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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -1,7 +1,9 @@
import { ScreenScrollView as ScrollView } from "../../components/ScreenScrollView";
import { SymbolView } from "../../components/AppSymbol";
import { AppText as Text } from "../../components/AppText";
import { AppText as Text, AppTextInput } from "../../components/AppText";
import {
PROVIDER_SEND_TURN_MAX_INPUT_CHARS,
UsageLimitContinuationPrompt,
type ResponseStreamingMode,
type ServerSettings,
type ServerSettingsPatch,
Expand All @@ -10,6 +12,7 @@ import {
type ProjectScopedServerSettingKey,
} from "@t3tools/contracts";
import { useRef, useState, type ComponentProps } from "react";
import * as Schema from "effect/Schema";
import { Alert, Platform, Pressable, View } from "react-native";
import { useSafeAreaInsets } from "react-native-safe-area-context";

Expand All @@ -33,6 +36,8 @@ import {
type ScopedMobileSettingsTarget,
} from "./settings-scoped-server";

const isUsageLimitContinuationPrompt = Schema.is(UsageLimitContinuationPrompt);

type SettingsPage = "new-threads" | "source-control" | "agent-behavior" | "maintenance";

const PAGE_TITLES: Record<SettingsPage, string> = {
Expand Down Expand Up @@ -365,6 +370,30 @@ function ServerSettingsDetail(props: { readonly page: SettingsPage }) {
onValueChange={(value) => write({ continueThreadsAfterServerUpdate: value })}
/>
</View>
<View className="border-t border-border-subtle">
<FanoutSwitchRow
icon="arrow.uturn.forward"
label="Continue after usage limits"
subtitle={
projectSelected
? "Environment-wide setting. Select All projects to change it."
: "Continue Codex threads when subscription usage is available again."
}
value={uniform("continueThreadsAfterUsageLimit")}
disabled={disabledFor("continueThreadsAfterUsageLimit")}
onValueChange={(value) => write({ continueThreadsAfterUsageLimit: value })}
/>
</View>
{uniform("continueThreadsAfterUsageLimit") === true ? (
<UsageLimitPromptEditor
key={`${displayTargets.map((target) => target.environment.environmentId).join(",")}:${projectSelected}:${uniform("usageLimitContinuationPrompt")}`}
value={uniform("usageLimitContinuationPrompt")}
disabled={disabledFor("usageLimitContinuationPrompt")}
onSave={(usageLimitContinuationPrompt) =>
write({ usageLimitContinuationPrompt })
}
/>
) : null}
</SettingsSection>
) : null}
</>
Expand All @@ -375,6 +404,52 @@ function ServerSettingsDetail(props: { readonly page: SettingsPage }) {
);
}

function UsageLimitPromptEditor(props: {
readonly value: string | null;
readonly disabled: boolean;
readonly onSave: (value: string) => void;
}) {
const [draft, setDraft] = useState<string | null>(null);
const [error, setError] = useState<string | null>(null);
return (
<View className="gap-2 border-t border-border-subtle p-4">
<Text className="text-base text-foreground">Continuation prompt</Text>
<AppTextInput
accessibilityLabel="Continuation prompt"
className="min-h-24 rounded-xl bg-card px-3 py-2 text-base text-foreground"
multiline
textAlignVertical="top"
value={draft ?? props.value ?? ""}
placeholder={props.value === null ? "Mixed" : undefined}
maxLength={PROVIDER_SEND_TURN_MAX_INPUT_CHARS}
editable={!props.disabled}
onChangeText={(text) => {
setDraft(text);
setError(null);
}}
onBlur={() => {
if (draft === null || props.disabled) return;
const prompt = draft.trim();
if (!isUsageLimitContinuationPrompt(prompt)) {
setError("Enter a continuation prompt.");
return;
}
setDraft(prompt);
if (prompt !== props.value) props.onSave(prompt);
}}
/>
{error ? (
<Text accessibilityRole="alert" className="text-sm text-destructive">
{error}
</Text>
) : null}
<Text className="text-sm text-foreground-muted">
The server must remain running to continue.
</Text>
</View>
);
}

function ChoiceRow(props: {
readonly label: string;
readonly description: string;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,7 @@ import { RuntimeReceiptBusTest } from "../src/orchestration/Layers/RuntimeReceip
import { OrchestrationReactorLive } from "../src/orchestration/Layers/OrchestrationReactor.ts";
import { ProviderCommandReactorLive } from "../src/orchestration/Layers/ProviderCommandReactor.ts";
import { ProviderRuntimeIngestionLive } from "../src/orchestration/Layers/ProviderRuntimeIngestion.ts";
import { UsageLimitContinuation } from "../src/orchestration/UsageLimitContinuation.ts";
import { CheckpointReactor } from "../src/orchestration/Services/CheckpointReactor.ts";
import { ProviderRuntimeIngestionService } from "../src/orchestration/Services/ProviderRuntimeIngestion.ts";
import {
Expand Down Expand Up @@ -268,6 +269,7 @@ export const makeOrchestrationIntegrationHarness = (

const persistenceLayer = makeSqlitePersistenceLive(dbPath);
const orchestrationLayer = OrchestrationEngineLive.pipe(
Layer.provide(ServerSettingsService.layerTest()),
Layer.provide(OrchestrationProjectionPipelineLive),
Layer.provide(OrchestrationEventStoreLive),
Layer.provide(OrchestrationCommandReceiptRepositoryLive),
Expand Down Expand Up @@ -388,6 +390,12 @@ export const makeOrchestrationIntegrationHarness = (
}),
),
Layer.provideMerge(runtimeIngestionLayer),
Layer.provideMerge(
Layer.mock(UsageLimitContinuation)({
start: () => Effect.void,
recordFailure: () => Effect.void,
}),
),
Layer.provideMerge(providerCommandReactorLayer),
Layer.provideMerge(checkpointReactorLayer),
Layer.provideMerge(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,7 @@ const stoppedBindingResumeCursor = {
const makePersistedRuntimeLayer = (dbPath: string) => {
const persistence = makeSqlitePersistenceLive(dbPath);
const orchestration = OrchestrationLayerLive.pipe(
Layer.provide(ServerSettings.layerTest()),
Layer.provideMerge(RepositoryIdentityResolver.layer),
Layer.provideMerge(persistence),
);
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/bin.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@ import * as ServerEnvironment from "./environment/ServerEnvironment.ts";
import * as ProjectionSnapshotQuery from "./orchestration/Services/ProjectionSnapshotQuery.ts";
import * as OrchestrationEngine from "./orchestration/Services/OrchestrationEngine.ts";
import { OrchestrationLayerLive } from "./orchestration/runtimeLayer.ts";
import { ServerSettingsService } from "./serverSettings.ts";
import { orchestrationHttpApiLayer } from "./orchestration/http.ts";
import * as ProjectCloneTracker from "./project/ProjectCloneTracker.ts";
import { layerConfig as SqlitePersistenceLayerLive } from "./persistence/Layers/Sqlite.ts";
Expand Down Expand Up @@ -128,6 +129,7 @@ const makeCliTestServerConfig = (baseDir: string) =>
const makeProjectPersistenceLayer = (config: ServerConfig.ServerConfig["Service"]) =>
Layer.mergeAll(
OrchestrationLayerLive.pipe(
Layer.provide(ServerSettingsService.layerTest()),
Layer.provideMerge(RepositoryIdentityResolver.layer),
Layer.provideMerge(SqlitePersistenceLayerLive),
),
Expand Down
3 changes: 3 additions & 0 deletions apps/server/src/cli/project.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,8 @@ import { FetchHttpClient, HttpClient, HttpClientError } from "effect/unstable/ht
import * as HttpApiClient from "effect/unstable/httpapi/HttpApiClient";

import * as EnvironmentAuth from "../auth/EnvironmentAuth.ts";
import * as ServerSecretStore from "../auth/ServerSecretStore.ts";
import * as ServerSettings from "../serverSettings.ts";

import * as ServerConfig from "../config.ts";
import * as OrchestrationEngine from "../orchestration/Services/OrchestrationEngine.ts";
Expand Down Expand Up @@ -200,6 +202,7 @@ const projectCommandUuid = Crypto.Crypto.pipe(
const ProjectCliRuntimeLive = Layer.mergeAll(
WorkspacePaths.layer,
OrchestrationLayerLive.pipe(
Layer.provide(ServerSettings.layer.pipe(Layer.provide(ServerSecretStore.layer))),
Layer.provideMerge(RepositoryIdentityResolver.layer),
Layer.provideMerge(SqlitePersistenceLayerLive),
),
Expand Down
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import { ServerSettingsService } from "../../serverSettings.ts";
// @effect-diagnostics nodeBuiltinImport:off
import * as NodeFS from "node:fs";
import * as NodeOS from "node:os";
Expand Down Expand Up @@ -323,6 +324,7 @@ describe("CheckpointReactor", () => {
options?.providerName ?? ProviderDriverKind.make("codex"),
);
const orchestrationLayer = OrchestrationEngineLive.pipe(
Layer.provide(ServerSettingsService.layerTest()),
Layer.provide(OrchestrationProjectionSnapshotQueryLive),
Layer.provide(ThreadBackgroundLiveness.layer),
Layer.provide(ThreadPlanProgress.layer),
Expand Down
12 changes: 12 additions & 0 deletions apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,7 @@ import {
} from "../Services/ProjectionPipeline.ts";
import { ProjectionSnapshotQuery } from "../Services/ProjectionSnapshotQuery.ts";
import { ServerConfig } from "../../config.ts";
import { ServerSettingsService } from "../../serverSettings.ts";

const asProjectId = (value: string): ProjectId => ProjectId.make(value);
const asMessageId = (value: string): MessageId => MessageId.make(value);
Expand All @@ -72,6 +73,7 @@ function makeOrchestrationLayer(
});
return Layer.mergeAll(
OrchestrationEngineLive.pipe(
Layer.provideMerge(ServerSettingsService.layerTest({ continueThreadsAfterUsageLimit: true })),
Layer.provide(OrchestrationProjectionSnapshotQueryLive),
Layer.provide(OrchestrationProjectionPipelineLive),
),
Expand Down Expand Up @@ -417,6 +419,7 @@ describe("OrchestrationEngine", () => {
let fullSnapshotReadCount = 0;

const layer = OrchestrationEngineLive.pipe(
Layer.provideMerge(ServerSettingsService.layerTest({ continueThreadsAfterUsageLimit: true })),
Layer.provide(
Layer.succeed(ProjectionSnapshotQuery, {
getUserInputActivity: () => Effect.die("unused"),
Expand Down Expand Up @@ -1511,6 +1514,9 @@ describe("OrchestrationEngine", () => {

const runtime = ManagedRuntime.make(
OrchestrationEngineLive.pipe(
Layer.provideMerge(
ServerSettingsService.layerTest({ continueThreadsAfterUsageLimit: true }),
),
Layer.provide(OrchestrationProjectionSnapshotQueryLive),
Layer.provide(ThreadBackgroundLiveness.layer),
Layer.provide(ThreadPlanProgress.layer),
Expand Down Expand Up @@ -1619,6 +1625,9 @@ describe("OrchestrationEngine", () => {

const runtime = ManagedRuntime.make(
OrchestrationEngineLive.pipe(
Layer.provideMerge(
ServerSettingsService.layerTest({ continueThreadsAfterUsageLimit: true }),
),
Layer.provide(OrchestrationProjectionSnapshotQueryLive),
Layer.provide(ThreadBackgroundLiveness.layer),
Layer.provide(ThreadPlanProgress.layer),
Expand Down Expand Up @@ -1768,6 +1777,9 @@ describe("OrchestrationEngine", () => {

const runtime = ManagedRuntime.make(
OrchestrationEngineLive.pipe(
Layer.provideMerge(
ServerSettingsService.layerTest({ continueThreadsAfterUsageLimit: true }),
),
Layer.provide(OrchestrationProjectionSnapshotQueryLive),
Layer.provide(ThreadBackgroundLiveness.layer),
Layer.provide(ThreadPlanProgress.layer),
Expand Down
46 changes: 45 additions & 1 deletion apps/server/src/orchestration/Layers/OrchestrationEngine.ts
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,8 @@ import {
type OrchestrationDispatchError,
type OrchestrationProjectorDecodeError,
} from "../Errors.ts";
import { decideOrchestrationCommand } from "../decider.ts";
import { cancelsUsageLimitContinuation, decideOrchestrationCommand } from "../decider.ts";
import { ServerSettingsService } from "../../serverSettings.ts";
import { createEmptyReadModel, projectEvent } from "../projector.ts";
import { OrchestrationProjectionPipeline } from "../Services/ProjectionPipeline.ts";
import { ProjectionSnapshotQuery } from "../Services/ProjectionSnapshotQuery.ts";
Expand Down Expand Up @@ -83,6 +84,7 @@ function commandToAggregateRef(command: OrchestrationCommand): {

const makeOrchestrationEngine = Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;
const settingsService = yield* ServerSettingsService;
const eventStore = yield* OrchestrationEventStore;
const commandReceiptRepository = yield* OrchestrationCommandReceiptRepository;
const projectionPipeline = yield* OrchestrationProjectionPipeline;
Expand Down Expand Up @@ -185,6 +187,48 @@ const makeOrchestrationEngine = Effect.gen(function* () {
});
}

if (
(envelope.command.type === "thread.turn.start" ||
envelope.command.type === "thread.session.set") &&
envelope.command.expectedUsageLimit !== undefined
) {
const command = envelope.command;
const canceled = yield* eventStore
.readAggregateRange({
aggregateKind: "thread",
aggregateId: command.threadId,
fromSequenceExclusive: envelope.command.expectedUsageLimit.snapshotSequence,
toSequenceInclusive: commandReadModel.snapshotSequence,
limit:
commandReadModel.snapshotSequence -
envelope.command.expectedUsageLimit.snapshotSequence,
})
.pipe(Stream.filter(cancelsUsageLimitContinuation), Stream.runHead);
if (Option.isSome(canceled)) {
return yield* new OrchestrationCommandInvariantError({
commandType: command.type,
detail: `thread ${command.threadId} changed before usage-limit continuation`,
});
}
if (command.type === "thread.turn.start") {
const settings = yield* settingsService.getSettings.pipe(
Effect.mapError(
() =>
new OrchestrationCommandInvariantError({
commandType: command.type,
detail: "Could not read usage-limit continuation settings.",
}),
),
);
if (!settings.continueThreadsAfterUsageLimit) {
return yield* new OrchestrationCommandInvariantError({
commandType: command.type,
detail: "Usage-limit continuation is disabled.",
});
}
}
}

// The decider compares the lookup inputs. Only recreation needs an
// event check, since it can reset a thread to the same field values.
if (
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import * as ThreadPullRequestReactor from "../ThreadPullRequestReactor.ts";
import { OrchestrationReactor } from "../Services/OrchestrationReactor.ts";
import { makeOrchestrationReactor } from "./OrchestrationReactor.ts";
import * as AgentAwarenessRelay from "../../relay/AgentAwarenessRelay.ts";
import { UsageLimitContinuation } from "../UsageLimitContinuation.ts";
import { StorageCleanup } from "../../storageCleanup.ts";

describe("OrchestrationReactor", () => {
Expand All @@ -32,6 +33,7 @@ describe("OrchestrationReactor", () => {

runtime = ManagedRuntime.make(
Layer.effect(OrchestrationReactor, makeOrchestrationReactor).pipe(
Layer.provide(Layer.mock(UsageLimitContinuation)({ start: () => Effect.void })),
Layer.provideMerge(
Layer.succeed(StorageCleanup, {
start: () => {
Expand Down
3 changes: 3 additions & 0 deletions apps/server/src/orchestration/Layers/OrchestrationReactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 UsageLimitContinuation from "../UsageLimitContinuation.ts";
import * as PullRequestSyncReactor from "../PullRequestSyncReactor.ts";
import * as ThreadPullRequestReactor from "../ThreadPullRequestReactor.ts";
import * as AgentAwarenessRelay from "../../relay/AgentAwarenessRelay.ts";
Expand All @@ -21,12 +22,14 @@ export const makeOrchestrationReactor = Effect.gen(function* () {
const checkpointReactor = yield* CheckpointReactor;
const threadDeletionReactor = yield* ThreadDeletionReactor;
const threadSettlementReactor = yield* ThreadSettlementReactor.ThreadSettlementReactor;
const usageLimitContinuation = yield* UsageLimitContinuation.UsageLimitContinuation;
const pullRequestSyncReactor = yield* PullRequestSyncReactor.PullRequestSyncReactor;
const threadPullRequestReactor = yield* ThreadPullRequestReactor.ThreadPullRequestReactor;
const agentAwarenessRelay = yield* AgentAwarenessRelay.AgentAwarenessRelay;
const storageCleanup = yield* StorageCleanup.StorageCleanup;

const start: OrchestrationReactorShape["start"] = Effect.fn("start")(function* () {
yield* usageLimitContinuation.start();
yield* providerRuntimeIngestion.start();
yield* providerCommandReactor.start();
yield* checkpointReactor.start();
Expand Down
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import { ServerSettingsService } from "../../serverSettings.ts";
import {
ApprovalRequestId,
CheckpointRef,
Expand Down Expand Up @@ -4355,6 +4356,7 @@ it.effect("restores pending turn-start metadata across projection pipeline resta

const engineLayer = it.layer(
OrchestrationEngineLive.pipe(
Layer.provide(ServerSettingsService.layerTest()),
Layer.provideMerge(OrchestrationProjectionSnapshotQueryLive),
Layer.provide(ThreadBackgroundLiveness.layer),
Layer.provide(ThreadPlanProgress.layer),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1052,6 +1052,15 @@ projectionSnapshotLayer("ProjectionSnapshotQuery", (it) => {
);
assert.equal(archivedShellSnapshot.threads[0]?.archivedAt, "2026-04-06T00:00:06.000Z");
assert.deepEqual(archivedShellSnapshot.threads[0]?.branchPullRequest, branchPullRequest);
assert.equal(
(yield* snapshotQuery.getThreadShellById(ThreadId.make("thread-archived")))._tag,
"None",
);
const archivedThread = yield* snapshotQuery.getThreadShellById(
ThreadId.make("thread-archived"),
{ includeArchived: true },
);
assert.equal(Option.getOrThrow(archivedThread).archivedAt, "2026-04-06T00:00:06.000Z");
const activeContext = yield* snapshotQuery.getThreadRuntimeContext(
ThreadId.make("thread-active"),
);
Expand Down Expand Up @@ -1086,6 +1095,12 @@ projectionSnapshotLayer("ProjectionSnapshotQuery", (it) => {
},
]);
assert.deepEqual((yield* snapshotQuery.getArchivedShellSnapshot()).threads, []);
assert.equal(
(yield* snapshotQuery.getThreadShellById(ThreadId.make("thread-archived"), {
includeArchived: true,
}))._tag,
"None",
);
yield* sql`
UPDATE projection_projects
SET deleted_at = '2026-04-06T00:00:10.000Z', updated_at = '2026-04-06T00:00:10.000Z'
Expand Down
Loading
Loading