feat(FN-1264): enforce heartbeat budget governance for trigger execution
- Gate heartbeat runs in executeHeartbeat when agents are over budget and skip timer runs above threshold - Pause agents with pauseReason "budget-exhausted" when run usage pushes them over budget - Budget-gate assignment and timer triggers in HeartbeatTriggerScheduler and include budgetStatus in assignment wake context - Expand agent-heartbeat test coverage for budget exhaustion, threshold behavior, pause transitions, and trigger gating
This commit is contained in:
@@ -1,6 +1,6 @@
|
||||
import { describe, it, expect, vi, beforeEach, afterEach } from "vitest";
|
||||
import { HeartbeatMonitor, HeartbeatTriggerScheduler, type AgentSession, type HeartbeatExecutionOptions, HEARTBEAT_SYSTEM_PROMPT } from "./agent-heartbeat.js";
|
||||
import type { AgentStore, AgentHeartbeatRun, TaskStore, TaskDetail, Agent, MessageStore, Message } from "@fusion/core";
|
||||
import type { AgentStore, AgentHeartbeatRun, TaskStore, TaskDetail, Agent, MessageStore, Message, AgentBudgetStatus } from "@fusion/core";
|
||||
|
||||
// Mock logger to suppress noise in test output
|
||||
vi.mock("./logger.js", () => {
|
||||
@@ -68,6 +68,21 @@ function createMessage(overrides: Partial<Message> = {}): Message {
|
||||
};
|
||||
}
|
||||
|
||||
function createBudgetStatus(overrides: Partial<AgentBudgetStatus> = {}): AgentBudgetStatus {
|
||||
return {
|
||||
agentId: "agent-001",
|
||||
currentUsage: 0,
|
||||
budgetLimit: null,
|
||||
usagePercent: null,
|
||||
thresholdPercent: null,
|
||||
isOverBudget: false,
|
||||
isOverThreshold: false,
|
||||
lastResetAt: null,
|
||||
nextResetAt: null,
|
||||
...overrides,
|
||||
};
|
||||
}
|
||||
|
||||
describe("HeartbeatMonitor", () => {
|
||||
let store: AgentStore;
|
||||
let monitor: HeartbeatMonitor;
|
||||
@@ -1013,6 +1028,7 @@ describe("HeartbeatMonitor", () => {
|
||||
};
|
||||
}),
|
||||
endHeartbeatRun: vi.fn().mockResolvedValue(undefined),
|
||||
getBudgetStatus: vi.fn().mockResolvedValue(createBudgetStatus()),
|
||||
getCachedAgent: vi.fn().mockReturnValue(null),
|
||||
} as unknown as AgentStore;
|
||||
}
|
||||
@@ -1525,6 +1541,132 @@ describe("HeartbeatMonitor", () => {
|
||||
expect(monitor.getTrackedAgents()).not.toContain("agent-001");
|
||||
});
|
||||
});
|
||||
|
||||
describe("Budget Governance", () => {
|
||||
it("skips heartbeat when agent is over budget (timer)", async () => {
|
||||
const budgetStatus = createBudgetStatus({
|
||||
currentUsage: 10000,
|
||||
budgetLimit: 10000,
|
||||
usagePercent: 100,
|
||||
thresholdPercent: 80,
|
||||
isOverBudget: true,
|
||||
isOverThreshold: true,
|
||||
});
|
||||
const store = createStoreWithAgentForExec();
|
||||
(store.getBudgetStatus as ReturnType<typeof vi.fn>).mockResolvedValue(budgetStatus);
|
||||
|
||||
const monitor = new HeartbeatMonitor({ store, taskStore: mockTaskStore, rootDir: "/tmp" });
|
||||
const result = await monitor.executeHeartbeat({ agentId: "agent-001", source: "timer" });
|
||||
|
||||
expect(result.status).toBe("completed");
|
||||
expect(result.resultJson).toMatchObject({ reason: "budget_exhausted", budgetStatus });
|
||||
expect(mockedCreateKbAgent).not.toHaveBeenCalled();
|
||||
expect(store.updateAgentState).not.toHaveBeenCalledWith("agent-001", "active");
|
||||
});
|
||||
|
||||
it("skips heartbeat when agent is over budget (on_demand)", async () => {
|
||||
const store = createStoreWithAgentForExec();
|
||||
(store.getBudgetStatus as ReturnType<typeof vi.fn>).mockResolvedValue(
|
||||
createBudgetStatus({ isOverBudget: true, isOverThreshold: true, usagePercent: 100 })
|
||||
);
|
||||
|
||||
const monitor = new HeartbeatMonitor({ store, taskStore: mockTaskStore, rootDir: "/tmp" });
|
||||
const result = await monitor.executeHeartbeat({ agentId: "agent-001", source: "on_demand" });
|
||||
|
||||
expect(result.resultJson).toMatchObject({ reason: "budget_exhausted" });
|
||||
expect(mockedCreateKbAgent).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("skips heartbeat when agent is over budget (assignment)", async () => {
|
||||
const store = createStoreWithAgentForExec();
|
||||
(store.getBudgetStatus as ReturnType<typeof vi.fn>).mockResolvedValue(
|
||||
createBudgetStatus({ isOverBudget: true, isOverThreshold: true, usagePercent: 100 })
|
||||
);
|
||||
|
||||
const monitor = new HeartbeatMonitor({ store, taskStore: mockTaskStore, rootDir: "/tmp" });
|
||||
const result = await monitor.executeHeartbeat({ agentId: "agent-001", source: "assignment" });
|
||||
|
||||
expect(result.resultJson).toMatchObject({ reason: "budget_exhausted" });
|
||||
expect(mockedCreateKbAgent).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("skips timer heartbeat when agent is over threshold but not over budget", async () => {
|
||||
const budgetStatus = createBudgetStatus({
|
||||
currentUsage: 850,
|
||||
budgetLimit: 1000,
|
||||
usagePercent: 85,
|
||||
thresholdPercent: 80,
|
||||
isOverBudget: false,
|
||||
isOverThreshold: true,
|
||||
});
|
||||
const store = createStoreWithAgentForExec();
|
||||
(store.getBudgetStatus as ReturnType<typeof vi.fn>).mockResolvedValue(budgetStatus);
|
||||
|
||||
const monitor = new HeartbeatMonitor({ store, taskStore: mockTaskStore, rootDir: "/tmp" });
|
||||
const result = await monitor.executeHeartbeat({ agentId: "agent-001", source: "timer" });
|
||||
|
||||
expect(result.resultJson).toMatchObject({ reason: "budget_threshold_exceeded", budgetStatus });
|
||||
expect(mockedCreateKbAgent).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("allows on_demand heartbeat when agent is over threshold", async () => {
|
||||
const store = createStoreWithAgentForExec();
|
||||
const mockSession = createMockAgentSession();
|
||||
mockedCreateKbAgent.mockResolvedValue({ session: mockSession as any });
|
||||
(store.getBudgetStatus as ReturnType<typeof vi.fn>).mockResolvedValue(
|
||||
createBudgetStatus({ isOverThreshold: true, usagePercent: 85, budgetLimit: 1000, thresholdPercent: 80 })
|
||||
);
|
||||
|
||||
const monitor = new HeartbeatMonitor({ store, taskStore: mockTaskStore, rootDir: "/tmp" });
|
||||
const result = await monitor.executeHeartbeat({ agentId: "agent-001", source: "on_demand" });
|
||||
|
||||
expect(result.status).toBe("completed");
|
||||
expect(mockedCreateKbAgent).toHaveBeenCalledOnce();
|
||||
});
|
||||
|
||||
it("allows assignment heartbeat when agent is over threshold", async () => {
|
||||
const store = createStoreWithAgentForExec();
|
||||
const mockSession = createMockAgentSession();
|
||||
mockedCreateKbAgent.mockResolvedValue({ session: mockSession as any });
|
||||
(store.getBudgetStatus as ReturnType<typeof vi.fn>).mockResolvedValue(
|
||||
createBudgetStatus({ isOverThreshold: true, usagePercent: 85, budgetLimit: 1000, thresholdPercent: 80 })
|
||||
);
|
||||
|
||||
const monitor = new HeartbeatMonitor({ store, taskStore: mockTaskStore, rootDir: "/tmp" });
|
||||
const result = await monitor.executeHeartbeat({ agentId: "agent-001", source: "assignment" });
|
||||
|
||||
expect(result.status).toBe("completed");
|
||||
expect(mockedCreateKbAgent).toHaveBeenCalledOnce();
|
||||
});
|
||||
|
||||
it("proceeds normally when agent is below threshold", async () => {
|
||||
const store = createStoreWithAgentForExec();
|
||||
const mockSession = createMockAgentSession();
|
||||
mockedCreateKbAgent.mockResolvedValue({ session: mockSession as any });
|
||||
(store.getBudgetStatus as ReturnType<typeof vi.fn>).mockResolvedValue(
|
||||
createBudgetStatus({ isOverBudget: false, isOverThreshold: false, usagePercent: 30, budgetLimit: 1000, thresholdPercent: 80 })
|
||||
);
|
||||
|
||||
const monitor = new HeartbeatMonitor({ store, taskStore: mockTaskStore, rootDir: "/tmp" });
|
||||
const result = await monitor.executeHeartbeat({ agentId: "agent-001", source: "timer" });
|
||||
|
||||
expect(result.status).toBe("completed");
|
||||
expect(mockedCreateKbAgent).toHaveBeenCalledOnce();
|
||||
});
|
||||
|
||||
it("proceeds normally when getBudgetStatus throws", async () => {
|
||||
const store = createStoreWithAgentForExec();
|
||||
const mockSession = createMockAgentSession();
|
||||
mockedCreateKbAgent.mockResolvedValue({ session: mockSession as any });
|
||||
(store.getBudgetStatus as ReturnType<typeof vi.fn>).mockRejectedValue(new Error("budget unavailable"));
|
||||
|
||||
const monitor = new HeartbeatMonitor({ store, taskStore: mockTaskStore, rootDir: "/tmp" });
|
||||
const result = await monitor.executeHeartbeat({ agentId: "agent-001", source: "timer" });
|
||||
|
||||
expect(result.status).toBe("completed");
|
||||
expect(mockedCreateKbAgent).toHaveBeenCalledOnce();
|
||||
});
|
||||
});
|
||||
});
|
||||
|
||||
// ── Task Creation Tracking Tests ──────────────────────────────────────
|
||||
@@ -1744,6 +1886,139 @@ describe("HeartbeatMonitor", () => {
|
||||
});
|
||||
});
|
||||
|
||||
describe("Budget Governance", () => {
|
||||
function createCompleteRunBudgetStore(options: {
|
||||
agent?: Partial<Agent>;
|
||||
budgetStatus?: AgentBudgetStatus;
|
||||
budgetStatusError?: Error;
|
||||
} = {}): AgentStore {
|
||||
const run: AgentHeartbeatRun = {
|
||||
id: "run-budget-001",
|
||||
agentId: "agent-001",
|
||||
startedAt: new Date().toISOString(),
|
||||
endedAt: null,
|
||||
status: "active",
|
||||
};
|
||||
const agent: Agent = {
|
||||
id: "agent-001",
|
||||
name: "Budget Agent",
|
||||
role: "executor",
|
||||
state: "running",
|
||||
taskId: "FN-001",
|
||||
createdAt: new Date().toISOString(),
|
||||
updatedAt: new Date().toISOString(),
|
||||
metadata: {},
|
||||
...options.agent,
|
||||
} as Agent;
|
||||
|
||||
return {
|
||||
getRunDetail: vi.fn().mockResolvedValue(run),
|
||||
saveRun: vi.fn().mockResolvedValue(undefined),
|
||||
endHeartbeatRun: vi.fn().mockResolvedValue(undefined),
|
||||
getAgent: vi.fn().mockResolvedValue(agent),
|
||||
updateAgent: vi.fn().mockResolvedValue(undefined),
|
||||
updateAgentState: vi.fn().mockResolvedValue(undefined),
|
||||
getBudgetStatus: options.budgetStatusError
|
||||
? vi.fn().mockRejectedValue(options.budgetStatusError)
|
||||
: vi.fn().mockResolvedValue(options.budgetStatus ?? createBudgetStatus()),
|
||||
} as unknown as AgentStore;
|
||||
}
|
||||
|
||||
it("pauses agent with budget-exhausted reason when run pushes usage over budget", async () => {
|
||||
const store = createCompleteRunBudgetStore({
|
||||
agent: { totalInputTokens: 950, totalOutputTokens: 0 },
|
||||
budgetStatus: createBudgetStatus({
|
||||
currentUsage: 1050,
|
||||
budgetLimit: 1000,
|
||||
usagePercent: 105,
|
||||
thresholdPercent: 80,
|
||||
isOverBudget: true,
|
||||
isOverThreshold: true,
|
||||
}),
|
||||
});
|
||||
const monitor = new HeartbeatMonitor({ store });
|
||||
|
||||
await monitor.completeRun("agent-001", "run-budget-001", {
|
||||
status: "completed",
|
||||
usageJson: { inputTokens: 0, outputTokens: 100, cachedTokens: 0 },
|
||||
});
|
||||
|
||||
expect(store.updateAgentState).toHaveBeenCalledWith("agent-001", "paused");
|
||||
expect(store.updateAgent).toHaveBeenCalledWith("agent-001", { pauseReason: "budget-exhausted" });
|
||||
expect(store.updateAgentState).not.toHaveBeenCalledWith("agent-001", "active");
|
||||
});
|
||||
|
||||
it("does not pause agent when below budget after run", async () => {
|
||||
const store = createCompleteRunBudgetStore({
|
||||
budgetStatus: createBudgetStatus({
|
||||
currentUsage: 700,
|
||||
budgetLimit: 1000,
|
||||
usagePercent: 70,
|
||||
thresholdPercent: 80,
|
||||
isOverBudget: false,
|
||||
isOverThreshold: false,
|
||||
}),
|
||||
});
|
||||
const monitor = new HeartbeatMonitor({ store });
|
||||
|
||||
await monitor.completeRun("agent-001", "run-budget-001", {
|
||||
status: "completed",
|
||||
usageJson: { inputTokens: 10, outputTokens: 50, cachedTokens: 0 },
|
||||
});
|
||||
|
||||
expect(store.updateAgentState).toHaveBeenCalledWith("agent-001", "active");
|
||||
expect(store.updateAgent).not.toHaveBeenCalledWith("agent-001", { pauseReason: "budget-exhausted" });
|
||||
});
|
||||
|
||||
it("does not pause agent when run fails (status=failed)", async () => {
|
||||
const store = createCompleteRunBudgetStore({
|
||||
budgetStatus: createBudgetStatus({ isOverBudget: true, isOverThreshold: true }),
|
||||
});
|
||||
const monitor = new HeartbeatMonitor({ store });
|
||||
|
||||
await monitor.completeRun("agent-001", "run-budget-001", {
|
||||
status: "failed",
|
||||
usageJson: { inputTokens: 10, outputTokens: 50, cachedTokens: 0 },
|
||||
stderrExcerpt: "failure",
|
||||
});
|
||||
|
||||
expect(store.getBudgetStatus).not.toHaveBeenCalled();
|
||||
expect(store.updateAgentState).toHaveBeenCalledWith("agent-001", "error");
|
||||
expect(store.updateAgent).not.toHaveBeenCalledWith("agent-001", { pauseReason: "budget-exhausted" });
|
||||
});
|
||||
|
||||
it("does not pause agent when run is terminated", async () => {
|
||||
const store = createCompleteRunBudgetStore({
|
||||
budgetStatus: createBudgetStatus({ isOverBudget: true, isOverThreshold: true }),
|
||||
});
|
||||
const monitor = new HeartbeatMonitor({ store });
|
||||
|
||||
await monitor.completeRun("agent-001", "run-budget-001", {
|
||||
status: "terminated",
|
||||
usageJson: { inputTokens: 10, outputTokens: 50, cachedTokens: 0 },
|
||||
});
|
||||
|
||||
expect(store.getBudgetStatus).not.toHaveBeenCalled();
|
||||
expect(store.updateAgentState).toHaveBeenCalledWith("agent-001", "terminated");
|
||||
expect(store.updateAgent).not.toHaveBeenCalledWith("agent-001", { pauseReason: "budget-exhausted" });
|
||||
});
|
||||
|
||||
it("does not pause agent when usageJson is undefined", async () => {
|
||||
const store = createCompleteRunBudgetStore({
|
||||
budgetStatus: createBudgetStatus({ isOverBudget: true, isOverThreshold: true }),
|
||||
});
|
||||
const monitor = new HeartbeatMonitor({ store });
|
||||
|
||||
await monitor.completeRun("agent-001", "run-budget-001", {
|
||||
status: "completed",
|
||||
});
|
||||
|
||||
expect(store.getBudgetStatus).not.toHaveBeenCalled();
|
||||
expect(store.updateAgentState).toHaveBeenCalledWith("agent-001", "active");
|
||||
expect(store.updateAgent).not.toHaveBeenCalledWith("agent-001", { pauseReason: "budget-exhausted" });
|
||||
});
|
||||
});
|
||||
|
||||
describe("clearRunState", () => {
|
||||
it("resets accumulated task state for an agent", async () => {
|
||||
const savedRuns: Map<string, AgentHeartbeatRun> = new Map();
|
||||
@@ -1817,6 +2092,7 @@ describe("HeartbeatTriggerScheduler", () => {
|
||||
callback = vi.fn().mockResolvedValue(undefined);
|
||||
store = {
|
||||
getActiveHeartbeatRun: vi.fn().mockResolvedValue(null),
|
||||
getBudgetStatus: vi.fn().mockResolvedValue(createBudgetStatus()),
|
||||
on: vi.fn(),
|
||||
off: vi.fn(),
|
||||
} as unknown as AgentStore;
|
||||
@@ -1996,6 +2272,60 @@ describe("HeartbeatTriggerScheduler", () => {
|
||||
|
||||
expect(callback).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("skips timer tick when agent is over budget", async () => {
|
||||
(store.getBudgetStatus as ReturnType<typeof vi.fn>).mockResolvedValue(
|
||||
createBudgetStatus({ isOverBudget: true, isOverThreshold: true, usagePercent: 100 })
|
||||
);
|
||||
|
||||
scheduler.registerAgent("agent-001", { heartbeatIntervalMs: 5000 });
|
||||
await vi.advanceTimersByTimeAsync(5000);
|
||||
|
||||
expect(callback).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("skips timer tick when agent is over threshold", async () => {
|
||||
(store.getBudgetStatus as ReturnType<typeof vi.fn>).mockResolvedValue(
|
||||
createBudgetStatus({
|
||||
budgetLimit: 1000,
|
||||
usagePercent: 85,
|
||||
thresholdPercent: 80,
|
||||
isOverBudget: false,
|
||||
isOverThreshold: true,
|
||||
})
|
||||
);
|
||||
|
||||
scheduler.registerAgent("agent-001", { heartbeatIntervalMs: 5000 });
|
||||
await vi.advanceTimersByTimeAsync(5000);
|
||||
|
||||
expect(callback).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("fires timer tick normally when below threshold", async () => {
|
||||
(store.getBudgetStatus as ReturnType<typeof vi.fn>).mockResolvedValue(
|
||||
createBudgetStatus({
|
||||
budgetLimit: 1000,
|
||||
usagePercent: 30,
|
||||
thresholdPercent: 80,
|
||||
isOverBudget: false,
|
||||
isOverThreshold: false,
|
||||
})
|
||||
);
|
||||
|
||||
scheduler.registerAgent("agent-001", { heartbeatIntervalMs: 5000 });
|
||||
await vi.advanceTimersByTimeAsync(5000);
|
||||
|
||||
expect(callback).toHaveBeenCalledOnce();
|
||||
});
|
||||
|
||||
it("fires timer tick when getBudgetStatus throws", async () => {
|
||||
(store.getBudgetStatus as ReturnType<typeof vi.fn>).mockRejectedValue(new Error("budget unavailable"));
|
||||
|
||||
scheduler.registerAgent("agent-001", { heartbeatIntervalMs: 5000 });
|
||||
await vi.advanceTimersByTimeAsync(5000);
|
||||
|
||||
expect(callback).toHaveBeenCalledOnce();
|
||||
});
|
||||
});
|
||||
|
||||
describe("stop clears all timers", () => {
|
||||
@@ -2082,6 +2412,75 @@ describe("HeartbeatTriggerScheduler", () => {
|
||||
expect(callback).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("blocks assignment trigger when agent is over budget", async () => {
|
||||
(eventStore as any).getBudgetStatus = vi.fn().mockResolvedValue(
|
||||
createBudgetStatus({
|
||||
agentId: "agent-test",
|
||||
isOverBudget: true,
|
||||
isOverThreshold: true,
|
||||
usagePercent: 100,
|
||||
budgetLimit: 1000,
|
||||
thresholdPercent: 80,
|
||||
})
|
||||
);
|
||||
|
||||
const agent = { id: "agent-test", name: "Test" } as import("@fusion/core").Agent;
|
||||
eventStore.emit("agent:assigned", agent, "FN-003");
|
||||
|
||||
await new Promise((resolve) => setTimeout(resolve, 10));
|
||||
|
||||
expect(callback).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("allows assignment trigger when agent is over threshold", async () => {
|
||||
const budgetStatus = createBudgetStatus({
|
||||
agentId: "agent-test",
|
||||
budgetLimit: 1000,
|
||||
usagePercent: 85,
|
||||
thresholdPercent: 80,
|
||||
isOverBudget: false,
|
||||
isOverThreshold: true,
|
||||
});
|
||||
(eventStore as any).getBudgetStatus = vi.fn().mockResolvedValue(budgetStatus);
|
||||
|
||||
const agent = { id: "agent-test", name: "Test" } as import("@fusion/core").Agent;
|
||||
eventStore.emit("agent:assigned", agent, "FN-003");
|
||||
|
||||
await new Promise((resolve) => setTimeout(resolve, 10));
|
||||
|
||||
expect(callback).toHaveBeenCalledOnce();
|
||||
expect(callback).toHaveBeenCalledWith("agent-test", "assignment", {
|
||||
taskId: "FN-003",
|
||||
wakeReason: "assignment",
|
||||
triggerDetail: "task-assigned",
|
||||
budgetStatus,
|
||||
});
|
||||
});
|
||||
|
||||
it("passes budgetStatus in WakeContext for assignment triggers", async () => {
|
||||
const budgetStatus = createBudgetStatus({
|
||||
agentId: "agent-test",
|
||||
budgetLimit: 1000,
|
||||
usagePercent: 45,
|
||||
thresholdPercent: 80,
|
||||
});
|
||||
(eventStore as any).getBudgetStatus = vi.fn().mockResolvedValue(budgetStatus);
|
||||
|
||||
const agent = { id: "agent-test", name: "Test" } as import("@fusion/core").Agent;
|
||||
eventStore.emit("agent:assigned", agent, "FN-005");
|
||||
|
||||
await new Promise((resolve) => setTimeout(resolve, 10));
|
||||
|
||||
expect(callback).toHaveBeenCalledWith(
|
||||
"agent-test",
|
||||
"assignment",
|
||||
expect.objectContaining({
|
||||
taskId: "FN-005",
|
||||
budgetStatus,
|
||||
}),
|
||||
);
|
||||
});
|
||||
|
||||
it("cleans up listener on unwatch", async () => {
|
||||
scheduler.unwatchAssignments();
|
||||
|
||||
|
||||
@@ -17,7 +17,7 @@
|
||||
* - onTerminated: Called when an unresponsive agent is terminated
|
||||
*/
|
||||
|
||||
import type { AgentStore, AgentHeartbeatRun, HeartbeatInvocationSource, AgentHeartbeatConfig, Message, MessageStore, TaskStore, TaskDetail } from "@fusion/core";
|
||||
import type { AgentStore, AgentHeartbeatRun, HeartbeatInvocationSource, AgentHeartbeatConfig, AgentBudgetStatus, Message, MessageStore, TaskStore, TaskDetail } from "@fusion/core";
|
||||
import type { ToolDefinition } from "@mariozechner/pi-coding-agent";
|
||||
import { Type, type Static } from "@mariozechner/pi-ai";
|
||||
import { createTaskCreateTool, createTaskLogTool, taskCreateParams } from "./agent-tools.js";
|
||||
@@ -308,24 +308,25 @@ export class HeartbeatMonitor {
|
||||
if (!run) return;
|
||||
|
||||
const tracked = this.trackedAgents.get(agentId);
|
||||
let completionResult = result;
|
||||
|
||||
// Merge accumulated task creations into resultJson
|
||||
const createdTasks = this.runCreatedTasks.get(agentId);
|
||||
const enrichedResultJson = createdTasks?.length
|
||||
? { ...result.resultJson, tasksCreated: createdTasks }
|
||||
: result.resultJson;
|
||||
? { ...completionResult.resultJson, tasksCreated: createdTasks }
|
||||
: completionResult.resultJson;
|
||||
|
||||
const completedRun: AgentHeartbeatRun = {
|
||||
...run,
|
||||
endedAt: new Date().toISOString(),
|
||||
status: result.status,
|
||||
exitCode: result.exitCode,
|
||||
status: completionResult.status,
|
||||
exitCode: completionResult.exitCode,
|
||||
sessionIdBefore: tracked?.sessionIdBefore,
|
||||
sessionIdAfter: result.sessionIdAfter,
|
||||
usageJson: result.usageJson,
|
||||
sessionIdAfter: completionResult.sessionIdAfter,
|
||||
usageJson: completionResult.usageJson,
|
||||
resultJson: enrichedResultJson,
|
||||
stdoutExcerpt: result.stdoutExcerpt,
|
||||
stderrExcerpt: result.stderrExcerpt,
|
||||
stdoutExcerpt: completionResult.stdoutExcerpt,
|
||||
stderrExcerpt: completionResult.stderrExcerpt,
|
||||
};
|
||||
|
||||
await this.store.saveRun(completedRun);
|
||||
@@ -334,13 +335,13 @@ export class HeartbeatMonitor {
|
||||
this.clearRunState(agentId);
|
||||
|
||||
// Update cumulative usage on agent
|
||||
if (result.usageJson) {
|
||||
if (completionResult.usageJson) {
|
||||
try {
|
||||
const agent = await this.store.getAgent(agentId);
|
||||
if (agent) {
|
||||
await this.store.updateAgent(agentId, {
|
||||
totalInputTokens: (agent.totalInputTokens ?? 0) + result.usageJson.inputTokens,
|
||||
totalOutputTokens: (agent.totalOutputTokens ?? 0) + result.usageJson.outputTokens,
|
||||
totalInputTokens: (agent.totalInputTokens ?? 0) + completionResult.usageJson.inputTokens,
|
||||
totalOutputTokens: (agent.totalOutputTokens ?? 0) + completionResult.usageJson.outputTokens,
|
||||
});
|
||||
}
|
||||
} catch {
|
||||
@@ -348,13 +349,29 @@ export class HeartbeatMonitor {
|
||||
}
|
||||
}
|
||||
|
||||
// Transition agent state based on result
|
||||
if (!result.skipStateTransition) {
|
||||
// Budget governance: pause agent if over budget after usage update
|
||||
if (completionResult.usageJson && completionResult.status !== "failed" && completionResult.status !== "terminated") {
|
||||
try {
|
||||
if (result.status === "failed") {
|
||||
const budgetStatus = await this.store.getBudgetStatus(agentId);
|
||||
if (budgetStatus.isOverBudget) {
|
||||
heartbeatLog.log(`Agent ${agentId} is over budget — pausing with reason "budget-exhausted"`);
|
||||
await this.store.updateAgentState(agentId, "paused");
|
||||
await this.store.updateAgent(agentId, { pauseReason: "budget-exhausted" });
|
||||
// Skip the normal state transition below since we already set the correct state
|
||||
completionResult = { ...completionResult, skipStateTransition: true };
|
||||
}
|
||||
} catch {
|
||||
// If budget check fails, proceed with normal state transition
|
||||
}
|
||||
}
|
||||
|
||||
// Transition agent state based on result
|
||||
if (!completionResult.skipStateTransition) {
|
||||
try {
|
||||
if (completionResult.status === "failed") {
|
||||
await this.store.updateAgentState(agentId, "error");
|
||||
await this.store.updateAgent(agentId, { lastError: result.stderrExcerpt ?? "Run failed" });
|
||||
} else if (result.status === "terminated") {
|
||||
await this.store.updateAgent(agentId, { lastError: completionResult.stderrExcerpt ?? "Run failed" });
|
||||
} else if (completionResult.status === "terminated") {
|
||||
await this.store.updateAgentState(agentId, "terminated");
|
||||
} else {
|
||||
// Completed successfully - back to active
|
||||
@@ -366,7 +383,7 @@ export class HeartbeatMonitor {
|
||||
}
|
||||
|
||||
// End the heartbeat run tracking
|
||||
await this.store.endHeartbeatRun(runId, result.status === "completed" ? "completed" : "terminated");
|
||||
await this.store.endHeartbeatRun(runId, completionResult.status === "completed" ? "completed" : "terminated");
|
||||
|
||||
this.onRunCompleted?.(agentId, completedRun);
|
||||
}
|
||||
@@ -476,6 +493,11 @@ export class HeartbeatMonitor {
|
||||
* 3. Work — run a lightweight agent session with readonly tools + task_create/task_log
|
||||
* 4. Exit — record results and complete the run
|
||||
*
|
||||
* Budget governance:
|
||||
* - Skip all triggers when the agent is over budget (`isOverBudget`)
|
||||
* - Skip timer triggers when over the warning threshold (`isOverThreshold`)
|
||||
* - Continue normal execution for critical triggers (assignment/on_demand) when only over threshold
|
||||
*
|
||||
* Per-agent execution is serialized via `withAgentStartLock` — concurrent calls
|
||||
* for the same agent wait for the previous run to complete.
|
||||
*
|
||||
@@ -501,6 +523,32 @@ export class HeartbeatMonitor {
|
||||
const run = await this.startRun(agentId, { source, triggerDetail, contextSnapshot });
|
||||
|
||||
try {
|
||||
// Budget governance: check if agent can run
|
||||
try {
|
||||
const budgetStatus = await this.store.getBudgetStatus(agentId);
|
||||
if (budgetStatus.isOverBudget) {
|
||||
heartbeatLog.log(`Agent ${agentId} budget exhausted — heartbeat skipped`);
|
||||
await this.completeRun(agentId, run.id, {
|
||||
status: "completed",
|
||||
resultJson: { reason: "budget_exhausted", budgetStatus },
|
||||
skipStateTransition: true,
|
||||
});
|
||||
return (await this.store.getRunDetail(agentId, run.id))!;
|
||||
}
|
||||
// Above threshold: only allow critical triggers (assignment, on_demand)
|
||||
if (budgetStatus.isOverThreshold && source === "timer") {
|
||||
heartbeatLog.log(`Agent ${agentId} over budget threshold (${budgetStatus.usagePercent}%) — timer heartbeat skipped`);
|
||||
await this.completeRun(agentId, run.id, {
|
||||
status: "completed",
|
||||
resultJson: { reason: "budget_threshold_exceeded", budgetStatus },
|
||||
skipStateTransition: true,
|
||||
});
|
||||
return (await this.store.getRunDetail(agentId, run.id))!;
|
||||
}
|
||||
} catch {
|
||||
// If getBudgetStatus fails (e.g., method not available), proceed without budget check
|
||||
}
|
||||
|
||||
// Resolve agent
|
||||
const agent = await this.store.getAgent(agentId);
|
||||
if (!agent) {
|
||||
@@ -838,6 +886,8 @@ export interface WakeContext {
|
||||
wakeReason: string;
|
||||
/** Detail about the specific trigger */
|
||||
triggerDetail: string;
|
||||
/** Budget governance status for the agent at trigger time */
|
||||
budgetStatus?: AgentBudgetStatus;
|
||||
/** Additional context (intervalMs, etc.) */
|
||||
[key: string]: unknown;
|
||||
}
|
||||
@@ -998,11 +1048,24 @@ export class HeartbeatTriggerScheduler {
|
||||
return;
|
||||
}
|
||||
|
||||
let budgetStatus: AgentBudgetStatus | undefined;
|
||||
// Budget governance: block even critical triggers when budget is fully exhausted
|
||||
try {
|
||||
budgetStatus = await this.store.getBudgetStatus(agent.id);
|
||||
if (budgetStatus.isOverBudget) {
|
||||
heartbeatLog.log(`Agent ${agent.id} budget exhausted — assignment trigger skipped`);
|
||||
return;
|
||||
}
|
||||
} catch {
|
||||
// If getBudgetStatus fails, proceed without budget check
|
||||
}
|
||||
|
||||
heartbeatLog.log(`Assignment trigger for ${agent.id} (task: ${taskId})`);
|
||||
await this.callback(agent.id, "assignment", {
|
||||
taskId,
|
||||
wakeReason: "assignment",
|
||||
triggerDetail: "task-assigned",
|
||||
...(budgetStatus && { budgetStatus }),
|
||||
});
|
||||
} catch (err) {
|
||||
heartbeatLog.error(`Assignment trigger error for ${agent.id}: ${err instanceof Error ? err.message : err}`);
|
||||
@@ -1039,6 +1102,21 @@ export class HeartbeatTriggerScheduler {
|
||||
return;
|
||||
}
|
||||
|
||||
// Budget governance: skip timer triggers for over-budget agents
|
||||
try {
|
||||
const budgetStatus = await this.store.getBudgetStatus(agentId);
|
||||
if (budgetStatus.isOverBudget) {
|
||||
heartbeatLog.log(`Agent ${agentId} budget exhausted — timer tick skipped`);
|
||||
return;
|
||||
}
|
||||
if (budgetStatus.isOverThreshold) {
|
||||
heartbeatLog.log(`Agent ${agentId} over budget threshold (${budgetStatus.usagePercent}%) — timer tick skipped`);
|
||||
return;
|
||||
}
|
||||
} catch {
|
||||
// If getBudgetStatus fails, proceed without budget check
|
||||
}
|
||||
|
||||
await this.callback(agentId, "timer", {
|
||||
wakeReason: "timer",
|
||||
triggerDetail: "scheduled",
|
||||
|
||||
Reference in New Issue
Block a user