package main // Notif Worker - notification subscription, dedup, filtering, watermark tracking. // // supervisor -> worker: // ["N_EVENT", subID, evJSON] // ["N_EOSE", subID] // ["N_SET_PUBKEY", pk] // ["N_SET_RELAYS", urlsJSON] // ["N_SET_MUTES", pksJSON] // ["N_LOAD_MORE"] // ["N_MARK_READ"] // ["N_SET_FILTER", filter] -- "all", "mentions", "reactions", "zaps" // ["N_READ_TS", ts] -- initial notifSt.notifReadTs from IDB // // worker -> supervisor (routed to Shell): // ["N_RENDER", evJSON, mode] -- mode: "prepend" or "append" // ["N_DOT", show] -- show=1 or 0 // ["N_EOSE_DONE", subID] // ["N_STATUS", count, notifSt.exhausted] // ["N_SUB", subID, filterJSON, urlsJSON] // ["N_CLOSE", subID] // ["N_FETCH_REFS", idsJSON] -- Shell should fetch these referenced events // ["N_STORE_READ_TS", ts] -- save notif-read-ts to IDB import ( "runtime" "git.smesh.lol/musiquay/web/common/accum" "git.smesh.lol/musiquay/web/common/helpers" "git.smesh.lol/musiquay/web/common/jsbridge/notif" "git.smesh.lol/musiquay/web/common/mw" "git.smesh.lol/nostr/pkg/core" ) // notifState is the notification // Package globals are immutable outside initialization, so it lives in // one self-mutating type reached through a package-level pointer. type notifState struct { myPK string relayURLs accum.Strings muteSet map[string]bool notifFilter string notifReadTs int64 notifSeen map[string]bool oldestTs int64 eventCount int32 exhausted bool emptyStreak int32 initLoad bool loading bool moreGot int32 moreTimer int32 notifMissing accum.Strings refTimer int32 } var notifSt *notifState func initState() { if notifSt != nil { return } // notifSt 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()) notifSt = ¬ifState{} notifSt.muteSet = map[string]bool{} notifSt.notifSeen = map[string]bool{} notifSt.notifFilter = "all" notifSt.initLoad = true runtime.SovereignRestoreArena(runtime.RootArena()) } func main() { initState() notif.WorkerOnMessage(handleMessage) } func handleMessage(msg string) { w := mw.New(msg) switch w.Str() { case "N_SET_PUBKEY": notifSt.myPK = w.Str() case "N_SET_RELAYS": notifSt.relayURLs.Set(parseStringArray(w.Raw())) case "N_SET_MUTES": pks := parseStringArray(w.Raw()) notifSt.muteSet = map[string]bool{} for _, pk := range pks { notifSt.muteSet[pk] = true } case "N_SET_FILTER": notifSt.notifFilter = w.Str() case "N_READ_TS": notifSt.notifReadTs = w.Num() case "N_EVENT": _ = w.Str() // subID consumed evJSON := w.Raw() handleNotifEvent(evJSON) case "N_EOSE": subID := w.Str() handleEOSE(subID) case "N_SUBSCRIBE": handleSubscribe() case "N_LOAD_MORE": handleLoadMore() case "N_MARK_READ": handleMarkRead() } } func handleNotifEvent(evJSON string) { ev := nostr.ParseEvent(evJSON) if ev == nil { return } if notifSt.notifSeen == nil { notifSt.notifSeen = map[string]bool{} } if notifSt.notifSeen[ev.ID] { return } notifSt.notifSeen[ev.ID] = true if notifSt.oldestTs == 0 || ev.CreatedAt < notifSt.oldestTs { notifSt.oldestTs = ev.CreatedAt } if ev.PubKey == notifSt.myPK { return } if notifSt.muteSet[notifSenderPK(ev)] { return } if repliesToMuted(ev) { return } // Queue missing referenced events if ev.Kind != 1 { refID := notifRefID(ev) if refID != "" { notifSt.notifMissing.Push(refID) } } // Check against read watermark for dot display if ev.CreatedAt > notifSt.notifReadTs { notif.WorkerPost(`["N_DOT",1]`) } if !passesFilter(ev) { return } mode := "append" if !notifSt.initLoad { mode = "prepend" } notif.WorkerPost(`["N_RENDER",` | evJSON | `,"` | mode | `"]`) notifSt.eventCount++ } func handleEOSE(subID string) { if subID == "ntf" { notifSt.initLoad = false if notifSt.notifMissing.Len() > 0 { notif.WorkerPost(`["N_FETCH_REFS",` | buildIDsJSON() | `]`) notifSt.notifMissing.Reset() } notif.WorkerPost(`["N_EOSE_DONE","ntf"]`) notif.WorkerPost(`["N_STATUS",` | helpers.Itoa(int64(notifSt.eventCount)) | `,0]`) } else if subID == "ntf-more" { if notifSt.moreTimer != 0 { notif.WorkerClearTimeout(notifSt.moreTimer); notifSt.moreTimer = 0 } notifSt.loading = false if notifSt.moreGot == 0 { notifSt.emptyStreak++ if notifSt.emptyStreak >= 3 { notifSt.exhausted = true } } else { notifSt.emptyStreak = 0 } ex := int64(0) if notifSt.exhausted { ex = 1 } notif.WorkerPost(`["N_STATUS",` | helpers.Itoa(int64(notifSt.eventCount)) | `,` | helpers.Itoa(ex) | `]`) notif.WorkerPost(`["N_EOSE_DONE","ntf-more"]`) notif.WorkerPost(`["N_CLOSE","ntf-more"]`) } } func handleSubscribe() { if notifSt.myPK == "" || notifSt.relayURLs.Len() == 0 { return } // Close any existing ntf sub, then open fresh. notif.WorkerPost(`["N_CLOSE","ntf"]`) filter := `{"kinds":[1,6,7,9735],"#p":[` | jstr(notifSt.myPK) | `],"limit":20}` notif.WorkerPost(`["N_SUB","ntf",` | filter | `,` | buildURLsJSON(notifSt.relayURLs.Slice()) | `]`) } func handleLoadMore() { if notifSt.loading || notifSt.exhausted || notifSt.oldestTs == 0 || notifSt.myPK == "" { return } notifSt.loading = true notifSt.moreGot = 0 filter := `{"kinds":[1,6,7,9735],"#p":[` | jstr(notifSt.myPK) | `],"until":` | helpers.Itoa(notifSt.oldestTs) | `,"limit":20}` notif.WorkerPost(`["N_SUB","ntf-more",` | filter | `,` | buildURLsJSON(notifSt.relayURLs.Slice()) | `]`) notifSt.moreTimer = notif.WorkerSetTimeout(10000, func() { notifSt.moreTimer = 0 notifSt.loading = false notif.WorkerPost(`["N_EOSE_DONE","ntf-more"]`) }) } func handleMarkRead() { ts := notif.WorkerNowSeconds() notifSt.notifReadTs = ts notif.WorkerPost(`["N_STORE_READ_TS",` | helpers.Itoa(ts) | `]`) notif.WorkerPost(`["N_DOT",0]`) } func passesFilter(ev *nostr.Event) (ok bool) { switch notifSt.notifFilter { case "mentions": return ev.Kind == 1 case "reactions": return ev.Kind == 7 case "zaps": return ev.Kind == 9735 default: return true } } func notifRefID(ev *nostr.Event) (s string) { switch ev.Kind { case 7, 9735, 6: tag := nostr.TagsGetFirst(ev.Tags, []byte("e")) if tag != nil { return string(nostr.TagValue(tag)) } } return "" } func notifSenderPK(ev *nostr.Event) (s string) { if ev.Kind == 9735 { pTag := nostr.TagsGetFirst(ev.Tags, []byte("P")) if pTag != nil { return string(nostr.TagValue(pTag)) } } return ev.PubKey } func repliesToMuted(ev *nostr.Event) (ok bool) { for _, tag := range nostr.TagsGetAll(ev.Tags, "p") { if v := string(nostr.TagValue(tag)); v != "" && notifSt.muteSet[v] { return true } } return false } // buildIDsJSON serialises the pending reference ids. It reads the holder // directly: the ids live in chunks owned by the holder's sovereign arena, and // the only slice the worker ever passes around is the one built here, in this // frame, which dies at return. func buildIDsJSON() (s string) { s = "[" for i := int32(0); i < notifSt.notifMissing.Len(); i++ { if i > 0 { s = s | "," } s = s | jstr(notifSt.notifMissing.At(i)) } return s | "]" } 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) }