diff --git a/apps/server/src/provider/Layers/AntigravityAdapter.test.ts b/apps/server/src/provider/Layers/AntigravityAdapter.test.ts index 4bd0bb3e1010..64904879dad1 100644 --- a/apps/server/src/provider/Layers/AntigravityAdapter.test.ts +++ b/apps/server/src/provider/Layers/AntigravityAdapter.test.ts @@ -589,9 +589,9 @@ it.layer(layer)("AntigravityAdapter", (it) => { }), ); - it.effect("waits for native cancellation before a steer changes the model", () => + it.effect("queues a follow-up turn until the in-flight prompt completes", () => Effect.gen(function* () { - const h = yield* makeHarness({ holdCancel: true }); + const h = yield* makeHarness(); yield* h.adapter.startSession({ threadId, cwd: process.cwd(), @@ -609,38 +609,34 @@ it.layer(layer)("AntigravityAdapter", (it) => { const second = yield* h.adapter .sendTurn({ threadId, - input: "Steer the turn", + input: "Follow-up prompt", modelSelection: { instanceId, model: nativeAlternative }, }) - .pipe(Effect.forkChild); - expect(yield* h.nextCancellation).toBe(1); - expect(h.calls.slice(marker)).toEqual(["cancel:1"]); - yield* h.emitNative({ - _tag: "ContentDelta", - text: "The first prompt stopped.", - rawPayload: {}, - }); - yield* Deferred.succeed(h.cancelRelease, undefined); + .pipe(Effect.forkChild({ startImmediately: true })); + + expect(h.hasActivePrompt()).toBe(true); + expect(h.calls.slice(marker)).toEqual([]); + + yield* Deferred.succeed(initialPrompt.result, { stopReason: "end_turn" }); + const firstResult = yield* Fiber.join(first); + const replacement = yield* h.nextPrompt; expect(replacement.content).toEqual([ - { type: "text", text: "Steer the turn" }, + { type: "text", text: "Follow-up prompt" }, { type: "text", text: expect.stringContaining(`Antigravity harness, as ${nativeAlternative}`), }, ]); expect(h.calls.slice(marker)).toEqual([ - "cancel:1", - "drained:1", `model:${nativeAlternative}`, "mode:default", "prompt:2", ]); yield* Deferred.succeed(replacement.result, { stopReason: "end_turn" }); - const [oldResult, newResult] = yield* Effect.all([Fiber.join(first), Fiber.join(second)]); - expect(oldResult.turnId).toBe(newResult.turnId); - yield* h.waitForEvent((event) => event.type === "turn.completed"); - expect(h.seen.filter((event) => event.type === "turn.completed")).toHaveLength(1); + const secondResult = yield* Fiber.join(second); + expect(firstResult.turnId).not.toBe(secondResult.turnId); + expect(h.seen.filter((event) => event.type === "turn.completed")).toHaveLength(2); expect((yield* h.adapter.listSessions())[0]).toMatchObject({ status: "ready", activeTurnId: undefined, @@ -649,6 +645,48 @@ it.layer(layer)("AntigravityAdapter", (it) => { }), ); + it.effect("queued follow-up turn without explicit model inherits updated session model", () => + Effect.gen(function* () { + const h = yield* makeHarness(); + yield* h.adapter.startSession({ + threadId, + cwd: process.cwd(), + runtimeMode: "approval-required", + }); + const first = yield* h.adapter + .sendTurn({ + threadId, + input: "First prompt", + modelSelection: { instanceId, model: nativeAlternative }, + }) + .pipe(Effect.forkChild); + const initialPrompt = yield* h.nextPrompt; + + const second = yield* h.adapter + .sendTurn({ + threadId, + input: "Second queued prompt without model", + }) + .pipe(Effect.forkChild({ startImmediately: true })); + + expect(h.hasActivePrompt()).toBe(true); + + yield* Deferred.succeed(initialPrompt.result, { stopReason: "end_turn" }); + yield* Fiber.join(first); + + const replacement = yield* h.nextPrompt; + expect(replacement.content).toEqual([ + { type: "text", text: "Second queued prompt without model" }, + { + type: "text", + text: expect.stringContaining(`Antigravity harness, as ${nativeAlternative}`), + }, + ]); + yield* Deferred.succeed(replacement.result, { stopReason: "end_turn" }); + yield* Fiber.join(second); + }), + ); + it.effect("rejects an unavailable steer model without cancelling current work", () => Effect.gen(function* () { const h = yield* makeHarness(); @@ -685,7 +723,7 @@ it.layer(layer)("AntigravityAdapter", (it) => { runtimeMode: "approval-required", }); const first = yield* h.adapter.sendTurn({ threadId, input: "First" }).pipe(Effect.forkChild); - yield* h.nextPrompt; + const prompt1 = yield* h.nextPrompt; h.controls.failModel = true; const failed = yield* h.adapter .sendTurn({ @@ -693,9 +731,12 @@ it.layer(layer)("AntigravityAdapter", (it) => { input: "Replacement", modelSelection: { instanceId, model: nativeAlternative }, }) - .pipe(Effect.exit); - expect(Exit.isFailure(failed)).toBe(true); + .pipe(Effect.forkChild); + yield* Deferred.succeed(prompt1.result, { stopReason: "end_turn" }); yield* Fiber.join(first); + yield* h.waitForEvent((event) => event.type === "turn.completed"); + const failedExit = yield* Fiber.join(failed).pipe(Effect.exit); + expect(Exit.isFailure(failedExit)).toBe(true); const ended = yield* h.waitForEvent((event) => event.type === "turn.completed"); expect(ended.payload.state).toBe("failed"); expect((yield* h.adapter.listSessions())[0]).toMatchObject({ @@ -1090,6 +1131,8 @@ it.layer(layer)("AntigravityAdapter", (it) => { const steering = yield* h.adapter .sendTurn({ threadId, input: "Change direction" }) .pipe(Effect.forkChild); + yield* Deferred.succeed(prompt.result, { stopReason: "end_turn" }); + yield* Fiber.join(sending); const replacement = yield* h.nextPrompt; yield* Deferred.succeed(replacement.result, { stopReason: "end_turn" }); yield* Fiber.join(steering); @@ -1101,12 +1144,7 @@ it.layer(layer)("AntigravityAdapter", (it) => { taskId: "trajectory:4", title: "Antigravity subagent batch", taskType: "subagent_batch", - status: - stop === "disconnect" - ? "failed" - : stop === "cancel" || stop === "steer" - ? "cancelled" - : "idle", + status: stop === "disconnect" ? "failed" : stop === "cancel" ? "cancelled" : "idle", }); if (stop === "disconnect") yield* h.waitForEvent((event) => event.type === "session.exited"); diff --git a/apps/server/src/provider/Layers/AntigravityAdapter.ts b/apps/server/src/provider/Layers/AntigravityAdapter.ts index 61a9b3c3a645..ce85560bfad2 100644 --- a/apps/server/src/provider/Layers/AntigravityAdapter.ts +++ b/apps/server/src/provider/Layers/AntigravityAdapter.ts @@ -198,6 +198,7 @@ interface SessionContext { readonly promptLock: Semaphore.Semaphore; readonly stopLock: Semaphore.Semaphore; readonly commandLock: Semaphore.Semaphore; + readonly turnQueue: Semaphore.Semaphore; readonly approvals: Map; readonly questions: Map; readonly commands: Map; @@ -871,6 +872,7 @@ export const makeAntigravityAdapter = Effect.fn("makeAntigravityAdapter")(functi promptLock: yield* Semaphore.make(1), stopLock: yield* Semaphore.make(1), commandLock: yield* Semaphore.make(1), + turnQueue: yield* Semaphore.make(1), approvals: new Map(), questions: new Map(), commands: new Map(), @@ -1026,10 +1028,28 @@ export const makeAntigravityAdapter = Effect.fn("makeAntigravityAdapter")(functi }); }).pipe(Effect.uninterruptible); - return yield* Effect.gen(function* () { - const launch = yield* context.promptLock.withPermit( + if (input.modelSelection?.model) { + const configOptions = yield* context.runtime.getConfigOptions; + const explicitModel = resolveAntigravityModel({ + configOptions, + model: input.modelSelection.model, + defaultModel: yield* options.defaultModel ?? Effect.succeed(undefined), + }); + const availableModels = antigravityModelOptions(configOptions); + if (explicitModel && !availableModels.some((option) => option.value === explicitModel)) { + return yield* mapAntigravityError( + input.threadId, + "session/prompt", + EffectAcpErrors.AcpRequestError.invalidParams( + `Antigravity model '${explicitModel}' is unavailable for this Google account. Select an available model.`, + ), + ); + } + } + + return yield* context.turnQueue + .withPermit( Effect.gen(function* () { - yield* requireSession(input.threadId); const requestedModel = input.modelSelection?.model ?? context.session.model; const configOptions = yield* context.runtime.getConfigOptions; const model = resolveAntigravityModel({ @@ -1039,127 +1059,135 @@ export const makeAntigravityAdapter = Effect.fn("makeAntigravityAdapter")(functi }); const availableModels = antigravityModelOptions(configOptions); if (model && !availableModels.some((option) => option.value === model)) { - return yield* EffectAcpErrors.AcpRequestError.invalidParams( - `Antigravity model '${model}' is unavailable for this Google account. Select an available model.`, + return yield* mapAntigravityError( + input.threadId, + "session/prompt", + EffectAcpErrors.AcpRequestError.invalidParams( + `Antigravity model '${model}' is unavailable for this Google account. Select an available model.`, + ), ); } - const turnId = context.activeTurnId ?? TurnId.make(yield* randomId); - const steering = context.activeTurnId !== undefined; - const turn: TurnIntent = { turnId, generation: ++context.generation, settled: false }; - intent = turn; - context.activeTurnId = turnId; - if (!steering) { - yield* emit({ - type: "turn.started", - ...(yield* stamp), + + const launch = yield* context.promptLock.withPermit( + Effect.gen(function* () { + yield* requireSession(input.threadId); + const turnId = TurnId.make(yield* randomId); + const turn: TurnIntent = { turnId, generation: ++context.generation, settled: false }; + intent = turn; + context.activeTurnId = turnId; + yield* emit({ + type: "turn.started", + ...(yield* stamp), + provider: PROVIDER, + threadId: input.threadId, + turnId, + payload: model ? { model } : {}, + }); + yield* applyAntigravityAcpModelSelection({ + runtime: context.runtime, + model, + mapError: (cause) => cause, + }); + yield* context.runtime.setMode( + antigravityPermissionMode(context.session.runtimeMode), + ); + context.session = { + ...context.session, + status: "running", + activeTurnId: turnId, + ...(model ? { model } : {}), + updatedAt: yield* nowIso, + }; + const dispatched = yield* Deferred.make(); + const fiber = yield* context.runtime + .prompt( + { + prompt: [ + ...prompt, + { + type: "text", + text: buildRuntimeInstructions({ harness: "Antigravity", model }), + }, + ], + }, + { dispatched }, + ) + .pipe(Effect.forkIn(context.scope)); + context.promptFiber = fiber; + // Fiber.join can skip a scope-close waiter when the child is interrupted. + // Unwrap the Exit after Fiber.await returns. + yield* Effect.raceFirst( + Deferred.await(dispatched), + Fiber.await(fiber).pipe( + Effect.flatMap((exit) => exit), + Effect.asVoid, + ), + ); + return { turn, fiber }; + }), + ); + const result = yield* Fiber.await(launch.fiber).pipe(Effect.flatMap((exit) => exit)); + yield* context.runtime.drainEvents; + if (context.stopped) { + return yield* new ProviderAdapterSessionClosedError({ provider: PROVIDER, threadId: input.threadId, - turnId, - payload: model ? { model } : {}, }); } - if (context.promptFiber) { - yield* cancelRequests(context); - yield* context.runtime.cancel; - yield* Fiber.await(context.promptFiber); - yield* finishSubagents(context, "cancelled"); - } - yield* applyAntigravityAcpModelSelection({ - runtime: context.runtime, - model, - mapError: (cause) => cause, - }); - yield* context.runtime.setMode(antigravityPermissionMode(context.session.runtimeMode)); - context.session = { - ...context.session, - status: "running", - activeTurnId: turnId, - ...(model ? { model } : {}), - updatedAt: yield* nowIso, - }; - const dispatched = yield* Deferred.make(); - const fiber = yield* context.runtime - .prompt( - { - prompt: [ - ...prompt, - { - type: "text", - text: buildRuntimeInstructions({ harness: "Antigravity", model }), - }, - ], - }, - { dispatched }, - ) - .pipe(Effect.forkIn(context.scope)); - context.promptFiber = fiber; - // Fiber.join can skip a scope-close waiter when the child is interrupted. - // Unwrap the Exit after Fiber.await returns. - yield* Effect.raceFirst( - Deferred.await(dispatched), - Fiber.await(fiber).pipe( - Effect.flatMap((exit) => exit), - Effect.asVoid, - ), + const record = context.turns.find((turn) => turn.id === launch.turn.turnId); + if (record) record.items.push(result); + else context.turns.push({ id: launch.turn.turnId, items: [result] }); + yield* context.promptLock.withPermit( + finishTurn(launch.turn, { + state: result.stopReason === "cancelled" ? "cancelled" : "completed", + stopReason: result.stopReason, + }), ); - return { turn, fiber }; - }), - ); - const result = yield* Fiber.await(launch.fiber).pipe(Effect.flatMap((exit) => exit)); - yield* context.runtime.drainEvents; - if (context.stopped) { - return yield* new ProviderAdapterSessionClosedError({ - provider: PROVIDER, - threadId: input.threadId, - }); - } - const record = context.turns.find((turn) => turn.id === launch.turn.turnId); - if (record) record.items.push(result); - else context.turns.push({ id: launch.turn.turnId, items: [result] }); - yield* context.promptLock.withPermit( - finishTurn(launch.turn, { - state: result.stopReason === "cancelled" ? "cancelled" : "completed", - stopReason: result.stopReason, + return { + threadId: input.threadId, + turnId: launch.turn.turnId, + resumeCursor: context.session.resumeCursor, + }; }), - ); - return { - threadId: input.threadId, - turnId: launch.turn.turnId, - resumeCursor: context.session.resumeCursor, - }; - }).pipe( - Effect.tapError((cause) => - isAntigravitySignInRequiredError(cause) - ? (options.onAuthRequired ?? Effect.void) - : Effect.void, - ), - Effect.mapError((cause) => - isAcpError(cause) ? mapAntigravityError(input.threadId, "session/prompt", cause) : cause, - ), - Effect.tapError((cause) => - Effect.suspend(() => - intent - ? context.promptLock.withPermit( - finishTurn(intent, { state: "failed", errorMessage: cause.message }), - ) + ) + .pipe( + Effect.tapError((cause) => + isAntigravitySignInRequiredError(cause) + ? (options.onAuthRequired ?? Effect.void) : Effect.void, ), - ), - Effect.onInterrupt(() => - context.promptLock.withPermit( - Effect.gen(function* () { - const turn = intent; - if (!turn || turn.settled || context.stopped || context.generation !== turn.generation) - return; - const promptFiber = context.promptFiber; - yield* cancelRequests(context); - yield* Effect.ignore(context.runtime.cancel); - if (promptFiber) yield* Fiber.interrupt(promptFiber); - yield* finishTurn(turn, { state: "cancelled", stopReason: "cancelled" }); - }), + Effect.mapError((cause) => + isAcpError(cause) ? mapAntigravityError(input.threadId, "session/prompt", cause) : cause, ), - ), - ); + Effect.tapError((cause) => + Effect.suspend(() => + intent + ? context.promptLock.withPermit( + finishTurn(intent, { state: "failed", errorMessage: cause.message }), + ) + : Effect.void, + ), + ), + Effect.onInterrupt(() => + context.promptLock.withPermit( + Effect.gen(function* () { + const turn = intent; + if ( + !turn || + turn.settled || + context.stopped || + context.generation !== turn.generation + ) + return; + const promptFiber = context.promptFiber; + yield* cancelRequests(context); + yield* Effect.ignore(context.runtime.cancel); + if (promptFiber) yield* Fiber.interrupt(promptFiber); + yield* finishTurn(turn, { state: "cancelled", stopReason: "cancelled" }); + }), + ), + ), + ); }); const interruptTurn: Adapter["interruptTurn"] = (threadId) =>