diff --git a/CHANGELOG.md b/CHANGELOG.md index 62d020041..12a67a727 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -11,6 +11,7 @@ This project adheres to [Semantic Versioning](http://semver.org/). - #3644, Make --dump-schema work with in-database pgrst.db_schemas setting - @wolfgangwalther - #3644, Show number of timezones in schema cache load report - @wolfgangwalther - #3644, List correct enum options in OpenApi output when multiple types with same name are present - @wolfgangwalther + - #3523, Fix schema cache loading retry without backoff - @steve-chavez ## [12.2.1] - 2024-06-27 diff --git a/src/PostgREST/App.hs b/src/PostgREST/App.hs index aeb43ffd2..8026ff66b 100644 --- a/src/PostgREST/App.hs +++ b/src/PostgREST/App.hs @@ -68,14 +68,14 @@ run appState = do observer $ AppStartObs prettyVersion - AppState.connectionWorker appState - Unix.installSignalHandlers (AppState.getMainThreadId appState) (AppState.connectionWorker appState) (AppState.reReadConfig False appState) + AppState.schemaCacheLoader appState -- Loads the initial SchemaCache + Unix.installSignalHandlers (AppState.getMainThreadId appState) (AppState.schemaCacheLoader appState) (AppState.readInDbConfig False appState) Listener.runListener appState Admin.runAdmin appState (serverSettings conf) - let app = postgrest configLogLevel appState (AppState.connectionWorker appState) + let app = postgrest configLogLevel appState (AppState.schemaCacheLoader appState) case configServerUnixSocket of Just path -> do diff --git a/src/PostgREST/AppState.hs b/src/PostgREST/AppState.hs index 391b1f0e5..11eebffc6 100644 --- a/src/PostgREST/AppState.hs +++ b/src/PostgREST/AppState.hs @@ -24,8 +24,8 @@ module PostgREST.AppState , putPgVersion , putIsListenerOn , usePool - , reReadConfig - , connectionWorker + , readInDbConfig + , schemaCacheLoader , getObserver , isLoaded , isPending @@ -85,39 +85,37 @@ data AuthResult = AuthResult data AppState = AppState -- | Database connection pool - { statePool :: SQL.Pool - -- | Database server version, will be updated by the connectionWorker - , statePgVersion :: IORef PgVersion - -- | No schema cache at the start. Will be filled in by the connectionWorker - , stateSchemaCache :: IORef (Maybe SchemaCache) + { statePool :: SQL.Pool + -- | Database server version + , statePgVersion :: IORef PgVersion + -- | Schema cache + , stateSchemaCache :: IORef (Maybe SchemaCache) -- | The schema cache status - , stateSCacheStatus :: IORef SchemaCacheStatus - -- | The connection status - , stateConnStatus :: IORef ConnectionStatus + , stateSCacheStatus :: IORef SchemaCacheStatus -- | State of the LISTEN channel - , stateIsListenerOn :: IORef Bool + , stateIsListenerOn :: IORef Bool -- | starts the connection worker with a debounce - , debouncedConnectionWorker :: IO () + , debouncedSCacheLoader :: IO () -- | Config that can change at runtime - , stateConf :: IORef AppConfig + , stateConf :: IORef AppConfig -- | Time used for verifying JWT expiration - , stateGetTime :: IO UTCTime + , stateGetTime :: IO UTCTime -- | Used for killing the main thread in case a subthread fails - , stateMainThreadId :: ThreadId + , stateMainThreadId :: ThreadId -- | Keeps track of the next delay for db connection retry - , stateNextDelay :: IORef Int + , stateNextDelay :: IORef Int -- | Keeps track of the next delay for the listener - , stateNextListenerDelay :: IORef Int + , stateNextListenerDelay :: IORef Int -- | JWT Cache - , jwtCache :: C.Cache ByteString AuthResult + , jwtCache :: C.Cache ByteString AuthResult -- | Network socket for REST API - , stateSocketREST :: NS.Socket + , stateSocketREST :: NS.Socket -- | Network socket for the admin UI - , stateSocketAdmin :: Maybe NS.Socket + , stateSocketAdmin :: Maybe NS.Socket -- | Observation handler - , stateObserver :: ObservationHandler - , stateLogger :: Logger.LoggerState - , stateMetrics :: Metrics.MetricsState + , stateObserver :: ObservationHandler + , stateLogger :: Logger.LoggerState + , stateMetrics :: Metrics.MetricsState } -- | Schema cache status @@ -126,15 +124,8 @@ data SchemaCacheStatus | SCPending deriving Eq --- | Current database connection status -data ConnectionStatus - = ConnEstablished - | ConnPending - deriving Eq - type AppSockets = (NS.Socket, Maybe NS.Socket) - init :: AppConfig -> IO AppState init conf@AppConfig{configLogLevel, configDbPoolSize} = do loggerState <- Logger.init @@ -153,7 +144,6 @@ initWithPool (sock, adminSock) pool conf loggerState metricsState observer = do <$> newIORef minimumPgVersion -- assume we're in a supported version when starting, this will be corrected on a later step <*> newIORef Nothing <*> newIORef SCPending - <*> newIORef ConnPending <*> newIORef False <*> pure (pure ()) <*> newIORef conf @@ -168,15 +158,15 @@ initWithPool (sock, adminSock) pool conf loggerState metricsState observer = do <*> pure loggerState <*> pure metricsState - debWorker <- + deb <- let decisecond = 100000 in mkDebounce defaultDebounceSettings - { debounceAction = internalConnectionWorker appState + { debounceAction = internalSchemaCacheLoad appState , debounceFreq = decisecond , debounceEdge = leadingEdge -- runs the worker at the start and the end } - return appState { debouncedConnectionWorker = debWorker} + return appState { debouncedSCacheLoader = deb} destroy :: AppState -> IO () destroy = destroyPool @@ -302,15 +292,12 @@ getSchemaCache = readIORef . stateSchemaCache putSchemaCache :: AppState -> Maybe SchemaCache -> IO () putSchemaCache appState = atomicWriteIORef (stateSchemaCache appState) -connectionWorker :: AppState -> IO () -connectionWorker = debouncedConnectionWorker +schemaCacheLoader :: AppState -> IO () +schemaCacheLoader = debouncedSCacheLoader getNextDelay :: AppState -> IO Int getNextDelay = readIORef . stateNextDelay -putNextDelay :: AppState -> Int -> IO () -putNextDelay = atomicWriteIORef . stateNextDelay - getNextListenerDelay :: AppState -> IO Int getNextListenerDelay = readIORef . stateNextListenerDelay @@ -338,167 +325,120 @@ getSocketAdmin = stateSocketAdmin getMainThreadId :: AppState -> ThreadId getMainThreadId = stateMainThreadId -getIsListenerOn :: AppState -> IO Bool -getIsListenerOn appState = do +isConnEstablished :: AppState -> IO Bool +isConnEstablished appState = do AppConfig{..} <- getConfig appState - if configDbChannelEnabled then + if configDbChannelEnabled then -- if the listener is enabled, we can be sure the connection is up readIORef $ stateIsListenerOn appState - else - pure True + else -- otherwise the only way to check the connection is to make a query + isRight <$> usePool appState (SQL.sql "SELECT 1") putIsListenerOn :: AppState -> Bool -> IO () putIsListenerOn = atomicWriteIORef . stateIsListenerOn -isConnEstablished :: AppState -> IO Bool -isConnEstablished x = do - conf <- getConfig x - if configDbChannelEnabled conf - then do -- if the listener is enabled, we can be sure the connection status is always up to date - st <- readIORef $ stateConnStatus x - return $ st == ConnEstablished - else -- otherwise the only way to check the connection is to make a query - isRight <$> usePool x (SQL.sql "SELECT 1") - isLoaded :: AppState -> IO Bool isLoaded x = do scacheStatus <- readIORef $ stateSCacheStatus x connEstablished <- isConnEstablished x - listenerOn <- getIsListenerOn x - return $ scacheStatus == SCLoaded && connEstablished && listenerOn + return $ scacheStatus == SCLoaded && connEstablished isPending :: AppState -> IO Bool isPending x = do scacheStatus <- readIORef $ stateSCacheStatus x - connStatus <- readIORef $ stateConnStatus x - listenerOn <- getIsListenerOn x - return $ scacheStatus == SCPending || connStatus == ConnPending || not listenerOn + connEstablished <- isConnEstablished x + return $ scacheStatus == SCPending || not connEstablished putSCacheStatus :: AppState -> SchemaCacheStatus -> IO () putSCacheStatus = atomicWriteIORef . stateSCacheStatus -putConnStatus :: AppState -> ConnectionStatus -> IO () -putConnStatus = atomicWriteIORef . stateConnStatus - getObserver :: AppState -> ObservationHandler getObserver = stateObserver --- | Load the SchemaCache by using a connection from the pool. -loadSchemaCache :: AppState -> IO SchemaCacheStatus -loadSchemaCache appState@AppState{stateObserver=observer} = do - conf@AppConfig{..} <- getConfig appState - (resultTime, result) <- - let transaction = if configDbPreparedStatements then SQL.transaction else SQL.unpreparedTransaction in - timeItT $ usePool appState (transaction SQL.ReadCommitted SQL.Read $ querySchemaCache conf) - case result of - Left e -> do - putSCacheStatus appState SCPending - putSchemaCache appState Nothing - observer $ SchemaCacheErrorObs e - return SCPending +internalSchemaCacheLoad :: AppState -> IO () +internalSchemaCacheLoad appState = do + AppConfig{..} <- getConfig appState + void $ retryingSchemaCacheLoad appState + -- We cannot retry reading the in-db config after it fails immediately, because it could have user errors. We just report the error and continue. + when configDbConfig $ readInDbConfig False appState - 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 - -- 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. - putSCacheStatus appState SCPending - putSchemaCache appState $ Just sCache - observer $ SchemaCacheQueriedObs resultTime - (t, _) <- timeItT $ observer $ SchemaCacheSummaryObs $ showSummary sCache - observer $ SchemaCacheLoadedObs t - putSCacheStatus appState SCLoaded - return SCLoaded - --- | The purpose of this worker is to obtain a healthy connection to pg and an --- up-to-date schema cache(SchemaCache). This method is meant to be called --- multiple times by the same thread, but does nothing if the previous --- invocation has not terminated. In all cases this method does not halt the --- calling thread, the work is performed in a separate thread. +-- | Try to load the schema cache and retry if it fails. -- --- Background thread that does the following : --- 1. Tries to connect to pg server and will keep trying until success. --- 2. Checks if the pg version is supported and if it's not it kills the main --- program. --- 3. Obtains the sCache. If this fails, it goes back to 1. -internalConnectionWorker :: AppState -> IO () -internalConnectionWorker appState@AppState{stateObserver=observer, stateMainThreadId=mainThreadId} = work +-- This is done by repeatedly: 1) flushing the pool, 2) querying the version and validating that the postgres version is supported by us, and 3) loading the schema cache. +-- It's necessary to flush the pool: +-- +-- + Because connections cache the pg catalog(see #2620) +-- + For rapid recovery. Otherwise, the pool idle or lifetime timeout would have to be reached for new healthy connections to be acquired. +retryingSchemaCacheLoad :: AppState -> IO (Maybe PgVersion, Maybe SchemaCache) +retryingSchemaCacheLoad appState@AppState{stateObserver=observer, stateMainThreadId=mainThreadId} = + retrying retryPolicy shouldRetry (\RetryStatus{rsIterNumber, rsPreviousDelay} -> do + when (rsIterNumber > 0) $ do + let delay = fromMaybe 0 rsPreviousDelay `div` oneSecondInUs + observer $ ConnectionRetryObs delay + putNextListenerDelay appState delay + + flushPool appState + (,) <$> qPgVersion <*> qSchemaCache + ) where - work = do + qPgVersion :: IO (Maybe PgVersion) + qPgVersion = do AppConfig{..} <- getConfig appState - observer DBConnectAttemptObs - connStatus <- establishConnection appState - case connStatus of - ConnPending -> - unless configDbPoolAutomaticRecovery $ do - observer ExitDBNoRecoveryObs - killThread mainThreadId - ConnEstablished -> do - actualPgVersion <- getPgVersion appState - when (actualPgVersion < minimumPgVersion) $ do - observer $ ExitUnsupportedPgVersion actualPgVersion minimumPgVersion - killThread mainThreadId - observer (DBConnectedObs $ pgvFullName actualPgVersion) - -- this could be fail because the connection drops, but the loadSchemaCache will pick the error and retry again - -- We cannot retry after it fails immediately, because db-pre-config could have user errors. We just log the error and continue. - when configDbConfig $ reReadConfig False appState - scStatus <- loadSchemaCache appState - case scStatus of - SCLoaded -> - -- do nothing and proceed if the load was successful - return () - SCPending -> - -- retry reloading the schema cache - work - --- | Repeatedly flush the pool, and check if a connection from the --- pool allows access to the PostgreSQL database. --- --- Releasing the pool is key for rapid recovery. Otherwise, the pool --- timeout would have to be reached for new healthy connections to be acquired. --- Which might not happen if the server is busy with requests. No idle --- connection, no pool timeout. --- --- It's also necessary to release the pool connections because they cache the pg catalog(see #2620) --- --- The connection tries are capped, but if the connection times out no error is --- thrown, just 'False' is returned. -establishConnection :: AppState -> IO ConnectionStatus -establishConnection appState@AppState{stateObserver=observer} = - retrying retryPolicy shouldRetry $ - const $ flushPool appState >> getConnectionStatus - where - getConnectionStatus :: IO ConnectionStatus - getConnectionStatus = do pgVersion <- usePool appState (queryPgVersion False) -- No need to prepare the query here, as the connection might not be established case pgVersion of Left e -> do - observer $ ConnectionPgVersionErrorObs e - putConnStatus appState ConnPending - return ConnPending - Right version -> do - putConnStatus appState ConnEstablished - putPgVersion appState version - return ConnEstablished + observer $ QueryPgVersionError e + unless configDbPoolAutomaticRecovery $ do + observer ExitDBNoRecoveryObs + killThread mainThreadId + return Nothing + Right actualPgVersion -> do + when (actualPgVersion < minimumPgVersion) $ do + observer $ ExitUnsupportedPgVersion actualPgVersion minimumPgVersion + killThread mainThreadId + observer $ DBConnectedObs $ pgvFullName actualPgVersion + putPgVersion appState actualPgVersion + return $ Just actualPgVersion - shouldRetry :: RetryStatus -> ConnectionStatus -> IO Bool - shouldRetry rs isConnSucc = do + qSchemaCache :: IO (Maybe SchemaCache) + qSchemaCache = do + conf@AppConfig{..} <- getConfig appState + (resultTime, result) <- + let transaction = if configDbPreparedStatements then SQL.transaction else SQL.unpreparedTransaction in + timeItT $ usePool appState (transaction SQL.ReadCommitted SQL.Read $ querySchemaCache conf) + case result of + Left e -> do + putSCacheStatus appState SCPending + putSchemaCache appState Nothing + observer $ SchemaCacheErrorObs e + return Nothing + + 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 + -- 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. + putSCacheStatus appState SCPending + putSchemaCache appState $ Just sCache + observer $ SchemaCacheQueriedObs resultTime + (t, _) <- timeItT $ observer $ SchemaCacheSummaryObs $ showSummary sCache + observer $ SchemaCacheLoadedObs t + putSCacheStatus appState SCLoaded + return $ Just sCache + + shouldRetry :: RetryStatus -> (Maybe PgVersion, Maybe SchemaCache) -> IO Bool + shouldRetry _ (pgVer, sCache) = do AppConfig{..} <- getConfig appState - let - delay = fromMaybe 0 (rsPreviousDelay rs) `div` oneSecondInUs - itShould = ConnPending == isConnSucc && configDbPoolAutomaticRecovery - when itShould $ observer $ ConnectionRetryObs delay - when itShould $ putNextDelay appState delay + let itShould = configDbPoolAutomaticRecovery && (isNothing pgVer || isNothing sCache) return itShould retryPolicy :: RetryPolicy retryPolicy = - let - delayMicroseconds = 32000000 -- 32 seconds - in + let delayMicroseconds = 32*oneSecondInUs {-32 seconds-} in capDelay delayMicroseconds $ exponentialBackoff oneSecondInUs - oneSecondInUs = 1000000 -- | One second in microseconds --- | Re-reads the config plus config options from the db -reReadConfig :: Bool -> AppState -> IO () -reReadConfig startingUp appState@AppState{stateObserver=observer} = do + oneSecondInUs = 1000000 -- one second in microseconds + +-- | Reads the in-db config and reads the config file again +readInDbConfig :: Bool -> AppState -> IO () +readInDbConfig startingUp appState@AppState{stateObserver=observer} = do AppConfig{..} <- getConfig appState pgVer <- getPgVersion appState dbSettings <- diff --git a/src/PostgREST/CLI.hs b/src/PostgREST/CLI.hs index c09e3b5f6..0ba71144b 100644 --- a/src/PostgREST/CLI.hs +++ b/src/PostgREST/CLI.hs @@ -42,10 +42,10 @@ main CLI{cliCommand, cliPath} = do AppState.destroy (\appState -> case cliCommand of CmdDumpConfig -> do - when configDbConfig $ AppState.reReadConfig True appState + when configDbConfig $ AppState.readInDbConfig True appState putStr . Config.toText =<< AppState.getConfig appState CmdDumpSchema -> do - when configDbConfig $ AppState.reReadConfig True appState + when configDbConfig $ AppState.readInDbConfig True appState putStrLn =<< dumpSchema appState CmdRun -> App.run appState) diff --git a/src/PostgREST/Listener.hs b/src/PostgREST/Listener.hs index 7438000e3..0c5c42d4e 100644 --- a/src/PostgREST/Listener.hs +++ b/src/PostgREST/Listener.hs @@ -54,8 +54,8 @@ retryingListen appState = do delay <- AppState.getNextListenerDelay appState when (delay > 1) $ do -- if we did a retry - -- assume we lost notifications, call the connection worker which will also reload the schema cache - AppState.connectionWorker appState + -- assume we lost notifications, refresh the schema cache + AppState.schemaCacheLoader appState -- reset the delay AppState.putNextListenerDelay appState 1 @@ -74,8 +74,8 @@ retryingListen appState = do handleNotification channel msg = if | BS.null msg -> observer (DBListenerGotSCacheMsg channel) >> cacheReloader | msg == "reload schema" -> observer (DBListenerGotSCacheMsg channel) >> cacheReloader - | msg == "reload config" -> observer (DBListenerGotConfigMsg channel) >> AppState.reReadConfig False appState + | msg == "reload config" -> observer (DBListenerGotConfigMsg channel) >> AppState.readInDbConfig False appState | otherwise -> pure () -- Do nothing if anything else than an empty message is sent cacheReloader = - AppState.connectionWorker appState + AppState.schemaCacheLoader appState diff --git a/src/PostgREST/Observation.hs b/src/PostgREST/Observation.hs index 28cdbf3ef..884d5f309 100644 --- a/src/PostgREST/Observation.hs +++ b/src/PostgREST/Observation.hs @@ -29,7 +29,6 @@ data Observation | AppStartObs ByteString | AppServerPortObs NS.PortNumber | AppServerUnixObs FilePath - | DBConnectAttemptObs | ExitUnsupportedPgVersion PgVersion PgVersion | ExitDBNoRecoveryObs | ExitDBFatalError ObsFatalError SQL.UsageError @@ -39,7 +38,6 @@ data Observation | SchemaCacheSummaryObs Text | SchemaCacheLoadedObs Double | ConnectionRetryObs Int - | ConnectionPgVersionErrorObs SQL.UsageError | DBListenStart Text | DBListenFail Text (Either SQL.ConnectionError (Either SomeException ())) | DBListenRetry Int @@ -50,6 +48,7 @@ data Observation | ConfigSucceededObs | QueryRoleSettingsErrorObs SQL.UsageError | QueryErrorCodeHighObs SQL.UsageError + | QueryPgVersionError SQL.UsageError | PoolAcqTimeoutObs SQL.UsageError | HasqlPoolObs SQL.Observation | PoolRequest @@ -69,8 +68,6 @@ observationMessage = \case "Listening on port " <> show port AppServerUnixObs sock -> "Listening on unix socket " <> show sock - DBConnectAttemptObs -> - "Attempting to connect to the database..." DBConnectedObs ver -> "Successfully connected to " <> ver ExitUnsupportedPgVersion pgVer minPgVer -> @@ -78,7 +75,7 @@ observationMessage = \case ExitDBNoRecoveryObs -> "Automatic recovery disabled, exiting." ExitDBFatalError ServerAuthError usageErr -> - jsonMessage usageErr + "Failed to establish a connection. " <> jsonMessage usageErr ExitDBFatalError ServerPgrstBug usageErr -> "This is probably a bug in PostgREST, please report it at https://github.com/PostgREST/postgrest/issues. " <> jsonMessage usageErr ExitDBFatalError ServerError42P05 usageErr -> @@ -95,8 +92,8 @@ observationMessage = \case "Schema cache loaded in " <> showMillis resultTime <> " milliseconds" ConnectionRetryObs delay -> "Attempting to reconnect to the database in " <> (show delay::Text) <> " seconds..." - ConnectionPgVersionErrorObs usageErr -> - jsonMessage usageErr + QueryPgVersionError usageErr -> + "Failed to query the PostgreSQL version. " <> jsonMessage usageErr DBListenStart channel -> do "Listening for notifications on the " <> show channel <> " channel" DBListenFail channel listenErr -> diff --git a/test/io/test_io.py b/test/io/test_io.py index 32919544a..d1dc347e4 100644 --- a/test/io/test_io.py +++ b/test/io/test_io.py @@ -778,7 +778,7 @@ def test_metrics_include_schema_cache_fails(defaultenv, metapostgrest): r'pgrst_schema_cache_loads_total{status="FAIL"} (\d+)', response.text ).group(1) ) - assert metrics > 3.0 + assert metrics == 1.0 reset_statement_timeout(metapostgrest, role)