diff --git a/apps/mobile/src/features/settings/SettingsServerControlsRouteScreen.tsx b/apps/mobile/src/features/settings/SettingsServerControlsRouteScreen.tsx index 8aff022910d0..750727847846 100644 --- a/apps/mobile/src/features/settings/SettingsServerControlsRouteScreen.tsx +++ b/apps/mobile/src/features/settings/SettingsServerControlsRouteScreen.tsx @@ -24,6 +24,7 @@ import { } from "./components/SettingsEnvironmentFilterHeader"; import { SettingsSection } from "./components/SettingsSection"; import { SettingsControlRow } from "./components/SettingsControlRow"; +import { AutoResumeMessageField } from "./components/AutoResumeMessageField"; import { SettingsSwitchRow } from "./components/SettingsSwitchRow"; import { SettingsProjectOverridesSection } from "./components/SettingsProjectOverridesSection"; import { useSettingsEnvironmentFilter } from "./settings-environment-filter"; @@ -393,6 +394,44 @@ function ServerSettingsDetail(props: { readonly page: SettingsPage }) { onValueChange={(value) => write({ enableAgentBrowserAccess: value })} /> + + write({ autoResumeAfterUsageLimit: value })} + /> + + + write({ autoResumeMessage: value })} + /> + + + + write({ autoResumeDisablesFastMode: value })} + /> + + ) : null} diff --git a/apps/mobile/src/features/settings/components/AutoResumeMessageField.tsx b/apps/mobile/src/features/settings/components/AutoResumeMessageField.tsx new file mode 100644 index 000000000000..1e8c7a2a470b --- /dev/null +++ b/apps/mobile/src/features/settings/components/AutoResumeMessageField.tsx @@ -0,0 +1,32 @@ +import { useState } from "react"; + +import { AppTextInput } from "../../../components/AppText"; + +/** Commits on blur or submit so each keystroke does not fan out a settings write. */ +export function AutoResumeMessageField(props: { + readonly value: string; + readonly placeholder: string; + readonly disabled: boolean; + readonly onCommit: (value: string) => void; +}) { + const [draft, setDraft] = useState(null); + const commit = () => { + const next = draft; + setDraft(null); + if (!props.disabled && next !== null && next !== props.value) props.onCommit(next); + }; + return ( + + ); +} diff --git a/apps/mobile/src/features/threads/ThreadDetailScreen.tsx b/apps/mobile/src/features/threads/ThreadDetailScreen.tsx index fe3c98e7f200..865424712cbc 100644 --- a/apps/mobile/src/features/threads/ThreadDetailScreen.tsx +++ b/apps/mobile/src/features/threads/ThreadDetailScreen.tsx @@ -93,6 +93,7 @@ import { ComposerFeedback } from "./ComposerFeedback"; import { ComposerUsageLimits } from "./ComposerUsageLimits"; import { PendingUserInputCard } from "./PendingUserInputCard"; import { ThreadCreationFailedCard } from "./ThreadCreationFailedCard"; +import { UsageLimitResumeCard } from "./UsageLimitResumeCard"; import { FLOATING_WORKING_CONTROL_COVERAGE, FloatingWorkingControl, @@ -322,6 +323,9 @@ export const ThreadDetailScreen = memo(function ThreadDetailScreen(props: Thread const windowHeight = useWindowDimensions().height; const navigationHeaderHeight = useContext(HeaderHeightContext) || insets.top + IOS_NAV_BAR_HEIGHT; const agentLabel = `${props.selectedThread.modelSelection.instanceId} agent`; + const usageLimit = props.selectedThread.usageLimit; + const usageLimitWithReset = + usageLimit?.resetsAt != null ? { ...usageLimit, resetsAt: usageLimit.resetsAt } : null; const selectedThreadKey = scopedThreadKey(props.environmentId, props.selectedThread.id); useReadAloudLifecycle({ environmentId: props.environmentId, threadId: props.selectedThread.id }); const composerEditorRef = useRef(null); @@ -1007,6 +1011,19 @@ export const ThreadDetailScreen = memo(function ThreadDetailScreen(props: Thread /> ) : null} + {usageLimitWithReset && activeUserInputRequestId === null ? ( + + + + ) : null} {props.creationState?.kind === "failed" ? ( + + {scheduled ? `Resumes automatically ${resetTime}.` : `The usage limit resets ${resetTime}.`} + + + void setAutoResume({ + environmentId: props.environmentId, + input: { threadId: props.threadId, scheduled: !scheduled }, + }) + } + /> + + ); +} diff --git a/apps/server/integration/OrchestrationEngineHarness.integration.ts b/apps/server/integration/OrchestrationEngineHarness.integration.ts index 0df54be2f701..35eee0b51950 100644 --- a/apps/server/integration/OrchestrationEngineHarness.integration.ts +++ b/apps/server/integration/OrchestrationEngineHarness.integration.ts @@ -66,6 +66,7 @@ import { } from "../src/orchestration/Services/OrchestrationEngine.ts"; import { ThreadDeletionReactor } from "../src/orchestration/Services/ThreadDeletionReactor.ts"; import * as ThreadSettlementReactor from "../src/orchestration/ThreadSettlementReactor.ts"; +import * as UsageLimitResumeReactor from "../src/orchestration/UsageLimitResumeReactor.ts"; import * as PullRequestSyncReactor from "../src/orchestration/PullRequestSyncReactor.ts"; import * as ThreadPullRequestReactor from "../src/orchestration/ThreadPullRequestReactor.ts"; import { OrchestrationReactor } from "../src/orchestration/Services/OrchestrationReactor.ts"; @@ -408,6 +409,12 @@ export const makeOrchestrationIntegrationHarness = ( drain: Effect.void, }), ), + Layer.provideMerge( + Layer.succeed(UsageLimitResumeReactor.UsageLimitResumeReactor, { + start: () => Effect.void, + drain: Effect.void, + }), + ), Layer.provideMerge( Layer.succeed(PullRequestSyncReactor.PullRequestSyncReactor, { start: () => Effect.void, diff --git a/apps/server/src/orchestration/Layers/OrchestrationReactor.test.ts b/apps/server/src/orchestration/Layers/OrchestrationReactor.test.ts index d25c442ac8b2..4644acbd56b1 100644 --- a/apps/server/src/orchestration/Layers/OrchestrationReactor.test.ts +++ b/apps/server/src/orchestration/Layers/OrchestrationReactor.test.ts @@ -10,6 +10,7 @@ import { ProviderCommandReactor } from "../Services/ProviderCommandReactor.ts"; import { ProviderRuntimeIngestionService } from "../Services/ProviderRuntimeIngestion.ts"; import { ThreadDeletionReactor } from "../Services/ThreadDeletionReactor.ts"; import * as ThreadSettlementReactor from "../ThreadSettlementReactor.ts"; +import * as UsageLimitResumeReactor from "../UsageLimitResumeReactor.ts"; import * as PullRequestSyncReactor from "../PullRequestSyncReactor.ts"; import * as ThreadPullRequestReactor from "../ThreadPullRequestReactor.ts"; import { OrchestrationReactor } from "../Services/OrchestrationReactor.ts"; @@ -95,6 +96,15 @@ describe("OrchestrationReactor", () => { drain: Effect.void, }), ), + Layer.provideMerge( + Layer.succeed(UsageLimitResumeReactor.UsageLimitResumeReactor, { + start: () => { + started.push("usage-limit-resume-reactor"); + return Effect.void; + }, + drain: Effect.void, + }), + ), Layer.provideMerge( Layer.succeed(PullRequestSyncReactor.PullRequestSyncReactor, { start: () => { @@ -128,6 +138,7 @@ describe("OrchestrationReactor", () => { "thread-deletion-reactor", "thread-pull-request-reactor", "thread-settlement-reactor", + "usage-limit-resume-reactor", "pull-request-sync-reactor", "agent-awareness-relay", "storage-cleanup", diff --git a/apps/server/src/orchestration/Layers/OrchestrationReactor.ts b/apps/server/src/orchestration/Layers/OrchestrationReactor.ts index da5cc337e420..99894dd37533 100644 --- a/apps/server/src/orchestration/Layers/OrchestrationReactor.ts +++ b/apps/server/src/orchestration/Layers/OrchestrationReactor.ts @@ -10,6 +10,7 @@ import { ProviderCommandReactor } from "../Services/ProviderCommandReactor.ts"; import { ProviderRuntimeIngestionService } from "../Services/ProviderRuntimeIngestion.ts"; import { ThreadDeletionReactor } from "../Services/ThreadDeletionReactor.ts"; import * as ThreadSettlementReactor from "../ThreadSettlementReactor.ts"; +import * as UsageLimitResumeReactor from "../UsageLimitResumeReactor.ts"; import * as PullRequestSyncReactor from "../PullRequestSyncReactor.ts"; import * as ThreadPullRequestReactor from "../ThreadPullRequestReactor.ts"; import * as AgentAwarenessRelay from "../../relay/AgentAwarenessRelay.ts"; @@ -21,6 +22,7 @@ export const makeOrchestrationReactor = Effect.gen(function* () { const checkpointReactor = yield* CheckpointReactor; const threadDeletionReactor = yield* ThreadDeletionReactor; const threadSettlementReactor = yield* ThreadSettlementReactor.ThreadSettlementReactor; + const usageLimitResumeReactor = yield* UsageLimitResumeReactor.UsageLimitResumeReactor; const pullRequestSyncReactor = yield* PullRequestSyncReactor.PullRequestSyncReactor; const threadPullRequestReactor = yield* ThreadPullRequestReactor.ThreadPullRequestReactor; const agentAwarenessRelay = yield* AgentAwarenessRelay.AgentAwarenessRelay; @@ -33,6 +35,7 @@ export const makeOrchestrationReactor = Effect.gen(function* () { yield* threadDeletionReactor.start(); yield* threadPullRequestReactor.start(); yield* threadSettlementReactor.start(); + yield* usageLimitResumeReactor.start(); yield* pullRequestSyncReactor.start(); yield* agentAwarenessRelay.start(); yield* storageCleanup.start(); diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts index b5e6cb0cdd54..c4ebb0daa2ae 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts @@ -630,6 +630,7 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti unsettledAt: null, snoozedUntil: null, snoozedAt: null, + usageLimit: null, pinnedAt: null, pinOrderKey: null, activeOrderKey: null, @@ -748,6 +749,35 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti return; } + case "thread.usage-limit-set": { + const existingRow = yield* projectionThreadRepository.getById({ + threadId: event.payload.threadId, + }); + if (Option.isNone(existingRow)) { + return; + } + yield* projectionThreadRepository.upsert({ + ...existingRow.value, + usageLimit: event.payload.usageLimit, + }); + return; + } + + case "thread.auto-resume-set": { + const existingRow = yield* projectionThreadRepository.getById({ + threadId: event.payload.threadId, + }); + const usageLimit = Option.isSome(existingRow) ? existingRow.value.usageLimit : null; + if (Option.isNone(existingRow) || usageLimit == null) { + return; + } + yield* projectionThreadRepository.upsert({ + ...existingRow.value, + usageLimit: { ...usageLimit, resumeScheduled: event.payload.scheduled }, + }); + return; + } + case "thread.pinned": { const existingRow = yield* projectionThreadRepository.getById({ threadId: event.payload.threadId, diff --git a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts index 843eb8343d84..4e68d7648f99 100644 --- a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts +++ b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts @@ -482,6 +482,7 @@ projectionSnapshotLayer("ProjectionSnapshotQuery", (it) => { unsettledAt: null, snoozedUntil: null, snoozedAt: null, + usageLimit: null, pinnedAt: "2026-02-24T00:00:01.000Z", pinOrderKey: "gm", activeOrderKey: "hq", @@ -608,6 +609,7 @@ projectionSnapshotLayer("ProjectionSnapshotQuery", (it) => { unsettledAt: null, snoozedUntil: null, snoozedAt: null, + usageLimit: null, pinnedAt: "2026-02-24T00:00:01.000Z", pinOrderKey: "gm", activeOrderKey: "hq", @@ -744,6 +746,7 @@ projectionSnapshotLayer("ProjectionSnapshotQuery", (it) => { projectId: asProjectId("project-1"), title: "Thread 1", titleState: null, + usageLimit: null, session: snapshot.threads[0]?.session ?? null, }); } diff --git a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts index 1e7058742e25..0a4894a1caba 100644 --- a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts +++ b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts @@ -31,6 +31,7 @@ import { ProjectId, ThreadLinkedPullRequest, ThreadTitleState, + ThreadUsageLimit, ThreadId, ThreadPullRequestSnapshot, ThreadPullRequestStack, @@ -132,6 +133,7 @@ const ProjectionThreadDbRowSchema = ProjectionThread.mapFields( titleState: Schema.NullOr(Schema.fromJsonString(ThreadTitleState)), linkedPullRequest: Schema.NullOr(Schema.fromJsonString(ThreadLinkedPullRequest)), branchPullRequest: Schema.NullOr(Schema.fromJsonString(ThreadLinkedPullRequest)), + usageLimit: Schema.NullOr(Schema.fromJsonString(ThreadUsageLimit)), }), ); const ProjectionThreadActivityDbRowSchema = ProjectionThreadActivity.mapFields( @@ -146,6 +148,7 @@ const ProjectionThreadActivityIdRowSchema = Schema.Struct({ const ProjectionThreadSessionDbRowSchema = ProjectionThreadSession; const ProjectionThreadRuntimeContextDbRowSchema = Schema.Struct({ titleState: Schema.NullOr(Schema.fromJsonString(ThreadTitleState)), + usageLimit: Schema.NullOr(Schema.fromJsonString(ThreadUsageLimit)), id: ThreadId, projectId: ProjectId, title: Schema.String, @@ -584,6 +587,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { unsettled_at AS "unsettledAt", snoozed_until AS "snoozedUntil", snoozed_at AS "snoozedAt", + usage_limit_json AS "usageLimit", pinned_at AS "pinnedAt", pin_order_key AS "pinOrderKey", active_order_key AS "activeOrderKey", @@ -625,6 +629,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { unsettled_at AS "unsettledAt", snoozed_until AS "snoozedUntil", snoozed_at AS "snoozedAt", + usage_limit_json AS "usageLimit", pinned_at AS "pinnedAt", pin_order_key AS "pinOrderKey", active_order_key AS "activeOrderKey", @@ -698,6 +703,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { unsettled_at AS "unsettledAt", snoozed_until AS "snoozedUntil", snoozed_at AS "snoozedAt", + usage_limit_json AS "usageLimit", pinned_at AS "pinnedAt", pin_order_key AS "pinOrderKey", active_order_key AS "activeOrderKey", @@ -1263,6 +1269,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { unsettled_at AS "unsettledAt", snoozed_until AS "snoozedUntil", snoozed_at AS "snoozedAt", + usage_limit_json AS "usageLimit", pinned_at AS "pinnedAt", pin_order_key AS "pinOrderKey", active_order_key AS "activeOrderKey", @@ -1291,6 +1298,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { threads.project_id AS "projectId", threads.title, threads.title_state_json AS "titleState", + threads.usage_limit_json AS "usageLimit", sessions.thread_id AS "threadId", sessions.status, sessions.provider_name AS "providerName", @@ -1313,6 +1321,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { projectId: row.projectId, title: row.title, titleState: row.titleState, + usageLimit: row.usageLimit, session: row.threadId === null ? null : row, })), ), @@ -2338,6 +2347,7 @@ pending_approval_requests AS ( unsettledAt: row.unsettledAt, snoozedUntil: row.snoozedUntil, snoozedAt: row.snoozedAt, + usageLimit: row.usageLimit ?? null, pinnedAt: row.pinnedAt, pinOrderKey: row.pinOrderKey ?? null, activeOrderKey: row.activeOrderKey ?? null, @@ -2583,6 +2593,7 @@ pending_approval_requests AS ( unsettledAt: row.unsettledAt, snoozedUntil: row.snoozedUntil, snoozedAt: row.snoozedAt, + usageLimit: row.usageLimit ?? null, pinnedAt: row.pinnedAt, pinOrderKey: row.pinOrderKey ?? null, activeOrderKey: row.activeOrderKey ?? null, @@ -2739,6 +2750,7 @@ pending_approval_requests AS ( unsettledAt: row.unsettledAt, snoozedUntil: row.snoozedUntil, snoozedAt: row.snoozedAt, + usageLimit: row.usageLimit ?? null, pinnedAt: row.pinnedAt, pinOrderKey: row.pinOrderKey ?? null, activeOrderKey: row.activeOrderKey ?? null, @@ -2902,6 +2914,7 @@ pending_approval_requests AS ( unsettledAt: row.unsettledAt, snoozedUntil: row.snoozedUntil, snoozedAt: row.snoozedAt, + usageLimit: row.usageLimit ?? null, pinnedAt: row.pinnedAt, pinOrderKey: row.pinOrderKey ?? null, activeOrderKey: row.activeOrderKey ?? null, @@ -3258,6 +3271,7 @@ pending_approval_requests AS ( unsettledAt: threadRow.value.unsettledAt, snoozedUntil: threadRow.value.snoozedUntil, snoozedAt: threadRow.value.snoozedAt, + usageLimit: threadRow.value.usageLimit ?? null, pinnedAt: threadRow.value.pinnedAt, pinOrderKey: threadRow.value.pinOrderKey ?? null, activeOrderKey: threadRow.value.activeOrderKey ?? null, @@ -3290,6 +3304,7 @@ pending_approval_requests AS ( projectId: row.projectId, title: row.title, titleState: row.titleState, + usageLimit: row.usageLimit, session: row.session === null ? null : mapSessionRow(row.session), })); }); @@ -3559,6 +3574,7 @@ pending_approval_requests AS ( unsettledAt: threadRow.value.unsettledAt, snoozedUntil: threadRow.value.snoozedUntil, snoozedAt: threadRow.value.snoozedAt, + usageLimit: threadRow.value.usageLimit ?? null, pinnedAt: threadRow.value.pinnedAt, pinOrderKey: threadRow.value.pinOrderKey ?? null, activeOrderKey: threadRow.value.activeOrderKey ?? null, diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index 0db70e491235..3542b112f1b1 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -54,6 +54,8 @@ import { import { projectActivityPayload } from "../ActivityPayloadProjection.ts"; import { forkParked } from "../../serverActivation.ts"; import { ServerSettingsService } from "../../serverSettings.ts"; +import { ProviderRegistry } from "../../provider/Services/ProviderRegistry.ts"; +import { latestExhaustedWindowResetAt } from "../../provider/providerUsageLimits.ts"; import { resolveProjectSettings } from "@t3tools/shared/projectSettings"; import { canReplaceThreadTitle } from "../threadTitles.ts"; @@ -1029,6 +1031,7 @@ const make = Effect.gen(function* () { const projectionThreadActivityRepository = yield* ProjectionThreadActivityRepository; const serverSettingsService = yield* ServerSettingsService; const checkpointStore = yield* CheckpointStore.CheckpointStore; + const providerRegistry = yield* Effect.serviceOption(ProviderRegistry); const providerCommandId = (event: ProviderRuntimeEvent, tag: string) => crypto.randomUUIDv4.pipe( Effect.map((uuid) => CommandId.make(`provider:${event.eventId}:${tag}:${uuid}`)), @@ -1755,6 +1758,25 @@ const make = Effect.gen(function* () { }, ); + const laterResetsAt = (current: string | null, next: string | null) => + current !== null && (next === null || Date.parse(current) > Date.parse(next)) ? current : next; + + const resolveUsageLimitResetsAt = ( + event: Extract, + ) => + Effect.gen(function* () { + const reported = event.payload.resetsAt; + if (reported !== undefined && Date.parse(reported) > Date.parse(event.createdAt)) { + return reported; + } + if (Option.isNone(providerRegistry) || event.providerInstanceId === undefined) { + return null; + } + const providers = yield* providerRegistry.value.getProviders; + const provider = providers.find((entry) => entry.instanceId === event.providerInstanceId); + return latestExhaustedWindowResetAt(provider?.usageLimits, event.createdAt) ?? null; + }); + const processRuntimeEvent = (event: ProviderRuntimeEvent) => Effect.gen(function* () { if ( @@ -1934,6 +1956,39 @@ const make = Effect.gen(function* () { } } + if (event.type === "turn.usage-limited") { + yield* orchestrationEngine.dispatch({ + type: "thread.usage-limit.set", + commandId: yield* providerCommandId(event, "usage-limit-set"), + threadId: thread.id, + usageLimit: { + reachedAt: thread.usageLimit?.reachedAt ?? now, + // A turn parked on several windows resumes only once the last one reopens. + resetsAt: laterResetsAt( + thread.usageLimit?.resetsAt ?? null, + yield* resolveUsageLimitResetsAt(event), + ), + }, + createdAt: now, + }); + } else if ( + thread.usageLimit != null && + ((event.type === "content.delta" && event.payload.streamKind === "assistant_text") || + (event.type === "turn.completed" && + shouldApplyThreadLifecycle && + normalizeRuntimeTurnState(event.payload.state) === "completed")) + ) { + // A parked turn the provider carried on by itself no longer needs a resume, + // and must not be interrupted by one. + yield* orchestrationEngine.dispatch({ + type: "thread.usage-limit.set", + commandId: yield* providerCommandId(event, "usage-limit-clear"), + threadId: thread.id, + usageLimit: null, + createdAt: now, + }); + } + const assistantDelta = event.type === "content.delta" && event.payload.streamKind === "assistant_text" ? event.payload.delta diff --git a/apps/server/src/orchestration/Schemas.ts b/apps/server/src/orchestration/Schemas.ts index 29468dc3f84e..aab73ae72a37 100644 --- a/apps/server/src/orchestration/Schemas.ts +++ b/apps/server/src/orchestration/Schemas.ts @@ -12,6 +12,8 @@ import { ThreadUnarchivedPayload as ContractsThreadUnarchivedPayloadSchema, ThreadUnsettledPayload as ContractsThreadUnsettledPayloadSchema, ThreadSnoozedPayload as ContractsThreadSnoozedPayloadSchema, + ThreadUsageLimitSetPayload as ContractsThreadUsageLimitSetPayloadSchema, + ThreadAutoResumeSetPayload as ContractsThreadAutoResumeSetPayloadSchema, ThreadUnsnoozedPayload as ContractsThreadUnsnoozedPayloadSchema, ThreadPinnedPayload as ContractsThreadPinnedPayloadSchema, ThreadUnpinnedPayload as ContractsThreadUnpinnedPayloadSchema, @@ -47,6 +49,8 @@ export const ThreadDeletedPayload = ContractsThreadDeletedPayloadSchema; export const ThreadUnarchivedPayload = ContractsThreadUnarchivedPayloadSchema; export const ThreadUnsettledPayload = ContractsThreadUnsettledPayloadSchema; export const ThreadSnoozedPayload = ContractsThreadSnoozedPayloadSchema; +export const ThreadUsageLimitSetPayload = ContractsThreadUsageLimitSetPayloadSchema; +export const ThreadAutoResumeSetPayload = ContractsThreadAutoResumeSetPayloadSchema; export const ThreadUnsnoozedPayload = ContractsThreadUnsnoozedPayloadSchema; export const ThreadPinnedPayload = ContractsThreadPinnedPayloadSchema; export const ThreadUnpinnedPayload = ContractsThreadUnpinnedPayloadSchema; diff --git a/apps/server/src/orchestration/Services/ProjectionSnapshotQuery.ts b/apps/server/src/orchestration/Services/ProjectionSnapshotQuery.ts index eac3ede9c1ee..024878961011 100644 --- a/apps/server/src/orchestration/Services/ProjectionSnapshotQuery.ts +++ b/apps/server/src/orchestration/Services/ProjectionSnapshotQuery.ts @@ -236,7 +236,10 @@ export interface ProjectionSnapshotQueryShape { threadId: ThreadId, ) => Effect.Effect< Option.Option< - Pick + Pick< + OrchestrationThreadShell, + "id" | "projectId" | "title" | "titleState" | "usageLimit" | "session" + > >, ProjectionRepositoryError >; diff --git a/apps/server/src/orchestration/UsageLimitResumePolicy.test.ts b/apps/server/src/orchestration/UsageLimitResumePolicy.test.ts new file mode 100644 index 000000000000..6807c152c667 --- /dev/null +++ b/apps/server/src/orchestration/UsageLimitResumePolicy.test.ts @@ -0,0 +1,146 @@ +import { describe, expect, it } from "vite-plus/test"; +import { + MessageId, + ProviderInstanceId, + ThreadId, + type ModelSelection, + type OrchestrationSession, +} from "@t3tools/contracts"; + +import { + autoResumeMessageId, + countTrailingAutoResumes, + resolveResumeAction, + shouldAutoScheduleResume, + withoutFastMode, +} from "./UsageLimitResumePolicy.ts"; + +const RESETS_AT = "2026-09-24T13:50:00.000Z"; + +const session = (status: OrchestrationSession["status"]): OrchestrationSession => ({ + threadId: ThreadId.make("thread-1"), + status, + providerName: "claudeAgent", + runtimeMode: "full-access", + activeTurnId: null, + lastError: null, + updatedAt: "2026-09-24T11:00:00.000Z", +}); + +const scheduled = { + reachedAt: "2026-09-24T11:00:00.000Z", + resetsAt: RESETS_AT, + resumeScheduled: true, +}; + +describe("resolveResumeAction", () => { + it("does nothing for threads without a scheduled resume", () => { + expect(resolveResumeAction({ usageLimit: null, session: null }, RESETS_AT)).toBeNull(); + expect( + resolveResumeAction( + { usageLimit: { ...scheduled, resumeScheduled: false }, session: session("error") }, + "2026-09-24T15:00:00.000Z", + ), + ).toBeNull(); + }); + + it("waits until a grace period after the reset", () => { + expect( + resolveResumeAction({ usageLimit: scheduled, session: session("error") }, RESETS_AT), + ).toBe("wait"); + expect( + resolveResumeAction( + { usageLimit: scheduled, session: session("error") }, + "2026-09-24T13:51:00.000Z", + ), + ).toBe("send"); + }); + + it("stops a turn still parked on the limit before sending", () => { + expect( + resolveResumeAction( + { usageLimit: scheduled, session: session("running") }, + "2026-09-24T14:00:00.000Z", + ), + ).toBe("interrupt"); + expect( + resolveResumeAction( + { usageLimit: scheduled, session: session("starting") }, + "2026-09-24T14:00:00.000Z", + ), + ).toBe("wait"); + }); +}); + +describe("shouldAutoScheduleResume", () => { + const limit = { ...scheduled, resumeScheduled: false }; + + it("schedules a known reset when the setting is on", () => { + expect( + shouldAutoScheduleResume({ enabled: true, usageLimit: limit, trailingAutoResumes: 0 }), + ).toBe(true); + expect( + shouldAutoScheduleResume({ enabled: false, usageLimit: limit, trailingAutoResumes: 0 }), + ).toBe(false); + expect( + shouldAutoScheduleResume({ + enabled: true, + usageLimit: { ...limit, resetsAt: null }, + trailingAutoResumes: 0, + }), + ).toBe(false); + }); + + it("retries once, then leaves the thread for the user", () => { + expect( + shouldAutoScheduleResume({ enabled: true, usageLimit: limit, trailingAutoResumes: 1 }), + ).toBe(true); + expect( + shouldAutoScheduleResume({ enabled: true, usageLimit: limit, trailingAutoResumes: 2 }), + ).toBe(false); + }); +}); + +describe("countTrailingAutoResumes", () => { + const user = (id: string) => ({ id: MessageId.make(id), role: "user" as const }); + const assistant = (id: string) => ({ id: MessageId.make(id), role: "assistant" as const }); + + it("counts only the unbroken run of resume messages at the end", () => { + expect( + countTrailingAutoResumes([ + user(autoResumeMessageId("a")), + user("manual"), + assistant("reply"), + user(autoResumeMessageId("b")), + assistant("reply-2"), + user(autoResumeMessageId("c")), + ]), + ).toBe(2); + expect(countTrailingAutoResumes([user(autoResumeMessageId("a")), user("manual")])).toBe(0); + }); +}); + +describe("withoutFastMode", () => { + const base = { instanceId: ProviderInstanceId.make("claudeAgent"), model: "claude-opus" }; + + it("turns fast mode off and drops a fast service tier", () => { + const selection: ModelSelection = { + ...base, + options: [ + { id: "fastMode", value: true }, + { id: "serviceTier", value: "fast" }, + { id: "effort", value: "high" }, + ], + }; + expect(withoutFastMode(selection).options).toEqual([ + { id: "fastMode", value: false }, + { id: "effort", value: "high" }, + ]); + }); + + it("returns selections without fast mode unchanged", () => { + const flex: ModelSelection = { ...base, options: [{ id: "serviceTier", value: "flex" }] }; + expect(withoutFastMode(flex)).toBe(flex); + expect(withoutFastMode(base)).toBe(base); + }); +}); diff --git a/apps/server/src/orchestration/UsageLimitResumePolicy.ts b/apps/server/src/orchestration/UsageLimitResumePolicy.ts new file mode 100644 index 000000000000..bd1d06e33b9f --- /dev/null +++ b/apps/server/src/orchestration/UsageLimitResumePolicy.ts @@ -0,0 +1,90 @@ +import { + MessageId, + type ModelSelection, + type OrchestrationMessage, + type OrchestrationThreadShell, + type ThreadUsageLimit, +} from "@t3tools/contracts"; + +/** Reset times are rounded by providers; resuming on the dot can hit the old window. */ +const RESUME_GRACE_MS = 60_000; + +/** Auto-scheduled resumes allowed in a row: the resume, then one retry. */ +const MAX_TRAILING_AUTO_RESUMES = 2; + +const AUTO_RESUME_MESSAGE_ID_PREFIX = "auto-resume:"; + +const DEFAULT_RESUME_TEXT = "go on"; + +/** Resume messages carry a marked id so the retry cap survives restarts. */ +export function autoResumeMessageId(uuid: string): MessageId { + return MessageId.make(`${AUTO_RESUME_MESSAGE_ID_PREFIX}${uuid}`); +} + +export function countTrailingAutoResumes( + messages: ReadonlyArray>, +): number { + let count = 0; + for (let index = messages.length - 1; index >= 0; index -= 1) { + const message = messages[index]!; + if (message.role !== "user") continue; + if (!message.id.startsWith(AUTO_RESUME_MESSAGE_ID_PREFIX)) break; + count += 1; + } + return count; +} + +export function shouldAutoScheduleResume(input: { + readonly enabled: boolean; + readonly usageLimit: ThreadUsageLimit; + readonly trailingAutoResumes: number; +}): boolean { + return ( + input.enabled && + input.usageLimit.resetsAt !== null && + !input.usageLimit.resumeScheduled && + input.trailingAutoResumes < MAX_TRAILING_AUTO_RESUMES + ); +} + +/** + * What a scheduled resume needs right now, or null when none is scheduled. + * A turn Claude parked on the limit is still running and must be stopped + * before the resume message can start a new one. + */ +export function resolveResumeAction( + thread: Pick, + now: string, +): "wait" | "interrupt" | "send" | null { + const usageLimit = thread.usageLimit; + if (!usageLimit?.resumeScheduled || usageLimit.resetsAt === null) return null; + if (!(Date.parse(now) >= Date.parse(usageLimit.resetsAt) + RESUME_GRACE_MS)) return "wait"; + switch (thread.session?.status) { + case "running": + return "interrupt"; + case "starting": + return "wait"; + default: + return "send"; + } +} + +export function autoResumeText(configured: string): string { + return configured.trim() || DEFAULT_RESUME_TEXT; +} + +/** Claude and Cursor use `fastMode`; Codex also expresses it as the `fast` service tier. */ +export function withoutFastMode(selection: ModelSelection): ModelSelection { + const isFast = (option: NonNullable[number]) => + (option.id === "fastMode" && option.value === true) || + (option.id === "serviceTier" && option.value === "fast"); + if (!selection.options?.some(isFast)) return selection; + return { + ...selection, + options: selection.options.flatMap((option) => { + if (option.id === "fastMode") return [{ ...option, value: false }]; + if (option.id === "serviceTier" && option.value === "fast") return []; + return [option]; + }), + }; +} diff --git a/apps/server/src/orchestration/UsageLimitResumeReactor.ts b/apps/server/src/orchestration/UsageLimitResumeReactor.ts new file mode 100644 index 000000000000..b5fe87689103 --- /dev/null +++ b/apps/server/src/orchestration/UsageLimitResumeReactor.ts @@ -0,0 +1,220 @@ +import { CommandId, type OrchestrationEvent, type ThreadId } from "@t3tools/contracts"; +import { makeDrainableWorker } from "@t3tools/shared/DrainableWorker"; +import * as Cause from "effect/Cause"; +import * as Context from "effect/Context"; +import * as Crypto from "effect/Crypto"; +import * as DateTime from "effect/DateTime"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; +import * as Schedule from "effect/Schedule"; +import type * as Scope from "effect/Scope"; +import * as Stream from "effect/Stream"; + +import * as ServerSettings from "../serverSettings.ts"; +import { forkParked } from "../serverActivation.ts"; +import * as OrchestrationEngine from "./Services/OrchestrationEngine.ts"; +import * as ProjectionSnapshotQuery from "./Services/ProjectionSnapshotQuery.ts"; +import { + autoResumeMessageId, + autoResumeText, + countTrailingAutoResumes, + resolveResumeAction, + shouldAutoScheduleResume, + withoutFastMode, +} from "./UsageLimitResumePolicy.ts"; + +/** + * Sends the configured resume message to threads stopped by a provider usage + * limit once the limit resets. Schedules live on the thread (`usageLimit`), so + * a restart picks them up on the first sweep. + */ +export class UsageLimitResumeReactor extends Context.Service< + UsageLimitResumeReactor, + { + readonly start: () => Effect.Effect; + readonly drain: Effect.Effect; + } +>()("t3/orchestration/UsageLimitResumeReactor") {} + +/** @public Service construction is part of the canonical Effect module API. */ +export const make = Effect.gen(function* () { + const engine = yield* OrchestrationEngine.OrchestrationEngineService; + const snapshots = yield* ProjectionSnapshotQuery.ProjectionSnapshotQuery; + const settingsService = yield* ServerSettings.ServerSettingsService; + const crypto = yield* Crypto.Crypto; + // One interrupt per stop: a parked turn that ignores it is left to the user. + const interruptedStops = new Set(); + // Filled from one full scan at start, then kept current from events, so the + // minute tick reads only threads with a pending resume. + const scheduledThreadIds = new Set(); + + const commandId = (tag: string, threadId: ThreadId) => + Effect.map(crypto.randomUUIDv4, (uuid) => + CommandId.make(`server:auto-resume:${tag}:${threadId}:${uuid}`), + ); + + const autoSchedule = Effect.fn("UsageLimitResumeReactor.autoSchedule")(function* ( + threadId: ThreadId, + ) { + const settings = yield* settingsService.getSettings; + if (!settings.autoResumeAfterUsageLimit) return; + const detail = yield* snapshots.getThreadDetailSnapshot(threadId, { turnLimit: 3 }); + if (Option.isNone(detail)) return; + const thread = detail.value.thread; + if ( + !thread.usageLimit || + !shouldAutoScheduleResume({ + enabled: true, + usageLimit: thread.usageLimit, + trailingAutoResumes: countTrailingAutoResumes(thread.messages), + }) + ) { + return; + } + yield* engine.dispatch({ + type: "thread.auto-resume.set", + commandId: yield* commandId("schedule", threadId), + threadId, + scheduled: true, + }); + }); + + const resume = Effect.fn("UsageLimitResumeReactor.resume")(function* (threadId: ThreadId) { + const shell = yield* snapshots.getThreadShellById(threadId); + if (Option.isNone(shell) || !shell.value.usageLimit?.resumeScheduled) { + scheduledThreadIds.delete(threadId); + return; + } + const thread = shell.value; + const now = DateTime.formatIso(yield* DateTime.now); + const action = resolveResumeAction(thread, now); + if (action === null || action === "wait" || !thread.usageLimit) return; + const stopKey = `${threadId}:${thread.usageLimit.reachedAt}`; + if (action === "interrupt") { + if (interruptedStops.has(stopKey)) return; + interruptedStops.add(stopKey); + yield* engine.dispatch({ + type: "thread.turn.interrupt", + commandId: yield* commandId("interrupt", threadId), + threadId, + ...(thread.session?.activeTurnId ? { turnId: thread.session.activeTurnId } : {}), + createdAt: now, + }); + return; + } + interruptedStops.delete(stopKey); + const settings = yield* settingsService.getSettings; + const modelSelection = settings.autoResumeDisablesFastMode + ? withoutFastMode(thread.modelSelection) + : thread.modelSelection; + if (modelSelection !== thread.modelSelection) { + // Saved on the thread so fast mode stays off after the resumed turn. + yield* engine.dispatch({ + type: "thread.meta.update", + commandId: yield* commandId("fast-mode-off", threadId), + threadId, + modelSelection, + }); + } + yield* engine.dispatch({ + type: "thread.turn.start", + commandId: yield* commandId("send", threadId), + threadId, + message: { + messageId: autoResumeMessageId(yield* crypto.randomUUIDv4), + role: "user", + text: autoResumeText(settings.autoResumeMessage), + attachments: [], + }, + modelSelection, + runtimeMode: thread.runtimeMode, + interactionMode: thread.interactionMode, + createdAt: now, + }); + }); + + const scan = Effect.fn("UsageLimitResumeReactor.scan")(function* () { + const snapshot = yield* snapshots.getShellSnapshot(); + for (const thread of snapshot.threads) { + if (thread.usageLimit?.resumeScheduled) scheduledThreadIds.add(thread.id); + } + }); + + const sweep = Effect.fn("UsageLimitResumeReactor.sweep")(function* () { + yield* Effect.forEach([...scheduledThreadIds], resume, { discard: true }); + }); + + type Job = + | { readonly kind: "sweep" } + | { readonly kind: "resume"; readonly threadId: ThreadId } + | { readonly kind: "auto-schedule"; readonly threadId: ThreadId }; + + const worker = yield* makeDrainableWorker((job: Job) => + (job.kind === "sweep" + ? sweep() + : job.kind === "resume" + ? resume(job.threadId) + : autoSchedule(job.threadId) + ).pipe( + Effect.catchCause((cause) => + Cause.hasInterruptsOnly(cause) + ? Effect.failCause(cause) + : Effect.logWarning("usage limit resume skipped", { + job, + cause: Cause.pretty(cause), + }), + ), + ), + ); + + const processEvent = (event: OrchestrationEvent) => { + switch (event.type) { + case "thread.usage-limit-set": + // Re-reports of a stop keep its reachedAt, so only a new stop is + // auto-scheduled and a Cancel on the current one sticks. + return event.payload.usageLimit !== null && + event.payload.usageLimit.reachedAt === event.occurredAt && + !event.payload.usageLimit.resumeScheduled + ? worker.enqueue({ kind: "auto-schedule", threadId: event.payload.threadId }) + : Effect.void; + case "thread.auto-resume-set": + if (!event.payload.scheduled) { + scheduledThreadIds.delete(event.payload.threadId); + return Effect.void; + } + scheduledThreadIds.add(event.payload.threadId); + return worker.enqueue({ kind: "resume", threadId: event.payload.threadId }); + case "thread.session-set": + if (!scheduledThreadIds.has(event.payload.threadId)) return Effect.void; + // The interrupted parked turn has settled; the resume can go out now. + return event.payload.session.status !== "running" && + event.payload.session.status !== "starting" + ? worker.enqueue({ kind: "resume", threadId: event.payload.threadId }) + : Effect.void; + } + return Effect.void; + }; + + const start: UsageLimitResumeReactor["Service"]["start"] = Effect.fn( + "UsageLimitResumeReactor.start", + )(function* () { + const events = yield* engine.subscribeDomainEvents; + yield* scan().pipe( + Effect.catchCause((cause) => + Effect.logWarning("usage limit resume scan failed", { cause: Cause.pretty(cause) }), + ), + ); + yield* forkParked( + Effect.gen(function* () { + yield* worker.enqueue({ kind: "sweep" }); + yield* worker.drain; + }).pipe(Effect.repeat(Schedule.spaced("1 minute")), Effect.asVoid), + ); + yield* forkParked(Stream.runForEach(events, processEvent)); + }); + + return { start, drain: worker.drain } satisfies UsageLimitResumeReactor["Service"]; +}); + +export const layer = Layer.effect(UsageLimitResumeReactor, make); diff --git a/apps/server/src/orchestration/decider.ts b/apps/server/src/orchestration/decider.ts index 0119c0e8599a..10bb2f02c557 100644 --- a/apps/server/src/orchestration/decider.ts +++ b/apps/server/src/orchestration/decider.ts @@ -725,6 +725,66 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand" }; } + case "thread.usage-limit.set": { + const thread = yield* requireThread({ + readModel, + command, + threadId: command.threadId, + }); + // A limit re-reported for the same stop (a parked Claude turn re-fires + // while its wait shrinks) keeps the user's resume choice. + const usageLimit = + command.usageLimit === null + ? null + : { + ...command.usageLimit, + resumeScheduled: thread.usageLimit?.resumeScheduled ?? false, + }; + return { + ...(yield* withEventBase({ + aggregateKind: "thread", + aggregateId: command.threadId, + occurredAt: command.createdAt, + commandId: command.commandId, + })), + type: "thread.usage-limit-set", + payload: { + threadId: command.threadId, + usageLimit, + }, + }; + } + + case "thread.auto-resume.set": { + const thread = yield* requireThreadNotArchived({ + readModel, + command, + threadId: command.threadId, + }); + if (command.scheduled && thread.usageLimit?.resetsAt == null) { + return yield* Effect.fail( + new OrchestrationCommandInvariantError({ + commandType: command.type, + detail: `thread ${command.threadId} has no usage limit with a known reset time to resume after`, + }), + ); + } + const occurredAt = yield* nowIso; + return { + ...(yield* withEventBase({ + aggregateKind: "thread", + aggregateId: command.threadId, + occurredAt, + commandId: command.commandId, + })), + type: "thread.auto-resume-set", + payload: { + threadId: command.threadId, + scheduled: command.scheduled, + }, + }; + } + case "thread.pin": { const thread = yield* requireThreadNotArchived({ readModel, @@ -1494,6 +1554,22 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand" }, }); } + // Any new turn, manual or the scheduled resume itself, spends the stop. + if (targetThread.usageLimit != null) { + lifecycleResetEvents.push({ + ...(yield* withEventBase({ + aggregateKind: "thread", + aggregateId: command.threadId, + occurredAt: command.createdAt, + commandId: command.commandId, + })), + type: "thread.usage-limit-set", + payload: { + threadId: command.threadId, + usageLimit: null, + }, + }); + } return [ ...lifecycleResetEvents, ...(userMessageEvent ? [userMessageEvent] : []), diff --git a/apps/server/src/orchestration/projector.test.ts b/apps/server/src/orchestration/projector.test.ts index c4e1996f1ddd..1d4730a8f6be 100644 --- a/apps/server/src/orchestration/projector.test.ts +++ b/apps/server/src/orchestration/projector.test.ts @@ -98,6 +98,7 @@ describe("orchestration projector", () => { unsettledAt: null, snoozedUntil: null, snoozedAt: null, + usageLimit: null, deletedAt: null, messages: [], proposedPlans: [], diff --git a/apps/server/src/orchestration/projector.ts b/apps/server/src/orchestration/projector.ts index 53013770b15b..de4446447d08 100644 --- a/apps/server/src/orchestration/projector.ts +++ b/apps/server/src/orchestration/projector.ts @@ -46,6 +46,8 @@ import { ThreadPullRequestSyncedPayload, ThreadPullRequestUnlinkedPayload, ThreadSnoozedPayload, + ThreadUsageLimitSetPayload, + ThreadAutoResumeSetPayload, ThreadUnpinnedPayload, ThreadUnarchivedPayload, ThreadUnsettledPayload, @@ -442,6 +444,7 @@ export function projectEvent( activeOrderKey: null, snoozedUntil: null, snoozedAt: null, + usageLimit: null, deletedAt: null, messages: [], activities: [], @@ -554,6 +557,31 @@ export function projectEvent( })), ); + case "thread.usage-limit-set": + return decodeForEvent(ThreadUsageLimitSetPayload, event.payload, event.type, "payload").pipe( + Effect.map((payload) => ({ + ...nextBase, + threads: updateThread(nextBase.threads, payload.threadId, { + usageLimit: payload.usageLimit, + }), + })), + ); + + case "thread.auto-resume-set": + return decodeForEvent(ThreadAutoResumeSetPayload, event.payload, event.type, "payload").pipe( + Effect.map((payload) => ({ + ...nextBase, + threads: nextBase.threads.map((thread) => + thread.id === payload.threadId && thread.usageLimit + ? { + ...thread, + usageLimit: { ...thread.usageLimit, resumeScheduled: payload.scheduled }, + } + : thread, + ), + })), + ); + case "thread.pinned": return decodeForEvent(ThreadPinnedPayload, event.payload, event.type, "payload").pipe( Effect.map((payload) => ({ diff --git a/apps/server/src/persistence/Layers/ProjectionThreads.ts b/apps/server/src/persistence/Layers/ProjectionThreads.ts index af36578f286e..ecf0620b3e40 100644 --- a/apps/server/src/persistence/Layers/ProjectionThreads.ts +++ b/apps/server/src/persistence/Layers/ProjectionThreads.ts @@ -12,7 +12,12 @@ import { ProjectionThreadRepository, type ProjectionThreadRepositoryShape, } from "../Services/ProjectionThreads.ts"; -import { ModelSelection, ThreadLinkedPullRequest, ThreadTitleState } from "@t3tools/contracts"; +import { + ModelSelection, + ThreadLinkedPullRequest, + ThreadTitleState, + ThreadUsageLimit, +} from "@t3tools/contracts"; const ProjectionThreadDbRow = ProjectionThread.mapFields( Struct.assign({ @@ -20,6 +25,7 @@ const ProjectionThreadDbRow = ProjectionThread.mapFields( titleState: Schema.NullOr(Schema.fromJsonString(ThreadTitleState)), linkedPullRequest: Schema.NullOr(Schema.fromJsonString(ThreadLinkedPullRequest)), branchPullRequest: Schema.NullOr(Schema.fromJsonString(ThreadLinkedPullRequest)), + usageLimit: Schema.NullOr(Schema.fromJsonString(ThreadUsageLimit)), }), ); @@ -51,6 +57,7 @@ const makeProjectionThreadRepository = Effect.gen(function* () { unsettled_at, snoozed_until, snoozed_at, + usage_limit_json, pinned_at, pin_order_key, active_order_key, @@ -83,6 +90,7 @@ const makeProjectionThreadRepository = Effect.gen(function* () { ${row.unsettledAt}, ${row.snoozedUntil}, ${row.snoozedAt}, + ${row.usageLimit == null ? null : JSON.stringify(row.usageLimit)}, ${row.pinnedAt}, ${row.pinOrderKey ?? null}, ${row.activeOrderKey ?? null}, @@ -115,6 +123,7 @@ const makeProjectionThreadRepository = Effect.gen(function* () { unsettled_at = excluded.unsettled_at, snoozed_until = excluded.snoozed_until, snoozed_at = excluded.snoozed_at, + usage_limit_json = excluded.usage_limit_json, pinned_at = excluded.pinned_at, pin_order_key = excluded.pin_order_key, active_order_key = excluded.active_order_key, @@ -154,6 +163,7 @@ const makeProjectionThreadRepository = Effect.gen(function* () { unsettled_at AS "unsettledAt", snoozed_until AS "snoozedUntil", snoozed_at AS "snoozedAt", + usage_limit_json AS "usageLimit", pinned_at AS "pinnedAt", pin_order_key AS "pinOrderKey", active_order_key AS "activeOrderKey", diff --git a/apps/server/src/persistence/Migrations.ts b/apps/server/src/persistence/Migrations.ts index 7a0fd6b26cb0..a39dce7dd763 100644 --- a/apps/server/src/persistence/Migrations.ts +++ b/apps/server/src/persistence/Migrations.ts @@ -66,6 +66,7 @@ import Migration0051 from "./Migrations/051_ProjectionThreadMessageContext.ts"; import Migration0052 from "./Migrations/052_ProjectionThreadTitleState.ts"; import Migration0053 from "./Migrations/053_PullRequestFilesViewed.ts"; import Migration0054 from "./Migrations/054_AuthSessionLastSeenAt.ts"; +import Migration0055 from "./Migrations/055_ProjectionThreadUsageLimit.ts"; /** * Migration loader with all migrations defined inline. @@ -132,6 +133,7 @@ const migrationEntries = [ [52, "ProjectionThreadTitleState", Migration0052], [53, "PullRequestFilesViewed", Migration0053], [54, "AuthSessionLastSeenAt", Migration0054], + [55, "ProjectionThreadUsageLimit", Migration0055], ] as const; export const migrationManifest = migrationEntries.map(([id, name]) => [id, name] as const); diff --git a/apps/server/src/persistence/Migrations/055_ProjectionThreadUsageLimit.ts b/apps/server/src/persistence/Migrations/055_ProjectionThreadUsageLimit.ts new file mode 100644 index 000000000000..ef39b2201637 --- /dev/null +++ b/apps/server/src/persistence/Migrations/055_ProjectionThreadUsageLimit.ts @@ -0,0 +1,13 @@ +import * as Effect from "effect/Effect"; +import * as SqlClient from "effect/unstable/sql/SqlClient"; + +export default Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + const columns = yield* sql<{ readonly name: string }>` + PRAGMA table_info(projection_threads) + `; + // Guarded so the column survives renumbering if upstream claims this id. + if (!columns.some((column) => column.name === "usage_limit_json")) { + yield* sql`ALTER TABLE projection_threads ADD COLUMN usage_limit_json TEXT`; + } +}); diff --git a/apps/server/src/persistence/Services/ProjectionThreads.ts b/apps/server/src/persistence/Services/ProjectionThreads.ts index 895fe596db24..45d7dcc8d079 100644 --- a/apps/server/src/persistence/Services/ProjectionThreads.ts +++ b/apps/server/src/persistence/Services/ProjectionThreads.ts @@ -16,6 +16,7 @@ import { RuntimeMode, ThreadLinkedPullRequest, ThreadTitleState, + ThreadUsageLimit, ThreadId, TurnId, } from "@t3tools/contracts"; @@ -47,6 +48,7 @@ export const ProjectionThread = Schema.Struct({ unsettledAt: Schema.NullOr(IsoDateTime), snoozedUntil: Schema.NullOr(IsoDateTime), snoozedAt: Schema.NullOr(IsoDateTime), + usageLimit: Schema.optional(Schema.NullOr(ThreadUsageLimit)), pinnedAt: Schema.NullOr(IsoDateTime), pinOrderKey: Schema.optional(Schema.NullOr(Schema.String)), activeOrderKey: Schema.optional(Schema.NullOr(Schema.String)), diff --git a/apps/server/src/provider/Layers/ClaudeAdapter.ts b/apps/server/src/provider/Layers/ClaudeAdapter.ts index 13645a33a05b..8722a2b2b75a 100644 --- a/apps/server/src/provider/Layers/ClaudeAdapter.ts +++ b/apps/server/src/provider/Layers/ClaudeAdapter.ts @@ -74,6 +74,7 @@ import { HostProcessIsExecutable } from "@t3tools/shared/hostProcess"; import * as Cause from "effect/Cause"; import * as Crypto from "effect/Crypto"; import * as DateTime from "effect/DateTime"; +import * as Option from "effect/Option"; import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; @@ -2533,6 +2534,28 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( }); }); + const emitTurnUsageLimited = Effect.fn("emitTurnUsageLimited")(function* ( + context: ClaudeSessionContext, + resetsAtEpochSeconds: number | undefined, + ) { + const turnState = context.turnState; + const stamp = yield* makeEventStamp(); + const resetsAt = + resetsAtEpochSeconds === undefined + ? Option.none() + : Option.map(DateTime.make(resetsAtEpochSeconds * 1000), DateTime.formatIso); + yield* offerRuntimeEvent({ + type: "turn.usage-limited", + eventId: stamp.eventId, + provider: PROVIDER, + createdAt: stamp.createdAt, + threadId: context.session.threadId, + ...(turnState ? { turnId: asCanonicalTurnId(turnState.turnId) } : {}), + payload: Option.isSome(resetsAt) ? { resetsAt: resetsAt.value } : {}, + providerRefs: nativeProviderRefs(context), + }); + }); + const emitThreadTokenUsage = Effect.fn("emitThreadTokenUsage")(function* ( context: ClaudeSessionContext, usage: ThreadTokenUsageSnapshot | undefined, @@ -3496,6 +3519,16 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( : undefined); const { status, errorMessage } = resultOutcome(message, failureHint); + // A rejected rate_limit_event already reported the stop with its reset; + // an assistant-text limit or a blocking_limit result carries none. + if ( + status === "failed" && + turn?.rejectedRateLimitTypes.size === 0 && + (turn.latestAssistantRateLimited || message.terminal_reason === "blocking_limit") + ) { + yield* emitTurnUsageLimited(context, undefined); + } + if (status === "failed") { yield* emitRuntimeError(context, errorMessage ?? "Claude turn failed."); } @@ -4136,6 +4169,7 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( names, ); yield* emitRuntimeWarning(context, notice, rateLimitInfo); + yield* emitTurnUsageLimited(context, rateLimitInfo.resetsAt); } } return; diff --git a/apps/server/src/provider/Layers/CodexAdapter.test.ts b/apps/server/src/provider/Layers/CodexAdapter.test.ts index 9f464bdaa177..3b9465840a17 100644 --- a/apps/server/src/provider/Layers/CodexAdapter.test.ts +++ b/apps/server/src/provider/Layers/CodexAdapter.test.ts @@ -2914,7 +2914,7 @@ usageLimitLayer("CodexAdapterLive usage limits", (it) => { Effect.gen(function* () { const { adapter, runtime } = yield* startUsageLimitRuntime(); const eventsFiber = yield* adapter.streamEvents.pipe( - Stream.take(5), + Stream.take(7), Stream.runCollect, Effect.forkChild, ); @@ -2946,8 +2946,10 @@ usageLimitLayer("CodexAdapterLive usage limits", (it) => { [ "account.rate-limits.updated", "runtime.error", + "turn.usage-limited", "turn.completed", "runtime.error", + "turn.usage-limited", "turn.completed", ], ); @@ -2967,7 +2969,7 @@ usageLimitLayer("CodexAdapterLive usage limits", (it) => { Effect.gen(function* () { const { adapter, runtime } = yield* startUsageLimitRuntime(); const eventsFiber = yield* adapter.streamEvents.pipe( - Stream.take(3), + Stream.take(4), Stream.runCollect, Effect.forkChild, ); @@ -2994,6 +2996,8 @@ usageLimitLayer("CodexAdapterLive usage limits", (it) => { completed?.payload.errorMessage, "Codex usage limit reached. The session limit resets in 3h 20m. Send the message again once the limit resets.", ); + const limited = events.find((event) => event.type === "turn.usage-limited"); + NodeAssert.equal(limited?.payload.resetsAt, "2026-01-01T03:20:00.000Z"); }), ); @@ -3001,7 +3005,7 @@ usageLimitLayer("CodexAdapterLive usage limits", (it) => { Effect.gen(function* () { const { adapter, runtime } = yield* startUsageLimitRuntime(); const eventsFiber = yield* adapter.streamEvents.pipe( - Stream.take(3), + Stream.take(4), Stream.runCollect, Effect.forkChild, ); @@ -3035,7 +3039,7 @@ usageLimitLayer("CodexAdapterLive usage limits", (it) => { Effect.gen(function* () { const { adapter, runtime } = yield* startUsageLimitRuntime(); const eventsFiber = yield* adapter.streamEvents.pipe( - Stream.take(2), + Stream.take(3), Stream.runCollect, Effect.forkChild, ); @@ -3053,8 +3057,10 @@ usageLimitLayer("CodexAdapterLive usage limits", (it) => { const expected = "Codex usage limit reached. Send the message again once the limit resets."; NodeAssert.deepStrictEqual( events.map((event) => event.type), - ["runtime.error", "turn.completed"], + ["runtime.error", "turn.usage-limited", "turn.completed"], ); + const limited = events.find((event) => event.type === "turn.usage-limited"); + NodeAssert.equal(limited?.payload.resetsAt, undefined); const runtimeError = events.find((event) => event.type === "runtime.error"); NodeAssert.equal(runtimeError?.payload.message, expected); const completed = events.find((event) => event.type === "turn.completed"); diff --git a/apps/server/src/provider/Layers/CodexAdapter.ts b/apps/server/src/provider/Layers/CodexAdapter.ts index 0ecc9693ab04..f9b457ed883f 100644 --- a/apps/server/src/provider/Layers/CodexAdapter.ts +++ b/apps/server/src/provider/Layers/CodexAdapter.ts @@ -75,6 +75,7 @@ import { type CodexRateLimitSnapshot, codexRateLimitsToUpdate, codexUsageLimitMessage, + codexUsageLimitResetsAt, mergeCodexRateLimits, } from "./codexUsageLimits.ts"; const isCodexAppServerProcessExitedError = Schema.is(CodexErrors.CodexAppServerProcessExitedError); @@ -2389,7 +2390,7 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( if (errorPayload?.error.codexErrorInfo === "usageLimitExceeded") return; } - let usageLimitError: ProviderRuntimeEvent | undefined; + let usageLimitEvents: ReadonlyArray = []; let usageLimitMessage: string | undefined; if (event.method === "turn/completed") { const completedPayload = readPayload( @@ -2402,15 +2403,23 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( : undefined; if (turnError?.codexErrorInfo === "usageLimitExceeded") { usageLimitMessage = codexUsageLimitMessage(rateLimits, event.createdAt); - usageLimitError = { - ...runtimeEventBase(event, event.threadId), - type: "runtime.error", - payload: { - message: usageLimitMessage, - class: "provider_error", - ...(turnError.message ? { detail: turnError.message } : {}), + const resetsAt = codexUsageLimitResetsAt(rateLimits, event.createdAt); + usageLimitEvents = [ + { + ...runtimeEventBase(event, event.threadId), + type: "runtime.error", + payload: { + message: usageLimitMessage, + class: "provider_error", + ...(turnError.message ? { detail: turnError.message } : {}), + }, + }, + { + ...runtimeEventBase(event, event.threadId), + type: "turn.usage-limited", + payload: resetsAt ? { resetsAt } : {}, }, - }; + ]; } } @@ -2444,9 +2453,7 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( } return runtimeEvent; }); - const runtimeEvents = usageLimitError - ? [usageLimitError, ...mappedEvents] - : mappedEvents; + const runtimeEvents = [...usageLimitEvents, ...mappedEvents]; if (runtimeEvents.length === 0) { yield* Effect.logDebug("ignoring unhandled Codex provider event", { method: event.method, diff --git a/apps/server/src/provider/Layers/GrokAdapter.ts b/apps/server/src/provider/Layers/GrokAdapter.ts index bec0fc39faa0..445c080e28c6 100644 --- a/apps/server/src/provider/Layers/GrokAdapter.ts +++ b/apps/server/src/provider/Layers/GrokAdapter.ts @@ -76,6 +76,7 @@ import { extractGrokPlanMarkdownFromToolCallData, extractXAiAskUserQuestions, extractXAiExitPlanMarkdown, + isXAiUsageLimitError, makeXAiAskUserQuestionCancelledResponse, makeXAiAskUserQuestionResponse, makeXAiExitPlanModeCapturedResponse, @@ -1823,7 +1824,23 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte Ref.set( promptFailureMessageRef, mapAcpToAdapterError(PROVIDER, input.threadId, "session/prompt", error).message, - ).pipe(Effect.andThen(prepared.acp.drainEvents)), + ).pipe( + Effect.andThen( + isXAiUsageLimitError(error) + ? Effect.flatMap(makeEventStamp(), (stamp) => + offerRuntimeEvent({ + type: "turn.usage-limited", + ...stamp, + provider: PROVIDER, + threadId: input.threadId, + turnId: prepared.turnId, + payload: {}, + }), + ).pipe(Effect.ignore) + : Effect.void, + ), + Effect.andThen(prepared.acp.drainEvents), + ), ), Effect.mapError((error) => mapAcpToAdapterError(PROVIDER, input.threadId, "session/prompt", error), diff --git a/apps/server/src/provider/Layers/codexUsageLimits.ts b/apps/server/src/provider/Layers/codexUsageLimits.ts index 74ce04a7895f..7edee5abe788 100644 --- a/apps/server/src/provider/Layers/codexUsageLimits.ts +++ b/apps/server/src/provider/Layers/codexUsageLimits.ts @@ -209,6 +209,31 @@ function codexUsageLimitNextStep(rateLimitReachedType: string | null | undefined } } +/** The exhausted window that has yet to reset, latest first, as of `atIso`. */ +function codexBlockingWindow(snapshot: CodexRateLimitSnapshot | undefined, atIso: string) { + const atMs = Date.parse(atIso); + const windows = snapshot && Number.isFinite(atMs) ? codexRateLimitsToWindows(snapshot) : []; + let blocking: + | { readonly kind: string; readonly resetsAt: string; readonly waitMs: number } + | undefined; + for (const window of windows) { + if (window.usedPercent < 100 || !window.resetsAt) continue; + const resetMs = Date.parse(window.resetsAt); + if (!Number.isFinite(resetMs) || resetMs <= atMs) continue; + if (blocking && resetMs - atMs <= blocking.waitMs) continue; + blocking = { kind: window.kind, resetsAt: window.resetsAt, waitMs: resetMs - atMs }; + } + return blocking; +} + +/** When the limit that stopped a turn at `atIso` resets, if the snapshot says. */ +export function codexUsageLimitResetsAt( + snapshot: CodexRateLimitSnapshot | undefined, + atIso: string, +): string | undefined { + return codexBlockingWindow(snapshot, atIso)?.resetsAt; +} + /** * The message a usage-limit stop shows instead of the provider sentence, which * on a Business workspace blames credits for a window that simply ran out. The @@ -219,16 +244,9 @@ export function codexUsageLimitMessage( snapshot: CodexRateLimitSnapshot | undefined, atIso: string, ): string { - const atMs = Date.parse(atIso); - const windows = snapshot && Number.isFinite(atMs) ? codexRateLimitsToWindows(snapshot) : []; - let reset = ""; - let latestResetMs = Number.NEGATIVE_INFINITY; - for (const window of windows) { - if (window.usedPercent < 100 || !window.resetsAt) continue; - const resetMs = Date.parse(window.resetsAt); - if (!Number.isFinite(resetMs) || resetMs <= atMs || resetMs <= latestResetMs) continue; - latestResetMs = resetMs; - reset = ` The ${window.kind} limit resets in ${formatCodexUsageLimitWait(resetMs - atMs)}.`; - } + const blocking = codexBlockingWindow(snapshot, atIso); + const reset = blocking + ? ` The ${blocking.kind} limit resets in ${formatCodexUsageLimitWait(blocking.waitMs)}.` + : ""; return `Codex usage limit reached.${reset}${codexUsageLimitNextStep(snapshot?.rateLimitReachedType)}`; } diff --git a/apps/server/src/provider/acp/XAiAcpExtension.ts b/apps/server/src/provider/acp/XAiAcpExtension.ts index 08bbd674e9b0..ec08c5311b23 100644 --- a/apps/server/src/provider/acp/XAiAcpExtension.ts +++ b/apps/server/src/provider/acp/XAiAcpExtension.ts @@ -32,6 +32,13 @@ const completedXAiPromptIdLimit = 128; const xAiStopReasonMissingMetaKey = "xAiStopReasonMissing"; const xAiRateLimitedErrorCode = -32003; +const isAcpRequestError = Schema.is(EffectAcpErrors.AcpRequestError); + +/** Whether a failed Grok prompt stopped on the account's usage limit. */ +export function isXAiUsageLimitError(error: unknown): boolean { + return isAcpRequestError(error) && error.code === xAiRateLimitedErrorCode; +} + const XAiAskUserQuestionOption = Schema.Struct({ label: Schema.String, description: Schema.optional(Schema.String), diff --git a/apps/server/src/provider/providerUsageLimits.ts b/apps/server/src/provider/providerUsageLimits.ts index ea8d0d1d029f..80ec9935328e 100644 --- a/apps/server/src/provider/providerUsageLimits.ts +++ b/apps/server/src/provider/providerUsageLimits.ts @@ -132,3 +132,23 @@ export function resolveUsageLimitsAfterProbe(input: { } return probed; } + +/** + * When a usage-limit stop says nothing about its reset, the exhausted window + * that resets last is what still blocks the account. Undefined when no + * exhausted window reports a future reset. + */ +export function latestExhaustedWindowResetAt( + limits: ServerProviderUsageLimits | undefined, + nowIso: string, +): string | undefined { + const nowMs = Date.parse(nowIso); + let latest: { readonly iso: string; readonly ms: number } | undefined; + for (const window of limits?.windows ?? []) { + if (window.usedPercent < 100 || !window.resetsAt) continue; + const ms = Date.parse(window.resetsAt); + if (!(ms > nowMs) || (latest && latest.ms >= ms)) continue; + latest = { iso: window.resetsAt, ms }; + } + return latest?.iso; +} diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index db7246810b2d..cd5c73f2a503 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -85,6 +85,7 @@ import { ProviderCommandReactorLive } from "./orchestration/Layers/ProviderComma import { CheckpointReactorLive } from "./orchestration/Layers/CheckpointReactor.ts"; import { ThreadDeletionReactorLive } from "./orchestration/Layers/ThreadDeletionReactor.ts"; import * as ThreadSettlementReactor from "./orchestration/ThreadSettlementReactor.ts"; +import * as UsageLimitResumeReactor from "./orchestration/UsageLimitResumeReactor.ts"; import * as StorageCleanup from "./storageCleanup.ts"; import * as PullRequestSyncReactor from "./orchestration/PullRequestSyncReactor.ts"; import * as ThreadPullRequestReactor from "./orchestration/ThreadPullRequestReactor.ts"; @@ -253,6 +254,7 @@ const ReactorLayerLive = Layer.empty.pipe( Layer.provideMerge(StorageCleanup.layer), Layer.provideMerge(ThreadDeletionReactorLive), Layer.provideMerge(ThreadSettlementReactor.layer), + Layer.provideMerge(UsageLimitResumeReactor.layer), Layer.provideMerge(PullRequestSyncReactor.layer), Layer.provideMerge(ThreadPullRequestReactor.layer), Layer.provideMerge(AgentAwarenessRelay.layer.pipe(Layer.provide(ServerSecretStore.layer))), diff --git a/apps/web/src/components/ChatView.tsx b/apps/web/src/components/ChatView.tsx index 6c0d4d6d7fdb..27e07e2739e6 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -398,6 +398,7 @@ import { shouldShowThreadErrorBanner, ThreadErrorBanner, } from "./chat/ThreadErrorBanner"; +import { UsageLimitResumeBanner } from "./chat/UsageLimitResumeBanner"; import type { ComposerBannerStackItem } from "./chat/ComposerBannerStack"; import { ComposerSurface } from "./chat/ComposerSurface"; import { @@ -9930,6 +9931,13 @@ export default function ChatView(props: ChatViewProps) { setThreadErrorBannerDismissTick((tick) => tick + 1); }} /> + {activeServerThread !== null && ( + + )} {/* Messages Wrapper */}
diff --git a/apps/web/src/components/chat/UsageLimitResumeBanner.tsx b/apps/web/src/components/chat/UsageLimitResumeBanner.tsx new file mode 100644 index 000000000000..376dceae68cf --- /dev/null +++ b/apps/web/src/components/chat/UsageLimitResumeBanner.tsx @@ -0,0 +1,56 @@ +import type { ClientSettings, EnvironmentId, ThreadId, ThreadUsageLimit } from "@t3tools/contracts"; +import { AlarmClockIcon } from "lucide-react"; +import { memo } from "react"; + +import { useClientSettings } from "../../hooks/useSettings"; +import { threadEnvironment } from "../../state/threads"; +import { useAtomCommand } from "../../state/use-atom-command"; +import { formatUpcomingTimestamp } from "../../timestampFormat"; +import { Alert, AlertAction, AlertDescription } from "../ui/alert"; +import { Button } from "../ui/button"; + +const selectTimestampFormat = (settings: ClientSettings) => settings.timestampFormat; + +/** Offers, or shows, the server-side resume of a thread stopped by a usage limit. */ +export const UsageLimitResumeBanner = memo(function UsageLimitResumeBanner({ + environmentId, + threadId, + usageLimit, +}: { + environmentId: EnvironmentId; + threadId: ThreadId; + usageLimit: ThreadUsageLimit | null | undefined; +}) { + const timestampFormat = useClientSettings(selectTimestampFormat); + const setAutoResume = useAtomCommand(threadEnvironment.setAutoResume); + if (!usageLimit?.resetsAt) return null; + const upcoming = formatUpcomingTimestamp(usageLimit.resetsAt, timestampFormat); + const resetTime = upcoming.startsWith("tomorrow") ? upcoming : `at ${upcoming}`; + const scheduled = usageLimit.resumeScheduled; + return ( +
+ + + + {scheduled + ? `Resumes automatically ${resetTime}.` + : `The usage limit resets ${resetTime}.`} + + + + + +
+ ); +}); diff --git a/apps/web/src/components/settings/SettingsPanels.tsx b/apps/web/src/components/settings/SettingsPanels.tsx index f37c9fb631c0..c5fc86d12ffd 100644 --- a/apps/web/src/components/settings/SettingsPanels.tsx +++ b/apps/web/src/components/settings/SettingsPanels.tsx @@ -605,6 +605,12 @@ export function useSettingsRestore(onRestored?: () => void) { DEFAULT_UNIFIED_SETTINGS.continueThreadsAfterServerUpdate ? ["Continue threads after restarts"] : []), + ...(settings.autoResumeAfterUsageLimit !== + DEFAULT_UNIFIED_SETTINGS.autoResumeAfterUsageLimit || + settings.autoResumeMessage !== DEFAULT_UNIFIED_SETTINGS.autoResumeMessage || + settings.autoResumeDisablesFastMode !== DEFAULT_UNIFIED_SETTINGS.autoResumeDisablesFastMode + ? ["Resume after usage limits"] + : []), ...(isBackgroundActivityDirty ? ["Background activity"] : []), ...(settings.defaultThreadEnvMode !== DEFAULT_UNIFIED_SETTINGS.defaultThreadEnvMode ? ["New thread mode"] @@ -676,6 +682,9 @@ export function useSettingsRestore(onRestored?: () => void) { settings.responseStreamingMode, settings.enableProviderUpdateChecks, settings.continueThreadsAfterServerUpdate, + settings.autoResumeAfterUsageLimit, + settings.autoResumeMessage, + settings.autoResumeDisablesFastMode, settings.sidebarAutoSettleAfterDays, settings.sidebarAutoSettleOnMerge, settings.sidebarProjectGroupingMode, @@ -782,6 +791,9 @@ export function useSettingsRestore(onRestored?: () => void) { responseStreamingMode: DEFAULT_UNIFIED_SETTINGS.responseStreamingMode, enableProviderUpdateChecks: DEFAULT_UNIFIED_SETTINGS.enableProviderUpdateChecks, continueThreadsAfterServerUpdate: DEFAULT_UNIFIED_SETTINGS.continueThreadsAfterServerUpdate, + autoResumeAfterUsageLimit: DEFAULT_UNIFIED_SETTINGS.autoResumeAfterUsageLimit, + autoResumeMessage: DEFAULT_UNIFIED_SETTINGS.autoResumeMessage, + autoResumeDisablesFastMode: DEFAULT_UNIFIED_SETTINGS.autoResumeDisablesFastMode, backgroundActivity: DEFAULT_UNIFIED_SETTINGS.backgroundActivity, backgroundActivityProfile: DEFAULT_UNIFIED_SETTINGS.backgroundActivityProfile, automaticGitFetchInterval: DEFAULT_UNIFIED_SETTINGS.automaticGitFetchInterval, @@ -2804,6 +2816,93 @@ export function GeneralSettingsPanel() { } /> + + updateSettings({ + autoResumeAfterUsageLimit: DEFAULT_UNIFIED_SETTINGS.autoResumeAfterUsageLimit, + }) + } + /> + ) : null + } + control={ + + updateSettings({ autoResumeAfterUsageLimit: Boolean(checked) }) + } + aria-label="Resume after usage limits" + /> + } + /> + + + updateSettings({ autoResumeMessage: DEFAULT_UNIFIED_SETTINGS.autoResumeMessage }) + } + /> + ) : null + } + control={ + updateSettings({ autoResumeMessage: next })} + placeholder={DEFAULT_UNIFIED_SETTINGS.autoResumeMessage} + aria-label="Resume message" + /> + } + /> + + + updateSettings({ + autoResumeDisablesFastMode: DEFAULT_UNIFIED_SETTINGS.autoResumeDisablesFastMode, + }) + } + /> + ) : null + } + control={ + + updateSettings({ autoResumeDisablesFastMode: Boolean(checked) }) + } + aria-label="Turn off fast mode when resuming" + /> + } + /> + ; export type UnsettleThreadInput = CommandInput<"thread.unsettle">; export type SnoozeThreadInput = CommandInput<"thread.snooze">; export type UnsnoozeThreadInput = CommandInput<"thread.unsnooze">; +export type SetThreadAutoResumeInput = CommandInput<"thread.auto-resume.set">; export type PinThreadInput = CommandInput<"thread.pin">; export type UnpinThreadInput = CommandInput<"thread.unpin">; export type ReorderPinnedThreadInput = CommandInput<"thread.pin.reorder">; @@ -206,6 +207,16 @@ export const unsnoozeThread: (input: UnsnoozeThreadInput) => CommandEffect = Eff }); }); +export const setThreadAutoResume: (input: SetThreadAutoResumeInput) => CommandEffect = Effect.fn( + "EnvironmentCommands.setThreadAutoResume", +)(function* (input) { + return yield* dispatch({ + ...input, + type: "thread.auto-resume.set", + commandId: yield* commandId(input), + }); +}); + export const pinThread: (input: PinThreadInput) => CommandEffect = Effect.fn( "EnvironmentCommands.pinThread", )(function* (input) { diff --git a/packages/client-runtime/src/state/threadCommands.ts b/packages/client-runtime/src/state/threadCommands.ts index 93e22cfd0c70..5fbb786a03d7 100644 --- a/packages/client-runtime/src/state/threadCommands.ts +++ b/packages/client-runtime/src/state/threadCommands.ts @@ -24,6 +24,7 @@ import { type RespondToThreadUserInputInput, type DismissThreadUserInputInput, type RevertThreadCheckpointInput, + type SetThreadAutoResumeInput, type SetThreadInteractionModeInput, type SetThreadRuntimeModeInput, type PinThreadInput, @@ -48,6 +49,7 @@ import { respondToThreadUserInput, dismissThreadUserInput, revertThreadCheckpoint, + setThreadAutoResume, setThreadInteractionMode, setThreadRuntimeMode, pinThread, @@ -76,6 +78,7 @@ export type { RespondToThreadUserInputInput, DismissThreadUserInputInput, RevertThreadCheckpointInput, + SetThreadAutoResumeInput, SetThreadInteractionModeInput, SetThreadRuntimeModeInput, PinThreadInput, @@ -152,6 +155,12 @@ export function createThreadEnvironmentAtoms( scheduler, concurrency, }), + setAutoResume: createEnvironmentCommand(runtime, { + label: "environment-data:commands:thread:set-auto-resume", + execute: (input: SetThreadAutoResumeInput) => setThreadAutoResume(input), + scheduler, + concurrency, + }), pin: createEnvironmentCommand(runtime, { label: "environment-data:commands:thread:pin", execute: (input: PinThreadInput) => pinThread(input), diff --git a/packages/client-runtime/src/state/threadDetail.ts b/packages/client-runtime/src/state/threadDetail.ts index 379985b71243..35b8f109197b 100644 --- a/packages/client-runtime/src/state/threadDetail.ts +++ b/packages/client-runtime/src/state/threadDetail.ts @@ -62,6 +62,7 @@ export function mergeEnvironmentThread( activeOrderKey: shell.activeOrderKey, snoozedUntil: shell.snoozedUntil, snoozedAt: shell.snoozedAt, + usageLimit: shell.usageLimit, pinnedAt: shell.pinnedAt, pinOrderKey: shell.pinOrderKey, session: shell.session, diff --git a/packages/client-runtime/src/state/threadReducer.ts b/packages/client-runtime/src/state/threadReducer.ts index 101bb34fba91..3d4ad2c7245d 100644 --- a/packages/client-runtime/src/state/threadReducer.ts +++ b/packages/client-runtime/src/state/threadReducer.ts @@ -133,6 +133,7 @@ export function applyThreadDetailEvent( activeOrderKey: null, snoozedUntil: null, snoozedAt: null, + usageLimit: null, deletedAt: null, pullRequests: [], messages: [], @@ -215,6 +216,23 @@ export function applyThreadDetailEvent( }, }; + case "thread.usage-limit-set": + return { + kind: "updated", + thread: { ...thread, usageLimit: event.payload.usageLimit }, + }; + + case "thread.auto-resume-set": + return thread.usageLimit + ? { + kind: "updated", + thread: { + ...thread, + usageLimit: { ...thread.usageLimit, resumeScheduled: event.payload.scheduled }, + }, + } + : { kind: "unchanged" }; + case "thread.pinned": return { kind: "updated", diff --git a/packages/contracts/src/orchestration.ts b/packages/contracts/src/orchestration.ts index 3e323e4964d5..de7322ed677a 100644 --- a/packages/contracts/src/orchestration.ts +++ b/packages/contracts/src/orchestration.ts @@ -790,6 +790,19 @@ export const ThreadPullRequestLink = Schema.Struct({ }); export type ThreadPullRequestLink = typeof ThreadPullRequestLink.Type; +/** + * A provider usage limit stopped this thread's latest turn. `resetsAt` is null + * when the provider reported no reset time; such a stop cannot be resumed on a + * timer. `resumeScheduled` asks the server to send the configured resume + * message once `resetsAt` passes. Cleared by the next turn start. + */ +export const ThreadUsageLimit = Schema.Struct({ + reachedAt: IsoDateTime, + resetsAt: Schema.NullOr(IsoDateTime), + resumeScheduled: Schema.Boolean, +}); +export type ThreadUsageLimit = typeof ThreadUsageLimit.Type; + export const OrchestrationThread = Schema.Struct({ id: ThreadId, projectId: ProjectId, @@ -826,6 +839,8 @@ export const OrchestrationThread = Schema.Struct({ // Optional so payloads from pre-snooze servers still decode. snoozedUntil: Schema.optional(Schema.NullOr(IsoDateTime)), snoozedAt: Schema.optional(Schema.NullOr(IsoDateTime)), + // Optional so payloads from pre-auto-resume servers still decode. + usageLimit: Schema.optional(Schema.NullOr(ThreadUsageLimit)), // Active pinned threads render in the pinned block. Settled and snoozed // threads remain in their respective shelves even when pinned. // Optional so payloads from pre-pinning servers still decode. @@ -905,6 +920,7 @@ export const OrchestrationThreadShell = Schema.Struct({ unsettledAt: Schema.optional(Schema.NullOr(IsoDateTime)), snoozedUntil: Schema.optional(Schema.NullOr(IsoDateTime)), snoozedAt: Schema.optional(Schema.NullOr(IsoDateTime)), + usageLimit: Schema.optional(Schema.NullOr(ThreadUsageLimit)), pinnedAt: Schema.optional(Schema.NullOr(IsoDateTime)), pinOrderKey: Schema.optional(Schema.NullOr(TrimmedNonEmptyString)), activeOrderKey: Schema.optional(Schema.NullOr(TrimmedNonEmptyString)), @@ -1191,6 +1207,13 @@ const ThreadUnsnoozeCommand = Schema.Struct({ reason: Schema.Literal("user"), }); +const ThreadAutoResumeSetCommand = Schema.Struct({ + type: Schema.Literal("thread.auto-resume.set"), + commandId: CommandId, + threadId: ThreadId, + scheduled: Schema.Boolean, +}); + const ThreadPinCommand = Schema.Struct({ type: Schema.Literal("thread.pin"), commandId: CommandId, @@ -1423,6 +1446,7 @@ const DispatchableClientOrchestrationCommand = Schema.Union([ ThreadUnsettleCommand, ThreadSnoozeCommand, ThreadUnsnoozeCommand, + ThreadAutoResumeSetCommand, ThreadPinCommand, ThreadUnpinCommand, ThreadPinReorderCommand, @@ -1456,6 +1480,7 @@ export const ClientOrchestrationCommand = Schema.Union([ ThreadUnsettleCommand, ThreadSnoozeCommand, ThreadUnsnoozeCommand, + ThreadAutoResumeSetCommand, ThreadPinCommand, ThreadUnpinCommand, ThreadPinReorderCommand, @@ -1643,7 +1668,19 @@ const ThreadPullRequestLinkSyncCommand = Schema.Struct({ stack: Schema.NullOr(ThreadPullRequestStack), }); +const ThreadUsageLimitSetCommand = Schema.Struct({ + type: Schema.Literal("thread.usage-limit.set"), + commandId: CommandId, + threadId: ThreadId, + // Null records that the limit no longer blocks the thread. + usageLimit: Schema.NullOr( + Schema.Struct({ reachedAt: IsoDateTime, resetsAt: Schema.NullOr(IsoDateTime) }), + ), + createdAt: IsoDateTime, +}); + const InternalOrchestrationCommand = Schema.Union([ + ThreadUsageLimitSetCommand, ThreadAutoSettleCommand, ThreadPullRequestSyncCommand, ThreadPullRequestLinkSyncCommand, @@ -1684,6 +1721,8 @@ export const OrchestrationEventType = Schema.Literals([ "thread.unsettled", "thread.snoozed", "thread.unsnoozed", + "thread.usage-limit-set", + "thread.auto-resume-set", "thread.pinned", "thread.unpinned", "thread.pin-reordered", @@ -1805,6 +1844,16 @@ export const ThreadUnsnoozedPayload = Schema.Struct({ updatedAt: IsoDateTime, }); +export const ThreadUsageLimitSetPayload = Schema.Struct({ + threadId: ThreadId, + usageLimit: Schema.NullOr(ThreadUsageLimit), +}); + +export const ThreadAutoResumeSetPayload = Schema.Struct({ + threadId: ThreadId, + scheduled: Schema.Boolean, +}); + export const ThreadPinnedPayload = Schema.Struct({ threadId: ThreadId, pinnedAt: IsoDateTime, @@ -2074,6 +2123,16 @@ export const OrchestrationEvent = Schema.Union([ type: Schema.Literal("thread.unsnoozed"), payload: ThreadUnsnoozedPayload, }), + Schema.Struct({ + ...EventBaseFields, + type: Schema.Literal("thread.usage-limit-set"), + payload: ThreadUsageLimitSetPayload, + }), + Schema.Struct({ + ...EventBaseFields, + type: Schema.Literal("thread.auto-resume-set"), + payload: ThreadAutoResumeSetPayload, + }), Schema.Struct({ ...EventBaseFields, type: Schema.Literal("thread.pinned"), diff --git a/packages/contracts/src/providerRuntime.ts b/packages/contracts/src/providerRuntime.ts index 309b61935485..e86db5142ad3 100644 --- a/packages/contracts/src/providerRuntime.ts +++ b/packages/contracts/src/providerRuntime.ts @@ -196,6 +196,7 @@ const ConfigWarningType = Schema.Literal("config.warning"); const DeprecationNoticeType = Schema.Literal("deprecation.notice"); const FilesPersistedType = Schema.Literal("files.persisted"); const ToolDeniedType = Schema.Literal("tool.denied"); +const TurnUsageLimitedType = Schema.Literal("turn.usage-limited"); const RuntimeWarningType = Schema.Literal("runtime.warning"); const RuntimeErrorType = Schema.Literal("runtime.error"); @@ -801,6 +802,16 @@ const ToolDeniedPayload = Schema.Struct({ }); export type ToolDeniedPayload = typeof ToolDeniedPayload.Type; +/** + * A provider usage limit stopped or parked the turn. Emitted alongside the + * adapter's own warning or failure text; `resetsAt` is omitted when the + * provider did not say when the limit resets. + */ +const TurnUsageLimitedPayload = Schema.Struct({ + resetsAt: Schema.optional(IsoDateTime), +}); +export type TurnUsageLimitedPayload = typeof TurnUsageLimitedPayload.Type; + const RuntimeWarningPayload = Schema.Struct({ message: TrimmedNonEmptyStringSchema, detail: Schema.optional(Schema.Unknown), @@ -1160,6 +1171,13 @@ const ProviderRuntimeToolDeniedEvent = Schema.Struct({ }); export type ProviderRuntimeToolDeniedEvent = typeof ProviderRuntimeToolDeniedEvent.Type; +const ProviderRuntimeTurnUsageLimitedEvent = Schema.Struct({ + ...ProviderRuntimeEventBase.fields, + type: TurnUsageLimitedType, + payload: TurnUsageLimitedPayload, +}); +export type ProviderRuntimeTurnUsageLimitedEvent = typeof ProviderRuntimeTurnUsageLimitedEvent.Type; + const ProviderRuntimeWarningEvent = Schema.Struct({ ...ProviderRuntimeEventBase.fields, type: RuntimeWarningType, @@ -1222,6 +1240,7 @@ export const ProviderRuntimeEventV2 = Schema.Union([ ProviderRuntimeDeprecationNoticeEvent, ProviderRuntimeFilesPersistedEvent, ProviderRuntimeToolDeniedEvent, + ProviderRuntimeTurnUsageLimitedEvent, ProviderRuntimeWarningEvent, ProviderRuntimeErrorEvent, ]); diff --git a/packages/contracts/src/settings.ts b/packages/contracts/src/settings.ts index 65aa113228fe..b2099e22c7ea 100644 --- a/packages/contracts/src/settings.ts +++ b/packages/contracts/src/settings.ts @@ -1107,6 +1107,14 @@ export const ServerSettings = Schema.Struct({ continueThreadsAfterServerUpdate: Schema.Boolean.pipe( Schema.withDecodingDefault(Effect.succeed(false)), ), + /** + * Resume a thread stopped by a provider usage limit once the limit resets, + * by sending `autoResumeMessage`. Threads can also opt in one stop at a time. + */ + autoResumeAfterUsageLimit: Schema.Boolean.pipe(Schema.withDecodingDefault(Effect.succeed(false))), + autoResumeMessage: Schema.String.pipe(Schema.withDecodingDefault(Effect.succeed("go on"))), + /** Send the resume without fast mode; nobody is waiting on it. */ + autoResumeDisablesFastMode: Schema.Boolean.pipe(Schema.withDecodingDefault(Effect.succeed(true))), /** * Whether agents may drive the in-app preview browser. Turning this off * withholds the MCP credential, so the `t3-code` server (and with it every @@ -1479,6 +1487,9 @@ export const ServerSettingsPatch = Schema.Struct({ responseStreamingMode: Schema.optionalKey(ResponseStreamingMode), enableProviderUpdateChecks: Schema.optionalKey(Schema.Boolean), continueThreadsAfterServerUpdate: Schema.optionalKey(Schema.Boolean), + autoResumeAfterUsageLimit: Schema.optionalKey(Schema.Boolean), + autoResumeMessage: Schema.optionalKey(Schema.String), + autoResumeDisablesFastMode: Schema.optionalKey(Schema.Boolean), enableAgentBrowserAccess: Schema.optionalKey(Schema.Boolean), projectAgentBrowserAccessOverrides: Schema.optionalKey( Schema.Record(ProjectId, Schema.NullOr(Schema.Boolean)),