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

it("does not republish another server's turn when a local dispatch fails", async () => {
const directory = await NodeFSP.mkdtemp(NodePath.join(NodeOS.tmpdir(), "t3-shared-db-"));
const databasePath = NodePath.join(directory, "state.sqlite");
const serverA = await createOrchestrationSystem(databasePath);
const serverB = await createOrchestrationSystem(databasePath);
const threadId = ThreadId.make("thread-shared");
const createdAt = now();
const sentinelCommandId = CommandId.make("cmd-shared-rename");
try {
await serverA.run(
serverA.engine.dispatch({
type: "project.create",
commandId: CommandId.make("cmd-shared-project-create"),
projectId: asProjectId("project-shared"),
title: "Shared Project",
workspaceRoot: "/tmp/project-shared",
createdAt,
}),
);
await serverA.run(
serverA.engine.dispatch({
type: "thread.create",
commandId: CommandId.make("cmd-shared-thread-create"),
threadId,
projectId: asProjectId("project-shared"),
title: "shared",
modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "gpt-5-codex" },
interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE,
runtimeMode: "approval-required",
branch: null,
worktreePath: null,
createdAt,
}),
);
await serverA.run(
serverA.engine.dispatch({
type: "thread.turn.start",
commandId: CommandId.make("cmd-shared-turn-start"),
threadId,
message: {
messageId: asMessageId("msg-shared"),
role: "user",
text: "hello",
attachments: [],
},
interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE,
runtimeMode: "approval-required",
createdAt,
}),
);

const published = await serverB.run(
Effect.gen(function* () {
const events = yield* serverB.engine.subscribeDomainEvents;
// B's command model is still empty, so this fails and reconciles.
yield* serverB.engine
.dispatch({
type: "thread.meta.update",
commandId: CommandId.make("cmd-shared-stale-rename"),
threadId,
title: "stale",
})
.pipe(Effect.flip);
yield* serverB.engine.dispatch({
type: "thread.meta.update",
commandId: sentinelCommandId,
threadId,
title: "renamed on B",
});
return yield* Stream.runCollect(
Stream.takeUntil(events, (event) => event.commandId === sentinelCommandId),
);
}).pipe(Effect.scoped),
);

expect(Array.from(published).map((event) => event.type)).toEqual(["thread.meta-updated"]);
} finally {
await serverA.dispose();
await serverB.dispose();
await NodeFSP.rm(directory, { recursive: true, force: true });
}
});

it("fails command dispatch when command invariants are violated", async () => {
const system = await createOrchestrationSystem();
const { engine } = system;
Expand Down
9 changes: 8 additions & 1 deletion apps/server/src/orchestration/Layers/OrchestrationEngine.ts
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,10 @@ const makeOrchestrationEngine = Effect.gen(function* () {

const processEnvelope = (envelope: CommandEnvelope): Effect.Effect<void> => {
const dispatchStartSequence = commandReadModel.snapshotSequence;
// Events this dispatch appended. Reconcile republishes only these: a
// shared state directory can contain events another server already
// handled, and republishing them starts a second provider turn.
const appendedEventIds = new Set<OrchestrationEvent["eventId"]>();
let processingStartedAtMs = 0;
const aggregateRef = commandToAggregateRef(envelope.command);
const baseMetricAttributes = {
Expand All @@ -127,7 +131,9 @@ const makeOrchestrationEngine = Effect.gen(function* () {
commandReadModel = yield* projectEventsOntoReadModel(commandReadModel, persistedEvents);

for (const persistedEvent of persistedEvents) {
yield* PubSub.publish(eventPubSub, persistedEvent);
if (appendedEventIds.has(persistedEvent.eventId)) {
yield* PubSub.publish(eventPubSub, persistedEvent);
}
}
});

Expand Down Expand Up @@ -279,6 +285,7 @@ const makeOrchestrationEngine = Effect.gen(function* () {

for (const nextEvent of eventBases) {
const savedEvent = yield* eventStore.append(nextEvent);
appendedEventIds.add(savedEvent.eventId);
nextCommandReadModel = yield* projectEvent(nextCommandReadModel, savedEvent);
const cleanup = yield* projectionPipeline.projectEventDeferred(savedEvent);
attachmentCleanups.push(cleanup);
Expand Down
Loading