141 lines
5.0 KiB
TypeScript
141 lines
5.0 KiB
TypeScript
|
|
/**
|
|||
|
|
* BridgeDeliveryLog — V4.8.4 cross-channel dedup.
|
|||
|
|
*
|
|||
|
|
* Records per-`(address, msgId)` "delivered via bridge push" timestamps so
|
|||
|
|
* the inbox-fetch route can filter out blobs the bridge has already pushed
|
|||
|
|
* to the recipient. The relay's contract becomes:
|
|||
|
|
*
|
|||
|
|
* one `Inbox.send` ⇒ one observable delivery on the recipient
|
|||
|
|
*
|
|||
|
|
* even when the recipient runs both a bridge subscription (WS / SSE) AND
|
|||
|
|
* the regular inbox-poll. Without it, bridge-push and inbox-poll are
|
|||
|
|
* independent paths against the same store and the recipient gets the
|
|||
|
|
* same envelope twice — bridge-first, then ~30 s later via the next poll
|
|||
|
|
* — tripping on already-consumed prekeys (`one-time prekey not found`)
|
|||
|
|
* or surfacing as duplicate `shade.receive` work.
|
|||
|
|
*
|
|||
|
|
* The log is in-memory per process and intentionally bounded: each entry
|
|||
|
|
* lives for `graceMs` (default 60 s, well past a typical `pollIntervalMs`
|
|||
|
|
* of 30 s). After grace, the entry is forgotten and inbox-poll falls back
|
|||
|
|
* to delivering the blob — that's the legitimate "bridge dropped the
|
|||
|
|
* frame, poll picked up" recovery path. If the recipient explicitly
|
|||
|
|
* acks the blob (HTTP `DELETE /v1/inbox/:addr/:msgId`), the blob is gone
|
|||
|
|
* from storage and the log entry is moot.
|
|||
|
|
*
|
|||
|
|
* Multi-bridge per address (e.g. WS + SSE redundancy on the same client,
|
|||
|
|
* or two devices sharing one signing key) is preserved: every bridge
|
|||
|
|
* connection still fetches + pushes the blob — each push records its own
|
|||
|
|
* timestamp — so each connected bridge gets the frame. Only the *poll*
|
|||
|
|
* fetch is filtered, not the bridge fetches themselves.
|
|||
|
|
*
|
|||
|
|
* @see Prism FR `cross-channel-duplicate-fanout-v4.8.2.md`.
|
|||
|
|
*/
|
|||
|
|
|
|||
|
|
const DEFAULT_GRACE_MS = 60_000;
|
|||
|
|
|
|||
|
|
export interface BridgeDeliveryLogOptions {
|
|||
|
|
/**
|
|||
|
|
* How long a `(address, msgId)` mark suppresses inbox-poll delivery.
|
|||
|
|
* Defaults to 60_000ms — twice the default `pollIntervalMs` of the
|
|||
|
|
* `@shade/inbox` orchestrator, so a poll cycle that races a bridge
|
|||
|
|
* push always sees the mark, but a stuck recipient still gets the
|
|||
|
|
* blob via poll within ~minutes.
|
|||
|
|
*/
|
|||
|
|
graceMs?: number;
|
|||
|
|
/**
|
|||
|
|
* Maximum entries per address. Bounds memory under a busy address.
|
|||
|
|
* Oldest entries (by recorded timestamp) are evicted first. Default
|
|||
|
|
* 8192 — comfortably above any realistic backlog.
|
|||
|
|
*/
|
|||
|
|
maxPerAddress?: number;
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
export class BridgeDeliveryLog {
|
|||
|
|
private readonly log = new Map<string, Map<string, number>>();
|
|||
|
|
private readonly graceMs: number;
|
|||
|
|
private readonly maxPerAddress: number;
|
|||
|
|
|
|||
|
|
constructor(options: BridgeDeliveryLogOptions = {}) {
|
|||
|
|
this.graceMs = options.graceMs ?? DEFAULT_GRACE_MS;
|
|||
|
|
this.maxPerAddress = options.maxPerAddress ?? 8192;
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
/** Mark `(address, msgId)` as bridge-delivered at `now`. */
|
|||
|
|
recordDelivered(address: string, msgId: string, now: number): void {
|
|||
|
|
let inner = this.log.get(address);
|
|||
|
|
if (!inner) {
|
|||
|
|
inner = new Map();
|
|||
|
|
this.log.set(address, inner);
|
|||
|
|
}
|
|||
|
|
inner.set(msgId, now);
|
|||
|
|
// Lazy cleanup: drop entries past 2× grace so the map stays bounded
|
|||
|
|
// without a separate timer. Bound by `maxPerAddress` as a fallback
|
|||
|
|
// for pathological burst scenarios.
|
|||
|
|
if (inner.size > this.maxPerAddress) {
|
|||
|
|
const cutoff = now - this.graceMs * 2;
|
|||
|
|
for (const [id, ts] of inner) {
|
|||
|
|
if (ts < cutoff) inner.delete(id);
|
|||
|
|
}
|
|||
|
|
// Still over cap? Drop the oldest.
|
|||
|
|
if (inner.size > this.maxPerAddress) {
|
|||
|
|
const sorted = Array.from(inner.entries()).sort((a, b) => a[1] - b[1]);
|
|||
|
|
const toDrop = sorted.slice(0, inner.size - this.maxPerAddress);
|
|||
|
|
for (const [id] of toDrop) inner.delete(id);
|
|||
|
|
}
|
|||
|
|
}
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
/**
|
|||
|
|
* Returns true if `(address, msgId)` was bridge-delivered within the
|
|||
|
|
* grace window.
|
|||
|
|
*/
|
|||
|
|
isRecentlyDelivered(address: string, msgId: string, now: number): boolean {
|
|||
|
|
const inner = this.log.get(address);
|
|||
|
|
if (!inner) return false;
|
|||
|
|
const ts = inner.get(msgId);
|
|||
|
|
if (ts === undefined) return false;
|
|||
|
|
if (now - ts > this.graceMs) {
|
|||
|
|
inner.delete(msgId); // tombstone the stale entry
|
|||
|
|
return false;
|
|||
|
|
}
|
|||
|
|
return true;
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
/**
|
|||
|
|
* Filter `blobs` down to those not currently in the bridge-delivered
|
|||
|
|
* grace window. Used by the inbox-fetch route to suppress duplicates.
|
|||
|
|
*/
|
|||
|
|
filterRecent<T extends { msgId: string }>(
|
|||
|
|
address: string,
|
|||
|
|
blobs: T[],
|
|||
|
|
now: number,
|
|||
|
|
): T[] {
|
|||
|
|
const inner = this.log.get(address);
|
|||
|
|
if (!inner || inner.size === 0) return blobs;
|
|||
|
|
return blobs.filter((b) => {
|
|||
|
|
const ts = inner.get(b.msgId);
|
|||
|
|
if (ts === undefined) return true;
|
|||
|
|
if (now - ts > this.graceMs) {
|
|||
|
|
inner.delete(b.msgId);
|
|||
|
|
return true;
|
|||
|
|
}
|
|||
|
|
return false;
|
|||
|
|
});
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
/** Drop the entry for `(address, msgId)`. Called from blob-delete paths. */
|
|||
|
|
forget(address: string, msgId: string): void {
|
|||
|
|
this.log.get(address)?.delete(msgId);
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
/** Drop every entry for `address`. Called from address-delete paths. */
|
|||
|
|
forgetAddress(address: string): void {
|
|||
|
|
this.log.delete(address);
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
/** Test-only inspection. */
|
|||
|
|
size(address: string): number {
|
|||
|
|
return this.log.get(address)?.size ?? 0;
|
|||
|
|
}
|
|||
|
|
}
|