Extract relay-agnostic forum fetchers from the stores

This commit is contained in:
dtonon 2026-08-14 22:27:04 +01:00
parent 54c0682dce
commit a7b5d8f223
16 changed files with 694 additions and 531 deletions

View file

@ -1,27 +1,14 @@
import type { AbstractRelay } from "@nostr/tools/abstract-relay";
import type { Event } from "@nostr/tools/core";
import type { Filter } from "@nostr/tools/filter";
import { loadNostrUser, type NostrUser } from "$lib/gadgets";
import { ensureForumRelay } from "$lib/relay";
import { ingestNostrUser } from "$lib/profiles.svelte";
import { relayQuery } from "$lib/forum/query";
import {
fetchThreadPage,
type SortMode,
type ThreadData,
} from "$lib/forum/threads";
const PAGE_SIZE = 30;
const WALK_LIMIT = 150; // Events per walk query (~5x PAGE_SIZE)
const WALK_MAX_ITERS = 12;
export type SortMode = "active" | "new";
export type ThreadData = {
id: string;
title: string;
labels: string[];
authorPubkey: string;
createdAt: number;
replyCount: number;
latestAt: number;
latestPubkey: string;
replierPubkeys: string[]; // unique reply authors, excl. OP, max 4
};
export type { SortMode, ThreadData };
let threads = $state<ThreadData[]>([]);
let profiles = $state<Record<string, NostrUser>>({});
@ -59,238 +46,29 @@ async function loadProfile(pubkey: string) {
ingestNostrUser(user);
}
function querySync(relay: AbstractRelay, filter: Filter): Promise<Event[]> {
return new Promise((resolve) => {
const events: Event[] = [];
const sub = relay.subscribe([filter], {
onevent(e) {
events.push(e);
},
oneose() {
sub.close();
resolve(events);
},
onclose() {
resolve(events);
},
});
});
}
function threadIdOf(e: Event): string | undefined {
return e.kind === 11 ? e.id : e.tags.find((t) => t[0] === "E")?.[1];
}
type SliceItem = {
id: string;
latestAt: number;
latestPubkey: string;
op?: Event;
};
// Walk the combined OP + reply stream newest-first, collecting up to `n` unique
// threads not already shown. The first event seen for a thread defines its
// activity timestamp. Returns the slice plus the cursor for the next call.
async function fetchActivitySlice(
relay: AbstractRelay,
groupId: string,
until: number,
exclude: Set<string>,
n: number,
): Promise<{ items: SliceItem[]; nextCursor: number | null; done: boolean }> {
const collected = new Map<string, SliceItem>();
let cur = until;
let done = false;
for (let i = 0; i < WALK_MAX_ITERS && collected.size < n; i++) {
const events = await querySync(relay, {
kinds: [11, 1111],
"#h": [groupId],
until: cur,
limit: WALK_LIMIT,
});
if (events.length === 0) {
done = true;
break;
}
events.sort((a, b) => b.created_at - a.created_at);
let oldest = cur;
for (const e of events) {
oldest = Math.min(oldest, e.created_at);
const id = threadIdOf(e);
if (!id || exclude.has(id)) continue;
const seen = collected.get(id);
if (seen) {
// Grab the OP from the stream when a thread was first seen via a reply,
// so we avoid the by-id backfill the group relay truncates
if (!seen.op && e.kind === 11) seen.op = e;
continue;
}
if (collected.size >= n) continue; // Page full; keep scanning for OPs
collected.set(id, {
id,
latestAt: e.created_at,
latestPubkey: e.pubkey,
op: e.kind === 11 ? e : undefined,
});
}
if (events.length < WALK_LIMIT) {
done = true;
break;
}
if (oldest >= cur) break; // No progress (single timestamp floods the window)
cur = oldest; // Inclusive; thread-level dedupe absorbs re-reads
}
const items = [...collected.values()].sort((a, b) => b.latestAt - a.latestAt);
const nextCursor = items.length > 0 ? items[items.length - 1].latestAt : null;
return { items, nextCursor, done };
}
// Chronological-by-creation slice: just OPs ordered by created_at. No reply
// data is needed to order them (stable cursor), keeping the path cheap; reply
// counts are still attached later via enrichment.
async function fetchNewSlice(
relay: AbstractRelay,
groupId: string,
until: number,
exclude: Set<string>,
n: number,
): Promise<{ items: SliceItem[]; nextCursor: number | null; done: boolean }> {
const ops = await querySync(relay, {
kinds: [11],
"#h": [groupId],
until,
limit: n + 10, // headroom for boundary OPs re-read at the inclusive cursor
});
ops.sort((a, b) => b.created_at - a.created_at);
const fresh = ops.filter((e) => !exclude.has(e.id));
const slice = fresh.slice(0, n);
const items: SliceItem[] = slice.map((op) => ({
id: op.id,
latestAt: op.created_at,
latestPubkey: op.pubkey,
op,
}));
const done = ops.length < n + 10;
const last = slice[slice.length - 1] ?? ops[ops.length - 1];
const nextCursor = last ? last.created_at : null;
return { items, nextCursor, done };
}
// Reply enrichment (exact counts + sampled repliers), bounded by the frozen
// snapshot. Isolated so the future creation-by-date view can skip it entirely.
async function enrichWithReplies(
relay: AbstractRelay,
groupId: string,
items: SliceItem[],
): Promise<Map<string, { count: number; repliers: string[] }>> {
const result = new Map<string, { count: number; repliers: string[] }>();
if (items.length === 0) return result;
const replies = await querySync(relay, {
kinds: [1111],
"#h": [groupId],
"#E": items.map((it) => it.id),
until: snapshotAt,
limit: 5000,
});
const acc = new Map<string, { count: number; pubkeys: Set<string> }>();
for (const r of replies) {
const root = r.tags.find((t) => t[0] === "E")?.[1];
if (!root) continue;
let a = acc.get(root);
if (!a) {
a = { count: 0, pubkeys: new Set() };
acc.set(root, a);
}
a.count++;
a.pubkeys.add(r.pubkey);
}
for (const it of items) {
const a = acc.get(it.id);
const repliers = a
? [...a.pubkeys].filter((p) => p !== it.op?.pubkey).slice(0, 4)
: [];
result.set(it.id, { count: a?.count ?? 0, repliers });
}
return result;
}
async function buildThreads(
relay: AbstractRelay,
groupId: string,
items: SliceItem[],
): Promise<ThreadData[]> {
// Backfill OPs that fell outside the activity window. The group relay
// truncates multi-id queries, so fetch each one on its own.
const missing = items.filter((it) => !it.op);
if (missing.length > 0) {
const fetched = await Promise.all(
missing.map((it) =>
querySync(relay, { kinds: [11], "#h": [groupId], ids: [it.id] }),
),
);
const byId = new Map<string, Event>();
for (const evs of fetched) for (const e of evs) byId.set(e.id, e);
for (const it of items) if (!it.op) it.op = byId.get(it.id);
}
const enriched = await enrichWithReplies(relay, groupId, items);
const out: ThreadData[] = [];
for (const it of items) {
const op = it.op;
if (!op) continue; // OP missing (deleted/unavailable) — drop the row
const e = enriched.get(it.id);
out.push({
id: it.id,
title: op.tags.find((t) => t[0] === "title")?.[1] ?? "(untitled)",
labels: op.tags.filter((t) => t[0] === "t" && t[1]).map((t) => t[1]),
authorPubkey: op.pubkey,
createdAt: op.created_at,
replyCount: e?.count ?? 0,
latestAt: it.latestAt,
latestPubkey: it.latestPubkey,
replierPubkeys: e?.repliers ?? [],
});
}
return out;
}
async function runLoad(append: boolean, groupId: string) {
const id = ++reqId;
if (append) loadingMore = true;
else loading = true;
const relay = await ensureForumRelay();
// Paged loads drive their own subscriptions on the shared connection
const q = relayQuery(await ensureForumRelay());
try {
const until = append ? (cursor ?? snapshotAt) : snapshotAt;
const exclude = new Set(threads.map((t) => t.id));
const slice =
sortMode === "new"
? await fetchNewSlice(relay, groupId, until, exclude, PAGE_SIZE)
: await fetchActivitySlice(relay, groupId, until, exclude, PAGE_SIZE);
const built = await buildThreads(relay, groupId, slice.items);
const page = await fetchThreadPage(q, groupId, {
sort: sortMode,
until: append ? (cursor ?? snapshotAt) : snapshotAt,
snapshotAt,
exclude: new Set(threads.map((t) => t.id)),
});
if (id !== reqId) return; // Superseded by a newer load — discard results
threads = append ? [...threads, ...built] : built;
cursor = slice.nextCursor ?? cursor;
exhausted = slice.done || built.length === 0;
for (const t of built) {
threads = append ? [...threads, ...page.threads] : page.threads;
cursor = page.nextCursor ?? cursor;
exhausted = page.done || page.threads.length === 0;
for (const t of page.threads) {
loadProfile(t.authorPubkey);
loadProfile(t.latestPubkey);
for (const p of t.replierPubkeys) loadProfile(p);
}
} finally {
// shared forum connection is long-lived — don't close it here
if (id === reqId) {
if (append) loadingMore = false;
else loading = false;