Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
139 changes: 52 additions & 87 deletions packages/registry-server/client/lib/client.js
Original file line number Diff line number Diff line change
Expand Up @@ -304,7 +304,7 @@ class QVACRegistryClient extends ReadyResource {
throw new Error(`Invalid options: ${typeof options}`)
}

let core, blobs
let core, blobs, blockStart, blockEnd, rangeDownload

try {
this.logger.info('Downloading model', { path, source })
Expand Down Expand Up @@ -344,13 +344,13 @@ class QVACRegistryClient extends ReadyResource {

const totalSize = model.blobBinding.byteLength

const rangeDownload = core.download({
rangeDownload = core.download({
start: model.blobBinding.blockOffset,
length: model.blobBinding.blockLength
})

const blockStart = model.blobBinding.blockOffset
const blockEnd = blockStart + model.blobBinding.blockLength
blockStart = model.blobBinding.blockOffset
blockEnd = blockStart + model.blobBinding.blockLength

let artifact
if (options.outputFile) {
Expand All @@ -369,38 +369,15 @@ class QVACRegistryClient extends ReadyResource {
)
artifact = { path: options.outputFile, totalSize }

rangeDownload.destroy()
await this._clearBlobBlocks(core, blockStart, blockEnd)
if (blobs) await blobs.close()
if (core) await core.close()
await this._releaseDownload(core, blobs, rangeDownload, blockStart, blockEnd)
} else {
const stream = blobs.createReadStream(model.blobBinding, {
wait: true,
timeout: options.timeout || 30000
})
artifact = { stream, totalSize }

const cleanup = async () => {
rangeDownload.destroy()
await this._clearBlobBlocks(core, blockStart, blockEnd)
if (blobs) {
try {
await blobs.close()
} catch (cleanupError) {
this.logger.warn('Error closing blob instance', { error: cleanupError.message })
}
}
if (core) {
try {
await core.close()
} catch (cleanupError) {
this.logger.warn('Error closing blob core', { error: cleanupError.message })
}
}
this.logger.debug('Blob resources closed after stream end')
}

stream.once('end', cleanup)
this._releaseOnStreamEnd(stream, core, blobs, rangeDownload, blockStart, blockEnd)
}

this.logger.info('Model downloaded successfully')
Expand All @@ -412,20 +389,7 @@ class QVACRegistryClient extends ReadyResource {
} catch (error) {
this.logger.error('Error downloading model', error)

if (blobs) {
try {
await blobs.close()
} catch (cleanupError) {
this.logger.warn('Error closing blob instance on error', { error: cleanupError.message })
}
}
if (core) {
try {
await core.close()
} catch (cleanupError) {
this.logger.warn('Error closing blob core on error', { error: cleanupError.message })
}
}
await this._releaseDownload(core, blobs, rangeDownload, blockStart, blockEnd)

throw error
}
Expand All @@ -449,7 +413,7 @@ class QVACRegistryClient extends ReadyResource {
throw new Error(`Invalid options: ${typeof options}`)
}

let core, blobs
let core, blobs, blockStart, blockEnd, rangeDownload

try {
this.logger.info('Downloading blob directly', {
Expand Down Expand Up @@ -487,10 +451,10 @@ class QVACRegistryClient extends ReadyResource {
}
const totalSize = blobBinding.byteLength

const blockStart = pointer.blockOffset
const blockEnd = blockStart + pointer.blockLength
blockStart = pointer.blockOffset
blockEnd = blockStart + pointer.blockLength

const rangeDownload = core.download({
rangeDownload = core.download({
start: pointer.blockOffset,
length: pointer.blockLength
})
Expand All @@ -510,38 +474,15 @@ class QVACRegistryClient extends ReadyResource {
)
artifact = { path: options.outputFile, totalSize }

rangeDownload.destroy()
await this._clearBlobBlocks(core, blockStart, blockEnd)
if (blobs) await blobs.close()
if (core) await core.close()
await this._releaseDownload(core, blobs, rangeDownload, blockStart, blockEnd)
} else {
const stream = blobs.createReadStream(pointer, {
wait: true,
timeout: options.timeout || 30000
})
artifact = { stream, totalSize }

const cleanup = async () => {
rangeDownload.destroy()
await this._clearBlobBlocks(core, blockStart, blockEnd)
if (blobs) {
try {
await blobs.close()
} catch (e) {
this.logger.warn('Error closing blob instance', { error: e.message })
}
}
if (core) {
try {
await core.close()
} catch (e) {
this.logger.warn('Error closing blob core', { error: e.message })
}
}
this.logger.debug('Blob resources closed after stream end')
}

stream.once('end', cleanup)
this._releaseOnStreamEnd(stream, core, blobs, rangeDownload, blockStart, blockEnd)
}

this.logger.info('Blob download complete (direct)')
Expand All @@ -550,26 +491,50 @@ class QVACRegistryClient extends ReadyResource {
} catch (error) {
this.logger.error('Error downloading blob directly', error)

if (blobs) {
try {
await blobs.close()
} catch (e) {
this.logger.warn('Error closing blob instance on error', { error: e.message })
}
}
if (core) {
try {
await core.close()
} catch (e) {
this.logger.warn('Error closing blob core on error', { error: e.message })
}
}
await this._releaseDownload(core, blobs, rangeDownload, blockStart, blockEnd)

throw error
}
}

async _clearBlobBlocks(core, start, end) {
_releaseOnStreamEnd (stream, core, blobs, rangeDownload, blockStart, blockEnd) {
let released = false

// 'close' also covers a destroyed or errored stream; on 'end' alone a
// cancelled stream download would never free its blocks.
const release = () => {
if (released) return
released = true
return this._releaseDownload(core, blobs, rangeDownload, blockStart, blockEnd)
.catch(e => this.logger.warn('Error releasing blob resources', { error: e.message }))
}

stream.once('end', release)
stream.once('close', release)
}

async _releaseDownload (core, blobs, rangeDownload, blockStart, blockEnd) {
// Stop replication before clearing to prevent blocks from being refetched.
if (rangeDownload) rangeDownload.destroy()

if (core && blockStart != null) {
await this._clearBlobBlocks(core, blockStart, blockEnd)
}
if (blobs) {
try { await blobs.close() } catch (e) {
this.logger.warn('Error closing blob instance', { error: e.message })
}
}
if (core) {
try { await core.close() } catch (e) {
this.logger.warn('Error closing blob core', { error: e.message })
}
}

this.logger.debug('Blob resources released')
}

async _clearBlobBlocks (core, start, end) {
try {
const cleared = await core.clear(start, end, { diff: true })
await core.compact()
Expand Down
Loading
Loading