Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
4de1078
fix(web): restore composer tasks after relaunch
sethwebster Aug 24, 2026
fbb148b
test(web): cover composer tasks after relaunch
sethwebster Aug 24, 2026
6122258
fix(web): restore composer tasks from active plan
sethwebster Aug 24, 2026
c00b865
docs: add composer task relaunch evidence
sethwebster Aug 24, 2026
8e06465
chore: keep PR evidence out of the repository
sethwebster Aug 24, 2026
6feb266
test(server): cover active plan outside detail limit
sethwebster Aug 24, 2026
571fe15
fix(server): retain active plan in bounded detail
sethwebster Aug 24, 2026
ff58d07
test(server): require indexed active plan lookup
sethwebster Aug 24, 2026
d4187ef
fix(server): index active plan hydration
sethwebster Aug 24, 2026
802ea26
refactor(web): share active task derivation
sethwebster Aug 24, 2026
b31d5ae
test: cover task restart wiring and query plan
sethwebster Aug 24, 2026
a3d222d
docs(user): explain composer task progress
sethwebster Aug 24, 2026
cf7673a
test(server): cover provider-aborted turn lifecycle
sethwebster Aug 24, 2026
52a3944
fix(server): settle provider-aborted turns
sethwebster Aug 24, 2026
fb3bd8e
test: cover active-turn divergence and untargeted abort
sethwebster Aug 24, 2026
7d4698c
fix(server): clear progress only for accepted aborts
sethwebster Aug 24, 2026
ec648cf
fix: bind restored tasks to the active session turn
sethwebster Aug 24, 2026
80ce111
test(server): cover abort finalization
sethwebster Aug 24, 2026
82f78bc
fix(server): finalize buffered content on abort
sethwebster Aug 24, 2026
1e8276a
fix(web): key task dismissal to the active turn
sethwebster Aug 24, 2026
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
138 changes: 138 additions & 0 deletions apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2405,6 +2405,144 @@ projectionSnapshotLayer("ProjectionSnapshotQuery windowed thread detail", (it) =
}),
);

it.effect("preserves the latest active-turn plan outside the activity hydration limit", () =>
Effect.gen(function* () {
yield* seedFanOutThread();
const snapshotQuery = yield* ProjectionSnapshotQuery;
const sql = yield* SqlClient.SqlClient;

yield* sql`DELETE FROM projection_thread_activities`;
yield* sql`
UPDATE projection_turns
SET state = 'running', completed_at = NULL
WHERE thread_id = 'thread-w' AND turn_id = 'turn-5'
`;
yield* sql`
INSERT INTO projection_thread_activities (
activity_id, thread_id, turn_id, tone, kind, summary, payload_json, sequence, created_at
)
VALUES (
'active-plan', 'thread-w', 'turn-5', 'info', 'turn.plan.updated', 'Plan updated',
'{"plan":[{"step":"Keep working","status":"inProgress"}]}', 1,
'2026-03-01T00:04:00.000Z'
)
`;
yield* sql`
WITH RECURSIVE activity_rows(sequence) AS (
SELECT 2
UNION ALL
SELECT sequence + 1 FROM activity_rows WHERE sequence < 501
)
INSERT INTO projection_thread_activities (
activity_id, thread_id, turn_id, tone, kind, summary, payload_json, sequence, created_at
)
SELECT
printf('activity-%04d', sequence),
'thread-w',
'turn-5',
'tool',
'tool.completed',
'ran tool',
printf('{"sequence":%d}', sequence),
sequence,
'2026-03-01T00:04:00.000Z'
FROM activity_rows
`;

const fullDetail = yield* snapshotQuery.getThreadDetailById(threadW);
assert.equal(fullDetail._tag, "Some");
if (fullDetail._tag === "Some") {
assert.equal(fullDetail.value.activities.length, 501);
assert.equal(fullDetail.value.activities[0]?.id, asEventId("active-plan"));
}

const windowedDetail = yield* snapshotQuery.getThreadDetailSnapshot(threadW, {
turnLimit: 2,
});
assert.equal(windowedDetail._tag, "Some");
if (windowedDetail._tag === "Some") {
assert.equal(windowedDetail.value.thread.activities.length, 501);
assert.equal(windowedDetail.value.thread.activities[0]?.id, asEventId("active-plan"));
}
}),
);

it.effect("pins the session's active-turn plan when latest turn projection diverges", () =>
Effect.gen(function* () {
yield* seedFanOutThread();
const snapshotQuery = yield* ProjectionSnapshotQuery;
const sql = yield* SqlClient.SqlClient;

yield* sql`DELETE FROM projection_thread_activities`;
yield* sql`
INSERT INTO projection_turns (
thread_id, turn_id, pending_message_id, state, requested_at, started_at, completed_at,
checkpoint_files_json
)
VALUES (
'thread-w', 'turn-active', NULL, 'running',
'2026-03-01T00:05:00.000Z', '2026-03-01T00:05:00.000Z', NULL, '[]'
)
`;
yield* sql`
INSERT INTO projection_thread_sessions (
thread_id, status, provider_name, provider_instance_id, runtime_mode,
active_turn_id, last_error, updated_at
)
VALUES (
'thread-w', 'running', 'codex', 'codex', 'full-access',
'turn-active', NULL, '2026-03-01T00:05:00.000Z'
)
`;
yield* sql`
INSERT INTO projection_thread_activities (
activity_id, thread_id, turn_id, tone, kind, summary, payload_json, sequence, created_at
)
VALUES
(
'old-plan', 'thread-w', 'turn-5', 'info', 'turn.plan.updated', 'Old plan',
'{"plan":[{"step":"Old turn","status":"inProgress"}]}', 1,
'2026-03-01T00:04:00.000Z'
),
(
'active-plan', 'thread-w', 'turn-active', 'info', 'turn.plan.updated', 'Active plan',
'{"plan":[{"step":"Active turn","status":"inProgress"}]}', 2,
'2026-03-01T00:05:00.000Z'
)
`;
yield* sql`
WITH RECURSIVE activity_rows(sequence) AS (
SELECT 3
UNION ALL
SELECT sequence + 1 FROM activity_rows WHERE sequence < 502
)
INSERT INTO projection_thread_activities (
activity_id, thread_id, turn_id, tone, kind, summary, payload_json, sequence, created_at
)
SELECT
printf('activity-%04d', sequence),
'thread-w',
'turn-active',
'tool',
'tool.completed',
'ran tool',
printf('{"sequence":%d}', sequence),
sequence,
'2026-03-01T00:05:00.000Z'
FROM activity_rows
`;

const detail = yield* snapshotQuery.getThreadDetailById(threadW);
assert.equal(detail._tag, "Some");
if (detail._tag === "Some") {
const ids = detail.value.activities.map((activity) => activity.id);
assert.equal(ids.length, 501);
assert.equal(ids.includes(asEventId("active-plan")), true);
assert.equal(ids.includes(asEventId("old-plan")), false);
}
}),
);

it.effect("a thread with no turns returns its content unwindowed on the first page", () =>
Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;
Expand Down
48 changes: 45 additions & 3 deletions apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1253,9 +1253,10 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () {
`,
});

// Blocking request payloads must remain available even if they predate the
// recent activity window. Each CTE returns at most one unresolved row per
// request, so the merge below stays bounded by actionable work.
// Blocking request payloads and the current turn's latest plan must remain
// available even if they predate the recent activity window. Each CTE
// returns at most one row per request or active plan, so the merge below
// stays bounded by actionable work.
const listPinnedThreadActivityRowsByThread = SqlSchema.findAll({
Request: ThreadIdLookupInput,
Result: ProjectionThreadActivityDbRowSchema,
Expand Down Expand Up @@ -1315,6 +1316,44 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () {
)
AND json_extract(activity.payload_json, '$.requestId') IS NOT NULL
),
restorable_turn_candidate AS (
SELECT
CASE
WHEN session.status = 'running' THEN session.active_turn_id
ELSE thread.latest_turn_id
END AS turn_id,
session.status AS session_status
FROM projection_threads AS thread
LEFT JOIN projection_thread_sessions AS session
ON session.thread_id = thread.thread_id
WHERE thread.thread_id = ${threadId}
),
latest_unsettled_turn AS (
SELECT candidate.turn_id
FROM restorable_turn_candidate AS candidate
LEFT JOIN projection_turns AS turn
ON turn.thread_id = ${threadId}
AND turn.turn_id = candidate.turn_id
WHERE candidate.turn_id IS NOT NULL
AND (
candidate.session_status = 'running'
OR turn.started_at IS NULL
OR turn.completed_at IS NULL
)
),
active_plan_activities AS (
SELECT activity.activity_id
FROM latest_unsettled_turn AS active_turn
CROSS JOIN projection_thread_activities AS activity
WHERE activity.thread_id = ${threadId}
AND activity.turn_id = active_turn.turn_id
AND activity.kind = 'turn.plan.updated'
ORDER BY
activity.sequence DESC,
activity.created_at DESC,
activity.activity_id DESC
LIMIT 1
),
pinned_activity_ids AS (
SELECT activity_id
FROM pending_approval_activities
Expand All @@ -1324,6 +1363,9 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () {
FROM user_input_lifecycle
WHERE request_order = 1
AND kind = 'user-input.requested'
UNION ALL
SELECT activity_id
FROM active_plan_activities
)
SELECT
activity.activity_id AS "activityId",
Expand Down
163 changes: 163 additions & 0 deletions apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -321,6 +321,7 @@ describe("ProviderRuntimeIngestion", () => {
engine,
dispatch,
readModel: () => Effect.runPromise(snapshotQuery.getSnapshot()),
readShell: () => Effect.runPromise(snapshotQuery.getShellSnapshot()),
emit: provider.emit,
setProviderSession: provider.setSession,
drain,
Expand Down Expand Up @@ -740,6 +741,168 @@ describe("ProviderRuntimeIngestion", () => {
);
});

it("settles the active turn when the provider aborts it", async () => {
const harness = await createHarness();
const turnId = asTurnId("turn-aborted");
const startedAt = "2026-01-01T00:00:00.000Z";
const abortedAt = "2026-01-01T00:00:01.000Z";

harness.emit({
type: "turn.started",
eventId: asEventId("evt-turn-started-aborted"),
provider: ProviderDriverKind.make("codex"),
createdAt: startedAt,
threadId: asThreadId("thread-1"),
turnId,
});

await harness.drain();
const runningReadModel = await harness.readModel();
const runningThread = runningReadModel.threads.find(
(entry) => entry.id === ThreadId.make("thread-1"),
);
expect(runningThread?.session?.status).toBe("running");
expect(runningThread?.session?.activeTurnId).toBe(turnId);

harness.emit({
type: "turn.plan.updated",
eventId: asEventId("evt-turn-plan-aborted"),
provider: ProviderDriverKind.make("codex"),
createdAt: "2026-01-01T00:00:00.250Z",
threadId: asThreadId("thread-1"),
turnId,
payload: { plan: [{ step: "Keep working", status: "inProgress" }] },
});
harness.emit({
type: "content.delta",
eventId: asEventId("evt-turn-aborted-assistant-delta"),
provider: ProviderDriverKind.make("codex"),
createdAt: "2026-01-01T00:00:00.250Z",
threadId: asThreadId("thread-1"),
turnId,
itemId: asItemId("item-turn-aborted"),
payload: {
streamKind: "assistant_text",
delta: "Partial answer before abort.",
},
});
harness.emit({
type: "turn.proposed.delta",
eventId: asEventId("evt-turn-aborted-plan-delta"),
provider: ProviderDriverKind.make("codex"),
createdAt: "2026-01-01T00:00:00.250Z",
threadId: asThreadId("thread-1"),
turnId,
payload: { delta: "## Partial plan before abort" },
});
await harness.drain();
const shellWithPlan = await harness.readShell();
expect(shellWithPlan.threads.find((entry) => entry.id === "thread-1")?.planProgress).toEqual({
step: "Keep working",
completedSteps: 0,
totalSteps: 1,
});

harness.emit({
type: "turn.aborted",
eventId: asEventId("evt-untargeted-turn-aborted"),
provider: ProviderDriverKind.make("codex"),
createdAt: "2026-01-01T00:00:00.375Z",
threadId: asThreadId("thread-1"),
payload: { reason: "An untargeted abort arrived." },
});
await harness.drain();
const shellAfterUntargetedAbort = await harness.readShell();
expect(
shellAfterUntargetedAbort.threads.find((entry) => entry.id === "thread-1")?.planProgress,
).toEqual({
step: "Keep working",
completedSteps: 0,
totalSteps: 1,
});

harness.emit({
type: "turn.aborted",
eventId: asEventId("evt-other-turn-aborted"),
provider: ProviderDriverKind.make("codex"),
createdAt: "2026-01-01T00:00:00.500Z",
threadId: asThreadId("thread-1"),
turnId: asTurnId("turn-other"),
payload: { reason: "A different turn aborted." },
});

await harness.drain();
const activeReadModel = await harness.readModel();
const activeThread = activeReadModel.threads.find(
(entry) => entry.id === ThreadId.make("thread-1"),
);
expect(activeThread?.session?.status).toBe("running");
expect(activeThread?.session?.activeTurnId).toBe(turnId);
const shellAfterOtherAbort = await harness.readShell();
expect(
shellAfterOtherAbort.threads.find((entry) => entry.id === "thread-1")?.planProgress,
).toEqual({
step: "Keep working",
completedSteps: 0,
totalSteps: 1,
});

harness.emit({
type: "turn.aborted",
eventId: asEventId("evt-turn-aborted"),
provider: ProviderDriverKind.make("codex"),
createdAt: abortedAt,
threadId: asThreadId("thread-1"),
turnId,
payload: { reason: "Provider aborted the turn." },
});

await harness.drain();
const readModel = await harness.readModel();
const thread = readModel.threads.find((entry) => entry.id === ThreadId.make("thread-1"));
expect(thread?.session?.status).toBe("interrupted");
expect(thread?.session?.activeTurnId).toBeNull();
expect(thread?.latestTurn).toMatchObject({
turnId,
state: "interrupted",
completedAt: abortedAt,
});
expect(
thread?.messages.find((message) => message.id === "assistant:item-turn-aborted"),
).toMatchObject({
text: "Partial answer before abort.",
streaming: false,
});
expect(
thread?.proposedPlans.find(
(proposedPlan) => proposedPlan.id === "plan:thread-1:turn:turn-aborted",
),
).toMatchObject({
planMarkdown: "## Partial plan before abort",
});
const shellAfterActiveAbort = await harness.readShell();
expect(
shellAfterActiveAbort.threads.find((entry) => entry.id === "thread-1")?.planProgress,
).toBeNull();

harness.emit({
type: "turn.aborted",
eventId: asEventId("evt-late-turn-aborted"),
provider: ProviderDriverKind.make("codex"),
createdAt: "2026-01-01T00:00:02.000Z",
threadId: asThreadId("thread-1"),
turnId,
payload: { reason: "Late duplicate abort." },
});
await harness.drain();
const afterLateAbort = await harness.readModel();
const threadAfterLateAbort = afterLateAbort.threads.find(
(entry) => entry.id === ThreadId.make("thread-1"),
);
expect(threadAfterLateAbort?.session?.status).toBe("interrupted");
expect(threadAfterLateAbort?.session?.activeTurnId).toBeNull();
});

it("accepts claude turn lifecycle when seeded thread id is a synthetic placeholder", async () => {
const harness = await createHarness();
const seededAt = "2026-01-01T00:00:00.000Z";
Expand Down
Loading
Loading