flotilla/src/app/feeds.ts

830 lines
25 KiB
TypeScript
Raw Normal View History

2026-08-30 00:42:07 +00:00
import {derived, get, readable, writable} from "svelte/store"
import type {Readable} from "svelte/store"
import {batch, call, int, ms, now, on, sleep, uniqBy, MONTH, YEAR} from "@welshman/lib"
import {
COMMENT,
DELETE,
EVENT_TIME,
NOTE,
addressTags,
compareEventsAsc,
getAddress,
getCommentFiltersForRoot,
2026-08-31 17:03:01 +00:00
getIdOrAddress,
getReplyFilters,
hexTags,
isReplaceableKind,
matchFilters,
tagSpec,
tagValue,
tagValues,
} from "@welshman/util"
2026-08-30 00:42:07 +00:00
import type {Maybe} from "@welshman/lib"
import type {Filter, TrustedEvent} from "@welshman/util"
2026-07-28 16:16:26 +00:00
import {mergeRepositoryUpdates} from "@welshman/net"
import type {RepositoryUpdate} from "@welshman/net"
2025-01-28 16:13:20 +00:00
import {createScroller} from "@lib/html"
2026-08-30 00:42:07 +00:00
import type {ScrollerOpts} from "@lib/html"
2025-02-06 00:26:22 +00:00
import {daysBetween} from "@lib/util"
import {EVENT_CONTEXT_KINDS, REACTION_KINDS} from "@app/content"
2026-07-28 16:16:26 +00:00
import {app, network} from "@app/core"
import {getEventsForUrl} from "@app/repository"
2024-12-11 00:38:22 +00:00
const noEvents: TrustedEvent[] = []
const mergeSorted = <T>(left: T[], right: T[], compare: (a: T, b: T) => number) => {
const merged: T[] = []
let i = 0
let j = 0
while (i < left.length && j < right.length) {
merged.push(compare(left[i], right[j]) <= 0 ? left[i++] : right[j++])
}
2026-09-21 16:48:35 +00:00
while (i < left.length) {
merged.push(left[i++])
}
while (j < right.length) {
merged.push(right[j++])
}
return merged
}
// Reactions, zaps and reports point at their subject with `e`/`a`. A NIP-22 comment instead
// points at its thread *root* with `E`/`A`, so filing it by those tags puts a whole thread in
// the root's bucket — which is the scope a reply count wants.
const getTargets = ({kind, tags}: TrustedEvent) =>
kind === COMMENT
? [...tagValues(hexTags("E"), tags), ...tagValues(addressTags("A"), tags)]
: [...tagValues(hexTags("e"), tags), ...tagValues(addressTags("a"), tags)]
// A related event may point at a replaceable one by either id or address, so both are keys
const getKeys = (event: TrustedEvent) =>
isReplaceableKind(event.kind) ? [event.id, getAddress(event)] : [event.id]
export const makeFeedContext = ({
relays,
withReplies = false,
}: {
relays: string[] | Promise<string[]>
withReplies?: boolean
}) => {
// Every feed loads NIP-22 comments. Kind 1 notes reply to each other with `e` tags instead,
// so a feed that renders those replies has to ask for them as well.
const contextKinds = withReplies ? [...EVENT_CONTEXT_KINDS, NOTE] : EVENT_CONTEXT_KINDS
const requestKinds = withReplies ? [...REACTION_KINDS, NOTE] : REACTION_KINDS
const {repository, tracker} = app.get()
const controller = new AbortController()
const targets = new Set<string>()
const eventsByTarget = new Map<string, TrustedEvent[]>()
const targetsByEventId = new Map<string, string[]>()
const subscribersByTarget = new Map<string, Set<(events: TrustedEvent[]) => void>>()
2026-08-31 17:03:01 +00:00
const deletedChecks = new Map<string, Set<() => void>>()
const addEvent = (event: TrustedEvent, touched: Set<string>) => {
// An event seen before its target was tracked stays unfiled, so that adding the target
// later can pick it up out of the repository
if (!targetsByEventId.has(event.id)) {
const eventTargets = getTargets(event).filter(target => targets.has(target))
if (eventTargets.length > 0) {
targetsByEventId.set(event.id, eventTargets)
for (const target of eventTargets) {
eventsByTarget.set(target, [...(eventsByTarget.get(target) || []), event])
touched.add(target)
}
}
}
}
const removeEvent = (id: string, touched: Set<string>) => {
const eventTargets = targetsByEventId.get(id)
if (eventTargets) {
targetsByEventId.delete(id)
for (const target of eventTargets) {
const events = eventsByTarget.get(target)!.filter(event => event.id !== id)
if (events.length > 0) {
eventsByTarget.set(target, events)
} else {
eventsByTarget.delete(target)
}
touched.add(target)
}
}
}
2026-08-31 17:03:01 +00:00
const notifyDeleted = (added: TrustedEvent[], removed: Set<string>) => {
if (deletedChecks.size > 0) {
if (removed.size > 0) {
for (const checks of deletedChecks.values()) {
2026-09-21 16:48:35 +00:00
for (const check of checks) {
check()
}
2026-08-31 17:03:01 +00:00
}
} else {
for (const event of added) {
for (const check of deletedChecks.get(getIdOrAddress(event)) || []) {
check()
}
}
}
}
}
const notify = (touched: Set<string>) => {
for (const target of touched) {
for (const subscriber of subscribersByTarget.get(target) || []) {
subscriber(eventsByTarget.get(target) || noEvents)
}
}
}
const loadFrom = async (urls: string[], events: TrustedEvent[]) => {
2026-09-21 17:58:14 +00:00
const context = await network.get().loadComplete({
relays: urls,
signal: controller.signal,
filters: [
...getReplyFilters(events, {kinds: requestKinds}),
...getCommentFiltersForRoot(events),
],
})
if (context.length > 0) {
2026-09-21 17:58:14 +00:00
network.get().loadLenient({
relays: urls,
signal: controller.signal,
filters: getReplyFilters(context, {kinds: [DELETE]}),
})
}
}
const loadContext = batch(100, async (events: TrustedEvent[]) => {
const touched = new Set<string>()
// What's already local — an earlier page, our own optimistic reactions — never comes
// through the update listener, so file it before asking the network for the rest
for (const event of repository.query([
...getReplyFilters(events, {kinds: contextKinds}),
...getCommentFiltersForRoot(events),
])) {
addEvent(event, touched)
}
notify(touched)
const urls = await relays
// A space relay holds the whole conversation for its own feed. A feed built out of other
// people's outboxes does not, so ask each event's own relays for its context too.
const eventsByRelay = new Map<string, TrustedEvent[]>()
for (const event of events) {
for (const url of tracker.getRelays(event.id)) {
if (!urls.includes(url)) {
eventsByRelay.set(url, [...(eventsByRelay.get(url) || []), event])
}
}
}
await Promise.all([
loadFrom(urls, events),
...Array.from(eventsByRelay, ([url, seenThere]) => loadFrom([url], seenThere)),
])
})
const unsubscribe = on(
repository,
"update",
batch(150, (updates: RepositoryUpdate[]) => {
const {added, removed} = mergeRepositoryUpdates(updates)
const touched = new Set<string>()
for (const event of added) {
if (contextKinds.includes(event.kind)) {
addEvent(event, touched)
}
}
for (const id of removed) {
removeEvent(id, touched)
}
notify(touched)
2026-08-31 17:03:01 +00:00
notifyDeleted(added, removed)
}),
)
// Track an event so its context gets requested with the rest of the batch
const add = (event: TrustedEvent) => {
if (!targets.has(event.id)) {
for (const key of getKeys(event)) {
targets.add(key)
}
loadContext(event)
}
}
const relatedForKey = (key: string) =>
readable(eventsByTarget.get(key) || noEvents, set => {
let subscribers = subscribersByTarget.get(key)
if (!subscribers) {
subscribers = new Set()
subscribersByTarget.set(key, subscribers)
}
subscribers.add(set)
set(eventsByTarget.get(key) || noEvents)
return () => {
subscribers.delete(set)
if (subscribers.size === 0) {
subscribersByTarget.delete(key)
}
}
})
return {
add,
2026-08-31 17:03:01 +00:00
deleted: (event: TrustedEvent) =>
readable(repository.isDeleted(event), set => {
const key = getIdOrAddress(event)
const check = () => set(repository.isDeleted(event))
let checks = deletedChecks.get(key)
if (!checks) {
checks = new Set()
deletedChecks.set(key, checks)
}
checks.add(check)
check()
return () => {
checks.delete(check)
if (checks.size === 0) {
deletedChecks.delete(key)
}
}
}),
related: (event: TrustedEvent): Readable<TrustedEvent[]> => {
add(event)
const [key, ...rest] = getKeys(event)
// A replaceable event collects both its buckets, and something tagging it by id and
// address at once lands in both
return rest.length > 0
? derived([relatedForKey(key), ...rest.map(relatedForKey)], buckets =>
uniqBy(event => event.id, buckets.flat()),
)
: relatedForKey(key)
},
cleanup: () => {
controller.abort()
unsubscribe()
},
}
}
export type FeedContext = ReturnType<typeof makeFeedContext>
2026-08-30 00:42:07 +00:00
// Keeps a feed's store in step with the repository: events that arrive later, events that get
// deleted, and events already held that have only now been seen on one of the feed's relays.
const syncFeed = ({
2026-06-29 23:28:11 +00:00
relays,
2025-10-06 18:23:19 +00:00
filters,
addEvents,
removeEvents,
2026-08-30 00:42:07 +00:00
requireRelay = true,
2025-01-28 16:13:20 +00:00
}: {
2026-06-29 23:28:11 +00:00
relays: string[]
2025-10-06 18:23:19 +00:00
filters: Filter[]
addEvents: (events: TrustedEvent[]) => void
removeEvents: (ids: Set<string>) => void
2026-08-30 00:42:07 +00:00
// Whether an event has to have been seen on one of `relays` to belong to this feed
requireRelay?: boolean
2025-01-28 16:13:20 +00:00
}) => {
2026-07-28 00:13:14 +00:00
const onTrackedId = batch(150, (ids: string[]) => {
const matching: TrustedEvent[] = []
for (const id of new Set(ids)) {
2026-07-28 16:16:26 +00:00
const event = app.get().repository.getEvent(id)
2026-07-28 00:13:14 +00:00
if (event && matchFilters(filters, event)) {
matching.push(event)
}
}
if (matching.length > 0) {
addEvents(matching)
2026-07-28 00:13:14 +00:00
}
})
2026-08-30 00:42:07 +00:00
return [
2026-04-10 18:09:26 +00:00
on(
2026-07-28 16:16:26 +00:00
app.get().repository,
2026-04-10 18:09:26 +00:00
"update",
batch(150, (updates: RepositoryUpdate[]) => {
const {added, removed} = mergeRepositoryUpdates(updates)
if (removed.size > 0) {
removeEvents(removed)
2026-04-10 18:09:26 +00:00
}
2026-02-13 22:51:56 +00:00
2026-04-10 18:09:26 +00:00
const matching = added.filter(
2026-06-29 23:28:11 +00:00
event =>
matchFilters(filters, event) &&
2026-08-30 00:42:07 +00:00
(!requireRelay || relays.some(url => app.get().tracker.getRelays(event.id).has(url))),
2026-04-10 18:09:26 +00:00
)
2026-04-10 18:09:26 +00:00
if (matching.length > 0) {
addEvents(matching)
2026-04-10 18:09:26 +00:00
}
}),
),
2026-07-28 16:16:26 +00:00
on(app.get().tracker, "add", (id: string, url: string) => {
2026-06-29 23:28:11 +00:00
if (relays.includes(url)) {
2026-07-28 00:13:14 +00:00
onTrackedId(id)
2026-02-13 23:18:46 +00:00
}
}),
]
2026-08-30 00:42:07 +00:00
}
// One direction of a feed. A span that comes back empty is a gap in the timeline, not the end of
// it — conflating the two either stops loading at the first gap or walks the whole history
// looking for the end of one.
export type FeedLoadState =
| {status: "idle"}
| {status: "loading"}
| {status: "searching"}
| {status: "exhausted"}
2026-09-01 17:29:14 +00:00
// One span of a feed's timeline. `complete` is whether the relays actually answered for it: a
// request the socket dropped reports nothing found, which is not the same thing as a span with
// nothing in it, and the two have to move the window differently.
export type FeedSpan = {found: number; complete: boolean; exhausted: boolean}
// How far a span covered is set by its least generous relay, and quiet relays answer first.
const SPAN_THRESHOLD = 0.8
// A relay that accepts a socket and then says nothing neither answers nor drops.
const SPAN_TIMEOUT = 3000
// Aborting a request resolves it with what arrived, the same way closing it does.
const spanSignal = (signal: AbortSignal) =>
AbortSignal.any([signal, AbortSignal.timeout(SPAN_TIMEOUT)])
2026-08-30 00:42:07 +00:00
// Empty spans to walk per trigger. Enough to cross a gap; not enough to reach the end of the
// history on a single request.
const SPANS_PER_TRIGGER = 3
// How many events one step back through the history asks for. Spans are a month wide so that a
// quiet feed finds its history in a few requests, which in a busy one is thousands of events —
// far more than anyone is about to read, and all of it queued ahead of what they are looking at.
// Asking for a page instead leaves the relays to answer with the events nearest the anchor,
// which is what NIP-01 promises a filter carrying a limit.
const PAGE_SIZE = 100
// A span that turns up nothing widens the next one, so walking a sparse history doesn't take
// dozens of round trips
const nextInterval = (interval: number, found: number) =>
found > 0 ? int(MONTH) : Math.round(interval * 1.5)
2026-08-30 00:42:07 +00:00
// Whether a request is actually in flight, which is not the same as whether more might exist. A
// list that already has what it needs shouldn't sit under a spinner just because it hasn't
// walked to the end of the history.
export const isFeedLoading = (state: Maybe<FeedLoadState>) =>
state?.status === "loading" || state?.status === "searching"
2026-09-01 17:29:14 +00:00
const makeFeedLoader = (load: () => Promise<FeedSpan>) => {
2026-08-30 00:42:07 +00:00
const state = writable<FeedLoadState>({status: "idle"})
let running = false
const run = async () => {
2026-09-21 16:48:35 +00:00
if (running || get(state).status === "exhausted") {
return
}
2026-08-30 00:42:07 +00:00
running = true
try {
for (let span = 0; span < SPANS_PER_TRIGGER; span++) {
state.set({status: span > 0 ? "searching" : "loading"})
2026-09-01 17:29:14 +00:00
const {found, complete, exhausted} = await load()
2026-08-30 00:42:07 +00:00
if (exhausted) {
state.set({status: "exhausted"})
return
}
2026-09-01 17:29:14 +00:00
// A span nobody answered for hasn't moved the window, so hold here rather than walking
// past it, and give the socket a moment before the scroller comes back around
if (!complete) {
await sleep(ms(3))
break
2026-08-30 00:42:07 +00:00
}
2026-09-01 17:29:14 +00:00
2026-09-21 16:48:35 +00:00
if (found > 0) {
break
}
2026-08-30 00:42:07 +00:00
}
} finally {
running = false
}
}
// A run covers a few spans and the scroller starts another one a moment later, so settling at
// the end of a run blinks the spinner once per page. What ends a load is the trigger going
// quiet — the list grown long enough that nothing more is wanted.
const settle = () => {
if (!running && isFeedLoading(get(state))) {
state.set({status: "idle"})
}
}
return {subscribe: state.subscribe, run, settle}
2026-08-30 00:42:07 +00:00
}
// A loader triggered by proximity to the end of a scroll container, which is how every list in
// the app pages. The container's orientation decides which direction `reverse` reaches, so the
// caller passes it — a reversed chat scrolls away from its origin to find older messages, an
// ordinary feed scrolls toward the end of its own content.
export const makeScrollLoader = (
element: HTMLElement,
2026-09-01 17:29:14 +00:00
load: () => Promise<FeedSpan>,
2026-08-30 00:42:07 +00:00
options: Partial<ScrollerOpts> = {},
) => {
const loader = makeFeedLoader(load)
const scroller = createScroller({
element,
delay: 300,
threshold: 5000,
...options,
onScroll: loader.run,
onSettle: loader.settle,
2026-08-30 00:42:07 +00:00
})
return {subscribe: loader.subscribe, stop: scroller.stop}
}
// Holds every event a view has loaded, sorted oldest to newest, and knows how to ask for the
// next span in either direction. It does not decide *when* to ask — the view does, because the
// view is what knows what is on screen.
export const makeFeed = ({
relays,
filters,
onEvent,
at = now(),
}: {
relays: string[]
filters: Filter[]
onEvent?: (event: TrustedEvent) => void
at?: number
}) => {
const controller = new AbortController()
const events = writable<TrustedEvent[]>([])
const seen = new Set<string>()
// Events from further back than the feed has reached. Everything the app does fills the same
// repository — a space-wide sync reconciling a month of every room at once, most of all — and
// putting each of those on screen as it lands is what walks a room backwards under the reader.
const held = new Map<string, TrustedEvent>()
// The span the relays have been asked about, which grows outward from the anchor. Each
// direction widens its own: what one of them is walking through says nothing about the other,
// and sharing a width lets a backward walk finding history reset the forward walk's growth.
2026-08-30 00:42:07 +00:00
let oldest = at
let newest = at
let olderInterval = int(MONTH)
let newerInterval = int(MONTH)
2026-08-30 00:42:07 +00:00
// How far back the feed has been answered for, which is what decides whether an event is ready
// to render or has to wait for the window to come and get it.
let reached = at
2026-08-30 00:42:07 +00:00
const insertEvents = (newEvents: Iterable<TrustedEvent>) => {
const added: TrustedEvent[] = []
for (const event of newEvents) {
if (!seen.has(event.id)) {
seen.add(event.id)
held.delete(event.id)
2026-08-30 00:42:07 +00:00
added.push(event)
onEvent?.(event)
}
}
if (added.length > 0) {
added.sort(compareEventsAsc)
2026-08-30 00:42:07 +00:00
events.update($events => mergeSorted($events, added, compareEventsAsc))
2026-08-30 00:42:07 +00:00
}
}
// The one door into the feed, and anything older than it has reached waits for the window.
const addEvents = (newEvents: TrustedEvent[]) => {
const ready: TrustedEvent[] = []
2025-02-04 22:06:05 +00:00
for (const event of newEvents) {
if (!seen.has(event.id) && !held.has(event.id)) {
if (event.created_at >= reached) {
ready.push(event)
} else {
held.set(event.id, event)
}
}
}
insertEvents(ready)
}
const removeEvents = (ids: Set<string>) => {
events.update($events => $events.filter(event => !ids.has(event.id)))
for (const id of ids) {
seen.delete(id)
held.delete(id)
}
}
// Take the feed back to `timestamp`, releasing everything that had arrived for the stretch
// between there and where it had got to
const reach = (timestamp: number) => {
if (timestamp < reached) {
reached = timestamp
const ready: TrustedEvent[] = []
for (const [id, event] of held) {
if (event.created_at >= reached) {
held.delete(id)
ready.push(event)
}
}
insertEvents(ready)
}
}
const unsubscribers = syncFeed({relays, filters, addEvents, removeEvents})
// One request per direction, reported per relay as well as in total: each relay answers a
// limit for itself, so a page only runs out where the relay that gave the least of it ran out.
const loadSpan = async (extension: Filter) => {
2026-09-01 17:29:14 +00:00
let complete = false
const pages = new Map<string, {count: number; lowest: number}>()
const countEvent = (event: TrustedEvent, url: string) => {
const page = pages.get(url)
if (page) {
page.count += 1
page.lowest = Math.min(page.lowest, event.created_at)
} else {
pages.set(url, {count: 1, lowest: event.created_at})
}
}
2026-08-30 00:42:07 +00:00
const found = await network.get().request({
2026-06-29 23:28:11 +00:00
relays,
2026-02-16 21:31:43 +00:00
autoClose: true,
signal: spanSignal(controller.signal),
threshold: SPAN_THRESHOLD,
filters: filters.map(filter => ({...filter, ...extension})),
onEvent: countEvent,
onDuplicate: countEvent,
2026-09-01 17:29:14 +00:00
onEose: () => {
complete = true
},
2026-02-16 21:31:43 +00:00
})
return {found, complete, pages}
2026-02-16 21:31:43 +00:00
}
2026-08-30 00:42:07 +00:00
// Ask for the next span in each direction. A span that comes back empty is normal while
// walking a sparse history, so the count is reported separately from whether there is any
// history left — a caller watching its list for changes would never hear about an empty one.
2026-09-01 17:29:14 +00:00
// The window only moves once the relays have answered: a request the socket dropped looks
// exactly like an empty span, and walking past it would leave a hole nothing goes back for.
const loadOlder = async (): Promise<FeedSpan> => {
2026-09-21 16:48:35 +00:00
if (oldest < now() - int(2, YEAR)) {
return {found: 0, complete: true, exhausted: true}
}
2026-02-16 21:31:43 +00:00
2026-08-30 00:42:07 +00:00
const until = oldest
const since = until - olderInterval
const {found, complete, pages} = await loadSpan({since, until, limit: PAGE_SIZE})
// A relay that answered with less than it was allowed has covered its whole span and holds
// nothing back, so the page runs out at the highest of the rest
let edge: Maybe<number>
for (const page of pages.values()) {
if (page.count >= PAGE_SIZE && (edge === undefined || page.lowest > edge)) {
edge = page.lowest
}
}
2026-02-16 21:31:43 +00:00
2026-09-01 17:29:14 +00:00
if (complete) {
olderInterval = nextInterval(olderInterval, found.length)
// The second the edge steps back is what stops a page that filled up inside one from being
// asked for over and over
oldest = edge === undefined ? since : Math.min(edge, until - 1)
2026-09-01 17:29:14 +00:00
}
2026-02-16 21:31:43 +00:00
// A span reaches further back than it covers, so its deepest events wait for the window.
addEvents(found)
reach(oldest)
return {found: found.length, complete, exhausted: false}
2026-08-30 00:42:07 +00:00
}
2025-01-28 16:13:20 +00:00
// A limit is answered with the newest events matching it, which reaches away from an anchor in
// the past rather than toward it, so this direction walks spans as it always has
2026-09-01 17:29:14 +00:00
const loadNewer = async (): Promise<FeedSpan> => {
2026-09-21 16:48:35 +00:00
if (newest >= now()) {
return {found: 0, complete: true, exhausted: true}
}
2026-02-16 21:31:43 +00:00
2026-08-30 00:42:07 +00:00
const since = newest
const until = Math.min(now(), since + newerInterval)
const {found, complete} = await loadSpan({since, until})
2025-01-28 16:13:20 +00:00
2026-09-01 17:29:14 +00:00
if (complete) {
newerInterval = nextInterval(newerInterval, found.length)
2026-09-01 17:29:14 +00:00
newest = until
}
2026-02-13 22:51:56 +00:00
addEvents(found)
return {found: found.length, complete, exhausted: false}
2026-08-30 00:42:07 +00:00
}
2025-01-28 16:13:20 +00:00
addEvents(relays.flatMap(url => Array.from(getEventsForUrl(url, filters))))
2025-01-28 16:13:20 +00:00
return {
events,
2026-08-30 00:42:07 +00:00
loadOlder,
loadNewer,
2025-01-28 16:13:20 +00:00
cleanup: () => {
2025-04-11 16:27:19 +00:00
controller.abort()
2026-02-13 23:18:46 +00:00
unsubscribers.forEach(call)
2025-01-28 16:13:20 +00:00
},
}
}
2026-08-30 00:42:07 +00:00
// Same split as makeFeed: it holds what has been loaded and knows how to reach further out in
// either direction, while the page decides when to ask. Calendar events are addressed by the
// days they cover rather than by when they were published, so the spans are date hashes.
2025-02-06 00:26:22 +00:00
export const makeCalendarFeed = ({
2026-06-29 23:28:11 +00:00
relays,
2025-10-06 18:23:19 +00:00
filters,
onEvent,
2025-02-06 00:26:22 +00:00
}: {
2026-06-29 23:28:11 +00:00
relays: string[]
2025-10-06 18:23:19 +00:00
filters: Filter[]
onEvent?: (event: TrustedEvent) => void
2025-02-06 00:26:22 +00:00
}) => {
2026-06-15 18:45:36 +00:00
const interval = int(5, MONTH)
2025-04-11 16:27:19 +00:00
const controller = new AbortController()
2026-06-29 23:28:11 +00:00
const seen = new Set<string>()
2026-08-30 00:42:07 +00:00
let oldest = now()
let newest = now()
2025-02-06 00:26:22 +00:00
2026-07-28 16:16:26 +00:00
const getStart = (event: TrustedEvent) => parseInt(tagValue(tagSpec("start"), event.tags) || "")
2025-02-06 01:05:41 +00:00
2026-07-28 16:16:26 +00:00
const getEnd = (event: TrustedEvent) => parseInt(tagValue(tagSpec("end"), event.tags) || "")
2025-02-06 01:05:41 +00:00
const compareByStart = (a: TrustedEvent, b: TrustedEvent) =>
getStart(a) - getStart(b) || compareEventsAsc(a, b)
2026-06-29 23:28:11 +00:00
const events = writable(
uniqBy(
e => e.id,
relays.flatMap(url => Array.from(getEventsForUrl(url, filters))),
).sort(compareByStart),
2026-06-29 23:28:11 +00:00
)
2025-02-06 01:05:41 +00:00
const insertEvents = (newEvents: TrustedEvent[]) => {
2026-06-29 23:28:11 +00:00
const valid = newEvents.filter(e => !isNaN(getStart(e)) && !isNaN(getEnd(e)) && !seen.has(e.id))
2026-09-21 16:48:35 +00:00
if (valid.length === 0) {
return
}
2025-02-06 01:05:41 +00:00
2026-07-28 00:13:14 +00:00
for (const event of valid) {
seen.add(event.id)
onEvent?.(event)
2026-07-28 00:13:14 +00:00
}
valid.sort(compareByStart)
2025-02-06 00:26:22 +00:00
2026-07-28 00:13:14 +00:00
events.update($events => {
// Calendar events are addressable, so a new version supersedes the old one
const superseded = new Set(valid.map(getAddress))
return mergeSorted(
$events.filter(e => !superseded.has(getAddress(e))),
valid,
compareByStart,
)
2025-02-06 00:26:22 +00:00
})
}
const removeEvents = (ids: Set<string>) => {
events.update($events => $events.filter(event => !ids.has(event.id)))
for (const id of ids) {
seen.delete(id)
}
}
2026-08-30 00:42:07 +00:00
// Calendar events are addressable and often relayed on from elsewhere, so this feed takes any
// matching event rather than only those seen on its own relays
const unsubscribers = syncFeed({
relays,
filters,
addEvents: insertEvents,
removeEvents,
2026-08-30 00:42:07 +00:00
requireRelay: false,
2026-07-28 00:13:14 +00:00
})
2026-08-30 00:42:07 +00:00
const loadTimeframe = async (since: number, until: number) => {
2026-09-01 17:29:14 +00:00
let complete = false
2026-08-30 00:42:07 +00:00
const found = await network.get().request({
2026-06-29 23:28:11 +00:00
relays,
2025-04-09 22:32:18 +00:00
autoClose: true,
signal: spanSignal(controller.signal),
threshold: SPAN_THRESHOLD,
2026-08-30 00:42:07 +00:00
filters: [{kinds: [EVENT_TIME], "#D": daysBetween(since, until).map(String)}],
2026-09-01 17:29:14 +00:00
onEose: () => {
complete = true
},
2025-02-06 00:26:22 +00:00
})
2026-09-01 17:29:14 +00:00
return {found: found.length, complete}
2025-02-06 01:05:41 +00:00
}
2026-09-01 17:29:14 +00:00
const loadOlder = async (): Promise<FeedSpan> => {
2026-09-21 16:48:35 +00:00
if (oldest < now() - int(2, YEAR)) {
return {found: 0, complete: true, exhausted: true}
}
2025-02-06 00:26:22 +00:00
2026-08-30 00:42:07 +00:00
const until = oldest
2026-09-01 17:29:14 +00:00
const since = until - interval
const {found, complete} = await loadTimeframe(since, until)
2025-02-06 00:26:22 +00:00
2026-09-01 17:29:14 +00:00
if (complete) {
oldest = since
}
2025-02-06 00:26:22 +00:00
2026-09-01 17:29:14 +00:00
return {found, complete, exhausted: false}
2026-08-30 00:42:07 +00:00
}
2025-02-06 00:26:22 +00:00
2026-09-01 17:29:14 +00:00
const loadNewer = async (): Promise<FeedSpan> => {
2026-09-21 16:48:35 +00:00
if (newest > now() + int(2, YEAR)) {
return {found: 0, complete: true, exhausted: true}
}
2025-02-06 00:26:22 +00:00
2026-08-30 00:42:07 +00:00
const since = newest
2026-09-01 17:29:14 +00:00
const until = since + interval
const {found, complete} = await loadTimeframe(since, until)
2026-08-30 00:42:07 +00:00
2026-09-01 17:29:14 +00:00
if (complete) {
newest = until
}
2026-08-30 00:42:07 +00:00
2026-09-01 17:29:14 +00:00
return {found, complete, exhausted: false}
2026-08-30 00:42:07 +00:00
}
2025-02-06 00:26:22 +00:00
return {
events,
2026-08-30 00:42:07 +00:00
loadOlder,
loadNewer,
// The month and week views jump to arbitrary ranges rather than scrolling through them, and
// wait on the request so they can show progress for the range on screen
load: loadTimeframe,
2025-02-06 00:26:22 +00:00
cleanup: () => {
2025-04-11 16:27:19 +00:00
controller.abort()
2026-02-13 23:18:46 +00:00
unsubscribers.forEach(call)
2025-02-06 00:26:22 +00:00
},
}
}