flotilla/src/app/feeds.ts

598 lines
15 KiB
TypeScript
Raw Normal View History

import {derived, readable, writable} from "svelte/store"
import type {Readable} from "svelte/store"
2026-07-28 00:13:14 +00:00
import {batch, between, call, int, now, on, sortBy, uniqBy, MONTH, YEAR} from "@welshman/lib"
import {
COMMENT,
DELETE,
EVENT_TIME,
addressTags,
getAddress,
getCommentFiltersForRoot,
getReplyFilters,
hexTags,
isReplaceableKind,
matchFilters,
tagSpec,
tagValue,
tagValues,
} from "@welshman/util"
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"
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[] = []
// 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}: {relays: string[] | Promise<string[]>}) => {
const {repository} = 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>>()
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)
}
}
}
const notify = (touched: Set<string>) => {
for (const target of touched) {
for (const subscriber of subscribersByTarget.get(target) || []) {
subscriber(eventsByTarget.get(target) || noEvents)
}
}
}
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: EVENT_CONTEXT_KINDS}),
...getCommentFiltersForRoot(events),
])) {
addEvent(event, touched)
}
notify(touched)
const urls = await relays
const context = await network.get().load({
relays: urls,
signal: controller.signal,
filters: [
...getReplyFilters(events, {kinds: REACTION_KINDS}),
...getCommentFiltersForRoot(events),
],
})
if (context.length > 0) {
network.get().load({
relays: urls,
signal: controller.signal,
filters: getReplyFilters(context, {kinds: [DELETE]}),
})
}
})
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 (EVENT_CONTEXT_KINDS.includes(event.kind)) {
addEvent(event, touched)
}
}
for (const id of removed) {
removeEvent(id, touched)
}
notify(touched)
}),
)
// 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,
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>
2025-01-28 16:13:20 +00:00
export const makeFeed = ({
2026-06-29 23:28:11 +00:00
relays,
2025-10-06 18:23:19 +00:00
filters,
2025-01-28 16:13:20 +00:00
element,
onEvent,
2026-02-16 21:31:43 +00:00
onBackwardExhausted,
onForwardExhausted,
at = now(),
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[]
2025-01-28 16:13:20 +00:00
element: HTMLElement
onEvent?: (event: TrustedEvent) => void
2026-02-16 21:31:43 +00:00
onBackwardExhausted?: () => void
onForwardExhausted?: () => void
at?: number
2025-01-28 16:13:20 +00:00
}) => {
2025-04-11 16:27:19 +00:00
const controller = new AbortController()
2026-02-13 22:51:56 +00:00
const events = writable<TrustedEvent[]>([])
2026-06-29 23:28:11 +00:00
const seen = new Set<string>()
let interval = int(6, MONTH)
let buffer: TrustedEvent[] = []
2026-02-16 21:31:43 +00:00
let backwardWindow = [at - interval, at]
let forwardWindow = [at, at + interval]
const insertIntoBuffer = (event: TrustedEvent) => {
for (let i = 0; i < buffer.length; i++) {
2026-04-10 19:40:28 +00:00
if (buffer[i].created_at < event.created_at) {
buffer.splice(i, 0, event)
return
}
}
buffer.push(event)
}
2026-02-13 22:51:56 +00:00
// Batch-insert events into the visible store with a single update
const insertEvents = (newEvents: Iterable<TrustedEvent>) => {
const visible: TrustedEvent[] = []
2026-02-13 22:51:56 +00:00
for (const event of newEvents) {
2026-06-29 23:28:11 +00:00
if (seen.has(event.id)) {
continue
}
seen.add(event.id)
if (between([backwardWindow[0], forwardWindow[1]], event.created_at)) {
visible.push(event)
} else {
insertIntoBuffer(event)
2026-02-13 22:51:56 +00:00
}
}
2026-02-13 22:51:56 +00:00
if (visible.length > 0) {
2026-07-28 00:13:14 +00:00
visible.sort((a, b) => a.created_at - b.created_at)
for (const event of visible) {
onEvent?.(event)
}
events.update($events => {
2026-07-28 00:13:14 +00:00
const merged: TrustedEvent[] = []
let i = 0
let j = 0
while (i < $events.length && j < visible.length) {
if ($events[i].created_at <= visible[j].created_at) {
merged.push($events[i++])
} else {
merged.push(visible[j++])
}
2026-02-13 22:51:56 +00:00
}
2026-06-29 23:28:11 +00:00
2026-07-28 00:13:14 +00:00
while (i < $events.length) merged.push($events[i++])
while (j < visible.length) merged.push(visible[j++])
return merged
})
2026-02-13 22:51:56 +00:00
}
}
// Buffered events are routed through insertEvents again, so forget we've seen
// them to let the window check run a second time
const drainBuffer = () => {
const drained = buffer.splice(0, 30)
for (const event of drained) {
seen.delete(event.id)
}
insertEvents(drained)
}
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) {
insertEvents(matching)
}
})
2026-02-13 23:18:46 +00:00
const unsubscribers = [
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) {
buffer = buffer.filter(e => !removed.has(e.id))
events.update($events => $events.filter(e => !removed.has(e.id)))
for (const id of removed) {
seen.delete(id)
}
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-07-28 16:16:26 +00:00
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) {
insertEvents(matching)
}
}),
),
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
}
}),
]
2025-02-04 22:06:05 +00:00
const loadTimeframe = async (since: number, until: number) => {
2026-07-28 16:16:26 +00:00
const events = await network.get().request({
2026-06-29 23:28:11 +00:00
relays,
2026-02-16 21:31:43 +00:00
autoClose: true,
signal: controller.signal,
filters: filters.map(filter => ({...filter, since, until})),
})
// If we found nothing, accelerate
if (events.length === 0) {
interval = Math.round(interval * 1.1)
} else {
2026-06-15 18:45:36 +00:00
interval = int(MONTH)
}
2026-02-16 21:31:43 +00:00
}
const backwardScroller = createScroller({
element,
delay: 300,
threshold: 5000,
2026-04-10 19:40:28 +00:00
onScroll: async () => {
2026-02-16 21:31:43 +00:00
const [since, until] = backwardWindow
backwardWindow = [since - interval, since]
drainBuffer()
2026-02-16 21:31:43 +00:00
if (until > now() - int(2, YEAR)) {
2026-04-10 19:40:28 +00:00
await loadTimeframe(since, until)
2026-02-16 21:31:43 +00:00
} else if (!buffer.some(e => e.created_at < at)) {
backwardScroller.stop()
onBackwardExhausted?.()
}
},
2025-01-28 16:13:20 +00:00
})
2026-02-16 21:31:43 +00:00
const forwardScroller = createScroller({
2025-01-28 16:13:20 +00:00
element,
2026-02-16 21:31:43 +00:00
reverse: true,
2025-01-28 16:13:20 +00:00
delay: 300,
2026-02-16 21:31:43 +00:00
threshold: 5000,
2026-04-10 19:40:28 +00:00
onScroll: async () => {
2026-02-16 21:31:43 +00:00
const [since, until] = forwardWindow
forwardWindow = [until, until + interval]
2025-01-28 16:13:20 +00:00
drainBuffer()
2026-02-13 22:51:56 +00:00
2026-02-16 21:31:43 +00:00
if (until < now()) {
2026-04-10 19:40:28 +00:00
await loadTimeframe(since, until)
2026-02-16 21:31:43 +00:00
} else if (!buffer.some(e => e.created_at > at)) {
forwardScroller.stop()
onForwardExhausted?.()
2026-02-13 22:51:56 +00:00
}
2025-01-28 16:13:20 +00:00
},
})
for (const url of relays) {
insertEvents(getEventsForUrl(url, filters))
}
2025-01-28 16:13:20 +00:00
return {
events,
cleanup: () => {
2025-04-11 16:27:19 +00:00
controller.abort()
2026-02-16 21:31:43 +00:00
forwardScroller.stop()
backwardScroller.stop()
2026-02-13 23:18:46 +00:00
unsubscribers.forEach(call)
2025-01-28 16:13:20 +00:00
},
}
}
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,
2025-02-06 00:26:22 +00:00
element,
onEvent,
2025-02-06 00:26:22 +00:00
onExhausted,
}: {
2026-06-29 23:28:11 +00:00
relays: string[]
2025-10-06 18:23:19 +00:00
filters: Filter[]
2025-02-06 00:26:22 +00:00
element: HTMLElement
onEvent?: (event: TrustedEvent) => void
2025-02-06 00:26:22 +00:00
onExhausted?: () => void
}) => {
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>()
2025-02-06 01:05:41 +00:00
let exhaustedScrollers = 0
let backwardWindow = [now() - interval, now()]
let forwardWindow = [now(), now() + interval]
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
2026-06-29 23:28:11 +00:00
const events = writable(
sortBy(
getStart,
uniqBy(
e => e.id,
relays.flatMap(url => Array.from(getEventsForUrl(url, filters))),
),
),
)
2025-02-06 01:05:41 +00:00
// Batch-insert calendar events into the store with a single update
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))
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
}
2026-07-28 00:13:14 +00:00
valid.sort((a, b) => getStart(a) - getStart(b))
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))
const kept = $events.filter(e => !superseded.has(getAddress(e)))
const merged: TrustedEvent[] = []
let i = 0
let j = 0
while (i < kept.length && j < valid.length) {
if (getStart(kept[i]) <= getStart(valid[j])) {
merged.push(kept[i++])
} else {
merged.push(valid[j++])
}
}
2026-06-29 23:28:11 +00:00
2026-07-28 00:13:14 +00:00
while (i < kept.length) merged.push(kept[i++])
while (j < valid.length) merged.push(valid[j++])
return merged
2025-02-06 00:26:22 +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) {
insertEvents(matching)
}
})
2026-02-13 23:18:46 +00:00
const unsubscribers = [
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) {
events.update($events => $events.filter(e => !removed.has(e.id)))
for (const id of removed) {
seen.delete(id)
}
2026-04-10 18:09:26 +00:00
}
2025-02-06 00:26:22 +00:00
2026-04-10 18:09:26 +00:00
const matching = added.filter(event => matchFilters(filters, event))
2026-04-10 18:09:26 +00:00
if (matching.length > 0) {
insertEvents(matching)
}
}),
),
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
}
}),
]
2025-02-06 00:26:22 +00:00
const loadTimeframe = (since: number, until: number) => {
const hashes = daysBetween(since, until).map(String)
2026-07-28 16:16:26 +00:00
network.get().request({
2026-06-29 23:28:11 +00:00
relays,
2025-04-09 22:32:18 +00:00
autoClose: true,
2025-10-06 18:23:19 +00:00
signal: controller.signal,
2025-02-06 00:26:22 +00:00
filters: [{kinds: [EVENT_TIME], "#D": hashes}],
})
}
2025-02-06 01:05:41 +00:00
const maybeExhausted = () => {
if (++exhaustedScrollers === 2) {
onExhausted?.()
}
}
2025-02-06 00:26:22 +00:00
const backwardScroller = createScroller({
element,
reverse: true,
onScroll: () => {
const [since, until] = backwardWindow
backwardWindow = [since - interval, since]
2025-02-06 00:26:22 +00:00
if (until > now() - int(2, YEAR)) {
loadTimeframe(since, until)
} else {
backwardScroller.stop()
2025-02-06 01:05:41 +00:00
maybeExhausted()
2025-02-06 00:26:22 +00:00
}
},
})
const forwardScroller = createScroller({
element,
onScroll: () => {
const [since, until] = forwardWindow
forwardWindow = [until, until + interval]
2025-02-06 00:26:22 +00:00
if (until < now() + int(2, YEAR)) {
loadTimeframe(since, until)
} else {
forwardScroller.stop()
2025-02-06 01:05:41 +00:00
maybeExhausted()
2025-02-06 00:26:22 +00:00
}
},
})
return {
events,
cleanup: () => {
2025-04-11 16:27:19 +00:00
controller.abort()
2026-02-13 23:18:46 +00:00
forwardScroller.stop()
backwardScroller.stop()
unsubscribers.forEach(call)
2025-02-06 00:26:22 +00:00
},
}
}