Add Experience reflect background worker with env-gated scheduling.
Start a periodic reflect loop from portal and MindSpace service bootstraps when EXPERIENCE_REFLECT_ENABLED=1, aggregating repeated task outcomes on a configurable interval. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -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,
|
||||||
|
};
|
||||||
|
}
|
||||||
@@ -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();
|
||||||
|
});
|
||||||
@@ -59,6 +59,7 @@ async function loadMemindModules(memindRoot) {
|
|||||||
userAuth,
|
userAuth,
|
||||||
agentRunner,
|
agentRunner,
|
||||||
experienceService,
|
experienceService,
|
||||||
|
experienceReflectWorker,
|
||||||
workspaceThumbnails,
|
workspaceThumbnails,
|
||||||
workspaceSync,
|
workspaceSync,
|
||||||
] = await Promise.all([
|
] = await Promise.all([
|
||||||
@@ -71,6 +72,7 @@ async function loadMemindModules(memindRoot) {
|
|||||||
importFromMemindRoot(memindRoot, 'user-auth.mjs'),
|
importFromMemindRoot(memindRoot, 'user-auth.mjs'),
|
||||||
importFromMemindRoot(memindRoot, 'mindspace-agent-runner.mjs'),
|
importFromMemindRoot(memindRoot, 'mindspace-agent-runner.mjs'),
|
||||||
importFromMemindRoot(memindRoot, 'experience-service.mjs'),
|
importFromMemindRoot(memindRoot, 'experience-service.mjs'),
|
||||||
|
importFromMemindRoot(memindRoot, 'experience-reflect-worker.mjs'),
|
||||||
importFromMemindRoot(memindRoot, 'mindspace-workspace-thumbnails.mjs'),
|
importFromMemindRoot(memindRoot, 'mindspace-workspace-thumbnails.mjs'),
|
||||||
importFromMemindRoot(memindRoot, 'mindspace-workspace-sync.mjs'),
|
importFromMemindRoot(memindRoot, 'mindspace-workspace-sync.mjs'),
|
||||||
]);
|
]);
|
||||||
@@ -84,6 +86,7 @@ async function loadMemindModules(memindRoot) {
|
|||||||
...userAuth,
|
...userAuth,
|
||||||
...agentRunner,
|
...agentRunner,
|
||||||
...experienceService,
|
...experienceService,
|
||||||
|
...experienceReflectWorker,
|
||||||
...workspaceThumbnails,
|
...workspaceThumbnails,
|
||||||
...workspaceSync,
|
...workspaceSync,
|
||||||
};
|
};
|
||||||
@@ -134,6 +137,7 @@ export async function bootstrapMindSpaceService({
|
|||||||
createSessionSnapshotService,
|
createSessionSnapshotService,
|
||||||
createUserAuth,
|
createUserAuth,
|
||||||
createExperienceService,
|
createExperienceService,
|
||||||
|
startExperienceReflectWorker,
|
||||||
resolveMindSpaceServerRuntimeOptions,
|
resolveMindSpaceServerRuntimeOptions,
|
||||||
startWorkspaceAssetSyncWatcher,
|
startWorkspaceAssetSyncWatcher,
|
||||||
startWorkspaceThumbnailWatcher,
|
startWorkspaceThumbnailWatcher,
|
||||||
@@ -178,6 +182,13 @@ export async function bootstrapMindSpaceService({
|
|||||||
if (runtime.experienceEnabled) {
|
if (runtime.experienceEnabled) {
|
||||||
experienceService = createExperienceService(pool, { productEventsPool: pool });
|
experienceService = createExperienceService(pool, { productEventsPool: pool });
|
||||||
}
|
}
|
||||||
|
const experienceReflectWorker = experienceService
|
||||||
|
? startExperienceReflectWorker({
|
||||||
|
experienceService,
|
||||||
|
env,
|
||||||
|
logger,
|
||||||
|
})
|
||||||
|
: { stop() {}, policy: { enabled: false } };
|
||||||
const agentRunner = createMindSpaceAgentRunner({
|
const agentRunner = createMindSpaceAgentRunner({
|
||||||
apiTarget,
|
apiTarget,
|
||||||
apiSecret,
|
apiSecret,
|
||||||
@@ -218,9 +229,11 @@ export async function bootstrapMindSpaceService({
|
|||||||
sessionSnapshotService,
|
sessionSnapshotService,
|
||||||
userAuth,
|
userAuth,
|
||||||
agentRunner,
|
agentRunner,
|
||||||
|
experienceReflectWorker,
|
||||||
adapter,
|
adapter,
|
||||||
backgroundJobs,
|
backgroundJobs,
|
||||||
async close() {
|
async close() {
|
||||||
|
experienceReflectWorker.stop?.();
|
||||||
await pool.end();
|
await pool.end();
|
||||||
},
|
},
|
||||||
};
|
};
|
||||||
|
|||||||
@@ -1,4 +1,5 @@
|
|||||||
import { createAssetGatewayConfigService } from '../asset-gateway.mjs';
|
import { createAssetGatewayConfigService } from '../asset-gateway.mjs';
|
||||||
|
import { startExperienceReflectWorker } from '../experience-reflect-worker.mjs';
|
||||||
import { createExperienceService } from '../experience-service.mjs';
|
import { createExperienceService } from '../experience-service.mjs';
|
||||||
import { createImageMakeAdminConfigService } from '../image-make-admin-config.mjs';
|
import { createImageMakeAdminConfigService } from '../image-make-admin-config.mjs';
|
||||||
import { createImageMakeClientFromEnv } from '../image-make-client.mjs';
|
import { createImageMakeClientFromEnv } from '../image-make-client.mjs';
|
||||||
@@ -53,6 +54,8 @@ export async function bootstrapPortalAgentServices({
|
|||||||
createMindSpaceImageGenerationServiceFn =
|
createMindSpaceImageGenerationServiceFn =
|
||||||
createMindSpaceImageGenerationService,
|
createMindSpaceImageGenerationService,
|
||||||
createWordFilterServiceFn = createWordFilterService,
|
createWordFilterServiceFn = createWordFilterService,
|
||||||
|
startExperienceReflectWorkerFn = startExperienceReflectWorker,
|
||||||
|
setIntervalFn = setInterval,
|
||||||
relayBootstrap = RELAY_BOOTSTRAP,
|
relayBootstrap = RELAY_BOOTSTRAP,
|
||||||
} = {}) {
|
} = {}) {
|
||||||
if (
|
if (
|
||||||
@@ -99,6 +102,15 @@ export async function bootstrapPortalAgentServices({
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const experienceReflectWorker = experienceService
|
||||||
|
? startExperienceReflectWorkerFn({
|
||||||
|
experienceService,
|
||||||
|
env,
|
||||||
|
logger,
|
||||||
|
setIntervalFn,
|
||||||
|
})
|
||||||
|
: { stop() {}, policy: { enabled: false } };
|
||||||
|
|
||||||
const mindSpaceAgentRunner =
|
const mindSpaceAgentRunner =
|
||||||
createMindSpaceAgentRunnerFn({
|
createMindSpaceAgentRunnerFn({
|
||||||
apiTarget,
|
apiTarget,
|
||||||
@@ -226,6 +238,7 @@ export async function bootstrapPortalAgentServices({
|
|||||||
return {
|
return {
|
||||||
mindSpaceAgentJobs,
|
mindSpaceAgentJobs,
|
||||||
experienceService,
|
experienceService,
|
||||||
|
experienceReflectWorker,
|
||||||
mindSpaceAgentRunner,
|
mindSpaceAgentRunner,
|
||||||
mindSpaceAudit,
|
mindSpaceAudit,
|
||||||
userDataSpaceBackfill,
|
userDataSpaceBackfill,
|
||||||
|
|||||||
Reference in New Issue
Block a user