Skip to content
Closed
72 changes: 68 additions & 4 deletions apps/server/src/orchestration-v2/Orchestrator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2370,13 +2370,14 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
}
if (
command.type === "thread.metadata.update" &&
command.expectedWorktreePath !== undefined &&
command.expectedWorktreePath !== thread.worktreePath
((command.expectedWorktreePath !== undefined &&
command.expectedWorktreePath !== thread.worktreePath) ||
(command.expectedBranch !== undefined && command.expectedBranch !== thread.branch))
) {
return yield* new OrchestratorDispatchError({
commandId: command.commandId,
commandType: command.type,
cause: `Thread ${command.threadId} worktree changed before the metadata update could be applied.`,
cause: `Thread ${command.threadId} workspace binding changed before the metadata update could be applied.`,
});
}
if (command.type === "thread.metadata.update" && command.expectedEmpty === true) {
Expand Down Expand Up @@ -5166,7 +5167,12 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
? {}
: { restartContinuationOfRunId: command.restartContinuationOfRunId }),
...(dispatchMode.type === "defer_start" && dispatchMode.workspaceStrategy !== undefined
? { workspacePreparation: dispatchMode.workspaceStrategy }
? {
workspacePreparation: dispatchMode.workspaceStrategy,
...(dispatchMode.workspaceStrategy.type === "worktree"
? { completedWorktreePath: null }
: {}),
}
: {}),
...wakeWorkStartedAt(projection.runs, command),
};
Expand Down Expand Up @@ -7564,6 +7570,52 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
});
}
const now = yield* DateTime.now;
const workspace = command.completedWorkspace;
if (
workspace !== undefined &&
(command.phase !== "setup" ||
state.run.workspacePreparation?.type !== "worktree" ||
projection.thread.worktreePath !== workspace.expectedWorktreePath ||
projection.thread.branch !== workspace.expectedBranch)
) {
return yield* new OrchestratorDispatchError({
commandId: command.commandId,
commandType: command.type,
cause: "The prepared run no longer owns the workspace binding.",
});
}
if (workspace !== undefined) {
yield* emit(
events,
command,
)({
type: "thread.metadata-updated",
threadId: command.threadId,
occurredAt: now,
payload: {
...projection.thread,
worktreePath: workspace.worktreePath,
branch: workspace.branch,
updatedAt: now,
},
});
}
if (workspace !== undefined || command.phase === "worktree") {
yield* emit(
events,
command,
)({
type: "run.updated",
threadId: command.threadId,
runId: state.run.id,
providerInstanceId: state.run.providerInstanceId,
occurredAt: now,
payload: {
...state.run,
completedWorktreePath: workspace?.worktreePath ?? null,
},
});
}
yield* emit(
events,
command,
Expand Down Expand Up @@ -7601,6 +7653,18 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
cause: `Run ${command.runId} is not awaiting workspace preparation.`,
});
}
if (
state.run.workspacePreparation?.type === "worktree" &&
state.run.completedWorktreePath !== undefined &&
(state.run.completedWorktreePath === null ||
state.run.completedWorktreePath !== projection.thread.worktreePath)
) {
return yield* new OrchestratorDispatchError({
commandId: command.commandId,
commandType: command.type,
cause: "The prepared run has no completed checkout at its current workspace binding.",
});
}
const now = yield* DateTime.now;
const resolvedRuntimePolicy = yield* runtimePolicy
.resolve({ thread: projection.thread, modelSelection: state.run.modelSelection })
Expand Down
Loading
Loading