Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions packages/cli/src/commands/start.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1586,6 +1586,9 @@ async function startServerCore(
agentId,
state.status as 'idle' | 'working' | 'offline' | 'error'
);
if (state.lastHeartbeat) {
await storage.agentRepo.updateLastHeartbeat(agentId, state.lastHeartbeat);
}
} catch (err) {
log.warn('Failed to persist agent state', { agentId, error: String(err) });
}
Expand Down
6 changes: 3 additions & 3 deletions packages/core/src/agent-manager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -472,7 +472,7 @@ export class AgentManager {
) => Promise<{ approved: boolean; comment?: string }>;
private stateChangeHandler?: (
agentId: string,
state: { status: string; tokensUsedToday: number; activeTaskIds: string[]; lastError?: string; lastErrorAt?: string; currentActivity?: AgentActivity }
state: { status: string; tokensUsedToday: number; activeTaskIds: string[]; lastError?: string; lastErrorAt?: string; currentActivity?: AgentActivity; lastHeartbeat?: string }
) => void;
private disabledChangeHandler?: (agentId: string, disabled: boolean) => void;
/** Grace timers for releasing scoped MCP processes after agent goes idle */
Expand Down Expand Up @@ -3335,7 +3335,7 @@ export class AgentManager {
*/
private buildStateChangeCallback(): (
agentId: string,
state: { status: string; tokensUsedToday: number; activeTaskIds: string[]; lastError?: string; lastErrorAt?: string; currentActivity?: AgentActivity }
state: { status: string; tokensUsedToday: number; activeTaskIds: string[]; lastError?: string; lastErrorAt?: string; currentActivity?: AgentActivity; lastHeartbeat?: string }
) => void {
return (agentId, state) => {
if (state.status === 'idle' && state.activeTaskIds.length === 0) {
Expand Down Expand Up @@ -3374,7 +3374,7 @@ export class AgentManager {
setStateChangeHandler(
handler: (
agentId: string,
state: { status: string; tokensUsedToday: number; activeTaskIds: string[]; lastError?: string; lastErrorAt?: string; currentActivity?: AgentActivity }
state: { status: string; tokensUsedToday: number; activeTaskIds: string[]; lastError?: string; lastErrorAt?: string; currentActivity?: AgentActivity; lastHeartbeat?: string }
) => void
): void {
this.stateChangeHandler = handler;
Expand Down
24 changes: 13 additions & 11 deletions packages/core/src/agent.ts
Original file line number Diff line number Diff line change
Expand Up @@ -76,7 +76,6 @@ import { detectEnvironment, type EnvironmentProfile } from './environment-profil
import { ToolSelector } from './tool-selector.js';
import {
scenarioToPack,
getReflexAllowlist,
formatEvictedToolCatalog,
COMMENT_RESPONSE_ALLOWED_TOOLS,
REQUIREMENT_ACTION_ALLOWED_TOOLS,
Expand Down Expand Up @@ -564,7 +563,7 @@ export class Agent {
private stateManager?: AgentStateManager;
private stateChangeCallback?: (
agentId: string,
state: { status: string; tokensUsedToday: number; activeTaskIds: string[]; lastError?: string; lastErrorAt?: string; currentActivity?: AgentActivity }
state: { status: string; tokensUsedToday: number; activeTaskIds: string[]; lastError?: string; lastErrorAt?: string; currentActivity?: AgentActivity; lastHeartbeat?: string }
) => void;
private memoryConsolidationTimer?: ReturnType<typeof setInterval>;
private loopDetector = new ToolLoopDetector();
Expand Down Expand Up @@ -3972,7 +3971,7 @@ export class Agent {
setStateChangeCallback(
cb: (
agentId: string,
state: { status: string; tokensUsedToday: number; activeTaskIds: string[]; lastError?: string; lastErrorAt?: string; currentActivity?: AgentActivity }
state: { status: string; tokensUsedToday: number; activeTaskIds: string[]; lastError?: string; lastErrorAt?: string; currentActivity?: AgentActivity; lastHeartbeat?: string }
) => void
): void {
this.stateChangeCallback = cb;
Expand All @@ -3987,6 +3986,7 @@ export class Agent {
lastError: this.state.lastError,
lastErrorAt: this.state.lastErrorAt,
currentActivity: this.getCurrentActivity(),
lastHeartbeat: this.state.lastHeartbeat,
});
}
}
Expand Down Expand Up @@ -8865,6 +8865,7 @@ export class Agent {
this.endActivity(skipActivityId, { success: true });
this.state.lastHeartbeat = new Date().toISOString();
this.metricsCollector.recordHeartbeat(true, true);
this.notifyStateChange(); // 心跳也是存活证明:落库 last_heartbeat,供所有观察者读取
return;
}

Expand Down Expand Up @@ -8910,6 +8911,7 @@ export class Agent {
this.endActivity(skipActivityId, { success: true });
this.state.lastHeartbeat = new Date().toISOString();
this.metricsCollector.recordHeartbeat(true, true);
this.notifyStateChange(); // 心跳存活证明落库
if (deepSleep) {
try {
const cur = (this as unknown as { heartbeatIntervalMs?: number }).heartbeatIntervalMs ?? 6 * 3600_000;
Expand Down Expand Up @@ -9063,10 +9065,10 @@ export class Agent {
selfEvolutionSection,
'',
'## Patrol Reminder',
'This is a lightweight patrol (reflex pack), not a work session.',
'- **Do**: `task_list`/`task_get` triage, `memory_save` one-line insights, `notify_user` / `request_user_input` when humans must act, `schedule_wakeup` for precise follow-ups.',
'- **Don\'t**: `task_create`, `requirement_propose`, write code, refactor, or deep analysis. Escalate or wake into a chat/task session instead.',
'- Use `set_heartbeat_interval` only if the patrol cadence itself is wrong.',
'You have the same toolset as a normal session — use it with the same judgement, scoped to patrol.',
'- **Scope**: follow YOUR HEARTBEAT.md checklist. Triage with `task_list`/`task_get`; leave one-line insights via `memory_save`; escalate decisions to humans via `notify_user`/`request_user_input`; use `schedule_wakeup` for precise follow-ups.',
'- **Don\u2019t drift into deep work in a routine heartbeat**: writing code / refactoring / editing files is allowed only for small, low-risk fixes the patrol genuinely requires; otherwise wake a proper chat/task session instead of doing heavy work here.',
'- **Cadence**: if the patrol interval itself is wrong, ASK THE HUMAN FIRST (e.g. `request_user_input` / `notify_user`), and only with their consent call `set_heartbeat_interval`.',
'',
'## Finishing Up',
'- Compare against your last heartbeat summary above. Skip unchanged items.',
Expand All @@ -9075,9 +9077,9 @@ export class Agent {
'- If nothing needs attention and no daily report is due, respond with exactly: HEARTBEAT_OK',
].join('\n');

// Reflex pack allowlist (AGENT-RUNTIME §2.2) — no package/goal/spawn fat.
const HEARTBEAT_ALLOWED_TOOLS = getReflexAllowlist(isManager);

// 工具集:与普通 session 同级(scenario 'heartbeat' → converse 包)。
// 不再传入 reflex 白名单锁死能力 —— 约束由 HEARTBEAT.md(巡检范围/红线)
// 与心跳间隔(成本)承担。工具迭代上限仍保留,作为单次巡检的成本护栏。
const HEARTBEAT_MAX_RETRIES = 3;
const HEARTBEAT_RETRY_BASE_MS = 3000;
let lastError: unknown;
Expand All @@ -9086,7 +9088,6 @@ export class Agent {
try {
const reply = await this.handleMessage(prompt, undefined, undefined, {
sessionId: heartbeatSessionId(this.id),
allowedTools: HEARTBEAT_ALLOWED_TOOLS,
scenario: 'heartbeat',
maxToolIterations: Agent.HEARTBEAT_MAX_TOOL_ITERATIONS,
});
Expand All @@ -9101,6 +9102,7 @@ export class Agent {

this.state.lastHeartbeat = new Date().toISOString();
this.metricsCollector.recordHeartbeat(true);
this.notifyStateChange(); // 心跳完成=存活证明:落库 last_heartbeat

if (dailyReportSection) {
this.memory.addEntry({
Expand Down
4 changes: 4 additions & 0 deletions packages/core/src/capability-packs.ts
Original file line number Diff line number Diff line change
Expand Up @@ -135,6 +135,10 @@ export function allowsWorkContextBoundTools(
export function scenarioToPack(scenario: string | undefined): CapabilityPack {
switch (scenario) {
case 'heartbeat':
// 心跳与普通 session 同级的能力(converse)。不再用小工具 reflex 包锁死——
// 约束回归「心跳描述文件(HEARTBEAT.md)+ 心跳间隔」:HEARTBEAT.md 定义巡检
// 范围与行为红线,间隔控制成本。见 docs/agent-liveness-redesign.md §三.2。
return 'converse';
case 'memory_consolidation':
case 'memory_flush':
case 'distillation':
Expand Down
4 changes: 1 addition & 3 deletions packages/core/src/context-engine.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1609,9 +1609,7 @@ export class ContextEngine {
lines.push('');
lines.push('The runtime skips heartbeat LLM turns while a human chat is focused or queued — you should not see a heartbeat mid-conversation. If you do run, keep it brief.');
lines.push('');
lines.push('**Tools (reflex pack only)**: `task_list`, `task_get`, `memory_save`/`memory_search`, `notify_user`, `request_user_input`, `schedule_wakeup`/`cancel_wakeup`, `set_heartbeat_interval`, `discover_tools`, `check_mailbox`, `file_read`, `agent_send_message`, `update_notebook`'
+ (extra?.isManager ? ', `team_status`' : '')
+ '. Do **not** call `task_create`, `requirement_propose`, `package_install`, or other execute-pack tools here.');
lines.push('**Tools (full session capability)**: heartbeat runs with the same toolset as a normal session — the boundaries come from your HEARTBEAT.md checklist (scope) and the heartbeat interval (cost), not a hardcoded tool subset. Triage and follow-up are welcome; for deep work, wake a real chat/task session instead of doing it inside a routine patrol. If the patrol cadence itself is wrong, ask the human first (`request_user_input` / `notify_user`) and only with their consent call `set_heartbeat_interval`.');
lines.push('');
lines.push('**Priority actions (in order):**');
lines.push('1. **Patrol board**: `task_list` / `task_get` for reviews due, failed, or stuck items. Note blockers; `notify_user` if a human must act.');
Expand Down
3 changes: 2 additions & 1 deletion packages/core/test/capability-packs.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,8 @@ describe('capability packs (AGENT-RUNTIME §2)', () => {
});

it('maps heartbeat/review scenarios', () => {
expect(scenarioToPack('heartbeat')).toBe('reflex');
// 心跳与普通 session 同级能力(converse)——约束由 HEARTBEAT.md + 间隔承担
expect(scenarioToPack('heartbeat')).toBe('converse');
expect(scenarioToPack('review')).toBe('govern');
expect(packToolDefBudget('reflex')).toBe(TOOL_DEF_BUDGET_REFLEX);
expect(packToolDefBudget('converse')).toBe(TOOL_DEF_BUDGET_CONVERSE);
Expand Down
7 changes: 4 additions & 3 deletions packages/core/test/context-scenario-matrix.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -149,10 +149,11 @@ describe('group chat is never rendered with 1:1-A2A semantics', () => {
});

describe('reflex scenarios restrict tools in the prompt (and the runtime enforces it)', () => {
it('heartbeat names the reflex pack and forbids execute-pack tools', async () => {
it('heartbeat runs with full session capability (boundaries via HEARTBEAT.md + interval)', async () => {
const { text } = await build({ scenario: 'heartbeat' });
expect(text).toContain('reflex pack only');
expect(text).toContain('task_create');
expect(text).toContain('full session capability');
expect(text).not.toContain('reflex pack only');
expect(text).not.toContain('Do **not** call `task_create`');
expect(text).toContain('HEARTBEAT_OK');
});

Expand Down
6 changes: 4 additions & 2 deletions packages/core/test/prompt-profiles.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -151,7 +151,7 @@ describe('prompt profiles (AGENT-RUNTIME §4)', () => {
expect(tight.truncated).toBe(true);
});

it('A-heartbeat-no-task-create-prompt: heartbeat mode omits create/propose tools', async () => {
it('A-heartbeat-full-capability-prompt: heartbeat runs with full session toolset, no reflex lock', async () => {
const engine = new ContextEngine();
const { text } = await engine.buildSystemPrompt({
agentId: 'agt_1',
Expand All @@ -162,7 +162,9 @@ describe('prompt profiles (AGENT-RUNTIME §4)', () => {
promptProfile: 'reflex',
});
expect(text).toContain('heartbeat mode');
expect(text).toMatch(/Do \*\*not\*\* call `task_create`/);
expect(text).toContain('full session capability');
expect(text).not.toContain('reflex pack only');
expect(text).not.toContain('Do **not** call `task_create`');
expect(text).not.toContain('You MAY create tasks via `task_create`');
expect(text).not.toContain('## Self-Evolution');
});
Expand Down
55 changes: 49 additions & 6 deletions packages/org-manager/src/agent-dirty-reconciler.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,19 @@ const log = createLogger('agent-dirty-reconciler');
/** 同一脏态两次兜底尝试的最短间隔:失败后 5 分钟再试,给 agent 自愈时间又不会永久放弃。 */
export const RETRY_AFTER_MS = 5 * 60_000;

/**
* trigger-heartbeat 连续兜底次数的上限:达到后仍未见效(agent 仍处 processing-like)
* 就停止自动触发并升级为 human-review。
*
* 背景(心跳风暴回归):stuck-working 的 agent 被触发一次心跳后,lastHeartbeat 会在
* heartbeatGraceMs(2 分钟)内保持新鲜 → 判定为「不脏」→ 重试 key 被释放;宽限一过又判脏 →
* 立刻再触发(绕过 RETRY_AFTER_MS)。结果:每 2 分钟触发一次恢复心跳、每个心跳跑一轮完整
* LLM 巡检并产生一条「定时心跳签到」记录 —— 一天数百条,属于自愈机制的反馈回路 bug。
* 修复双重保险:① 仅当 agent 真正脱离 processing-like 才释放重试状态;② 连续 N 次触发仍
* 未解决则放弃自动兜底、升级人工介入,杜绝无限心跳风暴。
*/
export const MAX_TRIGGER_HEARTBEAT_ATTEMPTS = 3;

/** 扫描所需的最小 agent live-state 视图(对应 agentManager.listAgents() 的字段)。 */
export interface AgentLiveView {
agentId: string;
Expand Down Expand Up @@ -67,6 +80,8 @@ export class AgentDirtyReconciler {
private timer?: ReturnType<typeof setInterval>;
/** 已兜底过的脏 key(agentId:recovery)→ 最近一次尝试时间。 */
private attempted = new Map<string, number>();
/** 同一脏 key 的连续兜底次数(用于 trigger-heartbeat 升级阈值)。 */
private consecutive = new Map<string, number>();

constructor(private opts: DirtyReconcilerOptions) {
this.cfg = { ...DEFAULT_DIRTY_CONFIG, ...opts.cfg };
Expand All @@ -77,6 +92,7 @@ export class AgentDirtyReconciler {
if (!this.cfg.enabled) return [];
const done: DirtyVerdict[] = [];
const seenKeys = new Set<string>();
const viewById = new Map(agents.map((a) => [a.agentId, a] as const));

for (const a of agents) {
const v = evaluateDirtyState(
Expand Down Expand Up @@ -104,21 +120,48 @@ export class AgentDirtyReconciler {
if (lastAt !== undefined && now - lastAt < RETRY_AFTER_MS) continue;
this.attempted.set(key, now);

let effective = v;
if (v.recovery === 'trigger-heartbeat') {
const n = (this.consecutive.get(key) ?? 0) + 1;
this.consecutive.set(key, n);
if (n > MAX_TRIGGER_HEARTBEAT_ATTEMPTS) {
// 连续多次心跳仍未见效 → 说明自愈对该 agent 无效,停止自动触发,升级人工介入。
effective = {
...v,
recovery: 'human-review',
reason: `${v.reason}(已连续 ${n - 1} 次触发恢复心跳仍未脱离 stuck-busy,停止自动兜底)`,
suggestions: [
'在 Agent 设置中手动重置/重启该 agent 以清除卡死的 working 状态',
'检查其真实任务是否早已结束,status 是否被错误钉在 working',
],
};
this.consecutive.delete(key);
}
}

// 可观测:写一条执行流事件(谁/何时/为何/建议)。
this.observe(a, v);
this.observe(a, effective);

if (v.recovery === 'human-review') {
if (this.opts.onNeedsHuman) void this.opts.onNeedsHuman(v);
else log.warn('Dirty agent needs human review', { agentId: a.agentId, verdict: v });
if (effective.recovery === 'human-review') {
if (this.opts.onNeedsHuman) void this.opts.onNeedsHuman(effective);
else log.warn('Dirty agent needs human review', { agentId: a.agentId, verdict: effective });
} else if (this.opts.recover) {
await this.opts.recover(v);
await this.opts.recover(effective);
}
// recover 默认无动作(纯观察 + 事件)——安全默认。
}

// 释放已恢复(不再脏)的 key,允许再次变脏时立即重新兜底。
// 注意:仅当 agent 真正脱离 processing-like(status != working 且无活动痕迹)才释放。
// 不能只凭「本次不脏」就释放 —— 心跳宽限窗口(heartbeatGraceMs=2min)内的暂时新鲜会
// 让 key 被提前释放,宽限一过又立刻重触发,绕过 RETRY_AFTER_MS 形成 2 分钟心跳风暴。
for (const k of [...this.attempted.keys()]) {
if (!seenKeys.has(k)) this.attempted.delete(k);
if (seenKeys.has(k)) continue;
const agentId = k.slice(0, k.lastIndexOf(':'));
const view = viewById.get(agentId);
if (view && (view.status === 'working' || view.currentActivity)) continue; // 仍 processing-like,非真恢复
this.attempted.delete(k);
this.consecutive.delete(k);
}
return done;
}
Expand Down
41 changes: 41 additions & 0 deletions packages/org-manager/test/agent-dirty-reconciler.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -100,4 +100,45 @@ describe('AgentDirtyReconciler — 脏态周期兜底(OB-3)', () => {
expect(out).toEqual([]);
expect(recover).not.toHaveBeenCalled();
});

it('回归:stuck-working 心跳风暴 —— 心跳宽限期内不释放重试 key,连续 N 次后升级 human-review', async () => {
const recover = vi.fn();
const onNeedsHuman = vi.fn();
const { reconciler } = make({ recover, onNeedsHuman });

// 模拟被顶死的 working:无活动、无任务、心跳每轮触发后短暂新鲜 2 分钟又变旧(宽限窗口)。
const stuckWorking = (heartbeatAgeMs: number): AgentLiveView => ({
agentId: 'storm',
status: 'working',
currentActivity: null,
activeTaskIds: [],
lastHeartbeat: new Date(NOW - heartbeatAgeMs).toISOString(),
});
const HBEAT_FRESH_MS = 60_000; // < 2min 宽限 → busy 判定新鲜

// t0:心跳已旧 → 判脏 → 触发恢复心跳(第 1 次)
await reconciler.scan([stuckWorking(3 * 60_000)], NOW);
expect(recover).toHaveBeenCalledTimes(1);

// t0+30s:心跳刚被触发 → 新鲜 → 不脏,但 status 仍是 working → key 不得释放
await reconciler.scan([stuckWorking(HBEAT_FRESH_MS)], NOW + 30_000);
// t0+2.5min:心跳再次变旧 → 又判脏,但距上次仅 2.5min < RETRY_AFTER_MS(5min) → 不重复触发
await reconciler.scan([stuckWorking(3 * 60_000)], NOW + 2.5 * 60_000);
expect(recover).toHaveBeenCalledTimes(1);

// t0+5min、+10min:超窗重试(第 2、3 次)
await reconciler.scan([stuckWorking(3 * 60_000)], NOW + 5 * 60_000);
await reconciler.scan([stuckWorking(3 * 60_000)], NOW + 10 * 60_000);
// 第 4 次越窗时仍未见效 → 升级 human-review,不再自动触发心跳
await reconciler.scan([stuckWorking(3 * 60_000)], NOW + 15 * 60_000);
expect(onNeedsHuman).toHaveBeenCalledTimes(1);

// 全程没有出现「每 2 分钟一次」的密集触发(对照修复前:每次宽限窗口过后都会立即重触发)
expect(recover.mock.calls.filter(([v]: any) => v.recovery === 'trigger-heartbeat')).toHaveLength(3);

// agent 真正恢复 idle 后 key 释放 → 未来再变脏可重新兜底
await reconciler.scan([{ ...stuckWorking(3 * 60_000), status: 'idle' }], NOW + 20 * 60_000);
await reconciler.scan([stuckWorking(3 * 60_000)], NOW + 21 * 60_000);
expect(recover).toHaveBeenCalledTimes(4);
});
});
5 changes: 5 additions & 0 deletions packages/storage/src/sqlite-storage.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1147,6 +1147,11 @@ export class SqliteAgentRepo {
.run(status, containerId ?? null, now(), id);
}

/** 记录 agent 最近一次心跳完成时间(存活证明的持久化事实源)。 */
updateLastHeartbeat(id: string, ts: string) {
this.db.prepare('UPDATE agents SET last_heartbeat = ?, updated_at = ? WHERE id = ?').run(ts, now(), id);
}

updateTokens(id: string, tokensUsed: number) {
this.db
.prepare('UPDATE agents SET tokens_used_today = ?, updated_at = ? WHERE id = ?')
Expand Down
Loading