diff --git a/apps/server/src/provider/Drivers/ClaudeHome.test.ts b/apps/server/src/provider/Drivers/ClaudeHome.test.ts index 70d95afa10be..0395690390f2 100644 --- a/apps/server/src/provider/Drivers/ClaudeHome.test.ts +++ b/apps/server/src/provider/Drivers/ClaudeHome.test.ts @@ -11,6 +11,7 @@ import * as Path from "effect/Path"; import { subscriptionRuntimeEnvironment } from "../../subscription-auth/runtime.ts"; import { SubscriptionAuthService } from "../../subscription-auth/service.ts"; import { + claudeSignedOutMessage, makeClaudeCapabilitiesCacheKey, makeClaudeContinuationGroupKey, makeClaudeEnvironment, @@ -44,6 +45,17 @@ it.layer(NodeServices.layer)("ClaudeHome", (it) => { }), ); + it("points the signed-out hint at the configured Claude home", () => { + expect(claudeSignedOutMessage({ configDir: undefined, cwd: "/synthetic" })).toContain( + "run `claude auth login`", + ); + const configDir = "/synthetic/Claude work's $literal"; + const message = claudeSignedOutMessage({ configDir, cwd: "/synthetic/project" }); + expect(message).toContain(`CLAUDE_CONFIG_DIR set to "${configDir}"`); + expect(message).not.toContain("CLAUDE_CONFIG_DIR="); + expect(message).toContain("then start a new thread"); + }); + it.effect("marks CLAUDE_CONFIG_DIR explicit so saved API keys do not replace the account", () => Effect.gen(function* () { const secretsDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "akeru-claude-home-")); diff --git a/apps/server/src/provider/Drivers/ClaudeHome.ts b/apps/server/src/provider/Drivers/ClaudeHome.ts index aec4a1fc7434..15edd465005c 100644 --- a/apps/server/src/provider/Drivers/ClaudeHome.ts +++ b/apps/server/src/provider/Drivers/ClaudeHome.ts @@ -3,10 +3,13 @@ import * as NodeOS from "node:os"; import type { ClaudeSettings } from "@t3tools/contracts"; import * as Effect from "effect/Effect"; import * as Path from "effect/Path"; +import * as Schema from "effect/Schema"; import { expandHomePath } from "../../pathExpansion.ts"; import { withExplicitEnvironmentKeys } from "../../subscription-auth/runtime.ts"; +const quotePath = Schema.encodeSync(Schema.fromJsonString(Schema.String)); + export const resolveClaudeHomePath = Effect.fn("resolveClaudeHomePath")(function* ( config: Pick, ): Effect.fn.Return { @@ -51,3 +54,18 @@ export const makeClaudeCapabilitiesCacheKey = Effect.fn("makeClaudeCapabilitiesC return `${config.binaryPath}\0${resolvedHomePath}\0${cwd ?? ""}`; }, ); + +/** + * Describe the spawned CLI's environment separately from the login command so + * paths remain literal on every shell, including relative inherited values. + */ +export const claudeSignedOutMessage = (input: { + readonly configDir: string | undefined; + readonly cwd: string; +}): string => { + const configuration = + input.configDir !== undefined + ? ` from ${quotePath(input.cwd)}, with CLAUDE_CONFIG_DIR set to ${quotePath(input.configDir)}` + : ""; + return `Claude could not authenticate. For subscription login, run \`claude auth login\` on this environment's machine${configuration}, then start a new thread. For API-key authentication, check this instance's configured credentials.`; +}; diff --git a/apps/server/src/provider/Layers/ClaudeAdapter.test.ts b/apps/server/src/provider/Layers/ClaudeAdapter.test.ts index 3b9ed57c619a..e19fe2c6da36 100644 --- a/apps/server/src/provider/Layers/ClaudeAdapter.test.ts +++ b/apps/server/src/provider/Layers/ClaudeAdapter.test.ts @@ -27,6 +27,7 @@ import * as Context from "effect/Context"; import * as Effect from "effect/Effect"; import * as Fiber from "effect/Fiber"; import * as Layer from "effect/Layer"; +import * as Clock from "effect/Clock"; import * as Random from "effect/Random"; import * as Schema from "effect/Schema"; import * as Stream from "effect/Stream"; @@ -156,6 +157,7 @@ function makeHarness(config?: { readonly baseDir?: string; readonly claudeConfig?: Partial; readonly instanceId?: ProviderInstanceId; + readonly environment?: ClaudeAdapterLiveOptions["environment"]; }) { const query = new FakeClaudeQuery(); let createInput: @@ -166,6 +168,7 @@ function makeHarness(config?: { | undefined; const adapterOptions: ClaudeAdapterLiveOptions = { + ...(config?.environment ? { environment: config.environment } : {}), ...(config?.instanceId ? { instanceId: config.instanceId } : {}), createQuery: (input) => { createInput = input; @@ -266,6 +269,28 @@ async function readFirstPromptMessage( const THREAD_ID = ThreadId.make("thread-claude-1"); const RESUME_THREAD_ID = ThreadId.make("thread-claude-resume"); +const encodeUnknownJsonString = Schema.encodeSync(Schema.fromJsonString(Schema.String)); + +const completedTurn = (runtimeEvents: ReadonlyArray) => { + const event = runtimeEvents[runtimeEvents.length - 1]; + assert.equal(event?.type, "turn.completed"); + assert(event?.type === "turn.completed"); + return event.payload; +}; + +const AUTH_FAILURE_ASSISTANT = { + type: "assistant", + session_id: "sdk-session-auth", + uuid: "assistant-auth", + parent_tool_use_id: null, + error: "authentication_failed", + is_api_error_message: true, + message: { + id: "assistant-message-auth", + model: "", + content: [{ type: "text", text: "Not logged in ยท Please run /login" }], + }, +} as unknown as SDKMessage; describe("ClaudeAdapterLive", () => { it.effect("passes a newly saved API key and endpoint to the Claude SDK", () => { @@ -1696,6 +1721,547 @@ describe("ClaudeAdapterLive", () => { ); }); + it.effect("completes with result usage without querying current context usage", () => { + const harness = makeHarness(); + let getContextUsageCalls = 0; + Object.assign(harness.query, { + getContextUsage: async () => { + getContextUsageCalls += 1; + return { + totalTokens: 999, + maxTokens: 200000, + isAutoCompactEnabled: true, + }; + }, + }); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const runtimeEventsFiber = yield* adapter.streamEvents.pipe( + Stream.takeUntil((event) => event.type === "turn.completed"), + Stream.runCollect, + Effect.forkChild, + ); + yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + yield* adapter.sendTurn({ + threadId: THREAD_ID, + input: "hello", + attachments: [], + }); + harness.query.emit({ + type: "assistant", + session_id: "sdk-session-result-usage", + uuid: "assistant-result-usage-1", + parent_tool_use_id: null, + message: { + id: "assistant-message-result-usage-1", + role: "assistant", + content: [], + usage: { input_tokens: 80, output_tokens: 20 }, + }, + } as unknown as SDKMessage); + harness.query.emit({ + type: "assistant", + session_id: "sdk-session-result-usage", + uuid: "assistant-result-usage-2", + parent_tool_use_id: null, + message: { + id: "assistant-message-result-usage-2", + role: "assistant", + content: [], + usage: { input_tokens: 180, output_tokens: 20 }, + }, + } as unknown as SDKMessage); + harness.query.emit({ + type: "result", + subtype: "success", + is_error: false, + duration_ms: 1234, + duration_api_ms: 1200, + num_turns: 1, + result: "done", + stop_reason: "end_turn", + session_id: "sdk-session-result-usage", + usage: { input_tokens: 400, output_tokens: 50 }, + modelUsage: { + "claude-opus-4-6": { contextWindow: 200000, maxOutputTokens: 64000 }, + }, + } as unknown as SDKMessage); + + const runtimeEvents = Array.from(yield* Fiber.join(runtimeEventsFiber)); + assert.equal(getContextUsageCalls, 0); + const usageEvent = runtimeEvents.find((event) => event.type === "thread.token-usage.updated"); + assert.equal(usageEvent?.type, "thread.token-usage.updated"); + if (usageEvent?.type === "thread.token-usage.updated") { + assert.deepEqual(usageEvent.payload.usage, { + usedTokens: 200, + lastUsedTokens: 200, + totalProcessedTokens: 450, + inputTokens: 180, + outputTokens: 20, + maxTokens: 200000, + }); + } + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + it.effect("preserves compacted usage when completion follows an older assistant frame", () => { + const harness = makeHarness(); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const runtimeEventsFiber = yield* adapter.streamEvents.pipe( + Stream.takeUntil((event) => event.type === "turn.completed"), + Stream.runCollect, + Effect.forkChild, + ); + yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + yield* adapter.sendTurn({ + threadId: THREAD_ID, + input: "hello", + attachments: [], + }); + harness.query.emit({ + type: "assistant", + session_id: "sdk-session-compacted-usage", + uuid: "assistant-compacted-usage", + parent_tool_use_id: null, + message: { + id: "assistant-message-compacted-usage", + role: "assistant", + content: [], + usage: { input_tokens: 180, output_tokens: 20 }, + }, + } as unknown as SDKMessage); + harness.query.emit({ + type: "system", + subtype: "compact_boundary", + compact_metadata: { pre_tokens: 200, post_tokens: 40 }, + session_id: "sdk-session-compacted-usage", + uuid: "compact-boundary-usage", + } as unknown as SDKMessage); + harness.query.emit({ + type: "result", + subtype: "success", + is_error: false, + duration_ms: 1234, + duration_api_ms: 1200, + num_turns: 2, + result: "done", + stop_reason: "end_turn", + session_id: "sdk-session-compacted-usage", + usage: { input_tokens: 400, output_tokens: 50 }, + modelUsage: { + "claude-opus-4-6": { contextWindow: 200000, maxOutputTokens: 64000 }, + }, + } as unknown as SDKMessage); + + const runtimeEvents = Array.from(yield* Fiber.join(runtimeEventsFiber)); + const finalUsageEvent = runtimeEvents.findLast( + (event) => event.type === "thread.token-usage.updated", + ); + assert.equal(finalUsageEvent?.type, "thread.token-usage.updated"); + if (finalUsageEvent?.type === "thread.token-usage.updated") { + assert.deepEqual(finalUsageEvent.payload.usage, { + usedTokens: 40, + lastUsedTokens: 200, + totalProcessedTokens: 450, + maxTokens: 200000, + }); + } + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + it.effect.each(["success", "error_during_execution"] as const)( + "preserves %s behavior for an unknown runtime terminal reason", + (subtype) => { + const harness = makeHarness(); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const completionFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.type === "turn.completed"), + Stream.runHead, + Effect.forkChild, + ); + const session = yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + yield* adapter.sendTurn({ threadId: session.threadId, input: "hello", attachments: [] }); + harness.query.emit({ + type: "result", + subtype, + is_error: subtype !== "success", + result: "", + errors: subtype === "success" ? [] : ["Provider error detail"], + stop_reason: null, + terminal_reason: "future_terminal_reason", + session_id: "sdk-session-future-reason", + uuid: "result-future-reason", + } as unknown as SDKMessage); + const completed = yield* Fiber.join(completionFiber); + assert.equal(completed._tag, "Some"); + if (completed._tag === "Some" && completed.value.type === "turn.completed") { + assert.equal( + completed.value.payload.state, + subtype === "success" ? "completed" : "failed", + ); + assert.equal( + completed.value.payload.errorMessage, + subtype === "success" ? undefined : "Provider error detail", + ); + } + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }, + ); + + it.effect.each([ + { + name: "an api_error terminal reason", + result: { subtype: "success", is_error: false, terminal_reason: "api_error", errors: [] }, + state: "failed", + errorMessage: /claude auth login/, + }, + { + name: "an is_error success with no terminal reason", + result: { subtype: "success", is_error: true, errors: [] }, + state: "failed", + errorMessage: /claude auth login/, + }, + { + name: "a terminal reason of its own", + result: { + subtype: "success", + is_error: false, + terminal_reason: "prompt_too_long", + errors: [], + }, + state: "failed", + errorMessage: /prompt exceeds the model's context window/, + }, + { + name: "a listed tool failure", + result: { + subtype: "error_during_execution", + is_error: true, + errors: ["Tool execution failed: EACCES"], + }, + state: "failed", + errorMessage: /EACCES/, + }, + { + name: "a user interrupt", + result: { + subtype: "error_during_execution", + is_error: true, + terminal_reason: "aborted_tools", + errors: [], + }, + state: "interrupted", + errorMessage: undefined, + }, + { + name: "a cancellation", + result: { subtype: "error_during_execution", is_error: true, errors: ["cancelled"] }, + state: "cancelled", + errorMessage: /cancelled/, + }, + ])( + "reports the real cause when an expired login is followed by $name", + ({ result, state, errorMessage }) => { + const harness = makeHarness(); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const runtimeEventsFiber = yield* adapter.streamEvents.pipe( + Stream.takeUntil((event) => event.type === "turn.completed"), + Stream.runCollect, + Effect.forkChild, + ); + const session = yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + yield* adapter.sendTurn({ threadId: session.threadId, input: "hello", attachments: [] }); + harness.query.emit(AUTH_FAILURE_ASSISTANT); + harness.query.emit({ + type: "result", + ...result, + session_id: "sdk-session-auth", + uuid: "result-auth", + } as unknown as SDKMessage); + const payload = completedTurn(Array.from(yield* Fiber.join(runtimeEventsFiber))); + assert.equal(payload.state, state); + if (errorMessage === undefined) { + assert.equal(payload.errorMessage, undefined); + } else { + assert.match(payload.errorMessage ?? "", errorMessage); + } + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }, + ); + + it.effect("fails a usage-limited turn with the limit it parked on", () => { + const harness = makeHarness(); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const runtimeEventsFiber = yield* adapter.streamEvents.pipe( + Stream.takeUntil((event) => event.type === "turn.completed"), + Stream.runCollect, + Effect.forkChild, + ); + const session = yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + yield* adapter.sendTurn({ threadId: session.threadId, input: "hello", attachments: [] }); + const nowMs = yield* Clock.currentTimeMillis; + harness.query.emit({ + type: "rate_limit_event", + rate_limit_info: { + status: "rejected", + rateLimitType: "five_hour", + resetsAt: Math.floor(nowMs / 1000) + 2 * 60 * 60, + }, + session_id: "sdk-session-limit", + uuid: "rate-limit-rejected", + } as unknown as SDKMessage); + harness.query.emit({ + type: "result", + subtype: "success", + is_error: false, + terminal_reason: "api_error", + errors: [], + session_id: "sdk-session-limit", + uuid: "result-limit", + } as unknown as SDKMessage); + const payload = completedTurn(Array.from(yield* Fiber.join(runtimeEventsFiber))); + assert.equal(payload.state, "failed"); + assert.equal( + payload.errorMessage, + "Claude usage limit reached. Send the message again once the limit resets.", + ); + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + const usageLimitMessage = + "Claude usage limit reached. Send the message again once the limit resets."; + const genericApiErrorMessage = "Claude gave up after repeated API errors."; + const rateLimitAssistant = { + type: "assistant", + session_id: "sdk-session-limit", + uuid: "assistant-limit", + parent_tool_use_id: null, + error: "rate_limit", + message: { + id: "assistant-message-limit", + model: "", + content: [{ type: "text", text: "You've hit your session limit" }], + }, + }; + const rateLimitResult = { + type: "result", + subtype: "success", + is_error: true, + terminal_reason: "api_error", + session_id: "sdk-session-limit", + uuid: "result-limit", + }; + + it.effect.each([ + { + name: "an assistant-only rate limit", + messages: [rateLimitAssistant], + expected: usageLimitMessage, + }, + { + name: "a normal parent response after a rate limit", + messages: [rateLimitAssistant, { ...rateLimitAssistant, error: undefined }], + expected: genericApiErrorMessage, + }, + { + name: "a server error after a rate limit", + messages: [rateLimitAssistant, { ...rateLimitAssistant, error: "server_error" }], + expected: genericApiErrorMessage, + }, + { + name: "a subagent rate limit", + messages: [{ ...rateLimitAssistant, parent_tool_use_id: "nested-tool" }], + expected: genericApiErrorMessage, + }, + { + name: "a subagent response after a parent rate limit", + messages: [ + rateLimitAssistant, + { ...rateLimitAssistant, error: undefined, parent_tool_use_id: "nested-tool" }, + ], + expected: usageLimitMessage, + }, + ])("classifies the terminal API failure after $name", ({ messages, expected }) => { + const harness = makeHarness(); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const eventsFiber = yield* adapter.streamEvents.pipe( + Stream.takeUntil((event) => event.type === "turn.completed"), + Stream.runCollect, + Effect.forkChild, + ); + const session = yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + yield* adapter.sendTurn({ threadId: session.threadId, input: "hello", attachments: [] }); + for (const [index, message] of messages.entries()) { + harness.query.emit({ ...message, uuid: `assistant-${index}` } as unknown as SDKMessage); + } + harness.query.emit(rateLimitResult as unknown as SDKMessage); + const events = Array.from(yield* Fiber.join(eventsFiber)); + const errors = events.filter((event) => event.type === "runtime.error"); + assert.equal(errors.length, 1); + assert.equal(errors[0]?.payload.message, expected); + assert.equal(completedTurn(events).state, "failed"); + assert.equal(completedTurn(events).errorMessage, expected); + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + it.effect("names repeated usage limits without carrying them into a later turn", () => { + const harness = makeHarness(); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const session = yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + for (const [index, expected] of [ + usageLimitMessage, + usageLimitMessage, + genericApiErrorMessage, + ].entries()) { + const eventsFiber = yield* adapter.streamEvents.pipe( + Stream.takeUntil((event) => event.type === "turn.completed"), + Stream.runCollect, + Effect.forkChild, + ); + yield* adapter.sendTurn({ threadId: session.threadId, input: "again", attachments: [] }); + if (index === 0) { + harness.query.emit({ + type: "rate_limit_event", + rate_limit_info: { status: "rejected", rateLimitType: "five_hour" }, + session_id: "sdk-session-limit", + uuid: "limit-rejected", + } as unknown as SDKMessage); + } + if (index < 2) { + harness.query.emit({ + ...rateLimitAssistant, + uuid: `assistant-limit-${index}`, + } as unknown as SDKMessage); + } + harness.query.emit({ + ...rateLimitResult, + uuid: `result-limit-${index}`, + } as unknown as SDKMessage); + const payload = completedTurn(Array.from(yield* Fiber.join(eventsFiber))); + assert.equal(payload.state, "failed"); + assert.equal(payload.errorMessage, expected); + } + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + it.effect.each([ + { homePath: "./synthetic config's $literal", inherited: undefined }, + { homePath: "", inherited: ".synthetic config's $literal" }, + { homePath: "", inherited: " /synthetic/path with edge spaces " }, + ])( + "reports the same Claude config and cwd used by the spawned query ($homePath, $inherited)", + ({ homePath, inherited }) => { + const harness = makeHarness({ + claudeConfig: { homePath }, + environment: { ...process.env, CLAUDE_CONFIG_DIR: inherited }, + }); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const eventsFiber = yield* adapter.streamEvents.pipe( + Stream.takeUntil((event) => event.type === "turn.completed"), + Stream.runCollect, + Effect.forkChild, + ); + const cwd = NodePath.resolve("/tmp/synthetic-audit-project"); + const session = yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + cwd, + }); + yield* adapter.sendTurn({ + threadId: session.threadId, + input: "synthetic", + attachments: [], + }); + harness.query.emit(AUTH_FAILURE_ASSISTANT); + harness.query.emit({ + type: "result", + subtype: "success", + is_error: false, + terminal_reason: "api_error", + errors: [], + session_id: "sdk-session-auth", + uuid: "result-auth", + } as unknown as SDKMessage); + const events = Array.from(yield* Fiber.join(eventsFiber)); + const completed = events.at(-1); + assert(completed?.type === "turn.completed"); + const actualQuery = harness.getLastCreateQueryInput(); + assert(actualQuery !== undefined); + const expectedConfigDir = homePath ? NodePath.resolve(homePath) : inherited; + assert(expectedConfigDir !== undefined); + assert.equal(actualQuery.options.env?.CLAUDE_CONFIG_DIR, expectedConfigDir); + assert.equal(actualQuery.options.cwd, cwd); + assert( + completed.payload.errorMessage?.includes( + `CLAUDE_CONFIG_DIR set to ${encodeUnknownJsonString(expectedConfigDir)}`, + ), + ); + assert(completed.payload.errorMessage?.includes(`from ${encodeUnknownJsonString(cwd)}`)); + assert(!completed.payload.errorMessage?.includes("CLAUDE_CONFIG_DIR=")); + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }, + ); + it.effect("interruptTurn settles live tasks and closes the provider session", () => { const harness = makeHarness(); return Effect.gen(function* () { @@ -1861,12 +2427,8 @@ describe("ClaudeAdapterLive", () => { ); }); - it.effect("keeps a resumed replacement session during slow stop cleanup", () => { + it.effect("keeps a resumed replacement session after interrupt cleanup", () => { const queries: FakeClaudeQuery[] = []; - let signalUsageStarted: () => void = () => undefined; - const usageStarted = new Promise((resolve) => { - signalUsageStarted = resolve; - }); const layer = Layer.effect( ClaudeAdapter, Effect.gen(function* () { @@ -1874,14 +2436,6 @@ describe("ClaudeAdapterLive", () => { return yield* makeClaudeAdapter(claudeConfig, { createQuery: () => { const query = new FakeClaudeQuery(); - if (queries.length === 0) { - Object.assign(query, { - getContextUsage: async () => { - signalUsageStarted(); - return await new Promise(() => undefined); - }, - }); - } queries.push(query); return query; }, @@ -1895,7 +2449,9 @@ describe("ClaudeAdapterLive", () => { return Effect.gen(function* () { const adapter = yield* ClaudeAdapter; - const runtimeEventsFiber = yield* Stream.take(adapter.streamEvents, 8).pipe( + const runtimeEventsFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.type.startsWith("session.")), + Stream.take(7), Stream.runCollect, Effect.forkChild, ); @@ -1910,10 +2466,7 @@ describe("ClaudeAdapterLive", () => { attachments: [], }); - const interruptFiber = yield* adapter - .interruptTurn(firstSession.threadId) - .pipe(Effect.forkChild); - yield* Effect.promise(() => usageStarted); + yield* adapter.interruptTurn(firstSession.threadId); assert.equal(queries[0]?.closeCalls, 1); const replacement = yield* adapter.startSession({ @@ -1922,8 +2475,6 @@ describe("ClaudeAdapterLive", () => { runtimeMode: "full-access", resumeCursor: firstSession.resumeCursor, }); - yield* TestClock.adjust("1 second"); - yield* Fiber.join(interruptFiber); const activeSessions = yield* adapter.listSessions(); const runtimeEvents = Array.from(yield* Fiber.join(runtimeEventsFiber)); @@ -1939,6 +2490,7 @@ describe("ClaudeAdapterLive", () => { "session.started", "session.configured", "session.state.changed", + "session.exited", "session.started", "session.configured", "session.state.changed", diff --git a/apps/server/src/provider/Layers/ClaudeAdapter.ts b/apps/server/src/provider/Layers/ClaudeAdapter.ts index fb7c435703cf..f7649ebaa677 100644 --- a/apps/server/src/provider/Layers/ClaudeAdapter.ts +++ b/apps/server/src/provider/Layers/ClaudeAdapter.ts @@ -14,7 +14,7 @@ import { type PermissionResult, type PermissionUpdate, type SDKMessage, - type SDKControlGetContextUsageResponse, + type SDKRateLimitInfo, type SDKResultMessage, type SettingSource, type SDKUserMessage, @@ -69,7 +69,6 @@ import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; import * as FileSystem from "effect/FileSystem"; import * as Fiber from "effect/Fiber"; -import * as Option from "effect/Option"; import * as Path from "effect/Path"; import * as Queue from "effect/Queue"; import * as Ref from "effect/Ref"; @@ -86,7 +85,7 @@ import { import * as McpProviderSession from "../../mcp/McpProviderSession.ts"; import { toClaudeMcpServers } from "../McpServerConfig.ts"; import { resolveClaudeSdkExecutablePath } from "../Drivers/ClaudeExecutable.ts"; -import { makeClaudeEnvironment } from "../Drivers/ClaudeHome.ts"; +import { claudeSignedOutMessage, makeClaudeEnvironment } from "../Drivers/ClaudeHome.ts"; import { getClaudeModelCapabilities, isClaudeUltracodeEffort, @@ -151,7 +150,34 @@ interface ClaudeTurnState { readonly assistantTextBlocks: Map; readonly assistantTextBlockOrder: Array; readonly capturedProposedPlanKeys: Set; + latestAssistantUsage: unknown | undefined; + compactedSinceLatestAssistantUsage: boolean; nextSyntheticAssistantBlockIndex: number; + authenticationFailureMessage: string | undefined; + rejectedRateLimitTypes: Set; + latestAssistantRateLimited: boolean; +} + +function createClaudeTurnState( + turnId: TurnId, + startedAt: string, + extra?: { readonly synthetic?: true }, +): ClaudeTurnState { + return { + turnId, + startedAt, + ...(extra?.synthetic ? { synthetic: true } : {}), + items: [], + assistantTextBlocks: new Map(), + assistantTextBlockOrder: [], + capturedProposedPlanKeys: new Set(), + latestAssistantUsage: undefined, + compactedSinceLatestAssistantUsage: false, + nextSyntheticAssistantBlockIndex: -1, + authenticationFailureMessage: undefined, + rejectedRateLimitTypes: new Set(), + latestAssistantRateLimited: false, + }; } interface AssistantTextBlockState { @@ -320,6 +346,7 @@ interface ClaudeSessionContext { lastKnownTotalProcessedTokens: number | undefined; lastAssistantUuid: string | undefined; lastThreadStartedId: string | undefined; + announcedUsageLimits: { turnId: string; keys: Set } | undefined; stopped: boolean; } @@ -327,7 +354,6 @@ interface ClaudeQueryRuntime extends AsyncIterable { readonly setModel: (model?: string) => Promise; readonly setPermissionMode: (mode: PermissionMode) => Promise; readonly setMaxThinkingTokens: (maxThinkingTokens: number | null) => Promise; - readonly getContextUsage?: () => Promise; readonly close: () => void; } @@ -421,16 +447,37 @@ function resultErrorsText(result: SDKResultMessage): string { : ""; } -/** - * First user-facing error from a non-success result. "[ede_diagnostic] ..." - * entries are CLI-internal telemetry (the CLI hides them from its own UI too), - * so they must never become the error banner. - */ -function resultUserFacingError(result: SDKResultMessage): string | undefined { - if (result.subtype === "success" || !Array.isArray(result.errors)) { - return undefined; +/** Failure text for structured terminal reasons, including success-tagged failures. */ +function terminalResultError( + reason: SDKResultMessage["terminal_reason"] | string | undefined, + failureHint?: string, +): string | undefined { + switch (reason) { + case "api_error": + return failureHint ?? "Claude gave up after repeated API errors."; + case "malformed_tool_use_exhausted": + return "Claude gave up after repeated malformed tool calls."; + case "budget_exhausted": + return "Claude stopped: the turn's token budget was exhausted."; + case "structured_output_retry_exhausted": + return "Claude could not produce the requested structured output."; + case "tool_deferred_unavailable": + return "Claude could not resume a deferred tool call: the tool is no longer available."; + case "turn_setup_failed": + return "Claude could not start the turn."; + case "blocking_limit": + return "Claude stopped: a usage limit blocked the request."; + case "rapid_refill_breaker": + return "Claude stopped: the context refilled too quickly after compaction."; + case "prompt_too_long": + return "Claude stopped: the prompt exceeds the model's context window."; + case "image_error": + return "Claude stopped: an image in the conversation could not be processed."; + case "model_error": + return "Claude stopped: the model returned an error."; + default: + return undefined; } - return result.errors.find((error) => !error.startsWith("[ede_diagnostic]")); } function isInterruptedResult(result: SDKResultMessage): boolean { @@ -458,6 +505,45 @@ function isInterruptedResult(result: SDKResultMessage): boolean { ); } +const CLAUDE_USAGE_LIMIT_WINDOWS = { + five_hour: "5-hour", + seven_day: "7-day", + seven_day_opus: "7-day Opus", + seven_day_sonnet: "7-day Sonnet", + overage: "overage", +} satisfies Record, string>; + +/** Beyond this the reset time is not credible, so the row ships without a wait. */ +const CLAUDE_USAGE_LIMIT_MAX_WAIT_MS = 30 * 24 * 60 * 60 * 1000; + +/** + * `resetsAt` is epoch seconds. The row states the remaining wait rather than a + * wall-clock time: this renders on the server, while the row is read on clients + * that may sit in another timezone and locale, and that carry their own + * timestamp preference. A wait reads the same everywhere. + */ +function describeClaudeUsageLimit(info: SDKRateLimitInfo, nowMs: number): string { + const label = info.rateLimitType ? CLAUDE_USAGE_LIMIT_WINDOWS[info.rateLimitType] : undefined; + const resetsAtMs = info.resetsAt === undefined ? undefined : info.resetsAt * 1000; + const waitMs = + resetsAtMs === undefined || !Number.isFinite(nowMs) ? undefined : resetsAtMs - nowMs; + const wait = + waitMs !== undefined && waitMs > 0 && waitMs <= CLAUDE_USAGE_LIMIT_MAX_WAIT_MS + ? formatClaudeUsageLimitWait(waitMs) + : undefined; + return `Claude usage limit reached. This turn is paused until the ${ + label ? `${label} ` : "" + }limit resets${wait ? ` in ${wait}` : ""}.`; +} + +function formatClaudeUsageLimitWait(waitMs: number): string { + const totalMinutes = Math.ceil(waitMs / 60_000); + const hours = Math.floor(totalMinutes / 60); + const minutes = totalMinutes % 60; + if (hours === 0) return `${totalMinutes}m`; + return minutes === 0 ? `${hours}h` : `${hours}h ${minutes}m`; +} + function asRuntimeItemId(value: string): RuntimeItemId { return RuntimeItemId.make(value); } @@ -614,20 +700,6 @@ function normalizeClaudeActiveTokenUsage( }); } -function normalizeClaudeContextUsageApiSnapshot( - value: SDKControlGetContextUsageResponse, - totalProcessedTokens?: number, -): ThreadTokenUsageSnapshot | undefined { - const autoCompactThreshold = finitePositiveInteger(value.autoCompactThreshold); - return makeClaudeTokenUsageSnapshot({ - activeTokens: value.totalTokens, - contextWindow: value.maxTokens, - ...(totalProcessedTokens !== undefined ? { totalProcessedTokens } : {}), - compactsAutomatically: value.isAutoCompactEnabled, - ...(autoCompactThreshold !== undefined ? { autoCompactThreshold } : {}), - }); -} - function compactBoundaryTokenUsageSnapshot( message: Record, contextWindow?: number, @@ -1359,19 +1431,49 @@ const buildUserMessageEffect = Effect.fn("buildUserMessageEffect")(function* ( return buildUserMessage({ sdkContent }); }); -function turnStatusFromResult(result: SDKResultMessage): ProviderRuntimeTurnStatus { - if (result.subtype === "success") { - return "completed"; - } +/** + * The CLI reports repeated 529 overload failures as a success-subtype result + * with api_error_status 529 and an empty error list; the status code is the + * only structured signal. + */ +function isOverloadedResult(result: SDKResultMessage): boolean { + return result.subtype === "success" && result.api_error_status === 529; +} - const errors = resultErrorsText(result); - if (isInterruptedResult(result)) { - return "interrupted"; - } - if (errors.includes("cancel")) { - return "cancelled"; - } - return "failed"; +/** Derives turn status and its error from the same provider result. */ +function resultOutcome( + result: SDKResultMessage, + failureHint?: string, +): { + status: ProviderRuntimeTurnStatus; + errorMessage: string | undefined; +} { + // A success result flagged is_error only fails when the turn already + // reported its cause (expired login, rejected usage window). + const successTaggedFailure = result.subtype === "success" && result.is_error === true; + const structuredError = isOverloadedResult(result) + ? "Claude API is overloaded (529). Try again shortly." + : (terminalResultError(result.terminal_reason, failureHint) ?? + (successTaggedFailure ? failureHint : undefined)); + // CLI diagnostic entries must not become the error banner. Success results + // carry no typed error list, but a success-tagged failure may still list one. + const listedErrors: ReadonlyArray = + "errors" in result && Array.isArray(result.errors) ? result.errors : []; + const listedError = + result.subtype === "success" && !successTaggedFailure + ? undefined + : listedErrors.find( + (error): error is string => + typeof error === "string" && !error.startsWith("[ede_diagnostic]"), + ); + const errorMessage = listedError || structuredError; + if (structuredError !== undefined) return { status: "failed", errorMessage }; + if (result.subtype === "success") return { status: "completed", errorMessage }; + if (isInterruptedResult(result)) return { status: "interrupted", errorMessage }; + return { + status: resultErrorsText(result).includes("cancel") ? "cancelled" : "failed", + errorMessage, + }; } function streamKindFromDeltaType(deltaType: string): ClaudeTextStreamKind { @@ -2137,29 +2239,6 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( }); }); - const queryCurrentContextUsage = Effect.fn("queryCurrentContextUsage")(function* ( - context: ClaudeSessionContext, - totalProcessedTokens?: number, - ) { - if (!context.query.getContextUsage) { - return undefined; - } - - const usage = yield* Effect.promise(async () => { - try { - return await context.query.getContextUsage?.(); - } catch { - return undefined; - } - }).pipe(Effect.timeoutOption("1 second")); - if (Option.isNone(usage) || !usage.value) { - return undefined; - } - - context.lastKnownContextWindow = usage.value.maxTokens; - return normalizeClaudeContextUsageApiSnapshot(usage.value, totalProcessedTokens); - }); - const emitProposedPlanCompleted = Effect.fn("emitProposedPlanCompleted")(function* ( context: ClaudeSessionContext, input: { @@ -2260,10 +2339,7 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( context.lastKnownTotalProcessedTokens = accumulatedTotalProcessedTokens; } - const contextUsageSnapshot = yield* queryCurrentContextUsage( - context, - accumulatedTotalProcessedTokens ?? context.lastKnownTotalProcessedTokens, - ); + // Avoid getContextUsage because its token-count fallback can make extra model requests. const resultUsageRecord = result?.usage && typeof result.usage === "object" && !Array.isArray(result.usage) ? (result.usage as Record) @@ -2285,24 +2361,31 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( accumulatedTotalProcessedTokens ?? context.lastKnownTotalProcessedTokens, ) : undefined; + const latestAssistantSnapshot = normalizeClaudeActiveTokenUsage( + context.turnState?.latestAssistantUsage, + maxTokens, + accumulatedTotalProcessedTokens ?? context.lastKnownTotalProcessedTokens, + ); const lastGoodUsage = context.lastKnownTokenUsage; const usageSnapshot: ThreadTokenUsageSnapshot | undefined = - contextUsageSnapshot ?? - (resultTotalOnly && lastGoodUsage - ? { - ...lastGoodUsage, - ...(typeof maxTokens === "number" && Number.isFinite(maxTokens) && maxTokens > 0 - ? { maxTokens } - : {}), - ...(typeof accumulatedTotalProcessedTokens === "number" && - Number.isFinite(accumulatedTotalProcessedTokens) && - accumulatedTotalProcessedTokens > lastGoodUsage.usedTokens - ? { - totalProcessedTokens: accumulatedTotalProcessedTokens, - } - : {}), - } - : resultIterationSnapshot) ?? + latestAssistantSnapshot ?? + (context.turnState?.compactedSinceLatestAssistantUsage + ? undefined + : resultTotalOnly && lastGoodUsage + ? { + ...lastGoodUsage, + ...(typeof maxTokens === "number" && Number.isFinite(maxTokens) && maxTokens > 0 + ? { maxTokens } + : {}), + ...(typeof accumulatedTotalProcessedTokens === "number" && + Number.isFinite(accumulatedTotalProcessedTokens) && + accumulatedTotalProcessedTokens > lastGoodUsage.usedTokens + ? { + totalProcessedTokens: accumulatedTotalProcessedTokens, + } + : {}), + } + : resultIterationSnapshot) ?? (lastGoodUsage ? { ...lastGoodUsage, @@ -2946,16 +3029,7 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( if (!context.turnState) { const turnId = TurnId.make(yield* randomUUIDv4); const startedAt = yield* nowIso; - context.turnState = { - turnId, - startedAt, - synthetic: true, - items: [], - assistantTextBlocks: new Map(), - assistantTextBlockOrder: [], - capturedProposedPlanKeys: new Set(), - nextSyntheticAssistantBlockIndex: -1, - }; + context.turnState = createClaudeTurnState(turnId, startedAt, { synthetic: true }); context.session = { ...context.session, status: "running", @@ -3013,7 +3087,28 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( } if (context.turnState) { + // Limited retries may only carry an assistant error, without a new window + // event. Later parent responses replace this evidence if the turn recovers. + context.turnState.latestAssistantRateLimited = message.error === "rate_limit"; + // The CLI can report authentication failure before ending the turn as a + // generic API error, so retain that evidence for the result fallback. + if (message.error === "authentication_failed") { + context.turnState.authenticationFailureMessage = claudeSignedOutMessage({ + configDir: claudeEnvironment.CLAUDE_CONFIG_DIR, + cwd: path.resolve(context.session.cwd ?? "."), + }); + } context.turnState.items.push(message.message); + if ( + normalizeClaudeActiveTokenUsage( + message.message.usage, + context.lastKnownContextWindow, + context.lastKnownTotalProcessedTokens, + ) + ) { + context.turnState.latestAssistantUsage = message.message.usage; + context.turnState.compactedSinceLatestAssistantUsage = false; + } yield* backfillAssistantTextBlocksFromSnapshot(context, message); } @@ -3029,8 +3124,13 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( return; } - const status = turnStatusFromResult(message); - const errorMessage = resultUserFacingError(message); + const turn = context.turnState; + const failureHint = + turn?.authenticationFailureMessage ?? + (turn && (turn.rejectedRateLimitTypes.size > 0 || turn.latestAssistantRateLimited) + ? "Claude usage limit reached. Send the message again once the limit resets." + : undefined); + const { status, errorMessage } = resultOutcome(message, failureHint); if (status === "failed") { yield* emitRuntimeError(context, errorMessage ?? "Claude turn failed."); @@ -3178,6 +3278,10 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( }); return; case "compact_boundary": + if (context.turnState) { + context.turnState.latestAssistantUsage = undefined; + context.turnState.compactedSinceLatestAssistantUsage = true; + } yield* emitThreadTokenUsage( context, compactBoundaryTokenUsageSnapshot( @@ -3575,6 +3679,45 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( rateLimits: message, }, }); + const rateLimitInfo = message.rate_limit_info; + if (!rateLimitInfo) return; + // A rejected window parks the turn inside the SDK: no further messages + // arrive and no result lands, so without a row the thread just spins. + // Warnings (allowed_warning) still have headroom and stay quiet, an + // account spending provisioned overage keeps running despite the reject, + // and between turns there is no turn to report as paused. + const overageAllowed = + rateLimitInfo.overageStatus === "allowed" || + rateLimitInfo.overageStatus === "allowed_warning" || + rateLimitInfo.isUsingOverage === true || + rateLimitInfo.overageInUse === true; + const blocked = rateLimitInfo.status === "rejected" && !overageAllowed; + const limitType = rateLimitInfo.rateLimitType ?? "unknown"; + const limitKey = `${limitType}:${rateLimitInfo.resetsAt ?? "unknown"}`; + if (context.turnState) { + // Current blocking evidence is independent of whether its warning has + // already been shown. A recovery can omit or advance the reset time; + // its window type remains stable without clearing another window. + if (blocked) context.turnState.rejectedRateLimitTypes.add(limitType); + else if ( + rateLimitInfo.status === "allowed" || + rateLimitInfo.status === "allowed_warning" || + overageAllowed + ) { + context.turnState.rejectedRateLimitTypes.delete(limitType); + } + } + if (blocked && context.turnState !== undefined) { + const turnId = context.turnState.turnId; + if (context.announcedUsageLimits?.turnId !== turnId) { + context.announcedUsageLimits = { turnId, keys: new Set() }; + } + if (!context.announcedUsageLimits.keys.has(limitKey)) { + context.announcedUsageLimits.keys.add(limitKey); + const notice = describeClaudeUsageLimit(rateLimitInfo, Date.parse(stamp.createdAt)); + yield* emitRuntimeWarning(context, notice, rateLimitInfo); + } + } return; } }); @@ -4441,6 +4584,7 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( lastKnownTotalProcessedTokens: undefined, lastAssistantUuid: resumeState?.resumeSessionAt, lastThreadStartedId: undefined, + announcedUsageLimits: undefined, stopped: false, }; yield* Ref.set(contextRef, context); @@ -4578,15 +4722,7 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( const turnId = steeringTurnState?.turnId ?? TurnId.make(yield* randomUUIDv4); if (steeringTurnState === null) { - const turnState: ClaudeTurnState = { - turnId, - startedAt: yield* nowIso, - items: [], - assistantTextBlocks: new Map(), - assistantTextBlockOrder: [], - capturedProposedPlanKeys: new Set(), - nextSyntheticAssistantBlockIndex: -1, - }; + const turnState: ClaudeTurnState = createClaudeTurnState(turnId, yield* nowIso); const updatedAt = yield* nowIso; context.turnState = turnState; diff --git a/docs/user/providers-claude.md b/docs/user/providers-claude.md index 681551707a0d..7b44a96a4914 100644 --- a/docs/user/providers-claude.md +++ b/docs/user/providers-claude.md @@ -24,6 +24,8 @@ real provider request succeeds. - Select **Reconnect** after an expired or revoked login. - Disconnect the account to remove its stored credential from this environment. +If a Claude chat fails because the login expired, Akeru names that cause instead of a generic API error and tells you to run `claude auth login` on the environment machine. A usage-limit failure says the limit was reached and to send the message again after it resets. + Sign-in state belongs to one environment. Connect Claude again on each separate Akeru server that should use the account.