import { randomUUID } from "node:crypto";
import { readFileSync, writeSync } from "node:fs";
import { join } from "node:path";
import {
ExplicitInternalActivationError,
runDirectoryFromHostContext,
type HostContext,
type RoleEnvelopeHost,
type RoleHost,
} from "./host-contracts.ts";
import { Value } from "typebox/value";
import { sitianReport } from "./sitian-facade.ts";
import { createSubmissionLedgerHost, sealAcceptedSubmission } from "./submission-ledger.ts";
import { registerFiledSubmissionTool, type FiledSubmissionBeforeAccept } from "./filed-submission.ts";
import type { RoleSubmissionDeclaration } from "./role-submission-declarations.ts";
import { activationTraceRecordSchema, namedActivationCause, type ActivationTraceRecord, type ActivationTraceWriter } from "./activation-trace.ts";
import { homeFromRunDirectory } from "./activation-ledger-topology.ts";
import {
durableSessionPointer,
resolveBookKeyFromGit,
} from "./activation-ledger.ts";
import { readPackageMaterial } from "./session-opening-materials.ts";
import { writeStderrJsonlRecord } from "./stderr-jsonl.ts";
import {
ENGINE_DETOUR_TOOL_NAME,
ENGINE_MODEL_FLAG_NAME,
resolveEngineModel,
resolveEngineName,
} from "./engine-detour.ts";
import { engineSessionMaterialFromOptions } from "./package-resources/engine-material.ts";
import { registerEngineDetourTool } from "./engine-detour-tool.ts";
import {
createReceiptDeliveryPolicy,
deliveryLimitFromConfig,
deliveryLimitFromEnv,
NO_RECEIPT_LIFECYCLE_ENTRY_TYPE,
priorReceiptContinuation,
RECEIPT_DELIVERY_REQUEST_ENTRY,
RECEIPT_REJECTION_ENTRY,
receiptAttemptPointer,
} from "./receipt-delivery-policy.ts";
import {
COLLECTOR_CONSTRUCTION_TOOLS,
COLLECTOR_REQUIRED_TOOLS,
COLLECTOR_TRANSPORT_FLAGS,
createCollectorRoleRuntime,
type CollectorActivation,
} from "./collector-role.ts";
import type { ComplianceDecision } from "./compliance-transport.ts";
import { createDoctorRoleRuntime } from "./doctor-role.ts";
import {
createNotaryRoleRuntime,
NOTARY_SESSION_BOUND_ENTRY,
projectNotaryBoundFromFlags,
readNotaryTicketFlag,
} from "./notary-role.ts";
import { NOTARY_SOURCE_RUN_FLAG, NOTARY_TICKET_FLAG } from "./notary-contracts.ts";
import {
COUNTERSIGN_TOOL_SPEC,
type CountersignRuntimeDependencies,
} from "./countersign-role.ts";
import {
GLEANER_LEFT_TOOL_SPEC,
type GleanerLeftRuntimeDependencies,
} from "./gleaner-left-role.ts";
import { GLEANER_LEFT_BASE_FLAG } from "./gleaner-left-contracts.ts";
import {
INSPECTOR_TOOL_SPEC,
type InspectorRuntimeDependencies,
} from "./inspector-role.ts";
import { INSPECTOR_SOURCE_RUN_FLAG } from "./inspector-contracts.ts";
import {
DIARIST_TOOL_SPEC,
type DiaristRuntimeDependencies,
} from "./diarist-role.ts";
import {
SECRETARIAT_OUTPUT_TOOL_SPEC,
type SecretariatRuntimeDependencies,
} from "./secretariat-role.ts";
import { SECRETARIAT_OUTPUT_TOOL_NAME } from "./secretariat-contracts.ts";
import {
projectDiaristSessions,
projectDiaristTicketSessions,
} from "./diarist-contracts.ts";
import { commitDiaristProjection } from "./diarist.ts";
import { TicketProvenanceInputError } from "./ticket-provenance.ts";
import { bindTicketNumberOnRunDirectory } from "./public-cli/invocation.ts";
import {
GATEKEEPER_TOOL_SPEC,
type GatekeeperRuntimeDependencies,
} from "./gatekeeper-role.ts";
import {
NAVIGATOR_TOOL_SPEC,
type NavigatorRuntimeDependencies,
} from "./navigator-role.ts";
import {
AUDITOR_TOOL_SPEC,
type AuditorRuntimeDependencies,
} from "./auditor-role.ts";
import { formatNavigatorReport, NAVIGATOR_EVENT_TYPE, NAVIGATOR_ROUTE_PLAYBOOK_FAILURE_ENTRY, navigatorSubjectKey, navigatorUnavailableError, subjectPath, type NavigatorAttendance, type NavigatorAttendanceOptions, type NavigatorEvent, type NavigatorPhase, type NavigatorReport, type NavigatorSettlement, type NavigatorSubjectProvenance, type NavigatorWorkContext } from "./navigator-attendance.ts";
import { loadNavigatorWorkBaseSuffix } from "./navigator-work-base.ts";
import {
buildNavigatorInfrastructureFailureFact,
classifyPackagedRoleTerminalResult,
packagedRoleOutputRejectionReason,
extractInfrastructureFailureEvidence,
NAVIGATOR_INVOCATION_ENTRY,
resolveLifecycleInvocationPrincipal,
} from "./navigator-invocation-identity.ts";
import { NAVIGATOR_POST_ROLE_GRACE_MS, raceNavigatorGrace } from "./public-cli/settlement.ts";
import { PACKAGED_ROLE_REGISTRY, isOfficerReviewSeat, packagedRoleActivationFlags, packagedRoleInputFlag, packagedRoleMetadata, packagedRoleOutputTool, packagedRolePhaseFlag, type PackagedRole } from "./packaged-role-registry.ts";
import {
createJudgeRoleRuntime,
} from "./judge-role.ts";
import {
createReviewerRoleRuntime,
type ReviewerActivation,
type ReviewerAdmittedInputs,
} from "./reviewer-role.ts";
import type { SubmissionGateNonPassResult } from "./gatekeeper-role.ts";
/**
* Private transport flag names/definitions for Reviewer admitted inputs.
* Shared activation envelope owns registration and decoding (ADR 0018).
*/
const REVIEWER_TRANSPORT_FLAGS = Object.freeze([
Object.freeze({
name: "ak-review-base",
definition: Object.freeze({
description: "Fixed base revision for the pinned review target",
type: "string" as const,
}),
}),
Object.freeze({
name: "ak-review-lens",
definition: Object.freeze({
description: "Single review lens: completeness or correctness",
type: "string" as const,
}),
}),
Object.freeze({
name: "ak-review-authority-refs",
definition: Object.freeze({
description: "JSON array of durable authority references projected as Skill --authority inputs",
type: "string" as const,
}),
}),
Object.freeze({
name: "ak-review-ticket-number",
definition: Object.freeze({
description: "Typed ticketNumber for Spec self-fetch primary path",
type: "string" as const,
}),
}),
] as const);
/** Notary private transport: optional court-diary ticket (ADR 0018 / 0075 — envelope-owned). */
const NOTARY_TRANSPORT_FLAGS = Object.freeze([
Object.freeze({
name: NOTARY_TICKET_FLAG.name,
definition: NOTARY_TICKET_FLAG.definition,
}),
] as const);
/** Gleaner-left private transport: comparison-base revision (ADR 0018 / #502). */
const GLEANER_LEFT_TRANSPORT_FLAGS = Object.freeze([
Object.freeze({
name: GLEANER_LEFT_BASE_FLAG.name,
definition: GLEANER_LEFT_BASE_FLAG.definition,
}),
] as const);
export const STATION_CHILD_FLAG = Object.freeze({
name: "ak-station-child",
definition: Object.freeze({
description: "Station-child role run (omit automatic navigator attendance; #840)",
type: "boolean" as const,
}),
} as const);
/**
* Decode private transport flags into frozen admitted inputs.
* Envelope-owned; necessary JSON decode only (public --authority-ref owns grammar).
*/
function decodeReviewerAdmittedInputs(getFlag: (name: string) => unknown): ReviewerAdmittedInputs {
let authorityRefs: readonly string[] | undefined;
const rawAuthorityRefs = getFlag("ak-review-authority-refs");
if (rawAuthorityRefs !== undefined) {
if (typeof rawAuthorityRefs !== "string") {
throw new Error("Reviewer authority refs transport error: flag value must be a string");
}
let parsed: unknown;
try {
parsed = JSON.parse(rawAuthorityRefs);
} catch (error) {
throw new Error(
`Reviewer authority refs transport error: JSON decode failed: ${errorText(error)}`,
);
}
if (!Array.isArray(parsed) || parsed.some((ref) => typeof ref !== "string")) {
throw new Error("Reviewer authority refs transport error: expected a JSON array of strings");
}
authorityRefs = Object.freeze(parsed as string[]);
}
let ticketNumber: number | undefined;
const rawTicketNumber = getFlag("ak-review-ticket-number");
// Shape-invalid flag values do not abort: omit typed candidate; branch→commit→degrade continues.
if (typeof rawTicketNumber === "string" && /^[1-9]\d*$/.test(rawTicketNumber)) {
ticketNumber = Number(rawTicketNumber);
}
const baseRevision = getFlag("ak-review-base");
if (typeof baseRevision !== "string" || !baseRevision.trim()) {
throw new Error("Reviewer role requires --ak-review-base");
}
const rawLens = getFlag("ak-review-lens");
if (rawLens !== "completeness" && rawLens !== "correctness") {
throw new Error("Reviewer role requires --ak-review-lens completeness|correctness");
}
return Object.freeze({
baseRevision,
lens: rawLens,
...(authorityRefs === undefined ? {} : { authorityRefs }),
...(ticketNumber === undefined ? {} : { ticketNumber }),
});
}
/**
* Parent system-prompt assembly for the shared activation envelope.
* Verification cadence lives in soul/quality-law only (ADR 0073: no machine copy).
*/
function assembleReviewerParentSystemPrompt(input: {
baseSystemPrompt: string;
soul: string;
}): string {
return [
input.baseSystemPrompt,
"",
"",
input.soul,
"",
].join("\n");
}
import {
CODER_OUTPUT_TOOL_NAME,
createCoderRoleRuntime,
createFixerRoleRuntime,
FIXER_FLAG_DEFINITIONS,
FIXER_OUTPUT_TOOL_NAME,
FIXER_PHASES,
} from "./worker-role.ts";
import { JUDGE_OUTPUT_TOOL_NAME } from "./package-contracts/judge-output.ts";
import { REVIEWER_OUTPUT_TOOL_NAME } from "./package-contracts/reviewer-output.ts";
import { DOCTOR_OUTPUT_TOOL_NAME } from "./doctor-contracts.ts";
import { MERGER_OUTPUT_TOOL_NAME } from "./merger-contracts.ts";
import { createMergerRoleRuntime, type MergerRoleDependencies } from "./merger-role.ts";
export {
buildNavigatorInfrastructureFailureFact,
classifyPackagedRoleTerminalResult,
extractInfrastructureFailureEvidence,
hasNavigatorInfrastructureFailureBase,
isAcceptedPackagedRoleTerminalResult,
isDurablePackagedRoleTerminalResult,
isNavigatorInfrastructureFailureFact,
NAVIGATOR_INFRASTRUCTURE_FAILURE_EVIDENCE_KEYS,
NAVIGATOR_INFRASTRUCTURE_FAILURE_KIND,
type NavigatorInfrastructureFailureFact,
type PackagedRoleTerminalClassification,
} from "./navigator-invocation-identity.ts";
export { activationTraceRecordSchema, namedActivationCause } from "./activation-trace.ts";
export type { ActivationTraceRecord, ActivationTraceWriter } from "./activation-trace.ts";
export {
ActivationGitRepositoryRequiredError,
ActivationLedgerError,
ActivationSessionFileMissingError,
activationBookDirectory,
durableSessionPointer,
resolveActivationLedgerHome,
resolveBookKeyFromGit,
} from "./activation-ledger.ts";
export type {
ActivationSessionManager,
ActivationSessionPointer,
} from "./activation-ledger.ts";
import {
NOTARY_OUTPUT_TOOL,
INSPECTOR_OUTPUT_TOOL,
GatekeeperDecisionError,
runGatekeeper,
gateOfficerForSubject,
} from "./gatekeeper-role.ts";
export {
NOTARY_OUTPUT_TOOL,
INSPECTOR_OUTPUT_TOOL,
GatekeeperDecisionError,
runGatekeeper,
gateOfficerForSubject,
};
export type { GatekeeperResult, GatekeeperSubject, SubmissionGateNonPassResult, GateOfficer, RunGatekeeperOptions } from "./gatekeeper-role.ts";
import { ParentQueueReaskError, unreadableDiscriminatorNotice } from "./submission-errors.ts";
import { REVIEW_QUEUE_STATUSES } from "./review-submission.ts";
import { isRecord, errorText } from "./unknown-value.ts";
export {
DOCTOR_EVIDENCE_TOOL_NAME,
DOCTOR_OUTPUT_TOOL_NAME,
} from "./doctor-role.ts";
export type { DoctorCase, DoctorCaseCost, DoctorSubmission, DoctorOutput, DoctorFinding } from "./doctor-contracts.ts";
export { validateDoctorSubmissionShape, validateDoctorOutput, DoctorEvidenceStore } from "./doctor-contracts.ts";
export { loadDoctorCase } from "./doctor-evidence.ts";
export {
JUDGE_OUTPUT_TOOL_NAME,
type JudgeVerdict,
} from "./judge-role.ts";
export { ENGINE_DETOUR_TOOL_NAME, AK_ROLE_ENGINE_ENV } from "./engine-detour.ts";
export {
REVIEWER_OUTPUT_TOOL_NAME,
type ReviewerIntent,
} from "./reviewer-role.ts";
export {
CODER_OUTPUT_TOOL_NAME,
FIXER_FLAG_DEFINITIONS,
FIXER_OUTPUT_TOOL_NAME,
FIXER_PHASES,
type CoderOutput,
type FixerOutput,
type WorkerOutput,
} from "./worker-role.ts";
export { fixerOutputSchema, validateFixerOutput } from "./package-contracts/fixer-output.ts";
export type { FixerBlocker, FixerClassResult, FixerPhase, FixerTestEvidence } from "./package-contracts/fixer-output.ts";
export { fixerPrerequisiteSchema, fixerPrerequisitesSchema, parseFixerPrerequisites, validateFixerPrerequisites } from "./package-contracts/fixer-packet.ts";
export type { FixerInvocationInput, FixerPrerequisite } from "./package-contracts/fixer-packet.ts";
export {
AUDITOR_SOUL_ROLES,
AK_ROLE_AUDITOR_SUBJECT_ENV,
loadAuditorSoul,
loadAuditorSoulFromSubjectInput,
resolveAuditorSubject,
} from "./auditor-soul.ts";
export type { AuditorSoulRole } from "./auditor-soul.ts";
export { JUDGE_AUDIT_TOOL_NAME, SOUL_AUDIT_TOOL_NAME } from "./judge-auditor.ts";
export { DOCTOR_AUDIT_TOOL_NAME, createPiDoctorAuditor } from "./doctor-auditor.ts";
export type { ComplianceDecision } from "./compliance-transport.ts";
export { COLLECTOR_OUTPUT_TOOL } from "./collector-role.ts";
export type { CollectorReceipt } from "./package-contracts/collector-output.ts";
export * from "./navigator-attendance.ts";
export { MERGER_INPUT_FLAG, createMergerRoleRuntime } from "./merger-role.ts";
export { MERGER_OUTPUT_TOOL_NAME, mergerInputSchema, mergerOutputSchema, validateMergerInput, validateMergerOutput } from "./merger-contracts.ts";
export type { MergerInput, MergerMaterial, MergerOutput } from "./merger-contracts.ts";
export { createProductionMergerGitState } from "./merger-git-state.ts";
export type { MergerGitState, ActiveMergerGitState } from "./merger-git-state.ts";
export type { MergerRoleDependencies } from "./merger-role.ts";
function activationStage(
role: PackagedRole,
activate: Record Promise>,
): { id: string; run(): Promise } {
const record = packagedRoleMetadata(role);
if (record === undefined) {
throw new Error(`Unsupported workflow role: ${String(role)}`);
}
return {
id: record.activationStage,
run: () => activate[role](),
};
}
function validateActivationTraceRecord(record: unknown): ActivationTraceRecord {
if (!Value.Check(activationTraceRecordSchema, record)) {
throw new TypeError("Activation trace record does not match its closed contract");
}
return record as ActivationTraceRecord;
}
async function emitActivationTrace(
writeTrace: (record: ActivationTraceRecord) => void | Promise,
record: unknown,
): Promise {
await writeTrace(validateActivationTraceRecord(record));
}
async function executeActivationStage(
role: string,
stage: { id: string; run(): Promise },
infrastructure: { clock(): string; writeTrace(record: ActivationTraceRecord): void | Promise },
): Promise {
try {
await stage.run();
} catch (activationError) {
try {
await emitActivationTrace(infrastructure.writeTrace, {
role,
stageId: stage.id,
status: "failed",
timestamp: infrastructure.clock(),
cause: namedActivationCause(activationError),
});
} catch (traceError) {
throw new AggregateError([activationError, traceError], `Activation stage ${stage.id} failed and its failure trace could not be emitted`);
}
throw activationError;
}
}
export function writeActivationTraceRecord(
record: ActivationTraceRecord,
write: typeof writeSync = writeSync,
): void {
writeStderrJsonlRecord(record, write);
}
export class ActivationBarrierError extends Error {
readonly code = "AK_ACTIVATION_NOT_ADMITTED";
constructor(role: unknown) {
super(`Workflow role ${String(role)} activation did not complete`);
this.name = "ActivationBarrierError";
}
}
export const WORKFLOW_ROLES = PACKAGED_ROLE_REGISTRY.map(({ role }) => role) as Array<(typeof PACKAGED_ROLE_REGISTRY)[number]["role"]>;
export const ROLE_FLAG = {
name: "ak-role",
definition: {
description: `Activate a packaged workflow role: ${WORKFLOW_ROLES.slice(0, -1).join(", ")}, or ${WORKFLOW_ROLES.at(-1)}`,
type: "string" as const,
},
} as const;
type NavigatorAttendanceDependency = NavigatorAttendance;
export type RoleRuntimeDependencies = {
/** Package root for packaged engine-note resolution (#879). */
packageRoot?: string;
/**
* Composition-root adapter table for nested gate officer summons (#969).
* Production leaves unset (default pi + packaged externals). Tests inject
* faux nested hosts so the public entry's reviewer summons use the same
* host resolution as production.
*/
hostAdapters?: readonly import("./public-cli/role-turn-host-resolution.ts").NamedRoleTurnHostAdapter[];
/** Non-identity opening materials delivered in the ordinary role brief. */
loadRoleReferenceMaterials?(role: PackagedRole): Promise;
/** One soul loader. The role argument selects the registry record. */
loadRoleSoul(role: PackagedRole): Promise;
loadFixPacket?(path: string): Promise;
loadCoderTask?(path: string): Promise;
loadNotarySourceRun?(path: string): Promise;
loadDoctorCase?(path: string): Promise;
loadMergerInput?(path: string): Promise;
createNavigatorAttendance?(options: { context: HostContext; role: string; phase: NavigatorPhase; subjectKey: string; subject: string; authority: string; contextError?: unknown; invocationId: string; deliveryRequestLimit?: number; onEvent: (event: import("./navigator-attendance.ts").NavigatorEvent, report: import("./navigator-attendance.ts").NavigatorReport) => void | Promise }): NavigatorAttendanceDependency | Promise;
loadNavigatorWorkContext?(options: { context: HostContext; role: string; phase: NavigatorPhase; getFlag?: (name: string) => unknown }): Promise;
activationClock?(): string;
activationTraceWriter?: (record: ActivationTraceRecord) => void | Promise;
};
function abortContext(ctx: { abort(): void }): void {
ctx.abort();
}
function failInfrastructure(error: unknown, ctx: { mode: string; abort(): void }): never {
abortContext(ctx);
if (ctx.mode === "print" || ctx.mode === "json") process.exitCode = 1;
throw error;
}
/**
* Envelope-owned pending infrastructure failure: closed fact + original Error evidence.
* tool_result projects once from this record; settlement consumes durable details as-is.
*/
type PendingInfrastructureFailure = {
readonly details: Record;
};
function buildPendingInfrastructureFailure(error: unknown): PendingInfrastructureFailure {
return {
details: {
...buildNavigatorInfrastructureFailureFact(),
...extractInfrastructureFailureEvidence(error),
},
};
}
function navigatorPhase(roleHost: RoleHost, role: string): NavigatorPhase {
const metadata = packagedRoleMetadata(role);
if (metadata === undefined || metadata.phases[0] === null) return null;
const phaseFlag = packagedRolePhaseFlag(role);
const requested = phaseFlag === undefined ? undefined : roleHost.getFlag(phaseFlag);
return requested === "apply" ? "apply" : "plan";
}
function navigatorOutputTool(role: string): string | undefined {
return packagedRoleOutputTool(role);
}
/** Status leaves that are not the shared `status` key, in registry order. */
function receiptStatusFromRegistry(details: Record): string | undefined {
for (const entry of PACKAGED_ROLE_REGISTRY) {
if (!("receiptStatusKey" in entry)) continue;
const value = details[entry.receiptStatusKey];
if (typeof value === "string") return value;
}
return undefined;
}
export function publicNavigatorSettlement(role: string, phase: NavigatorPhase, event: { toolName: string; isError?: unknown; details: unknown }): NavigatorSettlement | undefined {
// One shared classifier owns terminal discriminant (lifecycle + settlement + extractors).
if (event.toolName !== navigatorOutputTool(role)) return undefined;
const classification = classifyPackagedRoleTerminalResult(event);
if (classification.kind === "nonterminal") return undefined;
if (classification.kind === "infrastructure") {
return { kind: "role_infrastructure_failure", role, phase };
}
// accepted/human — project role/phase status; classifier already rejected infra/contradiction.
const details = isRecord(event.details)
? event.details as Record
: {};
const status = typeof details.status === "string"
? details.status
: receiptStatusFromRegistry(details);
if (status !== undefined && status === "escalate") {
return { kind: "human_decision", role, phase, status };
}
return { kind: "accepted", role, phase, ...(status === undefined ? {} : { status }) };
}
export async function projectClosedSubmissionLifecycle(
closed: import("./submission-ledger.ts").ClosedSubmission,
context: HostContext,
phase: NavigatorPhase,
recordAccepted: () => void,
settle: (settlement: NavigatorSettlement | undefined) => Promise,
): Promise {
recordAccepted();
const closure = {
toolName: navigatorOutputTool(closed.role)!,
isError: false,
details: closed.accepted,
};
let navigator: NavigatorEvent | undefined;
try {
navigator = await settle(publicNavigatorSettlement(closed.role, phase, closure));
} finally {
// A Navigator failure cannot erase an already accepted role closure.
// Missing attendance remains a typed unavailable at public settlement.
context.sessionManager.appendCustomEntry?.("ak-role-submission-closure", {
...closure,
...(navigator === undefined ? {} : { navigator }),
});
}
}
/**
* Shared registration envelope for filed officers (ADR 0018 / #572):
* activate, tool register, before_agent_start prompt, inventory check.
* Role module keeps label/soul/spec shape only; sole-final barrier is ledger-owned.
* Optional beforeAccept projects local submission facts before the tool returns.
*/
function createFiledOfficerRuntime(
roleHost: RoleHost,
spec: {
role: PackagedRole;
tool: RoleSubmissionDeclaration;
soulTag: string;
beforeAccept?: FiledSubmissionBeforeAccept;
},
dependencies: { loadSoul(): Promise },
) {
let soul: string | undefined;
let registered = false;
return {
async activate() {
const loaded = (await dependencies.loadSoul()).trim();
if (loaded.length === 0) throw new Error(`${spec.role} soul is empty`);
soul = loaded;
if (!registered) {
registered = true;
registerFiledSubmissionTool(roleHost, spec.tool, {
readyError: () => (soul === undefined ? `${spec.role} 职分未装载` : undefined),
...(spec.beforeAccept === undefined ? {} : { beforeAccept: spec.beforeAccept }),
});
roleHost.on("before_agent_start", (event) => {
if (soul === undefined) throw new Error(`${spec.role} 职分未装载`);
const tail = `\n\n<${spec.soulTag}_soul>\n${soul}\n${spec.soulTag}_soul>`;
return { systemPrompt: `${event.systemPrompt}${tail}` };
});
}
const all = roleHost.getAllTools().map((tool) => tool.name);
if (all.filter((name) => name === spec.tool.name).length !== 1) {
throw new Error(`${spec.role} required tool collision or missing: ${spec.tool.name}`);
}
},
};
}
function activationPublishesFlag(role: unknown, flagName: string): boolean {
if (typeof role !== "string") return false;
return packagedRoleActivationFlags(role).some((spec) => spec.flag === flagName);
}
export function createGleanerLeftRoleRuntime(
roleHost: RoleHost,
dependencies: GleanerLeftRuntimeDependencies,
) {
return createFiledOfficerRuntime(
roleHost,
{
role: "gleaner-left",
tool: GLEANER_LEFT_TOOL_SPEC,
soulTag: "gleaner-left",
},
dependencies,
);
}
function decodeGleanerLeftBase(getFlag: (name: string) => unknown): string {
const baseRevision = getFlag(GLEANER_LEFT_BASE_FLAG.name);
if (typeof baseRevision !== "string" || !baseRevision.trim()) {
throw new Error("Gleaner-left role requires --ak-gleaner-left-base");
}
return baseRevision;
}
export function createInspectorRoleRuntime(
roleHost: RoleHost,
dependencies: InspectorRuntimeDependencies,
) {
return createFiledOfficerRuntime(
roleHost,
{
role: "inspector",
tool: INSPECTOR_TOOL_SPEC,
soulTag: "inspector",
},
dependencies,
);
}
/** #639: direct Gatekeeper public seat on the shared filed-officer envelope. */
export function createGatekeeperRoleRuntime(
roleHost: RoleHost,
dependencies: GatekeeperRuntimeDependencies,
) {
return createFiledOfficerRuntime(
roleHost,
{
role: "gatekeeper",
tool: GATEKEEPER_TOOL_SPEC,
soulTag: "gatekeeper",
},
dependencies,
);
}
/** #639: direct Navigator public seat on the shared filed-officer envelope. */
export function createNavigatorRoleRuntime(
roleHost: RoleHost,
dependencies: NavigatorRuntimeDependencies,
) {
const base = createFiledOfficerRuntime(
roleHost,
{
role: "navigator",
tool: NAVIGATOR_TOOL_SPEC,
soulTag: "navigator",
},
dependencies,
);
let playbookBound = false;
return {
async activate() {
await base.activate();
if (playbookBound) return;
playbookBound = true;
// Standing system prompt, not a per-turn user message. Native read failure
// is the explanation (ADR 0061); no code-written sentence.
const loadRoutePlaybook = dependencies.loadRoutePlaybook
?? (() => readPackageMaterial("resources/navigator-route-playbook.md"));
let playbookRead: Promise | undefined;
roleHost.on("before_agent_start", async (event) => {
const basePrompt = typeof event.systemPrompt === "string" ? event.systemPrompt : "";
const parts: string[] = [];
try {
playbookRead ??= loadRoutePlaybook();
const content = await playbookRead;
dependencies.recordRoutePlaybookReadFailure?.(undefined);
if (content.trim() !== "") parts.push(content);
} catch (error) {
const message = errorText(error);
if (message.trim() !== "") {
dependencies.recordRoutePlaybookReadFailure?.(message);
parts.push(message);
}
}
const prompt = typeof event.prompt === "string" ? event.prompt : "";
const work = await loadNavigatorWorkBaseSuffix(prompt);
if (work !== undefined && work.trim() !== "") parts.push(work);
if (parts.length === 0) return;
const text = parts.join("\n\n");
return {
systemPrompt: basePrompt.trim() === "" ? text : `${basePrompt}\n\n${text}`,
};
});
},
};
}
/** #675: public 审刑院 seat on the shared filed-officer envelope. */
export function createAuditorRoleRuntime(
roleHost: RoleHost,
dependencies: AuditorRuntimeDependencies,
) {
const base = createFiledOfficerRuntime(
roleHost,
{
role: "auditor",
tool: AUDITOR_TOOL_SPEC,
soulTag: "auditor",
},
dependencies,
);
return {
async activate() {
await base.activate();
// Same tools whether nested or direct (#675): dossier tool always registered.
// Source run: only the shared --source-run input face (never own-run fallback).
const { createAuditorDossierTool, AUDITOR_DOSSIER_TOOL_NAME } =
await import("./auditor-dossier-tool.ts");
const { AK_ROLE_AUDITOR_SOURCE_RUN_ENV } = await import("./auditor-soul.ts");
const already = roleHost.getAllTools().some((tool) => tool.name === AUDITOR_DOSSIER_TOOL_NAME);
if (already) return;
const sourceRun =
typeof process.env[AK_ROLE_AUDITOR_SOURCE_RUN_ENV] === "string"
&& process.env[AK_ROLE_AUDITOR_SOURCE_RUN_ENV].trim() !== ""
? process.env[AK_ROLE_AUDITOR_SOURCE_RUN_ENV].trim()
: undefined;
roleHost.registerTool(createAuditorDossierTool(sourceRun) as never);
},
};
}
/**
* LLM typed court-target assertion from diarist output (ADR 0075 / #779).
* null/absent = true-unbound; positive safe integer = ticket N.
* Shape-only read of the typed field — not content judgment of free text.
* Any other shape fails honestly — never washes into unbound.
*/
function readDiaristTicketAssertion(
submitted: Record | undefined,
): { kind: "true-unbound" } | { kind: "ticket"; ticketNumber: number } | { kind: "invalid" } {
if (submitted === undefined || !("ticketNumber" in submitted)) {
return { kind: "true-unbound" };
}
const raw = submitted.ticketNumber;
if (raw === null) return { kind: "true-unbound" };
if (typeof raw === "number" && Number.isSafeInteger(raw) && raw >= 1) {
return { kind: "ticket", ticketNumber: raw };
}
if (typeof raw === "string" && /^[1-9]\d*$/.test(raw)) {
const n = Number(raw);
if (Number.isSafeInteger(n) && n >= 1) {
return { kind: "ticket", ticketNumber: n };
}
}
return { kind: "invalid" };
}
/** Shared durable coordinates for role tools that summon another public role. */
function readRoleRunCoordinates(ctx: HostContext, label: string): {
readonly runDirectory: string;
readonly projectRoot: string;
readonly home: string;
readonly admitted: Record;
} {
const runDirectory = runDirectoryFromHostContext(ctx);
if (runDirectory === undefined) throw new Error(`${label} requires AK_ROLE_RUN_DIR`);
const admittedPath = join(runDirectory, "admitted-request.json");
const admitted = JSON.parse(readFileSync(admittedPath, "utf8")) as Record;
if (typeof admitted.projectRoot !== "string" || admitted.projectRoot.trim() === "") {
throw new Error(`${label} admitted-request missing projectRoot (${admittedPath})`);
}
return {
runDirectory,
projectRoot: admitted.projectRoot,
home: homeFromRunDirectory(runDirectory),
admitted,
};
}
/** Run coordinates + optional pre-bound ticket from durable pages (#779). */
function readDiaristRunCoordinates(ctx: HostContext): {
readonly runDirectory: string;
readonly projectRoot: string;
readonly home: string;
readonly boundTicketNumber?: number;
} {
const coordinates = readRoleRunCoordinates(ctx, "diarist accept");
const bound =
typeof coordinates.admitted.ticketNumber === "number" &&
Number.isSafeInteger(coordinates.admitted.ticketNumber) &&
coordinates.admitted.ticketNumber >= 1
? coordinates.admitted.ticketNumber
: undefined;
return {
runDirectory: coordinates.runDirectory,
projectRoot: coordinates.projectRoot,
home: coordinates.home,
...(bound === undefined ? {} : { boundTicketNumber: bound }),
};
}
/** Plain-language re-ask when diarist bounds cannot be used (#901 / reask-not-explode). */
const DIARIST_BOUNDS_REASK =
"边界无法使用。多票请逐票重交 ticketSessions,单票重交 sessions;每卷 path + ranges,每端以原生 id 或本轮行号二选一指名。" as const;
/**
* #708 / #779 / #901: 起居郎 public seat on the shared filed-officer envelope.
* LLM judges ticket + dialogue bounds; mechanical layer reprojects the unique
* records.jsonl. Unusable bounds reask via ParentQueueReaskError. Machine
* facts never come from model self-report (锚定宪法).
*/
export function createDiaristRoleRuntime(
roleHost: RoleHost,
dependencies: DiaristRuntimeDependencies,
) {
return createFiledOfficerRuntime(
roleHost,
{
role: "diarist",
tool: DIARIST_TOOL_SPEC,
soulTag: "diarist",
beforeAccept: async ({ parameters, ctx }) => {
const submitted =
isRecord(parameters)
? (parameters as Record)
: undefined;
// Routing belongs to the public seam after this tool call has returned.
// Only completed receipts need local diary projection; other statuses
// pass through verbatim for that seam to route.
const status = submitted?.status;
if (status !== "completed") return parameters;
const assertion = readDiaristTicketAssertion(submitted);
const coords = readDiaristRunCoordinates(ctx);
// #836 7.3: pre-bound ticket is material for the LLM, not an override.
const ticketNumber = assertion.kind === "ticket" ? assertion.ticketNumber : undefined;
if (ticketNumber !== undefined) {
if (coords.boundTicketNumber === undefined) {
await bindTicketNumberOnRunDirectory(coords.runDirectory, ticketNumber);
}
}
// Strict-schema hosts emit ticketSessions: null for a single ticket.
const multiTicket = submitted?.ticketSessions != null;
const singleSessions = submitted && !multiTicket
? projectDiaristSessions(parameters)
: undefined;
const ticketSessions = submitted && multiTicket
? projectDiaristTicketSessions(submitted)
: singleSessions === undefined
? undefined
: [{ ticketNumber: ticketNumber ?? null, sessions: singleSessions }];
if (ticketSessions === undefined) {
throw new ParentQueueReaskError(DIARIST_BOUNDS_REASK);
}
try {
await commitDiaristProjection({
cwd: coords.projectRoot,
home: coords.home,
runDirectory: coords.runDirectory,
tickets: ticketSessions,
});
} catch (error) {
// Bound/session input failures → reask via typed identity (not message prefix).
// Unexpected infrastructure keeps its own identity (do not wash).
if (error instanceof ParentQueueReaskError) throw error;
if (error instanceof TicketProvenanceInputError) {
throw new ParentQueueReaskError(
`${DIARIST_BOUNDS_REASK}\n${error.message}`,
);
}
throw error;
}
return parameters;
},
},
dependencies,
);
}
/** Secretariat files before the public post-submission audit starts. */
export function createSecretariatRoleRuntime(
roleHost: RoleHost,
dependencies: SecretariatRuntimeDependencies,
) {
const base = createFiledOfficerRuntime(
roleHost,
{
role: "secretariat",
tool: SECRETARIAT_OUTPUT_TOOL_SPEC,
soulTag: "secretariat",
},
dependencies,
);
return {
async activate() {
await base.activate();
const packageRequired = [SECRETARIAT_OUTPUT_TOOL_NAME] as const;
const all = roleHost.getAllTools().map((tool) => tool.name);
for (const name of packageRequired) {
if (all.filter((item) => item === name).length !== 1) {
throw new Error(`secretariat required tool collision or missing: ${name}`);
}
}
// Host-neutral minimum package surface: declare package tools via
// setActiveTools without hardcoding host builtin names. Preserve any host
// surface already visible on getAllTools/getActiveTools so body rewrite
// (public gh path) stays reachable on hosts that expose it.
const priorActive = roleHost.getActiveTools();
const hostSurface = priorActive.length > 0 ? priorActive : all;
const nextActive = [
...new Set([...hostSurface, ...packageRequired]),
];
roleHost.setActiveTools(nextActive);
const active = roleHost.getActiveTools();
for (const name of packageRequired) {
if (!active.includes(name)) {
throw new Error(`secretariat failed to activate required tool ${name}`);
}
}
},
claimsEngineDetour: true,
};
}
export function createCountersignRoleRuntime(
roleHost: RoleHost,
dependencies: CountersignRuntimeDependencies,
) {
return createFiledOfficerRuntime(
roleHost,
{
role: "countersign",
tool: COUNTERSIGN_TOOL_SPEC,
soulTag: "countersign",
},
dependencies,
);
}
/**
* #959: last assistant text parts as navigator prose exit.
* Tool-call-only messages yield undefined — those already ride the tool path.
*/
function lastAssistantProse(
messages: readonly { role: string; content?: readonly { type: string; text?: string }[] }[],
): string | undefined {
for (let index = messages.length - 1; index >= 0; index -= 1) {
const message = messages[index];
if (message?.role !== "assistant" || !Array.isArray(message.content)) continue;
const texts: string[] = [];
for (const part of message.content) {
if (part.type === "text" && typeof part.text === "string" && part.text.length > 0) {
texts.push(part.text);
}
}
if (texts.length === 0) continue;
return texts.join("");
}
return undefined;
}
export function createRoleRuntimeExtension(
dependencies: RoleRuntimeDependencies,
/**
* Effective ceiling already resolved for this turn (#1132). Absent only on
* the Pi child, whose env is that same resolved value projected by the host.
*/
deliveryRequestLimit?: number,
): (envelopeHost: RoleEnvelopeHost) => void {
const deliveryLimit = deliveryRequestLimit === undefined
? deliveryLimitFromEnv(process.env)
: deliveryLimitFromConfig(deliveryRequestLimit);
return (envelopeHost) => {
let projectClosedSubmission: (closed: import("./submission-ledger.ts").ClosedSubmission, context: HostContext) => Promise = async () => {
throw new Error("角色终局投射接缝尚未初始化");
};
const roleHost = createSubmissionLedgerHost(
envelopeHost.host,
new Map(PACKAGED_ROLE_REGISTRY.reduce<(readonly [string, import("./public-cli/terminal.ts").TerminalRoleName | readonly import("./public-cli/terminal.ts").TerminalRoleName[]])[]>((entries, { role, outputTool }) => {
const existingIndex = entries.findIndex(([name]) => name === outputTool);
if (existingIndex < 0) entries.push([outputTool, role]);
else {
const previous = entries[existingIndex]![1];
entries[existingIndex] = [outputTool, [...(Array.isArray(previous) ? previous : [previous]), role]];
}
return entries;
}, [])),
failInfrastructure,
async (closed, context) => projectClosedSubmission(closed, context),
);
roleHost.registerFlag(ROLE_FLAG.name, ROLE_FLAG.definition);
// Reviewer transport flags: shared envelope owns registration (ADR 0018).
for (const flag of REVIEWER_TRANSPORT_FLAGS) {
roleHost.registerFlag(flag.name, flag.definition);
}
for (const flag of NOTARY_TRANSPORT_FLAGS) {
roleHost.registerFlag(flag.name, flag.definition);
}
roleHost.registerFlag(
INSPECTOR_SOURCE_RUN_FLAG.name,
INSPECTOR_SOURCE_RUN_FLAG.definition,
);
for (const flag of GLEANER_LEFT_TRANSPORT_FLAGS) {
roleHost.registerFlag(flag.name, flag.definition);
}
// Collector transport flags: shared envelope owns registration (ADR 0018 / #676 E).
for (const flag of COLLECTOR_TRANSPORT_FLAGS) {
roleHost.registerFlag(flag.name, flag.definition);
}
// Station-child identity (#840): omit navigator attendance. One flag.
roleHost.registerFlag(STATION_CHILD_FLAG.name, STATION_CHILD_FLAG.definition);
// Register model only. Pi never sets ak-engine — resolveEngineName must
// fall through to child-process env. An empty default would block that.
roleHost.registerFlag(ENGINE_MODEL_FLAG_NAME, {
description: "本次劳务引擎模型",
type: "string",
default: "",
});
let admitted = false;
let selectedRole: PackagedRole | undefined;
let roleReferenceMaterials = "";
/** Live Reviewer parent activation for envelope agent_start prompt assembly. */
let activeReviewerParent: ReviewerActivation | undefined;
let navigatorAttendance: NavigatorAttendanceDependency | undefined;
// #351: session-lifecycle owner for periodic OAuth refresh (orthogonal to role admission).
let pendingNavigatorPresentation: { event: import("./navigator-attendance.ts").NavigatorEvent; report: import("./navigator-attendance.ts").NavigatorReport } | undefined;
let navigatorActivation = 0;
let navigatorDeliveryClosed = false;
let pendingNavigatorSettlement: Promise | undefined;
let navigatorWorkContext: NavigatorWorkContext | undefined;
let navigatorSessionParent: string | undefined;
let navigatorCwd: string | undefined;
/** toolCallId → fact+evidence; one-shot projected onto durable tool_result (#475). */
const pendingInfrastructureFailures = new Map();
// Envelope-owned execute→tool_result bridge for submission non-pass (ADR 0018 / #525).
const pendingSubmissionNonPassByToolCallId = new Map();
let engineDetourRegistered = false;
// #288 primary-session thin adapter. The policy is the sole budget owner.
// #1132: one ceiling for this turn, closed over from the execution seam.
let receiptDelivery = createReceiptDeliveryPolicy(deliveryLimit);
let noReceiptRecorded = false;
/** Envelope-owned: abort/teardown attendance without re-blocking the parent court (#959). */
const disposeNavigatorAttendanceNonBlocking = (
attendance: NavigatorAttendanceDependency | undefined,
): void => {
if (attendance === undefined) return;
const recordDisposeFailure = (error: unknown): void => {
const diagnostic = errorText(error);
try {
sitianReport({
level: "event",
kind: "navigator-dispose-failure",
cwd: navigatorCwd,
sessionParent: navigatorSessionParent,
payload: { diagnostic },
source: "role-runtime",
});
} catch (recordError) {
try {
envelopeHost.appendEntry?.("ak-navigator-dispose-failure", {
diagnostic,
recordFailure: errorText(recordError),
});
} catch {
// Failure already diagnosed; recording must not create unhandled rejection (#959).
}
}
};
// Evaluate dispose inside try: Promise.resolve(dispose()) throws sync before .then attaches.
let pending: void | Promise;
try {
pending = attendance.dispose();
} catch (error) {
recordDisposeFailure(error);
return;
}
void Promise.resolve(pending).then(undefined, recordDisposeFailure);
};
const settleNavigatorProjection = async (settlement: NavigatorSettlement | undefined): Promise => {
const attendance = navigatorAttendance;
if (settlement === undefined || attendance === undefined) return;
const workContext = navigatorWorkContext;
const pending = (async () => {
// ADR 0052 / #959: every settlement that starts a post-role feed host round
// (accepted, human_decision, role_infrastructure_failure) shares one grace
// and the same honest unavailable projection — no unbounded parallel branch.
const settlePromise = attendance.settle(settlement);
// Attach catch immediately so a late rejection after grace timeout cannot
// surface as unhandledRejection / stale-ctx after session dispose (#675).
void settlePromise.catch(() => undefined);
const raced = await raceNavigatorGrace(settlePromise, NAVIGATOR_POST_ROLE_GRACE_MS);
if (raced.status !== "timeout") return;
navigatorDeliveryClosed = true;
if (pendingNavigatorPresentation === undefined) {
const report: NavigatorReport = {
disposition: "unavailable",
unavailableReason: "Navigator exceeded post-role delivery grace",
unavailableSource: "unknown",
unavailableCause: "unknown",
};
const event: NavigatorEvent = {
version: 1,
disposition: "unavailable",
invocationId: "post-role-grace-timeout",
role: settlement.role,
phase: settlement.phase,
subjectKey: workContext?.subjectKey ?? "",
unavailableReason: "Navigator exceeded post-role delivery grace",
unavailableSource: "unknown",
unavailableCause: "unknown",
};
pendingNavigatorPresentation = { event, report };
}
// 过时不候: do not await sidecar teardown — parent court must close.
disposeNavigatorAttendanceNonBlocking(attendance);
})();
pendingNavigatorSettlement = pending;
await pending;
return pendingNavigatorPresentation?.event;
};
projectClosedSubmission = async (closed, context) => projectClosedSubmissionLifecycle(
closed,
context,
navigatorPhase(roleHost, closed.role),
() => receiptDelivery.recordAccepted(),
settleNavigatorProjection,
);
roleHost.on("input", (_event) => {
const role = roleHost.getFlag(ROLE_FLAG.name);
if (role !== undefined && !admitted) return { action: "handled" as const };
return { action: "continue" as const };
});
// Reference law/guides use the existing typed reading-material channel;
// they are not rewritten into operator dialogue or the identity Soul.
roleHost.on("before_agent_start", () => roleReferenceMaterials === ""
? undefined
: { readingMaterial: { kind: "role-reference-materials", content: roleReferenceMaterials } });
roleHost.on("before_agent_start", async (event, ctx) => {
const role = roleHost.getFlag(ROLE_FLAG.name);
const prompt = event.prompt;
if (role === undefined) return;
if (!admitted || selectedRole !== role) {
failInfrastructure(new ActivationBarrierError(role), ctx);
}
// Notary session bound: envelope-owned lifecycle write (ADR 0018 / #582).
// Ticket flag register/read + session entry live here; role projects admitted bound only.
if (activationPublishesFlag(role, NOTARY_SOURCE_RUN_FLAG.name)) {
const bound = projectNotaryBoundFromFlags((name) => roleHost.getFlag(name));
if (bound !== undefined) {
ctx.sessionManager.appendCustomEntry?.(NOTARY_SESSION_BOUND_ENTRY, bound);
}
}
if (navigatorAttendance !== undefined && navigatorWorkContext !== undefined && navigatorWorkContext.contextError === undefined) {
// Flagged roles already have a concrete packet/task/case/review input.
// A bare Judge (and other bare packaged entrypoint) gets its concrete
// user task at this seam; do not copy the assembled system prompt.
// Replacement is keyed by typed subject provenance, never prose prefixes.
// Soft session_start placeholders (no materials yet) recover here; hard
// contextError from true loader failures stays poisoned and honest.
if (navigatorWorkContext.subjectProvenance === "placeholder") {
// Navigator's own input bootstrap, not a rewrite of another role's
// output payload (#836 targets the latter): the human's own raw
// prompt bytes, copied verbatim into Navigator's work context when
// nothing else has supplied a subject yet.
const subject = prompt.trim();
if (subject !== "") {
const root = subjectPath(ctx.sessionManager.getSessionDir(), ctx.cwd);
const subjectProvenance = "user_prompt" satisfies NavigatorSubjectProvenance;
const priorAuthority = navigatorWorkContext.authority;
const authority = typeof priorAuthority === "string" && priorAuthority.trim() !== ""
? priorAuthority
: subject;
navigatorWorkContext = {
subjectKey: navigatorSubjectKey(root, subject, subjectProvenance),
subject,
authority,
subjectProvenance,
};
await navigatorAttendance.setWorkContext(navigatorWorkContext);
}
}
}
navigatorAttendance?.prepare();
// Reviewer parent prompt assembly.
if (activeReviewerParent !== undefined && selectedRole === role) {
return {
systemPrompt: assembleReviewerParentSystemPrompt({
baseSystemPrompt: event.systemPrompt,
soul: activeReviewerParent.soul,
}),
};
}
// #676 E / J1: collector materials share this envelope hook (no parallel register).
if (activeCollector !== undefined && selectedRole === role) {
return {
systemPrompt: collectorBusiness.assembleMaterials(activeCollector, event.systemPrompt),
};
}
});
// #879: engine coordinates ride readingMaterial only when station-child
// officer dialogue replaced the ordinary transport prompt (no double fold).
roleHost.on("before_agent_start", () => {
const role = selectedRole ?? roleHost.getFlag(ROLE_FLAG.name);
if (typeof role !== "string" || !isOfficerReviewSeat(role)) return;
if (roleHost.getFlag(STATION_CHILD_FLAG.name) !== true) return;
const engine = resolveEngineName((name) => roleHost.getFlag(name));
if (engine === undefined || dependencies.packageRoot === undefined) return;
const engineModel = resolveEngineModel((name) => roleHost.getFlag(name));
const material = engineSessionMaterialFromOptions({
engine,
...(engineModel === undefined ? {} : { engineModel }),
packageRoot: dependencies.packageRoot,
});
if (material === undefined) return;
return {
readingMaterial: {
kind: "engine-session-material" as const,
name: material.name,
...(material.model === undefined ? {} : { model: material.model }),
...(material.materialPath === undefined ? {} : { materialPath: material.materialPath }),
},
};
});
roleHost.on("before_agent_start", () => {
const path = roleHost.getFlag(INSPECTOR_SOURCE_RUN_FLAG.name);
if (typeof path !== "string" || path.trim() === "") return;
return {
readingMaterial: {
kind: "inspector-parent-binding" as const,
sourceRunPath: path,
},
};
});
roleHost.on("tool_result", async (event, ctx) => {
const role = selectedRole;
if (role === undefined) return;
const pendingInfra = pendingInfrastructureFailures.get(event.toolCallId);
const isRoleInfrastructureFailure = pendingInfra !== undefined;
if (pendingInfra !== undefined) pendingInfrastructureFailures.delete(event.toolCallId);
// One-shot project fact + typed evidence so live settlement and durable session agree.
const infrastructureDetails = pendingInfra?.details;
const classified = infrastructureDetails === undefined
? event
: { ...event, details: infrastructureDetails };
const isOutputTool = event.toolName === navigatorOutputTool(role);
const outputClassification = isOutputTool ? classifyPackagedRoleTerminalResult(classified) : undefined;
if (isRoleInfrastructureFailure || outputClassification?.kind === "infrastructure") {
receiptDelivery.stopForInfrastructure();
} else if (isOutputTool) {
const reason = packagedRoleOutputRejectionReason(classified);
if (reason !== undefined) {
receiptDelivery.recordRejected(reason);
// Bind the observed rejection now, before agent_end or a host failure.
// This observation is not an exhausted lifecycle or a budget spend.
if (ctx.invocationScopeId !== undefined) {
envelopeHost.appendEntry(RECEIPT_REJECTION_ENTRY, {
invocationScopeId: ctx.invocationScopeId,
reason,
});
}
}
}
// Accepted/human terminal projection belongs exclusively to typed ledger
// closure. tool_result retains only infrastructure settlement.
const settlement = isRoleInfrastructureFailure || outputClassification?.kind === "infrastructure"
? publicNavigatorSettlement(role, navigatorPhase(roleHost, role), classified)
: undefined;
await settleNavigatorProjection(settlement);
// Persist typed infrastructure-failure fact onto the role session toolResult so
// exact-session restart shares the same durable completion classification.
if (infrastructureDetails !== undefined) {
return { isError: true };
}
// Submission non-pass: throw kept message text for the model; project the
// envelope-bound structured result onto session details at this tool_result seam.
const submissionNonPass = pendingSubmissionNonPassByToolCallId.get(event.toolCallId);
if (submissionNonPass !== undefined) {
pendingSubmissionNonPassByToolCallId.delete(event.toolCallId);
return { details: submissionNonPass, isError: true };
}
});
// Queue receipt delivery before `agent_settled`: that event means Pi has
// already decided no queued continuation will run, so a triggerTurn there is
// too late for print/json sessions. `agent_end` is the last production seam
// whose queued next turn is consumed before settlement.
roleHost.on("agent_end", async (event, ctx) => {
const lastMessage = event.messages.at(-1);
if (lastMessage?.role === "assistant"
&& (lastMessage.stopReason === "error" || lastMessage.stopReason === "aborted")) {
// Abort after an already-recorded receipt must not un-accept or催交.
if (receiptDelivery.nextAction() !== "accepted") {
receiptDelivery.stopForInfrastructure();
}
return;
}
const role = selectedRole ?? roleHost.getFlag(ROLE_FLAG.name);
// #959: navigator prose exit — final assistant text is the receipt.
// No typed-tool 催交; no JSON required. Tool path still wins when already accepted.
if (
typeof role === "string"
&& packagedRoleOutputTool(role) === NAVIGATOR_TOOL_SPEC.name
&& receiptDelivery.nextAction() !== "accepted"
) {
const prose = lastAssistantProse(event.messages);
if (prose !== undefined && prose.trim() !== "") {
const accepted = { prose };
await sealAcceptedSubmission({
context: ctx,
role: "navigator",
accepted,
toolCallId: `navigator-prose-exit:${randomUUID()}`,
});
await projectClosedSubmission(
{ role: "navigator", kind: "accepted", accepted },
ctx,
);
return;
}
// Attended with neither tool nor prose → honest no_receipt (not typed 催交).
// Close the budget without sending prompts; the fact stays the real count.
if (!noReceiptRecorded) {
const runPointer = runDirectoryFromHostContext(ctx);
if (runPointer !== undefined) {
noReceiptRecorded = true;
receiptDelivery.closeBudget();
const facts = receiptDelivery.facts({
runPointer,
attemptPointer: receiptAttemptPointer(runPointer, ctx.invocationScopeId),
});
envelopeHost.appendEntry(NO_RECEIPT_LIFECYCLE_ENTRY_TYPE, facts);
try {
sitianReport({
level: "event",
kind: "no-receipt-lifecycle",
cwd: ctx.cwd,
sessionParent: ctx.sessionManager.getSessionFile(),
payload: facts,
source: "role-runtime",
});
} catch {}
}
}
return;
}
if (receiptDelivery.nextAction() === "request-delivery") {
receiptDelivery.recordDeliveryRequest();
// Keep the package-owned continuation off the public input lifecycle:
// receipt delivery must not be mistaken for later caller input.
const scopeId = ctx.invocationScopeId?.trim() ?? "";
envelopeHost.appendEntry(
RECEIPT_DELIVERY_REQUEST_ENTRY,
scopeId.length > 0 ? { invocationScopeId: scopeId } : undefined,
);
try {
sitianReport({
level: "event",
kind: "receipt-delivery",
cwd: ctx.cwd,
sessionParent: ctx.sessionManager.getSessionFile(),
payload: { type: "ak-receipt-delivery-request" },
source: "role-runtime",
});
} catch {}
envelopeHost.sendMessage({
customType: "ak-receipt-delivery-prompt",
content: JSON.stringify(receiptDelivery.deliveryState()),
display: false,
}, { triggerTurn: true, deliverAs: "followUp" });
} else if (receiptDelivery.nextAction() === "no-receipt" && !noReceiptRecorded) {
const runPointer = runDirectoryFromHostContext(ctx);
if (runPointer !== undefined) {
noReceiptRecorded = true;
const facts = receiptDelivery.facts({
runPointer,
attemptPointer: receiptAttemptPointer(runPointer, ctx.invocationScopeId),
});
envelopeHost.appendEntry(NO_RECEIPT_LIFECYCLE_ENTRY_TYPE, facts);
try {
sitianReport({
level: "event",
kind: "no-receipt-lifecycle",
cwd: ctx.cwd,
sessionParent: ctx.sessionManager.getSessionFile(),
payload: facts,
source: "role-runtime",
});
} catch {}
}
}
});
roleHost.on("agent_settled", async () => {
if (pendingNavigatorSettlement !== undefined) {
await pendingNavigatorSettlement;
}
// Do not auto-drain Navigator on ordinary mid-turn agent_settled.
// Multi-turn roles (#162 coder) fire agent_settled after every tool turn;
// discarding a healthy/in-flight prepare there forces a cold prepare on the
// final accepted terminal and races the post-role grace. Keep preparation
// until an accepted/human/infrastructure tool_result settlement (or dispose).
// Provider death without any role tool_result leaves preparation held; the
// next explicit settlement or session end owns it — not mid-turn churn.
pendingNavigatorSettlement = undefined;
const presentation = pendingNavigatorPresentation;
pendingNavigatorPresentation = undefined;
if (presentation === undefined) return;
await envelopeHost.sendMessage({
customType: NAVIGATOR_EVENT_TYPE,
content: formatNavigatorReport(presentation.report),
display: true,
details: presentation.event,
}, { triggerTurn: false });
});
roleHost.on("session_shutdown", async () => {
// #351: stop OAuth keepalive first so shutdown yields zero further ticks.
envelopeHost.stopKeepalive();
// Flush any still-pending affirmative attendance before teardown.
// Grace-timeout paths normally emit on agent_settled; abort can skip that hook.
const presentation = pendingNavigatorPresentation;
pendingNavigatorPresentation = undefined;
try {
if (presentation !== undefined) {
// Required attendance persistence must reach its failure owner, not be
// discarded as cleanup or merged with an existing host report.
await envelopeHost.sendMessage({
customType: NAVIGATOR_EVENT_TYPE,
content: formatNavigatorReport(presentation.report),
display: true,
details: presentation.event,
}, { triggerTurn: false });
}
} finally {
// Same non-blocking dispose as post-role grace: awaiting here re-blocked the
// parent court for 43–270s after unavailable was already projected (#959 reopen).
const attendanceToDispose = navigatorAttendance;
navigatorAttendance = undefined;
pendingNavigatorSettlement = undefined;
disposeNavigatorAttendanceNonBlocking(attendanceToDispose);
pendingInfrastructureFailures.clear();
pendingSubmissionNonPassByToolCallId.clear();
}
});
const hostActions = {
failInfrastructure(error: unknown, ctx: HostContext, toolCallId?: string): never {
if (toolCallId !== undefined) {
pendingInfrastructureFailures.set(toolCallId, buildPendingInfrastructureFailure(error));
}
failInfrastructure(error, ctx);
},
bindSubmissionNonPass(toolCallId: string, result: SubmissionGateNonPassResult): void {
pendingSubmissionNonPassByToolCallId.set(toolCallId, result);
},
};
const requireRoleSoul = (role: PackagedRole): Promise => dependencies.loadRoleSoul(role);
const judge = createJudgeRoleRuntime(
roleHost,
{
loadSoul: () => requireRoleSoul("judge"),
},
);
const fixer = createFixerRoleRuntime(
roleHost,
{
loadSoul: () => requireRoleSoul("fixer"),
async loadPacket(path) {
if (dependencies.loadFixPacket === undefined) {
throw new Error("Fixer packet loader is not configured");
}
return dependencies.loadFixPacket(path);
},
},
hostActions,
{ unfinishedReasonBounceLimit: deliveryLimit },
);
const coder = createCoderRoleRuntime(
roleHost,
{
loadSoul: () => requireRoleSoul("coder"),
async loadTask(path) {
if (dependencies.loadCoderTask === undefined) {
throw new Error("Coder task loader is not configured");
}
return dependencies.loadCoderTask(path);
},
},
hostActions,
{ unfinishedReasonBounceLimit: deliveryLimit },
);
const reviewer = createReviewerRoleRuntime(
roleHost,
{
loadSoul: () => requireRoleSoul("reviewer"),
},
hostActions,
);
const doctor = createDoctorRoleRuntime(roleHost, {
loadSoul: () => requireRoleSoul("doctor"),
async loadCase(path) { if (!dependencies.loadDoctorCase) throw new Error("Doctor runtime dependencies are not configured"); return dependencies.loadDoctorCase(path); },
}, hostActions);
const notary = createNotaryRoleRuntime(roleHost, {
loadSoul: () => requireRoleSoul("notary"),
async loadSourceRunLocator(path) {
if (!dependencies.loadNotarySourceRun) throw new Error("Notary runtime dependencies are not configured");
return dependencies.loadNotarySourceRun(path);
},
}, hostActions);
const countersign = createCountersignRoleRuntime(roleHost, {
loadSoul: () => requireRoleSoul("countersign"),
});
const gleanerLeftRuntime = createGleanerLeftRoleRuntime(roleHost, {
loadSoul: () => requireRoleSoul("gleaner-left"),
});
const gleanerLeft = {
async activate() {
decodeGleanerLeftBase((name) => envelopeHost.host.getFlag(name));
return gleanerLeftRuntime.activate();
},
};
const inspector = createInspectorRoleRuntime(roleHost, {
loadSoul: () => requireRoleSoul("inspector"),
});
const gatekeeper = createGatekeeperRoleRuntime(roleHost, {
loadSoul: () => requireRoleSoul("gatekeeper"),
});
const navigator = createNavigatorRoleRuntime(roleHost, {
loadSoul: () => requireRoleSoul("navigator"),
recordRoutePlaybookReadFailure: (message) => {
envelopeHost.appendEntry(NAVIGATOR_ROUTE_PLAYBOOK_FAILURE_ENTRY, {
message: message ?? "",
});
},
});
const auditor = createAuditorRoleRuntime(roleHost, {
loadSoul: () => requireRoleSoul("auditor"),
});
const diarist = createDiaristRoleRuntime(
roleHost,
{
loadSoul: () => requireRoleSoul("diarist"),
},
);
const secretariat = createSecretariatRoleRuntime(roleHost, {
loadSoul: () => requireRoleSoul("secretariat"),
});
const merger = createMergerRoleRuntime(roleHost, {
loadSoul: () => requireRoleSoul("merger"),
async loadInput(path) { if (!dependencies.loadMergerInput) throw new Error("Merger runtime dependencies are not configured"); return dependencies.loadMergerInput(path); },
});
// #676 E / #1088: shared envelope owns collector lifecycle (mode/fork, host
// tool surface). Role supplies soul/materials + submission tool only; code
// no longer observes or merges GitHub findings.
let activeCollector: CollectorActivation | undefined;
const collectorBusiness = createCollectorRoleRuntime(
roleHost,
{
loadSoul: () => requireRoleSoul("collector"),
},
);
let collectorToolCallRegistered = false;
const collector = {
async activate(context: HostContext, event: { reason: string }) {
activeCollector = undefined;
if (context.mode !== "print" && context.mode !== "json") {
throw new Error(
`Collector supports only print or json mode (got ${context.mode})`,
);
}
if (event.reason === "fork" || event.reason === "reload") {
throw new Error(
`Collector does not support session_start reason ${event.reason}`,
);
}
const preExisting = roleHost.getAllTools();
const alreadyRegistered = COLLECTOR_REQUIRED_TOOLS.every((required) =>
preExisting.some((tool) => tool.name === required),
);
if (!alreadyRegistered) {
for (const required of COLLECTOR_REQUIRED_TOOLS) {
const prior = preExisting.filter((tool) => tool.name === required);
if (prior.length > 0) {
throw new Error(`Collector required tool name collision: ${required}`);
}
}
}
collectorBusiness.registerBusinessTools(() => activeCollector);
const allTools = roleHost.getAllTools();
for (const required of COLLECTOR_REQUIRED_TOOLS) {
const matches = allTools.filter((tool) => tool.name === required);
if (matches.length === 0) {
throw new Error(`Collector required tool missing: ${required}`);
}
if (matches.length > 1) {
throw new Error(`Collector required tool name collision: ${required}`);
}
}
// #1088: keep host CLI surface for self-collect; drop construction write/edit (ADR 0064).
const priorActive = roleHost.getActiveTools();
const hostSurface = (priorActive.length > 0 ? priorActive : allTools.map((tool) => tool.name))
.filter((name) => !(COLLECTOR_CONSTRUCTION_TOOLS as readonly string[]).includes(name));
const nextActive = [...new Set([...hostSurface, ...COLLECTOR_REQUIRED_TOOLS])];
roleHost.setActiveTools(nextActive);
const active = new Set(roleHost.getActiveTools());
for (const required of COLLECTOR_REQUIRED_TOOLS) {
if (!active.has(required)) {
throw new Error(`Collector failed to activate required tool ${required}`);
}
}
for (const name of COLLECTOR_CONSTRUCTION_TOOLS) {
if (active.has(name)) {
throw new Error(`Collector must not activate construction tool ${name}`);
}
}
if (!collectorToolCallRegistered) {
collectorToolCallRegistered = true;
roleHost.on("tool_call", (toolEvent) => {
const liveRole = roleHost.getFlag(ROLE_FLAG.name);
if (activeCollector === undefined || selectedRole === undefined || selectedRole !== liveRole) return;
return collectorBusiness.onToolCall(activeCollector, toolEvent);
});
}
activeCollector = await collectorBusiness.activate(context);
},
};
const clock = dependencies.activationClock ?? (() => new Date().toISOString());
const writeTrace = dependencies.activationTraceWriter ?? writeActivationTraceRecord;
roleHost.on("session_start", async (event, ctx) => {
admitted = false;
selectedRole = undefined;
roleReferenceMaterials = "";
activeReviewerParent = undefined;
activeCollector = undefined;
receiptDelivery = createReceiptDeliveryPolicy(deliveryLimit);
noReceiptRecorded = false;
pendingNavigatorPresentation = undefined;
navigatorActivation += 1;
navigatorDeliveryClosed = false;
pendingNavigatorSettlement = undefined;
pendingInfrastructureFailures.clear();
pendingSubmissionNonPassByToolCallId.clear();
navigatorWorkContext = undefined;
// #351: OAuth keepalive is orthogonal to --ak-role; start before role early-return
// so role-less sessions (and reload after shutdown stop) still keep tokens alive.
envelopeHost.startKeepalive(ctx);
const rawRole = roleHost.getFlag(ROLE_FLAG.name);
if (rawRole === undefined) return;
const entry = PACKAGED_ROLE_REGISTRY.find(({ role }) => role === rawRole);
if (entry === undefined) {
failInfrastructure(new Error(`Unsupported workflow role: ${String(rawRole)}`), ctx);
}
selectedRole = entry.role;
await navigatorAttendance?.dispose();
navigatorAttendance = undefined;
// One activation call per seat. Stage id still comes from the registry.
// Envelope owns Reviewer and Notary flag reads (ADR 0018).
const activateByRole = {
judge: () => judge.activate(),
fixer: () => fixer.activate(),
coder: () => coder.activate(ctx),
reviewer: async () => {
const admitted = decodeReviewerAdmittedInputs((name) => roleHost.getFlag(name));
activeReviewerParent = await reviewer.activate(ctx, admitted);
},
collector: () => collector.activate(ctx, event),
doctor: () => doctor.activate(),
notary: () => {
const ticketNumber = readNotaryTicketFlag(roleHost.getFlag(NOTARY_TICKET_FLAG.name));
return notary.activate(ticketNumber === undefined ? undefined : { ticketNumber });
},
countersign: () => countersign.activate(),
"gleaner-left": () => gleanerLeft.activate(),
inspector: () => inspector.activate(),
gatekeeper: () => gatekeeper.activate(),
navigator: () => navigator.activate(),
auditor: () => auditor.activate(),
diarist: () => diarist.activate(),
secretariat: () => secretariat.activate(),
merger: () => merger.activate(),
} satisfies Record Promise>;
try {
// Production topology only (ADR 0048/0049): no test-only ledger hooks.
// Admit durable session file first so lifecycle getEntries/appendEntry are truthful:
// deferred SM materialization must not wipe an in-memory principal marker.
// #855: two-face waiting.jsonl write removed — fail-closed book-key + session only.
resolveBookKeyFromGit(ctx.cwd);
durableSessionPointer(ctx.sessionManager);
const scopeId = ctx.invocationScopeId?.trim() ?? "";
if (scopeId.length > 0) {
// Same public call, new process: keep sends already issued and the
// rejection facts. A different manual call has a different scope.
receiptDelivery.continueIssued(
priorReceiptContinuation(ctx.sessionManager.getEntries(), scopeId),
);
}
// Station children (court diarist, inner-gate summons) omit navigator sidecar (#840).
// Top-level public entry legs still attend automatically.
// Navigator seat never re-attaches itself — prepare turns already ARE the navigator
// public activation (prevents summonPublicRole navigator ↔ attendance recursion).
const isStationChild = roleHost.getFlag(STATION_CHILD_FLAG.name) === true;
if (
dependencies.createNavigatorAttendance !== undefined
&& entry.outputTool !== NAVIGATOR_TOOL_SPEC.name
&& !isStationChild
) {
navigatorSessionParent = ctx.sessionManager.getSessionFile();
navigatorCwd = ctx.cwd;
let work: NavigatorWorkContext;
let contextError: unknown;
if (dependencies.loadNavigatorWorkContext === undefined) {
const fallbackSubjectKey = subjectPath(ctx.sessionManager.getSessionDir(), ctx.cwd);
contextError = new Error("Navigator work context loader is not configured");
work = { subjectKey: fallbackSubjectKey, subject: `work subject: ${fallbackSubjectKey}`, authority: "", subjectProvenance: "placeholder" };
} else {
try {
work = await dependencies.loadNavigatorWorkContext({
context: ctx,
role: entry.role,
phase: navigatorPhase(roleHost, entry.role),
getFlag: (name) => roleHost.getFlag(name),
});
contextError = work.contextError;
} catch (error) {
// Contract: README.md#Navigator-attendance — a failed context load continues with a typed placeholder work context; the original cause is retained in contextError for the typed unavailable report.
contextError = navigatorUnavailableError("context", error);
const fallbackSubjectKey = subjectPath(ctx.sessionManager.getSessionDir(), ctx.cwd);
work = { subjectKey: fallbackSubjectKey, subject: `work subject: ${fallbackSubjectKey}`, authority: "", subjectProvenance: "placeholder" };
}
}
navigatorWorkContext = { ...work, ...(contextError === undefined ? {} : { contextError }) };
// Shared envelope owns exact invocation principal from admitted session lifecycle.
// session_start is process activation: resume unfinished principal only when marker
// role/phase/subjectKey still match; mint for contradictory marker, malformed nearest,
// missing marker, or terminal already completed.
const sessionEntries = [...ctx.sessionManager.getEntries()];
const invocationPhase = navigatorPhase(roleHost, entry.role);
const lifecyclePrincipal = resolveLifecycleInvocationPrincipal(sessionEntries, {
role: entry.role,
phase: invocationPhase,
subjectKey: work.subjectKey,
});
const invocationId = lifecyclePrincipal.invocationId;
if (!lifecyclePrincipal.resume) {
const data = {
invocationId,
role: entry.role,
phase: invocationPhase,
subjectKey: work.subjectKey,
};
envelopeHost.appendEntry(NAVIGATOR_INVOCATION_ENTRY, data);
try {
sitianReport({
level: "event",
kind: "attendance",
cwd: ctx.cwd,
sessionParent: ctx.sessionManager.getSessionFile(),
payload: { type: NAVIGATOR_INVOCATION_ENTRY, ...data },
source: "role-runtime",
});
} catch {}
}
const activation = navigatorActivation;
navigatorAttendance = await dependencies.createNavigatorAttendance({
context: ctx,
role: entry.role,
phase: invocationPhase,
subjectKey: work.subjectKey,
subject: work.subject,
authority: work.authority,
invocationId,
deliveryRequestLimit: deliveryLimit,
...(contextError === undefined ? {} : { contextError }),
onEvent: (navigatorEvent, report) => {
if (activation === navigatorActivation && !navigatorDeliveryClosed) {
pendingNavigatorPresentation = { event: navigatorEvent, report };
}
},
});
// Concrete work context starts standby attendance (record only, no model).
// Placeholder subjects wait for before_agent_start (user prompt may replace the subject key).
if (
navigatorWorkContext.contextError === undefined &&
navigatorWorkContext.subjectProvenance !== "placeholder"
) {
navigatorAttendance.prepare();
}
}
await executeActivationStage(entry.role, activationStage(entry.role, activateByRole), { clock, writeTrace });
roleReferenceMaterials = await dependencies.loadRoleReferenceMaterials?.(entry.role) ?? "";
// #357 T2 / #378 / #380 / #391 / #818: any role+engine activation registers the package detour tool once.
// Gate is resolveEngineName (RoleHost flag → env fallback) — no per-engine execute branch; no role-module spawn.
if (!engineDetourRegistered) {
engineDetourRegistered = registerEngineDetourTool(roleHost, hostActions);
}
// Secretariat owns a declared active surface. The shared engine detour is
// registered after role activation, so include it here rather than leave
// a newly registered optional tool unreachable until a later reload.
if (secretariat.claimsEngineDetour === true && engineDetourRegistered) {
roleHost.setActiveTools([
...new Set([...roleHost.getActiveTools(), ENGINE_DETOUR_TOOL_NAME]),
]);
}
// Worker gates ①②: registry declares the worker seats (ADR 0066 / 0070).
// Parent session feeds #216 createRecordSession so baseline/bounce survive resume.
if ("worker" in entry && entry.worker === true) {
const workerArms = { coder, fixer } as const;
workerArms[entry.role].armSubmissionGate(ctx.cwd, ctx.sessionManager, ctx.invocationScopeId);
}
admitted = true;
} catch (error) {
failInfrastructure(error, ctx);
}
});
};
}