From 8a8ca65d5d35bef291a9fa2b1c250429ff5c09e7 Mon Sep 17 00:00:00 2001 From: Jon Staab Date: Fri, 25 Sep 2026 21:20:44 -0700 Subject: [PATCH] Keep the chats index and latest-event store on per-app Chats and LatestEvents plugins instead of hand-rebinding module state --- .agents/skills/flotilla-architecture/SKILL.md | 2 +- .agents/skills/flotilla-state/SKILL.md | 18 +- src/app/chats.ts | 208 ++++++++---------- src/app/components/Chat.svelte | 4 +- src/app/components/ProfileBadges.svelte | 4 +- src/app/notifications.ts | 4 +- src/app/repository.ts | 113 +++++----- 7 files changed, 174 insertions(+), 179 deletions(-) diff --git a/.agents/skills/flotilla-architecture/SKILL.md b/.agents/skills/flotilla-architecture/SKILL.md index c033fc67..706b6eba 100644 --- a/.agents/skills/flotilla-architecture/SKILL.md +++ b/.agents/skills/flotilla-architecture/SKILL.md @@ -81,7 +81,7 @@ only gate. cache - `sync.ts`: `syncApplicationData`, the background sync of user data, spaces and DMs - `settings.ts`: the `Settings` plugin over encrypted app data, plus notification settings -- `repository.ts`: `derive*` helpers over the current app's repository +- `repository.ts`: the `LatestEvents` plugin, each watched author's most recent event - `thunks.ts` (publish status by event id), `signer.ts` (signer request tracking) - `env.ts`: every `VITE_` value, parsed - `logger.ts` (log capture and sending), `analytics.ts` (Plausible pageviews), `device.ts` (a diff --git a/.agents/skills/flotilla-state/SKILL.md b/.agents/skills/flotilla-state/SKILL.md index caf82151..836cd9cc 100644 --- a/.agents/skills/flotilla-state/SKILL.md +++ b/.agents/skills/flotilla-state/SKILL.md @@ -88,8 +88,9 @@ The rule is about when a binding is made: - **Module scope, and anything outside the gate**, goes through `app`, `usePlugin`, `fromApp` or `deriveUserItem`. A module-level `app.get()` builds the first app before the policies - register, and a binding made that way keeps reading the discarded app after login. Long-lived - listeners re-bind on `app.subscribe`, as `chatsById` in `src/app/chats.ts`, + register, and a binding made that way keeps reading the discarded app after login. A + long-lived listener over the repository belongs on a plugin, as in `Chats` (`chats.ts`), so + it is dropped with the app. Other listeners re-bind on `app.subscribe`, as `syncCheckedRemote` in `notifications.ts` and the resync in the root layout do. - **Code under the gate** can bind at call time. `rooms.get().forUrl(url).$` inside `deriveUserRooms`, or `$app.use(X)` in a component's script, is fine there, because the app @@ -118,6 +119,8 @@ they cache is itself a rebinding store. `commandsByUrl` holds `fromApp` stores. | `Commands` (`commands.ts`) | `RelayScopedDerivedPlugin` | slash-command definitions, keyed per relay | | `HealthChecks` (`healthChecks.ts`) | none | a plain class over `IApp` exposing `Projection`s | | `SpaceManagement` (`management.ts`) | none | NIP-86 supported methods (5-minute TTL), admins and bans per URL | +| `Chats` (`chats.ts`) | none | DM conversations by chat id (`index`, `one`), and a fuzzy `search` | +| `LatestEvents` (`repository.ts`) | none | the most recent event held from an author (`forPubkey`) | Each is exposed with `usePlugin`. `Statuses` is the minimal shape: @@ -186,7 +189,7 @@ Every method returns a projection over the current app's repository: Components bind `$events.forUrl(url, filters).$`. A `derive*` function called from gated code binds `events.get().forUrl(...).$` at call time, and a module-scope store wraps it in `fromApp($app => $app.use(Events).byIdByUrl(filters).$)`, as `latestActivityByPath` does. -`src/app/repository.ts` holds only `deriveLatestEvent`. Plugin reads (`get`, `one`, `load`, +`src/app/repository.ts` holds only the `LatestEvents` plugin. Plugin reads (`get`, `one`, `load`, `index`) are documented in welshman-app. ### Free functions in an app module @@ -224,14 +227,15 @@ Rules that involve more than one plugin belong in these functions rather than in A derivation that every row subscribes to, or that joins large sets, is built by hand: -- `chatsById` (`chats.ts`) updates incrementally from repository `update` events rather than - re-querying. +- `Chats` (`chats.ts`) updates its `index` incrementally from repository `update` events rather + than re-querying. Its `search` reads profile names when it is rebuilt, so it also follows the + profiles index. - `thunksByEventId` (`thunks.ts`) indexes thunk history once, and hands back the previous array wherever an event's thunks are unchanged so rows don't churn. - `latestActivityByPath` (`notifications.ts`) joins chats, room lists, relay info, events and settings behind `throttled(1000, …)`. -- `deriveLatestEvent` (`repository.ts`) shares one repository listener across every watched - author. +- `LatestEvents.forPubkey` (`repository.ts`) shares one repository listener across every + watched author. ## Local and persisted state diff --git a/src/app/chats.ts b/src/app/chats.ts index 1823ecb5..c92995f5 100644 --- a/src/app/chats.ts +++ b/src/app/chats.ts @@ -1,13 +1,23 @@ import {derived, readable} from "svelte/store" -import {append, call, on, partition, remove, sort, sortBy, uniq, uniqBy} from "@welshman/lib" -import type {Override} from "@welshman/lib" -import {DELETE, PROFILE, getPow, hexTags, tagValues} from "@welshman/util" +import type {Readable} from "svelte/store" +import {append, identity, on, partition, remove, sort, sortBy, uniq, uniqBy} from "@welshman/lib" +import type {Maybe} from "@welshman/lib" +import {getPow, hexTags, tagValues} from "@welshman/util" import type {TrustedEvent} from "@welshman/util" import type {RepositoryUpdate, WrapItem} from "@welshman/net" import {makeDeriveItem, throttled} from "@welshman/store" -import {FollowLists, RelayMemberLists, RoomLists, Wot, WotScope, createSearch} from "@welshman/app" -import type {App} from "@welshman/app" -import {app, deriveUserItem, fromApp, profiles, user} from "@app/core" +import { + FollowLists, + Profiles, + RelayMemberLists, + RoomLists, + Wot, + WotScope, + createSearch, + projection, +} from "@welshman/app" +import type {IApp, Projection, Search} from "@welshman/app" +import {app, deriveUserItem, fromApp, usePlugin, user} from "@app/core" import {DM_KINDS} from "@app/content" import {userSettingsValues} from "@app/settings" @@ -16,7 +26,6 @@ export type Chat = { pubkeys: string[] messages: TrustedEvent[] last_activity: number - search_text: string } export const getChatPubkeys = (pubkeys: string[]) => sort(uniq(append(user.get().pubkey, pubkeys))) @@ -41,128 +50,103 @@ export const isUserMessage = (event: TrustedEvent) => { return event.pubkey === pubkey || tagValues(hexTags("p"), event.tags).includes(pubkey) } -export const chatsById = call(() => { - const chatsById = new Map() - const chatsByPubkey = new Map() +export class Chats { + private chatsById = new Map() - const displayProfile = (pubkey: string) => profiles.get().display(pubkey).get() + index: Projection> + one: (id: string) => Readable> + search: Readable> - const addSearchText = (chat: Override) => { - chat.search_text = - chat.pubkeys.length === 1 - ? displayProfile(chat.pubkeys[0]) + " note to self" - : remove(user.get().pubkey, chat.pubkeys).map(displayProfile).join(" ") + constructor(private readonly app: IApp) { + this.index = projection( + readable(this.chatsById, set => { + this.chatsById.clear() + this.addEvents(app.repository.query([{kinds: DM_KINDS}])) + set(this.chatsById) - return chat as Chat + return on(app.repository, "update", ({added, removed}: RepositoryUpdate) => { + const addedChats = this.addEvents(added) + const removedChats = this.removeEvents(removed) + + if (addedChats || removedChats) { + set(this.chatsById) + } + }) + }), + ) + + this.one = makeDeriveItem(this.index.$) + + // Profiles arrive after the messages that name them, so search text is read when the search is built. + this.search = derived( + throttled(1500, derived([this.index.$, app.use(Profiles).index.$], identity)), + ([$chatsById]) => + createSearch( + sortBy(c => -c.last_activity, Array.from($chatsById.values())), + { + getValue: (chat: Chat) => chat.id, + fuseOptions: {keys: [{name: "search_text", getFn: this.getSearchText}]}, + }, + ), + ) } - return readable(chatsById, set => { - const indexChatByPubkeys = (chat: Chat) => { - for (const pubkey of chat.pubkeys) { - chatsByPubkey.set(pubkey, uniq(append(chat.id, chatsByPubkey.get(pubkey) || []))) + private getSearchText = (chat: Chat) => { + const display = (pubkey: string) => this.app.use(Profiles).display(pubkey).get() + + return chat.pubkeys.length === 1 + ? display(chat.pubkeys[0]) + " note to self" + : remove(user.get().pubkey, chat.pubkeys).map(display).join(" ") + } + + private addEvents = (events: TrustedEvent[]) => { + let dirty = false + + for (const event of events) { + if (DM_KINDS.includes(event.kind) && isUserMessage(event)) { + const pubkeys = getChatPubkeysFromEvent(event) + const id = makeChatId(pubkeys) + const chat = this.chatsById.get(id) + const messages = sortBy( + e => -e.created_at, + uniqBy(e => e.id, append(event, chat?.messages || [])), + ) + const last_activity = Math.max(chat?.last_activity || 0, event.created_at) + + this.chatsById.set(id, {id, pubkeys, messages, last_activity}) + + dirty = true } } - const addEvents = (events: TrustedEvent[]) => { - let dirty = false - for (const event of events) { - if (DM_KINDS.includes(event.kind) && isUserMessage(event)) { - const pubkeys = getChatPubkeysFromEvent(event) - const id = makeChatId(pubkeys) - const chat = chatsById.get(id) - const messages = sortBy( - e => -e.created_at, - uniqBy(e => e.id, append(event, chat?.messages || [])), - ) - const last_activity = Math.max(chat?.last_activity || 0, event.created_at) - const updatedChat = addSearchText({id, pubkeys, messages, last_activity}) + return dirty + } - chatsById.set(id, updatedChat) - indexChatByPubkeys(updatedChat) + private removeEvents = (removed: Set) => { + let dirty = false - dirty = true + // `one` dedupes by reference, so an affected chat is replaced rather than mutated in place. + for (const [chatId, chat] of this.chatsById) { + const messages = chat.messages.filter(e => !removed.has(e.id)) + + if (messages.length !== chat.messages.length) { + if (messages.length > 0) { + this.chatsById.set(chatId, {...chat, messages}) + } else { + this.chatsById.delete(chatId) } - if (event.kind === PROFILE) { - for (const chatId of chatsByPubkey.get(event.pubkey) || []) { - const chat = chatsById.get(chatId) - - if (chat) { - addSearchText(chat) - dirty = true - } - } - } - } - - if (dirty) { - set(chatsById) + dirty = true } } - const removeEvents = (removed: Set) => { - let dirty = false + return dirty + } +} - // deriveChat dedupes by reference, so an affected chat is replaced rather than mutated in place. - for (const [chatId, chat] of chatsById) { - const messages = chat.messages.filter(e => !removed.has(e.id)) +export const chats = usePlugin(Chats) - if (messages.length !== chat.messages.length) { - if (messages.length > 0) { - chatsById.set(chatId, {...chat, messages}) - } else { - chatsById.delete(chatId) - } - - dirty = true - } - } - - if (dirty) { - set(chatsById) - } - } - - // A listener bound to `app.get().repository` would keep reading the app login discarded. - let repoUnsubscribe: (() => void) | undefined - - const bindRepository = ($app: App) => { - repoUnsubscribe?.() - - chatsById.clear() - chatsByPubkey.clear() - addEvents($app.repository.query([{kinds: [...DM_KINDS, DELETE, PROFILE]}])) - set(chatsById) - - repoUnsubscribe = on($app.repository, "update", ({added, removed}: RepositoryUpdate) => { - // Do this async so that profiles are populated - setTimeout(() => { - addEvents(added) - removeEvents(removed) - }, 200) - }) - } - - const unsubscribeApp = app.subscribe(bindRepository) - - return () => { - repoUnsubscribe?.() - unsubscribeApp() - } - }) -}) - -export const deriveChat = makeDeriveItem(chatsById) - -export const chatSearch = derived(throttled(1500, chatsById), $chatsByPubkey => { - return createSearch( - sortBy(c => -c.last_activity, Array.from($chatsByPubkey.values())), - { - getValue: (chat: Chat) => chat.id, - fuseOptions: {keys: ["search_text"]}, - }, - ) -}) +export const chatSearch = fromApp($app => $app.use(Chats).search) // Conversations and requests diff --git a/src/app/components/Chat.svelte b/src/app/components/Chat.svelte index 5b6f9dc7..ee543d2a 100644 --- a/src/app/components/Chat.svelte +++ b/src/app/components/Chat.svelte @@ -38,7 +38,7 @@ import ThunkToast from "@app/components/ThunkToast.svelte" import {app, deletes, router, user, wraps} from "@app/core" import {loadSendDelay} from "@app/settings" - import {deriveChat, makeChatId} from "@app/chats" + import {chats, makeChatId} from "@app/chats" import {makeFeedContext} from "@app/feeds" import {navigate, pushModal} from "@app/modal" import {prependParent} from "@app/rooms" @@ -55,7 +55,7 @@ const context = makeFeedContext({relays: $router.resolver.relays([userInbox()])}) onDestroy(context.cleanup) - const chat = deriveChat(chatId) + const chat = $chats.one(chatId) const others = remove($user.pubkey, pubkeys) const messagingRelayLists = $app.use(MessagingRelayLists).index.$ const missingRelayLists = $derived(others.filter(pk => !$messagingRelayLists.has(pk))) diff --git a/src/app/components/ProfileBadges.svelte b/src/app/components/ProfileBadges.svelte index 208d77a9..155ac0d3 100644 --- a/src/app/components/ProfileBadges.svelte +++ b/src/app/components/ProfileBadges.svelte @@ -5,7 +5,7 @@ import Button from "@lib/components/Button.svelte" import ProfileSpaces from "@app/components/ProfileSpaces.svelte" import {network, roomLists} from "@app/core" - import {deriveLatestEvent} from "@app/repository" + import {latestEvents} from "@app/repository" import {goToEvent} from "@app/routes" import {pushModal} from "@app/modal" @@ -16,7 +16,7 @@ const {pubkey, url}: Props = $props() - const latest = deriveLatestEvent(pubkey) + const latest = $latestEvents.forPubkey(pubkey) const spaceUrls = $roomLists.urls(pubkey).$ diff --git a/src/app/notifications.ts b/src/app/notifications.ts index 0cb2d821..db94ee09 100644 --- a/src/app/notifications.ts +++ b/src/app/notifications.ts @@ -20,7 +20,7 @@ import {app, fromApp, reader} from "@app/core" import {makeRoomPath, makeSpaceChatPath, makeChatPath, makeContentPath} from "@app/routes" import {CONTENT_KINDS, makeCommentFilter} from "@app/content" import {getIsMuted, notificationSettings, userSettingsValues} from "@app/settings" -import {chatsById} from "@app/chats" +import {Chats} from "@app/chats" import {dufflepud, DUFFLEPUD_URL, PLATFORM_RELAYS} from "@app/env" import {kv} from "@app/storage" @@ -232,7 +232,7 @@ export const latestActivityByPath = derived( derived( [ app, - chatsById, + fromApp($app => $app.use(Chats).index.$), fromApp($app => $app.use(Relays).index.$), fromApp($app => $app.use(RoomLists).index.$), fromApp( diff --git a/src/app/repository.ts b/src/app/repository.ts index a7061133..27f945de 100644 --- a/src/app/repository.ts +++ b/src/app/repository.ts @@ -4,72 +4,79 @@ import {first, on} from "@welshman/lib" import type {Maybe} from "@welshman/lib" import {sortEventsDesc} from "@welshman/util" import type {TrustedEvent} from "@welshman/util" -import {app} from "@app/core" +import type {RepositoryUpdate} from "@welshman/net" +import type {IApp} from "@welshman/app" +import {usePlugin} from "@app/core" -// The most recent event held from one author. +// The most recent event held from each watched author, sharing one repository listener across them. +export class LatestEvents { + private latestByPubkey = new Map>() + private subscribers = new Map) => void>>() + private unsubscriber: Maybe -const latestByPubkey = new Map>() + constructor(private readonly app: IApp) {} -const latestSubscribers = new Map) => void>>() + private read = (pubkey: string) => + first(sortEventsDesc(this.app.repository.query([{authors: [pubkey]}]))) -let latestUnsubscriber: Maybe + private onUpdate = ({added, removed}: RepositoryUpdate) => { + const touched = new Set() -const readLatest = (pubkey: string) => - first(sortEventsDesc(app.get().repository.query([{authors: [pubkey]}]))) + for (const event of added) { + if (this.subscribers.has(event.pubkey)) { + const current = this.latestByPubkey.get(event.pubkey) -export const deriveLatestEvent = (pubkey: string) => - readable>(undefined, set => { - let subscribers = latestSubscribers.get(pubkey) - - if (!subscribers) { - subscribers = new Set() - latestSubscribers.set(pubkey, subscribers) - latestByPubkey.set(pubkey, readLatest(pubkey)) + if (!current || event.created_at > current.created_at) { + this.latestByPubkey.set(event.pubkey, event) + touched.add(event.pubkey) + } + } } - subscribers.add(set) - set(latestByPubkey.get(pubkey)) + // A removal can take the very event a row is showing, so that author is looked up again. + for (const [author, event] of this.latestByPubkey) { + if (event && removed.has(event.id)) { + this.latestByPubkey.set(author, this.read(author)) + touched.add(author) + } + } - latestUnsubscriber ??= on(app.get().repository, "update", ({added, removed}) => { - const touched = new Set() + for (const author of touched) { + for (const subscriber of this.subscribers.get(author) || []) { + subscriber(this.latestByPubkey.get(author)) + } + } + } - for (const event of added) { - if (latestSubscribers.has(event.pubkey)) { - const current = latestByPubkey.get(event.pubkey) + forPubkey = (pubkey: string) => + readable>(undefined, set => { + let subscribers = this.subscribers.get(pubkey) - if (!current || event.created_at > current.created_at) { - latestByPubkey.set(event.pubkey, event) - touched.add(event.pubkey) + if (!subscribers) { + subscribers = new Set() + this.subscribers.set(pubkey, subscribers) + this.latestByPubkey.set(pubkey, this.read(pubkey)) + } + + subscribers.add(set) + set(this.latestByPubkey.get(pubkey)) + + this.unsubscriber ??= on(this.app.repository, "update", this.onUpdate) + + return () => { + subscribers.delete(set) + + if (subscribers.size === 0) { + this.subscribers.delete(pubkey) + this.latestByPubkey.delete(pubkey) + + if (this.subscribers.size === 0) { + this.unsubscriber?.() + this.unsubscriber = undefined } } } - - // A removal can take the very event a row is showing, so that author is looked up again. - for (const [author, event] of latestByPubkey) { - if (event && removed.has(event.id)) { - latestByPubkey.set(author, readLatest(author)) - touched.add(author) - } - } - - for (const author of touched) { - for (const subscriber of latestSubscribers.get(author) || []) { - subscriber(latestByPubkey.get(author)) - } - } }) +} - return () => { - subscribers.delete(set) - - if (subscribers.size === 0) { - latestSubscribers.delete(pubkey) - latestByPubkey.delete(pubkey) - - if (latestSubscribers.size === 0) { - latestUnsubscriber?.() - latestUnsubscriber = undefined - } - } - } - }) +export const latestEvents = usePlugin(LatestEvents)