diff --git a/apps/server/scripts/acp-mock-agent.ts b/apps/server/scripts/acp-mock-agent.ts index d60688ce09fb..1565de017707 100644 --- a/apps/server/scripts/acp-mock-agent.ts +++ b/apps/server/scripts/acp-mock-agent.ts @@ -14,6 +14,8 @@ import type * as AcpSchema from "effect-acp/schema"; const requestLogPath = process.env.T3_ACP_REQUEST_LOG_PATH; const exitLogPath = process.env.T3_ACP_EXIT_LOG_PATH; const emitToolCalls = process.env.T3_ACP_EMIT_TOOL_CALLS === "1"; +const emitToolCallsOnPermissionPhrase = + process.env.T3_ACP_EMIT_TOOL_CALLS_ON_PERMISSION_PHRASE === "1"; const emitInterleavedAssistantToolCalls = process.env.T3_ACP_EMIT_INTERLEAVED_ASSISTANT_TOOL_CALLS === "1"; const emitGenericToolPlaceholders = process.env.T3_ACP_EMIT_GENERIC_TOOL_PLACEHOLDERS === "1"; @@ -22,6 +24,7 @@ const emitXAiAskUserQuestion = process.env.T3_ACP_EMIT_XAI_ASK_USER_QUESTION === const emitCreatePlan = process.env.T3_ACP_EMIT_CREATE_PLAN === "1"; const emitUpdateTodos = process.env.T3_ACP_EMIT_UPDATE_TODOS === "1"; const failPrompt = process.env.T3_ACP_FAIL_PROMPT === "1"; +const hangPromptWithXAiComplete = process.env.T3_ACP_HANG_PROMPT_WITH_XAI_COMPLETE === "1"; const failSetConfigOption = process.env.T3_ACP_FAIL_SET_CONFIG_OPTION === "1"; const exitOnSetConfigOption = process.env.T3_ACP_EXIT_ON_SET_CONFIG_OPTION === "1"; const promptResponseText = process.env.T3_ACP_PROMPT_RESPONSE_TEXT; @@ -330,6 +333,9 @@ const program = Effect.gen(function* () { yield* agent.handlePrompt((request) => Effect.gen(function* () { const requestedSessionId = String(request.sessionId ?? sessionId); + const promptText = request.prompt + .map((block) => (block.type === "text" ? block.text : "")) + .join("\n"); if (failPrompt) { return yield* AcpError.AcpRequestError.invalidParams("Mock prompt failure", { @@ -338,6 +344,14 @@ const program = Effect.gen(function* () { }); } + if (hangPromptWithXAiComplete) { + yield* agent.client.extNotification("_x.ai/session/prompt_complete", { + sessionId: requestedSessionId, + stopReason: "end_turn", + }); + return yield* Effect.never; + } + if (emitInterleavedAssistantToolCalls) { const toolCallId = "tool-call-1"; @@ -388,7 +402,10 @@ const program = Effect.gen(function* () { return { stopReason: "end_turn" }; } - if (emitToolCalls) { + if ( + emitToolCalls && + (!emitToolCallsOnPermissionPhrase || promptText.includes("trigger permission")) + ) { const toolCallId = "tool-call-1"; yield* agent.client.sessionUpdate({ diff --git a/apps/server/src/access/ServerExposure.test.ts b/apps/server/src/access/ServerExposure.test.ts index 2787440ce9a5..5fb144f03c89 100644 --- a/apps/server/src/access/ServerExposure.test.ts +++ b/apps/server/src/access/ServerExposure.test.ts @@ -188,8 +188,7 @@ describe("ServerExposure", () => { ); it.effect("reports Tailscale Serve disabled when startup configure fails", () => { - const commands: Array<{ readonly command: string; readonly args: ReadonlyArray }> = - []; + const commands: Array<{ readonly command: string; readonly args: ReadonlyArray }> = []; return Effect.scoped( Effect.gen(function* () { @@ -232,14 +231,10 @@ describe("ServerExposure", () => { const exposure = yield* ServerExposure; const endpoints = yield* exposure.getAdvertisedEndpoints; - expect(endpoints.map((endpoint) => endpoint.httpBaseUrl)).toContain( - "http://127.0.0.1:3773/", - ); + expect(endpoints.map((endpoint) => endpoint.httpBaseUrl)).toContain("http://127.0.0.1:3773/"); expect(endpoints.some((endpoint) => endpoint.label === "This machine")).toBe(true); expect(endpoints.some((endpoint) => endpoint.label === "Tailscale HTTPS")).toBe(true); - expect( - endpoints.find((endpoint) => endpoint.label === "Tailscale HTTPS"), - ).toMatchObject({ + expect(endpoints.find((endpoint) => endpoint.label === "Tailscale HTTPS")).toMatchObject({ httpBaseUrl: "https://desktop.tail.ts.net/", status: "unavailable", }); @@ -300,4 +295,4 @@ describe("ServerExposure", () => { ), ); }); -}); \ No newline at end of file +}); diff --git a/apps/server/src/provider/Layers/GrokBuildAdapter.test.ts b/apps/server/src/provider/Layers/GrokBuildAdapter.test.ts index 9a497adbf636..6a0c0cbc5937 100644 --- a/apps/server/src/provider/Layers/GrokBuildAdapter.test.ts +++ b/apps/server/src/provider/Layers/GrokBuildAdapter.test.ts @@ -388,7 +388,10 @@ const grokPermissionAdapterTestLayer = it.layer( GrokBuildAdapter, Effect.gen(function* () { const wrapperPath = yield* Effect.promise(() => - makeMockAgentWrapper({ T3_ACP_EMIT_TOOL_CALLS: "1" }), + makeMockAgentWrapper({ + T3_ACP_EMIT_TOOL_CALLS: "1", + T3_ACP_EMIT_TOOL_CALLS_ON_PERMISSION_PHRASE: "1", + }), ); const settings = decodeGrokBuildSettings({ enabled: true, @@ -433,8 +436,14 @@ grokPermissionAdapterTestLayer("GrokBuildAdapter permissions", (it) => { Stream.runCollect, Effect.forkChild, ); + const turnCompletedFiber = yield* Stream.take(adapter.streamEvents, 20).pipe( + Stream.filter((event) => event.type === "turn.completed"), + Stream.take(1), + Stream.runCollect, + Effect.forkChild, + ); - yield* adapter.sendTurn({ + const turn = yield* adapter.sendTurn({ threadId, input: "trigger permission", attachments: [], @@ -443,15 +452,238 @@ grokPermissionAdapterTestLayer("GrokBuildAdapter permissions", (it) => { const openedEvents = Array.from(yield* Fiber.join(requestOpenedFiber)); const opened = openedEvents[0]; assert.isDefined(opened); - if (opened?.type === "request.opened" && opened.requestId) { - yield* adapter.respondToRequest( - threadId, - ApprovalRequestId.make(String(opened.requestId)), - "accept", - ); + + yield* adapter.interruptTurn(threadId, turn.turnId); + + const completedEvents = Array.from(yield* Fiber.join(turnCompletedFiber)); + const completed = completedEvents[0]; + assert.equal(completed?.type, "turn.completed"); + if (completed?.type === "turn.completed") { + assert.equal(completed.turnId, turn.turnId); + assert.equal(completed.payload.state, "cancelled"); + assert.equal(completed.payload.stopReason, "cancelled"); } - yield* adapter.interruptTurn(threadId); + const session = (yield* adapter.listSessions()).find( + (candidate) => candidate.threadId === threadId, + ); + assert.equal(session?.status, "ready"); + assert.isUndefined(session?.activeTurnId); + + yield* adapter.stopSession(threadId); + }), + ); + + it.effect("completes a later turn after interrupting a blocked permission turn", () => + Effect.gen(function* () { + const adapter = yield* GrokBuildAdapter; + const threadId = ThreadId.make("grok-permission-interrupt-next-turn-thread"); + yield* adapter.startSession({ + threadId, + provider: ProviderDriverKind.make("grok-build"), + cwd: process.cwd(), + runtimeMode: "approval-required", + modelSelection: { + instanceId: ProviderInstanceId.make("grok-build-test"), + model: "default", + }, + }); + + const requestOpenedFiber = yield* Stream.take(adapter.streamEvents, 20).pipe( + Stream.filter((event) => event.type === "request.opened"), + Stream.take(1), + Stream.runDrain, + Effect.forkChild, + ); + const completedFiber = yield* Stream.take(adapter.streamEvents, 40).pipe( + Stream.filter((event) => event.type === "turn.completed"), + Stream.take(2), + Stream.runCollect, + Effect.forkChild, + ); + + const interruptedTurn = yield* adapter.sendTurn({ + threadId, + input: "trigger permission", + attachments: [], + }); + + yield* Fiber.join(requestOpenedFiber); + yield* adapter.interruptTurn(threadId, interruptedTurn.turnId); + + const nextTurn = yield* adapter.sendTurn({ + threadId, + input: "plain follow-up after interrupt", + attachments: [], + }); + + const completedEvents = Array.from(yield* Fiber.join(completedFiber)); + assert.equal(completedEvents.length, 2); + const [cancelled, completed] = completedEvents; + assert.equal(cancelled?.type, "turn.completed"); + assert.equal(completed?.type, "turn.completed"); + if (cancelled?.type === "turn.completed") { + assert.equal(cancelled.turnId, interruptedTurn.turnId); + assert.equal(cancelled.payload.state, "cancelled"); + } + if (completed?.type === "turn.completed") { + assert.equal(completed.turnId, nextTurn.turnId); + assert.equal(completed.payload.state, "completed"); + assert.equal(completed.payload.stopReason, "end_turn"); + } + + const session = (yield* adapter.listSessions()).find( + (candidate) => candidate.threadId === threadId, + ); + assert.equal(session?.status, "ready"); + assert.isUndefined(session?.activeTurnId); + + yield* adapter.stopSession(threadId); + }), + ); +}); + +const grokPromptCompleteAdapterTestLayer = it.layer( + Layer.effect( + GrokBuildAdapter, + Effect.gen(function* () { + const wrapperPath = yield* Effect.promise(() => + makeMockAgentWrapper({ T3_ACP_HANG_PROMPT_WITH_XAI_COMPLETE: "1" }), + ); + const settings = decodeGrokBuildSettings({ + enabled: true, + command: wrapperPath, + args: [], + envJson: "{}", + customModels: [], + }); + return yield* makeGrokBuildAdapter(settings, { + instanceId: ProviderInstanceId.make("grok-build-test"), + }); + }), + ).pipe( + Layer.provideMerge( + ServerConfig.layerTest(process.cwd(), { + prefix: "t3code-grok-prompt-complete-test-", + }), + ), + Layer.provideMerge(NodeServices.layer), + ), +); + +grokPromptCompleteAdapterTestLayer("GrokBuildAdapter xAI prompt completion", (it) => { + it.effect("settles a turn from prompt_complete when session/prompt remains pending", () => + Effect.gen(function* () { + const adapter = yield* GrokBuildAdapter; + const threadId = ThreadId.make("grok-prompt-complete-thread"); + const turnCompletedFiber = yield* Stream.take(adapter.streamEvents, 20).pipe( + Stream.filter((event) => event.type === "turn.completed"), + Stream.take(1), + Stream.runCollect, + Effect.forkChild, + ); + + yield* adapter.startSession({ + threadId, + provider: ProviderDriverKind.make("grok-build"), + cwd: process.cwd(), + runtimeMode: "full-access", + modelSelection: { + instanceId: ProviderInstanceId.make("grok-build-test"), + model: "default", + }, + }); + + const turn = yield* adapter.sendTurn({ + threadId, + input: "complete through xai notification", + attachments: [], + }); + + const completedEvents = Array.from(yield* Fiber.join(turnCompletedFiber)); + const completed = completedEvents[0]; + assert.equal(completed?.type, "turn.completed"); + if (completed?.type === "turn.completed") { + assert.equal(completed.turnId, turn.turnId); + assert.equal(completed.payload.state, "completed"); + assert.equal(completed.payload.stopReason, "end_turn"); + } + + const session = (yield* adapter.listSessions()).find( + (candidate) => candidate.threadId === threadId, + ); + assert.equal(session?.status, "ready"); + assert.isUndefined(session?.activeTurnId); + + yield* adapter.stopSession(threadId); + }), + ); + + it.effect("completes a later turn after prompt_complete settlement", () => + Effect.gen(function* () { + const adapter = yield* GrokBuildAdapter; + const threadId = ThreadId.make("grok-prompt-complete-next-turn-thread"); + const firstCompletedFiber = yield* Stream.take(adapter.streamEvents, 40).pipe( + Stream.filter((event) => event.type === "turn.completed"), + Stream.take(1), + Stream.runCollect, + Effect.forkChild, + ); + + yield* adapter.startSession({ + threadId, + provider: ProviderDriverKind.make("grok-build"), + cwd: process.cwd(), + runtimeMode: "full-access", + modelSelection: { + instanceId: ProviderInstanceId.make("grok-build-test"), + model: "default", + }, + }); + + const firstTurn = yield* adapter.sendTurn({ + threadId, + input: "complete through xai notification", + attachments: [], + }); + + const firstCompletedEvents = Array.from(yield* Fiber.join(firstCompletedFiber)); + const firstCompleted = firstCompletedEvents[0]; + assert.equal(firstCompleted?.type, "turn.completed"); + if (firstCompleted?.type === "turn.completed") { + assert.equal(firstCompleted.turnId, firstTurn.turnId); + assert.equal(firstCompleted.payload.state, "completed"); + assert.equal(firstCompleted.payload.stopReason, "end_turn"); + } + + const secondCompletedFiber = yield* Stream.take(adapter.streamEvents, 40).pipe( + Stream.filter((event) => event.type === "turn.completed"), + Stream.take(1), + Stream.runCollect, + Effect.forkChild, + ); + + const nextTurn = yield* adapter.sendTurn({ + threadId, + input: "plain follow-up after prompt_complete", + attachments: [], + }); + + const secondCompletedEvents = Array.from(yield* Fiber.join(secondCompletedFiber)); + const secondCompleted = secondCompletedEvents[0]; + assert.equal(secondCompleted?.type, "turn.completed"); + if (secondCompleted?.type === "turn.completed") { + assert.equal(secondCompleted.turnId, nextTurn.turnId); + assert.equal(secondCompleted.payload.state, "completed"); + assert.equal(secondCompleted.payload.stopReason, "end_turn"); + } + + const session = (yield* adapter.listSessions()).find( + (candidate) => candidate.threadId === threadId, + ); + assert.equal(session?.status, "ready"); + assert.isUndefined(session?.activeTurnId); + yield* adapter.stopSession(threadId); }), ); diff --git a/apps/server/src/provider/Layers/GrokBuildAdapter.ts b/apps/server/src/provider/Layers/GrokBuildAdapter.ts index 876aa7ffc2f8..d250fcb23ef8 100644 --- a/apps/server/src/provider/Layers/GrokBuildAdapter.ts +++ b/apps/server/src/provider/Layers/GrokBuildAdapter.ts @@ -228,6 +228,15 @@ export function makeGrokBuildAdapter( ctx.stopped = true; yield* settlePendingApprovalsAsCancelled(ctx.pendingApprovals); yield* settlePendingAcpUserInputsAsCancelled(ctx.pendingUserInputs); + ctx.interruptedTurnIds.clear(); + ctx.activeTurnId = undefined; + ctx.promptsInFlight = 0; + const promptFibers = Array.from(ctx.promptFibers); + ctx.promptFibers.clear(); + yield* Effect.forEach(promptFibers, (fiber) => Fiber.interrupt(fiber), { + discard: true, + }); + yield* ctx.acp.cancel.pipe(Effect.ignore); if (ctx.notificationFiber) { yield* Fiber.interrupt(ctx.notificationFiber); } @@ -586,6 +595,8 @@ export function makeGrokBuildAdapter( pendingUserInputs, turns: [], activeTurnId: undefined, + interruptedTurnIds: new Set(), + promptFibers: new Set(), lastPlanFingerprint: undefined, currentModelId: mapGrokSlugToAcpModelId(input.modelSelection?.model), promptsInFlight: 0, @@ -684,13 +695,55 @@ export function makeGrokBuildAdapter( sendTurn, - interruptTurn: (threadId) => + interruptTurn: (threadId, turnId) => withThreadLock( threadId, Effect.gen(function* () { const ctx = yield* requireSession(threadId); + const targetTurnId = turnId ?? ctx.activeTurnId; + if (targetTurnId === undefined) { + yield* settlePendingApprovalsAsCancelled(ctx.pendingApprovals); + yield* settlePendingAcpUserInputsAsCancelled(ctx.pendingUserInputs); + yield* ctx.acp.cancel.pipe(Effect.ignore); + return; + } + const shouldCancelActiveTurn = ctx.activeTurnId === targetTurnId; + if (turnId !== undefined && !shouldCancelActiveTurn) { + return; + } + ctx.interruptedTurnIds.add(targetTurnId); yield* settlePendingApprovalsAsCancelled(ctx.pendingApprovals); yield* settlePendingAcpUserInputsAsCancelled(ctx.pendingUserInputs); + if (shouldCancelActiveTurn) { + ctx.activeTurnId = undefined; + ctx.promptsInFlight = 0; + const promptFibers = Array.from(ctx.promptFibers); + ctx.promptFibers.clear(); + yield* Effect.forEach(promptFibers, (fiber) => Fiber.interrupt(fiber), { + discard: true, + }); + const { + activeTurnId: _activeTurnId, + lastError: _lastError, + ...session + } = ctx.session; + ctx.session = { + ...session, + status: "ready", + updatedAt: yield* nowIso, + }; + yield* offerRuntimeEvent({ + type: "turn.completed", + ...(yield* makeEventStamp()), + provider: PROVIDER, + threadId, + turnId: targetTurnId, + payload: { + state: "cancelled", + stopReason: "cancelled", + }, + }); + } yield* ctx.acp.cancel.pipe(Effect.ignore); }), ), diff --git a/apps/server/src/provider/Layers/GrokBuildSendTurn.ts b/apps/server/src/provider/Layers/GrokBuildSendTurn.ts index fa3b6873784b..753305cc3594 100644 --- a/apps/server/src/provider/Layers/GrokBuildSendTurn.ts +++ b/apps/server/src/provider/Layers/GrokBuildSendTurn.ts @@ -8,6 +8,7 @@ import { TurnId, } from "@t3tools/contracts"; import * as Effect from "effect/Effect"; +import type * as Fiber from "effect/Fiber"; import * as Scope from "effect/Scope"; import type * as EffectAcpSchema from "effect-acp/schema"; @@ -40,6 +41,8 @@ export interface GrokBuildSendTurnContext { readonly promptCapabilities: GrokAcpPromptCapabilities; readonly turns: Array; activeTurnId: TurnId | undefined; + readonly interruptedTurnIds: Set; + readonly promptFibers: Set>; lastPlanFingerprint: string | undefined; readonly currentModelId: string | undefined; promptsInFlight: number; @@ -89,7 +92,12 @@ export function makeGrokBuildSendTurn(input: { context.promptsInFlight = Math.max(0, context.promptsInFlight - 1); }); const releasePrompt = (prepared: PreparedGrokBuildPrompt) => - releaseContextPrompt(prepared.context); + Effect.sync(() => { + if (prepared.context.interruptedTurnIds.has(prepared.turnId)) { + return; + } + prepared.context.promptsInFlight = Math.max(0, prepared.context.promptsInFlight - 1); + }); const preparePrompt = (request: Parameters[0]) => input.withThreadLock( @@ -201,6 +209,10 @@ export function makeGrokBuildSendTurn(input: { detail: "Grok Build session changed before the turn completed.", }); } + if (context.interruptedTurnIds.has(prepared.turnId)) { + context.interruptedTurnIds.delete(prepared.turnId); + return; + } const existingTurnRecord = context.turns.find((turn) => turn.id === prepared.turnId); const item = { prompt: prepared.promptBlocks, result }; @@ -249,6 +261,10 @@ export function makeGrokBuildSendTurn(input: { if (context !== prepared.context) { return; } + if (context.interruptedTurnIds.has(prepared.turnId)) { + context.interruptedTurnIds.delete(prepared.turnId); + return; + } if (context.promptsInFlight === 1 && context.activeTurnId === prepared.turnId) { context.activeTurnId = undefined; const { activeTurnId: _activeTurnId, ...session } = context.session; @@ -306,14 +322,26 @@ export function makeGrokBuildSendTurn(input: { ); }).pipe(Effect.tapCause(() => releasePrompt(prepared))); - yield* context.acp.prompt(payload).pipe( - Effect.tap((result) => - input.logNative(request.threadId, "session/prompt(response)", result), + const promptFiber = yield* input.withThreadLock( + request.threadId, + context.acp.prompt(payload).pipe( + Effect.tap((result) => + input.logNative(request.threadId, "session/prompt(response)", result), + ), + Effect.flatMap((result) => settleSuccess(request, prepared, result)), + Effect.catchCause(() => settleFailure(request, prepared)), + Effect.ensuring( + Effect.sync(() => { + context.promptFibers.delete(promptFiber); + }).pipe(Effect.andThen(releasePrompt(prepared))), + ), + Effect.forkIn(context.scope), + Effect.tap((fiber) => + Effect.sync(() => { + context.promptFibers.add(fiber); + }), + ), ), - Effect.flatMap((result) => settleSuccess(request, prepared, result)), - Effect.catchCause(() => settleFailure(request, prepared)), - Effect.ensuring(releasePrompt(prepared)), - Effect.forkIn(context.scope), ); return { diff --git a/apps/server/src/provider/acp/GrokAcpSupport.ts b/apps/server/src/provider/acp/GrokAcpSupport.ts index c1776f77e870..19010be29b71 100644 --- a/apps/server/src/provider/acp/GrokAcpSupport.ts +++ b/apps/server/src/provider/acp/GrokAcpSupport.ts @@ -1,6 +1,8 @@ import { type GrokBuildSettings } from "@t3tools/contracts"; import * as Effect from "effect/Effect"; +import * as Deferred from "effect/Deferred"; import * as Layer from "effect/Layer"; +import * as Ref from "effect/Ref"; import * as Scope from "effect/Scope"; import { ChildProcessSpawner } from "effect/unstable/process"; import type * as EffectAcpSchema from "effect-acp/schema"; @@ -12,6 +14,7 @@ import { type AcpSessionRuntimeShape, type AcpSpawnInput, } from "./AcpSessionRuntime.ts"; +import { parseXAiPromptCompleteResponse } from "./XAiAcpExtension.ts"; export const GROK_BUILD_RESUME_VERSION = 1; @@ -125,7 +128,125 @@ export const makeGrokBuildAcpRuntime = ( ), ), ); - return yield* Effect.service(AcpSessionRuntime).pipe(Effect.provide(acpContext)); + const runtime = yield* Effect.service(AcpSessionRuntime).pipe(Effect.provide(acpContext)); + const activeSessionIdRef = yield* Ref.make(undefined); + const promptCompleteWaitersRef = yield* Ref.make( + new Map>>(), + ); + const notificationResolvedWaitersRef = yield* Ref.make( + new Set>(), + ); + + yield* runtime.handleUnknownExtNotification((method, payload) => { + if (method !== "_x.ai/session/prompt_complete" && method !== "x.ai/session/prompt_complete") { + return Effect.void; + } + const completed = parseXAiPromptCompleteResponse(payload); + if (!completed) { + return Effect.void; + } + return Ref.modify(promptCompleteWaitersRef, (waiters) => { + const sessionWaiters = waiters.get(completed.sessionId) ?? []; + const [waiter, ...remainingWaiters] = sessionWaiters; + if (!waiter) { + return [Effect.void, waiters] as const; + } + const next = new Map(waiters); + if (remainingWaiters.length > 0) { + next.set(completed.sessionId, remainingWaiters); + } else { + next.delete(completed.sessionId); + } + return [ + Deferred.succeed(waiter, completed.response).pipe( + Effect.tap(() => + Ref.update(notificationResolvedWaitersRef, (resolved) => { + const nextResolved = new Set(resolved); + nextResolved.add(waiter); + return nextResolved; + }), + ), + Effect.asVoid, + ), + next, + ] as const; + }).pipe(Effect.flatten); + }); + + const start: AcpSessionRuntimeShape["start"] = () => + runtime.start().pipe(Effect.tap((started) => Ref.set(activeSessionIdRef, started.sessionId))); + + const settlePromptCompleteWaitersAsCancelled = Ref.modify( + promptCompleteWaitersRef, + (waiters) => [ + Effect.forEach( + Array.from(waiters.values()).flat(), + (waiter) => + Deferred.succeed(waiter, { + stopReason: "cancelled", + } satisfies EffectAcpSchema.PromptResponse), + { discard: true }, + ), + new Map>>(), + ], + ).pipe(Effect.flatten); + + const prompt: AcpSessionRuntimeShape["prompt"] = (payload) => + Effect.gen(function* () { + const sessionId = yield* Ref.get(activeSessionIdRef); + if (!sessionId) { + return yield* runtime.prompt(payload); + } + const waiter = yield* Deferred.make(); + // xAI prompt_complete notifications carry no turn id; waiters are FIFO per session. + yield* Ref.update(promptCompleteWaitersRef, (waiters) => { + const next = new Map(waiters); + next.set(sessionId, [...(next.get(sessionId) ?? []), waiter]); + return next; + }); + return yield* Effect.raceFirst(runtime.prompt(payload), Deferred.await(waiter)).pipe( + Effect.ensuring( + Effect.gen(function* () { + const completedViaNotification = (yield* Ref.get(notificationResolvedWaitersRef)).has( + waiter, + ); + if (completedViaNotification) { + yield* Ref.update(notificationResolvedWaitersRef, (resolved) => { + const nextResolved = new Set(resolved); + nextResolved.delete(waiter); + return nextResolved; + }); + yield* runtime.cancel.pipe(Effect.ignore); + } + yield* Ref.update(promptCompleteWaitersRef, (waiters) => { + const sessionWaiters = waiters.get(sessionId); + if (!sessionWaiters?.includes(waiter)) { + return waiters; + } + const remainingWaiters = sessionWaiters.filter((candidate) => candidate !== waiter); + const next = new Map(waiters); + if (remainingWaiters.length > 0) { + next.set(sessionId, remainingWaiters); + } else { + next.delete(sessionId); + } + return next; + }); + }), + ), + ); + }); + + const cancel: AcpSessionRuntimeShape["cancel"] = settlePromptCompleteWaitersAsCancelled.pipe( + Effect.andThen(runtime.cancel), + ); + + return { + ...runtime, + start, + prompt, + cancel, + } satisfies AcpSessionRuntimeShape; }); export function extractGrokAcpPromptCapabilities( diff --git a/apps/server/src/provider/acp/XAiAcpExtension.ts b/apps/server/src/provider/acp/XAiAcpExtension.ts index 113316e6e63a..47fed6b14644 100644 --- a/apps/server/src/provider/acp/XAiAcpExtension.ts +++ b/apps/server/src/provider/acp/XAiAcpExtension.ts @@ -1,5 +1,7 @@ import type { ProviderUserInputAnswers, UserInputQuestion } from "@t3tools/contracts"; +import * as Option from "effect/Option"; import * as Schema from "effect/Schema"; +import type * as EffectAcpSchema from "effect-acp/schema"; const XAiAskUserQuestionOption = Schema.Struct({ label: Schema.String, @@ -27,13 +29,33 @@ const XAiWrappedAskUserQuestionParams = Schema.Struct({ params: XAiAskUserQuestionParams, }); +const XAiPromptCompleteParams = Schema.Struct({ + sessionId: Schema.String, + stopReason: Schema.optional( + Schema.NullOr( + Schema.Literals(["end_turn", "max_tokens", "max_turn_requests", "refusal", "cancelled"]), + ), + ), +}); + +const XAiWrappedPromptCompleteParams = Schema.Struct({ + method: Schema.Literals(["x.ai/session/prompt_complete", "_x.ai/session/prompt_complete"]), + params: XAiPromptCompleteParams, +}); + export const XAiAskUserQuestionRequest = Schema.Union([ XAiAskUserQuestionParams, XAiWrappedAskUserQuestionParams, ]); +export const XAiPromptCompleteRequest = Schema.Union([ + XAiPromptCompleteParams, + XAiWrappedPromptCompleteParams, +]); type XAiAskUserQuestionRequestParams = typeof XAiAskUserQuestionParams.Type; type XAiAskUserQuestionRequest = typeof XAiAskUserQuestionRequest.Type; +type XAiPromptCompleteRequest = typeof XAiPromptCompleteRequest.Type; +const decodeXAiPromptCompleteRequest = Schema.decodeUnknownOption(XAiPromptCompleteRequest); const EMPTY_OPTIONS_CONTINUE_FALLBACK = [{ label: "OK", description: "Continue" }] as const; function trimmed(value: string | undefined): string | undefined { @@ -47,6 +69,31 @@ function unwrapAskUserQuestionParams( return "params" in params ? params.params : params; } +function unwrapPromptCompleteParams( + params: XAiPromptCompleteRequest, +): typeof XAiPromptCompleteParams.Type { + return "params" in params ? params.params : params; +} + +export function parseXAiPromptCompleteResponse(input: unknown): + | { + readonly sessionId: string; + readonly response: EffectAcpSchema.PromptResponse; + } + | undefined { + const decoded = decodeXAiPromptCompleteRequest(input); + if (Option.isNone(decoded)) { + return undefined; + } + const params = unwrapPromptCompleteParams(decoded.value); + return { + sessionId: params.sessionId, + response: { + stopReason: params.stopReason ?? "end_turn", + }, + }; +} + export function extractXAiAskUserQuestions( params: XAiAskUserQuestionRequest, ): ReadonlyArray { diff --git a/packages/tailscale/src/childProcessSpawnerTestHelpers.ts b/packages/tailscale/src/childProcessSpawnerTestHelpers.ts index 619a075e9422..3d050b34f20b 100644 --- a/packages/tailscale/src/childProcessSpawnerTestHelpers.ts +++ b/packages/tailscale/src/childProcessSpawnerTestHelpers.ts @@ -46,7 +46,9 @@ export function mockChildProcessSpawnerLayer( command: childProcess.command, args: childProcess.args, }); - return Effect.succeed(mockChildProcessHandle(handler(childProcess.command, childProcess.args))); + return Effect.succeed( + mockChildProcessHandle(handler(childProcess.command, childProcess.args)), + ); }), ); -} \ No newline at end of file +}