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
1 change: 1 addition & 0 deletions apps/server/src/environment/ServerEnvironment.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -171,6 +171,7 @@ it.layer(NodeServices.layer)("ServerEnvironmentLive", (it) => {
expect(second.capabilities.threadTitleRegeneration).toBe(true);
expect(second.capabilities.threadPullRequests).toBe(true);
expect(second.capabilities.threadPullRequestLinking).toBe(true);
expect(second.capabilities.guardedSessionStop).toBe(true);
expect(second.capabilities.agentActivityPublishing).toBe(false);
}),
);
Expand Down
1 change: 1 addition & 0 deletions apps/server/src/environment/ServerEnvironment.ts
Original file line number Diff line number Diff line change
Expand Up @@ -222,6 +222,7 @@ export const make = Effect.gen(function* () {
threadSettlement: true,
threadAutoSettlement: true,
threadRestartContinuation: true,
guardedSessionStop: true,
threadSnooze: true,
environmentThemes: true,
usageLimitSources: true,
Expand Down
93 changes: 93 additions & 0 deletions apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ import {
type OrchestrationCommand,
type OrchestrationEvent,
ProviderInstanceId,
ProviderDriverKind,
} from "@t3tools/contracts";
import * as NodeServices from "@effect/platform-node/NodeServices";
import { it as effectIt } from "@effect/vitest";
Expand Down Expand Up @@ -697,6 +698,98 @@ describe("OrchestrationEngine", () => {
}).pipe(Effect.provide(makeOrchestrationLayer())),
);

effectIt.effect("rejects a guarded session stop from a stale thread snapshot", () =>
Effect.gen(function* () {
const engine = yield* OrchestrationEngineService;
const backgroundLiveness = yield* ThreadBackgroundLiveness.ThreadBackgroundLivenessService;
const projectId = ProjectId.make("project-guarded-stop");
const threadId = ThreadId.make("thread-guarded-stop");
yield* engine.dispatch({
type: "project.create",
commandId: CommandId.make("cmd-guarded-stop-project"),
projectId,
title: "Project",
workspaceRoot: "/tmp/project-guarded-stop",
createdAt: now(),
});
yield* engine.dispatch({
type: "thread.create",
commandId: CommandId.make("cmd-guarded-stop-thread"),
threadId,
projectId,
title: "Thread",
modelSelection: {
instanceId: ProviderInstanceId.make("codex"),
model: "gpt-5-codex",
},
interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE,
runtimeMode: "full-access",
branch: null,
worktreePath: null,
createdAt: now(),
});
yield* engine.dispatch({
type: "thread.session.set",
commandId: CommandId.make("cmd-guarded-stop-session"),
threadId,
session: {
threadId,
status: "ready",
providerName: "codex",
providerInstanceId: ProviderInstanceId.make("codex"),
runtimeMode: "full-access",
activeTurnId: null,
lastError: null,
updatedAt: now(),
},
createdAt: now(),
});
const snapshotSequence = yield* engine.latestSequence;
yield* engine.dispatch({
type: "thread.meta.update",
commandId: CommandId.make("cmd-guarded-stop-update"),
threadId,
title: "Changed",
});

const error = yield* engine
.dispatch({
type: "thread.session.stop",
commandId: CommandId.make("cmd-guarded-stop-stale"),
threadId,
createdAt: now(),
onlyIfIdle: true,
snapshotSequence,
expectedProviderName: ProviderDriverKind.make("codex"),
})
.pipe(Effect.flip);

expect(error._tag).toBe("OrchestrationCommandInvariantError");

const freshSequence = yield* engine.latestSequence;
backgroundLiveness.recordTaskLiveness({
threadId,
taskId: "guarded-stop-background-task",
taskType: "subagent",
status: undefined,
kind: "started",
});
const backgroundError = yield* engine
.dispatch({
type: "thread.session.stop",
commandId: CommandId.make("cmd-guarded-stop-background"),
threadId,
createdAt: now(),
onlyIfIdle: true,
snapshotSequence: freshSequence,
expectedProviderName: ProviderDriverKind.make("codex"),
})
.pipe(Effect.flip);
expect(backgroundError._tag).toBe("OrchestrationCommandInvariantError");
backgroundLiveness.clearThreadLiveness(threadId);
}).pipe(Effect.provide(makeOrchestrationLayer())),
);

it("persists deterministic read models for repeated snapshot reads", async () => {
const createdAt = now();
const system = await createOrchestrationSystem();
Expand Down
26 changes: 26 additions & 0 deletions apps/server/src/orchestration/Layers/OrchestrationEngine.ts
Original file line number Diff line number Diff line change
Expand Up @@ -185,6 +185,22 @@ const makeOrchestrationEngine = Effect.gen(function* () {
});
}

if (
envelope.command.type === "thread.session.stop" &&
envelope.command.onlyIfIdle === true &&
envelope.command.snapshotSequence !== undefined &&
(yield* eventStore.hasEventAfter({
aggregateKind: "thread",
aggregateId: envelope.command.threadId,
sequenceExclusive: envelope.command.snapshotSequence,
}))
) {
return yield* new OrchestrationCommandInvariantError({
commandType: envelope.command.type,
detail: `thread ${envelope.command.threadId} changed before guarded session stop`,
});
}

// 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 All @@ -211,6 +227,16 @@ const makeOrchestrationEngine = Effect.gen(function* () {
detail: `thread ${envelope.command.threadId} has live background work`,
});
}
if (
envelope.command.type === "thread.session.stop" &&
envelope.command.onlyIfIdle === true &&
threadBackgroundLiveness.getThreadBackgroundLiveness(envelope.command.threadId) !== null
) {
return yield* new OrchestrationCommandInvariantError({
commandType: envelope.command.type,
detail: `thread ${envelope.command.threadId} has live background work`,
});
}

// New and moved projects do not carry a resolved identity in the event-derived
// command model. Legacy PR edits need it to identify the link they replace.
Expand Down
83 changes: 73 additions & 10 deletions apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -173,6 +173,7 @@ describe("ProviderCommandReactor", () => {
readonly unreadableHistory?: boolean;
readonly titleRegenerationCompletionDispatchFailures?: number;
readonly titleRegenerationBeforeStart?: "one" | "two";
readonly startReactor?: boolean;
readonly serverActivation?: Effect.Effect<void>;
readonly beforeReadySessionDispatch?: () => Effect.Effect<void>;
readonly beforeTurnStartDispatch?: () => Effect.Effect<void>;
Expand All @@ -191,6 +192,11 @@ describe("ProviderCommandReactor", () => {
createdBaseDirs.add(baseDir);
const { stateDir } = deriveServerPathsSync(baseDir, undefined);
createdStateDirs.add(stateDir);
const backgroundLiveness = ThreadBackgroundLiveness.make();
const backgroundLivenessLayer = Layer.succeed(
ThreadBackgroundLiveness.ThreadBackgroundLivenessService,
backgroundLiveness,
);
const runtimeEventPubSub = Effect.runSync(PubSub.unbounded<ProviderRuntimeEvent>());
const tryHandlePromptCommand = vi.fn<ProviderAuthService["Service"]["tryHandlePromptCommand"]>(
input?.tryHandlePromptCommandEffect ?? (() => Effect.succeed(false)),
Expand Down Expand Up @@ -400,7 +406,7 @@ describe("ProviderCommandReactor", () => {

const orchestrationLayer = OrchestrationEngineLive.pipe(
Layer.provide(OrchestrationProjectionSnapshotQueryLive),
Layer.provide(ThreadBackgroundLiveness.layer),
Layer.provide(backgroundLivenessLayer),
Layer.provide(ThreadPlanProgress.layer),
Layer.provide(OrchestrationProjectionPipelineLive),
Layer.provide(OrchestrationEventStoreLive),
Expand All @@ -409,7 +415,7 @@ describe("ProviderCommandReactor", () => {
Layer.provide(SqlitePersistenceMemory),
);
const projectionSnapshotLayer = OrchestrationProjectionSnapshotQueryLive.pipe(
Layer.provide(ThreadBackgroundLiveness.layer),
Layer.provide(backgroundLivenessLayer),
Layer.provide(ThreadPlanProgress.layer),
Layer.provide(RepositoryIdentityResolver.layer),
Layer.provide(SqlitePersistenceMemory),
Expand Down Expand Up @@ -488,6 +494,7 @@ describe("ProviderCommandReactor", () => {
}),
),
Layer.provideMerge(ServerSettingsService.layerTest()),
Layer.provideMerge(backgroundLivenessLayer),
Layer.provideMerge(SqlitePersistenceMemory),
Layer.provideMerge(ServerConfig.layerTest(process.cwd(), baseDir)),
Layer.provideMerge(NodeServices.layer),
Expand Down Expand Up @@ -579,14 +586,18 @@ describe("ProviderCommandReactor", () => {
}

scope = await Effect.runPromise(Scope.make("sequential"));
await Effect.runPromise(
reactor
.start()
.pipe(
Scope.provide(scope),
Effect.provideService(ServerActivation, input?.serverActivation),
),
);
const start = () =>
Effect.runPromise(
reactor
.start()
.pipe(
Scope.provide(scope!),
Effect.provideService(ServerActivation, input?.serverActivation),
),
);
if (input?.startReactor !== false) {
await start();
}
const drain = () => Effect.runPromise(reactor.drain);

return {
Expand Down Expand Up @@ -620,6 +631,8 @@ describe("ProviderCommandReactor", () => {
generateThreadTitle,
runtimeSessions,
stateDir,
backgroundLiveness,
start,
drain,
runEffect,
get titleRegenerationCompletionDispatchAttempts() {
Expand Down Expand Up @@ -4130,6 +4143,56 @@ describe("ProviderCommandReactor", () => {
}),
);

effectIt.effect("skips a guarded stop when background work appears before execution", () =>
Effect.gen(function* () {
const harness = yield* Effect.promise(() => createHarness({ startReactor: false }));
const threadId = ThreadId.make("thread-1");
const now = "2026-01-01T00:00:00.000Z";
yield* harness.engine.dispatch({
type: "thread.session.set",
commandId: CommandId.make("cmd-session-set-before-guarded-stop"),
threadId,
session: {
threadId,
status: "ready",
providerName: "codex",
providerInstanceId: ProviderInstanceId.make("codex"),
runtimeMode: "approval-required",
activeTurnId: null,
lastError: null,
updatedAt: now,
},
createdAt: now,
});
const snapshotSequence = yield* harness.engine.latestSequence;
yield* harness.engine.dispatch({
type: "thread.session.stop",
commandId: CommandId.make("cmd-guarded-stop-before-background"),
threadId,
createdAt: now,
onlyIfIdle: true,
snapshotSequence,
expectedProviderName: ProviderDriverKind.make("codex"),
});
harness.backgroundLiveness.recordTaskLiveness({
threadId,
taskId: "guarded-stop-background-task",
taskType: "subagent",
status: undefined,
kind: "started",
});

yield* Effect.promise(() => harness.start());
yield* Effect.promise(() => harness.drain());

expect(harness.stopSession).not.toHaveBeenCalled();
const thread = yield* harness.snapshotQuery
.getThreadShellById(threadId)
.pipe(Effect.map(Option.getOrThrow));
expect(thread.session?.status).toBe("ready");
}),
);

effectIt.effect("stops a ready provider session after automatic settlement", () =>
Effect.gen(function* () {
const sessionStopped = yield* Deferred.make<void>();
Expand Down
25 changes: 25 additions & 0 deletions apps/server/src/orchestration/Layers/ProviderCommandReactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,8 @@ import {
} from "../Services/ProviderCommandReactor.ts";
import { forkParked, ServerActivation } from "../../serverActivation.ts";
import { canReplaceThreadTitle, DEFAULT_THREAD_TITLE } from "../threadTitles.ts";
import { canStopThreadSessionIfIdle } from "../SessionStopPolicy.ts";
import { threadHasQueuedTurnStart } from "../ThreadSettlementPolicy.ts";
import {
resolveSourceControlWriterModelSelection,
ServerSettingsService,
Expand Down Expand Up @@ -1757,6 +1759,29 @@ const make = Effect.gen(function* () {
}

const now = event.payload.createdAt;
if (
event.payload.onlyIfIdle === true &&
!canStopThreadSessionIfIdle({
expectedProviderName: event.payload.expectedProviderName,
session: thread.session,
latestTurnState: thread.latestTurn?.state ?? null,
hasQueuedTurnStart: threadHasQueuedTurnStart(
thread,
DateTime.formatIso(yield* DateTime.now),
),
hasPendingRequests: thread.hasPendingApprovals || thread.hasPendingUserInput,
backgroundLiveness: thread.backgroundLiveness ?? null,
})
) {
if (thread.session !== null) {
yield* setThreadSession({
threadId: thread.id,
session: thread.session,
createdAt: DateTime.formatIso(yield* DateTime.now),
});
}
return;
}
const wasCompacting = compactingThreadIds.has(thread.id);
stoppingThreadIds.add(thread.id);
const clearStopping = Effect.sync(() => void stoppingThreadIds.delete(thread.id));
Expand Down
51 changes: 51 additions & 0 deletions apps/server/src/orchestration/SessionStopPolicy.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,51 @@
import { ProviderDriverKind, ProviderInstanceId, ThreadId, TurnId } from "@t3tools/contracts";
import { describe, expect, it } from "vite-plus/test";

import { canStopThreadSessionIfIdle } from "./SessionStopPolicy.ts";

const idle = {
expectedProviderName: ProviderDriverKind.make("codex"),
session: {
threadId: ThreadId.make("thread-1"),
status: "ready" as const,
providerName: "codex",
providerInstanceId: ProviderInstanceId.make("codex"),
runtimeMode: "full-access" as const,
activeTurnId: null,
lastError: null,
updatedAt: "2026-09-11T00:00:00.000Z",
},
latestTurnState: null,
hasQueuedTurnStart: false,
hasPendingRequests: false,
backgroundLiveness: null,
};

describe("guarded Codex session stop policy", () => {
it("allows the four non-working session states", () => {
for (const status of ["ready", "idle", "interrupted", "error"] as const) {
expect(canStopThreadSessionIfIdle({ ...idle, session: { ...idle.session, status } })).toBe(
true,
);
}
});

it("rejects every signal that work is active or belongs to another provider", () => {
const unsafe = [
{ ...idle, expectedProviderName: ProviderDriverKind.make("claudeAgent") },
{ ...idle, session: { ...idle.session, providerName: "claudeAgent" } },
{ ...idle, session: { ...idle.session, status: "starting" as const } },
{ ...idle, session: { ...idle.session, status: "running" as const } },
{ ...idle, session: { ...idle.session, activeTurnId: TurnId.make("turn-1") } },
{ ...idle, latestTurnState: "running" as const },
{ ...idle, hasQueuedTurnStart: true },
{ ...idle, hasPendingRequests: true },
{ ...idle, backgroundLiveness: "working" as const },
{ ...idle, backgroundLiveness: "monitoring" as const },
];

for (const input of unsafe) {
expect(canStopThreadSessionIfIdle(input)).toBe(false);
}
});
});
Loading
Loading