import { backfillConversationPackageArtifacts } from './mindspace-conversation-package-backfill.mjs'; import { createConversationPackagePublicHtmlHydrator } from './mindspace-conversation-package-public-html.mjs'; import { createMindSpaceLocalRuntimeServices } from './mindspace-local-runtime-services.mjs'; export function createMindSpaceLocalServerAdapter({ pool, h5Root, env = process.env, maxFileBytes, publicPageLimit, resolveUserIdForAgentSession, resolveWorkspaceRoot = null, resolveSessionSnapshot, registerPublicHtmlArtifactsForConversation, logger = console, setIntervalFn = setInterval, syncWorkspaceAssetsEnabled = true, } = {}) { let services; services = createMindSpaceLocalRuntimeServices({ pool, h5Root, env, maxFileBytes, publicPageLimit, resolveUserIdForAgentSession, resolveWorkspaceRoot, conversationPackageBackfill: ({ user, sessionId, storageRoot, registry }) => backfillConversationPackageArtifacts({ pool, registry, storageRoot, h5Root, user, sessionId, }), conversationPackagePublicHtmlHydrator: createConversationPackagePublicHtmlHydrator({ h5Root, resolveSessionSnapshot, registerPublicHtmlArtifactsForConversation, }), logger, syncWorkspaceAssets: syncWorkspaceAssetsEnabled ? (targetUserId, options) => services?.assetService?.syncWorkspaceAssets?.(targetUserId, options) : null, }); return { ...services, kind: 'local', implementationStatus: 'ready', assertReady() { return true; }, startPublicationCleanup(intervalMs = 60 * 1000) { if (!services.publicationService || typeof setIntervalFn !== 'function') return null; const timer = setIntervalFn(async () => { try { const result = await services.publicationService.cleanupExpiredUnconfirmedPublications(); if (result.cleaned > 0) { logger.log?.( `[Publication Cleanup] Auto-privatized ${result.cleaned} expired unconfirmed publications`, ); } } catch (error) { logger.error?.('[Publication Cleanup Error]', error instanceof Error ? error.message : error); } }, intervalMs); timer?.unref?.(); return timer; }, startAgentWorker({ enabled = false, concurrency = 1, pollMs = 1000, staleMs = 5 * 60 * 1000, runner, } = {}) { if (!enabled || !services.agentJobService || !runner || typeof setIntervalFn !== 'function') return null; let inFlight = 0; let draining = false; const drainQueue = async () => { if (draining) return; draining = true; try { while (inFlight < concurrency) { const claim = await services.agentJobService.claimNextJob(); if (!claim) break; inFlight += 1; void runner .runJob(claim.jobId, claim) .catch((error) => { logger.error?.('Agent worker job failed:', error); }) .finally(() => { inFlight -= 1; }); } } catch (error) { logger.error?.('Agent worker drain failed:', error); } finally { draining = false; } }; const workerTimer = setIntervalFn(() => { void drainQueue(); }, pollMs); const reaperTimer = setIntervalFn(() => { void services.agentJobService .reapStaleJobs(staleMs) .then((reaped) => { if (reaped > 0) { logger.warn?.(`Agent worker reaped ${reaped} stale running job(s)`); } }) .catch((error) => { logger.error?.('Agent worker reaper failed:', error); }); }, Math.min(staleMs, 60_000)); workerTimer?.unref?.(); reaperTimer?.unref?.(); logger.log?.(`Agent job worker enabled (concurrency=${concurrency}, poll=${pollMs}ms)`); return { workerTimer, reaperTimer }; }, startBackgroundJobs({ publicationCleanupIntervalMs = 60 * 1000, agentWorker, agentRunner, workspaceMaintenanceEnabled = false, startWorkspaceThumbnailWatcher, startWorkspaceAssetSyncWatcher, publishRoot, syncUserWorkspaceByDirKey, expireStaleUploadsIntervalMs = 5 * 60 * 1000, } = {}) { const backgroundJobs = { publicationCleanup: this.startPublicationCleanup(publicationCleanupIntervalMs), agentWorker: this.startAgentWorker({ ...agentWorker, runner: agentRunner, }), workspaceMaintenance: null, }; if (!workspaceMaintenanceEnabled) { logger.log?.('Workspace maintenance daemons disabled (MEMIND_WORKSPACE_MAINTENANCE=0)'); return backgroundJobs; } startWorkspaceThumbnailWatcher?.(publishRoot); startWorkspaceAssetSyncWatcher?.({ publishRoot, syncUserWorkspaceByDirKey, }); void services.assetService?.expireStaleUploads?.().catch(() => {}); backgroundJobs.workspaceMaintenance = { expireStaleUploadsTimer: setIntervalFn(() => { void services.assetService?.expireStaleUploads?.().catch(() => {}); }, expireStaleUploadsIntervalMs), }; backgroundJobs.workspaceMaintenance.expireStaleUploadsTimer?.unref?.(); return backgroundJobs; }, }; }