/** * ACPX turn adapters. Modern runtimes can expose startTurn directly; legacy * runtimes that only stream runTurn events are adapted to the newer contract. */ import { toErrorObject } from "openclaw/plugin-sdk/error-runtime"; import { createDeferred } from "openclaw/plugin-sdk/extension-shared"; import type { AcpRuntime, AcpRuntimeEvent, AcpRuntimeTurn, AcpRuntimeTurnInput, AcpRuntimeTurnResult, } from "../runtime-api.js"; function isCancellationStopReason(stopReason: string | undefined): boolean { return stopReason === "cancel" || stopReason === "cancelled" || stopReason === "manual-cancel"; } class LegacyRunTurnEventQueue { private readonly items: AcpRuntimeEvent[] = []; private readonly waits: Array<{ resolve: (value: AcpRuntimeEvent | null) => void; reject: (error: unknown) => void; }> = []; private closed = false; private error: unknown; push(item: AcpRuntimeEvent): void { if (this.closed) { return; } const waiter = this.waits.shift(); if (waiter) { waiter.resolve(item); return; } this.items.push(item); } clear(): void { this.items.length = 0; } close(): void { if (this.closed) { return; } this.closed = true; for (const waiter of this.waits.splice(0)) { waiter.resolve(null); } } fail(error: unknown): void { if (this.closed) { return; } this.error = error; this.closed = true; for (const waiter of this.waits.splice(0)) { waiter.reject(error); } } private async next(): Promise { const item = this.items.shift(); if (item) { return item; } if (this.error) { throw toErrorObject(this.error, "Non-Error thrown"); } if (this.closed) { return null; } return await new Promise((resolve, reject) => { this.waits.push({ resolve, reject }); }); } async *iterate(): AsyncIterable { for (;;) { const item = await this.next(); if (!item) { return; } yield item; } } } function legacyRunTurnAsStartTurn(runtime: AcpRuntime, input: AcpRuntimeTurnInput): AcpRuntimeTurn { const result = createDeferred(); result.promise.catch(() => {}); const queue = new LegacyRunTurnEventQueue(); let resultSettled = false; const settleResult = (next: AcpRuntimeTurnResult) => { if (resultSettled) { return; } resultSettled = true; result.resolve(next); }; void (async () => { try { for await (const event of runtime.runTurn(input)) { if (event.type === "done") { // Legacy runTurn events omit result.status but preserve stopReason, so infer // cancellation here instead of silently converting it to success. settleResult({ status: event.status ?? (isCancellationStopReason(event.stopReason) ? "cancelled" : "completed"), ...(event.stopReason ? { stopReason: event.stopReason } : {}), }); continue; } if (event.type === "error") { settleResult({ status: "failed", error: { message: event.message, ...(event.code ? { code: event.code } : {}), ...(event.detailCode ? { detailCode: event.detailCode } : {}), ...(event.retryable === undefined ? {} : { retryable: event.retryable }), }, }); continue; } queue.push(event); } settleResult({ status: "failed", error: { code: "ACP_TURN_FAILED", message: "ACP turn ended without a terminal done event.", }, }); } catch (error) { result.reject(error); queue.fail(error); return; } queue.close(); })(); return { requestId: input.requestId, events: queue.iterate(), result: result.promise, async cancel(inputArgs) { await runtime.cancel({ handle: input.handle, reason: inputArgs?.reason }); }, async closeStream() { queue.clear(); queue.close(); }, }; } /** Start an ACP turn, adapting legacy runTurn-only runtimes when needed. */ function startRuntimeTurn(runtime: AcpRuntime, input: AcpRuntimeTurnInput): AcpRuntimeTurn { return runtime.startTurn?.(input) ?? legacyRunTurnAsStartTurn(runtime, input); } /** Start an ACP turn through a lazy runtime resolver. */ export function lazyStartRuntimeTurn( resolveRuntime: () => Promise, input: AcpRuntimeTurnInput, ): AcpRuntimeTurn { const turnPromise = resolveRuntime().then((runtime) => startRuntimeTurn(runtime, input)); return { requestId: input.requestId, events: { async *[Symbol.asyncIterator]() { yield* (await turnPromise).events; }, }, result: turnPromise.then((turn) => turn.result), cancel(inputArgs) { return turnPromise.then((turn) => turn.cancel(inputArgs)); }, closeStream(inputArgs) { return turnPromise.then((turn) => turn.closeStream(inputArgs)); }, }; }