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 = &notifState{}
  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