import type { DockerHost } from './classes.host.js'; import type { TLabels } from './interfaces/label.js'; import type { IDockerResourceFileTarget } from './interfaces/resource.js'; import type { IDockerGlobalJobCompletionProof, IDockerGlobalJobCompletionProofOptions, IDockerGlobalJobTaskCompletionEvidence, IDockerServiceConfigReference, IDockerServiceSpecMode, IDockerServiceSecretReference, IDockerServiceStopProofOptions, IDockerServiceStoppedImagePinOptions, IDockerServiceTask, IServiceCreationDescriptor, } from './interfaces/service.js'; import { DockerResource } from './classes.base.js'; import { DockerConfig } from './classes.config.js'; import { DockerImage } from './classes.image.js'; import { DockerSecret } from './classes.secret.js'; import { assertDockerNodeId, assertDockerResponseStatus, assertDockerServiceId, assertDockerMutableImageReference, assertNonemptyDockerString, formatDockerResponseBody, } from './helpers.docker.js'; import { logger } from './logger.js'; const isTargetedSecretReference = ( secretArg: IServiceCreationDescriptor['secrets'][number], ): secretArg is IDockerServiceSecretReference => typeof secretArg === 'object' && 'secret' in secretArg; const isTargetedConfigReference = ( configArg: NonNullable[number], ): configArg is IDockerServiceConfigReference => typeof configArg === 'object' && 'config' in configArg; const isRecord = (valueArg: unknown): valueArg is Record => ( typeof valueArg === 'object' && valueArg !== null && !Array.isArray(valueArg) ); // SwarmKit agent end states prove that execution has ended. Orchestrator // remove/orphaned states can still refer to shutting-down or unreachable tasks. const observedTerminalTaskStates = new Set([ 'complete', 'shutdown', 'failed', 'rejected', ]); const desiredTerminalTaskStates = new Set([ ...observedTerminalTaskStates, 'remove', 'orphaned', ]); const canonicalJson = (valueArg: unknown): string => { const normalize = (entryArg: unknown): unknown => { if (Array.isArray(entryArg)) return entryArg.map(normalize); if (isRecord(entryArg)) { return Object.fromEntries( Object.entries(entryArg) .sort(([leftArg], [rightArg]) => leftArg.localeCompare(rightArg)) .map(([keyArg, valueArg]) => [keyArg, normalize(valueArg)]), ); } return entryArg; }; return JSON.stringify(normalize(valueArg)); }; const isImmutableRepositoryDigest = (referenceArg: string): boolean => { const match = /^(?[^\s@]+)@sha256:[a-f0-9]{64}$/u.exec(referenceArg); const repository = match?.groups?.repository; if (!repository) return false; const lastSlashIndex = repository.lastIndexOf('/'); return repository.lastIndexOf(':') <= lastSlashIndex; }; const hasVerifiedRepoDigest = DockerImage.prototype.hasVerifiedRepoDigest; const hasPulledMutableReference = DockerImage.prototype.hasPulledMutableReference; const createDockerServiceMode = ( modeArg: IServiceCreationDescriptor['mode'], ): IDockerServiceSpecMode | undefined => { if (modeArg === undefined) return undefined; if (!isRecord(modeArg)) throw new TypeError('Docker service mode must be an object'); if (modeArg.type === 'replicated') { if (canonicalJson(Object.keys(modeArg).sort()) !== canonicalJson(['replicas', 'type'])) { throw new TypeError('Docker replicated service mode must contain only type and replicas'); } if (!Number.isSafeInteger(modeArg.replicas) || modeArg.replicas < 0) { throw new TypeError('Docker replicated service replicas must be a non-negative safe integer'); } return { Replicated: { Replicas: modeArg.replicas } }; } if (modeArg.type === 'global-job') { if (canonicalJson(Object.keys(modeArg)) !== canonicalJson(['type'])) { throw new TypeError('Docker global-job service mode must contain only type'); } return { GlobalJob: {} }; } throw new TypeError('Docker service mode must be replicated or global-job'); }; const createDockerServicePlacement = ( placementArg: IServiceCreationDescriptor['placement'], ): { Constraints: string[] } | undefined => { if (placementArg === undefined) return undefined; if (!isRecord(placementArg) || canonicalJson(Object.keys(placementArg)) !== canonicalJson(['constraints']) || !Array.isArray(placementArg.constraints) || placementArg.constraints.length === 0) { throw new TypeError( 'Docker service placement must contain only a non-empty array of non-empty control-free constraints', ); } const constraints = Array.from(placementArg.constraints); if (constraints.some((constraintArg) => ( typeof constraintArg !== 'string' || constraintArg.length === 0 || /[\u0000-\u001f\u007f]/.test(constraintArg) ))) { throw new TypeError( 'Docker service placement must contain only a non-empty array of non-empty control-free constraints', ); } return { Constraints: constraints }; }; const createDockerServiceArgs = ( argsArg: IServiceCreationDescriptor['args'], ): string[] | undefined => { if (argsArg === undefined) return undefined; if (!Array.isArray(argsArg) || argsArg.length === 0) { throw new TypeError('Docker service args must be a non-empty array of strings'); } const args = Array.from(argsArg); if (args.some((argumentArg) => typeof argumentArg !== 'string')) { throw new TypeError('Docker service args must be a non-empty array of strings'); } return args; }; const createDockerServiceRestartPolicy = ( policyArg: IServiceCreationDescriptor['restartPolicy'], ): { Condition: 'none' } | undefined => { if (policyArg === undefined) return undefined; if (!isRecord(policyArg) || canonicalJson(Object.keys(policyArg)) !== canonicalJson(['condition']) || policyArg.condition !== 'none') { throw new TypeError( 'Docker service restart policy must contain only condition set to none', ); } return { Condition: 'none' }; }; const readDockerJobIterationIndex = (iterationArg: unknown): number | undefined => { if (!isRecord(iterationArg) || Object.keys(iterationArg).some((keyArg) => keyArg !== 'Index')) { return undefined; } if (iterationArg.Index === undefined) return 0; return Number.isSafeInteger(iterationArg.Index) && (iterationArg.Index as number) >= 0 ? iterationArg.Index as number : undefined; }; const hasExactDockerGlobalJobPlacement = ( placementArg: unknown, constraintsArg: string[], ): boolean => ( isRecord(placementArg) && Object.keys(placementArg).every((keyArg) => ( keyArg === 'Constraints' || keyArg === 'Platforms' )) && canonicalJson(placementArg.Constraints) === canonicalJson(constraintsArg) ); const hasNonretryingDockerRestartPolicy = (policyArg: unknown): boolean => ( isRecord(policyArg) && Object.keys(policyArg).every((keyArg) => [ 'Condition', 'Delay', 'MaxAttempts', 'Window', ].includes(keyArg)) && policyArg.Condition === 'none' ); const readExactWritableDockerBindMounts = ( mountsArg: unknown, ): Array<{ hostFsPath: string; containerFsPath: string }> | undefined => { if (!Array.isArray(mountsArg)) return undefined; const mounts: Array<{ hostFsPath: string; containerFsPath: string }> = []; const targets = new Set(); for (const mountArg of mountsArg) { if (!isRecord(mountArg) || Object.keys(mountArg).some((keyArg) => ![ 'Target', 'Source', 'Type', 'ReadOnly', 'Consistency', 'BindOptions', ].includes(keyArg)) || mountArg.Type !== 'bind' || typeof mountArg.Source !== 'string' || typeof mountArg.Target !== 'string' || (mountArg.ReadOnly !== undefined && mountArg.ReadOnly !== false) || (mountArg.Consistency !== undefined && mountArg.Consistency !== 'default') || (mountArg.BindOptions !== undefined && (!isRecord(mountArg.BindOptions) || Object.keys(mountArg.BindOptions).length !== 0))) { return undefined; } if (targets.has(mountArg.Target)) return undefined; targets.add(mountArg.Target); mounts.push({ hostFsPath: mountArg.Source, containerFsPath: mountArg.Target, }); } return mounts.sort((leftArg, rightArg) => ( leftArg.containerFsPath.localeCompare(rightArg.containerFsPath) || leftArg.hostFsPath.localeCompare(rightArg.hostFsPath) )); }; const resolveDockerServiceImageReference = ( descriptorArg: IServiceCreationDescriptor, ): string | undefined => { const immutableReference = descriptorArg.immutableImageReference; const mutableReference = descriptorArg.mutableImageReference; if (immutableReference !== undefined && mutableReference !== undefined) { throw new TypeError( 'Docker service image cannot use both immutable and mutable explicit references', ); } if (immutableReference === undefined && mutableReference === undefined) return undefined; if (immutableReference !== undefined && !isImmutableRepositoryDigest(immutableReference)) { throw new TypeError( 'Docker immutable service image reference must be a complete repository@sha256 reference', ); } if (mutableReference !== undefined) assertDockerMutableImageReference(mutableReference); if (!(descriptorArg.image instanceof DockerImage)) { throw new TypeError( `Docker ${immutableReference === undefined ? 'mutable' : 'immutable'} service image reference requires a verified DockerImage instance`, ); } if (immutableReference !== undefined) { if (!hasVerifiedRepoDigest.call(descriptorArg.image, immutableReference)) { throw new Error( 'Docker immutable service image reference does not match the DockerImage verification evidence', ); } return immutableReference; } if (mutableReference === undefined) { throw new Error('Docker mutable service image reference is unexpectedly absent'); } if (!hasPulledMutableReference.call(descriptorArg.image, mutableReference)) { throw new TypeError( 'Docker mutable service image reference does not match the DockerImage pull evidence', ); } return mutableReference; }; export class DockerService extends DockerResource { // STATIC (Internal - prefixed with _ to indicate internal use) /** * Internal: Get all services * Public API: Use dockerHost.listServices() instead */ public static async _list(dockerHost: DockerHost) { const services: DockerService[] = []; const response = await dockerHost.request('GET', '/services'); for (const serviceObject of response.body) { const dockerService = new DockerService(dockerHost); Object.assign(dockerService, serviceObject); services.push(dockerService); } return services; } /** * Internal: Get service by ID through direct inspection. */ public static async _fromId( dockerHostArg: DockerHost, idArg: string, ): Promise { assertNonemptyDockerString(idArg, 'Docker service ID'); const response = await dockerHostArg.request( 'GET', `/services/${encodeURIComponent(idArg)}`, ); if (response.statusCode === 404) { return undefined; } assertDockerResponseStatus(`Docker service inspect "${idArg}"`, 200, response); const service = new DockerService(dockerHostArg); Object.assign(service, response.body); return service; } /** * Internal: Get service by name * Public API: Use dockerHost.getServiceByName(name) instead */ public static async _fromName( dockerHost: DockerHost, networkName: string, ): Promise { const allServices = await DockerService._list(dockerHost); const wantedService = allServices.find((service) => { return service.Spec.Name === networkName; }); if (!wantedService) { throw new Error(`Service not found: ${networkName}`); } return wantedService; } /** * Internal: Create a service * Public API: Use dockerHost.createService(descriptor) instead */ public static async _create( dockerHost: DockerHost, serviceCreationDescriptor: IServiceCreationDescriptor, ): Promise { const serviceMode = createDockerServiceMode(serviceCreationDescriptor.mode); const servicePlacement = createDockerServicePlacement(serviceCreationDescriptor.placement); const serviceArgs = createDockerServiceArgs(serviceCreationDescriptor.args); const serviceRestartPolicy = createDockerServiceRestartPolicy( serviceCreationDescriptor.restartPolicy, ); const immutableImageReference = resolveDockerServiceImageReference(serviceCreationDescriptor); logger.log( 'info', `now creating service ${serviceCreationDescriptor.name}`, ); const requestedFileNames = [ ...serviceCreationDescriptor.secrets.map((secretArg) => isTargetedSecretReference(secretArg) ? secretArg.file.name : 'secret.json', ), ...(serviceCreationDescriptor.configs || []).map((configArg) => isTargetedConfigReference(configArg) ? configArg.file.name : 'config.json', ), ]; const seenFileNames = new Set(); for (const fileName of requestedFileNames) { if (seenFileNames.has(fileName)) { throw new Error(`Duplicate Docker resource file target: ${fileName}`); } seenFileNames.add(fileName); } // Resolve image (support both string and DockerImage instance) let imageInstance: DockerImage; if (typeof serviceCreationDescriptor.image === 'string') { const foundImage = await DockerImage._fromReference(dockerHost, serviceCreationDescriptor.image); if (!foundImage) { throw new Error(`Image not found: ${serviceCreationDescriptor.image}`); } imageInstance = foundImage; } else { imageInstance = serviceCreationDescriptor.image; } const serviceVersion = await imageInstance.getVersion(); const labels: TLabels = { ...serviceCreationDescriptor.labels, version: serviceVersion, }; const mounts: Array<{ /** * the target inside the container */ Target: string; /** * The Source from which to mount the data (Volume or host path) */ Source: string; Type: 'bind' | 'volume' | 'tmpfs' | 'npipe'; ReadOnly: boolean; Consistency: 'default' | 'consistent' | 'cached' | 'delegated'; }> = []; if (serviceCreationDescriptor.accessHostDockerSock) { mounts.push({ Target: '/var/run/docker.sock', Source: '/var/run/docker.sock', Consistency: 'default', ReadOnly: false, Type: 'bind', }); } if ( serviceCreationDescriptor.resources && serviceCreationDescriptor.resources.volumeMounts ) { for (const volumeMount of serviceCreationDescriptor.resources .volumeMounts) { mounts.push({ Target: volumeMount.containerFsPath, Source: volumeMount.hostFsPath, Consistency: 'default', ReadOnly: volumeMount.readOnly ?? false, Type: 'bind', }); } } // Resolve networks (support both string[] and DockerNetwork[]) const networkArray: Array<{ Target: string; Aliases: string[]; }> = []; for (const network of serviceCreationDescriptor.networks) { // Skip null networks (can happen if network creation fails) if (!network) { logger.log('warn', 'Skipping null network in service creation'); continue; } // Resolve network name const networkName = typeof network === 'string' ? network : network.Name; networkArray.push({ Target: networkName, Aliases: [serviceCreationDescriptor.networkAlias], }); } const ports: Array<{ Protocol: string; PublishedPort: number; TargetPort: number }> = []; for (const port of serviceCreationDescriptor.ports) { // "host:container" with an optional "/udp" or "/tcp" suffix (compose syntax) const [portPart, protocolPart] = port.split('/'); const portArray = portPart.split(':'); const hostPort = portArray[0]; const containerPort = portArray[1]; ports.push({ Protocol: protocolPart === 'udp' ? 'udp' : 'tcp', PublishedPort: parseInt(hostPort, 10), TargetPort: parseInt(containerPort, 10), }); } const secretArray: Array<{ File: { Name: string; UID: string; GID: string; Mode: number }; SecretID: string; SecretName: string; }> = []; for (const secretReference of serviceCreationDescriptor.secrets) { let secret: string | DockerSecret; let fileTarget: IDockerResourceFileTarget | undefined; if (isTargetedSecretReference(secretReference)) { secret = secretReference.secret; fileTarget = secretReference.file; } else { secret = secretReference; } let secretInstance: DockerSecret; if (typeof secret === 'string') { const foundSecret = await DockerSecret._fromName(dockerHost, secret); if (!foundSecret) { throw new Error(`Secret not found: ${secret}`); } secretInstance = foundSecret; } else { secretInstance = secret; } secretArray.push({ File: { Name: fileTarget?.name ?? 'secret.json', UID: fileTarget?.uid ?? '33', GID: fileTarget?.gid ?? '33', Mode: fileTarget?.mode ?? 0o600, }, SecretID: secretInstance.ID, SecretName: secretInstance.Spec.Name, }); } const configArray: Array<{ File: { Name: string; UID: string; GID: string; Mode: number }; ConfigID: string; ConfigName: string; }> = []; for (const configReference of serviceCreationDescriptor.configs || []) { let config: string | DockerConfig; let fileTarget: IDockerResourceFileTarget | undefined; if (isTargetedConfigReference(configReference)) { config = configReference.config; fileTarget = configReference.file; } else { config = configReference; } let configInstance: DockerConfig; if (typeof config === 'string') { const foundConfig = await DockerConfig._fromName(dockerHost, config); if (!foundConfig) { throw new Error(`Config not found: ${config}`); } configInstance = foundConfig; } else { configInstance = config; } configArray.push({ File: { Name: fileTarget?.name ?? 'config.json', UID: fileTarget?.uid ?? '0', GID: fileTarget?.gid ?? '0', Mode: fileTarget?.mode ?? 0o444, }, ConfigID: configInstance.ID, ConfigName: configInstance.Spec.Name, }); } // lets configure limits const memoryLimitMB = serviceCreationDescriptor.resources?.memorySizeMB ?? 1000; const limits = { MemoryBytes: memoryLimitMB * 1000000, }; // Swarm HealthConfig durations are nanoseconds const healthcheck = serviceCreationDescriptor.healthcheck ? { Test: serviceCreationDescriptor.healthcheck.test, ...(serviceCreationDescriptor.healthcheck.intervalMs !== undefined ? { Interval: serviceCreationDescriptor.healthcheck.intervalMs * 1_000_000 } : {}), ...(serviceCreationDescriptor.healthcheck.timeoutMs !== undefined ? { Timeout: serviceCreationDescriptor.healthcheck.timeoutMs * 1_000_000 } : {}), ...(serviceCreationDescriptor.healthcheck.startPeriodMs !== undefined ? { StartPeriod: serviceCreationDescriptor.healthcheck.startPeriodMs * 1_000_000 } : {}), ...(serviceCreationDescriptor.healthcheck.retries !== undefined ? { Retries: serviceCreationDescriptor.healthcheck.retries } : {}), } : undefined; const response = await dockerHost.request('POST', '/services/create', { Name: serviceCreationDescriptor.name, TaskTemplate: { ContainerSpec: { Image: immutableImageReference ?? imageInstance.RepoTags[0], Labels: labels, Secrets: secretArray, ...(configArray.length > 0 ? { Configs: configArray } : {}), Mounts: mounts, ...(serviceArgs ? { Args: serviceArgs } : {}), ...(healthcheck ? { Healthcheck: healthcheck } : {}), /* DNSConfig: { Nameservers: ['1.1.1.1'] } */ }, UpdateConfig: { Parallelism: 0, Delay: 0, FailureAction: 'pause', Monitor: 15000000000, MaxFailureRatio: 0.15, }, ForceUpdate: 1, Resources: { Limits: limits, }, ...(serviceRestartPolicy ? { RestartPolicy: serviceRestartPolicy } : {}), ...(servicePlacement ? { Placement: servicePlacement } : {}), Networks: networkArray, LogDriver: { Name: 'json-file', Options: { 'max-file': '3', 'max-size': '10M', }, }, }, ...(serviceMode ? { Mode: serviceMode } : {}), Labels: labels, EndpointSpec: { Ports: ports, }, }); assertDockerResponseStatus('Docker service create', 201, response); const createdId = response.body?.ID; if (typeof createdId !== 'string' || createdId.length === 0) { throw new Error( `Docker service create failed: expected response body.ID to be a nonempty string; response body: ${formatDockerResponseBody(response.body)}`, ); } const createdService = await DockerService._fromId(dockerHost, createdId); if (!createdService) { throw new Error( `Docker service create failed: direct inspect did not find created ID "${createdId}"`, ); } return createdService; } // INSTANCE PROPERTIES // Note: dockerHost (not dockerHostRef) for consistency with base class public ID!: string; public Version!: { Index: number }; public CreatedAt!: string; public UpdatedAt!: string; public Spec!: { Name: string; Labels: TLabels; TaskTemplate: { ContainerSpec: { Image: string; Isolation: string; Command?: string[]; Args?: string[]; Mounts?: Array<{ Target?: string; Source?: string; Type?: string; ReadOnly?: boolean; Consistency?: string; BindOptions?: Record; }>; Secrets: Array<{ File: { Name: string; UID: string; GID: string; Mode: number; }; SecretID: string; SecretName: string; }>; Configs?: Array<{ File: { Name: string; UID: string; GID: string; Mode: number; }; ConfigID: string; ConfigName: string; }>; }; ForceUpdate: 0; RestartPolicy?: { Condition?: string; Delay?: number; MaxAttempts?: number; Window?: number; }; Placement?: { Constraints?: string[]; Platforms?: Array<{ Architecture?: string; OS?: string; }>; }; Networks: Array<{ Target: string; Aliases: string[]; }>; }; Mode: IDockerServiceSpecMode; }; public JobStatus?: { JobIteration?: { Index?: number; }; LastExecution?: string; }; public Endpoint!: { Spec: {}; VirtualIPs: [any[]] }; constructor(dockerHostArg: DockerHost) { super(dockerHostArg); } // INSTANCE METHODS /** * Refreshes this service's state from the Docker daemon */ public async refresh(): Promise { const updated = await DockerService._fromName(this.dockerHost, this.Spec.Name); if (updated) { Object.assign(this, updated); } } /** * Removes this service from the Docker daemon */ public async remove() { await this.dockerHost.request('DELETE', `/services/${this.ID}`); } /** * Re-reads service data from Docker engine * @deprecated Use refresh() instead */ public async reReadFromDockerEngine() { const dockerData = await this.dockerHost.request( 'GET', `/services/${this.ID}`, ); if (dockerData.statusCode >= 300) { throw new Error(`Failed to refresh Docker service ${this.ID}: ${JSON.stringify(dockerData.body)}`); } if (dockerData.body) { Object.assign(this, dockerData.body); } } /** * Lists tasks for this swarm service. */ public async listTasks(optionsArg: { onlyRunning?: boolean } = {}): Promise { const serviceIdentifier = this.ID || this.Spec?.Name; assertNonemptyDockerString(serviceIdentifier, 'Docker service task filter ID'); const filters = encodeURIComponent(JSON.stringify({ service: [serviceIdentifier] })); const response = await this.dockerHost.request('GET', `/tasks?filters=${filters}`); assertDockerResponseStatus(`Docker service tasks "${serviceIdentifier}"`, 200, response); if (!Array.isArray(response.body)) { throw new Error( `Docker service tasks "${serviceIdentifier}" failed: expected an array; response body: ${formatDockerResponseBody(response.body)}`, ); } const tasks: IDockerServiceTask[] = response.body; if (!optionsArg.onlyRunning) { return tasks; } return tasks.filter((taskArg) => { return taskArg.DesiredState === 'running' && taskArg.Status?.State === 'running'; }); } /** * Proves one exact GlobalJob execution completed once on every authoritative target node. */ public async proveGlobalJobCompletion( optionsArg: IDockerGlobalJobCompletionProofOptions, ): Promise { assertDockerServiceId(this.ID); const options = DockerService.normalizeGlobalJobProofOptions(optionsArg); const first = await this.dockerHost.getServiceById(this.ID); if (!first) { throw new Error(`Docker global job ${this.ID} is unavailable for completion proof`); } DockerService.requireGlobalJobProofSpec(first, options, this.ID); const tasks = await first.listTasks(); const second = await this.dockerHost.getServiceById(this.ID); if (!second) { throw new Error(`Docker global job ${this.ID} disappeared during completion proof`); } DockerService.requireGlobalJobProofSpec(second, options, this.ID); if (first.Version.Index !== second.Version.Index || canonicalJson(first.Spec) !== canonicalJson(second.Spec) || readDockerJobIterationIndex(first.JobStatus?.JobIteration) !== readDockerJobIterationIndex(second.JobStatus?.JobIteration)) { throw new Error(`Docker global job ${this.ID} changed during completion proof`); } const expectedNodeIds = new Set(options.expectedTargetNodeIds); const evidence: IDockerGlobalJobTaskCompletionEvidence[] = []; const seenNodeIds = new Set(); for (const task of tasks) { assertDockerServiceId(task.ID, 'Docker global job task ID'); if (task.ServiceID !== this.ID) { throw new Error(`Docker global job ${this.ID} task ownership could not be proven`); } const taskIteration = readDockerJobIterationIndex(task.JobIteration); if (taskIteration === undefined) { throw new Error(`Docker global job ${this.ID} has a task without a valid job iteration`); } if (taskIteration > options.expectedJobIterationIndex) { throw new Error(`Docker global job ${this.ID} has a task from a future job iteration`); } if (taskIteration !== options.expectedJobIterationIndex) { if (typeof task.DesiredState !== 'string' || !desiredTerminalTaskStates.has(task.DesiredState) || typeof task.Status?.State !== 'string' || !observedTerminalTaskStates.has(task.Status.State)) { throw new Error(`Docker global job ${this.ID} has an unproven historical task`); } continue; } DockerService.requireGlobalJobTaskSpec(task, options, this.ID); assertDockerNodeId(task.NodeID, 'Docker global job task node ID'); if (!expectedNodeIds.has(task.NodeID) || seenNodeIds.has(task.NodeID)) { throw new Error( `Docker global job ${this.ID} current tasks do not exactly cover target nodes`, ); } if (task.DesiredState !== 'complete' || task.Status?.State !== 'complete' || task.Status.ContainerStatus?.ExitCode !== 0) { throw new Error(`Docker global job ${this.ID} has an unproven current task`); } seenNodeIds.add(task.NodeID); evidence.push({ taskId: task.ID, nodeId: task.NodeID, jobIterationIndex: options.expectedJobIterationIndex, desiredState: 'complete', state: 'complete', exitCode: 0, }); } if (seenNodeIds.size !== expectedNodeIds.size) { throw new Error( `Docker global job ${this.ID} current tasks do not exactly cover target nodes`, ); } evidence.sort((leftArg, rightArg) => leftArg.nodeId.localeCompare(rightArg.nodeId)); return { serviceId: this.ID, serviceVersionIndex: options.expectedServiceVersionIndex, jobIterationIndex: options.expectedJobIterationIndex, imageReference: options.expectedImageReference, args: options.expectedArgs, writableBindMounts: options.expectedWritableBindMounts, placementConstraints: options.expectedPlacementConstraints, targetNodeIds: [...options.expectedTargetNodeIds].sort(), tasks: evidence, }; } private static normalizeGlobalJobProofOptions( optionsArg: IDockerGlobalJobCompletionProofOptions, ): IDockerGlobalJobCompletionProofOptions { if (!isRecord(optionsArg) || canonicalJson(Object.keys(optionsArg).sort()) !== canonicalJson([ 'expectedArgs', 'expectedImageReference', 'expectedJobIterationIndex', 'expectedPlacementConstraints', 'expectedServiceVersionIndex', 'expectedTargetNodeIds', 'expectedWritableBindMounts', 'requiredServiceLabels', ])) { throw new TypeError('Docker global job proof options must use the exact schema'); } if (!Number.isSafeInteger(optionsArg.expectedServiceVersionIndex) || optionsArg.expectedServiceVersionIndex < 0) { throw new TypeError( 'Docker global job proof expectedServiceVersionIndex must be a non-negative safe integer', ); } if (!Number.isSafeInteger(optionsArg.expectedJobIterationIndex) || optionsArg.expectedJobIterationIndex < 0) { throw new TypeError( 'Docker global job proof expectedJobIterationIndex must be a non-negative safe integer', ); } if (typeof optionsArg.expectedImageReference !== 'string' || !isImmutableRepositoryDigest(optionsArg.expectedImageReference)) { throw new TypeError( 'Docker global job proof image must be a complete immutable repository digest', ); } const args = createDockerServiceArgs(optionsArg.expectedArgs); if (!args) throw new TypeError('Docker global job proof args must be present'); if (!Array.isArray(optionsArg.expectedWritableBindMounts) || optionsArg.expectedWritableBindMounts.length === 0) { throw new TypeError( 'Docker global job proof writable bind mounts must be a non-empty array', ); } const mountTargets = new Set(); const writableBindMounts = optionsArg.expectedWritableBindMounts.map((mountArg) => { if (!isRecord(mountArg) || canonicalJson(Object.keys(mountArg).sort()) !== canonicalJson(['containerFsPath', 'hostFsPath']) || typeof mountArg.hostFsPath !== 'string' || typeof mountArg.containerFsPath !== 'string' || !mountArg.hostFsPath.startsWith('/') || !mountArg.containerFsPath.startsWith('/') || /[\u0000-\u001f\u007f]/.test(mountArg.hostFsPath) || /[\u0000-\u001f\u007f]/.test(mountArg.containerFsPath)) { throw new TypeError( 'Docker global job proof writable bind mounts must contain exact absolute control-free paths', ); } if (mountTargets.has(mountArg.containerFsPath)) { throw new TypeError('Docker global job proof writable bind mount targets must be unique'); } mountTargets.add(mountArg.containerFsPath); return { hostFsPath: mountArg.hostFsPath, containerFsPath: mountArg.containerFsPath, }; }).sort((leftArg, rightArg) => ( leftArg.containerFsPath.localeCompare(rightArg.containerFsPath) || leftArg.hostFsPath.localeCompare(rightArg.hostFsPath) )); const placement = createDockerServicePlacement({ constraints: optionsArg.expectedPlacementConstraints, }); if (!placement) { throw new TypeError('Docker global job proof placement constraints must be present'); } if (new Set(placement.Constraints).size !== placement.Constraints.length) { throw new TypeError('Docker global job proof placement constraints must be unique'); } if (!isRecord(optionsArg.requiredServiceLabels) || Object.keys(optionsArg.requiredServiceLabels).length === 0 || Object.values(optionsArg.requiredServiceLabels).some( (valueArg) => typeof valueArg !== 'string', )) { throw new TypeError( 'Docker global job proof requiredServiceLabels must contain string values', ); } if (!Array.isArray(optionsArg.expectedTargetNodeIds) || optionsArg.expectedTargetNodeIds.length === 0) { throw new TypeError('Docker global job proof target node IDs must be a non-empty array'); } const targetNodeIds = optionsArg.expectedTargetNodeIds.map((nodeIdArg) => { assertDockerNodeId(nodeIdArg, 'Docker global job proof target node ID'); return nodeIdArg; }); if (new Set(targetNodeIds).size !== targetNodeIds.length) { throw new TypeError('Docker global job proof target node IDs must be unique'); } return { expectedServiceVersionIndex: optionsArg.expectedServiceVersionIndex, expectedJobIterationIndex: optionsArg.expectedJobIterationIndex, expectedImageReference: optionsArg.expectedImageReference, expectedArgs: args as [string, ...string[]], expectedWritableBindMounts: writableBindMounts, expectedPlacementConstraints: placement.Constraints, requiredServiceLabels: { ...optionsArg.requiredServiceLabels }, expectedTargetNodeIds: targetNodeIds, }; } private static requireGlobalJobProofSpec( serviceArg: DockerService, optionsArg: IDockerGlobalJobCompletionProofOptions, serviceIdArg: string, ): void { if (serviceArg.ID !== serviceIdArg) { throw new Error('Docker global job completion proof requires an exact service ID'); } const spec = serviceArg.Spec; const taskTemplate = isRecord(spec?.TaskTemplate) ? spec.TaskTemplate : undefined; const containerSpec = taskTemplate && isRecord(taskTemplate.ContainerSpec) ? taskTemplate.ContainerSpec : undefined; const actualMounts = readExactWritableDockerBindMounts(containerSpec?.Mounts); const actualLabels = isRecord(spec?.Labels) ? spec.Labels : {}; if (!Number.isSafeInteger(serviceArg.Version?.Index) || serviceArg.Version.Index !== optionsArg.expectedServiceVersionIndex || canonicalJson(spec?.Mode) !== canonicalJson({ GlobalJob: {} }) || !taskTemplate || !containerSpec || containerSpec.Image !== optionsArg.expectedImageReference || containerSpec.Command !== undefined || canonicalJson(containerSpec.Args) !== canonicalJson(optionsArg.expectedArgs) || canonicalJson(actualMounts) !== canonicalJson(optionsArg.expectedWritableBindMounts) || !hasNonretryingDockerRestartPolicy(taskTemplate.RestartPolicy) || !hasExactDockerGlobalJobPlacement( taskTemplate.Placement, optionsArg.expectedPlacementConstraints, ) || readDockerJobIterationIndex(serviceArg.JobStatus?.JobIteration) !== optionsArg.expectedJobIterationIndex) { throw new Error(`Docker global job ${serviceArg.ID} specification could not be proven`); } for (const [key, value] of Object.entries(optionsArg.requiredServiceLabels)) { if (actualLabels[key] !== value) { throw new Error(`Docker global job ${serviceArg.ID} label proof failed for "${key}"`); } } } private static requireGlobalJobTaskSpec( taskArg: IDockerServiceTask, optionsArg: IDockerGlobalJobCompletionProofOptions, serviceIdArg: string, ): void { const taskSpec = isRecord(taskArg.Spec) ? taskArg.Spec : undefined; const containerSpec = taskSpec && isRecord(taskSpec.ContainerSpec) ? taskSpec.ContainerSpec : undefined; const actualMounts = readExactWritableDockerBindMounts(containerSpec?.Mounts); if (!taskSpec || !containerSpec || containerSpec.Image !== optionsArg.expectedImageReference || containerSpec.Command !== undefined || canonicalJson(containerSpec.Args) !== canonicalJson(optionsArg.expectedArgs) || canonicalJson(actualMounts) !== canonicalJson(optionsArg.expectedWritableBindMounts) || !hasNonretryingDockerRestartPolicy(taskSpec.RestartPolicy) || !hasExactDockerGlobalJobPlacement( taskSpec.Placement, optionsArg.expectedPlacementConstraints, )) { throw new Error(`Docker global job ${serviceIdArg} current task specification is unproven`); } } /** * Scales an exact replicated service to zero and proves every exact-ID task is terminal. * A vanished service is accepted only after its exact-ID task inventory is also terminal. * Optional version and service-label preconditions are checked on the fresh exact-ID * inspection before either an update or already-stopped proof is accepted. */ public async stopAndProveStopped( optionsArg: IDockerServiceStopProofOptions = {}, ): Promise { assertDockerServiceId(this.ID); const timeoutMs = optionsArg.timeoutMs ?? 30_000; const pollIntervalMs = optionsArg.pollIntervalMs ?? 250; if (!Number.isSafeInteger(timeoutMs) || timeoutMs < 1) { throw new TypeError('Docker service stop timeoutMs must be a positive safe integer'); } if (!Number.isSafeInteger(pollIntervalMs) || pollIntervalMs < 1) { throw new TypeError('Docker service stop pollIntervalMs must be a positive safe integer'); } if ( optionsArg.expectedVersionIndex !== undefined && (!Number.isSafeInteger(optionsArg.expectedVersionIndex) || optionsArg.expectedVersionIndex < 0) ) { throw new TypeError( 'Docker service stop expectedVersionIndex must be a non-negative safe integer', ); } if (optionsArg.requiredLabels !== undefined) { if ( !isRecord(optionsArg.requiredLabels) || Object.values(optionsArg.requiredLabels).some((valueArg) => typeof valueArg !== 'string') ) { throw new TypeError('Docker service stop requiredLabels must contain only string values'); } } const current = await this.dockerHost.getServiceById(this.ID); let updateError: unknown; if (current) { const currentSpec = DockerService.requireReplicatedStopSpec(current); if ( optionsArg.expectedVersionIndex !== undefined && currentSpec.version !== optionsArg.expectedVersionIndex ) { throw new Error( `Docker service stop "${this.ID}" version precondition failed: expected ${optionsArg.expectedVersionIndex}, actual ${currentSpec.version}`, ); } const currentLabels = isRecord(currentSpec.spec.Labels) ? currentSpec.spec.Labels : {}; for (const [key, value] of Object.entries(optionsArg.requiredLabels ?? {})) { if (currentLabels[key] !== value) { throw new Error( `Docker service stop "${this.ID}" label precondition failed for "${key}"`, ); } } if (currentSpec.replicas !== 0) { try { const response = await this.dockerHost.request( 'POST', `/services/${encodeURIComponent(this.ID)}/update?version=${currentSpec.version}`, { ...currentSpec.spec, Mode: { ...currentSpec.mode, Replicated: { ...currentSpec.replicated, Replicas: 0, }, }, }, ); assertDockerResponseStatus(`Docker service stop "${this.ID}"`, 200, response); } catch (error) { updateError = error; } } } const deadline = Date.now() + timeoutMs; let proofError: unknown; while (true) { try { if (await this.hasExactStoppedProof()) { return; } proofError = undefined; } catch (error) { proofError = error; } const remainingMs = deadline - Date.now(); if (remainingMs <= 0) { break; } await new Promise((resolveArg) => { setTimeout(resolveArg, Math.min(pollIntervalMs, remainingMs)); }); } const errors = [updateError, proofError].filter( (errorArg) => errorArg !== undefined, ); if (errors.length > 0) { throw new AggregateError( errors, `Docker service ${this.ID} could not be proven stopped`, ); } throw new Error(`Docker service ${this.ID} remains active after stop`); } /** * Pins a stopped replicated service to an immutable repository digest. */ public async pinStoppedImage( optionsArg: IDockerServiceStoppedImagePinOptions, ): Promise { assertDockerServiceId(this.ID); assertNonemptyDockerString( optionsArg.expectedImageReference, 'Expected Docker service image reference', ); if (!isImmutableRepositoryDigest(optionsArg.imageReference)) { throw new TypeError( 'Docker service image pin requires a complete immutable repository digest', ); } const current = await this.dockerHost.getServiceById(this.ID); if (!current) { throw new Error(`Docker service ${this.ID} is unavailable for image pinning`); } const currentSpec = DockerService.requireStoppedImagePinSpec( current, optionsArg.expectedImageReference, ); DockerService.assertExactTerminalTasks( await current.listTasks(), this.ID, 'before image pinning', ); const expectedSpec = { ...currentSpec.spec, TaskTemplate: { ...currentSpec.taskTemplate, ContainerSpec: { ...currentSpec.containerSpec, Image: optionsArg.imageReference, }, }, }; let updateError: unknown; try { const response = await this.dockerHost.request( 'POST', `/services/${encodeURIComponent(this.ID)}/update?version=${currentSpec.version}`, expectedSpec, ); assertDockerResponseStatus(`Docker service image pin "${this.ID}"`, 200, response); } catch (error) { updateError = error; } let proofError: unknown; try { const updated = await this.dockerHost.getServiceById(this.ID); if (!updated) { throw new Error(`Docker service ${this.ID} disappeared after image pinning`); } DockerService.requireStoppedImagePinSpec(updated, optionsArg.imageReference); if ( !Number.isSafeInteger(updated.Version?.Index) || updated.Version.Index <= currentSpec.version || canonicalJson(updated.Spec) !== canonicalJson(expectedSpec) ) { throw new Error( `Docker service ${this.ID} specification changed during image pinning`, ); } DockerService.assertExactTerminalTasks( await updated.listTasks(), this.ID, 'after image pinning', ); return updated; } catch (error) { proofError = error; } const errors = [updateError, proofError].filter( (errorArg) => errorArg !== undefined, ); throw new AggregateError( errors, `Docker service ${this.ID} image pin could not be proven`, ); } private static requireReplicatedStopSpec(serviceArg: DockerService): { version: number; spec: Record; mode: Record; replicated: Record; replicas: number; } { const version = serviceArg.Version?.Index; const spec = serviceArg.Spec; const mode = isRecord(spec?.Mode) ? spec.Mode : undefined; const replicated = isRecord(mode?.Replicated) ? mode.Replicated : undefined; const replicas = replicated?.Replicas; if ( !Number.isSafeInteger(version) || (version as number) < 0 || !isRecord(spec) || !mode || !replicated || !Number.isSafeInteger(replicas) || (replicas as number) < 0 ) { throw new Error('Docker service stop requires an exact replicated service specification'); } return { version: version as number, spec, mode, replicated, replicas: replicas as number, }; } private static requireStoppedImagePinSpec( serviceArg: DockerService, expectedImageReferenceArg: string, ): { version: number; spec: Record; taskTemplate: Record; containerSpec: Record; } { const stoppedSpec = DockerService.requireReplicatedStopSpec(serviceArg); const taskTemplate = isRecord(stoppedSpec.spec.TaskTemplate) ? stoppedSpec.spec.TaskTemplate : undefined; const containerSpec = taskTemplate && isRecord(taskTemplate.ContainerSpec) ? taskTemplate.ContainerSpec : undefined; if ( stoppedSpec.replicas !== 0 || !taskTemplate || !containerSpec || containerSpec.Image !== expectedImageReferenceArg ) { throw new Error( 'Docker service image pin requires the exact stopped service specification', ); } return { version: stoppedSpec.version, spec: stoppedSpec.spec, taskTemplate, containerSpec, }; } private static assertExactTerminalTasks( tasksArg: IDockerServiceTask[], serviceIdArg: string, phaseArg: string, ): void { for (const task of tasksArg) { if ( task.ServiceID !== serviceIdArg || typeof task.Status?.State !== 'string' || !observedTerminalTaskStates.has(task.Status.State) ) { throw new Error( `Docker service ${serviceIdArg} has an unproven task ${phaseArg}`, ); } } } private static async listTasksByExactServiceId( dockerHostArg: DockerHost, serviceIdArg: string, ): Promise { const filters = encodeURIComponent(JSON.stringify({ service: [serviceIdArg] })); const response = await dockerHostArg.request('GET', `/tasks?filters=${filters}`); assertDockerResponseStatus(`Docker service tasks "${serviceIdArg}"`, 200, response); if (!Array.isArray(response.body)) { throw new Error( `Docker service tasks "${serviceIdArg}" failed: expected an array; response body: ${formatDockerResponseBody(response.body)}`, ); } return response.body as IDockerServiceTask[]; } private async hasExactStoppedProof(): Promise { const current = await this.dockerHost.getServiceById(this.ID); if (current && DockerService.requireReplicatedStopSpec(current).replicas !== 0) { return false; } const tasks = await DockerService.listTasksByExactServiceId(this.dockerHost, this.ID); return tasks.every((taskArg) => { const state = taskArg.Status?.State; return typeof state === 'string' && observedTerminalTaskStates.has(state); }); } /** * Gets image IDs for currently running tasks of this service. */ public async getRunningTaskImageIds(): Promise { const runningTasks = await this.listTasks({ onlyRunning: true }); const imageIds: string[] = []; for (const task of runningTasks) { const imageId = await this.getTaskContainerImageId(task); if (imageId) { imageIds.push(imageId); } } return imageIds; } /** * Checks if this service needs an image update. */ public async needsUpdate(desiredImageArg: DockerImage): Promise { await this.reReadFromDockerEngine(); const desiredImage = desiredImageArg; const desiredVersion = desiredImage.Labels?.version; const currentVersion = this.Spec?.Labels?.version; if (desiredVersion && currentVersion && desiredVersion !== currentVersion) { return true; } const desiredImageId = DockerService.normalizeImageId( desiredImage.Id || (desiredImage as DockerImage & { ID?: string }).ID, ); if (!desiredImageId) { return false; } const runningTasks = await this.listTasks({ onlyRunning: true }); if (runningTasks.length === 0) { return true; } for (const runningTask of runningTasks) { const runningImageId = DockerService.normalizeImageId( await this.getTaskContainerImageId(runningTask), ); if (!runningImageId || runningImageId !== desiredImageId) { return true; } } return false; } private async getTaskContainerImageId(taskArg: IDockerServiceTask): Promise { const directImageId = taskArg.Status?.ContainerStatus?.ImageID; if (directImageId) { return directImageId; } const containerId = taskArg.Status?.ContainerStatus?.ContainerID; if (!containerId) { return undefined; } const containerResponse = await this.dockerHost.request( 'GET', `/containers/${encodeURIComponent(containerId)}/json`, ); if (containerResponse.statusCode >= 300) { return undefined; } const containerBody = containerResponse.body || {}; return typeof containerBody.Image === 'string' ? containerBody.Image : typeof containerBody.ImageID === 'string' ? containerBody.ImageID : undefined; } private static normalizeImageId(imageIdArg?: string): string | undefined { if (!imageIdArg) { return undefined; } return imageIdArg.replace(/^sha256:/, ''); } }