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
139 changes: 139 additions & 0 deletions apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3259,6 +3259,145 @@ it.layer(makeProjectionPipelinePrefixedTestLayer("t3-pending-turn-terminal-test-
assert.deepEqual(pendingRows, [{ messageId: "new-message" }]);
}),
);

it.effect("promotes a turn start deferred behind compaction when the session restores", () =>
Effect.gen(function* () {
const projectionPipeline = yield* OrchestrationProjectionPipeline;
const eventStore = yield* OrchestrationEventStore;
const sql = yield* SqlClient.SqlClient;
const threadId = ThreadId.make("thread-deferred-compaction-turn");
const compactMessageId = MessageId.make("message-deferred-compact");
const queuedMessageId = MessageId.make("message-deferred-queued");
const compactAt = "2026-02-26T16:00:00.000Z";
const queuedAt = "2026-02-26T16:00:01.000Z";

for (const [index, message] of [
{ id: compactMessageId, text: "/compact", createdAt: compactAt },
{ id: queuedMessageId, text: "send after compact", createdAt: queuedAt },
].entries()) {
yield* eventStore.append({
type: "thread.message-sent",
eventId: EventId.make(`evt-deferred-message-${index}`),
aggregateKind: "thread",
aggregateId: threadId,
occurredAt: message.createdAt,
commandId: CommandId.make(`cmd-deferred-message-${index}`),
causationEventId: null,
correlationId: CorrelationId.make(`cmd-deferred-message-${index}`),
metadata: {},
payload: {
threadId,
messageId: message.id,
role: "user",
text: message.text,
attachments: [],
turnId: null,
streaming: false,
createdAt: message.createdAt,
updatedAt: message.createdAt,
},
});
yield* eventStore.append({
type: "thread.turn-start-requested",
eventId: EventId.make(`evt-deferred-turn-${index}`),
aggregateKind: "thread",
aggregateId: threadId,
occurredAt: message.createdAt,
commandId: CommandId.make(`cmd-deferred-turn-${index}`),
causationEventId: null,
correlationId: CorrelationId.make(`cmd-deferred-turn-${index}`),
metadata: {},
payload: {
threadId,
messageId: message.id,
...(index === 1
? {
sourceProposedPlan: {
threadId: ThreadId.make("thread-deferred-plan-source"),
planId: "plan-deferred",
},
}
: {}),
runtimeMode: "full-access",
createdAt: message.createdAt,
},
});
}
yield* eventStore.append({
type: "thread.activity-appended",
eventId: EventId.make("evt-deferred-compaction-complete"),
aggregateKind: "thread",
aggregateId: threadId,
occurredAt: "2026-02-26T16:00:02.000Z",
commandId: CommandId.make("cmd-deferred-compaction-complete"),
causationEventId: null,
correlationId: CorrelationId.make("cmd-deferred-compaction-complete"),
metadata: {},
payload: {
threadId,
activity: {
id: EventId.make("activity-deferred-compaction-complete"),
tone: "info",
kind: "context-compaction",
summary: "Context compacted",
payload: { requestId: compactMessageId },
turnId: null,
createdAt: "2026-02-26T16:00:02.000Z",
},
},
});
yield* eventStore.append({
type: "thread.session-set",
eventId: EventId.make("evt-deferred-session-ready"),
aggregateKind: "thread",
aggregateId: threadId,
occurredAt: "2026-02-26T16:00:03.000Z",
commandId: CommandId.make("server:provider-session-set:deferred"),
causationEventId: null,
correlationId: CorrelationId.make("server:provider-session-set:deferred"),
metadata: {},
payload: {
threadId,
session: {
threadId,
status: "ready",
providerName: "codex",
runtimeMode: "full-access",
activeTurnId: null,
lastError: null,
updatedAt: "2026-02-26T16:00:03.000Z",
},
},
});

yield* projectionPipeline.bootstrap;

const pendingRows = yield* sql<{
readonly messageId: string;
readonly sourceThreadId: string | null;
readonly sourcePlanId: string | null;
readonly requestedAt: string;
}>`
SELECT
pending_message_id AS "messageId",
source_proposed_plan_thread_id AS "sourceThreadId",
source_proposed_plan_id AS "sourcePlanId",
requested_at AS "requestedAt"
FROM projection_turns
WHERE thread_id = ${threadId}
AND turn_id IS NULL
AND state = 'pending'
`;
assert.deepEqual(pendingRows, [
{
messageId: queuedMessageId,
sourceThreadId: "thread-deferred-plan-source",
sourcePlanId: "plan-deferred",
requestedAt: queuedAt,
},
]);
}),
);
},
);

Expand Down
66 changes: 41 additions & 25 deletions apps/server/src/orchestration/Layers/ProjectionPipeline.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import {
ApprovalRequestId,
type ChatAttachment,
type MessageId,
type OrchestrationEvent,
type OrchestrationSessionStatus,
ThreadId,
Expand Down Expand Up @@ -503,6 +504,7 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti
const fileSystem = yield* FileSystem.FileSystem;
const path = yield* Path.Path;
const serverConfig = yield* ServerConfig;
const compactRequestIds = new Map<ThreadId, MessageId>();

const applyProjectsProjection: ProjectorDefinition["apply"] = Effect.fn(
"applyProjectsProjection",
Expand Down Expand Up @@ -1263,35 +1265,35 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti
)(function* (event, _attachmentSideEffects) {
switch (event.type) {
case "thread.created":
compactRequestIds.delete(event.payload.threadId);
yield* projectionTurnRepository.deleteByThreadId({
threadId: event.payload.threadId,
});
return;

case "thread.turn-start-requested": {
const pendingTurnStart = yield* projectionTurnRepository.getPendingTurnStartByThreadId({
threadId: event.payload.threadId,
});
if (Option.isSome(pendingTurnStart)) {
const pendingMessage = yield* projectionThreadMessageRepository.getByMessageId({
messageId: pendingTurnStart.value.messageId,
});
if (
Option.isSome(pendingMessage) &&
pendingMessage.value.role === "user" &&
(pendingMessage.value.attachments?.length ?? 0) === 0 &&
pendingMessage.value.text.trim().toLowerCase() === "/compact"
) {
return;
}
}
yield* projectionTurnRepository.replacePendingTurnStart({
const nextPendingTurnStart = {
threadId: event.payload.threadId,
messageId: event.payload.messageId,
sourceProposedPlanThreadId: event.payload.sourceProposedPlan?.threadId ?? null,
sourceProposedPlanId: event.payload.sourceProposedPlan?.planId ?? null,
requestedAt: event.payload.createdAt,
};
const requestedMessage = yield* projectionThreadMessageRepository.getByMessageId({
messageId: event.payload.messageId,
});
const isCompactRequest =
Option.isSome(requestedMessage) &&
requestedMessage.value.role === "user" &&
(requestedMessage.value.attachments?.length ?? 0) === 0 &&
requestedMessage.value.text.trim().toLowerCase() === "/compact";
if (isCompactRequest) {
if (compactRequestIds.has(event.payload.threadId)) {
return;
}
compactRequestIds.set(event.payload.threadId, event.payload.messageId);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Medium Layers/ProjectionPipeline.ts:1294

A failed replacePendingTurnStart leaves compactRequestIds populated even though the database transaction rolled back, so the next /compact for that thread returns at the duplicate check without creating a pending-start row. Move the map update until after replacePendingTurnStart succeeds.

🤖 Copy this AI Prompt to have your agent fix this:
In file @apps/server/src/orchestration/Layers/ProjectionPipeline.ts around line 1294:

A failed `replacePendingTurnStart` leaves `compactRequestIds` populated even though the database transaction rolled back, so the next `/compact` for that thread returns at the duplicate check without creating a pending-start row. Move the map update until after `replacePendingTurnStart` succeeds.

}
yield* projectionTurnRepository.replacePendingTurnStart(nextPendingTurnStart);
return;
}

Expand Down Expand Up @@ -1328,16 +1330,29 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti
case "thread.session-set": {
const turnId = event.payload.session.activeTurnId;
if (turnId === null || event.payload.session.status !== "running") {
if (
(event.payload.session.status === "ready" &&
event.commandId?.startsWith("server:provider-session-set:") === true) ||
const restoredReadySession =
event.payload.session.status === "ready" &&
event.commandId?.startsWith("server:provider-session-set:") === true;
const terminalSession =
event.payload.session.status === "error" ||
event.payload.session.status === "stopped" ||
event.payload.session.status === "interrupted"
) {
yield* projectionTurnRepository.deletePendingTurnStartByThreadId({
threadId: event.payload.threadId,
});
event.payload.session.status === "interrupted";
if (restoredReadySession || terminalSession) {
const pendingTurnStart =
yield* projectionTurnRepository.getPendingTurnStartByThreadId({
threadId: event.payload.threadId,
});
const compactRequestId = compactRequestIds.get(event.payload.threadId);
const pendingTurnBelongsToCompaction =
compactRequestId === undefined ||

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Medium Layers/ProjectionPipeline.ts:1347

After a restart, a restoring ready session deletes a queued message's pending turn-start row, so the queued turn loses its pending-message/proposed-plan metadata and is no longer projected as starting. Because compactRequestIds is in-memory, compactRequestId === undefined must not be treated as proof that the pending row belongs to /compact; only delete it when a known compaction request matches (while terminal sessions can still clear it unconditionally).

🤖 Copy this AI Prompt to have your agent fix this:
In file @apps/server/src/orchestration/Layers/ProjectionPipeline.ts around line 1347:

After a restart, a restoring `ready` session deletes a queued message's pending turn-start row, so the queued turn loses its pending-message/proposed-plan metadata and is no longer projected as `starting`. Because `compactRequestIds` is in-memory, `compactRequestId === undefined` must not be treated as proof that the pending row belongs to `/compact`; only delete it when a known compaction request matches (while terminal sessions can still clear it unconditionally).

Option.isNone(pendingTurnStart) ||
pendingTurnStart.value.messageId === compactRequestId;
if (terminalSession || pendingTurnBelongsToCompaction) {
yield* projectionTurnRepository.deletePendingTurnStartByThreadId({
threadId: event.payload.threadId,
});
}
compactRequestIds.delete(event.payload.threadId);
}
// Leaving the "running" session status is the turn-end signal:
// settle still-running turns so their duration reflects the whole
Expand Down Expand Up @@ -1461,6 +1476,7 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti
yield* projectionTurnRepository.deletePendingTurnStartByThreadId({
threadId: event.payload.threadId,
});
compactRequestIds.delete(event.payload.threadId);
return;
}

Expand Down
Loading
Loading