Skip to content
5 changes: 4 additions & 1 deletion desktop/scripts/verify-service-contracts.ts
Original file line number Diff line number Diff line change
Expand Up @@ -136,7 +136,10 @@ const QUOTED_RE = /"([^"]+)"/g
const IDENT_RE = /\b([A-Za-z_]\w*)\b/g
const NOTIFY_LITERAL_RE = /(?:notify|emit|publish|broadcast)\(\s*"([^"]+)"/gi
const NOTIFY_PREFIX_CONCAT_RE = /(?:notify|emit|publish|broadcast)\(\s*"([^"]+)"\s*\+/i
const NOTIFY_IDENT_RE = /(?:notify|emit|publish|broadcast)\(\s*([A-Za-z_]\w*)\s*,/gi
// Broadcast(ctx, frame) is HTTP fan-out, not a JSON-RPC notification. Keep
// literal Broadcast methods above, but do not treat an identifier argument to
// Broadcast as an unresolved JSON-RPC method.
const NOTIFY_IDENT_RE = /(?:notify|emit|publish)\(\s*([A-Za-z_]\w*)\s*,/gi

/** A `| `nvpair-ui-broker` | 0.37.0 | … |` row — forbidden in hand-maintained docs. */
const DOC_COMPONENT_ROW_RE = /^\|\s*`([a-z0-9-]+)`\s*\|\s*(\d+\.\d+\.\d+)\s*\|/gm
Expand Down
16 changes: 12 additions & 4 deletions services/nvpair-workload-manager/cluster_mtls_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -81,17 +81,25 @@ func setupNode(t *testing.T, certPEM, keyPEM []byte, pins map[string][]byte) str
return dir
}

// newPinnedPeerDirs builds two mutually pinned nodes. Tests that construct a
// Manager need selfDir because NewManager opens its own mesh from that path.
func newPinnedPeerDirs(t *testing.T) (selfDir, peerDir string) {
t.Helper()
selfCert, selfKey := genLeaf(t, "uuid-self")
peerCert, peerKey := genLeaf(t, "uuid-peer")
selfDir = setupNode(t, selfCert, selfKey, map[string][]byte{"uuid-peer": peerCert})
peerDir = setupNode(t, peerCert, peerKey, map[string][]byte{"uuid-self": selfCert})
return selfDir, peerDir
}

// newPinnedPeerMeshes builds the two sides of a minimal two-node cluster: self
// pins peer and peer pins self, so a test can drive the real inter-node path
// (cluster mTLS from a pinned caller) rather than a plaintext shortcut. There is
// no plaintext shortcut to take — the interface is mTLS unconditionally — so every
// receiver test goes through here.
func newPinnedPeerMeshes(t *testing.T) (self, peer *clustertrust.Mesh) {
t.Helper()
selfCert, selfKey := genLeaf(t, "uuid-self")
peerCert, peerKey := genLeaf(t, "uuid-peer")
selfDir := setupNode(t, selfCert, selfKey, map[string][]byte{"uuid-peer": peerCert})
peerDir := setupNode(t, peerCert, peerKey, map[string][]byte{"uuid-self": selfCert})
selfDir, peerDir := newPinnedPeerDirs(t)
self, peer = clustertrust.Open(selfDir), clustertrust.Open(peerDir)
if !self.Clustered() || !peer.Clustered() {
t.Fatal("a dir holding a keypair and a pin must read as clustered")
Expand Down
61 changes: 47 additions & 14 deletions services/nvpair-workload-manager/dedup.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,10 +13,10 @@ import (
// session-scoped volume at ~dozen-node scale (spec §5).
const defaultDedupCapacity = 10000

// dedupIndex is a bounded LRU set of keys. It answers a single question:
// "have I seen this key before?" and records it if not, evicting the
// least-recently-seen key once capacity is exceeded. It is safe for
// concurrent use — the inter-node HTTP handler runs one goroutine per
// dedupIndex is a bounded LRU set of successfully emitted keys. It records a
// key only after the broker emit succeeds, and serializes concurrent emits for
// the same key. Completed keys are evicted least-recently-used first. It is
// safe for concurrent use — the inter-node HTTP handler runs one goroutine per
// request.
//
// Keys are opaque strings built by the caller: lifecycle events key on
Expand All @@ -33,6 +33,7 @@ type dedupIndex struct {
capacity int
ll *list.List // front = most recently seen
items map[string]*list.Element // key -> element in ll
inFlight map[string]chan struct{} // one broker emit at a time per key
}

func newDedupIndex(capacity int) *dedupIndex {
Expand All @@ -43,21 +44,54 @@ func newDedupIndex(capacity int) *dedupIndex {
capacity: capacity,
ll: list.New(),
items: make(map[string]*list.Element, capacity),
inFlight: make(map[string]chan struct{}),
}
}

// seenOrAdd returns true if the key was already present (a duplicate). On a
// first sighting it records the key and returns false. Either way the key is
// promoted to most-recently-seen.
func (d *dedupIndex) seenOrAdd(key string) bool {
d.mu.Lock()
defer d.mu.Unlock()
// emitOnce serializes the check, broker emit, and record for one key. Other
// keys can emit concurrently. If emit fails, the key stays absent and a waiting
// request can retry it. A completed key is reported as a duplicate.
func (d *dedupIndex) emitOnce(key string, emit func() error) (bool, error) {
for {
d.mu.Lock()
if el, ok := d.items[key]; ok {
d.ll.MoveToFront(el)
d.mu.Unlock()
return true, nil
}
if done, ok := d.inFlight[key]; ok {
d.mu.Unlock()
<-done
continue
}
done := make(chan struct{})
d.inFlight[key] = done
d.mu.Unlock()

if el, ok := d.items[key]; ok {
d.ll.MoveToFront(el)
return true
// Release waiters even if an emitter panics and net/http recovers the
// request. A panicking emit has not completed successfully.
err := func() (err error) {
completed := false
defer func() {
d.mu.Lock()
if completed && err == nil {
d.addLocked(key)
}
delete(d.inFlight, key)
close(done)
d.mu.Unlock()
}()
err = emit()
completed = true
return err
}()
return false, err
}
}

// addLocked records a new key and evicts the least-recently-seen key past
// capacity. The caller holds d.mu and has checked that key is absent.
func (d *dedupIndex) addLocked(key string) {
el := d.ll.PushFront(key)
d.items[key] = el
if d.ll.Len() > d.capacity {
Expand All @@ -67,7 +101,6 @@ func (d *dedupIndex) seenOrAdd(key string) bool {
delete(d.items, oldest.Value.(string))
}
}
return false
}

// keyLifecycle builds the dedup key for a lifecycle event. Workload.id is only
Expand Down
Loading
Loading