diff --git a/scheduled-task-executor.mjs b/scheduled-task-executor.mjs index c860142..fecdf4f 100644 --- a/scheduled-task-executor.mjs +++ b/scheduled-task-executor.mjs @@ -1,7 +1,13 @@ import crypto from 'node:crypto'; +import fs from 'node:fs'; import path from 'node:path'; import { buildChatSkillPrompt, SCHEDULED_TASK_AUTOMATION_SKILL_NAME } from './chat-skills.mjs'; -import { releaseMaterializedPageDeliveryContracts } from './mindspace-delivery-contract.mjs'; +import { + markPageDeliveryContractReady, + normalizeDeliveryRelativePath, + preparePageDeliveryContract, + releaseMaterializedPageDeliveryContracts, +} from './mindspace-delivery-contract.mjs'; import { collectOwnPublicHtmlRelativePaths, materializeMissingPublicHtmlWrites, @@ -83,6 +89,56 @@ const SCHEDULED_TASK_CLARIFICATION_PATTERNS = [ /缺(?:少|失)/u, ]; +const DEFAULT_SCHEDULED_TASK_DELIVERY_RETRY_DELAYS_MS = [ + 250, + 1_000, + 3_000, + 5_000, + 10_000, + 30_000, + 60_000, +]; + +export function resolveScheduledTaskDeliveryRetryDelaysMs( + env = process.env, +) { + const raw = String( + env.H5_SCHEDULED_TASK_DELIVERY_RETRY_DELAYS_MS ?? '', + ).trim(); + if (!raw) return [...DEFAULT_SCHEDULED_TASK_DELIVERY_RETRY_DELAYS_MS]; + const parsed = raw + .split(',') + .map((value) => Number(value.trim())) + .filter((value) => Number.isFinite(value) && value >= 0); + return parsed.length > 0 + ? parsed + : [...DEFAULT_SCHEDULED_TASK_DELIVERY_RETRY_DELAYS_MS]; +} + +export function extractPublicHtmlPathsFromText(text) { + const paths = new Set(); + const normalized = String(text ?? ''); + for (const match of normalized.matchAll( + /public\/[^\s"'<>]+\.html/gi, + )) { + paths.add(match[0].replace(/\\/g, '/')); + } + for (const match of normalized.matchAll( + /\/MindSpace\/[^/\s"'<>]+\/(public\/[^\s"'<>]+\.html)/gi, + )) { + paths.add(match[1].replace(/\\/g, '/')); + } + return [...paths]; +} + +export function deliveryTextPromisesPublicHtml(text) { + const normalized = String(text ?? '').trim(); + if (!normalized) return false; + if (extractPublicHtmlPathsFromText(normalized).length > 0) return true; + return /\/MindSpace\/[^/\s"'<>]+/i.test(normalized) + && /\.html/i.test(normalized); +} + export function looksLikeScheduledTaskNonDelivery(text, { readyPaths = [] } = {}) { if (Array.isArray(readyPaths) && readyPaths.length > 0) return false; const normalized = String(text ?? '').trim(); @@ -90,6 +146,7 @@ export function looksLikeScheduledTaskNonDelivery(text, { readyPaths = [] } = {} if (SCHEDULED_TASK_CLARIFICATION_PATTERNS.some((pattern) => pattern.test(normalized))) { return true; } + if (deliveryTextPromisesPublicHtml(normalized)) return true; if (/https?:\/\//i.test(normalized)) return false; if (/public\/[^\s]+\.html/i.test(normalized)) return false; if (/页面链接/u.test(normalized)) return false; @@ -100,18 +157,35 @@ export function looksLikeScheduledTaskNonDelivery(text, { readyPaths = [] } = {} return normalized.length < 12; } -export async function finalizeScheduledTaskPageDelivery({ - pool, +async function refreshScheduledTaskMessages({ userId, sessionId, + tkmindProxy, + sessionSnapshotService, +}) { + if (typeof tkmindProxy?.fetchSessionConversationForUser === 'function') { + return tkmindProxy + .fetchSessionConversationForUser(userId, sessionId) + .catch(() => []); + } + if (typeof sessionSnapshotService?.get === 'function') { + const snapshot = await sessionSnapshotService.get(sessionId).catch(() => null); + return snapshot?.messages ?? snapshot?.conversation?.messages ?? []; + } + return []; +} + +function collectScheduledTaskPageRelativePaths({ messages, publishDir, - currentUser = null, - logger = console, -} = {}) { - if (!pool || !userId || !publishDir) return []; - - const materialized = materializeMissingPublicHtmlWrites({ messages, publishDir }); + currentUser, + userId, + deliveryText = '', +}) { + const materialized = materializeMissingPublicHtmlWrites({ + messages, + publishDir, + }); const relativePaths = new Set( collectOwnPublicHtmlRelativePaths({ messages, @@ -121,6 +195,35 @@ export async function finalizeScheduledTaskPageDelivery({ skipped: materialized.skipped, }), ); + for (const relativePath of extractPublicHtmlPathsFromText(deliveryText)) { + relativePaths.add(relativePath); + } + return { + materialized, + relativePaths: [...relativePaths], + }; +} + +export async function finalizeScheduledTaskPageDelivery({ + pool, + userId, + sessionId, + messages, + publishDir, + currentUser = null, + deliveryText = '', + logger = console, +} = {}) { + if (!pool || !userId || !publishDir) return []; + + const { relativePaths } = collectScheduledTaskPageRelativePaths({ + messages, + publishDir, + currentUser, + userId, + deliveryText, + }); + const pathsToRelease = new Set(relativePaths); if (sessionId) { const [rows] = await pool.query( @@ -130,14 +233,31 @@ export async function finalizeScheduledTaskPageDelivery({ [userId], ); for (const row of rows ?? []) { - if (row?.workspace_relative_path) relativePaths.add(row.workspace_relative_path); + if (row?.workspace_relative_path) { + pathsToRelease.add(row.workspace_relative_path); + } } } + for (const relativePath of pathsToRelease) { + await preparePageDeliveryContract({ + pool, + userId, + requestId: sessionId, + relativePath, + pgRequired: false, + }).catch((error) => { + logger.warn?.( + `[ScheduledTask] prepare delivery contract failed for ${relativePath}:`, + error, + ); + }); + } + const readyPaths = await releaseMaterializedPageDeliveryContracts({ pool, userId, - relativePaths: [...relativePaths], + relativePaths: [...pathsToRelease], allowPgRequired: true, }).catch((error) => { logger.warn?.('[ScheduledTask] release delivery contracts failed:', error); @@ -153,6 +273,117 @@ export async function finalizeScheduledTaskPageDelivery({ return readyPaths; } +export async function awaitScheduledTaskPageDelivery({ + pool, + userId, + sessionId, + messages, + publishDir, + deliveryText = '', + tkmindProxy = null, + sessionSnapshotService = null, + retryDelaysMs = resolveScheduledTaskDeliveryRetryDelaysMs(), + sleepFn = (delayMs) => new Promise((resolve) => { + setTimeout(resolve, delayMs); + }), + logger = console, +} = {}) { + let currentMessages = Array.isArray(messages) ? messages : []; + let readyPaths = []; + const attempts = [0, ...retryDelaysMs]; + + for (let attempt = 0; attempt < attempts.length; attempt += 1) { + if (attempt > 0) { + await sleepFn(attempts[attempt]); + currentMessages = await refreshScheduledTaskMessages({ + userId, + sessionId, + tkmindProxy, + sessionSnapshotService, + }); + } + readyPaths = await finalizeScheduledTaskPageDelivery({ + pool, + userId, + sessionId, + messages: currentMessages, + publishDir, + deliveryText, + logger, + }).catch((error) => { + logger.warn?.('[ScheduledTask] finalize page delivery failed:', error); + return []; + }); + + const promisesHtml = deliveryTextPromisesPublicHtml(deliveryText) + || collectScheduledTaskPageRelativePaths({ + messages: currentMessages, + publishDir, + userId, + deliveryText, + }).relativePaths.length > 0; + if (!promisesHtml || readyPaths.length > 0) { + break; + } + logger.warn?.('[ScheduledTask] page delivery still preparing', { + userId, + sessionId, + attempt: attempt + 1, + maxAttempts: attempts.length, + }); + } + + return { + messages: currentMessages, + readyPaths, + }; +} + +export async function reconcileStuckStaticPageDeliveryContracts({ + pool, + h5Root, + limit = 20, + logger = console, +} = {}) { + if (!pool || !h5Root) return []; + const [rows] = await pool.query( + `SELECT user_id, workspace_relative_path + FROM h5_page_delivery_contracts + WHERE status = 'preparing' AND data_mode = 'static' + ORDER BY updated_at ASC + LIMIT ?`, + [Math.max(1, Number(limit) || 20)], + ); + const released = []; + for (const row of rows ?? []) { + const relativePath = normalizeDeliveryRelativePath( + row.workspace_relative_path, + ); + if (!relativePath) continue; + const filePath = path.join( + h5Root, + 'MindSpace', + row.user_id, + relativePath, + ); + if (!fs.existsSync(filePath)) continue; + if ( + await markPageDeliveryContractReady({ + pool, + userId: row.user_id, + relativePath, + }) + ) { + released.push(relativePath); + logger.info?.('[ScheduledTask] reconciled static delivery contract', { + userId: row.user_id, + relativePath, + }); + } + } + return released; +} + export async function executeScheduledTask(task, { userAuth, tkmindProxy, @@ -204,30 +435,34 @@ export async function executeScheduledTask(task, { { timeoutMs }, ); - let messages = []; - if (typeof tkmindProxy.fetchSessionConversationForUser === 'function') { - messages = await tkmindProxy.fetchSessionConversationForUser(task.userId, sessionId).catch(() => []); - } else if (typeof sessionSnapshotService?.get === 'function') { - const snapshot = await sessionSnapshotService.get(sessionId).catch(() => null); - messages = snapshot?.messages ?? snapshot?.conversation?.messages ?? []; - } + let messages = await refreshScheduledTaskMessages({ + userId: task.userId, + sessionId, + tkmindProxy, + sessionSnapshotService, + }); const publishDir = h5Root && task.userId ? path.join(h5Root, 'MindSpace', task.userId) : null; - const readyPaths = await finalizeScheduledTaskPageDelivery({ - pool, - userId: task.userId, - sessionId, - messages, - publishDir, - logger, - }).catch((error) => { - logger.warn?.('[ScheduledTask] finalize page delivery failed:', error); - return []; - }); - let deliveryText = extractScheduledTaskDeliveryText(messages, task); + const deliveryResult = publishDir + ? await awaitScheduledTaskPageDelivery({ + pool, + userId: task.userId, + sessionId, + messages, + publishDir, + deliveryText, + tkmindProxy, + sessionSnapshotService, + logger, + }) + : { messages, readyPaths: [] }; + messages = deliveryResult.messages; + const readyPaths = deliveryResult.readyPaths; + + deliveryText = extractScheduledTaskDeliveryText(messages, task); if ( readyPaths.length > 0 && !/https?:\/\//i.test(deliveryText) diff --git a/scheduled-task-executor.test.mjs b/scheduled-task-executor.test.mjs index 94bc4ee..f9b1108 100644 --- a/scheduled-task-executor.test.mjs +++ b/scheduled-task-executor.test.mjs @@ -1,10 +1,18 @@ import assert from 'node:assert/strict'; +import fs from 'node:fs/promises'; +import os from 'node:os'; +import path from 'node:path'; import test from 'node:test'; import { buildScheduledTaskExecutionPrompt, + deliveryTextPromisesPublicHtml, + extractPublicHtmlPathsFromText, extractScheduledTaskDeliveryText, + finalizeScheduledTaskPageDelivery, formatScheduledTaskDeliveryMessage, looksLikeScheduledTaskNonDelivery, + reconcileStuckStaticPageDeliveryContracts, + resolveScheduledTaskDeliveryRetryDelaysMs, } from './scheduled-task-executor.mjs'; test('buildScheduledTaskExecutionPrompt includes task spec and automation marker', () => { @@ -32,6 +40,12 @@ test('looksLikeScheduledTaskNonDelivery detects clarification replies', () => { looksLikeScheduledTaskNonDelivery('页面已生成:https://example.com/news.html'), false, ); + assert.equal( + looksLikeScheduledTaskNonDelivery( + '页面已生成:https://m.tkmind.cn/MindSpace/user-1/public/daily-news.html', + ), + true, + ); assert.equal( looksLikeScheduledTaskNonDelivery('好的', { readyPaths: ['public/news.html'] }), false, @@ -41,6 +55,91 @@ test('looksLikeScheduledTaskNonDelivery detects clarification replies', () => { assert.equal(looksLikeScheduledTaskNonDelivery('好的'), true); }); +test('extractPublicHtmlPathsFromText parses MindSpace public links', () => { + assert.deepEqual( + extractPublicHtmlPathsFromText( + '链接:https://m.tkmind.cn/MindSpace/u/public/daily-news-0816.html', + ), + ['public/daily-news-0816.html'], + ); +}); + +test('deliveryTextPromisesPublicHtml detects page delivery replies', () => { + assert.equal( + deliveryTextPromisesPublicHtml('今日摘要:天气不错'), + false, + ); + assert.equal( + deliveryTextPromisesPublicHtml('public/daily-news.html 已生成'), + true, + ); +}); + +test('resolveScheduledTaskDeliveryRetryDelaysMs reads env override', () => { + assert.deepEqual( + resolveScheduledTaskDeliveryRetryDelaysMs({ + H5_SCHEDULED_TASK_DELIVERY_RETRY_DELAYS_MS: '0,100,250', + }), + [0, 100, 250], + ); +}); + +test('finalizeScheduledTaskPageDelivery prepares and releases static contracts', async () => { + const calls = []; + const pool = { + async query(sql, params) { + calls.push({ sql, params }); + if (sql.includes('FROM h5_page_delivery_contracts')) { + return [[{ workspace_relative_path: 'public/news.html' }]]; + } + if (sql.includes('SELECT id, data_mode, status')) { + return [[{ id: 'c1', data_mode: 'static', status: 'preparing' }]]; + } + if (sql.includes("SET status = 'ready'")) { + return [{ affectedRows: 1 }]; + } + return [[]]; + }, + }; + const readyPaths = await finalizeScheduledTaskPageDelivery({ + pool, + userId: 'user-1', + sessionId: 'session-1', + messages: [], + publishDir: '/tmp/publish', + deliveryText: 'public/news.html 已生成', + logger: { warn() {}, info() {} }, + }); + assert.deepEqual(readyPaths, ['public/news.html']); + assert.ok(calls.some((call) => call.sql.includes('INSERT INTO h5_page_delivery_contracts'))); +}); + +test('reconcileStuckStaticPageDeliveryContracts releases materialized static pages', async () => { + const dir = await fs.mkdtemp(path.join(os.tmpdir(), 'scheduled-task-')); + const userId = 'user-1'; + const relativePath = 'public/news.html'; + const filePath = path.join(dir, 'MindSpace', userId, relativePath); + await fs.mkdir(path.dirname(filePath), { recursive: true }); + await fs.writeFile(filePath, '', 'utf8'); + const pool = { + async query(sql) { + if (sql.includes('FROM h5_page_delivery_contracts')) { + return [[{ user_id: userId, workspace_relative_path: relativePath }]]; + } + if (sql.includes("SET status = 'ready'")) { + return [{ affectedRows: 1 }]; + } + return [[]]; + }, + }; + const released = await reconcileStuckStaticPageDeliveryContracts({ + pool, + h5Root: dir, + logger: { info() {} }, + }); + assert.deepEqual(released, [relativePath]); +}); + test('extractScheduledTaskDeliveryText reads last assistant message', () => { const text = extractScheduledTaskDeliveryText([ { role: 'user', content: [{ type: 'text', text: 'hi' }] }, diff --git a/scheduled-task-worker.mjs b/scheduled-task-worker.mjs index 923ddf0..f4a1038 100644 --- a/scheduled-task-worker.mjs +++ b/scheduled-task-worker.mjs @@ -2,6 +2,7 @@ import { executeScheduledTask, formatScheduledTaskDeliveryMessage, looksLikeScheduledTaskNonDelivery, + reconcileStuckStaticPageDeliveryContracts, } from './scheduled-task-executor.mjs'; export function startScheduledTaskWorker({ @@ -67,6 +68,15 @@ export function startScheduledTaskWorker({ if (running || stopped) return; running = true; try { + if (pool && h5Root) { + await reconcileStuckStaticPageDeliveryContracts({ + pool, + h5Root, + logger, + }).catch((err) => { + logger.warn?.('Scheduled task delivery reconcile failed:', err); + }); + } const dueTasks = await scheduledTaskService.listDueTasks({ limit: 10 }); for (const candidate of dueTasks) { const task = await scheduledTaskService.lockTask(candidate.id); diff --git a/scheduled-task-worker.test.mjs b/scheduled-task-worker.test.mjs index 829a206..0f0efbd 100644 --- a/scheduled-task-worker.test.mjs +++ b/scheduled-task-worker.test.mjs @@ -54,6 +54,7 @@ test('scheduled task worker executes due task and completes one-shot task', asyn sessionId: 'session-1', requestId: 'req-1', deliveryText: '页面:https://example.com/news.html', + readyPaths: [], }), logger: { warn() {}, info() {} }, runOnStart: false, @@ -224,6 +225,7 @@ test('scheduled task worker keeps success when wechat notification fails', async sessionId: 'session-4', requestId: 'req-4', deliveryText: '页面已生成:https://example.com/news.html', + readyPaths: [], }), logger: { warn() {} }, runOnStart: false, @@ -234,3 +236,56 @@ test('scheduled task worker keeps success when wechat notification fails', async assert.deepEqual(calls, ['notify:web', 'success:task-4']); }); + +test('scheduled task worker marks failure when page link is not deliverable yet', async () => { + const calls = []; + const task = { + id: 'task-5', + userId: 'user-5', + title: '每日新闻页', + recurrence: 'daily', + notifyChannel: 'web', + attempts: 1, + }; + const worker = startScheduledTaskWorker({ + intervalMs: 60_000, + userAuth: { id: 'user-auth' }, + tkmindProxy: { id: 'proxy' }, + scheduledTaskService: { + async listDueTasks() { + return [task]; + }, + async lockTask() { + return task; + }, + async markTaskRunning(input) { + return input; + }, + async markTaskSucceeded() { + calls.push('success'); + }, + async markTaskFailed(input, err) { + calls.push(`failed:${err.code}`); + return { ...input, status: 'failed', lastError: err.message }; + }, + }, + scheduleService: { + async createUserNotification(input) { + calls.push(`notify:${input.notificationType}`); + }, + }, + executeTask: async () => ({ + sessionId: 'session-5', + requestId: 'req-5', + deliveryText: '页面:https://m.tkmind.cn/MindSpace/user-5/public/daily-news.html', + readyPaths: [], + }), + logger: { warn() {} }, + runOnStart: false, + }); + + await worker.runOnce(); + worker.stop(); + + assert.deepEqual(calls, ['failed:SCHEDULED_TASK_NON_DELIVERY', 'notify:scheduled_task_failed']); +});