Files
memind/scripts/agent-run-worker.mjs
T
john 1d165bc6e3 feat(goose): complete v1.49 phase3 closeout gates and context fusion plan
Expand Goose v1.49 smoke coverage (memory chat, portal resume, page e2e,
multiturn provider), add canary memory policy lock, refresh baselines, and
ignore one-off evidence artifacts. Document headroom-based context runtime
fusion plan; include auth, scheduled-task, and wechat intent fixes on branch.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-09 18:11:39 +08:00

327 lines
12 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 { createManagedChatIntentRouter } from '../chat-intent-router.mjs';
import { createConversationMemoryService } from '../conversation-memory.mjs';
import { createDbPool } from '../db.mjs';
import { createEpisodicMemoryService } from '../episodic-memory.mjs';
import { createLlmProviderService } from '../llm-providers.mjs';
import { createMemoryV2AdminConfigService } from '../memory-v2-admin-config.mjs';
import { createManagedMemoryV2Runtime } from '../memory-v2-runtime.mjs';
import { createTkmindProxy } from '../tkmind-proxy.mjs';
import { createToolGateway } from '../tool-gateway.mjs';
import { createUserAuth } from '../user-auth.mjs';
import { createSessionSnapshotService } from '../session-snapshot.mjs';
import { createDirectChatService } from '../direct-chat-service.mjs';
import { createSessionAccess } from '../session-broker.mjs';
import { createOrchestratorAdminConfigService } from '../services/orchestrator/admin-config.mjs';
import { createWorkflowShadowObserver } from '../services/orchestrator/shadow-observer.mjs';
import { createExperienceService } from '../experience-service.mjs';
import {
applyGooseV149CanaryBlockEnv,
resolveGooseApiTargetsFromEnv,
} from './goose-v149-canary.mjs';
import { loadMemindEnvFiles } from './memind-runtime-profile.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 gooseCanary = resolveGooseApiTargetsFromEnv(process.env);
if (gooseCanary?.targets?.length) {
return [...gooseCanary.targets];
}
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 isEnabledFlag(value) {
return ['1', 'true', 'yes', 'on'].includes(
String(value ?? '').trim().toLowerCase(),
);
}
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 <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'));
}
if (process.env.MEMIND_ENV_FILE) {
loadEnvFile(process.env.MEMIND_ENV_FILE);
}
loadMemindEnvFiles(root, process.env);
applyGooseV149CanaryBlockEnv(process.env, root);
const workerApiTargets = parseTargets();
if (resolveGooseApiTargetsFromEnv(process.env)?.mode === 'all') {
console.log(
`[agent-run-worker] Goose v1.49 canary mode=all primary=${workerApiTargets[0] ?? ''}`,
);
}
// Bundled worker lives under scripts/; MCP path resolution must not use that
// directory as the portal runtime root (see resolveBundledMcpServerPath).
if (!String(process.env.MEMIND_PORTAL_H5_ROOT ?? '').trim()) {
process.env.MEMIND_PORTAL_H5_ROOT = root;
}
const args = parseArgs(process.argv.slice(2));
if (args.help) {
printHelp();
process.exit(0);
}
async function bootstrapWorker() {
const pool = createDbPool();
const workflowShadowEnabled = isEnabledFlag(
process.env.MEMIND_ORCHESTRATOR_SHADOW_OBSERVATION_ENABLED,
);
const workflowShadowObserver = workflowShadowEnabled
? createWorkflowShadowObserver({
configService: createOrchestratorAdminConfigService(pool),
serviceToken: process.env.MEMIND_ORCHESTRATOR_SERVICE_TOKEN,
})
: null;
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 apiTargets = workerApiTargets.length ? workerApiTargets : parseTargets();
const apiTarget = apiTargets[0] ?? process.env.TKMIND_API_TARGET ?? 'https://127.0.0.1:18006';
const llmProviderService = createLlmProviderService(pool, {
apiTarget,
apiTargets,
apiSecret: process.env.TKMIND_SERVER__SECRET_KEY ?? 'local-dev-secret',
});
const memoryV2ConfigService = createMemoryV2AdminConfigService(pool);
const getEffectiveEnv = async () => {
const state = await memoryV2ConfigService.getRuntimeState().catch(() => null);
return {
...process.env,
...(state?.overrides ?? {}),
};
};
const conversationMemoryService = createConversationMemoryService(pool, {
llmProviderService,
getEffectiveEnv,
});
const memoryV2 = await createManagedMemoryV2Runtime({
legacyMemoryService: conversationMemoryService,
configService: memoryV2ConfigService,
mysqlPool: pool,
});
const episodicMemoryService = createEpisodicMemoryService(pool, {
getEffectiveEnv,
logger: console,
});
await episodicMemoryService.ensureSchema().catch((err) => {
console.warn(
'[agent-run-worker] episodic memory schema setup degraded:',
err instanceof Error ? err.message : err,
);
});
const sessionSnapshotService = createSessionSnapshotService(pool, {
conversationMemoryService,
memoryV2,
episodicMemoryService,
});
const sessionAccess = createSessionAccess({ userAuth, enabled: false });
const directChatService = createDirectChatService({
userAuth,
sessionAccess,
llmProviderService,
sessionSnapshotService,
memoryV2,
conversationMemoryService,
episodicMemoryService,
});
const chatIntentRouter = createManagedChatIntentRouter({
llmProviderService,
memoryV2,
conversationMemoryService,
episodicMemoryService,
configService: memoryV2ConfigService,
});
const toolGateway = createToolGateway({ llmProviderService });
const tkmindProxy = createTkmindProxy({
apiTarget,
apiTargets,
apiSecret: process.env.TKMIND_SERVER__SECRET_KEY ?? 'local-dev-secret',
userAuth,
});
const experienceService = createExperienceService(pool);
const gateway = createAgentRunGateway({
pool,
userAuth,
sessionAccess,
tkmindProxy,
toolGateway,
directChatService,
llmProviderService,
sessionSnapshotService,
conversationMemoryService,
chatIntentRouter,
experienceService,
conversationMemoryService,
autoDispatch: false,
maxConcurrentRuns: Number(process.env.MEMIND_AGENT_RUN_QUEUE_CONCURRENCY ?? 1),
runTimeoutMs: Number(process.env.MEMIND_AGENT_RUN_TIMEOUT_MS ?? 15 * 60 * 1000),
observeWorkflowRun: workflowShadowObserver,
observeWorkflowValidation:
workflowShadowObserver?.observeValidation ?? null,
enforcePageDataWorkflowValidation:
workflowShadowEnabled
&& isEnabledFlag(
process.env.MEMIND_ORCHESTRATOR_PAGE_DATA_VALIDATION_GATE_ENABLED,
),
});
return { pool, gateway };
}
const { pool, gateway } = await bootstrapWorker();
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);
}