Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
161 changes: 87 additions & 74 deletions src/features/relay/live.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import { getLogger, setLogLevel } from "../developer/logging";
import { assert, afterEach, expect, it, vi } from "vitest";
import { assert, afterEach, beforeAll, describe, expect, it, vi } from "vitest";
import {
createLiveAdmission,
liveChannels,
Expand Down Expand Up @@ -205,82 +205,95 @@ it("batches joined background interests without rebalance and fences retired bat
expect(vi.getTimerCount()).toBe(0);
});

it("keeps quiet-channel unread evidence when another filter fills its replay allowance", async () => {
vi.useFakeTimers();
const h = setup([]);
const relay = keypair(),
author = keypair();
const wire = scriptedTransport(h.key.pubkey, relay.pubkey);
const owner = createRelaySession({
...wire.transport,
subscribe(callbacks) {
h.callbacks.receive.mockImplementation(callbacks.receive);
h.callbacks.state.mockImplementation(callbacks.state);
return h.owner;
},
describe("per-channel replay allowance", () => {
const now = 1_700_000_100;
let hot: ReturnType<typeof message>[];
beforeAll(() => {
// Signing also verifies each signature. Keep large fixture construction
// separate from the behavior budget; socket delivery still verifies all 501
// distinct received events through the real live/session path.
const author = keypair();
hot = Array.from({ length: 501 }, (_, i) =>
message(author, "hot", `hot-${i}`, now - 1),
);
});
try {
const ids = ["hot", "quiet"];
h.callbacks.receive(ids.map((id) => roster(relay, id, [h.key.pubkey])));
await h.first.auth();
const batch = h.first.requests().find((r) => scopeOf(r).includes("quiet"));
assert.exists(batch);
const incoming = vi.fn();
owner.session.subscribeIncoming(incoming);
const now = Math.floor(Date.now() / 1000);
const quiet = signed(author, {
kind: 9,
content: "quiet mention",
created_at: now - 100,
tags: [
["h", "quiet"],
["p", h.key.pubkey],
],
it("keeps quiet-channel unread evidence when another filter fills its replay allowance", async () => {
vi.useFakeTimers({ now: now * 1000 });
const h = setup([]);
const relay = keypair(),
author = keypair();
const wire = scriptedTransport(h.key.pubkey, relay.pubkey);
const owner = createRelaySession({
...wire.transport,
subscribe(callbacks) {
h.callbacks.receive.mockImplementation(callbacks.receive);
h.callbacks.state.mockImplementation(callbacks.state);
return h.owner;
},
});
const history = [
quiet,
...Array.from({ length: 501 }, (_, i) =>
message(author, "hot", `hot-${i}`, now - 1),
),
];
// Relay's existing OR contract applies each filter's limit separately.
// The former multi-h filter loses quiet to the 500 newer hot events.
const delivered = new Set<string>();
for (const filter of filtersOf(batch)) {
const matches = history
.filter(
(e) =>
filter.kinds.includes(e.kind) &&
e.created_at >= filter.since &&
e.tags.some(
([k, v]) =>
k === "h" && v !== undefined && filter["#h"]?.includes(v),
),
)
.sort((a, b) => b.created_at - a.created_at)
.slice(0, filter.limit);
for (const event of matches)
if (!delivered.has(event.id)) {
delivered.add(event.id);
await h.first.receive(["EVENT", batch[1], event]);
}
try {
const ids = ["hot", "quiet"];
h.callbacks.receive(ids.map((id) => roster(relay, id, [h.key.pubkey])));
await h.first.auth();
const batch = h.first
.requests()
.find((r) => scopeOf(r).includes("quiet"));
assert.exists(batch);
const incoming = vi.fn();
owner.session.subscribeIncoming(incoming);
const quiet = signed(author, {
kind: 9,
content: "quiet mention",
created_at: now - 100,
tags: [
["h", "quiet"],
["p", h.key.pubkey],
],
});
const history = [quiet, ...hot];
// Relay's existing OR contract applies each filter's limit separately.
// The former multi-h filter loses quiet to the 500 newer hot events.
const delivered = new Set<string>();
for (const filter of filtersOf(batch)) {
const matches = history
.filter(
(e) =>
filter.kinds.includes(e.kind) &&
e.created_at >= filter.since &&
e.tags.some(
([k, v]) =>
k === "h" && v !== undefined && filter["#h"]?.includes(v),
),
)
.sort((a, b) => b.created_at - a.created_at)
.slice(0, filter.limit);
for (const event of matches)
if (!delivered.has(event.id)) {
delivered.add(event.id);
await h.first.receive(["EVENT", batch[1], event]);
}
}
await h.first.receive(["EOSE", batch[1]]);
expect(delivered.has(quiet.id)).toBe(true);
expect(delivered.size).toBe(501);
expect(
owner.session.unread.snapshot({ kind: "channel", channelId: "hot" }),
).toMatchObject({ observedCount: 500, attentionCount: 0 });
expect(
owner.session.unread.snapshot({ kind: "channel", channelId: "quiet" }),
).toMatchObject({ observedCount: 1, attentionCount: 1 });
expect(incoming).not.toHaveBeenCalled();
// Aggregate replay evidence stays conservative for every logical channel.
expect(
h.callbacks.state.mock.lastCall?.[0].routes
.filter((r) => r.channelId)
.every((r) => r.replay === "limited"),
).toBe(true);
} finally {
owner.dispose();
}
await h.first.receive(["EOSE", batch[1]]);
expect(delivered.has(quiet.id)).toBe(true);
expect(
owner.session.unread.snapshot({ kind: "channel", channelId: "quiet" }),
).toMatchObject({ observedCount: 1, attentionCount: 1 });
expect(incoming).not.toHaveBeenCalled();
// Aggregate replay evidence stays conservative for every logical channel.
expect(
h.callbacks.state.mock.lastCall?.[0].routes
.filter((r) => r.channelId)
.every((r) => r.replay === "limited"),
).toBe(true);
} finally {
owner.dispose();
}
expect(vi.getTimerCount()).toBe(0);
expect(vi.getTimerCount()).toBe(0);
});
});

it.each(["single", "chained"])(
Expand Down
Loading