/** * Collector business tools, soul/material assembly, and result projection. * Registration flags, activation barrier, resource assembly (no-skills / no-context), * sole-final seal, mode/fork gates, tool surface narrowing, and failInfrastructure * exit live on the shared envelope (ADR 0018 / #676 E). Role keeps materials, * business tools, ledger facts, and projection only. */ import type { RoleHost, HostContext, HostToolResult } from "./host-contracts.ts"; import type { Static } from "typebox"; import { emptyCollectorManifest, loadCollectorManifest, parseCollectorPrNumber, parseCollectorRepository, type CollectorManifest, type CollectorRepository, } from "./collector-config.ts"; import { COLLECTOR_DEFAULT_WAIT_WINDOW_MS, createSystemCollectorClock, type CollectorClock, } from "./collector-evidence.ts"; import { createGhApiRunner, listPullRequestNumbersByTicket, type CollectorGitHubTransport, } from "./collector-github.ts"; import { createCollectorHandbookStore, resolveCollectorHandbookRoot, type CollectorHandbookStore, } from "./collector-handbook.ts"; import { COLLECTOR_BIND_TARGET_TOOL, COLLECTOR_HANDBOOK_WRITE_TOOL, COLLECTOR_OBSERVE_TOOL, COLLECTOR_OPEN_WAIT_WINDOW_TOOL, COLLECTOR_OUTPUT_TOOL, COLLECTOR_READ_TOOL, COLLECTOR_REQUEST_TOOL, COLLECTOR_WAIT_TOOL, projectEvidenceEntryView, type CollectorConfigState, type CollectorLedger, } from "./collector-ledger.ts"; import { collectorBindTargetArgsSchema, collectorHandbookWriteArgsSchema, collectorObserveArgsSchema, collectorOpenWaitWindowArgsSchema, collectorOutputArgsSchema, collectorReadArgsSchema, collectorRequestArgsSchema, collectorWaitArgsSchema, } from "./collector-tool-schemas.ts"; import { COLLECTOR_ACCEPTED_TEXT } from "./package-contracts/collector-output.ts"; import { CorrectableSubmissionError, isCorrectableExecuteError } from "./submission-correctable-error.ts"; import { CollectorUnknownEvidenceError } from "./collector-identity.ts"; export { CollectorUnknownEvidenceError } from "./collector-identity.ts"; /** #676 A: role-chosen target could not bind uniquely — correctable, not host failure. */ export class CollectorTargetBindError extends CorrectableSubmissionError { constructor(message: string) { super(message); this.name = "CollectorTargetBindError"; } } export { COLLECTOR_BIND_TARGET_TOOL, COLLECTOR_HANDBOOK_WRITE_TOOL, COLLECTOR_OBSERVE_TOOL, COLLECTOR_OPEN_WAIT_WINDOW_TOOL, COLLECTOR_OUTPUT_TOOL, COLLECTOR_READ_TOOL, COLLECTOR_REQUEST_TOOL, COLLECTOR_WAIT_TOOL, }; export const COLLECTOR_REQUIRED_TOOLS = [ COLLECTOR_BIND_TARGET_TOOL, COLLECTOR_OBSERVE_TOOL, COLLECTOR_READ_TOOL, COLLECTOR_REQUEST_TOOL, COLLECTOR_OPEN_WAIT_WINDOW_TOOL, COLLECTOR_WAIT_TOOL, COLLECTOR_HANDBOOK_WRITE_TOOL, COLLECTOR_OUTPUT_TOOL, ] as const; /** Envelope-owned transport flags (ADR 0018 / #676 E). */ export const COLLECTOR_TRANSPORT_FLAGS = Object.freeze([ Object.freeze({ name: "ak-collector-repo", definition: Object.freeze({ description: "GitHub owner/repo target for Collector (github.com only; conservative ASCII grammar).", type: "string" as const, }), }), Object.freeze({ name: "ak-collector-pr", definition: Object.freeze({ description: "Optional positive safe-integer pull request number for Collector. Omit when the role will bind from task materials.", type: "string" as const, }), }), Object.freeze({ name: "ak-collector-request-manifest", definition: Object.freeze({ description: "Path to the Collector v1 request manifest JSON file.", type: "string" as const, }), }), Object.freeze({ name: "ak-collector-wait-ms", definition: Object.freeze({ description: `Collector wait-window duration in milliseconds (default ${COLLECTOR_DEFAULT_WAIT_WINDOW_MS}). Caller-configurable; opens at a work step, not session start (#678 D4).`, type: "string" as const, }), }), ] as const); const observeSchema = collectorObserveArgsSchema; const readSchema = collectorReadArgsSchema; const requestSchema = collectorRequestArgsSchema; const openWaitWindowSchema = collectorOpenWaitWindowArgsSchema; const waitSchema = collectorWaitArgsSchema; const bindSchema = collectorBindTargetArgsSchema; const handbookWriteSchema = collectorHandbookWriteArgsSchema; const outputSchema = collectorOutputArgsSchema; type RequestParams = Static; type ReadParams = Static; type OpenWaitWindowParams = Static; type WaitParams = Static; type BindParams = Static; type HandbookWriteParams = Static; type OutputParams = Static; function parseWaitWindowMsFlag(raw: unknown): number { if (raw === undefined || raw === null || raw === "") return COLLECTOR_DEFAULT_WAIT_WINDOW_MS; if (typeof raw === "number" && Number.isSafeInteger(raw) && raw >= 1) return raw; if (typeof raw === "string" && /^[1-9]\d*$/.test(raw.trim())) { const value = Number(raw.trim()); if (Number.isSafeInteger(value) && value >= 1) return value; } throw new Error("Collector --ak-collector-wait-ms must be a positive safe-integer millisecond string"); } export type CollectorRoleDependencies = { loadSoul(): Promise; createTransport(): CollectorGitHubTransport; createClock?(): CollectorClock; createLedger(config: CollectorConfigState, clock: CollectorClock, ctx: HostContext): CollectorLedger; /** Optional packaged seed for first-use general handbook (#677). */ loadHandbookSeed?(): Promise; }; export type CollectorRoleHostActions = { failInfrastructure(error: unknown, ctx: HostContext, toolCallId?: string): never; }; export type CollectorActivation = { soul: string; repository: CollectorRepository; manifest: CollectorManifest; ledger: CollectorLedger; transport: CollectorGitHubTransport; clock: CollectorClock; handbook: CollectorHandbookStore; /** Loaded at activate; refreshed after handbook write for same-session materials. */ handbookView: { general: string; repo: string; generalSource: "book" | "seed" | "empty"; repoSource: "book" | "empty"; generalPath: string; repoPath: string; }; }; function buildMethodContext(activation: CollectorActivation): string { const pr = activation.ledger.config.prNumber; const handbook = activation.handbookView; const lines = [ "", `host: github.com`, `repository: ${activation.repository.canonical}`, `prNumber: ${pr === undefined ? "未绑定" : String(pr)}`, `waitWindowMs: ${String(activation.ledger.config.waitWindowMs)}`, `requests: ${JSON.stringify(activation.manifest.requests.map((request) => ({ id: request.id })))}`, `handbookGeneralSource: ${handbook.generalSource}`, `handbookRepoSource: ${handbook.repoSource}`, "", ]; // Opaque working memory only — no directional instructions (ADR 0073). // JSON + `<` → \u003c keeps bodies from forging the delivery close tag (#677). if (handbook.general.length > 0 || handbook.repo.length > 0) { const payload = JSON.stringify({ general: handbook.general, repo: handbook.repo, generalSource: handbook.generalSource, repoSource: handbook.repoSource, repository: activation.repository.canonical, }).replaceAll("<", "\\u003c"); lines.push("", "", payload, ""); } return lines.join("\n"); } function parsePositiveTicket(raw: unknown, label: string): number | undefined { if (raw === undefined || raw === null) return undefined; if (typeof raw === "number" && Number.isSafeInteger(raw) && raw >= 1) return raw; if (typeof raw === "string" && /^[1-9]\d*$/.test(raw.trim())) return Number(raw.trim()); throw new CollectorTargetBindError(`通进司绑定 ${label} 须为正安全整数`); } /** * Shared envelope installs business tools and material callback after lifecycle gates. * Role module does not self-hook events, setActiveTools, or mode/fork checks (ADR 0018 / #676 E). */ export function createCollectorRoleRuntime( pi: RoleHost, dependencies: CollectorRoleDependencies, hostActions: CollectorRoleHostActions, ): { /** Business activation only: soul, flags→config, ledger. Envelope owns lifecycle gates. */ activate(ctx: HostContext): Promise; /** Business material assembly — envelope hangs this on before_agent_start. */ assembleMaterials(activation: CollectorActivation, baseSystemPrompt: string): string; /** Business operational bookkeeping callbacks for shared tool_call / tool_result seams. */ onToolCall( activation: CollectorActivation, event: { toolName: string; toolCallId: string }, ): { block: true; reason: string } | undefined; onToolResult(activation: CollectorActivation, event: { toolCallId: string }): void; /** Register business tools once. Envelope owns uniqueness + setActiveTools. */ registerBusinessTools(getActivation: () => CollectorActivation | undefined): void; } { let toolsRegistered = false; return { async activate(ctx) { const soul = (await dependencies.loadSoul()).trim(); if (soul.length === 0) throw new Error("Collector soul is empty"); // Flags registered by the shared envelope (COLLECTOR_TRANSPORT_FLAGS). const repoFlag = pi.getFlag("ak-collector-repo"); const prFlag = pi.getFlag("ak-collector-pr"); const requestManifestFlag = pi.getFlag("ak-collector-request-manifest"); const waitMsFlag = pi.getFlag("ak-collector-wait-ms"); if (typeof repoFlag !== "string" || repoFlag.trim().length === 0) { throw new Error("Collector requires --ak-collector-repo"); } const repository = parseCollectorRepository(repoFlag); // #676 A: PR may be unbound — role binds via ak_collector_bind_target from materials. let prNumber: number | undefined; if (typeof prFlag === "string" && prFlag.trim().length > 0) { prNumber = parseCollectorPrNumber(prFlag); } else if (typeof prFlag === "number") { prNumber = parseCollectorPrNumber(prFlag); } const manifest = typeof requestManifestFlag === "string" && requestManifestFlag.trim().length > 0 ? await loadCollectorManifest(requestManifestFlag) : emptyCollectorManifest(); const waitWindowMs = parseWaitWindowMsFlag(waitMsFlag); const clock = dependencies.createClock?.() ?? createSystemCollectorClock(); const transport = dependencies.createTransport(); const ledger = dependencies.createLedger( { repository, prNumber, manifest, waitWindowMs }, clock, ctx, ); // #677: handbook under admitted book topology (session path → books//collector-handbook). const sessionPath = ctx.sessionManager?.getSessionFile?.() ?? ctx.sessionManager?.getSessionDir?.(); if (typeof sessionPath !== "string" || sessionPath.length === 0) { throw new Error("通进司手册要求 session 路径位于 books//"); } const placement = resolveCollectorHandbookRoot(sessionPath); const seedGeneral = dependencies.loadHandbookSeed === undefined ? undefined : (await dependencies.loadHandbookSeed()).trim(); const handbook = createCollectorHandbookStore({ ledgerHome: placement.ledgerHome, handbookRoot: placement.root, repositoryCanonical: repository.canonical, ...(seedGeneral === undefined || seedGeneral.length === 0 ? {} : { seedGeneral }), }); const handbookView = await handbook.read(); return { soul, repository, manifest, ledger, transport, clock, handbook, handbookView, }; }, assembleMaterials(activation, baseSystemPrompt) { return [ baseSystemPrompt, "", "", activation.soul, "", "", buildMethodContext(activation), ].join("\n"); }, onToolCall(activation, event) { // Business ledger facts only — allowed-tool / mode gates live on the envelope. if (activation.ledger.fatal) { return { block: true, reason: activation.ledger.fatalReason ?? "通进司致命状态", }; } if (event.toolName === COLLECTOR_OUTPUT_TOOL) { activation.ledger.beginOperational(COLLECTOR_OUTPUT_TOOL, event.toolCallId); } return undefined; }, onToolResult(activation, event) { activation.ledger.completeOperational(event.toolCallId); }, registerBusinessTools(getActivation) { if (toolsRegistered) return; toolsRegistered = true; pi.registerTool({ name: COLLECTOR_BIND_TARGET_TOOL, label: "通进司认票绑定", description: "角色判定任务材料后绑定本仓唯一 PR 目标。可提交 prNumber 或 issueNumber(线上关联唯一 PR);显式 --pr 已绑定时无需再调。多义或无法确定时会正确驳回,要求调用方明确 --pr。", promptSnippet: "绑定角色判定的 issue/PR 目标", parameters: bindSchema, async execute(toolCallId: string, params: BindParams, _signal: AbortSignal | undefined, _onUpdate: unknown, ctx: HostContext) { const activation = getActivation(); if (activation === undefined) throw new Error("通进司未激活"); try { activation.ledger.beginOperational(COLLECTOR_BIND_TARGET_TOOL, toolCallId); const prNumber = parsePositiveTicket(params.prNumber, "prNumber"); const issueNumber = parsePositiveTicket(params.issueNumber, "issueNumber"); if (prNumber === undefined && issueNumber === undefined) { throw new CollectorTargetBindError( "通进司绑定须由角色判定 prNumber 与/或 issueNumber", ); } let bound = prNumber; let associated: readonly number[] = []; if (issueNumber !== undefined) { associated = await listPullRequestNumbersByTicket(createGhApiRunner(), { owner: activation.repository.owner, repo: activation.repository.repo, ticketNumber: issueNumber, }); if (bound === undefined && associated.length === 1) { bound = associated[0]!; } } if (bound === undefined) { return { content: [{ type: "text" as const, text: JSON.stringify({ bound: false, repository: activation.repository.canonical, prNumber: prNumber ?? null, issueNumber: issueNumber ?? null, associated, }), }], details: { bound: false, repository: activation.repository.canonical, prNumber: prNumber ?? null, issueNumber: issueNumber ?? null, associated, }, }; } activation.ledger.bindTarget(bound); activation.ledger.completeOperational(toolCallId); return { content: [{ type: "text" as const, text: `目标已绑定:${activation.repository.canonical}#${bound}`, }], details: { bound: true, repository: activation.repository.canonical, prNumber: bound, issueNumber: issueNumber ?? null, associated, }, }; } catch (error) { if (isCorrectableExecuteError(error)) throw error; hostActions.failInfrastructure(error, ctx, toolCallId); } finally { try { activation.ledger.completeOperational(toolCallId); } catch { // already completed or not begun } } }, }); pi.registerTool({ name: COLLECTOR_OBSERVE_TOOL, label: "通进司观察", description: "抓取配置目标的完整 GitHub PR 证据,存不可变快照入卷。正文在上下文中只给头部摘录加指针;需要头部之外的正文时,用 ak_collector_read 按 evidenceId 开卷;findings 的拆分与归类由你在交件时完成。目标未绑定前须先 ak_collector_bind_target。", promptSnippet: "抓取配置目标 PR 证据", parameters: observeSchema, async execute(toolCallId: string, _params: unknown, signal: AbortSignal | undefined, _onUpdate: unknown, ctx: HostContext) { const activation = getActivation(); if (activation === undefined) throw new Error("通进司未激活"); try { activation.ledger.beginOperational(COLLECTOR_OBSERVE_TOOL, toolCallId); const { snapshot, contextView } = await activation.ledger.observe( activation.transport, activation.clock, signal, ); activation.ledger.completeOperational(toolCallId); return { content: [{ type: "text" as const, text: JSON.stringify(contextView), }], details: contextView, }; } catch (error) { if (isCorrectableExecuteError(error)) throw error; hostActions.failInfrastructure(error, ctx, toolCallId); } }, }); pi.registerTool({ name: COLLECTOR_READ_TOOL, label: "通进司开卷", description: "按 evidenceId 开卷读取一条已观测材料的全量正文与指针;只在观察头部摘录不足以判读时调用。", promptSnippet: "按指针开卷读材料", parameters: readSchema, async execute(toolCallId: string, params: ReadParams, _signal: AbortSignal | undefined, _onUpdate: unknown, ctx: HostContext) { const activation = getActivation(); if (activation === undefined) throw new Error("通进司未激活"); try { activation.ledger.beginOperational(COLLECTOR_READ_TOOL, toolCallId); const record = activation.ledger.getEvidence(params.evidenceId); if ( record === undefined || (record.kind !== "review" && record.kind !== "issue_comment" && record.kind !== "review_comment" && record.kind !== "reaction") || typeof record.body !== "string" ) { throw new CollectorUnknownEvidenceError(params.evidenceId); } const material = projectEvidenceEntryView(record); activation.ledger.completeOperational(toolCallId); return { content: [{ type: "text" as const, text: JSON.stringify(material), }], details: material, }; } catch (error) { if (isCorrectableExecuteError(error)) throw error; hostActions.failInfrastructure(error, ctx, toolCallId); } }, }); pi.registerTool({ name: COLLECTOR_REQUEST_TOOL, label: "通进司请求", description: "在所引最新快照 HEAD 发一次请求。requestId 可取配置清单,或角色依手册/现场判定的稳定 id;后者须同时提供 body。", promptSnippet: "发一次评审请求", parameters: requestSchema, async execute(toolCallId: string, params: RequestParams, signal: AbortSignal | undefined, _onUpdate: unknown, ctx: HostContext) { const activation = getActivation(); if (activation === undefined) throw new Error("通进司未激活"); try { activation.ledger.beginOperational(COLLECTOR_REQUEST_TOOL, toolCallId); const details = await activation.ledger.request( { requestId: params.requestId, snapshotId: params.snapshotId, ...(typeof params.body === "string" ? { body: params.body } : {}), }, activation.transport, activation.clock, signal, ); activation.ledger.completeOperational(toolCallId); return { content: [{ type: "text" as const, text: `请求尝试已记录:request ${params.requestId}`, }], details, }; } catch (error) { if (isCorrectableExecuteError(error)) throw error; hostActions.failInfrastructure(error, ctx, toolCallId); } }, }); pi.registerTool({ name: COLLECTOR_HANDBOOK_WRITE_TOOL, label: "通进司手册写入", description: "写入通用手册或当前仓库差异全文(整份替换)。", promptSnippet: "更新 bot 手册", parameters: handbookWriteSchema, async execute(toolCallId: string, params: HandbookWriteParams, _signal: AbortSignal | undefined, _onUpdate: unknown, ctx: HostContext) { const activation = getActivation(); if (activation === undefined) throw new Error("通进司未激活"); try { activation.ledger.beginOperational(COLLECTOR_HANDBOOK_WRITE_TOOL, toolCallId); // scope/body shape authority = collectorHandbookWriteArgsSchema (host parameters). const details = await activation.handbook.write(params.scope, params.body); activation.handbookView = await activation.handbook.read(); activation.ledger.completeOperational(toolCallId); return { content: [{ type: "text" as const, text: `手册已写入:${details.scope}`, }], details: { scope: details.scope, path: details.path, byteLength: details.byteLength, generalSource: activation.handbookView.generalSource, repoSource: activation.handbookView.repoSource, }, }; } catch (error) { if (isCorrectableExecuteError(error)) throw error; hostActions.failInfrastructure(error, ctx, toolCallId); } finally { try { activation.ledger.completeOperational(toolCallId); } catch { // already completed or not begun } } }, }); pi.registerTool({ name: COLLECTOR_OPEN_WAIT_WINDOW_TOOL, label: "通进司开启等待窗", description: "在工作步骤上开启等待窗:新建并自动触发的 PR 传创建成功时刻;已有 PR 在触发阶段结束后省略 startedAt(=现在)。窗长由使用方配置,默认十分钟;开启后不因 PR 更新重置。", promptSnippet: "按工作步骤开启等待窗", parameters: openWaitWindowSchema, async execute(toolCallId: string, params: OpenWaitWindowParams, _signal: AbortSignal | undefined, _onUpdate: unknown, ctx: HostContext) { const activation = getActivation(); if (activation === undefined) throw new Error("通进司未激活"); try { activation.ledger.beginOperational(COLLECTOR_OPEN_WAIT_WINDOW_TOOL, toolCallId); let startedAt: Date | undefined; if (typeof params.startedAt === "string" && params.startedAt.trim().length > 0) { const parsed = new Date(params.startedAt); if (Number.isNaN(parsed.getTime())) { throw new Error("通进司等待窗 startedAt 须为有效 ISO 时间"); } startedAt = parsed; } activation.ledger.openWaitWindow( activation.clock, startedAt === undefined ? undefined : { startedAt }, ); activation.ledger.completeOperational(toolCallId); return { content: [{ type: "text" as const, text: "等待窗已开启", }], details: { activationTime: activation.ledger.activationTime?.toISOString(), deadlineTime: activation.ledger.deadlineTime?.toISOString(), waitWindowMs: activation.ledger.config.waitWindowMs, }, }; } catch (error) { if (isCorrectableExecuteError(error)) throw error; hostActions.failInfrastructure(error, ctx, toolCallId); } finally { try { activation.ledger.completeOperational(toolCallId); } catch { // already completed or not begun } } }, }); pi.registerTool({ name: COLLECTOR_WAIT_TOOL, label: "通进司等待", description: "再观察前等待;实际睡眠不超过剩余等待窗。", promptSnippet: "等待窗内等待", parameters: waitSchema, async execute(toolCallId: string, params: WaitParams, signal: AbortSignal | undefined, _onUpdate: unknown, ctx: HostContext) { const activation = getActivation(); if (activation === undefined) throw new Error("通进司未激活"); try { activation.ledger.beginOperational(COLLECTOR_WAIT_TOOL, toolCallId); const details = await activation.ledger.wait( params, activation.clock, signal, ); activation.ledger.completeOperational(toolCallId); return { content: [{ type: "text" as const, text: `已等待 ${String((details as { effectiveMs: number }).effectiveMs)}ms`, }], details, }; } catch (error) { if (isCorrectableExecuteError(error)) throw error; hostActions.failInfrastructure(error, ctx, toolCallId); } }, }); pi.registerTool({ name: COLLECTOR_OUTPUT_TOOL, label: "通进司输出", description: "观察完成后提交原交卷。正常完工提交空对象 {}(如需报 finding,填 findings 指针数组);仅在基础设施真实失败时才可填 infrastructureFailure,无失败时必须省略该字段。", promptSnippet: "提交通进司回执", parameters: outputSchema, async execute(toolCallId: string, params: OutputParams, _signal: AbortSignal | undefined, _onUpdate: unknown, ctx: HostContext) { const activation = getActivation(); if (activation === undefined) throw new Error("通进司未激活"); try { activation.ledger.beginOperational(COLLECTOR_OUTPUT_TOOL, toolCallId); activation.ledger.recordOutputCandidate(); return { content: [{ type: "text" as const, text: COLLECTOR_ACCEPTED_TEXT, }], details: params, terminate: true as const, }; } catch (error) { if ( error instanceof Error && (error as { collectorFatal?: boolean }).collectorFatal === true ) { hostActions.failInfrastructure(error, ctx, toolCallId); } throw error; } finally { activation.ledger.completeOperational(toolCallId); } }, }); }, }; }