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
19 changes: 14 additions & 5 deletions apps/mobile/src/features/threads/ThreadDetailScreen.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -173,9 +173,15 @@ export interface ThreadDetailScreenProps {
* The server has not created this thread yet. "preparing" runs while the
* queued creation is delivered (a worktree may be checking out); "failed"
* is a rejected creation whose content went back to the project draft.
* `serverOwned` marks a thread the server already created that is still
* setting up its worktree; follow-ups queue on the server behind that setup.
*/
readonly creationState:
| { readonly kind: "preparing"; readonly preparingWorktree: boolean }
| {
readonly kind: "preparing";
readonly preparingWorktree: boolean;
readonly serverOwned: boolean;
}
| { readonly kind: "failed"; readonly reason: string; readonly onEditTask: () => void }
| null;
readonly activePendingApproval: PendingApproval | null;
Expand Down Expand Up @@ -1439,11 +1445,14 @@ export const ThreadDetailScreen = memo(function ThreadDetailScreen(props: Thread
canStopThread={props.canStopThread}
environmentId={props.environmentId}
projectCwd={props.threadCwd ?? props.projectWorkspaceRoot}
// Follow-ups typed during setup wait in the draft: queueing
// them against a thread id the server may still reject
// would strand them in the outbox.
// Follow-ups typed before the server creates the thread wait in
// the draft: queueing them against a thread id the server may
// still reject would strand them in the outbox.
sendBlockedReason={
props.creationState?.kind === "preparing" ? "Starting the task…" : null
props.creationState?.kind === "preparing" &&
!props.creationState.serverOwned
? "Starting the task…"
: null
}
draftKey={props.composerDraftKey ?? undefined}
followUpBehavior={props.followUpBehavior}
Expand Down
5 changes: 4 additions & 1 deletion apps/mobile/src/features/threads/ThreadRouteScreen.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -889,7 +889,9 @@ function ThreadRouteContent(
]);
const creationState = ((): ThreadDetailScreenProps["creationState"] => {
if (selectedThreadCreation === null) {
return awaitingBootstrapTurn ? { kind: "preparing", preparingWorktree: true } : null;
return awaitingBootstrapTurn
? { kind: "preparing", preparingWorktree: true, serverOwned: true }
: null;
}
if (selectedThreadCreation.outcome?.kind === "failed") {
return {
Expand All @@ -901,6 +903,7 @@ function ThreadRouteContent(
return {
kind: "preparing",
preparingWorktree: selectedThreadCreation.message.creation?.workspaceMode === "worktree",
serverOwned: false,
};
})();
if (!environmentId || !threadId) {
Expand Down
18 changes: 16 additions & 2 deletions apps/mobile/src/state/use-thread-composer-state.ts
Original file line number Diff line number Diff line change
Expand Up @@ -677,8 +677,20 @@ export function useThreadComposerState() {
alternateModifier: followUpOverride !== undefined && followUpOverride !== followUpBehavior,
activeTurnDefault: followUpBehavior,
});
const followUpDispatchMode =
followUpAction === "auto" ? null : followUpAction === "queue" ? "queue" : "auto";
// A send while the worktree is still being set up waits behind that setup,
// so the server can hold it if the setup ends without a worktree.
// Before detail loads, the shell's latest run can be a queued follow-up
// while its runtime still reports the setup.
const setupRunning = [selectedThreadActivityRun?.status, selectedThreadRuntime?.status].some(
(status) => status === "preparing" || status === "starting",
);
const followUpDispatchMode = setupRunning
? "queue"
: followUpAction === "auto"
? null
: followUpAction === "queue"
? "queue"
: "auto";

const metadata = makeQueuedMessageMetadata();
const messageId = MessageId.make(metadata.messageId);
Expand Down Expand Up @@ -735,7 +747,9 @@ export function useThreadComposerState() {
saveQueuedRunEdit,
selectedEnvironmentRuntime?.connectionState,
selectedEnvironmentRuntime?.serverConfig,
selectedThreadActivityRun?.status,
selectedThreadCreation,
selectedThreadRuntime?.status,
selectedThreadShell,
uploadThreadFeedback,
],
Expand Down
65 changes: 54 additions & 11 deletions apps/server/src/orchestration-v2/Orchestrator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1266,7 +1266,10 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
]);
});

const startNextQueuedRun = (threadId: ThreadId, options?: { readonly failedRunId?: RunId }) =>
const startNextQueuedRun = (
threadId: ThreadId,
options?: { readonly failedRunId?: RunId; readonly endedRunId?: RunId },
) =>
Effect.gen(function* () {
// Every terminal run checks the queue. Only a deliverable queued run
// needs the transcript for provider handoff and legacy import context.
Expand Down Expand Up @@ -1301,15 +1304,22 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
// Hold the queue so the user decides when to resume it. Validation
// failures (setup, unsupported handoff) belong to that message alone,
// and a message queued for another provider is how users recover.
// A message queued behind a worktree that was cancelled or never created
// would run in the project checkout instead, so it waits for the user.
const failedRun = latestExecutedRun(projection.runs);
const failureClass =
failedRun?.id === options?.failedRunId
? latestRootProviderFailure(failedRun, projection.turnItems)?.class
: undefined;
const worktreeMissing =
failedRun?.id === options?.endedRunId &&
Comment thread
coderabbitai[bot] marked this conversation as resolved.
failedRun?.workspacePreparation?.type === "worktree" &&
projection.thread.worktreePath === null;
if (
failureClass !== undefined &&
failureClass !== "validation_error" &&
failedRun?.providerInstanceId === queuedRun.providerInstanceId
worktreeMissing ||
Comment thread
coderabbitai[bot] marked this conversation as resolved.
(failureClass !== undefined &&
failureClass !== "validation_error" &&
failedRun?.providerInstanceId === queuedRun.providerInstanceId)
) {
const now = yield* DateTime.now;
yield* writeSystemEvents(
Expand Down Expand Up @@ -4852,9 +4862,22 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
(candidate) => candidate.id === projection.thread.activeProviderThreadId,
);
const activeRun = projection.runs.find(isBlockingRun);
// A message queued during worktree setup can arrive after that setup failed
// without a worktree. It waits, held, like the messages queued before the
// failure, instead of starting in the project checkout.
const latestRun = latestExecutedRun(projection.runs);
const worktreeMissingRun =
activeRun === undefined &&
dispatchMode.type === "queue_after_active" &&
(latestRun?.status === "failed" || latestRun?.status === "interrupted") &&
latestRun.workspacePreparation?.type === "worktree" &&
projection.thread.worktreePath === null
? latestRun
: undefined;
const queueBehindRun = activeRun ?? worktreeMissingRun;
const pendingMergeBackTransfers = pendingMergeBackTransfersForThread(projection);
const shouldQueue =
activeRun !== undefined &&
queueBehindRun !== undefined &&
(dispatchMode.type === "defer_start" ||
dispatchMode.type === "start_immediately" ||
dispatchMode.type === "queue_after_active");
Expand All @@ -4869,13 +4892,13 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
const queueProviderThread =
activeProviderThread ??
projection.providerThreads.find(
(candidate) => candidate.id === activeRun.providerThreadId,
(candidate) => candidate.id === queueBehindRun.providerThreadId,
);
if (queueProviderThread === undefined) {
return yield* new OrchestratorDispatchError({
commandId: command.commandId,
commandType: command.type,
cause: `Active run ${activeRun.id} has no provider thread for queued dispatch.`,
cause: `Active run ${queueBehindRun.id} has no provider thread for queued dispatch.`,
});
}
const now = yield* DateTime.now;
Expand Down Expand Up @@ -4929,7 +4952,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
const attemptId = idAllocator.derive.runAttempt({ runId, attemptOrdinal: 1 });
const rootNodeId = idAllocator.derive.rootNode({ runId });
const checkpointScope =
activeRun.status === "preparing"
queueBehindRun.status === "preparing" || worktreeMissingRun !== undefined
? null
: yield* runtimePolicy.resolve({ thread: projection.thread, modelSelection }).pipe(
Effect.flatMap((resolvedRuntimePolicy) =>
Expand Down Expand Up @@ -4966,7 +4989,8 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
rootNodeId,
activeAttemptId: attemptId,
status: "queued",
...(projection.runs.some(
...(worktreeMissingRun !== undefined ||
projection.runs.some(
(candidate) => candidate.status === "queued" && candidate.queueHeld === true,
)
? { queueHeld: true }
Expand Down Expand Up @@ -8273,6 +8297,20 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
type: "run.updated",
payload: { ...state.run, status: "preparing", completedAt: null },
});
// Messages held when this preparation failed follow the run again. A hold
// from a later run (a stop, say) is not this failure's to release.
const failureOwnsQueue = latestExecutedRun(projection.runs)?.id === state.run.id;
for (const run of projection.runs) {
if (!failureOwnsQueue || run.status !== "queued" || run.queueHeld !== true) continue;
yield* emitEvent({
type: "run.updated",
threadId: command.threadId,
runId: run.id,
providerInstanceId: run.providerInstanceId,
occurredAt: now,
payload: { ...run, queueHeld: false },
});
}
});

/**
Expand Down Expand Up @@ -10889,8 +10927,13 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
threadId,
startNextQueuedRun(
threadId,
stored.event.type === "run.updated" && stored.event.payload.status === "failed"
? { failedRunId: stored.event.payload.id }
stored.event.type === "run.updated"
? {
endedRunId: stored.event.payload.id,
...(stored.event.payload.status === "failed"
? { failedRunId: stored.event.payload.id }
: {}),
}
: undefined,
),
)
Expand Down
156 changes: 156 additions & 0 deletions apps/server/src/orchestration-v2/ThreadLaunchService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1422,6 +1422,162 @@ it.effect("retries a failed workspace preparation on the same run", () => {
}).pipe(Effect.provide(harness.layer));
});

it.effect.each(["fails", "is interrupted"] as const)(
"holds a message queued during setup when the workspace preparation %s",
(ending) =>
Effect.gen(function* () {
const fetchEntered = yield* Deferred.make<void>();
const allowFetch = yield* Deferred.make<void>();
let fetchFailures = 1;
const harness = makeHarness({
fetchRemote: () =>
fetchFailures-- <= 0
? Effect.void
: Deferred.succeed(fetchEntered, undefined).pipe(
Effect.andThen(Deferred.await(allowFetch)),
Effect.andThen(
Effect.fail(
new GitCommandError({
operation: "GitVcsDriver.fetchRemote",
command: "git",
cwd: project.workspaceRoot,
detail: "Git could not reach the remote.",
exitCode: 128,
}),
),
),
),
});
yield* Effect.gen(function* () {
const launches = yield* ThreadLaunch.ThreadLaunchService;
const threads = yield* ThreadManagement.ThreadManagementService;
const launched = yield* launches.launch(
launchInput({
command: `command:launch:queued-setup-${ending}`,
thread: `thread:launch:queued-setup-${ending}`,
message: "First message",
workspace: { type: "worktree", baseRef: "main", startFromOrigin: true },
}),
);
yield* Deferred.await(fetchEntered);
const queued = yield* threads.sendToThread({
projectId,
commandId: CommandId.make(`command:launch:queued-setup-${ending}:follow-up`),
threadId: launched.threadId,
messageId: MessageId.make(`message:launch:queued-setup-${ending}:follow-up`),
text: "Sent during setup",
attachments: [],
mode: "queue",
createdBy: "user",
creationSource: "web",
});
assert.equal(queued.delivery, "queued");
if (ending === "fails") {
yield* Deferred.succeed(allowFetch, undefined);
} else {
// An agent stop does not ask to hold the queue.
yield* threads.dispatch({
type: "run.interrupt",
commandId: CommandId.make("command:launch:queued-setup-interrupt"),
threadId: launched.threadId,
runId: launched.projection.runs[0]!.id,
});
}

// Without its worktree the follow-up would start in the project checkout.
yield* threads.streamStoredEventsFrom({ threadId: launched.threadId }).pipe(
Stream.filter(
(stored) =>
stored.event.type === "run.updated" && stored.event.payload.queueHeld === true,
),
Stream.runHead,
);
const projection = yield* threads.getThreadProjection(launched.threadId);
assert.equal(projection.runs.at(-1)?.status, "queued");
assert.equal(projection.thread.worktreePath, null);

if (ending === "fails") {
// A successful retry lets the held message follow the first turn again.
yield* launches.retryPreparation({
commandId: CommandId.make("command:launch:queued-setup-retry"),
threadId: launched.threadId,
runId: launched.projection.runs[0]!.id,
});
const retried = yield* threads.getThreadProjection(launched.threadId);
assert.equal(retried.runs.at(-1)?.queueHeld, false);
}
}).pipe(Effect.provide(harness.layer));
}),
);

it.effect.each(["failed", "interrupted"] as const)(
"holds a message queued during setup that arrives after the setup %s",
(ending) => {
let fetchEntered: Deferred.Deferred<void> | null = null;
const harness = makeHarness({
fetchRemote: () =>
ending === "interrupted"
? Deferred.succeed(fetchEntered!, undefined).pipe(Effect.andThen(Effect.never))
: Effect.fail(
new GitCommandError({
operation: "GitVcsDriver.fetchRemote",
command: "git",
cwd: project.workspaceRoot,
detail: "Git could not reach the remote.",
exitCode: 128,
}),
),
});
return Effect.gen(function* () {
fetchEntered = yield* Deferred.make<void>();
const launches = yield* ThreadLaunch.ThreadLaunchService;
const threads = yield* ThreadManagement.ThreadManagementService;
const launched = yield* launches.launch(
launchInput({
command: `command:launch:late-queued-follow-up-${ending}`,
thread: `thread:launch:late-queued-follow-up-${ending}`,
message: "First message",
workspace: { type: "worktree", baseRef: "main", startFromOrigin: true },
}),
);
if (ending === "interrupted") {
yield* Deferred.await(fetchEntered);
// An agent stop does not ask to hold the queue.
yield* threads.dispatch({
type: "run.interrupt",
commandId: CommandId.make("command:launch:late-queued-follow-up:interrupt"),
threadId: launched.threadId,
runId: launched.projection.runs[0]!.id,
});
}
yield* threads.streamStoredEventsFrom({ threadId: launched.threadId }).pipe(
Stream.filter(
(stored) => stored.event.type === "run.updated" && stored.event.payload.status === ending,
),
Stream.runHead,
);

// Queued while setup ran, delivered only after it failed: it must not start
// in the project checkout.
const queued = yield* threads.sendToThread({
projectId,
commandId: CommandId.make(`command:launch:late-queued-follow-up-${ending}:send`),
threadId: launched.threadId,
messageId: MessageId.make(`message:launch:late-queued-follow-up-${ending}:send`),
text: "Sent during setup",
attachments: [],
mode: "queue",
createdBy: "user",
creationSource: "web",
});
assert.equal(queued.delivery, "queued");
const projection = yield* threads.getThreadProjection(launched.threadId);
assert.equal(projection.runs.at(-1)?.status, "queued");
assert.equal(projection.runs.at(-1)?.queueHeld, true);
}).pipe(Effect.provide(harness.layer));
},
);

it.effect("a retry reuses a recorded worktree without undoing its branch rename", () => {
let setupFailures = 1;
const harness = makeHarness({
Expand Down
Loading
Loading