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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 14 additions & 1 deletion apps/mobile/src/features/settings/SettingsThreadsRouteScreen.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -86,7 +86,9 @@ function AutoSettleSettingsRows() {
return null;
}

const writeToAll = (patch: Partial<AutoSettleSettings>) => {
const writeToAll = (
patch: Partial<AutoSettleSettings> & { autoResumeLimitedThreads?: boolean },
) => {
if (writeInFlight.current) return;
const writes = planMobileScopedSettingsPatch(syncTargets, projectSelected, patch);
if (writes.length === 0) return;
Expand Down Expand Up @@ -164,6 +166,17 @@ function AutoSettleSettingsRows() {
onClear={clearProjectOverrides}
/>
) : null}
{!projectSelected ? (
<SettingsSection title="Usage limits">
<SettingsSwitchRow
icon="clock"
label="Auto-resume limited threads"
value={referenceSettings.autoResumeLimitedThreads}
disabled={disabled}
onValueChange={(value) => writeToAll({ autoResumeLimitedThreads: value })}
/>
</SettingsSection>
) : null}
<SettingsSection title="Auto-settle">
<SettingsSwitchRow
icon="arrow.triangle.branch"
Expand Down
6 changes: 6 additions & 0 deletions apps/mobile/src/features/threads/ThreadDetailScreen.tsx
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import { UsageLimitRecoveryCard } from "./UsageLimitRecoveryCard";
import { useNavigation } from "@react-navigation/native";
import type { WorktreeSetupCardProps } from "./worktree-setup-card";
import type { ComposerTextPaste } from "../../native/T3ComposerEditor.types";
Expand Down Expand Up @@ -1093,6 +1094,11 @@ export const ThreadDetailScreen = memo(function ThreadDetailScreen(props: Thread
/>
</Animated.View>
) : null}
<UsageLimitRecoveryCard
key={props.selectedThread.latestRun?.runId}
thread={props.selectedThread}
environmentId={props.environmentId}
/>
{props.feedbackSubmissions.map((submission) => (
<ComposerFeedback
key={submission.id}
Expand Down
77 changes: 77 additions & 0 deletions apps/mobile/src/features/threads/UsageLimitRecoveryCard.tsx
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
import { squashAtomCommandFailure } from "@t3tools/client-runtime/state/runtime";
import type { EnvironmentThreadShell } from "@t3tools/client-runtime/state/shell";
import type { EnvironmentId } from "@t3tools/contracts";
import * as DateTime from "effect/DateTime";
import { useState } from "react";
import { Pressable, View } from "react-native";
import { AppText as Text } from "../../components/AppText";
import { threadEnvironment } from "../../state/threads";
import { useAtomCommand } from "../../state/use-atom-command";

export function UsageLimitRecoveryCard({
thread,
environmentId,
}: {
thread: EnvironmentThreadShell;
environmentId: EnvironmentId;
}) {
const updateMetadata = useAtomCommand(threadEnvironment.updateMetadata);
const [pending, setPending] = useState(false);
const [error, setError] = useState<string | null>(null);
const resetAt = thread.runtime?.usageLimitResetAt ?? null;
const canSchedule =
resetAt !== null &&
Date.parse(resetAt) > Date.parse(thread.latestRun?.completedAt ?? thread.updatedAt);
const runId = thread.latestRun?.runId;
const recovery = thread.limitRecovery;
const scheduled =
recovery?.runId === runId && recovery?.resetAt === resetAt && recovery?.autoResume;
if (
thread.runtime?.status !== "failed" ||
thread.runtime.lastErrorClass !== "usage_limit" ||
!runId
)
return null;
async function toggle() {
if (!resetAt || !runId || !canSchedule) return;
setPending(true);
setError(null);
try {
const result = await updateMetadata({
environmentId,
input: { threadId: thread.id, limitRecovery: { runId, resetAt, autoResume: !scheduled } },
});
if (result._tag === "Failure") throw squashAtomCommandFailure(result);
} catch (cause) {
setError(cause instanceof Error ? cause.message : "Could not change limit recovery.");
} finally {
setPending(false);
}
}
return (
<View className="mx-3 mb-2 gap-2 rounded-xl border border-warning-foreground/25 bg-background p-3">
<Text className="text-sm text-warning-foreground">
{resetAt
? `Usage limit resets ${DateTime.toDateUtc(DateTime.makeUnsafe(resetAt)).toLocaleString()}.`
: "The provider did not report a reset time. Retry manually when your limit is available."}
</Text>
{canSchedule ? (
<Pressable
accessibilityRole="button"
disabled={pending}
onPress={() => void toggle()}
className="self-start rounded-lg bg-subtle px-3 py-2 active:opacity-70"
>
<Text className="text-sm text-foreground">
{scheduled ? "Cancel auto-resume" : "Resume at reset"}
</Text>
</Pressable>
) : null}
{error ? (
<Text accessibilityRole="alert" className="text-sm text-destructive">
{error}
</Text>
) : null}
</View>
);
}
76 changes: 76 additions & 0 deletions apps/server/src/orchestration-v2/Orchestrator.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import { latestRootProviderFailure } from "@t3tools/shared/orchestrationV2ThreadError";
import { threadPullRequestsOf } from "@t3tools/shared/threadPullRequests";
import {
normalizeThreadPullRequestKey,
Expand Down Expand Up @@ -70,6 +71,7 @@ import {
import {
applyToProjection,
emptyProjection,
threadShellFromProjection,
isTurnItemAtOrBeforeRun,
ProjectionStoreV2,
type ProjectionCheckpointContext,
Expand Down Expand Up @@ -2225,6 +2227,28 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
: null;

const now = yield* DateTime.now;
if (command.type === "thread.metadata.update" && command.limitRecovery != null) {
const projection = yield* loadProjectionForCommand(command);
const run = projection.runs.at(-1) ?? null;
const failure = latestRootProviderFailure(run, projection.turnItems);
if (
thread.archivedAt !== null ||
thread.settledOverride === "settled" ||
run?.id !== command.limitRecovery.runId ||
failure?.class !== "usage_limit" ||
failure.resetAt !== command.limitRecovery.resetAt ||
Date.parse(command.limitRecovery.resetAt) <=
DateTime.toEpochMillis(run.completedAt ?? run.requestedAt) ||
projection.runtimeRequests.some((request) => request.status === "pending") ||
projection.runs.some((candidate) => candidate.status === "queued")
) {
return yield* new OrchestratorDispatchError({
commandId: command.commandId,
commandType: command.type,
cause: "The provider limit changed before recovery could be configured.",
});
}
}
let snoozedUntil: DateTime.Utc | null = null;
if (command.type === "thread.snooze") {
const projection = yield* loadProjectionForCommand(command);
Expand Down Expand Up @@ -2382,6 +2406,14 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
return {
...thread,
...(command.title === undefined ? {} : { title: command.title }),
...(command.limitRecovery === undefined
? {}
: {
limitRecovery:
command.limitRecovery === null
? null
: { ...command.limitRecovery, requestId: command.commandId },
}),
...(command.branch === undefined ? {} : { branch: command.branch }),
...(command.worktreePath === undefined ? {} : { worktreePath: command.worktreePath }),
...(command.linkedPullRequest === undefined
Expand Down Expand Up @@ -3754,6 +3786,50 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
) =>
Effect.gen(function* () {
let projection = yield* getProjectionWithPendingEvents(command.threadId, events);
if (command.usageLimitContinuationOfRunId !== undefined) {
const run = projection.runs.at(-1) ?? null;
const failure = latestRootProviderFailure(run, projection.turnItems);
const recovery = projection.thread.limitRecovery;
const now = yield* DateTime.now;
if (
run?.id !== command.usageLimitContinuationOfRunId ||
failure?.class !== "usage_limit" ||
threadShellFromProjection(projection).lastErrorClass !== "usage_limit" ||
!recovery?.autoResume ||
recovery.requestId !== command.usageLimitRecoveryRequestId ||
recovery.runId !== run.id ||
recovery.resetAt !== failure.resetAt ||
Date.parse(recovery.resetAt) > DateTime.toEpochMillis(now) ||
projection.thread.archivedAt !== null ||
projection.thread.deletedAt !== null ||
projection.thread.settledOverride === "settled" ||
projection.thread.providerInstanceId !== run.providerInstanceId ||
projection.runtimeRequests.some((request) => request.status === "pending") ||
(projection.thread.snoozedUntil != null &&
DateTime.toEpochMillis(projection.thread.snoozedUntil) > DateTime.toEpochMillis(now))
) {
yield* emit(
events,
command,
)({
type: "thread.metadata-updated",
threadId: command.threadId,
occurredAt: now,
payload: {
...projection.thread,
...(recovery?.autoResume &&
recovery.requestId === command.usageLimitRecoveryRequestId &&
recovery.runId === command.usageLimitContinuationOfRunId &&
projection.thread.snoozedUntil != null &&
DateTime.toEpochMillis(projection.thread.snoozedUntil) > DateTime.toEpochMillis(now)
? { limitRecovery: { ...recovery, requestId: command.commandId } }
: {}),
},
});
return;
}
}

if (command.restartContinuationOfRunId !== undefined) {
const source = projection.runs.find((run) => run.id === command.restartContinuationOfRunId);
if (
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/orchestration-v2/ProjectionStore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1274,6 +1274,7 @@ export function threadShellFromProjection(
pinOrderKey: projection.thread.pinOrderKey ?? null,
lastVisitedAt: projection.thread.lastVisitedAt,
titleRegeneration: projection.thread.titleRegeneration ?? null,
limitRecovery: projection.thread.limitRecovery ?? null,
deletedAt: projection.thread.deletedAt,
};
}
Expand Down Expand Up @@ -1495,6 +1496,7 @@ function shellFromState(input: {
pinOrderKey: input.state.thread.pinOrderKey ?? null,
lastVisitedAt: input.state.thread.lastVisitedAt,
titleRegeneration: input.state.thread.titleRegeneration ?? null,
limitRecovery: input.state.thread.limitRecovery ?? null,
deletedAt: input.state.thread.deletedAt,
};
}
Expand Down
109 changes: 109 additions & 0 deletions apps/server/src/orchestration-v2/UsageLimitRecoveryWorker.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,109 @@
import {
CommandId,
MessageId,
type OrchestrationV2ThreadShell,
type OrchestrationV2Command,
} from "@t3tools/contracts";
import * as DateTime from "effect/DateTime";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as Schedule from "effect/Schedule";
import * as ServerSettings from "../serverSettings.ts";
import * as ProjectionStore from "./ProjectionStore.ts";
import * as ThreadManagement from "./ThreadManagementService.ts";

/** The persisted run and reset form the identity of one recovery opportunity. */
export function limitRecoveryCommand(
thread: OrchestrationV2ThreadShell,
autoResume: boolean,
nowMs: number,
): OrchestrationV2Command | null {
if (
thread.status !== "failed" ||
thread.lastErrorClass !== "usage_limit" ||
!thread.latestRunId ||
!thread.usageLimitResetAt ||
thread.archivedAt !== null ||
thread.settledOverride === "settled" ||
thread.pendingRuntimeRequest !== null
)
return null;
const resetMs = Date.parse(thread.usageLimitResetAt);
// An already-expired window reported with a fresh failure cannot start a retry loop.
if (
!Number.isFinite(resetMs) ||
resetMs <= DateTime.toEpochMillis(thread.latestRunCompletedAt ?? thread.updatedAt)
)
return null;
const identity = `${thread.id}:${thread.latestRunId}:${resetMs}`;
const recovery = thread.limitRecovery;
if (recovery?.runId !== thread.latestRunId || recovery.resetAt !== thread.usageLimitResetAt) {
if (!autoResume) return null;
return {
type: "thread.metadata.update",
commandId: CommandId.make(`limit-arm:${identity}`),
threadId: thread.id,
limitRecovery: { runId: thread.latestRunId, resetAt: thread.usageLimitResetAt, autoResume },
};
}
if (
!recovery.autoResume ||
resetMs > nowMs ||
(thread.snoozedUntil != null && DateTime.toEpochMillis(thread.snoozedUntil) > nowMs)
)
return null;
const deliveryIdentity = `${identity}:${recovery.requestId ?? "legacy"}`;
return {
type: "message.dispatch",
commandId: CommandId.make(`limit-resume:${deliveryIdentity}`),
messageId: MessageId.make(`limit-resume:${deliveryIdentity}`),
Comment thread
coderabbitai[bot] marked this conversation as resolved.
threadId: thread.id,
usageLimitContinuationOfRunId: thread.latestRunId,
...(recovery.requestId === undefined
? {}
: { usageLimitRecoveryRequestId: recovery.requestId }),
text: "Continue where you left off.",
attachments: [],
dispatchMode: { type: "start_immediately" },
createdBy: "user",
creationSource: "server",
};
}

const makeSweep = Effect.gen(function* () {
const projections = yield* ProjectionStore.ProjectionStoreV2;
const threads = yield* ThreadManagement.ThreadManagementService;
const settings = yield* ServerSettings.ServerSettingsService;
return Effect.fn("UsageLimitRecoveryWorker.sweep")(function* () {
const preferences = yield* settings.getSettings;
const snapshot = yield* projections.getShellSnapshot();
const nowMs = DateTime.toEpochMillis(yield* DateTime.now);
for (const thread of snapshot.threads) {
const command = limitRecoveryCommand(thread, preferences.autoResumeLimitedThreads, nowMs);
if (command === null) continue;
yield* threads.dispatch(command).pipe(
Effect.catchCause((cause) =>
Effect.logWarning("orchestration-v2.limit-recovery.dispatch-failed", {
threadId: thread.id,
cause,
}),
),
);
}
});
});

// The schedule is derived from persisted failures and thread recovery choices,
// so restarts need no timer restoration and disconnected clients need not run it.
export const workerLive = Layer.effectDiscard(
Effect.gen(function* () {
const sweep = yield* makeSweep;
yield* sweep().pipe(
Effect.catchCause((cause) =>
Effect.logWarning("orchestration-v2.limit-recovery.sweep-failed", { cause }),
),
Effect.repeat(Schedule.spaced("30 seconds")),
Effect.forkScoped,
);
}),
);
Loading
Loading