149 lines
5.0 KiB
JavaScript
149 lines
5.0 KiB
JavaScript
#!/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,
|
|
pollMs: Number(process.env.MEMIND_AGENT_RUN_WORKER_POLL_MS ?? 1000),
|
|
limit: Number(process.env.MEMIND_AGENT_RUN_WORKER_BATCH_SIZE ?? 1),
|
|
};
|
|
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 === '--poll-ms') args.pollMs = Number(argv[++i] ?? args.pollMs);
|
|
else if (item === '--limit') args.limit = Number(argv[++i] ?? args.limit);
|
|
else if (item === '--help' || item === '-h') args.help = true;
|
|
}
|
|
return args;
|
|
}
|
|
|
|
function printHelp() {
|
|
console.log([
|
|
'Usage:',
|
|
' node scripts/agent-run-worker.mjs [--once] [--status] [--poll-ms 1000] [--limit 1]',
|
|
'',
|
|
'Notes:',
|
|
' - Processes existing h5_agent_runs queued/retryable rows through the Tool Gateway queue.',
|
|
' - Does not create agent runs by itself.',
|
|
' - 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 userAuth = createUserAuth(pool, {
|
|
usersRoot: process.env.H5_USERS_ROOT ?? path.join(root, 'users'),
|
|
h5Root: root,
|
|
});
|
|
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() {
|
|
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: 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));
|
|
} else if (args.once) {
|
|
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);
|
|
}
|