diff --git a/apps/server/src/provider/Drivers/AntigravityDriver.ts b/apps/server/src/provider/Drivers/AntigravityDriver.ts new file mode 100644 index 000000000000..e9778779e331 --- /dev/null +++ b/apps/server/src/provider/Drivers/AntigravityDriver.ts @@ -0,0 +1,155 @@ +import { AntigravitySettings, ProviderDriverKind, type ServerProvider } from "@t3tools/contracts"; +import * as Crypto from "effect/Crypto"; +import * as Effect from "effect/Effect"; +import * as FileSystem from "effect/FileSystem"; +import * as Path from "effect/Path"; +import * as Schema from "effect/Schema"; +import * as Scope from "effect/Scope"; +import { HttpClient } from "effect/unstable/http"; +import { ChildProcessSpawner } from "effect/unstable/process"; + +import * as BackgroundPolicy from "../../background/BackgroundPolicy.ts"; +import { ServerConfig } from "../../config.ts"; +import { ServerSettingsService } from "../../serverSettings.ts"; +import { ProviderDriverError } from "../Errors.ts"; +import { makeAntigravityAdapter } from "../Layers/AntigravityAdapter.ts"; +import { + buildInitialAntigravityProviderSnapshot, + checkAntigravityProviderStatus, +} from "../Layers/AntigravityProvider.ts"; +import { ProviderEventLoggers } from "../Layers/ProviderEventLoggers.ts"; +import { makeManagedServerProvider } from "../makeManagedServerProvider.ts"; +import { + defaultProviderContinuationIdentity, + type ProviderDriver, + type ProviderInstance, +} from "../ProviderDriver.ts"; +import type { ServerProviderDraft } from "../providerSnapshot.ts"; +import { mergeProviderInstanceEnvironment } from "../ProviderInstanceEnvironment.ts"; +import { + makeManualOnlyProviderMaintenanceCapabilities, + makeStaticProviderMaintenanceResolver, + resolveProviderMaintenanceCapabilitiesEffect, +} from "../providerMaintenance.ts"; +import { + haveProviderSnapshotSettingsChanged, + makeProviderSnapshotSettingsSource, + type ProviderSnapshotSettings, +} from "../providerUpdateSettings.ts"; + +const decodeAntigravitySettings = Schema.decodeSync(AntigravitySettings); + +const DRIVER_KIND = ProviderDriverKind.make("antigravity"); +const UPDATE = makeStaticProviderMaintenanceResolver( + makeManualOnlyProviderMaintenanceCapabilities({ + provider: DRIVER_KIND, + packageName: null, + }), +); + +export type AntigravityDriverEnv = + | BackgroundPolicy.BackgroundPolicy + | ChildProcessSpawner.ChildProcessSpawner + | Crypto.Crypto + | FileSystem.FileSystem + | HttpClient.HttpClient + | Path.Path + | ProviderEventLoggers + | ServerConfig + | ServerSettingsService + | Scope.Scope; + +const withInstanceIdentity = + (input: { + readonly instanceId: ProviderInstance["instanceId"]; + readonly displayName: string | undefined; + readonly accentColor: string | undefined; + readonly continuationGroupKey: string; + }) => + (snapshot: ServerProviderDraft): ServerProvider => ({ + ...snapshot, + instanceId: input.instanceId, + driver: DRIVER_KIND, + ...(input.displayName ? { displayName: input.displayName } : {}), + ...(input.accentColor ? { accentColor: input.accentColor } : {}), + continuation: { groupKey: input.continuationGroupKey }, + }); + +export const AntigravityDriver: ProviderDriver = { + driverKind: DRIVER_KIND, + metadata: { + displayName: "Antigravity", + supportsMultipleInstances: true, + }, + configSchema: AntigravitySettings, + defaultConfig: (): AntigravitySettings => decodeAntigravitySettings({}), + create: ({ instanceId, displayName, accentColor, environment, enabled, config }) => + Effect.gen(function* () { + const crypto = yield* Crypto.Crypto; + const spawner = yield* ChildProcessSpawner.ChildProcessSpawner; + const serverSettings = yield* ServerSettingsService; + const processEnv = mergeProviderInstanceEnvironment(environment); + const continuationIdentity = defaultProviderContinuationIdentity({ + driverKind: DRIVER_KIND, + instanceId, + }); + const stampIdentity = withInstanceIdentity({ + instanceId, + displayName, + accentColor, + continuationGroupKey: continuationIdentity.continuationKey, + }); + const effectiveConfig = { ...config, enabled } satisfies AntigravitySettings; + const maintenanceCapabilities = yield* resolveProviderMaintenanceCapabilitiesEffect(UPDATE, { + binaryPath: effectiveConfig.binaryPath, + env: processEnv, + }); + + const adapter = yield* makeAntigravityAdapter(effectiveConfig, { + environment: processEnv, + instanceId, + }); + + const checkProvider = checkAntigravityProviderStatus(effectiveConfig, processEnv).pipe( + Effect.map(stampIdentity), + Effect.provideService(Crypto.Crypto, crypto), + Effect.provideService(ChildProcessSpawner.ChildProcessSpawner, spawner), + ); + + const snapshotSettings = makeProviderSnapshotSettingsSource(effectiveConfig, serverSettings); + const snapshot = yield* makeManagedServerProvider< + ProviderSnapshotSettings + >({ + maintenanceCapabilities, + getSettings: snapshotSettings.getSettings, + streamSettings: snapshotSettings.streamSettings, + haveSettingsChanged: haveProviderSnapshotSettingsChanged, + initialSnapshot: (settings) => + buildInitialAntigravityProviderSnapshot(settings.provider).pipe( + Effect.map(stampIdentity), + ), + checkProvider, + }).pipe( + Effect.mapError( + (cause) => + new ProviderDriverError({ + driver: DRIVER_KIND, + instanceId, + detail: `Failed to build Antigravity snapshot: ${cause.message ?? String(cause)}`, + cause, + }), + ), + ); + + return { + instanceId, + driverKind: DRIVER_KIND, + continuationIdentity, + displayName, + accentColor, + enabled, + snapshot, + adapter, + } satisfies ProviderInstance; + }), +}; diff --git a/apps/server/src/provider/Layers/AntigravityAdapter.test.ts b/apps/server/src/provider/Layers/AntigravityAdapter.test.ts new file mode 100644 index 000000000000..fd837f5aa87a --- /dev/null +++ b/apps/server/src/provider/Layers/AntigravityAdapter.test.ts @@ -0,0 +1,364 @@ +import * as NodeServices from "@effect/platform-node/NodeServices"; +import { describe, expect, it } from "@effect/vitest"; +import * as Deferred from "effect/Deferred"; +import * as Duration from "effect/Duration"; +import * as Effect from "effect/Effect"; +import * as FileSystem from "effect/FileSystem"; +import * as Path from "effect/Path"; +import * as Schema from "effect/Schema"; +import * as Stream from "effect/Stream"; +import { + AntigravitySettings, + ProviderInstanceId, + ThreadId, + type ProviderRuntimeEvent, +} from "@t3tools/contracts"; + +import { makeAntigravityAdapter } from "./AntigravityAdapter.ts"; + +const decodeAntigravitySettings = Schema.decodeSync(AntigravitySettings); + +it.layer(NodeServices.layer)("AntigravityAdapter", (it) => { + it.effect("manages session lifecycle: start, list, has, stop", () => + Effect.scoped( + Effect.gen(function* () { + const adapter = yield* makeAntigravityAdapter( + decodeAntigravitySettings({ enabled: true }), + { instanceId: ProviderInstanceId.make("antigravity") }, + ); + + const threadId = ThreadId.make("thread_test_123"); + const session = yield* adapter.startSession({ + threadId, + cwd: "/tmp", + runtimeMode: "full-access", + }); + + expect(session.threadId).toBe(threadId); + expect(session.status).toBe("ready"); + expect(yield* adapter.hasSession(threadId)).toBe(true); + + const sessions = yield* adapter.listSessions(); + expect(sessions.length).toBe(1); + expect(sessions[0].threadId).toBe(threadId); + + yield* adapter.stopSession(threadId); + expect(yield* adapter.hasSession(threadId)).toBe(false); + }), + ), + ); + + it.effect("executes bounded read-only turn, streams deltas and tool events, and completes", () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const dir = yield* fs.makeTempDirectoryScoped({ prefix: "t3code-agy-adapter-" }); + const mockAgyPath = path.join(dir, "agy"); + + const line1 = JSON.stringify({ + event: "init", + conversation_id: "conv_agy_test_456", + init: { + cwd: dir, + model: "gemini-3.7-flash-high", + permission_mode: "request-review", + tools: ["run_command"], + }, + }); + const line2 = JSON.stringify({ + event: "step_update", + step_update: { + step_index: 1, + step_type: "agent_response", + state: "ACTIVE", + text_delta: "Inspecting repository safely...", + }, + }); + const line3 = JSON.stringify({ + event: "step_update", + step_update: { + step_index: 2, + step_type: "tool", + state: "ACTIVE", + tool_name: "run_command", + tool_info: { + name: "run_command", + parameters: { CommandLine: "echo READ_ONLY_PROBE" }, + }, + }, + }); + const line4 = JSON.stringify({ + event: "step_update", + step_update: { + step_index: 2, + step_type: "tool", + state: "DONE", + tool_name: "run_command", + duration_seconds: 0.12, + tool_info: { + name: "run_command", + parameters: { CommandLine: "echo READ_ONLY_PROBE" }, + output: "READ_ONLY_PROBE", + }, + }, + }); + const line5 = JSON.stringify({ + event: "result", + result: { + conversation_id: "conv_agy_test_456", + status: "SUCCESS", + response: "Inspection complete. Output: READ_ONLY_PROBE", + duration_seconds: 0.5, + num_turns: 1, + usage: { input_tokens: 50, output_tokens: 15 }, + }, + }); + + yield* fs.writeFileString( + mockAgyPath, + [ + "#!/bin/sh", + `cat << 'RAW_EOF'`, + line1, + line2, + line3, + line4, + line5, + `RAW_EOF`, + "exit 0", + "", + ].join("\n"), + ); + yield* fs.chmod(mockAgyPath, 0o755); + + const adapter = yield* makeAntigravityAdapter( + decodeAntigravitySettings({ enabled: true, binaryPath: mockAgyPath }), + { instanceId: ProviderInstanceId.make("antigravity") }, + ); + + const threadId = ThreadId.make("thread_turn_test"); + yield* adapter.startSession({ + threadId, + cwd: dir, + runtimeMode: "full-access", + }); + + const runtimeEvents: ProviderRuntimeEvent[] = []; + const turnCompleted = yield* Deferred.make(); + + yield* Stream.runForEach(adapter.streamEvents, (event) => + Effect.sync(() => { + runtimeEvents.push(event); + }).pipe( + Effect.andThen( + event.type === "turn.completed" || event.type === "turn.aborted" + ? Deferred.succeed(turnCompleted, undefined) + : Effect.void, + ), + ), + ).pipe(Effect.forkScoped); + + const turnResult = yield* adapter.sendTurn({ + threadId, + input: "Run safe inspection", + interactionMode: "plan", + }); + + expect(turnResult.threadId).toBe(threadId); + expect(turnResult.turnId).toBeDefined(); + + yield* Deferred.await(turnCompleted); + expect(runtimeEvents.length).toBeGreaterThanOrEqual(4); + + // 1. Content delta + const delta = runtimeEvents.find((e) => e.type === "content.delta"); + expect(delta).toBeDefined(); + if (delta && delta.type === "content.delta") { + expect(delta.payload.delta).toBe("Inspecting repository safely..."); + } + + // 2. Tool started and completed + const toolStarted = runtimeEvents.find( + (e) => e.type === "item.started" && (e.payload as any)?.itemType === "tool_use", + ); + expect(toolStarted).toBeDefined(); + expect((toolStarted?.payload as any)?.title).toBe("run_command"); + + const toolCompleted = runtimeEvents.find( + (e) => e.type === "item.completed" && (e.payload as any)?.itemType === "tool_use", + ); + expect(toolCompleted).toBeDefined(); + expect((toolCompleted?.payload as any)?.data?.output).toContain("READ_ONLY_PROBE"); + + // 3. Turn completed + const completed = runtimeEvents.find((e) => e.type === "turn.completed"); + expect(completed).toBeDefined(); + if (completed && completed.type === "turn.completed") { + expect(completed.payload.state).toBe("completed"); + } + }), + ), + ); + + it.effect("emits turn.aborted when interruptTurn runs against an active turn", () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const dir = yield* fs.makeTempDirectoryScoped({ prefix: "t3code-agy-interrupt-" }); + const mockAgyPath = path.join(dir, "agy"); + + // Exit quickly after one line — interrupt races the short window while + // activeProcess is set. Abort is emitted by interruptTurn itself. + yield* fs.writeFileString( + mockAgyPath, + [ + "#!/bin/sh", + 'printf \'%s\\n\' \'{"event":"init","conversation_id":"conv_int_1","init":{"cwd":"/tmp","model":"gemini-3.7-flash-high","permission_mode":"request-review","tools":[]}}\'', + 'printf \'%s\\n\' \'{"event":"result","result":{"conversation_id":"conv_int_1","status":"SUCCESS","response":"done"}}\'', + "exit 0", + "", + ].join("\n"), + ); + yield* fs.chmod(mockAgyPath, 0o755); + + const adapter = yield* makeAntigravityAdapter( + decodeAntigravitySettings({ enabled: true, binaryPath: mockAgyPath }), + { instanceId: ProviderInstanceId.make("antigravity") }, + ); + + const threadId = ThreadId.make("thread_interrupt_test"); + yield* adapter.startSession({ + threadId, + cwd: dir, + runtimeMode: "full-access", + }); + + const runtimeEvents: ProviderRuntimeEvent[] = []; + yield* Stream.runForEach(adapter.streamEvents, (event) => + Effect.sync(() => { + runtimeEvents.push(event); + }), + ).pipe(Effect.forkScoped); + + const turnResult = yield* adapter.sendTurn({ + threadId, + input: "Test interrupt handling", + }); + expect(turnResult.turnId).toBeDefined(); + + yield* adapter.interruptTurn(threadId, turnResult.turnId); + + // Either interrupt emitted abort, or the turn already completed — both + // prove the adapter did not suppress the terminal event. + const terminal = runtimeEvents.filter( + (event) => event.type === "turn.aborted" || event.type === "turn.completed", + ); + expect(terminal.length).toBeGreaterThanOrEqual(1); + + yield* adapter.stopSession(threadId); + expect(yield* adapter.hasSession(threadId)).toBe(false); + }), + ), + ); + + it.effect("maps tool ERROR steps to failed item.completed events", () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const dir = yield* fs.makeTempDirectoryScoped({ prefix: "t3code-agy-tool-error-" }); + const mockAgyPath = path.join(dir, "agy"); + + const lines = [ + JSON.stringify({ + event: "init", + conversation_id: "conv_err", + init: { + cwd: dir, + model: "gemini-3.7-flash-high", + permission_mode: "default", + tools: [], + }, + }), + JSON.stringify({ + event: "step_update", + step_update: { + step_index: 1, + step_type: "tool", + state: "ACTIVE", + tool_name: "run_command", + tool_info: { name: "run_command", parameters: { CommandLine: "false" } }, + }, + }), + JSON.stringify({ + event: "step_update", + step_update: { + step_index: 1, + step_type: "tool", + state: "ERROR", + tool_name: "run_command", + tool_info: { name: "run_command", error: "exit 1" }, + }, + }), + JSON.stringify({ + event: "result", + result: { + conversation_id: "conv_err", + status: "SUCCESS", + response: "tool failed but turn continued", + }, + }), + ]; + + yield* fs.writeFileString( + mockAgyPath, + ["#!/bin/sh", "cat << 'RAW_EOF'", ...lines, "RAW_EOF", "exit 0", ""].join("\n"), + ); + yield* fs.chmod(mockAgyPath, 0o755); + + const adapter = yield* makeAntigravityAdapter( + decodeAntigravitySettings({ + enabled: true, + binaryPath: mockAgyPath, + launchArgs: "--sandbox", + effort: "medium", + }), + { instanceId: ProviderInstanceId.make("antigravity") }, + ); + + const threadId = ThreadId.make("thread_tool_error"); + yield* adapter.startSession({ + threadId, + cwd: dir, + runtimeMode: "approval-required", + }); + + const runtimeEvents: ProviderRuntimeEvent[] = []; + const turnCompleted = yield* Deferred.make(); + yield* Stream.runForEach(adapter.streamEvents, (event) => + Effect.sync(() => { + runtimeEvents.push(event); + }).pipe( + Effect.andThen( + event.type === "turn.completed" || event.type === "turn.aborted" + ? Deferred.succeed(turnCompleted, undefined) + : Effect.void, + ), + ), + ).pipe(Effect.forkScoped); + + yield* adapter.sendTurn({ threadId, input: "force tool error" }); + yield* Deferred.await(turnCompleted); + + const failed = runtimeEvents.find( + (e) => + e.type === "item.completed" && (e.payload as { status?: string }).status === "failed", + ); + expect(failed).toBeDefined(); + expect((failed?.payload as { title?: string }).title).toBe("run_command"); + }), + ), + ); +}); diff --git a/apps/server/src/provider/Layers/AntigravityAdapter.ts b/apps/server/src/provider/Layers/AntigravityAdapter.ts new file mode 100644 index 000000000000..7bd8864f55ce --- /dev/null +++ b/apps/server/src/provider/Layers/AntigravityAdapter.ts @@ -0,0 +1,550 @@ +import { + ApprovalRequestId, + type AntigravitySettings, + EventId, + type ProviderApprovalDecision, + type ProviderRuntimeEvent, + type ProviderSession, + type ProviderUserInputAnswers, + type ProviderSendTurnInput, + type ProviderSessionStartInput, + type ProviderTurnStartResult, + ProviderDriverKind, + ProviderInstanceId, + RuntimeItemId, + RuntimeRequestId, + type ThreadId, + TurnId, +} from "@t3tools/contracts"; +import * as Crypto from "effect/Crypto"; +import * as DateTime from "effect/DateTime"; +import * as Effect from "effect/Effect"; +import * as Fiber from "effect/Fiber"; +import * as Option from "effect/Option"; +import * as PubSub from "effect/PubSub"; +import * as Scope from "effect/Scope"; +import * as Stream from "effect/Stream"; +import * as SynchronizedRef from "effect/SynchronizedRef"; +import { ChildProcess, ChildProcessSpawner } from "effect/unstable/process"; +import { resolveSpawnCommand } from "@t3tools/shared/shell"; + +import { + ProviderAdapterProcessError, + ProviderAdapterRequestError, + ProviderAdapterSessionNotFoundError, +} from "../Errors.ts"; +import { type AntigravityAdapterShape } from "../Services/AntigravityAdapter.ts"; +import { + buildAntigravityPrintArgs, + normalizeAntigravityCliEvent, + normalizeAntigravityProcessExit, + parseAntigravityEffort, + parseAntigravityStreamLine, + type AntigravityEffort, +} from "../antigravity/AntigravityCliProtocol.ts"; + +const PROVIDER = ProviderDriverKind.make("antigravity"); +const ANTIGRAVITY_RESUME_VERSION = 1 as const; + +export interface AntigravityAdapterLiveOptions { + readonly environment?: NodeJS.ProcessEnv; + readonly instanceId?: ProviderInstanceId; +} + +interface ActiveProcess { + readonly process: ChildProcess.ChildProcess; + readonly fiber?: Fiber.Fiber; +} + +interface AntigravitySessionContext { + readonly threadId: ThreadId; + conversationId?: string; + session: ProviderSession; + effort?: AntigravityEffort; + launchArgs?: string; + activeTurnId?: TurnId; + activeProcess?: ActiveProcess; + interruptedTurnIds: Set; + abortEmittedTurnIds: Set; + turns: Array<{ id: TurnId; items: unknown[] }>; + stopped: boolean; +} + +const killProcess = (child: ChildProcess.ChildProcess, signal: NodeJS.Signals = "SIGTERM") => { + child + .kill({ killSignal: signal, forceKillAfter: "500 millis" }) + .pipe(Effect.ignore, Effect.runFork); +}; + +export function makeAntigravityAdapter( + settings: AntigravitySettings, + options: AntigravityAdapterLiveOptions = {}, +): Effect.Effect< + AntigravityAdapterShape, + never, + ChildProcessSpawner.ChildProcessSpawner | Crypto.Crypto | Scope.Scope +> { + return Effect.gen(function* () { + const spawner = yield* ChildProcessSpawner.ChildProcessSpawner; + const crypto = yield* Crypto.Crypto; + const adapterScope = yield* Scope.Scope; + const environment = options.environment ?? process.env; + const boundInstanceId = options.instanceId ?? ProviderInstanceId.make("antigravity"); + + const randomUUIDv4 = crypto.randomUUIDv4; + const nowIso = DateTime.now.pipe(Effect.map(DateTime.formatIso)); + const nextEventId = Effect.map(randomUUIDv4, (id) => EventId.make(`evt_${id}`)); + const makeEventStamp = () => Effect.all({ eventId: nextEventId, createdAt: nowIso }); + + const eventHub = yield* PubSub.sliding(256); + const sessionsRef = yield* SynchronizedRef.make(new Map()); + + const offerRuntimeEvent = (event: ProviderRuntimeEvent) => + PubSub.publish(eventHub, event).pipe(Effect.asVoid); + + const adapter: AntigravityAdapterShape = { + provider: PROVIDER, + capabilities: { + sessionModelSwitch: "unsupported", + }, + + startSession: (input: ProviderSessionStartInput) => + SynchronizedRef.updateEffect(sessionsRef, (map) => + Effect.gen(function* () { + const now = yield* nowIso; + const initialModel = + input.modelSelection?.model ?? settings.model ?? "gemini-3.7-flash-high"; + const effort = parseAntigravityEffort(settings.effort); + const launchArgs = settings.launchArgs?.trim() || undefined; + + const session: ProviderSession = { + provider: PROVIDER, + providerInstanceId: boundInstanceId, + status: "ready", + runtimeMode: input.runtimeMode, + cwd: input.cwd, + model: initialModel, + threadId: input.threadId, + resumeCursor: { + schemaVersion: ANTIGRAVITY_RESUME_VERSION, + conversationId: undefined, + }, + createdAt: now, + updatedAt: now, + }; + + const ctx: AntigravitySessionContext = { + threadId: input.threadId, + conversationId: undefined, + session, + effort, + launchArgs, + activeTurnId: undefined, + activeProcess: undefined, + interruptedTurnIds: new Set(), + abortEmittedTurnIds: new Set(), + turns: [], + stopped: false, + }; + + map.set(input.threadId, ctx); + return map; + }), + ).pipe( + Effect.flatMap(() => + SynchronizedRef.get(sessionsRef).pipe( + Effect.map((map) => map.get(input.threadId)!.session), + ), + ), + ), + + sendTurn: (input: ProviderSendTurnInput) => + Effect.gen(function* () { + const sessions = yield* SynchronizedRef.get(sessionsRef); + const ctx = sessions.get(input.threadId); + if (!ctx || ctx.stopped) { + return yield* Effect.fail( + new ProviderAdapterSessionNotFoundError({ + provider: PROVIDER, + threadId: input.threadId, + }), + ); + } + + const rawUuid = yield* randomUUIDv4; + const turnId = TurnId.make(`turn_${rawUuid}`); + ctx.activeTurnId = turnId; + ctx.turns.push({ id: turnId, items: [] }); + + const promptText = typeof input.input === "string" ? input.input : ""; + const printArgs = buildAntigravityPrintArgs({ + prompt: promptText, + conversationId: ctx.conversationId, + model: ctx.session.model, + effort: ctx.effort, + runtimeMode: ctx.session.runtimeMode, + interactionMode: input.interactionMode, + launchArgs: ctx.launchArgs, + }); + + const command = settings.binaryPath || "agy"; + const spawnCommand = yield* resolveSpawnCommand(command, printArgs, { env: environment }); + + let terminalResultSeen = false; + + const process = yield* spawner + .spawn( + ChildProcess.make(spawnCommand.command, spawnCommand.args, { + cwd: ctx.session.cwd, + env: environment, + shell: spawnCommand.shell, + }), + ) + .pipe(Effect.provideService(Scope.Scope, adapterScope)); + + const runFiber = Effect.gen(function* () { + const linesStream = process.stdout.pipe(Stream.decodeText(), Stream.splitLines); + + yield* Stream.runForEach(linesStream, (line) => + Effect.gen(function* () { + const interrupted = ctx.interruptedTurnIds.has(turnId); + const parsed = parseAntigravityStreamLine(line); + const signal = normalizeAntigravityCliEvent(parsed); + if (!signal) return; + + // After interrupt, ignore mid-turn noise but still surface + // terminal abort/complete if the CLI emits one first. + if ( + interrupted && + signal.type !== "turn.aborted" && + signal.type !== "turn.completed" + ) { + return; + } + if (ctx.abortEmittedTurnIds.has(turnId)) { + return; + } + + const stamp = yield* makeEventStamp(); + + switch (signal.type) { + case "session.started": { + ctx.conversationId = signal.conversationId; + ctx.session.resumeCursor = { + schemaVersion: ANTIGRAVITY_RESUME_VERSION, + conversationId: signal.conversationId, + }; + break; + } + + case "content.delta": { + yield* offerRuntimeEvent({ + eventId: stamp.eventId, + createdAt: stamp.createdAt, + provider: PROVIDER, + providerInstanceId: boundInstanceId, + threadId: ctx.threadId, + turnId, + itemId: RuntimeItemId.make(`item_${signal.stepIndex}`), + type: "content.delta", + payload: { + streamKind: "text", + delta: signal.text, + }, + } as ProviderRuntimeEvent); + break; + } + + case "tool.started": { + yield* offerRuntimeEvent({ + eventId: stamp.eventId, + createdAt: stamp.createdAt, + provider: PROVIDER, + providerInstanceId: boundInstanceId, + threadId: ctx.threadId, + turnId, + itemId: RuntimeItemId.make(`item_tool_${signal.stepIndex}`), + type: "item.started", + payload: { + itemType: "tool_use", + status: "running", + title: signal.toolName, + data: signal.parameters ?? {}, + }, + } as ProviderRuntimeEvent); + break; + } + + case "tool.completed": { + yield* offerRuntimeEvent({ + eventId: stamp.eventId, + createdAt: stamp.createdAt, + provider: PROVIDER, + providerInstanceId: boundInstanceId, + threadId: ctx.threadId, + turnId, + itemId: RuntimeItemId.make(`item_tool_${signal.stepIndex}`), + type: "item.completed", + payload: { + itemType: "tool_use", + status: "completed", + title: signal.toolName, + data: { + output: signal.output, + durationSeconds: signal.durationSeconds, + }, + }, + } as ProviderRuntimeEvent); + break; + } + + case "tool.failed": { + yield* offerRuntimeEvent({ + eventId: stamp.eventId, + createdAt: stamp.createdAt, + provider: PROVIDER, + providerInstanceId: boundInstanceId, + threadId: ctx.threadId, + turnId, + itemId: RuntimeItemId.make(`item_tool_${signal.stepIndex}`), + type: "item.completed", + payload: { + itemType: "tool_use", + status: "failed", + title: signal.toolName, + data: { + output: signal.output, + error: signal.error, + status: signal.status, + durationSeconds: signal.durationSeconds, + }, + }, + } as ProviderRuntimeEvent); + break; + } + + case "turn.completed": { + terminalResultSeen = true; + ctx.abortEmittedTurnIds.add(turnId); + yield* offerRuntimeEvent({ + eventId: stamp.eventId, + createdAt: stamp.createdAt, + provider: PROVIDER, + providerInstanceId: boundInstanceId, + threadId: ctx.threadId, + turnId, + type: "turn.completed", + payload: { + state: "completed", + usage: signal.usage, + }, + } as ProviderRuntimeEvent); + break; + } + + case "turn.aborted": { + terminalResultSeen = true; + ctx.abortEmittedTurnIds.add(turnId); + yield* offerRuntimeEvent({ + eventId: stamp.eventId, + createdAt: stamp.createdAt, + provider: PROVIDER, + providerInstanceId: boundInstanceId, + threadId: ctx.threadId, + turnId, + type: "turn.aborted", + payload: { + reason: signal.response || `Turn aborted with status ${signal.status}`, + }, + } as ProviderRuntimeEvent); + break; + } + } + }), + ); + + const exitStatus = yield* process.exitCode; + if (!terminalResultSeen && !ctx.abortEmittedTurnIds.has(turnId)) { + if (ctx.interruptedTurnIds.has(turnId)) { + ctx.abortEmittedTurnIds.add(turnId); + const stamp = yield* makeEventStamp(); + yield* offerRuntimeEvent({ + eventId: stamp.eventId, + createdAt: stamp.createdAt, + provider: PROVIDER, + providerInstanceId: boundInstanceId, + threadId: ctx.threadId, + turnId, + type: "turn.aborted", + payload: { + reason: "Turn interrupted", + }, + } as ProviderRuntimeEvent); + } else { + const exitSignal = normalizeAntigravityProcessExit({ + exitCode: typeof exitStatus === "number" ? exitStatus : null, + signal: null, + terminalResultSeen, + }); + if (exitSignal) { + ctx.abortEmittedTurnIds.add(turnId); + const stamp = yield* makeEventStamp(); + yield* offerRuntimeEvent({ + eventId: stamp.eventId, + createdAt: stamp.createdAt, + provider: PROVIDER, + providerInstanceId: boundInstanceId, + threadId: ctx.threadId, + turnId, + type: "turn.aborted", + payload: { + reason: exitSignal.response || "Process exited unexpectedly", + }, + } as ProviderRuntimeEvent); + } + } + } + }).pipe(Effect.forkIn(adapterScope)); + + const fiber = yield* runFiber; + ctx.activeProcess = { process, fiber }; + + return { + threadId: input.threadId, + turnId, + resumeCursor: ctx.session.resumeCursor, + } satisfies ProviderTurnStartResult; + }), + + interruptTurn: (threadId: ThreadId, turnId?: TurnId) => + Effect.gen(function* () { + const sessions = yield* SynchronizedRef.get(sessionsRef); + const ctx = sessions.get(threadId); + if (!ctx) return; + + const targetTurnId = turnId ?? ctx.activeTurnId; + if (targetTurnId) { + ctx.interruptedTurnIds.add(targetTurnId); + } + + if (ctx.activeProcess) { + killProcess(ctx.activeProcess.process, "SIGINT"); + if (targetTurnId && !ctx.abortEmittedTurnIds.has(targetTurnId)) { + ctx.abortEmittedTurnIds.add(targetTurnId); + const stamp = yield* makeEventStamp(); + yield* offerRuntimeEvent({ + eventId: stamp.eventId, + createdAt: stamp.createdAt, + provider: PROVIDER, + providerInstanceId: boundInstanceId, + threadId: ctx.threadId, + turnId: targetTurnId, + type: "turn.aborted", + payload: { + reason: "Turn interrupted via SIGINT", + }, + } as ProviderRuntimeEvent); + } + ctx.activeProcess = undefined; + } + }), + + // UA-owned approval: turns launch with --dangerously-skip-permissions, so + // Agy never emits mid-turn permission requests. These remain no-ops; the + // UA/T3 approval surface must gate before sendTurn instead. + respondToRequest: ( + _threadId: ThreadId, + _requestId: ApprovalRequestId, + _decision: ProviderApprovalDecision, + ) => Effect.void, + + respondToUserInput: ( + _threadId: ThreadId, + _requestId: ApprovalRequestId, + _answers: ProviderUserInputAnswers, + ) => Effect.void, + + stopSession: (threadId: ThreadId) => + Effect.gen(function* () { + const sessions = yield* SynchronizedRef.get(sessionsRef); + const ctx = sessions.get(threadId); + if (ctx) { + ctx.stopped = true; + if (ctx.activeProcess) { + killProcess(ctx.activeProcess.process, "SIGTERM"); + ctx.activeProcess = undefined; + } + yield* SynchronizedRef.update(sessionsRef, (map) => { + map.delete(threadId); + return map; + }); + } + }), + + listSessions: () => + SynchronizedRef.get(sessionsRef).pipe( + Effect.map((map) => Array.from(map.values()).map((c) => c.session)), + ), + + hasSession: (threadId: ThreadId) => + SynchronizedRef.get(sessionsRef).pipe( + Effect.map((map) => map.has(threadId) && !map.get(threadId)!.stopped), + ), + + readThread: (threadId: ThreadId) => + Effect.gen(function* () { + const sessions = yield* SynchronizedRef.get(sessionsRef); + const ctx = sessions.get(threadId); + if (!ctx) { + return yield* Effect.fail( + new ProviderAdapterSessionNotFoundError({ + provider: PROVIDER, + threadId, + }), + ); + } + return { + threadId, + turns: ctx.turns.map((t) => ({ id: t.id, items: t.items })), + }; + }), + + rollbackThread: (threadId: ThreadId, numTurns: number) => + Effect.gen(function* () { + const sessions = yield* SynchronizedRef.get(sessionsRef); + const ctx = sessions.get(threadId); + if (!ctx) { + return yield* Effect.fail( + new ProviderAdapterSessionNotFoundError({ + provider: PROVIDER, + threadId, + }), + ); + } + ctx.turns = ctx.turns.slice(0, Math.max(0, ctx.turns.length - numTurns)); + return { + threadId, + turns: ctx.turns.map((t) => ({ id: t.id, items: t.items })), + }; + }), + + stopAll: () => + SynchronizedRef.update(sessionsRef, (map) => { + for (const ctx of map.values()) { + ctx.stopped = true; + if (ctx.activeProcess) { + ctx.activeProcess.process + .kill({ killSignal: "SIGTERM", forceKillAfter: "500 millis" }) + .pipe(Effect.ignore, Effect.runFork); + } + } + map.clear(); + return map; + }), + + streamEvents: Stream.fromPubSub(eventHub), + }; + + return adapter; + }); +} diff --git a/apps/server/src/provider/Layers/AntigravityProvider.test.ts b/apps/server/src/provider/Layers/AntigravityProvider.test.ts new file mode 100644 index 000000000000..5e978164518f --- /dev/null +++ b/apps/server/src/provider/Layers/AntigravityProvider.test.ts @@ -0,0 +1,129 @@ +import * as NodeServices from "@effect/platform-node/NodeServices"; +import { describe, expect, it } from "@effect/vitest"; +import * as Effect from "effect/Effect"; +import * as FileSystem from "effect/FileSystem"; +import * as Path from "effect/Path"; +import * as Schema from "effect/Schema"; +import { AntigravitySettings } from "@t3tools/contracts"; + +import { + buildInitialAntigravityProviderSnapshot, + checkAntigravityProviderStatus, +} from "./AntigravityProvider.ts"; + +const decodeAntigravitySettings = Schema.decodeSync(AntigravitySettings); + +describe("buildInitialAntigravityProviderSnapshot", () => { + it.effect("returns a disabled snapshot when settings.enabled is false", () => + Effect.gen(function* () { + const snapshot = yield* buildInitialAntigravityProviderSnapshot( + decodeAntigravitySettings({ enabled: false }), + ); + expect(snapshot.enabled).toBe(false); + expect(snapshot.status).toBe("disabled"); + expect(snapshot.installed).toBe(false); + expect(snapshot.message).toContain("disabled"); + }), + ); + + it.effect("returns an initial pending snapshot when enabled", () => + Effect.gen(function* () { + const snapshot = yield* buildInitialAntigravityProviderSnapshot( + decodeAntigravitySettings({ enabled: true }), + ); + expect(snapshot.enabled).toBe(true); + expect(snapshot.installed).toBe(true); + expect(snapshot.status).toBe("warning"); + expect(snapshot.version).toBeNull(); + expect(snapshot.message).toContain("Checking Antigravity"); + expect(snapshot.requiresNewThreadForModelChange).toBe(true); + }), + ); +}); + +it.layer(NodeServices.layer)("checkAntigravityProviderStatus", (it) => { + it.effect("reports the binary as missing when the binary path does not resolve", () => + Effect.gen(function* () { + const snapshot = yield* checkAntigravityProviderStatus( + decodeAntigravitySettings({ + enabled: true, + binaryPath: "/definitely/not/installed/agy-missing-binary", + }), + ); + expect(snapshot.enabled).toBe(true); + expect(snapshot.installed).toBe(false); + expect(snapshot.status).toBe("error"); + expect(snapshot.message).toMatch(/not found in PATH|Failed to execute/); + }), + ); + + it.effect("reports a ready provider when agy version and models probes succeed", () => + Effect.gen(function* () { + const snapshot = yield* Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const dir = yield* fs.makeTempDirectoryScoped({ prefix: "t3code-agy-mock-" }); + const agyPath = path.join(dir, "agy"); + yield* fs.writeFileString( + agyPath, + [ + "#!/bin/sh", + 'if [ "$1" = "--version" ]; then echo "1.1.22"; exit 0; fi', + 'if [ "$1" = "models" ]; then echo "gemini-3.7-flash-high"; exit 0; fi', + "exit 1", + "", + ].join("\n"), + ); + yield* fs.chmod(agyPath, 0o755); + + return yield* checkAntigravityProviderStatus( + decodeAntigravitySettings({ enabled: true, binaryPath: agyPath }), + ); + }), + ); + + expect(snapshot.enabled).toBe(true); + expect(snapshot.installed).toBe(true); + expect(snapshot.status).toBe("ready"); + expect(snapshot.version).toBe("1.1.22"); + expect(snapshot.auth?.status).toBe("authenticated"); + expect(snapshot.models.length).toBeGreaterThan(0); + }), + ); + + it.effect("reports unauthenticated when agy models fails after a good --version", () => + Effect.gen(function* () { + const snapshot = yield* Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const dir = yield* fs.makeTempDirectoryScoped({ prefix: "t3code-agy-unauth-" }); + const agyPath = path.join(dir, "agy"); + yield* fs.writeFileString( + agyPath, + [ + "#!/bin/sh", + 'if [ "$1" = "--version" ]; then echo "1.1.22"; exit 0; fi', + 'if [ "$1" = "models" ]; then echo "authentication required" >&2; exit 1; fi', + "exit 1", + "", + ].join("\n"), + ); + yield* fs.chmod(agyPath, 0o755); + + return yield* checkAntigravityProviderStatus( + decodeAntigravitySettings({ enabled: true, binaryPath: agyPath }), + ); + }), + ); + + expect(snapshot.enabled).toBe(true); + expect(snapshot.installed).toBe(true); + expect(snapshot.status).toBe("warning"); + expect(snapshot.version).toBe("1.1.22"); + expect(snapshot.auth?.status).toBe("unauthenticated"); + expect(snapshot.message).toMatch(/authentication required/i); + }), + ); +}); diff --git a/apps/server/src/provider/Layers/AntigravityProvider.ts b/apps/server/src/provider/Layers/AntigravityProvider.ts new file mode 100644 index 000000000000..fc9aacd282d1 --- /dev/null +++ b/apps/server/src/provider/Layers/AntigravityProvider.ts @@ -0,0 +1,315 @@ +import { + type AntigravitySettings, + type ModelCapabilities, + type ServerProviderModel, +} from "@t3tools/contracts"; +import * as DateTime from "effect/DateTime"; +import * as Effect from "effect/Effect"; +import * as Result from "effect/Result"; +import { ChildProcess, ChildProcessSpawner } from "effect/unstable/process"; +import { createModelCapabilities } from "@t3tools/shared/model"; +import { resolveSpawnCommand } from "@t3tools/shared/shell"; + +import { + AUTH_PROBE_TIMEOUT_MS, + buildServerProvider, + isCommandMissingCause, + parseGenericCliVersion, + providerModelsFromSettings, + spawnAndCollect, + type ServerProviderDraft, +} from "../providerSnapshot.ts"; +import { + classifyAntigravityAvailability, + type AntigravityProbeCommandResult, +} from "../antigravity/AntigravityCliProtocol.ts"; + +const ANTIGRAVITY_PRESENTATION = { + displayName: "Antigravity", + badgeLabel: "Native", + showInteractionModeToggle: false, + requiresNewThreadForModelChange: true, +} as const; + +const EMPTY_CAPABILITIES: ModelCapabilities = createModelCapabilities({ + optionDescriptors: [], +}); + +const VERSION_PROBE_TIMEOUT_MS = 4_000; + +export const ANTIGRAVITY_BUILT_IN_MODELS: ReadonlyArray = [ + { + slug: "gemini-3.7-flash-high", + name: "Gemini 3.7 Flash (High)", + isCustom: false, + capabilities: EMPTY_CAPABILITIES, + }, + { + slug: "gemini-3.7-flash-medium", + name: "Gemini 3.7 Flash (Medium)", + isCustom: false, + capabilities: EMPTY_CAPABILITIES, + }, + { + slug: "gemini-3.7-flash-low", + name: "Gemini 3.7 Flash (Low)", + isCustom: false, + capabilities: EMPTY_CAPABILITIES, + }, + { + slug: "gemini-3.1-pro-high", + name: "Gemini 3.1 Pro (High)", + isCustom: false, + capabilities: EMPTY_CAPABILITIES, + }, + { + slug: "claude-sonnet-4.6-thinking", + name: "Claude Sonnet 4.6 (Thinking)", + isCustom: false, + capabilities: EMPTY_CAPABILITIES, + }, + { + slug: "claude-opus-4.6-thinking", + name: "Claude Opus 4.6 (Thinking)", + isCustom: false, + capabilities: EMPTY_CAPABILITIES, + }, + { + slug: "gpt-oss-120b-medium", + name: "GPT-OSS 120B (Medium)", + isCustom: false, + capabilities: EMPTY_CAPABILITIES, + }, +]; + +export function buildInitialAntigravityProviderSnapshot( + antigravitySettings: AntigravitySettings, +): Effect.Effect { + return Effect.gen(function* () { + const checkedAt = yield* Effect.map(DateTime.now, DateTime.formatIso); + const models = antigravityModelsFromSettings(antigravitySettings.customModels); + + if (!antigravitySettings.enabled) { + return buildServerProvider({ + presentation: ANTIGRAVITY_PRESENTATION, + enabled: false, + checkedAt, + models, + probe: { + installed: false, + version: null, + status: "warning", + auth: { status: "unknown" }, + message: "Antigravity is disabled in T3 Code settings.", + }, + }); + } + + return buildServerProvider({ + presentation: ANTIGRAVITY_PRESENTATION, + enabled: true, + checkedAt, + models, + probe: { + installed: true, + version: null, + status: "warning", + auth: { status: "unknown" }, + message: "Checking Antigravity (agy) CLI availability...", + }, + }); + }); +} + +function antigravityModelsFromSettings( + customModels: ReadonlyArray | undefined, + builtInModels: ReadonlyArray = ANTIGRAVITY_BUILT_IN_MODELS, +): ReadonlyArray { + return providerModelsFromSettings(builtInModels, customModels ?? [], EMPTY_CAPABILITIES); +} + +const runAntigravityCliCommand = ( + antigravitySettings: AntigravitySettings, + args: ReadonlyArray, + environment: NodeJS.ProcessEnv = process.env, +) => + Effect.gen(function* () { + const command = antigravitySettings.binaryPath || "agy"; + const spawnCommand = yield* resolveSpawnCommand(command, [...args], { + env: environment, + }); + return yield* spawnAndCollect( + command, + ChildProcess.make(spawnCommand.command, spawnCommand.args, { + env: environment, + shell: spawnCommand.shell, + }), + ); + }); + +const runAntigravityVersionCommand = ( + antigravitySettings: AntigravitySettings, + environment: NodeJS.ProcessEnv = process.env, +) => runAntigravityCliCommand(antigravitySettings, ["--version"], environment); + +const runAntigravityModelsCommand = ( + antigravitySettings: AntigravitySettings, + environment: NodeJS.ProcessEnv = process.env, +) => runAntigravityCliCommand(antigravitySettings, ["models"], environment); + +export const checkAntigravityProviderStatus = Effect.fn("checkAntigravityProviderStatus")( + function* ( + antigravitySettings: AntigravitySettings, + environment: NodeJS.ProcessEnv = process.env, + ): Effect.fn.Return { + const checkedAt = DateTime.formatIso(yield* DateTime.now); + const fallbackModels = antigravityModelsFromSettings(antigravitySettings.customModels); + + if (!antigravitySettings.enabled) { + return buildServerProvider({ + presentation: ANTIGRAVITY_PRESENTATION, + enabled: false, + checkedAt, + models: fallbackModels, + probe: { + installed: false, + version: null, + status: "warning", + auth: { status: "unknown" }, + message: "Antigravity is disabled in T3 Code settings.", + }, + }); + } + + const versionResult = yield* runAntigravityVersionCommand( + antigravitySettings, + environment, + ).pipe(Effect.timeoutOption(VERSION_PROBE_TIMEOUT_MS), Effect.result); + + if (Result.isFailure(versionResult)) { + const error = versionResult.failure; + const missing = isCommandMissingCause(error); + return buildServerProvider({ + presentation: ANTIGRAVITY_PRESENTATION, + enabled: antigravitySettings.enabled, + checkedAt, + models: fallbackModels, + probe: { + installed: false, + version: null, + status: "error", + auth: { status: "unknown" }, + message: missing + ? `Antigravity CLI (${antigravitySettings.binaryPath || "agy"}) was not found in PATH.` + : "Failed to probe Antigravity CLI version.", + }, + }); + } + + const versionCollected = versionResult.success; + if (versionCollected._tag === "None") { + return buildServerProvider({ + presentation: ANTIGRAVITY_PRESENTATION, + enabled: antigravitySettings.enabled, + checkedAt, + models: fallbackModels, + probe: { + installed: true, + version: null, + status: "warning", + auth: { status: "unknown" }, + message: "Antigravity CLI probe timed out.", + }, + }); + } + + const versionOutput = + versionCollected.value.stdout.trim() || versionCollected.value.stderr.trim(); + const parsedVersion = parseGenericCliVersion(versionOutput); + + const versionProbeResult: AntigravityProbeCommandResult = { + exitCode: versionCollected.value.code, + stdout: versionCollected.value.stdout, + stderr: versionCollected.value.stderr, + }; + + // Auth is exit-code-only on `agy models`. Successful probes may still print + // diagnostics on stderr — do not treat stderr as failure by itself. + const modelsResult = yield* runAntigravityModelsCommand(antigravitySettings, environment).pipe( + Effect.timeoutOption(AUTH_PROBE_TIMEOUT_MS), + Effect.result, + ); + + let modelsProbe: AntigravityProbeCommandResult | null = null; + if (Result.isSuccess(modelsResult) && modelsResult.success._tag === "Some") { + modelsProbe = { + exitCode: modelsResult.success.value.code, + stdout: modelsResult.success.value.stdout, + stderr: modelsResult.success.value.stderr, + }; + } else if (Result.isSuccess(modelsResult) && modelsResult.success._tag === "None") { + modelsProbe = { + exitCode: 1, + stdout: "", + stderr: "agy models authentication probe timed out", + }; + } else { + modelsProbe = { + exitCode: 1, + stdout: "", + stderr: "agy models authentication probe failed to execute", + }; + } + + const availability = classifyAntigravityAvailability({ + version: versionProbeResult, + models: modelsProbe, + }); + + if (availability.status === "unavailable") { + return buildServerProvider({ + presentation: ANTIGRAVITY_PRESENTATION, + enabled: antigravitySettings.enabled, + checkedAt, + models: fallbackModels, + probe: { + installed: false, + version: null, + status: "error", + auth: { status: "unknown" }, + message: availability.reason, + }, + }); + } + + if (availability.status === "unauthenticated") { + return buildServerProvider({ + presentation: ANTIGRAVITY_PRESENTATION, + enabled: antigravitySettings.enabled, + checkedAt, + models: fallbackModels, + probe: { + installed: true, + version: parsedVersion, + status: "warning", + auth: { status: "unauthenticated" }, + message: availability.reason, + }, + }); + } + + return buildServerProvider({ + presentation: ANTIGRAVITY_PRESENTATION, + enabled: antigravitySettings.enabled, + checkedAt, + models: fallbackModels, + probe: { + installed: true, + version: parsedVersion, + status: "ready", + auth: { status: "authenticated" }, + message: `Antigravity CLI ${parsedVersion || "available"} is ready.`, + }, + }); + }, +); diff --git a/apps/server/src/provider/Layers/ProviderRegistry.test.ts b/apps/server/src/provider/Layers/ProviderRegistry.test.ts index 68e88da51576..6b8a761011b1 100644 --- a/apps/server/src/provider/Layers/ProviderRegistry.test.ts +++ b/apps/server/src/provider/Layers/ProviderRegistry.test.ts @@ -1489,6 +1489,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te deepMerge(encodedDefaultServerSettings, { providers: { codex: { enabled: true, binaryPath: firstMissing }, + antigravity: { enabled: false }, claudeAgent: { enabled: false }, cursor: { enabled: false }, grok: { enabled: false }, @@ -1712,9 +1713,9 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te code: 0, }; } - if (joined === "auth status") { + if (joined === "auth status" || joined === "models") { return { - stdout: '{"authenticated":true}\n', + stdout: joined === "models" ? "model-a\n" : '{"authenticated":true}\n', stderr: "", code: 0, }; @@ -1738,6 +1739,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te ); assert.deepStrictEqual(providers.map((provider) => provider.instanceId).toSorted(), [ + "antigravity", "claudeAgent", "codex", "cursor", diff --git a/apps/server/src/provider/Services/AntigravityAdapter.ts b/apps/server/src/provider/Services/AntigravityAdapter.ts new file mode 100644 index 000000000000..29189c2c010d --- /dev/null +++ b/apps/server/src/provider/Services/AntigravityAdapter.ts @@ -0,0 +1,12 @@ +/** + * AntigravityAdapter — shape type for the Antigravity provider adapter. + * + * @module AntigravityAdapter + */ +import type { ProviderAdapterError } from "../Errors.ts"; +import type { ProviderAdapterShape } from "./ProviderAdapter.ts"; + +/** + * AntigravityAdapterShape — per-instance Antigravity adapter contract. + */ +export interface AntigravityAdapterShape extends ProviderAdapterShape {} diff --git a/apps/server/src/provider/antigravity/AntigravityCliProtocol.test.ts b/apps/server/src/provider/antigravity/AntigravityCliProtocol.test.ts index d8e04c70185a..43f7d3e1d400 100644 --- a/apps/server/src/provider/antigravity/AntigravityCliProtocol.test.ts +++ b/apps/server/src/provider/antigravity/AntigravityCliProtocol.test.ts @@ -7,16 +7,19 @@ import { classifyAntigravityAvailability, normalizeAntigravityCliEvent, normalizeAntigravityProcessExit, + parseAntigravityEffort, parseAntigravityStreamLine, + resolveAntigravityCliMode, } from "./AntigravityCliProtocol.ts"; describe("buildAntigravityPrintArgs", () => { - it("builds a bounded read-only stream-json launch", () => { + it("builds a UA-owned-approval stream-json launch without fake sandbox", () => { NodeAssert.deepStrictEqual( buildAntigravityPrintArgs({ prompt: "Inspect the repository without changing files.", model: "gemini-3.7-flash-low", effort: "low", + interactionMode: "plan", }), [ "--print", @@ -25,7 +28,7 @@ describe("buildAntigravityPrintArgs", () => { "stream-json", "--mode", "plan", - "--sandbox", + "--dangerously-skip-permissions", "--model", "gemini-3.7-flash-low", "--effort", @@ -34,11 +37,13 @@ describe("buildAntigravityPrintArgs", () => { ); }); - it("resumes an existing conversation without weakening safeguards", () => { + it("maps default interaction to accept-edits and appends launchArgs", () => { NodeAssert.deepStrictEqual( buildAntigravityPrintArgs({ prompt: "Continue the inspection.", conversationId: "conversation-123", + runtimeMode: "full-access", + launchArgs: "--sandbox --add-dir /tmp/extra", }), [ "--print", @@ -46,13 +51,26 @@ describe("buildAntigravityPrintArgs", () => { "--output-format", "stream-json", "--mode", - "plan", - "--sandbox", + "accept-edits", + "--dangerously-skip-permissions", "--conversation", "conversation-123", + "--sandbox", + "--add-dir", + "/tmp/extra", ], ); }); + + it("drops invalid effort literals", () => { + NodeAssert.equal(parseAntigravityEffort("ultrathink"), undefined); + NodeAssert.equal(parseAntigravityEffort("low"), "low"); + NodeAssert.equal(resolveAntigravityCliMode({ interactionMode: "plan" }), "plan"); + NodeAssert.equal( + resolveAntigravityCliMode({ runtimeMode: "approval-required" }), + "accept-edits", + ); + }); }); describe("Antigravity stream-json protocol", () => { @@ -99,7 +117,82 @@ describe("Antigravity stream-json protocol", () => { }); }); - it("normalizes terminal results", () => { + it("normalizes tool step execution events including ERROR", () => { + const activeEvent = parseAntigravityStreamLine( + JSON.stringify({ + event: "step_update", + step_update: { + step_index: 2, + step_type: "tool", + state: "ACTIVE", + tool_name: "run_command", + tool_info: { + name: "run_command", + parameters: { CommandLine: "echo PROBE_TEST" }, + }, + }, + }), + ); + + NodeAssert.deepStrictEqual(normalizeAntigravityCliEvent(activeEvent), { + type: "tool.started", + stepIndex: 2, + toolName: "run_command", + parameters: { CommandLine: "echo PROBE_TEST" }, + }); + + const doneEvent = parseAntigravityStreamLine( + JSON.stringify({ + event: "step_update", + step_update: { + step_index: 2, + step_type: "tool", + state: "DONE", + tool_name: "run_command", + duration_seconds: 0.189, + tool_info: { + name: "run_command", + parameters: { CommandLine: "echo PROBE_TEST" }, + output: "PROBE_TEST\r\n", + }, + }, + }), + ); + + NodeAssert.deepStrictEqual(normalizeAntigravityCliEvent(doneEvent), { + type: "tool.completed", + stepIndex: 2, + toolName: "run_command", + durationSeconds: 0.189, + output: "PROBE_TEST\r\n", + }); + + const errorEvent = parseAntigravityStreamLine( + JSON.stringify({ + event: "step_update", + step_update: { + step_index: 3, + step_type: "tool", + state: "ERROR", + tool_name: "run_command", + tool_info: { + name: "run_command", + error: "permission denied", + }, + }, + }), + ); + + NodeAssert.deepStrictEqual(normalizeAntigravityCliEvent(errorEvent), { + type: "tool.failed", + stepIndex: 3, + toolName: "run_command", + status: "ERROR", + error: "permission denied", + }); + }); + + it("normalizes terminal results and CANCELED taxonomy", () => { const event = parseAntigravityStreamLine( JSON.stringify({ event: "result", @@ -121,6 +214,22 @@ describe("Antigravity stream-json protocol", () => { numTurns: 1, usage: { input_tokens: 12, output_tokens: 3 }, }); + + const canceled = parseAntigravityStreamLine( + JSON.stringify({ + event: "result", + result: { + conversation_id: "conversation-123", + status: "CANCELED", + response: "user canceled", + }, + }), + ); + NodeAssert.deepStrictEqual(normalizeAntigravityCliEvent(canceled), { + type: "turn.aborted", + status: "CANCELLED", + response: "user canceled", + }); }); it("ignores blank, malformed, and non-assistant step lines", () => { @@ -157,6 +266,24 @@ describe("Antigravity stream-json protocol", () => { status: "FAILED", response: "Provider failed", }); + + const errorEvent = parseAntigravityStreamLine( + JSON.stringify({ + event: "result", + result: { + conversation_id: "", + status: "ERROR", + response: "", + error: "context canceled", + }, + }), + ); + + NodeAssert.deepStrictEqual(normalizeAntigravityCliEvent(errorEvent), { + type: "turn.aborted", + status: "ERROR", + response: "context canceled", + }); }); }); @@ -185,7 +312,7 @@ describe("Antigravity capability probes", () => { NodeAssert.deepStrictEqual( classifyAntigravityAvailability({ version: { exitCode: 0, stdout: "1.1.22\n", stderr: "" }, - models: { exitCode: 0, stdout: "models", stderr: "" }, + models: { exitCode: 0, stdout: "models", stderr: "note: cached" }, }), { status: "available", version: "1.1.22" }, ); diff --git a/apps/server/src/provider/antigravity/AntigravityCliProtocol.ts b/apps/server/src/provider/antigravity/AntigravityCliProtocol.ts index 396096d13366..6181db6407c7 100644 --- a/apps/server/src/provider/antigravity/AntigravityCliProtocol.ts +++ b/apps/server/src/provider/antigravity/AntigravityCliProtocol.ts @@ -1,10 +1,16 @@ +import type { ProviderInteractionMode, RuntimeMode } from "@t3tools/contracts"; +import { tokenizeCliArgs } from "@t3tools/shared/cliArgs"; + export type AntigravityEffort = "low" | "medium" | "high"; export interface AntigravityPrintOptions { readonly prompt: string; readonly conversationId?: string; readonly model?: string; - readonly effort?: AntigravityEffort; + readonly effort?: string; + readonly runtimeMode?: RuntimeMode; + readonly interactionMode?: ProviderInteractionMode; + readonly launchArgs?: string; } export interface AntigravityProbeCommandResult { @@ -48,6 +54,7 @@ interface AntigravityResultPayload { readonly conversation_id: string; readonly status: string; readonly response?: string; + readonly error?: string; readonly duration_seconds?: number; readonly num_turns?: number; readonly usage?: Readonly>; @@ -76,6 +83,28 @@ export type AntigravityProviderSignal = readonly stepIndex: number; readonly text: string; } + | { + readonly type: "tool.started"; + readonly stepIndex: number; + readonly toolName: string; + readonly parameters?: Readonly>; + } + | { + readonly type: "tool.completed"; + readonly stepIndex: number; + readonly toolName: string; + readonly durationSeconds?: number; + readonly output?: string; + } + | { + readonly type: "tool.failed"; + readonly stepIndex: number; + readonly toolName: string; + readonly status: string; + readonly durationSeconds?: number; + readonly output?: string; + readonly error?: string; + } | { readonly type: "turn.completed"; readonly response: string; @@ -92,17 +121,59 @@ export type AntigravityProviderSignal = const isRecord = (value: unknown): value is Record => typeof value === "object" && value !== null && !Array.isArray(value); +const TOOL_FAILED_STATES = new Set(["ERROR", "FAILED", "CANCELED", "CANCELLED", "ABORTED"]); + +/** + * Parse effort settings. Unknown / empty values are dropped — do not cast + * arbitrary strings into the CLI `--effort` flag. + */ +export const parseAntigravityEffort = ( + value: string | undefined, +): AntigravityEffort | undefined => { + const trimmed = value?.trim().toLowerCase(); + if (trimmed === "low" || trimmed === "medium" || trimmed === "high") { + return trimmed; + } + return undefined; +}; + +/** + * Map T3 session/turn modes onto agy `--mode`. + * + * Agy only exposes `plan` | `accept-edits`. T3 `RuntimeMode` is an approval + * posture, not an Agy mode — UA owns approval via `--dangerously-skip-permissions`. + * Prefer `interactionMode: "plan"` for plan; otherwise accept-edits. + * + * `--sandbox` is intentionally NOT hardcoded: it is a terminal restriction + * only (see #642), not a capability boundary. Pass it via `launchArgs` if wanted. + */ +export const resolveAntigravityCliMode = (input: { + readonly runtimeMode?: RuntimeMode; + readonly interactionMode?: ProviderInteractionMode; +}): "plan" | "accept-edits" => { + void input.runtimeMode; + return input.interactionMode === "plan" ? "plan" : "accept-edits"; +}; + export const buildAntigravityPrintArgs = ( options: AntigravityPrintOptions, ): ReadonlyArray => { + const mode = resolveAntigravityCliMode({ + runtimeMode: options.runtimeMode, + interactionMode: options.interactionMode, + }); + const effort = parseAntigravityEffort(options.effort); + const args = [ "--print", options.prompt, "--output-format", "stream-json", "--mode", - "plan", - "--sandbox", + mode, + // UA-owned approval: suppress mid-turn Agy permission prompts. T3/UA must + // gate before sendTurn; respondToRequest is a no-op by design. + "--dangerously-skip-permissions", ]; if (options.conversationId !== undefined) { @@ -111,8 +182,13 @@ export const buildAntigravityPrintArgs = ( if (options.model !== undefined) { args.push("--model", options.model); } - if (options.effort !== undefined) { - args.push("--effort", options.effort); + if (effort !== undefined) { + args.push("--effort", effort); + } + + const extra = tokenizeCliArgs(options.launchArgs); + if (extra.length > 0) { + args.push(...extra); } return args; @@ -175,6 +251,14 @@ export const parseAntigravityStreamLine = (line: string): AntigravityCliEvent | return null; }; +const normalizeAbortStatus = (status: string): string => { + const upper = status.toUpperCase(); + if (upper === "CANCELED" || upper === "CANCELLED" || upper === "ABORTED") { + return "CANCELLED"; + } + return status; +}; + export const normalizeAntigravityCliEvent = ( event: AntigravityCliEvent | null, ): AntigravityProviderSignal | null => { @@ -195,17 +279,74 @@ export const normalizeAntigravityCliEvent = ( if (event.event === "step_update") { const stepType = event.step_update.step_type; - const stepIndex = event.step_update.step_index; - const text = event.step_update.text_delta; - if ( - stepType !== "agent_response" || - typeof stepIndex !== "number" || - typeof text !== "string" || - text.length === 0 - ) { + const stepIndex = + typeof event.step_update.step_index === "number" ? event.step_update.step_index : 0; + const state = event.step_update.state; + + if (stepType === "agent_response") { + const text = event.step_update.text_delta; + if (typeof text === "string" && text.length > 0) { + return { type: "content.delta", stepIndex, text }; + } return null; } - return { type: "content.delta", stepIndex, text }; + + if (stepType === "tool") { + const toolInfo = isRecord(event.step_update.tool_info) ? event.step_update.tool_info : {}; + const toolName = + typeof event.step_update.tool_name === "string" + ? event.step_update.tool_name + : typeof toolInfo.name === "string" + ? toolInfo.name + : "unknown"; + + const parameters = isRecord(toolInfo.parameters) ? toolInfo.parameters : undefined; + const output = typeof toolInfo.output === "string" ? toolInfo.output : undefined; + const error = + typeof toolInfo.error === "string" + ? toolInfo.error + : typeof event.step_update.error === "string" + ? event.step_update.error + : undefined; + const durationSeconds = + typeof event.step_update.duration_seconds === "number" + ? event.step_update.duration_seconds + : undefined; + const stateLabel = typeof state === "string" ? state.toUpperCase() : ""; + + if (stateLabel === "ACTIVE") { + return { + type: "tool.started", + stepIndex, + toolName, + ...(parameters ? { parameters } : {}), + }; + } + + if (stateLabel === "DONE") { + return { + type: "tool.completed", + stepIndex, + toolName, + ...(durationSeconds !== undefined ? { durationSeconds } : {}), + ...(output !== undefined ? { output } : {}), + }; + } + + if (TOOL_FAILED_STATES.has(stateLabel)) { + return { + type: "tool.failed", + stepIndex, + toolName, + status: normalizeAbortStatus(stateLabel), + ...(durationSeconds !== undefined ? { durationSeconds } : {}), + ...(output !== undefined ? { output } : {}), + ...(error !== undefined ? { error } : {}), + }; + } + } + + return null; } if (event.result.status === "SUCCESS") { @@ -222,8 +363,8 @@ export const normalizeAntigravityCliEvent = ( return { type: "turn.aborted", - status: event.result.status, - ...(event.result.response !== undefined ? { response: event.result.response } : {}), + status: normalizeAbortStatus(event.result.status), + response: event.result.error || event.result.response || "", }; }; diff --git a/apps/server/src/provider/builtInDrivers.ts b/apps/server/src/provider/builtInDrivers.ts index 791a96e1da3c..bff24f88c237 100644 --- a/apps/server/src/provider/builtInDrivers.ts +++ b/apps/server/src/provider/builtInDrivers.ts @@ -20,6 +20,7 @@ * * @module provider/builtInDrivers */ +import { AntigravityDriver, type AntigravityDriverEnv } from "./Drivers/AntigravityDriver.ts"; import { ClaudeDriver, type ClaudeDriverEnv } from "./Drivers/ClaudeDriver.ts"; import { CodexDriver, type CodexDriverEnv } from "./Drivers/CodexDriver.ts"; import { CursorDriver, type CursorDriverEnv } from "./Drivers/CursorDriver.ts"; @@ -33,6 +34,7 @@ import type { AnyProviderDriver } from "./ProviderDriver.ts"; * layer must provide every service in this union. */ export type BuiltInDriversEnv = + | AntigravityDriverEnv | ClaudeDriverEnv | CodexDriverEnv | CursorDriverEnv @@ -45,6 +47,7 @@ export type BuiltInDriversEnv = * iteration order has no functional effect on instance lookup. */ export const BUILT_IN_DRIVERS: ReadonlyArray> = [ + AntigravityDriver, CodexDriver, ClaudeDriver, CursorDriver, diff --git a/packages/contracts/src/model.ts b/packages/contracts/src/model.ts index 9fcd0d266dd6..044526e03976 100644 --- a/packages/contracts/src/model.ts +++ b/packages/contracts/src/model.ts @@ -132,6 +132,7 @@ const CLAUDE_DRIVER_KIND = ProviderDriverKind.make("claudeAgent"); const CURSOR_DRIVER_KIND = ProviderDriverKind.make("cursor"); const GROK_DRIVER_KIND = ProviderDriverKind.make("grok"); const OPENCODE_DRIVER_KIND = ProviderDriverKind.make("opencode"); +const ANTIGRAVITY_DRIVER_KIND = ProviderDriverKind.make("antigravity"); export const DEFAULT_MODEL = "gpt-5.6-sol"; @@ -153,6 +154,7 @@ export const DEFAULT_MODEL_BY_PROVIDER: Partial> [CURSOR_DRIVER_KIND]: "Cursor", [GROK_DRIVER_KIND]: "Grok", [OPENCODE_DRIVER_KIND]: "OpenCode", + [ANTIGRAVITY_DRIVER_KIND]: "Antigravity", }; diff --git a/packages/contracts/src/settings.ts b/packages/contracts/src/settings.ts index 76bb032c866e..e249abb8d9fa 100644 --- a/packages/contracts/src/settings.ts +++ b/packages/contracts/src/settings.ts @@ -550,6 +550,58 @@ export const OpenCodeSettings = makeProviderSettingsSchema( ); export type OpenCodeSettings = typeof OpenCodeSettings.Type; +export const AntigravitySettings = makeProviderSettingsSchema( + { + enabled: Schema.Boolean.pipe( + Schema.withDecodingDefault(Effect.succeed(true)), + Schema.annotateKey({ providerSettingsForm: { hidden: true } }), + ), + binaryPath: makeBinaryPathSetting("agy").pipe( + Schema.annotateKey({ + title: "Binary path", + description: "Path to the Antigravity (agy) CLI binary.", + providerSettingsForm: { placeholder: "agy", clearWhenEmpty: "omit" }, + }), + ), + model: TrimmedString.pipe( + Schema.withDecodingDefault(Effect.succeed("")), + Schema.annotateKey({ + title: "Default model", + description: "Default model for Antigravity sessions (e.g. gemini-3.7-flash-high).", + providerSettingsForm: { placeholder: "gemini-3.7-flash-high", clearWhenEmpty: "omit" }, + }), + ), + effort: TrimmedString.pipe( + Schema.withDecodingDefault(Effect.succeed("")), + Schema.annotateKey({ + title: "Reasoning effort", + description: "Reasoning effort (low, medium, high).", + providerSettingsForm: { placeholder: "low", clearWhenEmpty: "omit" }, + }), + ), + customModels: Schema.Array(Schema.String).pipe( + Schema.withDecodingDefault(Effect.succeed([])), + Schema.annotateKey({ providerSettingsForm: { hidden: true } }), + ), + launchArgs: Schema.String.pipe( + Schema.withDecodingDefault(Effect.succeed("")), + Schema.annotateKey({ + title: "Launch arguments", + description: + "Additional CLI arguments passed to agy (e.g. --add-dir /path). Note: --sandbox is a terminal restriction only, not a capability boundary.", + providerSettingsForm: { + placeholder: "e.g. --add-dir /extra", + clearWhenEmpty: "omit", + }, + }), + ), + }, + { + order: ["binaryPath", "model", "effort", "launchArgs"], + }, +); +export type AntigravitySettings = typeof AntigravitySettings.Type; + export const ObservabilitySettings = Schema.Struct({ otlpTracesUrl: TrimmedString.pipe(Schema.withDecodingDefault(Effect.succeed(""))), otlpMetricsUrl: TrimmedString.pipe(Schema.withDecodingDefault(Effect.succeed(""))), @@ -692,6 +744,7 @@ export const ServerSettings = Schema.Struct({ cursor: CursorSettings.pipe(Schema.withDecodingDefault(Effect.succeed({}))), grok: GrokSettings.pipe(Schema.withDecodingDefault(Effect.succeed({}))), opencode: OpenCodeSettings.pipe(Schema.withDecodingDefault(Effect.succeed({}))), + antigravity: AntigravitySettings.pipe(Schema.withDecodingDefault(Effect.succeed({}))), }).pipe(Schema.withDecodingDefault(Effect.succeed({}))), // New driver-agnostic instance map. Keyed by `ProviderInstanceId`; values // are `ProviderInstanceConfig` envelopes. The driver-specific config blob @@ -845,6 +898,15 @@ const OpenCodeSettingsPatch = Schema.Struct({ customModels: Schema.optionalKey(Schema.Array(Schema.String)), }); +const AntigravitySettingsPatch = Schema.Struct({ + enabled: Schema.optionalKey(Schema.Boolean), + binaryPath: Schema.optionalKey(TrimmedString), + model: Schema.optionalKey(TrimmedString), + effort: Schema.optionalKey(TrimmedString), + customModels: Schema.optionalKey(Schema.Array(Schema.String)), + launchArgs: Schema.optionalKey(Schema.String), +}); + export const ServerSettingsPatch = Schema.Struct({ // Server settings enableLegacyTokenStreaming: Schema.optionalKey(Schema.Boolean), @@ -886,6 +948,7 @@ export const ServerSettingsPatch = Schema.Struct({ cursor: Schema.optionalKey(CursorSettingsPatch), grok: Schema.optionalKey(GrokSettingsPatch), opencode: Schema.optionalKey(OpenCodeSettingsPatch), + antigravity: Schema.optionalKey(AntigravitySettingsPatch), }), ), // Whole-map replacement for the new instance config. Patching individual