Keep daily morning greetings scheduled after WeChat 45015 failures.
A terminal deferred delivery was marking the reminder failed without creating the next daily run, which silently killed active 早安问候 subscriptions. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -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);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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']);
|
||||
});
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
Reference in New Issue
Block a user