import * as plugins from './plugins.js'; import type * as interfaces from './interfaces/index.js'; import type { DockerHost, IHijackedStreamingResponse } from './classes.host.js'; import { DockerResource } from './classes.base.js'; import { assertDockerContainerId, assertDockerImageId, assertDockerNetworkId, assertDockerResponseStatus, assertDockerResponseStatusOneOf, assertNonemptyDockerString, assertNonnegativeInteger, assertNonRootDockerUser, buildDockerFilterRoute, formatDockerResponseBody, rollbackFailedDockerCreation, } from './helpers.docker.js'; import { parseDockerStartedAtNanoseconds } from './helpers.timestamp.js'; const defaultExecTimeoutMs = 30_000; const defaultExecMaxOutputBytes = 1024 * 1024; const maxDockerLogErrorBytes = 64 * 1024; const dockerNamedVolumePattern = /^[A-Za-z0-9][A-Za-z0-9_.-]*$/; const dockerContainerNamePattern = /^[A-Za-z0-9][A-Za-z0-9_.-]{0,127}$/; const containerRuntimeStatuses = new Set([ 'created', 'running', 'paused', 'restarting', 'removing', 'exited', 'dead', ]); interface ICollectedExecOutput { stdout: Buffer; stderr: Buffer; } export class DockerContainer extends DockerResource { public static async _list( dockerHostArg: DockerHost, optionsArg: interfaces.IContainerListOptions = {}, ): Promise { const includeAll = optionsArg.all ?? true; let route = `/containers/json?all=${includeAll ? 'true' : 'false'}`; route = buildDockerFilterRoute(route, optionsArg.filters); const response = await dockerHostArg.request('GET', route); assertDockerResponseStatus('Docker container list', 200, response); if (!Array.isArray(response.body)) { throw new Error( `Docker container list failed: expected response body to be an array; response body: ${formatDockerResponseBody(response.body)}`, ); } return response.body.map( (containerObjectArg: Record) => new DockerContainer(dockerHostArg, containerObjectArg), ); } public static async _fromId( dockerHostArg: DockerHost, containerIdArg: string, optionsArg: interfaces.IContainerReadOptions = {}, ): Promise { assertDockerContainerId(containerIdArg); const requestOptions: interfaces.IContainerReadOptions = { ...(optionsArg.timeoutMs === undefined ? {} : { timeoutMs: optionsArg.timeoutMs }), ...(optionsArg.signal === undefined ? {} : { signal: optionsArg.signal }), }; const response = await dockerHostArg.request( 'GET', `/containers/${encodeURIComponent(containerIdArg)}/json`, {}, requestOptions, ); if (response.statusCode === 404) { return undefined; } assertDockerResponseStatus( `Docker container inspect "${containerIdArg}"`, 200, response, ); const container = new DockerContainer(dockerHostArg, response.body); if (container.Id !== containerIdArg) { throw new Error( `Docker container inspect "${containerIdArg}" failed: inspected Id "${container.Id}" does not equal the requested immutable ID`, ); } return container; } public static async _create( dockerHostArg: DockerHost, descriptorArg: interfaces.IContainerCreationDescriptor, ): Promise { DockerContainer.validateCreationDescriptor(descriptorArg); if (descriptorArg.allowRootOnRootless) { const info = await dockerHostArg.info(); if (!info.SecurityOptions.includes('name=rootless')) { throw new Error('Docker container root user requires a verified rootless daemon'); } } const exposedPorts: Record> = {}; const portBindings: Record< string, Array<{ HostIp: interfaces.TLoopbackHostIp; HostPort: string }> > = {}; for (const portBinding of descriptorArg.portBindings || []) { const portKey = `${portBinding.containerPort}/${portBinding.protocol || 'tcp'}`; exposedPorts[portKey] = {}; (portBindings[portKey] ||= []).push({ HostIp: portBinding.hostIp, HostPort: portBinding.hostPort === undefined ? '' : String(portBinding.hostPort), }); } const mounts: Array> = []; for (const bindMount of descriptorArg.bindMounts || []) { mounts.push({ Type: 'bind', Source: bindMount.source, Target: bindMount.target, ReadOnly: bindMount.readOnly ?? false, ...(bindMount.createMountpoint === undefined ? {} : { BindOptions: { CreateMountpoint: bindMount.createMountpoint } }), }); } for (const volumeMount of descriptorArg.namedVolumeMounts || []) { mounts.push({ Type: 'volume', Source: volumeMount.source, Target: volumeMount.target, ReadOnly: volumeMount.readOnly ?? false, }); } for (const tmpfsMount of descriptorArg.tmpfsMounts || []) { mounts.push({ Type: 'tmpfs', Target: tmpfsMount.target, ReadOnly: false, TmpfsOptions: { ...(tmpfsMount.sizeBytes === undefined ? {} : { SizeBytes: tmpfsMount.sizeBytes }), ...(tmpfsMount.mode === undefined ? {} : { Mode: tmpfsMount.mode }), }, }); } const endpointsConfig: Record = {}; for (const endpoint of descriptorArg.networkEndpoints || []) { endpointsConfig[endpoint.networkId] = { Aliases: endpoint.aliases || [], }; } const response = await dockerHostArg.request( 'POST', `/containers/create?name=${encodeURIComponent(descriptorArg.name)}`, { Image: descriptorArg.imageId, User: descriptorArg.user, Cmd: descriptorArg.command, ...(descriptorArg.entrypoint === undefined ? {} : { Entrypoint: descriptorArg.entrypoint }), ...(descriptorArg.workingDirectory === undefined ? {} : { WorkingDir: descriptorArg.workingDirectory }), ...(descriptorArg.tty === undefined ? {} : { Tty: descriptorArg.tty }), ...(descriptorArg.openStdin === undefined ? {} : { OpenStdin: descriptorArg.openStdin }), ...(descriptorArg.stdinOnce === undefined ? {} : { StdinOnce: descriptorArg.stdinOnce }), ...(descriptorArg.stopSignal === undefined ? {} : { StopSignal: descriptorArg.stopSignal }), ...(descriptorArg.stopTimeout === undefined ? {} : { StopTimeout: descriptorArg.stopTimeout }), ...(descriptorArg.hostname === undefined ? {} : { Hostname: descriptorArg.hostname }), Env: Object.entries(descriptorArg.env || {}).map( ([keyArg, valueArg]) => `${keyArg}=${valueArg}`, ), Labels: descriptorArg.labels || {}, ExposedPorts: exposedPorts, ...(descriptorArg.healthcheck ? { Healthcheck: DockerContainer.createHealthcheck(descriptorArg.healthcheck) } : {}), HostConfig: { ReadonlyRootfs: descriptorArg.readOnlyRootFilesystem ?? false, Mounts: mounts, PortBindings: portBindings, ...(descriptorArg.init === undefined ? {} : { Init: descriptorArg.init }), ...(descriptorArg.memoryBytes === undefined ? {} : { Memory: descriptorArg.memoryBytes }), ...(descriptorArg.memorySwapBytes === undefined ? {} : { MemorySwap: descriptorArg.memorySwapBytes }), ...(descriptorArg.nanoCpus === undefined ? {} : { NanoCpus: descriptorArg.nanoCpus }), ...(descriptorArg.pidsLimit === undefined ? {} : { PidsLimit: descriptorArg.pidsLimit }), ...(descriptorArg.shmSize === undefined ? {} : { ShmSize: descriptorArg.shmSize }), ...(descriptorArg.logDriver === undefined ? {} : { LogConfig: { Type: descriptorArg.logDriver } }), ...(descriptorArg.networkMode === undefined ? {} : { NetworkMode: descriptorArg.networkMode }), }, ...(descriptorArg.networkMode === 'none' ? {} : { NetworkingConfig: { EndpointsConfig: endpointsConfig, }, }), }, ); assertDockerResponseStatus('Docker container create', 201, response); const createdId = response.body?.Id; try { assertDockerContainerId(createdId, 'Docker container create response body.Id'); } catch (error) { const verificationError = new Error( `Docker container create failed: ${error instanceof Error ? error.message : formatDockerResponseBody(error)}; response body: ${formatDockerResponseBody(response.body)}`, ); return await rollbackFailedDockerCreation( `Docker container "${descriptorArg.name}"`, verificationError, async () => { const rollbackResponse = await dockerHostArg.request( 'DELETE', `/containers/${encodeURIComponent(descriptorArg.name)}?force=true&v=true`, ); assertDockerResponseStatusOneOf( `Docker container rollback remove "${descriptorArg.name}"`, [204, 404], rollbackResponse, ); }, ); } try { const createdContainer = await DockerContainer._fromId( dockerHostArg, createdId, ); if (!createdContainer) { throw new Error( `Docker container create failed: direct inspect did not find created ID "${createdId}"`, ); } if (createdContainer.Image !== descriptorArg.imageId) { throw new Error( `Docker container create failed: inspected image ID "${createdContainer.Image}" does not equal requested image ID "${descriptorArg.imageId}"`, ); } createdContainer.ImageReference = descriptorArg.imageReference; return createdContainer; } catch (verificationError) { return await rollbackFailedDockerCreation( `Docker container "${createdId}"`, verificationError, async () => { const rollbackResponse = await dockerHostArg.request( 'DELETE', `/containers/${encodeURIComponent(createdId)}?force=true&v=true`, ); assertDockerResponseStatusOneOf( `Docker container rollback remove "${createdId}"`, [204, 404], rollbackResponse, ); }, ); } } private static validateCreationDescriptor( descriptorArg: interfaces.IContainerCreationDescriptor, ): void { assertNonemptyDockerString(descriptorArg.name, 'Docker container name'); assertDockerImageId(descriptorArg.imageId, 'Docker container image ID'); if (descriptorArg.allowRootOnRootless === undefined) { assertNonRootDockerUser(descriptorArg.user, 'Docker container user'); } else { if (descriptorArg.allowRootOnRootless !== true) { throw new TypeError('Docker container allowRootOnRootless must be true when present'); } assertNonemptyDockerString(descriptorArg.user, 'Docker container root user'); if (!/^(?:root|0)(?::(?:root|0))?$/i.test(descriptorArg.user)) { throw new TypeError('Docker container root opt-in requires an explicit root user'); } } if (!Array.isArray(descriptorArg.command) || descriptorArg.command.length === 0) { throw new TypeError('Docker container command must be a nonempty argv array'); } for (const commandPart of descriptorArg.command) { assertNonemptyDockerString(commandPart, 'Docker container command argument'); } if (descriptorArg.entrypoint !== undefined) { if (!Array.isArray(descriptorArg.entrypoint) || descriptorArg.entrypoint.length === 0) { throw new TypeError('Docker container entrypoint must be a nonempty argv array'); } for (const entrypointPart of descriptorArg.entrypoint) { assertNonemptyDockerString(entrypointPart, 'Docker container entrypoint argument'); } } if (descriptorArg.workingDirectory !== undefined && !descriptorArg.workingDirectory.startsWith('/')) { throw new TypeError('Docker container workingDirectory must be an absolute path'); } for (const [description, value] of [ ['tty', descriptorArg.tty], ['openStdin', descriptorArg.openStdin], ['stdinOnce', descriptorArg.stdinOnce], ['init', descriptorArg.init], ] as const) { if (value !== undefined && typeof value !== 'boolean') { throw new TypeError(`Docker container ${description} must be a boolean`); } } if (descriptorArg.stdinOnce && !descriptorArg.openStdin) { throw new TypeError('Docker container stdinOnce requires openStdin'); } if (descriptorArg.stopSignal !== undefined) { assertNonemptyDockerString(descriptorArg.stopSignal, 'Docker container stopSignal'); } if (descriptorArg.hostname !== undefined) { assertNonemptyDockerString(descriptorArg.hostname, 'Docker container hostname'); if (!/^[A-Za-z0-9](?:[A-Za-z0-9.-]*[A-Za-z0-9])?$/.test(descriptorArg.hostname)) { throw new TypeError('Docker container hostname must be a valid hostname'); } } if (descriptorArg.stopTimeout !== undefined) { assertNonnegativeInteger(descriptorArg.stopTimeout, 'Docker container stopTimeout'); } for (const [description, value] of [ ['memoryBytes', descriptorArg.memoryBytes], ['memorySwapBytes', descriptorArg.memorySwapBytes], ['nanoCpus', descriptorArg.nanoCpus], ['pidsLimit', descriptorArg.pidsLimit], ['shmSize', descriptorArg.shmSize], ] as const) { if (value !== undefined) { assertNonnegativeInteger(value, `Docker container ${description}`); if (value === 0) { throw new TypeError(`Docker container ${description} must be greater than zero`); } } } if (descriptorArg.memorySwapBytes !== undefined && descriptorArg.memoryBytes === undefined) { throw new TypeError('Docker container memorySwapBytes requires memoryBytes'); } if (descriptorArg.memorySwapBytes !== undefined && descriptorArg.memoryBytes !== undefined && descriptorArg.memorySwapBytes < descriptorArg.memoryBytes) { throw new TypeError('Docker container memorySwapBytes must be at least memoryBytes'); } if (descriptorArg.logDriver !== undefined && descriptorArg.logDriver !== 'none') { throw new TypeError('Docker container logDriver only supports none'); } for (const [key, value] of Object.entries(descriptorArg.env || {})) { assertNonemptyDockerString(key, 'Docker container environment key'); if (key.includes('=')) { throw new TypeError('Docker container environment keys must not contain "="'); } if (typeof value !== 'string') { throw new TypeError(`Docker container environment value for "${key}" must be a string`); } } const mountTargets = new Set(); for (const mount of descriptorArg.bindMounts || []) { if (!mount.source.startsWith('/') || !mount.target.startsWith('/')) { throw new TypeError('Docker bind mount source and target must be absolute paths'); } if (mount.createMountpoint !== undefined && typeof mount.createMountpoint !== 'boolean') { throw new TypeError('Docker bind createMountpoint must be a boolean'); } DockerContainer.recordUniqueValue( mountTargets, mount.target, 'Docker mount target', ); } for (const mount of descriptorArg.namedVolumeMounts || []) { assertNonemptyDockerString( mount.source, 'Docker named volume source', ); if ( mount.source !== mount.source.trim() || !dockerNamedVolumePattern.test(mount.source) ) { throw new TypeError( 'Docker named volume source must be a canonical volume name starting with an alphanumeric character and containing only alphanumeric characters, underscore, period, or hyphen', ); } if (!mount.target.startsWith('/')) { throw new TypeError('Docker named volume target must be an absolute path'); } DockerContainer.recordUniqueValue( mountTargets, mount.target, 'Docker mount target', ); } for (const mount of descriptorArg.tmpfsMounts || []) { if (!mount.target.startsWith('/')) { throw new TypeError('Docker tmpfs mount target must be an absolute path'); } DockerContainer.recordUniqueValue( mountTargets, mount.target, 'Docker mount target', ); if (mount.sizeBytes !== undefined) { assertNonnegativeInteger(mount.sizeBytes, 'Docker tmpfs sizeBytes'); } if (mount.mode !== undefined) { assertNonnegativeInteger(mount.mode, 'Docker tmpfs mode'); } } const networkIds = new Set(); for (const endpoint of descriptorArg.networkEndpoints || []) { assertDockerNetworkId(endpoint.networkId, 'Docker network endpoint ID'); DockerContainer.recordUniqueValue( networkIds, endpoint.networkId, 'Docker network endpoint ID', ); for (const alias of endpoint.aliases || []) { assertNonemptyDockerString(alias, 'Docker network alias'); } } if ( descriptorArg.networkMode === 'none' && ( (descriptorArg.networkEndpoints?.length || 0) > 0 || (descriptorArg.portBindings?.length || 0) > 0 ) ) { throw new TypeError( 'Docker networkless containers cannot declare network endpoints or port bindings', ); } for (const binding of descriptorArg.portBindings || []) { DockerContainer.assertPort(binding.containerPort, 'Docker container port'); if (binding.hostPort !== undefined) { DockerContainer.assertPort(binding.hostPort, 'Docker host port'); } if (binding.hostIp !== '127.0.0.1' && binding.hostIp !== '::1') { throw new TypeError( 'Docker standalone port bindings must use loopback 127.0.0.1 or ::1', ); } } if (descriptorArg.healthcheck) { DockerContainer.createHealthcheck(descriptorArg.healthcheck); } } private static recordUniqueValue( seenArg: Set, valueArg: string, descriptionArg: string, ): void { if (seenArg.has(valueArg)) { throw new TypeError(`${descriptionArg} "${valueArg}" is duplicated`); } seenArg.add(valueArg); } private static assertPort(portArg: number, descriptionArg: string): void { if (!Number.isSafeInteger(portArg) || portArg < 1 || portArg > 65535) { throw new TypeError(`${descriptionArg} must be an integer from 1 through 65535`); } } private static createHealthcheck( healthcheckArg: interfaces.IContainerHealthcheck, ): Record { const [testKind] = healthcheckArg.test; const validTest = (testKind === 'NONE' && healthcheckArg.test.length === 1) || (testKind === 'CMD' && healthcheckArg.test.length >= 2) || (testKind === 'CMD-SHELL' && healthcheckArg.test.length === 2); if (!validTest) { throw new TypeError('Docker Healthcheck.test is not a valid exact Docker healthcheck argv'); } const durationFields: Array<[string, number | undefined]> = [ ['Interval', healthcheckArg.intervalMs], ['Timeout', healthcheckArg.timeoutMs], ['StartPeriod', healthcheckArg.startPeriodMs], ['StartInterval', healthcheckArg.startIntervalMs], ]; const result: Record = { Test: healthcheckArg.test }; for (const [dockerField, milliseconds] of durationFields) { if (milliseconds === undefined) { continue; } assertNonnegativeInteger(milliseconds, `Docker Healthcheck ${dockerField} milliseconds`); if (milliseconds > Number.MAX_SAFE_INTEGER / 1_000_000) { throw new TypeError(`Docker Healthcheck ${dockerField} milliseconds is too large`); } result[dockerField] = milliseconds * 1_000_000; } if (healthcheckArg.retries !== undefined) { assertNonnegativeInteger(healthcheckArg.retries, 'Docker Healthcheck retries'); result.Retries = healthcheckArg.retries; } return result; } public Id!: string; public Names!: string[]; public Name?: string; public Image!: string; public ImageID?: string; public ImageReference?: string; public Command?: string; public Created?: number | string; public Ports?: interfaces.TPorts; public Labels?: interfaces.TLabels; public State?: string | Record; public Status?: string; public Config?: Record; public HostConfig?: Record; public NetworkSettings?: Record; public Mounts?: unknown[]; private readonly immutableId: string; constructor( dockerHostArg: DockerHost, dockerContainerObjectArg: Record, ) { super(dockerHostArg); assertDockerContainerId( dockerContainerObjectArg.Id, 'Docker container response Id', ); this.immutableId = dockerContainerObjectArg.Id; Object.assign(this, dockerContainerObjectArg); } public async refresh(): Promise { const updated = await DockerContainer._fromId( this.dockerHost, this.immutableId, ); if (!updated) { throw new Error( `Docker container refresh failed: container "${this.immutableId}" was not found`, ); } Object.assign(this, updated); this.Id = this.immutableId; } public async inspect( optionsArg: interfaces.IContainerReadOptions = {}, ): Promise> { const requestOptions: interfaces.IContainerReadOptions = { ...(optionsArg.timeoutMs === undefined ? {} : { timeoutMs: optionsArg.timeoutMs }), ...(optionsArg.signal === undefined ? {} : { signal: optionsArg.signal }), }; const response = await this.dockerHost.request( 'GET', `/containers/${encodeURIComponent(this.immutableId)}/json`, {}, requestOptions, ); assertDockerResponseStatus( `Docker container inspect "${this.immutableId}"`, 200, response, ); assertDockerContainerId( response.body?.Id, `Docker container inspect "${this.immutableId}" response Id`, ); if (response.body.Id !== this.immutableId) { throw new Error( `Docker container inspect "${this.immutableId}" failed: inspected Id "${response.body.Id}" does not equal the immutable container ID`, ); } Object.assign(this, response.body); this.Id = this.immutableId; return response.body; } /** * Inspect this immutable container ID and validate the state fields used for runtime recovery. * Missing or malformed OOM evidence is an error, never an implicit `false`. */ public async inspectState( optionsArg: interfaces.IContainerReadOptions = {}, ): Promise { const inspection = await this.inspect(optionsArg); return this.parseRuntimeState(inspection); } /** Inspect one exact-ID snapshot and retain the current run's full nanosecond start time. */ public async inspectRunState( optionsArg: interfaces.IContainerReadOptions = {}, ): Promise { const inspection = await this.inspect(optionsArg); const runtime = this.parseRuntimeState(inspection); const state = inspection.State as Record; const startedAt = state.StartedAt; const startedAtUnixNano = parseDockerStartedAtNanoseconds(startedAt).toString(); return { ...runtime, StartedAt: startedAt as string, startedAtUnixNano, }; } private parseRuntimeState(inspectionArg: Record): interfaces.IContainerRuntimeState { const state: unknown = inspectionArg.State; if (typeof state !== 'object' || state === null || Array.isArray(state)) { throw new Error(`Docker container inspect state "${this.immutableId}" is missing or malformed`); } const status: unknown = Reflect.get(state, 'Status'); const running: unknown = Reflect.get(state, 'Running'); const oomKilled: unknown = Reflect.get(state, 'OOMKilled'); const pid: unknown = Reflect.get(state, 'Pid'); const exitCode: unknown = Reflect.get(state, 'ExitCode'); if ( typeof status !== 'string' || !containerRuntimeStatuses.has(status as interfaces.IContainerRuntimeState['Status']) || typeof running !== 'boolean' || typeof oomKilled !== 'boolean' || !Number.isSafeInteger(pid) || (pid as number) < 0 || !Number.isSafeInteger(exitCode) ) { throw new Error(`Docker container inspect state "${this.immutableId}" is missing or malformed`); } return { Status: status as interfaces.IContainerRuntimeState['Status'], Running: running, OOMKilled: oomKilled, Pid: pid as number, ExitCode: exitCode as number, }; } /** Renames this exact container and verifies the result through its immutable ID. */ public async rename(nextNameArg: string): Promise { assertNonemptyDockerString(nextNameArg, 'Docker container name'); if ( nextNameArg !== nextNameArg.trim() || !dockerContainerNamePattern.test(nextNameArg) ) { throw new TypeError( 'Docker container name must start with an alphanumeric character and contain only alphanumeric characters, underscore, period, or hyphen', ); } const response = await this.dockerHost.request( 'POST', `/containers/${encodeURIComponent(this.immutableId)}/rename?name=${encodeURIComponent(nextNameArg)}`, ); assertDockerResponseStatus( `Docker container rename "${this.immutableId}"`, 204, response, ); await this.refresh(); const inspectedName = typeof this.Name === 'string' ? this.Name.replace(/^\/+/, '') : ''; if (inspectedName !== nextNameArg) { throw new Error( `Docker container rename "${this.immutableId}" failed: inspected name "${inspectedName}" does not equal requested name "${nextNameArg}"`, ); } } public async start(): Promise { const response = await this.dockerHost.request( 'POST', `/containers/${encodeURIComponent(this.immutableId)}/start`, ); assertDockerResponseStatusOneOf( `Docker container start "${this.immutableId}"`, [204, 304], response, ); await this.refresh(); return response.statusCode === 204 ? 'started' : 'already-running'; } public async stop( optionsArg: interfaces.IContainerStopOptions = {}, ): Promise { const queryParams = new URLSearchParams(); if (optionsArg.timeoutSeconds !== undefined) { assertNonnegativeInteger( optionsArg.timeoutSeconds, 'Docker container stop timeoutSeconds', ); queryParams.set('t', String(optionsArg.timeoutSeconds)); } if (optionsArg.signal !== undefined) { assertNonemptyDockerString(optionsArg.signal, 'Docker container stop signal'); queryParams.set('signal', optionsArg.signal); } const query = queryParams.size > 0 ? `?${queryParams.toString()}` : ''; const response = await this.dockerHost.request( 'POST', `/containers/${encodeURIComponent(this.immutableId)}/stop${query}`, ); assertDockerResponseStatusOneOf( `Docker container stop "${this.immutableId}"`, [204, 304], response, ); await this.refresh(); return response.statusCode === 204 ? 'stopped' : 'already-stopped'; } /** Wait indefinitely by default; pass a signal or deadline to cancel. */ public async wait( conditionArg: interfaces.TContainerWaitCondition = 'not-running', optionsArg: interfaces.IContainerReadOptions = {}, ): Promise { if (conditionArg !== 'not-running' && conditionArg !== 'next-exit' && conditionArg !== 'removed') { throw new TypeError('Docker container wait condition must be not-running, next-exit, or removed'); } const timeoutMs = optionsArg.timeoutMs ?? 0; assertNonnegativeInteger(timeoutMs, 'Docker container wait timeoutMs'); const response = await this.dockerHost.request( 'POST', `/containers/${encodeURIComponent(this.immutableId)}/wait?condition=${encodeURIComponent(conditionArg)}`, {}, { timeoutMs, signal: optionsArg.signal }, ); assertDockerResponseStatus(`Docker container wait "${this.immutableId}"`, 200, response); const body: unknown = response.body; if (typeof body !== 'object' || body === null) { throw new Error('Docker container wait failed: expected an integer StatusCode'); } const exitCode: unknown = Reflect.get(body, 'StatusCode'); if (typeof exitCode !== 'number' || !Number.isSafeInteger(exitCode)) { throw new Error('Docker container wait failed: expected an integer StatusCode'); } const waitError: unknown = Reflect.get(body, 'Error'); if (waitError === undefined || waitError === null) { return { exitCode }; } if (typeof waitError !== 'object' || waitError === null) { throw new Error('Docker container wait failed: expected Error.Message to be a string'); } const errorMessage: unknown = Reflect.get(waitError, 'Message'); if (typeof errorMessage !== 'string') { throw new Error('Docker container wait failed: expected Error.Message to be a string'); } return { exitCode, errorMessage }; } /** Send a signal to this exact container. */ public async kill(signalArg = 'SIGKILL'): Promise { assertNonemptyDockerString(signalArg, 'Docker container kill signal'); const response = await this.dockerHost.request( 'POST', `/containers/${encodeURIComponent(this.immutableId)}/kill?signal=${encodeURIComponent(signalArg)}`, ); assertDockerResponseStatus(`Docker container kill "${this.immutableId}"`, 204, response); await this.refresh(); } /** Resize this container's TTY in rows and columns. */ public async resize(heightArg: number, widthArg: number): Promise { DockerContainer.assertTerminalSize(heightArg, widthArg); const response = await this.dockerHost.request( 'POST', `/containers/${encodeURIComponent(this.immutableId)}/resize?h=${heightArg}&w=${widthArg}`, ); assertDockerResponseStatus(`Docker container resize "${this.immutableId}"`, 200, response); } /** Resize an exec after proving it belongs to this exact container. */ public async resizeExec(execIdArg: string, heightArg: number, widthArg: number): Promise { assertNonemptyDockerString(execIdArg, 'Docker exec ID'); DockerContainer.assertTerminalSize(heightArg, widthArg); await this.inspectExec(execIdArg); const response = await this.dockerHost.request( 'POST', `/exec/${encodeURIComponent(execIdArg)}/resize?h=${heightArg}&w=${widthArg}`, ); assertDockerResponseStatus(`Docker exec resize "${execIdArg}"`, 200, response); } private static assertTerminalSize(heightArg: number, widthArg: number): void { if (!Number.isSafeInteger(heightArg) || heightArg <= 0 || !Number.isSafeInteger(widthArg) || widthArg <= 0) { throw new TypeError('Docker TTY height and width must be positive safe integers'); } } public async remove( optionsArg: interfaces.IContainerRemoveOptions = {}, ): Promise { const queryParams = new URLSearchParams(); if (optionsArg.force) { queryParams.set('force', 'true'); } if (optionsArg.removeAnonymousVolumes) { queryParams.set('v', 'true'); } const query = queryParams.size > 0 ? `?${queryParams.toString()}` : ''; const response = await this.dockerHost.request( 'DELETE', `/containers/${encodeURIComponent(this.immutableId)}${query}`, ); assertDockerResponseStatusOneOf( `Docker container remove "${this.immutableId}"`, [204, 404], response, ); return response.statusCode === 204 ? 'removed' : 'already-removed'; } public async logs( optionsArg: interfaces.IContainerLogsOptions = {}, ): Promise { const queryParams = new URLSearchParams(); queryParams.set('stdout', optionsArg.stdout === false ? '0' : '1'); queryParams.set('stderr', optionsArg.stderr === false ? '0' : '1'); if (optionsArg.timestamps) queryParams.set('timestamps', '1'); if (optionsArg.tail !== undefined) queryParams.set('tail', String(optionsArg.tail)); if (optionsArg.since !== undefined) queryParams.set('since', String(optionsArg.since)); const response = await this.dockerHost.requestStreaming( 'GET', `/containers/${encodeURIComponent(this.immutableId)}/logs?${queryParams.toString()}`, undefined, undefined, { timeoutMs: defaultExecTimeoutMs }, ); if (!(response instanceof plugins.stream.Readable)) { throw new Error( `Docker container logs "${this.immutableId}" failed: actual status ${response.statusCode}; response body: ${response.body}`, ); } const statusCode = Reflect.get(response, 'statusCode'); if (statusCode !== 200) { const errorBody = await DockerContainer.collectStream( response, maxDockerLogErrorBytes, `Docker container logs "${this.immutableId}" error response`, ); throw new Error( `Docker container logs "${this.immutableId}" failed: actual status ${String(statusCode)}; response body: ${errorBody.toString('utf8')}`, ); } const demultiplexed = DockerContainer.createDemultiplexedStream(response); const logs = await DockerContainer.collectStream( demultiplexed, undefined, `Docker container logs "${this.immutableId}"`, ); return logs.toString('utf8'); } public async stats( optionsArg: interfaces.IContainerStatsOptions = {}, ): Promise { const queryParams = new URLSearchParams(); queryParams.set('stream', optionsArg.stream ? '1' : '0'); if (optionsArg.oneShot) queryParams.set('one-shot', '1'); const response = await this.dockerHost.request( 'GET', `/containers/${encodeURIComponent(this.immutableId)}/stats?${queryParams.toString()}`, ); assertDockerResponseStatus( `Docker container stats "${this.immutableId}"`, 200, response, ); return response.body; } public async streamLogs( optionsArg: interfaces.IContainerStreamLogsOptions = {}, ): Promise { const queryParams = new URLSearchParams(); queryParams.set('stdout', optionsArg.stdout === false ? '0' : '1'); queryParams.set('stderr', optionsArg.stderr === false ? '0' : '1'); queryParams.set('follow', '1'); if (optionsArg.timestamps) queryParams.set('timestamps', '1'); if (optionsArg.tail !== undefined) queryParams.set('tail', String(optionsArg.tail)); if (optionsArg.since !== undefined) queryParams.set('since', String(optionsArg.since)); const response = await this.dockerHost.requestStreaming( 'GET', `/containers/${encodeURIComponent(this.immutableId)}/logs?${queryParams.toString()}`, undefined, undefined, { timeoutMs: 0 }, ); if (!(response instanceof plugins.stream.Readable)) { throw new Error( `Docker container log stream "${this.immutableId}" failed: actual status ${response.statusCode}; response body: ${response.body}`, ); } if (!optionsArg.demux) { return response; } return DockerContainer.createDemultiplexedStream(response); } private static createDemultiplexedStream( rawStreamArg: plugins.stream.Readable, ): plugins.stream.Readable { let pending = Buffer.alloc(0); const demuxTransform = new plugins.stream.Transform({ transform(chunkArg: Buffer, _encodingArg, callbackArg) { pending = pending.length ? Buffer.concat([pending, chunkArg]) : Buffer.from(chunkArg); while (pending.length >= 8) { const validHeader = (pending[0] === 1 || pending[0] === 2) && pending[1] === 0 && pending[2] === 0 && pending[3] === 0; if (!validHeader) { this.push(pending); pending = Buffer.alloc(0); break; } const frameLength = pending.readUInt32BE(4); if (pending.length < frameLength + 8) { break; } this.push(pending.subarray(8, frameLength + 8)); pending = pending.subarray(frameLength + 8); } callbackArg(); }, flush(callbackArg) { if (pending.length > 0) { this.push(pending); } callbackArg(); }, }); rawStreamArg.on('error', (errorArg) => demuxTransform.destroy(errorArg)); demuxTransform.on('close', () => rawStreamArg.destroy()); return rawStreamArg.pipe(demuxTransform); } private static async collectStream( streamArg: plugins.stream.Readable, maxBytesArg: number | undefined, contextArg: string, ): Promise { const chunks: Buffer[] = []; let totalBytes = 0; try { for await (const chunk of streamArg) { const buffer = Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk); totalBytes += buffer.length; if (maxBytesArg !== undefined && totalBytes > maxBytesArg) { throw new Error(`${contextArg} exceeded ${maxBytesArg} bytes`); } chunks.push(buffer); } return Buffer.concat(chunks, totalBytes); } finally { streamArg.destroy(); } } public async attach( optionsArg: interfaces.IContainerAttachOptions = {}, ): Promise<{ stream: plugins.stream.Duplex; close: () => Promise; }> { const timeoutMs = optionsArg.timeoutMs ?? defaultExecTimeoutMs; assertNonnegativeInteger(timeoutMs, 'Docker container attach timeoutMs'); if (optionsArg.signal?.aborted) { throw DockerContainer.createAttachAbortError(optionsArg.signal); } const queryParams = new URLSearchParams(); queryParams.set('stream', optionsArg.stream === false ? '0' : '1'); queryParams.set('stdin', optionsArg.stdin ? '1' : '0'); queryParams.set('stdout', optionsArg.stdout === false ? '0' : '1'); queryParams.set('stderr', optionsArg.stderr === false ? '0' : '1'); if (optionsArg.logs) queryParams.set('logs', '1'); if (optionsArg.detachKeys !== undefined) { assertNonemptyDockerString(optionsArg.detachKeys, 'Docker container attach detachKeys'); queryParams.set('detachKeys', optionsArg.detachKeys); } const response = await this.dockerHost.requestHijackedStreaming( 'POST', `/containers/${encodeURIComponent(this.immutableId)}/attach?${queryParams.toString()}`, {}, { timeoutMs, signal: optionsArg.signal }, ); try { assertDockerResponseStatusOneOf( `Docker container attach "${this.immutableId}"`, [101, 200], { statusCode: response.statusCode, body: '' }, ); } catch (error) { response.stream.destroy(); try { await response.close(); } catch (cleanupError) { throw new AggregateError( [error, cleanupError], `Docker container attach "${this.immutableId}" failed and its hijacked stream cleanup also failed`, ); } throw error; } let closePromise: Promise | undefined; const close = async (): Promise => { if (!closePromise) { optionsArg.signal?.removeEventListener('abort', handleTerminalEvent); response.stream.off('end', handleTerminalEvent); response.stream.off('close', handleTerminalEvent); response.stream.off('error', handleTerminalEvent); closePromise = Promise.resolve().then(async () => { response.stream.destroy(); await response.close(); }); } await closePromise; }; const handleTerminalEvent = () => { void close().catch(() => {}); }; optionsArg.signal?.addEventListener('abort', handleTerminalEvent, { once: true }); response.stream.once('end', handleTerminalEvent); response.stream.once('close', handleTerminalEvent); response.stream.once('error', handleTerminalEvent); if (optionsArg.signal?.aborted) { await close(); throw DockerContainer.createAttachAbortError(optionsArg.signal); } return { stream: response.stream, close }; } private static createAttachAbortError(signalArg: AbortSignal): Error { const reason = signalArg.reason; const suffix = reason === undefined ? '' : `: ${reason instanceof Error ? reason.message : String(reason)}`; const error = new Error(`Docker container attach was aborted${suffix}`); error.name = 'AbortError'; return error; } /** * Runs a bounded argv command, collects multiplexed stdout/stderr, inspects * the exact exit status, and closes the hijacked transport on every path. */ public async exec( commandArg: interfaces.TContainerCommand, optionsArg: interfaces.IContainerExecOptions = {}, ): Promise { if (!Array.isArray(commandArg) || commandArg.length === 0) { throw new TypeError('Docker exec command must be a nonempty argv array'); } for (const commandPart of commandArg) { assertNonemptyDockerString(commandPart, 'Docker exec command argument'); } const timeoutMs = optionsArg.timeoutMs ?? defaultExecTimeoutMs; const maxOutputBytes = optionsArg.maxOutputBytes ?? defaultExecMaxOutputBytes; if (!Number.isSafeInteger(timeoutMs) || timeoutMs <= 0) { throw new TypeError('Docker exec timeoutMs must be a positive safe integer'); } if (!Number.isSafeInteger(maxOutputBytes) || maxOutputBytes <= 0) { throw new TypeError('Docker exec maxOutputBytes must be a positive safe integer'); } if (optionsArg.user !== undefined) { assertNonRootDockerUser(optionsArg.user, 'Docker exec user override'); } const startedAt = Date.now(); const createResponse = await DockerContainer.awaitExecStep( this.dockerHost.request( 'POST', `/containers/${encodeURIComponent(this.immutableId)}/exec`, { Cmd: commandArg, AttachStdin: false, AttachStdout: true, AttachStderr: true, Tty: false, Env: Object.entries(optionsArg.env || {}).map( ([keyArg, valueArg]) => `${keyArg}=${valueArg}`, ), ...(optionsArg.workingDirectory === undefined ? {} : { WorkingDir: optionsArg.workingDirectory }), ...(optionsArg.user === undefined ? {} : { User: optionsArg.user }), }, { timeoutMs: DockerContainer.remainingExecTime( `for container "${this.immutableId}"`, startedAt, timeoutMs, ), }, ), `for container "${this.immutableId}"`, startedAt, timeoutMs, ); assertDockerResponseStatus( `Docker exec create for container "${this.immutableId}"`, 201, createResponse, ); const execId = createResponse.body?.Id; if (typeof execId !== 'string' || execId.length === 0) { throw new Error( `Docker exec create for container "${this.immutableId}" failed: expected response body.Id to be a nonempty string; response body: ${formatDockerResponseBody(createResponse.body)}`, ); } let startResponse: IHijackedStreamingResponse | undefined; let operationError: unknown; try { startResponse = await this.dockerHost.requestHijackedStreaming( 'POST', `/exec/${encodeURIComponent(execId)}/start`, { Detach: false, Tty: false }, { timeoutMs: DockerContainer.remainingExecTime( execId, startedAt, timeoutMs, ), }, ); assertDockerResponseStatusOneOf( `Docker exec start "${execId}"`, [101, 200], { statusCode: startResponse.statusCode, body: '' }, ); const output = await DockerContainer.collectExecOutput( execId, startResponse, maxOutputBytes, DockerContainer.remainingExecTime(execId, startedAt, timeoutMs), timeoutMs, ); const inspect = await this.waitForExecCompletion( execId, startedAt, timeoutMs, ); return { execId, stdout: output.stdout.toString('utf8'), stderr: output.stderr.toString('utf8'), exitCode: inspect.ExitCode, inspect, }; } catch (error) { operationError = error; throw error; } finally { if (startResponse) { startResponse.stream.destroy(); try { await startResponse.close(); } catch (cleanupError) { if (operationError !== undefined) { throw new AggregateError( [operationError, cleanupError], `Docker exec "${execId}" failed and its hijacked stream cleanup also failed`, ); } throw cleanupError; } } } } /** * Starts a caller-owned bidirectional exec session without changing the * collected, bounded semantics of exec(). */ public async execInteractive( commandArg: interfaces.TContainerCommand, optionsArg: interfaces.IContainerInteractiveExecOptions = {}, ): Promise { if (!Array.isArray(commandArg) || commandArg.length === 0) { throw new TypeError('Docker interactive exec command must be a nonempty argv array'); } for (const commandPart of commandArg) { assertNonemptyDockerString(commandPart, 'Docker interactive exec command argument'); } for (const [key, value] of Object.entries(optionsArg.env || {})) { assertNonemptyDockerString(key, 'Docker interactive exec environment key'); if (key.includes('=')) { throw new TypeError('Docker interactive exec environment keys must not contain "="'); } if (typeof value !== 'string') { throw new TypeError( `Docker interactive exec environment value for "${key}" must be a string`, ); } } if (optionsArg.workingDirectory !== undefined) { assertNonemptyDockerString( optionsArg.workingDirectory, 'Docker interactive exec working directory', ); } if (optionsArg.user !== undefined) { assertNonRootDockerUser(optionsArg.user, 'Docker interactive exec user override'); } const timeoutMs = optionsArg.timeoutMs ?? defaultExecTimeoutMs; if (!Number.isSafeInteger(timeoutMs) || timeoutMs <= 0) { throw new TypeError( 'Docker interactive exec timeoutMs must be a positive safe integer', ); } if (optionsArg.signal?.aborted) { throw DockerContainer.createInteractiveExecAbortError(optionsArg.signal); } const tty = optionsArg.tty ?? true; if (optionsArg.detachKeys !== undefined) { assertNonemptyDockerString(optionsArg.detachKeys, 'Docker interactive exec detachKeys'); } if (optionsArg.consoleSize !== undefined) { if (!tty || !Array.isArray(optionsArg.consoleSize) || optionsArg.consoleSize.length !== 2) { throw new TypeError('Docker interactive exec consoleSize requires a TTY and [height, width]'); } DockerContainer.assertTerminalSize(optionsArg.consoleSize[0], optionsArg.consoleSize[1]); } const createResponse = await this.dockerHost.request( 'POST', `/containers/${encodeURIComponent(this.immutableId)}/exec`, { Cmd: commandArg, AttachStdin: true, AttachStdout: true, AttachStderr: true, Tty: tty, ...(optionsArg.detachKeys === undefined ? {} : { DetachKeys: optionsArg.detachKeys }), ...(optionsArg.consoleSize === undefined ? {} : { ConsoleSize: optionsArg.consoleSize }), Env: Object.entries(optionsArg.env || {}).map( ([keyArg, valueArg]) => `${keyArg}=${valueArg}`, ), ...(optionsArg.workingDirectory === undefined ? {} : { WorkingDir: optionsArg.workingDirectory }), ...(optionsArg.user === undefined ? {} : { User: optionsArg.user }), }, { timeoutMs }, ); assertDockerResponseStatus( `Docker interactive exec create for container "${this.immutableId}"`, 201, createResponse, ); const execId = createResponse.body?.Id; if (typeof execId !== 'string' || execId.length === 0) { throw new Error( `Docker interactive exec create for container "${this.immutableId}" failed: expected response body.Id to be a nonempty string; response body: ${formatDockerResponseBody(createResponse.body)}`, ); } const startResponse = await this.dockerHost.requestHijackedStreaming( 'POST', `/exec/${encodeURIComponent(execId)}/start`, { Detach: false, Tty: tty, ...(optionsArg.consoleSize === undefined ? {} : { ConsoleSize: optionsArg.consoleSize }), }, { timeoutMs, signal: optionsArg.signal }, ); try { assertDockerResponseStatusOneOf( `Docker interactive exec start "${execId}"`, [101, 200], { statusCode: startResponse.statusCode, body: '' }, ); } catch (error) { startResponse.stream.destroy(); try { await startResponse.close(); } catch (cleanupError) { throw new AggregateError( [error, cleanupError], `Docker interactive exec "${execId}" start failed and its hijacked stream cleanup also failed`, ); } throw error; } let closePromise: Promise | undefined; const close = async (): Promise => { if (!closePromise) { optionsArg.signal?.removeEventListener('abort', handleTerminalEvent); startResponse.stream.off('end', handleTerminalEvent); startResponse.stream.off('close', handleTerminalEvent); startResponse.stream.off('error', handleTerminalEvent); closePromise = Promise.resolve().then(async () => { startResponse.stream.destroy(); await startResponse.close(); }); } await closePromise; }; const handleTerminalEvent = () => { void close().catch(() => {}); }; optionsArg.signal?.addEventListener('abort', handleTerminalEvent, { once: true }); startResponse.stream.once('end', handleTerminalEvent); startResponse.stream.once('close', handleTerminalEvent); startResponse.stream.once('error', handleTerminalEvent); if (optionsArg.signal?.aborted) { await close(); throw DockerContainer.createInteractiveExecAbortError(optionsArg.signal); } return { execId, stream: startResponse.stream, close, inspect: async () => await this.inspectExec(execId), }; } private static createInteractiveExecAbortError(signalArg: AbortSignal): Error { const reason = signalArg.reason; const suffix = reason === undefined ? '' : `: ${reason instanceof Error ? reason.message : String(reason)}`; const error = new Error(`Docker interactive exec was aborted${suffix}`); error.name = 'AbortError'; return error; } private async inspectExec(execIdArg: string): Promise { const response = await this.dockerHost.request( 'GET', `/exec/${encodeURIComponent(execIdArg)}/json`, ); assertDockerResponseStatus(`Docker exec inspect "${execIdArg}"`, 200, response); if (typeof response.body?.Running !== 'boolean') { throw new Error( `Docker exec inspect "${execIdArg}" failed: expected response body.Running to be boolean; response body: ${formatDockerResponseBody(response.body)}`, ); } if (response.body.ID !== execIdArg) { throw new Error( `Docker exec inspect "${execIdArg}" failed: inspected exec ID does not equal the requested exact exec ID; response body: ${formatDockerResponseBody(response.body)}`, ); } if (response.body.ContainerID !== this.immutableId) { throw new Error( `Docker exec inspect "${execIdArg}" failed: inspected container ID does not equal the immutable container ID; response body: ${formatDockerResponseBody(response.body)}`, ); } if (!response.body.Running && !Number.isSafeInteger(response.body.ExitCode)) { throw new Error( `Docker exec inspect "${execIdArg}" failed: expected response body.ExitCode to be an integer; response body: ${formatDockerResponseBody(response.body)}`, ); } return response.body; } private static remainingExecTime( execIdArg: string, startedAtArg: number, timeoutMsArg: number, ): number { const remaining = timeoutMsArg - (Date.now() - startedAtArg); if (remaining <= 0) { throw new Error(`Docker exec "${execIdArg}" timed out after ${timeoutMsArg}ms`); } return remaining; } private static async awaitExecStep( promiseArg: Promise, execDescriptionArg: string, startedAtArg: number, timeoutMsArg: number, ): Promise { const remainingMs = DockerContainer.remainingExecTime( execDescriptionArg, startedAtArg, timeoutMsArg, ); return await new Promise((resolveArg, rejectArg) => { let settled = false; const timeout = setTimeout(() => { settled = true; rejectArg( new Error( `Docker exec ${execDescriptionArg} timed out after ${timeoutMsArg}ms`, ), ); }, remainingMs); void promiseArg.then( (valueArg) => { if (settled) { return; } settled = true; clearTimeout(timeout); resolveArg(valueArg); }, (errorArg) => { if (settled) { return; } settled = true; clearTimeout(timeout); rejectArg(errorArg); }, ); }); } private async waitForExecCompletion( execIdArg: string, startedAtArg: number, timeoutMsArg: number, ): Promise { while (true) { DockerContainer.remainingExecTime(execIdArg, startedAtArg, timeoutMsArg); const response = await DockerContainer.awaitExecStep( this.dockerHost.request( 'GET', `/exec/${encodeURIComponent(execIdArg)}/json`, {}, { timeoutMs: DockerContainer.remainingExecTime( execIdArg, startedAtArg, timeoutMsArg, ), }, ), `"${execIdArg}"`, startedAtArg, timeoutMsArg, ); assertDockerResponseStatus(`Docker exec inspect "${execIdArg}"`, 200, response); const running = response.body?.Running; const exitCode = response.body?.ExitCode; if (typeof running !== 'boolean') { throw new Error( `Docker exec inspect "${execIdArg}" failed: expected response body.Running to be boolean; response body: ${formatDockerResponseBody(response.body)}`, ); } if (!running) { if (!Number.isSafeInteger(exitCode)) { throw new Error( `Docker exec inspect "${execIdArg}" failed: expected response body.ExitCode to be an integer; response body: ${formatDockerResponseBody(response.body)}`, ); } return response.body; } await new Promise((resolve) => { setTimeout( resolve, Math.min( 10, DockerContainer.remainingExecTime( execIdArg, startedAtArg, timeoutMsArg, ), ), ); }); } } private static async collectExecOutput( execIdArg: string, responseArg: IHijackedStreamingResponse, maxOutputBytesArg: number, remainingTimeoutMsArg: number, configuredTimeoutMsArg: number, ): Promise { return await new Promise((resolve, reject) => { const stdoutChunks: Buffer[] = []; const stderrChunks: Buffer[] = []; let pendingHeader = Buffer.alloc(0); let frameBytesRemaining = 0; let frameStreamType: 1 | 2 = 1; let combinedOutputBytes = 0; let settled = false; const contentTypeHeader = responseArg.headers['content-type']; const contentType = Array.isArray(contentTypeHeader) ? contentTypeHeader.join(';') : contentTypeHeader || ''; let streamMode: 'raw' | 'multiplexed' | 'unknown' = contentType.includes('raw-stream') ? 'raw' : contentType.includes('multiplexed-stream') ? 'multiplexed' : 'unknown'; const cleanupListeners = () => { responseArg.stream.off('data', onData); responseArg.stream.off('end', onEnd); responseArg.stream.off('close', onEnd); responseArg.stream.off('error', onError); clearTimeout(timeout); }; const fail = (errorArg: Error) => { if (settled) return; settled = true; cleanupListeners(); reject(errorArg); }; const appendPayload = (streamTypeArg: 1 | 2, payloadArg: Buffer) => { if (payloadArg.length > maxOutputBytesArg - combinedOutputBytes) { fail( new Error( `Docker exec "${execIdArg}" exceeded max combined output of ${maxOutputBytesArg} bytes`, ), ); return false; } combinedOutputBytes += payloadArg.length; (streamTypeArg === 2 ? stderrChunks : stdoutChunks).push( Buffer.from(payloadArg), ); return true; }; const consumeChunk = (chunkArg: Buffer) => { let chunkOffset = 0; while (!settled && chunkOffset < chunkArg.length) { if (frameBytesRemaining > 0) { const payloadLength = Math.min( frameBytesRemaining, chunkArg.length - chunkOffset, ); const payload = chunkArg.subarray( chunkOffset, chunkOffset + payloadLength, ); if (!appendPayload(frameStreamType, payload)) { return; } frameBytesRemaining -= payloadLength; chunkOffset += payloadLength; continue; } const headerBytesNeeded = 8 - pendingHeader.length; const availableBytes = chunkArg.length - chunkOffset; if (availableBytes < headerBytesNeeded) { const headerFragment = chunkArg.subarray(chunkOffset); pendingHeader = pendingHeader.length ? Buffer.concat( [pendingHeader, headerFragment], pendingHeader.length + headerFragment.length, ) : Buffer.from(headerFragment); return; } let header: Buffer; if (pendingHeader.length === 0) { header = chunkArg.subarray(chunkOffset, chunkOffset + 8); } else { const headerFragment = chunkArg.subarray( chunkOffset, chunkOffset + headerBytesNeeded, ); header = Buffer.concat([pendingHeader, headerFragment], 8); pendingHeader = Buffer.alloc(0); } chunkOffset += headerBytesNeeded; const validHeader = (header[0] === 1 || header[0] === 2) && header[1] === 0 && header[2] === 0 && header[3] === 0; if (streamMode === 'unknown') { streamMode = validHeader ? 'multiplexed' : 'raw'; } if (streamMode === 'raw') { if (!appendPayload(1, header)) { return; } if (chunkOffset < chunkArg.length) { appendPayload(1, chunkArg.subarray(chunkOffset)); } return; } if (!validHeader) { fail(new Error(`Docker exec "${execIdArg}" returned a malformed multiplexed stream`)); return; } const frameLength = header.readUInt32BE(4); if (frameLength > maxOutputBytesArg - combinedOutputBytes) { fail( new Error( `Docker exec "${execIdArg}" exceeded max combined output of ${maxOutputBytesArg} bytes`, ), ); return; } frameStreamType = header[0] === 2 ? 2 : 1; frameBytesRemaining = frameLength; } }; const onData = (chunkArg: Buffer | Uint8Array | string) => { const chunk = typeof chunkArg === 'string' ? Buffer.from(chunkArg) : Buffer.from(chunkArg); if (streamMode === 'raw') { appendPayload(1, chunk); return; } consumeChunk(chunk); }; const onEnd = () => { if (settled) return; if (streamMode === 'unknown' && pendingHeader.length > 0) { streamMode = 'raw'; if (!appendPayload(1, pendingHeader)) return; pendingHeader = Buffer.alloc(0); } if ( streamMode === 'multiplexed' && (pendingHeader.length !== 0 || frameBytesRemaining !== 0) ) { fail(new Error(`Docker exec "${execIdArg}" returned a truncated multiplexed stream`)); return; } settled = true; cleanupListeners(); resolve({ stdout: Buffer.concat(stdoutChunks), stderr: Buffer.concat(stderrChunks), }); }; const onError = (errorArg: Error) => fail(errorArg); const timeout = setTimeout(() => { fail( new Error( `Docker exec "${execIdArg}" timed out after ${configuredTimeoutMsArg}ms`, ), ); }, remainingTimeoutMsArg); responseArg.stream.on('data', onData); responseArg.stream.once('end', onEnd); responseArg.stream.once('close', onEnd); responseArg.stream.once('error', onError); }); } }