import type Database from "better-sqlite3"; import type { InboundDelivery, TransportClient } from "../transport/types.js"; import { type AgentBackend } from "../agent/types.js"; import { type AgentBackendType } from "../config.js"; import { type UpgradeSafenessSource } from "./auto-update.js"; import { type StableSystemContextOptions } from "../memory/inject.js"; import { type ScheduleAgentCommand, type ScheduleAgentCommandResult } from "./schedule-command.js"; import { type GoalFinishCommand } from "./goal.js"; import type { BackendCapability } from "../agent/backend-capability.js"; import type { EngineLifecycle } from "../engine-lifecycle.js"; import { type CollabTurnDecision } from "./collab-loop.js"; export { resolveUpdateCommandCwd } from "../update-command.js"; export declare const SHELL_COMMAND_TIMEOUT_MS = 300000; /** Bot 身份信息,由外部传入 */ export interface BotIdentity { /** Bot 显示名称(如 "CowBot",从平台 API 获取或 config 指定) */ name: string; /** IM 平台标识(如 "feishu") */ platform: string; /** Bot 在平台上的唯一标识(用于 DB 中的 bot 用户记录);启动前尚未解析时为空。 */ platformBotId?: string; /** 此 Bot 过去使用过的本地占位 platform_id,只用于启动时清洗旧数据。 */ legacyPlatformBotIds?: readonly string[]; /** 主模型 ID(可选,覆盖 backend 默认值) */ model?: string; /** 推理强度运行时选择(low/medium/high/xhigh/max),backend 支持时生效 */ effort?: string; } export declare class Pipeline { private db; private transport; private agent; /** 引擎默认 backend(配置值)。`/agent` 只改当前 scope,不改这个默认。 */ private backendType; private backendResolver?; private backends; private backendLoads; private getAvailableBackends; private getBackendCapabilities; private queue; private responseSender; private runtimeState; private runManager; private engineLifecycle?; /** 自动升级是否启用:唯一来源是服务配置文件的 autoUpdate 布尔值。 */ private isAutoUpdateEnabled; private effectiveAutoUpdateConfig; private botIdentity; private log; /** 每个 chat 的当前 agent session */ private chatSessions; /** 异步创建中、尚未加入 chatSessions 的 session */ private sessionCreations; /** chatId → platformChatId 映射 */ private platformChatIds; /** chatId → userId 映射 */ private chatUserIds; /** bot 的内部用户 ID */ private botUserId; /** 每个 session scope 连续被 Bot 触发并跑 Agent 的回合数(人入站清零) */ private botTurnCounts; /** 当前进程里本群被 at 过的 Bot;补交棒时与数据库里的 Bot 发言合并。 */ private chatMentionedBotIds; /** admin 角色映射:userId → role */ private adminRoles; /** 运行中的独立 task(agentSession.id → RunningTask) */ private runningTasks; /** processIndependentSession 的活跃计数:覆盖从入口到清理完成的完整生命周期,供优雅关闭等待。 */ private independentRunCount; /** shell 命令历史(admin 专用) */ private shellHistory; private readonly MAX_SHELL_HISTORY; /** agent 工作目录 */ private workingDirectory; /** 数据库路径(传递给 agent 子进程) */ private dbPath; /** 已处理的消息 ID 去重集合(有上限防内存泄漏) */ private processedMsgIds; private static readonly MAX_PROCESSED_IDS; /** chatId → triggerPlatformMsgId,暂存触发消息 ID */ private triggerMsgIds; /** chatId → transition promise,session 切换期间后续消息先挂起 */ private sessionTransitionLocks; /** 全局过渡(重启等)阻塞所有 chat;/agent 只用 per-scope 锁 */ private globalSessionTransition?; /** chatId → transition 期间暂存的后续消息 */ private pendingTransitionMessages; /** 已加过 Pin 的消息,避免重复加 reaction */ private pinnedMsgIds; /** 已加过 Get 的消息,避免重复加 reaction */ private processingMsgIds; /** 协作消息的最终反应;同一 Bot 对同一条协议消息只保留一个表情。 */ private collabReactionEmojis; private collabReactionTasks; /** 每个 backend 的模型配置快照,切换时保存/恢复 */ private backendModelCache; /** Watchdog 定时器 */ private watchdogTimer; /** 已发送 compact 通知的 session(避免重复通知) */ private compactNotifiedSessions; /** chatId → 上次从 backend 看到的 compact 次数 */ private lastCompactCounts; /** chatId 集合:下一条发给 agent 的消息需要注入 compact 恢复提醒 */ private pendingCompactRecovery; /** 创建新 agent session 前刷新 AGENTS.md / CLAUDE.md 等静态上下文文件 */ private refreshAgentContextFiles?; private stableContextOptions; private archiveHome; private activeScheduleAgentCommands; /** 当前协作回合的能力令牌和动作;只存在于本进程当前 Agent 回合。 */ private activeCollabTurns; /** chatId → 主会话调度令牌:随主 Agent 环境注入,独立 session 拿不到,防跨进程身份借用。 */ private chatScheduleTokens; /** token → scopeKey:IPC 反查,避免平台 chat_id 串话题。 */ private tokenToScope; /** 正在占用主聊天队列的 Loop 回合;取消命令用它精确中止对应 run。 */ private activeLoopRuns; /** chatId → 未结束的 Goal(纯内存;重启即断)。 */ private activeGoals; constructor(db: Database.Database, transport: TransportClient, agent: AgentBackend, botIdentity: BotIdentity, workingDirectory: string, dbPath: string, bufferMs: number, backendType?: AgentBackendType, backendResolver?: (type: AgentBackendType) => Promise, getAvailableBackends?: () => string[], refreshAgentContextFiles?: () => void, stableContextOptions?: StableSystemContextOptions, _legacyRestartConfig?: unknown, _legacyAutoUpdateCoordinator?: boolean, archiveHome?: string, getBackendCapabilities?: () => BackendCapability[] | Promise, _legacyAutoUpdateConfig?: unknown, _legacyConfigPath?: string, _legacyOnAutoUpdateConfigChanged?: () => void, engineLifecycle?: EngineLifecycle); private createAgentSession; /** 启动 Engine 管道;Transport 入站入口由装配层连接。 */ start(): Promise; private resolveBotIdentity; private runStartupPlatformProbes; private requireBotPlatformId; /** 停止管道:清除队列计时器 */ stop(): void; /** 优雅关闭:cancel 所有活跃 session,清理资源(DB 中 session 保持 active,下次启动恢复) */ shutdown(): Promise; /** 是否有正在处理的主会话或独立 session(Cron/task)——优雅关闭等待用 */ hasBusyChats(): boolean; /** 检查用户是否为 admin 或 owner */ isAdmin(userId: string): boolean; /** 检查用户是否为 owner */ isOwner(userId: string): boolean; /** 调度写入口。身份取当前 Agent 回合,不信任 Session 创建时固定的环境变量。 */ executeScheduleAgentCommand(chatId: string, command: ScheduleAgentCommand, token?: string, scope?: { scopeKey?: string; threadId?: string; replyToMsgId?: string; }): Promise; /** * 协作 Agent 的唯一写入口。动作只写入当前回合的内存上下文, * 真实出站消息仍由 process() 统一发送。 */ executeCollabTurn(chatId: string, decision: CollabTurnDecision, token?: string, scope?: { scopeKey?: string; threadId?: string; replyToMsgId?: string; }): Promise<{ output: string; }>; /** 获取 bot 用户 ID */ getBotUserId(): string | null; private platformChatIdToScopeKey; private resolveOutboundScopeKey; /** IPC 发送:仅主 Agent 当前回合(schedule token 匹配)才挂到用户消息下 */ private currentTurnReplyTarget; /** 话题内最后一次平台消息 ID;跨进程 IPC 没有 schedule token 时作为回复锚点。 */ private latestThreadPlatformMsgId; private sendPreferringReply; /** 通过 IPC 发送消息到指定 chat */ sendToChat(platformChatId: string, text: string, scheduleToken?: string, scope?: { scopeKey?: string; threadId?: string; replyToMsgId?: string; }): Promise; /** 通过 IPC 发送卡片到指定 chat */ sendCardToChat(platformChatId: string, header: string, content: string, scheduleToken?: string, scope?: { scopeKey?: string; threadId?: string; replyToMsgId?: string; }): Promise; /** 卡片失败时,完整飞书 at 改走文本,避免对方 Bot 收不到。 */ private sendCardKeepingAt; private loadMentionUsers; private rememberMentionedBots; private chatBotMentionCandidates; private isBotMention; private isCurrentBotMention; private isCurrentBotPlatformId; /** * 协作状态保存的是消息中出现的参与者 ID。Bot 身份迁移后,这个 ID 可能是旧 ID; * 执行和动作校验必须继续使用同一个参与者 ID,不能混用当前的新 ID。 */ private currentCollabParticipantId; /** finish 时提醒协作启动消息的提问者;找不到根消息则宁可少 @,不猜身份。 */ private collabRequester; private firstBotMention; private backendSupportsCollab; private sameCollabProjection; private isCurrentCollabProjection; private hasCollabRuntimeEvent; private isCollabRunForState; /** 安装新的显式协作起点,并取消仍在执行的旧协作回合。 */ private installCollabStart; /** * 为当前入站消息准备协作状态。 * 启动消息和后续协议消息都会 @ 全部参与者;@ 列表第一个 Bot 才是唯一执行者。 */ private prepareCollabTurn; /** 活跃协作期间,所有人类消息都是补充材料;只有 /stop 才会打断协作。 */ private isActiveHumanCollabContext; /** 多 Bot 起点先确认当前 Bot 的线上接收权限,提示独立发送,不混入协作报告。 */ private warnAboutCollabBotAtPermission; /** 转发读取错误由适配器按外层消息合并,此处保留其回复锚点。 */ private warnAboutMessageReadError; /** * 取本 Bot 尚未见过、且发生在本次协作启动之后的人类补充消息。 * 每台设备各自持久化 seen 标记,下一棒无论在哪台设备都能看到同一批消息一次。 */ private buildCollabSupplementContext; /** * 飞书晚投递的协作事件不能照普通过期消息直接抛弃,也不能直接执行。 * 用话题历史重建后,只有这条消息仍是当前链最后一个有效触发点才允许入队。 */ private recoverLateCollabTurn; /** * 全员协议消息只更新投影;@ 列表首位等于自身时,才把本轮放入 Agent 队列。 */ private prepareProtocolCollabTurn; /** * 设备离线或刚启动时,本地没有投影也可以从飞书话题历史重建。 * 当前入站消息尚未入库,因此显式补在历史末尾参与重放。 */ private recoverProtocolCollabState; private applyProtocolCollabTurn; private toCollabHistoryMessage; private readCollabHistory; private extractCollabMentions; private blockCollabState; private stopCollabState; /** 参与者没有用全员协议交棒时,不能让旧协作投影继续占住话题。 */ private stopBrokenCollabFromBot; private prepareOutboundText; private sendPreparedFinalResponse; /** 内置命令卡片:有当前消息就 reply 到话题,始终不主动 create 新话题。 */ private sendBuiltinCard; private noteHumanInbound; private noteBotCollabTurn; /** 通过 IPC 发送文件到指定 chat */ sendFileToChat(platformChatId: string, filePath: string, scheduleToken?: string, scope?: { scopeKey?: string; threadId?: string; replyToMsgId?: string; }): Promise; private markQueuedMessage; private moveMessageToProcessing; /** * 协作反应按消息串行发送。完成优先于思考中和处理中;升级时先移除旧表情, * 避免同一 Bot 在同一轮留下多个状态标记。 */ private setCollabReaction; /** * 协作协议会 @ 全部参与 Bot。非当前轮次的 Bot 以思考中确认已同步, * 但不进入队列,也不表示正在执行。 */ private acknowledgeObservedCollabMessage; /** * 注入 prompt 到 agent pipeline(用于 cron 等内部触发场景)。 * 和用户消息走相同的 queue → process → agent 链路。 */ injectPrompt(chatId: string, userId: string, text: string): void; private markRuntimeRun; private persistRuntimeEvent; private markUnfinishedRuntimeRunsFailedByRestart; private ensureChatMetadata; private isStrictTopicChat; private isIsolatedTopicChat; /** 这个话题已经请过本 Bot:内存队列/session、带 session_key 的本 Bot 回复,或根帖已被本引擎处理。 */ private topicThreadBelongsToBot; private resolveMessageScope; /** Standard Transport entrypoint. Platform events are persisted before reaching this method. */ handleInbound(delivery: InboundDelivery): Promise; private handleMessage; /** Store a bot-sent message in DB */ private storeBotResponse; /** Store message without triggering agent (for group chat non-targeted messages) */ private storeMessageOnly; private persistInboundMessage; /** Check if a platform message was sent by the bot */ private isMessageFromBot; /** 被回复的那条:`正文`,找不到则空。 */ private buildReplyQuoted; /** Persist a user's admin role in both memory and DB */ private setAdminRole; /** Remove admin from both memory and DB (cannot remove owner) */ private removeAdmin; /** Restore admin users from local DB only. This must stay fast during startup. */ private restoreAdminsFromDb; /** Detect platform app creator in the background. */ private detectAppCreatorAdmin; /** * 内置命令拦截:匹配 /xxx 格式的消息,命中则直接处理并返回 true。 * //xxx 视为强制透传给 agent,本地不拦截。 * 未命中返回 false,消息继续走 agent 流程。 * * 分发顺序: * 1. 内置命令 switch(含 /loop、/cron 的查看、帮助和取消) * 2. 管理员 shell 命令(tryShellCommand) * 3. return false → 转发给 agent */ private handleBuiltinCommand; private isBuiltinCommand; private handleScheduleBuiltinCommand; private sendScheduleBuiltinList; private cancelRunningCronSessions; private cancelActiveLoopRun; /** //xxx 表示强制透传给 agent,实际发送时去掉一个前缀 / */ private normalizeUserTextForAgent; /** 回复文本:有 msgId 时引用回复,否则直接发送,并存入 DB */ private replyText; /** * /list:列出所有运行中的会话(主 session + 独立 task),含最近日志。 */ private sendRunningList; private buildLoopStatusSection; private stopAllTasks; private sendStatus; /** /status:与 watchdog 一致,取 agent activity.lastActiveAt */ private getLatestAgentOutputAt; /** /goal:创建或查询 Goal(纯内存,重启即断)。 */ /** /goal 内置命令:只处理查询(无参)。带参创建由 hybrid 翻译层转发给 Agent(nbt goal start)。 */ private handleGoalCommand; /** nbt goal finish:Agent 显式结束请求(令牌 + Run 一致性校验;三条件结算在回合收尾时做)。 */ executeGoalFinishCommand(chatId: string, command: GoalFinishCommand, scheduleToken?: string, scope?: { scopeKey?: string; threadId?: string; }): Promise<{ output: string; }>; /** nbt goal start:Agent 主动进入 Goal 模式。当前回合计入第 1 轮(process 检测 startRunId 后由 runGoalLoop 接管)。 */ executeGoalStartCommand(chatId: string, objective: string, scheduleToken?: string, scope?: { scopeKey?: string; threadId?: string; }): Promise<{ output: string; }>; /** * nbt goal progress:中间轮静默记录进展(不发送 IM)。 * content = 本次步骤(一两句话,保留最近 N 条);status = 全局进展状态(覆盖式:任务整体进行到哪、还剩什么)。 */ executeGoalProgressCommand(chatId: string, content: string, status?: string, scope?: { scopeKey?: string; threadId?: string; }): Promise<{ output: string; }>; /** nbt restart --wake:重启完成后注入主会话任务(在原上下文触发 Agent 回合)。 */ executeWakeCommand(chatId: string, prompt: string, scope?: { scopeKey?: string; threadId?: string; replyToMsgId?: string; }): Promise<{ output: string; }>; /** 构建 Goal 每轮注入的引导(目标原文 + 检查引导 + 进展:全局状态与最近步骤,防遗忘)。 */ private buildGoalTurnPrompt; /** Goal 主循环:同一个 Run 内连续多轮执行,直到 Agent 调用 finish 或保护触发。 */ private runGoalLoop; /** 处理一轮 Goal 回合结果:停止/结算(交付)/未结束时落库统计。返回 true = Goal 已结束。 */ private consumeGoalTurn; /** 回合收尾(主对话与 Goal 共用):assistant 正文落库 + session 统计(turn_count +1 等)。返回落库消息 ID。 */ private recordAgentTurn; /** Goal 结束清理:状态与 Run 收尾(队列释放由 process 返回后 queue 接管)。 */ private cleanupGoal; /** 结算 Goal 并记录结果(纯内存;发送成功与否影响 outcome 落库摘要)。 */ private finishGoal; /** 发送 Goal 最终正文(唯一一次 IM 交付,卡片 + 引用触发消息 + 汇总 + footer;与常规交付同一降级链);失败时返回 false。 */ private deliverGoalFinalResponse; /** * 执行定时任务:创建独立 session,发送 prompt,结果用 ⏰ header 卡片发送,完成后归档。 * 不走用户消息队列,不干扰当前对话 session。 */ processCronJob(chatId: string, userId: string, prompt: string, description: string, cronJobId?: number, claimToken?: string, threadId?: string): Promise; reportCronJobFailure(chatId: string, description: string, error: string, paused: boolean, threadId?: string, replyToMsgId?: string): Promise; /** Loop Scheduler 只投递 ID;任务内容在真正轮到该 chat 时重新从 DB 读取。 */ enqueueLoopJob(loopJobId: number): void; /** 进程恢复:用户主 session 懒加载,有人说话 / Loop / wake 时再 attach。 */ recover(): Promise; /** Loop 是 Engine 生成的内部续接事件,不伪装成一条新的用户消息。 */ private buildLoopContinuationPrompt; private buildCollabActionRetryPrompt; private processIndependentSession; private isCronClaimCurrent; /** /new:归档当前 session,让下一条消息自然创建新 session。 */ private resetSession; private startSessionTransition; private startGlobalSessionTransition; private stopActiveRunForSessionTransition; private waitForSessionCreation; private enqueuePendingTransitionMessage; private drainPendingTransitionMessages; /** * /admin 命令:管理员列表/添加/移除。 * - /admin → 显示管理员列表 * - /admin add @某人 → 添加管理员(需要 @ mention) * - /admin remove @某人 → 移除管理员 */ private handleAdminCommand; /** * /agent 命令:查看或切换当前对话的 agent 整包(backend + 该 backend 自己的 model/effort)。 * 当前会话写覆盖;私聊同时更新 Bot 默认。没覆盖的跟着 Bot 默认走。 */ private handleAgentCommand; private resetScopeAgentConfig; /** * /model 命令:查看或切换模型。 * - /model → 显示当前模型 + 可选列表 * - /model → 切主模型 * - /model reset → 恢复为配置初始值 */ private handleModelCommand; /** * /effort 命令:查看或切换推理强度。 * - /effort → 显示当前 effort + 可选值 + 当前 backend 是否支持 * - /effort → 切换(low/medium/high/xhigh/max) * - /effort reset → 恢复 backend 默认 */ private handleEffortCommand; /** * /tz:查看或切换展示时区。时刻仍按 UTC 存储。 * - /tz → 显示当前时区 * - /tz 东京 或 /tz 把时区改成纽约 → 切换并写入配置 * - /tz reset → 恢复默认北京时间 */ private handleTimezoneCommand; /** /tz 和 agent(nbt timezone set)共用:解析后写入配置并切换展示时区。 */ setEngineTimezone(raw: string): string; /** /autoupdate:查看/开关自动升级。 */ private handleAutoUpdateCommand; private persistAutoUpdateSetting; /** 当前会话覆盖 ?? Bot 默认 ?? 配置文件。都是整包。私聊写入时同时更新 Bot 默认。 */ private resolveScopeConfig; private materializeScopeConfig; private resolveFallbackConfig; private enginePackage; private loadBotDefaultPackage; /** 某个 backend 自己的整包:用它上次的 model/effort,绝不借用别的 backend 的模型。 */ private packageForBackend; private persistScopeConfig; private persistBotDefault; private clearBotDefault; private resolveScopeBackendType; private isP2pScope; /** 用户这句话在话题里才回话题;话题群没有主时间线,缺 thread 也必须 reply。 */ private userSpokeInThread; private hasOwnScopeConfig; private rememberChatUser; /** Loop/wake 重启后 map 是空的,要从 session 或私聊行找回说话的人。Bot 发言不记成配置主人。 */ private resolveChatUserId; private ensureBackend; private backendForType; private backendForSession; private backendForScope; private resolveScopeModels; private setScopeModels; /** 同步当前 scope 已存在的 backend session;各 backend resume 时会从 session 对象读取 model/effort。 */ private updateActiveScopeSessionModels; /** 构建模型候选列表:初始配置 → 历史 → 运行时新增,顺序稳定 */ private buildModelCandidates; /** 解析模型参数:数字当编号,否则当名字 */ private resolveModelArg; /** 记录模型使用历史 */ private recordModelHistory; /** 显示模型列表卡片 */ private sendModelList; /** 发送 Agent 命令卡片回复。话题里必须带 threadId 落库,否则 nbt messages 会显示成主群。 */ private sendAgentCard; /** 发送 /help 卡片 */ private sendHelpCard; private handleAwakeCommand; /** 管理员命令:Windows 优先 pwsh、回退 powershell.exe,Unix 保持平台默认 shell。 */ private tryShellCommand; private recordShellHistory; private sendShellHistory; private handleUpdate; /** 自动升级状态/帮助行(供 /update 默认展示,单卡片合并)。 */ private buildAutoUpdateStatusLines; /** 向 Engine 暴露本 Bot 的只读空闲状态;Pipeline 不执行升级决策。 */ getUpgradeSafenessSources(): UpgradeSafenessSource[]; /** 请求 Engine 启动 restart worker;Pipeline 只负责定位通知会话和展示结果。 */ triggerRestart(opts?: { platformChatId?: string; chatId?: string; scopeKey?: string; threadId?: string; updateVersion?: string; silent?: boolean; replyToMsgId?: string; }): void; private process; private getOrCreateSession; private createChatSession; /** 新开 agent 时首条消息的动态上下文暂存(native resume 不灌) */ private pendingMessageContext; /** stable context 暂存(needsStableUserPrefix 时挂到首条 user 消息) */ private pendingStableContext; /** 新 session 首条消息提醒暂存 */ private pendingNewSessionReminder; private clearChatRuntimeState; private queueNewAgentMessageContext; private buildSessionProfile; private buildStableSystemContext; private updateCompactRecoveryState; private archiveSession; private lookupActiveUserSession; private resolveNativeAgentSessionId; private markSessionEnded; private archiveTranscript; private cancelChat; private maybeUnloadScope; private resolveWatchdogTarget; /** 向指定 chat 发送系统通知(不走 pipeline 队列) */ private sendWatchdogNotification; private sendWatchdogCard; /** 发送 compact 中提示(仅通知一次) */ notifyCompacting(chatId: string): void; /** Watchdog 主循环:检测 idle session,按策略通知或自动 kill */ private runIdleWatchdogSafely; private runIdleWatchdog; } export declare function formatShellExecError(cwd: string, cmd: string, err: unknown): string;