diff --git a/extensions/tlon/src/monitor/index.ts b/extensions/tlon/src/monitor/index.ts index db7f91f53a1b..251b8a2907f5 100644 --- a/extensions/tlon/src/monitor/index.ts +++ b/extensions/tlon/src/monitor/index.ts @@ -36,6 +36,7 @@ import { mergeUniqueStrings, shouldMigrateTlonSetting, } from "./settings-helpers.js"; +import { createActiveSnapshotTracker, createParticipatedThreadTracker } from "./tracking.js"; import { asRecord, formatErrorMessage, readString } from "./utils.js"; import { extractMessageText, @@ -160,8 +161,8 @@ export async function monitorTlonProvider(opts: MonitorTlonOpts = {}): Promise(); + // Track recent threads we've participated in so replies can omit a mention. + const participatedThreads = createParticipatedThreadTracker(); // Track DM senders per session to detect shared sessions (security warning) const dmSendersBySession = new Map>(); @@ -894,8 +895,8 @@ export async function monitorTlonProvider(opts: MonitorTlonOpts = {}): Promise(); + // Track processed DM invites only while they remain in the active /v3 snapshot. + const processedDmInvites = createActiveSnapshotTracker(); const handleChatFirehose = async ( event: unknown, @@ -904,9 +905,16 @@ export async function monitorTlonProvider(opts: MonitorTlonOpts = {}): Promise normalizeShip(invite.ship || "")).filter(Boolean), + ); + + for (const ship of ships) { + if (processedDmInvites.has(ship)) { continue; } @@ -926,11 +934,9 @@ export async function monitorTlonProvider(opts: MonitorTlonOpts = {}): Promise { + it("evicts the least recently used thread at the configured limit", () => { + const tracker = createParticipatedThreadTracker(3); + tracker.add("oldest"); + tracker.add("refreshed"); + tracker.add("recent"); + + expect(tracker.has("refreshed")).toBe(true); + tracker.add("newest"); + + expect(tracker.has("oldest")).toBe(false); + expect(tracker.has("refreshed")).toBe(true); + expect(tracker.has("recent")).toBe(true); + expect(tracker.has("newest")).toBe(true); + }); +}); + +describe("createActiveSnapshotTracker", () => { + it("forgets processed keys after they leave the active snapshot", () => { + const tracker = createActiveSnapshotTracker(); + expect(tracker.beginSnapshot(["active", "removed"])).toEqual(new Set(["active", "removed"])); + tracker.add("active"); + tracker.add("removed"); + + tracker.beginSnapshot(["active"]); + expect(tracker.has("active")).toBe(true); + expect(tracker.has("removed")).toBe(false); + + tracker.beginSnapshot(["active", "removed"]); + expect(tracker.has("removed")).toBe(false); + }); + + it("does not impose a count cap on the authoritative active snapshot", () => { + const tracker = createActiveSnapshotTracker(); + const keys = Array.from({ length: 2_001 }, (_, index) => `invite-${index}`); + const active = tracker.beginSnapshot(keys); + for (const key of active) { + tracker.add(key); + } + + expect(keys.every((key) => tracker.has(key))).toBe(true); + }); +}); diff --git a/extensions/tlon/src/monitor/tracking.ts b/extensions/tlon/src/monitor/tracking.ts new file mode 100644 index 000000000000..8adda3d098b6 --- /dev/null +++ b/extensions/tlon/src/monitor/tracking.ts @@ -0,0 +1,40 @@ +// Tlon monitor module owns bounded and snapshot-scoped identifier tracking. +import { createDedupeCache } from "../../runtime-api.js"; + +const TLON_PARTICIPATED_THREAD_LIMIT = 2_000; + +export function createParticipatedThreadTracker(limit = TLON_PARTICIPATED_THREAD_LIMIT) { + const cache = createDedupeCache({ ttlMs: 0, maxSize: limit }); + + return { + add: (parentId: string) => { + cache.check(parentId); + }, + has: (parentId: string) => { + if (!cache.peek(parentId)) { + return false; + } + // Mention-free replies refresh recency before older participation is evicted. + cache.check(parentId); + return true; + }, + }; +} + +export function createActiveSnapshotTracker() { + const processed = new Set(); + + return { + beginSnapshot: (keys: Iterable): ReadonlySet => { + const active = new Set(keys); + for (const key of processed) { + if (!active.has(key)) { + processed.delete(key); + } + } + return active; + }, + has: (key: string) => processed.has(key), + add: (key: string) => processed.add(key), + }; +}