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
6 changes: 6 additions & 0 deletions .github/upstream-candidates.json
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,12 @@
"sourceSha": "f7eaa00b99e67a9c1caf09e2fe6ff1f536f66b91",
"status": "active",
"purpose": "Show answered provider questions in the web thread timeline"
},
{
"upstreamPr": 4245,
"sourceSha": "6be48eb238cf310a60cc9571240601e774c7d381",
"status": "active",
"purpose": "Queue follow-up messages server-side during active turns with explicit steer controls"
}
]
}
2 changes: 2 additions & 0 deletions apps/mobile/src/lib/threadActivity.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,8 @@ function makeThread(
archivedAt: null,
deletedAt: null,
messages: [],
queuedMessages: [],
pendingTurnStart: null,
proposedPlans: [],
activities: [],
checkpoints: [],
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -150,6 +150,8 @@ describe("OrchestrationEngine", () => {
settledAt: null,
deletedAt: null,
messages: [],
queuedMessages: [],
pendingTurnStart: null,
proposedPlans: [],
activities: [],
checkpoints: [],
Expand All @@ -162,6 +164,8 @@ describe("OrchestrationEngine", () => {
threads: projectionSnapshot.threads.map((thread) => ({
...thread,
messages: [],
queuedMessages: [],
pendingTurnStart: null,
proposedPlans: [],
activities: [],
checkpoints: [],
Expand Down
109 changes: 98 additions & 11 deletions apps/server/src/orchestration/Layers/ProjectionPipeline.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ import {
type ProjectionThreadProposedPlan,
ProjectionThreadProposedPlanRepository,
} from "../../persistence/Services/ProjectionThreadProposedPlans.ts";
import { ProjectionQueuedMessageRepository } from "../../persistence/Services/ProjectionQueuedMessages.ts";
import { ProjectionThreadSessionRepository } from "../../persistence/Services/ProjectionThreadSessions.ts";
import {
type ProjectionTurn,
Expand All @@ -40,6 +41,7 @@ import { ProjectionStateRepositoryLive } from "../../persistence/Layers/Projecti
import { ProjectionThreadActivityRepositoryLive } from "../../persistence/Layers/ProjectionThreadActivities.ts";
import { ProjectionThreadMessageRepositoryLive } from "../../persistence/Layers/ProjectionThreadMessages.ts";
import { ProjectionThreadProposedPlanRepositoryLive } from "../../persistence/Layers/ProjectionThreadProposedPlans.ts";
import { ProjectionQueuedMessageRepositoryLive } from "../../persistence/Layers/ProjectionQueuedMessages.ts";
import { ProjectionThreadSessionRepositoryLive } from "../../persistence/Layers/ProjectionThreadSessions.ts";
import { ProjectionTurnRepositoryLive } from "../../persistence/Layers/ProjectionTurns.ts";
import { ProjectionThreadRepositoryLive } from "../../persistence/Layers/ProjectionThreads.ts";
Expand All @@ -59,6 +61,7 @@ export const ORCHESTRATION_PROJECTOR_NAMES = {
projects: "projection.projects",
threads: "projection.threads",
threadMessages: "projection.thread-messages",
queuedMessages: "projection.queued-messages",
threadProposedPlans: "projection.thread-proposed-plans",
threadActivities: "projection.thread-activities",
threadSessions: "projection.thread-sessions",
Expand Down Expand Up @@ -329,7 +332,7 @@ function retainProjectionProposedPlansAfterRevert(

function collectThreadAttachmentRelativePaths(
threadId: string,
messages: ReadonlyArray<ProjectionThreadMessage>,
messages: ReadonlyArray<Pick<ProjectionThreadMessage, "attachments">>,
): Set<string> {
const threadSegment = toSafeThreadAttachmentSegment(threadId);
if (!threadSegment) {
Expand Down Expand Up @@ -476,6 +479,7 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti
const projectionThreadRepository = yield* ProjectionThreadRepository;
const projectionThreadMessageRepository = yield* ProjectionThreadMessageRepository;
const projectionThreadProposedPlanRepository = yield* ProjectionThreadProposedPlanRepository;
const projectionQueuedMessageRepository = yield* ProjectionQueuedMessageRepository;
const projectionThreadActivityRepository = yield* ProjectionThreadActivityRepository;
const projectionThreadSessionRepository = yield* ProjectionThreadSessionRepository;
const projectionTurnRepository = yield* ProjectionTurnRepository;
Expand Down Expand Up @@ -782,6 +786,8 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti
}

case "thread.message-sent":
case "thread.message-queued":
case "thread.queued-message-removed":
case "thread.proposed-plan-upserted":
case "thread.activity-appended":
case "thread.approval-response-requested":
Expand Down Expand Up @@ -942,9 +948,17 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti
yield* Effect.forEach(keptRows, projectionThreadMessageRepository.upsert, {
concurrency: 1,
}).pipe(Effect.asVoid);
// Queued messages survive a revert (they are not part of the
// timeline), so their attachment files must survive pruning too.
const queuedRows = yield* projectionQueuedMessageRepository.listByThreadId({
threadId: event.payload.threadId,
});
attachmentSideEffects.prunedThreadRelativePaths.set(
event.payload.threadId,
collectThreadAttachmentRelativePaths(event.payload.threadId, keptRows),
collectThreadAttachmentRelativePaths(event.payload.threadId, [
...keptRows,
...queuedRows,
]),
);
return;
}
Expand All @@ -954,6 +968,67 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti
}
});

const applyQueuedMessagesProjection: ProjectorDefinition["apply"] = Effect.fn(
"applyQueuedMessagesProjection",
)(function* (event, attachmentSideEffects) {
switch (event.type) {
case "thread.message-queued":
yield* projectionQueuedMessageRepository.upsert({
messageId: event.payload.messageId,
threadId: event.payload.threadId,
text: event.payload.text,
attachments: event.payload.attachments,
modelSelection: event.payload.modelSelection ?? null,
sourceProposedPlanThreadId: event.payload.sourceProposedPlan?.threadId ?? null,
sourceProposedPlanId: event.payload.sourceProposedPlan?.planId ?? null,
queuedAt: event.payload.queuedAt,
});
return;

case "thread.queued-message-removed": {
// A user removal orphans the removed message's attachment files —
// prune to what the timeline and remaining queue still reference.
// Dispatch removals keep everything: the same attachments re-enter
// the timeline via the paired thread.message-sent.
const removedQueuedMessage =
event.payload.reason === "user"
? (yield* projectionQueuedMessageRepository.listByThreadId({
threadId: event.payload.threadId,
})).find((entry) => entry.messageId === event.payload.messageId)
: undefined;
yield* projectionQueuedMessageRepository.deleteByMessageId({
threadId: event.payload.threadId,
messageId: event.payload.messageId,
});
if (removedQueuedMessage && (removedQueuedMessage.attachments?.length ?? 0) > 0) {
const retainedMessageRows = yield* projectionThreadMessageRepository.listByThreadId({
threadId: event.payload.threadId,
});
const retainedQueuedRows = yield* projectionQueuedMessageRepository.listByThreadId({
threadId: event.payload.threadId,
});
attachmentSideEffects.prunedThreadRelativePaths.set(
event.payload.threadId,
collectThreadAttachmentRelativePaths(event.payload.threadId, [
...retainedMessageRows,
...retainedQueuedRows,
]),
);
}
return;
}

case "thread.deleted":
yield* projectionQueuedMessageRepository.deleteByThreadId({
threadId: event.payload.threadId,
});
return;

default:
return;
}
});

const applyThreadProposedPlansProjection: ProjectorDefinition["apply"] = Effect.fn(
"applyThreadProposedPlansProjection",
)(function* (event, _attachmentSideEffects) {
Expand Down Expand Up @@ -1093,19 +1168,22 @@ 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 === "error" ||
event.payload.session.status === "stopped" ||
event.payload.session.status === "interrupted"
) {
yield* projectionTurnRepository.deletePendingTurnStartByThreadId({
threadId: event.payload.threadId,
});
}
// Leaving the "running" session status is the turn-end signal:
// settle still-running turns so their duration reflects the whole
// turn rather than the last assistant message.
const settledTurnState = settledTurnStateForSessionStatus(event.payload.session.status);
// Any settled status abandons an unadopted pending turn start —
// including "ready": a mid-turn steer re-arms the pending row
// without a fresh adoption, and a stale row would block queue
// drains (and re-arm the read model's pendingTurnStart flag on
// restart hydration). Ready-with-genuinely-pending starts never
// reach this projection: ingestion maps that shape to
// "starting" before dispatching the session set.
if (settledTurnState !== null) {
yield* projectionTurnRepository.deletePendingTurnStartByThreadId({
threadId: event.payload.threadId,
});
}
if (settledTurnState === null) {
return;
}
Expand Down Expand Up @@ -1541,6 +1619,14 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti
name: ORCHESTRATION_PROJECTOR_NAMES.projects,
apply: applyProjectsProjection,
},
// queuedMessages must bootstrap before threadMessages: the revert
// handler in threadMessages reads the queued-message projection to
// retain queued attachments, so on replay that table has to be
// populated first or the prune deletes files still referenced.
{
name: ORCHESTRATION_PROJECTOR_NAMES.queuedMessages,
apply: applyQueuedMessagesProjection,
},
{
name: ORCHESTRATION_PROJECTOR_NAMES.threadMessages,
apply: applyThreadMessagesProjection,
Expand Down Expand Up @@ -1671,6 +1757,7 @@ export const OrchestrationProjectionPipelineLive = Layer.effect(
Layer.provideMerge(ProjectionThreadRepositoryLive),
Layer.provideMerge(ProjectionThreadMessageRepositoryLive),
Layer.provideMerge(ProjectionThreadProposedPlanRepositoryLive),
Layer.provideMerge(ProjectionQueuedMessageRepositoryLive),
Layer.provideMerge(ProjectionThreadActivityRepositoryLive),
Layer.provideMerge(ProjectionThreadSessionRepositoryLive),
Layer.provideMerge(ProjectionTurnRepositoryLive),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -325,6 +325,8 @@ projectionSnapshotLayer("ProjectionSnapshotQuery", (it) => {
updatedAt: "2026-02-24T00:00:05.000Z",
},
],
queuedMessages: [],
pendingTurnStart: null,
proposedPlans: [
{
id: "plan-1",
Expand Down
Loading