Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 39 additions & 0 deletions apps/server/scripts/acp-mock-agent.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ const emitInterleavedAssistantToolCalls =
const emitGenericToolPlaceholders = process.env.T3_ACP_EMIT_GENERIC_TOOL_PLACEHOLDERS === "1";
const emitAskQuestion = process.env.T3_ACP_EMIT_ASK_QUESTION === "1";
const emitXAiAskUserQuestion = process.env.T3_ACP_EMIT_XAI_ASK_USER_QUESTION === "1";
const emitXAiSubagent = process.env.T3_ACP_EMIT_XAI_SUBAGENT === "1";
const emitXAiExitPlanMode = process.env.T3_ACP_EMIT_XAI_EXIT_PLAN_MODE === "1";
const emitXAiPlanMdWrite = process.env.T3_ACP_EMIT_XAI_PLAN_MD_WRITE === "1";
const emitXAiPromptCompleteThenHang = process.env.T3_ACP_EMIT_XAI_PROMPT_COMPLETE_THEN_HANG === "1";
Expand Down Expand Up @@ -543,6 +544,44 @@ const program = Effect.gen(function* () {
return yield* Effect.never;
}

if (emitXAiSubagent) {
writeJsonRpcNotification("_x.ai/session/update", {
sessionId: requestedSessionId,
update: {
sessionUpdate: "subagent_spawned",
subagent_id: "grok-subagent-1",
subagent_type: "explore",
description: "Search the codebase",
model: "grok-4.6",
},
});
writeJsonRpcNotification("_x.ai/session/update", {
sessionId: requestedSessionId,
update: {
sessionUpdate: "subagent_progress",
subagent_id: "grok-subagent-1",
subagent_type: "explore",
description: "Search the codebase",
status: "running",
last_tool_name: "grep",
tokens_used: 80,
},
});
writeJsonRpcNotification("_x.ai/session/update", {
sessionId: requestedSessionId,
update: {
sessionUpdate: "subagent_finished",
subagent_id: "grok-subagent-1",
subagent_type: "explore",
description: "Search the codebase",
status: "completed",
output: "Found 3 call sites.",
tokens_used: 120,
duration_ms: 400,
},
});
}

if (emitXAiRateLimitThenHang) {
writeJsonRpcNotification("_x.ai/session/prompt_complete", {
sessionId: requestedSessionId,
Expand Down
53 changes: 53 additions & 0 deletions apps/server/src/provider/Layers/GrokAdapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2085,4 +2085,57 @@ it.layer(grokAdapterTestLayer)("GrokAdapterLive", (it) => {
// hang until the suite timeout instead of failing here.
}).pipe(TestClock.withLive),
);

it.effect("maps Grok ACP subagent notifications onto the Agents panel task stream", () =>
Effect.gen(function* () {
const threadId = ThreadId.make("grok-subagent-panel");
const wrapperPath = yield* Effect.promise(() =>
makeMockGrokWrapper({ T3_ACP_EMIT_XAI_SUBAGENT: "1" }),
);
const adapter = yield* makeTestAdapter(wrapperPath);
const started =
yield* Deferred.make<Extract<ProviderRuntimeEvent, { type: "task.started" }>>();
const completed =
yield* Deferred.make<Extract<ProviderRuntimeEvent, { type: "task.completed" }>>();
const eventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => {
if (String(event.threadId) !== String(threadId)) {
return Effect.void;
}
if (event.type === "task.started") {
return Deferred.succeed(started, event).pipe(Effect.ignore);
}
if (event.type === "task.completed") {
return Deferred.succeed(completed, event).pipe(Effect.ignore);
}
return Effect.void;
}).pipe(Effect.forkChild);

yield* adapter.startSession({
threadId,
provider: ProviderDriverKind.make("grok"),
cwd: process.cwd(),
runtimeMode: "full-access",
});
yield* adapter.sendTurn({
threadId,
input: "delegate the search",
attachments: [],
});

const startedEvent = yield* Deferred.await(started);
const completedEvent = yield* Deferred.await(completed);
assert.equal(startedEvent.payload.taskId, "grok-subagent-1");
assert.equal(startedEvent.payload.taskType, "subagent");
assert.equal(startedEvent.payload.role, "explore");
assert.equal(startedEvent.payload.title, "Search the codebase");
assert.equal(startedEvent.payload.timelineBypass, true);
assert.equal(startedEvent.raw?.method, "_x.ai/session/update");
assert.equal(completedEvent.payload.taskId, "grok-subagent-1");
assert.equal(completedEvent.payload.status, "completed");
assert.equal(completedEvent.payload.summary, "Found 3 call sites.");

yield* Fiber.interrupt(eventsFiber);
yield* adapter.stopSession(threadId);
}),
);
});
153 changes: 153 additions & 0 deletions apps/server/src/provider/Layers/GrokAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import {
ProviderDriverKind,
ProviderInstanceId,
RuntimeRequestId,
RuntimeTaskId,
type ThreadId,
TurnId,
} from "@t3tools/contracts";
Expand Down Expand Up @@ -67,6 +68,18 @@ import {
normalizeGrokReasoningEffort,
resolveGrokAcpBaseModelId,
} from "../acp/GrokAcpSupport.ts";
import {
applyGrokSubagentUpdate,
applyGrokWorkflowUpdate,
emptyGrokSubagentTrackState,
GROK_SESSION_NOTIFICATION_METHODS,
parseXAiSubagentUpdate,
parseXAiWorkflowUpdated,
XAiSessionNotification,
type GrokSessionNotificationMethod,
type GrokSubagentTrackState,
type GrokTaskEventSpec,
} from "../acp/GrokAcpSubagents.ts";
import {
extractGrokPlanMarkdownFromToolCallData,
extractXAiAskUserQuestions,
Expand Down Expand Up @@ -162,6 +175,14 @@ interface GrokSessionContext {
promptResponsesReady: number;
currentModelId: string | undefined;
currentReasoningEffort: string | undefined;
/**
* Grok `agent()` / workflow ACP extras mapped onto task.* so the Agents
* panel can list them the same way Claude Task and Codex collab children
* already do.
*/
workflowTrack: GrokSubagentTrackState;
/** First-seen spawn turn per task, reused after the parent turn settles. */
readonly taskSpawnTurnIds: Map<string, TurnId>;
stopped: boolean;
}

Expand Down Expand Up @@ -399,6 +420,45 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte
const offerRuntimeEvent = (event: ProviderRuntimeEvent) =>
PubSub.publish(runtimeEventPubSub, event).pipe(Effect.asVoid);

const emitGrokTaskSpecs = (input: {
readonly ctx: GrokSessionContext;
readonly method: string;
readonly payload: unknown;
readonly specs: ReadonlyArray<GrokTaskEventSpec>;
}) =>
Effect.forEach(
input.specs,
(spec) =>
Effect.gen(function* () {
const taskIdValue = spec.payload.taskId;
if (typeof taskIdValue !== "string" || taskIdValue.length === 0) {
return;
}
const taskId = RuntimeTaskId.make(taskIdValue);
let turnId = input.ctx.taskSpawnTurnIds.get(taskIdValue);
if (turnId === undefined) {
turnId = resolveNotificationTurnId(input.ctx);
if (turnId !== undefined) {
input.ctx.taskSpawnTurnIds.set(taskIdValue, turnId);
}
}
yield* offerRuntimeEvent({
type: spec.type,
...(yield* makeEventStamp()),
provider: PROVIDER,
threadId: input.ctx.threadId,
...(turnId !== undefined ? { turnId } : {}),
payload: { ...spec.payload, taskId },
raw: {
source: "acp.grok.extension",
method: input.method,
payload: input.payload,
},
} as ProviderRuntimeEvent);
}),
{ discard: true },
);

const getThreadSemaphore = (threadId: string) =>
SynchronizedRef.modifyEffect(threadLocksRef, (current) => {
const existing: Option.Option<Semaphore.Semaphore> = Option.fromNullishOr(
Expand Down Expand Up @@ -506,6 +566,54 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte
return ctx && turnId !== undefined ? signalTurnLiveness(ctx, turnId) : Effect.void;
};

const refreshGrokSessionTurnLiveness = (ctx: GrokSessionContext) =>
Effect.gen(function* () {
const turnId = resolveNotificationTurnId(ctx);
if (turnId === undefined || ctx.livenessTurnId !== turnId) {
return;
}
// Subagent/workflow ticks prove the provider is still working even
// when the parent ACP stream is quiet. Stamp the watchdog clock
// before signalling; a wake without a fresh timestamp cancels a
// turn that has already sat near turnInactivityTimeoutMs.
ctx.lastTurnActivityAtNanos = yield* Clock.monotonicTimeNanos;
yield* signalTurnLiveness(ctx, turnId);
});

const applyGrokSessionNotification = (
ctx: GrokSessionContext,
method: GrokSessionNotificationMethod,
params: unknown,
) =>
Effect.gen(function* () {
const workflow = parseXAiWorkflowUpdated(params);
if (workflow) {
const applied = applyGrokWorkflowUpdate(ctx.workflowTrack, workflow);
ctx.workflowTrack = applied.state;
yield* emitGrokTaskSpecs({
ctx,
method,
payload: params,
specs: applied.events,
});
yield* refreshGrokSessionTurnLiveness(ctx);
return;
}
const subagent = parseXAiSubagentUpdate(params);
if (!subagent) {
return;
}
const applied = applyGrokSubagentUpdate(ctx.workflowTrack, subagent);
ctx.workflowTrack = applied.state;
yield* emitGrokTaskSpecs({
ctx,
method,
payload: params,
specs: applied.events,
});
yield* refreshGrokSessionTurnLiveness(ctx);
});

const resumeSessionTurnLiveness = Effect.fn("GrokAdapter.resumeSessionTurnLiveness")(function* (
threadId: ThreadId,
turnId: TurnId | undefined,
Expand Down Expand Up @@ -1018,6 +1126,12 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte
}),
),
);
const pendingSessionNotifications: Array<{
readonly method: GrokSessionNotificationMethod;
readonly params: unknown;
}> = [];
let sessionNotificationsReady = false;
const sessionNotificationLock = yield* Semaphore.make(1);
const started = yield* Effect.gen(function* () {
yield* Effect.forEach(
["x.ai/ask_user_question", "_x.ai/ask_user_question"] as const,
Expand Down Expand Up @@ -1121,6 +1235,34 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte
),
{ discard: true },
);
// Grok Build streams agent() / workflow runs over a private ACP
// channel. Claude maps Task/workflow_progress onto task.*; Codex
// maps collabAgent/* the same way. Keep that seam: parse here,
// emit canonical events, never a third Agents UI.
yield* Effect.forEach(
GROK_SESSION_NOTIFICATION_METHODS,
(method) =>
acp.handleExtNotification(method, XAiSessionNotification, (params) =>
mapAcpCallbackFailure(
Effect.gen(function* () {
yield* logNative(input.threadId, method, params);
if (!sessionNotificationsReady) {
pendingSessionNotifications.push({ method, params });
return;
}
const liveCtx = sessions.get(input.threadId);
if (!liveCtx) {
pendingSessionNotifications.push({ method, params });
return;
}
yield* sessionNotificationLock.withPermit(
applyGrokSessionNotification(liveCtx, method, params),
);
}),
),
),
{ discard: true },
);
yield* acp.handleRequestPermission((params) =>
mapAcpCallbackFailure(
Effect.gen(function* () {
Expand Down Expand Up @@ -1290,6 +1432,8 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte
requestedStartReasoningEffort !== undefined
? normalizeGrokReasoningEffort(requestedStartReasoningEffort)
: currentStartReasoningEffort,
workflowTrack: emptyGrokSubagentTrackState(),
taskSpawnTurnIds: new Map(),
stopped: false,
};

Expand Down Expand Up @@ -1433,6 +1577,15 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte
sessions.set(input.threadId, ctx);
yield* runTurnLivenessWatchdog(ctx).pipe(Effect.forkIn(ctx.scope), Effect.asVoid);
sessionScopeTransferred = true;
yield* sessionNotificationLock.withPermit(
Effect.gen(function* () {
const batch = pendingSessionNotifications.splice(0);
sessionNotificationsReady = true;
yield* Effect.forEach(batch, (pending) =>
applyGrokSessionNotification(ctx, pending.method, pending.params),
);
}),
);

yield* offerRuntimeEvent({
type: "session.started",
Expand Down
Loading
Loading