diff --git a/bun.lock b/bun.lock index f16c9ff..60b32a3 100644 --- a/bun.lock +++ b/bun.lock @@ -4,6 +4,10 @@ "workspaces": { "": { "name": "t3layer", + "dependencies": { + "@t3tools/runtime-client": "https://github.com/EtanHey/t3code/releases/download/runtime-client-v0.0.31-rpc.2/t3tools-runtime-client-0.0.31-rpc.2.tgz", + "effect": "4.0.0-beta.102", + }, "devDependencies": { "@types/bun": "1.3.11", "typescript": "~6.0.3", @@ -11,14 +15,58 @@ }, }, "packages": { + "@msgpackr-extract/msgpackr-extract-darwin-arm64": ["@msgpackr-extract/msgpackr-extract-darwin-arm64@3.0.4", "", { "os": "darwin", "cpu": "arm64" }, "sha512-LCkGo6JDfaBhgST7UpPWgNgLINpcpabaHfyz5OBx75nUYxBsaEPxjnyNjWpeb/xBup/682QnBfRBy2/LvPutZQ=="], + + "@msgpackr-extract/msgpackr-extract-darwin-x64": ["@msgpackr-extract/msgpackr-extract-darwin-x64@3.0.4", "", { "os": "darwin", "cpu": "x64" }, "sha512-zExlW9zUJKZH/tOtVMttwjKa4Xm/3KcNjnE3dPN92uCktwavMxpgCA3MoJK/DOnTWsQgo224OaST27/mPNAf+w=="], + + "@msgpackr-extract/msgpackr-extract-linux-arm": ["@msgpackr-extract/msgpackr-extract-linux-arm@3.0.4", "", { "os": "linux", "cpu": "arm" }, "sha512-Tg3yX65f5GbtXLkrYEHE5oibZG9epyYWas7FogTTEJeDEF9JlXJzKgXaNhT3UXlTOeA+AfZpYZYZ0uPj7Cfquw=="], + + "@msgpackr-extract/msgpackr-extract-linux-arm64": ["@msgpackr-extract/msgpackr-extract-linux-arm64@3.0.4", "", { "os": "linux", "cpu": "arm64" }, "sha512-dgX0P/9wGPJeHFBG+ZmhgE6bmtMt7NP5CRBGyyktpopdk/mW4POnrpQsSLtKI1dwpc+pPLuXHDh6vvskyQE/sw=="], + + "@msgpackr-extract/msgpackr-extract-linux-x64": ["@msgpackr-extract/msgpackr-extract-linux-x64@3.0.4", "", { "os": "linux", "cpu": "x64" }, "sha512-8TNXMEjJc3QEy7R/x1INhgiU+XakDAFUzBhaz7+Rbrs8NH5UQeHQxxmzsSBJGyV6I1jW79undiQm8tOI+D+8FQ=="], + + "@msgpackr-extract/msgpackr-extract-win32-x64": ["@msgpackr-extract/msgpackr-extract-win32-x64@3.0.4", "", { "os": "win32", "cpu": "x64" }, "sha512-CmCXPQrkbwExx3j946/PtHWHbYJiCRBRDl4BlkRQcJB/YOwQxJRTpoo7aTsortjgoJ1x7opzTSxn7C+ASSLVjQ=="], + + "@standard-schema/spec": ["@standard-schema/spec@1.1.0", "", {}, "sha512-l2aFy5jALhniG5HgqrD6jXLi/rUWrKvqN/qJx6yoJsgKhblVd+iqqU4RCXavm/jPityDo5TCvKMnpjKnOriy0w=="], + + "@t3tools/runtime-client": ["@t3tools/runtime-client@https://github.com/EtanHey/t3code/releases/download/runtime-client-v0.0.31-rpc.2/t3tools-runtime-client-0.0.31-rpc.2.tgz", { "peerDependencies": { "effect": "4.0.0-beta.102" } }, "sha512-+GRWl2WZxg9MJfQNKDc+BsnGAj8gD8zmlwWPa9Qe/TbEV3v/HRsoTZ8uZfuJX/m4YJnARIGzkhONN4eg9QeqTw=="], + "@types/bun": ["@types/bun@1.3.11", "", { "dependencies": { "bun-types": "1.3.11" } }, "sha512-5vPne5QvtpjGpsGYXiFyycfpDF2ECyPcTSsFBMa0fraoxiQyMJ3SmuQIGhzPg2WJuWxVBoxWJ2kClYTcw/4fAg=="], "@types/node": ["@types/node@26.1.2", "", { "dependencies": { "undici-types": "~8.3.0" } }, "sha512-Vu4a5UFA9rIIFJ7rB/Vaafh9lrCQszopTCx6KjFboXTGQbPNasehVR5TEiithSDGyd1DEiUByggTZsg8jukeIg=="], "bun-types": ["bun-types@1.3.11", "", { "dependencies": { "@types/node": "*" } }, "sha512-1KGPpoxQWl9f6wcZh57LvrPIInQMn2TQ7jsgxqpRzg+l0QPOFvJVH7HmvHo/AiPgwXy+/Thf6Ov3EdVn1vOabg=="], + "detect-libc": ["detect-libc@2.1.2", "", {}, "sha512-Btj2BOOO83o3WyH59e8MgXsxEQVcarkUOpEYrubB0urwnN10yQ364rsiByU11nZlqWYZm05i/of7io4mzihBtQ=="], + + "effect": ["effect@4.0.0-beta.102", "", { "dependencies": { "@standard-schema/spec": "^1.1.0", "fast-check": "^4.9.0", "find-my-way-ts": "^0.1.6", "ini": "^7.0.0", "kubernetes-types": "^1.30.0", "msgpackr": "^2.0.4", "multipasta": "^0.2.8", "toml": "^4.1.2", "uuid": "^14.0.1", "yaml": "^2.9.0" } }, "sha512-z8Y+Q76Hh/kjLFZrXu8tGn6e+tDsg45R+UHhxd190pXxD53OGwf/G/zDxXTkse4HJ5mobNZfitLfUCp4fMvu6w=="], + + "fast-check": ["fast-check@4.9.0", "", { "dependencies": { "pure-rand": "^8.0.0" } }, "sha512-7ms6T7SybUev/PQITciI0yLM2pOSFy5zpG8Ty7tQofcVaQUvrMXp6CBwqF6fThLCLOrfBtuHAtwq6Yu4XPCllg=="], + + "find-my-way-ts": ["find-my-way-ts@0.1.6", "", {}, "sha512-a85L9ZoXtNAey3Y6Z+eBWW658kO/MwR7zIafkIUPUMf3isZG0NCs2pjW2wtjxAKuJPxMAsHUIP4ZPGv0o5gyTA=="], + + "ini": ["ini@7.0.0", "", {}, "sha512-ifK0CgjALofS5bkrcTy4RaQ9Vx2Knf/eLeIO+NaswQEpH1UblrtTSCIvN71qQDMq0PeQ/SSPojvEJp9vvvfr+w=="], + + "kubernetes-types": ["kubernetes-types@1.30.0", "", {}, "sha512-Dew1okvhM/SQcIa2rcgujNndZwU8VnSapDgdxlYoB84ZlpAD43U6KLAFqYo17ykSFGHNPrg0qry0bP+GJd9v7Q=="], + + "msgpackr": ["msgpackr@2.0.5", "", { "optionalDependencies": { "msgpackr-extract": "^3.0.4" } }, "sha512-cef05H/dSYpLpqp3sj/qyZh5vhUYCalnaLO7j1yOmpsR0y/XwLVtK7r5gn+U/F7CTEfMowcGhlUQJDLcLf7jcA=="], + + "msgpackr-extract": ["msgpackr-extract@3.0.4", "", { "dependencies": { "node-gyp-build-optional-packages": "5.2.2" }, "optionalDependencies": { "@msgpackr-extract/msgpackr-extract-darwin-arm64": "3.0.4", "@msgpackr-extract/msgpackr-extract-darwin-x64": "3.0.4", "@msgpackr-extract/msgpackr-extract-linux-arm": "3.0.4", "@msgpackr-extract/msgpackr-extract-linux-arm64": "3.0.4", "@msgpackr-extract/msgpackr-extract-linux-x64": "3.0.4", "@msgpackr-extract/msgpackr-extract-win32-x64": "3.0.4" }, "bin": { "download-msgpackr-prebuilds": "bin/download-prebuilds.js" } }, "sha512-4kmO/MdyUIkLIvTPr8VHLil4AtoKIoniWPIEk5+CDy0xnWC84azhSFmuJ7PxZdsYtiP5kEeQsORAVIeMgxT+Hw=="], + + "multipasta": ["multipasta@0.2.8", "", {}, "sha512-ZPWuMKyv0cSO29f7hozp+k6+crZbQijV8ipMvxNxRf2SwtYGTX1ZX89Kd20VV4H9Znonx+EQn+iy1wGQsJ+b+Q=="], + + "node-gyp-build-optional-packages": ["node-gyp-build-optional-packages@5.2.2", "", { "dependencies": { "detect-libc": "^2.0.1" }, "bin": { "node-gyp-build-optional-packages": "bin.js", "node-gyp-build-optional-packages-optional": "optional.js", "node-gyp-build-optional-packages-test": "build-test.js" } }, "sha512-s+w+rBWnpTMwSFbaE0UXsRlg7hU4FjekKU4eyAih5T8nJuNZT1nNsskXpxmeqSK9UzkBl6UgRlnKc8hz8IEqOw=="], + + "pure-rand": ["pure-rand@8.4.2", "", {}, "sha512-vvuOGgcuPJAirlHvuQw1TrOiw7ptaIXXmIbNuiNOY6lNGJJH49PQ1Kj4nd783nPdQhQdicgOjVI2yI/9BD6/Ng=="], + + "toml": ["toml@4.3.0", "", {}, "sha512-lVb8X9BsPVuH0M4BKeS91tXAmJvCjQ5UIyAbQFaxkKGyUFK2RPkhwaFSQH8vbpl1d23eu/IBH+dwVMHWaq9A5A=="], + "typescript": ["typescript@6.0.3", "", { "bin": { "tsc": "bin/tsc", "tsserver": "bin/tsserver" } }, "sha512-y2TvuxSZPDyQakkFRPZHKFm+KKVqIisdg9/CZwm9ftvKXLP8NRWj38/ODjNbr43SsoXqNuAisEf1GdCxqWcdBw=="], "undici-types": ["undici-types@8.3.0", "", {}, "sha512-j375ScV60dom+YkPFIfTLcOiPxkN/buHz5GobjLhixFuANaNs3C9l4GmrWqejgXWJ7BbJcFYpTEUkS1Ge8bpZQ=="], + + "uuid": ["uuid@14.0.1", "", { "bin": { "uuid": "dist-node/bin/uuid" } }, "sha512-6ZxzVpzDXDa3bJWaHilVayA+BH/1zmxCJoVgvmqJnid/gPoKHxUrS/aC/T6LGQtNHT+XHG9fXPJB4d+IrU30Ew=="], + + "yaml": ["yaml@2.9.0", "", { "bin": { "yaml": "bin.mjs" } }, "sha512-2AvhNX3mb8zd6Zy7INTtSpl1F15HW6Wnqj0srWlkKLcpYl/gMIMJiyuGq2KeI2YFxUPjdlB+3Lc10seMLtL4cA=="], } } diff --git a/package.json b/package.json index eccbd13..0544f33 100644 --- a/package.json +++ b/package.json @@ -10,5 +10,9 @@ "devDependencies": { "@types/bun": "1.3.11", "typescript": "~6.0.3" + }, + "dependencies": { + "@t3tools/runtime-client": "https://github.com/EtanHey/t3code/releases/download/runtime-client-v0.0.31-rpc.2/t3tools-runtime-client-0.0.31-rpc.2.tgz", + "effect": "4.0.0-beta.102" } } diff --git a/src/config.ts b/src/config.ts index 23b5138..2d0a4eb 100644 --- a/src/config.ts +++ b/src/config.ts @@ -5,6 +5,7 @@ export interface ConfigInput { readonly effort?: unknown; readonly contextWindow?: unknown; readonly runtimeMode?: unknown; + readonly interactionMode?: unknown; readonly [key: string]: unknown; } @@ -15,6 +16,7 @@ export interface ExperimentConfig { readonly effort: "high"; readonly contextWindow: "1m"; readonly runtimeMode: "full-access"; + readonly interactionMode: "default"; } export function createConfig(input: ConfigInput): ExperimentConfig { @@ -25,6 +27,7 @@ export function createConfig(input: ConfigInput): ExperimentConfig { "effort", "contextWindow", "runtimeMode", + "interactionMode", ]); for (const field of Object.keys(input)) { @@ -75,6 +78,13 @@ export function createConfig(input: ConfigInput): ExperimentConfig { throw new TypeError("runtimeMode must be full-access"); } + if (input.interactionMode === undefined) { + throw new TypeError("interactionMode is required"); + } + if (input.interactionMode !== "default") { + throw new TypeError("interactionMode must be default"); + } + return Object.freeze({ baseUrl: input.baseUrl, provider: input.provider, @@ -82,5 +92,6 @@ export function createConfig(input: ConfigInput): ExperimentConfig { effort: input.effort, contextWindow: input.contextWindow, runtimeMode: input.runtimeMode, + interactionMode: input.interactionMode, }); } diff --git a/src/facade.ts b/src/facade.ts index 3f5644f..c239f39 100644 --- a/src/facade.ts +++ b/src/facade.ts @@ -1,3 +1,8 @@ +import type { ExperimentConfig } from "./config"; + +type RuntimeMode = ExperimentConfig["runtimeMode"]; +type InteractionMode = ExperimentConfig["interactionMode"]; + export interface NativeProject { readonly projectId: string; readonly workspaceRoot: string; @@ -32,8 +37,8 @@ export interface NativeStartThreadInput { readonly title: string; readonly message: string; readonly modelSelection: ModelSelection; - readonly runtimeMode: string; - readonly interactionMode: string; + readonly runtimeMode: RuntimeMode; + readonly interactionMode: InteractionMode; readonly branch: string | null; readonly worktreePath: string | null; readonly createdAt: string; @@ -55,6 +60,8 @@ export interface NativeStartTurnInput { readonly threadId: string; readonly messageId: string; readonly message: string; + readonly runtimeMode: RuntimeMode; + readonly interactionMode: InteractionMode; readonly createdAt: string; readonly attachments: readonly []; } @@ -87,9 +94,9 @@ export interface NativeThreadObservation { export interface ModelSelection { readonly instanceId: string; readonly model: string; - readonly options: ReadonlyArray<{ + readonly options?: ReadonlyArray<{ readonly id: string; - readonly value: string; + readonly value: string | boolean; }>; } @@ -98,8 +105,8 @@ export interface SpawnInput { readonly title: string; readonly message: string; readonly modelSelection: ModelSelection; - readonly runtimeMode: string; - readonly interactionMode: string; + readonly runtimeMode: RuntimeMode; + readonly interactionMode: InteractionMode; readonly branch: string | null; readonly worktreePath: string | null; } @@ -201,7 +208,10 @@ export class FacadeError extends Error { } } -export interface FacadeOptions { +export interface FacadeOptions extends Pick< + ExperimentConfig, + "runtimeMode" | "interactionMode" +> { readonly id?: () => string; readonly now?: () => string; readonly evidence?: (record: Readonly>) => void; @@ -384,10 +394,7 @@ function assertNonEmptyTerminal(snapshot: NativeThreadSnapshot): void { } } -export function createT3Facade( - runtime: NativeRuntime, - options: FacadeOptions = {}, -) { +export function createT3Facade(runtime: NativeRuntime, options: FacadeOptions) { const id = options.id ?? defaultId; const now = options.now ?? (() => new Date().toISOString()); const evidence = options.evidence ?? (() => undefined); @@ -483,7 +490,7 @@ export function createT3Facade( modelSelection: { instanceId: input.modelSelection.instanceId, model: input.modelSelection.model, - optionCount: input.modelSelection.options.length, + optionCount: input.modelSelection.options?.length ?? 0, }, runtimeMode: input.runtimeMode, interactionMode: input.interactionMode, @@ -570,6 +577,8 @@ export function createT3Facade( threadId: agentId, messageId, message, + runtimeMode: options.runtimeMode, + interactionMode: options.interactionMode, createdAt, attachments: [], }; @@ -578,6 +587,8 @@ export function createT3Facade( commandId, threadId: agentId, messageId, + runtimeMode: options.runtimeMode, + interactionMode: options.interactionMode, createdAt, attachments: 0, messageBytes: new TextEncoder().encode(message).byteLength, diff --git a/src/nativeRuntime.ts b/src/nativeRuntime.ts new file mode 100644 index 0000000..0430760 --- /dev/null +++ b/src/nativeRuntime.ts @@ -0,0 +1,1157 @@ +import { + ClientOrchestrationCommand, + EnvironmentId, + ORCHESTRATION_WS_METHODS, + OrchestrationSubscribeThreadInput, + applyShellStreamEvent, + applyThreadDetailEvent, + makeRpcSessionFactory, + type OrchestrationShellSnapshot, + type OrchestrationShellStreamItem, + type OrchestrationSubscribeShellInput, + type OrchestrationThread, + type OrchestrationThreadShell, + type OrchestrationThreadStreamItem, + type RuntimeClientRpcSessionFactory, +} from "@t3tools/runtime-client"; +import * as Cause from "effect/Cause"; +import * as Effect from "effect/Effect"; +import * as Exit from "effect/Exit"; +import * as Option from "effect/Option"; +import * as Schema from "effect/Schema"; +import * as Scope from "effect/Scope"; +import * as Stream from "effect/Stream"; +import * as Socket from "effect/unstable/socket/Socket"; +import { + AmbiguousDispatchError, + type NativeCreateProjectInput, + type NativeProject, + type NativeRuntime, + type NativeStartThreadInput, + type NativeStartTurnInput, + type NativeThreadObservation, + type NativeThreadSnapshot, +} from "./facade"; + +export type NativeRuntimeAdapterErrorCode = + | "authorization" + | "version_mismatch" + | "command_rejected" + | "transport_unavailable" + | "projection_invalid"; + +export class NativeRuntimeAdapterError extends Error { + readonly code: NativeRuntimeAdapterErrorCode; + + constructor(code: NativeRuntimeAdapterErrorCode) { + super(code); + this.name = "NativeRuntimeAdapterError"; + this.code = code; + } +} + +export interface RuntimeClientSession { + readonly dispatchCommand: ( + command: ClientOrchestrationCommand, + ) => Promise<{ readonly sequence: number }>; + readonly subscribeShell: ( + input: OrchestrationSubscribeShellInput, + ) => AsyncIterable; + readonly subscribeThread: ( + input: OrchestrationSubscribeThreadInput, + ) => AsyncIterable; + readonly close: () => Promise; +} + +export interface RuntimeClientSessionFactory { + readonly connect: (connection: { + readonly environmentId: string; + readonly label: string; + readonly socketUrl: string; + readonly timeoutMs?: number; + }) => Promise; +} + +export interface T3NativeRuntimeOptions { + readonly environmentId: string; + readonly label: string; + /** + * Acquires a one-use authorized WebSocket URL. The adapter passes the value + * directly into a new scoped session and never retains or reports it. + */ + readonly acquireSocketUrl: () => Promise; + readonly sessionFactory?: RuntimeClientSessionFactory; + readonly connectionTimeoutMs?: number; + readonly alignmentTimeoutMs?: number; +} + +interface VersionedDetail { + readonly sequence: number; + readonly thread: OrchestrationThread; +} + +interface VersionedShell { + readonly sequence: number; + readonly snapshot: OrchestrationShellSnapshot; + readonly origin: "seed" | "snapshot" | "event"; +} + +interface ReconciledThreadState { + readonly observation: NativeThreadObservation; + readonly detail: OrchestrationThread; + readonly shellSnapshot: OrchestrationShellSnapshot; +} + +interface ReconcileThreadOptions { + readonly emitAfterSequence?: number; + readonly resumeFromSequence?: number; + readonly seed?: ReconciledThreadState; + readonly alignmentTimeoutMs?: number; +} + +type TaggedNext = + | { + readonly source: "detail"; + readonly result: IteratorResult; + } + | { + readonly source: "shell"; + readonly result: IteratorResult; + } + | { + readonly source: "detail" | "shell"; + readonly error: unknown; + }; + +const DEFAULT_CONNECTION_TIMEOUT_MS = 15_000; +const DEFAULT_ALIGNMENT_TIMEOUT_MS = 5_000; +const MAX_TIMER_TIMEOUT_MS = 2_147_483_647; + +function positiveBound(value: number | undefined, fallback: number): number { + const resolved = value ?? fallback; + if ( + !Number.isSafeInteger(resolved) || + resolved <= 0 || + resolved > MAX_TIMER_TIMEOUT_MS + ) { + throw new NativeRuntimeAdapterError("transport_unavailable"); + } + return resolved; +} + +function remainingMillis(deadline: number): number { + const remaining = deadline - Date.now(); + if (remaining <= 0) { + throw new NativeRuntimeAdapterError("transport_unavailable"); + } + return remaining; +} + +// Initiates teardown synchronously, then detaches it. Callers may rely on the +// attempt having started, but not on cleanup having completed successfully. +function startBestEffortCleanup(cleanup: (() => unknown) | undefined): void { + if (cleanup === undefined) return; + try { + void Promise.resolve(cleanup()).catch(() => undefined); + } catch { + // Cleanup is initiated but deliberately neither awaited nor reported as + // complete; teardown failures must not mask or delay the request result. + } +} + +function withTransportTimeout( + promise: Promise, + timeoutMs: number, + onLateResolve?: (value: T) => void | Promise, + onTimeout?: () => void | Promise, +): Promise { + return new Promise((resolve, reject) => { + let settled = false; + const timer = setTimeout(() => { + settled = true; + startBestEffortCleanup(onTimeout); + reject(new NativeRuntimeAdapterError("transport_unavailable")); + }, timeoutMs); + promise.then( + (value) => { + if (settled) { + if (onLateResolve !== undefined) { + startBestEffortCleanup(() => onLateResolve(value)); + } + return; + } + settled = true; + clearTimeout(timer); + resolve(value); + }, + (error) => { + if (settled) return; + settled = true; + clearTimeout(timer); + reject(error); + }, + ); + }); +} + +function taggedNext( + source: "detail" | "shell", + next: Promise>, +): Promise { + return next.then( + (result) => ({ source, result }) as TaggedNext, + (error) => ({ source, error }), + ); +} + +function adapterError( + error: unknown, + fallback: NativeRuntimeAdapterErrorCode, +): NativeRuntimeAdapterError { + if (error instanceof NativeRuntimeAdapterError) return error; + const tagged = error as { + readonly _tag?: unknown; + readonly reason?: unknown; + }; + if (tagged?._tag === "EnvironmentAuthorizationError") { + return new NativeRuntimeAdapterError("authorization"); + } + if (tagged?._tag === "ConnectionBlockedError") { + return new NativeRuntimeAdapterError( + tagged.reason === "version_mismatch" + ? "version_mismatch" + : "authorization", + ); + } + if (tagged?._tag === "OrchestrationDispatchCommandError") { + return new NativeRuntimeAdapterError("command_rejected"); + } + return new NativeRuntimeAdapterError(fallback); +} + +function dispatchError(error: unknown): Error { + const mapped = adapterError(error, "transport_unavailable"); + if ( + mapped.code === "authorization" || + mapped.code === "version_mismatch" || + mapped.code === "command_rejected" + ) { + return mapped; + } + return new AmbiguousDispatchError(); +} + +async function runEffect( + effect: Effect.Effect, + signal?: AbortSignal, +): Promise { + const exit = await Effect.runPromiseExit( + effect, + signal === undefined ? undefined : { signal }, + ); + if (Exit.isSuccess(exit)) return exit.value; + const error = Cause.findErrorOption(exit.cause); + if (Option.isSome(error)) throw error.value; + throw new NativeRuntimeAdapterError("transport_unavailable"); +} + +function loadRuntimeClientFactory(): Promise { + return runEffect( + makeRpcSessionFactory.pipe( + Effect.provideService( + Socket.WebSocketConstructor, + (url, protocols) => new globalThis.WebSocket(url, protocols), + ), + ), + ); +} + +export function createDefaultSessionFactory( + loadFactory: () => Promise = loadRuntimeClientFactory, + defaultTimeoutMs = DEFAULT_CONNECTION_TIMEOUT_MS, +): RuntimeClientSessionFactory { + let factoryPromise: Promise | undefined; + const getFactory = () => { + if (factoryPromise === undefined) { + const attempt = loadFactory().catch((error) => { + if (factoryPromise === attempt) factoryPromise = undefined; + throw error; + }); + factoryPromise = attempt; + } + return factoryPromise; + }; + + return { + async connect(connection) { + const timeoutMs = positiveBound( + connection.timeoutMs, + positiveBound(defaultTimeoutMs, DEFAULT_CONNECTION_TIMEOUT_MS), + ); + const deadline = Date.now() + timeoutMs; + const factoryAttempt = getFactory(); + let factory: RuntimeClientRpcSessionFactory; + try { + factory = await withTransportTimeout( + factoryAttempt, + remainingMillis(deadline), + ); + } catch (error) { + if (factoryPromise === factoryAttempt) factoryPromise = undefined; + throw error; + } + const scope = await runEffect(Scope.make()); + let closed = false; + const close = async () => { + if (closed) return; + closed = true; + await runEffect(Scope.close(scope, Exit.succeed(undefined))); + }; + try { + const connectAbort = new AbortController(); + const session = await withTransportTimeout( + runEffect( + factory + .connect({ + environmentId: Schema.decodeUnknownSync(EnvironmentId)( + connection.environmentId, + ), + label: connection.label, + socketUrl: connection.socketUrl, + }) + .pipe(Effect.provideService(Scope.Scope, scope)), + connectAbort.signal, + ), + remainingMillis(deadline), + undefined, + () => connectAbort.abort(), + ); + const readyAbort = new AbortController(); + await withTransportTimeout( + runEffect(session.ready, readyAbort.signal), + remainingMillis(deadline), + undefined, + () => readyAbort.abort(), + ); + return { + dispatchCommand: (command) => + runEffect( + session.client[ORCHESTRATION_WS_METHODS.dispatchCommand](command), + ), + subscribeShell: (input) => + Stream.toAsyncIterable( + session.client[ORCHESTRATION_WS_METHODS.subscribeShell](input), + ), + subscribeThread: (input) => + Stream.toAsyncIterable( + session.client[ORCHESTRATION_WS_METHODS.subscribeThread](input), + ), + close, + }; + } catch (error) { + startBestEffortCleanup(close); + throw adapterError(error, "transport_unavailable"); + } + }, + }; +} + +function closeQuietly(session: RuntimeClientSession): void { + startBestEffortCleanup(() => session.close()); +} + +function latestVersionAt( + versions: readonly T[], + sequence: number, +): T | undefined { + for (let index = versions.length - 1; index >= 0; index -= 1) { + const version = versions[index]; + if (version !== undefined && version.sequence <= sequence) return version; + } + return undefined; +} + +function pruneVersionsThrough( + versions: T[], + sequence: number, +): void { + let retainedIndex = 0; + for (let index = 1; index < versions.length; index += 1) { + if (versions[index]!.sequence > sequence) break; + retainedIndex = index; + } + if (retainedIndex > 0) versions.splice(0, retainedIndex); +} + +function sameValue(left: unknown, right: unknown): boolean { + return JSON.stringify(left) === JSON.stringify(right); +} + +function shellCanAdvanceDetail( + detail: OrchestrationThread, + shell: OrchestrationThreadShell, +): boolean { + return ( + detail.id === shell.id && + detail.projectId === shell.projectId && + detail.title === shell.title && + sameValue(detail.modelSelection, shell.modelSelection) && + detail.runtimeMode === shell.runtimeMode && + detail.interactionMode === shell.interactionMode && + detail.branch === shell.branch && + detail.worktreePath === shell.worktreePath && + detail.archivedAt === shell.archivedAt && + detail.settledOverride === shell.settledOverride && + detail.settledAt === shell.settledAt && + sameValue(detail.latestTurn, shell.latestTurn) && + (detail.session === null || + sameValue( + { + status: detail.session.status, + providerName: detail.session.providerName, + providerInstanceId: detail.session.providerInstanceId, + runtimeMode: detail.session.runtimeMode, + activeTurnId: detail.session.activeTurnId, + lastError: detail.session.lastError, + }, + shell.session === null + ? null + : { + status: shell.session.status, + providerName: shell.session.providerName, + providerInstanceId: shell.session.providerInstanceId, + runtimeMode: shell.session.runtimeMode, + activeTurnId: shell.session.activeTurnId, + lastError: shell.session.lastError, + }, + )) + ); +} + +function shellDeltaIsPendingOnly( + previous: OrchestrationThreadShell, + current: OrchestrationThreadShell, +): boolean { + const { + updatedAt: _previousUpdatedAt, + hasPendingApprovals: previousPendingApprovals, + hasPendingUserInput: previousPendingInput, + hasActionableProposedPlan: previousActionablePlan, + ...previousStructural + } = previous; + const { + updatedAt: _currentUpdatedAt, + hasPendingApprovals: currentPendingApprovals, + hasPendingUserInput: currentPendingInput, + hasActionableProposedPlan: currentActionablePlan, + ...currentStructural + } = current; + const pendingChanged = + previousPendingApprovals !== currentPendingApprovals || + previousPendingInput !== currentPendingInput || + previousActionablePlan !== currentActionablePlan; + return pendingChanged && sameValue(previousStructural, currentStructural); +} + +function toNativeSnapshot( + sequence: number, + detail: OrchestrationThread, + shell: OrchestrationThreadShell, +): NativeThreadSnapshot { + const latestTurn = detail.latestTurn; + const userMessage = + latestTurn === null + ? undefined + : [...detail.messages] + .reverse() + .find( + (message) => + message.role === "user" && message.turnId === latestTurn.turnId, + ); + const assistantMessage = + latestTurn?.assistantMessageId == null + ? undefined + : detail.messages.find( + (message) => message.id === latestTurn.assistantMessageId, + ); + const session = detail.session ?? shell.session; + + return { + threadId: detail.id, + projectId: detail.projectId, + snapshotSequence: sequence, + session: { + status: session?.status ?? "unknown", + activeTurnId: session?.activeTurnId ?? null, + }, + latestTurn: + latestTurn === null + ? null + : { + turnId: latestTurn.turnId, + status: latestTurn.state, + ...(userMessage === undefined + ? {} + : { userMessageId: userMessage.id }), + assistantMessage: + assistantMessage === undefined + ? null + : { + content: assistantMessage.text, + streaming: assistantMessage.streaming, + }, + }, + pendingApproval: shell.hasPendingApprovals ? true : null, + pendingInput: shell.hasPendingUserInput ? true : null, + }; +} + +function updateDetail( + versions: VersionedDetail[], + item: OrchestrationThreadStreamItem, +): { readonly synchronized: boolean; readonly deleted: boolean } { + if (item.kind === "synchronized") { + return { synchronized: true, deleted: false }; + } + if (item.kind === "snapshot") { + const latest = versions.at(-1); + if ( + latest === undefined || + item.snapshot.snapshotSequence > latest.sequence + ) { + versions.push({ + sequence: item.snapshot.snapshotSequence, + thread: item.snapshot.thread, + }); + } + return { synchronized: false, deleted: false }; + } + const latest = versions.at(-1); + if (latest === undefined) { + throw new NativeRuntimeAdapterError("projection_invalid"); + } + if (item.event.sequence <= latest.sequence) { + return { synchronized: false, deleted: false }; + } + const reduced = applyThreadDetailEvent(latest.thread, item.event); + if (reduced.kind === "deleted") { + return { synchronized: false, deleted: true }; + } + versions.push({ + sequence: item.event.sequence, + thread: reduced.kind === "updated" ? reduced.thread : latest.thread, + }); + return { synchronized: false, deleted: false }; +} + +function updateShell( + versions: VersionedShell[], + item: OrchestrationShellStreamItem, +): boolean { + if (item.kind === "synchronized") return true; + if (item.kind === "snapshot") { + const latest = versions.at(-1); + if ( + latest === undefined || + item.snapshot.snapshotSequence > latest.sequence + ) { + versions.push({ + sequence: item.snapshot.snapshotSequence, + snapshot: item.snapshot, + origin: "snapshot", + }); + } + return false; + } + const latest = versions.at(-1); + if (latest === undefined) { + throw new NativeRuntimeAdapterError("projection_invalid"); + } + if (item.sequence <= latest.sequence) return false; + const snapshot = applyShellStreamEvent(latest.snapshot, item); + versions.push({ sequence: item.sequence, snapshot, origin: "event" }); + return false; +} + +async function catchUpDetailThrough( + session: RuntimeClientSession, + threadId: string, + afterSequence: number, + versions: VersionedDetail[], + deadline: number, +): Promise<{ readonly deleted: boolean }> { + const iterator = session + .subscribeThread( + Schema.decodeUnknownSync(OrchestrationSubscribeThreadInput)({ + threadId, + afterSequence, + requestCompletionMarker: true, + }), + ) + [Symbol.asyncIterator](); + try { + while (true) { + let next: IteratorResult; + try { + next = await withTransportTimeout( + iterator.next(), + remainingMillis(deadline), + ); + } catch (error) { + throw adapterError(error, "transport_unavailable"); + } + if (next.done) { + throw new NativeRuntimeAdapterError("transport_unavailable"); + } + const update = updateDetail(versions, next.value); + if (update.deleted) return { deleted: true }; + if (update.synchronized) return { deleted: false }; + } + } finally { + startBestEffortCleanup(() => iterator.return?.()); + } +} + +async function catchUpShellThrough( + session: RuntimeClientSession, + afterSequence: number, + versions: VersionedShell[], + deadline: number, +): Promise<{ readonly snapshotSequence: number | undefined }> { + const iterator = session + .subscribeShell({ + afterSequence, + requestCompletionMarker: true, + }) + [Symbol.asyncIterator](); + let snapshotSequence: number | undefined; + try { + while (true) { + let next: IteratorResult; + try { + next = await withTransportTimeout( + iterator.next(), + remainingMillis(deadline), + ); + } catch (error) { + throw adapterError(error, "transport_unavailable"); + } + if (next.done) { + throw new NativeRuntimeAdapterError("transport_unavailable"); + } + if (next.value.kind === "snapshot") { + snapshotSequence = Math.max( + snapshotSequence ?? -1, + next.value.snapshot.snapshotSequence, + ); + } + if (updateShell(versions, next.value)) { + return { snapshotSequence }; + } + } + } finally { + startBestEffortCleanup(() => iterator.return?.()); + } +} + +async function alignInitialVersions( + session: RuntimeClientSession, + threadId: string, + detailVersions: VersionedDetail[], + shellVersions: VersionedShell[], + timeoutMs: number, +): Promise<{ readonly sequence: number; readonly deleted: boolean }> { + const deadline = Date.now() + timeoutMs; + let detailValidatedThrough = detailVersions.at(-1)!.sequence; + let shellValidatedThrough = shellVersions.at(-1)!.sequence; + let targetSequence = Math.max(detailValidatedThrough, shellValidatedThrough); + + while ( + detailValidatedThrough < targetSequence || + shellValidatedThrough < targetSequence + ) { + if (detailValidatedThrough < targetSequence) { + const result = await catchUpDetailThrough( + session, + threadId, + detailValidatedThrough, + detailVersions, + deadline, + ); + if (result.deleted) { + return { sequence: targetSequence, deleted: true }; + } + detailValidatedThrough = targetSequence; + } + if (shellValidatedThrough < targetSequence) { + const result = await catchUpShellThrough( + session, + shellValidatedThrough, + shellVersions, + deadline, + ); + shellValidatedThrough = targetSequence; + if ( + result.snapshotSequence !== undefined && + result.snapshotSequence > targetSequence + ) { + targetSequence = result.snapshotSequence; + shellValidatedThrough = result.snapshotSequence; + } + } + } + + return { sequence: targetSequence, deleted: false }; +} + +async function* reconcileThread( + session: RuntimeClientSession, + threadId: string, + options: ReconcileThreadOptions = {}, +): AsyncIterable { + const resumeFromSequence = options.resumeFromSequence; + const detailIterator = session + .subscribeThread( + Schema.decodeUnknownSync(OrchestrationSubscribeThreadInput)({ + threadId, + ...(resumeFromSequence === undefined + ? {} + : { afterSequence: resumeFromSequence }), + requestCompletionMarker: true, + }), + ) + [Symbol.asyncIterator](); + const shellIterator = session + .subscribeShell({ + ...(resumeFromSequence === undefined + ? {} + : { afterSequence: resumeFromSequence }), + requestCompletionMarker: true, + }) + [Symbol.asyncIterator](); + const detailVersions: VersionedDetail[] = + options.seed === undefined + ? [] + : [ + { + sequence: options.seed.observation.sequence, + thread: options.seed.detail, + }, + ]; + const shellVersions: VersionedShell[] = + options.seed === undefined + ? [] + : [ + { + sequence: options.seed.observation.sequence, + snapshot: { + ...options.seed.shellSnapshot, + snapshotSequence: options.seed.observation.sequence, + }, + origin: "seed", + }, + ]; + let detailSynchronized = false; + let shellSynchronized = false; + let detailDeleted = false; + let detailFailure: NativeRuntimeAdapterError | undefined; + let lastEmitted = options.emitAfterSequence ?? -1; + let initialAlignmentComplete = options.seed !== undefined; + let alignedInitialSequence: number | undefined; + let detailDone = false; + let shellDone = false; + let detailNext: Promise | undefined = taggedNext( + "detail", + detailIterator.next(), + ); + let shellNext: Promise | undefined = taggedNext( + "shell", + shellIterator.next(), + ); + + try { + while (!detailDone || !shellDone) { + const pending = [detailNext, shellNext].filter( + (entry): entry is Promise => entry !== undefined, + ); + if (pending.length === 0) return; + const next = await Promise.race(pending); + if ("error" in next) { + if (next.source === "shell") { + throw adapterError(next.error, "transport_unavailable"); + } + detailFailure = adapterError(next.error, "transport_unavailable"); + detailDone = true; + detailNext = undefined; + } else if (next.source === "detail") { + detailNext = undefined; + if (next.result.done) { + detailDone = true; + } else { + const update = updateDetail(detailVersions, next.result.value); + if (update.synchronized) detailSynchronized = true; + if (update.deleted) detailDeleted = true; + detailNext = taggedNext("detail", detailIterator.next()); + } + } else { + shellNext = undefined; + if (next.result.done) { + shellDone = true; + } else { + if (updateShell(shellVersions, next.result.value)) { + shellSynchronized = true; + } + shellNext = taggedNext("shell", shellIterator.next()); + } + } + + if (shellSynchronized && shellVersions.length > 0) { + const targetExists = shellVersions + .at(-1)! + .snapshot.threads.some((candidate) => candidate.id === threadId); + if (!targetExists) return; + if (detailFailure !== undefined) throw detailFailure; + if (detailDone && !detailSynchronized) { + throw new NativeRuntimeAdapterError("transport_unavailable"); + } + } + + if ( + !detailSynchronized || + !shellSynchronized || + detailDeleted || + detailVersions.length === 0 || + shellVersions.length === 0 + ) { + continue; + } + if (!initialAlignmentComplete) { + const alignment = await alignInitialVersions( + session, + threadId, + detailVersions, + shellVersions, + positiveBound( + options.alignmentTimeoutMs, + DEFAULT_ALIGNMENT_TIMEOUT_MS, + ), + ); + if (alignment.deleted) return; + alignedInitialSequence = alignment.sequence; + initialAlignmentComplete = true; + } + + let detail: VersionedDetail; + let shell: VersionedShell; + let commonSequence: number; + if (alignedInitialSequence !== undefined) { + commonSequence = alignedInitialSequence; + const alignedDetail = latestVersionAt(detailVersions, commonSequence); + const alignedShell = latestVersionAt(shellVersions, commonSequence); + if (alignedDetail === undefined || alignedShell === undefined) { + throw new NativeRuntimeAdapterError("projection_invalid"); + } + detail = alignedDetail; + shell = alignedShell; + } else { + const latestDetail = detailVersions.at(-1)!; + const latestShell = shellVersions.at(-1)!; + const latestShellThread = latestShell.snapshot.threads.find( + (candidate) => candidate.id === threadId, + ); + if (latestShellThread === undefined) return; + if (latestDetail.sequence === latestShell.sequence) { + commonSequence = latestDetail.sequence; + detail = latestDetail; + shell = latestShell; + } else { + const overlappingStateMatches = shellCanAdvanceDetail( + latestDetail.thread, + latestShellThread, + ); + const previousShell = shellVersions + .at(-2) + ?.snapshot.threads.find((candidate) => candidate.id === threadId); + const provenPendingOnlyShellAdvance = + latestShell.sequence > latestDetail.sequence && + latestShell.origin === "event" && + previousShell !== undefined && + shellDeltaIsPendingOnly(previousShell, latestShellThread); + if (!overlappingStateMatches || !provenPendingOnlyShellAdvance) { + continue; + } + // After initial replay alignment, carry detail forward only for a + // shell delta proven to change pending/actionable flags and nothing + // else. Any other skew waits for the matching canonical stream. + commonSequence = latestShell.sequence; + detail = latestDetail; + shell = latestShell; + } + } + if (commonSequence <= lastEmitted) { + alignedInitialSequence = undefined; + continue; + } + const shellThread = shell.snapshot.threads.find( + (candidate) => candidate.id === threadId, + ); + if (shellThread === undefined || detail.thread.id !== threadId) return; + + const snapshot = toNativeSnapshot( + commonSequence, + detail.thread, + shellThread, + ); + const shellSnapshot = { + ...shell.snapshot, + snapshotSequence: commonSequence, + }; + pruneVersionsThrough(detailVersions, commonSequence); + pruneVersionsThrough(shellVersions, commonSequence); + lastEmitted = commonSequence; + alignedInitialSequence = undefined; + yield { + observation: { sequence: commonSequence, snapshot }, + detail: detail.thread, + shellSnapshot, + }; + } + if (detailFailure !== undefined) throw detailFailure; + } finally { + startBestEffortCleanup(() => detailIterator.return?.()); + startBestEffortCleanup(() => shellIterator.return?.()); + } +} + +async function readShellSnapshot( + session: RuntimeClientSession, +): Promise { + const iterator = session + .subscribeShell({ requestCompletionMarker: true }) + [Symbol.asyncIterator](); + let snapshot: OrchestrationShellSnapshot | undefined; + try { + while (true) { + let next: IteratorResult; + try { + next = await iterator.next(); + } catch (error) { + throw adapterError(error, "transport_unavailable"); + } + if (next.done) { + throw new NativeRuntimeAdapterError("transport_unavailable"); + } + if (next.value.kind === "synchronized") { + if (snapshot === undefined) { + throw new NativeRuntimeAdapterError("projection_invalid"); + } + return snapshot; + } + if (next.value.kind === "snapshot") { + snapshot = next.value.snapshot; + } else if (snapshot !== undefined) { + snapshot = applyShellStreamEvent(snapshot, next.value); + } else { + throw new NativeRuntimeAdapterError("projection_invalid"); + } + } + } finally { + startBestEffortCleanup(() => iterator.return?.()); + } +} + +function projectCommand( + input: NativeCreateProjectInput, +): ClientOrchestrationCommand { + return decodeClientCommand({ + type: "project.create", + commandId: input.commandId, + projectId: input.projectId, + title: input.title, + workspaceRoot: input.workspaceRoot, + createWorkspaceRootIfMissing: input.createWorkspaceRootIfMissing, + defaultModelSelection: input.defaultModelSelection, + createdAt: input.createdAt, + }); +} + +function spawnCommand( + input: NativeStartThreadInput, +): ClientOrchestrationCommand { + return decodeClientCommand({ + type: "thread.turn.start", + commandId: input.commandId, + threadId: input.threadId, + message: { + messageId: input.messageId, + role: "user", + text: input.message, + attachments: [], + }, + modelSelection: input.modelSelection, + runtimeMode: input.runtimeMode, + interactionMode: input.interactionMode, + bootstrap: { + createThread: { + projectId: input.projectId, + title: input.title, + modelSelection: input.modelSelection, + runtimeMode: input.runtimeMode, + interactionMode: input.interactionMode, + branch: input.branch, + worktreePath: input.worktreePath, + createdAt: input.createdAt, + }, + }, + createdAt: input.createdAt, + }); +} + +function turnCommand(input: NativeStartTurnInput): ClientOrchestrationCommand { + return decodeClientCommand({ + type: "thread.turn.start", + commandId: input.commandId, + threadId: input.threadId, + message: { + messageId: input.messageId, + role: "user", + text: input.message, + attachments: [], + }, + runtimeMode: input.runtimeMode, + interactionMode: input.interactionMode, + createdAt: input.createdAt, + }); +} + +function decodeClientCommand(input: unknown): ClientOrchestrationCommand { + try { + return Schema.decodeUnknownSync(ClientOrchestrationCommand)(input); + } catch { + throw new NativeRuntimeAdapterError("command_rejected"); + } +} + +export function createT3NativeRuntime( + options: T3NativeRuntimeOptions, +): NativeRuntime { + const connectionTimeoutMs = positiveBound( + options.connectionTimeoutMs, + DEFAULT_CONNECTION_TIMEOUT_MS, + ); + const alignmentTimeoutMs = positiveBound( + options.alignmentTimeoutMs, + DEFAULT_ALIGNMENT_TIMEOUT_MS, + ); + const ownsSessionFactory = options.sessionFactory === undefined; + const sessionFactory = + options.sessionFactory ?? + createDefaultSessionFactory(undefined, connectionTimeoutMs); + + const openSession = async (): Promise => { + const deadline = Date.now() + connectionTimeoutMs; + try { + const socketUrl = await withTransportTimeout( + Promise.resolve().then(options.acquireSocketUrl), + remainingMillis(deadline), + ); + const connectPromise = sessionFactory.connect({ + environmentId: options.environmentId, + label: options.label, + socketUrl, + timeoutMs: remainingMillis(deadline), + }); + if (ownsSessionFactory) return await connectPromise; + return await withTransportTimeout( + connectPromise, + remainingMillis(deadline), + closeQuietly, + ); + } catch (error) { + throw adapterError(error, "transport_unavailable"); + } + }; + + const dispatch = async ( + command: ClientOrchestrationCommand, + ): Promise<{ readonly sequence: number }> => { + const session = await openSession(); + try { + return await session.dispatchCommand(command); + } catch (error) { + throw dispatchError(error); + } finally { + closeQuietly(session); + } + }; + + return { + async listProjects(): Promise { + const session = await openSession(); + try { + const snapshot = await readShellSnapshot(session); + return snapshot.projects.map((entry) => ({ + projectId: entry.id, + workspaceRoot: entry.workspaceRoot, + })); + } finally { + closeQuietly(session); + } + }, + createProject: (input) => dispatch(projectCommand(input)), + startThread: (input) => dispatch(spawnCommand(input)), + startTurn: (input) => dispatch(turnCommand(input)), + async getThread(threadId): Promise { + const session = await openSession(); + try { + for await (const state of reconcileThread(session, threadId, { + alignmentTimeoutMs, + })) { + return state.observation.snapshot; + } + return undefined; + } finally { + closeQuietly(session); + } + }, + async *subscribeThread( + threadId, + input, + ): AsyncIterable { + const session = await openSession(); + try { + if (input.afterSequence === undefined) { + for await (const state of reconcileThread(session, threadId, { + alignmentTimeoutMs, + })) { + yield state.observation; + } + return; + } + + let initial: ReconciledThreadState | undefined; + for await (const state of reconcileThread(session, threadId, { + alignmentTimeoutMs, + })) { + initial = state; + break; + } + if (initial === undefined) return; + if (initial.observation.sequence > input.afterSequence) { + yield initial.observation; + } + const resumeFromSequence = initial.observation.sequence; + for await (const state of reconcileThread(session, threadId, { + seed: initial, + resumeFromSequence, + emitAfterSequence: Math.max(input.afterSequence, resumeFromSequence), + alignmentTimeoutMs, + })) { + yield state.observation; + } + } finally { + closeQuietly(session); + } + }, + }; +} diff --git a/test/config.test.ts b/test/config.test.ts index 33d463c..980b1f6 100644 --- a/test/config.test.ts +++ b/test/config.test.ts @@ -8,6 +8,7 @@ const VALID_CONFIG = { effort: "high", contextWindow: "1m", runtimeMode: "full-access", + interactionMode: "default", } as const; describe("createConfig", () => { @@ -29,6 +30,7 @@ describe("createConfig", () => { "effort", "contextWindow", "runtimeMode", + "interactionMode", ] as const)("rejects a missing %s", (field) => { const input: Record = { ...VALID_CONFIG }; delete input[field]; @@ -37,11 +39,16 @@ describe("createConfig", () => { }); test.each([ - ["baseUrl", "http://127.0.0.1:3774", "baseUrl must be http://127.0.0.1:3773"], + [ + "baseUrl", + "http://127.0.0.1:3774", + "baseUrl must be http://127.0.0.1:3773", + ], ["provider", "claudeDesktop", "provider must be claudeAgent"], ["effort", "medium", "effort must be high"], ["contextWindow", "200k", "contextWindow must be 1m"], ["runtimeMode", "approval-required", "runtimeMode must be full-access"], + ["interactionMode", "plan", "interactionMode must be default"], ] as const)("rejects invalid %s", (field, value, message) => { expect(() => createConfig({ diff --git a/test/facade.contract.test.ts b/test/facade.contract.test.ts new file mode 100644 index 0000000..22c79ad --- /dev/null +++ b/test/facade.contract.test.ts @@ -0,0 +1,193 @@ +import { describe, expect, test } from "bun:test"; +import { createConfig } from "../src/config"; +import { createT3Facade } from "../src/facade"; + +const FACADE_CONFIG = createConfig({ + baseUrl: "http://127.0.0.1:3773", + provider: "claudeAgent", + model: "claude-opus-5", + effort: "high", + contextWindow: "1m", + runtimeMode: "full-access", + interactionMode: "default", +}); + +function makeRuntime( + startThreadInputs: unknown[], + startTurnInputs: unknown[] = [], +) { + return { + async listProjects() { + return [{ projectId: "project-1", workspaceRoot: "/work/app" }]; + }, + async createProject() { + return { sequence: 1 }; + }, + async startThread(input: unknown) { + startThreadInputs.push(input); + return { sequence: 2 }; + }, + async startTurn(input: unknown) { + startTurnInputs.push(input); + return { sequence: 3 }; + }, + async getThread(threadId: string) { + return { + threadId, + projectId: "project-1", + snapshotSequence: 2, + session: { status: "ready", activeTurnId: null }, + latestTurn: + startThreadInputs.length === 0 + ? null + : { + turnId: "turn-1", + status: "completed", + userMessageId: "message-1", + assistantMessage: { + content: "complete", + streaming: false, + }, + }, + pendingApproval: null, + pendingInput: null, + }; + }, + async *subscribeThread() { + return; + }, + }; +} + +describe("canonical pre-adapter contracts", () => { + test("accepts omitted model options and records an evidence count of zero", async () => { + const startThreadInputs: unknown[] = []; + const evidence: unknown[] = []; + const ids = ["thread-1", "command-1", "message-1"]; + const facade = createT3Facade(makeRuntime(startThreadInputs), { + ...FACADE_CONFIG, + id: () => ids.shift()!, + now: () => "2026-07-31T10:00:00.000Z", + evidence: (record) => evidence.push(record), + }); + + await facade.spawn({ + workspaceRoot: "/work/app", + title: "worker", + message: "task", + modelSelection: { + instanceId: "codex", + model: "gpt-5.6-sol", + }, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + }); + + expect(startThreadInputs[0]).toMatchObject({ + modelSelection: { + instanceId: "codex", + model: "gpt-5.6-sol", + }, + }); + expect( + ( + startThreadInputs[0] as { + readonly modelSelection: Record; + } + ).modelSelection, + ).not.toHaveProperty("options"); + expect(evidence[0]).toMatchObject({ + modelSelection: { + instanceId: "codex", + model: "gpt-5.6-sol", + optionCount: 0, + }, + }); + }); + + test("accepts exact boolean model option values without exposing them", async () => { + const startThreadInputs: unknown[] = []; + const evidence: unknown[] = []; + const ids = ["thread-1", "command-1", "message-1"]; + const facade = createT3Facade(makeRuntime(startThreadInputs), { + ...FACADE_CONFIG, + id: () => ids.shift()!, + now: () => "2026-07-31T10:00:00.000Z", + evidence: (record) => evidence.push(record), + }); + + await facade.spawn({ + workspaceRoot: "/work/app", + title: "worker", + message: "task", + modelSelection: { + instanceId: "codex", + model: "gpt-5.6-sol", + options: [{ id: "fastMode", value: true }], + }, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + }); + + expect(startThreadInputs[0]).toMatchObject({ + modelSelection: { + options: [{ id: "fastMode", value: true }], + }, + }); + expect(evidence[0]).toMatchObject({ + modelSelection: { optionCount: 1 }, + }); + expect(JSON.stringify(evidence)).not.toContain("fastMode"); + expect( + ( + evidence[0] as { + readonly modelSelection: Record; + } + ).modelSelection, + ).not.toHaveProperty("options"); + }); + + test("includes the configured modes in follow-up dispatch and evidence", async () => { + const startTurnInputs: unknown[] = []; + const evidence: unknown[] = []; + const ids = ["command-2", "message-2"]; + const facade = createT3Facade(makeRuntime([], startTurnInputs), { + ...FACADE_CONFIG, + id: () => ids.shift()!, + now: () => "2026-07-31T10:05:00.000Z", + evidence: (record) => evidence.push(record), + }); + + await facade.send("thread-1", "follow-up"); + + expect(startTurnInputs).toEqual([ + { + commandId: "command-2", + threadId: "thread-1", + messageId: "message-2", + message: "follow-up", + runtimeMode: "full-access", + interactionMode: "default", + createdAt: "2026-07-31T10:05:00.000Z", + attachments: [], + }, + ]); + expect(evidence).toEqual([ + { + operation: "send", + commandId: "command-2", + threadId: "thread-1", + messageId: "message-2", + runtimeMode: "full-access", + interactionMode: "default", + createdAt: "2026-07-31T10:05:00.000Z", + attachments: 0, + messageBytes: 9, + }, + ]); + }); +}); diff --git a/test/facade.send.test.ts b/test/facade.send.test.ts index 6b5ee7c..154ad12 100644 --- a/test/facade.send.test.ts +++ b/test/facade.send.test.ts @@ -1,6 +1,11 @@ import { describe, expect, test } from "bun:test"; import { AmbiguousDispatchError, createT3Facade } from "../src/facade"; +const DISPATCH_MODES = { + runtimeMode: "full-access", + interactionMode: "default", +} as const; + describe("send", () => { test("reuses the thread with fresh IDs and no bootstrap payload", async () => { const calls: unknown[] = []; @@ -36,6 +41,7 @@ describe("send", () => { }; const ids = ["command-2", "message-2"]; const facade = createT3Facade(runtime, { + ...DISPATCH_MODES, id: () => ids.shift()!, now: () => "2026-07-30T18:05:00.000Z", evidence: (record) => evidence.push(record), @@ -49,6 +55,8 @@ describe("send", () => { threadId: "thread-1", messageId: "message-2", message: "secret follow-up", + runtimeMode: "full-access", + interactionMode: "default", createdAt: "2026-07-30T18:05:00.000Z", attachments: [], }, @@ -67,6 +75,8 @@ describe("send", () => { commandId: "command-2", threadId: "thread-1", messageId: "message-2", + runtimeMode: "full-access", + interactionMode: "default", createdAt: "2026-07-30T18:05:00.000Z", attachments: 0, messageBytes: 16, @@ -108,6 +118,7 @@ describe("send", () => { }; const ids = ["command-multibyte", "message-multibyte"]; const facade = createT3Facade(runtime, { + ...DISPATCH_MODES, id: () => ids.shift()!, now: () => "2026-07-31T00:00:00.000Z", evidence: (record) => evidence.push(record), @@ -122,6 +133,8 @@ describe("send", () => { commandId: "command-multibyte", threadId: "thread-1", messageId: "message-multibyte", + runtimeMode: "full-access", + interactionMode: "default", createdAt: "2026-07-31T00:00:00.000Z", attachments: 0, messageBytes: 6, @@ -166,7 +179,7 @@ describe("send", () => { return; }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); const result = facade.send("thread-requested", "follow-up"); const error = await result.catch((reason: unknown) => reason); @@ -203,7 +216,7 @@ describe("send", () => { return; }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); const result = facade.send("thread-requested", "follow-up"); @@ -268,6 +281,7 @@ describe("send", () => { }; const ids = ["command-2", "message-2"]; const facade = createT3Facade(runtime, { + ...DISPATCH_MODES, id: () => ids.shift()!, now: () => "2026-07-30T18:05:00.000Z", }); @@ -323,6 +337,7 @@ describe("send", () => { }; const ids = ["command-stable", "message-stable"]; const facade = createT3Facade(runtime, { + ...DISPATCH_MODES, id: () => ids.shift()!, now: () => "2026-07-31T00:00:00.000Z", }); @@ -336,6 +351,8 @@ describe("send", () => { threadId: "thread-1", messageId: "message-stable", message: "follow-up", + runtimeMode: "full-access", + interactionMode: "default", createdAt: "2026-07-31T00:00:00.000Z", attachments: [], }); @@ -396,6 +413,7 @@ describe("send", () => { }; const ids = ["command-stable", "message-stable"]; const facade = createT3Facade(runtime, { + ...DISPATCH_MODES, id: () => ids.shift()!, now: () => "2026-07-31T00:00:00.000Z", }); @@ -450,6 +468,7 @@ describe("send", () => { }; const ids = ["command-stable", "message-stable"]; const facade = createT3Facade(runtime, { + ...DISPATCH_MODES, id: () => ids.shift()!, now: () => "2026-07-31T00:00:00.000Z", }); @@ -506,6 +525,7 @@ describe("send", () => { }; const ids = ["command-stable", "message-stable"]; const facade = createT3Facade(runtime, { + ...DISPATCH_MODES, id: () => ids.shift()!, now: () => "2026-07-31T00:00:00.000Z", }); @@ -613,6 +633,7 @@ describe("send", () => { }; const ids = ["command-stable", "message-stable"]; const facade = createT3Facade(runtime, { + ...DISPATCH_MODES, id: () => ids.shift()!, now: () => "2026-07-31T00:00:00.000Z", }); @@ -662,6 +683,7 @@ describe("send", () => { }; const ids = ["command-stable", "message-stable"]; const facade = createT3Facade(runtime, { + ...DISPATCH_MODES, id: () => ids.shift()!, now: () => "2026-07-31T00:00:00.000Z", }); @@ -713,6 +735,7 @@ describe("send", () => { }, }; const facade = createT3Facade(runtime, { + ...DISPATCH_MODES, id: () => ids.shift()!, evidence: (record) => evidence.push(record), }); @@ -765,7 +788,7 @@ describe("send", () => { return; }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); await expect(facade.send("thread-1", "duplicate")).rejects.toMatchObject({ code: "turn_error", diff --git a/test/facade.spawn.test.ts b/test/facade.spawn.test.ts index cb39b3e..0028083 100644 --- a/test/facade.spawn.test.ts +++ b/test/facade.spawn.test.ts @@ -1,6 +1,11 @@ import { describe, expect, test } from "bun:test"; import { AmbiguousDispatchError, createT3Facade } from "../src/facade"; +const DISPATCH_MODES = { + runtimeMode: "full-access", + interactionMode: "default", +} as const; + describe("spawn", () => { test("discovers the project by exact workspace root and starts one explicit atomic turn", async () => { const calls: Array<{ @@ -49,6 +54,7 @@ describe("spawn", () => { }; const ids = ["thread-1", "command-1", "message-1"]; const facade = createT3Facade(runtime, { + ...DISPATCH_MODES, id: () => ids.shift()!, now: () => "2026-07-30T18:00:00.000Z", evidence: (record) => evidence.push(record), @@ -162,6 +168,7 @@ describe("spawn", () => { }; const ids = ["thread-multibyte", "command-multibyte", "message-multibyte"]; const facade = createT3Facade(runtime, { + ...DISPATCH_MODES, id: () => ids.shift()!, now: () => "2026-07-31T00:00:00.000Z", evidence: (record) => evidence.push(record), @@ -246,6 +253,7 @@ describe("spawn", () => { }; const ids = ["thread-1", "command-1", "message-1"]; const facade = createT3Facade(runtime, { + ...DISPATCH_MODES, id: () => ids.shift()!, now: () => "2026-07-31T00:00:00.000Z", evidence: (record) => evidence.push(record), @@ -326,6 +334,7 @@ describe("spawn", () => { "message-1", ]; const facade = createT3Facade(runtime, { + ...DISPATCH_MODES, id: () => ids.shift()!, now: () => "2026-07-30T18:00:00.000Z", }); @@ -422,6 +431,7 @@ describe("spawn", () => { "message-stable", ]; const facade = createT3Facade(runtime, { + ...DISPATCH_MODES, id: () => ids.shift()!, now: () => "2026-07-31T00:00:00.000Z", }); @@ -501,6 +511,7 @@ describe("spawn", () => { }; const ids = ["thread-stable", "command-stable", "message-stable"]; const facade = createT3Facade(runtime, { + ...DISPATCH_MODES, id: () => ids.shift()!, now: () => "2026-07-30T18:00:00.000Z", }); @@ -566,6 +577,7 @@ describe("spawn", () => { }; const ids = ["thread-stable", "command-stable", "message-stable"]; const facade = createT3Facade(runtime, { + ...DISPATCH_MODES, id: () => ids.shift()!, now: () => "2026-07-31T00:00:00.000Z", }); @@ -619,6 +631,7 @@ describe("spawn", () => { }; const ids = ["thread-stable", "command-stable", "message-stable"]; const facade = createT3Facade(runtime, { + ...DISPATCH_MODES, id: () => ids.shift()!, now: () => "2026-07-31T00:00:00.000Z", }); @@ -707,6 +720,7 @@ describe("spawn", () => { }; const ids = ["thread-stable", "command-stable", "message-stable"]; const facade = createT3Facade(runtime, { + ...DISPATCH_MODES, id: () => ids.shift()!, now: () => "2026-07-31T00:00:00.000Z", }); @@ -800,6 +814,7 @@ describe("spawn", () => { }; const ids = ["thread-stable", "command-stable", "message-stable"]; const facade = createT3Facade(runtime, { + ...DISPATCH_MODES, id: () => ids.shift()!, now: () => "2026-07-31T00:00:00.000Z", }); diff --git a/test/facade.wait.test.ts b/test/facade.wait.test.ts index e34d911..b591cc0 100644 --- a/test/facade.wait.test.ts +++ b/test/facade.wait.test.ts @@ -1,6 +1,11 @@ import { describe, expect, test } from "bun:test"; import { FacadeError, createT3Facade, type AgentEvent } from "../src/facade"; +const DISPATCH_MODES = { + runtimeMode: "full-access", + interactionMode: "default", +} as const; + async function collect( iterable: AsyncIterable, ): Promise { @@ -96,7 +101,7 @@ describe("wait", () => { throw new Error("wait consumed past completion"); }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); const events = []; for await (const event of facade.wait("thread-1", { @@ -184,7 +189,7 @@ describe("wait", () => { yield { sequence: 70, snapshot: invalid }; }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); const iterator = facade .wait("thread-1", { kind: "terminal", @@ -279,7 +284,7 @@ describe("wait", () => { }; }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); const iterator = facade .wait("thread-1", { kind: "terminal", @@ -337,7 +342,7 @@ describe("wait", () => { return; }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); const result = collect( facade.wait("thread-requested", { @@ -377,7 +382,7 @@ describe("wait", () => { return; }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); const result = collect( facade.wait("thread-requested", { @@ -438,7 +443,7 @@ describe("wait", () => { yield { sequence: 26, snapshot: divergent }; }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); const result = collect( facade.wait("thread-requested", { @@ -489,7 +494,7 @@ describe("wait", () => { return; }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); const result = collect( facade.wait("thread-1", { @@ -543,7 +548,7 @@ describe("wait", () => { throw new Error("missing assistant must fail before subscription"); }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); const result = collect( facade.wait("thread-1", { @@ -602,7 +607,7 @@ describe("wait", () => { throw new Error("pending state must stop before subscription"); }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); const events = await collect( facade.wait("thread-1", { @@ -652,7 +657,7 @@ describe("wait", () => { throw new Error("pending approval must stop before subscription"); }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); const events = await collect( facade.wait("thread-1", { @@ -705,7 +710,7 @@ describe("wait", () => { throw new Error("failed state must stop before subscription"); }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); const events = await collect( facade.wait("thread-1", { @@ -756,7 +761,7 @@ describe("wait", () => { yield { sequence: 33, snapshot: stopped }; }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); const events = await collect( facade.wait("thread-1", { @@ -818,7 +823,7 @@ describe("wait", () => { yield { sequence: 41, snapshot: completed }; }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); const result = collect( facade.wait("thread-1", { @@ -870,7 +875,7 @@ describe("wait", () => { throw new Error("late terminal lookup must not subscribe"); }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); const startedAt = performance.now(); const result = collect( @@ -922,7 +927,7 @@ describe("wait", () => { yield { sequence: 43, snapshot: running }; }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); const startedAt = performance.now(); const result = collect( @@ -976,7 +981,7 @@ describe("wait", () => { throw new Error(`transport rejected ${credential}`); }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); const result = collect( facade.wait("thread-1", { @@ -1030,7 +1035,7 @@ describe("wait", () => { return; }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); const result = collect( facade.wait("thread-1", { @@ -1081,7 +1086,7 @@ describe("wait", () => { return; }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); const result = collect( facade.wait("thread-1", { @@ -1138,7 +1143,7 @@ describe("wait", () => { ); }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); const result = collect( facade.wait("thread-1", { @@ -1191,7 +1196,7 @@ describe("wait", () => { throw new Error("interrupted state must stop before subscription"); }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); const events = await collect( facade.wait("thread-1", { @@ -1239,7 +1244,7 @@ describe("wait", () => { throw new Error("stopped state must stop before subscription"); }, }; - const facade = createT3Facade(runtime); + const facade = createT3Facade(runtime, DISPATCH_MODES); const events = await collect( facade.wait("thread-1", { diff --git a/test/native-runtime-adapter.test.ts b/test/native-runtime-adapter.test.ts new file mode 100644 index 0000000..bb4d251 --- /dev/null +++ b/test/native-runtime-adapter.test.ts @@ -0,0 +1,1523 @@ +import { describe, expect, test } from "bun:test"; +import type { + ClientOrchestrationCommand, + OrchestrationEvent, + OrchestrationProjectShell, + OrchestrationShellSnapshot, + OrchestrationShellStreamItem, + OrchestrationThread, + OrchestrationThreadShell, + OrchestrationThreadStreamItem, + RuntimeClientRpcSessionFactory, +} from "@t3tools/runtime-client"; +import * as Effect from "effect/Effect"; +import { + createDefaultSessionFactory, + createT3NativeRuntime, + type RuntimeClientSession, + type RuntimeClientSessionFactory, +} from "../src/nativeRuntime"; + +const NOW = "2026-07-31T00:00:00.000Z"; +const MODEL = { + instanceId: "codex", + model: "gpt-5.6-sol", + options: [{ id: "reasoningEffort", value: "high" }], +} as const; +const BOOLEAN_MODEL = { + instanceId: "codex", + model: "gpt-5.6-sol", + options: [{ id: "fastMode", value: true }], +} as const; + +function project( + id = "project-1", + workspaceRoot = "/repo", +): OrchestrationProjectShell { + return { + id, + title: "Project", + workspaceRoot, + defaultModelSelection: MODEL, + scripts: [], + createdAt: NOW, + updatedAt: NOW, + } as unknown as OrchestrationProjectShell; +} + +function thread( + sequence: number, + overrides: Partial = {}, +): OrchestrationThread { + return { + id: "thread-1", + projectId: "project-1", + title: "Worker", + modelSelection: MODEL, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + latestTurn: { + turnId: "turn-1", + state: "running", + requestedAt: NOW, + startedAt: NOW, + completedAt: null, + assistantMessageId: null, + }, + createdAt: NOW, + updatedAt: `${NOW.slice(0, -5)}${String(sequence).padStart(3, "0")}Z`, + archivedAt: null, + settledOverride: null, + settledAt: null, + deletedAt: null, + messages: [ + { + id: "message-user", + role: "user", + text: "work", + turnId: "turn-1", + streaming: false, + createdAt: NOW, + updatedAt: NOW, + }, + ], + proposedPlans: [], + activities: [], + checkpoints: [], + session: { + threadId: "thread-1", + status: "running", + providerName: "codex", + providerInstanceId: "codex", + runtimeMode: "full-access", + activeTurnId: "turn-1", + lastError: null, + updatedAt: NOW, + }, + ...overrides, + } as unknown as OrchestrationThread; +} + +function shellThread( + sequence: number, + options: { + readonly pendingApproval?: boolean; + readonly pendingInput?: boolean; + readonly status?: "running" | "ready"; + } = {}, +): OrchestrationThreadShell { + const status = options.status ?? "running"; + return { + id: "thread-1", + projectId: "project-1", + title: "Worker", + modelSelection: MODEL, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + latestTurn: { + turnId: "turn-1", + state: status === "ready" ? "completed" : "running", + requestedAt: NOW, + startedAt: NOW, + completedAt: status === "ready" ? NOW : null, + assistantMessageId: status === "ready" ? "message-assistant" : null, + }, + createdAt: NOW, + updatedAt: `${NOW.slice(0, -5)}${String(sequence).padStart(3, "0")}Z`, + archivedAt: null, + settledOverride: null, + settledAt: null, + session: { + threadId: "thread-1", + status, + providerName: "codex", + providerInstanceId: "codex", + runtimeMode: "full-access", + activeTurnId: status === "running" ? "turn-1" : null, + lastError: null, + updatedAt: NOW, + }, + latestUserMessageAt: NOW, + hasPendingApprovals: options.pendingApproval ?? false, + hasPendingUserInput: options.pendingInput ?? false, + hasActionableProposedPlan: false, + } as unknown as OrchestrationThreadShell; +} + +function shellSnapshot( + sequence: number, + threads: readonly OrchestrationThreadShell[] = [], +): OrchestrationShellSnapshot { + return { + snapshotSequence: sequence, + projects: [project()], + threads: [...threads], + updatedAt: NOW, + }; +} + +function stream( + items: readonly T[], + onReturn?: () => void, +): AsyncIterable { + return { + [Symbol.asyncIterator]() { + let index = 0; + return { + async next() { + if (index >= items.length) return { done: true, value: undefined }; + return { done: false, value: items[index++]! }; + }, + async return() { + onReturn?.(); + return { done: true, value: undefined }; + }, + }; + }, + }; +} + +function scriptedStream( + script: ReadonlyArray<{ + readonly delayMs: number; + readonly item: T; + }>, + onReturn?: () => void, +): AsyncIterable { + return { + [Symbol.asyncIterator]() { + let index = 0; + return { + async next() { + const entry = script[index++]; + if (entry === undefined) return { done: true, value: undefined }; + if (entry.delayMs > 0) { + await new Promise((resolve) => setTimeout(resolve, entry.delayMs)); + } + return { done: false, value: entry.item }; + }, + async return() { + onReturn?.(); + return { done: true, value: undefined }; + }, + }; + }, + }; +} + +type PromiseOutcome = + | { readonly kind: "resolved"; readonly value: T } + | { readonly kind: "rejected"; readonly error: unknown } + | { readonly kind: "pending" }; + +async function outcomeWithin( + promise: Promise, + timeoutMs = 100, +): Promise> { + return Promise.race([ + promise.then( + (value): PromiseOutcome => ({ kind: "resolved", value }), + (error): PromiseOutcome => ({ kind: "rejected", error }), + ), + new Promise>((resolve) => { + setTimeout(() => resolve({ kind: "pending" }), timeoutMs); + }), + ]); +} + +describe("T3 native runtime adapter", () => { + test("maps project and turn operations onto canonical RPC commands", async () => { + const commands: unknown[] = []; + const connections: Array<{ + readonly environmentId: string; + readonly label: string; + readonly socketUrl: string; + }> = []; + let closes = 0; + const sessionFactory: RuntimeClientSessionFactory = { + async connect(connection) { + connections.push(connection); + return { + async dispatchCommand(command) { + commands.push(command); + return { sequence: commands.length }; + }, + subscribeShell: () => + stream([ + { kind: "snapshot", snapshot: shellSnapshot(4) }, + { kind: "synchronized" }, + ]), + subscribeThread: () => stream([]), + async close() { + closes += 1; + }, + }; + }, + }; + let socketAcquisitions = 0; + const runtime = createT3NativeRuntime({ + environmentId: "environment-1", + label: "MacBook Pro", + acquireSocketUrl: async () => + `ws://127.0.0.1/socket?token=${++socketAcquisitions}`, + sessionFactory, + }); + + expect(await runtime.listProjects()).toEqual([ + { projectId: "project-1", workspaceRoot: "/repo" }, + ]); + expect( + await runtime.createProject({ + commandId: "command-project", + projectId: "project-new", + title: "New Project", + workspaceRoot: "/new", + createWorkspaceRootIfMissing: false, + defaultModelSelection: MODEL, + createdAt: NOW, + }), + ).toEqual({ sequence: 1 }); + expect( + await runtime.startThread({ + commandId: "command-spawn", + projectId: "project-1", + threadId: "thread-1", + messageId: "message-user", + title: "Worker", + message: "Do the work", + modelSelection: BOOLEAN_MODEL, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + createdAt: NOW, + attachments: [], + }), + ).toEqual({ sequence: 2 }); + expect( + await runtime.startTurn({ + commandId: "command-send", + threadId: "thread-1", + messageId: "message-follow-up", + message: "Continue", + runtimeMode: "full-access", + interactionMode: "default", + createdAt: NOW, + attachments: [], + }), + ).toEqual({ sequence: 3 }); + + expect(commands).toEqual([ + { + type: "project.create", + commandId: "command-project", + projectId: "project-new", + title: "New Project", + workspaceRoot: "/new", + createWorkspaceRootIfMissing: false, + defaultModelSelection: MODEL, + createdAt: NOW, + }, + { + type: "thread.turn.start", + commandId: "command-spawn", + threadId: "thread-1", + message: { + messageId: "message-user", + role: "user", + text: "Do the work", + attachments: [], + }, + modelSelection: BOOLEAN_MODEL, + runtimeMode: "full-access", + interactionMode: "default", + bootstrap: { + createThread: { + projectId: "project-1", + title: "Worker", + modelSelection: BOOLEAN_MODEL, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + createdAt: NOW, + }, + }, + createdAt: NOW, + }, + { + type: "thread.turn.start", + commandId: "command-send", + threadId: "thread-1", + message: { + messageId: "message-follow-up", + role: "user", + text: "Continue", + attachments: [], + }, + runtimeMode: "full-access", + interactionMode: "default", + createdAt: NOW, + }, + ]); + expect(connections).toHaveLength(4); + expect(socketAcquisitions).toBe(4); + expect(closes).toBe(4); + }); + + test("sanitizes unsent command projection failures without opening a session", async () => { + const secret = "projection-super-secret"; + const invalidCreatedAt = "not-an-iso-date"; + let connects = 0; + const runtime = createT3NativeRuntime({ + environmentId: "environment-1", + label: "MacBook Pro", + acquireSocketUrl: async () => "ws://127.0.0.1/ephemeral", + sessionFactory: { + async connect() { + connects += 1; + throw new Error("must not connect"); + }, + }, + }); + const invalidCalls = [ + () => + runtime.createProject({ + commandId: "", + projectId: "project-new", + title: secret, + workspaceRoot: "/new", + createWorkspaceRootIfMissing: false, + defaultModelSelection: MODEL, + createdAt: invalidCreatedAt, + }), + () => + runtime.startThread({ + commandId: "", + projectId: "project-1", + threadId: "thread-1", + messageId: "message-user", + title: "Worker", + message: secret, + modelSelection: MODEL, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + createdAt: invalidCreatedAt, + attachments: [], + }), + () => + runtime.startTurn({ + commandId: "", + threadId: "thread-1", + messageId: "message-follow-up", + message: secret, + runtimeMode: "full-access", + interactionMode: "default", + createdAt: invalidCreatedAt, + attachments: [], + }), + ]; + + for (const call of invalidCalls) { + const outcome = await outcomeWithin(Promise.resolve().then(call)); + expect(outcome).toMatchObject({ + kind: "rejected", + error: { + name: "NativeRuntimeAdapterError", + code: "command_rejected", + }, + }); + if (outcome.kind !== "rejected") { + throw new Error("expected rejected command projection"); + } + expect(String(outcome.error)).toBe( + "NativeRuntimeAdapterError: command_rejected", + ); + expect((outcome.error as Error).message).toBe("command_rejected"); + expect(String(outcome.error)).not.toContain(secret); + expect(String(outcome.error)).not.toContain("ParseError"); + expect((outcome.error as Error).name).not.toBe("AmbiguousDispatchError"); + } + expect(connects).toBe(0); + }); + + test("reconciles a pending-only shell snapshot ahead of a synchronized detail snapshot", async () => { + const subscriptionInputs: unknown[] = []; + let closes = 0; + const session: RuntimeClientSession = { + async dispatchCommand() { + return { sequence: 1 }; + }, + subscribeShell(input) { + subscriptionInputs.push(["shell", input]); + return stream([ + { + kind: "snapshot", + snapshot: shellSnapshot(11, [ + shellThread(11, { pendingApproval: true }), + ]), + }, + { kind: "synchronized" }, + ]); + }, + subscribeThread(input) { + subscriptionInputs.push(["thread", input]); + return stream([ + { + kind: "snapshot", + snapshot: { snapshotSequence: 10, thread: thread(10) }, + }, + { kind: "synchronized" }, + ]); + }, + async close() { + closes += 1; + }, + }; + const runtime = createT3NativeRuntime({ + environmentId: "environment-1", + label: "MacBook Pro", + acquireSocketUrl: async () => "ws://127.0.0.1/ephemeral", + sessionFactory: { connect: async () => session }, + }); + + const snapshot = await runtime.getThread("thread-1"); + + expect(snapshot).toMatchObject({ + threadId: "thread-1", + projectId: "project-1", + snapshotSequence: 11, + session: { status: "running", activeTurnId: "turn-1" }, + latestTurn: { + turnId: "turn-1", + status: "running", + userMessageId: "message-user", + }, + pendingApproval: true, + pendingInput: null, + }); + expect(subscriptionInputs).toHaveLength(3); + expect(subscriptionInputs).toContainEqual([ + "shell", + { requestCompletionMarker: true }, + ]); + expect(subscriptionInputs).toContainEqual([ + "thread", + { threadId: "thread-1", requestCompletionMarker: true }, + ]); + expect(subscriptionInputs).toContainEqual([ + "thread", + { + threadId: "thread-1", + afterSequence: 10, + requestCompletionMarker: true, + }, + ]); + expect(closes).toBe(1); + }); + + test("reconciles a synchronized detail snapshot ahead of a compatible shell snapshot", async () => { + const runtime = createT3NativeRuntime({ + environmentId: "environment-1", + label: "MacBook Pro", + acquireSocketUrl: async () => "ws://127.0.0.1/ephemeral", + sessionFactory: { + async connect() { + return { + async dispatchCommand() { + return { sequence: 1 }; + }, + subscribeShell: () => + stream([ + { + kind: "snapshot", + snapshot: shellSnapshot(10, [shellThread(10)]), + }, + { kind: "synchronized" }, + ]), + subscribeThread: () => + stream([ + { + kind: "snapshot", + snapshot: { snapshotSequence: 11, thread: thread(11) }, + }, + { kind: "synchronized" }, + ]), + async close() {}, + }; + }, + }, + }); + + expect(await runtime.getThread("thread-1")).toMatchObject({ + threadId: "thread-1", + snapshotSequence: 11, + session: { status: "running", activeTurnId: "turn-1" }, + pendingApproval: null, + pendingInput: null, + }); + }); + + test("returns undefined for a synchronized missing shell thread even if detail lookup rejects", async () => { + let closes = 0; + const missingDetail: AsyncIterable = { + [Symbol.asyncIterator]() { + return { + async next(): Promise> { + throw { _tag: "OrchestrationGetSnapshotError" }; + }, + }; + }, + }; + const session: RuntimeClientSession = { + async dispatchCommand() { + return { sequence: 1 }; + }, + subscribeShell() { + return scriptedStream([ + { + delayMs: 5, + item: { kind: "snapshot", snapshot: shellSnapshot(10) }, + }, + { delayMs: 0, item: { kind: "synchronized" } }, + ]); + }, + subscribeThread() { + return missingDetail; + }, + async close() { + closes += 1; + }, + }; + const runtime = createT3NativeRuntime({ + environmentId: "environment-1", + label: "MacBook Pro", + acquireSocketUrl: async () => "ws://127.0.0.1/ephemeral", + sessionFactory: { connect: async () => session }, + }); + + expect(await runtime.getThread("missing-thread")).toBeUndefined(); + expect(closes).toBe(1); + }); + + test("does not swallow a detail failure when shell absence was never synchronized", async () => { + const failingDetail: AsyncIterable = { + [Symbol.asyncIterator]() { + return { + async next(): Promise> { + throw { _tag: "OrchestrationGetSnapshotError" }; + }, + }; + }, + }; + const runtime = createT3NativeRuntime({ + environmentId: "environment-1", + label: "MacBook Pro", + acquireSocketUrl: async () => "ws://127.0.0.1/ephemeral", + sessionFactory: { + async connect() { + return { + async dispatchCommand() { + return { sequence: 1 }; + }, + subscribeShell: () => stream([]), + subscribeThread: () => failingDetail, + async close() {}, + }; + }, + }, + }); + + await expect(runtime.getThread("thread-1")).rejects.toMatchObject({ + name: "NativeRuntimeAdapterError", + code: "transport_unavailable", + }); + }); + + test("waits for a same-sequence detail update when a detail-relevant shell update arrives first", async () => { + const readyEvent = { + ...eventBase(11, "thread.session-set"), + payload: { + threadId: "thread-1", + session: { + ...thread(11).session!, + status: "ready", + activeTurnId: null, + updatedAt: NOW, + }, + }, + } as OrchestrationEvent; + const session: RuntimeClientSession = { + async dispatchCommand() { + return { sequence: 1 }; + }, + subscribeShell() { + return scriptedStream([ + { + delayMs: 0, + item: { + kind: "snapshot", + snapshot: shellSnapshot(10, [shellThread(10)]), + }, + }, + { delayMs: 0, item: { kind: "synchronized" } }, + { + delayMs: 1, + item: { + kind: "thread-upserted", + sequence: 11, + thread: shellThread(11, { status: "ready" }), + }, + }, + ]); + }, + subscribeThread() { + return scriptedStream([ + { + delayMs: 0, + item: { + kind: "snapshot", + snapshot: { snapshotSequence: 10, thread: thread(10) }, + }, + }, + { delayMs: 0, item: { kind: "synchronized" } }, + { delayMs: 15, item: { kind: "event", event: readyEvent } }, + ]); + }, + async close() {}, + }; + const runtime = createT3NativeRuntime({ + environmentId: "environment-1", + label: "MacBook Pro", + acquireSocketUrl: async () => "ws://127.0.0.1/ephemeral", + sessionFactory: { connect: async () => session }, + }); + + const observations = []; + for await (const observation of runtime.subscribeThread("thread-1", {})) { + observations.push(observation); + if (observation.sequence === 11) break; + } + + expect(observations.map(({ sequence }) => sequence)).toEqual([10, 11]); + expect(observations[1]?.snapshot.session).toEqual({ + status: "ready", + activeTurnId: null, + }); + }); + + test("does not let a shell-first same-sequence upsert suppress a detail-only message update", async () => { + const assistantEvent = { + ...eventBase(11, "thread.message-sent"), + payload: { + threadId: "thread-1", + messageId: "message-assistant", + role: "assistant", + text: "AB", + turnId: "turn-1", + streaming: true, + createdAt: NOW, + updatedAt: NOW, + }, + } as OrchestrationEvent; + const session: RuntimeClientSession = { + async dispatchCommand() { + return { sequence: 1 }; + }, + subscribeShell: () => + scriptedStream([ + { + delayMs: 0, + item: { + kind: "snapshot", + snapshot: shellSnapshot(10, [shellThread(10)]), + }, + }, + { delayMs: 0, item: { kind: "synchronized" } }, + { + delayMs: 1, + item: { + kind: "thread-upserted", + sequence: 11, + thread: shellThread(11), + }, + }, + ]), + subscribeThread: () => + scriptedStream([ + { + delayMs: 0, + item: { + kind: "snapshot", + snapshot: { snapshotSequence: 10, thread: thread(10) }, + }, + }, + { delayMs: 0, item: { kind: "synchronized" } }, + { delayMs: 15, item: { kind: "event", event: assistantEvent } }, + ]), + async close() {}, + }; + const runtime = createT3NativeRuntime({ + environmentId: "environment-1", + label: "MacBook Pro", + acquireSocketUrl: async () => "ws://127.0.0.1/ephemeral", + sessionFactory: { connect: async () => session }, + }); + + const observations = []; + for await (const observation of runtime.subscribeThread("thread-1", {})) { + observations.push(observation); + if (observation.sequence === 11) break; + } + + expect(observations.map(({ sequence }) => sequence)).toEqual([10, 11]); + expect( + observations[1]?.snapshot.latestTurn?.assistantMessage, + ).toMatchObject({ + content: "AB", + streaming: true, + }); + }); + + test("replays an initially lagging detail stream before emitting a newer shell snapshot sequence", async () => { + const assistantEvent = { + ...eventBase(11, "thread.message-sent"), + payload: { + threadId: "thread-1", + messageId: "message-assistant", + role: "assistant", + text: "AB", + turnId: "turn-1", + streaming: true, + createdAt: NOW, + updatedAt: NOW, + }, + } as OrchestrationEvent; + const detailInputs: unknown[] = []; + const session: RuntimeClientSession = { + async dispatchCommand() { + return { sequence: 1 }; + }, + subscribeShell: () => + stream([ + { + kind: "snapshot", + snapshot: shellSnapshot(11, [shellThread(11)]), + }, + { kind: "synchronized" }, + ]), + subscribeThread(input) { + detailInputs.push(input); + if (input.afterSequence === 10) { + return stream([ + { kind: "event", event: assistantEvent }, + { kind: "synchronized" }, + ]); + } + return scriptedStream([ + { + delayMs: 0, + item: { + kind: "snapshot", + snapshot: { snapshotSequence: 10, thread: thread(10) }, + }, + }, + { delayMs: 0, item: { kind: "synchronized" } }, + { delayMs: 25, item: { kind: "event", event: assistantEvent } }, + ]); + }, + async close() {}, + }; + const runtime = createT3NativeRuntime({ + environmentId: "environment-1", + label: "MacBook Pro", + acquireSocketUrl: async () => "ws://127.0.0.1/ephemeral", + sessionFactory: { connect: async () => session }, + }); + + const snapshot = await runtime.getThread("thread-1"); + + expect(snapshot?.snapshotSequence).toBe(11); + expect(snapshot?.latestTurn?.assistantMessage).toMatchObject({ + content: "AB", + streaming: true, + }); + expect(detailInputs).toContainEqual({ + threadId: "thread-1", + afterSequence: 10, + requestCompletionMarker: true, + }); + }); + + test("times out a hanging initial alignment replay without awaiting a hanging iterator return", async () => { + let replayReturns = 0; + let detailReturns = 0; + let shellReturns = 0; + const hangingReplay: AsyncIterable = { + [Symbol.asyncIterator]() { + return { + next: () => + new Promise>( + () => undefined, + ), + return() { + replayReturns += 1; + return new Promise>( + () => undefined, + ); + }, + }; + }, + }; + const session: RuntimeClientSession = { + async dispatchCommand() { + return { sequence: 1 }; + }, + subscribeShell: () => + stream( + [ + { + kind: "snapshot", + snapshot: shellSnapshot(11, [shellThread(11)]), + }, + { kind: "synchronized" }, + ], + () => { + shellReturns += 1; + }, + ), + subscribeThread(input) { + if (input.afterSequence === 10) return hangingReplay; + return stream( + [ + { + kind: "snapshot", + snapshot: { snapshotSequence: 10, thread: thread(10) }, + }, + { kind: "synchronized" }, + ], + () => { + detailReturns += 1; + }, + ); + }, + async close() {}, + }; + const runtime = createT3NativeRuntime({ + environmentId: "environment-1", + label: "MacBook Pro", + alignmentTimeoutMs: 10, + acquireSocketUrl: async () => "ws://127.0.0.1/ephemeral", + sessionFactory: { connect: async () => session }, + }); + + const outcome = await outcomeWithin(runtime.getThread("thread-1")); + + expect(outcome).toMatchObject({ + kind: "rejected", + error: { + name: "NativeRuntimeAdapterError", + code: "transport_unavailable", + }, + }); + expect(replayReturns).toBe(1); + expect(detailReturns).toBe(1); + expect(shellReturns).toBe(1); + }); + + test("returns missing when aligned shell replay removes the target thread", async () => { + const session: RuntimeClientSession = { + async dispatchCommand() { + return { sequence: 1 }; + }, + subscribeShell(input) { + if (input.afterSequence === 10) { + return stream([ + { + kind: "snapshot", + snapshot: shellSnapshot(12, []), + }, + { kind: "synchronized" }, + ]); + } + return stream([ + { + kind: "snapshot", + snapshot: shellSnapshot(10, [shellThread(10)]), + }, + { kind: "synchronized" }, + ]); + }, + subscribeThread(input) { + if (input.afterSequence === 11) { + return stream([{ kind: "synchronized" }]); + } + return stream([ + { + kind: "snapshot", + snapshot: { snapshotSequence: 11, thread: thread(11) }, + }, + { kind: "synchronized" }, + ]); + }, + async close() {}, + }; + const runtime = createT3NativeRuntime({ + environmentId: "environment-1", + label: "MacBook Pro", + acquireSocketUrl: async () => "ws://127.0.0.1/ephemeral", + sessionFactory: { connect: async () => session }, + }); + + expect(await runtime.getThread("thread-1")).toBeUndefined(); + }); + + test("terminates fallback alignment after one shell catch-up and one detail catch-up", async () => { + let shellSubscriptions = 0; + let detailSubscriptions = 0; + let iteratorReturns = 0; + const session: RuntimeClientSession = { + async dispatchCommand() { + return { sequence: 1 }; + }, + subscribeShell(input) { + shellSubscriptions += 1; + if (input.afterSequence === 10) { + return stream( + [ + { + kind: "snapshot", + snapshot: shellSnapshot(20, [shellThread(20)]), + }, + { kind: "synchronized" }, + ], + () => { + iteratorReturns += 1; + }, + ); + } + return stream( + [ + { + kind: "snapshot", + snapshot: shellSnapshot(10, [shellThread(10)]), + }, + { kind: "synchronized" }, + ], + () => { + iteratorReturns += 1; + }, + ); + }, + subscribeThread(input) { + detailSubscriptions += 1; + if (input.afterSequence === 11) { + return stream([{ kind: "synchronized" }], () => { + iteratorReturns += 1; + }); + } + return stream( + [ + { + kind: "snapshot", + snapshot: { snapshotSequence: 11, thread: thread(11) }, + }, + { kind: "synchronized" }, + ], + () => { + iteratorReturns += 1; + }, + ); + }, + async close() {}, + }; + const runtime = createT3NativeRuntime({ + environmentId: "environment-1", + label: "MacBook Pro", + acquireSocketUrl: async () => "ws://127.0.0.1/ephemeral", + sessionFactory: { connect: async () => session }, + }); + + expect(await runtime.getThread("thread-1")).toMatchObject({ + threadId: "thread-1", + snapshotSequence: 20, + }); + expect(shellSubscriptions).toBe(2); + expect(detailSubscriptions).toBe(2); + expect(iteratorReturns).toBe(4); + }); + + test("reconciles interleaved reducers at a common monotonic watermark and preserves pending precedence", async () => { + let shellReturns = 0; + let detailReturns = 0; + let closes = 0; + const assistantEvent = { + ...eventBase(12, "thread.message-sent"), + payload: { + threadId: "thread-1", + messageId: "message-assistant", + role: "assistant", + text: "Done", + turnId: "turn-1", + streaming: false, + createdAt: NOW, + updatedAt: NOW, + }, + } as OrchestrationEvent; + const readyEvent = { + ...eventBase(13, "thread.session-set"), + payload: { + threadId: "thread-1", + session: { + threadId: "thread-1", + status: "ready", + providerName: "codex", + providerInstanceId: "codex", + runtimeMode: "full-access", + activeTurnId: null, + lastError: null, + updatedAt: NOW, + }, + }, + } as OrchestrationEvent; + const session: RuntimeClientSession = { + async dispatchCommand() { + return { sequence: 1 }; + }, + subscribeShell() { + return scriptedStream( + [ + { + delayMs: 5, + item: { + kind: "snapshot", + snapshot: shellSnapshot(10, [shellThread(10)]), + }, + }, + { delayMs: 0, item: { kind: "synchronized" } }, + { + delayMs: 5, + item: { + kind: "thread-upserted", + sequence: 11, + thread: shellThread(11, { pendingApproval: true }), + }, + }, + { + delayMs: 5, + item: { + kind: "thread-upserted", + sequence: 11, + thread: shellThread(11, { pendingApproval: true }), + }, + }, + { + delayMs: 40, + item: { + kind: "thread-upserted", + sequence: 13, + thread: shellThread(13, { status: "ready" }), + }, + }, + ], + () => { + shellReturns += 1; + }, + ); + }, + subscribeThread() { + return scriptedStream( + [ + { + delayMs: 0, + item: { + kind: "snapshot", + snapshot: { snapshotSequence: 10, thread: thread(10) }, + }, + }, + { delayMs: 0, item: { kind: "synchronized" } }, + { delayMs: 20, item: { kind: "event", event: assistantEvent } }, + { delayMs: 20, item: { kind: "event", event: readyEvent } }, + ], + () => { + detailReturns += 1; + }, + ); + }, + async close() { + closes += 1; + }, + }; + const runtime = createT3NativeRuntime({ + environmentId: "environment-1", + label: "MacBook Pro", + acquireSocketUrl: async () => "ws://127.0.0.1/ephemeral", + sessionFactory: { connect: async () => session }, + }); + + const observations = []; + for await (const observation of runtime.subscribeThread("thread-1", {})) { + observations.push(observation); + if (observation.sequence === 13) break; + } + + expect(observations.map(({ sequence }) => sequence)).toEqual([10, 11, 13]); + expect(observations[1]?.snapshot).toMatchObject({ + snapshotSequence: 11, + pendingApproval: true, + pendingInput: null, + latestTurn: { status: "running", assistantMessage: null }, + }); + expect(observations[2]?.snapshot).toMatchObject({ + snapshotSequence: 13, + pendingApproval: null, + pendingInput: null, + session: { status: "ready", activeTurnId: null }, + latestTurn: { + status: "completed", + assistantMessage: { content: "Done", streaming: false }, + }, + }); + expect(shellReturns).toBe(1); + expect(detailReturns).toBe(1); + expect(closes).toBe(1); + }); + + test("requests synchronized subscription markers, filters the resume boundary, and sanitizes connection failures", async () => { + const inputs: unknown[] = []; + let closes = 0; + const healthy: RuntimeClientSession = { + async dispatchCommand() { + return { sequence: 1 }; + }, + subscribeShell(input) { + inputs.push(["shell", input]); + if (input.afterSequence === 8) { + return stream([ + { + kind: "thread-upserted", + sequence: 9, + thread: shellThread(9), + }, + { kind: "synchronized" }, + ]); + } + return stream([ + { kind: "snapshot", snapshot: shellSnapshot(8, [shellThread(8)]) }, + { kind: "synchronized" }, + ]); + }, + subscribeThread(input) { + inputs.push(["thread", input]); + if (input.afterSequence === 8) { + return stream([ + { + kind: "event", + event: { + ...eventBase(9, "thread.session-set"), + payload: { + threadId: "thread-1", + session: { + ...thread(9).session!, + updatedAt: NOW, + }, + }, + } as OrchestrationEvent, + }, + { kind: "synchronized" }, + ]); + } + return stream([ + { + kind: "snapshot", + snapshot: { snapshotSequence: 8, thread: thread(8) }, + }, + { kind: "synchronized" }, + ]); + }, + async close() { + closes += 1; + }, + }; + const runtime = createT3NativeRuntime({ + environmentId: "environment-1", + label: "MacBook Pro", + acquireSocketUrl: async () => "ws://127.0.0.1/ephemeral", + sessionFactory: { connect: async () => healthy }, + }); + + const iterator = runtime + .subscribeThread("thread-1", { afterSequence: 8 }) + [Symbol.asyncIterator](); + expect((await iterator.next()).value?.sequence).toBe(9); + await iterator.return?.(); + + expect(inputs).toHaveLength(4); + expect(inputs).toContainEqual(["shell", { requestCompletionMarker: true }]); + expect(inputs).toContainEqual([ + "thread", + { threadId: "thread-1", requestCompletionMarker: true }, + ]); + expect(inputs).toContainEqual([ + "shell", + { afterSequence: 8, requestCompletionMarker: true }, + ]); + expect(inputs).toContainEqual([ + "thread", + { + threadId: "thread-1", + afterSequence: 8, + requestCompletionMarker: true, + }, + ]); + expect(closes).toBe(1); + + const secret = "ws://127.0.0.1/socket?authorization=super-secret"; + const broken = createT3NativeRuntime({ + environmentId: "environment-1", + label: "MacBook Pro", + acquireSocketUrl: async () => secret, + sessionFactory: { + async connect() { + throw new Error(`could not dial ${secret}`); + }, + }, + }); + let failure: unknown; + try { + await broken.listProjects(); + } catch (error) { + failure = error; + } + expect(failure).toMatchObject({ + name: "NativeRuntimeAdapterError", + code: "transport_unavailable", + }); + expect(String(failure)).not.toContain("super-secret"); + expect(JSON.stringify(failure)).not.toContain("super-secret"); + }); + + test("bounds socket acquisition and injected session connection hangs", async () => { + let socketFactoryConnects = 0; + const socketHang = createT3NativeRuntime({ + environmentId: "environment-1", + label: "MacBook Pro", + connectionTimeoutMs: 10, + acquireSocketUrl: () => new Promise(() => undefined), + sessionFactory: { + async connect() { + socketFactoryConnects += 1; + throw new Error("must not connect"); + }, + }, + }); + const socketOutcome = await outcomeWithin(socketHang.listProjects()); + expect(socketOutcome).toMatchObject({ + kind: "rejected", + error: { + name: "NativeRuntimeAdapterError", + code: "transport_unavailable", + }, + }); + expect(socketFactoryConnects).toBe(0); + + const connectHang = createT3NativeRuntime({ + environmentId: "environment-1", + label: "MacBook Pro", + connectionTimeoutMs: 10, + acquireSocketUrl: async () => "ws://127.0.0.1/ephemeral", + sessionFactory: { + connect: () => new Promise(() => undefined), + }, + }); + const connectOutcome = await outcomeWithin(connectHang.listProjects()); + expect(connectOutcome).toMatchObject({ + kind: "rejected", + error: { + name: "NativeRuntimeAdapterError", + code: "transport_unavailable", + }, + }); + }); + + test("closes a late injected session exactly once after the caller times out", async () => { + let resolveConnect: ((session: RuntimeClientSession) => void) | undefined; + let closes = 0; + const lateSession: RuntimeClientSession = { + async dispatchCommand() { + return { sequence: 1 }; + }, + subscribeShell: () => stream([]), + subscribeThread: () => stream([]), + async close() { + closes += 1; + }, + }; + const runtime = createT3NativeRuntime({ + environmentId: "environment-1", + label: "MacBook Pro", + connectionTimeoutMs: 10, + acquireSocketUrl: async () => "ws://127.0.0.1/ephemeral", + sessionFactory: { + connect: () => + new Promise((resolve) => { + resolveConnect = resolve; + }), + }, + }); + + const outcome = await outcomeWithin(runtime.listProjects()); + expect(outcome).toMatchObject({ + kind: "rejected", + error: { + name: "NativeRuntimeAdapterError", + code: "transport_unavailable", + }, + }); + + resolveConnect?.(lateSession); + await new Promise((resolve) => setTimeout(resolve, 10)); + expect(closes).toBe(1); + }); + + test("rejects timeout values larger than the platform timer limit", () => { + const options = { + environmentId: "environment-1", + label: "MacBook Pro", + acquireSocketUrl: async () => "ws://127.0.0.1/ephemeral", + }; + + for (const oversized of [ + { connectionTimeoutMs: 2_147_483_648 }, + { alignmentTimeoutMs: 2_147_483_648 }, + ]) { + expect(() => createT3NativeRuntime({ ...options, ...oversized })).toThrow( + expect.objectContaining({ + name: "NativeRuntimeAdapterError", + code: "transport_unavailable", + }), + ); + } + }); + + test("bounds hanging default factory effects without awaiting their cleanup finalizers", async () => { + const connection = { + environmentId: "environment-1", + label: "MacBook Pro", + socketUrl: "ws://127.0.0.1/ephemeral", + }; + let connectScopeCloses = 0; + const connectHang = createDefaultSessionFactory( + async () => + ({ + connect: () => + Effect.gen(function* () { + yield* Effect.addFinalizer(() => + Effect.sync(() => { + connectScopeCloses += 1; + }).pipe(Effect.andThen(Effect.never)), + ); + return yield* Effect.never; + }), + }) as unknown as RuntimeClientRpcSessionFactory, + 10, + ); + const connectOutcome = await outcomeWithin(connectHang.connect(connection)); + expect(connectOutcome).toMatchObject({ + kind: "rejected", + error: { + name: "NativeRuntimeAdapterError", + code: "transport_unavailable", + }, + }); + await new Promise((resolve) => setTimeout(resolve, 10)); + expect(connectScopeCloses).toBe(1); + + let readyInterrupts = 0; + let readyScopeCloses = 0; + const readyHang = createDefaultSessionFactory( + async () => + ({ + connect: () => + Effect.gen(function* () { + yield* Effect.addFinalizer(() => + Effect.sync(() => { + readyScopeCloses += 1; + }).pipe(Effect.andThen(Effect.never)), + ); + return { + ready: Effect.never.pipe( + Effect.ensuring( + Effect.sync(() => { + readyInterrupts += 1; + }).pipe(Effect.andThen(Effect.never)), + ), + ), + client: {}, + }; + }), + }) as unknown as RuntimeClientRpcSessionFactory, + 10, + ); + const readyOutcome = await outcomeWithin(readyHang.connect(connection)); + expect(readyOutcome).toMatchObject({ + kind: "rejected", + error: { + name: "NativeRuntimeAdapterError", + code: "transport_unavailable", + }, + }); + await new Promise((resolve) => setTimeout(resolve, 10)); + expect(readyInterrupts).toBe(1); + expect(readyScopeCloses).toBe(1); + }); + + test("retries runtime-client factory initialization after a rejected attempt", async () => { + let loads = 0; + const sessionFactory = createDefaultSessionFactory(async () => { + loads += 1; + if (loads === 1) throw new Error("temporary initialization failure"); + return { + connect: () => + Effect.succeed({ + ready: Effect.succeed(undefined), + client: {}, + }), + } as unknown as RuntimeClientRpcSessionFactory; + }); + const connection = { + environmentId: "environment-1", + label: "MacBook Pro", + socketUrl: "ws://127.0.0.1/ephemeral", + }; + + await expect(sessionFactory.connect(connection)).rejects.toThrow( + "temporary initialization failure", + ); + const session = await sessionFactory.connect(connection); + await session.close(); + + expect(loads).toBe(2); + }); +}); + +function eventBase(sequence: number, type: string) { + return { + sequence, + eventId: `event-${sequence}`, + aggregateKind: "thread", + aggregateId: "thread-1", + occurredAt: NOW, + commandId: null, + causationEventId: null, + correlationId: null, + metadata: {}, + type, + }; +}