Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 11 additions & 8 deletions services/nvpair-cluster-manager/truststore.go
Original file line number Diff line number Diff line change
Expand Up @@ -207,8 +207,8 @@ func (ts *TrustStore) Pin(pin *TrustedPin) error {
}
// Identical re-pin: fold in any newly-seen endorsements so the
// gossiped trust web thickens (idempotent on the cert itself).
err := ts.mergeEndorsementsLocked(pin.NodeUUID, pin.Endorsements)
changed = err == nil
merged, err := ts.mergeEndorsementsLocked(pin.NodeUUID, pin.Endorsements)
changed = merged
return err
}
return fmt.Errorf("uuid %s already pinned to a different certificate; remove it first to re-pin", pin.NodeUUID)
Expand Down Expand Up @@ -239,10 +239,10 @@ func (ts *TrustStore) writePinLocked(pin *TrustedPin) error {
// mergeEndorsementsLocked unions incoming endorsements into an existing pin
// (dedup by signer+signature) and persists if anything new was added. Caller
// holds ts.mu. A missing pin is a silent no-op.
func (ts *TrustStore) mergeEndorsementsLocked(uuid string, incoming []Endorsement) error {
func (ts *TrustStore) mergeEndorsementsLocked(uuid string, incoming []Endorsement) (bool, error) {
pin, ok := ts.pins[uuid]
if !ok || len(incoming) == 0 {
return nil
return false, nil
}
seen := make(map[string]struct{}, len(pin.Endorsements))
for _, e := range pin.Endorsements {
Expand All @@ -260,9 +260,12 @@ func (ts *TrustStore) mergeEndorsementsLocked(uuid string, incoming []Endorsemen
added = true
}
if !added {
return nil
return false, nil
}
return ts.writePinLocked(updated)
if err := ts.writePinLocked(updated); err != nil {
return false, err
}
return true, nil
}

func endorsementKey(e Endorsement) string {
Expand All @@ -279,8 +282,8 @@ func (ts *TrustStore) AddEndorsements(uuid string, endorsements []Endorsement) e
defer ts.announce(&changed)
ts.mu.Lock()
defer ts.mu.Unlock()
err := ts.mergeEndorsementsLocked(uuid, endorsements)
changed = err == nil
merged, err := ts.mergeEndorsementsLocked(uuid, endorsements)
changed = merged
return err
}

Expand Down
228 changes: 221 additions & 7 deletions services/nvpair-cluster-manager/truststore_announce_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,13 @@
package main

import (
"bytes"
"errors"
"os"
"path/filepath"
"reflect"
"sync"
"sync/atomic"
"testing"
"time"
)
Expand Down Expand Up @@ -44,8 +51,8 @@ func testPin(t *testing.T, uuid string) *TrustedPin {
// announcement exists to prevent — so the hook lives on the store rather than at
// the ~19 call sites that pin and unpin peers.
//
// Pinning, removing, forgetting, and a display-name update all mutate what is on
// disk and must each announce exactly once.
// Pinning, removing, and a display-name update mutate disk state; forgetting
// intentionally changes only live authorization. Each must announce once.
func TestTrustStoreAnnouncesEveryMutation(t *testing.T) {
ts, count := newAnnouncingStore(t)
const uuid = "principal-peer"
Expand Down Expand Up @@ -80,6 +87,210 @@ func TestTrustStoreAnnouncesEveryMutation(t *testing.T) {
}
}

// assertStoredEndorsements checks the live snapshot and a separately loaded
// store. It is also safe to call from an onChange callback in a joined worker.
func assertStoredEndorsements(t *testing.T, ts *TrustStore, uuid string, want []Endorsement) {
t.Helper()
pin, ok := ts.Get(uuid)
if !ok || !reflect.DeepEqual(pin.Endorsements, want) {
t.Errorf("live endorsements = %+v, want %+v", pin, want)
}
reloaded, err := newTrustStore(filepath.Dir(ts.dir))
if err != nil {
t.Errorf("reload trust store: %v", err)
return
}
pin, ok = reloaded.Get(uuid)
if !ok || !reflect.DeepEqual(pin.Endorsements, want) {
t.Errorf("reloaded endorsements = %+v, want %+v", pin, want)
}
}

func TestTrustStoreAnnouncesNewEndorsementsAfterPersistence(t *testing.T) {
for _, identicalPin := range []bool{false, true} {
name := "AddEndorsements"
if identicalPin {
name = "IdenticalPin"
}
t.Run(name, func(t *testing.T) {
ts, _ := newAnnouncingStore(t)
pin := testPin(t, "principal-peer")
first := Endorsement{By: "trusted-peer", Sig: "signature-1"}
second := Endorsement{By: "trusted-peer", SigV2: "signature-2"}
pin.Endorsements = []Endorsement{first}
if err := ts.Pin(pin); err != nil {
t.Fatalf("pin: %v", err)
}
want := []Endorsement{first, second}
calls := 0
ts.SetOnChange(func() {
calls++
// Acquiring Get's read lock here also witnesses that the
// mutation lock was released before announcing the change.
assertStoredEndorsements(t, ts, pin.NodeUUID, want)
})
merge := func(batch []Endorsement) error {
if identicalPin {
updated := *pin
updated.Endorsements = batch
return ts.Pin(&updated)
}
return ts.AddEndorsements(pin.NodeUUID, batch)
}
// Mix an existing endorsement, a new endorsement, and an
// in-batch duplicate. One operation causes one announcement.
batch := []Endorsement{first, second, second}
if err := merge(batch); err != nil {
t.Fatalf("merge: %v", err)
}
if calls != 1 {
t.Fatalf("announcements after merge = %d, want 1", calls)
}
assertStoredEndorsements(t, ts, pin.NodeUUID, want)
before, err := os.ReadFile(ts.pinPath(pin.NodeUUID))
if err != nil {
t.Fatal(err)
}
for _, noOp := range [][]Endorsement{batch, nil} {
if err := merge(noOp); err != nil {
t.Fatalf("no-op merge: %v", err)
}
if calls != 1 {
t.Fatalf("announcements after no-op = %d, want 1", calls)
}
}
after, err := os.ReadFile(ts.pinPath(pin.NodeUUID))
if err != nil || !bytes.Equal(before, after) {
t.Fatalf("no-op changed disk contents: %v", err)
}
assertStoredEndorsements(t, ts, pin.NodeUUID, want)
})
}
}

func TestTrustStoreStaysSilentWhenEndorsementWriteFails(t *testing.T) {
for _, identicalPin := range []bool{false, true} {
name := "AddEndorsements"
if identicalPin {
name = "IdenticalPin"
}
t.Run(name, func(t *testing.T) {
ts, count := newAnnouncingStore(t)
pin := testPin(t, "principal-peer")
first := Endorsement{By: "trusted-peer", Sig: "signature-1"}
second := Endorsement{By: "trusted-peer", SigV2: "signature-2"}
pin.Endorsements = []Endorsement{first}
if err := ts.Pin(pin); err != nil {
t.Fatalf("pin: %v", err)
}
before, err := os.ReadFile(ts.pinPath(pin.NodeUUID))
if err != nil {
t.Fatal(err)
}
beforeCount := count()
// Fail the final replace, after the temporary file was written.
// The existing pin must remain intact on disk and in memory.
originalRename := renameFile
writeErr := errors.New("injected endorsement replace failure")
renameFile = func(_, _ string) error { return writeErr }
t.Cleanup(func() { renameFile = originalRename })
merge := func() error {
if identicalPin {
updated := *pin
updated.Endorsements = []Endorsement{second}
return ts.Pin(&updated)
}
return ts.AddEndorsements(pin.NodeUUID, []Endorsement{second})
}
if err := merge(); !errors.Is(err, writeErr) {
t.Fatalf("merge error = %v, want injected replace failure", err)
}
if count() != beforeCount {
t.Fatalf("announcements after failed write = %d, want %d", count(), beforeCount)
}
after, err := os.ReadFile(ts.pinPath(pin.NodeUUID))
if err != nil || !bytes.Equal(before, after) {
t.Fatalf("failed write changed disk contents: %v", err)
}
assertStoredEndorsements(t, ts, pin.NodeUUID, []Endorsement{first})
entries, err := os.ReadDir(ts.dir)
if err != nil || len(entries) != 1 || entries[0].Name() != pin.NodeUUID+".json" {
t.Fatalf("failed write left temporary residue: entries=%v err=%v", entries, err)
}
renameFile = originalRename
if err := merge(); err != nil {
t.Fatalf("retry after storage recovery: %v", err)
}
if count() != beforeCount+1 {
t.Fatalf("announcements after retry = %d, want %d", count(), beforeCount+1)
}
assertStoredEndorsements(t, ts, pin.NodeUUID, []Endorsement{first, second})
})
}
}

func TestTrustStoreConcurrentDuplicateEndorsementsAnnounceOnce(t *testing.T) {
ts, err := newTrustStore(t.TempDir())
if err != nil {
t.Fatal(err)
}
pin := testPin(t, "principal-peer")
if err := ts.Pin(pin); err != nil {
t.Fatal(err)
}
endorsement := Endorsement{By: "trusted-peer", SigV2: "signature-1"}
want := []Endorsement{endorsement}
var calls atomic.Int32
ts.SetOnChange(func() {
calls.Add(1)
assertStoredEndorsements(t, ts, pin.NodeUUID, want)
})
const workers = 16
start := make(chan struct{})
errs := make(chan error, workers)
var wg sync.WaitGroup
for i := range workers {
wg.Add(1)
go func() {
defer wg.Done()
<-start
if i%2 == 0 {
updated := *pin
updated.Endorsements = want
errs <- ts.Pin(&updated)
} else {
errs <- ts.AddEndorsements(pin.NodeUUID, want)
}
}()
}
close(start)
wg.Wait()
close(errs)
for err := range errs {
if err != nil {
t.Errorf("concurrent merge: %v", err)
}
}
if got := calls.Load(); got != 1 {
t.Errorf("announcements for concurrent identical submissions = %d, want 1", got)
}
assertStoredEndorsements(t, ts, pin.NodeUUID, want)
}

func TestTrustStoreMissingEndorsementTargetStaysSilent(t *testing.T) {
ts, count := newAnnouncingStore(t)
if err := ts.AddEndorsements("principal-stranger", []Endorsement{{By: "trusted-peer", SigV2: "signature-1"}}); err != nil {
t.Fatal(err)
}
if count() != 0 || len(ts.List()) != 0 {
t.Fatalf("missing-target merge changed live state: announcements=%d pins=%v", count(), ts.List())
}
entries, err := os.ReadDir(ts.dir)
if err != nil || len(entries) != 0 {
t.Fatalf("missing-target merge changed disk state: entries=%v err=%v", entries, err)
}
}

// TestTrustStoreStaysSilentWhenNothingChanged keeps the announcement meaningful.
// The scanner answers it by walking its whole directory, and the broker relays
// it, so a store that announced on every call — including the idempotent re-pin
Expand All @@ -96,11 +307,14 @@ func TestTrustStoreStaysSilentWhenNothingChanged(t *testing.T) {
before := count()

// An identical re-pin folds in no new endorsements and rewrites nothing.
if err := ts.Pin(testPin(t, uuid)); err == nil {
// A fresh leaf for the same uuid is a DIFFERENT certificate, which the
// store refuses rather than silently re-pinning; either way it must not
// announce a change it did not make.
t.Log("re-pin with a new certificate was accepted")
// Reuse the exact pin: generating another fixture would mint a different
// certificate and exercise the key-rotation rejection path instead.
if err := ts.Pin(pin); err != nil {
t.Fatalf("identical re-pin: %v", err)
}
// An empty endorsement merge is also a no-op and must stay silent.
if err := ts.AddEndorsements(uuid, nil); err != nil {
t.Fatalf("empty endorsement merge: %v", err)
}
// A rename to the values already stored changes nothing.
if ok, err := ts.UpdateIdentity(uuid, uuid, uuid); err != nil || ok {
Expand Down
2 changes: 1 addition & 1 deletion services/versions.json
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@
"nvpair-node-settings": "1.0.4",
"nvpair-ui-broker": "0.40.2",
"nvpair-engine-manager": "0.17.4",
"nvpair-cluster-manager": "1.1.4",
"nvpair-cluster-manager": "1.1.5",
"nvpair-job-scheduler": "0.4.1",
"nvpair-tui": "0.7.2"
}
Expand Down