import crypto from 'node:crypto'; import fs from 'node:fs'; import path from 'node:path'; import { createAgentRunGateway } from '../agent-run-gateway.mjs'; import { isDirectChatSessionId } from '../direct-chat-service.mjs'; import { cancelSessionActiveRequest, quiesceSessionStdioExtensions, } from '../session-runtime-lifecycle.mjs'; import { createOrchestratorAdminConfigService } from '../services/orchestrator/admin-config.mjs'; import { createWorkflowShadowObserver } from '../services/orchestrator/shadow-observer.mjs'; import { createTkmindProxy } from '../tkmind-proxy.mjs'; import { createToolGateway } from '../tool-gateway.mjs'; import { isPassiveCanaryRuntime } from './portal-runtime-role.mjs'; function isEnabledFlag(value, fallback = '') { return ['1', 'true', 'yes', 'on'].includes( String(value ?? fallback).trim().toLowerCase(), ); } export function startPortalAgentRunRecoveryLoop( agentRunGateway, { env = process.env, logger = console, setIntervalFn = setInterval, } = {}, ) { if ( typeof agentRunGateway?.dispatchQueuedRuns !== 'function' || typeof agentRunGateway?.recoverStaleRunningRuns !== 'function' ) { return null; } const configuredInterval = Number( env.MEMIND_AGENT_RUN_RECOVERY_INTERVAL_MS ?? 30_000, ); const intervalMs = Number.isFinite(configuredInterval) ? Math.max(5_000, configuredInterval) : 30_000; const configuredHeartbeatStaleMs = Number( env.MEMIND_AGENT_RUN_HEARTBEAT_STALE_MS ?? 90_000, ); const heartbeatStaleMs = Number.isFinite(configuredHeartbeatStaleMs) ? Math.max(intervalMs * 2, configuredHeartbeatStaleMs) : Math.max(intervalMs * 2, 90_000); let sweepActive = false; const sweep = async () => { if (sweepActive) return; sweepActive = true; try { const heartbeatRecovery = await agentRunGateway.recoverStaleRunningRuns({ staleMs: heartbeatStaleMs, limit: 20, dryRun: false, reason: 'worker_heartbeat_stale_recovery', }); const result = await agentRunGateway.dispatchQueuedRuns(); const recovered = Number(heartbeatRecovery?.recovered ?? 0) + Number(result?.staleRecovery?.recovered ?? 0); if (recovered > 0) { logger.warn( `[AgentRun] recovered ${recovered} stale run(s)`, ); } } catch (error) { logger.warn( '[AgentRun] background recovery sweep failed:', error instanceof Error ? error.message : error, ); } finally { sweepActive = false; } }; void sweep(); const timer = setIntervalFn(sweep, intervalMs); timer?.unref?.(); return timer; } export function createPortalRunDeliverablesValidator({ getWorkspacePublicationDelivery = () => null, } = {}) { if ( typeof getWorkspacePublicationDelivery !== 'function' ) { throw new Error( 'createPortalRunDeliverablesValidator requires MindSpace delivery access', ); } return async function validateRunDeliverables( input, ) { const service = getWorkspacePublicationDelivery(); if ( !service || typeof service.validateRunDeliverables !== 'function' ) { throw new Error( 'MindSpace 交付验证服务未启用', ); } return service.validateRunDeliverables( input, ); }; } export function bootstrapPortalGatewayServices({ pool, h5Root, env = process.env, apiTarget, apiTargets, apiSecret, userAuth, sessionAccess, sessionStreamStore, llmProviderService, subscriptionService, sessionSnapshotService, conversationMemoryService, memoryV2, systemDisclosurePolicyService, mindSpaceAssets, getWorkspacePublicationDelivery = () => null, directChatService, chatIntentRouter, syncUserGeneratedPages, isSessionPageDeliveryActive, createTkmindProxyFn = createTkmindProxy, createToolGatewayFn = createToolGateway, createAgentRunGatewayFn = createAgentRunGateway, createOrchestratorAdminConfigServiceFn = createOrchestratorAdminConfigService, createWorkflowShadowObserverFn = createWorkflowShadowObserver, createRunDeliverablesValidatorFn = createPortalRunDeliverablesValidator, startAgentRunRecoveryLoopFn = startPortalAgentRunRecoveryLoop, } = {}) { if ( !pool || !h5Root || !userAuth || !sessionAccess || !llmProviderService || typeof syncUserGeneratedPages !== 'function' || typeof isSessionPageDeliveryActive !== 'function' ) { throw new Error( 'bootstrapPortalGatewayServices requires gateway dependencies', ); } // GOOSED PROXY BOUNDARY: H5 chat → goosed unique entry. const tkmindProxy = createTkmindProxyFn({ apiTarget, apiTargets, apiSecret, userAuth, sessionAccess, sessionStreamStore, llmProviderService, subscriptionService, sessionSnapshotService, conversationMemoryService, memoryV2, systemDisclosurePolicyService, localFetchAsset: mindSpaceAssets ? async (userId, assetId) => { const { asset, bodyBase64 } = await mindSpaceAssets .readAssetContent( userId, assetId, ); return { buffer: Buffer.from( bodyBase64, 'base64', ), mimeType: asset.mimeType, }; } : null, }); const toolGateway = createToolGatewayFn({ llmProviderService, }); const validateRunDeliverables = createRunDeliverablesValidatorFn({ getWorkspacePublicationDelivery, }); const workflowShadowEnabled = isEnabledFlag( env.MEMIND_ORCHESTRATOR_SHADOW_OBSERVATION_ENABLED, ); const workflowShadowObserver = workflowShadowEnabled ? createWorkflowShadowObserverFn({ configService: createOrchestratorAdminConfigServiceFn(pool), serviceToken: env.MEMIND_ORCHESTRATOR_SERVICE_TOKEN, logger: console, }) : null; let runtimeRoot; try { runtimeRoot = fs.realpathSync(h5Root); } catch { runtimeRoot = path.resolve(h5Root); } const workerIdentity = { workerId: String( env.MEMIND_AGENT_RUN_WORKER_ID ?? '', ).trim() || `${process.pid}:${crypto.randomUUID()}`, runtimeRoot, buildId: String( env.MEMIND_RUNTIME_BUILD_ID ?? env.MEMIND_RELEASE_ID ?? env.GIT_COMMIT ?? 'dev', ).trim() || 'dev', }; const agentRunGateway = createAgentRunGatewayFn({ pool, userAuth, sessionAccess, tkmindProxy, toolGateway, directChatService, systemDisclosurePolicyService, chatIntentRouter, sessionSnapshotService, conversationMemoryService, observeWorkflowRun: workflowShadowObserver, observeWorkflowValidation: workflowShadowObserver?.observeValidation ?? null, enforcePageDataWorkflowValidation: workflowShadowEnabled && isEnabledFlag( env.MEMIND_ORCHESTRATOR_PAGE_DATA_VALIDATION_GATE_ENABLED, ), observePersonalMemoryOnSuccess: async ({ userId, sessionId, userMessage, }) => { if (!memoryV2?.observePersonalMemory) return; await memoryV2.observePersonalMemory({ userId, sessionId, messages: [userMessage], }); }, syncUserPagesOnSuccess: async ({ userId, sessionId, runStartedAtMs, }) => syncUserGeneratedPages(userId, { sessionId, sinceMs: runStartedAtMs, }), isSessionExternallyBusy: ({ sessionId }) => isSessionPageDeliveryActive(sessionId), validateRunDeliverables, cancelSessionOnRetry: async ({ sessionId, requestId, }) => { if (isDirectChatSessionId(sessionId)) { return { cancelled: false, skipped: true }; } const target = await tkmindProxy.resolveTarget(sessionId); const sessionApiFetch = (pathname, init) => tkmindProxy.apiFetchTo( target, pathname, init, ); return cancelSessionActiveRequest( sessionApiFetch, sessionId, requestId, ); }, quiesceSessionOnTerminal: async ({ sessionId, requestId, status, }) => { if (isDirectChatSessionId(sessionId)) { return { removed: [], skipped: true }; } const target = await tkmindProxy.resolveTarget(sessionId); const sessionApiFetch = (pathname, init) => tkmindProxy.apiFetchTo( target, pathname, init, ); const cancellation = status === 'failed' ? await cancelSessionActiveRequest( sessionApiFetch, sessionId, requestId, ) : { cancelled: false, skipped: true }; const quiesced = await quiesceSessionStdioExtensions( sessionApiFetch, sessionId, ); return { ...quiesced, cancellation }; }, autoDispatch: isEnabledFlag( env.MEMIND_AGENT_RUN_AUTODISPATCH, '1', ), maxConcurrentRuns: Number( env.MEMIND_AGENT_RUN_QUEUE_CONCURRENCY ?? 1, ), runTimeoutMs: Number( env.MEMIND_AGENT_RUN_TIMEOUT_MS ?? 15 * 60 * 1000, ), maxConcurrentShadowObservations: Number( env.MEMIND_ORCHESTRATOR_SHADOW_MAX_CONCURRENCY ?? 2, ), maxQueuedShadowObservations: Number( env.MEMIND_ORCHESTRATOR_SHADOW_MAX_QUEUE ?? 100, ), workerIdentity, }); const agentRunRecoveryTimer = isPassiveCanaryRuntime(env) ? null : startAgentRunRecoveryLoopFn( agentRunGateway, { env, logger: console }, ); return { tkmindProxy, toolGateway, agentRunGateway, agentRunRecoveryTimer, validateRunDeliverables, }; }