package main // 有界去重表 —— 三个 Node 插件里 `lib/bounded.js` 的 Go 对应物。 // // **不能共用那个文件**(homeagent 是 Go 子进程插件),但要解决的问题完全一样: // `deliveredMails` 是「这封邮件我处理过吗」的记忆,键来自 SSE 事件流, // 而插件跟着 homed 长期活着 —— 邮件数单调增长,键却从来没有出口。 // // # 为什么 Go 侧不做 LRU // // Node 的 `Map` 保证插入顺序,所以那边「删掉再插入」就等于「移到队尾」, // LRU 几乎免费。Go 的 map **不保证遍历顺序**,做 LRU 要额外维护一个链表。 // // 这里不值得:`deliveredMails` 防的两种重复(SSE 重放、心跳与建连之间的窗口) // 都发生在秒到分钟级,先进先出(丢最早插入的)与丢最久未访问的在这个场景下 // 没有可观察的差别。而跨进程、跨天的去重本来就由 `ledger`(落盘,14 天保留期) // 负责,这张表只是同进程内的快速路径。 // // # 为什么不是「攒满就整表清空」 // // 整表清空会在那一刻把**全部**记忆丢掉,于是紧接着到达的 SSE 重放会被当成 // 新邮件全部重投一遍 —— 一次性放大成一批重复投递。FIFO 每次只丢最老的一条, // 而最老的那条恰好是最不可能再出现的。 // maxTrackedMails 是同进程内已投递邮件 id 的记忆上限。 // // 与 Node 侧的 MAX_TRACKED_MAILS 取同一个数:SSE 重放最多回放服务端环形缓冲的 // 500 条事件,一次补拉最多 5 封(catchupLimit)。2000 是三个数量级的余量, // 内存代价约 200KB。 const maxTrackedMails = 2000 // boundedIDSet 是一个带 FIFO 上限的字符串集合。 // // 非并发安全:调用方(plugin.go)已经用 sseMu 保护着它, // 自带一把锁只会让「到底该拿哪把锁」变得含糊。 type boundedIDSet struct { limit int seen map[string]struct{} // order 记录插入顺序,用来知道该丢谁。 // // 用 slice 而不是 container/list:上限只有 2000,切片头部推进的代价 // (一次 append + 一个下标)远小于链表节点的分配开销。 order []string // head 是 order 里第一个仍然有效的下标。丢弃时只推进它,不做 order[1:] —— // 后者每次都要搬移整个底层数组。 head int // evicted 累计淘汰条数,观测用。 evicted int } func newBoundedIDSet(limit int) *boundedIDSet { // 上限非法时回落到 1 而不是 panic:这张表是优化项,配错了应当退化成 // 「只记得最后一条」(多几次重复投递),而不是让插件起不来。 if limit < 1 { limit = 1 } return &boundedIDSet{ limit: limit, seen: make(map[string]struct{}, limit), order: make([]string, 0, limit), } } // has 报告这个 id 是否已经记住过。 func (s *boundedIDSet) has(id string) bool { if s == nil || id == "" { return false } _, ok := s.seen[id] return ok } // add 记住一个 id,并在超限时丢掉最早插入的那些。 // // 返回值是「这次调用**新加入**了吗」—— 调用方常常想在一次操作里同时完成 // 「查重」与「登记」,分两步做需要两次加锁或一段不必要的临界区。 func (s *boundedIDSet) add(id string) bool { if s == nil || id == "" { return false } if _, ok := s.seen[id]; ok { // 已存在时**不**移到队尾:FIFO 语义下位置由首次插入决定。 return false } s.seen[id] = struct{}{} s.order = append(s.order, id) for len(s.order)-s.head > s.limit { oldest := s.order[s.head] s.order[s.head] = "" // 断引用,让字符串可回收 s.head++ delete(s.seen, oldest) s.evicted++ } // 前缀攒到一半以上时压实一次,否则 order 的底层数组会随插入次数无限增长 // —— 那正是这张表本来要修的病,只是换了个地方。 if s.head > 0 && s.head >= len(s.order)/2 { s.compact() } return true } // compact 把 order 重排到从 0 开始,丢掉已淘汰的前缀。 // // 容量固定为 `limit*2` 而不是 `cap(live)+limit`:后者里的 `cap(live)` 是 // 原切片剩下的容量,而 append 会不断扩容 —— 于是每次压实都把上一轮扩大后的 // 容量继承下去,底层数组仍然单调增长(实测灌 1 万条后 cap 到 1130)。 // 那正是这张表本来要修的病,只是从 map 换到了切片上。 // // limit*2 刚好是下一次触发压实的长度(head 走到 limit 时 len == 2*limit), // 于是 append 在两次压实之间不会扩容。 func (s *boundedIDSet) compact() { live := s.order[s.head:] fresh := make([]string, len(live), s.limit*2) copy(fresh, live) s.order = fresh s.head = 0 } // size 是当前记住的条数,观测与测试用。 func (s *boundedIDSet) size() int { if s == nil { return 0 } return len(s.seen) }