#!/usr/bin/env node import fs from 'node:fs'; import path from 'node:path'; import { fileURLToPath } from 'node:url'; import { createAgentRunGateway } from '../agent-run-gateway.mjs'; import { createDbPool } from '../db.mjs'; import { createLlmProviderService } from '../llm-providers.mjs'; import { createTkmindProxy } from '../tkmind-proxy.mjs'; import { createToolGateway } from '../tool-gateway.mjs'; import { createUserAuth } from '../user-auth.mjs'; const root = path.join(path.dirname(fileURLToPath(import.meta.url)), '..'); function loadEnvFile(filePath) { if (!fs.existsSync(filePath)) return; for (const line of fs.readFileSync(filePath, 'utf8').split('\n')) { const trimmed = line.trim(); if (!trimmed || trimmed.startsWith('#')) continue; const eq = trimmed.indexOf('='); if (eq < 0) continue; const key = trimmed.slice(0, eq).trim(); const value = trimmed.slice(eq + 1).trim(); if (!process.env[key]) process.env[key] = value; } } function parseTargets() { const csv = String(process.env.TKMIND_API_TARGETS ?? '').trim(); if (csv) { return csv.split(',').map((item) => item.trim()).filter(Boolean); } return [process.env.TKMIND_API_TARGET ?? 'https://127.0.0.1:18006']; } function parseArgs(argv) { const args = { once: false, status: false, recoverStale: false, applyRecovery: false, runId: '', pollMs: Number(process.env.MEMIND_AGENT_RUN_WORKER_POLL_MS ?? 1000), limit: Number(process.env.MEMIND_AGENT_RUN_WORKER_BATCH_SIZE ?? 1), staleMs: Number(process.env.MEMIND_AGENT_RUN_STALE_RUNNING_MS ?? process.env.MEMIND_AGENT_RUN_TIMEOUT_MS ?? 15 * 60 * 1000), }; for (let i = 0; i < argv.length; i += 1) { const item = argv[i]; if (item === '--once') args.once = true; else if (item === '--status') args.status = true; else if (item === '--recover-stale') args.recoverStale = true; else if (item === '--apply-recovery') { args.recoverStale = true; args.applyRecovery = true; } else if (item === '--run-id') args.runId = String(argv[++i] ?? '').trim(); else if (item === '--poll-ms') args.pollMs = Number(argv[++i] ?? args.pollMs); else if (item === '--limit') args.limit = Number(argv[++i] ?? args.limit); else if (item === '--stale-ms') args.staleMs = Number(argv[++i] ?? args.staleMs); else if (item === '--help' || item === '-h') args.help = true; } return args; } function printHelp() { console.log([ 'Usage:', ' node scripts/agent-run-worker.mjs [--once] [--status] [--run-id ] [--poll-ms 1000] [--limit 1]', ' node scripts/agent-run-worker.mjs --recover-stale [--apply-recovery] [--stale-ms 900000]', '', 'Notes:', ' - Processes existing h5_agent_runs queued/retryable rows through the Tool Gateway queue.', ' - Regular dispatch automatically fails stale running rows older than MEMIND_AGENT_RUN_TIMEOUT_MS.', ' - --recover-stale is dry-run by default; --apply-recovery marks stale running rows failed.', ' - Does not create agent runs by itself.', ' - --run-id dispatches exactly one known run and avoids scanning the global queue.', ' - Set MEMIND_AGENT_RUN_AUTODISPATCH=0 on Portal only when this worker is ready to take over.', ].join('\n')); } loadEnvFile(process.env.MEMIND_ENV_FILE || path.join(root, '.env')); loadEnvFile(path.join(root, '.env.local')); const args = parseArgs(process.argv.slice(2)); if (args.help) { printHelp(); process.exit(0); } const pool = createDbPool(); const baseUserAuth = createUserAuth(pool, { usersRoot: process.env.H5_USERS_ROOT ?? path.join(root, 'users'), h5Root: root, }); const overrideWorkingDir = String(process.env.MEMIND_AGENT_RUN_WORKDIR_OVERRIDE ?? '').trim(); const overrideWorkingDirUserId = String(process.env.MEMIND_AGENT_RUN_WORKDIR_USER_ID ?? '').trim(); const userAuth = overrideWorkingDir ? { ...baseUserAuth, async resolveWorkingDir(userId) { if (!overrideWorkingDirUserId || overrideWorkingDirUserId === userId) { return path.resolve(overrideWorkingDir); } return baseUserAuth.resolveWorkingDir(userId); }, } : baseUserAuth; const llmProviderService = createLlmProviderService(pool, { apiTarget: process.env.TKMIND_API_TARGET ?? 'https://127.0.0.1:18006', apiSecret: process.env.TKMIND_SERVER__SECRET_KEY ?? 'local-dev-secret', }); const toolGateway = createToolGateway({ llmProviderService }); const apiTargets = parseTargets(); const tkmindProxy = createTkmindProxy({ apiTarget: apiTargets[0], apiTargets, apiSecret: process.env.TKMIND_SERVER__SECRET_KEY ?? 'local-dev-secret', userAuth, }); const gateway = createAgentRunGateway({ pool, userAuth, tkmindProxy, toolGateway, autoDispatch: false, maxConcurrentRuns: Number(process.env.MEMIND_AGENT_RUN_QUEUE_CONCURRENCY ?? 1), runTimeoutMs: Number(process.env.MEMIND_AGENT_RUN_TIMEOUT_MS ?? 15 * 60 * 1000), }); let stopping = false; process.on('SIGINT', () => { stopping = true; }); process.on('SIGTERM', () => { stopping = true; }); async function runOnce() { let result = null; if (args.runId) { gateway.dispatchRun(args.runId); result = { considered: 1, dispatched: 1, runId: args.runId, }; } else { result = await gateway.dispatchQueuedRuns({ limit: args.limit }); } await waitForIdle(); console.log(JSON.stringify({ ok: true, mode: args.status ? 'status' : 'dispatch', runId: args.runId || null, dispatched: args.status ? 0 : result.dispatched, queue: await gateway.getQueueStatus(), })); } async function runStaleRecovery() { const result = await gateway.recoverStaleRunningRuns({ staleMs: args.staleMs, limit: args.limit, dryRun: !args.applyRecovery, reason: args.applyRecovery ? 'manual_stale_recovery' : 'manual_stale_recovery_dry_run', }); console.log(JSON.stringify({ ok: true, mode: args.applyRecovery ? 'recover-stale-apply' : 'recover-stale-dry-run', recovery: result, queue: await gateway.getQueueStatus(), }, null, 2)); } 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)); } else if (args.recoverStale) { await runStaleRecovery(); } else if (args.once || args.runId) { await runOnce(); } else { console.log(JSON.stringify({ ok: true, worker: 'agent-run-worker', pollMs: args.pollMs, limit: args.limit })); while (!stopping) { await runOnce(); 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); }