package main // Feed Worker - subscription management, dedup, filtering, and pagination for the note feed. // // Shell forwards verified feed events (subID "feed" / "feed-more") to this worker. // This worker deduplicates, applies mute/follow filters, and sends render commands back to Shell. // // supervisor -> worker (from Shell via feed.Send): // ["F_EVENT", subID, evJSON] -- feed or feed-more event arrived // ["F_EOSE", subID] -- EOSE for feed or feed-more sub // ["F_SET_MODE", mode] -- "follows", "relays", or relay URL // ["F_SET_RELAYS", urlsJSON] // ["F_SET_FOLLOWS", pksJSON] -- current follow list (pubkeys) // ["F_SET_MUTES", pksJSON] -- current mute list (pubkeys) // ["F_SET_PUBKEY", pk] // ["F_LOAD_MORE"] -- user clicked "load more" // ["F_REFRESH"] -- user triggered refresh // // worker -> supervisor (to Shell via F_* prefix routing): // ["F_RENDER", evJSON, now] -- now=1: render immediately, 0: buffer // ["F_NEW"] -- new buffered events available (pulse) // ["F_EOSE_DONE", subID] -- signal Shell to update UI after EOSE // ["F_STATUS", count, exhausted] -- event count and exhaustion state // ["F_SUB", subID, filterJSON, urlsJSON] // ["F_CLOSE", subID] import ( "runtime" "git.smesh.lol/musiquay/web/common/accum" "git.smesh.lol/musiquay/web/common/helpers" "git.smesh.lol/musiquay/web/common/jsbridge/feed" "git.smesh.lol/musiquay/web/common/mw" "git.smesh.lol/nostr/pkg/core" ) // feedState is the feed // Package globals are immutable outside initialization, so it lives in // one self-mutating type reached through a package-level pointer. type feedState struct { feedMode string relayURLs accum.Strings followList accum.Strings followSet map[string]bool muteSet map[string]bool myPK string seenEvents map[string]bool oldestFeedTs int64 eventCount int32 feedInitialLoad bool feedExhausted bool feedEmptyStreak int32 feedLoading bool feedMoreGot int32 feedMoreTimer int32 } var feedSt *feedState func initState() { if feedSt != nil { return } // feedSt and everything it holds live as long as this worker: build them // in the root arena, not in this call's frame arena which dies on // return. runtime.SovereignSetArena(runtime.RootArena()) feedSt = &feedState{} feedSt.seenEvents = map[string]bool{} feedSt.followSet = map[string]bool{} feedSt.muteSet = map[string]bool{} feedSt.feedInitialLoad = true runtime.SovereignRestoreArena(runtime.RootArena()) } func main() { initState() feed.WorkerOnMessage(handleMessage) } func handleMessage(msg string) { w := mw.New(msg) switch w.Str() { case "F_EVENT": subID := w.Str() evJSON := w.Raw() handleFeedEvent(subID, evJSON) case "F_EOSE": subID := w.Str() handleEOSE(subID) case "F_SET_MODE": feedSt.feedMode = w.Str() case "F_SET_RELAYS": feedSt.relayURLs.Set(parseStringArray(w.Raw())) case "F_SET_FOLLOWS": feedSt.followList.Set(parseStringArray(w.Raw())) feedSt.followSet = map[string]bool{} for i := int32(0); i < feedSt.followList.Len(); i++ { feedSt.followSet[feedSt.followList.At(i)] = true } case "F_SET_MUTES": pks := parseStringArray(w.Raw()) feedSt.muteSet = map[string]bool{} for _, pk := range pks { feedSt.muteSet[pk] = true } case "F_SET_PUBKEY": feedSt.myPK = w.Str() case "F_LOAD_MORE": handleLoadMore() case "F_REFRESH": handleRefresh() } } func handleFeedEvent(subID, evJSON string) { ev := nostr.ParseEvent(evJSON) if ev == nil { return } if feedSt.seenEvents[ev.ID] { return } feedSt.seenEvents[ev.ID] = true if !feedPassesFilter(ev) { return } if feedSt.muteSet[ev.PubKey] { return } if repliesToMuted(ev) { return } if (ev.Kind == 1 || ev.Kind == 1111) && looksLikeJSONSpam(ev.Content) { return } if feedSt.oldestFeedTs == 0 || ev.CreatedAt < feedSt.oldestFeedTs { feedSt.oldestFeedTs = ev.CreatedAt } feedSt.eventCount++ if subID == "feed-more" { feedSt.feedMoreGot++ feed.WorkerPost(`["F_RENDER",` | evJSON | `,1]`) return } // Events for the initial feed subscription are the feed: render them // immediately. feedSt.feedInitialLoad cannot decide this, because the feed has two // event sources (the local cache query and the relay subscription) and // whichever EOSE arrives first clears the flag: when the local query's // empty EOSE won the race, the relay's stored events arrived afterwards, // were buffered as "new posts", and the feed looked empty. feed.WorkerPost(`["F_RENDER",` | evJSON | `,1]`) } func handleEOSE(subID string) { if subID == "feed" { feedSt.feedInitialLoad = false feed.WorkerPost(`["F_EOSE_DONE","feed"]`) feed.WorkerPost(`["F_STATUS",` | helpers.Itoa(int64(feedSt.eventCount)) | `,0]`) } else if subID == "feed-more" { if feedSt.feedMoreTimer != 0 { feed.WorkerClearTimeout(feedSt.feedMoreTimer) feedSt.feedMoreTimer = 0 } feedSt.feedLoading = false if feedSt.feedMoreGot == 0 { feedSt.feedEmptyStreak++ if feedSt.feedEmptyStreak >= 3 { feedSt.feedExhausted = true } } else { feedSt.feedEmptyStreak = 0 } ex := int64(0) if feedSt.feedExhausted { ex = 1 } feed.WorkerPost(`["F_STATUS",` | helpers.Itoa(int64(feedSt.eventCount)) | `,` | helpers.Itoa(ex) | `]`) feed.WorkerPost(`["F_EOSE_DONE","feed-more"]`) feed.WorkerPost(`["F_CLOSE","feed-more"]`) } } func handleRefresh() { feedSt.seenEvents = map[string]bool{} feedSt.oldestFeedTs = 0 feedSt.eventCount = 0 feedSt.feedInitialLoad = true feedSt.feedExhausted = false feedSt.feedEmptyStreak = 0 feed.WorkerPost(`["F_CLOSE","feed"]`) feed.WorkerPost(`["F_CLOSE","feed-more"]`) subscribe() } func handleLoadMore() { if feedSt.feedLoading || feedSt.feedExhausted || feedSt.oldestFeedTs == 0 { return } feedSt.feedLoading = true feedSt.feedMoreGot = 0 filter := buildFeedFilter(20) until := `,"until":` | helpers.Itoa(feedSt.oldestFeedTs) | `}` filter = filter[:len(filter)-1] | until urls := feedRelays() feed.WorkerPost(`["F_SUB","feed-more",` | filter | `,` | buildURLsJSON(urls) | `]`) feedSt.feedMoreTimer = feed.WorkerSetTimeout(10000, func() { feedSt.feedMoreTimer = 0 feedSt.feedLoading = false feed.WorkerPost(`["F_EOSE_DONE","feed-more"]`) }) } func subscribe() { if feedSt.myPK == "" || feedSt.relayURLs.Len() == 0 { return } filter := buildFeedFilter(20) urls := feedRelays() feed.WorkerPost(`["F_SUB","feed",` | filter | `,` | buildURLsJSON(urls) | `]`) } func buildFeedFilter(limit int32) (s string) { if feedSt.feedMode == "follows" && feedSt.followList.Len() > 0 { authors := jstr(feedSt.myPK) for i := int32(0); i < feedSt.followList.Len(); i++ { pk := feedSt.followList.At(i) if pk != feedSt.myPK { authors = authors | "," | jstr(pk) } } return `{"kinds":[1,6,7,1111],"authors":[` | authors | `],"limit":` | helpers.Itoa(int64(limit)) | `}` } return `{"kinds":[1,6,7,1111],"limit":` | helpers.Itoa(int64(limit)) | `}` } func feedPassesFilter(ev *nostr.Event) (ok bool) { if feedSt.feedMode != "follows" { return true } if feedSt.followList.Len() == 0 { return true } if ev.PubKey == feedSt.myPK { return true } return feedSt.followSet[ev.PubKey] } func feedRelays() (ss []string) { if feedSt.feedMode != "" && feedSt.feedMode != "follows" && feedSt.feedMode != "relays" { return []string{feedSt.feedMode} } return feedSt.relayURLs.Slice() } func repliesToMuted(ev *nostr.Event) (ok bool) { for _, tag := range nostr.TagsGetAll(ev.Tags, "p") { if v := string(nostr.TagValue(tag)); v != "" && feedSt.muteSet[v] { return true } } return false } func looksLikeJSONSpam(content string) (ok bool) { i := 0 for i < len(content) { c := content[i] if c == ' ' || c == '\n' || c == '\r' || c == '\t' { i++ } else { break } } if i >= len(content) { return false } open := content[i] return open == '{' || open == '[' } func buildURLsJSON(urls []string) (s string) { s := "[" for i, u := range urls { if i > 0 { s = s | "," } s = s | jstr(u) } return s | "]" } func parseStringArray(json string) (ss []string) { var out []string i := 0 for i < len(json) && json[i] != '[' { i++ } if i >= len(json) { return nil } i++ for { for i < len(json) && (json[i] == ' ' || json[i] == ',' || json[i] == '\n') { i++ } if i >= len(json) || json[i] == ']' { break } if json[i] != '"' { break } i++ start := i for i < len(json) && json[i] != '"' { if json[i] == '\\' { i++ } i++ } if i >= len(json) { break } out = push(out, json[start:i]) i++ } return out } func jstr(s string) (sv string) { return helpers.JsonString(s) }