diff --git a/.azure-pipelines/templates/IsolationSession.TestBundle.Build.Job.yml b/.azure-pipelines/templates/IsolationSession.TestBundle.Build.Job.yml index c1360029a..86ec655cc 100644 --- a/.azure-pipelines/templates/IsolationSession.TestBundle.Build.Job.yml +++ b/.azure-pipelines/templates/IsolationSession.TestBundle.Build.Job.yml @@ -72,9 +72,9 @@ jobs: path: s - task: NodeTool@1 - displayName: Use Node.js 22 + displayName: Use Node.js 24.21.0 inputs: - version: '22.x' + version: '24.21.0' - script: node "$(Pipeline.Workspace)/lock-check/scripts/versioning/check-package-lock-integrity.js" displayName: 'Check package locks use npmjs and SHA-512 integrity' diff --git a/.azure-pipelines/templates/Package.Lock.Check.Job.yml b/.azure-pipelines/templates/Package.Lock.Check.Job.yml index ae31001fc..eb3a48a8d 100644 --- a/.azure-pipelines/templates/Package.Lock.Check.Job.yml +++ b/.azure-pipelines/templates/Package.Lock.Check.Job.yml @@ -12,9 +12,9 @@ jobs: - checkout: self - task: NodeTool@1 - displayName: Use Node.js 20 + displayName: Use Node.js 24.21.0 inputs: - version: '20.x' + version: '24.21.0' - script: node scripts/versioning/check-package-lock-integrity.js displayName: Check package locks use npmjs and SHA-512 integrity diff --git a/.azure-pipelines/templates/Package.NpmSdk.Job.yml b/.azure-pipelines/templates/Package.NpmSdk.Job.yml index 480a42f97..a9ed86c1a 100644 --- a/.azure-pipelines/templates/Package.NpmSdk.Job.yml +++ b/.azure-pipelines/templates/Package.NpmSdk.Job.yml @@ -46,9 +46,9 @@ jobs: - template: Cargo.Setup.Public.yml@self - task: UseNode@1 - displayName: Use Node.js 20 + displayName: Use Node.js 24.21.0 inputs: - version: '20.x' + version: '24.21.0' - task: Cache@2 displayName: Cache npm diff --git a/.azure-pipelines/templates/SDK.Integration.Test.Job.yml b/.azure-pipelines/templates/SDK.Integration.Test.Job.yml index b2b32dbc6..9ceddc063 100644 --- a/.azure-pipelines/templates/SDK.Integration.Test.Job.yml +++ b/.azure-pipelines/templates/SDK.Integration.Test.Job.yml @@ -51,9 +51,9 @@ jobs: - checkout: self - task: UseNode@1 - displayName: Use Node.js 20 + displayName: Use Node.js 24.21.0 inputs: - version: '20.x' + version: '24.21.0' - task: Cache@2 displayName: Cache npm diff --git a/.azure-pipelines/templates/SDK.Unit.Test.Job.yml b/.azure-pipelines/templates/SDK.Unit.Test.Job.yml index bc1c8904f..3b1a7e614 100644 --- a/.azure-pipelines/templates/SDK.Unit.Test.Job.yml +++ b/.azure-pipelines/templates/SDK.Unit.Test.Job.yml @@ -36,9 +36,9 @@ jobs: - checkout: self - task: UseNode@1 - displayName: Use Node.js 20 + displayName: Use Node.js 24.21.0 inputs: - version: '20.x' + version: '24.21.0' - task: Cache@2 displayName: Cache npm diff --git a/.github/workflows/Package.Lock.Check.Job.yml b/.github/workflows/Package.Lock.Check.Job.yml index a1f03aa01..15e738d78 100644 --- a/.github/workflows/Package.Lock.Check.Job.yml +++ b/.github/workflows/Package.Lock.Check.Job.yml @@ -15,7 +15,7 @@ jobs: - uses: actions/setup-node@820762786026740c76f36085b0efc47a31fe5020 # v7.0.0 with: - node-version: 20 + node-version: 24.21.0 - name: Check package locks use npmjs and SHA-512 integrity run: node scripts/versioning/check-package-lock-integrity.js diff --git a/.github/workflows/Package.NpmSdk.Job.yml b/.github/workflows/Package.NpmSdk.Job.yml index 6c8730f5e..76b464b1a 100644 --- a/.github/workflows/Package.NpmSdk.Job.yml +++ b/.github/workflows/Package.NpmSdk.Job.yml @@ -15,7 +15,7 @@ jobs: - uses: actions/setup-node@820762786026740c76f36085b0efc47a31fe5020 # v7.0.0 with: - node-version: 20 + node-version: 24.21.0 cache: npm cache-dependency-path: sdk/node/package-lock.json diff --git a/.github/workflows/SDK.Integration.Test.Job.yml b/.github/workflows/SDK.Integration.Test.Job.yml index fc8495e37..5bdc1adc3 100644 --- a/.github/workflows/SDK.Integration.Test.Job.yml +++ b/.github/workflows/SDK.Integration.Test.Job.yml @@ -37,7 +37,7 @@ jobs: - uses: actions/setup-node@820762786026740c76f36085b0efc47a31fe5020 # v7.0.0 with: - node-version: 20 + node-version: 24.21.0 cache: npm cache-dependency-path: sdk/node/tests/integration/package-lock.json diff --git a/.github/workflows/SDK.Unit.Test.Job.yml b/.github/workflows/SDK.Unit.Test.Job.yml index fec809016..79e53af92 100644 --- a/.github/workflows/SDK.Unit.Test.Job.yml +++ b/.github/workflows/SDK.Unit.Test.Job.yml @@ -23,7 +23,7 @@ jobs: - uses: actions/setup-node@820762786026740c76f36085b0efc47a31fe5020 # v7.0.0 with: - node-version: 20 + node-version: 24.21.0 cache: npm cache-dependency-path: sdk/node/package-lock.json diff --git a/.github/workflows/Versioning.Checks.Job.yml b/.github/workflows/Versioning.Checks.Job.yml index b6a31ee73..9611b97b4 100644 --- a/.github/workflows/Versioning.Checks.Job.yml +++ b/.github/workflows/Versioning.Checks.Job.yml @@ -19,7 +19,7 @@ jobs: - uses: actions/setup-node@820762786026740c76f36085b0efc47a31fe5020 # v7.0.0 with: - node-version: 20 + node-version: 24.21.0 - name: Setup Rust toolchain run: rustup update stable diff --git a/docs/architecture.md b/docs/architecture.md index e580d996b..e735044a2 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -102,22 +102,25 @@ flowchart LR typescript["TypeScript SDK"] cli["CLI"] csharp["C# SDK"] - rust["Rust SDK"] + rust["Rust SDK (mxc-sdk crate)"] executor["Platform executor"] ffi["mxc_ffi"] - sdk["mxc-sdk"] engine["mxc_engine"] typescript --> executor + typescript -. "internal native streaming" .-> ffi cli --> executor - csharp --> ffi --> sdk - rust --> sdk + csharp --> ffi + ffi -->|"existing SDK APIs"| rust + ffi -->|"callback streaming"| engine executor --> engine - sdk --> engine + rust --> engine ``` -`mxc_ffi` is the C ABI used by the C# SDK. Its generated C# P/Invoke file is -created during the C# build. Generated TypeScript wire types come from the +`mxc_ffi` is the C ABI used by the C# SDK and by internal event-loop binding +paths. The callback streaming entry points call `mxc_engine` directly and are +not included in the generated C# P/Invoke surface. The generated C# P/Invoke +file is created during the C# build. Generated TypeScript wire types come from the schema tooling. ## Tests diff --git a/docs/isolation-session/state-aware-rust.md b/docs/isolation-session/state-aware-rust.md index cee34727c..237223a16 100644 --- a/docs/isolation-session/state-aware-rust.md +++ b/docs/isolation-session/state-aware-rust.md @@ -30,7 +30,7 @@ The Rust SDK (`mxc-sdk`) and the C ABI over it (`mxc_ffi`), each with an |---|---|---| | provision / start / stop / deprovision | `wxc-exec --config …` | `mxc_sdk::run_state_aware_json`, `mxc_state_aware` | | exec, attached to the caller's stdio | `wxc-exec --config …` | `mxc_sdk::exec_attached`, `mxc_state_aware_exec_attached` | -| exec, caller drives the pipes | *(no CLI equivalent)* | `mxc_sdk::exec_sandbox`, `mxc_state_aware_exec` | +| exec, caller drives the pipes | *(no CLI equivalent)* | `mxc_sdk::exec_sandbox`, `mxc_state_aware_exec`, `mxc_io_state_aware_exec_callback` | Requirements on an in-process caller: diff --git a/sdk/node/README.md b/sdk/node/README.md index f4ac77785..5c8be80aa 100644 --- a/sdk/node/README.md +++ b/sdk/node/README.md @@ -51,6 +51,10 @@ child.on('close', (code) => console.log('exit:', code)); When a new schema graduates or a new backend ships, update only this block. --> +**Node.js:** `24.21.0` or later. On Windows, native pipe streaming uses the +`windowsHandle` option added to `fs.ReadStream` and `fs.WriteStream` in +Node.js 24.21.0. + **Policy / config schema versions:** | Version | Status | Schema file | diff --git a/sdk/node/package-lock.json b/sdk/node/package-lock.json index 092a6fe31..89b84c0ef 100644 --- a/sdk/node/package-lock.json +++ b/sdk/node/package-lock.json @@ -14,14 +14,14 @@ "semver": "^7.7.4" }, "devDependencies": { - "@types/node": "^20.10.0", + "@types/node": "^24.13.4", "@types/semver": "^7.7.1", "node-gyp": "^12.2.0", "rimraf": "^6.1.3", "typescript": "^5.3.3" }, "engines": { - "node": ">=18.0.0" + "node": ">=24.21.0" } }, "node_modules/@gar/promise-retry": { @@ -376,13 +376,13 @@ } }, "node_modules/@types/node": { - "version": "20.19.31", - "resolved": "https://registry.npmjs.org/@types/node/-/node-20.19.31.tgz", - "integrity": "sha512-5jsi0wpncvTD33Sh1UCgacK37FFwDn+EG7wCmEvs62fCvBL+n8/76cAYDok21NF6+jaVWIqKwCZyX7Vbu8eB3A==", + "version": "24.13.4", + "resolved": "https://ms-feed-25.pkgs.visualstudio.com/1es-public/_packaging/npm-public/npm/registry/@types/node/-/node-24.13.4.tgz", + "integrity": "sha1-5U90zq5oOL/1eoFLfTH1BCoPkHQ=", "dev": true, "license": "MIT", "dependencies": { - "undici-types": "~6.21.0" + "undici-types": "~7.18.0" } }, "node_modules/@types/semver": { @@ -1121,9 +1121,9 @@ } }, "node_modules/undici-types": { - "version": "6.21.0", - "resolved": "https://registry.npmjs.org/undici-types/-/undici-types-6.21.0.tgz", - "integrity": "sha512-iwDZqg0QAGrg9Rav5H4n0M64c3mkR59cJ6wQp+7C4nI0gsmExaedaYLNO44eT4AtBBwjbTiGPMlt2Md0T9H9JQ==", + "version": "7.18.2", + "resolved": "https://ms-feed-25.pkgs.visualstudio.com/1es-public/_packaging/npm-public/npm/registry/undici-types/-/undici-types-7.18.2.tgz", + "integrity": "sha1-KTV6iee3ykrvO/D9P9DNc4hCKek=", "dev": true, "license": "MIT" }, diff --git a/sdk/node/package.json b/sdk/node/package.json index 9bf807ac0..04772a2eb 100644 --- a/sdk/node/package.json +++ b/sdk/node/package.json @@ -23,7 +23,7 @@ "watch": "tsc --watch", "clean": "rimraf dist", "test": "npm run test:unit", - "test:unit": "npm run build:test-unit && node --test dist-tests/tests/unit/sandbox.test.js dist-tests/tests/unit/policy.test.js dist-tests/tests/unit/logger.test.js dist-tests/tests/unit/errors.test.js dist-tests/tests/unit/state-aware-types.test.js dist-tests/tests/unit/state-aware.test.js dist-tests/tests/unit/platform.test.js dist-tests/tests/unit/native-library.test.js dist-tests/tests/unit/binding-request.test.js dist-tests/tests/unit/binding-run.test.js dist-tests/tests/unit/inprocess-run.test.js dist-tests/tests/unit/telemetry.test.js dist-tests/tests/unit/default-consent-protocol-runner.test.js dist-tests/tests/unit/wire-conformance.test.js dist-tests/tests/unit/wire-conformance-state-aware.test.js", + "test:unit": "npm run build:test-unit && node --test dist-tests/tests/unit/sandbox.test.js dist-tests/tests/unit/policy.test.js dist-tests/tests/unit/logger.test.js dist-tests/tests/unit/errors.test.js dist-tests/tests/unit/state-aware-types.test.js dist-tests/tests/unit/state-aware.test.js dist-tests/tests/unit/platform.test.js dist-tests/tests/unit/native-library.test.js dist-tests/tests/unit/binding-request.test.js dist-tests/tests/unit/binding-run.test.js dist-tests/tests/unit/streaming-binding.test.js dist-tests/tests/unit/inprocess-run.test.js dist-tests/tests/unit/sandbox-process.test.js dist-tests/tests/unit/telemetry.test.js dist-tests/tests/unit/default-consent-protocol-runner.test.js dist-tests/tests/unit/wire-conformance.test.js dist-tests/tests/unit/wire-conformance-state-aware.test.js", "test:integration": "cd tests/integration && npm install && npm run build && npm test", "prepublishOnly": "npm run build", "typecheck:integration": "cd tests/integration && npx tsc --noEmit -p tsconfig.json", @@ -40,7 +40,7 @@ "author": "Microsoft Corporation", "license": "MIT", "devDependencies": { - "@types/node": "^20.10.0", + "@types/node": "^24.13.4", "@types/semver": "^7.7.1", "node-gyp": "^12.2.0", "rimraf": "^6.1.3", @@ -52,6 +52,6 @@ "semver": "^7.7.4" }, "engines": { - "node": ">=18.0.0" + "node": ">=24.21.0" } } \ No newline at end of file diff --git a/sdk/node/src/bindings/streaming.ts b/sdk/node/src/bindings/streaming.ts new file mode 100644 index 000000000..678ca2607 --- /dev/null +++ b/sdk/node/src/bindings/streaming.ts @@ -0,0 +1,477 @@ +// Copyright (c) Microsoft Corporation. +// Licensed under the MIT License. + +// Native stdio ownership is transferred once to Node. Rust retains only +// process lifecycle control, so bytes never cross Koffi callbacks. + +import * as fs from 'node:fs'; +import * as net from 'node:net'; +import * as os from 'node:os'; +import type { Readable, Writable } from 'node:stream'; +import koffi, { type KoffiFunc } from 'koffi'; +import { + _createMxcSandboxProcess, + type MxcSandboxProcess, + type NativeLifecycleDriver, + type NativeLifecycleStatus, +} from '../sandbox-process.js'; +import { loadMxcFfi, type MxcNativeLibrary } from '../native-library.js'; +import type { RequestSpec } from './request.js'; +import { bindNativeFunction } from './native-function.js'; +import { + AbiErrorDetailType, + decodeString, + nativeStatusError, + parseStringArray, + type AbiErrorDetail, +} from './native-error.js'; + +type Pointer = unknown; +type NativeHandle = number | bigint; +type NativeLibraryHandle = MxcNativeLibrary['handle']; +type NativeFreeCompletion = (error: Error | null) => void; + +interface AbiNativeStdio { + stdin_handle: NativeHandle; + stdout_handle: NativeHandle; + stderr_handle: NativeHandle; +} + +const AbiCoordinator = koffi.opaque('MxcNodeLifecycleCoordinator'); +const AbiNativeStdioType = koffi.struct('MxcNodeNativeStdio', { + stdin_handle: 'intptr_t', + stdout_handle: 'intptr_t', + stderr_handle: 'intptr_t', +}); + +export interface _StreamingNativeFacade { + spawn( + request: string, + outHandle: Pointer[], + error: AbiErrorDetail, + ): number; + id(handle: Pointer): number; + takeNativeStdio(handle: Pointer, stdio: AbiNativeStdio): number; + closeNativePipe(handle: NativeHandle): void; + tryWait( + handle: Pointer, + outExit: number[], + outRunning: number[], + outTimedOut: number[], + ): number; + requestKill(handle: Pointer): number; + requestShutdown(handle: Pointer): number; + warningsJson(handle: Pointer, out: Pointer[]): number; + outputMetadataJson(handle: Pointer, out: Pointer[]): number; + free(handle: Pointer, completion: NativeFreeCompletion): void; + freeError(error: AbiErrorDetail): void; + freeString(value: Pointer): void; +} + +export interface _NativeStreamFactory { + readonly platform: NodeJS.Platform; + readable(handle: NativeHandle): Readable; + writable(handle: NativeHandle): Writable; +} + +const MINIMUM_NODE_VERSION = [24, 21, 0] as const; + +function windowsHandleOptions( + handle: NativeHandle, +): { autoClose: true; windowsHandle: bigint } { + return { + autoClose: true, + windowsHandle: typeof handle === 'bigint' ? handle : BigInt(handle), + }; +} + +function unixFd(handle: NativeHandle): number { + const fd = Number(handle); + if (!Number.isSafeInteger(fd) || fd < 0) { + throw new Error(`native runtime returned invalid file descriptor ${handle}`); + } + return fd; +} + +const nodeStreamFactory: _NativeStreamFactory = { + platform: os.platform(), + readable(handle) { + if (this.platform === 'win32') { + return fs.createReadStream( + '', + windowsHandleOptions(handle) as unknown as + Parameters[1], + ); + } + return new net.Socket({ + fd: unixFd(handle), + readable: true, + writable: false, + }); + }, + writable(handle) { + if (this.platform === 'win32') { + return fs.createWriteStream( + '', + windowsHandleOptions(handle) as unknown as + Parameters[1], + ); + } + return new net.Socket({ + fd: unixFd(handle), + readable: false, + writable: true, + }); + }, +}; + +function bindCoordinatorFunctions( + handle: NativeLibraryHandle, +): _StreamingNativeFacade { + const pointer = koffi.pointer(AbiCoordinator); + const free = bindNativeFunction void>>( + handle, + { + symbol: 'mxc_io_free', + result: 'void', + parameters: [pointer], + }, + ); + return { + spawn: bindNativeFunction(handle, { + symbol: 'mxc_io_spawn_request', + result: 'int32_t', + parameters: [ + 'const char *', + koffi.out(koffi.pointer(AbiCoordinator, 2)), + koffi.out(koffi.pointer(AbiErrorDetailType)), + ], + }), + id: bindNativeFunction(handle, { + symbol: 'mxc_io_id', + result: 'uint32_t', + parameters: [pointer], + }), + takeNativeStdio: bindNativeFunction(handle, { + symbol: 'mxc_io_take_native_stdio', + result: 'int32_t', + parameters: [pointer, koffi.out(koffi.pointer(AbiNativeStdioType))], + }), + closeNativePipe: bindNativeFunction(handle, { + symbol: 'mxc_native_pipe_close', + result: 'void', + parameters: ['intptr_t'], + }), + tryWait: bindNativeFunction(handle, { + symbol: 'mxc_io_try_wait', + result: 'int32_t', + parameters: [ + pointer, + koffi.out(koffi.pointer('int32_t')), + koffi.out(koffi.pointer('int32_t')), + koffi.out(koffi.pointer('int32_t')), + ], + }), + requestKill: bindNativeFunction(handle, { + symbol: 'mxc_io_request_kill', + result: 'int32_t', + parameters: [pointer], + }), + requestShutdown: bindNativeFunction(handle, { + symbol: 'mxc_io_request_shutdown', + result: 'int32_t', + parameters: [pointer], + }), + warningsJson: bindNativeFunction(handle, { + symbol: 'mxc_io_warnings_json', + result: 'int32_t', + parameters: [pointer, koffi.out(koffi.pointer('char', 2))], + }), + outputMetadataJson: bindNativeFunction(handle, { + symbol: 'mxc_io_output_metadata_json', + result: 'int32_t', + parameters: [pointer, koffi.out(koffi.pointer('char', 2))], + }), + free(coordinator, completion) { + free.async(coordinator, completion); + }, + freeError: bindNativeFunction(handle, { + symbol: 'mxc_error_detail_free', + result: 'void', + parameters: [koffi.pointer(AbiErrorDetailType)], + }), + freeString: bindNativeFunction(handle, { + symbol: 'mxc_string_free', + result: 'void', + parameters: ['char *'], + }), + } as _StreamingNativeFacade; +} + +let sharedNative: _StreamingNativeFacade | undefined; + +function getNative(): _StreamingNativeFacade { + return sharedNative ??= bindCoordinatorFunctions(loadMxcFfi().handle); +} + +function throwIfFailed(status: number, message: string): void { + if (status !== 0) throw nativeStatusError(status, {}, message); +} + +function isMissingHandle( + handle: NativeHandle, + platform: NodeJS.Platform, +): boolean { + return platform === 'win32' + ? handle === 0 || handle === 0n + : handle === -1 || handle === -1n; +} + +function destroyStream(stream: Readable | Writable | null): void { + if (stream !== null && !stream.destroyed) stream.destroy(); +} + +function freeCoordinatorAsync( + native: _StreamingNativeFacade, + handle: Pointer, +): Promise { + return new Promise((resolve, reject) => { + native.free(handle, (error) => { + if (error === null) resolve(); + else reject(error); + }); + }); +} + +class KoffiLifecycleDriver implements NativeLifecycleDriver { + private freePromise: Promise | undefined; + private shutdownRequested = false; + + constructor( + private readonly native: _StreamingNativeFacade, + private readonly handle: Pointer, + readonly id: number, + readonly standardInput: Writable | null, + readonly standardOutput: Readable | null, + readonly standardError: Readable | null, + ) {} + + poll(): NativeLifecycleStatus { + const exit = [0]; + const running = [1]; + const timedOut = [0]; + throwIfFailed( + this.native.tryWait( + this.handle, + exit, + running, + timedOut, + ), + 'polling sandbox process failed', + ); + return { + exitCode: exit[0], + running: running[0] !== 0, + timedOut: timedOut[0] !== 0, + }; + } + + warnings(): readonly string[] { + return parseStringArray( + this.readOwnedJson(this.native.warningsJson), + ); + } + + outputMetadata(): unknown | undefined { + const json = this.readOwnedJson(this.native.outputMetadataJson); + return json === undefined ? undefined : JSON.parse(json); + } + + kill(): void { + throwIfFailed( + this.native.requestKill(this.handle), + 'killing sandbox process failed', + ); + } + + shutdown(): void { + if (this.shutdownRequested) return; + throwIfFailed( + this.native.requestShutdown(this.handle), + 'shutting down sandbox process failed', + ); + this.shutdownRequested = true; + } + + free(): Promise { + return this.freePromise ??= freeCoordinatorAsync(this.native, this.handle); + } + + private readOwnedJson( + read: (handle: Pointer, out: Pointer[]) => number, + ): string | undefined { + const out: Pointer[] = [null]; + throwIfFailed(read(this.handle, out), 'reading sandbox process data failed'); + try { + return decodeString(out[0]); + } finally { + if (out[0] !== null) this.native.freeString(out[0]); + } + } +} + +function adoptEndpoint( + factory: _NativeStreamFactory, + handle: NativeHandle, + writable: boolean, +): Readable | Writable | null { + if (isMissingHandle(handle, factory.platform)) return null; + + return writable ? factory.writable(handle) : factory.readable(handle); +} + +function closeUnadopted( + native: _StreamingNativeFacade, + factory: _NativeStreamFactory, + handles: NativeHandle[], +): void { + for (const handle of handles) { + if (!isMissingHandle(handle, factory.platform)) { + native.closeNativePipe(handle); + } + } +} + +function beginFailedSpawnCleanup( + native: _StreamingNativeFacade, + handle: Pointer, +): void { + try { + native.requestShutdown(handle); + } finally { + void freeCoordinatorAsync(native, handle).catch(() => {}); + } +} + +/** @internal Builds a driver with injectable native and stream facades. */ +export function _spawnStreamingDriverForTest( + request: RequestSpec, + native: _StreamingNativeFacade, + factory: _NativeStreamFactory, +): NativeLifecycleDriver { + const outHandle: Pointer[] = [null]; + const error = {} as AbiErrorDetail; + const status = native.spawn(JSON.stringify(request), outHandle, error); + if (status !== 0) { + try { + throw nativeStatusError(status, error); + } finally { + native.freeError(error); + } + } + + const handle = outHandle[0]; + if (handle === null || handle === undefined) { + throw new Error('native runtime returned a null lifecycle handle'); + } + + let input: Writable | null = null; + let output: Readable | null = null; + let errorOutput: Readable | null = null; + try { + const stdio = {} as AbiNativeStdio; + throwIfFailed( + native.takeNativeStdio(handle, stdio), + 'taking native stdio failed', + ); + const remaining = [ + stdio.stdin_handle, + stdio.stdout_handle, + stdio.stderr_handle, + ]; + try { + const invalid = factory.platform === 'win32' ? 0 : -1; + input = adoptEndpoint(factory, remaining[0], true) as Writable | null; + remaining[0] = invalid; + output = adoptEndpoint( + factory, + remaining[1], + false, + ) as Readable | null; + remaining[1] = invalid; + errorOutput = adoptEndpoint( + factory, + remaining[2], + false, + ) as Readable | null; + remaining[2] = invalid; + } catch (error) { + closeUnadopted(native, factory, remaining); + throw error; + } + + const id = native.id(handle); + return new KoffiLifecycleDriver( + native, + handle, + id, + input, + output, + errorOutput, + ); + } catch (error) { + destroyStream(input); + destroyStream(output); + destroyStream(errorOutput); + beginFailedSpawnCleanup(native, handle); + throw error; + } +} + +/** @internal Tests whether a Node.js version supports native Windows handles. */ +export function _isSupportedNodeVersionForTest(version: string): boolean { + const current = version + .split('.') + .slice(0, 3) + .map((component) => Number.parseInt(component, 10)); + const difference = current.findIndex( + (component, index) => component !== MINIMUM_NODE_VERSION[index], + ); + return difference === -1 || + current[difference] > MINIMUM_NODE_VERSION[difference]; +} + +function ensureSupportedNodeVersion(): void { + if (!_isSupportedNodeVersionForTest(process.versions.node)) { + throw new Error( + `native stdio requires Node.js ${MINIMUM_NODE_VERSION.join('.')} or newer; ` + + `current version is ${process.versions.node}`, + ); + } +} + +function spawnDriver(request: RequestSpec): NativeLifecycleDriver { + ensureSupportedNodeVersion(); + return _spawnStreamingDriverForTest(request, getNative(), nodeStreamFactory); +} + +export function spawnBindingSandboxProcess( + request: RequestSpec, +): MxcSandboxProcess { + const driver = spawnDriver(request); + try { + return _createMxcSandboxProcess(driver); + } catch (error) { + destroyStream(driver.standardInput); + destroyStream(driver.standardOutput); + destroyStream(driver.standardError); + try { + driver.shutdown(); + } catch { + // Preserve the construction failure while still releasing native state. + } finally { + void driver.free().catch(() => {}); + } + throw error; + } +} diff --git a/sdk/node/src/sandbox-process.ts b/sdk/node/src/sandbox-process.ts new file mode 100644 index 000000000..2f85c5bbb --- /dev/null +++ b/sdk/node/src/sandbox-process.ts @@ -0,0 +1,267 @@ +// Copyright (c) Microsoft Corporation. +// Licensed under the MIT License. + +import type { Readable, Writable } from 'node:stream'; + +export interface SandboxWaitResult { + exitCode: number; + timedOut: boolean; +} + +export interface NativeLifecycleStatus extends SandboxWaitResult { + running: boolean; +} + +export interface NativeLifecycleDriver { + readonly id: number; + readonly standardInput: Writable | null; + readonly standardOutput: Readable | null; + readonly standardError: Readable | null; + poll(): NativeLifecycleStatus; + warnings(): readonly string[]; + outputMetadata(): unknown | undefined; + kill(): void; + shutdown(): void; + free(): Promise; +} + +const POLL_INTERVAL_MS = 10; + +function asError(error: unknown): Error { + return error instanceof Error ? error : new Error(String(error)); +} + +function destroyStream(stream: Readable | Writable | null): void { + if (stream !== null && !stream.destroyed) stream.destroy(); +} + +/** + * A sandbox process whose stdio is backed by native Node streams. + * + * Rust retains lifecycle ownership; Node owns the transferred stdio + * descriptors and provides their normal buffering and backpressure. + */ +export class MxcSandboxProcess { + readonly id: number; + + private readonly input: Writable | null; + private readonly output: Readable | null; + private readonly errorOutput: Readable | null; + private inputTaken = false; + private outputTaken = false; + private errorTaken = false; + private outputDrained = false; + private errorDrained = false; + private disposed = false; + private terminal = false; + private cleanupStarted = false; + private warningsValue: readonly string[]; + private metadataValue: unknown | undefined; + private readonly cleanups: Array<() => void> = []; + private resolveWait!: (result: SandboxWaitResult) => void; + private rejectWait!: (error: Error) => void; + private readonly waitPromise: Promise; + private pollTimer: NodeJS.Timeout | undefined; + + private constructor(private readonly driver: NativeLifecycleDriver) { + this.id = driver.id; + this.input = driver.standardInput; + this.output = driver.standardOutput; + this.errorOutput = driver.standardError; + this.warningsValue = driver.warnings(); + this.waitPromise = new Promise((resolve, reject) => { + this.resolveWait = resolve; + this.rejectWait = reject; + }); + void this.waitPromise.catch(() => {}); + this.input?.on('error', (error) => void this.finish(undefined, asError(error))); + this.output?.on('error', (error) => void this.finish(undefined, asError(error))); + this.errorOutput?.on('error', (error) => void this.finish(undefined, asError(error))); + this.poll(); + } + + static create(driver: NativeLifecycleDriver): MxcSandboxProcess { + return new MxcSandboxProcess(driver); + } + + get standardInput(): Writable | null { + this.throwIfDisposed(); + this.inputTaken = true; + return this.input; + } + + get standardOutput(): Readable | null { + this.throwIfDisposed(); + if (this.outputDrained) { + throw new Error('standard output is being drained internally'); + } + this.outputTaken = true; + return this.output; + } + + get standardError(): Readable | null { + this.throwIfDisposed(); + if (this.errorDrained) { + throw new Error('standard error is being drained internally'); + } + this.errorTaken = true; + return this.errorOutput; + } + + get warnings(): readonly string[] { + return this.warningsValue; + } + + get outputMetadata(): unknown | undefined { + return this.metadataValue; + } + + waitAsync(): Promise { + this.throwIfDisposed(); + if (!this.inputTaken && this.input !== null && !this.input.destroyed) { + this.input.end(); + } + if (!this.outputTaken) this.drainOutput(); + if (!this.errorTaken) this.drainError(); + return this.waitPromise; + } + + kill(): void { + this.throwIfDisposed(); + if (!this.terminal) this.driver.kill(); + } + + dispose(): void { + if (this.disposed) return; + this.disposed = true; + let firstError: Error | undefined; + if (!this.terminal) { + try { + this.driver.shutdown(); + } catch (error) { + firstError = asError(error); + } + } + this.stopPolling(); + destroyStream(this.input); + destroyStream(this.output); + destroyStream(this.errorOutput); + const cleanupError = this.runCleanups(); + firstError ??= cleanupError; + void this.driver.free().catch(() => {}); + if (!this.terminal) { + this.terminal = true; + this.rejectWait(new Error('sandbox process was disposed before completion')); + } + if (firstError !== undefined) throw firstError; + } + + /** @internal Registers state-aware cleanup tied to process completion. */ + _registerCleanup(cleanup: () => void): void { + if (this.cleanupStarted) { + cleanup(); + return; + } + this.cleanups.push(cleanup); + } + + private poll(): void { + if (this.disposed || this.terminal) return; + let status: NativeLifecycleStatus; + try { + status = this.driver.poll(); + } catch (error) { + void this.finish(undefined, asError(error)); + return; + } + if (!status.running) { + void this.finish({ + exitCode: status.exitCode, + timedOut: status.timedOut, + }); + return; + } + this.pollTimer = setTimeout(() => this.poll(), POLL_INTERVAL_MS); + } + + private async finish( + result: SandboxWaitResult | undefined, + initialError?: Error, + ): Promise { + if (this.terminal) return; + this.terminal = true; + this.stopPolling(); + + let error = initialError; + if (error === undefined) { + try { + this.warningsValue = this.driver.warnings(); + this.metadataValue = this.driver.outputMetadata(); + } catch (cause) { + error = asError(cause); + } + } + const cleanupError = this.runCleanups(); + error ??= cleanupError; + + try { + this.driver.shutdown(); + } catch (cause) { + error ??= asError(cause); + } + try { + await this.driver.free(); + } catch (cause) { + error ??= asError(cause); + } + + if (error !== undefined) { + this.rejectWait(error); + } else { + this.resolveWait(result!); + } + } + + private runCleanups(): Error | undefined { + if (this.cleanupStarted) return undefined; + this.cleanupStarted = true; + let firstError: Error | undefined; + while (this.cleanups.length > 0) { + const cleanup = this.cleanups.shift()!; + try { + cleanup(); + } catch (error) { + firstError ??= asError(error); + } + } + return firstError; + } + + private stopPolling(): void { + if (this.pollTimer !== undefined) { + clearTimeout(this.pollTimer); + this.pollTimer = undefined; + } + } + + private drainOutput(): void { + this.outputDrained = true; + this.output?.resume(); + } + + private drainError(): void { + this.errorDrained = true; + this.errorOutput?.resume(); + } + + private throwIfDisposed(): void { + if (this.disposed) throw new Error('sandbox process has been disposed'); + } +} + +/** @internal Creates a process around a test or native lifecycle driver. */ +export function _createMxcSandboxProcess( + driver: NativeLifecycleDriver, +): MxcSandboxProcess { + return MxcSandboxProcess.create(driver); +} diff --git a/sdk/node/tests/integration/native-streaming.test.ts b/sdk/node/tests/integration/native-streaming.test.ts new file mode 100644 index 000000000..d7ca0b0ed --- /dev/null +++ b/sdk/node/tests/integration/native-streaming.test.ts @@ -0,0 +1,111 @@ +// Copyright (c) Microsoft Corporation. +// Licensed under the MIT License. + +import assert from 'node:assert'; +import { once } from 'node:events'; +import os from 'node:os'; +import path from 'node:path'; +import { describe, it } from 'node:test'; +import { pathToFileURL } from 'node:url'; +import type { Readable } from 'node:stream'; +import type { ContainerConfig } from '@microsoft/mxc-sdk'; +import { + debugSpawnOptions, + getSdkPackageRoot, + isLinuxBubblewrap, + sandboxSkipReason, + sdk, + supportedVersions, +} from './test-helpers.js'; + +interface NativeSandbox { + readonly standardOutput: Readable | null; + readonly standardError: Readable | null; + waitAsync(): Promise<{ exitCode: number; timedOut: boolean }>; +} + +interface RequestModule { + prepareRequestSpec( + config: ContainerConfig, + options?: { experimental?: boolean }, + ): unknown; +} + +interface StreamingModule { + spawnBindingSandboxProcess(request: unknown): NativeSandbox; +} + +const platformSupport = sdk.getPlatformSupport(); +const schemaVersion = supportedVersions.at(-1)!; +const skipReason = + sandboxSkipReason ?? + (!platformSupport.isSupported ? `Platform not supported: ${platformSupport.reason}` : undefined) ?? + (os.platform() === 'linux' && !isLinuxBubblewrap + ? 'Native streaming requires Bubblewrap on Linux' + : undefined); + +describe(`Internal native streaming (schema ${schemaVersion})`, { skip: skipReason }, () => { + it('delivers output before the sandbox exits', { timeout: 30000 }, async () => { + const packageRoot = getSdkPackageRoot(); + const requestModule = await import(pathToFileURL( + path.join(packageRoot, 'dist', 'bindings', 'request.js'), + ).href) as RequestModule; + const streamingModule = await import(pathToFileURL( + path.join(packageRoot, 'dist', 'bindings', 'streaming.js'), + ).href) as StreamingModule; + + const command = os.platform() === 'win32' + ? 'powershell.exe -NoProfile -Command "Write-Output STREAM_FIRST; ' + + 'Start-Sleep -Milliseconds 500; Write-Output STREAM_SECOND; ' + + '[Console]::Error.WriteLine(\'STREAM_ERROR\')"' + : 'sh -c "printf \'STREAM_FIRST\\n\'; sleep 1; ' + + 'printf \'STREAM_SECOND\\n\'; printf \'STREAM_ERROR\\n\' >&2"'; + const policy = { + version: schemaVersion.raw, + ...(os.platform() === 'win32' ? { ui: { allowWindows: true } } : {}), + }; + const config = sdk.createConfigFromPolicy(policy); + config.process!.commandLine = command; + const request = requestModule.prepareRequestSpec(config, { + experimental: debugSpawnOptions.experimental, + }); + const sandbox = streamingModule.spawnBindingSandboxProcess(request); + const standardOutput = sandbox.standardOutput; + const standardError = sandbox.standardError; + assert.ok(standardOutput, 'streaming stdout should be available'); + assert.ok(standardError, 'streaming stderr should be available'); + + let stdout = ''; + let stderr = ''; + let resolveFirstChunk: (() => void) | undefined; + const firstChunk = new Promise((resolve) => { + resolveFirstChunk = resolve; + }); + standardOutput.on('data', (data: Buffer) => { + stdout += data.toString(); + if (stdout.includes('STREAM_FIRST')) { + resolveFirstChunk?.(); + } + }); + standardError.on('data', (data: Buffer) => { + stderr += data.toString(); + }); + + let completed = false; + const wait = sandbox.waitAsync().then((result) => { + completed = true; + return result; + }); + const outputEnded = once(standardOutput, 'end'); + const errorEnded = once(standardError, 'end'); + await firstChunk; + assert.strictEqual(completed, false, 'first output should arrive before process completion'); + + const result = await wait; + await Promise.all([outputEnded, errorEnded]); + assert.strictEqual(result.exitCode, 0, stderr); + assert.ok(stdout.includes('STREAM_FIRST')); + assert.ok(stdout.includes('STREAM_SECOND')); + assert.ok(stderr.includes('STREAM_ERROR')); + }); +}); diff --git a/sdk/node/tests/integration/package-lock.json b/sdk/node/tests/integration/package-lock.json index 4037804bc..64ed15dd7 100644 --- a/sdk/node/tests/integration/package-lock.json +++ b/sdk/node/tests/integration/package-lock.json @@ -8,36 +8,325 @@ "name": "mxc-integration-tests", "version": "0.0.0", "dependencies": { - "@microsoft/mxc-sdk": "file:../../" + "@microsoft/mxc-sdk": "file:../.." }, "devDependencies": { - "@types/node": "^20.10.0", + "@types/node": "^24.13.4", "@types/semver": "^7.7.1", "rimraf": "^6.1.3", "semver": "^7.7.4", "typescript": "^5.3.3" } }, + "node_modules/@koromix/koffi-android-arm64": { + "version": "3.2.1", + "resolved": "https://ms-feed-25.pkgs.visualstudio.com/1es-public/_packaging/npm-public/npm/registry/@koromix/koffi-android-arm64/-/koffi-android-arm64-3.2.1.tgz", + "integrity": "sha1-Lc55jTZp0re1bpwCqKaidKvWx8Y=", + "cpu": [ + "arm64" + ], + "license": "MIT", + "optional": true, + "os": [ + "android" + ], + "funding": { + "url": "https://liberapay.com/Koromix" + } + }, + "node_modules/@koromix/koffi-android-x64": { + "version": "3.2.1", + "resolved": "https://ms-feed-25.pkgs.visualstudio.com/1es-public/_packaging/npm-public/npm/registry/@koromix/koffi-android-x64/-/koffi-android-x64-3.2.1.tgz", + "integrity": "sha1-TQr8uvZgQjEKwbOs8sWMb/6XN0M=", + "cpu": [ + "x64" + ], + "license": "MIT", + "optional": true, + "os": [ + "android" + ], + "funding": { + "url": "https://liberapay.com/Koromix" + } + }, + "node_modules/@koromix/koffi-darwin-arm64": { + "version": "3.2.1", + "resolved": "https://ms-feed-25.pkgs.visualstudio.com/1es-public/_packaging/npm-public/npm/registry/@koromix/koffi-darwin-arm64/-/koffi-darwin-arm64-3.2.1.tgz", + "integrity": "sha1-IP1RZvBGbrn1lplQbLeK8cUgfQk=", + "cpu": [ + "arm64" + ], + "license": "MIT", + "optional": true, + "os": [ + "darwin" + ], + "funding": { + "url": "https://liberapay.com/Koromix" + } + }, + "node_modules/@koromix/koffi-darwin-x64": { + "version": "3.2.1", + "resolved": "https://ms-feed-25.pkgs.visualstudio.com/1es-public/_packaging/npm-public/npm/registry/@koromix/koffi-darwin-x64/-/koffi-darwin-x64-3.2.1.tgz", + "integrity": "sha1-A+TMtBoFV58aKcKB5IJiEzF89so=", + "cpu": [ + "x64" + ], + "license": "MIT", + "optional": true, + "os": [ + "darwin" + ], + "funding": { + "url": "https://liberapay.com/Koromix" + } + }, + "node_modules/@koromix/koffi-freebsd-arm64": { + "version": "3.2.1", + "resolved": "https://ms-feed-25.pkgs.visualstudio.com/1es-public/_packaging/npm-public/npm/registry/@koromix/koffi-freebsd-arm64/-/koffi-freebsd-arm64-3.2.1.tgz", + "integrity": "sha1-0+Tp8iGTKchvN7woVoXriVpSgb0=", + "cpu": [ + "arm64" + ], + "license": "MIT", + "optional": true, + "os": [ + "freebsd" + ], + "funding": { + "url": "https://liberapay.com/Koromix" + } + }, + "node_modules/@koromix/koffi-freebsd-ia32": { + "version": "3.2.1", + "resolved": "https://ms-feed-25.pkgs.visualstudio.com/1es-public/_packaging/npm-public/npm/registry/@koromix/koffi-freebsd-ia32/-/koffi-freebsd-ia32-3.2.1.tgz", + "integrity": "sha1-hphgO4/QJ8Hs9tw6rX303Bbv398=", + "cpu": [ + "ia32" + ], + "license": "MIT", + "optional": true, + "os": [ + "freebsd" + ], + "funding": { + "url": "https://liberapay.com/Koromix" + } + }, + "node_modules/@koromix/koffi-freebsd-x64": { + "version": "3.2.1", + "resolved": "https://ms-feed-25.pkgs.visualstudio.com/1es-public/_packaging/npm-public/npm/registry/@koromix/koffi-freebsd-x64/-/koffi-freebsd-x64-3.2.1.tgz", + "integrity": "sha1-hROyr0V+Jtl+gNNPY8vwBjhzBYY=", + "cpu": [ + "x64" + ], + "license": "MIT", + "optional": true, + "os": [ + "freebsd" + ], + "funding": { + "url": "https://liberapay.com/Koromix" + } + }, + "node_modules/@koromix/koffi-linux-arm": { + "version": "3.2.1", + "resolved": "https://ms-feed-25.pkgs.visualstudio.com/1es-public/_packaging/npm-public/npm/registry/@koromix/koffi-linux-arm/-/koffi-linux-arm-3.2.1.tgz", + "integrity": "sha1-U8kNnV+aiTlB862CivSHs3lejuE=", + "cpu": [ + "arm" + ], + "license": "MIT", + "optional": true, + "os": [ + "linux" + ], + "funding": { + "url": "https://liberapay.com/Koromix" + } + }, + "node_modules/@koromix/koffi-linux-arm64": { + "version": "3.2.1", + "resolved": "https://ms-feed-25.pkgs.visualstudio.com/1es-public/_packaging/npm-public/npm/registry/@koromix/koffi-linux-arm64/-/koffi-linux-arm64-3.2.1.tgz", + "integrity": "sha1-ddSlPuDQYfoVnMG+eugLBNEAkg8=", + "cpu": [ + "arm64" + ], + "license": "MIT", + "optional": true, + "os": [ + "linux" + ], + "funding": { + "url": "https://liberapay.com/Koromix" + } + }, + "node_modules/@koromix/koffi-linux-ia32": { + "version": "3.2.1", + "resolved": "https://ms-feed-25.pkgs.visualstudio.com/1es-public/_packaging/npm-public/npm/registry/@koromix/koffi-linux-ia32/-/koffi-linux-ia32-3.2.1.tgz", + "integrity": "sha1-TvuhuJjnO2GNIqngyfPp0R192rQ=", + "cpu": [ + "ia32" + ], + "license": "MIT", + "optional": true, + "os": [ + "linux" + ], + "funding": { + "url": "https://liberapay.com/Koromix" + } + }, + "node_modules/@koromix/koffi-linux-loong64": { + "version": "3.2.1", + "resolved": "https://ms-feed-25.pkgs.visualstudio.com/1es-public/_packaging/npm-public/npm/registry/@koromix/koffi-linux-loong64/-/koffi-linux-loong64-3.2.1.tgz", + "integrity": "sha1-Kg2NCvf2bpXiVxEeX1TsIQjYs7g=", + "cpu": [ + "loong64" + ], + "license": "MIT", + "optional": true, + "os": [ + "linux" + ], + "funding": { + "url": "https://liberapay.com/Koromix" + } + }, + "node_modules/@koromix/koffi-linux-riscv64": { + "version": "3.2.1", + "resolved": "https://ms-feed-25.pkgs.visualstudio.com/1es-public/_packaging/npm-public/npm/registry/@koromix/koffi-linux-riscv64/-/koffi-linux-riscv64-3.2.1.tgz", + "integrity": "sha1-1VWiDsNvfBkM8W/vEFWkKtkYPU0=", + "cpu": [ + "riscv64" + ], + "license": "MIT", + "optional": true, + "os": [ + "linux" + ], + "funding": { + "url": "https://liberapay.com/Koromix" + } + }, + "node_modules/@koromix/koffi-linux-x64": { + "version": "3.2.1", + "resolved": "https://ms-feed-25.pkgs.visualstudio.com/1es-public/_packaging/npm-public/npm/registry/@koromix/koffi-linux-x64/-/koffi-linux-x64-3.2.1.tgz", + "integrity": "sha1-I3tGQNc4A+tny7pDe1Eu2HkzgD0=", + "cpu": [ + "x64" + ], + "license": "MIT", + "optional": true, + "os": [ + "linux" + ], + "funding": { + "url": "https://liberapay.com/Koromix" + } + }, + "node_modules/@koromix/koffi-openbsd-ia32": { + "version": "3.2.1", + "resolved": "https://ms-feed-25.pkgs.visualstudio.com/1es-public/_packaging/npm-public/npm/registry/@koromix/koffi-openbsd-ia32/-/koffi-openbsd-ia32-3.2.1.tgz", + "integrity": "sha1-QzemOyiO/9HqwdOHbXlK4rfuPRg=", + "cpu": [ + "ia32" + ], + "license": "MIT", + "optional": true, + "os": [ + "openbsd" + ], + "funding": { + "url": "https://liberapay.com/Koromix" + } + }, + "node_modules/@koromix/koffi-openbsd-x64": { + "version": "3.2.1", + "resolved": "https://ms-feed-25.pkgs.visualstudio.com/1es-public/_packaging/npm-public/npm/registry/@koromix/koffi-openbsd-x64/-/koffi-openbsd-x64-3.2.1.tgz", + "integrity": "sha1-xZycfkCVP/t+MvjY+8sZ/zio7wM=", + "cpu": [ + "x64" + ], + "license": "MIT", + "optional": true, + "os": [ + "openbsd" + ], + "funding": { + "url": "https://liberapay.com/Koromix" + } + }, + "node_modules/@koromix/koffi-win32-arm64": { + "version": "3.2.1", + "resolved": "https://ms-feed-25.pkgs.visualstudio.com/1es-public/_packaging/npm-public/npm/registry/@koromix/koffi-win32-arm64/-/koffi-win32-arm64-3.2.1.tgz", + "integrity": "sha1-kVoYz1dHK16cbXFoE7w5mDekUU4=", + "cpu": [ + "arm64" + ], + "license": "MIT", + "optional": true, + "os": [ + "win32" + ], + "funding": { + "url": "https://liberapay.com/Koromix" + } + }, + "node_modules/@koromix/koffi-win32-ia32": { + "version": "3.2.1", + "resolved": "https://ms-feed-25.pkgs.visualstudio.com/1es-public/_packaging/npm-public/npm/registry/@koromix/koffi-win32-ia32/-/koffi-win32-ia32-3.2.1.tgz", + "integrity": "sha1-n60joxMqlAcm//MI5bKXVRepoys=", + "cpu": [ + "ia32" + ], + "license": "MIT", + "optional": true, + "os": [ + "win32" + ], + "funding": { + "url": "https://liberapay.com/Koromix" + } + }, + "node_modules/@koromix/koffi-win32-x64": { + "version": "3.2.1", + "resolved": "https://ms-feed-25.pkgs.visualstudio.com/1es-public/_packaging/npm-public/npm/registry/@koromix/koffi-win32-x64/-/koffi-win32-x64-3.2.1.tgz", + "integrity": "sha1-nTgKsOqLz7oW9P7aLJ6jYw/1Yg8=", + "cpu": [ + "x64" + ], + "license": "MIT", + "optional": true, + "os": [ + "win32" + ], + "funding": { + "url": "https://liberapay.com/Koromix" + } + }, "node_modules/@microsoft/mxc-sdk": { "version": "0.8.0", "resolved": "file:../..", "license": "MIT", "dependencies": { + "koffi": "^3.2.1", "node-pty": "^1.2.0-beta.12", "semver": "^7.7.4" }, "engines": { - "node": ">=18.0.0" + "node": ">=24.21.0" } }, "node_modules/@types/node": { - "version": "20.19.43", - "resolved": "https://registry.npmjs.org/@types/node/-/node-20.19.43.tgz", - "integrity": "sha512-6oYBAi5ikg4Pl+kGsoYtawUMBT2zZMCvPNF7pVLnHZfd1zf38DRiWn/gT01RYCdUqkv7Fhr+C9ot4/tb+2sVvA==", + "version": "24.13.4", + "resolved": "https://ms-feed-25.pkgs.visualstudio.com/1es-public/_packaging/npm-public/npm/registry/@types/node/-/node-24.13.4.tgz", + "integrity": "sha1-5U90zq5oOL/1eoFLfTH1BCoPkHQ=", "dev": true, "license": "MIT", "dependencies": { - "undici-types": "~6.21.0" + "undici-types": "~7.18.0" } }, "node_modules/@types/semver": { @@ -88,6 +377,36 @@ "url": "https://github.com/sponsors/isaacs" } }, + "node_modules/koffi": { + "version": "3.2.1", + "resolved": "https://ms-feed-25.pkgs.visualstudio.com/1es-public/_packaging/npm-public/npm/registry/koffi/-/koffi-3.2.1.tgz", + "integrity": "sha1-lPmhEdssUTzdchscXYwk2DjCziU=", + "hasInstallScript": true, + "license": "MIT", + "funding": { + "url": "https://liberapay.com/Koromix" + }, + "optionalDependencies": { + "@koromix/koffi-android-arm64": "3.2.1", + "@koromix/koffi-android-x64": "3.2.1", + "@koromix/koffi-darwin-arm64": "3.2.1", + "@koromix/koffi-darwin-x64": "3.2.1", + "@koromix/koffi-freebsd-arm64": "3.2.1", + "@koromix/koffi-freebsd-ia32": "3.2.1", + "@koromix/koffi-freebsd-x64": "3.2.1", + "@koromix/koffi-linux-arm": "3.2.1", + "@koromix/koffi-linux-arm64": "3.2.1", + "@koromix/koffi-linux-ia32": "3.2.1", + "@koromix/koffi-linux-loong64": "3.2.1", + "@koromix/koffi-linux-riscv64": "3.2.1", + "@koromix/koffi-linux-x64": "3.2.1", + "@koromix/koffi-openbsd-ia32": "3.2.1", + "@koromix/koffi-openbsd-x64": "3.2.1", + "@koromix/koffi-win32-arm64": "3.2.1", + "@koromix/koffi-win32-ia32": "3.2.1", + "@koromix/koffi-win32-x64": "3.2.1" + } + }, "node_modules/lru-cache": { "version": "11.5.1", "resolved": "https://registry.npmjs.org/lru-cache/-/lru-cache-11.5.1.tgz", @@ -211,9 +530,9 @@ } }, "node_modules/undici-types": { - "version": "6.21.0", - "resolved": "https://registry.npmjs.org/undici-types/-/undici-types-6.21.0.tgz", - "integrity": "sha512-iwDZqg0QAGrg9Rav5H4n0M64c3mkR59cJ6wQp+7C4nI0gsmExaedaYLNO44eT4AtBBwjbTiGPMlt2Md0T9H9JQ==", + "version": "7.18.2", + "resolved": "https://ms-feed-25.pkgs.visualstudio.com/1es-public/_packaging/npm-public/npm/registry/undici-types/-/undici-types-7.18.2.tgz", + "integrity": "sha1-KTV6iee3ykrvO/D9P9DNc4hCKek=", "dev": true, "license": "MIT" } diff --git a/sdk/node/tests/integration/package.json b/sdk/node/tests/integration/package.json index bd261a314..97efe06ef 100644 --- a/sdk/node/tests/integration/package.json +++ b/sdk/node/tests/integration/package.json @@ -9,10 +9,10 @@ "test": "node run-tests.js" }, "dependencies": { - "@microsoft/mxc-sdk": "file:../../" + "@microsoft/mxc-sdk": "file:../.." }, "devDependencies": { - "@types/node": "^20.10.0", + "@types/node": "^24.13.4", "@types/semver": "^7.7.1", "rimraf": "^6.1.3", "semver": "^7.7.4", diff --git a/sdk/node/tests/unit/sandbox-process.test.ts b/sdk/node/tests/unit/sandbox-process.test.ts new file mode 100644 index 000000000..feec357cf --- /dev/null +++ b/sdk/node/tests/unit/sandbox-process.test.ts @@ -0,0 +1,177 @@ +// Copyright (c) Microsoft Corporation. +// Licensed under the MIT License. + +import assert from 'node:assert'; +import { once } from 'node:events'; +import { PassThrough } from 'node:stream'; +import { describe, it } from 'node:test'; +import { + _createMxcSandboxProcess, + type NativeLifecycleDriver, + type NativeLifecycleStatus, +} from '../../src/sandbox-process.js'; + +class FakeDriver implements NativeLifecycleDriver { + readonly standardInput = new PassThrough(); + readonly standardOutput = new PassThrough(); + readonly standardError = new PassThrough(); + warningValues = ['initial warning']; + metadata: unknown = undefined; + status: NativeLifecycleStatus = { + exitCode: 0, + running: true, + timedOut: false, + }; + pollError: Error | undefined; + killCount = 0; + shutdownCount = 0; + freeCount = 0; + + constructor(readonly id = 17) {} + + poll(): NativeLifecycleStatus { + if (this.pollError !== undefined) throw this.pollError; + return this.status; + } + + warnings(): readonly string[] { + return this.warningValues; + } + + outputMetadata(): unknown { + return this.metadata; + } + + kill(): void { + this.killCount += 1; + } + + shutdown(): void { + this.shutdownCount += 1; + } + + async free(): Promise { + this.freeCount += 1; + } + + complete(exitCode = 0, timedOut = false): void { + this.status = { exitCode, running: false, timedOut }; + } +} + +describe('native sandbox process', () => { + it('exposes the transferred Node streams directly', async () => { + const driver = new FakeDriver(); + const proc = _createMxcSandboxProcess(driver); + assert.strictEqual(proc.standardInput, driver.standardInput); + assert.strictEqual(proc.standardOutput, driver.standardOutput); + assert.strictEqual(proc.standardError, driver.standardError); + + const data = once(proc.standardOutput!, 'data'); + driver.standardOutput.write('native output'); + assert.strictEqual((await data)[0].toString(), 'native output'); + + proc.dispose(); + }); + + it('closes untaken stdin and drains untaken output while waiting', async () => { + const driver = new FakeDriver(); + const proc = _createMxcSandboxProcess(driver); + const inputEnded = once(driver.standardInput, 'finish'); + const wait = proc.waitAsync(); + + driver.standardOutput.end(Buffer.alloc(128 * 1024)); + driver.standardError.end(Buffer.alloc(128 * 1024)); + driver.complete(4); + + assert.deepStrictEqual(await wait, { exitCode: 4, timedOut: false }); + await inputEnded; + assert.strictEqual(driver.standardOutput.readableFlowing, true); + assert.strictEqual(driver.standardError.readableFlowing, true); + }); + + it('refreshes warnings and metadata after terminal completion', async () => { + const driver = new FakeDriver(); + const proc = _createMxcSandboxProcess(driver); + assert.deepStrictEqual(proc.warnings, ['initial warning']); + driver.warningValues = ['cleanup warning']; + driver.metadata = { source: 'native' }; + driver.complete(9, true); + + assert.deepStrictEqual(await proc.waitAsync(), { + exitCode: 9, + timedOut: true, + }); + assert.deepStrictEqual(proc.warnings, ['cleanup warning']); + assert.deepStrictEqual(proc.outputMetadata, { source: 'native' }); + assert.strictEqual(driver.shutdownCount, 1); + assert.strictEqual(driver.freeCount, 1); + }); + + it('does not destroy caller-owned output when the process exits', async () => { + const driver = new FakeDriver(); + const proc = _createMxcSandboxProcess(driver); + const output = proc.standardOutput!; + output.pause(); + driver.standardOutput.write('trailing output'); + driver.complete(); + + await proc.waitAsync(); + + assert.strictEqual(output.destroyed, false); + assert.strictEqual(output.read()?.toString(), 'trailing output'); + }); + + it('forwards kill while the process is running', () => { + const driver = new FakeDriver(); + const proc = _createMxcSandboxProcess(driver); + + proc.kill(); + + assert.strictEqual(driver.killCount, 1); + proc.dispose(); + }); + + it('rejects wait when lifecycle polling fails', async () => { + const driver = new FakeDriver(); + const proc = _createMxcSandboxProcess(driver); + driver.pollError = new Error('poll failed'); + + await assert.rejects(proc.waitAsync(), /poll failed/); + assert.strictEqual(driver.shutdownCount, 1); + assert.strictEqual(driver.freeCount, 1); + }); + + it('runs all registered cleanup and preserves the first failure', async () => { + const driver = new FakeDriver(); + const proc = _createMxcSandboxProcess(driver); + let secondRan = false; + proc._registerCleanup(() => { + throw new Error('cleanup failed'); + }); + proc._registerCleanup(() => { + secondRan = true; + }); + driver.complete(); + + await assert.rejects(proc.waitAsync(), /cleanup failed/); + assert.strictEqual(secondRan, true); + }); + + it('disposes streams and native lifecycle ownership exactly once', async () => { + const driver = new FakeDriver(); + const proc = _createMxcSandboxProcess(driver); + const wait = proc.waitAsync(); + + proc.dispose(); + proc.dispose(); + + await assert.rejects(wait, /disposed before completion/); + assert.strictEqual(driver.shutdownCount, 1); + assert.strictEqual(driver.freeCount, 1); + assert.strictEqual(driver.standardInput.destroyed, true); + assert.strictEqual(driver.standardOutput.destroyed, true); + assert.strictEqual(driver.standardError.destroyed, true); + assert.throws(() => proc.standardOutput, /disposed/); + }); +}); diff --git a/sdk/node/tests/unit/streaming-binding.test.ts b/sdk/node/tests/unit/streaming-binding.test.ts new file mode 100644 index 000000000..af0ef7a27 --- /dev/null +++ b/sdk/node/tests/unit/streaming-binding.test.ts @@ -0,0 +1,235 @@ +// Copyright (c) Microsoft Corporation. +// Licensed under the MIT License. + +import assert from 'node:assert'; +import { PassThrough, type Readable, type Writable } from 'node:stream'; +import { describe, it } from 'node:test'; +import { + _isSupportedNodeVersionForTest, + _spawnStreamingDriverForTest, + type _NativeStreamFactory, + type _StreamingNativeFacade, +} from '../../src/bindings/streaming.js'; + +class FakeNative implements _StreamingNativeFacade { + readonly handle = {}; + readonly closedHandles: Array = []; + readonly freedStrings: unknown[] = []; + spawnStatus = 0; + takeStatus = 0; + shutdownCount = 0; + freeCount = 0; + freeErrorCount = 0; + stdinHandle: number | bigint = 11; + stdoutHandle: number | bigint = 12; + stderrHandle: number | bigint = 13; + + spawn(_request: string, outHandle: unknown[], _error: unknown): number { + outHandle[0] = this.handle; + return this.spawnStatus; + } + + id(): number { + return 23; + } + + takeNativeStdio(_handle: unknown, stdio: { + stdin_handle: number | bigint; + stdout_handle: number | bigint; + stderr_handle: number | bigint; + }): number { + stdio.stdin_handle = this.stdinHandle; + stdio.stdout_handle = this.stdoutHandle; + stdio.stderr_handle = this.stderrHandle; + return this.takeStatus; + } + + closeNativePipe(handle: number | bigint): void { + this.closedHandles.push(handle); + } + + tryWait( + _handle: unknown, + exit: number[], + running: number[], + timedOut: number[], + ): number { + exit[0] = 0; + running[0] = 1; + timedOut[0] = 0; + return 0; + } + + requestKill(): number { + return 0; + } + + requestShutdown(): number { + this.shutdownCount += 1; + return 0; + } + + warningsJson(_handle: unknown, out: unknown[]): number { + out[0] = null; + return 0; + } + + outputMetadataJson(_handle: unknown, out: unknown[]): number { + out[0] = null; + return 0; + } + + free(_handle: unknown, completion: (error: Error | null) => void): void { + this.freeCount += 1; + queueMicrotask(() => completion(null)); + } + + freeError(): void { + this.freeErrorCount += 1; + } + + freeString(value: unknown): void { + this.freedStrings.push(value); + } +} + +class FakeStreams implements _NativeStreamFactory { + readonly readableHandles: Array = []; + readonly writableHandles: Array = []; + failHandle: number | bigint | undefined; + + constructor(readonly platform: NodeJS.Platform = 'linux') {} + + readable(handle: number | bigint): Readable { + if (handle === this.failHandle) { + throw new Error('readable construction failed'); + } + this.readableHandles.push(handle); + return new PassThrough(); + } + + writable(handle: number | bigint): Writable { + if (handle === this.failHandle) { + throw new Error('writable construction failed'); + } + this.writableHandles.push(handle); + return new PassThrough(); + } +} + +describe('native streaming binding ownership', () => { + it('passes raw Windows handles directly to the stream factory', async () => { + const native = new FakeNative(); + native.stdinHandle = 0x100000001n; + native.stdoutHandle = 0x100000002n; + native.stderrHandle = 0x100000003n; + const streams = new FakeStreams('win32'); + + const driver = _spawnStreamingDriverForTest( + {} as never, + native, + streams, + ); + + assert.deepStrictEqual(streams.writableHandles, [0x100000001n]); + assert.deepStrictEqual( + streams.readableHandles, + [0x100000002n, 0x100000003n], + ); + assert.deepStrictEqual(native.closedHandles, []); + await driver.free(); + }); + + it('adopts each endpoint once with the correct direction', async () => { + const native = new FakeNative(); + const streams = new FakeStreams(); + + const driver = _spawnStreamingDriverForTest( + {} as never, + native, + streams, + ); + + assert.strictEqual(driver.id, 23); + assert.deepStrictEqual(streams.writableHandles, [11]); + assert.deepStrictEqual(streams.readableHandles, [12, 13]); + assert.deepStrictEqual(native.closedHandles, []); + await driver.free(); + await driver.free(); + assert.strictEqual(native.freeCount, 1); + }); + + it('treats platform-specific missing endpoint sentinels as absent', async () => { + const unixNative = new FakeNative(); + unixNative.stdinHandle = -1; + unixNative.stderrHandle = -1n; + const unixStreams = new FakeStreams(); + const unixDriver = _spawnStreamingDriverForTest( + {} as never, + unixNative, + unixStreams, + ); + + assert.strictEqual(unixDriver.standardInput, null); + assert.strictEqual(unixDriver.standardError, null); + assert.deepStrictEqual(unixStreams.readableHandles, [12]); + await unixDriver.free(); + + const windowsNative = new FakeNative(); + windowsNative.stdinHandle = 0n; + windowsNative.stderrHandle = 0; + const windowsStreams = new FakeStreams('win32'); + const windowsDriver = _spawnStreamingDriverForTest( + {} as never, + windowsNative, + windowsStreams, + ); + + assert.strictEqual(windowsDriver.standardInput, null); + assert.strictEqual(windowsDriver.standardError, null); + assert.deepStrictEqual(windowsStreams.readableHandles, [12]); + await windowsDriver.free(); + }); + + it('closes the failed and remaining handles when stream construction fails', async () => { + const native = new FakeNative(); + const streams = new FakeStreams(); + streams.failHandle = 12; + + assert.throws( + () => _spawnStreamingDriverForTest({} as never, native, streams), + /readable construction failed/, + ); + + assert.deepStrictEqual(native.closedHandles, [12, 13]); + await new Promise((resolve) => setImmediate(resolve)); + assert.strictEqual(native.freeCount, 1); + }); + + it('frees native error detail when spawn fails', () => { + const native = new FakeNative(); + native.spawnStatus = 12; + + assert.throws( + () => _spawnStreamingDriverForTest( + {} as never, + native, + new FakeStreams(), + ), + /native runtime failed/, + ); + + assert.strictEqual(native.freeErrorCount, 1); + assert.strictEqual(native.freeCount, 0); + }); +}); + +describe('native streaming Node version support', () => { + it('requires Node 24.21.0 or newer', () => { + assert.strictEqual(_isSupportedNodeVersionForTest('24.20.9'), false); + assert.strictEqual(_isSupportedNodeVersionForTest('24.21.0'), true); + assert.strictEqual(_isSupportedNodeVersionForTest('24.21.1'), true); + assert.strictEqual(_isSupportedNodeVersionForTest('25.0.0'), true); + assert.strictEqual(_isSupportedNodeVersionForTest('23.99.99'), false); + }); +}); diff --git a/src/Cargo.lock b/src/Cargo.lock index 044adf05e..b9c8e8455 100644 --- a/src/Cargo.lock +++ b/src/Cargo.lock @@ -1541,6 +1541,7 @@ version = "0.8.0" dependencies = [ "appcontainer_common", "bwrap_common", + "crossbeam-channel", "getrandom 0.2.17", "hyperlight_common", "isolation_session_bindings", @@ -1567,6 +1568,7 @@ version = "0.8.0" dependencies = [ "csbindgen", "mxc-sdk", + "mxc_engine", "semver", "serde", "serde_ignored", diff --git a/src/Cargo.toml b/src/Cargo.toml index 615231007..d56ae8a66 100644 --- a/src/Cargo.toml +++ b/src/Cargo.toml @@ -69,6 +69,7 @@ bwrap_common = { path = "backends/bubblewrap/common" } chrono = { version = "0.4", default-features = false, features = ["std", "clock"] } cidr = "0.3.2" clap = { version = "4", features = ["derive"] } +crossbeam-channel = "0.5" flatbuffers = "25" getrandom = "0.2" hyperlight_common = { path = "backends/hyperlight/common" } diff --git a/src/backends/bubblewrap/common/src/bwrap_runner.rs b/src/backends/bubblewrap/common/src/bwrap_runner.rs index 941a736fb..a9f87cb88 100644 --- a/src/backends/bubblewrap/common/src/bwrap_runner.rs +++ b/src/backends/bubblewrap/common/src/bwrap_runner.rs @@ -27,6 +27,7 @@ use std::collections::HashSet; use std::fmt::Write as FmtWrite; +use std::os::fd::AsFd; use std::os::unix::process::CommandExt; use std::path::{Component, Path, PathBuf}; use std::process::{Child, ChildStdin, Command, Stdio}; @@ -38,8 +39,8 @@ use wxc_common::logger::Logger; use wxc_common::models::{ExecutionRequest, ScriptResponse}; use wxc_common::sandbox_process::{ boxed_closer, cancel_and_join_discard, group_kill, spawn_discard, take_boxed_read, - take_boxed_write, wait_with_timeout, SandboxBackend, SandboxProcess, StdioMode, StreamCloser, - WaitError, + take_boxed_write, wait_with_timeout, NativeStdio, SandboxBackend, SandboxProcess, StdioMode, + StreamCloser, WaitError, }; use wxc_common::unix_proxy_coordinator::UnixProxyCoordinator; use wxc_common::validator::{ @@ -768,6 +769,38 @@ impl BubblewrapSandboxProcess { } impl SandboxProcess for BubblewrapSandboxProcess { + fn take_native_stdio(&mut self) -> std::io::Result> { + let stdio = NativeStdio { + stdin: self + .inner + .stdin + .as_ref() + .map(|stream| stream.as_fd().try_clone_to_owned()) + .transpose()?, + stdout: self + .inner + .stdout + .as_ref() + .map(InterruptibleReader::try_clone_owned_fd) + .transpose()?, + stderr: self + .inner + .stderr + .as_ref() + .map(InterruptibleReader::try_clone_owned_fd) + .transpose()?, + }; + if stdio.is_empty() { + return Ok(None); + } + self.inner.stdin.take(); + self.inner.stdout.take(); + self.inner.stderr.take(); + self.inner.stdout_canceller.take(); + self.inner.stderr_canceller.take(); + Ok(Some(stdio)) + } + fn take_stdin(&mut self) -> Option> { take_boxed_write(&mut self.inner.stdin) } diff --git a/src/backends/isolation_session/common/src/sandbox.rs b/src/backends/isolation_session/common/src/sandbox.rs index 724a73c9d..71e858655 100644 --- a/src/backends/isolation_session/common/src/sandbox.rs +++ b/src/backends/isolation_session/common/src/sandbox.rs @@ -312,6 +312,10 @@ impl SandboxProcess for OneShotSandboxProcess { self.inner.take_stdin() } + fn stdin_closer(&self) -> Option> { + self.inner.stdin_closer() + } + fn take_stdout(&mut self) -> Option> { self.inner.take_stdout() } diff --git a/src/backends/process_container/common/src/appcontainer_runner.rs b/src/backends/process_container/common/src/appcontainer_runner.rs index 34d9c8338..b65e5e4e0 100644 --- a/src/backends/process_container/common/src/appcontainer_runner.rs +++ b/src/backends/process_container/common/src/appcontainer_runner.rs @@ -59,7 +59,7 @@ use wxc_common::process_util::{ }; use wxc_common::sandbox_process::{ boxed_closer, cancel_and_join_discard, spawn_discard, take_boxed_read, take_boxed_write, - SandboxBackend, SandboxProcess, StdioMode, StreamCloser, + NativeStdio, SandboxBackend, SandboxProcess, StdioMode, StreamCloser, }; use wxc_common::script_runner::get_timeout_milliseconds; use wxc_common::validator::{validate_network_policy_support, NetworkPolicySupport}; @@ -1865,6 +1865,7 @@ struct AppContainerSandboxProcess { filesystem_mode: FilesystemMode, preserve_policy: bool, timeout_ms: u32, + coordinator_timed_out: bool, teardown_result: Option>, /// Live guarded WPR capture session, moved from the `SpawnedChild`. /// Stopped and analyzed in `run_teardown` once the child has exited and @@ -1933,6 +1934,7 @@ impl AppContainerSandboxProcess { filesystem_mode, preserve_policy: request.lifecycle.preserve_policy, timeout_ms: child.timeout_ms, + coordinator_timed_out: false, teardown_result: None, capture_session: child.capture_session.take(), capture_output_path: child.capture_output_path.take(), @@ -1975,6 +1977,24 @@ impl AppContainerSandboxProcess { } } + fn timeout_result(&mut self) -> std::io::Result { + wxc_common::telemetry::log_process_event( + &self.identity, + self.pid, + wxc_common::telemetry::ProcessEvent::TimedOut(self.timeout_ms as u64), + ); + if self.audit_enabled() { + let record = self + .audit(AuditEventName::ProcessTimedOut) + .u64("timeout_ms", self.timeout_ms as u64); + self.audit_logger.log_audit_event(&record); + } + Err(std::io::Error::new( + std::io::ErrorKind::TimedOut, + format!("script timed out after {}ms", self.timeout_ms), + )) + } + fn run_teardown(&mut self, allow_trace_transfer: bool) -> std::io::Result<()> { if let Some(result) = &self.teardown_result { return result.clone().map_err(std::io::Error::other); @@ -2078,6 +2098,35 @@ impl SandboxProcess for AppContainerSandboxProcess { self.output_metadata.as_ref() } + fn take_native_stdio(&mut self) -> std::io::Result> { + let stdio = NativeStdio { + stdin: self + .stdin + .as_ref() + .map(|stream| stream.try_clone_owned_handle()) + .transpose()?, + stdout: self + .stdout + .as_ref() + .map(|stream| stream.try_clone_owned_handle()) + .transpose()?, + stderr: self + .stderr + .as_ref() + .map(|stream| stream.try_clone_owned_handle()) + .transpose()?, + }; + if stdio.is_empty() { + return Ok(None); + } + self.stdin.take(); + self.stdout.take(); + self.stderr.take(); + self.stdout_canceller.take(); + self.stderr_canceller.take(); + Ok(Some(stdio)) + } + fn take_stdin(&mut self) -> Option> { take_boxed_write(&mut self.stdin) } @@ -2147,6 +2196,11 @@ impl SandboxProcess for AppContainerSandboxProcess { } } + fn kill_for_timeout(&mut self) -> std::io::Result<()> { + self.coordinator_timed_out = true; + self.kill() + } + fn wait(&mut self) -> std::io::Result { // Close our copy of any not-taken stdin so the child sees EOF and can // exit reliably (an interactive command would otherwise block waiting @@ -2163,6 +2217,8 @@ impl SandboxProcess for AppContainerSandboxProcess { let mut code: u32 = 0; if unsafe { GetExitCodeProcess(self.process.get(), &mut code) }.is_err() { Err(std::io::Error::other("GetExitCodeProcess failed")) + } else if self.coordinator_timed_out { + self.timeout_result() } else { let exit_code = code as i32; wxc_common::telemetry::log_process_event( @@ -2179,23 +2235,7 @@ impl SandboxProcess for AppContainerSandboxProcess { Ok(exit_code) } } - WAIT_TIMEOUT => { - wxc_common::telemetry::log_process_event( - &self.identity, - self.pid, - wxc_common::telemetry::ProcessEvent::TimedOut(self.timeout_ms as u64), - ); - if self.audit_enabled() { - let record = self - .audit(AuditEventName::ProcessTimedOut) - .u64("timeout_ms", self.timeout_ms as u64); - self.audit_logger.log_audit_event(&record); - } - Err(std::io::Error::new( - std::io::ErrorKind::TimedOut, - format!("script timed out after {}ms", self.timeout_ms), - )) - } + WAIT_TIMEOUT => self.timeout_result(), _ => Err(std::io::Error::other("WaitForSingleObject failed")), }; diff --git a/src/backends/process_container/common/src/base_container_runner.rs b/src/backends/process_container/common/src/base_container_runner.rs index 571b1b4f6..01913c767 100644 --- a/src/backends/process_container/common/src/base_container_runner.rs +++ b/src/backends/process_container/common/src/base_container_runner.rs @@ -65,7 +65,7 @@ use wxc_common::process_util::{ }; use wxc_common::sandbox_process::{ boxed_closer, cancel_and_join_discard, spawn_discard, take_boxed_read, take_boxed_write, - SandboxBackend, SandboxProcess, StdioMode, StreamCloser, + NativeStdio, SandboxBackend, SandboxProcess, StdioMode, StreamCloser, }; use wxc_common::script_runner::get_timeout_milliseconds; use wxc_common::string_util; @@ -1556,6 +1556,7 @@ struct BaseContainerSandboxProcess { stdout_canceller: Option, stderr_canceller: Option, timeout_ms: u32, + coordinator_timed_out: bool, preserve_policy: bool, identity: String, proxy_coordinator: ProxyCoordinator, @@ -1607,6 +1608,7 @@ impl BaseContainerSandboxProcess { stdout_canceller, stderr_canceller, timeout_ms: child.timeout_ms, + coordinator_timed_out: false, preserve_policy: child.preserve_policy, identity: sanitize_identity(&std::mem::take(&mut child.identity)).to_string(), proxy_coordinator: std::mem::take(&mut child.proxy_coordinator), @@ -1859,6 +1861,24 @@ impl BaseContainerSandboxProcess { } } + fn timeout_result(&mut self) -> std::io::Result { + wxc_common::telemetry::log_process_event( + &self.identity, + self.pid, + wxc_common::telemetry::ProcessEvent::TimedOut(self.timeout_ms as u64), + ); + if self.audit_enabled() { + let record = self + .audit(AuditEventName::ProcessTimedOut) + .u64("timeout_ms", self.timeout_ms as u64); + self.audit_logger.log_audit_event(&record); + } + Err(std::io::Error::new( + std::io::ErrorKind::TimedOut, + format!("script timed out after {}ms", self.timeout_ms), + )) + } + fn terminate_and_reap(&mut self) -> std::io::Result<()> { self.kill_process_tree()?; unsafe { @@ -2034,6 +2054,35 @@ impl SandboxProcess for BaseContainerSandboxProcess { self.output_metadata.as_ref() } + fn take_native_stdio(&mut self) -> std::io::Result> { + let stdio = NativeStdio { + stdin: self + .stdin + .as_ref() + .map(|stream| stream.try_clone_owned_handle()) + .transpose()?, + stdout: self + .stdout + .as_ref() + .map(|stream| stream.try_clone_owned_handle()) + .transpose()?, + stderr: self + .stderr + .as_ref() + .map(|stream| stream.try_clone_owned_handle()) + .transpose()?, + }; + if stdio.is_empty() { + return Ok(None); + } + self.stdin.take(); + self.stdout.take(); + self.stderr.take(); + self.stdout_canceller.take(); + self.stderr_canceller.take(); + Ok(Some(stdio)) + } + fn take_stdin(&mut self) -> Option> { take_boxed_write(&mut self.stdin) } @@ -2081,6 +2130,11 @@ impl SandboxProcess for BaseContainerSandboxProcess { self.kill_process_tree() } + fn kill_for_timeout(&mut self) -> std::io::Result<()> { + self.coordinator_timed_out = true; + self.kill_process_tree() + } + fn wait(&mut self) -> std::io::Result { // Close our copy of any not-taken stdin so the child sees EOF and can // exit reliably (an interactive command would otherwise block waiting @@ -2097,6 +2151,8 @@ impl SandboxProcess for BaseContainerSandboxProcess { let mut code: u32 = 0; if unsafe { GetExitCodeProcess(self.process.get(), &mut code) }.is_err() { Err(std::io::Error::other("GetExitCodeProcess failed")) + } else if self.coordinator_timed_out { + self.timeout_result() } else { let exit_code = code as i32; wxc_common::telemetry::log_process_event( @@ -2113,23 +2169,7 @@ impl SandboxProcess for BaseContainerSandboxProcess { Ok(exit_code) } } - WAIT_TIMEOUT => { - wxc_common::telemetry::log_process_event( - &self.identity, - self.pid, - wxc_common::telemetry::ProcessEvent::TimedOut(self.timeout_ms as u64), - ); - if self.audit_enabled() { - let record = self - .audit(AuditEventName::ProcessTimedOut) - .u64("timeout_ms", self.timeout_ms as u64); - self.audit_logger.log_audit_event(&record); - } - Err(std::io::Error::new( - std::io::ErrorKind::TimedOut, - format!("script timed out after {}ms", self.timeout_ms), - )) - } + WAIT_TIMEOUT => self.timeout_result(), _ => Err(std::io::Error::other("WaitForSingleObject failed")), }; diff --git a/src/backends/process_container/common/src/dispatcher.rs b/src/backends/process_container/common/src/dispatcher.rs index 04a173860..83a73dea6 100644 --- a/src/backends/process_container/common/src/dispatcher.rs +++ b/src/backends/process_container/common/src/dispatcher.rs @@ -778,6 +778,10 @@ impl SandboxProcess for DaclGuardedProcess { self.inner.take_stdin() } + fn stdin_closer(&self) -> Option> { + self.inner.stdin_closer() + } + fn take_stdout(&mut self) -> Option> { self.inner.take_stdout() } @@ -798,6 +802,10 @@ impl SandboxProcess for DaclGuardedProcess { self.inner.kill() } + fn kill_for_timeout(&mut self) -> std::io::Result<()> { + self.inner.kill_for_timeout() + } + fn wait(&mut self) -> std::io::Result { self.inner.wait() } @@ -1389,6 +1397,9 @@ mod tests { } fn kill(&mut self) -> std::io::Result<()> { self.killed = true; + Err(std::io::Error::other("ordinary kill")) + } + fn kill_for_timeout(&mut self) -> std::io::Result<()> { Ok(()) } fn wait(&mut self) -> std::io::Result { @@ -1415,7 +1426,11 @@ mod tests { ); assert!(matches!(guarded.wait(), Ok(7)), "wait() must delegate"); assert!(guarded.take_stdin().is_none(), "take_stdin() must delegate"); - assert!(guarded.kill().is_ok(), "kill() must delegate"); + assert!(guarded.kill().is_err(), "kill() must delegate"); + assert!( + guarded.kill_for_timeout().is_ok(), + "kill_for_timeout() must delegate" + ); assert!( guarded.output_metadata().is_some(), "output_metadata() must delegate" diff --git a/src/backends/seatbelt/common/src/seatbelt_runner.rs b/src/backends/seatbelt/common/src/seatbelt_runner.rs index 82b8cd4b8..d74de0d40 100644 --- a/src/backends/seatbelt/common/src/seatbelt_runner.rs +++ b/src/backends/seatbelt/common/src/seatbelt_runner.rs @@ -23,6 +23,7 @@ use std::ffi::{CStr, CString}; use std::fmt::Write as FmtWrite; use std::fs; +use std::os::fd::AsFd; use std::os::unix::process::CommandExt; use std::path::{Path, PathBuf}; use std::process::{Command, Stdio}; @@ -33,8 +34,8 @@ use wxc_common::logger::Logger; use wxc_common::models::{ExecutionRequest, LaunchMethod, ProxyAddress, ScriptResponse}; use wxc_common::sandbox_process::{ boxed_closer, cancel_and_join_discard, group_kill, spawn_discard, take_boxed_read, - take_boxed_write, wait_with_timeout, SandboxBackend, SandboxProcess, StdioMode, StreamCloser, - WaitError, + take_boxed_write, wait_with_timeout, NativeStdio, SandboxBackend, SandboxProcess, StdioMode, + StreamCloser, WaitError, }; use wxc_common::unix_proxy_coordinator::UnixProxyCoordinator; use wxc_common::validator::{ @@ -495,6 +496,35 @@ impl SeatbeltSandboxProcess { } impl SandboxProcess for SeatbeltSandboxProcess { + fn take_native_stdio(&mut self) -> std::io::Result> { + let stdio = NativeStdio { + stdin: self + .stdin + .as_ref() + .map(|stream| stream.as_fd().try_clone_to_owned()) + .transpose()?, + stdout: self + .stdout + .as_ref() + .map(InterruptibleReader::try_clone_owned_fd) + .transpose()?, + stderr: self + .stderr + .as_ref() + .map(InterruptibleReader::try_clone_owned_fd) + .transpose()?, + }; + if stdio.is_empty() { + return Ok(None); + } + self.stdin.take(); + self.stdout.take(); + self.stderr.take(); + self.stdout_canceller.take(); + self.stderr_canceller.take(); + Ok(Some(stdio)) + } + fn take_stdin(&mut self) -> Option> { take_boxed_write(&mut self.stdin) } diff --git a/src/core/mxc-sdk/src/sandbox.rs b/src/core/mxc-sdk/src/sandbox.rs index de1413c03..9f7d6a221 100644 --- a/src/core/mxc-sdk/src/sandbox.rs +++ b/src/core/mxc-sdk/src/sandbox.rs @@ -10,7 +10,7 @@ use std::io::{Read, Write}; pub use wxc_common::models::{ CaptureDenialsErrorOutput, CaptureDenialsOutput, SandboxOutputMetadata, }; -use wxc_common::sandbox_process::{SandboxProcess, StreamCloser as InnerCloser}; +use wxc_common::sandbox_process::{NativeStdio, SandboxProcess, StreamCloser as InnerCloser}; /// The outcome of waiting on a [`Sandbox`] (see [`Sandbox::wait`]). /// @@ -105,6 +105,11 @@ impl Sandbox { self.inner.take_stderr() } + #[doc(hidden)] + pub fn take_native_stdio(&mut self) -> std::io::Result> { + self.inner.take_native_stdio() + } + /// A [`StreamCloser`] that unblocks a parked blocking read on stdout without /// killing the child. `None` if stdout was not piped. pub fn stdout_closer(&self) -> Option { diff --git a/src/core/mxc_engine/Cargo.toml b/src/core/mxc_engine/Cargo.toml index 67f406107..44408a2fa 100644 --- a/src/core/mxc_engine/Cargo.toml +++ b/src/core/mxc_engine/Cargo.toml @@ -12,6 +12,7 @@ path = "src/lib.rs" [dependencies] wxc_common.workspace = true mxc_config_contract.workspace = true +crossbeam-channel.workspace = true serde = { workspace = true } serde_json = { workspace = true } # MicroVM runner — used by the Windows and Linux run-to-completion bodies under diff --git a/src/core/mxc_engine/src/io_coordinator.rs b/src/core/mxc_engine/src/io_coordinator.rs new file mode 100644 index 000000000..b61edd42a --- /dev/null +++ b/src/core/mxc_engine/src/io_coordinator.rs @@ -0,0 +1,612 @@ +// Copyright (c) Microsoft Corporation. +// Licensed under the MIT License. + +//! Native process lifecycle coordination for event-loop language bindings. +//! +//! Stdio ownership is transferred to the language runtime as native endpoints. +//! This module retains only process lifecycle responsibilities: timeout +//! enforcement, kill, reaping, warnings, metadata, and terminal status. + +use std::panic::{catch_unwind, AssertUnwindSafe}; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::{Arc, Mutex, MutexGuard}; +use std::thread::{self, JoinHandle}; +use std::time::{Duration, Instant}; + +use crossbeam_channel::{unbounded, Receiver, RecvTimeoutError, Sender}; +use wxc_common::models::SandboxOutputMetadata; +use wxc_common::sandbox_process::{NativeStdio, SandboxProcess}; + +use crate::{spawn, Error, ErrorCode, SandboxRequest}; + +const CONTROL_POLL_INTERVAL: Duration = Duration::from_millis(10); +const MAX_CONSECUTIVE_PROCESS_ERRORS: usize = 100; + +/// A non-blocking snapshot of a coordinated process. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct IoProcessStatus { + /// Whether the process is still running. + pub running: bool, + /// The terminal exit code. Meaningful only when `running` is false. + pub exit_code: i32, + /// Whether the configured deadline terminated the process. + pub timed_out: bool, +} + +/// Failures produced by lifecycle operations after a sandbox has spawned. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum IoCoordinatorError { + /// The coordinator no longer accepts commands. + Closed, + /// The underlying sandbox operation failed. + Backend, +} + +#[derive(Debug, Clone, Copy)] +struct TerminalResult { + exit_code: i32, + timed_out: bool, +} + +#[derive(Default)] +struct ProcessState { + result: Option, + terminal_failed: bool, + error_reported: bool, + warnings: Vec, + output_metadata: Option, +} + +#[derive(Clone, Copy)] +enum ControlCommand { + Kill, + Shutdown, +} + +struct SharedCoordinator { + id: u32, + process: Mutex, + stdio: Mutex>, + control: Sender, + kill_requested: AtomicBool, + shutdown_requested: AtomicBool, +} + +impl SharedCoordinator { + fn request_shutdown(&self) { + if !self.shutdown_requested.swap(true, Ordering::AcqRel) { + let _ = self.control.send(ControlCommand::Shutdown); + } + } +} + +/// Coordinates a sandbox process without transporting its stdio bytes. +pub struct IoCoordinator { + shared: Arc, + worker: Mutex>>, +} + +impl std::fmt::Debug for IoCoordinator { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter + .debug_struct("IoCoordinator") + .field("id", &self.shared.id) + .finish_non_exhaustive() + } +} + +impl Drop for IoCoordinator { + fn drop(&mut self) { + self.shared.request_shutdown(); + if let Some(worker) = lock_unpoisoned(&self.worker).take() { + let _ = worker.join(); + } + } +} + +impl IoCoordinator { + /// Return the child process identifier. + pub fn id(&self) -> u32 { + self.shared.id + } + + /// Transfer the native stdio endpoints exactly once. + pub fn take_native_stdio(&self) -> Option { + lock_unpoisoned(&self.shared.stdio).take() + } + + /// Poll process completion without blocking. + pub fn poll_process(&self) -> Result { + let state = lock_unpoisoned(&self.shared.process); + if let Some(result) = state.result { + return Ok(IoProcessStatus { + running: false, + exit_code: result.exit_code, + timed_out: result.timed_out, + }); + } + if state.terminal_failed { + return Err(IoCoordinatorError::Backend); + } + Ok(IoProcessStatus { + running: true, + exit_code: 0, + timed_out: false, + }) + } + + /// Report whether process monitoring reached a terminal state. + pub fn process_is_terminal(&self) -> bool { + let state = lock_unpoisoned(&self.shared.process); + state.result.is_some() || state.terminal_failed + } + + /// Report whether the native lifecycle worker has exited and been joined. + pub fn workers_finished(&self) -> bool { + let mut worker = lock_unpoisoned(&self.worker); + if worker.as_ref().is_some_and(|worker| !worker.is_finished()) { + return false; + } + if let Some(worker) = worker.take() { + let _ = worker.join(); + } + true + } + + /// Queue a process-tree kill. + pub fn request_kill(&self) -> Result<(), IoCoordinatorError> { + if self.process_is_terminal() || self.shared.kill_requested.swap(true, Ordering::AcqRel) { + return Ok(()); + } + self.shared.control.send(ControlCommand::Kill).or_else(|_| { + self.shared.kill_requested.store(false, Ordering::Release); + if self.process_is_terminal() { + Ok(()) + } else { + Err(IoCoordinatorError::Closed) + } + }) + } + + /// Request process shutdown. + pub fn request_shutdown(&self) { + self.shared.request_shutdown(); + } + + /// Return the latest sandbox warnings. + pub fn warnings(&self) -> Vec { + lock_unpoisoned(&self.shared.process).warnings.clone() + } + + /// Return structured output metadata after terminal completion. + pub fn output_metadata(&self) -> Option { + lock_unpoisoned(&self.shared.process) + .output_metadata + .clone() + } +} + +/// Spawn a sandbox under native lifecycle coordination. +pub fn spawn_io(request: &SandboxRequest) -> Result { + let timeout_ms = (request.inner.script_timeout > 0).then_some(request.inner.script_timeout); + coordinate_io(spawn(request)?, timeout_ms) +} + +/// Wrap an already-spawned process under native lifecycle coordination. +pub fn coordinate_io( + mut process: Box, + timeout_ms: Option, +) -> Result { + let stdio = process + .take_native_stdio() + .map_err(|error| Error::new(ErrorCode::BackendError, error.to_string()))? + .filter(|stdio| !stdio.is_empty()) + .ok_or_else(|| { + Error::new( + ErrorCode::UnsupportedContainment, + "the selected backend does not expose transferable native stdio endpoints", + ) + })?; + Ok(start_coordinator(process, stdio, timeout_ms)) +} + +fn start_coordinator( + process: Box, + stdio: NativeStdio, + timeout_ms: Option, +) -> IoCoordinator { + let id = process.id(); + let warnings = process.warnings(); + let (control_tx, control_rx) = unbounded(); + let shared = Arc::new(SharedCoordinator { + id, + process: Mutex::new(ProcessState { + warnings, + ..ProcessState::default() + }), + stdio: Mutex::new(Some(stdio)), + control: control_tx, + kill_requested: AtomicBool::new(false), + shutdown_requested: AtomicBool::new(false), + }); + let worker_shared = Arc::clone(&shared); + let panic_shared = Arc::clone(&shared); + let worker = thread::spawn(move || { + let mut process = process; + if catch_unwind(AssertUnwindSafe(|| { + run_control( + &mut process, + timeout_ms.filter(|timeout| *timeout > 0), + control_rx, + worker_shared, + ); + })) + .is_err() + { + let mut state = lock_unpoisoned(&panic_shared.process); + state.terminal_failed = true; + state + .warnings + .push("sandbox lifecycle coordinator thread panicked".to_string()); + } + }); + IoCoordinator { + shared, + worker: Mutex::new(Some(worker)), + } +} + +fn run_control( + process: &mut Box, + timeout_ms: Option, + commands: Receiver, + shared: Arc, +) { + let started_at = Instant::now(); + let mut consecutive_process_errors = 0; + loop { + if shared.shutdown_requested.load(Ordering::Acquire) { + if shutdown_process(process, &shared) { + return; + } + if record_consecutive_process_error(&shared, &mut consecutive_process_errors) { + fail_process_control(&shared); + return; + } + thread::sleep(CONTROL_POLL_INTERVAL); + continue; + } + + match commands.recv_timeout(CONTROL_POLL_INTERVAL) { + Ok(ControlCommand::Kill) => { + let finished = match process.kill() { + Ok(()) => { + consecutive_process_errors = 0; + false + } + Err(_) => { + let finished = finish_after_kill_race(&mut **process, &shared); + if !finished + && record_consecutive_process_error( + &shared, + &mut consecutive_process_errors, + ) + { + fail_process_control(&shared); + return; + } + finished + } + }; + shared.kill_requested.store(false, Ordering::Release); + if finished { + return; + } + } + Ok(ControlCommand::Shutdown) | Err(RecvTimeoutError::Disconnected) => { + shared.shutdown_requested.store(true, Ordering::Release); + continue; + } + Err(RecvTimeoutError::Timeout) => {} + } + + match process.try_wait() { + Ok(Some(_)) => { + finish_process(&mut **process, &shared, false); + return; + } + Err(error) if error.kind() == std::io::ErrorKind::TimedOut => { + finish_process(&mut **process, &shared, true); + return; + } + Err(_) => { + if record_consecutive_process_error(&shared, &mut consecutive_process_errors) { + fail_process_control(&shared); + return; + } + } + Ok(None) => consecutive_process_errors = 0, + } + + if timeout_ms + .is_some_and(|timeout| started_at.elapsed() >= Duration::from_millis(timeout.into())) + { + let finished = match process.try_wait() { + Ok(Some(_)) => { + finish_process(&mut **process, &shared, false); + true + } + Err(error) if error.kind() == std::io::ErrorKind::TimedOut => { + finish_process(&mut **process, &shared, true); + true + } + Ok(None) => { + consecutive_process_errors = 0; + terminate_for_timeout(&mut **process, &shared) + } + Err(_) => terminate_for_timeout(&mut **process, &shared), + }; + if finished { + return; + } + if record_consecutive_process_error(&shared, &mut consecutive_process_errors) { + fail_process_control(&shared); + return; + } + } + } +} + +fn terminate_for_timeout(process: &mut dyn SandboxProcess, shared: &SharedCoordinator) -> bool { + if process.kill_for_timeout().is_err() { + if finish_after_kill_race(process, shared) { + true + } else { + record_process_error(shared); + false + } + } else { + finish_process(process, shared, true); + true + } +} + +fn shutdown_process(process: &mut Box, shared: &SharedCoordinator) -> bool { + if process.kill().is_ok() { + finish_process(&mut **process, shared, false); + true + } else if finish_after_kill_race(&mut **process, shared) { + true + } else { + record_process_error(shared); + false + } +} + +fn finish_after_kill_race(process: &mut dyn SandboxProcess, shared: &SharedCoordinator) -> bool { + match process.try_wait() { + Ok(Some(_)) => { + finish_process(process, shared, false); + true + } + Err(error) if error.kind() == std::io::ErrorKind::TimedOut => { + finish_process(process, shared, true); + true + } + Ok(None) | Err(_) => false, + } +} + +fn finish_process( + process: &mut dyn SandboxProcess, + shared: &SharedCoordinator, + forced_timeout: bool, +) { + let result = process.wait(); + let warnings = process.warnings(); + let output_metadata = process.output_metadata().cloned(); + let mut state = lock_unpoisoned(&shared.process); + state.warnings = warnings; + state.output_metadata = output_metadata; + match result { + Ok(exit_code) => { + state.result = Some(TerminalResult { + exit_code, + timed_out: forced_timeout, + }); + } + Err(error) if error.kind() == std::io::ErrorKind::TimedOut => { + state.result = Some(TerminalResult { + exit_code: -1, + timed_out: true, + }); + } + Err(_) => state.terminal_failed = true, + } +} + +fn record_process_error(shared: &SharedCoordinator) { + let mut state = lock_unpoisoned(&shared.process); + if !state.error_reported { + state.error_reported = true; + state + .warnings + .push("sandbox lifecycle monitoring encountered a transient backend error".to_string()); + } +} + +fn record_consecutive_process_error( + shared: &SharedCoordinator, + consecutive_process_errors: &mut usize, +) -> bool { + record_process_error(shared); + *consecutive_process_errors += 1; + *consecutive_process_errors >= MAX_CONSECUTIVE_PROCESS_ERRORS +} + +fn fail_process_control(shared: &SharedCoordinator) { + let mut state = lock_unpoisoned(&shared.process); + state.terminal_failed = true; + state + .warnings + .push("sandbox lifecycle coordinator stopped after repeated backend errors".to_string()); +} + +fn lock_unpoisoned(mutex: &Mutex) -> MutexGuard<'_, T> { + mutex + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) +} + +#[cfg(test)] +mod tests { + use super::*; + use std::collections::VecDeque; + use std::sync::atomic::AtomicUsize; + + enum PollStep { + Running, + Exited(i32), + } + + struct ScriptedProcess { + polls: VecDeque, + exit_code: i32, + kill_count: Arc, + timeout_kill_count: Arc, + } + + impl SandboxProcess for ScriptedProcess { + fn take_stdin(&mut self) -> Option> { + None + } + + fn take_stdout(&mut self) -> Option> { + None + } + + fn take_stderr(&mut self) -> Option> { + None + } + + fn try_wait(&mut self) -> std::io::Result> { + if self.exit_code == -1 { + return Ok(Some(-1)); + } + Ok(match self.polls.pop_front().unwrap_or(PollStep::Running) { + PollStep::Running => None, + PollStep::Exited(code) => { + self.exit_code = code; + Some(code) + } + }) + } + + fn wait(&mut self) -> std::io::Result { + Ok(self.exit_code) + } + + fn id(&self) -> u32 { + 42 + } + + fn kill(&mut self) -> std::io::Result<()> { + self.kill_count.fetch_add(1, Ordering::Relaxed); + self.exit_code = -1; + Ok(()) + } + + fn kill_for_timeout(&mut self) -> std::io::Result<()> { + self.timeout_kill_count.fetch_add(1, Ordering::Relaxed); + self.exit_code = -1; + Ok(()) + } + } + + fn empty_stdio() -> NativeStdio { + NativeStdio { + stdin: None, + stdout: None, + stderr: None, + } + } + + fn wait_for_terminal(coordinator: &IoCoordinator) -> IoProcessStatus { + for _ in 0..200 { + if let Ok(status) = coordinator.poll_process() { + if !status.running { + return status; + } + } + thread::sleep(Duration::from_millis(1)); + } + panic!("coordinator did not reach a terminal state"); + } + + #[test] + fn native_stdio_is_transferred_once() { + let process = Box::new(ScriptedProcess { + polls: VecDeque::from([PollStep::Exited(0)]), + exit_code: 0, + kill_count: Arc::new(AtomicUsize::new(0)), + timeout_kill_count: Arc::new(AtomicUsize::new(0)), + }); + let coordinator = start_coordinator(process, empty_stdio(), None); + + assert!(coordinator.take_native_stdio().is_some()); + assert!(coordinator.take_native_stdio().is_none()); + } + + #[test] + fn terminal_exit_is_reported() { + let process = Box::new(ScriptedProcess { + polls: VecDeque::from([PollStep::Running, PollStep::Exited(23)]), + exit_code: 0, + kill_count: Arc::new(AtomicUsize::new(0)), + timeout_kill_count: Arc::new(AtomicUsize::new(0)), + }); + let coordinator = start_coordinator(process, empty_stdio(), None); + + assert_eq!( + wait_for_terminal(&coordinator), + IoProcessStatus { + running: false, + exit_code: 23, + timed_out: false, + } + ); + } + + #[test] + fn kill_is_serialized_through_the_control_thread() { + let kill_count = Arc::new(AtomicUsize::new(0)); + let process = Box::new(ScriptedProcess { + polls: VecDeque::new(), + exit_code: 0, + kill_count: Arc::clone(&kill_count), + timeout_kill_count: Arc::new(AtomicUsize::new(0)), + }); + let coordinator = start_coordinator(process, empty_stdio(), None); + + coordinator.request_kill().expect("kill should be accepted"); + assert_eq!(wait_for_terminal(&coordinator).exit_code, -1); + assert_eq!(kill_count.load(Ordering::Relaxed), 1); + } + + #[test] + fn timeout_uses_timeout_kill_and_marks_the_result() { + let timeout_kill_count = Arc::new(AtomicUsize::new(0)); + let process = Box::new(ScriptedProcess { + polls: VecDeque::new(), + exit_code: 0, + kill_count: Arc::new(AtomicUsize::new(0)), + timeout_kill_count: Arc::clone(&timeout_kill_count), + }); + let coordinator = start_coordinator(process, empty_stdio(), Some(1)); + + let status = wait_for_terminal(&coordinator); + assert_eq!(status.exit_code, -1); + assert!(status.timed_out); + assert_eq!(timeout_kill_count.load(Ordering::Relaxed), 1); + } +} diff --git a/src/core/mxc_engine/src/lib.rs b/src/core/mxc_engine/src/lib.rs index a082a8f47..46a237a2c 100644 --- a/src/core/mxc_engine/src/lib.rs +++ b/src/core/mxc_engine/src/lib.rs @@ -20,6 +20,8 @@ //! port of the SDK's `createConfigFromPolicy`), for the host's native //! containment or an explicitly selected [`Containment`] backend. //! - [`spawn`] — spawn a streaming [`SandboxProcess`] handle for a request. +//! - [`spawn_io`] / [`coordinate_io`] — transfer native stdio ownership and +//! retain lifecycle control for event-loop language bindings. //! - [`run`] / [`resolve_runner`] (Windows) — run-to-completion backend //! selection and execution. //! - [`run_state_aware`] — state-aware lifecycle backend resolution + dispatch. @@ -34,6 +36,7 @@ mod dispatch; mod error; #[cfg(target_os = "windows")] mod guarded_capture; +mod io_coordinator; mod platform; pub mod policy; mod probe; @@ -44,6 +47,9 @@ mod state_aware; mod verbose_telemetry; pub use error::{Error, ErrorCode}; +pub use io_coordinator::{ + coordinate_io, spawn_io, IoCoordinator, IoCoordinatorError, IoProcessStatus, +}; #[cfg(all(target_os = "windows", feature = "isolation_session"))] pub use platform::isolation_session_available; pub use platform::{platform_support, BubblewrapNetworkSupport, PlatformSupport, ProxyEnforcement}; @@ -518,6 +524,10 @@ impl SandboxProcess for TelemetryProcess { self.inner.take_stdin() } + fn stdin_closer(&self) -> Option> { + self.inner.stdin_closer() + } + fn take_stdout(&mut self) -> Option> { self.inner.take_stdout() } @@ -563,6 +573,17 @@ impl SandboxProcess for TelemetryProcess { result } + fn kill_for_timeout(&mut self) -> std::io::Result<()> { + let result = self.inner.kill_for_timeout(); + if result.is_ok() { + self.emit(&Err(std::io::Error::new( + std::io::ErrorKind::TimedOut, + "sandbox execution timed out", + ))); + } + result + } + fn wait(&mut self) -> std::io::Result { let result = self.inner.wait(); self.emit(&result); @@ -624,6 +645,10 @@ impl SandboxProcess for ProcessWithWarnings { self.inner.take_stdin() } + fn stdin_closer(&self) -> Option> { + self.inner.stdin_closer() + } + fn take_stdout(&mut self) -> Option> { self.inner.take_stdout() } @@ -644,6 +669,10 @@ impl SandboxProcess for ProcessWithWarnings { self.inner.kill() } + fn kill_for_timeout(&mut self) -> std::io::Result<()> { + self.inner.kill_for_timeout() + } + fn wait(&mut self) -> std::io::Result { self.inner.wait() } @@ -778,6 +807,12 @@ mod telemetry_process_tests { assert_eq!(killed.wait().unwrap(), 0); assert!(!killed.active); + let mut timed_out_by_coordinator = wrapped(TryWaitResult::Running); + timed_out_by_coordinator.kill_for_timeout().unwrap(); + assert!(!timed_out_by_coordinator.active); + assert_eq!(timed_out_by_coordinator.wait().unwrap(), 0); + assert!(!timed_out_by_coordinator.active); + let mut exited = wrapped(TryWaitResult::Exited(7)); assert_eq!(exited.try_wait().unwrap(), Some(7)); assert!(!exited.active); diff --git a/src/core/wxc_common/src/exec_stream.rs b/src/core/wxc_common/src/exec_stream.rs index 4dcabf895..c9ad54996 100644 --- a/src/core/wxc_common/src/exec_stream.rs +++ b/src/core/wxc_common/src/exec_stream.rs @@ -49,7 +49,9 @@ use std::sync::{Arc, Mutex}; use std::thread::JoinHandle; use crate::mxc_error::MxcError; -use crate::sandbox_process::{boxed_closer, cancel_and_join_discard, SandboxProcess, StreamCloser}; +use crate::sandbox_process::{ + boxed_closer, cancel_and_join_discard, NativeStdio, OwnedPipe, SandboxProcess, StreamCloser, +}; use crate::state_aware_backend::{ExecHandle, ExecOutcome, PipeHandle}; /// The platform's closer for a cancellable read — fired to make an in-flight @@ -69,9 +71,22 @@ type ReadStream = (Box, StreamCanceller); /// A pipe reaches EOF only once every write handle closes, so dropping the /// caller's duplicate is not enough: the backend keeps its own. Dropping this /// closes both, in that order. +type StdinCloseCallback = Box; + +#[derive(Clone)] +struct StdinCloser(Arc>>); + +impl StreamCloser for StdinCloser { + fn close(&self) { + if let Some(close) = self.0.lock().unwrap_or_else(|e| e.into_inner()).take() { + close(); + } + } +} + struct StdinWriter { writer: Option>, - backend_closer: Option>, + backend_closer: Option, } impl Write for StdinWriter { @@ -93,8 +108,8 @@ impl Write for StdinWriter { impl Drop for StdinWriter { fn drop(&mut self) { drop(self.writer.take()); - if let Some(close) = self.backend_closer.take() { - close(); + if let Some(closer) = self.backend_closer.take() { + closer.close(); } } } @@ -104,6 +119,14 @@ struct PreparedStreams { stdout: Option, stderr: Option, stdin: Option>, + stdin_closer: Option, +} + +#[derive(Clone, Copy)] +struct NativeStdioSource { + stdin: Option, + stdout: Option, + stderr: Option, } /// A streaming [`SandboxProcess`] backed by a state-aware [`ExecHandle`]. @@ -111,6 +134,7 @@ pub struct ExecSandboxProcess { stdout: Option>, stderr: Option>, stdin: Option>, + stdin_closer: Option, /// Closers for the two readable streams, kept whether or not the caller /// takes them: [`wait`](SandboxProcess::wait) fires one to end its own /// safety-drain, and [`stdout_closer`](SandboxProcess::stdout_closer) hands @@ -124,6 +148,9 @@ pub struct ExecSandboxProcess { /// Kills the process tree. Taken by the first [`kill`](SandboxProcess::kill) /// or by `Drop`. terminator: Option Result<(), MxcError> + Send>>, + /// Backend-owned handles retained only as sources for caller-owned + /// duplicates. The backend continues to own and close the originals. + native_stdio: Option, /// The waiter's outcome once joined, so repeat waits are idempotent. Holds /// the outcome rather than a code because a timeout has no code. exit: Option, @@ -186,6 +213,12 @@ impl ExecSandboxProcess { )); } + let native_stdio = NativeStdioSource { + stdin: native_pipe_source(stdin), + stdout: native_pipe_source(stdout), + stderr: native_pipe_source(stderr), + }; + let stdin_closer = stdin_closer.map(|close| StdinCloser(Arc::new(Mutex::new(Some(close))))); let streams = wrap_cancellable_read_checked(stdout, "stdout").and_then(|out| { let err = wrap_cancellable_read_checked(stderr, "stderr")?; let input = wrap_write_checked(stdin, "stdin")?; @@ -195,13 +228,16 @@ impl ExecSandboxProcess { stdin: input.map(|writer| { Box::new(StdinWriter { writer: Some(writer), - backend_closer: stdin_closer, + backend_closer: stdin_closer.clone(), }) as Box }), + stdin_closer, }) }); - Self::from_prepared_streams(streams, waiter, terminator) + let mut process = Self::from_prepared_streams(streams, waiter, terminator)?; + process.native_stdio = Some(native_stdio); + Ok(process) } /// Build the handle from already-classified streams. @@ -219,6 +255,7 @@ impl ExecSandboxProcess { stdout, stderr, stdin, + stdin_closer, } = match streams { Ok(streams) => streams, Err(error) => { @@ -278,10 +315,12 @@ impl ExecSandboxProcess { stdout, stderr, stdin, + stdin_closer, stdout_canceller, stderr_canceller, waiter: Some(waiter_thread), terminator: Some(terminator), + native_stdio: None, exit: None, kill_refused: false, }) @@ -452,10 +491,42 @@ fn spawn_discard_checked( } impl SandboxProcess for ExecSandboxProcess { + fn take_native_stdio(&mut self) -> std::io::Result> { + let Some(source) = self.native_stdio else { + return Ok(None); + }; + if (source.stdin.is_some() && self.stdin.is_none()) + || (source.stdout.is_some() && self.stdout.is_none()) + || (source.stderr.is_some() && self.stderr.is_none()) + { + return Err(std::io::Error::other( + "native stdio must be taken before taking individual streams", + )); + } + + let stdio = NativeStdio { + stdin: duplicate_native_pipe(source.stdin, "stdin")?, + stdout: duplicate_native_pipe(source.stdout, "stdout")?, + stderr: duplicate_native_pipe(source.stderr, "stderr")?, + }; + self.native_stdio.take(); + self.stdin.take(); + self.stdout.take(); + self.stderr.take(); + self.stdin_closer.take(); + self.stdout_canceller.take(); + self.stderr_canceller.take(); + Ok(Some(stdio)) + } + fn take_stdin(&mut self) -> Option> { self.stdin.take() } + fn stdin_closer(&self) -> Option> { + boxed_closer(&self.stdin_closer) + } + fn take_stdout(&mut self) -> Option> { self.stdout.take() } @@ -654,6 +725,31 @@ fn dup_handle_to_file(handle: PipeHandle) -> Option { dup_handle_to_owned(handle).map(std::fs::File::from) } +#[cfg(target_os = "windows")] +fn native_pipe_source(handle: PipeHandle) -> Option { + (!is_null_pipe(handle)).then_some(handle.0 as isize) +} + +#[cfg(target_os = "windows")] +fn duplicate_native_pipe( + source: Option, + stream: &str, +) -> std::io::Result> { + use std::os::windows::io::BorrowedHandle; + let Some(raw) = source else { + return Ok(None); + }; + // SAFETY: `NativeStdioSource` retains backend-owned handles only while the + // backend process object is alive; this call duplicates, never adopts, it. + let borrowed = unsafe { BorrowedHandle::borrow_raw(raw as _) }; + borrowed.try_clone_to_owned().map(Some).map_err(|error| { + std::io::Error::new( + error.kind(), + format!("failed to duplicate the exec {stream} pipe handle: {error}"), + ) + }) +} + #[cfg(not(target_os = "windows"))] fn wrap_read(handle: PipeHandle) -> Option> { dup_fd_to_file(handle).map(|f| Box::new(f) as Box) @@ -750,6 +846,38 @@ fn dup_fd_to_file(handle: PipeHandle) -> Option { borrowed.try_clone_to_owned().ok().map(std::fs::File::from) } +#[cfg(not(target_os = "windows"))] +fn native_pipe_source(handle: PipeHandle) -> Option { + (!is_null_pipe(handle)).then_some(handle as isize) +} + +#[cfg(not(target_os = "windows"))] +fn duplicate_native_pipe( + source: Option, + stream: &str, +) -> std::io::Result> { + use std::os::fd::BorrowedFd; + let Some(raw) = source else { + return Ok(None); + }; + let handle = i32::try_from(raw).map_err(|_| { + std::io::Error::new( + std::io::ErrorKind::InvalidInput, + format!("the exec {stream} pipe fd is out of range"), + ) + })?; + // SAFETY: `NativeStdioSource` retains backend-owned descriptors only while + // the backend process object is alive; this call duplicates, never adopts, + // the descriptor. + let borrowed = unsafe { BorrowedFd::borrow_raw(handle) }; + borrowed.try_clone_to_owned().map(Some).map_err(|error| { + std::io::Error::new( + error.kind(), + format!("failed to duplicate the exec {stream} pipe fd: {error}"), + ) + }) +} + #[cfg(test)] mod tests { use super::*; diff --git a/src/core/wxc_common/src/interruptible_reader.rs b/src/core/wxc_common/src/interruptible_reader.rs index 78d6e0137..e0913f035 100644 --- a/src/core/wxc_common/src/interruptible_reader.rs +++ b/src/core/wxc_common/src/interruptible_reader.rs @@ -17,7 +17,7 @@ //! and sets a flag so later reads short-circuit to EOF. use std::io::{self, Read}; -use std::os::fd::{AsRawFd, FromRawFd, OwnedFd, RawFd}; +use std::os::fd::{AsFd, AsRawFd, FromRawFd, OwnedFd, RawFd}; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::Arc; @@ -124,6 +124,10 @@ impl InterruptibleReader { pub fn canceller(&self) -> ReadCanceller { ReadCanceller(Arc::clone(&self.state)) } + + pub fn try_clone_owned_fd(&self) -> io::Result { + self.fd.as_fd().try_clone_to_owned() + } } /// Wrap an optional child pipe end into an [`InterruptibleReader`] plus a diff --git a/src/core/wxc_common/src/process_util.rs b/src/core/wxc_common/src/process_util.rs index ea494017c..02bd76ed5 100644 --- a/src/core/wxc_common/src/process_util.rs +++ b/src/core/wxc_common/src/process_util.rs @@ -2,6 +2,7 @@ // Licensed under the MIT License. use std::os::windows::ffi::OsStringExt; +use std::os::windows::io::{BorrowedHandle, OwnedHandle as StdOwnedHandle}; use std::path::PathBuf; use std::sync::{Arc, Mutex, Weak}; use std::time::Duration; @@ -69,6 +70,10 @@ impl PipeReader { pub fn new(mut handle: OwnedHandle) -> Self { Self(SendOwnedHandle::take(&mut handle)) } + + pub fn try_clone_owned_handle(&self) -> std::io::Result { + self.0.try_clone_owned_handle() + } } impl std::io::Read for PipeReader { @@ -157,6 +162,10 @@ impl InterruptiblePipeReader { pub fn canceller(&self) -> PipeReadCanceller { PipeReadCanceller(Arc::downgrade(&self.0)) } + + pub fn try_clone_owned_handle(&self) -> std::io::Result { + self.0.handle.try_clone_owned_handle() + } } impl std::io::Read for InterruptiblePipeReader { @@ -253,6 +262,10 @@ impl PipeWriter { pub fn new(mut handle: OwnedHandle) -> Self { Self(SendOwnedHandle::take(&mut handle)) } + + pub fn try_clone_owned_handle(&self) -> std::io::Result { + self.0.try_clone_owned_handle() + } } impl std::io::Write for PipeWriter { @@ -300,6 +313,13 @@ impl SendOwnedHandle { pub fn get(&self) -> HANDLE { HANDLE(self.0 as *mut core::ffi::c_void) } + + fn try_clone_owned_handle(&self) -> std::io::Result { + // SAFETY: this wrapper owns a valid process-wide handle for the + // duration of the borrow; `try_clone_to_owned` duplicates it. + let borrowed = unsafe { BorrowedHandle::borrow_raw(self.get().0) }; + borrowed.try_clone_to_owned() + } } impl Drop for SendOwnedHandle { diff --git a/src/core/wxc_common/src/sandbox_process.rs b/src/core/wxc_common/src/sandbox_process.rs index 3cde6232f..f31e1b575 100644 --- a/src/core/wxc_common/src/sandbox_process.rs +++ b/src/core/wxc_common/src/sandbox_process.rs @@ -22,6 +22,29 @@ use crate::models::{ExecutionRequest, FailurePhase, SandboxOutputMetadata, Scrip use crate::script_runner::ScriptRunner; use crate::validator::{validate_common, validate_network_policy_support, NetworkPolicySupport}; +#[cfg(unix)] +pub type OwnedPipe = std::os::fd::OwnedFd; +#[cfg(windows)] +pub type OwnedPipe = std::os::windows::io::OwnedHandle; + +/// Owned native endpoints for a sandbox process. +/// +/// Taking these endpoints transfers stream ownership to the caller. The +/// process retains only lifecycle control; its normal `take_*` methods return +/// `None` afterward. +#[derive(Debug)] +pub struct NativeStdio { + pub stdin: Option, + pub stdout: Option, + pub stderr: Option, +} + +impl NativeStdio { + pub fn is_empty(&self) -> bool { + self.stdin.is_none() && self.stdout.is_none() && self.stderr.is_none() + } +} + /// A handle to a running sandboxed process. /// /// Modelled on [`std::process::Child`]: the caller may `take_*` the std @@ -85,10 +108,24 @@ pub trait SandboxProcess: Send { None } + /// Transfer owned native stdio endpoints to the caller. + /// + /// Backends that cannot expose OS pipe endpoints return `Ok(None)`. + fn take_native_stdio(&mut self) -> std::io::Result> { + Ok(None) + } + /// Take ownership of the child's stdin so the caller can write to it. /// Returns `None` if already taken. Drop the writer to send EOF. fn take_stdin(&mut self) -> Option>; + /// A closer that interrupts the stdin stream returned by + /// [`take_stdin`](SandboxProcess::take_stdin), including an in-flight + /// blocking write. The default returns `None`. + fn stdin_closer(&self) -> Option> { + None + } + /// Take ownership of the child's stdout for live reading. Returns `None` /// if already taken. A taken stream is **not** drained by /// [`wait`](SandboxProcess::wait). @@ -120,6 +157,15 @@ pub trait SandboxProcess: Send { /// [`wait`](SandboxProcess::wait). fn kill(&mut self) -> std::io::Result<()>; + /// Request termination because the execution deadline elapsed. + /// + /// The default uses the same process-tree termination primitive as + /// [`kill`](SandboxProcess::kill). Wrappers may override this to preserve + /// timeout-specific reporting while delegating the actual termination. + fn kill_for_timeout(&mut self) -> std::io::Result<()> { + self.kill() + } + /// Block until the child exits (honouring the request's `scriptTimeout`, /// where `0` means wait forever) and return its exit code. /// diff --git a/src/ffi/mxc_ffi/Cargo.toml b/src/ffi/mxc_ffi/Cargo.toml index 0e37150ea..00261f111 100644 --- a/src/ffi/mxc_ffi/Cargo.toml +++ b/src/ffi/mxc_ffi/Cargo.toml @@ -14,6 +14,7 @@ crate-type = ["cdylib", "staticlib", "lib"] [dependencies] mxc-sdk = { workspace = true } +mxc_engine = { workspace = true } serde = { workspace = true, features = ["derive"] } serde_ignored = "0.1" serde_json = { workspace = true } diff --git a/src/ffi/mxc_ffi/build.rs b/src/ffi/mxc_ffi/build.rs index e8134dba5..fcc9c4ef6 100644 --- a/src/ffi/mxc_ffi/build.rs +++ b/src/ffi/mxc_ffi/build.rs @@ -16,6 +16,7 @@ fn main() { println!("cargo:rerun-if-changed=src/lib.rs"); println!("cargo:rerun-if-changed=src/error_detail.rs"); + println!("cargo:rerun-if-changed=src/io_coordinator.rs"); println!("cargo:rerun-if-changed=src/streaming.rs"); println!("cargo:rerun-if-changed=src/state_aware.rs"); println!("cargo:rerun-if-changed=build.rs"); diff --git a/src/ffi/mxc_ffi/src/io_coordinator.rs b/src/ffi/mxc_ffi/src/io_coordinator.rs new file mode 100644 index 000000000..c803bd498 --- /dev/null +++ b/src/ffi/mxc_ffi/src/io_coordinator.rs @@ -0,0 +1,467 @@ +// Copyright (c) Microsoft Corporation. +// Licensed under the MIT License. + +//! C translation layer for native stdio and process lifecycle coordination. + +use std::ffi::{c_char, CStr}; +use std::panic::{catch_unwind, AssertUnwindSafe}; +use std::ptr; + +use mxc_engine::{IoCoordinator, IoCoordinatorError}; + +use crate::{ + alloc_cstring, request, status_from_error_code, MxcErrorDetail, MXC_STATUS_BACKEND_ERROR, + MXC_STATUS_INVALID_UTF8, MXC_STATUS_NULL_ARGUMENT, MXC_STATUS_PANIC, MXC_STATUS_SUCCESS, +}; + +/// Opaque lifecycle coordinator for event-loop language bindings. +pub struct MxcIoCoordinator { + inner: IoCoordinator, +} + +/// Caller-owned native stdio endpoints. +/// +/// Values are Win32 `HANDLE`s on Windows and file descriptors on Unix. +/// Absent endpoints are `0` on Windows and `-1` on Unix. +#[repr(C)] +pub struct MxcNativeStdio { + pub stdin_handle: isize, + pub stdout_handle: isize, + pub stderr_handle: isize, +} + +impl MxcNativeStdio { + const fn invalid() -> Self { + #[cfg(target_os = "windows")] + const INVALID: isize = 0; + #[cfg(not(target_os = "windows"))] + const INVALID: isize = -1; + + Self { + stdin_handle: INVALID, + stdout_handle: INVALID, + stderr_handle: INVALID, + } + } +} + +fn coordinator_ref<'a>(handle: *mut MxcIoCoordinator) -> Option<&'a MxcIoCoordinator> { + if handle.is_null() { + None + } else { + // SAFETY: callers of the public FFI functions guarantee a live handle. + Some(unsafe { &*handle }) + } +} + +fn catch_status(operation: &str, body: impl FnOnce() -> i32) -> i32 { + catch_unwind(AssertUnwindSafe(body)).unwrap_or_else(|panic| { + crate::report_panic(operation, &*panic); + MXC_STATUS_PANIC + }) +} + +/// Spawn a one-shot sandbox with native stdio and lifecycle coordination. +/// +/// # Safety +/// - `request_json_utf8` must be null or valid NUL-terminated UTF-8. +/// - `out_handle` must point to writable pointer storage holding no live handle. +/// - `out_error` must be null or point to writable fresh error storage. +#[no_mangle] +pub unsafe extern "C" fn mxc_io_spawn_request( + request_json_utf8: *const c_char, + out_handle: *mut *mut MxcIoCoordinator, + out_error: *mut MxcErrorDetail, +) -> i32 { + unsafe { + spawn_coordinator( + "mxc_io_spawn_request", + request_json_utf8, + out_handle, + out_error, + |request_json| { + let request = request::build_request_from_json(request_json) + .map_err(|error| sdk_error_detail(&error))?; + mxc_engine::spawn_io(&request).map_err(|error| sdk_error_detail(&error)) + }, + ) + } +} + +/// Execute a state-aware request with native stdio and lifecycle coordination. +/// +/// # Safety +/// Pointer requirements are identical to [`mxc_io_spawn_request`]. +#[no_mangle] +pub unsafe extern "C" fn mxc_io_state_aware_exec( + request_json_utf8: *const c_char, + experimental: i32, + out_handle: *mut *mut MxcIoCoordinator, + out_error: *mut MxcErrorDetail, +) -> i32 { + unsafe { + spawn_coordinator( + "mxc_io_state_aware_exec", + request_json_utf8, + out_handle, + out_error, + |request_json| { + let process = mxc_engine::exec_state_aware_json(request_json, experimental != 0) + .map_err(|error| sdk_error_detail(&error))?; + mxc_engine::coordinate_io(process, None).map_err(|error| sdk_error_detail(&error)) + }, + ) + } +} + +unsafe fn spawn_coordinator( + operation: &str, + request_json_utf8: *const c_char, + out_handle: *mut *mut MxcIoCoordinator, + out_error: *mut MxcErrorDetail, + spawn: impl FnOnce(&str) -> Result, +) -> i32 { + if !out_handle.is_null() { + // SAFETY: caller-guaranteed writable pointer-sized storage. + unsafe { *out_handle = ptr::null_mut() }; + } + if !out_error.is_null() { + // SAFETY: caller-guaranteed writable storage for one fresh detail. + unsafe { ptr::write(out_error, MxcErrorDetail::none()) }; + } + if out_handle.is_null() { + return MXC_STATUS_NULL_ARGUMENT; + } + + let outcome = catch_unwind(AssertUnwindSafe(|| { + let request_json = if request_json_utf8.is_null() { + return Err(( + MXC_STATUS_NULL_ARGUMENT, + MxcErrorDetail::from_message("request JSON pointer is null"), + )); + } else { + // SAFETY: caller contract guarantees a valid C string. + unsafe { CStr::from_ptr(request_json_utf8) } + .to_str() + .map_err(|_| { + ( + MXC_STATUS_INVALID_UTF8, + MxcErrorDetail::from_message("request JSON is not UTF-8"), + ) + })? + }; + spawn(request_json) + })) + .unwrap_or_else(|panic| { + crate::report_panic(operation, &*panic); + Err(( + MXC_STATUS_PANIC, + MxcErrorDetail::from_message("the mxc engine panicked"), + )) + }); + + match outcome { + Ok(coordinator) => { + // SAFETY: `out_handle` is non-null and writable. + unsafe { + *out_handle = Box::into_raw(Box::new(MxcIoCoordinator { inner: coordinator })) + }; + MXC_STATUS_SUCCESS + } + Err((status, mut detail)) => { + if out_error.is_null() { + detail.free_strings(); + } else { + // SAFETY: caller-guaranteed writable fresh storage. + unsafe { *out_error = detail }; + } + status + } + } +} + +fn sdk_error_detail(error: &mxc_sdk::Error) -> (i32, MxcErrorDetail) { + ( + status_from_error_code(error.code), + MxcErrorDetail::from_error(error), + ) +} + +/// Return the child process identifier, or zero for an invalid handle. +/// +/// # Safety +/// `handle` must be null or a live coordinator handle. +#[no_mangle] +pub unsafe extern "C" fn mxc_io_id(handle: *mut MxcIoCoordinator) -> u32 { + catch_unwind(AssertUnwindSafe(|| { + coordinator_ref(handle).map_or(0, |coordinator| coordinator.inner.id()) + })) + .unwrap_or_else(|panic| { + crate::report_panic("mxc_io_id", &*panic); + 0 + }) +} + +/// Transfer native stdio endpoints from the coordinator exactly once. +/// +/// # Safety +/// - `handle` must be null or a live coordinator handle. +/// - `out_stdio` must point to writable storage for one [`MxcNativeStdio`]. +#[no_mangle] +pub unsafe extern "C" fn mxc_io_take_native_stdio( + handle: *mut MxcIoCoordinator, + out_stdio: *mut MxcNativeStdio, +) -> i32 { + if out_stdio.is_null() { + return MXC_STATUS_NULL_ARGUMENT; + } + // SAFETY: caller-guaranteed writable storage. + unsafe { ptr::write(out_stdio, MxcNativeStdio::invalid()) }; + catch_status("mxc_io_take_native_stdio", || { + let Some(coordinator) = coordinator_ref(handle) else { + return MXC_STATUS_NULL_ARGUMENT; + }; + let Some(stdio) = coordinator.inner.take_native_stdio() else { + return MXC_STATUS_BACKEND_ERROR; + }; + let result = MxcNativeStdio { + stdin_handle: native_pipe_into_raw(stdio.stdin), + stdout_handle: native_pipe_into_raw(stdio.stdout), + stderr_handle: native_pipe_into_raw(stdio.stderr), + }; + // SAFETY: non-null writable output per the caller contract. + unsafe { ptr::write(out_stdio, result) }; + MXC_STATUS_SUCCESS + }) +} + +#[cfg(target_os = "windows")] +fn native_pipe_into_raw(pipe: Option) -> isize { + use std::os::windows::io::IntoRawHandle; + pipe.map_or(0, |handle| handle.into_raw_handle() as isize) +} + +#[cfg(not(target_os = "windows"))] +fn native_pipe_into_raw(pipe: Option) -> isize { + use std::os::fd::IntoRawFd; + pipe.map_or(-1, |fd| fd.into_raw_fd() as isize) +} + +/// Close a transferred endpoint that the caller could not adopt. +/// +/// # Safety +/// `handle` must be an owned endpoint returned by +/// [`mxc_io_take_native_stdio`] and must not have been closed or adopted. +#[no_mangle] +pub unsafe extern "C" fn mxc_native_pipe_close(handle: isize) { + if let Err(panic) = catch_unwind(AssertUnwindSafe(|| close_native_pipe(handle))) { + crate::report_panic("mxc_native_pipe_close", &*panic); + } +} + +#[cfg(target_os = "windows")] +fn close_native_pipe(handle: isize) { + use std::os::windows::io::{FromRawHandle, OwnedHandle}; + if handle != 0 { + // SAFETY: caller transfers one live owned handle to this function. + drop(unsafe { OwnedHandle::from_raw_handle(handle as _) }); + } +} + +#[cfg(not(target_os = "windows"))] +fn close_native_pipe(handle: isize) { + use std::os::fd::{FromRawFd, OwnedFd}; + if let Ok(fd) = i32::try_from(handle) { + if fd >= 0 { + // SAFETY: caller transfers one live owned descriptor to this function. + drop(unsafe { OwnedFd::from_raw_fd(fd) }); + } + } +} + +/// Poll process completion without blocking. +/// +/// # Safety +/// All output pointers must be non-null and writable. +#[no_mangle] +pub unsafe extern "C" fn mxc_io_try_wait( + handle: *mut MxcIoCoordinator, + out_exit: *mut i32, + out_running: *mut i32, + out_timed_out: *mut i32, +) -> i32 { + if handle.is_null() || out_exit.is_null() || out_running.is_null() || out_timed_out.is_null() { + return MXC_STATUS_NULL_ARGUMENT; + } + // SAFETY: caller-guaranteed writable outputs. + unsafe { + *out_exit = 0; + *out_running = 1; + *out_timed_out = 0; + } + catch_status("mxc_io_try_wait", || { + let Some(coordinator) = coordinator_ref(handle) else { + return MXC_STATUS_NULL_ARGUMENT; + }; + match coordinator.inner.poll_process() { + Ok(status) => { + // SAFETY: non-null writable outputs per the caller contract. + unsafe { + *out_exit = status.exit_code; + *out_running = i32::from(status.running); + *out_timed_out = i32::from(status.timed_out); + } + MXC_STATUS_SUCCESS + } + Err(IoCoordinatorError::Closed | IoCoordinatorError::Backend) => { + MXC_STATUS_BACKEND_ERROR + } + } + }) +} + +/// Queue a process-tree kill and return immediately. +/// +/// # Safety +/// `handle` must be null or a live coordinator handle. +#[no_mangle] +pub unsafe extern "C" fn mxc_io_request_kill(handle: *mut MxcIoCoordinator) -> i32 { + catch_status("mxc_io_request_kill", || { + coordinator_ref(handle).map_or(MXC_STATUS_NULL_ARGUMENT, |coordinator| { + coordinator + .inner + .request_kill() + .map_or(MXC_STATUS_BACKEND_ERROR, |()| MXC_STATUS_SUCCESS) + }) + }) +} + +/// Request process shutdown without releasing the caller's handle. +/// +/// # Safety +/// `handle` must be null or a live coordinator handle. +#[no_mangle] +pub unsafe extern "C" fn mxc_io_request_shutdown(handle: *mut MxcIoCoordinator) -> i32 { + catch_status("mxc_io_request_shutdown", || { + let Some(coordinator) = coordinator_ref(handle) else { + return MXC_STATUS_NULL_ARGUMENT; + }; + coordinator.inner.request_shutdown(); + MXC_STATUS_SUCCESS + }) +} + +unsafe fn copy_owned_json( + handle: *mut MxcIoCoordinator, + out_json_utf8: *mut *mut c_char, + value: impl FnOnce(&IoCoordinator) -> Option>, +) -> i32 { + if !out_json_utf8.is_null() { + // SAFETY: caller-guaranteed writable pointer-sized storage. + unsafe { *out_json_utf8 = ptr::null_mut() }; + } + if handle.is_null() || out_json_utf8.is_null() { + return MXC_STATUS_NULL_ARGUMENT; + } + let Some(coordinator) = coordinator_ref(handle) else { + return MXC_STATUS_NULL_ARGUMENT; + }; + if let Some(json) = value(&coordinator.inner) { + // SAFETY: caller-guaranteed writable pointer-sized storage. + unsafe { *out_json_utf8 = alloc_cstring(&json) }; + } + MXC_STATUS_SUCCESS +} + +/// Return the latest warnings JSON without blocking. +/// +/// # Safety +/// `out_json_utf8` must be null or point to writable pointer-sized storage. +#[no_mangle] +pub unsafe extern "C" fn mxc_io_warnings_json( + handle: *mut MxcIoCoordinator, + out_json_utf8: *mut *mut c_char, +) -> i32 { + catch_status("mxc_io_warnings_json", || { + // SAFETY: forwarded caller contract. + unsafe { + copy_owned_json(handle, out_json_utf8, |coordinator| { + serde_json::to_vec(&coordinator.warnings()).ok() + }) + } + }) +} + +/// Return terminal output metadata JSON without blocking. +/// +/// # Safety +/// `out_json_utf8` must be null or point to writable pointer-sized storage. +#[no_mangle] +pub unsafe extern "C" fn mxc_io_output_metadata_json( + handle: *mut MxcIoCoordinator, + out_json_utf8: *mut *mut c_char, +) -> i32 { + catch_status("mxc_io_output_metadata_json", || { + // SAFETY: forwarded caller contract. + unsafe { + copy_owned_json(handle, out_json_utf8, |coordinator| { + coordinator + .output_metadata() + .and_then(|metadata| serde_json::to_vec(&metadata).ok()) + }) + } + }) +} + +/// Request shutdown and release the coordinator. +/// +/// This call may block while the native lifecycle worker exits. Event-loop +/// bindings should invoke it through a native worker pool. +/// +/// # Safety +/// `handle` must be null or a live, not-yet-freed coordinator handle. +#[no_mangle] +pub unsafe extern "C" fn mxc_io_free(handle: *mut MxcIoCoordinator) { + if handle.is_null() { + return; + } + if let Err(panic) = catch_unwind(AssertUnwindSafe(|| { + // SAFETY: live unique handle produced by `Box::into_raw`. + drop(unsafe { Box::from_raw(handle) }); + })) { + crate::report_panic("mxc_io_free", &*panic); + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn null_handle_accessors_fail_closed() { + let mut exit = 7; + let mut running = 7; + let mut timed_out = 7; + // SAFETY: null handles are explicitly accepted and writable outputs + // are provided. + let status = + unsafe { mxc_io_try_wait(ptr::null_mut(), &mut exit, &mut running, &mut timed_out) }; + assert_eq!(status, MXC_STATUS_NULL_ARGUMENT); + } + + #[test] + fn native_stdio_output_is_initialized_before_handle_validation() { + let mut stdio = MxcNativeStdio { + stdin_handle: 5, + stdout_handle: 6, + stderr_handle: 7, + }; + // SAFETY: null handles are explicitly accepted and `stdio` is writable. + let status = unsafe { mxc_io_take_native_stdio(ptr::null_mut(), &mut stdio) }; + assert_eq!(status, MXC_STATUS_NULL_ARGUMENT); + let invalid = MxcNativeStdio::invalid(); + assert_eq!(stdio.stdin_handle, invalid.stdin_handle); + assert_eq!(stdio.stdout_handle, invalid.stdout_handle); + assert_eq!(stdio.stderr_handle, invalid.stderr_handle); + } +} diff --git a/src/ffi/mxc_ffi/src/lib.rs b/src/ffi/mxc_ffi/src/lib.rs index f51d86155..f481e5e58 100644 --- a/src/ffi/mxc_ffi/src/lib.rs +++ b/src/ffi/mxc_ffi/src/lib.rs @@ -11,6 +11,10 @@ //! subset this SDK can launch. //! - **Streaming** (`streaming` module) — [`mxc_spawn_request`] accepts the //! same binding request and returns an opaque live handle. +//! - **Event-loop streaming** (`io_coordinator` module) — +//! [`mxc_io_spawn_request`] and [`mxc_io_state_aware_exec`] transfer native +//! stdio endpoints to the caller while retaining timeout, kill, wait, +//! metadata, and teardown control on a native lifecycle thread. //! - **State-aware lifecycle** (`state_aware` module) — [`mxc_state_aware`] //! drives the envelope phases (provision / start / stop / deprovision), and //! [`mxc_state_aware_exec`] runs the exec phase as a live streaming handle @@ -57,10 +61,12 @@ use std::sync::OnceLock; use mxc_sdk::{available_backends, platform_support, run, ErrorCode, SandboxRequest, WaitOutcome}; mod error_detail; +mod io_coordinator; mod request; mod state_aware; mod streaming; pub use error_detail::*; +pub use io_coordinator::*; pub use state_aware::*; pub use streaming::*;