feat: add external agent run worker entry
This commit is contained in:
@@ -30,6 +30,14 @@ function createFakePool() {
|
||||
}
|
||||
return [[...counts].map(([status, count]) => ({ status, count }))];
|
||||
}
|
||||
if (sql.includes('SELECT id') && sql.includes("status IN ('queued', 'retryable')")) {
|
||||
const limit = Number(params[0] ?? 1);
|
||||
return [[...runs.values()]
|
||||
.filter((row) => ['queued', 'retryable'].includes(row.status))
|
||||
.sort((a, b) => Number(a.updated_at) - Number(b.updated_at))
|
||||
.slice(0, limit)
|
||||
.map((row) => ({ id: row.id }))];
|
||||
}
|
||||
if (sql.includes('INSERT INTO h5_agent_runs')) {
|
||||
const [
|
||||
id,
|
||||
@@ -338,3 +346,77 @@ test('agent run queue status reports active database and local queue state', asy
|
||||
assert.equal(status.pendingDispatches, 0);
|
||||
assert.equal(status.statusCounts.queued, 1);
|
||||
});
|
||||
|
||||
test('external worker dispatches queued runs through the same queue controls', async () => {
|
||||
const pool = createFakePool();
|
||||
const submitted = [];
|
||||
const gateway = createAgentRunGateway({
|
||||
pool,
|
||||
userAuth: {},
|
||||
tkmindProxy: {
|
||||
async startSessionForUser() {
|
||||
return { id: 'session-worker' };
|
||||
},
|
||||
async submitSessionReplyForUser(_userId, _sessionId, requestId) {
|
||||
submitted.push(requestId);
|
||||
},
|
||||
},
|
||||
autoDispatch: false,
|
||||
retryDelaysMs: [],
|
||||
maxConcurrentRuns: 1,
|
||||
});
|
||||
|
||||
const run = await gateway.createRun('user-1', {
|
||||
requestId: 'req-worker',
|
||||
userMessage: { role: 'user', content: [] },
|
||||
});
|
||||
assert.equal(pool.runs.get(run.id).status, 'queued');
|
||||
|
||||
const result = await gateway.dispatchQueuedRuns({ limit: 10 });
|
||||
assert.equal(result.dispatched, 1);
|
||||
await waitFor(() => pool.runs.get(run.id)?.status === 'succeeded');
|
||||
assert.deepEqual(submitted, ['req-worker']);
|
||||
});
|
||||
|
||||
test('external worker does not dispatch more runs when local queue is full', async () => {
|
||||
const pool = createFakePool();
|
||||
const release = [];
|
||||
const gateway = createAgentRunGateway({
|
||||
pool,
|
||||
userAuth: {},
|
||||
tkmindProxy: {
|
||||
async startSessionForUser() {
|
||||
return { id: `session-${release.length + 1}` };
|
||||
},
|
||||
async submitSessionReplyForUser() {
|
||||
await new Promise((resolve) => release.push(resolve));
|
||||
},
|
||||
},
|
||||
autoDispatch: false,
|
||||
retryDelaysMs: [],
|
||||
maxConcurrentRuns: 1,
|
||||
});
|
||||
|
||||
const run1 = await gateway.createRun('user-1', {
|
||||
requestId: 'req-full-1',
|
||||
userMessage: { role: 'user', content: [] },
|
||||
});
|
||||
const run2 = await gateway.createRun('user-1', {
|
||||
requestId: 'req-full-2',
|
||||
userMessage: { role: 'user', content: [] },
|
||||
});
|
||||
|
||||
try {
|
||||
const first = await gateway.dispatchQueuedRuns({ limit: 10 });
|
||||
assert.equal(first.dispatched, 1);
|
||||
await waitFor(() => pool.runs.get(run1.id)?.status === 'running');
|
||||
assert.equal((await gateway.getQueueStatus()).inFlight, 1);
|
||||
const second = await gateway.dispatchQueuedRuns({ limit: 10 });
|
||||
assert.equal(second.dispatched, 0);
|
||||
assert.equal(pool.runs.get(run2.id).status, 'queued');
|
||||
} finally {
|
||||
while (release.length > 0) release.shift()();
|
||||
}
|
||||
await waitFor(() => pool.runs.get(run1.id)?.status === 'succeeded');
|
||||
assert.equal(pool.runs.get(run2.id).status, 'queued');
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user