export interface QueuedMessage { chatId: string; /** 队列键:非隔离时 === chatId;隔离时 c5#omt_aaa。 */ scopeKey?: string; /** 飞书话题 ID(仅隔离时使用)。 */ threadId?: string; /** 本轮要回复的 platform message ID(Loop/Cron/wake 等无当前用户消息时使用)。 */ replyToMsgId?: string; /** 即使未隔离也必须禁止 create 新话题。 */ strict?: boolean; text: string; timestamp: number; platformMsgId?: string; /** 发送者短标签(如 "U2"),用于多条消息合并时生成 YAML 格式 */ senderLabel?: string; /** 发送者内部 ID,用于群聊 speaker 注入 */ senderId?: string; /** 消息在 DB 中的 ID,用于 runtime 事件关联 */ dbMsgId?: number; /** 触发来源:用户消息(默认)或内部续接事件 */ triggerKind?: "user" | "loop_continuation" | "restart_wake"; /** triggerKind 为 loop_continuation 时携带的 Loop Job ID(内部事件,不写库) */ loopJobId?: number; /** 原始用户文本是否以 /loop 或 /cron 开头;回复渲染后仍保留调度意图。 */ scheduleCommand?: boolean; /** 这条消息是否是严格多 Bot 协作回合;普通群消息不能借用协作状态。 */ collabTurn?: boolean; } type ProcessFn = (scopeKey: string, mergedText: string, messages: QueuedMessage[], signal: AbortSignal) => Promise; export interface QueueSnapshot { buffer: QueuedMessage[]; pending: QueuedMessage[]; busy: boolean; } export declare class MessageQueue { private queues; private processFn; private startFn; private finishFn; private idleFn; private bufferMs; private pendingFn; private stateFn; private discardFn; private stopped; constructor(bufferMs?: number); /** 注册消息处理函数 */ onProcess(fn: ProcessFn): void; /** 注册启动门禁;返回 false 时保持 buffer,不进入 busy。 */ onStart(fn: (scopeKey: string) => boolean): void; /** 注册 process 结束回调;返回的 sibling scopeKey 将由队列立即尝试启动。 */ onFinish(fn: (scopeKey: string) => string[]): void; onIdle(fn: (scopeKey: string) => void): void; /** 注册 pending 通知函数(消息进入等待队列时立即回调) */ onPending(fn: (msg: QueuedMessage) => void): void; /** 注册队列状态变更通知,用于外部同步观测状态 */ onStateChange(fn: (chatId: string, snapshot: QueueSnapshot) => void): void; /** 注册显式丢弃回调;服务关闭时保留队列供持久层恢复,不调用此回调。 */ onDiscard(fn: (messages: QueuedMessage[]) => void): void; /** 推入一条新消息,返回是否进入 pending 队列 */ push(msg: QueuedMessage): boolean; /** 停止队列,清除所有计时器 */ stop(): void; isStopped(): boolean; /** 清空指定 chat 的等待队列(buffer + pending),返回丢弃的消息数。 */ drain(scopeKey: string, shouldDiscard?: (message: QueuedMessage) => boolean): number; /** 获取指定 chat 的待处理消息数(buffer + pending) */ pendingCount(scopeKey: string): number; /** 是否有正在处理的任务 */ hasBusyChats(): boolean; /** 指定 chat 是否正在处理 */ isBusy(scopeKey: string): boolean; /** 取消指定 chat 正在进行的 process 调用 */ cancel(scopeKey: string): boolean; scopeKeys(): string[]; flushScope(scopeKey: string): Promise; private getQueue; private resetBufferTimer; /** 标记某 chat 处理完成,检查后续队列 */ private processNext; private flush; private logDeferred; private emitState; } export {};