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())) +} 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() {