Skip to content
Open
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
11 changes: 11 additions & 0 deletions apps/server/scripts/acp-mock-agent.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ const emitGenericToolPlaceholders = process.env.T3_ACP_EMIT_GENERIC_TOOL_PLACEHO
const emitAskQuestion = process.env.T3_ACP_EMIT_ASK_QUESTION === "1";
const emitXAiAskUserQuestion = process.env.T3_ACP_EMIT_XAI_ASK_USER_QUESTION === "1";
const emitXAiPromptCompleteThenHang = process.env.T3_ACP_EMIT_XAI_PROMPT_COMPLETE_THEN_HANG === "1";
const emitXAiRateLimitThenHang = process.env.T3_ACP_EMIT_XAI_RATE_LIMIT_THEN_HANG === "1";
const emitForeignSessionUpdates = process.env.T3_ACP_EMIT_FOREIGN_SESSION_UPDATES === "1";
const hangPromptForever = process.env.T3_ACP_HANG_PROMPT_FOREVER === "1";
const hangFirstPromptForever = process.env.T3_ACP_HANG_FIRST_PROMPT_FOREVER === "1";
Expand Down Expand Up @@ -522,6 +523,16 @@ const program = Effect.gen(function* () {
return yield* Effect.never;
}

if (emitXAiRateLimitThenHang) {
writeJsonRpcNotification("_x.ai/session/prompt_complete", {
sessionId: requestedSessionId,
promptId: promptIdFromRequestMeta(request) ?? "mock-xai-rate-limit-prompt-1",
stopReason: "rate_limit",
agentResult: null,
});
return yield* Effect.never;
}

if (emitXAiPromptCompleteThenHang) {
writeJsonRpcNotification("session/update", {
sessionId: requestedSessionId,
Expand Down
58 changes: 58 additions & 0 deletions apps/server/src/provider/Layers/GrokAdapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -943,6 +943,64 @@ it.layer(grokAdapterTestLayer)("GrokAdapterLive", (it) => {
}),
);

it.effect("surfaces Grok usage limits without clearing the selected model", () =>
Effect.gen(function* () {
const threadId = ThreadId.make("grok-usage-limit-error");
const wrapperPath = yield* Effect.promise(() =>
makeMockGrokWrapper({
T3_ACP_EMIT_XAI_RATE_LIMIT_THEN_HANG: "1",
}),
);
const adapter = yield* makeTestAdapter(wrapperPath);
const runtimeEvents: ProviderRuntimeEvent[] = [];
const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) =>
Effect.sync(() => {
runtimeEvents.push(event);
}),
).pipe(Effect.forkChild);

yield* adapter.startSession({
threadId,
provider: ProviderDriverKind.make("grok"),
cwd: process.cwd(),
runtimeMode: "full-access",
modelSelection: { instanceId: ProviderInstanceId.make("grok"), model: "grok-build" },
});

const error = yield* Effect.flip(
adapter.sendTurn({
threadId,
input: "hit the usage limit",
attachments: [],
}),
);
const readySessions = yield* adapter.listSessions();
const readySession = readySessions.find((session) => session.threadId === threadId);
const terminalEvents = runtimeEvents.filter(
(event) => event.type === "turn.completed" && event.threadId === threadId,
);

assert.equal(error._tag, "ProviderAdapterRequestError");
assert.include(error.message, "Grok usage limit reached. Try again later.");
assert.equal(readySession?.status, "ready");
assert.equal(readySession?.model, "grok-build");
assert.isUndefined(readySession?.activeTurnId);
assert.lengthOf(terminalEvents, 1);
const [terminalEvent] = terminalEvents;
assert.equal(terminalEvent?.type, "turn.completed");
if (terminalEvent?.type === "turn.completed") {
assert.equal(terminalEvent.payload.state, "failed");
assert.include(
terminalEvent.payload.errorMessage ?? "",
"Grok usage limit reached. Try again later.",
);
}

yield* Fiber.interrupt(runtimeEventsFiber);
yield* adapter.stopSession(threadId);
}),
);

it.effect("ignores replayed session/load updates when resuming a Grok session", () =>
Effect.gen(function* () {
const threadId = ThreadId.make("grok-load-replay-filter");
Expand Down
21 changes: 21 additions & 0 deletions apps/server/src/provider/acp/XAiAcpExtension.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -299,6 +299,27 @@ describe("XAiAcpExtension", () => {
}).pipe(Effect.scoped, Effect.provide(NodeServices.layer)),
);

it.effect("fails a hung standard prompt from an xAI rate-limit completion", () =>
Effect.gen(function* () {
const runtime = yield* makePromptCompletionRuntime({
T3_ACP_EMIT_XAI_RATE_LIMIT_THEN_HANG: "1",
});
yield* runtime.start();

const error = yield* Effect.flip(
runtime.prompt({
prompt: [{ type: "text", text: "hi" }],
}),
);

expect(error).toMatchObject({
_tag: "AcpRequestError",
code: -32003,
errorMessage: "Grok usage limit reached. Try again later.",
});
}).pipe(Effect.scoped, Effect.provide(NodeServices.layer)),
);

it.effect("ignores stale xAI completion from an already settled prompt", () =>
Effect.gen(function* () {
const runtime = yield* makePromptCompletionRuntime({
Expand Down
48 changes: 44 additions & 4 deletions apps/server/src/provider/acp/XAiAcpExtension.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ import * as Deferred from "effect/Deferred";
import * as Effect from "effect/Effect";
import * as Ref from "effect/Ref";
import * as Schema from "effect/Schema";
import * as EffectAcpErrors from "effect-acp/errors";
import type * as EffectAcpSchema from "effect-acp/schema";

import type * as AcpSessionRuntime from "./AcpSessionRuntime.ts";
Expand All @@ -19,11 +20,15 @@ type XAiPromptCompleteNotification = typeof XAiPromptCompleteNotification.Type;
interface PendingXAiPromptCompletion {
readonly sessionId: string;
readonly promptId: string;
readonly deferred: Deferred.Deferred<EffectAcpSchema.PromptResponse>;
readonly deferred: Deferred.Deferred<
EffectAcpSchema.PromptResponse,
EffectAcpErrors.AcpRequestError
>;
}

const completedXAiPromptIdLimit = 128;
const xAiStopReasonMissingMetaKey = "xAiStopReasonMissing";
const xAiRateLimitedErrorCode = -32003;

const XAiAskUserQuestionOption = Schema.Struct({
label: Schema.String,
Expand Down Expand Up @@ -278,7 +283,7 @@ const registerXAiPromptCompletionFallback = (
sessionId: string,
promptId: string,
) =>
Deferred.make<EffectAcpSchema.PromptResponse>().pipe(
Deferred.make<EffectAcpSchema.PromptResponse, EffectAcpErrors.AcpRequestError>().pipe(
Effect.tap((deferred) =>
Ref.update(pendingRef, (pending) => [...pending, { sessionId, promptId, deferred }]),
),
Expand All @@ -287,7 +292,7 @@ const registerXAiPromptCompletionFallback = (

const unregisterXAiPromptCompletionFallback = (
pendingRef: Ref.Ref<ReadonlyArray<PendingXAiPromptCompletion>>,
deferred: Deferred.Deferred<EffectAcpSchema.PromptResponse>,
deferred: Deferred.Deferred<EffectAcpSchema.PromptResponse, EffectAcpErrors.AcpRequestError>,
) => Ref.update(pendingRef, (pending) => pending.filter((entry) => entry.deferred !== deferred));

const abortPendingPromptCompletions = (
Expand Down Expand Up @@ -358,13 +363,48 @@ const resolveXAiPromptCompletionFallback = ({
return [Effect.void, pending] as const;
}
return [
Deferred.succeed(entry.deferred, promptResponseFromXAi(notification)).pipe(Effect.asVoid),
settleXAiPromptCompletion(entry.deferred, notification),
[...pending.slice(0, index), ...pending.slice(index + 1)],
] as const;
}).pipe(Effect.flatten);
}),
);

const settleXAiPromptCompletion = (
deferred: Deferred.Deferred<EffectAcpSchema.PromptResponse, EffectAcpErrors.AcpRequestError>,
notification: XAiPromptCompleteNotification,
) => {
if (notification.stopReason === "rate_limit") {
return Deferred.fail(
deferred,
new EffectAcpErrors.AcpRequestError({
code: xAiRateLimitedErrorCode,
errorMessage: "Grok usage limit reached. Try again later.",
}),
).pipe(Effect.asVoid);
}
if (notification.stopReason === "error") {
return Deferred.fail(
deferred,
EffectAcpErrors.AcpRequestError.internalError(
xAiAgentResultMessage(notification.agentResult) ?? "Grok prompt failed.",
),
).pipe(Effect.asVoid);
}
return Deferred.succeed(deferred, promptResponseFromXAi(notification)).pipe(Effect.asVoid);
};

function xAiAgentResultMessage(value: unknown): string | undefined {
if (typeof value === "string") {
return trimmed(value);
}
if (value === null || typeof value !== "object") {
return undefined;
}
const message = "message" in value ? value.message : undefined;
return typeof message === "string" ? trimmed(message) : undefined;
}

const rememberCompletedXAiPromptId = (
completedPromptIdsRef: Ref.Ref<ReadonlyArray<string>>,
response: EffectAcpSchema.PromptResponse,
Expand Down
Loading