feat(FN-3795): add WhatsApp dedupe retention pruning and reports plugin sca
Merges the reports plugin scaffold (FN-3778, FN-3790) with notification service improvements across the engine, plus mailbox UI enhancements; adds WhatsApp message deduplication retention pruning (FN-3795) with tests and documentation; introduces /tasks deep-link routing with theme-data URL resoluti Fusion-Task-Id: FN-3795
This commit is contained in:
5
.changeset/fn-3795-whatsapp-dedupe-prune.md
Normal file
5
.changeset/fn-3795-whatsapp-dedupe-prune.md
Normal file
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@runfusion/fusion": patch
|
||||
---
|
||||
|
||||
Add a `dedupeRetentionDays` setting to the WhatsApp chat plugin (default 7 days) and prune old `whatsapp_chat_dedupe` rows on each inbound message to prevent unbounded dedupe-table growth.
|
||||
@@ -20,6 +20,7 @@ No Meta Cloud app, webhook URL, verify token, or Graph API credentials are requi
|
||||
- `allowedSenders`: allowed WhatsApp JIDs or E.164 digits.
|
||||
- `agentSystemPrompt`: system prompt for replies.
|
||||
- `historyTurnLimit`: persisted turn window (default `40`).
|
||||
- `dedupeRetentionDays`: replay-protection retention window for inbound message IDs (default `7` days). Rows older than this are pruned lazily whenever a new inbound message is processed.
|
||||
|
||||
## Routes
|
||||
|
||||
|
||||
@@ -1,5 +1,39 @@
|
||||
import { describe, expect, it, vi } from "vitest";
|
||||
import plugin, { splitMessageForWhatsapp } from "../index.js";
|
||||
import { describe, expect, it } from "vitest";
|
||||
import plugin, { ensureSchema, getDedupeRetentionDays, markProcessed, splitMessageForWhatsapp, wasProcessed } from "../index.js";
|
||||
|
||||
function createInMemoryDb() {
|
||||
const dedupe = new Map<string, { sender: string; receivedAt: string }>();
|
||||
|
||||
return {
|
||||
exec(_sql: string) {},
|
||||
prepare(sql: string) {
|
||||
return {
|
||||
get: (...args: unknown[]) => {
|
||||
if (sql.includes("FROM whatsapp_chat_dedupe") && sql.includes("messageId = ?")) {
|
||||
const row = dedupe.get(args[0] as string);
|
||||
return row ? { found: 1, ...row } : undefined;
|
||||
}
|
||||
return undefined;
|
||||
},
|
||||
run: (...args: unknown[]) => {
|
||||
if (sql.includes("INSERT INTO whatsapp_chat_dedupe")) {
|
||||
dedupe.set(args[0] as string, {
|
||||
sender: args[1] as string,
|
||||
receivedAt: args[2] as string,
|
||||
});
|
||||
}
|
||||
if (sql.includes("DELETE FROM whatsapp_chat_dedupe WHERE receivedAt < ?")) {
|
||||
const cutoff = args[0] as string;
|
||||
for (const [id, row] of dedupe.entries()) {
|
||||
if (row.receivedAt < cutoff) dedupe.delete(id);
|
||||
}
|
||||
}
|
||||
},
|
||||
};
|
||||
},
|
||||
_dedupe: dedupe,
|
||||
};
|
||||
}
|
||||
|
||||
describe("whatsapp plugin", () => {
|
||||
it("registers schema init hook", () => {
|
||||
@@ -19,6 +53,7 @@ describe("whatsapp plugin", () => {
|
||||
expect(Object.keys(schema).sort()).toEqual([
|
||||
"agentSystemPrompt",
|
||||
"allowedSenders",
|
||||
"dedupeRetentionDays",
|
||||
"historyTurnLimit",
|
||||
"pairingMode",
|
||||
"pairingPhoneNumber",
|
||||
@@ -31,3 +66,58 @@ describe("whatsapp plugin", () => {
|
||||
expect(chunks[0].length).toBeLessThanOrEqual(4096);
|
||||
});
|
||||
});
|
||||
|
||||
describe("markProcessed retention", () => {
|
||||
it("prunes rows older than retention and keeps recent rows", () => {
|
||||
const db = createInMemoryDb();
|
||||
ensureSchema(db as any);
|
||||
const now = Date.now();
|
||||
|
||||
db.prepare("INSERT INTO whatsapp_chat_dedupe(messageId, sender, receivedAt) VALUES(?, ?, ?)").run(
|
||||
"old-id",
|
||||
"sender",
|
||||
new Date(now - 30 * 86_400_000).toISOString(),
|
||||
);
|
||||
db.prepare("INSERT INTO whatsapp_chat_dedupe(messageId, sender, receivedAt) VALUES(?, ?, ?)").run(
|
||||
"recent-id",
|
||||
"sender",
|
||||
new Date(now - 3_600_000).toISOString(),
|
||||
);
|
||||
|
||||
markProcessed(db as any, "new-id", "sender", 7);
|
||||
|
||||
const oldRow = db.prepare("SELECT 1 as found FROM whatsapp_chat_dedupe WHERE messageId = ?").get("old-id") as { found?: number } | undefined;
|
||||
const recentRow = db.prepare("SELECT 1 as found FROM whatsapp_chat_dedupe WHERE messageId = ?").get("recent-id") as { found?: number } | undefined;
|
||||
expect(Boolean(oldRow?.found)).toBe(false);
|
||||
expect(Boolean(recentRow?.found)).toBe(true);
|
||||
expect(wasProcessed(db as any, "new-id")).toBe(true);
|
||||
});
|
||||
|
||||
it("keeps entries inside retention window", () => {
|
||||
const db = createInMemoryDb();
|
||||
ensureSchema(db as any);
|
||||
|
||||
db.prepare("INSERT INTO whatsapp_chat_dedupe(messageId, sender, receivedAt) VALUES(?, ?, ?)").run(
|
||||
"one-day-old-id",
|
||||
"sender",
|
||||
new Date(Date.now() - 86_400_000).toISOString(),
|
||||
);
|
||||
|
||||
markProcessed(db as any, "new-id", "sender", 7);
|
||||
|
||||
const oneDayOld = db.prepare("SELECT 1 as found FROM whatsapp_chat_dedupe WHERE messageId = ?").get("one-day-old-id") as { found?: number } | undefined;
|
||||
expect(Boolean(oneDayOld?.found)).toBe(true);
|
||||
});
|
||||
|
||||
it("parses dedupeRetentionDays safely", () => {
|
||||
expect(getDedupeRetentionDays({})).toBe(7);
|
||||
expect(getDedupeRetentionDays({ dedupeRetentionDays: undefined })).toBe(7);
|
||||
expect(getDedupeRetentionDays({ dedupeRetentionDays: null })).toBe(7);
|
||||
expect(getDedupeRetentionDays({ dedupeRetentionDays: 0 })).toBe(7);
|
||||
expect(getDedupeRetentionDays({ dedupeRetentionDays: -3 })).toBe(7);
|
||||
expect(getDedupeRetentionDays({ dedupeRetentionDays: "foo" })).toBe(7);
|
||||
expect(getDedupeRetentionDays({ dedupeRetentionDays: Number.POSITIVE_INFINITY })).toBe(7);
|
||||
expect(getDedupeRetentionDays({ dedupeRetentionDays: 14 })).toBe(14);
|
||||
expect(getDedupeRetentionDays({ dedupeRetentionDays: 3.7 })).toBe(3);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -3,7 +3,17 @@ import type { PluginContext } from "@fusion/plugin-sdk";
|
||||
import pino from "pino";
|
||||
import qrcode from "qrcode";
|
||||
import { clearAuthState, createPluginDbAuthState } from "./auth-state.js";
|
||||
import { getAllowedSenders, getHistoryTurnLimit, loadHistory, markProcessed, saveHistory, wasProcessed, type ChatTurn, type PluginDb } from "./index.js";
|
||||
import {
|
||||
getAllowedSenders,
|
||||
getDedupeRetentionDays,
|
||||
getHistoryTurnLimit,
|
||||
loadHistory,
|
||||
markProcessed,
|
||||
saveHistory,
|
||||
wasProcessed,
|
||||
type ChatTurn,
|
||||
type PluginDb,
|
||||
} from "./index.js";
|
||||
|
||||
export type ConnectionStatus = {
|
||||
state: "starting" | "awaiting-qr" | "awaiting-code" | "connected" | "disconnected" | "error";
|
||||
@@ -181,7 +191,7 @@ export class WhatsAppConnection {
|
||||
if (allowedSenders.size === 0 || (!allowedSenders.has(sender) && !allowedSenders.has(jid))) continue;
|
||||
if (wasProcessed(this.db, messageId)) continue;
|
||||
|
||||
markProcessed(this.db, messageId, sender);
|
||||
markProcessed(this.db, messageId, sender, getDedupeRetentionDays(this.ctx.settings));
|
||||
|
||||
try {
|
||||
const history = loadHistory(this.db, sender);
|
||||
|
||||
@@ -4,6 +4,7 @@ import { WhatsAppConnection } from "./connection.js";
|
||||
import { generateReply } from "./reply.js";
|
||||
|
||||
const DEFAULT_HISTORY_TURN_LIMIT = 40;
|
||||
const DEFAULT_DEDUPE_RETENTION_DAYS = 7;
|
||||
|
||||
export type ChatTurn = { role: "user" | "assistant"; text: string; createdAt: string };
|
||||
|
||||
@@ -39,6 +40,12 @@ const settingsSchema: Record<string, PluginSettingSchema> = {
|
||||
label: "History Turn Limit",
|
||||
defaultValue: DEFAULT_HISTORY_TURN_LIMIT,
|
||||
},
|
||||
dedupeRetentionDays: {
|
||||
type: "number",
|
||||
label: "Dedupe Retention (days)",
|
||||
description: "How long inbound message IDs are kept for replay protection. Older rows are pruned on each inbound message.",
|
||||
defaultValue: DEFAULT_DEDUPE_RETENTION_DAYS,
|
||||
},
|
||||
};
|
||||
|
||||
const connections = new Map<string, WhatsAppConnection>();
|
||||
@@ -62,6 +69,14 @@ export function getHistoryTurnLimit(settings: Record<string, unknown>): number {
|
||||
return Math.floor(value);
|
||||
}
|
||||
|
||||
export function getDedupeRetentionDays(settings: Record<string, unknown>): number {
|
||||
const value = settings.dedupeRetentionDays;
|
||||
if (typeof value !== "number" || !Number.isFinite(value) || value <= 0) {
|
||||
return DEFAULT_DEDUPE_RETENTION_DAYS;
|
||||
}
|
||||
return Math.floor(value);
|
||||
}
|
||||
|
||||
export function splitMessageForWhatsapp(text: string): string[] {
|
||||
return WhatsAppConnection.splitMessageForWhatsapp(text);
|
||||
}
|
||||
@@ -121,8 +136,16 @@ export function wasProcessed(db: PluginDb, messageId: string): boolean {
|
||||
return Boolean(row?.found);
|
||||
}
|
||||
|
||||
export function markProcessed(db: PluginDb, messageId: string, sender: string): void {
|
||||
db.prepare("INSERT INTO whatsapp_chat_dedupe(messageId, sender, receivedAt) VALUES(?, ?, ?)").run(messageId, sender, new Date().toISOString());
|
||||
export function markProcessed(
|
||||
db: PluginDb,
|
||||
messageId: string,
|
||||
sender: string,
|
||||
retentionDays: number = DEFAULT_DEDUPE_RETENTION_DAYS,
|
||||
): void {
|
||||
const now = new Date().toISOString();
|
||||
const cutoff = new Date(Date.now() - retentionDays * 86_400_000).toISOString();
|
||||
db.prepare("DELETE FROM whatsapp_chat_dedupe WHERE receivedAt < ?").run(cutoff);
|
||||
db.prepare("INSERT INTO whatsapp_chat_dedupe(messageId, sender, receivedAt) VALUES(?, ?, ?)").run(messageId, sender, now);
|
||||
}
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user