Files
memind/mindspace-agent-jobs.test.mjs
John 2e14873f2d Initial commit: Memind H5 portal with MindSpace, Plaza, and agent jobs.
Track application source and tests; exclude local env, user workspaces, and runtime data via .gitignore.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-06-15 15:04:43 -07:00

401 lines
13 KiB
JavaScript

import assert from 'node:assert/strict';
import test from 'node:test';
import { agentJobInternals, createAgentJobService } from './mindspace-agent-jobs.mjs';
function createMockPool(state) {
const runQuery = async (sql, params = []) => {
if (sql.includes('FROM h5_agent_jobs j') && sql.includes('idempotency_key')) {
const row = state.jobs.find(
(job) => job.user_id === params[0] && job.idempotency_key === params[1],
);
if (!row) return [[]];
const category = state.categories.find((item) => item.id === row.output_category_id);
return [[{ ...row, output_category_code: category?.category_code ?? null }]];
}
if (sql.includes('FROM h5_space_categories') && sql.includes("category_code = 'draft'")) {
const row = state.categories.find(
(item) => item.user_id === params[0] && item.category_code === 'draft',
);
return [row ? [row] : []];
}
if (sql.includes('FROM h5_space_categories') && sql.includes('WHERE id = ? AND user_id = ?')) {
const row = state.categories.find(
(item) => item.id === params[0] && item.user_id === params[1],
);
return [row ? [row] : []];
}
if (sql.includes('FROM h5_assets a') && sql.includes('a.id IN')) {
const userId = params[0];
const ids = new Set(params.slice(1));
return [[
...state.assets
.filter((asset) => asset.user_id === userId && ids.has(asset.id) && asset.status !== 'deleted')
.map((asset) => ({
...asset,
scan_status:
state.assetVersions.find((version) => version.id === asset.current_version_id)?.scan_status ??
'passed',
})),
]];
}
if (sql.includes('INSERT INTO h5_agent_jobs')) {
state.jobs.push({
id: params[0],
user_id: params[1],
job_type: params[2],
instruction: params[3],
permission_scope: params[4],
user_context_json: params[5],
output_category_id: params[6],
output_type: params[7],
status: 'queued',
idempotency_key: params[8],
progress_json: params[9],
queued_at: params[10],
expires_at: params[11],
updated_at: params[12],
max_output_bytes: params[13],
result_page_id: null,
result_asset_id: null,
error_code: null,
error_message: null,
started_at: null,
heartbeat_at: null,
completed_at: null,
job_token_hash: null,
});
return [[]];
}
if (sql.includes('INSERT INTO h5_agent_job_assets')) {
state.jobAssets.push({
id: params[0],
job_id: params[1],
asset_id: params[2],
asset_version_id: params[3],
permission: 'read',
created_at: params[4],
});
return [[]];
}
if (sql.includes('FROM h5_agent_jobs j') && sql.includes('WHERE j.id = ?')) {
const row = state.jobs.find(
(job) => job.id === params[0] && (params[1] == null || job.user_id === params[1]),
);
if (!row) return [[]];
const category = state.categories.find((item) => item.id === row.output_category_id);
return [[{ ...row, output_category_code: category?.category_code ?? null }]];
}
if (sql.includes('FROM h5_agent_job_assets ja') && sql.includes('ORDER BY ja.created_at ASC')) {
const jobId = params[0];
return [[
...state.jobAssets
.filter((binding) => binding.job_id === jobId)
.map((binding) => {
const asset = state.assets.find((item) => item.id === binding.asset_id);
const version = state.assetVersions.find((item) => item.id === binding.asset_version_id);
return {
asset_id: binding.asset_id,
asset_version_id: binding.asset_version_id,
permission: binding.permission,
display_name: asset?.display_name ?? '',
mime_type: asset?.mime_type ?? '',
status: asset?.status ?? 'ready',
scan_status: version?.scan_status ?? 'passed',
};
}),
]];
}
if (sql.includes('WHERE id = ? AND user_id = ?') && sql.includes('FROM h5_agent_jobs')) {
const row = state.jobs.find((job) => job.id === params[0] && job.user_id === params[1]);
return [row ? [row] : []];
}
if (sql.includes("SET status = 'running'")) {
const row = state.jobs.find((job) => job.id === params[6]);
row.status = 'running';
row.started_at ??= params[0];
row.heartbeat_at = params[1];
row.updated_at = params[2];
row.expires_at = params[3];
row.progress_json = params[4];
row.job_token_hash = params[5];
return [[]];
}
if (sql.includes('SET heartbeat_at = ?')) {
const row = state.jobs.find((job) => job.id === params[4]);
row.heartbeat_at = params[0];
row.updated_at = params[1];
row.expires_at = params[2];
row.progress_json = params[3];
return [[]];
}
if (sql.includes("SET status = 'completed'")) {
const row = state.jobs.find((job) => job.id === params[4]);
row.status = 'completed';
row.result_page_id = params[0];
row.completed_at = params[1];
row.updated_at = params[2];
row.heartbeat_at = null;
row.job_token_hash = null;
row.progress_json = params[3];
return [[]];
}
if (sql.includes("SET status = 'failed'")) {
const row = state.jobs.find((job) => job.id === params[5]);
row.status = 'failed';
row.error_code = params[0];
row.error_message = params[1];
row.completed_at = params[2];
row.updated_at = params[3];
row.heartbeat_at = null;
row.job_token_hash = null;
row.progress_json = params[4];
return [[]];
}
if (sql.includes("SET status = 'queued'")) {
const row = state.jobs.find((job) => job.id === params[3]);
row.status = 'queued';
row.progress_json = params[0];
row.started_at = null;
row.completed_at = null;
row.error_code = null;
row.error_message = null;
row.result_page_id = null;
row.result_asset_id = null;
row.job_token_hash = null;
row.heartbeat_at = null;
row.expires_at = params[1];
row.updated_at = params[2];
return [[]];
}
if (sql.includes("SET status = 'cancelled'")) {
const row = state.jobs.find((job) => job.id === params[4]);
row.status = 'cancelled';
row.completed_at = params[0];
row.expires_at = params[1];
row.updated_at = params[2];
row.job_token_hash = null;
row.heartbeat_at = null;
row.progress_json = params[3];
return [[]];
}
if (sql.includes('FROM h5_agent_job_assets ja') && sql.includes('ja.asset_id = ?')) {
const binding = state.jobAssets.find(
(item) => item.job_id === params[0] && item.asset_id === params[1],
);
if (!binding) return [[]];
const asset = state.assets.find((item) => item.id === binding.asset_id);
const version = state.assetVersions.find((item) => item.id === binding.asset_version_id);
return [[{
id: asset.id,
display_name: asset.display_name,
mime_type: asset.mime_type,
status: asset.status,
asset_version_id: binding.asset_version_id,
storage_key: version.storage_key,
scan_status: version.scan_status,
}]];
}
throw new Error(`Unhandled SQL: ${sql}`);
};
return {
query: runQuery,
async getConnection() {
return {
query: runQuery,
async beginTransaction() {},
async commit() {},
async rollback() {},
release() {},
};
},
};
}
test('normalizeJobInput validates required fields', () => {
const result = agentJobInternals.normalizeJobInput({
jobType: 'generate_page',
instruction: ' 生成周报 ',
outputType: 'page_draft',
allowedAssetIds: ['a1', 'a1', 'a2'],
});
assert.equal(result.instruction, '生成周报');
assert.deepEqual(result.allowedAssetIds, ['a1', 'a2']);
assert.throws(
() =>
agentJobInternals.normalizeJobInput({
jobType: 'unknown',
instruction: 'x',
outputType: 'page_draft',
allowedAssetIds: ['a1'],
}),
{ code: 'invalid_agent_job_input' },
);
});
test('agent job lifecycle creates, claims, heartbeats, and completes into a page draft', async () => {
const state = {
categories: [{ id: 'cat-draft', user_id: 'user-1', category_code: 'draft' }],
assets: [{
id: 'asset-1',
user_id: 'user-1',
current_version_id: 'ver-1',
display_name: '日报.xlsx',
mime_type: 'application/vnd.openxmlformats-officedocument.spreadsheetml.sheet',
status: 'ready',
}],
assetVersions: [{
id: 'ver-1',
asset_id: 'asset-1',
storage_key: 'users/user-1/assets/asset-1/versions/ver-1',
scan_status: 'passed',
}],
jobs: [],
jobAssets: [],
};
const createdPages = [];
let seq = 0;
const service = createAgentJobService(createMockPool(state), {
idFactory: () => `id-${++seq}`,
nowFactory: () => 1_700_000_000_000 + seq,
pageService: {
async createFromAgent(userId, input, source) {
createdPages.push({ userId, input, source });
return { id: 'page-1' };
},
},
});
const job = await service.createJob('user-1', {
jobType: 'generate_page',
instruction: '根据日报生成周报页面',
outputType: 'page_draft',
allowedAssetIds: ['asset-1'],
});
assert.equal(job.status, 'queued');
assert.equal(job.assets[0].assetId, 'asset-1');
const claim = await service.claimJob(job.id);
assert.equal(claim.allowedAssets[0].assetId, 'asset-1');
assert.equal(state.jobs[0].status, 'running');
const heartbeat = await service.heartbeat(job.id, claim.jobToken, {
stage: 'ai_analysis',
message: '正在分析资产内容',
});
assert.equal(heartbeat.progress.stage, 'ai_analysis');
const completed = await service.completeJob(job.id, claim.jobToken, {
title: '项目周报',
summary: '本周进展',
content: '# 项目周报\n\n一切正常。',
sourceAssetIds: ['asset-1'],
});
assert.equal(completed.status, 'completed');
assert.equal(completed.resultPageId, 'page-1');
assert.equal(createdPages[0].userId, 'user-1');
assert.equal(createdPages[0].source.jobId, job.id);
assert.deepEqual(createdPages[0].source.assetIds, ['asset-1']);
});
test('retry resets a retryable failed job back to queued', async () => {
const state = {
categories: [{ id: 'cat-draft', user_id: 'user-1', category_code: 'draft' }],
assets: [],
assetVersions: [],
jobs: [{
id: 'job-1',
user_id: 'user-1',
job_type: 'generate_page',
instruction: 'retry me',
permission_scope: '{}',
user_context_json: '{}',
output_category_id: 'cat-draft',
output_type: 'page_draft',
status: 'failed',
idempotency_key: null,
progress_json: '{"stage":"failed"}',
queued_at: 1,
expires_at: 2,
updated_at: 3,
max_output_bytes: 1024,
result_page_id: null,
result_asset_id: null,
error_code: 'worker_crashed',
error_message: 'crashed',
started_at: 4,
heartbeat_at: null,
completed_at: 5,
job_token_hash: null,
}],
jobAssets: [],
};
const service = createAgentJobService(createMockPool(state), {
nowFactory: () => 42,
pageService: { async createFromAgent() { return { id: 'page-x' }; } },
});
const retried = await service.retryJob('user-1', 'job-1');
assert.equal(retried.status, 'queued');
assert.equal(retried.errorCode, null);
});
test('listJobs returns paginated items with total count', async () => {
const state = {
categories: [{ id: 'cat-draft', user_id: 'user-1', category_code: 'draft' }],
jobs: Array.from({ length: 3 }, (_, index) => ({
id: `job-${index}`,
user_id: 'user-1',
job_type: 'generate_page',
instruction: `task-${index}`,
permission_scope: '{}',
user_context_json: '{}',
output_category_id: 'cat-draft',
output_type: 'page_draft',
output_category_code: 'draft',
status: 'completed',
idempotency_key: `key-${index}`,
progress_json: '{"stage":"completed"}',
queued_at: 100 - index,
started_at: 100 - index,
completed_at: 100 - index,
expires_at: null,
updated_at: 100 - index,
max_output_bytes: 1024,
result_page_id: null,
result_asset_id: null,
error_code: null,
error_message: null,
heartbeat_at: null,
job_token_hash: null,
})),
jobAssets: [],
assets: [],
assetVersions: [],
};
const pool = {
async query(sql, params = []) {
if (sql.includes('COUNT(*) AS total FROM h5_agent_jobs')) {
return [[{ total: state.jobs.length }]];
}
if (sql.includes('FROM h5_agent_jobs j') && sql.includes('ORDER BY')) {
const offset = Number(sql.match(/OFFSET (\d+)/)?.[1] ?? 0);
const limit = Number(sql.match(/LIMIT (\d+)/)?.[1] ?? 20);
const rows = state.jobs.slice(offset, offset + limit);
return [rows];
}
if (sql.includes('FROM h5_agent_job_assets ja')) {
return [[]];
}
return [[]];
},
};
const service = createAgentJobService(pool, { nowFactory: () => 100 });
const page = await service.listJobs('user-1', { limit: 2, offset: 1 });
assert.equal(page.total, 3);
assert.equal(page.items.length, 2);
assert.equal(page.offset, 1);
assert.equal(page.hasMore, false);
});