Skip to content
Closed
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
63 changes: 63 additions & 0 deletions pkg/mediorum/server/blob_presence_store.go
Original file line number Diff line number Diff line change
Expand Up @@ -276,6 +276,69 @@ func (ss *MediorumServer) savePresenceStore(ctx context.Context, bucket *blob.Bu
// caller must enumerate the bucket instead.
var errPresenceStoreNotReady = errors.New("presence store not ready")

// PresenceStoreStatus is the outcome of the last per-cycle decision about
// where presence comes from, published for the health endpoint.
//
// The rows in blob_presence survive restarts; what does not is the decision to
// read them. Every cycle re-runs presenceStoreReady, and a failed gate falls
// back to a full walk with the reason logged once at cycle start -- so from
// outside, a node that walked looks identical whether the store is disabled,
// was never fully enumerated, or failed its liveness sample. This is the
// answer to "why is it walking again".
type PresenceStoreStatus struct {
// Enabled is the operator setting (OPENAUDIO_PRESENCE_STORE_ENABLED).
Enabled bool `json:"enabled"`
// Used reports whether the current (or most recent) cycle resolves
// presence per batch from the store rather than enumerating buckets.
Used bool `json:"used"`
// Reason says why Used is false. Empty when the store is in use.
Reason string `json:"reason"`
CheckedAt time.Time `json:"checkedAt"`
}

func (ss *MediorumServer) publishPresenceStoreStatus(used bool, reason string) {
ss.presenceStore.Store(&PresenceStoreStatus{
Enabled: ss.Config.PresenceStoreEnabled,
Used: used,
Reason: reason,
CheckedAt: time.Now().UTC(),
})
}

// presenceStoreStatus returns the last published status, or nil before the
// first repair cycle of this process has decided.
func (ss *MediorumServer) presenceStoreStatus() *PresenceStoreStatus {
return ss.presenceStore.Load()
}

// presenceSourceForCycle decides whether a repair cycle reads presence per
// batch from the durable store, and records that decision where an operator
// can see it. Cleanup never qualifies: it is the ground-truth pass, and
// reading a table instead of the filesystem would defeat it.
//
// The fallback is logged at Warn only when the operator turned the store on.
// With it off -- the default -- enumerating is the expected path and a Warn
// every cycle would be noise.
func (ss *MediorumServer) presenceSourceForCycle(ctx context.Context, cleanupMode bool) bool {
if cleanupMode {
ss.publishPresenceStoreStatus(false, "cleanup cycle always enumerates its buckets")
return false
}
if err := ss.presenceStoreReady(ctx); err != nil {
if ss.Config.PresenceStoreEnabled {
ss.logger.Warn("presence store enabled but not usable this cycle; enumerating buckets",
zap.Error(err))
} else {
ss.logger.Debug("presence store not usable this cycle; enumerating buckets",
zap.Error(err))
}
ss.publishPresenceStoreStatus(false, err.Error())
return false
}
ss.publishPresenceStoreStatus(true, "")
return true
}

// presenceStoreReady reports whether every bucket this cycle will consult can
// be served from the store. It is all-or-nothing on purpose: a mixed
// file://-plus-cloud node falls back to enumerating everything, which is what
Expand Down
64 changes: 64 additions & 0 deletions pkg/mediorum/server/blob_presence_store_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -226,3 +226,67 @@ func TestVerifyPresenceStoreLivenessRejectsMissingDir(t *testing.T) {
ss.verifyPresenceStoreLiveness(ctx, ss.bucket, filepath.Join(t.TempDir(), "not-mounted")),
errPresenceStoreNotReady)
}

// The rows survive restarts; the decision to read them is remade every cycle.
// From outside, a node that walked looks the same whatever gate failed, so the
// decision publishes its reason -- and warns only when the operator turned the
// store on, since with it off enumerating is simply the expected path.
func TestPresenceSourceForCycleReportsWhyItWalked(t *testing.T) {
ctx := context.Background()
ss := testNetwork[0]
enablePresenceStore(t, ss)

ss.Config.PresenceStoreEnabled = false
assert.False(t, ss.presenceSourceForCycle(ctx, false))
st := ss.presenceStoreStatus()
require.NotNil(t, st)
assert.False(t, st.Enabled)
assert.False(t, st.Used)
assert.Contains(t, st.Reason, "disabled")
ss.Config.PresenceStoreEnabled = true

assert.False(t, ss.presenceSourceForCycle(ctx, true), "cleanup never reads the store")
st = ss.presenceStoreStatus()
assert.True(t, st.Enabled)
assert.False(t, st.Used)
assert.Contains(t, st.Reason, "cleanup")

assert.False(t, ss.presenceSourceForCycle(ctx, false))
st = ss.presenceStoreStatus()
assert.False(t, st.Used)
assert.Contains(t, st.Reason, "never been walked")

_, err := ss.buildRepairPresenceIndex(ctx)
require.NoError(t, err)
assert.True(t, ss.presenceSourceForCycle(ctx, false),
"a completed enumeration should make the store readable")
st = ss.presenceStoreStatus()
assert.True(t, st.Used)
assert.Empty(t, st.Reason)
assert.False(t, st.CheckedAt.IsZero())
}

// Archive eviction and relocation delete through dropFromBucket. A row left
// behind for a blob that is gone is exactly the drift the liveness sample
// exists to catch, so enough of them would send every later cycle back to a
// full walk.
func TestDropFromBucketForgetsPresence(t *testing.T) {
ctx := context.Background()
ss := testNetwork[0]
enablePresenceStore(t, ss)

const key = "zzz/dropped-key"
ss.recordBlobPresent(ss.bucket, key, 7)
index, err := ss.presenceForCIDs(ctx, []string{key})
require.NoError(t, err)
_, ok := index.Lookup(key, ss.bucket)
require.True(t, ok)

// The blob itself was never written; NotFound is benign for the delete.
require.NoError(t, ss.dropFromBucket(ctx, ss.bucket, key))

index, err = ss.presenceForCIDs(ctx, []string{key})
require.NoError(t, err)
_, ok = index.Lookup(key, ss.bucket)
assert.False(t, ok, "a bucket-scoped delete must forget presence")
}
10 changes: 1 addition & 9 deletions pkg/mediorum/server/repair.go
Original file line number Diff line number Diff line change
Expand Up @@ -273,15 +273,7 @@ func (ss *MediorumServer) runRepair(ctx context.Context, tracker *RepairTracker)
// also what populates the store, so it stays the fallback for every case
// the store cannot serve -- including, by default, all of them.
var cycleIndex *repairPresenceIndex
usePerBatchPresence := false
if !tracker.CleanupMode {
if err := ss.presenceStoreReady(ctx); err != nil {
ss.logger.Debug("presence store not usable this cycle; enumerating buckets",
zap.Error(err))
} else {
usePerBatchPresence = true
}
}
usePerBatchPresence := ss.presenceSourceForCycle(ctx, tracker.CleanupMode)
if usePerBatchPresence {
ss.logger.Info("resolving presence per batch from the durable store")
tracker.Counters["presence_from_store"] = 1
Expand Down
1 change: 1 addition & 0 deletions pkg/mediorum/server/replicate.go
Original file line number Diff line number Diff line change
Expand Up @@ -214,6 +214,7 @@ func (ss *MediorumServer) dropFromBucket(ctx context.Context, b *blob.Bucket, ke
return err
}
ss.knownPresent.Remove(ss.presenceCacheKey(key, b))
ss.forgetBlobPresent(b, key)
return nil
}

Expand Down
5 changes: 5 additions & 0 deletions pkg/mediorum/server/serve_health.go
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,10 @@ type HealthData struct {
// repair's first checkpoint, so /internal/logs/repair still shows the
// previous run as the latest and has no row for the current one.
PresenceWalk *PresenceWalkProgress `json:"presenceWalk"`
// PresenceStore is the last cycle's decision about whether presence came
// from the durable store, with the reason when it did not. Nil until the
// first repair cycle of this process has decided.
PresenceStore *PresenceStoreStatus `json:"presenceStore"`
}

func (ss *MediorumServer) getHealth() HealthData {
Expand Down Expand Up @@ -143,6 +147,7 @@ func (ss *MediorumServer) getHealth() HealthData {
TranscodeQueueLength: len(ss.transcodeWork),
TranscodeStats: ss.getTranscodeStats(),
PresenceWalk: ss.presenceWalkProgress(),
PresenceStore: ss.presenceStoreStatus(),
}
}

Expand Down
4 changes: 4 additions & 0 deletions pkg/mediorum/server/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -268,6 +268,10 @@ type MediorumServer struct {
// first checkpoint, so nothing else reports it while it is happening.
presenceWalk atomic.Pointer[presenceWalkCounter]

// presenceStore is the last per-cycle decision about whether presence is
// read from the durable store, or nil before the first cycle decides.
presenceStore atomic.Pointer[PresenceStoreStatus]

StartedAt time.Time
Config MediorumConfig

Expand Down