Files
memind/mindspace-local-server-adapter.mjs
john 7f8d692d16 fix(agent): recover stale runs and improve new-user OA delivery
Fix DEV logout cookie clearing, materialize selected MindSpace OA assets before agent runs, and recover zombie runs from synced workspace pages. Add client run wait timeout, harness retry limits, page-edit asset forwarding, and logout/john2 scenario tests.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-07-10 20:05:52 +08:00

166 lines
5.4 KiB
JavaScript

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;
},
};
}