import * as plugins from './plugins.js'; import * as paths from './paths.js'; import type * as interfaces from './interfaces/index.js'; import { DockerContainer } from './classes.container.js'; import { DockerConfig } from './classes.config.js'; import { DockerNetwork } from './classes.network.js'; import { DockerService } from './classes.service.js'; import { DockerSecret } from './classes.secret.js'; import { logger } from './logger.js'; import { DockerImageStore, type IDockerImageArchiveReceipt, type IDockerImageStoreOperationOptions } from './classes.imagestore.js'; import { DockerImage } from './classes.image.js'; import { DockerVolume } from './classes.volume.js'; import { assertDockerResponseStatus, assertDockerServiceId, assertNonnegativeInteger } from './helpers.docker.js'; import { parseDockerUnixNanoseconds } from './helpers.timestamp.js'; export interface IDockerInfo { /** Daemon security feature strings, including `name=rootless` when enabled. */ SecurityOptions: string[]; } export interface IDockerHostConstructorOptions { socketPath?: string; imageStoreDir?: string; /** Set false to prevent host lifecycle methods from creating image-store files. */ enableImageStore?: boolean; } export interface IHijackedStreamingResponse { stream: plugins.stream.Duplex; close: () => Promise; statusCode: number; headers: plugins.http.IncomingHttpHeaders; } export interface IHijackedStreamingRequestOptions { /** Handshake deadline in milliseconds. Defaults to 30000; pass 0 to disable. */ timeoutMs?: number; /** Cancels the pending handshake and closes its underlying request or socket. */ signal?: AbortSignal; } export interface IDockerRequestOptions { /** * Socket idle timeout in milliseconds. Defaults to 30000. * Pass 0 for Docker operations such as image pulls that may remain silent * while the daemon downloads or extracts a large layer. */ timeoutMs?: number; /** Registry credentials scoped exclusively to this request. */ registryAuth?: interfaces.TDockerRegistryAuth; /** Cancels the request and any response body consumption. */ signal?: AbortSignal; } export interface IDockerStreamingRequestOptions { /** * Socket idle timeout in milliseconds. Defaults to 30000. * Pass 0 for long-lived streams such as events or followed logs. */ timeoutMs?: number; /** Cancels a pending stream open. */ signal?: AbortSignal; } export interface IDockerStreamingErrorResponse { statusCode: number; body: string; headers: Record; } interface IDenoUnixConnectOptions { transport: 'unix'; path: string; } interface IDenoRuntime { connect(optionsArg: IDenoUnixConnectOptions): Promise; } interface IDenoUnixConnection { read(bufferArg: Uint8Array): Promise; write(bufferArg: Uint8Array): Promise; closeWrite(): Promise; close(): void; } const isDenoRuntime = (valueArg: unknown): valueArg is IDenoRuntime => typeof valueArg === 'object' && valueArg !== null && typeof Reflect.get(valueArg, 'connect') === 'function'; const isDenoUnixConnection = (valueArg: unknown): valueArg is IDenoUnixConnection => typeof valueArg === 'object' && valueArg !== null && typeof Reflect.get(valueArg, 'read') === 'function' && typeof Reflect.get(valueArg, 'write') === 'function' && typeof Reflect.get(valueArg, 'closeWrite') === 'function' && typeof Reflect.get(valueArg, 'close') === 'function'; const isDockerEvent = (valueArg: unknown): valueArg is interfaces.IDockerEvent => typeof valueArg === 'object' && valueArg !== null && !Array.isArray(valueArg); const getErrorCode = (errorArg: unknown): string | undefined => { if (typeof errorArg !== 'object' || errorArg === null) { return undefined; } const codeValue: unknown = Reflect.get(errorArg, 'code'); return typeof codeValue === 'string' ? codeValue : undefined; }; const asEventError = (errorArg: unknown): Error => errorArg instanceof Error ? errorArg : new Error(String(errorArg)); const buildDockerEventPath = (filtersArg?: Record, sinceArg?: number): string => { const params = new URLSearchParams(); if (filtersArg) params.append('filters', JSON.stringify(filtersArg)); if (sinceArg !== undefined) params.append('since', `${sinceArg}`); const query = params.toString(); return query ? `/events?${query}` : '/events'; }; const parseDockerEventLine = (lineArg: string): interfaces.IDockerEvent | undefined => { const parsedEvent: unknown = JSON.parse(lineArg); if (!isDockerEvent(parsedEvent)) return undefined; // JSON.parse has already rounded Docker's signed 64-bit timeNano. Recover the // original token without changing the numeric fields shipped to existing users. delete parsedEvent.timeNanoExact; try { const losslessEvent = plugins.losslessJson.parse(lineArg, undefined, { onDuplicateKey: ({ newValue }) => newValue, }); if (isDockerEvent(losslessEvent)) { const rawNano = losslessEvent.timeNano; if (plugins.losslessJson.isLosslessNumber(rawNano)) { const exact = rawNano.toString(); parseDockerUnixNanoseconds(exact, 'Docker event timeNano'); parsedEvent.timeNanoExact = exact; } } } catch { // Generic events still flow; the timed parser rejects missing exact proof. } return parsedEvent; }; /** Owns one stream's parser and listeners. Disposal prevents further callbacks. */ const consumeDockerEventStream = ( streamArg: plugins.smartstream.stream.Readable, onEventArg: (eventArg: interfaces.IDockerEvent) => void, onGoneArg: (errorArg?: Error) => void, onMalformedArg: (errorArg: Error) => void, ): (() => void) => { let remainder = ''; let disposed = false; const decoder = new plugins.stringDecoder.StringDecoder('utf8'); const cleanupListeners = () => { streamArg.off('data', handleData); streamArg.off('error', handleError); streamArg.off('end', handleEnd); streamArg.off('close', handleClose); }; const dispose = () => { if (disposed) return; disposed = true; cleanupListeners(); if (!streamArg.destroyed) streamArg.destroy(); }; const handleGone = (errorArg?: Error) => { if (disposed) return; dispose(); onGoneArg(errorArg); }; const handleData = (dataArg: Buffer | string) => { if (disposed) return; remainder += typeof dataArg === 'string' ? dataArg : decoder.write(dataArg); const lines = remainder.split('\n'); remainder = lines.pop() || ''; for (const line of lines) { if (disposed) return; const trimmed = line.trim(); if (!trimmed) continue; try { const event = parseDockerEventLine(trimmed); if (event) onEventArg(event); } catch (errorArg) { onMalformedArg(asEventError(errorArg)); } } }; const handleError = (errorArg: unknown) => handleGone(asEventError(errorArg)); const handleEnd = () => handleGone(); const handleClose = () => handleGone(); streamArg.on('data', handleData); streamArg.on('error', handleError); streamArg.on('end', handleEnd); streamArg.on('close', handleClose); return dispose; }; export class DockerHost { public options: IDockerHostConstructorOptions; /** * the path where the docker sock can be found */ public socketPath: string; private imageStore?: DockerImageStore; public smartBucket?: plugins.smartbucket.SmartBucket; private readonly storageLifecycleStack = new plugins.lik.AsyncExecutionStack(); private stopRequested = false; /** * the constructor to instantiate a new docker sock instance * @param pathArg */ constructor(optionsArg: IDockerHostConstructorOptions) { this.options = { ...{ enableImageStore: true, imageStoreDir: plugins.path.join( paths.nogitDir, 'temp-docker-image-store', ), }, ...optionsArg, }; let pathToUse: string; if (optionsArg.socketPath) { pathToUse = optionsArg.socketPath; } else if (process.env.DOCKER_HOST) { pathToUse = process.env.DOCKER_HOST; } else if (process.env.CI) { pathToUse = 'http://docker:2375/'; } else { pathToUse = 'http://unix:/var/run/docker.sock:'; } if (pathToUse.startsWith('unix:///')) { pathToUse = pathToUse.replace('unix://', 'http://unix:'); } if (pathToUse.endsWith('.sock')) { pathToUse = pathToUse.replace('.sock', '.sock:'); } console.log(`using docker sock at ${pathToUse}`); this.socketPath = pathToUse; if (this.options.enableImageStore) { this.imageStore = new DockerImageStore({ bucketDir: null!, localDirPath: this.options.imageStoreDir!, }); } } public async start() { await this.imageStore?.start(); } public async stop(): Promise { this.stopRequested = true; return await this.storageLifecycleStack.getExclusiveExecutionSlot( async () => { const cleanupErrors: unknown[] = []; try { await this.imageStore?.stop(); } catch (error) { cleanupErrors.push(error); } const smartBucket = this.smartBucket; if (smartBucket) { try { await smartBucket.close(); } catch (error) { cleanupErrors.push(error); } finally { if (this.smartBucket === smartBucket) { this.smartBucket = undefined; } } } if (cleanupErrors.length > 0) { throw new AggregateError(cleanupErrors, 'DockerHost cleanup failed'); } }, ); } /** * Ping the Docker daemon to check if it's running and accessible * @returns Promise that resolves if Docker is available, rejects otherwise * @throws Error if Docker ping fails */ public async ping(): Promise { const response = await this.request('GET', '/_ping'); if (response.statusCode !== 200) { throw new Error(`Docker ping failed with status ${response.statusCode}`); } } /** Read daemon security features afresh before a rootful container request. */ public async info(): Promise { const response = await this.request('GET', '/info'); assertDockerResponseStatus('Docker info', 200, response); const body: unknown = response.body; if (typeof body !== 'object' || body === null || Array.isArray(body)) { throw new Error('Docker info failed: expected an object response'); } const securityOptions: unknown = Reflect.get(body, 'SecurityOptions'); if (!Array.isArray(securityOptions) || !securityOptions.every((option) => typeof option === 'string')) { throw new Error('Docker info failed: expected SecurityOptions to be a string array'); } return { SecurityOptions: securityOptions }; } /** * Get Docker daemon version information * @returns Version info including Docker version, API version, OS, architecture, etc. */ public async getVersion(): Promise<{ Version: string; ApiVersion: string; MinAPIVersion?: string; GitCommit: string; GoVersion: string; Os: string; Arch: string; KernelVersion: string; BuildTime?: string; }> { const response = await this.request('GET', '/version'); return response.body; } // ============== // NETWORKS - Public Factory API // ============== /** * Lists all networks */ public async listNetworks(optionsArg: interfaces.INetworkListOptions = {}) { return await DockerNetwork._list(this, optionsArg); } /** Gets a network through direct immutable-ID inspection. */ public async getNetworkById(networkIdArg: string) { return await DockerNetwork._fromId(this, networkIdArg); } /** * Gets a network by name */ public async getNetworkByName(networkNameArg: string) { return await DockerNetwork._fromName(this, networkNameArg); } /** * Creates a network */ public async createNetwork( descriptor: interfaces.INetworkCreationDescriptor, ) { return await DockerNetwork._create(this, descriptor); } // ============== // VOLUMES - Public Factory API // ============== public async listVolumes(optionsArg: interfaces.IVolumeListOptions = {}) { return await DockerVolume._list(this, optionsArg); } public async getVolumeByName(volumeNameArg: string) { return await DockerVolume._fromName(this, volumeNameArg); } public async createVolume(descriptorArg: interfaces.IVolumeCreationDescriptor) { return await DockerVolume._create(this, descriptorArg); } // ============== // CONTAINERS - Public Factory API // ============== /** * Lists all containers */ public async listContainers(optionsArg: interfaces.IContainerListOptions = {}) { return await DockerContainer._list(this, optionsArg); } /** * Gets a container by ID * Returns undefined if container does not exist */ public async getContainerById( containerId: string, optionsArg: interfaces.IContainerReadOptions = {}, ): Promise { return await DockerContainer._fromId(this, containerId, optionsArg); } /** * Creates a container */ public async createContainer( descriptor: interfaces.IContainerCreationDescriptor, ) { return await DockerContainer._create(this, descriptor); } // ============== // SERVICES - Public Factory API // ============== /** * Lists all services */ public async listServices() { return await DockerService._list(this); } /** * Gets a service by name */ public async getServiceByName(serviceName: string) { return await DockerService._fromName(this, serviceName); } /** * Gets a service by its exact Docker service ID. * Returns undefined only when Docker reports that exact ID as absent. */ public async getServiceById(serviceId: string) { assertDockerServiceId(serviceId); const service = await DockerService._fromId(this, serviceId); if (service && service.ID !== serviceId) { throw new Error(`Docker service inspect returned a different ID for ${serviceId}`); } return service; } /** * Creates a service */ public async createService( descriptor: interfaces.IServiceCreationDescriptor, ) { return await DockerService._create(this, descriptor); } // ============== // IMAGES - Public Factory API // ============== /** * Lists all images */ public async listImages(optionsArg: interfaces.IImageListOptions = {}) { return await DockerImage._list(this, optionsArg); } /** * Gets an image by a complete local tag or digest reference. */ public async getImageByReference(referenceArg: string) { return await DockerImage._fromReference(this, referenceArg); } /** * Gets an image through its immutable local ID. */ public async getImageById(idArg: string) { return await DockerImage._fromId(this, idArg); } /** Pulls and verifies an exact registry image. */ public async pullImage(descriptorArg: interfaces.IImagePullDescriptor) { return await DockerImage._pull(this, descriptorArg); } /** Pulls an explicitly tagged mutable application image as an inspected DockerImage. */ public async pullMutableImage(descriptorArg: interfaces.IImageMutablePullDescriptor) { return await DockerImage._pullMutable(this, descriptorArg); } /** * Creates an image from a tar stream */ public async createImageFromTarStream( tarStream: plugins.smartstream.stream.Readable, descriptor: interfaces.IImageCreationDescriptor, ) { return await DockerImage._createFromTarStream(this, { creationObject: descriptor, tarStream: tarStream, }); } /** * Prune unused images * @param options Optional filters (dangling, until, label) * @returns Object with deleted images and space reclaimed */ public async pruneImages(options?: { dangling?: boolean; filters?: Record; }): Promise<{ ImagesDeleted: Array<{ Untagged?: string; Deleted?: string }>; SpaceReclaimed: number; }> { const filters: Record = options?.filters || {}; // Add dangling filter if specified if (options?.dangling !== undefined) { filters.dangling = [options.dangling.toString()]; } let route = '/images/prune'; if (filters && Object.keys(filters).length > 0) { route += `?filters=${encodeURIComponent(JSON.stringify(filters))}`; } const response = await this.request('POST', route); return response.body; } /** * Builds an image from a Dockerfile */ public async buildImage(imageTag: string) { return await DockerImage._build(this, imageTag); } // ============== // SECRETS - Public Factory API // ============== /** * Lists all secrets */ public async listSecrets() { return await DockerSecret._list(this); } /** * Gets a secret by name */ public async getSecretByName(secretName: string) { return await DockerSecret._fromName(this, secretName); } /** * Gets a secret by ID */ public async getSecretById(secretId: string) { return await DockerSecret._fromId(this, secretId); } /** * Creates a secret */ public async createSecret( descriptor: interfaces.ISecretCreationDescriptor, ) { return await DockerSecret._create(this, descriptor); } // ============== // CONFIGS - Public Factory API // ============== /** * Lists all configs */ public async listConfigs() { return await DockerConfig._list(this); } /** * Gets a config by ID */ public async getConfigById(configId: string) { return await DockerConfig._fromId(this, configId); } /** * Gets a config by name */ public async getConfigByName(configName: string) { return await DockerConfig._fromName(this, configName); } /** * Creates a config */ public async createConfig( descriptor: interfaces.IConfigCreationDescriptor, ) { return await DockerConfig._create(this, descriptor); } // ============== // IMAGE STORE - Public API // ============== /** * Atomically stores or replaces the conventional image archive in SmartBucket. */ public async storeImage( imageName: string, tarStream: plugins.smartstream.stream.Readable, ): Promise { return await this.getEnabledImageStore().storeImage(imageName, tarStream); } /** Creates a unique archive with an exact receipt for the repacked stored bytes. */ public storeImageArchive( imageName: string, tarStream: plugins.stream.Readable, options: IDockerImageStoreOperationOptions = {}, ): Promise { return this.getEnabledImageStore().storeImageArchive(imageName, tarStream, options); } /** * Retrieves an image archive from SmartBucket; the caller owns the returned stream. */ public async retrieveImage( imageName: string, ): Promise { return await this.getEnabledImageStore().getImage(imageName); } /** * Open one supervised event connection. Resolves only after Docker sends HTTP 200 headers * and event listeners are installed. There is no automatic reconnect: any loss is reported * once, so a runtime owner can stop the workload before establishing a new monitor. */ public async openEventMonitor( optionsArg: interfaces.IDockerEventMonitorOptions, ): Promise { if (optionsArg.signal?.aborted) { throw asEventError(optionsArg.signal.reason ?? new Error('Docker event monitor open aborted.')); } const openTimeoutMs = optionsArg.openTimeoutMs ?? 30_000; if (!Number.isSafeInteger(openTimeoutMs) || openTimeoutMs < 1 || openTimeoutMs > 2_147_483_647) { throw new TypeError('Docker event monitor openTimeoutMs must be a positive timer-safe integer.'); } const openController = new AbortController(); const abortPending = () => openController.abort( optionsArg.signal?.reason ?? new Error('Docker event monitor open aborted.'), ); optionsArg.signal?.addEventListener('abort', abortPending, { once: true }); if (optionsArg.signal?.aborted) abortPending(); const openTimer = setTimeout(() => { openController.abort(new Error('Docker event monitor HTTP readiness timed out.')); }, openTimeoutMs); openTimer.unref(); let stream: plugins.smartstream.stream.Readable; try { stream = await this.openDockerEventStream( optionsArg.filters, undefined, openController.signal, true, ); } catch (errorArg) { if (openController.signal.aborted) throw asEventError(openController.signal.reason); throw errorArg; } finally { clearTimeout(openTimer); optionsArg.signal?.removeEventListener('abort', abortPending); } if (openController.signal.aborted) { stream.destroy(); throw asEventError(openController.signal.reason); } let status: interfaces.IDockerEventMonitor['status'] = 'connected'; let disposeStream: () => void = () => stream.destroy(); const onAbort = () => monitor.close(); const monitor: interfaces.IDockerEventMonitor = { get status() { return status; }, close: () => { if (status === 'closed') return; status = 'closed'; optionsArg.signal?.removeEventListener('abort', onAbort); disposeStream(); }, }; const lose = (errorArg: Error) => { if (status !== 'connected') return; if (optionsArg.signal?.aborted) { monitor.close(); return; } status = 'lost'; optionsArg.signal?.removeEventListener('abort', onAbort); disposeStream(); try { optionsArg.onLost(errorArg); } catch (callbackError) { console.log(callbackError); } }; disposeStream = consumeDockerEventStream( stream, optionsArg.onEvent, (errorArg) => lose(errorArg ?? new Error('Docker event monitor stream ended.')), lose, ); optionsArg.signal?.addEventListener('abort', onAbort, { once: true }); if (optionsArg.signal?.aborted) { monitor.close(); throw asEventError(optionsArg.signal.reason ?? new Error('Docker event monitor open aborted.')); } if (stream.destroyed || stream.readableEnded) { lose(new Error('Docker event monitor stream closed before readiness.')); throw new Error('Docker event monitor stream closed before readiness.'); } return monitor; } private async openDockerEventStream( filtersArg: Record | undefined, sinceArg: number | undefined, signalArg: AbortSignal | undefined, requireHttp200Arg: boolean, ): Promise { // The stream may be idle indefinitely; disable the ordinary socket idle timeout. const response = await this.requestStreaming( 'GET', buildDockerEventPath(filtersArg, sinceArg), undefined, undefined, { timeoutMs: 0, signal: signalArg }, ); if (!(response instanceof plugins.smartstream.stream.Readable)) { throw new Error( `Docker event stream failed: actual status ${response.statusCode}; response body: ${response.body}`, ); } if (requireHttp200Arg && Reflect.get(response, 'statusCode') !== 200) { const statusCode: unknown = Reflect.get(response, 'statusCode'); response.destroy(); throw new Error(`Docker event monitor requires HTTP 200; received ${String(statusCode)}.`); } return response; } /** * Streams docker engine events as parsed objects. * * Events arrive newline-delimited; lines are buffered across chunks so * bursts (multiple events per chunk) and frames split across chunks parse * correctly. The observable completes when the daemon closes the stream — * unless `reconnect` keeps it alive. * * @param optionsArg.filters server-side event filters, e.g. * `{ type: ['container'], event: ['start', 'die'] }` — unwanted events * are never delivered or parsed. * @param optionsArg.reconnect re-establish the stream with backoff * (1s..30s) when it ends or errors, passing `since` so events from the * gap are replayed. Duplicates around the reconnect boundary are * possible; consumers should be idempotent. The stream then only ends * on unsubscribe. */ public async getEventObservable( optionsArg: interfaces.IDockerEventObservableOptions = {}, ): Promise> { return new plugins.rxjs.Observable((observer) => { let unsubscribed = false; let currentStream: plugins.smartstream.stream.Readable | undefined; let currentStreamCleanup: (() => void) | undefined; let pendingOpenController: AbortController | undefined; let reconnectTimer: ReturnType | undefined; let lastEventTime: number | undefined; let reconnectDelayMs = 1000; const scheduleReconnect = (errorArg?: Error) => { if (unsubscribed) { return; } if (!optionsArg.reconnect) { if (errorArg) { observer.error(errorArg); } else { observer.complete(); } return; } const delay = reconnectDelayMs; reconnectDelayMs = Math.min(reconnectDelayMs * 2, 30000); reconnectTimer = setTimeout(() => { reconnectTimer = undefined; if (!unsubscribed) { void openAndConsume(lastEventTime); } }, delay); }; const consume = (nodeStream: plugins.smartstream.stream.Readable) => { currentStream = nodeStream; const dispose = consumeDockerEventStream( nodeStream, (eventArg) => { if (typeof eventArg.time === 'number') lastEventTime = eventArg.time; reconnectDelayMs = 1000; observer.next(eventArg); }, (errorArg) => { if (currentStreamCleanup === dispose) currentStreamCleanup = undefined; if (unsubscribed) return; if (currentStream === nodeStream) { currentStream = undefined; } scheduleReconnect(getErrorCode(errorArg) === 'ECONNRESET' ? undefined : errorArg); }, (errorArg) => console.log(errorArg), ); currentStreamCleanup = dispose; }; const openAndConsume = async (sinceArg?: number): Promise => { const openController = new AbortController(); pendingOpenController = openController; try { const nextStream = await this.openDockerEventStream( optionsArg.filters, sinceArg, openController.signal, false, ); if (pendingOpenController === openController) { pendingOpenController = undefined; } if (unsubscribed) { nextStream.destroy(); return; } consume(nextStream); } catch (error) { if (pendingOpenController === openController) { pendingOpenController = undefined; } if (unsubscribed || openController.signal.aborted) { return; } scheduleReconnect( error instanceof Error ? error : new Error(String(error)), ); } }; void openAndConsume(); return () => { unsubscribed = true; if (reconnectTimer) clearTimeout(reconnectTimer); reconnectTimer = undefined; pendingOpenController?.abort(); pendingOpenController = undefined; currentStreamCleanup?.(); currentStream?.destroy(); currentStream = undefined; }; }); } /** * activates docker swarm */ public async activateSwarm(addvertisementIpArg?: string) { // determine advertisement address let addvertisementIp: string = ''; if (addvertisementIpArg) { addvertisementIp = addvertisementIpArg; } else { try { const smartnetworkInstance = new plugins.smartnetwork.SmartNetwork(); const defaultGateway = await smartnetworkInstance.getDefaultGateway(); if (defaultGateway) { addvertisementIp = defaultGateway.ipv4.address; } } catch (err) { // Failed to determine default gateway (e.g. in Deno without --allow-run) // Docker will auto-detect the advertise address } } const response = await this.request('POST', '/swarm/init', { ListenAddr: '0.0.0.0:2377', AdvertiseAddr: addvertisementIp, DataPathPort: 4789, DefaultAddrPool: ['10.10.0.0/8', '20.20.0.0/8'], SubnetSize: 24, ForceNewCluster: false, }); if (response.statusCode === 200) { logger.log('info', 'created Swam succesfully'); } else { logger.log('error', 'could not initiate swarm'); } } /** * fire a request */ public async request( methodArg: string, routeArg: string, dataArg = {}, optionsArg?: IDockerRequestOptions, ) { const requestUrl = `${this.socketPath}${routeArg}`; const timeoutMs = optionsArg?.timeoutMs ?? 30000; // Build the request using the fluent API const smartRequest = plugins.smartrequest.SmartRequest.create() .url(requestUrl) .header('Content-Type', 'application/json') .header('Host', 'docker.sock') .options( timeoutMs > 0 ? { keepAlive: false, signal: optionsArg?.signal } : // Setting keepAlive to either true or false installs an HTTP // agent with its own short idle timeout. Omit it entirely for // operations whose response can legitimately be silent. { signal: optionsArg?.signal }, ); if (optionsArg?.registryAuth) { smartRequest.header( 'X-Registry-Auth', this.encodeRegistryAuth(optionsArg.registryAuth), ); } if (timeoutMs > 0) { smartRequest.timeout(timeoutMs); } // Add body for methods that support it if (dataArg && Object.keys(dataArg).length > 0) { smartRequest.json(dataArg); } // Execute the request based on method let response; switch (methodArg.toUpperCase()) { case 'GET': response = await smartRequest.get(); break; case 'POST': response = await smartRequest.post(); break; case 'PUT': response = await smartRequest.put(); break; case 'DELETE': response = await smartRequest.delete(); break; default: throw new Error(`Unsupported HTTP method: ${methodArg}`); } // Parse the response body based on content type let body; const contentType = response.headers['content-type'] || ''; // Docker's streaming endpoints (like /images/create) return newline-delimited JSON // which can't be parsed as a single JSON object const isStreamingEndpoint = routeArg.includes('/images/create') || routeArg.includes('/images/load') || routeArg.includes('/build'); if (contentType.includes('application/json') && !isStreamingEndpoint) { body = await response.json(); } else { body = await response.text(); // Try to parse as JSON if it looks like JSON and is not a streaming response if ( !isStreamingEndpoint && body && (body.startsWith('{') || body.startsWith('[')) ) { try { body = JSON.parse(body); } catch { // Keep as text if parsing fails } } } // Create a response object compatible with existing code const legacyResponse = { statusCode: response.status, body: body, headers: response.headers, }; if (response.status !== 200) { console.log(body); } return legacyResponse; } public async requestStreaming( methodArg: string, routeArg: string, readStream?: plugins.smartstream.stream.Readable, jsonData?: Record, optionsArg: IDockerStreamingRequestOptions = {}, ): Promise { const requestUrl = `${this.socketPath}${routeArg}`; const timeoutMs = optionsArg.timeoutMs ?? 30000; // Build the request using the fluent API const smartRequest = plugins.smartrequest.SmartRequest.create() .url(requestUrl) .header('Content-Type', 'application/json') .header('Host', 'docker.sock') .options( timeoutMs > 0 ? { keepAlive: false, autoDrain: true, signal: optionsArg.signal } : // No keepAlive key: with keepAlive set (true or false) smartrequest // attaches an agentkeepalive agent whose 8s idle socket timeout // would destroy long-lived follow streams; without it, no agent // and no idle timeout are applied. { autoDrain: true, signal: optionsArg.signal }, ); if (timeoutMs > 0) { smartRequest.timeout(timeoutMs); } // If we have JSON data, add it to the request if (jsonData && Object.keys(jsonData).length > 0) { smartRequest.json(jsonData); } // If we have a readStream, use the new stream method with logging if (readStream) { let counter = 0; const smartduplex = new plugins.smartstream.SmartDuplex({ writeFunction: async (chunkArg) => { if (counter % 1000 === 0) { console.log(`posting chunk ${counter}`); } counter++; return chunkArg; }, }); // Pipe through the logging duplex stream const loggedStream = readStream.pipe(smartduplex); // Use the new stream method to stream the data smartRequest.stream(loggedStream, 'application/octet-stream'); } // Execute the request based on method let response: plugins.smartrequest.ICoreResponse; switch (methodArg.toUpperCase()) { case 'GET': response = await smartRequest.get(); break; case 'POST': response = await smartRequest.post(); break; case 'PUT': response = await smartRequest.put(); break; case 'DELETE': response = await smartRequest.delete(); break; default: throw new Error(`Unsupported HTTP method: ${methodArg}`); } console.log(response.status); // For streaming responses, get the web stream const webStream = response.stream(); if (!webStream) { // If no stream is available, consume the body as text const body = await response.text(); console.log(body); // Return a compatible response object return { statusCode: response.status, body: body, headers: response.headers, }; } // Convert web ReadableStream to Node.js stream for backward compatibility const nodeStream = plugins.smartstream.nodewebhelpers.convertWebReadableToNodeReadable(webStream); // Add a default error handler to prevent unhandled 'error' events from crashing the process. // Callers that attach their own 'error' listener will still receive the event. nodeStream.on('error', () => {}); // Retain the legacy runtime properties without weakening the public type. Object.defineProperties(nodeStream, { statusCode: { value: response.status, enumerable: true }, body: { value: '', enumerable: true }, }); return nodeStream; } public async requestHijackedStreaming( methodArg: string, routeArg: string, jsonData: Record = {}, optionsArg: IHijackedStreamingRequestOptions = {}, ): Promise { const timeoutMs = optionsArg.timeoutMs ?? 30_000; assertNonnegativeInteger(timeoutMs, 'Docker hijack request timeoutMs'); if (optionsArg.signal?.aborted) { throw this.createHijackAbortError(methodArg, routeArg, optionsArg.signal); } const body = JSON.stringify(jsonData); const headers: Record = { 'Content-Type': 'application/json', 'Content-Length': Buffer.byteLength(body), 'Connection': 'Upgrade', 'Upgrade': 'tcp', 'Host': 'docker.sock', }; if (this.socketPath.startsWith('http://unix:')) { return await this.requestHijackedStreamingOverRawSocket( methodArg, routeArg, body, headers, timeoutMs, optionsArg.signal, ); } const { requestModule, options } = this.getNodeRequestOptions(methodArg, routeArg, headers); return await new Promise((resolve, reject) => { let settled = false; let request: plugins.http.ClientRequest; let activeResponse: plugins.http.IncomingMessage | undefined; let activeSocket: plugins.stream.Duplex | undefined; let timeout: ReturnType | undefined; const cleanupHandshake = () => { if (timeout) { clearTimeout(timeout); timeout = undefined; } optionsArg.signal?.removeEventListener('abort', handleAbort); request.off('response', handleResponse); request.off('upgrade', handleUpgrade); }; const fail = (errorArg: Error) => { if (settled) { return; } settled = true; request.destroy(); activeResponse?.destroy(); activeSocket?.destroy(); cleanupHandshake(); reject(errorArg); }; const succeed = (responseArg: IHijackedStreamingResponse) => { if (settled) { responseArg.stream.destroy(); void responseArg.close().catch(() => {}); return; } settled = true; cleanupHandshake(); resolve(responseArg); }; const handleResponse = (response: plugins.http.IncomingMessage) => { activeResponse = response; if ((response.statusCode || 0) >= 400) { this.collectErrorResponse(response).then((bodyText) => { fail(new Error(`Docker hijack request failed with HTTP ${response.statusCode}: ${bodyText}`)); }).catch((errorArg: Error) => fail(errorArg)); return; } if (!response.socket) { fail(new Error('Docker hijack response did not include a socket')); return; } succeed({ stream: this.createDuplexForHijackedResponse(response, response.socket), close: async () => { response.destroy(); response.socket?.destroy(); }, statusCode: response.statusCode || 0, headers: response.headers, }); }; const handleUpgrade = ( response: plugins.http.IncomingMessage, socket: plugins.stream.Duplex, head: Buffer, ) => { activeSocket = socket; if (head.length > 0) { socket.unshift(head); } succeed({ stream: socket, close: async () => { socket.destroy(); }, statusCode: response.statusCode || 0, headers: response.headers, }); }; const handleRequestError = (errorArg: Error) => fail(errorArg); const handleRequestClose = () => { request.off('error', handleRequestError); }; const handleAbort = () => { if (optionsArg.signal) { fail(this.createHijackAbortError(methodArg, routeArg, optionsArg.signal)); } }; request = requestModule.request(options); request.once('response', handleResponse); request.once('upgrade', handleUpgrade); request.once('error', handleRequestError); request.once('close', handleRequestClose); if (timeoutMs > 0) { timeout = setTimeout(() => { fail( new Error( `Docker hijack request ${methodArg.toUpperCase()} ${routeArg} timed out after ${timeoutMs}ms`, ), ); }, timeoutMs); } optionsArg.signal?.addEventListener('abort', handleAbort, { once: true }); if (optionsArg.signal?.aborted) { handleAbort(); } if (!settled) { try { request.end(body); } catch (error) { fail(error instanceof Error ? error : new Error(String(error))); } } }); } private createHijackAbortError( methodArg: string, routeArg: string, signalArg: AbortSignal, ): Error { const reason = signalArg.reason; const reasonSuffix = reason === undefined ? '' : `: ${reason instanceof Error ? reason.message : String(reason)}`; const error = new Error( `Docker hijack request ${methodArg.toUpperCase()} ${routeArg} was aborted${reasonSuffix}`, ); error.name = 'AbortError'; return error; } private async requestHijackedStreamingOverRawSocket( methodArg: string, routeArg: string, bodyArg: string, headersArg: Record, timeoutMsArg: number, signalArg?: AbortSignal, ): Promise { const socketPath = this.socketPath.slice('http://unix:'.length, -1); const denoGlobal: unknown = Object.getOwnPropertyDescriptor(globalThis, 'Deno')?.value; if (isDenoRuntime(denoGlobal)) { return await this.requestHijackedStreamingOverDenoUnixSocket( socketPath, methodArg, routeArg, bodyArg, headersArg, timeoutMsArg, signalArg, ); } const socket = plugins.net.connect(socketPath); const requestHead = [ `${methodArg.toUpperCase()} ${routeArg} HTTP/1.1`, ...Object.entries(headersArg).map(([key, value]) => `${key}: ${value}`), '', bodyArg, ].join('\r\n'); return await new Promise((resolve, reject) => { let settled = false; let responseBuffer = Buffer.alloc(0); let dockerErrorStatus: number | undefined; const dockerErrorChunks: Buffer[] = []; let timeout: ReturnType | undefined; const cleanupHandshake = () => { if (timeout) { clearTimeout(timeout); timeout = undefined; } signalArg?.removeEventListener('abort', handleAbort); socket.off('connect', handleConnect); socket.off('data', handleData); socket.off('end', handleSocketEnd); }; const fail = (errorArg: Error) => { if (settled) { return; } settled = true; socket.destroy(); cleanupHandshake(); reject(errorArg); }; const succeed = (responseArg: IHijackedStreamingResponse) => { if (settled) { responseArg.stream.destroy(); void responseArg.close().catch(() => {}); return; } settled = true; cleanupHandshake(); resolve(responseArg); }; const handleData = (chunkArg: Buffer) => { if (dockerErrorStatus !== undefined) { dockerErrorChunks.push(Buffer.from(chunkArg)); return; } responseBuffer = Buffer.concat([responseBuffer, chunkArg]); const headerEndIndex = responseBuffer.indexOf('\r\n\r\n'); if (headerEndIndex === -1) { return; } const headerBuffer = responseBuffer.subarray(0, headerEndIndex); const bodyHead = responseBuffer.subarray(headerEndIndex + 4); const { statusCode, headers } = this.parseRawHttpResponseHeaders(headerBuffer.toString('utf8')); responseBuffer = Buffer.alloc(0); if (statusCode >= 400) { dockerErrorStatus = statusCode; if (bodyHead.length > 0) { dockerErrorChunks.push(Buffer.from(bodyHead)); } socket.end(); return; } succeed({ stream: this.createDuplexForRawSocket(socket, bodyHead), close: async () => { socket.destroy(); }, statusCode, headers, }); }; const handleConnect = () => { try { socket.write(requestHead); } catch (error) { fail(error instanceof Error ? error : new Error(String(error))); } }; const handleSocketError = (errorArg: Error) => fail(errorArg); const handleSocketEnd = () => { if (dockerErrorStatus !== undefined) { fail( new Error( `Docker hijack request failed with HTTP ${dockerErrorStatus}: ${Buffer.concat(dockerErrorChunks).toString('utf8')}`, ), ); return; } fail( new Error( `Docker hijack connection closed before response headers were received for ${methodArg.toUpperCase()} ${routeArg}`, ), ); }; const handleSocketClose = () => { socket.off('error', handleSocketError); if (!settled) { handleSocketEnd(); } }; const handleAbort = () => { if (signalArg) { fail(this.createHijackAbortError(methodArg, routeArg, signalArg)); } }; socket.once('connect', handleConnect); socket.on('data', handleData); socket.once('error', handleSocketError); socket.once('end', handleSocketEnd); socket.once('close', handleSocketClose); if (timeoutMsArg > 0) { timeout = setTimeout(() => { fail( new Error( `Docker hijack request ${methodArg.toUpperCase()} ${routeArg} timed out after ${timeoutMsArg}ms`, ), ); }, timeoutMsArg); } signalArg?.addEventListener('abort', handleAbort, { once: true }); if (signalArg?.aborted) { handleAbort(); } }); } private async requestHijackedStreamingOverDenoUnixSocket( socketPathArg: string, methodArg: string, routeArg: string, bodyArg: string, headersArg: Record, timeoutMsArg: number, signalArg?: AbortSignal, ): Promise { const denoGlobal: unknown = Object.getOwnPropertyDescriptor(globalThis, 'Deno')?.value; if (!isDenoRuntime(denoGlobal)) { throw new Error('Deno.connect is unavailable for Docker Unix-socket access'); } const requestHead = [ `${methodArg.toUpperCase()} ${routeArg} HTTP/1.1`, ...Object.entries(headersArg).map(([key, value]) => `${key}: ${value}`), '', bodyArg, ].join('\r\n'); let conn: IDenoUnixConnection | undefined; let cancelled = false; let handshakeComplete = false; let timeout: ReturnType | undefined; let rejectCancellation: (errorArg: Error) => void = () => {}; const cancellationPromise = new Promise((_resolve, reject) => { rejectCancellation = reject; }); const cleanupHandshake = () => { if (timeout) { clearTimeout(timeout); timeout = undefined; } signalArg?.removeEventListener('abort', handleAbort); }; const cancelHandshake = (errorArg: Error) => { if (cancelled || handshakeComplete) { return; } cancelled = true; if (conn) { try { conn.close(); } catch {} } cleanupHandshake(); rejectCancellation(errorArg); }; const handleAbort = () => { if (signalArg) { cancelHandshake(this.createHijackAbortError(methodArg, routeArg, signalArg)); } }; if (timeoutMsArg > 0) { timeout = setTimeout(() => { cancelHandshake( new Error( `Docker hijack request ${methodArg.toUpperCase()} ${routeArg} timed out after ${timeoutMsArg}ms`, ), ); }, timeoutMsArg); } signalArg?.addEventListener('abort', handleAbort, { once: true }); if (signalArg?.aborted) { handleAbort(); } const connectPromise = denoGlobal.connect({ transport: 'unix', path: socketPathArg, }); void connectPromise.then( (lateConnArg) => { if (cancelled && isDenoUnixConnection(lateConnArg)) { try { lateConnArg.close(); } catch {} } }, () => {}, ); try { const connectedValue = await Promise.race([connectPromise, cancellationPromise]); if (!isDenoUnixConnection(connectedValue)) { throw new Error('Deno.connect returned an invalid Unix connection'); } conn = connectedValue; await Promise.race([ conn.write(new TextEncoder().encode(requestHead)), cancellationPromise, ]); let responseBuffer = Buffer.alloc(0); const readBuffer = new Uint8Array(65536); while (true) { const bytesRead = await Promise.race([ conn.read(readBuffer), cancellationPromise, ]); if (bytesRead === null) { throw new Error('Docker hijack connection closed before response headers were received'); } responseBuffer = Buffer.concat([ responseBuffer, Buffer.from(readBuffer.subarray(0, bytesRead)), ]); const headerEndIndex = responseBuffer.indexOf('\r\n\r\n'); if (headerEndIndex === -1) { continue; } const headerBuffer = responseBuffer.subarray(0, headerEndIndex); const bodyHead = responseBuffer.subarray(headerEndIndex + 4); const { statusCode, headers } = this.parseRawHttpResponseHeaders(headerBuffer.toString('utf8')); if (statusCode >= 400) { throw new Error(`Docker hijack request failed with HTTP ${statusCode}: ${bodyHead.toString('utf8')}`); } handshakeComplete = true; cleanupHandshake(); const stream = this.createDuplexForDenoConn(conn, bodyHead); return { stream, close: async () => { stream.destroy(); }, statusCode, headers, }; } } catch (error) { if (!cancelled && conn) { try { conn.close(); } catch {} } cleanupHandshake(); throw error; } } private getNodeRequestOptions( methodArg: string, routeArg: string, headersArg: Record, ): { requestModule: typeof plugins.http | typeof plugins.https; options: plugins.http.RequestOptions | plugins.https.RequestOptions; } { const options: plugins.http.RequestOptions | plugins.https.RequestOptions = { method: methodArg.toUpperCase(), headers: headersArg, }; if (this.socketPath.startsWith('http://unix:')) { options.socketPath = this.socketPath.slice('http://unix:'.length, -1); options.path = routeArg; return { requestModule: plugins.http, options }; } const requestUrl = new URL(routeArg, this.socketPath); options.protocol = requestUrl.protocol; options.hostname = requestUrl.hostname; options.port = requestUrl.port; options.path = `${requestUrl.pathname}${requestUrl.search}`; return { requestModule: requestUrl.protocol === 'https:' ? plugins.https : plugins.http, options, }; } private createDuplexForHijackedResponse( responseArg: plugins.http.IncomingMessage, writableArg: plugins.stream.Duplex, ): plugins.stream.Duplex { const duplex = new plugins.stream.Duplex({ write(chunkArg, encodingArg, callbackArg) { writableArg.write(chunkArg, encodingArg, callbackArg); }, read() {}, destroy(errorArg, callbackArg) { responseArg.destroy(); writableArg.destroy(); callbackArg(errorArg || null); }, }); responseArg.on('data', (chunkArg) => duplex.push(chunkArg)); responseArg.on('end', () => duplex.push(null)); responseArg.on('error', (errorArg) => duplex.destroy(errorArg)); writableArg.on('error', (errorArg) => duplex.destroy(errorArg)); duplex.on('finish', () => writableArg.end()); return duplex; } private createDuplexForRawSocket( socketArg: plugins.net.Socket, bodyHeadArg: Buffer, ): plugins.stream.Duplex { const duplex = new plugins.stream.Duplex({ write(chunkArg, encodingArg, callbackArg) { socketArg.write(chunkArg, encodingArg, callbackArg); }, read() {}, destroy(errorArg, callbackArg) { socketArg.destroy(); callbackArg(errorArg || null); }, }); if (bodyHeadArg.length > 0) { duplex.push(bodyHeadArg); } socketArg.on('data', (chunkArg) => duplex.push(chunkArg)); socketArg.on('end', () => duplex.push(null)); socketArg.on('error', (errorArg) => duplex.destroy(errorArg)); duplex.on('finish', () => socketArg.end()); return duplex; } private createDuplexForDenoConn( connArg: IDenoUnixConnection, bodyHeadArg: Buffer, ): plugins.stream.Duplex { let closed = false; const closeConn = () => { if (closed) { return; } closed = true; try { connArg.close(); } catch {} }; const duplex = new plugins.stream.Duplex({ async write(chunkArg, encodingArg, callbackArg) { try { const chunkBuffer = Buffer.isBuffer(chunkArg) ? chunkArg : Buffer.from(chunkArg, encodingArg); await connArg.write(new Uint8Array(chunkBuffer)); callbackArg(); } catch (error) { callbackArg(error as Error); } }, read() {}, destroy(errorArg, callbackArg) { closeConn(); callbackArg(errorArg || null); }, }); if (bodyHeadArg.length > 0) { duplex.push(bodyHeadArg); } const readLoop = async () => { const readBuffer = new Uint8Array(65536); try { while (!duplex.destroyed) { const bytesRead = await connArg.read(readBuffer); if (bytesRead === null) { break; } if (bytesRead > 0) { duplex.push(Buffer.from(readBuffer.subarray(0, bytesRead))); } } duplex.push(null); } catch (error) { if (!closed && !duplex.destroyed) { duplex.destroy(error as Error); } } }; void readLoop(); duplex.on('finish', () => { void connArg.closeWrite().catch((errorArg) => { if (!closed && !duplex.destroyed) { duplex.destroy(errorArg); } }); }); return duplex; } private parseRawHttpResponseHeaders(headerTextArg: string): { statusCode: number; headers: plugins.http.IncomingHttpHeaders; } { const [statusLine, ...headerLines] = headerTextArg.split('\r\n'); const statusCode = Number(statusLine.split(' ')[1]); const headers: plugins.http.IncomingHttpHeaders = {}; for (const headerLine of headerLines) { const separatorIndex = headerLine.indexOf(':'); if (separatorIndex === -1) { continue; } const key = headerLine.slice(0, separatorIndex).trim().toLowerCase(); const value = headerLine.slice(separatorIndex + 1).trim(); headers[key] = value; } return { statusCode, headers }; } private async collectErrorResponse(responseArg: plugins.http.IncomingMessage): Promise { const chunks: Buffer[] = []; for await (const chunk of responseArg) { chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk)); } return Buffer.concat(chunks).toString('utf8'); } /** * add s3 storage * @param optionsArg */ public async addS3Storage(optionsArg: plugins.tsclass.storage.IS3Descriptor) { if (this.stopRequested) { throw new Error( 'Cannot configure Docker image-store S3 storage after DockerHost.stop() has begun', ); } const bucketName = optionsArg.bucketName; if (!bucketName) { throw new Error('bucketName is required'); } const imageStore = this.getEnabledImageStore(); return await this.storageLifecycleStack.getExclusiveExecutionSlot( async () => { if (this.stopRequested) { throw new Error( 'Cannot configure Docker image-store S3 storage after DockerHost.stop() has begun', ); } const candidate = new plugins.smartbucket.SmartBucket(optionsArg); const closeCandidateAndThrow = async ( operationErrorArg: unknown, messageArg: string, ): Promise => { try { await candidate.close(); } catch (cleanupError) { throw new AggregateError( [operationErrorArg, cleanupError], `${messageArg}; candidate cleanup also failed`, ); } throw operationErrorArg; }; let wantedDirectory: plugins.smartbucket.Directory; try { const bucket = await candidate.getBucketByName(bucketName); wantedDirectory = await bucket.getBaseDirectory(); if (optionsArg.directoryPath) { wantedDirectory = await wantedDirectory.getSubDirectoryByName( optionsArg.directoryPath, ); } } catch (error) { return await closeCandidateAndThrow( error, 'Docker image-store S3 setup failed', ); } if (this.stopRequested) { return await closeCandidateAndThrow( new Error( 'Cannot configure Docker image-store S3 storage after DockerHost.stop() has begun', ), 'Docker image-store S3 setup was superseded by host shutdown', ); } const previous = this.smartBucket; if (previous) { let previousCloseFailed = false; let previousCloseError: unknown; try { await previous.close(); } catch (error) { previousCloseFailed = true; previousCloseError = error; } finally { if (this.smartBucket === previous) { this.smartBucket = undefined; } } if (previousCloseFailed) { return await closeCandidateAndThrow( previousCloseError, 'Docker image-store S3 replacement failed while closing the prior client', ); } } if (this.stopRequested) { return await closeCandidateAndThrow( new Error( 'Cannot configure Docker image-store S3 storage after DockerHost.stop() has begun', ), 'Docker image-store S3 replacement was superseded by host shutdown', ); } imageStore.options.bucketDir = wantedDirectory; this.smartBucket = candidate; }, ); } private getEnabledImageStore(): DockerImageStore { if (!this.imageStore) { throw new Error('Docker image store is disabled for this host'); } return this.imageStore; } private encodeRegistryAuth(authArg: interfaces.TDockerRegistryAuth): string { if ('identitytoken' in authArg) { if (!authArg.identitytoken) { throw new TypeError('Docker registry identity token must be nonempty'); } } else if (!authArg.username || !authArg.password || !authArg.serveraddress) { throw new TypeError( 'Docker registry credentials require nonempty username, password, and serveraddress', ); } return plugins.buffer.Buffer.from( JSON.stringify(authArg), 'utf8', ).toString('base64url'); } }