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
88 changes: 81 additions & 7 deletions apps/server/src/orchestration-v2/Orchestrator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2205,6 +2205,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
interactionMode: command.interactionMode,
branch: command.branch,
worktreePath: command.worktreePath,
workspaceBindingId: command.commandId,
activeProviderThreadId: null,
lineage: {
parentThreadId: null,
Expand Down Expand Up @@ -2471,13 +2472,16 @@ 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) ||
(command.expectedWorkspaceBindingId !== undefined &&
command.expectedWorkspaceBindingId !== (thread.workspaceBindingId ?? null)))
) {
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 @@ -2937,6 +2941,9 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
: {}),
...(command.branch === undefined ? {} : { branch: command.branch }),
...(command.worktreePath === undefined ? {} : { worktreePath: command.worktreePath }),
...(command.branch === undefined && command.worktreePath === undefined
? {}
: { workspaceBindingId: command.commandId }),
...(command.linkedPullRequest === undefined
? {}
: {
Expand Down Expand Up @@ -5318,7 +5325,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 @@ -7742,7 +7754,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
});

const loadProjectionForCommand = <K extends ProjectionRecordField>(
command: OrchestrationV2Command,
command: OrchestrationV2ServerCommand,
fields: ReadonlyArray<K>,
filter?: ProjectionRecordFilter,
) =>
Expand All @@ -7756,7 +7768,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio

const preparedRunState = (
command: Extract<
OrchestrationV2Command,
OrchestrationV2ServerCommand,
{
readonly type:
| "prepared-run.release"
Expand Down Expand Up @@ -7797,7 +7809,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
};

const dispatchPreparedRunProgress = (
command: Extract<OrchestrationV2Command, { readonly type: "prepared-run.progress" }>,
command: Extract<OrchestrationV2ServerCommand, { readonly type: "prepared-run.progress" }>,
events: Ref.Ref<Array<OrchestrationV2DomainEvent>>,
) =>
Effect.gen(function* () {
Expand All @@ -7815,6 +7827,56 @@ 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 ||
(workspace.expectedWorkspaceBindingId !== undefined &&
workspace.expectedWorkspaceBindingId !==
(projection.thread.workspaceBindingId ?? null)))
) {
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,
workspaceBindingId: command.commandId,
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 @@ -7852,6 +7914,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
6 changes: 6 additions & 0 deletions apps/server/src/orchestration-v2/ProjectionStore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1447,6 +1447,9 @@ export function threadShellFromProjection(
interactionMode: projection.thread.interactionMode,
branch: projection.thread.branch,
worktreePath: projection.thread.worktreePath,
...(projection.thread.workspaceBindingId === undefined
? {}
: { workspaceBindingId: projection.thread.workspaceBindingId }),
pullRequests: threadPullRequestsOf(projection.thread),
...(projection.thread.linkedPullRequest === undefined
? {}
Expand Down Expand Up @@ -1715,6 +1718,9 @@ function shellFromState(input: {
interactionMode: input.state.thread.interactionMode,
branch: input.state.thread.branch,
worktreePath: input.state.thread.worktreePath,
...(input.state.thread.workspaceBindingId === undefined
? {}
: { workspaceBindingId: input.state.thread.workspaceBindingId }),
pullRequests: threadPullRequestsOf(input.state.thread),
...(input.state.thread.linkedPullRequest === undefined
? {}
Expand Down
Loading
Loading