Skip to content
Open
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
190 changes: 189 additions & 1 deletion apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ import {
type OrchestrationEventStoreShape,
} from "../../persistence/Services/OrchestrationEventStore.ts";
import * as RepositoryIdentityResolver from "../../project/RepositoryIdentityResolver.ts";
import { decideOrchestrationCommand } from "../decider.ts";
import { OrchestrationEngineLive } from "./OrchestrationEngine.ts";
import { OrchestrationProjectionPipelineLive } from "./ProjectionPipeline.ts";
import { OrchestrationProjectionSnapshotQueryLive } from "./ProjectionSnapshotQuery.ts";
Expand All @@ -54,8 +55,9 @@ async function createOrchestrationSystem() {
Layer.provide(OrchestrationProjectionPipelineLive),
),
OrchestrationProjectionSnapshotQueryLive,
OrchestrationProjectionPipelineLive,
).pipe(
Layer.provide(OrchestrationEventStoreLive),
Layer.provideMerge(OrchestrationEventStoreLive),
Layer.provide(OrchestrationCommandReceiptRepositoryLive),
Layer.provide(RepositoryIdentityResolver.layer),
Layer.provide(SqlitePersistenceMemory),
Expand All @@ -65,8 +67,15 @@ async function createOrchestrationSystem() {
const runtime = ManagedRuntime.make(orchestrationLayer);
const engine = await runtime.runPromise(Effect.service(OrchestrationEngineService));
const snapshotQuery = await runtime.runPromise(Effect.service(ProjectionSnapshotQuery));
const eventStore = await runtime.runPromise(Effect.service(OrchestrationEventStore));
const projectionPipeline = await runtime.runPromise(
Effect.service(OrchestrationProjectionPipeline),
);
return {
engine,
eventStore,
projectionPipeline,
getCommandReadModel: () => runtime.runPromise(snapshotQuery.getCommandReadModel()),
readModel: () => runtime.runPromise(snapshotQuery.getSnapshot()),
run: <A, E>(effect: Effect.Effect<A, E>) => runtime.runPromise(effect),
dispose: () => runtime.dispose(),
Expand Down Expand Up @@ -420,6 +429,185 @@ describe("OrchestrationEngine", () => {
await system.dispose();
});

it("dispatches commands against threads appended by an external writer", async () => {
const system = await createOrchestrationSystem();
const { engine, eventStore, projectionPipeline } = system;
const createdAt = now();
const projectId = asProjectId("project-external-writer");
const threadAId = ThreadId.make("thread-engine");
const threadBId = ThreadId.make("thread-external");

await system.run(
engine.dispatch({
type: "project.create",
commandId: CommandId.make("cmd-external-project-create"),
projectId,
title: "External Writer Project",
workspaceRoot: "/tmp/project-external-writer",
defaultModelSelection: {
instanceId: ProviderInstanceId.make("codex"),
model: "gpt-5-codex",
},
createdAt,
}),
);
await system.run(
engine.dispatch({
type: "thread.create",
commandId: CommandId.make("cmd-external-thread-a-create"),
threadId: threadAId,
projectId,
title: "engine thread",
modelSelection: {
instanceId: ProviderInstanceId.make("codex"),
model: "gpt-5-codex",
},
interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE,
runtimeMode: "approval-required",
branch: null,
worktreePath: null,
createdAt,
}),
);

const externalWriterReadModel = await system.getCommandReadModel();
await system.run(
Effect.gen(function* () {
const eventBase = yield* decideOrchestrationCommand({
command: {
type: "thread.create",
commandId: CommandId.make("cmd-external-thread-b-create"),
threadId: threadBId,
projectId,
title: "external thread",
modelSelection: {
instanceId: ProviderInstanceId.make("codex"),
model: "gpt-5-codex",
},
interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE,
runtimeMode: "approval-required",
branch: null,
worktreePath: null,
createdAt,
},
readModel: externalWriterReadModel,
});
const eventBases = Array.isArray(eventBase) ? eventBase : [eventBase];
for (const event of eventBases) {
const savedEvent = yield* eventStore.append(event);
yield* projectionPipeline.projectEvent(savedEvent);
}
}).pipe(Effect.provide(NodeServices.layer)),
);

await system.run(
engine.dispatch({
type: "thread.meta.update",
commandId: CommandId.make("cmd-external-watermark-advance"),
threadId: threadAId,
title: "engine thread updated",
}),
);

await system.run(
engine.dispatch({
type: "thread.delete",
commandId: CommandId.make("cmd-external-delete"),
threadId: threadBId,
}),
);

const events = await system.run(
Stream.runCollect(engine.readEvents(0)).pipe(
Effect.map((chunk): OrchestrationEvent[] => Array.from(chunk)),
),
);
expect(
events.some((event) => event.type === "thread.deleted" && event.aggregateId === threadBId),
).toBe(true);

await system.dispose();
});

it("dispatches thread creation into a project appended by an external writer", async () => {
const system = await createOrchestrationSystem();
const { engine, eventStore, projectionPipeline } = system;
const createdAt = now();
const projectId = asProjectId("project-external-import");
const threadId = ThreadId.make("thread-external-import");

const externalWriterReadModel = await system.getCommandReadModel();
await system.run(
Effect.gen(function* () {
const eventBase = yield* decideOrchestrationCommand({
command: {
type: "project.create",
commandId: CommandId.make("cmd-external-project-create"),
projectId,
title: "External Import Project",
workspaceRoot: "/tmp/project-external-import",
defaultModelSelection: {
instanceId: ProviderInstanceId.make("codex"),
model: "gpt-5-codex",
},
createdAt,
},
readModel: externalWriterReadModel,
});
const eventBases = Array.isArray(eventBase) ? eventBase : [eventBase];
for (const event of eventBases) {
const savedEvent = yield* eventStore.append(event);
yield* projectionPipeline.projectEvent(savedEvent);
}
}).pipe(Effect.provide(NodeServices.layer)),
);

await system.run(
engine.dispatch({
type: "project.create",
commandId: CommandId.make("cmd-watermark-project-create"),
projectId: asProjectId("project-watermark"),
title: "Watermark Project",
workspaceRoot: "/tmp/project-watermark",
defaultModelSelection: {
instanceId: ProviderInstanceId.make("codex"),
model: "gpt-5-codex",
},
createdAt,
}),
);

await system.run(
engine.dispatch({
type: "thread.create",
commandId: CommandId.make("cmd-thread-create-in-external-project"),
threadId,
projectId,
title: "thread in imported project",
modelSelection: {
instanceId: ProviderInstanceId.make("codex"),
model: "gpt-5-codex",
},
interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE,
runtimeMode: "approval-required",
branch: null,
worktreePath: null,
createdAt,
}),
);

const events = await system.run(
Stream.runCollect(engine.readEvents(0)).pipe(
Effect.map((chunk): OrchestrationEvent[] => Array.from(chunk)),
),
);
expect(
events.some((event) => event.type === "thread.created" && event.aggregateId === threadId),
).toBe(true);

await system.dispose();
});

it("streams persisted domain events in order", async () => {
const system = await createOrchestrationSystem();
const { engine } = system;
Expand Down
53 changes: 49 additions & 4 deletions apps/server/src/orchestration/Layers/OrchestrationEngine.ts
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@ import {
type OrchestrationDispatchError,
type OrchestrationProjectorDecodeError,
} from "../Errors.ts";
import { findProjectById, findThreadById } from "../commandInvariants.ts";
import { decideOrchestrationCommand } from "../decider.ts";
import { createEmptyReadModel, projectEvent } from "../projector.ts";
import { OrchestrationProjectionPipeline } from "../Services/ProjectionPipeline.ts";
Expand All @@ -56,10 +57,11 @@ interface CommandEnvelope {
startedAtMs: number;
}

function commandToAggregateRef(command: OrchestrationCommand): {
readonly aggregateKind: "project" | "thread";
readonly aggregateId: ProjectId | ThreadId;
} {
function commandToAggregateRef(
command: OrchestrationCommand,
):
| { readonly aggregateKind: "project"; readonly aggregateId: ProjectId }
| { readonly aggregateKind: "thread"; readonly aggregateId: ThreadId } {
switch (command.type) {
case "project.create":
case "project.meta.update":
Expand Down Expand Up @@ -150,6 +152,49 @@ const makeOrchestrationEngine = Effect.gen(function* () {
});
}

// External writers (e.g. the `t3 import` CLI) append events and projections
// to the shared database from another process, so an aggregate missing from
// the in-memory model may still exist in SQL; refresh before rejecting.
const mustExistAggregate =
envelope.command.type === "project.create"
? undefined
: envelope.command.type === "thread.create"
? {
aggregateKind: "project" as const,
aggregateId: envelope.command.projectId,
}
: aggregateRef;
const aggregateMissingFromMemory =
mustExistAggregate !== undefined &&
(mustExistAggregate.aggregateKind === "thread"
? findThreadById(commandReadModel, mustExistAggregate.aggregateId) === undefined
: findProjectById(commandReadModel, mustExistAggregate.aggregateId) === undefined);
if (aggregateMissingFromMemory) {
const mustExistAggregateId = mustExistAggregate.aggregateId;
yield* Effect.gen(function* () {
const refreshed = yield* projectionSnapshotQuery.getCommandReadModel();
commandReadModel =
refreshed.snapshotSequence >= commandReadModel.snapshotSequence
? refreshed
: {
...refreshed,
snapshotSequence: commandReadModel.snapshotSequence,
updatedAt: commandReadModel.updatedAt,
};
}).pipe(
Effect.catch(() =>
Effect.logWarning(
"failed to refresh orchestration read model before command dispatch",
).pipe(
Effect.annotateLogs({
commandId: envelope.command.commandId,
aggregateId: mustExistAggregateId,
}),
),
),
);
}

const eventBase = yield* decideOrchestrationCommand({
command: envelope.command,
readModel: commandReadModel,
Expand Down
Loading