feat: finalize mindspace service extraction phase a
This commit is contained in:
@@ -0,0 +1,157 @@
|
||||
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,
|
||||
resolveSessionSnapshot,
|
||||
registerPublicHtmlArtifactsForConversation,
|
||||
logger = console,
|
||||
setIntervalFn = setInterval,
|
||||
} = {}) {
|
||||
const services = createMindSpaceLocalRuntimeServices({
|
||||
pool,
|
||||
h5Root,
|
||||
env,
|
||||
maxFileBytes,
|
||||
publicPageLimit,
|
||||
resolveUserIdForAgentSession,
|
||||
conversationPackageBackfill: ({ user, sessionId, storageRoot, registry }) =>
|
||||
backfillConversationPackageArtifacts({
|
||||
pool,
|
||||
registry,
|
||||
storageRoot,
|
||||
h5Root,
|
||||
user,
|
||||
sessionId,
|
||||
}),
|
||||
conversationPackagePublicHtmlHydrator: createConversationPackagePublicHtmlHydrator({
|
||||
h5Root,
|
||||
resolveSessionSnapshot,
|
||||
registerPublicHtmlArtifactsForConversation,
|
||||
}),
|
||||
logger,
|
||||
});
|
||||
|
||||
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;
|
||||
},
|
||||
};
|
||||
}
|
||||
Reference in New Issue
Block a user