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
239 changes: 239 additions & 0 deletions apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2127,3 +2127,242 @@ describe("OrchestrationEngine", () => {
await system.dispose();
});
});

describe("temporary side conversations", () => {
it("freezes streaming context, persists the source link, and keeps or discards without altering the source", async () => {
const directory = await NodeFSP.mkdtemp(NodePath.join(NodeOS.tmpdir(), "t3-side-"));
const databasePath = NodePath.join(directory, "state.sqlite");
let system = await createOrchestrationSystem(databasePath);
const projectId = ProjectId.make("side-project");
const sourceId = ThreadId.make("side-source");
const sideId = ThreadId.make("side-temporary");
const dispatch = (command: OrchestrationCommand) => system.run(system.engine.dispatch(command));
try {
await dispatch({
type: "project.create",
commandId: CommandId.make("sp"),
projectId,
title: "Side test",
workspaceRoot: directory,
createdAt: now(),
});
await dispatch({
type: "thread.create",
commandId: CommandId.make("st"),
threadId: sourceId,
projectId,
title: "Running main task",
modelSelection: {
instanceId: ProviderInstanceId.make("codex"),
model: "gpt-6-astra",
options: [{ id: "reasoningEffort", value: "medium" }],
},
runtimeMode: "full-access",
interactionMode: "default",
branch: null,
worktreePath: null,
createdAt: now(),
});
await dispatch({
type: "thread.message.assistant.delta",
commandId: CommandId.make("sd1"),
threadId: sourceId,
messageId: MessageId.make("partial"),
delta: "So far: cobalt",
createdAt: now(),
});
await dispatch({
type: "thread.session.set",
commandId: CommandId.make("side-main-running"),
threadId: sourceId,
session: {
threadId: sourceId,
status: "running",
providerName: "codex",
runtimeMode: "full-access",
activeTurnId: TurnId.make("main-active-turn"),
lastError: null,
updatedAt: now(),
},
createdAt: now(),
});
const before = Option.getOrThrow(await system.readThread(sourceId));
const command: OrchestrationCommand = {
type: "thread.side.create",
commandId: CommandId.make("ss"),
threadId: sideId,
sourceThreadId: sourceId,
createdAt: now(),
};
await dispatch(command);
await dispatch(command);
const side = Option.getOrThrow(await system.readThread(sideId));
expect(side.sideChatOf).toBe(sourceId);
expect(side.messages).toHaveLength(1);
expect(side.messages[0]).toMatchObject({
text: "So far: cobalt\n[This answer was still in progress when the side chat opened.]",
streaming: false,
turnId: null,
});
expect(side.session).toBeNull();
expect(side.latestTurn).toBeNull();
expect(side.modelSelection).toEqual(before.modelSelection);
expect(Option.getOrThrow(await system.readThread(sourceId))).toEqual(before);
await dispatch({
type: "thread.message.assistant.delta",
commandId: CommandId.make("sd2"),
threadId: sourceId,
messageId: MessageId.make("partial"),
delta: ". Later: amber",
createdAt: now(),
});
expect(Option.getOrThrow(await system.readThread(sideId)).messages).toEqual(side.messages);
const laterSideId = ThreadId.make("later-side-snapshot");
await dispatch({
...command,
commandId: CommandId.make("later-side-create"),
threadId: laterSideId,
});
expect(Option.getOrThrow(await system.readThread(laterSideId)).messages[0]?.text).toContain(
"So far: cobalt. Later: amber",
);
expect(Option.getOrThrow(await system.readThread(sideId)).messages[0]?.text).not.toContain(
"Later: amber",
);
await system.dispose();
system = await createOrchestrationSystem(databasePath);
expect(Option.getOrThrow(await system.readThread(sideId)).sideChatOf).toBe(sourceId);
await dispatch({
type: "thread.meta.update",
commandId: CommandId.make("sk"),
threadId: sideId,
sideChatOf: null,
title: "Kept investigation",
});
const keptTitleState = {
source: "manual",
version: CommandId.make("sk"),
needsRefinement: false,
};
expect(Option.getOrThrow(await system.readThread(sideId))).toMatchObject({
sideChatOf: null,
title: "Kept investigation",
titleState: keptTitleState,
});
// Keeping and renaming must survive SQL projection reload together.
await system.dispose();
system = await createOrchestrationSystem(databasePath);
expect(Option.getOrThrow(await system.readThread(sideId))).toMatchObject({
sideChatOf: null,
title: "Kept investigation",
titleState: keptTitleState,
});
expect(Option.getOrThrow(await system.readThread(laterSideId)).sideChatOf).toBe(sourceId);
const continuedSource = Option.getOrThrow(await system.readThread(sourceId));
await dispatch({ type: "thread.delete", commandId: CommandId.make("sx"), threadId: sideId });
expect(Option.isNone(await system.readThread(sideId))).toBe(true);
expect(Option.getOrThrow(await system.readThread(sourceId))).toEqual(continuedSource);
} finally {
await system.dispose();
await NodeFSP.rm(directory, { recursive: true, force: true });
}
});
});

it("persists inline context through side snapshots and projection restart", async () => {
const directory = await NodeFSP.mkdtemp(NodePath.join(NodeOS.tmpdir(), "t3-context-branch-"));
const databasePath = NodePath.join(directory, "state.sqlite");
let system = await createOrchestrationSystem(databasePath);
const projectId = ProjectId.make("context-project");
const threadId = ThreadId.make("context-source");
const targetId = ThreadId.make("context-branch");
const createdAt = "2026-01-01T00:00:00.000Z";
const context = {
version: 1,
records: [
{
version: 1,
contextId: "terminal-retained" as never,
kind: "terminal",
label: "Terminal 1",
terminalId: "default",
terminalLabel: "Terminal 1",
lineStart: 1,
lineEnd: 1,
text: "retained inline payload",
},
],
} as const;
const dispatch = (command: OrchestrationCommand) => system.run(system.engine.dispatch(command));
try {
await dispatch({
type: "project.create",
commandId: CommandId.make("cp"),
projectId,
title: "Context",
workspaceRoot: directory,
createdAt,
});
await dispatch({
type: "thread.create",
commandId: CommandId.make("ct"),
threadId,
projectId,
title: "Context",
modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "gpt-5" },
runtimeMode: "full-access",
interactionMode: "default",
branch: null,
worktreePath: null,
createdAt,
});
await dispatch({
type: "thread.turn.start",
commandId: CommandId.make("cu"),
threadId,
message: {
messageId: MessageId.make("user"),
role: "user",
text: "Read [terminal](t3-context://v1/terminal/terminal-retained)",
attachments: [],
context,
},
runtimeMode: "full-access",
interactionMode: "default",
createdAt,
});
await dispatch({
type: "thread.message.assistant.delta",
commandId: CommandId.make("ca"),
threadId,
messageId: MessageId.make("answer"),
delta: "Read it",
createdAt: "2026-01-01T00:00:01.000Z",
});
await dispatch({
type: "thread.message.assistant.complete",
commandId: CommandId.make("cc"),
threadId,
messageId: MessageId.make("answer"),
createdAt: "2026-01-01T00:00:01.000Z",
});
await dispatch({
type: "thread.side.create",
commandId: CommandId.make("cb"),
threadId: targetId,
sourceThreadId: threadId,
createdAt: "2026-01-01T00:00:01.000Z",
});
expect(Option.getOrThrow(await system.readThread(targetId)).messages[0]?.context).toEqual(
context,
);
await system.dispose();
system = await createOrchestrationSystem(databasePath);
expect(Option.getOrThrow(await system.readThread(targetId)).messages[0]?.context).toEqual(
context,
);
} finally {
await system.dispose();
await NodeFSP.rm(directory, { recursive: true, force: true });
}
});
31 changes: 31 additions & 0 deletions apps/server/src/orchestration/Layers/OrchestrationEngine.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
import * as FileSystem from "effect/FileSystem";
import { ServerConfig } from "../../config.ts";
import type {
OrchestrationClientOrigin,
OrchestrationEvent,
Expand Down Expand Up @@ -40,6 +42,7 @@ import {
type OrchestrationDispatchError,
type OrchestrationProjectorDecodeError,
} from "../Errors.ts";
import { snapshotSideChat } from "../SideChatSnapshot.ts";
import { decideOrchestrationCommand } from "../decider.ts";
import { createEmptyReadModel, projectEvent } from "../projector.ts";
import { OrchestrationProjectionPipeline } from "../Services/ProjectionPipeline.ts";
Expand Down Expand Up @@ -83,6 +86,8 @@ function commandToAggregateRef(command: OrchestrationCommand): {

const makeOrchestrationEngine = Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;
const sideChatFileSystem = yield* FileSystem.FileSystem;
const sideChatServerConfig = yield* Effect.serviceOption(ServerConfig);
const eventStore = yield* OrchestrationEventStore;
const commandReceiptRepository = yield* OrchestrationCommandReceiptRepository;
const projectionPipeline = yield* OrchestrationProjectionPipeline;
Expand Down Expand Up @@ -242,9 +247,35 @@ const makeOrchestrationEngine = Effect.gen(function* () {
envelope.command.type === "thread.user-input.dismiss"
? yield* projectionSnapshotQuery.getUserInputActivity(envelope.command)
: Option.none();
const sideChatSource =
envelope.command.type === "thread.side.create"
? yield* Effect.gen(function* () {
const command = envelope.command;
if (command.type !== "thread.side.create") return undefined;
const source = yield* projectionSnapshotQuery.getThreadDetailById(
command.sourceThreadId,
{ activityKinds: [] },
);
if (Option.isNone(source))
return yield* new OrchestrationCommandInvariantError({
commandType: command.type,
detail: "The source conversation no longer exists.",
});
if (Option.isNone(sideChatServerConfig))
return yield* new OrchestrationCommandInvariantError({
commandType: command.type,
detail: "Side chats need the server attachment store.",
});
return yield* snapshotSideChat(source.value, command.threadId).pipe(
Effect.provideService(FileSystem.FileSystem, sideChatFileSystem),
Effect.provideService(ServerConfig, sideChatServerConfig.value),
);
})
: undefined;
const eventBase = yield* decideOrchestrationCommand({
command: envelope.command,
readModel: commandReadModel,
...(sideChatSource !== undefined ? { sideChatSource } : {}),
...(Option.isSome(userInputActivity)
? { userInputActivity: userInputActivity.value }
: {}),
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/orchestration/Layers/ProjectionPipeline.ts
Original file line number Diff line number Diff line change
Expand Up @@ -613,6 +613,7 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti
yield* projectionThreadRepository.upsert({
threadId: event.payload.threadId,
projectId: event.payload.projectId,
sideChatOf: event.payload.sideChatOf ?? null,
title: event.payload.title,
modelSelection: event.payload.modelSelection,
runtimeMode: event.payload.runtimeMode,
Expand Down Expand Up @@ -806,6 +807,7 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti
}
yield* projectionThreadRepository.upsert({
...existingRow.value,
...(event.payload.sideChatOf === null ? { sideChatOf: null } : {}),
...(event.payload.title !== undefined ? { title: event.payload.title } : {}),
...(event.payload.activeOrderKey !== undefined
? { activeOrderKey: event.payload.activeOrderKey }
Expand Down
Loading
Loading