diff --git a/.runtime/portal/scripts/agent-run-worker.mjs b/.runtime/portal/scripts/agent-run-worker.mjs index 20fd3ad..1c16130 100755 --- a/.runtime/portal/scripts/agent-run-worker.mjs +++ b/.runtime/portal/scripts/agent-run-worker.mjs @@ -286,7 +286,7 @@ function createAgentRunGateway({ error_message: message, completed_at: retryable ? null : nowMs() }); - if (retryable) { + if (retryable && autoDispatch) { setTimeout(() => dispatchRun(runId), retryDelaysMs[nextAttempt - 1]); } } @@ -6888,13 +6888,27 @@ process.on("SIGTERM", () => { }); async function runOnce() { const result = await gateway.dispatchQueuedRuns({ limit: args.limit }); + await waitForIdle(); console.log(JSON.stringify({ ok: true, mode: args.status ? "status" : "dispatch", dispatched: args.status ? 0 : result.dispatched, - queue: result.queue + queue: await gateway.getQueueStatus() })); } +async function waitForIdle() { + const startedAt = Date.now(); + const maxWaitMs = Number(process.env.MEMIND_AGENT_RUN_WORKER_DRAIN_WAIT_MS ?? 20 * 60 * 1e3); + while (!stopping) { + const status = await gateway.getQueueStatus(); + if (status.inFlight === 0 && status.pendingDispatches === 0) return; + if (Date.now() - startedAt > maxWaitMs) { + throw new Error(`agent run worker did not drain within ${maxWaitMs}ms`); + } + await new Promise((resolve) => setTimeout(resolve, 500)); + } +} +var exitCode = 0; try { if (args.status) { console.log(JSON.stringify({ ok: true, queue: await gateway.getQueueStatus() }, null, 2)); @@ -6907,7 +6921,11 @@ try { await new Promise((resolve) => setTimeout(resolve, Math.max(200, args.pollMs))); } } +} catch (err) { + exitCode = 1; + console.error(err instanceof Error ? err.message : String(err)); } finally { await pool.end().catch(() => { }); + if (args.status || args.once) process.exit(exitCode); } diff --git a/.runtime/portal/server.mjs b/.runtime/portal/server.mjs index 24c8632..6084dec 100644 --- a/.runtime/portal/server.mjs +++ b/.runtime/portal/server.mjs @@ -6783,7 +6783,7 @@ function createAgentRunGateway({ error_message: message, completed_at: retryable ? null : nowMs() }); - if (retryable) { + if (retryable && autoDispatch) { setTimeout(() => dispatchRun(runId), retryDelaysMs[nextAttempt - 1]); } } diff --git a/agent-run-gateway.mjs b/agent-run-gateway.mjs index 8299dda..8e3efde 100644 --- a/agent-run-gateway.mjs +++ b/agent-run-gateway.mjs @@ -310,7 +310,7 @@ export function createAgentRunGateway({ error_message: message, completed_at: retryable ? null : nowMs(), }); - if (retryable) { + if (retryable && autoDispatch) { setTimeout(() => dispatchRun(runId), retryDelaysMs[nextAttempt - 1]); } } diff --git a/scripts/agent-run-worker.mjs b/scripts/agent-run-worker.mjs index 8391f1e..c0ca476 100644 --- a/scripts/agent-run-worker.mjs +++ b/scripts/agent-run-worker.mjs @@ -96,14 +96,29 @@ process.on('SIGTERM', () => { stopping = true; }); async function runOnce() { const result = await gateway.dispatchQueuedRuns({ limit: args.limit }); + await waitForIdle(); console.log(JSON.stringify({ ok: true, mode: args.status ? 'status' : 'dispatch', dispatched: args.status ? 0 : result.dispatched, - queue: result.queue, + queue: await gateway.getQueueStatus(), })); } +async function waitForIdle() { + const startedAt = Date.now(); + const maxWaitMs = Number(process.env.MEMIND_AGENT_RUN_WORKER_DRAIN_WAIT_MS ?? 20 * 60 * 1000); + while (!stopping) { + const status = await gateway.getQueueStatus(); + if (status.inFlight === 0 && status.pendingDispatches === 0) return; + if (Date.now() - startedAt > maxWaitMs) { + throw new Error(`agent run worker did not drain within ${maxWaitMs}ms`); + } + await new Promise((resolve) => setTimeout(resolve, 500)); + } +} + +let exitCode = 0; try { if (args.status) { console.log(JSON.stringify({ ok: true, queue: await gateway.getQueueStatus() }, null, 2)); @@ -116,6 +131,10 @@ try { await new Promise((resolve) => setTimeout(resolve, Math.max(200, args.pollMs))); } } +} catch (err) { + exitCode = 1; + console.error(err instanceof Error ? err.message : String(err)); } finally { await pool.end().catch(() => {}); + if (args.status || args.once) process.exit(exitCode); }