refactor: provide AppState infrastructure to wait for schema cache load

This commit replaces ioRef based implementation of schema cache status tracking to MVar based, so that it is possible to wait for schema cache loading.

Waiting for schema cache loading is necessary to implement zero-downtime upgrades  with SO_REUSEPORT, where listening on a socket must wait for schema cache loading.
This commit is contained in:
Michał Kłeczek
2026-03-12 09:04:56 -05:00
committed by Steve Chavez
parent d5694672a4
commit e741c1bca7
+27 -18
View File
@@ -84,7 +84,7 @@ data AppState = AppState
-- | Schema cache -- | Schema cache
, stateSchemaCache :: IORef (Maybe SchemaCache) , stateSchemaCache :: IORef (Maybe SchemaCache)
-- | The schema cache status -- | The schema cache status
, stateSCacheStatus :: IORef SchemaCacheStatus , stateSCacheStatus :: SchemaCacheStatus
-- | State of the LISTEN channel -- | State of the LISTEN channel
, stateIsListenerOn :: IORef Bool , stateIsListenerOn :: IORef Bool
-- | starts the connection worker with a debounce -- | starts the connection worker with a debounce
@@ -111,11 +111,11 @@ data AppState = AppState
, stateMetrics :: Metrics.MetricsState , stateMetrics :: Metrics.MetricsState
} }
-- | Schema cache status -- | Schema cache status.
data SchemaCacheStatus -- Empty means pending and full means loaded.
= SCLoaded newtype SchemaCacheStatus = SchemaCacheStatus
| SCPending { getSCStatusMVar :: MVar ()
deriving Eq }
type AppSockets = (NS.Socket, Maybe NS.Socket) type AppSockets = (NS.Socket, Maybe NS.Socket)
@@ -138,7 +138,7 @@ initWithPool (sock, adminSock) pool conf loggerState metricsState observer = do
appState <- AppState pool appState <- AppState pool
<$> newIORef minimumPgVersion -- assume we're in a supported version when starting, this will be corrected on a later step <$> newIORef minimumPgVersion -- assume we're in a supported version when starting, this will be corrected on a later step
<*> newIORef Nothing <*> newIORef Nothing
<*> newIORef SCPending <*> newSchemaCacheStatus
<*> newIORef False <*> newIORef False
<*> pure (pure ()) <*> pure (pure ())
<*> newIORef conf <*> newIORef conf
@@ -335,18 +335,15 @@ putIsListenerOn = atomicWriteIORef . stateIsListenerOn
isLoaded :: AppState -> IO Bool isLoaded :: AppState -> IO Bool
isLoaded x = do isLoaded x = do
scacheStatus <- readIORef $ stateSCacheStatus x scacheLoaded <- isSchemaCacheLoaded x
connEstablished <- isConnEstablished x connEstablished <- isConnEstablished x
return $ scacheStatus == SCLoaded && connEstablished return $ scacheLoaded && connEstablished
isPending :: AppState -> IO Bool isPending :: AppState -> IO Bool
isPending x = do isPending x = do
scacheStatus <- readIORef $ stateSCacheStatus x scacheLoaded <- isSchemaCacheLoaded x
connEstablished <- isConnEstablished x connEstablished <- isConnEstablished x
return $ scacheStatus == SCPending || not connEstablished return $ not scacheLoaded || not connEstablished
putSCacheStatus :: AppState -> SchemaCacheStatus -> IO ()
putSCacheStatus = atomicWriteIORef . stateSCacheStatus
getObserver :: AppState -> ObservationHandler getObserver :: AppState -> ObservationHandler
getObserver = stateObserver getObserver = stateObserver
@@ -405,19 +402,19 @@ retryingSchemaCacheLoad appState@AppState{stateObserver=observer, stateMainThrea
timeItT $ usePool appState (transaction SQL.ReadCommitted SQL.Read $ querySchemaCache conf) timeItT $ usePool appState (transaction SQL.ReadCommitted SQL.Read $ querySchemaCache conf)
case result of case result of
Left e -> do Left e -> do
putSCacheStatus appState SCPending markSchemaCachePending appState
putSchemaCache appState Nothing putSchemaCache appState Nothing
observer $ SchemaCacheErrorObs configDbSchemas configDbExtraSearchPath e observer $ SchemaCacheErrorObs configDbSchemas configDbExtraSearchPath e
return Nothing return Nothing
Right sCache -> do Right sCache -> do
-- IMPORTANT: While the pending schema cache state starts from running the above querySchemaCache, only at this stage we block API requests due to the usage of an -- IMPORTANT: While the pending schema cache state starts from running the above querySchemaCache, only at this stage we block API requests due to the usage of an
-- IORef on putSchemaCache. This is why SCacheStatus is put at SCPending here to signal the Admin server (using isPending) that we're on a recovery state. -- IORef on putSchemaCache. This is why schema cache status is marked as pending here to signal the Admin server (using isPending) that we're on a recovery state.
putSCacheStatus appState SCPending markSchemaCachePending appState
putSchemaCache appState $ Just sCache putSchemaCache appState $ Just sCache
observer $ SchemaCacheQueriedObs resultTime observer $ SchemaCacheQueriedObs resultTime
observer . uncurry SchemaCacheLoadedObs =<< timeItT (evaluate $ showSummary sCache) observer . uncurry SchemaCacheLoadedObs =<< timeItT (evaluate $ showSummary sCache)
putSCacheStatus appState SCLoaded markSchemaCacheLoaded appState
return $ Just sCache return $ Just sCache
shouldRetry :: RetryStatus -> (Maybe PgVersion, Maybe SchemaCache) -> IO Bool shouldRetry :: RetryStatus -> (Maybe PgVersion, Maybe SchemaCache) -> IO Bool
@@ -433,6 +430,18 @@ retryingSchemaCacheLoad appState@AppState{stateObserver=observer, stateMainThrea
oneSecondInUs = 1000000 -- one second in microseconds oneSecondInUs = 1000000 -- one second in microseconds
newSchemaCacheStatus :: IO SchemaCacheStatus
newSchemaCacheStatus = SchemaCacheStatus <$> newEmptyMVar
markSchemaCachePending :: AppState -> IO ()
markSchemaCachePending = void . tryTakeMVar . getSCStatusMVar . stateSCacheStatus
markSchemaCacheLoaded :: AppState -> IO ()
markSchemaCacheLoaded = void . (`tryPutMVar` ()) . getSCStatusMVar . stateSCacheStatus
isSchemaCacheLoaded :: AppState -> IO Bool
isSchemaCacheLoaded = fmap not . isEmptyMVar . getSCStatusMVar . stateSCacheStatus
-- | Reads the in-db config and reads the config file again -- | Reads the in-db config and reads the config file again
-- | We don't retry reading the in-db config after it fails immediately, because it could have user errors. We just report the error and continue. -- | We don't retry reading the in-db config after it fails immediately, because it could have user errors. We just report the error and continue.
readInDbConfig :: Bool -> AppState -> IO () readInDbConfig :: Bool -> AppState -> IO ()