Keep the chats index and latest-event store on per-app Chats and LatestEvents plugins instead of hand-rebinding module state
This commit is contained in:
parent
66af2a201e
commit
8a8ca65d5d
7 changed files with 174 additions and 179 deletions
|
|
@ -81,7 +81,7 @@ only gate.
|
||||||
cache
|
cache
|
||||||
- `sync.ts`: `syncApplicationData`, the background sync of user data, spaces and DMs
|
- `sync.ts`: `syncApplicationData`, the background sync of user data, spaces and DMs
|
||||||
- `settings.ts`: the `Settings` plugin over encrypted app data, plus notification settings
|
- `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)
|
- `thunks.ts` (publish status by event id), `signer.ts` (signer request tracking)
|
||||||
- `env.ts`: every `VITE_` value, parsed
|
- `env.ts`: every `VITE_` value, parsed
|
||||||
- `logger.ts` (log capture and sending), `analytics.ts` (Plausible pageviews), `device.ts` (a
|
- `logger.ts` (log capture and sending), `analytics.ts` (Plausible pageviews), `device.ts` (a
|
||||||
|
|
|
||||||
|
|
@ -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`
|
- **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
|
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
|
register, and a binding made that way keeps reading the discarded app after login. A
|
||||||
listeners re-bind on `app.subscribe`, as `chatsById` in `src/app/chats.ts`,
|
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.
|
`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
|
- **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
|
`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 |
|
| `Commands` (`commands.ts`) | `RelayScopedDerivedPlugin` | slash-command definitions, keyed per relay |
|
||||||
| `HealthChecks` (`healthChecks.ts`) | none | a plain class over `IApp` exposing `Projection`s |
|
| `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 |
|
| `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:
|
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
|
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
|
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.
|
`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.
|
`index`) are documented in welshman-app.
|
||||||
|
|
||||||
### Free functions in an app module
|
### 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:
|
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
|
- `Chats` (`chats.ts`) updates its `index` incrementally from repository `update` events rather
|
||||||
re-querying.
|
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
|
- `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.
|
wherever an event's thunks are unchanged so rows don't churn.
|
||||||
- `latestActivityByPath` (`notifications.ts`) joins chats, room lists, relay info, events and
|
- `latestActivityByPath` (`notifications.ts`) joins chats, room lists, relay info, events and
|
||||||
settings behind `throttled(1000, …)`.
|
settings behind `throttled(1000, …)`.
|
||||||
- `deriveLatestEvent` (`repository.ts`) shares one repository listener across every watched
|
- `LatestEvents.forPubkey` (`repository.ts`) shares one repository listener across every
|
||||||
author.
|
watched author.
|
||||||
|
|
||||||
## Local and persisted state
|
## Local and persisted state
|
||||||
|
|
||||||
|
|
|
||||||
208
src/app/chats.ts
208
src/app/chats.ts
|
|
@ -1,13 +1,23 @@
|
||||||
import {derived, readable} from "svelte/store"
|
import {derived, readable} from "svelte/store"
|
||||||
import {append, call, on, partition, remove, sort, sortBy, uniq, uniqBy} from "@welshman/lib"
|
import type {Readable} from "svelte/store"
|
||||||
import type {Override} from "@welshman/lib"
|
import {append, identity, on, partition, remove, sort, sortBy, uniq, uniqBy} from "@welshman/lib"
|
||||||
import {DELETE, PROFILE, getPow, hexTags, tagValues} from "@welshman/util"
|
import type {Maybe} from "@welshman/lib"
|
||||||
|
import {getPow, hexTags, tagValues} from "@welshman/util"
|
||||||
import type {TrustedEvent} from "@welshman/util"
|
import type {TrustedEvent} from "@welshman/util"
|
||||||
import type {RepositoryUpdate, WrapItem} from "@welshman/net"
|
import type {RepositoryUpdate, WrapItem} from "@welshman/net"
|
||||||
import {makeDeriveItem, throttled} from "@welshman/store"
|
import {makeDeriveItem, throttled} from "@welshman/store"
|
||||||
import {FollowLists, RelayMemberLists, RoomLists, Wot, WotScope, createSearch} from "@welshman/app"
|
import {
|
||||||
import type {App} from "@welshman/app"
|
FollowLists,
|
||||||
import {app, deriveUserItem, fromApp, profiles, user} from "@app/core"
|
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 {DM_KINDS} from "@app/content"
|
||||||
import {userSettingsValues} from "@app/settings"
|
import {userSettingsValues} from "@app/settings"
|
||||||
|
|
||||||
|
|
@ -16,7 +26,6 @@ export type Chat = {
|
||||||
pubkeys: string[]
|
pubkeys: string[]
|
||||||
messages: TrustedEvent[]
|
messages: TrustedEvent[]
|
||||||
last_activity: number
|
last_activity: number
|
||||||
search_text: string
|
|
||||||
}
|
}
|
||||||
|
|
||||||
export const getChatPubkeys = (pubkeys: string[]) => sort(uniq(append(user.get().pubkey, pubkeys)))
|
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)
|
return event.pubkey === pubkey || tagValues(hexTags("p"), event.tags).includes(pubkey)
|
||||||
}
|
}
|
||||||
|
|
||||||
export const chatsById = call(() => {
|
export class Chats {
|
||||||
const chatsById = new Map<string, Chat>()
|
private chatsById = new Map<string, Chat>()
|
||||||
const chatsByPubkey = new Map<string, string[]>()
|
|
||||||
|
|
||||||
const displayProfile = (pubkey: string) => profiles.get().display(pubkey).get()
|
index: Projection<Map<string, Chat>>
|
||||||
|
one: (id: string) => Readable<Maybe<Chat>>
|
||||||
|
search: Readable<Search<string, Chat>>
|
||||||
|
|
||||||
const addSearchText = (chat: Override<Chat, {search_text?: string}>) => {
|
constructor(private readonly app: IApp) {
|
||||||
chat.search_text =
|
this.index = projection(
|
||||||
chat.pubkeys.length === 1
|
readable(this.chatsById, set => {
|
||||||
? displayProfile(chat.pubkeys[0]) + " note to self"
|
this.chatsById.clear()
|
||||||
: remove(user.get().pubkey, chat.pubkeys).map(displayProfile).join(" ")
|
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 => {
|
private getSearchText = (chat: Chat) => {
|
||||||
const indexChatByPubkeys = (chat: Chat) => {
|
const display = (pubkey: string) => this.app.use(Profiles).display(pubkey).get()
|
||||||
for (const pubkey of chat.pubkeys) {
|
|
||||||
chatsByPubkey.set(pubkey, uniq(append(chat.id, chatsByPubkey.get(pubkey) || [])))
|
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[]) => {
|
return dirty
|
||||||
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})
|
|
||||||
|
|
||||||
chatsById.set(id, updatedChat)
|
private removeEvents = (removed: Set<string>) => {
|
||||||
indexChatByPubkeys(updatedChat)
|
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) {
|
dirty = true
|
||||||
for (const chatId of chatsByPubkey.get(event.pubkey) || []) {
|
|
||||||
const chat = chatsById.get(chatId)
|
|
||||||
|
|
||||||
if (chat) {
|
|
||||||
addSearchText(chat)
|
|
||||||
dirty = true
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
if (dirty) {
|
|
||||||
set(chatsById)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
const removeEvents = (removed: Set<string>) => {
|
return dirty
|
||||||
let dirty = false
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// deriveChat dedupes by reference, so an affected chat is replaced rather than mutated in place.
|
export const chats = usePlugin(Chats)
|
||||||
for (const [chatId, chat] of chatsById) {
|
|
||||||
const messages = chat.messages.filter(e => !removed.has(e.id))
|
|
||||||
|
|
||||||
if (messages.length !== chat.messages.length) {
|
export const chatSearch = fromApp($app => $app.use(Chats).search)
|
||||||
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"]},
|
|
||||||
},
|
|
||||||
)
|
|
||||||
})
|
|
||||||
|
|
||||||
// Conversations and requests
|
// Conversations and requests
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -38,7 +38,7 @@
|
||||||
import ThunkToast from "@app/components/ThunkToast.svelte"
|
import ThunkToast from "@app/components/ThunkToast.svelte"
|
||||||
import {app, deletes, router, user, wraps} from "@app/core"
|
import {app, deletes, router, user, wraps} from "@app/core"
|
||||||
import {loadSendDelay} from "@app/settings"
|
import {loadSendDelay} from "@app/settings"
|
||||||
import {deriveChat, makeChatId} from "@app/chats"
|
import {chats, makeChatId} from "@app/chats"
|
||||||
import {makeFeedContext} from "@app/feeds"
|
import {makeFeedContext} from "@app/feeds"
|
||||||
import {navigate, pushModal} from "@app/modal"
|
import {navigate, pushModal} from "@app/modal"
|
||||||
import {prependParent} from "@app/rooms"
|
import {prependParent} from "@app/rooms"
|
||||||
|
|
@ -55,7 +55,7 @@
|
||||||
const context = makeFeedContext({relays: $router.resolver.relays([userInbox()])})
|
const context = makeFeedContext({relays: $router.resolver.relays([userInbox()])})
|
||||||
|
|
||||||
onDestroy(context.cleanup)
|
onDestroy(context.cleanup)
|
||||||
const chat = deriveChat(chatId)
|
const chat = $chats.one(chatId)
|
||||||
const others = remove($user.pubkey, pubkeys)
|
const others = remove($user.pubkey, pubkeys)
|
||||||
const messagingRelayLists = $app.use(MessagingRelayLists).index.$
|
const messagingRelayLists = $app.use(MessagingRelayLists).index.$
|
||||||
const missingRelayLists = $derived(others.filter(pk => !$messagingRelayLists.has(pk)))
|
const missingRelayLists = $derived(others.filter(pk => !$messagingRelayLists.has(pk)))
|
||||||
|
|
|
||||||
|
|
@ -5,7 +5,7 @@
|
||||||
import Button from "@lib/components/Button.svelte"
|
import Button from "@lib/components/Button.svelte"
|
||||||
import ProfileSpaces from "@app/components/ProfileSpaces.svelte"
|
import ProfileSpaces from "@app/components/ProfileSpaces.svelte"
|
||||||
import {network, roomLists} from "@app/core"
|
import {network, roomLists} from "@app/core"
|
||||||
import {deriveLatestEvent} from "@app/repository"
|
import {latestEvents} from "@app/repository"
|
||||||
import {goToEvent} from "@app/routes"
|
import {goToEvent} from "@app/routes"
|
||||||
import {pushModal} from "@app/modal"
|
import {pushModal} from "@app/modal"
|
||||||
|
|
||||||
|
|
@ -16,7 +16,7 @@
|
||||||
|
|
||||||
const {pubkey, url}: Props = $props()
|
const {pubkey, url}: Props = $props()
|
||||||
|
|
||||||
const latest = deriveLatestEvent(pubkey)
|
const latest = $latestEvents.forPubkey(pubkey)
|
||||||
|
|
||||||
const spaceUrls = $roomLists.urls(pubkey).$
|
const spaceUrls = $roomLists.urls(pubkey).$
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -20,7 +20,7 @@ import {app, fromApp, reader} from "@app/core"
|
||||||
import {makeRoomPath, makeSpaceChatPath, makeChatPath, makeContentPath} from "@app/routes"
|
import {makeRoomPath, makeSpaceChatPath, makeChatPath, makeContentPath} from "@app/routes"
|
||||||
import {CONTENT_KINDS, makeCommentFilter} from "@app/content"
|
import {CONTENT_KINDS, makeCommentFilter} from "@app/content"
|
||||||
import {getIsMuted, notificationSettings, userSettingsValues} from "@app/settings"
|
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 {dufflepud, DUFFLEPUD_URL, PLATFORM_RELAYS} from "@app/env"
|
||||||
import {kv} from "@app/storage"
|
import {kv} from "@app/storage"
|
||||||
|
|
||||||
|
|
@ -232,7 +232,7 @@ export const latestActivityByPath = derived(
|
||||||
derived(
|
derived(
|
||||||
[
|
[
|
||||||
app,
|
app,
|
||||||
chatsById,
|
fromApp($app => $app.use(Chats).index.$),
|
||||||
fromApp($app => $app.use(Relays).index.$),
|
fromApp($app => $app.use(Relays).index.$),
|
||||||
fromApp($app => $app.use(RoomLists).index.$),
|
fromApp($app => $app.use(RoomLists).index.$),
|
||||||
fromApp(
|
fromApp(
|
||||||
|
|
|
||||||
|
|
@ -4,72 +4,79 @@ import {first, on} from "@welshman/lib"
|
||||||
import type {Maybe} from "@welshman/lib"
|
import type {Maybe} from "@welshman/lib"
|
||||||
import {sortEventsDesc} from "@welshman/util"
|
import {sortEventsDesc} from "@welshman/util"
|
||||||
import type {TrustedEvent} 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<string, Maybe<TrustedEvent>>()
|
||||||
|
private subscribers = new Map<string, Set<(event: Maybe<TrustedEvent>) => void>>()
|
||||||
|
private unsubscriber: Maybe<Unsubscriber>
|
||||||
|
|
||||||
const latestByPubkey = new Map<string, Maybe<TrustedEvent>>()
|
constructor(private readonly app: IApp) {}
|
||||||
|
|
||||||
const latestSubscribers = new Map<string, Set<(event: Maybe<TrustedEvent>) => void>>()
|
private read = (pubkey: string) =>
|
||||||
|
first(sortEventsDesc(this.app.repository.query([{authors: [pubkey]}])))
|
||||||
|
|
||||||
let latestUnsubscriber: Maybe<Unsubscriber>
|
private onUpdate = ({added, removed}: RepositoryUpdate) => {
|
||||||
|
const touched = new Set<string>()
|
||||||
|
|
||||||
const readLatest = (pubkey: string) =>
|
for (const event of added) {
|
||||||
first(sortEventsDesc(app.get().repository.query([{authors: [pubkey]}])))
|
if (this.subscribers.has(event.pubkey)) {
|
||||||
|
const current = this.latestByPubkey.get(event.pubkey)
|
||||||
|
|
||||||
export const deriveLatestEvent = (pubkey: string) =>
|
if (!current || event.created_at > current.created_at) {
|
||||||
readable<Maybe<TrustedEvent>>(undefined, set => {
|
this.latestByPubkey.set(event.pubkey, event)
|
||||||
let subscribers = latestSubscribers.get(pubkey)
|
touched.add(event.pubkey)
|
||||||
|
}
|
||||||
if (!subscribers) {
|
}
|
||||||
subscribers = new Set()
|
|
||||||
latestSubscribers.set(pubkey, subscribers)
|
|
||||||
latestByPubkey.set(pubkey, readLatest(pubkey))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
subscribers.add(set)
|
// A removal can take the very event a row is showing, so that author is looked up again.
|
||||||
set(latestByPubkey.get(pubkey))
|
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}) => {
|
for (const author of touched) {
|
||||||
const touched = new Set<string>()
|
for (const subscriber of this.subscribers.get(author) || []) {
|
||||||
|
subscriber(this.latestByPubkey.get(author))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
for (const event of added) {
|
forPubkey = (pubkey: string) =>
|
||||||
if (latestSubscribers.has(event.pubkey)) {
|
readable<Maybe<TrustedEvent>>(undefined, set => {
|
||||||
const current = latestByPubkey.get(event.pubkey)
|
let subscribers = this.subscribers.get(pubkey)
|
||||||
|
|
||||||
if (!current || event.created_at > current.created_at) {
|
if (!subscribers) {
|
||||||
latestByPubkey.set(event.pubkey, event)
|
subscribers = new Set()
|
||||||
touched.add(event.pubkey)
|
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 () => {
|
export const latestEvents = usePlugin(LatestEvents)
|
||||||
subscribers.delete(set)
|
|
||||||
|
|
||||||
if (subscribers.size === 0) {
|
|
||||||
latestSubscribers.delete(pubkey)
|
|
||||||
latestByPubkey.delete(pubkey)
|
|
||||||
|
|
||||||
if (latestSubscribers.size === 0) {
|
|
||||||
latestUnsubscriber?.()
|
|
||||||
latestUnsubscriber = undefined
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
})
|
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue