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 { scanWorkspaceFilesForProhibitedBrowserStorage } from '../mindspace-browser-storage-policy.mjs'; import { evaluatePageDataHtmlContent } from '../mindspace-page-data-finish-guard.mjs'; import { normalizeWorkspaceRelativePath } from '../mindspace-pages.mjs'; import { resolveMindSpaceUserPublishDir } from '../mindspace-runtime-config.mjs'; import { policyAllowsAction } from '../page-access-policy.mjs'; import { detectPageDataDatasetUsageFromHtml } from '../page-data-html-detect.mjs'; import { readPageAccessPolicy } from '../page-data-policy-store.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({ h5Root, resolveMindSpaceUserPublishDirFn = resolveMindSpaceUserPublishDir, normalizeWorkspaceRelativePathFn = normalizeWorkspaceRelativePath, evaluatePageDataHtmlContentFn = evaluatePageDataHtmlContent, readPageAccessPolicyFn = readPageAccessPolicy, detectPageDataDatasetUsageFromHtmlFn = detectPageDataDatasetUsageFromHtml, policyAllowsActionFn = policyAllowsAction, scanWorkspaceFilesForProhibitedBrowserStorageFn = scanWorkspaceFilesForProhibitedBrowserStorage, resolvePathFn = path.resolve, pathSeparator = path.sep, existsSyncFn = fs.existsSync, readFileSyncFn = fs.readFileSync, } = {}) { if (!h5Root) { throw new Error( 'createPortalRunDeliverablesValidator requires h5Root', ); } return async function validateRunDeliverables({ userId, deliverables, }) { const publishDir = resolveMindSpaceUserPublishDirFn(h5Root, { id: userId, }); const resolvedPublishDir = resolvePathFn(publishDir); const pageDataErrors = []; for (const page of deliverables?.pages ?? []) { const relativePath = normalizeWorkspaceRelativePathFn( page.workspaceRelativePath, ); if (!relativePath?.startsWith('public/')) continue; const filePath = resolvePathFn( publishDir, relativePath, ); if ( !filePath.startsWith( `${resolvedPublishDir}${pathSeparator}`, ) || !existsSyncFn(filePath) ) { continue; } const html = readFileSyncFn(filePath, 'utf8'); const evaluation = evaluatePageDataHtmlContentFn(html, { relativePath, }); if (!evaluation.usesPageDataApi) continue; for (const issue of evaluation.issues) { pageDataErrors.push({ code: issue, message: `${relativePath} Page Data HTML 不可交付:${issue}`, }); } const policy = page.pageId ? readPageAccessPolicyFn( publishDir, page.pageId, ) : null; for (const [dataset, actions] of detectPageDataDatasetUsageFromHtmlFn(html)) { for (const action of ['read', 'insert']) { if ( actions?.[action] && !policyAllowsActionFn( policy, dataset, action, ) ) { pageDataErrors.push({ code: 'page_data_policy_action_missing', message: `${relativePath} 的 ${dataset}.${action} 未获最终 policy 授权或 dataset 已关闭`, }); } } } } const violations = scanWorkspaceFilesForProhibitedBrowserStorageFn({ publishDir, relativePaths: (deliverables?.pages ?? []) .map((page) => page.workspaceRelativePath) .filter(Boolean), }); return { errors: [ ...pageDataErrors, ...violations.map((violation) => ({ code: 'browser_storage_forbidden', message: `${violation.relativePath} 使用 ${violation.apis.join(', ')}`, })), ], }; }; } export function bootstrapPortalGatewayServices({ pool, h5Root, env = process.env, apiTarget, apiTargets, apiSecret, userAuth, sessionAccess, sessionStreamStore, llmProviderService, subscriptionService, sessionSnapshotService, conversationMemoryService, memoryV2, systemDisclosurePolicyService, mindSpaceAssets, directChatService, chatIntentRouter, syncUserGeneratedPages, isSessionPageDeliveryActive, createTkmindProxyFn = createTkmindProxy, createToolGatewayFn = createToolGateway, createAgentRunGatewayFn = createAgentRunGateway, createOrchestratorAdminConfigServiceFn = createOrchestratorAdminConfigService, createWorkflowShadowObserverFn = createWorkflowShadowObserver, createRunDeliverablesValidatorFn = createPortalRunDeliverablesValidator, startAgentRunRecoveryLoopFn = startPortalAgentRunRecoveryLoop, readAssetFileFn = fs.promises.readFile, } = {}) { 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, path: assetPath } = await mindSpaceAssets.readAsset( userId, assetId, ); const buffer = await readAssetFileFn(assetPath); return { buffer, mimeType: asset.mimeType, }; } : null, }); const toolGateway = createToolGatewayFn({ llmProviderService, }); const validateRunDeliverables = createRunDeliverablesValidatorFn({ h5Root }); 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, }; }