diff --git a/experience-reflect-worker.mjs b/experience-reflect-worker.mjs new file mode 100644 index 0000000..9c58924 --- /dev/null +++ b/experience-reflect-worker.mjs @@ -0,0 +1,130 @@ +import { DEFAULT_REFLECT_MIN_GROUP_SIZE } from './experience-reflect.mjs'; + +function readFlag(value, fallback = false) { + if (value == null || value === '') return fallback; + const normalized = String(value).trim().toLowerCase(); + if (['1', 'true', 'yes', 'on'].includes(normalized)) return true; + if (['0', 'false', 'no', 'off'].includes(normalized)) return false; + return fallback; +} + +function bounded(value, fallback, min, max) { + const parsed = Number(value); + if (!Number.isFinite(parsed)) return fallback; + return Math.min(max, Math.max(min, parsed)); +} + +export function resolveExperienceReflectPolicy(env = process.env) { + const intervalHours = Math.round( + bounded(env.EXPERIENCE_REFLECT_INTERVAL_HOURS, 6, 1, 168), + ); + return { + enabled: readFlag(env.EXPERIENCE_REFLECT_ENABLED, false), + intervalMs: intervalHours * 3_600_000, + minGroupSize: Math.round( + bounded(env.EXPERIENCE_REFLECT_MIN_GROUP_SIZE, DEFAULT_REFLECT_MIN_GROUP_SIZE, 2, 20), + ), + scopes: String(env.EXPERIENCE_REFLECT_SCOPES ?? 'global') + .split(/[\s,]+/u) + .map((scope) => scope.trim()) + .filter(Boolean) + .slice(0, 20), + runOnStart: readFlag(env.EXPERIENCE_REFLECT_RUN_ON_START, true), + }; +} + +export async function runExperienceReflectOnce({ + experienceService, + policy = resolveExperienceReflectPolicy(), + logger = console, +} = {}) { + if (!experienceService || typeof experienceService.reflect !== 'function') { + return { ok: false, skipped: true, reason: 'experience_service_unavailable' }; + } + if (!policy.enabled) { + return { ok: true, skipped: true, reason: 'disabled' }; + } + + const scopes = policy.scopes.length > 0 ? policy.scopes : ['global']; + const results = []; + let created = 0; + let archived = 0; + + for (const scope of scopes) { + const result = await experienceService.reflect({ + scope, + minGroupSize: policy.minGroupSize, + }); + results.push({ scope, ...result }); + created += Number(result?.created ?? 0); + archived += Number(result?.archived ?? 0); + } + + if (created > 0 || archived > 0) { + logger.log?.( + `[ExperienceReflect] created=${created} archived=${archived} scopes=${scopes.join(',')}`, + ); + } + + return { + ok: true, + skipped: false, + created, + archived, + scopes, + results, + }; +} + +export function startExperienceReflectWorker({ + experienceService, + env = process.env, + logger = console, + setIntervalFn = setInterval, +} = {}) { + const policy = resolveExperienceReflectPolicy(env); + if (!policy.enabled || !experienceService?.reflect) { + return { stop() {}, policy }; + } + + let stopped = false; + let running = false; + + const runOnce = async () => { + if (stopped || running) return; + running = true; + try { + await runExperienceReflectOnce({ + experienceService, + policy, + logger, + }); + } catch (error) { + logger.warn?.( + '[ExperienceReflect] worker run failed:', + error instanceof Error ? error.message : error, + ); + } finally { + running = false; + } + }; + + if (policy.runOnStart) { + void runOnce(); + } + + const timer = setIntervalFn(runOnce, policy.intervalMs); + timer?.unref?.(); + logger.log?.( + `[ExperienceReflect] worker enabled (interval=${policy.intervalMs}ms, minGroup=${policy.minGroupSize}, scopes=${policy.scopes.join(',') || 'global'})`, + ); + + return { + policy, + stop() { + stopped = true; + clearInterval(timer); + }, + runOnce, + }; +} diff --git a/experience-reflect-worker.test.mjs b/experience-reflect-worker.test.mjs new file mode 100644 index 0000000..b55d7c1 --- /dev/null +++ b/experience-reflect-worker.test.mjs @@ -0,0 +1,93 @@ +import assert from 'node:assert/strict'; +import test from 'node:test'; +import { + resolveExperienceReflectPolicy, + runExperienceReflectOnce, + startExperienceReflectWorker, +} from './experience-reflect-worker.mjs'; + +test('resolveExperienceReflectPolicy defaults to disabled with 6h interval', () => { + const policy = resolveExperienceReflectPolicy({}); + assert.equal(policy.enabled, false); + assert.equal(policy.intervalMs, 6 * 3_600_000); + assert.equal(policy.minGroupSize, 3); + assert.deepEqual(policy.scopes, ['global']); +}); + +test('resolveExperienceReflectPolicy parses env overrides', () => { + const policy = resolveExperienceReflectPolicy({ + EXPERIENCE_REFLECT_ENABLED: '1', + EXPERIENCE_REFLECT_INTERVAL_HOURS: '12', + EXPERIENCE_REFLECT_MIN_GROUP_SIZE: '4', + EXPERIENCE_REFLECT_SCOPES: 'global, team-a', + EXPERIENCE_REFLECT_RUN_ON_START: '0', + }); + assert.equal(policy.enabled, true); + assert.equal(policy.intervalMs, 12 * 3_600_000); + assert.equal(policy.minGroupSize, 4); + assert.deepEqual(policy.scopes, ['global', 'team-a']); + assert.equal(policy.runOnStart, false); +}); + +test('runExperienceReflectOnce skips when disabled', async () => { + const result = await runExperienceReflectOnce({ + experienceService: { reflect: async () => ({ created: 1 }) }, + policy: resolveExperienceReflectPolicy({}), + }); + assert.equal(result.skipped, true); + assert.equal(result.reason, 'disabled'); +}); + +test('runExperienceReflectOnce reflects configured scopes', async () => { + const calls = []; + const result = await runExperienceReflectOnce({ + experienceService: { + reflect: async (input) => { + calls.push(input); + return { ok: true, created: 1, archived: 3 }; + }, + }, + policy: { + enabled: true, + minGroupSize: 3, + scopes: ['global', 'team-a'], + intervalMs: 1000, + runOnStart: true, + }, + }); + assert.equal(result.created, 2); + assert.equal(result.archived, 6); + assert.deepEqual(calls, [ + { scope: 'global', minGroupSize: 3 }, + { scope: 'team-a', minGroupSize: 3 }, + ]); +}); + +test('startExperienceReflectWorker schedules periodic runs', async () => { + let ticks = 0; + const timers = []; + const worker = startExperienceReflectWorker({ + experienceService: { + reflect: async () => { + ticks += 1; + return { ok: true, created: 0, archived: 0 }; + }, + }, + env: { + EXPERIENCE_REFLECT_ENABLED: '1', + EXPERIENCE_REFLECT_RUN_ON_START: '0', + }, + setIntervalFn: (fn, ms) => { + assert.equal(ms, 6 * 3_600_000); + timers.push(fn); + return { unref() {} }; + }, + }); + assert.equal(worker.policy.enabled, true); + assert.equal(ticks, 0); + await worker.runOnce(); + assert.equal(ticks, 1); + await timers[0](); + assert.equal(ticks, 2); + worker.stop(); +}); diff --git a/mindspace-service/mindspace-service-bootstrap.mjs b/mindspace-service/mindspace-service-bootstrap.mjs index 307b8a7..e8f57cc 100644 --- a/mindspace-service/mindspace-service-bootstrap.mjs +++ b/mindspace-service/mindspace-service-bootstrap.mjs @@ -59,6 +59,7 @@ async function loadMemindModules(memindRoot) { userAuth, agentRunner, experienceService, + experienceReflectWorker, workspaceThumbnails, workspaceSync, ] = await Promise.all([ @@ -71,6 +72,7 @@ async function loadMemindModules(memindRoot) { importFromMemindRoot(memindRoot, 'user-auth.mjs'), importFromMemindRoot(memindRoot, 'mindspace-agent-runner.mjs'), importFromMemindRoot(memindRoot, 'experience-service.mjs'), + importFromMemindRoot(memindRoot, 'experience-reflect-worker.mjs'), importFromMemindRoot(memindRoot, 'mindspace-workspace-thumbnails.mjs'), importFromMemindRoot(memindRoot, 'mindspace-workspace-sync.mjs'), ]); @@ -84,6 +86,7 @@ async function loadMemindModules(memindRoot) { ...userAuth, ...agentRunner, ...experienceService, + ...experienceReflectWorker, ...workspaceThumbnails, ...workspaceSync, }; @@ -134,6 +137,7 @@ export async function bootstrapMindSpaceService({ createSessionSnapshotService, createUserAuth, createExperienceService, + startExperienceReflectWorker, resolveMindSpaceServerRuntimeOptions, startWorkspaceAssetSyncWatcher, startWorkspaceThumbnailWatcher, @@ -178,6 +182,13 @@ export async function bootstrapMindSpaceService({ if (runtime.experienceEnabled) { experienceService = createExperienceService(pool, { productEventsPool: pool }); } + const experienceReflectWorker = experienceService + ? startExperienceReflectWorker({ + experienceService, + env, + logger, + }) + : { stop() {}, policy: { enabled: false } }; const agentRunner = createMindSpaceAgentRunner({ apiTarget, apiSecret, @@ -218,9 +229,11 @@ export async function bootstrapMindSpaceService({ sessionSnapshotService, userAuth, agentRunner, + experienceReflectWorker, adapter, backgroundJobs, async close() { + experienceReflectWorker.stop?.(); await pool.end(); }, }; diff --git a/server/portal-agent-services-bootstrap.mjs b/server/portal-agent-services-bootstrap.mjs index d39d402..4f6ddf6 100644 --- a/server/portal-agent-services-bootstrap.mjs +++ b/server/portal-agent-services-bootstrap.mjs @@ -1,4 +1,5 @@ import { createAssetGatewayConfigService } from '../asset-gateway.mjs'; +import { startExperienceReflectWorker } from '../experience-reflect-worker.mjs'; import { createExperienceService } from '../experience-service.mjs'; import { createImageMakeAdminConfigService } from '../image-make-admin-config.mjs'; import { createImageMakeClientFromEnv } from '../image-make-client.mjs'; @@ -53,6 +54,8 @@ export async function bootstrapPortalAgentServices({ createMindSpaceImageGenerationServiceFn = createMindSpaceImageGenerationService, createWordFilterServiceFn = createWordFilterService, + startExperienceReflectWorkerFn = startExperienceReflectWorker, + setIntervalFn = setInterval, relayBootstrap = RELAY_BOOTSTRAP, } = {}) { if ( @@ -99,6 +102,15 @@ export async function bootstrapPortalAgentServices({ } } + const experienceReflectWorker = experienceService + ? startExperienceReflectWorkerFn({ + experienceService, + env, + logger, + setIntervalFn, + }) + : { stop() {}, policy: { enabled: false } }; + const mindSpaceAgentRunner = createMindSpaceAgentRunnerFn({ apiTarget, @@ -226,6 +238,7 @@ export async function bootstrapPortalAgentServices({ return { mindSpaceAgentJobs, experienceService, + experienceReflectWorker, mindSpaceAgentRunner, mindSpaceAudit, userDataSpaceBackfill,