main.mx raw
1 package main
2
3 // Profile Worker - owns all author profile data and subscription management.
4 //
5 // Data authority: profSt.authorNames/Pics/Content/Ts/Relays/Follows/Mutes (all profiles).
6 // Shell keeps: profSt.authorNames/Pics view cache (small, for sync render), followSet/muteSet
7 // (sync filter derived sets), myFollows/myMutes (own lists for action operations).
8 //
9 // Wire: supervisor <-> worker
10 // supervisor -> worker:
11 // ["P_RESOLVE", pk] -- request fetch if not already fetched
12 // ["P_RESOLVE_FORCE", pk] -- re-fetch unconditionally
13 // ["P_EVENT", evJSON] -- kind 0/3/10002/10000 event from ap-* or prof sub
14 // ["P_EOSE", subID] -- ap-* EOSE (retry logic)
15 // ["P_SET_RELAYS", urlsJSON]
16 // ["P_SET_PUBKEY", pk]
17 // ["P_IDB_LOADED", pksJSON] -- Shell finished IDB load; mark these pks fetched
18 //
19 // worker -> supervisor (Shell or relay-proxy):
20 // ["P_RESOLVED", pk, name, pic] -- kind 0 resolved
21 // ["P_CONTENT", pk, contentJSON] -- full kind 0 content
22 // ["P_FOLLOWS", pk, pksJSON] -- kind 3 follow list
23 // ["P_MUTES", pk, pksJSON] -- kind 10000 mute list
24 // ["P_RELAYS", pk, urlsJSON] -- kind 10002 relay list
25 // ["P_SUB", subID, filterJSON, urlsJSON]
26 // ["P_CLOSE", subID]
27 // ["P_RETRY_REQ"] -- request Shell's list of still-missing pks
28
29 import (
30 "runtime"
31 "git.smesh.lol/musiquay/web/common/accum"
32 "git.smesh.lol/musiquay/web/common/helpers"
33 "git.smesh.lol/musiquay/web/common/jsbridge/profile"
34 "git.smesh.lol/musiquay/web/common/mw"
35 "git.smesh.lol/nostr/pkg/core"
36 )
37
38
39 const rotateAuthorMax = 500
40
41 // retryMax is the per-pk attempt cap. Each attempt waits sweepIntervalMs.
42 const retryMax = 6
43 const sweepIntervalMs = 15000
44
45 var discoveryRelays []string
46
47 // profileState is the profile
48 // Package globals are immutable outside initialization, so it lives in
49 // one self-mutating type reached through a package-level pointer.
50 type profileState struct {
51 // Operational state (fetch queue management)
52 fetchedK0 map[string]bool
53 fetchedK10k map[string]bool
54 resolvedSet map[string]bool // pk -> kind 0 resolved with a name
55 requestedSet map[string]int32 // pk -> retry attempts so far (capped at retryMax)
56 authorSubPK map[string]string
57 fetchQueue accum.Strings
58 fetchTimer int32
59 retryRound int32
60 retryTimer int32
61 sweepTimer int32
62 subCounter int32
63 relayURLs accum.Strings
64 myPK string
65
66 // Relay discovery (two-gen rotation with author caches)
67 authorRelays map[string][]string
68 authorRelaysOld map[string][]string
69 relayFreq map[string]int32
70 relayFreqOld map[string]int32
71
72 // Profile data (canonical authority - two-generation rotation)
73 authorNames map[string]string
74 authorNamesOld map[string]string
75 authorPics map[string]string
76 authorPicsOld map[string]string
77 authorContent map[string]string
78 authorContentOld map[string]string
79 authorTs map[string]int64
80 authorTsOld map[string]int64
81 authorFollows map[string][]string
82 authorMutes map[string][]string
83 authorCacheCount int32
84 }
85
86 var profSt *profileState
87
88 func initState() {
89 if profSt != nil {
90 return
91 }
92 // profSt and everything it holds live as long as this worker: build them
93 // in the root arena, not in this call's frame arena which dies on
94 // return.
95 runtime.SovereignSetArena(runtime.RootArena())
96 profSt = &profileState{}
97 profSt.fetchedK0 = map[string]bool{}
98 profSt.fetchedK10k = map[string]bool{}
99 profSt.resolvedSet = map[string]bool{}
100 profSt.requestedSet = map[string]int32{}
101 profSt.authorSubPK = map[string]string{}
102 profSt.authorRelays = map[string][]string{}
103 profSt.authorRelaysOld = map[string][]string{}
104 profSt.relayFreq = map[string]int32{}
105 profSt.relayFreqOld = map[string]int32{}
106 profSt.authorNames = map[string]string{}
107 profSt.authorNamesOld = map[string]string{}
108 profSt.authorPics = map[string]string{}
109 profSt.authorPicsOld = map[string]string{}
110 profSt.authorContent = map[string]string{}
111 profSt.authorContentOld = map[string]string{}
112 profSt.authorTs = map[string]int64{}
113 profSt.authorTsOld = map[string]int64{}
114 profSt.authorFollows = map[string][]string{}
115 profSt.authorMutes = map[string][]string{}
116 runtime.SovereignRestoreArena(runtime.RootArena())
117 }
118
119 func initProfileGlobals() {
120 discoveryRelays = []string{
121 "wss://purplepag.es",
122 "wss://relay.damus.io",
123 "wss://nos.lol",
124 }
125 }
126
127 func main() {
128 initProfileGlobals()
129
130 initState()
131 profile.WorkerOnMessage(handleMessage)
132 scheduleSweep()
133 }
134
135 // scheduleSweep arms the periodic retry timer. Every sweepIntervalMs the
136 // worker re-fetches any pks in profSt.requestedSet that aren't in profSt.resolvedSet, until
137 // each pk hits retryMax. This is independent of EOSE-based retry and works
138 // even when subs were closed without producing the expected event.
139 func scheduleSweep() {
140 if profSt.sweepTimer != 0 {
141 profile.WorkerClearTimeout(profSt.sweepTimer)
142 }
143 profSt.sweepTimer = profile.WorkerSetTimeout(int32(sweepIntervalMs), func() {
144 profSt.sweepTimer = 0
145 runSweep()
146 scheduleSweep()
147 })
148 }
149
150 func runSweep() {
151 var missing []string
152 for pk, attempts := range profSt.requestedSet {
153 if profSt.resolvedSet[pk] {
154 continue
155 }
156 if attempts >= retryMax {
157 continue
158 }
159 missing = push(missing, pk)
160 }
161 if len(missing) == 0 {
162 return
163 }
164 for _, pk := range missing {
165 profSt.requestedSet[pk] = profSt.requestedSet[pk] + 1
166 profSt.fetchedK0[pk] = false
167 }
168 proxy := buildProxyURLs(topRelays(8))
169 const batchSize = 100
170 for i := 0; i < len(missing); i += batchSize {
171 end := i + batchSize
172 if end > len(missing) {
173 end = len(missing)
174 }
175 chunk := missing[i:end]
176 authors := buildAuthorsJSON(chunk)
177 limit := helpers.Itoa(int64(len(chunk) * 2))
178 profSt.subCounter++
179 emitSub("ap-sweep-" | helpers.Itoa(int64(profSt.subCounter)),
180 `{"authors":` | authors | `,"kinds":[0,10002],"limit":` | limit | `}`, proxy)
181 profSt.subCounter++
182 emitSub("ap-d-" | helpers.Itoa(int64(profSt.subCounter)),
183 `{"authors":` | authors | `,"kinds":[0],"limit":` | helpers.Itoa(int64(len(chunk))) | `}`, profSt.relayURLs.Slice())
184 }
185 }
186
187 func rotateAuthorCaches() {
188 profSt.authorNamesOld = profSt.authorNames
189 profSt.authorNames = map[string]string{}
190 profSt.authorPicsOld = profSt.authorPics
191 profSt.authorPics = map[string]string{}
192 profSt.authorContentOld = profSt.authorContent
193 profSt.authorContent = map[string]string{}
194 profSt.authorTsOld = profSt.authorTs
195 profSt.authorTs = map[string]int64{}
196 profSt.authorRelaysOld = profSt.authorRelays
197 profSt.authorRelays = map[string][]string{}
198 profSt.relayFreqOld = profSt.relayFreq
199 profSt.relayFreq = map[string]int32{}
200 // Clear fetch tracking maps on rotation. If we kept these populated while
201 // the data caches were thrown out, the worker would refuse to re-fetch
202 // authors whose data was rotated away ("we already fetched X, but the
203 // cache is empty for X" → permanent blank).
204 profSt.fetchedK0 = map[string]bool{}
205 profSt.fetchedK10k = map[string]bool{}
206 profSt.resolvedSet = map[string]bool{}
207 profSt.requestedSet = map[string]int32{}
208 profSt.authorCacheCount = 0
209 }
210
211 func handleMessage(msg string) {
212 w := mw.New(msg)
213 switch w.Str() {
214 case "P_RESOLVE":
215 pk := w.Str()
216 queueProfileFetch(pk)
217 case "P_RESOLVE_FORCE":
218 pk := w.Str()
219 profSt.fetchedK0[pk] = false
220 queueProfileFetch(pk)
221 case "P_EVENT":
222 evJSON := w.Raw()
223 ev := nostr.ParseEvent(evJSON)
224 if ev != nil {
225 applyEvent(ev)
226 }
227 case "P_EOSE":
228 subID := w.Str()
229 handleEOSE(subID)
230 case "P_SET_RELAYS":
231 profSt.relayURLs.Set(parseStringArray(w.Raw()))
232 case "P_SET_PUBKEY":
233 profSt.myPK = w.Str()
234 case "P_IDB_LOADED":
235 // Shell finished loading profiles from musiquay-kv IDB. Mark them fetched AND
236 // resolved so Profile Worker doesn't waste relay bandwidth re-fetching them.
237 pksJSON := w.Raw()
238 pks := parseStringArray(pksJSON)
239 for _, pk := range pks {
240 profSt.fetchedK0[pk] = true
241 profSt.resolvedSet[pk] = true
242 }
243 case "P_RETRY_PROVIDE":
244 // Shell sent list of pending pubkeys without resolved names. Re-queue
245 // them, the periodic sweep + queue flush will pick them up. We bypass
246 // per-pk attempt cap here because this is a user-driven retry signal
247 // (tab switch, refresh button) not an automatic background retry.
248 pksJSON := w.Raw()
249 pks := parseStringArray(pksJSON)
250 for _, pk := range pks {
251 if profSt.resolvedSet[pk] {
252 continue
253 }
254 profSt.fetchedK0[pk] = false
255 if _, ok2 := profSt.requestedSet[pk]; ok2 {
256 profSt.requestedSet[pk] = 0 // reset attempts on user-driven retry
257 } else {
258 profSt.requestedSet[pk] = 0
259 }
260 queueProfileFetch(pk)
261 }
262 }
263 }
264
265 func applyEvent(ev *nostr.Event) {
266 switch ev.Kind {
267 case 0:
268 applyKind0(ev)
269 case 3:
270 applyKind3(ev)
271 case 10000:
272 applyKind10000(ev)
273 case 10002:
274 applyKind10002(ev)
275 }
276 }
277
278 func applyKind0(ev *nostr.Event) {
279 prevTs, hasPrev := profSt.authorTs[ev.PubKey]
280 if hasPrev && ev.CreatedAt <= prevTs {
281 return
282 }
283 if _, exists := profSt.authorContent[ev.PubKey]; !exists {
284 profSt.authorCacheCount++
285 if profSt.authorCacheCount > rotateAuthorMax {
286 rotateAuthorCaches()
287 }
288 }
289 profSt.authorTs[ev.PubKey] = ev.CreatedAt
290 profSt.authorContent[ev.PubKey] = ev.Content
291
292 name := helpers.JsonGetString(ev.Content, "name")
293 if len(name) == 0 {
294 name = helpers.JsonGetString(ev.Content, "display_name")
295 }
296 pic := helpers.JsonGetString(ev.Content, "picture")
297 if len(name) > 0 {
298 profSt.authorNames[ev.PubKey] = name
299 profSt.resolvedSet[ev.PubKey] = true
300 }
301 if len(pic) > 0 {
302 profSt.authorPics[ev.PubKey] = pic
303 }
304
305 // Emit P_RESOLVED so Shell can fill note header placeholders.
306 profile.WorkerPost(`["P_RESOLVED",` | jstr(ev.PubKey) | `,` | jstr(name) | `,` | jstr(pic) | `]`)
307 // Emit P_CONTENT so Shell can update profile page and zap buttons.
308 if len(ev.Content) > 0 {
309 profile.WorkerPost(`["P_CONTENT",` | jstr(ev.PubKey) | `,` | helpers.JsonString(ev.Content) | `]`)
310 }
311 }
312
313 func applyKind3(ev *nostr.Event) {
314 var pks []string
315 for _, tag := range nostr.TagsGetAll(ev.Tags, "p") {
316 if v := nostr.TagValue(tag); v != "" {
317 pks = push(pks, v)
318 }
319 }
320 profSt.authorFollows[ev.PubKey] = pks
321 profile.WorkerPost(`["P_FOLLOWS",` | jstr(ev.PubKey) | `,` | buildAuthorsJSON(pks) | `]`)
322 }
323
324 func applyKind10000(ev *nostr.Event) {
325 var pks []string
326 for _, tag := range nostr.TagsGetAll(ev.Tags, "p") {
327 if v := nostr.TagValue(tag); v != "" {
328 pks = push(pks, v)
329 }
330 }
331 profSt.authorMutes[ev.PubKey] = pks
332 profile.WorkerPost(`["P_MUTES",` | jstr(ev.PubKey) | `,` | buildAuthorsJSON(pks) | `]`)
333 }
334
335 func applyKind10002(ev *nostr.Event) {
336 tags := nostr.TagsGetAll(ev.Tags, "r")
337 if tags == nil {
338 return
339 }
340 var urls []string
341 for _, tag := range tags {
342 u := string(nostr.TagValue(tag))
343 if u != "" {
344 urls = push(urls, u)
345 profSt.relayFreq[u] = profSt.relayFreq[u] + 1
346 }
347 }
348 if len(urls) > 0 {
349 profSt.authorRelays[ev.PubKey] = urls
350 profile.WorkerPost(`["P_RELAYS",` | jstr(ev.PubKey) | `,` | buildURLsJSON(urls) | `]`)
351 }
352 }
353
354 func queueProfileFetch(pk string) {
355 if len(pk) != 64 {
356 return
357 }
358 // Track every requested pk so the periodic sweep can retry until
359 // resolved or per-pk attempt cap is hit.
360 if _, ok := profSt.requestedSet[pk]; !ok {
361 profSt.requestedSet[pk] = 0
362 }
363 if profSt.fetchedK0[pk] {
364 return
365 }
366 profSt.fetchedK0[pk] = true
367 profSt.fetchQueue.Push(pk)
368 if profSt.fetchTimer != 0 {
369 profile.WorkerClearTimeout(profSt.fetchTimer)
370 }
371 profSt.fetchTimer = profile.WorkerSetTimeout(300, func() {
372 profSt.fetchTimer = 0
373 flushFetchQueue()
374 })
375 }
376
377 func flushFetchQueue() {
378 if profSt.fetchQueue.Len() == 0 {
379 return
380 }
381 queue := profSt.fetchQueue.Slice()
382 profSt.fetchQueue.Reset()
383 proxy := buildProxyURLs(nil)
384 const batchSize = 100
385 for i := 0; i < len(queue); i += batchSize {
386 end := i + batchSize
387 if end > len(queue) {
388 end = len(queue)
389 }
390 chunk := queue[i:end]
391 authors := buildAuthorsJSON(chunk)
392 limit := helpers.Itoa(int64(len(chunk)))
393 profSt.subCounter++
394 emitSub("ap-batch-q-" | helpers.Itoa(int64(profSt.subCounter)),
395 `{"authors":` | authors | `,"kinds":[0],"limit":` | limit | `}`, proxy)
396 profSt.subCounter++
397 emitSub("ap-d-" | helpers.Itoa(int64(profSt.subCounter)),
398 `{"authors":` | authors | `,"kinds":[0],"limit":` | limit | `}`, profSt.relayURLs.Slice())
399 }
400 }
401
402 func fetchAuthorProfile(pk string) {
403 if profSt.fetchedK0[pk] {
404 return
405 }
406 profSt.fetchedK0[pk] = true
407 profSt.subCounter++
408 subID := "ap-" | helpers.Itoa(int64(profSt.subCounter))
409 profSt.authorSubPK[subID] = pk
410 emitSub(subID, `{"authors":[` | jstr(pk) | `],"kinds":[0,3,10002,10000],"limit":6}`, buildProxy(pk))
411 }
412
413 func handleEOSE(subID string) {
414 // "ap-batch-" prefix matches both "ap-batch-q-N" (queue flush) and
415 // "ap-batch-N" (legacy P_RETRY_PROVIDE re-batch). The periodic sweep
416 // now handles per-pk retries with its own cap, so EOSE-triggered retry
417 // is no longer needed. Keeping it here would double-fetch and the
418 // global profSt.retryRound cap caused unresolved pks to permanently stick.
419 if len(subID) > 9 && subID[:9] == "ap-batch-" {
420 return
421 }
422 if len(subID) > 9 && subID[:9] == "ap-sweep-" {
423 return
424 }
425 if len(subID) > 3 && subID[:3] == "ap-" {
426 pk, ok := profSt.authorSubPK[subID]
427 if !ok {
428 return
429 }
430 delete(profSt.authorSubPK, subID)
431 if !profSt.resolvedSet[pk] {
432 rels, ok2 := profSt.authorRelays[pk]
433 if !ok2 {
434 rels, ok2 = profSt.authorRelaysOld[pk]
435 }
436 if ok2 && len(rels) > 0 && !profSt.fetchedK10k[pk] {
437 profSt.fetchedK10k[pk] = true
438 profSt.fetchedK0[pk] = false
439 fetchAuthorProfile(pk)
440 }
441 }
442 }
443 }
444
445 func emitSub(subID, filterJSON string, urls []string) {
446 profile.WorkerPost(`["P_SUB",` | jstr(subID) | `,` | filterJSON | `,` | buildURLsJSON(urls) | `]`)
447 }
448
449 func buildProxy(pk string) (ss []string) {
450 out := []string{:len(discoveryRelays)}
451 copy(out, discoveryRelays)
452 for i := int32(0); i < profSt.relayURLs.Len(); i++ {
453 out = appendUnique(out, profSt.relayURLs.At(i))
454 }
455 if rels, ok2 := profSt.authorRelays[pk]; ok2 {
456 for _, r := range rels {
457 out = appendUnique(out, r)
458 }
459 } else if oldRels, oldOK := profSt.authorRelaysOld[pk]; oldOK {
460 for _, r := range oldRels {
461 out = appendUnique(out, r)
462 }
463 }
464 for _, r := range topRelays(4) {
465 out = appendUnique(out, r)
466 }
467 return out
468 }
469
470 func buildProxyURLs(extra []string) (ss []string) {
471 out := []string{:len(discoveryRelays)}
472 copy(out, discoveryRelays)
473 for i := int32(0); i < profSt.relayURLs.Len(); i++ {
474 out = appendUnique(out, profSt.relayURLs.At(i))
475 }
476 for _, u := range extra {
477 out = appendUnique(out, u)
478 }
479 return out
480 }
481
482 func topRelays(n int32) (ss []string) {
483 type kv struct{ url string; count int32 }
484 var all []kv
485 for url, count := range profSt.relayFreq {
486 all = push(all, kv{url, count})
487 }
488 for i := 0; i < len(all); i++ {
489 for j := i + 1; j < len(all); j++ {
490 if all[j].count > all[i].count {
491 all[i], all[j] = all[j], all[i]
492 }
493 }
494 }
495 var result []string
496 for i := 0; i < len(all) && i < n; i++ {
497 result = push(result, all[i].url)
498 }
499 return result
500 }
501
502 func appendUnique(list []string, val string) (ss []string) {
503 for _, v := range list {
504 if v == val {
505 return list
506 }
507 }
508 return push(list, val)
509 }
510
511 func buildAuthorsJSON(pks []string) (s string) {
512 s := "["
513 for i, pk := range pks {
514 if i > 0 { s = s | "," }
515 s = s | jstr(pk)
516 }
517 return s | "]"
518 }
519
520 func buildURLsJSON(urls []string) (s string) {
521 s := "["
522 for i, u := range urls {
523 if i > 0 { s = s | "," }
524 s = s | jstr(u)
525 }
526 return s | "]"
527 }
528
529 func parseStringArray(json string) (ss []string) {
530 var out []string
531 i := 0
532 for i < len(json) && json[i] != '[' { i++ }
533 if i >= len(json) { return nil }
534 i++
535 for {
536 for i < len(json) && (json[i] == ' ' || json[i] == ',' || json[i] == '\n') { i++ }
537 if i >= len(json) || json[i] == ']' { break }
538 if json[i] != '"' { break }
539 i++
540 start := i
541 for i < len(json) && json[i] != '"' {
542 if json[i] == '\\' { i++ }
543 i++
544 }
545 if i >= len(json) { break }
546 out = push(out, json[start:i])
547 i++
548 }
549 return out
550 }
551
552 func jstr(s string) (sv string) { return helpers.JsonString(s) }
553
554