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

it("deletes archived project threads before deleting the project", async () => {
const system = await createOrchestrationSystem();
const { engine } = system;
const createdAt = now();

await system.run(
engine.dispatch({
type: "project.create",
commandId: CommandId.make("cmd-project-delete-archived-create"),
projectId: asProjectId("project-delete-archived"),
title: "Archived Delete Project",
workspaceRoot: "/tmp/project-delete-archived",
defaultModelSelection: {
provider: "codex",
model: "gpt-5-codex",
},
createdAt,
}),
);
await system.run(
engine.dispatch({
type: "thread.create",
commandId: CommandId.make("cmd-thread-delete-archived-create"),
threadId: ThreadId.make("thread-delete-archived"),
projectId: asProjectId("project-delete-archived"),
title: "Archive me first",
modelSelection: {
provider: "codex",
model: "gpt-5-codex",
},
interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE,
runtimeMode: "full-access",
branch: null,
worktreePath: null,
createdAt,
}),
);
await system.run(
engine.dispatch({
type: "thread.archive",
commandId: CommandId.make("cmd-thread-delete-archived-archive"),
threadId: ThreadId.make("thread-delete-archived"),
}),
);
await system.run(
engine.dispatch({
type: "project.delete",
commandId: CommandId.make("cmd-project-delete-archived"),
projectId: asProjectId("project-delete-archived"),
}),
);

const events = await system.run(
Stream.runCollect(engine.readEvents(0)).pipe(
Effect.map((chunk): OrchestrationEvent[] => Array.from(chunk)),
),
);
expect(events.map((event) => event.type)).toEqual([
"project.created",
"thread.created",
"thread.archived",
"thread.deleted",
"project.deleted",
]);

const readModel = await system.run(engine.getReadModel());
expect(
readModel.projects.find((project) => project.id === "project-delete-archived")?.deletedAt,
).not.toBeNull();
expect(
readModel.threads.find((thread) => thread.id === "thread-delete-archived")?.deletedAt,
).not.toBeNull();

await system.dispose();
});

it("rejects project deletion when active threads still exist", async () => {
const system = await createOrchestrationSystem();
const { engine } = system;
const createdAt = now();

await system.run(
engine.dispatch({
type: "project.create",
commandId: CommandId.make("cmd-project-delete-active-create"),
projectId: asProjectId("project-delete-active"),
title: "Active Delete Project",
workspaceRoot: "/tmp/project-delete-active",
defaultModelSelection: {
provider: "codex",
model: "gpt-5-codex",
},
createdAt,
}),
);
await system.run(
engine.dispatch({
type: "thread.create",
commandId: CommandId.make("cmd-thread-delete-active-create"),
threadId: ThreadId.make("thread-delete-active"),
projectId: asProjectId("project-delete-active"),
title: "Still active",
modelSelection: {
provider: "codex",
model: "gpt-5-codex",
},
interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE,
runtimeMode: "full-access",
branch: null,
worktreePath: null,
createdAt,
}),
);

await expect(
system.run(
engine.dispatch({
type: "project.delete",
commandId: CommandId.make("cmd-project-delete-active"),
projectId: asProjectId("project-delete-active"),
}),
),
).rejects.toThrow("still has active threads");

await system.dispose();
});

it("replays append-only events from sequence", async () => {
const system = await createOrchestrationSystem();
const { engine } = system;
Expand Down
20 changes: 20 additions & 0 deletions apps/server/src/orchestration/commandInvariants.ts
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,26 @@ export function requireProjectAbsent(input: {
);
}

export function requireProjectDeletionArchivedThreads(input: {
readonly readModel: OrchestrationReadModel;
readonly command: OrchestrationCommand;
readonly projectId: ProjectId;
}): Effect.Effect<ReadonlyArray<OrchestrationThread>, OrchestrationCommandInvariantError> {
const projectThreads = listThreadsByProjectId(input.readModel, input.projectId).filter(
(thread) => thread.deletedAt === null,
);
const activeThreads = projectThreads.filter((thread) => thread.archivedAt === null);
if (activeThreads.length > 0) {
return Effect.fail(
invariantError(
input.command.type,
`Project '${input.projectId}' still has active threads and cannot be deleted.`,
),
);
}
return Effect.succeed(projectThreads);
}

export function requireThread(input: {
readonly readModel: OrchestrationReadModel;
readonly command: OrchestrationCommand;
Expand Down
52 changes: 39 additions & 13 deletions apps/server/src/orchestration/decider.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import { Effect } from "effect";
import { OrchestrationCommandInvariantError } from "./Errors.ts";
import {
requireProject,
requireProjectDeletionArchivedThreads,
requireProjectAbsent,
requireThread,
requireThreadArchived,
Expand Down Expand Up @@ -119,20 +120,45 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand"
command,
projectId: command.projectId,
});
const occurredAt = nowIso();
return {
...withEventBase({
aggregateKind: "project",
aggregateId: command.projectId,
occurredAt,
commandId: command.commandId,
const archivedThreads = yield* requireProjectDeletionArchivedThreads({
readModel,
command,
projectId: command.projectId,
});
return [
...archivedThreads.map((thread) => {
const occurredAt = nowIso();
const eventBase = withEventBase({
aggregateKind: "thread",
aggregateId: thread.id,
occurredAt,
commandId: command.commandId,
});
return Object.assign({}, eventBase, {
type: "thread.deleted" as const,
payload: {
threadId: thread.id,
deletedAt: occurredAt,
},
});
}),
type: "project.deleted",
payload: {
projectId: command.projectId,
deletedAt: occurredAt,
},
};
(() => {
const occurredAt = nowIso();
return {
...withEventBase({
aggregateKind: "project",
aggregateId: command.projectId,
occurredAt,
commandId: command.commandId,
}),
type: "project.deleted" as const,
payload: {
projectId: command.projectId,
deletedAt: occurredAt,
},
};
})(),
];
}

case "thread.create": {
Expand Down
Loading
Loading