Files
memind/agent-run-deliverable-check.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

182 lines
5.3 KiB
JavaScript

export const RECOVERABLE_FINISH_ERROR_CODES = new Set([
'AGENT_RUN_TIMEOUT',
'SESSION_REPLY_TIMEOUT',
'SESSION_REPLY_INCOMPLETE',
'AGENT_RUN_STALE_RECOVERY',
]);
export const SESSION_FINISHED_STALE_GRACE_MS = 90 * 1000;
export const WORKSPACE_DELIVERABLE_LOOKBACK_MS = 60 * 1000;
export function isRecoverableFinishError(error) {
const code = String(error?.code ?? '').trim();
return RECOVERABLE_FINISH_ERROR_CODES.has(code);
}
export function hasRecoverableSessionDeliverables(summary) {
const pageCount = Number(summary?.pageCount ?? 0);
const publicationCount = Number(summary?.publicationCount ?? 0);
return pageCount > 0 || publicationCount > 0;
}
function mapDeliverableRows(rows) {
const pages = (rows ?? []).map((row) => ({
pageId: String(row.page_id ?? '').trim(),
title: String(row.title ?? '').trim(),
publicationId: row.publication_id ? String(row.publication_id) : null,
publicationStatus: row.publication_status ? String(row.publication_status) : null,
publicUrl: row.public_url ? String(row.public_url) : null,
})).filter((row) => row.pageId);
const publicationCount = pages.filter((row) => row.publicationId).length;
return {
pageCount: pages.length,
publicationCount,
pages,
};
}
export function mergeDeliverableSummaries(...summaries) {
const pagesById = new Map();
for (const summary of summaries) {
for (const page of summary?.pages ?? []) {
if (!page?.pageId) continue;
pagesById.set(page.pageId, page);
}
}
const pages = [...pagesById.values()];
return {
pageCount: pages.length,
publicationCount: pages.filter((row) => row.publicationId).length,
pages,
};
}
export async function detectSessionDeliverables(pool, userId, sessionId) {
if (!pool?.query || !userId || !sessionId) {
return { pageCount: 0, publicationCount: 0, pages: [] };
}
const [rows] = await pool.query(
`SELECT p.id AS page_id,
p.title,
pr.id AS publication_id,
pr.status AS publication_status,
pr.public_url
FROM h5_page_records p
LEFT JOIN h5_publish_records pr
ON pr.page_id = p.id
AND pr.user_id = p.user_id
AND pr.status = 'online'
WHERE p.user_id = ?
AND p.source_session_id = ?
AND p.status <> 'deleted'
ORDER BY p.created_at ASC`,
[userId, sessionId],
);
return mapDeliverableRows(rows);
}
export async function detectWorkspaceSyncedDeliverables(pool, userId, { sinceMs = null } = {}) {
if (!pool?.query || !userId) {
return { pageCount: 0, publicationCount: 0, pages: [] };
}
const sinceClause = sinceMs != null ? 'AND p.updated_at >= ?' : '';
const params = sinceMs != null ? [userId, sinceMs] : [userId];
const [rows] = await pool.query(
`SELECT p.id AS page_id,
p.title,
pr.id AS publication_id,
pr.status AS publication_status,
pr.public_url
FROM h5_page_records p
JOIN h5_page_versions pv ON pv.id = p.current_version_id
LEFT JOIN h5_publish_records pr
ON pr.page_id = p.id
AND pr.user_id = p.user_id
AND pr.status = 'online'
WHERE p.user_id = ?
AND p.status <> 'deleted'
${sinceClause}
AND (
p.workspace_relative_path LIKE 'public/%.html'
OR JSON_UNQUOTE(JSON_EXTRACT(pv.source_snapshot_json, '$.relative_path')) LIKE 'public/%.html'
)
AND (
JSON_EXTRACT(pv.source_snapshot_json, '$.auto_synced') = true
OR JSON_UNQUOTE(JSON_EXTRACT(pv.source_snapshot_json, '$.content_mode')) = 'static_html'
)
ORDER BY p.updated_at DESC`,
params,
);
return mapDeliverableRows(rows);
}
export async function prepareAndDetectSessionDeliverables({
pool,
userId,
sessionId,
runStartedAtMs = null,
prepareDeliverables = null,
} = {}) {
if (typeof prepareDeliverables === 'function') {
await prepareDeliverables({ userId, sessionId }).catch(() => {});
}
const sinceMs = runStartedAtMs == null
? null
: Math.max(0, Number(runStartedAtMs) - WORKSPACE_DELIVERABLE_LOOKBACK_MS);
const bySession = await detectSessionDeliverables(pool, userId, sessionId);
const byWorkspace = await detectWorkspaceSyncedDeliverables(pool, userId, { sinceMs });
return mergeDeliverableSummaries(bySession, byWorkspace);
}
export async function tryRecoverRunFromDeliverables({
pool,
userId,
sessionId,
error = null,
requireRecoverableError = false,
prepareDeliverables = null,
runStartedAtMs = null,
} = {}) {
if (requireRecoverableError && !isRecoverableFinishError(error)) {
return null;
}
const deliverables = await prepareAndDetectSessionDeliverables({
pool,
userId,
sessionId,
runStartedAtMs,
prepareDeliverables,
});
if (!hasRecoverableSessionDeliverables(deliverables)) {
return null;
}
return {
deliverables,
originalError: error
? {
code: error?.code ?? null,
message: error instanceof Error ? error.message : String(error ?? ''),
}
: null,
};
}
export async function tryRecoverRunFromSessionDeliverables({
pool,
userId,
sessionId,
error,
prepareDeliverables = null,
runStartedAtMs = null,
} = {}) {
return tryRecoverRunFromDeliverables({
pool,
userId,
sessionId,
error,
requireRecoverableError: true,
prepareDeliverables,
runStartedAtMs,
});
}