main.mx raw
1 package main
2
3 // Notif Worker - notification subscription, dedup, filtering, watermark tracking.
4 //
5 // supervisor -> worker:
6 // ["N_EVENT", subID, evJSON]
7 // ["N_EOSE", subID]
8 // ["N_SET_PUBKEY", pk]
9 // ["N_SET_RELAYS", urlsJSON]
10 // ["N_SET_MUTES", pksJSON]
11 // ["N_LOAD_MORE"]
12 // ["N_MARK_READ"]
13 // ["N_SET_FILTER", filter] -- "all", "mentions", "reactions", "zaps"
14 // ["N_READ_TS", ts] -- initial notifSt.notifReadTs from IDB
15 //
16 // worker -> supervisor (routed to Shell):
17 // ["N_RENDER", evJSON, mode] -- mode: "prepend" or "append"
18 // ["N_DOT", show] -- show=1 or 0
19 // ["N_EOSE_DONE", subID]
20 // ["N_STATUS", count, notifSt.exhausted]
21 // ["N_SUB", subID, filterJSON, urlsJSON]
22 // ["N_CLOSE", subID]
23 // ["N_FETCH_REFS", idsJSON] -- Shell should fetch these referenced events
24 // ["N_STORE_READ_TS", ts] -- save notif-read-ts to IDB
25
26 import (
27 "runtime"
28 "git.smesh.lol/musiquay/web/common/accum"
29 "git.smesh.lol/musiquay/web/common/helpers"
30 "git.smesh.lol/musiquay/web/common/jsbridge/notif"
31 "git.smesh.lol/musiquay/web/common/mw"
32 "git.smesh.lol/nostr/pkg/core"
33 )
34
35
36 // notifState is the notification
37 // Package globals are immutable outside initialization, so it lives in
38 // one self-mutating type reached through a package-level pointer.
39 type notifState struct {
40 myPK string
41 relayURLs accum.Strings
42 muteSet map[string]bool
43 notifFilter string
44 notifReadTs int64
45 notifSeen map[string]bool
46 oldestTs int64
47 eventCount int32
48 exhausted bool
49 emptyStreak int32
50 initLoad bool
51 loading bool
52 moreGot int32
53 moreTimer int32
54 notifMissing accum.Strings
55 refTimer int32
56 }
57
58 var notifSt *notifState
59
60 func initState() {
61 if notifSt != nil {
62 return
63 }
64 // notifSt and everything it holds live as long as this worker: build them
65 // in the root arena, not in this call's frame arena which dies on
66 // return.
67 runtime.SovereignSetArena(runtime.RootArena())
68 notifSt = ¬ifState{}
69 notifSt.muteSet = map[string]bool{}
70 notifSt.notifSeen = map[string]bool{}
71 notifSt.notifFilter = "all"
72 notifSt.initLoad = true
73 runtime.SovereignRestoreArena(runtime.RootArena())
74 }
75
76 func main() {
77 initState()
78 notif.WorkerOnMessage(handleMessage)
79 }
80
81 func handleMessage(msg string) {
82 w := mw.New(msg)
83 switch w.Str() {
84 case "N_SET_PUBKEY":
85 notifSt.myPK = w.Str()
86 case "N_SET_RELAYS":
87 notifSt.relayURLs.Set(parseStringArray(w.Raw()))
88 case "N_SET_MUTES":
89 pks := parseStringArray(w.Raw())
90 notifSt.muteSet = map[string]bool{}
91 for _, pk := range pks { notifSt.muteSet[pk] = true }
92 case "N_SET_FILTER":
93 notifSt.notifFilter = w.Str()
94 case "N_READ_TS":
95 notifSt.notifReadTs = w.Num()
96 case "N_EVENT":
97 _ = w.Str() // subID consumed
98 evJSON := w.Raw()
99 handleNotifEvent(evJSON)
100 case "N_EOSE":
101 subID := w.Str()
102 handleEOSE(subID)
103 case "N_SUBSCRIBE":
104 handleSubscribe()
105 case "N_LOAD_MORE":
106 handleLoadMore()
107 case "N_MARK_READ":
108 handleMarkRead()
109 }
110 }
111
112 func handleNotifEvent(evJSON string) {
113 ev := nostr.ParseEvent(evJSON)
114 if ev == nil { return }
115 if notifSt.notifSeen == nil { notifSt.notifSeen = map[string]bool{} }
116 if notifSt.notifSeen[ev.ID] { return }
117 notifSt.notifSeen[ev.ID] = true
118 if notifSt.oldestTs == 0 || ev.CreatedAt < notifSt.oldestTs { notifSt.oldestTs = ev.CreatedAt }
119 if ev.PubKey == notifSt.myPK { return }
120 if notifSt.muteSet[notifSenderPK(ev)] { return }
121 if repliesToMuted(ev) { return }
122
123 // Queue missing referenced events
124 if ev.Kind != 1 {
125 refID := notifRefID(ev)
126 if refID != "" {
127 notifSt.notifMissing.Push(refID)
128 }
129 }
130
131 // Check against read watermark for dot display
132 if ev.CreatedAt > notifSt.notifReadTs {
133 notif.WorkerPost(`["N_DOT",1]`)
134 }
135
136 if !passesFilter(ev) { return }
137
138 mode := "append"
139 if !notifSt.initLoad { mode = "prepend" }
140 notif.WorkerPost(`["N_RENDER",` | evJSON | `,"` | mode | `"]`)
141 notifSt.eventCount++
142 }
143
144 func handleEOSE(subID string) {
145 if subID == "ntf" {
146 notifSt.initLoad = false
147 if notifSt.notifMissing.Len() > 0 {
148 notif.WorkerPost(`["N_FETCH_REFS",` | buildIDsJSON() | `]`)
149 notifSt.notifMissing.Reset()
150 }
151 notif.WorkerPost(`["N_EOSE_DONE","ntf"]`)
152 notif.WorkerPost(`["N_STATUS",` | helpers.Itoa(int64(notifSt.eventCount)) | `,0]`)
153 } else if subID == "ntf-more" {
154 if notifSt.moreTimer != 0 { notif.WorkerClearTimeout(notifSt.moreTimer); notifSt.moreTimer = 0 }
155 notifSt.loading = false
156 if notifSt.moreGot == 0 {
157 notifSt.emptyStreak++
158 if notifSt.emptyStreak >= 3 { notifSt.exhausted = true }
159 } else {
160 notifSt.emptyStreak = 0
161 }
162 ex := int64(0)
163 if notifSt.exhausted { ex = 1 }
164 notif.WorkerPost(`["N_STATUS",` | helpers.Itoa(int64(notifSt.eventCount)) | `,` | helpers.Itoa(ex) | `]`)
165 notif.WorkerPost(`["N_EOSE_DONE","ntf-more"]`)
166 notif.WorkerPost(`["N_CLOSE","ntf-more"]`)
167 }
168 }
169
170 func handleSubscribe() {
171 if notifSt.myPK == "" || notifSt.relayURLs.Len() == 0 {
172 return
173 }
174 // Close any existing ntf sub, then open fresh.
175 notif.WorkerPost(`["N_CLOSE","ntf"]`)
176 filter := `{"kinds":[1,6,7,9735],"#p":[` | jstr(notifSt.myPK) | `],"limit":20}`
177 notif.WorkerPost(`["N_SUB","ntf",` | filter | `,` | buildURLsJSON(notifSt.relayURLs.Slice()) | `]`)
178 }
179
180 func handleLoadMore() {
181 if notifSt.loading || notifSt.exhausted || notifSt.oldestTs == 0 || notifSt.myPK == "" { return }
182 notifSt.loading = true
183 notifSt.moreGot = 0
184 filter := `{"kinds":[1,6,7,9735],"#p":[` | jstr(notifSt.myPK) | `],"until":` | helpers.Itoa(notifSt.oldestTs) | `,"limit":20}`
185 notif.WorkerPost(`["N_SUB","ntf-more",` | filter | `,` | buildURLsJSON(notifSt.relayURLs.Slice()) | `]`)
186 notifSt.moreTimer = notif.WorkerSetTimeout(10000, func() {
187 notifSt.moreTimer = 0
188 notifSt.loading = false
189 notif.WorkerPost(`["N_EOSE_DONE","ntf-more"]`)
190 })
191 }
192
193 func handleMarkRead() {
194 ts := notif.WorkerNowSeconds()
195 notifSt.notifReadTs = ts
196 notif.WorkerPost(`["N_STORE_READ_TS",` | helpers.Itoa(ts) | `]`)
197 notif.WorkerPost(`["N_DOT",0]`)
198 }
199
200 func passesFilter(ev *nostr.Event) (ok bool) {
201 switch notifSt.notifFilter {
202 case "mentions":
203 return ev.Kind == 1
204 case "reactions":
205 return ev.Kind == 7
206 case "zaps":
207 return ev.Kind == 9735
208 default:
209 return true
210 }
211 }
212
213 func notifRefID(ev *nostr.Event) (s string) {
214 switch ev.Kind {
215 case 7, 9735, 6:
216 tag := nostr.TagsGetFirst(ev.Tags, []byte("e"))
217 if tag != nil { return string(nostr.TagValue(tag)) }
218 }
219 return ""
220 }
221
222 func notifSenderPK(ev *nostr.Event) (s string) {
223 if ev.Kind == 9735 {
224 pTag := nostr.TagsGetFirst(ev.Tags, []byte("P"))
225 if pTag != nil { return string(nostr.TagValue(pTag)) }
226 }
227 return ev.PubKey
228 }
229
230 func repliesToMuted(ev *nostr.Event) (ok bool) {
231 for _, tag := range nostr.TagsGetAll(ev.Tags, "p") {
232 if v := string(nostr.TagValue(tag)); v != "" && notifSt.muteSet[v] { return true }
233 }
234 return false
235 }
236
237 // buildIDsJSON serialises the pending reference ids. It reads the holder
238 // directly: the ids live in chunks owned by the holder's sovereign arena, and
239 // the only slice the worker ever passes around is the one built here, in this
240 // frame, which dies at return.
241 func buildIDsJSON() (s string) {
242 s = "["
243 for i := int32(0); i < notifSt.notifMissing.Len(); i++ {
244 if i > 0 { s = s | "," }
245 s = s | jstr(notifSt.notifMissing.At(i))
246 }
247 return s | "]"
248 }
249
250 func buildURLsJSON(urls []string) (s string) {
251 s := "["
252 for i, u := range urls {
253 if i > 0 { s = s | "," }
254 s = s | jstr(u)
255 }
256 return s | "]"
257 }
258
259 func parseStringArray(json string) (ss []string) {
260 var out []string
261 i := 0
262 for i < len(json) && json[i] != '[' { i++ }
263 if i >= len(json) { return nil }
264 i++
265 for {
266 for i < len(json) && (json[i] == ' ' || json[i] == ',' || json[i] == '\n') { i++ }
267 if i >= len(json) || json[i] == ']' { break }
268 if json[i] != '"' { break }
269 i++
270 start := i
271 for i < len(json) && json[i] != '"' {
272 if json[i] == '\\' { i++ }
273 i++
274 }
275 if i >= len(json) { break }
276 out = push(out, json[start:i])
277 i++
278 }
279 return out
280 }
281
282 func jstr(s string) (sv string) { return helpers.JsonString(s) }
283