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
2 changes: 2 additions & 0 deletions apps/server/src/orchestration/Schemas.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import {
ThreadSeededWorkItemsUpsertedPayload as ContractsThreadSeededWorkItemsUpsertedPayloadSchema,
ThreadSeededWorkItemWritebackRequestedPayload as ContractsThreadSeededWorkItemWritebackRequestedPayloadSchema,
ThreadArchivedPayload as ContractsThreadArchivedPayloadSchema,
ThreadArchivedAndNewCreatedPayload as ContractsThreadArchivedAndNewCreatedPayloadSchema,
ThreadMetaUpdatedPayload as ContractsThreadMetaUpdatedPayloadSchema,
ThreadRuntimeModeSetPayload as ContractsThreadRuntimeModeSetPayloadSchema,
ThreadInteractionModeSetPayload as ContractsThreadInteractionModeSetPayloadSchema,
Expand Down Expand Up @@ -37,6 +38,7 @@ export const ThreadSeededWorkItemsUpsertedPayload =
export const ThreadSeededWorkItemWritebackRequestedPayload =
ContractsThreadSeededWorkItemWritebackRequestedPayloadSchema;
export const ThreadArchivedPayload = ContractsThreadArchivedPayloadSchema;
export const ThreadArchivedAndNewCreatedPayload = ContractsThreadArchivedAndNewCreatedPayloadSchema;
export const ThreadMetaUpdatedPayload = ContractsThreadMetaUpdatedPayloadSchema;
export const ThreadRuntimeModeSetPayload = ContractsThreadRuntimeModeSetPayloadSchema;
export const ThreadInteractionModeSetPayload = ContractsThreadInteractionModeSetPayloadSchema;
Expand Down
232 changes: 232 additions & 0 deletions apps/server/src/orchestration/decider.archiveAndNew.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,232 @@
import {
CommandId,
DEFAULT_PROVIDER_INTERACTION_MODE,
ProjectId,
ProviderInstanceId,
ThreadId,
type OrchestrationCommand,
type OrchestrationReadModel,
} from "@t3tools/contracts";
import { assert, describe, it } from "@effect/vitest";
import * as Crypto from "effect/Crypto";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import { NodeCrypto } from "@effect/platform-node";

import { decideOrchestrationCommand } from "./decider.ts";
import { createEmptyReadModel, projectEvent } from "./projector.ts";

const asCommandId = (value: string): CommandId => CommandId.make(value);
const asProjectId = (value: string): ProjectId => ProjectId.make(value);
const asThreadId = (value: string): ThreadId => ThreadId.make(value);

function seedThreadEffect(readModel: OrchestrationReadModel, overrides?: {
archivedAt?: string | null;
}): Effect.Effect<OrchestrationReadModel, never, Crypto.Crypto> {
return Effect.gen(function* () {
const now = "2026-01-01T00:00:00.000Z";
let model = readModel;

if (!model.projects.some((p) => p.id === asProjectId("project-1"))) {
model = yield* projectEvent(model, {
sequence: model.snapshotSequence + 1,
eventId: CommandId.make(crypto.randomUUID()) as unknown as never,
aggregateKind: "project",
aggregateId: asProjectId("project-1"),
type: "project.created",
occurredAt: now,
commandId: asCommandId("cmd-project-create"),
causationEventId: null,
correlationId: asCommandId("cmd-project-create"),
metadata: {},
payload: {
projectId: asProjectId("project-1"),
title: "Project 1",
workspaceRoot: "/tmp/project-1",
defaultModelSelection: null,
scripts: [],
createdAt: now,
updatedAt: now,
},
} as never);
}

model = yield* projectEvent(model, {
sequence: model.snapshotSequence + 1,
eventId: CommandId.make(crypto.randomUUID()) as unknown as never,
aggregateKind: "thread",
aggregateId: asThreadId("thread-1"),
type: "thread.created",
occurredAt: now,
commandId: asCommandId("cmd-thread-create"),
causationEventId: null,
correlationId: asCommandId("cmd-thread-create"),
metadata: {},
payload: {
threadId: asThreadId("thread-1"),
projectId: asProjectId("project-1"),
title: "Test Thread",
modelSelection: {
instanceId: ProviderInstanceId.make("codex"),
model: "gpt-5-codex",
},
interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE,
runtimeMode: "full-access",
branch: "feat/my-branch",
worktreePath: "/tmp/worktree-thread-1",
createdAt: overrides?.archivedAt ?? now,
updatedAt: now,
},
} as never);

if (overrides?.archivedAt) {
model = yield* projectEvent(model, {
sequence: model.snapshotSequence + 1,
eventId: CommandId.make(crypto.randomUUID()) as unknown as never,
aggregateKind: "thread",
aggregateId: asThreadId("thread-1"),
type: "thread.archived",
occurredAt: overrides.archivedAt,
commandId: asCommandId("cmd-archive"),
causationEventId: null,
correlationId: asCommandId("cmd-archive"),
metadata: {},
payload: {
threadId: asThreadId("thread-1"),
archivedAt: overrides.archivedAt,
updatedAt: overrides.archivedAt,
},
} as never);
}

return model;
});
}

describe("thread.archive-and-new decider", () => {
it.effect("emits thread.archived-and-new-created and thread.session-stop-requested", () =>
Effect.gen(function* () {
const readModel = yield* seedThreadEffect(
createEmptyReadModel("2026-01-01T00:00:00.000Z"),
);

const result = yield* decideOrchestrationCommand({
command: {
type: "thread.archive-and-new",
commandId: asCommandId("cmd-archive-new"),
threadId: asThreadId("thread-1"),
newThreadId: asThreadId("thread-new"),
createdAt: "2026-01-02T00:00:00.000Z",
} as Extract<OrchestrationCommand, { type: "thread.archive-and-new" }>,
readModel,
});

const events = Array.isArray(result) ? result : [result];
assert.strictEqual(events.length, 2);
assert.strictEqual(events[0]?.type, "thread.archived-and-new-created");
assert.strictEqual(events[1]?.type, "thread.session-stop-requested");

if (events[0]?.type === "thread.archived-and-new-created") {
assert.strictEqual(events[0].payload.archivedThreadId, "thread-1");
assert.strictEqual(events[0].payload.newThreadId, "thread-new");
assert.strictEqual(events[0].payload.createdAt, "2026-01-02T00:00:00.000Z");
}

if (events[1]?.type === "thread.session-stop-requested") {
assert.strictEqual(events[1].payload.threadId, "thread-1");
}
}).pipe(Effect.provide(NodeCrypto.layer)));

it.effect("rejects when thread does not exist", () =>
Effect.gen(function* () {
const readModel = createEmptyReadModel("2026-01-01T00:00:00.000Z");

const result = yield* Effect.exit(
decideOrchestrationCommand({
command: {
type: "thread.archive-and-new",
commandId: asCommandId("cmd-archive-new"),
threadId: asThreadId("thread-missing"),
newThreadId: asThreadId("thread-new"),
createdAt: "2026-01-02T00:00:00.000Z",
} as Extract<OrchestrationCommand, { type: "thread.archive-and-new" }>,
readModel,
}),
);

assert.strictEqual(result._tag, "Failure");
}).pipe(Effect.provide(NodeCrypto.layer)));

it.effect("rejects when thread is already archived", () =>
Effect.gen(function* () {
const readModel = yield* seedThreadEffect(
createEmptyReadModel("2026-01-01T00:00:00.000Z"),
{ archivedAt: "2026-01-01T00:30:00.000Z" },
);

const result = yield* Effect.exit(
decideOrchestrationCommand({
command: {
type: "thread.archive-and-new",
commandId: asCommandId("cmd-archive-new"),
threadId: asThreadId("thread-1"),
newThreadId: asThreadId("thread-new"),
createdAt: "2026-01-02T00:00:00.000Z",
} as Extract<OrchestrationCommand, { type: "thread.archive-and-new" }>,
readModel,
}),
);

assert.strictEqual(result._tag, "Failure");
}).pipe(Effect.provide(NodeCrypto.layer)));

it.effect("rejects when new thread ID already exists", () =>
Effect.gen(function* () {
const readModel = yield* seedThreadEffect(
createEmptyReadModel("2026-01-01T00:00:00.000Z"),
);

const modelWithCollision = yield* projectEvent(readModel, {
sequence: readModel.snapshotSequence + 1,
eventId: CommandId.make(crypto.randomUUID()) as unknown as never,
aggregateKind: "thread",
aggregateId: asThreadId("thread-new"),
type: "thread.created",
occurredAt: "2026-01-01T00:01:00.000Z",
commandId: asCommandId("cmd-collision-create"),
causationEventId: null,
correlationId: asCommandId("cmd-collision-create"),
metadata: {},
payload: {
threadId: asThreadId("thread-new"),
projectId: asProjectId("project-1"),
title: "Collision Thread",
modelSelection: {
instanceId: ProviderInstanceId.make("codex"),
model: "gpt-5-codex",
},
interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE,
runtimeMode: "full-access",
branch: null,
worktreePath: null,
createdAt: "2026-01-01T00:01:00.000Z",
updatedAt: "2026-01-01T00:01:00.000Z",
},
} as never);

const result = yield* Effect.exit(
decideOrchestrationCommand({
command: {
type: "thread.archive-and-new",
commandId: asCommandId("cmd-archive-new"),
threadId: asThreadId("thread-1"),
newThreadId: asThreadId("thread-new"),
createdAt: "2026-01-02T00:00:00.000Z",
} as Extract<OrchestrationCommand, { type: "thread.archive-and-new" }>,
readModel: modelWithCollision,
}),
);

assert.strictEqual(result._tag, "Failure");
}).pipe(Effect.provide(NodeCrypto.layer)));
});
43 changes: 43 additions & 0 deletions apps/server/src/orchestration/decider.ts
Original file line number Diff line number Diff line change
Expand Up @@ -740,6 +740,49 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand"
};
}

case "thread.archive-and-new": {
const thread = yield* requireThreadNotArchived({
readModel,
command,
threadId: command.threadId,
});
yield* requireThreadAbsent({
readModel,
command,
threadId: command.newThreadId,
});
const occurredAt = command.createdAt;
return [
{
...(yield* withEventBase({
aggregateKind: "thread",
aggregateId: command.threadId,
occurredAt,
commandId: command.commandId,
})),
type: "thread.archived-and-new-created",
payload: {
archivedThreadId: command.threadId,
newThreadId: command.newThreadId,
createdAt: occurredAt,
},
},
{
...(yield* withEventBase({
aggregateKind: "thread",
aggregateId: command.threadId,
occurredAt,
commandId: command.commandId,
})),
type: "thread.session-stop-requested",
payload: {
threadId: command.threadId,
createdAt: occurredAt,
},
},
];
}

case "thread.unarchive": {
yield* requireThreadArchived({
readModel,
Expand Down
Loading
Loading