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