diff --git a/schedule-reminder-worker.mjs b/schedule-reminder-worker.mjs index 91722ee..5c0aca5 100644 --- a/schedule-reminder-worker.mjs +++ b/schedule-reminder-worker.mjs @@ -106,6 +106,13 @@ export function startScheduleReminderWorker({ errorMessage: err instanceof Error ? err.message : String(err), }).catch(() => {}); await scheduleService.markReminderFailed(reminder, err, { maxAttempts }); + if (Number(reminder.attempts ?? 0) >= maxAttempts) { + try { + await scheduleService.scheduleNextDailyReminder?.(reminder); + } catch (scheduleErr) { + logger.warn?.('Schedule next daily reminder failed:', scheduleErr); + } + } } } diff --git a/schedule-reminder-worker.test.mjs b/schedule-reminder-worker.test.mjs index a414fc6..4507773 100644 --- a/schedule-reminder-worker.test.mjs +++ b/schedule-reminder-worker.test.mjs @@ -434,6 +434,9 @@ test('schedule reminder worker retries when wechat delivery is deferred', async async markReminderSent() { calls.push('sent'); }, + async scheduleNextDailyReminder() { + calls.push('next'); + }, async markReminderFailed() { calls.push('failed'); }, @@ -465,3 +468,68 @@ test('schedule reminder worker retries when wechat delivery is deferred', async assert.deepEqual(calls, ['log:failed', 'failed']); }); + +test('schedule reminder worker still schedules next daily reminder after terminal deferred failure', async () => { + const calls = []; + const reminder = { + id: 'rem-deferred-terminal', + userId: 'user-1', + itemId: 'item-morning', + remindAt: Date.now() - 1000, + channel: 'wechat', + attempts: 5, + }; + const worker = startScheduleReminderWorker({ + intervalMs: 60_000, + maxAttempts: 5, + scheduleService: { + async listDueReminders() { + return [reminder]; + }, + async lockReminder() { + return reminder; + }, + async buildReminderText() { + return '☀️ 早安'; + }, + async createUserNotification() {}, + async logDelivery(input) { + calls.push(`log:${input.status}`); + }, + async markReminderSent() { + calls.push('sent'); + }, + async scheduleNextDailyReminder() { + calls.push('next'); + }, + async markReminderFailed() { + calls.push('failed'); + }, + async markReminderCancelled() {}, + async listDueDigestSubscriptions() { + return []; + }, + async lockDigestSubscription() { + return null; + }, + async listDueBalanceAlerts() { + return []; + }, + async lockBalanceAlert() { + return null; + }, + }, + notificationDispatcher: { + async sendScheduleNotification() { + return { sent: false, deferred: true, errcode: 45015 }; + }, + }, + logger: { warn() {} }, + runOnStart: false, + }); + + await worker.runOnce(); + worker.stop(); + + assert.deepEqual(calls, ['log:failed', 'failed', 'next']); +}); diff --git a/scripts/repair-morning-greeting-chain-103.mjs b/scripts/repair-morning-greeting-chain-103.mjs new file mode 100644 index 0000000..aa5f312 --- /dev/null +++ b/scripts/repair-morning-greeting-chain-103.mjs @@ -0,0 +1,154 @@ +#!/usr/bin/env node +/** + * 恢复生产早安问候日更链:给已 active 但没有 pending/locked 提醒的订阅补上下一次推送。 + * + * node scripts/repair-morning-greeting-chain-103.mjs + * node scripts/repair-morning-greeting-chain-103.mjs --apply + */ +import { execSync } from 'node:child_process'; +import crypto from 'node:crypto'; +import mysql from 'mysql2/promise'; +import { nextDailyRunAt } from '../schedule-time.mjs'; +import { + LEGACY_MORNING_SOURCE, + pushSubscriptionMetadataSource, +} from '../wechat/push-subscription-catalog.mjs'; + +const HOST = process.env.STUDIO_HOST || 'john@114.85.107.50'; +const REMOTE_ENV = '/Users/john/Project/Memind/.env'; +const apply = process.argv.includes('--apply'); +const morningSources = [LEGACY_MORNING_SOURCE, pushSubscriptionMetadataSource('morning')]; + +function loadProdDatabaseUrl() { + if (process.env.DATABASE_URL && process.env.REPAIR_USE_LOCAL_DB === '1') { + return process.env.DATABASE_URL; + } + const raw = execSync( + `ssh -o BatchMode=yes -o ConnectTimeout=15 ${HOST} "grep '^DATABASE_URL=' ${REMOTE_ENV}"`, + { encoding: 'utf8' }, + ).trim(); + const value = raw.replace(/^DATABASE_URL=/, '').trim(); + if ( + (value.startsWith('"') && value.endsWith('"')) + || (value.startsWith("'") && value.endsWith("'")) + ) { + return value.slice(1, -1); + } + return value; +} + +function fmt(ms) { + return new Date(Number(ms)).toLocaleString('zh-CN', { timeZone: 'Asia/Shanghai' }); +} + +function parseMetadata(raw) { + if (raw == null || raw === '') return {}; + if (typeof raw === 'object') return raw; + try { + return JSON.parse(raw); + } catch { + return {}; + } +} + +async function main() { + const databaseUrl = loadProdDatabaseUrl(); + if (!databaseUrl) { + throw new Error('缺少 DATABASE_URL'); + } + const pool = mysql.createPool({ uri: databaseUrl, connectionLimit: 2 }); + const now = Date.now(); + const created = []; + const skipped = []; + + try { + const [items] = await pool.query( + `SELECT i.id, i.user_id, i.timezone, i.metadata_json, u.username + FROM h5_schedule_items i + JOIN h5_users u ON u.id = i.user_id + WHERE i.deleted_at IS NULL + AND i.status = 'active' + AND JSON_UNQUOTE(JSON_EXTRACT(i.metadata_json, '$.source')) IN (?, ?)`, + morningSources, + ); + + for (const item of items) { + const metadata = parseMetadata(item.metadata_json); + const hour = Number(metadata.dailyHour); + const minute = Number(metadata.dailyMinute ?? 0); + const [pending] = await pool.query( + `SELECT id, remind_at, status + FROM h5_schedule_reminders + WHERE item_id = ? AND status IN ('pending', 'locked') + ORDER BY remind_at ASC + LIMIT 1`, + [item.id], + ); + if (pending[0]) { + skipped.push({ + username: item.username, + reason: 'already-scheduled', + status: pending[0].status, + remindAt: fmt(pending[0].remind_at), + }); + continue; + } + if (metadata.recurrence !== 'daily' || !Number.isInteger(hour) || hour < 0 || hour > 23) { + skipped.push({ username: item.username, reason: 'invalid-schedule', hour, minute }); + continue; + } + if (!Number.isInteger(minute) || minute < 0 || minute > 59) { + skipped.push({ username: item.username, reason: 'invalid-minute', hour, minute }); + continue; + } + + const remindAt = nextDailyRunAt({ + hour, + minute, + timezone: item.timezone || 'Asia/Shanghai', + now: now + 1000, + }); + const row = { + username: item.username, + hour, + minute, + remindAt: fmt(remindAt), + }; + if (!apply) { + created.push({ ...row, dryRun: true }); + continue; + } + const id = crypto.randomUUID(); + await pool.query( + `INSERT INTO h5_schedule_reminders + (id, user_id, item_id, remind_at, offset_minutes, channel, status, attempts, + last_error, locked_until, sent_at, created_at, updated_at) + VALUES (?, ?, ?, ?, NULL, 'wechat', 'pending', 0, NULL, NULL, NULL, ?, ?) + ON DUPLICATE KEY UPDATE + offset_minutes = VALUES(offset_minutes), + status = 'pending', + attempts = 0, + last_error = NULL, + locked_until = NULL, + sent_at = NULL, + updated_at = VALUES(updated_at)`, + [id, item.user_id, item.id, remindAt, now, now], + ); + created.push(row); + } + + console.log(JSON.stringify({ + ok: true, + apply, + created, + skipped, + }, null, 2)); + } finally { + await pool.end(); + } +} + +main().catch((error) => { + console.error(error instanceof Error ? error.message : error); + process.exit(1); +});