Skip to content
Merged
Show file tree
Hide file tree
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
22 changes: 19 additions & 3 deletions src/features/relay/store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -870,6 +870,7 @@ export function createChannelStore(
epoch++;
hydration = undefined;
warmCandidates.clear();
warmEligible.clear();
warmPreferred = [];
media.dispose();
media = createMediaPreparation();
Expand All @@ -883,6 +884,9 @@ export function createChannelStore(
* for foreground slots and is dropped wholesale when the session resets. */
let warmPreferred: readonly string[] = [];
const warmCandidates = new Set<string>();
// Eligibility outlives head-cache retention. Consuming a candidate (including
// failure or yielding to demand) must not requeue it on the next preview update.
let warmEligible = new Set<string>();
function nextWarmId(): string | undefined {
for (const channelId of warmPreferred)
if (warmCandidates.has(channelId)) return channelId;
Expand Down Expand Up @@ -948,9 +952,16 @@ export function createChannelStore(
if (disposed || !transport || list.status !== "ready") return;
const starred = new Set(preferred);
warmPreferred = preferred;
for (const channel of list.channels)
if (!channel.archived || starred.has(channel.id))
warmCandidates.add(channel.id);
const eligible = new Set(
list.channels
.filter((channel) => !channel.archived || starred.has(channel.id))
.map((channel) => channel.id),
);
for (const id of warmCandidates)
if (!eligible.has(id)) warmCandidates.delete(id);
for (const id of eligible)
if (!warmEligible.has(id)) warmCandidates.add(id);
warmEligible = eligible;
void drainWarm();
},
ensure(channelId: string) {
Expand Down Expand Up @@ -1231,6 +1242,11 @@ export function createChannelStore(
const hadHydration = hydration !== undefined;
epoch++;
hydration = undefined;
for (const id of warmEligible)
if (!authorized(id)) {
warmEligible.delete(id);
warmCandidates.delete(id);
}
media.dispose();
media = createMediaPreparation();
for (const controller of controllers) controller.abort();
Expand Down
293 changes: 293 additions & 0 deletions src/features/relay/warm-lifecycle.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,293 @@
import { afterEach, expect, it, vi } from "vitest";
import { createRelaySession } from "./session";
import type { LiveCallbacks } from "./live";
import { ReadError } from "./errors";
import {
bounds,
flush,
keypair,
message,
metadata,
profile,
roster,
scriptedTransport,
signed,
} from "./testing";

const relay = keypair(),
viewer = keypair();
const owners: ReturnType<typeof createRelaySession>[] = [];
afterEach(() => {
for (const owner of owners.splice(0)) owner.dispose();
});
const head = (id: string, content = `Preview ${id}`) => [
message(viewer, id, content, 1_700_000_000),
bounds(relay, id, "head", { has_more: false, next_cursor: null }),
];
const discovery = (ids: readonly string[]) =>
ids.flatMap((id) => [
roster(relay, id, [viewer.pubkey]),
metadata(relay, id, id),
]);
async function setup(
ids: string[],
options: Parameters<typeof createRelaySession>[1] = {},
) {
const h = scriptedTransport(viewer.pubkey, relay.pubkey);
let live!: LiveCallbacks;
let starred: readonly string[] = [];
const transport = {
...h.transport,
subscribe(callbacks: LiveCallbacks) {
live = callbacks;
return { update() {}, retry() {}, dispose() {} };
},
decodeSidebarPreferences: async () => ({
sections: [],
assignments: {},
starred,
}),
query: vi.fn((...args: Parameters<typeof h.transport.query>) =>
args[0].some((filter) => filter.kinds?.includes(0))
? Promise.resolve([profile(viewer, { name: "Viewer" })])
: args[0].some((filter) => filter.kinds?.includes(30078))
? Promise.resolve([])
: h.transport.query(...args),
),
};
const store = createRelaySession(transport, {
prepared: true,
warm: true,
...options,
});
owners.push(store);
const queries = store.session.channels;
queries.ensureList();
h.next().respond(discovery(ids));
await flush();
return {
...h,
store,
queries,
transport,
live,
async setStars(ids: readonly string[]) {
starred = ids;
await store.session.sidebarPreferences.refresh();
},
};
}

it.each([
["entry", { maxHeads: 2 }],
["byte", { maxHeadBytes: 2500 }],
] as const)(
"finishes populated warming beyond the %s budget without requeueing evicted heads",
async (_budget, options) => {
let clock = Date.now();
const h = await setup(["a", "b", "c"], { ...options, now: () => clock });
for (const id of ["a", "b", "c"]) {
expect(h.pending).toHaveLength(1);
const request = h.next();
expect(request.filters[0]).toMatchObject({
"#h": [id],
top_level: true,
limit: 20,
});
request.respond(head(id));
await flush();
}
expect(h.store.diagnostics().heads.entries).toBeLessThan(3);
// Every response and its preview notification has settled; no live events or
// reconnects are needed to provoke the old self-replenishing queue.
expect(h.pending).toHaveLength(0);
clock += 61_000; // Expiry is not a fresh optional obligation either.
h.queries.refresh?.("c");
h.next().respond(head("c", "Changed preview"));
await flush();
expect(
h.queries.list().channels.find((channel) => channel.id === "c")?.preview,
).toBe("Changed preview");
expect(h.pending).toHaveLength(0);
// Eviction is not a new warm obligation, but opening that channel still reads.
h.queries.ensure("a");
expect(h.next().filters[0]?.["#h"]).toEqual(["a"]);
},
);

it("finishes a 65-channel populated roster with the production cache limits", async () => {
const ids = Array.from(
{ length: 65 },
(_, index) => `channel-${String(index).padStart(2, "0")}`,
);
const h = await setup(ids);
for (const id of ids) {
expect(h.pending).toHaveLength(1);
const request = h.next();
expect(request.filters[0]?.["#h"]).toEqual([id]);
request.respond(head(id));
await flush();
}
expect(h.store.diagnostics().heads.entries).toBe(64);
expect(h.pending).toHaveLength(0);
});

it("warms newly eligible channels and updates queued starred-first ordering without restarting completed work", async () => {
const h = await setup(["a", "b", "c"], { maxHeads: 1 });
const first = h.next();
await h.setStars(["c"]);
first.respond(head("a"));
await flush();
expect(h.pending[0]?.filters[0]?.["#h"]).toEqual(["c"]);
h.next().respond(head("c"));
await flush();
expect(h.pending[0]?.filters[0]?.["#h"]).toEqual(["b"]);
h.next().respond(head("b"));
await flush();
expect(h.pending).toHaveLength(0);
h.queries.refreshList?.();
h.next().respond(discovery(["a", "b", "c", "d"]));
await flush();
expect(h.pending).toHaveLength(1);
expect(h.pending[0]?.filters[0]?.["#h"]).toEqual(["d"]);
h.next().respond(head("d"));
await flush();
expect(h.pending).toHaveLength(0);
});

it("does not automatically retry a failed optional head when another preview changes", async () => {
const h = await setup(["a", "b"]);
h.next().fail(new Error("Head unavailable"));
await vi.waitFor(() => expect(h.pending).toHaveLength(1));
expect(h.pending[0]?.filters[0]?.["#h"]).toEqual(["b"]);
h.next().respond(head("b"));
await flush();
expect(h.pending).toHaveLength(0);
h.queries.ensure("a");
expect(h.pending[0]?.filters[0]?.["#h"]).toEqual(["a"]);
h.next().respond(head("a"));
await flush();
expect(h.queries.window("a").status).toBe("ready");
});

it.each(["clear", "dispose"])(
"does not resurrect queued warming after in-flight %s",
async (boundary) => {
const h = await setup(["a", "b"]);
const first = h.next();
if (boundary === "clear") await h.store.clearCache();
else h.store.dispose();
expect(first.signal?.aborted).toBe(true);
first.respond(head("a"));
await flush();
expect(h.pending).toHaveLength(0);
expect(h.store.diagnostics().heads.entries).toBe(0);
if (boundary === "clear") {
await h.store.session.sidebarPreferences.ensure();
expect(h.pending[0]?.filters[0]?.["#h"]).toEqual(["a"]);
h.next().respond(head("a"));
await flush();
h.next().respond(head("b"));
await flush();
expect(h.pending).toHaveLength(0);
}
},
);

it.each(["a", "b"])(
"drops revoked in-flight/queued work and rewarms %s after regrant",
async (revoked) => {
const h = await setup(["a", "b", "c"]);
const first = h.next();
// Revocation cancels all reads; neither that response nor preview changes
// may restart the consumed candidate. Other queued eligible work can finish.
const removed = roster(relay, revoked, [], 1_700_000_001);
h.live.receive([removed]);
expect(first.signal?.aborted).toBe(true);
first.respond(head("a"));
await flush();
for (const id of revoked === "a" ? ["b", "c"] : ["c"]) {
expect(h.pending[0]?.filters[0]?.["#h"]).toEqual([id]);
h.next().respond(head(id));
await flush();
}
expect(h.pending).toHaveLength(0);
expect(
h.queries.list().channels.some((channel) => channel.id === revoked),
).toBe(false);
const added = roster(relay, revoked, [viewer.pubkey], 1_700_000_002);
h.live.receive([added]);
expect(h.pending[0]?.filters[0]?.["#h"]).toEqual([revoked]);
h.next().respond(head(revoked));
await flush();
expect(h.pending).toHaveLength(0);
},
);

it("warms an archived channel only after it becomes starred", async () => {
const h = await setup([]);
h.queries.refreshList?.();
h.next().respond([
roster(relay, "a", [viewer.pubkey]),
signed(relay, {
kind: 39000,
content: JSON.stringify({ name: "a", archived: true }),
tags: [
["d", "a"],
["name", "a"],
["archived", "true"],
],
}),
]);
await flush();
expect(h.pending).toHaveLength(0);
await h.setStars(["a"]);
expect(h.pending[0]?.filters[0]?.["#h"]).toEqual(["a"]);
h.next().respond(head("a"));
await flush();
expect(h.pending).toHaveLength(0);
});

it("forgets warming eligibility after authoritative roster denial, then permits regrant", async () => {
const h = await setup(["a", "b"]);
const first = h.next();
h.queries.refreshList?.();
h.next().fail(new ReadError("denied", "Roster denied"));
await vi.waitFor(() => expect(h.queries.list().status).toBe("error"));
expect(first.signal?.aborted).toBe(true);
first.respond(head("a"));
await flush();
expect(h.pending).toHaveLength(0);
expect(h.store.diagnostics().heads.entries).toBe(0);
h.live.receive([roster(relay, "a", [viewer.pubkey], 1_700_000_002)]);
await vi.waitFor(() => expect(h.pending).toHaveLength(1));
expect(h.pending[0]?.filters[0]?.["#h"]).toEqual(["a"]);
h.next().respond(head("a"));
await flush();
expect(h.pending).toHaveLength(0);
});

it("drops an archived queued channel, then warms it after unarchiving", async () => {
const h = await setup(["a", "b"]);
const first = h.next();
h.live.receive([
signed(relay, {
kind: 39000,
content: JSON.stringify({ name: "b" }),
created_at: 1_700_000_001,
tags: [
["d", "b"],
["name", "b"],
["archived", "true"],
],
}),
]);
first.respond(head("a"));
await flush();
expect(h.pending).toHaveLength(0);
h.live.receive([metadata(relay, "b", "b", 1_700_000_002)]);
expect(h.pending[0]?.filters[0]?.["#h"]).toEqual(["b"]);
h.next().respond(head("b"));
await flush();
expect(h.pending).toHaveLength(0);
});
20 changes: 20 additions & 0 deletions tests/browser/channel-opening.spec.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,16 @@ test("cold opening bypasses held DM labels; warm switching paints within 100ms w
app.relay.holdProfiles(app.participants);
app.relay.holdEose("alpha");
app.relay.holdEose("beta");
// Establish a cold target deterministically: let the actual label read own
// the background slot before preferences release the roster warmer. Hold the
// host decoder, not a reader slot; demand and all production owners stay real.
const preferences = Promise.withResolvers();
let preferencesPending = false;
await page.route("**/sidebar-preferences", async (route) => {
preferencesPending = true;
await preferences.promise;
await route.fallback();
});
// A visible pre-establishment head is not a warm, verified cache. A signed
// missed message proves the catch-up reached the UI, not merely the broker.
const establish = async (channel) => {
Expand All @@ -56,6 +66,14 @@ test("cold opening bypasses held DM labels; warm switching paints within 100ms w
await establish("alpha");
await expect.poll(() => labelReads().length).toBe(1);
expect(labelReads()[0].filter.authors).toHaveLength(500);
expect(app.report.profileHolds.some((held) => held.pending)).toBe(true);
await expect.poll(() => preferencesPending).toBe(true);
const decoded = page.waitForResponse(
(response) =>
response.url().endsWith("/sidebar-preferences") && response.ok(),
);
preferences.resolve();
await decoded;
// The actual label hook has >1,000 missing participants. Keep its profile
// response held through cold opening; do not bypass that production caller.
expect(heads(app, "beta")).toHaveLength(0);
Expand Down Expand Up @@ -158,8 +176,10 @@ test("cold opening bypasses held DM labels; warm switching paints within 100ms w
expect(app.report.profileHolds.some((held) => held.pending)).toBe(true);
expect(app.report.profileHolds.some((held) => held.aborted)).toBe(false);
} finally {
preferences.resolve();
app.relay.releaseEose("alpha");
app.relay.releaseEose("beta");
app.relay.releaseProfiles();
await page.unrouteAll({ behavior: "wait" });
}
});
Loading