import { runRanAfter } from "@t3tools/shared/orchestrationV2ThreadError";
import { resolveProjectSettings } from "@t3tools/shared/projectSettings";
import {
CommandId,
MessageId,
type OrchestrationV2Run,
type RunId,
type ThreadId,
} from "@t3tools/contracts";
import * as Effect from "effect/Effect";
import type { ProjectionRuntimeRecoveryState } from "./ProjectionStore.ts";
import * as ServerSettings from "../serverSettings.ts";
import { isNativeMaintenanceCommand } from "./Orchestrator.ts";
import * as ThreadManagementService from "./ThreadManagementService.ts";
import {
isRestartNoteSource,
restartCancelledBackgroundWorkNote,
restartContinuationNote,
} from "./RestartBackgroundNote.ts";
const CONTINUE_PROMPT = "Continue where you left off.";
/**
* Resume only an unfinished root turn. Leftover background work is cleaned up
* separately and reported on the next user turn; it must not wake a settled run.
*/
export function restartContinuationRun(
projection: Pick<
ProjectionRuntimeRecoveryState,
"thread" | "runs" | "providerThreads" | "providerSessions" | "providerTurns"
>,
): OrchestrationV2Run | undefined {
if (projection.thread.archivedAt !== null || projection.thread.deletedAt !== null) return;
// Queued runs never started; recovery holds them behind the cut run.
const run = projection.runs.reduce<OrchestrationV2Run | undefined>(
(latest, candidate) =>
candidate.status !== "queued" && (!latest || runRanAfter(candidate, latest))
? candidate
: latest,
undefined,
);
if (!run) return;
const preparedContinuation =
run.status === "starting" && run.restartContinuationOfRunId !== undefined;
if (run.status !== "running" && !preparedContinuation) return;
const liveTurnRequired = !preparedContinuation;
if (projection.thread.providerInstanceId !== run.providerInstanceId) return;
const providerThread = projection.providerThreads.find(
(thread) => thread.id === run.providerThreadId,
);
if (
!providerThread ||
providerThread.appThreadId !== projection.thread.id ||
providerThread.ownerNodeId !== null ||
providerThread.providerInstanceId !== run.providerInstanceId ||
providerThread.nativeThreadRef?.nativeId == null ||
providerThread.nativeThreadRef.strength !== "strong" ||
providerThread.nativeThreadRef.driver !== providerThread.driver ||
(liveTurnRequired && providerThread.status !== "active") ||
providerThread.status === "closed" ||
providerThread.status === "archived"
)
return;
const session = projection.providerSessions.find(
(candidate) => candidate.id === providerThread.providerSessionId,
);
// Most adapters keep a live session "ready" through its turns, so only a
// stopped or failed session rules out a live turn.
if (
session === undefined ||
session.providerInstanceId !== run.providerInstanceId ||
session.driver !== providerThread.driver ||
(liveTurnRequired && (session.status === "stopped" || session.status === "error"))
)
return;
if (
liveTurnRequired &&
!projection.providerTurns.some(
(turn) =>
turn.providerThreadId === providerThread.id &&
turn.runAttemptId === run.activeAttemptId &&
turn.status === "running",
)
)
return;
return run;
}
export const continueRestartedRun = Effect.fn("RestartContinuation.continueRestartedRun")(
function* (input: { readonly threadId: ThreadId; readonly sourceRunId: RunId }) {
const settings = yield* ServerSettings.ServerSettingsService;
const enabled = yield* settings.getSettings.pipe(Effect.orElseSucceed(() => null));
if (!enabled) return;
const threads = yield* ThreadManagementService.ThreadManagementService;
const messageId = MessageId.make(`message:restart-continuation:${input.sourceRunId}`);
const projection = yield* threads.getThreadRecords(
input.threadId,
["messages", "runs", "providerTurns", "attempts"],
{ messageIds: [messageId] },
);
if (
!resolveProjectSettings(enabled, projection.thread.projectId).settings
.continueThreadsAfterServerUpdate
)
return;
if (projection.thread.archivedAt !== null || projection.thread.deletedAt !== null) return;
if (projection.messages.some((message) => message.id === messageId)) return;
const source = projection.runs.find((run) => run.id === input.sourceRunId);
// Pending effects from older versions may target settled background work,
// including waiting runs that reconciliation subsequently cancelled.
if (
!source ||
source.status !== "cancelled" ||
isRestartNoteSource(source, projection.providerTurns)
)
return;
// A user submission after reconciliation takes precedence over an automatic
// prompt. Queued runs never started and stay held behind this one.
if (
projection.runs.some(
(run) => run.id !== source.id && run.status !== "queued" && runRanAfter(run, source),
)
)
return;
if (projection.thread.providerInstanceId !== source.providerInstanceId) return;
const sourceRecords = yield* threads.getThreadRecords(
input.threadId,
["messages", "turnItems"],
{
messageIds: [source.userMessageId],
turnItemRunIds: [source.id],
turnItemTypes: ["run_interrupt_request"],
},
);
// The user asked this run to stop before the restart cut it.
if (
sourceRecords.turnItems.some(
(item) => item.runId === source.id && item.type === "run_interrupt_request",
)
)
return;
const sourceMessage = sourceRecords.messages.find(
(message) => message.id === source.userMessageId,
);
if (sourceMessage !== undefined && isNativeMaintenanceCommand(sourceMessage)) return;
const note = restartContinuationNote(
source,
projection.runs,
projection.providerTurns,
projection.attempts,
);
const noteText =
note.work.length === 0 ? undefined : restartCancelledBackgroundWorkNote(note.work);
yield* threads.dispatch({
type: "message.dispatch",
commandId: CommandId.make(`command:restart-continuation:${input.sourceRunId}`),
threadId: input.threadId,
messageId,
text:
noteText === undefined
? CONTINUE_PROMPT
: note.settled
? noteText
: `${noteText}\n\n${CONTINUE_PROMPT}`,
attachments: [],
modelSelection: source.modelSelection,
dispatchMode: { type: "start_immediately" },
createdBy: "agent",
creationSource: "server",
restartContinuationOfRunId: input.sourceRunId,
});
},
// A delegated child this declined to continue still owes its parent a
// result. Once a continuation run exists this is a no-op; that run settles it.
(effect, input) =>
effect.pipe(
Effect.andThen(
Effect.gen(function* () {
const threads = yield* ThreadManagementService.ThreadManagementService;
yield* threads.recoverDelegatedTask(input.threadId, input.sourceRunId);
}),
),
),
);