feat: add worker heartbeat for agent runs
This commit is contained in:
+131
-13
@@ -9,6 +9,13 @@ function createFakePool() {
|
||||
const runs = new Map();
|
||||
const events = [];
|
||||
|
||||
const latestHeartbeatAt = (runId) => {
|
||||
const timestamps = events
|
||||
.filter((event) => event.runId === runId && event.eventType === 'worker_heartbeat')
|
||||
.map((event) => Number(event.createdAt ?? 0));
|
||||
return timestamps.length ? Math.max(...timestamps) : null;
|
||||
};
|
||||
|
||||
return {
|
||||
runs,
|
||||
events,
|
||||
@@ -33,24 +40,45 @@ function createFakePool() {
|
||||
}
|
||||
return [[...counts].map(([status, count]) => ({ status, count }))];
|
||||
}
|
||||
if (sql.includes('SELECT MIN(started_at) AS oldest_started_at')) {
|
||||
const started = [...runs.values()]
|
||||
.filter((row) => row.status === 'running' && row.started_at != null)
|
||||
.map((row) => Number(row.started_at));
|
||||
return [[{ oldest_started_at: started.length ? Math.min(...started) : null }]];
|
||||
}
|
||||
if (sql.includes('SELECT id, request_id, started_at, updated_at, attempts')) {
|
||||
const [cutoff, limit = 1] = params;
|
||||
if (sql.includes('h.latest_heartbeat_at') && sql.includes("WHERE r.status = 'running'") && sql.includes('LIMIT 1')) {
|
||||
return [[...runs.values()]
|
||||
.filter((row) => row.status === 'running' && row.started_at != null && Number(row.started_at) <= Number(cutoff))
|
||||
.sort((a, b) => Number(a.started_at) - Number(b.started_at))
|
||||
.slice(0, Number(limit))
|
||||
.filter((row) => row.status === 'running')
|
||||
.map((row) => ({
|
||||
id: row.id,
|
||||
request_id: row.request_id,
|
||||
started_at: row.started_at,
|
||||
updated_at: row.updated_at,
|
||||
attempts: row.attempts,
|
||||
latest_heartbeat_at: latestHeartbeatAt(row.id),
|
||||
}))
|
||||
.sort((a, b) => Number(a.latest_heartbeat_at ?? a.started_at ?? 0) - Number(b.latest_heartbeat_at ?? b.started_at ?? 0))
|
||||
.slice(0, 1)];
|
||||
}
|
||||
if (sql.includes('SELECT COUNT(*) AS count') && sql.includes('worker_heartbeat')) {
|
||||
return [[{
|
||||
count: [...runs.values()]
|
||||
.filter((row) => row.status === 'running' && latestHeartbeatAt(row.id) == null)
|
||||
.length,
|
||||
}]];
|
||||
}
|
||||
if (sql.includes('SELECT') && sql.includes('latest_heartbeat_at') && sql.includes('COALESCE(h.latest_heartbeat_at, r.started_at) <= ?')) {
|
||||
const [cutoff, limit = 1] = params;
|
||||
return [[...runs.values()]
|
||||
.map((row) => ({ row, heartbeatAt: latestHeartbeatAt(row.id) }))
|
||||
.filter(({ row, heartbeatAt }) => (
|
||||
row.status === 'running' &&
|
||||
row.started_at != null &&
|
||||
Number(heartbeatAt ?? row.started_at) <= Number(cutoff)
|
||||
))
|
||||
.sort((a, b) => Number(a.heartbeatAt ?? a.row.started_at) - Number(b.heartbeatAt ?? b.row.started_at))
|
||||
.slice(0, Number(limit))
|
||||
.map(({ row, heartbeatAt }) => ({
|
||||
id: row.id,
|
||||
request_id: row.request_id,
|
||||
started_at: row.started_at,
|
||||
updated_at: row.updated_at,
|
||||
attempts: row.attempts,
|
||||
latest_heartbeat_at: heartbeatAt,
|
||||
}))];
|
||||
}
|
||||
if (sql.includes('SELECT id') && sql.includes("status IN ('queued', 'retryable')")) {
|
||||
@@ -107,10 +135,15 @@ function createFakePool() {
|
||||
});
|
||||
return [{ affectedRows: 1 }];
|
||||
}
|
||||
if (sql.includes("WHERE id = ? AND status = 'running' AND started_at <= ?")) {
|
||||
if (sql.includes("WHERE id = ?") && sql.includes("status = 'running'") && sql.includes('worker_heartbeat')) {
|
||||
const [errorMessage, updatedAt, completedAt, id, cutoff] = params;
|
||||
const row = runs.get(id);
|
||||
if (!row || row.status !== 'running' || Number(row.started_at) > Number(cutoff)) {
|
||||
if (
|
||||
!row ||
|
||||
row.status !== 'running' ||
|
||||
row.started_at == null ||
|
||||
Number(latestHeartbeatAt(id) ?? row.started_at) > Number(cutoff)
|
||||
) {
|
||||
return [{ affectedRows: 0 }];
|
||||
}
|
||||
Object.assign(row, {
|
||||
@@ -598,6 +631,7 @@ test('external worker dispatches queued runs through the same queue controls', a
|
||||
|
||||
const result = await gateway.dispatchQueuedRuns({ limit: 10 });
|
||||
assert.equal(result.dispatched, 1);
|
||||
await waitFor(() => pool.events.some((event) => event.runId === run.id && event.eventType === 'worker_heartbeat'));
|
||||
await waitFor(() => pool.runs.get(run.id)?.status === 'succeeded');
|
||||
assert.deepEqual(submitted, ['req-worker']);
|
||||
});
|
||||
@@ -710,3 +744,87 @@ test('stale running recovery marks old running rows failed with an event', async
|
||||
true,
|
||||
);
|
||||
});
|
||||
|
||||
test('stale running recovery ignores old running rows with fresh heartbeat', async () => {
|
||||
const pool = createFakePool();
|
||||
const gateway = createAgentRunGateway({
|
||||
pool,
|
||||
userAuth: {},
|
||||
tkmindProxy: {},
|
||||
autoDispatch: false,
|
||||
runTimeoutMs: 1000,
|
||||
});
|
||||
|
||||
const run = await gateway.createRun('user-1', {
|
||||
requestId: 'req-stale-heartbeat-fresh',
|
||||
userMessage: { role: 'user', content: [] },
|
||||
});
|
||||
Object.assign(pool.runs.get(run.id), {
|
||||
status: 'running',
|
||||
attempts: 1,
|
||||
started_at: Date.now() - 5000,
|
||||
updated_at: Date.now() - 5000,
|
||||
});
|
||||
pool.events.push({
|
||||
id: 'heartbeat-1',
|
||||
runId: run.id,
|
||||
eventType: 'worker_heartbeat',
|
||||
dataJson: JSON.stringify({ attempt: 1 }),
|
||||
createdAt: Date.now(),
|
||||
});
|
||||
|
||||
const result = await gateway.recoverStaleRunningRuns({ staleMs: 1000, dryRun: false });
|
||||
|
||||
assert.equal(result.considered, 0);
|
||||
assert.equal(result.recovered, 0);
|
||||
assert.equal(pool.runs.get(run.id).status, 'running');
|
||||
});
|
||||
|
||||
test('queue status reports running heartbeat age and missing heartbeat count', async () => {
|
||||
const pool = createFakePool();
|
||||
const gateway = createAgentRunGateway({
|
||||
pool,
|
||||
userAuth: {},
|
||||
tkmindProxy: {},
|
||||
autoDispatch: false,
|
||||
runTimeoutMs: 1000,
|
||||
heartbeatMs: 250,
|
||||
});
|
||||
|
||||
const runWithHeartbeat = await gateway.createRun('user-1', {
|
||||
requestId: 'req-heartbeat-status-1',
|
||||
userMessage: { role: 'user', content: [] },
|
||||
});
|
||||
Object.assign(pool.runs.get(runWithHeartbeat.id), {
|
||||
status: 'running',
|
||||
attempts: 1,
|
||||
started_at: Date.now() - 5000,
|
||||
updated_at: Date.now() - 5000,
|
||||
});
|
||||
pool.events.push({
|
||||
id: 'heartbeat-status-1',
|
||||
runId: runWithHeartbeat.id,
|
||||
eventType: 'worker_heartbeat',
|
||||
dataJson: JSON.stringify({ attempt: 1 }),
|
||||
createdAt: Date.now() - 100,
|
||||
});
|
||||
|
||||
const runWithoutHeartbeat = await gateway.createRun('user-1', {
|
||||
requestId: 'req-heartbeat-status-2',
|
||||
userMessage: { role: 'user', content: [] },
|
||||
});
|
||||
Object.assign(pool.runs.get(runWithoutHeartbeat.id), {
|
||||
status: 'running',
|
||||
attempts: 1,
|
||||
started_at: Date.now() - 2000,
|
||||
updated_at: Date.now() - 2000,
|
||||
});
|
||||
|
||||
const status = await gateway.getQueueStatus();
|
||||
|
||||
assert.equal(status.heartbeatMs, 250);
|
||||
assert.equal(status.runningWithoutHeartbeatCount, 1);
|
||||
assert.equal(status.latestRunningRun.id, runWithoutHeartbeat.id);
|
||||
assert.equal(status.oldestRunningHeartbeatAt, null);
|
||||
assert.ok(status.oldestRunningHeartbeatAgeMs >= 1900);
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user