From d4cfae95605293572cf8ed4bb953ffa678c9a57e Mon Sep 17 00:00:00 2001 From: vchaindz Date: Fri, 4 Sep 2026 19:37:28 +0200 Subject: [PATCH 1/2] fix(database): drop the db-less ref left behind by a failed open allocDB caches a placeholder &dbRef{count: 1} before openDB runs. When the open failed, Get called Release, which only decremented the refcount and left a dbRef{db: nil, count: 0} in dbCache for the life of the process. That single entry poisoned every reader that consults the cache: - GetState returned "unable to get state" forever and never reached the m.Get fallback below it that would have retried the open; - IsActive reported a database that never opened as active; - GetOptionsByIndex returned a nil *Options, which lazyDB.Path() and lazyDB.MaxResultSize() dereference; - CloseAll needed a nil guard (bfa78b34) to avoid panicking on shutdown. A transient 403 from S3 during the first open of two databases was therefore terminal: the client state path (CurrentState -> lazyDB.CurrentState -> GetState) short-circuits before m.Get, so nothing re-drove the open and no further storage traffic was attempted for over three hours until a restart. Release now pops the ref when the open failed and nobody picked it up, so a dbRef with a nil db is never observable. The pop lives in Release rather than a parallel cleanup path so that the pendingClose hand-off still runs when Close races a failed open; without that, a parked ref is orphaned and its database never closed. Get also unlocks db.mtx before Release acquires m.mtx. allocDB takes m.mtx then db.mtx and returns holding db.mtx, so the old ordering inverted the two: a second caller inside allocDB for the same index held m.mtx and blocked on db.mtx while the failing open held db.mtx and waited on m.mtx. Since allocDB takes m.mtx first, that wedged Get for every database, not just the failing one. The refcount always reaches zero on this path, so the inversion was reached on every failed open; only a concurrent caller was needed to close the cycle. TestDBManagerFailedOpenDoesNotDeadlock hangs without this change. GetState now falls through to m.Get instead of the sentinel error, surfacing the real storage error and re-driving the open, rate-limited by openRetryInterval so a permanently broken database is not hammered. A recorded failure outranks any state cached by a previous successful open, otherwise a database that fails to reopen would serve a stale state and never retry. GetOptionsByIndex collapses to dbInfo.getOptions(): both of its branches returned the same options, and the only distinct effect of the cache branch was the nil. OpenDB and NewDB additionally leaked the store on the sql.NewEngine, document.NewEngine and NewTxHolderPool error paths, where only InitIndexing closed it. A leaked store would make the retries this change enables fail. Fixes #2155 Signed-off-by: vchaindz --- pkg/database/database.go | 7 +- pkg/database/db_manager.go | 182 ++++++++++++++------ pkg/database/db_manager_test.go | 286 ++++++++++++++++++++++++++++++-- 3 files changed, 416 insertions(+), 59 deletions(-) diff --git a/pkg/database/database.go b/pkg/database/database.go index 7d68da4734..6e9b47354f 100644 --- a/pkg/database/database.go +++ b/pkg/database/database.go @@ -242,6 +242,7 @@ func OpenDB( dbi.sqlEngine, err = sql.NewEngine(dbi.st, sqlOpts) if err != nil { + dbi.st.Close() dbi.Logger.Errorf("unable to load sql-engine for database '%s' {replica = %v}. %v", dbName, opts.replica, err) return nil, err } @@ -250,12 +251,14 @@ func OpenDB( dbi.documentEngine, err = document.NewEngine(dbi.st, document.DefaultOptions().WithPrefix([]byte{DocumentPrefix})) if err != nil { - return nil, err + dbi.st.Close() + return nil, logErr(dbi.Logger, "unable to load document-engine: %s", err) } dbi.Logger.Infof("document-engine ready for database '%s' {replica = %v}", dbName, opts.replica) txPool, err := dbi.st.NewTxHolderPool(opts.readTxPoolSize, false) if err != nil { + dbi.st.Close() return nil, logErr(dbi.Logger, "unable to create tx pool: %s", err) } dbi.txPool = txPool @@ -369,6 +372,7 @@ func NewDB(dbName string, multidbHandler sql.MultiDBHandler, opts *Options, log dbi.sqlEngine, err = sql.NewEngine(dbi.st, sqlOpts) if err != nil { + dbi.st.Close() dbi.Logger.Errorf("unable to load sql-engine for database '%s' {replica = %v}. %v", dbName, opts.replica, err) return nil, err } @@ -376,6 +380,7 @@ func NewDB(dbName string, multidbHandler sql.MultiDBHandler, opts *Options, log dbi.documentEngine, err = document.NewEngine(dbi.st, document.DefaultOptions().WithPrefix([]byte{DocumentPrefix})) if err != nil { + dbi.st.Close() return nil, logErr(dbi.Logger, "Unable to open database: %s", err) } dbi.Logger.Infof("document-engine ready for database '%s' {replica = %v}", dbName, opts.replica) diff --git a/pkg/database/db_manager.go b/pkg/database/db_manager.go index d0549cf875..c8bece1e83 100644 --- a/pkg/database/db_manager.go +++ b/pkg/database/db_manager.go @@ -18,7 +18,6 @@ package database import ( "context" - "fmt" "sync" "sync/atomic" "time" @@ -39,13 +38,32 @@ type DBManager struct { databases []*dbInfo dbIndex map[string]int - mtx sync.Mutex - waitCond *sync.Cond + mtx sync.Mutex + waitCond *sync.Cond pendingClose map[int]*dbRef // refs removed from cache but still in use closed bool } +// openRetryInterval bounds how often GetState re-drives a failed open. Without +// it, every CurrentState call on a permanently broken database would hammer the +// underlying (possibly remote) storage. +const openRetryInterval = 5 * time.Second + +// openFailure records a failed attempt to open a database. +type openFailure struct { + err error + at time.Time +} + +// DatabaseOpenFailure reports a database whose most recent open attempt failed +// and which has not been opened successfully since. +type DatabaseOpenFailure struct { + Name string + Err error + At time.Time +} + type dbInfo struct { mtx sync.Mutex @@ -55,6 +73,23 @@ type dbInfo struct { name string deleted bool closed bool + + // lastOpenFail is read by readiness probes while an open is in flight, so it + // is accessed atomically rather than under mtx: mtx is held for the entire + // duration of an open, which is exactly when a probe wants an answer. + lastOpenFail atomic.Pointer[openFailure] +} + +func (db *dbInfo) recordOpenFailure(err error) { + db.lastOpenFail.Store(&openFailure{err: err, at: time.Now()}) +} + +func (db *dbInfo) clearOpenFailure() { + db.lastOpenFail.Store(nil) +} + +func (db *dbInfo) openFailure() *openFailure { + return db.lastOpenFail.Load() } func (db *dbInfo) cacheInfo(s *schema.ImmutableState, opts *Options) { @@ -87,6 +122,7 @@ func (db *dbInfo) close() error { return store.ErrAlreadyClosed } db.closed = true + db.clearOpenFailure() return nil } @@ -94,6 +130,11 @@ func (db *dbInfo) close() error { type dbRef struct { db DB count uint32 + + // opened is set to 1 once db has been assigned. It lets readers determine + // whether the ref backs a real database without reading db itself, which is + // only safe under the owning dbInfo's mutex. + opened uint32 } type OpenDBFunc func(name string, opts *Options) (DB, error) @@ -131,7 +172,9 @@ func createCache(m *DBManager, capacity int) *cache.Cache { // Moreover, since the reference cannot be altered after it has been set, // there is not need to acquire the database lock. if ref.db == nil { - m.logger.Errorf("db not initialised during eviction") + // Benign: a ref whose open failed can be evicted by another database's + // allocDB before Release pops it. + m.logger.Debugf("db not initialised during eviction") return } @@ -151,6 +194,7 @@ func createCache(m *DBManager, capacity int) *cache.Cache { m.databases[i].cacheInfo(state, opts) } ref.db = nil + atomic.StoreUint32(&ref.opened, 0) }) return c } @@ -164,6 +208,7 @@ func (m *DBManager) Put(dbName string, opts *Options, closed bool) int { ref.deleted = false ref.closed = closed ref.opts = opts + ref.clearOpenFailure() return idx } @@ -190,17 +235,31 @@ func (m *DBManager) Get(idx int) (DB, error) { if err != nil { return nil, err } - defer db.mtx.Unlock() + // NOTE: allocDB returns with db.mtx held. It is unlocked explicitly rather + // than deferred because the failed-open path must release it *before* + // Release acquires m.mtx: allocDB takes m.mtx then db.mtx, so holding + // db.mtx while waiting on m.mtx deadlocks against a concurrent allocDB. if ref.db == nil { d, err := m.openDB(db.name, db.opts) if err != nil { + db.recordOpenFailure(err) + db.mtx.Unlock() + + // Drops the db-less ref from the cache so a later access retries the + // open instead of observing a permanently poisoned entry. m.Release(idx) return nil, err } ref.db = d + atomic.StoreUint32(&ref.opened, 1) + db.clearOpenFailure() } - return ref.db, nil + + opened := ref.db + db.mtx.Unlock() + + return opened, nil } func (m *DBManager) allocDB(idx int, db *dbInfo) (*dbRef, error) { @@ -250,27 +309,32 @@ func (m *DBManager) Release(idx int) { return } - if atomic.AddUint32(&ref.count, ^uint32(0)) == 0 { - m.signal() - // If Close() deferred the underlying close to us, finalize it now. - m.mtx.Lock() - if pc := m.pendingClose[idx]; pc == ref { - delete(m.pendingClose, idx) - m.mtx.Unlock() - if ref.db != nil { - ref.db.Close() - } - return - } - m.mtx.Unlock() + if atomic.AddUint32(&ref.count, ^uint32(0)) != 0 { + return } -} -func (m *DBManager) signal() { - m.mtx.Lock() - defer m.mtx.Unlock() + var toClose DB + m.mtx.Lock() + if pc := m.pendingClose[idx]; pc == ref { + // Close() deferred the underlying close to us; finalize it below. + delete(m.pendingClose, idx) + toClose = ref.db + } else if v, err := m.dbCache.Get(idx); err == nil && + v.(*dbRef) == ref && + atomic.LoadUint32(&ref.count) == 0 && + atomic.LoadUint32(&ref.opened) == 0 { + // The open failed and nobody picked this ref up in the meantime. Drop it, + // so a dbRef with a nil db is never observable: leaving it cached is what + // made a transient storage error terminal until restart. + m.dbCache.Pop(idx) + } m.waitCond.Signal() + m.mtx.Unlock() + + if toClose != nil { + toClose.Close() + } } func (m *DBManager) Has(name string) bool { @@ -319,17 +383,6 @@ func (m *DBManager) GetOptionsByIndex(idx int) *Options { if !has { return nil } - - ref, err := m.dbCache.Get(idx) - if err == nil { - dbInfo.mtx.Lock() - defer dbInfo.mtx.Unlock() - - if dbRef := ref.(*dbRef); dbRef != nil && dbRef.db != nil { - return dbInfo.opts - } - return nil - } return dbInfo.getOptions() } @@ -339,20 +392,25 @@ func (m *DBManager) GetState(idx int) (*schema.ImmutableState, error) { return nil, ErrDatabaseNotExists } - ref, err := m.dbCache.Get(idx) - if err == nil { + if ref, err := m.dbCache.Get(idx); err == nil { dbInfo.mtx.Lock() - defer dbInfo.mtx.Unlock() - if dbRef := ref.(*dbRef); dbRef != nil && dbRef.db != nil { + defer dbInfo.mtx.Unlock() return dbRef.db.CurrentState() } - // this condition should never happen - return nil, fmt.Errorf("unable to get state") + // The open behind this ref failed, or is still in flight. Fall through and + // let m.Get drive it. The unlock is explicit: m.Get -> allocDB acquires + // this same mutex, so deferring it here would self-deadlock. + dbInfo.mtx.Unlock() } - s := dbInfo.getState() - if s != nil { + if failure := dbInfo.openFailure(); failure != nil { + // A recorded failure outranks any state cached by a previous successful + // open, otherwise a database that fails to *re*open never retries. + if time.Since(failure.at) < openRetryInterval { + return nil, failure.err + } + } else if s := dbInfo.getState(); s != nil { return s, nil } @@ -495,9 +553,9 @@ func (m *DBManager) CloseAll(ctx context.Context) error { return nil } - // ref.db may be nil if the database was never successfully opened - // (e.g. openDB failed in Get and left a db-less ref in the cache). - // Guard here as every other close site in this file does. + // Defensive: Release pops the ref left behind by a failed open, but + // that ref stays briefly observable in between. Guard as every other + // close site in this file does. if ref.db != nil { ref.db.Close() } @@ -511,7 +569,37 @@ func (m *DBManager) CloseAll(ctx context.Context) error { return nil } +// IsActive reports whether the database is currently open. A database whose +// open failed is not active, even while its ref is still cached. func (m *DBManager) IsActive(idx int) bool { - _, err := m.dbCache.Get(idx) - return err == nil + v, err := m.dbCache.Get(idx) + if err != nil { + return false + } + + ref, _ := v.(*dbRef) + return ref != nil && atomic.LoadUint32(&ref.opened) == 1 +} + +// OpenFailures returns every database whose most recent open attempt failed and +// which has not been opened successfully since. It is the basis of the +// readiness signal: a database that was never accessed has never attempted an +// open and is not reported here. +func (m *DBManager) OpenFailures() []DatabaseOpenFailure { + m.dbMutex.RLock() + dbs := make([]*dbInfo, len(m.databases)) + copy(dbs, m.databases) + m.dbMutex.RUnlock() + + // NOTE: no dbInfo mutex is taken here. db.mtx is held for the whole duration + // of an open, which is precisely when a probe needs an answer. A closed or + // deleted database cannot carry a recorded failure: close() and Put() both + // clear it, and allocDB refuses to open a deleted database. + var failures []DatabaseOpenFailure + for _, db := range dbs { + if failure := db.openFailure(); failure != nil { + failures = append(failures, DatabaseOpenFailure{Name: db.name, Err: failure.err, At: failure.at}) + } + } + return failures } diff --git a/pkg/database/db_manager_test.go b/pkg/database/db_manager_test.go index 35b8253de0..bd9e3d2e96 100644 --- a/pkg/database/db_manager_test.go +++ b/pkg/database/db_manager_test.go @@ -27,6 +27,7 @@ import ( "testing" "time" + "github.com/codenotary/immudb/embedded/cache" "github.com/codenotary/immudb/embedded/logger" "github.com/codenotary/immudb/embedded/sql" "github.com/codenotary/immudb/embedded/store" @@ -493,10 +494,10 @@ func TestDBManagerPutUpdate(t *testing.T) { } func TestDBManagerCloseAllAfterFailedOpen(t *testing.T) { - // A database whose open fails leaves a db-less ref in the cache: Get() has - // allocDB Put() a &dbRef{count: 1}, openDB then fails and Release() only - // decrements the count without removing the entry. CloseAll must skip such a - // ref instead of dereferencing its nil db. + // allocDB caches a placeholder ref before openDB runs. When the open fails, + // Release must drop that ref: a cached dbRef with a nil db used to poison + // GetState and IsActive until restart. CloseAll keeps its nil guard as + // defence in depth, since the ref stays briefly observable in between. openErr := fmt.Errorf("open failed") openDB := func(name string, opts *Options) (DB, error) { return nil, openErr @@ -508,15 +509,278 @@ func TestDBManagerCloseAllAfterFailedOpen(t *testing.T) { _, err := manager.Get(0) require.ErrorIs(t, err, openErr) - // The failed open must have left a ref with a nil db in the cache. - v, err := manager.dbCache.Get(0) - require.NoError(t, err) - ref, _ := v.(*dbRef) - require.NotNil(t, ref) - require.Nil(t, ref.db) - require.Zero(t, atomic.LoadUint32(&ref.count)) + _, err = manager.dbCache.Get(0) + require.ErrorIs(t, err, cache.ErrKeyNotFound, "failed open must not leave a db-less ref cached") + + require.False(t, manager.IsActive(0), "a database that never opened is not active") require.NotPanics(t, func() { require.NoError(t, manager.CloseAll(context.Background())) }) } + +// TestDBManagerFailedOpenDoesNotDeadlock guards the lock-order inversion that +// used to wedge the whole manager: allocDB returns holding db.mtx while having +// taken m.mtx first, and the failure path then wanted m.mtx again via Release. +// A concurrent allocDB for the same index closed the cycle. This hangs if Get +// reacquires m.mtx before releasing db.mtx. +func TestDBManagerFailedOpenDoesNotDeadlock(t *testing.T) { + openErr := fmt.Errorf("open failed") + openDB := func(name string, opts *Options) (DB, error) { + return nil, openErr + } + + manager := NewDBManager(openDB, 5, logger.NewMemoryLogger()) + manager.Put("testdb", DefaultOptions(), false) + + const n = 50 + + done := make(chan struct{}) + go func() { + defer close(done) + + var wg sync.WaitGroup + wg.Add(n) + for i := 0; i < n; i++ { + go func() { + defer wg.Done() + _, err := manager.Get(0) + require.ErrorIs(t, err, openErr) + }() + } + wg.Wait() + }() + + select { + case <-done: + case <-time.After(30 * time.Second): + t.Fatal("concurrent failed opens deadlocked") + } + + _, err := manager.dbCache.Get(0) + require.ErrorIs(t, err, cache.ErrKeyNotFound) +} + +// TestDBManagerRetriesAfterFailedOpen is the acceptance criterion for #2155: a +// transient storage error during an open must not be terminal. +func TestDBManagerRetriesAfterFailedOpen(t *testing.T) { + openErr := fmt.Errorf("transient storage failure") + + var fail atomic.Bool + fail.Store(true) + + var calls atomic.Uint32 + openDB := func(name string, opts *Options) (DB, error) { + calls.Add(1) + if fail.Load() { + return nil, openErr + } + return &mockDB{name: name}, nil + } + + manager := NewDBManager(openDB, 5, logger.NewMemoryLogger()) + manager.Put("testdb", DefaultOptions(), false) + + _, err := manager.Get(0) + require.ErrorIs(t, err, openErr) + + // GetState must surface the real open error, not the old "unable to get + // state" sentinel, and must not serve it forever. + _, err = manager.GetState(0) + require.ErrorIs(t, err, openErr) + + info, has := manager.getDB(0) + require.True(t, has) + require.NotNil(t, info.openFailure()) + + // Age the recorded failure past the retry backoff. + failure := info.openFailure() + info.lastOpenFail.Store(&openFailure{err: failure.err, at: failure.at.Add(-2 * openRetryInterval)}) + + fail.Store(false) + + // GetState alone re-drives the open: this is what makes the once-a-minute + // metrics loop a self-healing retry driver. + state, err := manager.GetState(0) + require.NoError(t, err) + require.NotNil(t, state) + + require.True(t, manager.IsActive(0)) + require.Nil(t, info.openFailure()) + require.Empty(t, manager.OpenFailures()) + require.Greater(t, calls.Load(), uint32(1)) + + require.NoError(t, manager.CloseAll(context.Background())) +} + +func TestDBManagerConcurrentFailedThenSuccessfulOpen(t *testing.T) { + openErr := fmt.Errorf("open failed") + + var fail atomic.Bool + fail.Store(true) + + openDB := func(name string, opts *Options) (DB, error) { + if fail.Load() { + return nil, openErr + } + return &mockDB{name: name}, nil + } + + manager := NewDBManager(openDB, 5, logger.NewMemoryLogger()) + manager.Put("testdb", DefaultOptions(), false) + + var wg sync.WaitGroup + wg.Add(20) + for i := 0; i < 20; i++ { + go func(i int) { + defer wg.Done() + + if i == 10 { + fail.Store(false) + } + if db, err := manager.Get(0); err == nil { + require.NotNil(t, db) + manager.Release(0) + } + }(i) + } + wg.Wait() + + fail.Store(false) + + db, err := manager.Get(0) + require.NoError(t, err) + require.NotNil(t, db) + manager.Release(0) + + v, err := manager.dbCache.Get(0) + require.NoError(t, err) + ref := v.(*dbRef) + require.NotNil(t, ref.db) + require.Zero(t, atomic.LoadUint32(&ref.count)) + require.Equal(t, 1, manager.dbCache.EntriesCount()) + + require.NoError(t, manager.CloseAll(context.Background())) +} + +// TestDBManagerCloseDuringFailedOpen guards the pendingClose hand-off: if Close +// parks the ref while the failing open is still counted, the last Release must +// drain it rather than orphan it. +func TestDBManagerCloseDuringFailedOpen(t *testing.T) { + openErr := fmt.Errorf("open failed") + + opening := make(chan struct{}) + release := make(chan struct{}) + openDB := func(name string, opts *Options) (DB, error) { + close(opening) + <-release + return nil, openErr + } + + manager := NewDBManager(openDB, 5, logger.NewMemoryLogger()) + manager.Put("testdb", DefaultOptions(), false) + + getDone := make(chan struct{}) + go func() { + defer close(getDone) + + _, err := manager.Get(0) + require.ErrorIs(t, err, openErr) + }() + + // allocDB has returned and the opener holds db.mtx. + <-opening + + closeDone := make(chan struct{}) + go func() { + defer close(closeDone) + + // Blocks on db.mtx until the failing open lets go of it. + _ = manager.Close(0) + }() + + time.Sleep(50 * time.Millisecond) + close(release) + + for _, c := range []chan struct{}{getDone, closeDone} { + select { + case <-c: + case <-time.After(30 * time.Second): + t.Fatal("failed open racing Close deadlocked") + } + } + + manager.mtx.Lock() + pending := len(manager.pendingClose) + manager.mtx.Unlock() + require.Zero(t, pending, "pendingClose must not be orphaned") + + _, err := manager.dbCache.Get(0) + require.ErrorIs(t, err, cache.ErrKeyNotFound) +} + +func TestDBManagerGetOptionsByIndexAfterFailedOpen(t *testing.T) { + openErr := fmt.Errorf("open failed") + openDB := func(name string, opts *Options) (DB, error) { + return nil, openErr + } + + manager := NewDBManager(openDB, 5, logger.NewMemoryLogger()) + manager.Put("testdb", DefaultOptions(), false) + + _, err := manager.Get(0) + require.ErrorIs(t, err, openErr) + + // A nil *Options here is what lazyDB.Path() and MaxResultSize() dereference. + opts := manager.GetOptionsByIndex(0) + require.NotNil(t, opts) + + db := &lazyDB{m: manager, idx: 0} + require.NotPanics(t, func() { _ = db.Path() }) +} + +func TestDBManagerOpenFailures(t *testing.T) { + openErr := fmt.Errorf("open failed") + + var fail atomic.Bool + fail.Store(true) + + openDB := func(name string, opts *Options) (DB, error) { + if fail.Load() { + return nil, openErr + } + return &mockDB{name: name}, nil + } + + manager := NewDBManager(openDB, 5, logger.NewMemoryLogger()) + manager.Put("okdb", DefaultOptions(), false) + manager.Put("faileddb", DefaultOptions(), false) + manager.Put("untoucheddb", DefaultOptions(), false) + + require.Empty(t, manager.OpenFailures(), "a database that was never accessed is not a failure") + + _, err := manager.Get(1) + require.ErrorIs(t, err, openErr) + + failures := manager.OpenFailures() + require.Len(t, failures, 1) + require.Equal(t, "faileddb", failures[0].Name) + require.ErrorIs(t, failures[0].Err, openErr) + require.False(t, failures[0].At.IsZero()) + + fail.Store(false) + + _, err = manager.Get(0) + require.NoError(t, err) + manager.Release(0) + + require.Len(t, manager.OpenFailures(), 1, "an unrelated successful open must not clear it") + + _, err = manager.Get(1) + require.NoError(t, err) + manager.Release(1) + + require.Empty(t, manager.OpenFailures(), "a successful open clears the recorded failure") + + require.NoError(t, manager.CloseAll(context.Background())) +} From e0d5b8d42fbcb34e89ad116d964fd81f21b39d62 Mon Sep 17 00:00:00 2001 From: vchaindz Date: Fri, 4 Sep 2026 19:38:23 +0200 Subject: [PATCH 2/2] feat(server): report databases that failed to open on /readyz Health returns Status: true unconditionally and DatabaseHealth reports request statistics, so neither can tell an orchestrator that a database is unusable. DatabaseListV2 is not an answer either: Loaded is !IsClosed(), and closed is set only by Close() or PutClosed(), never by a failed open, so a database whose open failed still reports Loaded: true. The consequence is that when a transient 403 from S3 left two databases unopened (#2155), the pod stayed 1/1 Running and the liveness probe passed for the whole outage. It surfaced instead as a crashloop in the client services, pointing on-call at the consumers rather than at immudb. The metrics server already serves /initz, /readyz and /livez, all of them unconditional 200 stubs. /readyz now returns 503 and names each database whose most recent open attempt failed and which has not opened successfully since. /livez and /initz stay unconditional: liveness should not depend on storage reachability. A database that was never accessed has never attempted an open and is not reported, so registering databases lazily at startup does not fail readiness. The failure record is read atomically rather than under the dbInfo mutex, which is held for the entire duration of an open and so is precisely unavailable when a probe needs an answer. IsActive is corrected in passing: it only checked that the cache lookup succeeded, which was true for a ref whose open had failed. It now requires a ref that actually holds an open database. Closes #2156 Signed-off-by: vchaindz --- pkg/database/types.go | 7 ++++ pkg/server/metrics.go | 28 +++++++++++++++- pkg/server/metrics_funcs.go | 22 +++++++++++++ pkg/server/metrics_funcs_test.go | 56 ++++++++++++++++++++++++++++++++ pkg/server/metrics_test.go | 40 +++++++++++++++++++++++ pkg/server/server.go | 1 + 6 files changed, 153 insertions(+), 1 deletion(-) diff --git a/pkg/database/types.go b/pkg/database/types.go index fb0ef0203d..d1c38b2d7b 100644 --- a/pkg/database/types.go +++ b/pkg/database/types.go @@ -29,6 +29,7 @@ type DatabaseList interface { Length() int Resize(n int) CloseAll(ctx context.Context) error + OpenFailures() []DatabaseOpenFailure } type databaseList struct { @@ -100,3 +101,9 @@ func (d *databaseList) CloseAll(ctx context.Context) error { func (d *databaseList) Resize(n int) { d.m.Resize(n) } + +// OpenFailures returns the databases whose most recent open attempt failed and +// which have not been opened successfully since. +func (d *databaseList) OpenFailures() []DatabaseOpenFailure { + return d.m.OpenFailures() +} diff --git a/pkg/server/metrics.go b/pkg/server/metrics.go index cd6805a7fb..b19f8ef573 100644 --- a/pkg/server/metrics.go +++ b/pkg/server/metrics.go @@ -21,6 +21,7 @@ import ( "crypto/tls" "encoding/json" "expvar" + "fmt" "net" "net/http" "net/http/pprof" @@ -412,6 +413,7 @@ func StartMetrics( computeDBEntries func() map[string]float64, computeLoadedDBSize func() float64, computeSessionCount func() float64, + computeReadiness func() (bool, []string), addPProf bool, ) *http.Server { Metrics.WithUptimeCounter(uptimeCounter) @@ -449,7 +451,7 @@ func StartMetrics( mux.HandleFunc("/debug/pprof/trace", corsHandlerFunc(pprof.Trace)) } mux.HandleFunc("/initz", corsHandlerFunc(ImmudbHealthHandlerFunc())) - mux.HandleFunc("/readyz", corsHandlerFunc(ImmudbHealthHandlerFunc())) + mux.HandleFunc("/readyz", corsHandlerFunc(ImmudbReadinessHandlerFunc(computeReadiness))) mux.HandleFunc("/livez", corsHandlerFunc(ImmudbHealthHandlerFunc())) mux.HandleFunc("/version", corsHandlerFunc(ImmudbVersionHandlerFunc)) @@ -485,6 +487,30 @@ func ImmudbHealthHandlerFunc() http.HandlerFunc { } } +// ImmudbReadinessHandlerFunc serves /readyz. Unlike /livez and /initz, which +// only report that the process is up, it fails while a database that attempted +// an open is not open, so an orchestrator can act on it. +func ImmudbReadinessHandlerFunc(computeReadiness func() (bool, []string)) http.HandlerFunc { + return func(w http.ResponseWriter, r *http.Request) { + if computeReadiness == nil { + w.WriteHeader(http.StatusOK) + return + } + + ready, reasons := computeReadiness() + if ready { + w.WriteHeader(http.StatusOK) + return + } + + w.Header().Set("Content-Type", "text/plain; charset=utf-8") + w.WriteHeader(http.StatusServiceUnavailable) + for _, reason := range reasons { + fmt.Fprintln(w, reason) + } + } +} + func ImmudbVersionHandlerFunc(w http.ResponseWriter, r *http.Request) { writeJSONResponse(w, r, 200, &Version) } diff --git a/pkg/server/metrics_funcs.go b/pkg/server/metrics_funcs.go index a191a15fff..7efc57a9a0 100644 --- a/pkg/server/metrics_funcs.go +++ b/pkg/server/metrics_funcs.go @@ -141,6 +141,28 @@ func (s *ImmuServer) metricFuncComputeLoadedDBSize() float64 { return float64(s.dbList.Length()) } +// metricFuncComputeReadiness reports whether every database that attempted an +// open actually opened. A transient remote-storage failure during an open used +// to be invisible to orchestrators: the process stayed healthy while the +// database was unusable. +func (s *ImmuServer) metricFuncComputeReadiness() (bool, []string) { + if s.dbList == nil { + return true, nil + } + + failures := s.dbList.OpenFailures() + if len(failures) == 0 { + return true, nil + } + + reasons := make([]string, 0, len(failures)) + for _, f := range failures { + reasons = append(reasons, fmt.Sprintf( + "database %s failed to open at %s: %v", f.Name, f.At.UTC().Format(time.RFC3339), f.Err)) + } + return false, reasons +} + func (s *ImmuServer) metricFuncComputeSessionCount() float64 { if s.SessManager == nil { s.Logger.Warningf( diff --git a/pkg/server/metrics_funcs_test.go b/pkg/server/metrics_funcs_test.go index 096d70b780..6a2d55dee1 100644 --- a/pkg/server/metrics_funcs_test.go +++ b/pkg/server/metrics_funcs_test.go @@ -53,6 +53,10 @@ func (dbm dbMock) GetOptions() *database.Options { return database.DefaultOptions() } +func (dbm dbMock) Size() (uint64, error) { + return 0, nil +} + func (dbm dbMock) GetName() string { if dbm.getNameF != nil { return dbm.getNameF() @@ -117,6 +121,58 @@ func TestMetricFuncComputeDBEntries(t *testing.T) { s.metricFuncComputeDBEntries() } +func TestMetricFuncComputeReadiness(t *testing.T) { + openErr := fmt.Errorf("403 Forbidden") + + failOpen := true + dbList := database.NewDatabaseList(database.NewDBManager( + func(name string, opts *database.Options) (database.DB, error) { + if failOpen { + return nil, openErr + } + return &dbMock{getNameF: func() string { return name }}, nil + }, 10, logger.NewMemoryLogger())) + + s := ImmuServer{dbList: dbList} + + ready, reasons := s.metricFuncComputeReadiness() + require.True(t, ready, "no databases registered") + require.Empty(t, reasons) + + dbList.Put("audit", database.DefaultOptions()) + + ready, reasons = s.metricFuncComputeReadiness() + require.True(t, ready, "a database that was never accessed has not failed to open") + require.Empty(t, reasons) + + db, err := dbList.GetByIndex(0) + require.NoError(t, err) + + _, err = db.CurrentState() + require.ErrorIs(t, err, openErr) + + ready, reasons = s.metricFuncComputeReadiness() + require.False(t, ready, "a database that failed to open must fail readiness") + require.Len(t, reasons, 1) + require.Contains(t, reasons[0], "audit") + require.Contains(t, reasons[0], "403 Forbidden") + + failOpen = false + + // Size() goes through DBManager.Get, which re-drives the open. + _, err = db.Size() + require.NoError(t, err) + + ready, reasons = s.metricFuncComputeReadiness() + require.True(t, ready, "a successful open clears the readiness failure") + require.Empty(t, reasons) + + // nil db list must not report unready + s.dbList = nil + ready, _ = s.metricFuncComputeReadiness() + require.True(t, ready) +} + func TestMetricFuncServerUptimeCounter(t *testing.T) { s := ImmuServer{} s.metricFuncServerUptimeCounter() diff --git a/pkg/server/metrics_test.go b/pkg/server/metrics_test.go index e3f808b0ac..b7a5ad5394 100644 --- a/pkg/server/metrics_test.go +++ b/pkg/server/metrics_test.go @@ -48,6 +48,7 @@ func TestStartMetricsHTTP(t *testing.T) { func() map[string]float64 { return make(map[string]float64) }, func() float64 { return 1.0 }, func() float64 { return 2.0 }, + func() (bool, []string) { return true, nil }, false, ) time.Sleep(200 * time.Millisecond) @@ -73,6 +74,7 @@ func TestStartMetricsHTTPS(t *testing.T) { func() map[string]float64 { return make(map[string]float64) }, func() float64 { return 1.0 }, func() float64 { return 2.0 }, + func() (bool, []string) { return true, nil }, false, ) time.Sleep(200 * time.Millisecond) @@ -107,6 +109,7 @@ func TestStartMetricsFail(t *testing.T) { func() map[string]float64 { return make(map[string]float64) }, func() float64 { return 1.0 }, func() float64 { return 2.0 }, + func() (bool, []string) { return true, nil }, false, ) time.Sleep(200 * time.Millisecond) @@ -293,6 +296,43 @@ func TestImmudbHealthHandlerFunc(t *testing.T) { require.Equal(t, http.StatusOK, rr.Code) } +func TestImmudbReadinessHandlerFunc(t *testing.T) { + t.Run("ready", func(t *testing.T) { + req, err := http.NewRequest("GET", "/readyz", nil) + require.NoError(t, err) + + rr := httptest.NewRecorder() + corsHandlerFunc(ImmudbReadinessHandlerFunc(func() (bool, []string) { + return true, nil + })).ServeHTTP(rr, req) + + require.Equal(t, http.StatusOK, rr.Code) + }) + + t.Run("not ready reports the databases that failed to open", func(t *testing.T) { + req, err := http.NewRequest("GET", "/readyz", nil) + require.NoError(t, err) + + rr := httptest.NewRecorder() + corsHandlerFunc(ImmudbReadinessHandlerFunc(func() (bool, []string) { + return false, []string{"database audit failed to open: 403 Forbidden"} + })).ServeHTTP(rr, req) + + require.Equal(t, http.StatusServiceUnavailable, rr.Code) + require.Contains(t, rr.Body.String(), "database audit failed to open") + }) + + t.Run("no readiness func stays healthy", func(t *testing.T) { + req, err := http.NewRequest("GET", "/readyz", nil) + require.NoError(t, err) + + rr := httptest.NewRecorder() + corsHandlerFunc(ImmudbReadinessHandlerFunc(nil)).ServeHTTP(rr, req) + + require.Equal(t, http.StatusOK, rr.Code) + }) +} + func TestImmudbVersionHandlerFunc(t *testing.T) { // test OPTIONS /version req, err := http.NewRequest("OPTIONS", "/version", nil) diff --git a/pkg/server/server.go b/pkg/server/server.go index 126452d2c4..10076f65bb 100644 --- a/pkg/server/server.go +++ b/pkg/server/server.go @@ -349,6 +349,7 @@ func (s *ImmuServer) Start() (err error) { s.metricFuncComputeDBEntries, s.metricFuncComputeLoadedDBSize, s.metricFuncComputeSessionCount, + s.metricFuncComputeReadiness, s.Options.PProf) defer func() {