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
50 changes: 50 additions & 0 deletions apps/server/src/orchestration/Layers/ProjectionPipeline.ts
Original file line number Diff line number Diff line change
Expand Up @@ -603,6 +603,8 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti
interactionMode: event.payload.interactionMode,
branch: event.payload.branch,
worktreePath: event.payload.worktreePath,
forkedFromThreadId: null,
forkedFromTurnId: null,
latestTurnId: null,
createdAt: event.payload.createdAt,
updatedAt: event.payload.updatedAt,
Expand All @@ -615,6 +617,33 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti
});
return;

case "thread.forked":
yield* projectionThreadRepository.upsert({
threadId: event.payload.threadId,
projectId: event.payload.projectId,
title: event.payload.title,
modelSelection: event.payload.modelSelection,
runtimeMode: event.payload.runtimeMode,
interactionMode: event.payload.interactionMode,
branch: event.payload.branch,
worktreePath: event.payload.worktreePath,
forkedFromThreadId: event.payload.forkedFrom.threadId,
forkedFromTurnId: event.payload.forkedFrom.turnId,
latestTurnId: null,
createdAt: event.payload.createdAt,
updatedAt: event.payload.updatedAt,
archivedAt: null,
latestUserMessageAt:
event.payload.inheritedMessages
.toReversed()
.find((message) => message.role === "user")?.createdAt ?? null,
pendingApprovalCount: 0,
pendingUserInputCount: 0,
hasActionableProposedPlan: 0,
deletedAt: null,
});
return;

case "thread.archived": {
const existingRow = yield* projectionThreadRepository.getById({
threadId: event.payload.threadId,
Expand Down Expand Up @@ -811,6 +840,27 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti
"applyThreadMessagesProjection",
)(function* (event, attachmentSideEffects) {
switch (event.type) {
case "thread.forked":

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.

🟠 High Layers/ProjectionPipeline.ts:843

The thread.forked message projection copies inherited attachments unchanged via [...message.attachments], so the forked thread's attachment IDs still contain the source thread's segment. When the source thread is deleted or reverted past those messages, runAttachmentSideEffects deletes files by that source segment, destroying attachment files that the fork still references. Consider materializing independent attachments for the fork (as thread.message-sent does via materializeAttachmentsForProjection) or making cleanup account for cross-thread references.

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

The `thread.forked` message projection copies inherited `attachments` unchanged via `[...message.attachments]`, so the forked thread's attachment IDs still contain the source thread's segment. When the source thread is deleted or reverted past those messages, `runAttachmentSideEffects` deletes files by that source segment, destroying attachment files that the fork still references. Consider materializing independent attachments for the fork (as `thread.message-sent` does via `materializeAttachmentsForProjection`) or making cleanup account for cross-thread references.

yield* Effect.forEach(
event.payload.inheritedMessages,
(message) =>
projectionThreadMessageRepository.upsert({
messageId: message.id,
threadId: event.payload.threadId,
turnId: message.turnId,
role: message.role,
text: message.text,
...(message.attachments !== undefined
? { attachments: [...message.attachments] }
: {}),
isStreaming: message.streaming,
createdAt: message.createdAt,
updatedAt: message.updatedAt,
}),
{ concurrency: 1, discard: true },
);
return;

case "thread.message-sent": {
const existingMessage = yield* projectionThreadMessageRepository.getByMessageId({
messageId: event.payload.messageId,
Expand Down
23 changes: 23 additions & 0 deletions apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts
Original file line number Diff line number Diff line change
Expand Up @@ -224,6 +224,15 @@ function mapSessionRow(
};
}

function mapForkedFrom(row: Schema.Schema.Type<typeof ProjectionThreadDbRowSchema>) {
return row.forkedFromThreadId == null
? null
: {
threadId: row.forkedFromThreadId,
turnId: row.forkedFromTurnId ?? null,
};
}

function mapProjectShellRow(
row: Schema.Schema.Type<typeof ProjectionProjectDbRowSchema>,
repositoryIdentity: OrchestrationProject["repositoryIdentity"],
Expand Down Expand Up @@ -330,6 +339,8 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () {
interaction_mode AS "interactionMode",
branch,
worktree_path AS "worktreePath",
forked_from_thread_id AS "forkedFromThreadId",
forked_from_turn_id AS "forkedFromTurnId",
latest_turn_id AS "latestTurnId",
created_at AS "createdAt",
updated_at AS "updatedAt",
Expand Down Expand Up @@ -358,6 +369,8 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () {
interaction_mode AS "interactionMode",
branch,
worktree_path AS "worktreePath",
forked_from_thread_id AS "forkedFromThreadId",
forked_from_turn_id AS "forkedFromTurnId",
latest_turn_id AS "latestTurnId",
created_at AS "createdAt",
updated_at AS "updatedAt",
Expand Down Expand Up @@ -388,6 +401,8 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () {
interaction_mode AS "interactionMode",
branch,
worktree_path AS "worktreePath",
forked_from_thread_id AS "forkedFromThreadId",
forked_from_turn_id AS "forkedFromTurnId",
latest_turn_id AS "latestTurnId",
created_at AS "createdAt",
updated_at AS "updatedAt",
Expand Down Expand Up @@ -750,6 +765,8 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () {
interaction_mode AS "interactionMode",
branch,
worktree_path AS "worktreePath",
forked_from_thread_id AS "forkedFromThreadId",
forked_from_turn_id AS "forkedFromTurnId",
latest_turn_id AS "latestTurnId",
created_at AS "createdAt",
updated_at AS "updatedAt",
Expand Down Expand Up @@ -1182,6 +1199,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () {
interactionMode: row.interactionMode,
branch: row.branch,
worktreePath: row.worktreePath,
forkedFrom: mapForkedFrom(row),
latestTurn: latestTurnByThread.get(row.threadId) ?? null,
createdAt: row.createdAt,
updatedAt: row.updatedAt,
Expand Down Expand Up @@ -1380,6 +1398,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () {
interactionMode: row.interactionMode,
branch: row.branch,
worktreePath: row.worktreePath,
forkedFrom: mapForkedFrom(row),
latestTurn: latestTurnByThread.get(row.threadId) ?? null,
createdAt: row.createdAt,
updatedAt: row.updatedAt,
Expand Down Expand Up @@ -1509,6 +1528,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () {
interactionMode: row.interactionMode,
branch: row.branch,
worktreePath: row.worktreePath,
forkedFrom: mapForkedFrom(row),
latestTurn: latestTurnByThread.get(row.threadId) ?? null,
createdAt: row.createdAt,
updatedAt: row.updatedAt,
Expand Down Expand Up @@ -1643,6 +1663,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () {
interactionMode: row.interactionMode,
branch: row.branch,
worktreePath: row.worktreePath,
forkedFrom: mapForkedFrom(row),
latestTurn: latestTurnByThread.get(row.threadId) ?? null,
createdAt: row.createdAt,
updatedAt: row.updatedAt,
Expand Down Expand Up @@ -1883,6 +1904,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () {
interactionMode: threadRow.value.interactionMode,
branch: threadRow.value.branch,
worktreePath: threadRow.value.worktreePath,
forkedFrom: mapForkedFrom(threadRow.value),
latestTurn: Option.isSome(latestTurnRow) ? mapLatestTurn(latestTurnRow.value) : null,
createdAt: threadRow.value.createdAt,
updatedAt: threadRow.value.updatedAt,
Expand Down Expand Up @@ -1977,6 +1999,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () {
interactionMode: threadRow.value.interactionMode,
branch: threadRow.value.branch,
worktreePath: threadRow.value.worktreePath,
forkedFrom: mapForkedFrom(threadRow.value),
latestTurn: Option.isSome(latestTurnRow) ? mapLatestTurn(latestTurnRow.value) : null,
createdAt: threadRow.value.createdAt,
updatedAt: threadRow.value.updatedAt,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@ import { OrchestrationEngineLive } from "./OrchestrationEngine.ts";
import { OrchestrationProjectionPipelineLive } from "./ProjectionPipeline.ts";
import { OrchestrationProjectionSnapshotQueryLive } from "./ProjectionSnapshotQuery.ts";
import {
findCompletedTurnIndex,
providerErrorLabel,
providerErrorLabelFromInstanceHint,
ProviderCommandReactorLive,
Expand All @@ -66,6 +67,24 @@ const asApprovalRequestId = (value: string): ApprovalRequestId => ApprovalReques
const asMessageId = (value: string): MessageId => MessageId.make(value);
const asTurnId = (value: string): TurnId => TurnId.make(value);

describe("findCompletedTurnIndex", () => {
it("maps a selected T3 turn to its provider response position", () => {
const first = asTurnId("turn-1");
const second = asTurnId("turn-2");
expect(
findCompletedTurnIndex(
[
{ role: "assistant", turnId: first, streaming: false },
{ role: "assistant", turnId: first, streaming: false },
{ role: "assistant", turnId: second, streaming: true },
{ role: "assistant", turnId: second, streaming: false },
],
second,
),
).toBe(1);
});
});

const deriveServerPathsSync = (baseDir: string, devUrl: URL | undefined) =>
Effect.runSync(deriveServerPaths(baseDir, devUrl).pipe(Effect.provide(NodeServices.layer)));

Expand Down
55 changes: 51 additions & 4 deletions apps/server/src/orchestration/Layers/ProviderCommandReactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -88,6 +88,32 @@ const HANDLED_TURN_START_KEY_TTL = Duration.minutes(30);
const DEFAULT_RUNTIME_MODE: RuntimeMode = "full-access";
const DEFAULT_THREAD_TITLE = "New thread";

export function findCompletedTurnIndex(
messages: ReadonlyArray<{
readonly role: string;
readonly turnId: TurnId | null;
readonly streaming: boolean;
}>,
sourceTurnId: TurnId,
): number | undefined {
const completedTurnIds: TurnId[] = [];
for (const message of messages) {
if (
message.role !== "assistant" ||
message.turnId === null ||
message.streaming ||
completedTurnIds.some((turnId) => turnId === message.turnId)
) {
continue;
}
completedTurnIds.push(message.turnId);
if (message.turnId === sourceTurnId) {
return completedTurnIds.length - 1;
}
}
return undefined;
}

export function providerErrorLabel(value: string | undefined): string {
const normalized = value?.trim();
return normalized && normalized.length > 0 ? normalized : "unknown";
Expand Down Expand Up @@ -471,19 +497,40 @@ const make = Effect.gen(function* () {
projects: project ? [project] : [],
});

const startProviderSession = (input?: {
const startProviderSession = Effect.fn("startProviderSession")(function* (input?: {
readonly resumeCursor?: unknown;
readonly provider?: ProviderDriverKind;
}) =>
providerService.startSession(threadId, {
readonly includeForkSource?: boolean;
}) {
const forkedFrom = input?.includeForkSource === true ? thread.forkedFrom : null;
const forkSource =
forkedFrom != null
? yield* Effect.gen(function* () {
const sourceThread = yield* resolveThread(forkedFrom.threadId);
const sourceTurnId = forkedFrom.turnId;
const sourceTurnIndex =
sourceThread != null && sourceTurnId !== null
? findCompletedTurnIndex(sourceThread.messages, sourceTurnId)
: undefined;
return {
threadId: forkedFrom.threadId,
...(sourceTurnId !== null ? { sourceTurnId } : {}),
...(sourceTurnIndex !== undefined ? { sourceTurnIndex } : {}),
};
})
: undefined;

return yield* providerService.startSession(threadId, {
threadId,
...(preferredProvider ? { provider: preferredProvider } : {}),
providerInstanceId: desiredInstanceId,
...(effectiveCwd ? { cwd: effectiveCwd } : {}),
modelSelection: desiredModelSelection,
...(input?.resumeCursor !== undefined ? { resumeCursor: input.resumeCursor } : {}),
...(forkSource !== undefined ? { forkFrom: forkSource } : {}),
runtimeMode: desiredRuntimeMode,
});
});

const bindSessionToThread = (session: ProviderSession) =>
Effect.gen(function* () {
Expand Down Expand Up @@ -578,7 +625,7 @@ const make = Effect.gen(function* () {
return restartedSession.threadId;
}

const startedSession = yield* startProviderSession(undefined);
const startedSession = yield* startProviderSession({ includeForkSource: true });
yield* bindSessionToThread(startedSession);
return startedSession.threadId;
});
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/orchestration/Schemas.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ import {
ProjectMetaUpdatedPayload as ContractsProjectMetaUpdatedPayloadSchema,
ProjectDeletedPayload as ContractsProjectDeletedPayloadSchema,
ThreadCreatedPayload as ContractsThreadCreatedPayloadSchema,
ThreadForkedPayload as ContractsThreadForkedPayloadSchema,
ThreadArchivedPayload as ContractsThreadArchivedPayloadSchema,
ThreadMetaUpdatedPayload as ContractsThreadMetaUpdatedPayloadSchema,
ThreadRuntimeModeSetPayload as ContractsThreadRuntimeModeSetPayloadSchema,
Expand All @@ -28,6 +29,7 @@ export const ProjectMetaUpdatedPayload = ContractsProjectMetaUpdatedPayloadSchem
export const ProjectDeletedPayload = ContractsProjectDeletedPayloadSchema;

export const ThreadCreatedPayload = ContractsThreadCreatedPayloadSchema;
export const ThreadForkedPayload = ContractsThreadForkedPayloadSchema;
export const ThreadArchivedPayload = ContractsThreadArchivedPayloadSchema;
export const ThreadMetaUpdatedPayload = ContractsThreadMetaUpdatedPayloadSchema;
export const ThreadRuntimeModeSetPayload = ContractsThreadRuntimeModeSetPayloadSchema;
Expand Down
Loading
Loading