From 50c3ba5c5da72bb758a4112363ba2fe1c0e968ea Mon Sep 17 00:00:00 2001 From: rnhmjoj Date: Wed, 5 Jun 2024 07:39:01 +0200 Subject: [PATCH 01/39] =?UTF-8?q?Fix=20build=20with=20mtl=20=E2=89=A5=202.?= =?UTF-8?q?3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit mtl stopped (re)exporting some functions, causing a build error. --- selda/src/Database/Selda/Backend/Internal.hs | 3 ++- selda/src/Database/Selda/Generic.hs | 3 ++- selda/src/Database/Selda/SQL/Print.hs | 3 ++- selda/src/Database/Selda/SqlRow.hs | 4 ++-- selda/src/Database/Selda/Unsafe.hs | 3 ++- 5 files changed, 10 insertions(+), 6 deletions(-) diff --git a/selda/src/Database/Selda/Backend/Internal.hs b/selda/src/Database/Selda/Backend/Internal.hs index b52fafa..c187dfe 100644 --- a/selda/src/Database/Selda/Backend/Internal.hs +++ b/selda/src/Database/Selda/Backend/Internal.hs @@ -40,11 +40,12 @@ import Database.Selda.SQL.Print.Config import Database.Selda.Types (TableName, ColName) import Data.Int (Int64) import Control.Concurrent ( newMVar, putMVar, takeMVar, MVar ) +import Control.Monad (when) import Control.Monad.Catch ( Exception, bracket, MonadCatch, MonadMask, MonadThrow(..) ) import Control.Monad.IO.Class ( MonadIO(..) ) import Control.Monad.Reader - ( MonadTrans(..), when, ReaderT(..), MonadReader(ask) ) + ( MonadTrans(..), ReaderT(..), MonadReader(ask) ) import Data.Dynamic ( Typeable, Dynamic ) import qualified Data.IntMap as M import Data.IORef diff --git a/selda/src/Database/Selda/Generic.hs b/selda/src/Database/Selda/Generic.hs index 9c18423..82a3fb9 100644 --- a/selda/src/Database/Selda/Generic.hs +++ b/selda/src/Database/Selda/Generic.hs @@ -8,8 +8,9 @@ module Database.Selda.Generic ( Relational, Generic , tblCols, params, def, gNew, gRow ) where +import Control.Monad (liftM2) import Control.Monad.State - ( liftM2, MonadState(put, get), evalState, State ) + ( MonadState(put, get), evalState, State ) import Data.Dynamic ( Typeable ) import Data.Text as Text (Text, pack) diff --git a/selda/src/Database/Selda/SQL/Print.hs b/selda/src/Database/Selda/SQL/Print.hs index cae9192..46d9509 100644 --- a/selda/src/Database/Selda/SQL/Print.hs +++ b/selda/src/Database/Selda/SQL/Print.hs @@ -16,8 +16,9 @@ import qualified Database.Selda.SQL.Print.Config as Cfg import Database.Selda.SqlType ( Lit(LJust, LNull), SqlTypeRep ) import Database.Selda.Types ( TableName, ColName, fromColName, fromTableName ) +import Control.Monad (liftM2) import Control.Monad.State - ( liftM2, MonadState(get, put), runState, State ) + ( MonadState(get, put), runState, State ) import Data.List ( group, sort ) import Data.Text (Text) import qualified Data.Text as Text diff --git a/selda/src/Database/Selda/SqlRow.hs b/selda/src/Database/Selda/SqlRow.hs index 4a52633..f8d153c 100644 --- a/selda/src/Database/Selda/SqlRow.hs +++ b/selda/src/Database/Selda/SqlRow.hs @@ -7,9 +7,9 @@ module Database.Selda.SqlRow , GSqlRow , runResultReader, next ) where +import Control.Monad (liftM2) import Control.Monad.State.Strict - ( liftM2, - StateT(StateT), + ( StateT(StateT), MonadState(state, get), State, evalState ) diff --git a/selda/src/Database/Selda/Unsafe.hs b/selda/src/Database/Selda/Unsafe.hs index 641674a..3ffac3e 100644 --- a/selda/src/Database/Selda/Unsafe.hs +++ b/selda/src/Database/Selda/Unsafe.hs @@ -9,8 +9,9 @@ module Database.Selda.Unsafe , QueryFragment, inj, injLit, rawName, rawExp, rawStm, rawQuery, rawQuery1 ) where import Control.Exception (throw) +import Control.Monad (void) import Control.Monad.State.Strict - ( MonadIO(liftIO), void, MonadState(put, get) ) + ( MonadIO(liftIO), MonadState(put, get) ) import Database.Selda.Backend.Internal ( SqlType(mkLit, sqlType), MonadSelda, From 603f81d8e674240ab6a102a2093e2beb287606b6 Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Wed, 10 Jul 2024 15:10:39 +0200 Subject: [PATCH 02/39] up time constraint to <1.15 --- selda-build-tools.cabal | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/selda-build-tools.cabal b/selda-build-tools.cabal index 000354d..61a12ff 100644 --- a/selda-build-tools.cabal +++ b/selda-build-tools.cabal @@ -10,7 +10,7 @@ executable selda-changelog main-is: ChangeLog.hs build-depends: base >=4.8 && <5, - time >=1.5 && <1.10, + time >=1.5 && <1.15, filepath >=1.4 && <1.5, process >=1.5 && <1.7 default-language: Haskell2010 From b19a7735bea43bca6dc896e730648eedd539f532 Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Wed, 10 Jul 2024 15:12:22 +0200 Subject: [PATCH 03/39] up filepath constraint to <1.6 --- selda-build-tools.cabal | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/selda-build-tools.cabal b/selda-build-tools.cabal index 61a12ff..6f2a8ce 100644 --- a/selda-build-tools.cabal +++ b/selda-build-tools.cabal @@ -11,6 +11,6 @@ executable selda-changelog build-depends: base >=4.8 && <5, time >=1.5 && <1.15, - filepath >=1.4 && <1.5, + filepath >=1.4 && <1.6, process >=1.5 && <1.7 default-language: Haskell2010 From 1f8555e281488106982b71f3fa2ca86fa79839d1 Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Wed, 10 Jul 2024 15:15:15 +0200 Subject: [PATCH 04/39] up time constraint to <1.15 --- selda/selda.cabal | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/selda/selda.cabal b/selda/selda.cabal index 83a6fbf..fc73d74 100644 --- a/selda/selda.cabal +++ b/selda/selda.cabal @@ -76,7 +76,7 @@ library , exceptions >=0.8 && <0.11 , mtl >=2.0 && <2.4 , text >=1.0 && <2.1 - , time >=1.5 && <1.13 + , time >=1.5 && <1.15 , containers >=0.4 && <0.7 , random >=1.1 && <1.3 , uuid-types >=1.0 && <1.1 From 83f9fa5353249f17ebb7d8fc195b7a363cf57e0b Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Wed, 10 Jul 2024 15:26:13 +0200 Subject: [PATCH 05/39] up time constraint to <1.15 --- selda-postgresql/selda-postgresql.cabal | 2 +- selda-sqlite/selda-sqlite.cabal | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/selda-postgresql/selda-postgresql.cabal b/selda-postgresql/selda-postgresql.cabal index 2fd0f2a..9e55301 100644 --- a/selda-postgresql/selda-postgresql.cabal +++ b/selda-postgresql/selda-postgresql.cabal @@ -37,7 +37,7 @@ library build-depends: postgresql-binary >=0.12 && <0.13 , postgresql-libpq >=0.9 && <0.10 - , time >=1.5 && <1.13 + , time >=1.5 && <1.15 , uuid-types >=1.0 && <1.1 hs-source-dirs: src diff --git a/selda-sqlite/selda-sqlite.cabal b/selda-sqlite/selda-sqlite.cabal index 288092a..bb07b47 100644 --- a/selda-sqlite/selda-sqlite.cabal +++ b/selda-sqlite/selda-sqlite.cabal @@ -33,7 +33,7 @@ library , direct-sqlite >=2.2 && <2.4 , directory >=1.2.2 && <1.4 , exceptions >=0.8 && <0.11 - , time >=1.5 && <1.13 + , time >=1.5 && <1.15 , uuid-types >=1.0 && <1.1 hs-source-dirs: src From c9f22a74dbfd2c36e641ba5d8946f45231e79162 Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Wed, 10 Jul 2024 15:31:26 +0200 Subject: [PATCH 06/39] bytestring constraint to <0.13 --- selda-json/selda-json.cabal | 2 +- selda-postgresql/selda-postgresql.cabal | 2 +- selda-sqlite/selda-sqlite.cabal | 2 +- selda-tests/selda-tests.cabal | 2 +- selda/selda.cabal | 2 +- 5 files changed, 5 insertions(+), 5 deletions(-) diff --git a/selda-json/selda-json.cabal b/selda-json/selda-json.cabal index 69ac5fc..0b58ee6 100644 --- a/selda-json/selda-json.cabal +++ b/selda-json/selda-json.cabal @@ -19,7 +19,7 @@ library build-depends: aeson >=1.0 && <2.2 , base >=4.9 && <5 - , bytestring >=0.10 && <0.12 + , bytestring >=0.10 && <0.13 , selda >=0.4 && <0.6 , text >=1.0 && <2.1 hs-source-dirs: src diff --git a/selda-postgresql/selda-postgresql.cabal b/selda-postgresql/selda-postgresql.cabal index 9e55301..438f385 100644 --- a/selda-postgresql/selda-postgresql.cabal +++ b/selda-postgresql/selda-postgresql.cabal @@ -28,7 +28,7 @@ library CPP build-depends: base >=4.9 && <5 - , bytestring >=0.9 && <0.12 + , bytestring >=0.9 && <0.13 , exceptions >=0.8 && <0.11 , selda >=0.5 && <0.6 , selda-json >=0.1 && <0.2 diff --git a/selda-sqlite/selda-sqlite.cabal b/selda-sqlite/selda-sqlite.cabal index bb07b47..032fcec 100644 --- a/selda-sqlite/selda-sqlite.cabal +++ b/selda-sqlite/selda-sqlite.cabal @@ -29,7 +29,7 @@ library , text >=1.0 && <2.1 if !flag(haste) build-depends: - bytestring >=0.10 && <0.12 + bytestring >=0.10 && <0.13 , direct-sqlite >=2.2 && <2.4 , directory >=1.2.2 && <1.4 , exceptions >=0.8 && <0.11 diff --git a/selda-tests/selda-tests.cabal b/selda-tests/selda-tests.cabal index 40ed106..542b914 100644 --- a/selda-tests/selda-tests.cabal +++ b/selda-tests/selda-tests.cabal @@ -34,7 +34,7 @@ test-suite selda-testsuite build-depends: aeson , base >=4.8 && <5 - , bytestring >=0.10 && <0.12 + , bytestring >=0.10 && <0.13 , directory >=1.2 && <1.4 , exceptions >=0.8 && <0.11 , HUnit >=1.4 && <1.7 diff --git a/selda/selda.cabal b/selda/selda.cabal index fc73d74..56f7155 100644 --- a/selda/selda.cabal +++ b/selda/selda.cabal @@ -72,7 +72,7 @@ library FlexibleContexts build-depends: base >=4.10 && <5 - , bytestring >=0.10 && <0.12 + , bytestring >=0.10 && <0.13 , exceptions >=0.8 && <0.11 , mtl >=2.0 && <2.4 , text >=1.0 && <2.1 From 6695b1cc52f525e2ce16eac281f091fcda62627b Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Mon, 28 Oct 2024 13:01:21 +0100 Subject: [PATCH 07/39] text constraint to <2.2 --- selda/selda.cabal | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/selda/selda.cabal b/selda/selda.cabal index 56f7155..1612047 100644 --- a/selda/selda.cabal +++ b/selda/selda.cabal @@ -75,7 +75,7 @@ library , bytestring >=0.10 && <0.13 , exceptions >=0.8 && <0.11 , mtl >=2.0 && <2.4 - , text >=1.0 && <2.1 + , text >=1.0 && <2.2 , time >=1.5 && <1.15 , containers >=0.4 && <0.7 , random >=1.1 && <1.3 From 8f677dd2b387ca18fd68e3f3b8bf3e250ae1893c Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Mon, 28 Oct 2024 13:01:55 +0100 Subject: [PATCH 08/39] text constraint to <2.2 --- selda-sqlite/selda-sqlite.cabal | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/selda-sqlite/selda-sqlite.cabal b/selda-sqlite/selda-sqlite.cabal index 032fcec..73d3536 100644 --- a/selda-sqlite/selda-sqlite.cabal +++ b/selda-sqlite/selda-sqlite.cabal @@ -26,7 +26,7 @@ library build-depends: base >=4.9 && <5 , selda >=0.5 && <0.6 - , text >=1.0 && <2.1 + , text >=1.0 && <2.2 if !flag(haste) build-depends: bytestring >=0.10 && <0.13 From 2d474614b30865f3bed63253c1557b5fd352f84f Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Thu, 7 Nov 2024 15:21:26 +0100 Subject: [PATCH 09/39] support containers 0.7 --- selda/selda.cabal | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/selda/selda.cabal b/selda/selda.cabal index 1612047..d25122c 100644 --- a/selda/selda.cabal +++ b/selda/selda.cabal @@ -77,7 +77,7 @@ library , mtl >=2.0 && <2.4 , text >=1.0 && <2.2 , time >=1.5 && <1.15 - , containers >=0.4 && <0.7 + , containers >=0.4 && <0.8 , random >=1.1 && <1.3 , uuid-types >=1.0 && <1.1 hs-source-dirs: From b8fcc0a56d9b4de85f726af4237755d8608fa659 Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Mon, 7 Jul 2025 13:16:16 +0200 Subject: [PATCH 10/39] wip streaming --- .../src/Database/Selda/PostgreSQL.hs | 17 +++---- selda-sqlite/src/Database/Selda/SQLite.hs | 45 +++++++++++++++---- selda/src/Database/Selda/Backend/Internal.hs | 7 ++- selda/src/Database/Selda/Frontend.hs | 6 ++- 4 files changed, 57 insertions(+), 18 deletions(-) diff --git a/selda-postgresql/src/Database/Selda/PostgreSQL.hs b/selda-postgresql/src/Database/Selda/PostgreSQL.hs index 79f83ec..da14956 100644 --- a/selda-postgresql/src/Database/Selda/PostgreSQL.hs +++ b/selda-postgresql/src/Database/Selda/PostgreSQL.hs @@ -200,14 +200,15 @@ pgPPConfig = defPPConfig pgBackend :: Connection -- ^ PostgreSQL connection object. -> SeldaBackend PG pgBackend c = SeldaBackend - { runStmt = \q ps -> right <$> pgQueryRunner c False q ps - , runStmtWithPK = \q ps -> left <$> pgQueryRunner c True q ps - , prepareStmt = pgPrepare c - , runPrepared = pgRun c - , getTableInfo = pgGetTableInfo c . rawTableName - , backendId = PostgreSQL - , ppConfig = pgPPConfig - , closeConnection = \_ -> finish c + { runStmt = \q ps -> right <$> pgQueryRunner c False q ps + , runStmtStreaming = undefined + , runStmtWithPK = \q ps -> left <$> pgQueryRunner c True q ps + , prepareStmt = pgPrepare c + , runPrepared = pgRun c + , getTableInfo = pgGetTableInfo c . rawTableName + , backendId = PostgreSQL + , ppConfig = pgPPConfig + , closeConnection = \_ -> finish c , disableForeignKeys = disableFKs c } where diff --git a/selda-sqlite/src/Database/Selda/SQLite.hs b/selda-sqlite/src/Database/Selda/SQLite.hs index 352ff34..55e4874 100644 --- a/selda-sqlite/src/Database/Selda/SQLite.hs +++ b/selda-sqlite/src/Database/Selda/SQLite.hs @@ -66,14 +66,15 @@ withSQLite file m = bracket (sqliteOpen file) seldaClose (runSeldaT m) -- Proceed with extreme caution. sqliteBackend :: Database -> SeldaBackend SQLite sqliteBackend db = SeldaBackend - { runStmt = \q ps -> snd <$> sqliteQueryRunner db q ps - , runStmtWithPK = \q ps -> fst <$> sqliteQueryRunner db q ps - , prepareStmt = \_ _ -> sqlitePrepare db - , runPrepared = sqliteRunPrepared db - , getTableInfo = sqliteGetTableInfo db . fromTableName - , ppConfig = defPPConfig {ppMaxInsertParams = Just 999} - , backendId = SQLite - , closeConnection = \conn -> do + { runStmt = \q ps -> snd <$> sqliteQueryRunner db q ps + , runStmtStreaming = \q ps -> snd <$> sqliteQueryRunnerStreaming db q ps + , runStmtWithPK = \q ps -> fst <$> sqliteQueryRunner db q ps + , prepareStmt = \_ _ -> sqlitePrepare db + , runPrepared = sqliteRunPrepared db + , getTableInfo = sqliteGetTableInfo db . fromTableName + , ppConfig = defPPConfig {ppMaxInsertParams = Just 999} + , backendId = SQLite + , closeConnection = \conn -> do stmts <- allStmts conn flip mapM_ stmts $ \(_, stm) -> do finalize $ fromDyn stm (error "BUG: non-statement SQLite statement") @@ -198,6 +199,24 @@ sqliteRunStmt db stm params = do cs <- changes db return (fromIntegral rid, (cs, [map fromSqlData r | r <- rows])) +sqliteQueryRunnerStreaming :: Database -> QueryRunner (Generator st) +sqliteQueryRunnerStreaming db qry params = do + eres <- try $ do + stm <- prepare db qry + sqliteRunStmtStreaming db stm params `finally` do + finalize stm + case eres of + Left e@(SQLError{}) -> throwM (SqlError (show e)) + Right res -> return res + +sqliteRunStmtStreaming :: Database -> Statement -> [Param] -> IO (Generator st) +sqliteRunStmtStreaming db stm params = do + bind stm [toSqlData p | Param p <- params] + rows <- streamRows stm [] + _rid <- lastInsertRowId db + !cs <- changes db + return rows -- this should be roughly correct + getRows :: Statement -> [[SQLData]] -> IO [[SQLData]] getRows s acc = do res <- step s @@ -208,6 +227,16 @@ getRows s acc = do _ -> do return $ reverse acc +streamRows :: Statement -> [[SQLData]] -> Generator st +streamRows s acc = \st -> do + res <- step s + case res of + Row -> do + cs <- columns s + yield cs + -- _ -> + -- what do we do instead of returning all accumulated values? + toSqlData :: Lit a -> SQLData toSqlData (LInt32 i) = SQLInteger $ fromIntegral i toSqlData (LInt64 i) = SQLInteger $ fromIntegral i diff --git a/selda/src/Database/Selda/Backend/Internal.hs b/selda/src/Database/Selda/Backend/Internal.hs index c187dfe..07bd205 100644 --- a/selda/src/Database/Selda/Backend/Internal.hs +++ b/selda/src/Database/Selda/Backend/Internal.hs @@ -211,11 +211,16 @@ tableInfo t = TableInfo ] ] +type Generator st = IO (st -> Maybe (st, [SqlValue])) + -- | A collection of functions making up a Selda backend. -data SeldaBackend b = SeldaBackend +data SeldaBackend b st = SeldaBackend { -- | Execute an SQL statement. runStmt :: Text -> [Param] -> IO (Int, [[SqlValue]]) + -- | Stream the result. + , runStmtStreaming :: Text -> [Param] -> Generator st + -- | Execute an SQL statement and return the last inserted primary key, -- where the primary key is auto-incrementing. -- Backends must take special care to make this thread-safe. diff --git a/selda/src/Database/Selda/Frontend.hs b/selda/src/Database/Selda/Frontend.hs index 247ae6d..8b1aa0b 100644 --- a/selda/src/Database/Selda/Frontend.hs +++ b/selda/src/Database/Selda/Frontend.hs @@ -15,7 +15,7 @@ import Database.Selda.Backend.Internal Param, SeldaT, MonadSelda(..), - SeldaBackend(runStmtWithPK, disableForeignKeys, ppConfig, runStmt), + SeldaBackend(runStmtWithPK, disableForeignKeys, ppConfig, runStmt, runStmtStreaming), QueryRunner, SeldaError(SqlError), withBackend ) @@ -58,6 +58,10 @@ import Control.Monad.IO.Class ( MonadIO(..) ) query :: (MonadSelda m, Result a) => Query (Backend m) a -> m [Res a] query q = withBackend (flip queryWith q . runStmt) +-- | Run a query within a Selda monad like `query` and stream the results. +queryStreaming :: (MonadSelda m, Result a) => Query (Backend m) a -> m (Generator st) +queryStreaming withBackend (flip _ q . runStmtStreaming) -- TODO: roll own queryWith to use here + -- | Perform the given query, and insert the result into the given table. -- Returns the number of inserted rows. queryInto :: (MonadSelda m, Relational a) From 677934580c64da8572817ccd0132e254002e1c86 Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Mon, 7 Jul 2025 15:47:35 +0200 Subject: [PATCH 11/39] wip streaming --- selda-sqlite/src/Database/Selda/SQLite.hs | 27 ++++++++++---------- selda/src/Database/Selda/Backend/Internal.hs | 4 +-- 2 files changed, 15 insertions(+), 16 deletions(-) diff --git a/selda-sqlite/src/Database/Selda/SQLite.hs b/selda-sqlite/src/Database/Selda/SQLite.hs index 55e4874..6e407d6 100644 --- a/selda-sqlite/src/Database/Selda/SQLite.hs +++ b/selda-sqlite/src/Database/Selda/SQLite.hs @@ -67,7 +67,7 @@ withSQLite file m = bracket (sqliteOpen file) seldaClose (runSeldaT m) sqliteBackend :: Database -> SeldaBackend SQLite sqliteBackend db = SeldaBackend { runStmt = \q ps -> snd <$> sqliteQueryRunner db q ps - , runStmtStreaming = \q ps -> snd <$> sqliteQueryRunnerStreaming db q ps + , runStmtStreaming = \q ps -> sqliteQueryRunnerStreaming db q ps , runStmtWithPK = \q ps -> fst <$> sqliteQueryRunner db q ps , prepareStmt = \_ _ -> sqlitePrepare db , runPrepared = sqliteRunPrepared db @@ -199,20 +199,20 @@ sqliteRunStmt db stm params = do cs <- changes db return (fromIntegral rid, (cs, [map fromSqlData r | r <- rows])) -sqliteQueryRunnerStreaming :: Database -> QueryRunner (Generator st) -sqliteQueryRunnerStreaming db qry params = do +sqliteQueryRunnerStreaming :: MonadIO m, Monoid a => (st -> m a) -> Database -> QueryRunner (m [SqlValue]) +sqliteQueryRunnerStreaming f db qry params = do eres <- try $ do stm <- prepare db qry - sqliteRunStmtStreaming db stm params `finally` do + sqliteRunStmtStreaming f db stm params `finally` do finalize stm case eres of Left e@(SQLError{}) -> throwM (SqlError (show e)) Right res -> return res -sqliteRunStmtStreaming :: Database -> Statement -> [Param] -> IO (Generator st) -sqliteRunStmtStreaming db stm params = do +sqliteRunStmtStreaming :: MonadIO m, Monoid a => (st -> m a) -> Database -> Statement -> [Param] -> IO (m [SqlValue]) +sqliteRunStmtStreaming f db stm params = do bind stm [toSqlData p | Param p <- params] - rows <- streamRows stm [] + rows <- streamRows f stm [] _rid <- lastInsertRowId db !cs <- changes db return rows -- this should be roughly correct @@ -227,15 +227,16 @@ getRows s acc = do _ -> do return $ reverse acc -streamRows :: Statement -> [[SQLData]] -> Generator st -streamRows s acc = \st -> do +-- TODO: should return `Generator [SqlData]` +streamRows :: MonadIO m, Monoid a => (st -> m a) -> Statement -> [[SQLData]] -> IO (m [SqlValue]) +streamRows f s acc = do res <- step s case res of Row -> do - cs <- columns s - yield cs - -- _ -> - -- what do we do instead of returning all accumulated values? + cs <- columns s -- :: IO [SQLData] + streamRows s `mappend` (f $ fromSqlData cs) `mappend` acc + Done -> do + return $ reverse acc toSqlData :: Lit a -> SQLData toSqlData (LInt32 i) = SQLInteger $ fromIntegral i diff --git a/selda/src/Database/Selda/Backend/Internal.hs b/selda/src/Database/Selda/Backend/Internal.hs index 07bd205..8bc84bc 100644 --- a/selda/src/Database/Selda/Backend/Internal.hs +++ b/selda/src/Database/Selda/Backend/Internal.hs @@ -211,15 +211,13 @@ tableInfo t = TableInfo ] ] -type Generator st = IO (st -> Maybe (st, [SqlValue])) - -- | A collection of functions making up a Selda backend. data SeldaBackend b st = SeldaBackend { -- | Execute an SQL statement. runStmt :: Text -> [Param] -> IO (Int, [[SqlValue]]) -- | Stream the result. - , runStmtStreaming :: Text -> [Param] -> Generator st + , runStmtStreaming :: MonadIO m => Text -> [Param] -> IO (m [SqlValue]) -- | Execute an SQL statement and return the last inserted primary key, -- where the primary key is auto-incrementing. From c6385914fadb3ce2f48190db8c06e4682250ebd2 Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Mon, 7 Jul 2025 22:10:23 +0200 Subject: [PATCH 12/39] wip streaming --- selda-postgresql/src/Database/Selda/PostgreSQL.hs | 2 +- selda-sqlite/src/Database/Selda/SQLite.hs | 15 +++++++-------- selda/src/Database/Selda/Backend/Internal.hs | 2 +- selda/src/Database/Selda/Frontend.hs | 14 ++++++++++++-- 4 files changed, 21 insertions(+), 12 deletions(-) diff --git a/selda-postgresql/src/Database/Selda/PostgreSQL.hs b/selda-postgresql/src/Database/Selda/PostgreSQL.hs index da14956..cd18efb 100644 --- a/selda-postgresql/src/Database/Selda/PostgreSQL.hs +++ b/selda-postgresql/src/Database/Selda/PostgreSQL.hs @@ -201,7 +201,7 @@ pgBackend :: Connection -- ^ PostgreSQL connection object. -> SeldaBackend PG pgBackend c = SeldaBackend { runStmt = \q ps -> right <$> pgQueryRunner c False q ps - , runStmtStreaming = undefined + , runStmtStreaming = \f q ps -> undefined , runStmtWithPK = \q ps -> left <$> pgQueryRunner c True q ps , prepareStmt = pgPrepare c , runPrepared = pgRun c diff --git a/selda-sqlite/src/Database/Selda/SQLite.hs b/selda-sqlite/src/Database/Selda/SQLite.hs index 6e407d6..f4c4a30 100644 --- a/selda-sqlite/src/Database/Selda/SQLite.hs +++ b/selda-sqlite/src/Database/Selda/SQLite.hs @@ -67,7 +67,7 @@ withSQLite file m = bracket (sqliteOpen file) seldaClose (runSeldaT m) sqliteBackend :: Database -> SeldaBackend SQLite sqliteBackend db = SeldaBackend { runStmt = \q ps -> snd <$> sqliteQueryRunner db q ps - , runStmtStreaming = \q ps -> sqliteQueryRunnerStreaming db q ps + , runStmtStreaming = \f q ps -> snd <$> sqliteQueryRunnerStreaming f db q ps , runStmtWithPK = \q ps -> fst <$> sqliteQueryRunner db q ps , prepareStmt = \_ _ -> sqlitePrepare db , runPrepared = sqliteRunPrepared db @@ -199,7 +199,7 @@ sqliteRunStmt db stm params = do cs <- changes db return (fromIntegral rid, (cs, [map fromSqlData r | r <- rows])) -sqliteQueryRunnerStreaming :: MonadIO m, Monoid a => (st -> m a) -> Database -> QueryRunner (m [SqlValue]) +sqliteQueryRunnerStreaming :: MonadIO m, Monoid a => (SqlValue -> m a) -> Database -> QueryRunner (Int64, (Int, m [SqlValue])) sqliteQueryRunnerStreaming f db qry params = do eres <- try $ do stm <- prepare db qry @@ -209,13 +209,13 @@ sqliteQueryRunnerStreaming f db qry params = do Left e@(SQLError{}) -> throwM (SqlError (show e)) Right res -> return res -sqliteRunStmtStreaming :: MonadIO m, Monoid a => (st -> m a) -> Database -> Statement -> [Param] -> IO (m [SqlValue]) +sqliteRunStmtStreaming :: MonadIO m, Monoid a => (SqlValue -> m a) -> Database -> Statement -> [Param] -> IO (Int64, (Int, m [SqlValue])) sqliteRunStmtStreaming f db stm params = do bind stm [toSqlData p | Param p <- params] rows <- streamRows f stm [] - _rid <- lastInsertRowId db - !cs <- changes db - return rows -- this should be roughly correct + rid <- lastInsertRowId db + cs <- changes db + return (fromIntegral rid, (cs, rows)) getRows :: Statement -> [[SQLData]] -> IO [[SQLData]] getRows s acc = do @@ -227,8 +227,7 @@ getRows s acc = do _ -> do return $ reverse acc --- TODO: should return `Generator [SqlData]` -streamRows :: MonadIO m, Monoid a => (st -> m a) -> Statement -> [[SQLData]] -> IO (m [SqlValue]) +streamRows :: MonadIO m, Monoid a => (SqlValue -> m a) -> Statement -> [[SQLData]] -> IO (m [SqlValue]) streamRows f s acc = do res <- step s case res of diff --git a/selda/src/Database/Selda/Backend/Internal.hs b/selda/src/Database/Selda/Backend/Internal.hs index 8bc84bc..638d2c7 100644 --- a/selda/src/Database/Selda/Backend/Internal.hs +++ b/selda/src/Database/Selda/Backend/Internal.hs @@ -217,7 +217,7 @@ data SeldaBackend b st = SeldaBackend runStmt :: Text -> [Param] -> IO (Int, [[SqlValue]]) -- | Stream the result. - , runStmtStreaming :: MonadIO m => Text -> [Param] -> IO (m [SqlValue]) + , runStmtStreaming :: MonadIO m => (SqlValue -> m a) -> Text -> [Param] -> IO (m [SqlValue]) -- | Execute an SQL statement and return the last inserted primary key, -- where the primary key is auto-incrementing. diff --git a/selda/src/Database/Selda/Frontend.hs b/selda/src/Database/Selda/Frontend.hs index 8b1aa0b..cea3336 100644 --- a/selda/src/Database/Selda/Frontend.hs +++ b/selda/src/Database/Selda/Frontend.hs @@ -59,8 +59,8 @@ query :: (MonadSelda m, Result a) => Query (Backend m) a -> m [Res a] query q = withBackend (flip queryWith q . runStmt) -- | Run a query within a Selda monad like `query` and stream the results. -queryStreaming :: (MonadSelda m, Result a) => Query (Backend m) a -> m (Generator st) -queryStreaming withBackend (flip _ q . runStmtStreaming) -- TODO: roll own queryWith to use here +queryStream :: (MonadSelda m, MonadIO m2, Result a, Monoid a) => (m2 b -> m2 a) -> (SqlValue -> m2 a) -> Query (Backend m) a -> m (m2 a) +queryStream ms f q = withBackend (flip queryStreamWith ms q . runStmtStreaming f) -- | Perform the given query, and insert the result into the given table. -- Returns the number of inserted rows. @@ -293,10 +293,20 @@ queryWith run q = withBackend $ \b -> do res <- fmap snd . liftIO . uncurry run $ compileWith (ppConfig b) q return $ mkResults (Proxy :: Proxy a) res +-- | Build the final result from streaming result columns. +queryStreamWith :: forall m a. (MonadSelda m, MonadIO m2, Result a, Monoid a) + => (m2 b -> m2 a) -> QueryRunner (Int, m2 [SqlValue]) -> Query (Backend m) a -> m (m2 (Res a)) +queryStreamWith ms run q = withBackend $ \b -> do + res <- fmap snd . liftIO . uncurry run $ compileWith (ppConfig b) q + return $ mkResultsStreaming ms (Proxy :: Proxy a) res -- FIXME: as res is the generator, we ought not just use ordinary map (which mkResults does) + -- | Generate the final result of a query from a list of untyped result rows. mkResults :: Result a => Proxy a -> [[SqlValue]] -> [Res a] mkResults p = map (buildResult p) +mkResultsStreaming :: (MonadIO m, Result a) => (m b -> m a) -> Proxy a -> m [SqlValue] -> m (Res a) +mkResultsStreaming ms p = ms (buildResult p) + {-# INLINE exec #-} -- | Execute a statement without a result. exec :: MonadSelda m => Text -> [Param] -> m Int From 986cbf0cc5e7b6e57eeb731ceae3e7d463e80f3b Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Wed, 9 Jul 2025 16:55:56 +0200 Subject: [PATCH 13/39] wip streaming; define Generator type synonyme to factor out type constraint from constructor level --- selda-json/selda-json.cabal | 4 +-- selda-postgresql/selda-postgresql.cabal | 2 +- selda/src/Database/Selda/Backend/Internal.hs | 32 +++++++++++--------- 3 files changed, 20 insertions(+), 18 deletions(-) diff --git a/selda-json/selda-json.cabal b/selda-json/selda-json.cabal index 0b58ee6..e6ab4ed 100644 --- a/selda-json/selda-json.cabal +++ b/selda-json/selda-json.cabal @@ -17,10 +17,10 @@ library exposed-modules: Database.Selda.JSON build-depends: - aeson >=1.0 && <2.2 + aeson >=1.0 && <2.3 , base >=4.9 && <5 , bytestring >=0.10 && <0.13 , selda >=0.4 && <0.6 - , text >=1.0 && <2.1 + , text >=1.0 && <2.2 hs-source-dirs: src default-language: Haskell2010 diff --git a/selda-postgresql/selda-postgresql.cabal b/selda-postgresql/selda-postgresql.cabal index 438f385..edab9d1 100644 --- a/selda-postgresql/selda-postgresql.cabal +++ b/selda-postgresql/selda-postgresql.cabal @@ -32,7 +32,7 @@ library , exceptions >=0.8 && <0.11 , selda >=0.5 && <0.6 , selda-json >=0.1 && <0.2 - , text >=1.0 && <2.1 + , text >=1.0 && <2.2 if !flag(haste) build-depends: postgresql-binary >=0.12 && <0.13 diff --git a/selda/src/Database/Selda/Backend/Internal.hs b/selda/src/Database/Selda/Backend/Internal.hs index 638d2c7..c9836ba 100644 --- a/selda/src/Database/Selda/Backend/Internal.hs +++ b/selda/src/Database/Selda/Backend/Internal.hs @@ -108,9 +108,9 @@ data SeldaStmt = SeldaStmt , stmtParams :: ![Either Int Param] } -data SeldaConnection b = SeldaConnection +data SeldaConnection b st = SeldaConnection { -- | The backend used by the current connection. - connBackend :: !(SeldaBackend b) + connBackend :: !(SeldaBackend b st) -- | A string uniquely identifying the database used by this connection. -- This could be, for instance, a PostgreSQL connection @@ -130,7 +130,7 @@ data SeldaConnection b = SeldaConnection -- | Create a new Selda connection for the given backend and database -- identifier string. -newConnection :: MonadIO m => SeldaBackend b -> Text -> m (SeldaConnection b) +newConnection :: MonadIO m => SeldaBackend b st -> Text -> m (SeldaConnection b st) newConnection back dbid = liftIO $ SeldaConnection back dbid <$> newIORef M.empty <*> newIORef False @@ -138,7 +138,7 @@ newConnection back dbid = -- | Get all statements and their corresponding identifiers for the current -- connection. -allStmts :: SeldaConnection b -> IO [(StmtID, Dynamic)] +allStmts :: SeldaConnection b st -> IO [(StmtID, Dynamic)] allStmts = fmap (map (\(k, v) -> (StmtID k, stmtHandle v)) . M.toList) . readIORef . connStmts @@ -211,13 +211,15 @@ tableInfo t = TableInfo ] ] +type Generator st = IO (st -> Maybe (st, [SqlValue])) + -- | A collection of functions making up a Selda backend. data SeldaBackend b st = SeldaBackend { -- | Execute an SQL statement. runStmt :: Text -> [Param] -> IO (Int, [[SqlValue]]) -- | Stream the result. - , runStmtStreaming :: MonadIO m => (SqlValue -> m a) -> Text -> [Param] -> IO (m [SqlValue]) + , runStmtStreaming :: Text -> [Param] -> Generator st -- | Execute an SQL statement and return the last inserted primary key, -- where the primary key is auto-incrementing. @@ -239,7 +241,7 @@ data SeldaBackend b st = SeldaBackend , ppConfig :: PPConfig -- | Close the currently open connection. - , closeConnection :: SeldaConnection b -> IO () + , closeConnection :: SeldaConnection b st -> IO () -- | Unique identifier for this backend. , backendId :: BackendID @@ -267,7 +269,7 @@ class MonadIO m => MonadSelda m where -- Selda computation. -- Thus, the computation must take care never to return or otherwise -- access the connection after returning. - withConnection :: (SeldaConnection (Backend m) -> m a) -> m a + withConnection :: (SeldaConnection (Backend m) st -> m a) -> m a -- | Perform the given computation as a transaction. -- Implementations must ensure that subsequent calls to 'withConnection' @@ -277,30 +279,30 @@ class MonadIO m => MonadSelda m where transact = id -- | Get the backend in use by the computation. -withBackend :: MonadSelda m => (SeldaBackend (Backend m) -> m a) -> m a +withBackend :: MonadSelda m => (SeldaBackend (Backend m) st -> m a) -> m a withBackend m = withConnection (m . connBackend) -- | Monad transformer adding Selda SQL capabilities. -newtype SeldaT b m a = S {unS :: ReaderT (SeldaConnection b) m a} +newtype SeldaT b st m a = S {unS :: ReaderT (SeldaConnection b st) m a} deriving ( Functor, Applicative, Monad, MonadIO , MonadThrow, MonadCatch, MonadMask , MonadFail ) -instance (MonadIO m, MonadMask m) => MonadSelda (SeldaT b m) where - type Backend (SeldaT b m) = b +instance (MonadIO m, MonadMask m) => MonadSelda (SeldaT b st m) where + type Backend (SeldaT b st m) = b withConnection m = S ask >>= m -instance MonadTrans (SeldaT b) where +instance MonadTrans (SeldaT b st) where lift = S . lift -- | The simplest form of Selda computation; 'SeldaT' specialized to 'IO'. -type SeldaM b = SeldaT b IO +type SeldaM b st = SeldaT b st IO -- | Run a Selda transformer. Backends should use this to implement their -- @withX@ functions. runSeldaT :: (MonadIO m, MonadMask m) - => SeldaT b m a - -> SeldaConnection b + => SeldaT b st m a + -> SeldaConnection b st -> m a runSeldaT m c = bracket (liftIO $ takeMVar (connLock c)) From 699207d719d0cf0ab9ee34f4eb3681ad4e5405a6 Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Fri, 11 Jul 2025 17:40:43 +0200 Subject: [PATCH 14/39] wip streaming --- selda/src/Database/Selda/Backend/Internal.hs | 47 +++++++++++--------- 1 file changed, 26 insertions(+), 21 deletions(-) diff --git a/selda/src/Database/Selda/Backend/Internal.hs b/selda/src/Database/Selda/Backend/Internal.hs index c9836ba..7ee6550 100644 --- a/selda/src/Database/Selda/Backend/Internal.hs +++ b/selda/src/Database/Selda/Backend/Internal.hs @@ -1,4 +1,4 @@ -{-# LANGUAGE GeneralizedNewtypeDeriving, CPP, TypeFamilies #-} +{-# LANGUAGE GeneralizedNewtypeDeriving, CPP, TypeFamilies, ScopedTypeVariables #-} -- | Internal backend API. -- Using anything exported from this module may or may not invalidate any -- safety guarantees made by Selda; use at your own peril. @@ -108,9 +108,9 @@ data SeldaStmt = SeldaStmt , stmtParams :: ![Either Int Param] } -data SeldaConnection b st = SeldaConnection +data SeldaConnection b = SeldaConnection { -- | The backend used by the current connection. - connBackend :: !(SeldaBackend b st) + connBackend :: !(SeldaBackend b) -- | A string uniquely identifying the database used by this connection. -- This could be, for instance, a PostgreSQL connection @@ -130,7 +130,7 @@ data SeldaConnection b st = SeldaConnection -- | Create a new Selda connection for the given backend and database -- identifier string. -newConnection :: MonadIO m => SeldaBackend b st -> Text -> m (SeldaConnection b st) +newConnection :: MonadIO m => SeldaBackend b -> Text -> m (SeldaConnection b s) newConnection back dbid = liftIO $ SeldaConnection back dbid <$> newIORef M.empty <*> newIORef False @@ -138,7 +138,7 @@ newConnection back dbid = -- | Get all statements and their corresponding identifiers for the current -- connection. -allStmts :: SeldaConnection b st -> IO [(StmtID, Dynamic)] +allStmts :: SeldaConnection b -> IO [(StmtID, Dynamic)] allStmts = fmap (map (\(k, v) -> (StmtID k, stmtHandle v)) . M.toList) . readIORef . connStmts @@ -211,15 +211,16 @@ tableInfo t = TableInfo ] ] -type Generator st = IO (st -> Maybe (st, [SqlValue])) +type Generator b = IO (b-> Maybe (b, [SqlValue])) -- | A collection of functions making up a Selda backend. -data SeldaBackend b st = SeldaBackend +data SeldaBackend b + = SeldaBackend { -- | Execute an SQL statement. runStmt :: Text -> [Param] -> IO (Int, [[SqlValue]]) -- | Stream the result. - , runStmtStreaming :: Text -> [Param] -> Generator st + , runStmtStreaming :: Text -> [Param] -> Generator b -- | Execute an SQL statement and return the last inserted primary key, -- where the primary key is auto-incrementing. @@ -241,7 +242,7 @@ data SeldaBackend b st = SeldaBackend , ppConfig :: PPConfig -- | Close the currently open connection. - , closeConnection :: SeldaConnection b st -> IO () + , closeConnection :: SeldaConnection b -> IO () -- | Unique identifier for this backend. , backendId :: BackendID @@ -260,8 +261,11 @@ data SeldaBackend b st = SeldaBackend class MonadIO m => MonadSelda m where {-# MINIMAL withConnection #-} + -- FIXME: how to define functional dependency `| st -> b` + type StreamingState b + -- | Type of database backend used by @m@. - type Backend m + type Backend m st -- | Pass a Selda connection to the given computation and execute it. -- After the computation finishes, @withConnection@ is free to do anything @@ -269,7 +273,7 @@ class MonadIO m => MonadSelda m where -- Selda computation. -- Thus, the computation must take care never to return or otherwise -- access the connection after returning. - withConnection :: (SeldaConnection (Backend m) st -> m a) -> m a + withConnection :: (SeldaConnection (Backend m st) -> m ()) -> m () -- | Perform the given computation as a transaction. -- Implementations must ensure that subsequent calls to 'withConnection' @@ -279,30 +283,31 @@ class MonadIO m => MonadSelda m where transact = id -- | Get the backend in use by the computation. -withBackend :: MonadSelda m => (SeldaBackend (Backend m) st -> m a) -> m a +withBackend :: MonadSelda m => (SeldaBackend (Backend m st) -> m ()) -> m () withBackend m = withConnection (m . connBackend) -- | Monad transformer adding Selda SQL capabilities. -newtype SeldaT b st m a = S {unS :: ReaderT (SeldaConnection b st) m a} +newtype SeldaT b m a = S {unS :: ReaderT (SeldaConnection b) m a} deriving ( Functor, Applicative, Monad, MonadIO , MonadThrow, MonadCatch, MonadMask , MonadFail ) -instance (MonadIO m, MonadMask m) => MonadSelda (SeldaT b st m) where - type Backend (SeldaT b st m) = b - withConnection m = S ask >>= m - -instance MonadTrans (SeldaT b st) where +instance (MonadIO m, MonadMask m) => MonadSelda (SeldaT b m) where + type Backend (SeldaT b m) st = b + withConnection m = S ask >>= _ + -- (SeldaConnection (Backend m) -> m ()) + +instance MonadTrans (SeldaT b) where lift = S . lift -- | The simplest form of Selda computation; 'SeldaT' specialized to 'IO'. -type SeldaM b st = SeldaT b st IO +type SeldaM b = SeldaT b IO -- | Run a Selda transformer. Backends should use this to implement their -- @withX@ functions. runSeldaT :: (MonadIO m, MonadMask m) - => SeldaT b st m a - -> SeldaConnection b st + => SeldaT b m a + -> SeldaConnection b -> m a runSeldaT m c = bracket (liftIO $ takeMVar (connLock c)) From 6a2c0dbdd152d337e95bc7b3444ce78b3d11a135 Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Fri, 11 Jul 2025 21:31:53 +0200 Subject: [PATCH 15/39] wip streaming --- selda/src/Database/Selda/Backend.hs | 2 ++ selda/src/Database/Selda/Backend/Internal.hs | 18 +++++++++--------- selda/src/Database/Selda/Frontend.hs | 13 +++++++------ 3 files changed, 18 insertions(+), 15 deletions(-) diff --git a/selda/src/Database/Selda/Backend.hs b/selda/src/Database/Selda/Backend.hs index 8e0bc8d..34dec64 100644 --- a/selda/src/Database/Selda/Backend.hs +++ b/selda/src/Database/Selda/Backend.hs @@ -3,6 +3,7 @@ module Database.Selda.Backend ( MonadSelda (..), SeldaT, SeldaM, SeldaError (..) , StmtID, BackendID (..), QueryRunner, SeldaBackend (..), SeldaConnection + , Generator , SqlValue (..) , IndexMethod (..) , Param (..), ColAttr (..), AutoIncType (..) @@ -32,6 +33,7 @@ import Database.Selda.Backend.Internal SeldaT, MonadSelda(..), SeldaBackend(..), + Generator, ColumnInfo(..), TableInfo(..), SeldaConnection(connClosed, connBackend), diff --git a/selda/src/Database/Selda/Backend/Internal.hs b/selda/src/Database/Selda/Backend/Internal.hs index 7ee6550..71919f7 100644 --- a/selda/src/Database/Selda/Backend/Internal.hs +++ b/selda/src/Database/Selda/Backend/Internal.hs @@ -5,6 +5,7 @@ module Database.Selda.Backend.Internal ( StmtID (..), BackendID (..) , QueryRunner, SeldaBackend (..), SeldaConnection (..), SeldaStmt (..) + , Generator , MonadSelda (..), SeldaT (..), SeldaM , SeldaError (..) , Param (..), Lit (..), ColAttr (..), AutoIncType (..) @@ -130,7 +131,7 @@ data SeldaConnection b = SeldaConnection -- | Create a new Selda connection for the given backend and database -- identifier string. -newConnection :: MonadIO m => SeldaBackend b -> Text -> m (SeldaConnection b s) +newConnection :: MonadIO m => SeldaBackend b -> Text -> m (SeldaConnection b) newConnection back dbid = liftIO $ SeldaConnection back dbid <$> newIORef M.empty <*> newIORef False @@ -211,7 +212,7 @@ tableInfo t = TableInfo ] ] -type Generator b = IO (b-> Maybe (b, [SqlValue])) +type Generator b = IO (b -> Maybe (b, [SqlValue])) -- | A collection of functions making up a Selda backend. data SeldaBackend b @@ -262,10 +263,10 @@ class MonadIO m => MonadSelda m where {-# MINIMAL withConnection #-} -- FIXME: how to define functional dependency `| st -> b` - type StreamingState b + type StreamingState m -- | Type of database backend used by @m@. - type Backend m st + type Backend m -- | Pass a Selda connection to the given computation and execute it. -- After the computation finishes, @withConnection@ is free to do anything @@ -273,7 +274,7 @@ class MonadIO m => MonadSelda m where -- Selda computation. -- Thus, the computation must take care never to return or otherwise -- access the connection after returning. - withConnection :: (SeldaConnection (Backend m st) -> m ()) -> m () + withConnection :: (SeldaConnection (Backend m) -> m a) -> m a -- | Perform the given computation as a transaction. -- Implementations must ensure that subsequent calls to 'withConnection' @@ -283,7 +284,7 @@ class MonadIO m => MonadSelda m where transact = id -- | Get the backend in use by the computation. -withBackend :: MonadSelda m => (SeldaBackend (Backend m st) -> m ()) -> m () +withBackend :: MonadSelda m => (SeldaBackend (Backend m) -> m a) -> m a withBackend m = withConnection (m . connBackend) -- | Monad transformer adding Selda SQL capabilities. @@ -293,9 +294,8 @@ newtype SeldaT b m a = S {unS :: ReaderT (SeldaConnection b) m a} ) instance (MonadIO m, MonadMask m) => MonadSelda (SeldaT b m) where - type Backend (SeldaT b m) st = b - withConnection m = S ask >>= _ - -- (SeldaConnection (Backend m) -> m ()) + type Backend (SeldaT b m) = b + withConnection m = S ask >>= m instance MonadTrans (SeldaT b) where lift = S . lift diff --git a/selda/src/Database/Selda/Frontend.hs b/selda/src/Database/Selda/Frontend.hs index cea3336..6aaaa7e 100644 --- a/selda/src/Database/Selda/Frontend.hs +++ b/selda/src/Database/Selda/Frontend.hs @@ -18,6 +18,7 @@ import Database.Selda.Backend.Internal SeldaBackend(runStmtWithPK, disableForeignKeys, ppConfig, runStmt, runStmtStreaming), QueryRunner, SeldaError(SqlError), + Generator, withBackend ) import Database.Selda.Column ( Row, Col ) import Database.Selda.Compile @@ -294,18 +295,18 @@ queryWith run q = withBackend $ \b -> do return $ mkResults (Proxy :: Proxy a) res -- | Build the final result from streaming result columns. -queryStreamWith :: forall m a. (MonadSelda m, MonadIO m2, Result a, Monoid a) - => (m2 b -> m2 a) -> QueryRunner (Int, m2 [SqlValue]) -> Query (Backend m) a -> m (m2 (Res a)) -queryStreamWith ms run q = withBackend $ \b -> do +queryStreamWith :: forall m a. (MonadSelda m, Result a, Monoid a) + => QueryRunner (Int, Generator [SqlValue]) -> Query (Backend m) a -> m (Generator (Res a)) +queryStreamWith run q = withBackend $ \b -> do res <- fmap snd . liftIO . uncurry run $ compileWith (ppConfig b) q - return $ mkResultsStreaming ms (Proxy :: Proxy a) res -- FIXME: as res is the generator, we ought not just use ordinary map (which mkResults does) + return $ _ -- mkResultsStreaming ms (Proxy :: Proxy a) res -- FIXME: as res is the generator, we ought not just use ordinary map (which mkResults does) -- | Generate the final result of a query from a list of untyped result rows. mkResults :: Result a => Proxy a -> [[SqlValue]] -> [Res a] mkResults p = map (buildResult p) -mkResultsStreaming :: (MonadIO m, Result a) => (m b -> m a) -> Proxy a -> m [SqlValue] -> m (Res a) -mkResultsStreaming ms p = ms (buildResult p) +mkResultsStreaming :: (MonadIO m, Result a) => Proxy a -> m [SqlValue] -> m (Generator (Res a)) +mkResultsStreaming p xs = return $ \ms -> ms (buildResult p) xs {-# INLINE exec #-} -- | Execute a statement without a result. From 84547944a5138da5b580daa7168152e0b05db4cd Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Mon, 14 Jul 2025 16:40:48 +0200 Subject: [PATCH 16/39] SQLite streaming compiles; next: write tests --- selda-sqlite/src/Database/Selda/SQLite.hs | 27 ++++++++++---------- selda/src/Database/Selda/Backend.hs | 2 -- selda/src/Database/Selda/Backend/Internal.hs | 12 +++------ selda/src/Database/Selda/Frontend.hs | 19 +++++++------- 4 files changed, 25 insertions(+), 35 deletions(-) diff --git a/selda-sqlite/src/Database/Selda/SQLite.hs b/selda-sqlite/src/Database/Selda/SQLite.hs index f4c4a30..648e508 100644 --- a/selda-sqlite/src/Database/Selda/SQLite.hs +++ b/selda-sqlite/src/Database/Selda/SQLite.hs @@ -199,7 +199,15 @@ sqliteRunStmt db stm params = do cs <- changes db return (fromIntegral rid, (cs, [map fromSqlData r | r <- rows])) -sqliteQueryRunnerStreaming :: MonadIO m, Monoid a => (SqlValue -> m a) -> Database -> QueryRunner (Int64, (Int, m [SqlValue])) +sqliteRunStmtStreaming :: (Monad m, Monoid (m [SqlValue])) => ([SqlValue] -> m [SqlValue]) -> Database -> Statement -> [Param] -> IO (Int64, (Int, m [SqlValue])) +sqliteRunStmtStreaming f db stm params = do + bind stm [toSqlData p | Param p <- params] + rows <- streamRows f stm mempty + rid <- lastInsertRowId db + cs <- changes db + return (fromIntegral rid, (cs, rows)) + +sqliteQueryRunnerStreaming :: (Monad m, Monoid (m [SqlValue])) => ([SqlValue] -> m [SqlValue]) -> Database -> QueryRunner (Int64, (Int, m [SqlValue])) sqliteQueryRunnerStreaming f db qry params = do eres <- try $ do stm <- prepare db qry @@ -209,14 +217,6 @@ sqliteQueryRunnerStreaming f db qry params = do Left e@(SQLError{}) -> throwM (SqlError (show e)) Right res -> return res -sqliteRunStmtStreaming :: MonadIO m, Monoid a => (SqlValue -> m a) -> Database -> Statement -> [Param] -> IO (Int64, (Int, m [SqlValue])) -sqliteRunStmtStreaming f db stm params = do - bind stm [toSqlData p | Param p <- params] - rows <- streamRows f stm [] - rid <- lastInsertRowId db - cs <- changes db - return (fromIntegral rid, (cs, rows)) - getRows :: Statement -> [[SQLData]] -> IO [[SQLData]] getRows s acc = do res <- step s @@ -227,15 +227,14 @@ getRows s acc = do _ -> do return $ reverse acc -streamRows :: MonadIO m, Monoid a => (SqlValue -> m a) -> Statement -> [[SQLData]] -> IO (m [SqlValue]) +streamRows :: (Monad m, Monoid (m [SqlValue])) => ([SqlValue]-> m [SqlValue]) -> Statement -> m [SqlValue]-> IO (m [SqlValue]) streamRows f s acc = do res <- step s case res of Row -> do - cs <- columns s -- :: IO [SQLData] - streamRows s `mappend` (f $ fromSqlData cs) `mappend` acc - Done -> do - return $ reverse acc + cs <- map fromSqlData <$> liftIO (columns s :: IO [SQLData]) + streamRows f s (acc `mappend` f cs) + _ -> return acc toSqlData :: Lit a -> SQLData toSqlData (LInt32 i) = SQLInteger $ fromIntegral i diff --git a/selda/src/Database/Selda/Backend.hs b/selda/src/Database/Selda/Backend.hs index 34dec64..8e0bc8d 100644 --- a/selda/src/Database/Selda/Backend.hs +++ b/selda/src/Database/Selda/Backend.hs @@ -3,7 +3,6 @@ module Database.Selda.Backend ( MonadSelda (..), SeldaT, SeldaM, SeldaError (..) , StmtID, BackendID (..), QueryRunner, SeldaBackend (..), SeldaConnection - , Generator , SqlValue (..) , IndexMethod (..) , Param (..), ColAttr (..), AutoIncType (..) @@ -33,7 +32,6 @@ import Database.Selda.Backend.Internal SeldaT, MonadSelda(..), SeldaBackend(..), - Generator, ColumnInfo(..), TableInfo(..), SeldaConnection(connClosed, connBackend), diff --git a/selda/src/Database/Selda/Backend/Internal.hs b/selda/src/Database/Selda/Backend/Internal.hs index 71919f7..8798d43 100644 --- a/selda/src/Database/Selda/Backend/Internal.hs +++ b/selda/src/Database/Selda/Backend/Internal.hs @@ -1,11 +1,10 @@ -{-# LANGUAGE GeneralizedNewtypeDeriving, CPP, TypeFamilies, ScopedTypeVariables #-} +{-# LANGUAGE GeneralizedNewtypeDeriving, CPP, TypeFamilies, ScopedTypeVariables, RankNTypes #-} -- | Internal backend API. -- Using anything exported from this module may or may not invalidate any -- safety guarantees made by Selda; use at your own peril. module Database.Selda.Backend.Internal ( StmtID (..), BackendID (..) , QueryRunner, SeldaBackend (..), SeldaConnection (..), SeldaStmt (..) - , Generator , MonadSelda (..), SeldaT (..), SeldaM , SeldaError (..) , Param (..), Lit (..), ColAttr (..), AutoIncType (..) @@ -19,6 +18,7 @@ module Database.Selda.Backend.Internal , runSeldaT, withBackend ) where import Data.List (nub) +import Data.Proxy ( Proxy(..) ) import Database.Selda.SQL (Param (..)) import Database.Selda.SqlType ( SqlValue(..), @@ -212,16 +212,13 @@ tableInfo t = TableInfo ] ] -type Generator b = IO (b -> Maybe (b, [SqlValue])) - -- | A collection of functions making up a Selda backend. data SeldaBackend b = SeldaBackend { -- | Execute an SQL statement. runStmt :: Text -> [Param] -> IO (Int, [[SqlValue]]) - -- | Stream the result. - , runStmtStreaming :: Text -> [Param] -> Generator b + , runStmtStreaming :: forall m. (Monad m, MonadIO m, Monoid (m [SqlValue])) => ([SqlValue] -> m [SqlValue]) -> Text -> [Param] -> IO (Int, m [SqlValue]) -- | Execute an SQL statement and return the last inserted primary key, -- where the primary key is auto-incrementing. @@ -262,9 +259,6 @@ data SeldaBackend b class MonadIO m => MonadSelda m where {-# MINIMAL withConnection #-} - -- FIXME: how to define functional dependency `| st -> b` - type StreamingState m - -- | Type of database backend used by @m@. type Backend m diff --git a/selda/src/Database/Selda/Frontend.hs b/selda/src/Database/Selda/Frontend.hs index 6aaaa7e..ad745a4 100644 --- a/selda/src/Database/Selda/Frontend.hs +++ b/selda/src/Database/Selda/Frontend.hs @@ -2,7 +2,7 @@ -- | API for running Selda operations over databases. module Database.Selda.Frontend ( Result, Res, MonadIO (..), MonadSelda (..), SeldaT, OnError (..) - , query, queryInto + , query, queryInto, queryStream , insert, insert_, insertWithPK, tryInsert, insertWhen, insertUnless , update, update_, upsert , deleteFrom, deleteFrom_ @@ -18,7 +18,6 @@ import Database.Selda.Backend.Internal SeldaBackend(runStmtWithPK, disableForeignKeys, ppConfig, runStmt, runStmtStreaming), QueryRunner, SeldaError(SqlError), - Generator, withBackend ) import Database.Selda.Column ( Row, Col ) import Database.Selda.Compile @@ -60,8 +59,8 @@ query :: (MonadSelda m, Result a) => Query (Backend m) a -> m [Res a] query q = withBackend (flip queryWith q . runStmt) -- | Run a query within a Selda monad like `query` and stream the results. -queryStream :: (MonadSelda m, MonadIO m2, Result a, Monoid a) => (m2 b -> m2 a) -> (SqlValue -> m2 a) -> Query (Backend m) a -> m (m2 a) -queryStream ms f q = withBackend (flip queryStreamWith ms q . runStmtStreaming f) +queryStream :: (MonadSelda m, MonadIO m, Monad m, Foldable m, Result a, Monoid (m (Res a)), Monoid (m [SqlValue])) => ([SqlValue] -> m [SqlValue]) -> Query (Backend m) a -> m (Res a) +queryStream f q = withBackend $ \b -> queryWithStream f (runStmtStreaming b f) q -- | Perform the given query, and insert the result into the given table. -- Returns the number of inserted rows. @@ -295,18 +294,18 @@ queryWith run q = withBackend $ \b -> do return $ mkResults (Proxy :: Proxy a) res -- | Build the final result from streaming result columns. -queryStreamWith :: forall m a. (MonadSelda m, Result a, Monoid a) - => QueryRunner (Int, Generator [SqlValue]) -> Query (Backend m) a -> m (Generator (Res a)) -queryStreamWith run q = withBackend $ \b -> do +queryWithStream :: forall m a. (MonadSelda m, Result a, Monad m, Foldable m, Monoid (m (Res a))) + => ([SqlValue] -> m [SqlValue]) -> QueryRunner (Int, m [SqlValue]) -> Query (Backend m) a -> m (Res a) -- m (m1 (Res a)) -- m1 is generator +queryWithStream f run q = withBackend $ \b -> do res <- fmap snd . liftIO . uncurry run $ compileWith (ppConfig b) q - return $ _ -- mkResultsStreaming ms (Proxy :: Proxy a) res -- FIXME: as res is the generator, we ought not just use ordinary map (which mkResults does) + mkResultsStream f (Proxy :: Proxy a) res -- | Generate the final result of a query from a list of untyped result rows. mkResults :: Result a => Proxy a -> [[SqlValue]] -> [Res a] mkResults p = map (buildResult p) -mkResultsStreaming :: (MonadIO m, Result a) => Proxy a -> m [SqlValue] -> m (Generator (Res a)) -mkResultsStreaming p xs = return $ \ms -> ms (buildResult p) xs +mkResultsStream :: (Monad m, Monoid (m (Res a)), Result a, Foldable m) => ([SqlValue] -> m [SqlValue]) -> Proxy a -> m [SqlValue] -> m (Res a) +mkResultsStream f p = foldl (\acc x -> acc `mappend` (buildResult p <$> f x)) mempty {-# INLINE exec #-} -- | Execute a statement without a result. From 4c8db387d6916ed297b47a1c12f3514f0dece5b5 Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Mon, 14 Jul 2025 21:03:48 +0200 Subject: [PATCH 17/39] wip streaming PostgreSQL --- .../src/Database/Selda/PostgreSQL.hs | 40 ++++++++++++++++++- 1 file changed, 39 insertions(+), 1 deletion(-) diff --git a/selda-postgresql/src/Database/Selda/PostgreSQL.hs b/selda-postgresql/src/Database/Selda/PostgreSQL.hs index cd18efb..70d5b9c 100644 --- a/selda-postgresql/src/Database/Selda/PostgreSQL.hs +++ b/selda-postgresql/src/Database/Selda/PostgreSQL.hs @@ -201,7 +201,7 @@ pgBackend :: Connection -- ^ PostgreSQL connection object. -> SeldaBackend PG pgBackend c = SeldaBackend { runStmt = \q ps -> right <$> pgQueryRunner c False q ps - , runStmtStreaming = \f q ps -> undefined + , runStmtStreaming = \f q ps -> right <$> pgQueryRunnerStream f c False q ps , runStmtWithPK = \q ps -> left <$> pgQueryRunner c True q ps , prepareStmt = pgPrepare c , runPrepared = pgRun c @@ -374,6 +374,20 @@ pgQueryRunner c return_lastid q ps = do getLastId res = (maybe 0 id . fmap readInt64) <$> getvalue res 0 0 +pgQueryRunnerStream :: (Monad m, Monoid (m [SqlValue])) => ([SqlValue] -> m [SqlValue]) -> Connection -> Bool -> T.Text -> [Param] -> IO (Either Int64 (Int, m [SqlValue])) +pgQueryRunnerStream f c return_lastid q ps = do + mres <- execParams c (encodeUtf8 q') [fromSqlValue p | Param p <- ps] Binary + unlessError c errmsg mres $ \res -> do + if return_lastid + then Left <$> getLastId res + else Right <$> streamRows f res + where + errmsg = "error executing query `" ++ T.unpack q' ++ "'" + q' | return_lastid = q <> " RETURNING LASTVAL();" + | otherwise = q + + getLastId res = (maybe 0 id . fmap readInt64) <$> getvalue res 0 0 + pgRun :: Connection -> Dynamic -> [Param] -> IO (Int, [[SqlValue]]) pgRun c hdl ps = do let Just sid = fromDynamic hdl :: Maybe StmtID @@ -385,6 +399,17 @@ pgRun c hdl ps = do Just (_, val, fmt) -> Just (val, fmt) Nothing -> Nothing +pgRunStream :: (Monad m, Monoid (m [SqlValue])) => ([SqlValue] -> m [SqlValue]) -> Connection -> Dynamic -> [Param] -> IO (Int, m [SqlValue]) +pgRunStream f c hdl ps = do + let Just sid = fromDynamic hdl :: Maybe StmtID + mres <- execPrepared c (BS.pack $ show sid) (map mkParam ps) Binary + unlessError c errmsg mres $ streamRows f + where + errmsg = "error executing prepared statement" + mkParam (Param p) = case fromSqlValue p of + Just (_, val, fmt) -> Just (val, fmt) + Nothing -> Nothing + -- | Get all rows from a result. getRows :: Result -> IO (Int, [[SqlValue]]) getRows res = do @@ -400,6 +425,19 @@ getRows res = do where bsToPositiveInt = BS.foldl' (\a x -> a*10+fromIntegral x-48) 0 +streamRows :: (Monad m, Monoid (m [SqlValue])) => ([SqlValue] -> m [SqlValue]) -> Result -> IO (Int, m [SqlValue]) +streamRows f res = do + rows <- ntuples res + cols <- nfields res + types <- ftype res [0..cols-1] + affected <- cmdTuples res + result <- mconcat $ traverse (getRow res types cols) [0..rows-1] + pure $ case affected of + Just "" -> (0, result) + Just s -> (bsToPositiveInt s, result) + _ -> (0, result) + where + bsToPositiveInt = BS.foldl' (\a x -> a*10+fromIntegral x-48) 0 -- | Get all columns for the given row. getRow :: Result -> [Oid] -> Column -> Row -> IO [SqlValue] From defdadc4dda7d7177f852ddab62b55a79d71620c Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Thu, 24 Jul 2025 12:26:47 +0200 Subject: [PATCH 18/39] export queryStream --- selda-sqlite/src/Database/Selda/SQLite.hs | 6 ++++++ selda/src/Database/Selda.hs | 3 ++- 2 files changed, 8 insertions(+), 1 deletion(-) diff --git a/selda-sqlite/src/Database/Selda/SQLite.hs b/selda-sqlite/src/Database/Selda/SQLite.hs index 648e508..07586e3 100644 --- a/selda-sqlite/src/Database/Selda/SQLite.hs +++ b/selda-sqlite/src/Database/Selda/SQLite.hs @@ -3,6 +3,7 @@ module Database.Selda.SQLite ( SQLite , withSQLite + , withSQLiteStreaming , sqliteOpen, seldaClose , sqliteBackend ) where @@ -56,6 +57,11 @@ sqliteBackend _ = error "sqliteBackend called in JS context" #else withSQLite file m = bracket (sqliteOpen file) seldaClose (runSeldaT m) +--withSQLiteStreaming :: +withSQLiteStreaming :: (MonadMask m, MonadIO m, Monoid (m b)) => FilePath -> SeldaT SQLite m b -> m b +withSQLiteStreaming file m = bracket (sqliteOpen file) seldaClose (runSeldaT m) +-- S.yield :: Monad m => a -> S.Stream (S.Of a) m () + -- | Create a Selda backend using an already open database handle. -- This is useful for situations where you want to use some SQLite-specific -- functionality alongside Selda. diff --git a/selda/src/Database/Selda.hs b/selda/src/Database/Selda.hs index c32846f..1a04af8 100644 --- a/selda/src/Database/Selda.hs +++ b/selda/src/Database/Selda.hs @@ -27,7 +27,7 @@ module Database.Selda , SeldaT, SeldaM , Relational, Only (..), The (..) , Table (tableName), Query, Row, Col, Res, Result - , query, queryInto + , query, queryStream, queryInto , transaction, withoutForeignKeyEnforcement , newUuid @@ -156,6 +156,7 @@ import Database.Selda.FieldSelectors import Database.Selda.Frontend ( MonadIO(..), query, + queryStream, queryInto, insert, tryInsert, From 248c4dc7830400613a8d0ea3bde7a09229ff328d Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Thu, 24 Jul 2025 14:45:08 +0200 Subject: [PATCH 19/39] get queryStream right --- selda-sqlite/src/Database/Selda/SQLite.hs | 5 ++--- selda/src/Database/Selda/Frontend.hs | 8 ++++---- 2 files changed, 6 insertions(+), 7 deletions(-) diff --git a/selda-sqlite/src/Database/Selda/SQLite.hs b/selda-sqlite/src/Database/Selda/SQLite.hs index 07586e3..9952464 100644 --- a/selda-sqlite/src/Database/Selda/SQLite.hs +++ b/selda-sqlite/src/Database/Selda/SQLite.hs @@ -57,10 +57,9 @@ sqliteBackend _ = error "sqliteBackend called in JS context" #else withSQLite file m = bracket (sqliteOpen file) seldaClose (runSeldaT m) ---withSQLiteStreaming :: -withSQLiteStreaming :: (MonadMask m, MonadIO m, Monoid (m b)) => FilePath -> SeldaT SQLite m b -> m b -withSQLiteStreaming file m = bracket (sqliteOpen file) seldaClose (runSeldaT m) -- S.yield :: Monad m => a -> S.Stream (S.Of a) m () +withSQLiteStreaming :: (MonadMask m, MonadIO m, Monad m, Monoid (m b)) => FilePath -> SeldaT SQLite m b -> m b +withSQLiteStreaming file m = bracket (sqliteOpen file) seldaClose (runSeldaT m) -- | Create a Selda backend using an already open database handle. -- This is useful for situations where you want to use some SQLite-specific diff --git a/selda/src/Database/Selda/Frontend.hs b/selda/src/Database/Selda/Frontend.hs index ad745a4..10aea5c 100644 --- a/selda/src/Database/Selda/Frontend.hs +++ b/selda/src/Database/Selda/Frontend.hs @@ -59,7 +59,7 @@ query :: (MonadSelda m, Result a) => Query (Backend m) a -> m [Res a] query q = withBackend (flip queryWith q . runStmt) -- | Run a query within a Selda monad like `query` and stream the results. -queryStream :: (MonadSelda m, MonadIO m, Monad m, Foldable m, Result a, Monoid (m (Res a)), Monoid (m [SqlValue])) => ([SqlValue] -> m [SqlValue]) -> Query (Backend m) a -> m (Res a) +queryStream :: (MonadSelda m, MonadIO m2, Monad m2, Foldable m2, Result a, Monoid (m2 (Res a)), Monoid (m2 [SqlValue])) => ([SqlValue] -> m2 [SqlValue]) -> Query (Backend m) a -> m (m2 (Res a)) queryStream f q = withBackend $ \b -> queryWithStream f (runStmtStreaming b f) q -- | Perform the given query, and insert the result into the given table. @@ -294,11 +294,11 @@ queryWith run q = withBackend $ \b -> do return $ mkResults (Proxy :: Proxy a) res -- | Build the final result from streaming result columns. -queryWithStream :: forall m a. (MonadSelda m, Result a, Monad m, Foldable m, Monoid (m (Res a))) - => ([SqlValue] -> m [SqlValue]) -> QueryRunner (Int, m [SqlValue]) -> Query (Backend m) a -> m (Res a) -- m (m1 (Res a)) -- m1 is generator +queryWithStream :: forall m a m1. (MonadSelda m, Result a, Monad m1, Foldable m1, Monoid (m1 (Res a))) + => ([SqlValue] -> m1 [SqlValue]) -> QueryRunner (Int, m1 [SqlValue]) -> Query (Backend m) a -> m (m1 (Res a)) -- m (m1 (Res a)) -- m1 is generator queryWithStream f run q = withBackend $ \b -> do res <- fmap snd . liftIO . uncurry run $ compileWith (ppConfig b) q - mkResultsStream f (Proxy :: Proxy a) res + return $ mkResultsStream f (Proxy :: Proxy a) res -- | Generate the final result of a query from a list of untyped result rows. mkResults :: Result a => Proxy a -> [[SqlValue]] -> [Res a] From e7ff2d2d28a3266baba428346ee1b6b4d09a1740 Mon Sep 17 00:00:00 2001 From: Mirek Kratochvil Date: Thu, 24 Jul 2025 15:57:07 +0200 Subject: [PATCH 20/39] try to simplify the streamy interface --- selda/src/Database/Selda/Backend/Internal.hs | 2 +- selda/src/Database/Selda/Frontend.hs | 17 ++++------------- 2 files changed, 5 insertions(+), 14 deletions(-) diff --git a/selda/src/Database/Selda/Backend/Internal.hs b/selda/src/Database/Selda/Backend/Internal.hs index 8798d43..623458c 100644 --- a/selda/src/Database/Selda/Backend/Internal.hs +++ b/selda/src/Database/Selda/Backend/Internal.hs @@ -218,7 +218,7 @@ data SeldaBackend b { -- | Execute an SQL statement. runStmt :: Text -> [Param] -> IO (Int, [[SqlValue]]) - , runStmtStreaming :: forall m. (Monad m, MonadIO m, Monoid (m [SqlValue])) => ([SqlValue] -> m [SqlValue]) -> Text -> [Param] -> IO (Int, m [SqlValue]) + , runStmtStreaming :: forall m r . (MonadIO m, Monoid r) => Text -> [Param] -> ([[SqlValue]] -> m r) -> m r -- | Execute an SQL statement and return the last inserted primary key, -- where the primary key is auto-incrementing. diff --git a/selda/src/Database/Selda/Frontend.hs b/selda/src/Database/Selda/Frontend.hs index ad745a4..aafe566 100644 --- a/selda/src/Database/Selda/Frontend.hs +++ b/selda/src/Database/Selda/Frontend.hs @@ -2,7 +2,7 @@ -- | API for running Selda operations over databases. module Database.Selda.Frontend ( Result, Res, MonadIO (..), MonadSelda (..), SeldaT, OnError (..) - , query, queryInto, queryStream + , query, queryInto, forQuery , insert, insert_, insertWithPK, tryInsert, insertWhen, insertUnless , update, update_, upsert , deleteFrom, deleteFrom_ @@ -59,8 +59,9 @@ query :: (MonadSelda m, Result a) => Query (Backend m) a -> m [Res a] query q = withBackend (flip queryWith q . runStmt) -- | Run a query within a Selda monad like `query` and stream the results. -queryStream :: (MonadSelda m, MonadIO m, Monad m, Foldable m, Result a, Monoid (m (Res a)), Monoid (m [SqlValue])) => ([SqlValue] -> m [SqlValue]) -> Query (Backend m) a -> m (Res a) -queryStream f q = withBackend $ \b -> queryWithStream f (runStmtStreaming b f) q +forQuery :: forall m a r . (MonadSelda m, Result a, Monoid r) => Query (Backend m) a -> (Res a -> m r) -> m r +forQuery q k = withBackend $ \b -> + uncurry (runStmtStreaming b) (ppConfig b `compileWith` q) $ fmap mconcat . traverse k . mkResults (Proxy :: Proxy a) -- | Perform the given query, and insert the result into the given table. -- Returns the number of inserted rows. @@ -293,20 +294,10 @@ queryWith run q = withBackend $ \b -> do res <- fmap snd . liftIO . uncurry run $ compileWith (ppConfig b) q return $ mkResults (Proxy :: Proxy a) res --- | Build the final result from streaming result columns. -queryWithStream :: forall m a. (MonadSelda m, Result a, Monad m, Foldable m, Monoid (m (Res a))) - => ([SqlValue] -> m [SqlValue]) -> QueryRunner (Int, m [SqlValue]) -> Query (Backend m) a -> m (Res a) -- m (m1 (Res a)) -- m1 is generator -queryWithStream f run q = withBackend $ \b -> do - res <- fmap snd . liftIO . uncurry run $ compileWith (ppConfig b) q - mkResultsStream f (Proxy :: Proxy a) res - -- | Generate the final result of a query from a list of untyped result rows. mkResults :: Result a => Proxy a -> [[SqlValue]] -> [Res a] mkResults p = map (buildResult p) -mkResultsStream :: (Monad m, Monoid (m (Res a)), Result a, Foldable m) => ([SqlValue] -> m [SqlValue]) -> Proxy a -> m [SqlValue] -> m (Res a) -mkResultsStream f p = foldl (\acc x -> acc `mappend` (buildResult p <$> f x)) mempty - {-# INLINE exec #-} -- | Execute a statement without a result. exec :: MonadSelda m => Text -> [Param] -> m Int From 16322d80391da4ecf28f9d95c01da68bd7ffa633 Mon Sep 17 00:00:00 2001 From: Mirek Kratochvil Date: Thu, 24 Jul 2025 16:16:49 +0200 Subject: [PATCH 21/39] sqlite backend seems to compile (wat) --- selda-sqlite/src/Database/Selda/SQLite.hs | 23 +++++++++----------- selda/src/Database/Selda/Backend/Internal.hs | 2 +- selda/src/Database/Selda/Frontend.hs | 2 +- 3 files changed, 12 insertions(+), 15 deletions(-) diff --git a/selda-sqlite/src/Database/Selda/SQLite.hs b/selda-sqlite/src/Database/Selda/SQLite.hs index 648e508..eeb322b 100644 --- a/selda-sqlite/src/Database/Selda/SQLite.hs +++ b/selda-sqlite/src/Database/Selda/SQLite.hs @@ -67,7 +67,7 @@ withSQLite file m = bracket (sqliteOpen file) seldaClose (runSeldaT m) sqliteBackend :: Database -> SeldaBackend SQLite sqliteBackend db = SeldaBackend { runStmt = \q ps -> snd <$> sqliteQueryRunner db q ps - , runStmtStreaming = \f q ps -> snd <$> sqliteQueryRunnerStreaming f db q ps + , runStmtStreaming = sqliteQueryRunnerStreaming db , runStmtWithPK = \q ps -> fst <$> sqliteQueryRunner db q ps , prepareStmt = \_ _ -> sqlitePrepare db , runPrepared = sqliteRunPrepared db @@ -199,20 +199,17 @@ sqliteRunStmt db stm params = do cs <- changes db return (fromIntegral rid, (cs, [map fromSqlData r | r <- rows])) -sqliteRunStmtStreaming :: (Monad m, Monoid (m [SqlValue])) => ([SqlValue] -> m [SqlValue]) -> Database -> Statement -> [Param] -> IO (Int64, (Int, m [SqlValue])) -sqliteRunStmtStreaming f db stm params = do - bind stm [toSqlData p | Param p <- params] - rows <- streamRows f stm mempty - rid <- lastInsertRowId db - cs <- changes db - return (fromIntegral rid, (cs, rows)) +sqliteRunStmtStreaming :: (MonadIO m, MonadMask m, MonadIO m, Monoid r) => Database -> Statement -> [Param] -> ([[SqlValue]] -> m r) -> m r +sqliteRunStmtStreaming db stm params k = do + liftIO $ bind stm [toSqlData p | Param p <- params] + liftIO (getRows stm []) >>= k . map (map fromSqlData) -sqliteQueryRunnerStreaming :: (Monad m, Monoid (m [SqlValue])) => ([SqlValue] -> m [SqlValue]) -> Database -> QueryRunner (Int64, (Int, m [SqlValue])) -sqliteQueryRunnerStreaming f db qry params = do +sqliteQueryRunnerStreaming :: (MonadIO m, MonadMask m, MonadIO m, Monoid r) => Database -> Text -> [Param] -> ([[SqlValue]] -> m r) -> m r +sqliteQueryRunnerStreaming db qry params k = do eres <- try $ do - stm <- prepare db qry - sqliteRunStmtStreaming f db stm params `finally` do - finalize stm + stm <- liftIO $ prepare db qry + sqliteRunStmtStreaming db stm params k `finally` do + liftIO $ finalize stm case eres of Left e@(SQLError{}) -> throwM (SqlError (show e)) Right res -> return res diff --git a/selda/src/Database/Selda/Backend/Internal.hs b/selda/src/Database/Selda/Backend/Internal.hs index 623458c..3dabbc1 100644 --- a/selda/src/Database/Selda/Backend/Internal.hs +++ b/selda/src/Database/Selda/Backend/Internal.hs @@ -218,7 +218,7 @@ data SeldaBackend b { -- | Execute an SQL statement. runStmt :: Text -> [Param] -> IO (Int, [[SqlValue]]) - , runStmtStreaming :: forall m r . (MonadIO m, Monoid r) => Text -> [Param] -> ([[SqlValue]] -> m r) -> m r + , runStmtStreaming :: forall m r . (MonadIO m, MonadMask m, Monoid r) => Text -> [Param] -> ([[SqlValue]] -> m r) -> m r -- | Execute an SQL statement and return the last inserted primary key, -- where the primary key is auto-incrementing. diff --git a/selda/src/Database/Selda/Frontend.hs b/selda/src/Database/Selda/Frontend.hs index aafe566..04a330a 100644 --- a/selda/src/Database/Selda/Frontend.hs +++ b/selda/src/Database/Selda/Frontend.hs @@ -59,7 +59,7 @@ query :: (MonadSelda m, Result a) => Query (Backend m) a -> m [Res a] query q = withBackend (flip queryWith q . runStmt) -- | Run a query within a Selda monad like `query` and stream the results. -forQuery :: forall m a r . (MonadSelda m, Result a, Monoid r) => Query (Backend m) a -> (Res a -> m r) -> m r +forQuery :: forall m a r . (MonadSelda m, MonadMask m, Result a, Monoid r) => Query (Backend m) a -> (Res a -> m r) -> m r forQuery q k = withBackend $ \b -> uncurry (runStmtStreaming b) (ppConfig b `compileWith` q) $ fmap mconcat . traverse k . mkResults (Proxy :: Proxy a) From 2bc9792a5b30b118bcb98d594e54b0f93219ee24 Mon Sep 17 00:00:00 2001 From: Mirek Kratochvil Date: Thu, 24 Jul 2025 18:24:18 +0200 Subject: [PATCH 22/39] this should stream properly now --- selda-sqlite/src/Database/Selda/SQLite.hs | 65 +++++++++++++++-------- 1 file changed, 42 insertions(+), 23 deletions(-) diff --git a/selda-sqlite/src/Database/Selda/SQLite.hs b/selda-sqlite/src/Database/Selda/SQLite.hs index eeb322b..f4a157b 100644 --- a/selda-sqlite/src/Database/Selda/SQLite.hs +++ b/selda-sqlite/src/Database/Selda/SQLite.hs @@ -199,21 +199,6 @@ sqliteRunStmt db stm params = do cs <- changes db return (fromIntegral rid, (cs, [map fromSqlData r | r <- rows])) -sqliteRunStmtStreaming :: (MonadIO m, MonadMask m, MonadIO m, Monoid r) => Database -> Statement -> [Param] -> ([[SqlValue]] -> m r) -> m r -sqliteRunStmtStreaming db stm params k = do - liftIO $ bind stm [toSqlData p | Param p <- params] - liftIO (getRows stm []) >>= k . map (map fromSqlData) - -sqliteQueryRunnerStreaming :: (MonadIO m, MonadMask m, MonadIO m, Monoid r) => Database -> Text -> [Param] -> ([[SqlValue]] -> m r) -> m r -sqliteQueryRunnerStreaming db qry params k = do - eres <- try $ do - stm <- liftIO $ prepare db qry - sqliteRunStmtStreaming db stm params k `finally` do - liftIO $ finalize stm - case eres of - Left e@(SQLError{}) -> throwM (SqlError (show e)) - Right res -> return res - getRows :: Statement -> [[SQLData]] -> IO [[SQLData]] getRows s acc = do res <- step s @@ -223,15 +208,49 @@ getRows s acc = do getRows s (cs : acc) _ -> do return $ reverse acc +sqliteQueryRunnerStreaming :: + (MonadIO m, MonadMask m, MonadIO m, Monoid r) + => Database + -> Text + -> [Param] + -> ([[SqlValue]] -> m r) + -> m r +sqliteQueryRunnerStreaming db qry params k = do + eres <- + try $ do + stm <- liftIO $ prepare db qry + sqliteRunStmtStreaming db stm params k `finally` do + liftIO $ finalize stm + case eres of + Left e@(SQLError {}) -> throwM (SqlError (show e)) + Right res -> return res -streamRows :: (Monad m, Monoid (m [SqlValue])) => ([SqlValue]-> m [SqlValue]) -> Statement -> m [SqlValue]-> IO (m [SqlValue]) -streamRows f s acc = do - res <- step s - case res of - Row -> do - cs <- map fromSqlData <$> liftIO (columns s :: IO [SQLData]) - streamRows f s (acc `mappend` f cs) - _ -> return acc +sqliteRunStmtStreaming :: + (MonadIO m, MonadMask m, MonadIO m, Monoid r) + => Database + -> Statement + -> [Param] + -> ([[SqlValue]] -> m r) + -> m r +sqliteRunStmtStreaming db stm params k = do + liftIO $ bind stm [toSqlData p | Param p <- params] + streamRows stm 1024 $ k . map (map fromSqlData) + +streamRows :: + (MonadIO m, Monoid r) => Statement -> Int -> ([[SQLData]] -> m r) -> m r +streamRows s n k = cont mempty + where + cont r = go r [] 0 + go r acc i + | i < n = do + res <- liftIO $ step s + case res of + Row -> do + cs <- liftIO $ columns s + go r (cs : acc) (succ i) + _ -> mappend r <$> send acc >>= cont + | otherwise = mappend r <$> send acc >>= cont + send = k . reverse toSqlData :: Lit a -> SQLData toSqlData (LInt32 i) = SQLInteger $ fromIntegral i From 19cacf79faa16735c9dfd8892665333801d5df02 Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Thu, 24 Jul 2025 19:22:21 +0200 Subject: [PATCH 23/39] export forQuery! --- selda/src/Database/Selda.hs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/selda/src/Database/Selda.hs b/selda/src/Database/Selda.hs index c32846f..f70b5de 100644 --- a/selda/src/Database/Selda.hs +++ b/selda/src/Database/Selda.hs @@ -27,7 +27,7 @@ module Database.Selda , SeldaT, SeldaM , Relational, Only (..), The (..) , Table (tableName), Query, Row, Col, Res, Result - , query, queryInto + , query, queryInto, forQuery , transaction, withoutForeignKeyEnforcement , newUuid @@ -157,6 +157,7 @@ import Database.Selda.Frontend ( MonadIO(..), query, queryInto, + forQuery, insert, tryInsert, upsert, From d6b760799b593a705f0b7a63f652a7314979603f Mon Sep 17 00:00:00 2001 From: Mirek Kratochvil Date: Thu, 24 Jul 2025 20:59:57 +0200 Subject: [PATCH 24/39] fix: do not cycle forever, it's indeed what we desire least --- selda-sqlite/src/Database/Selda/SQLite.hs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/selda-sqlite/src/Database/Selda/SQLite.hs b/selda-sqlite/src/Database/Selda/SQLite.hs index f4a157b..6d7354f 100644 --- a/selda-sqlite/src/Database/Selda/SQLite.hs +++ b/selda-sqlite/src/Database/Selda/SQLite.hs @@ -248,7 +248,7 @@ streamRows s n k = cont mempty Row -> do cs <- liftIO $ columns s go r (cs : acc) (succ i) - _ -> mappend r <$> send acc >>= cont + _ -> mappend r <$> send acc | otherwise = mappend r <$> send acc >>= cont send = k . reverse From d62b8391bf7422a238db0c220d23f725719adc4b Mon Sep 17 00:00:00 2001 From: Mirek Kratochvil Date: Thu, 24 Jul 2025 21:00:15 +0200 Subject: [PATCH 25/39] export forQuery --- selda/src/Database/Selda.hs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/selda/src/Database/Selda.hs b/selda/src/Database/Selda.hs index c32846f..f70b5de 100644 --- a/selda/src/Database/Selda.hs +++ b/selda/src/Database/Selda.hs @@ -27,7 +27,7 @@ module Database.Selda , SeldaT, SeldaM , Relational, Only (..), The (..) , Table (tableName), Query, Row, Col, Res, Result - , query, queryInto + , query, queryInto, forQuery , transaction, withoutForeignKeyEnforcement , newUuid @@ -157,6 +157,7 @@ import Database.Selda.Frontend ( MonadIO(..), query, queryInto, + forQuery, insert, tryInsert, upsert, From 1feaced9d616d72d323bdcf7459f7ca081e1f88f Mon Sep 17 00:00:00 2001 From: Mirek Kratochvil Date: Thu, 24 Jul 2025 21:37:11 +0200 Subject: [PATCH 26/39] make the postgres backend compile again --- .../src/Database/Selda/PostgreSQL.hs | 97 +++++++++++-------- 1 file changed, 57 insertions(+), 40 deletions(-) diff --git a/selda-postgresql/src/Database/Selda/PostgreSQL.hs b/selda-postgresql/src/Database/Selda/PostgreSQL.hs index 70d5b9c..61e2cb7 100644 --- a/selda-postgresql/src/Database/Selda/PostgreSQL.hs +++ b/selda-postgresql/src/Database/Selda/PostgreSQL.hs @@ -19,7 +19,7 @@ import Control.Monad.Catch import Control.Monad.IO.Class #ifndef __HASTE__ -import Control.Monad (void) +import Control.Monad (void, unless) import qualified Data.ByteString as BS (foldl') import qualified Data.ByteString.Char8 as BS (pack, unpack) import Data.Dynamic @@ -201,7 +201,7 @@ pgBackend :: Connection -- ^ PostgreSQL connection object. -> SeldaBackend PG pgBackend c = SeldaBackend { runStmt = \q ps -> right <$> pgQueryRunner c False q ps - , runStmtStreaming = \f q ps -> right <$> pgQueryRunnerStream f c False q ps + , runStmtStreaming = \q ps k -> pgQueryRunnerStream c q ps k , runStmtWithPK = \q ps -> left <$> pgQueryRunner c True q ps , prepareStmt = pgPrepare c , runPrepared = pgRun c @@ -374,20 +374,6 @@ pgQueryRunner c return_lastid q ps = do getLastId res = (maybe 0 id . fmap readInt64) <$> getvalue res 0 0 -pgQueryRunnerStream :: (Monad m, Monoid (m [SqlValue])) => ([SqlValue] -> m [SqlValue]) -> Connection -> Bool -> T.Text -> [Param] -> IO (Either Int64 (Int, m [SqlValue])) -pgQueryRunnerStream f c return_lastid q ps = do - mres <- execParams c (encodeUtf8 q') [fromSqlValue p | Param p <- ps] Binary - unlessError c errmsg mres $ \res -> do - if return_lastid - then Left <$> getLastId res - else Right <$> streamRows f res - where - errmsg = "error executing query `" ++ T.unpack q' ++ "'" - q' | return_lastid = q <> " RETURNING LASTVAL();" - | otherwise = q - - getLastId res = (maybe 0 id . fmap readInt64) <$> getvalue res 0 0 - pgRun :: Connection -> Dynamic -> [Param] -> IO (Int, [[SqlValue]]) pgRun c hdl ps = do let Just sid = fromDynamic hdl :: Maybe StmtID @@ -399,17 +385,6 @@ pgRun c hdl ps = do Just (_, val, fmt) -> Just (val, fmt) Nothing -> Nothing -pgRunStream :: (Monad m, Monoid (m [SqlValue])) => ([SqlValue] -> m [SqlValue]) -> Connection -> Dynamic -> [Param] -> IO (Int, m [SqlValue]) -pgRunStream f c hdl ps = do - let Just sid = fromDynamic hdl :: Maybe StmtID - mres <- execPrepared c (BS.pack $ show sid) (map mkParam ps) Binary - unlessError c errmsg mres $ streamRows f - where - errmsg = "error executing prepared statement" - mkParam (Param p) = case fromSqlValue p of - Just (_, val, fmt) -> Just (val, fmt) - Nothing -> Nothing - -- | Get all rows from a result. getRows :: Result -> IO (Int, [[SqlValue]]) getRows res = do @@ -425,20 +400,62 @@ getRows res = do where bsToPositiveInt = BS.foldl' (\a x -> a*10+fromIntegral x-48) 0 -streamRows :: (Monad m, Monoid (m [SqlValue])) => ([SqlValue] -> m [SqlValue]) -> Result -> IO (Int, m [SqlValue]) -streamRows f res = do - rows <- ntuples res - cols <- nfields res - types <- ftype res [0..cols-1] - affected <- cmdTuples res - result <- mconcat $ traverse (getRow res types cols) [0..rows-1] - pure $ case affected of - Just "" -> (0, result) - Just s -> (bsToPositiveInt s, result) - _ -> (0, result) +pgQueryRunnerStream :: + (MonadIO m, MonadMask m, Monoid r) + => Connection + -> T.Text + -> [Param] + -> ([[SqlValue]] -> m r) + -> m r +pgQueryRunnerStream c q ps k = do + qsent <- + liftIO + $ sendQueryParams c (encodeUtf8 q) [fromSqlValue p | Param p <- ps] Binary + unless qsent . throwM $ DbError "sendQueryParams failed" + msent <- liftIO $ setSingleRowMode c -- TODO: chunked mode, but the library doesn't wrap setChunkedRowsMode + unless msent . throwM $ DbError "setSingleRowMode failed" + streamRows c k mempty + +-- TODO connect this to the frontend +pgRunStream :: + (MonadIO m, MonadMask m, Monoid r) + => Connection + -> Dynamic + -> [Param] + -> ([[SqlValue]] -> m r) + -> m r +pgRunStream c hdl ps k = do + let Just sid = fromDynamic hdl :: Maybe StmtID + qsent <- + liftIO $ sendQueryPrepared c (BS.pack $ show sid) (map mkParam ps) Binary + unless qsent . throwM $ DbError "sendQueryParams failed" + msent <- liftIO $ setSingleRowMode c + unless msent . throwM $ DbError "setSingleRowMode failed" + streamRows c k mempty where - bsToPositiveInt = BS.foldl' (\a x -> a*10+fromIntegral x-48) 0 - + mkParam (Param p) = + case fromSqlValue p of + Just (_, val, fmt) -> Just (val, fmt) + Nothing -> Nothing + +streamRows :: + (MonadIO m, MonadMask m, Monoid r) + => Connection + -> ([[SqlValue]] -> m r) + -> r + -> m r +streamRows c k r = do + mres <- liftIO $ getResult c + case mres of + Just res -> do + result <- + liftIO $ do + rows <- ntuples res + cols <- nfields res + types <- mapM (ftype res) [0 .. cols - 1] + mapM (getRow res types cols) [0 .. rows - 1] + k result >>= streamRows c k . mappend r + Nothing -> pure r -- | Get all columns for the given row. getRow :: Result -> [Oid] -> Column -> Row -> IO [SqlValue] getRow res types cols row = do From 1b177f6c246f38ebaca8237df6f08eb3ee7d6532 Mon Sep 17 00:00:00 2001 From: Mirek Kratochvil Date: Thu, 24 Jul 2025 21:38:40 +0200 Subject: [PATCH 27/39] mind the failure --- selda-postgresql/src/Database/Selda/PostgreSQL.hs | 3 +++ 1 file changed, 3 insertions(+) diff --git a/selda-postgresql/src/Database/Selda/PostgreSQL.hs b/selda-postgresql/src/Database/Selda/PostgreSQL.hs index 61e2cb7..9cceaf7 100644 --- a/selda-postgresql/src/Database/Selda/PostgreSQL.hs +++ b/selda-postgresql/src/Database/Selda/PostgreSQL.hs @@ -446,6 +446,8 @@ streamRows :: -> m r streamRows c k r = do mres <- liftIO $ getResult c + --TODO: how does one actually catch errors here? + --(see caution note at https://www.postgresql.org/docs/17/libpq-single-row-mode.html ) case mres of Just res -> do result <- @@ -456,6 +458,7 @@ streamRows c k r = do mapM (getRow res types cols) [0 .. rows - 1] k result >>= streamRows c k . mappend r Nothing -> pure r + -- | Get all columns for the given row. getRow :: Result -> [Oid] -> Column -> Row -> IO [SqlValue] getRow res types cols row = do From 6726c253dc9dfa217f444e525cbcc2d4175836a6 Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Thu, 31 Jul 2025 18:10:18 +0200 Subject: [PATCH 28/39] fix test dependencies --- selda-tests/selda-tests.cabal | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/selda-tests/selda-tests.cabal b/selda-tests/selda-tests.cabal index 542b914..9a01077 100644 --- a/selda-tests/selda-tests.cabal +++ b/selda-tests/selda-tests.cabal @@ -33,15 +33,15 @@ test-suite selda-testsuite Tests.Validation build-depends: aeson - , base >=4.8 && <5 + , base >=4.10 && <5 , bytestring >=0.10 && <0.13 , directory >=1.2 && <1.4 , exceptions >=0.8 && <0.11 , HUnit >=1.4 && <1.7 , selda , selda-json - , text >=1.1 && <2.1 - , time >=1.4 && <1.13 + , text >=1.1 && <2.2 + , time >=1.4 && <1.15 , random >=1.1 && <1.3 , uuid-types >=1.0 && <1.1 if flag(postgres) From 2a9fd3a200a967315708150e9343f433a9eb7cf4 Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Thu, 31 Jul 2025 18:11:15 +0200 Subject: [PATCH 29/39] document test command --- README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/README.md b/README.md index b0c2ea9..7fa1b94 100644 --- a/README.md +++ b/README.md @@ -80,7 +80,7 @@ If you want to contribute code, please consult the following checklist before sending a pull request: * Does the code build with a recent version of GHC? -* Do all the tests pass? +* Do all the tests pass? Run `make test` in the root of the repo. * Have you added any tests covering your code? If you want to contribute code but don't really know where to begin, From 228d04de6da1a82adc74c62cac347bc466daac54 Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Thu, 31 Jul 2025 18:20:35 +0200 Subject: [PATCH 30/39] add test for streaming --- selda-tests/test/Tests/Query.hs | 8 ++++++++ 1 file changed, 8 insertions(+) diff --git a/selda-tests/test/Tests/Query.hs b/selda-tests/test/Tests/Query.hs index 1cffaa3..ac3eb4f 100644 --- a/selda-tests/test/Tests/Query.hs +++ b/selda-tests/test/Tests/Query.hs @@ -67,6 +67,7 @@ queryTests run = test , "nonNull" ~: run nonNullYieldsEmptyResult , "rawQuery1" ~: run rawQuery1Works , "rawQuery" ~: run rawQueryWorks + , "rawQueryStreaming" ~: run rawQueryStreamingWorks , "union" ~: run unionWorks , "union discards dupes" ~: run unionDiscardsDupes , "union works for whole rows" ~: run unionWorksForWholeRows @@ -620,6 +621,13 @@ rawQueryWorks = do let correct = [p | p <- peopleItems, name p == "Link"] assEq "wrong name list returned" correct ppl +rawQueryStreamingWorks = do + let q = rawQuery ["name", "age", "pet", "cash"] + ("SELECT * FROM people WHERE name = " <> injLit ("Link"::Text)) + ppl <- forQuery q $ pure . (:[]) + let correct = [p | p <- peopleItems, name p == "Link"] + assEq "wrong name list returned" correct ppl + unionWorks = assQueryEq "wrong name list returned" correct $ do let ppl = pName `from` select people pets = (pPet `from` select people) >>= nonNull From ee03232ca453feccfb63bec014c4b2c9a32d5075 Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Thu, 31 Jul 2025 18:22:04 +0200 Subject: [PATCH 31/39] add ghc 9.12.2 to github CI --- .github/workflows/run-tests.yml | 2 +- selda/selda.cabal | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/.github/workflows/run-tests.yml b/.github/workflows/run-tests.yml index a6c336e..c3bc86f 100644 --- a/.github/workflows/run-tests.yml +++ b/.github/workflows/run-tests.yml @@ -10,7 +10,7 @@ jobs: build: strategy: matrix: - ghc: ['8.8.2', '8.10.1', '9.2.5', '9.4.4'] + ghc: ['8.8.2', '8.10.1', '9.2.5', '9.4.4', '9.12.2'] runs-on: ubuntu-latest steps: - uses: actions/checkout@v2 diff --git a/selda/selda.cabal b/selda/selda.cabal index d25122c..76a39fc 100644 --- a/selda/selda.cabal +++ b/selda/selda.cabal @@ -18,7 +18,7 @@ maintainer: anton@ekblad.cc category: Database build-type: Simple cabal-version: >=1.10 -tested-with: GHC == 8.8.2, GHC == 8.10.1, GHC == 9.2.5, GHC == 9.4.4 +tested-with: GHC == 8.8.2, GHC == 8.10.1, GHC == 9.2.5, GHC == 9.4.4, GHC == 9.12.2 source-repository head type: git From 14e55660a601beaa6c9d8edbd71a03ab4037b80d Mon Sep 17 00:00:00 2001 From: Mirek Kratochvil Date: Fri, 1 Aug 2025 09:43:37 +0200 Subject: [PATCH 32/39] clean up errors, update TODOs --- selda-postgresql/src/Database/Selda/PostgreSQL.hs | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/selda-postgresql/src/Database/Selda/PostgreSQL.hs b/selda-postgresql/src/Database/Selda/PostgreSQL.hs index 9cceaf7..ad8dd35 100644 --- a/selda-postgresql/src/Database/Selda/PostgreSQL.hs +++ b/selda-postgresql/src/Database/Selda/PostgreSQL.hs @@ -412,11 +412,11 @@ pgQueryRunnerStream c q ps k = do liftIO $ sendQueryParams c (encodeUtf8 q) [fromSqlValue p | Param p <- ps] Binary unless qsent . throwM $ DbError "sendQueryParams failed" - msent <- liftIO $ setSingleRowMode c -- TODO: chunked mode, but the library doesn't wrap setChunkedRowsMode + msent <- liftIO $ setSingleRowMode c -- TODO: chunked mode, but the library doesn't wrap setChunkedRowsMode. See https://github.com/haskellari/postgresql-libpq/pull/79 for progress. unless msent . throwM $ DbError "setSingleRowMode failed" streamRows c k mempty --- TODO connect this to the frontend +-- TODO this is prepared for streaming of prepared statements (not in frontend yet) pgRunStream :: (MonadIO m, MonadMask m, Monoid r) => Connection @@ -446,8 +446,6 @@ streamRows :: -> m r streamRows c k r = do mres <- liftIO $ getResult c - --TODO: how does one actually catch errors here? - --(see caution note at https://www.postgresql.org/docs/17/libpq-single-row-mode.html ) case mres of Just res -> do result <- @@ -456,8 +454,10 @@ streamRows c k r = do cols <- nfields res types <- mapM (ftype res) [0 .. cols - 1] mapM (getRow res types cols) [0 .. rows - 1] - k result >>= streamRows c k . mappend r - Nothing -> pure r + if null result + then pure r + else k result >>= streamRows c k . mappend r + Nothing -> throwM $ DbError "streaming getResult failed" -- | Get all columns for the given row. getRow :: Result -> [Oid] -> Column -> Row -> IO [SqlValue] From 32f84173f65c5a9ffad20fc3af1c2b3f227fa7c8 Mon Sep 17 00:00:00 2001 From: Benjamin Weber Date: Mon, 25 Aug 2025 17:26:11 +0200 Subject: [PATCH 33/39] fix missing args in runStmtStreaming --- selda-sqlite/src/Database/Selda/SQLite.hs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/selda-sqlite/src/Database/Selda/SQLite.hs b/selda-sqlite/src/Database/Selda/SQLite.hs index ad0eeb3..5db79a4 100644 --- a/selda-sqlite/src/Database/Selda/SQLite.hs +++ b/selda-sqlite/src/Database/Selda/SQLite.hs @@ -72,7 +72,7 @@ withSQLiteStreaming file m = bracket (sqliteOpen file) seldaClose (runSeldaT m) sqliteBackend :: Database -> SeldaBackend SQLite sqliteBackend db = SeldaBackend { runStmt = \q ps -> snd <$> sqliteQueryRunner db q ps - , runStmtStreaming = sqliteQueryRunnerStreaming db + , runStmtStreaming = \q ps -> sqliteQueryRunnerStreaming db q ps , runStmtWithPK = \q ps -> fst <$> sqliteQueryRunner db q ps , prepareStmt = \_ _ -> sqlitePrepare db , runPrepared = sqliteRunPrepared db From 02c429735ca8461b0ac298f5e18d308037072a85 Mon Sep 17 00:00:00 2001 From: Mirek Kratochvil Date: Fri, 24 Oct 2025 11:19:07 +0200 Subject: [PATCH 34/39] remove trailing blanks here and there --- selda/src/Database/Selda/Backend/Internal.hs | 2 +- website/compile.hs | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/selda/src/Database/Selda/Backend/Internal.hs b/selda/src/Database/Selda/Backend/Internal.hs index ec76dc3..04154f2 100644 --- a/selda/src/Database/Selda/Backend/Internal.hs +++ b/selda/src/Database/Selda/Backend/Internal.hs @@ -291,7 +291,7 @@ newtype SeldaT b m a = S {unS :: ReaderT (SeldaConnection b) m a} instance (MonadIO m, MonadMask m) => MonadSelda (SeldaT b m) where type Backend (SeldaT b m) = b withConnection m = S ask >>= m - + instance MonadTrans (SeldaT b) where lift = S . lift diff --git a/website/compile.hs b/website/compile.hs index 22efbab..e504727 100644 --- a/website/compile.hs +++ b/website/compile.hs @@ -148,7 +148,7 @@ addToSiteMap Page{..} = do | pageSubdirectory == "." = "" | otherwise = pageSubdirectory pageUrl = base subdir pageFileName - + sitemapFile :: FilePath sitemapFile = siteDirectory "sitemap.txt" From 677d20a0e3df038c5a6c6386cdc5b4676f3f7cc8 Mon Sep 17 00:00:00 2001 From: Mirek Kratochvil Date: Fri, 24 Oct 2025 11:24:32 +0200 Subject: [PATCH 35/39] document runStmtStreaming, add prep-stmt analog --- selda/src/Database/Selda/Backend/Internal.hs | 13 ++++++++++++- 1 file changed, 12 insertions(+), 1 deletion(-) diff --git a/selda/src/Database/Selda/Backend/Internal.hs b/selda/src/Database/Selda/Backend/Internal.hs index 04154f2..28cc15b 100644 --- a/selda/src/Database/Selda/Backend/Internal.hs +++ b/selda/src/Database/Selda/Backend/Internal.hs @@ -219,7 +219,12 @@ data SeldaBackend b { -- | Execute an SQL statement. runStmt :: Text -> [Param] -> IO (Int, [[SqlValue]]) - , runStmtStreaming :: forall m r . (MonadIO m, MonadMask m, Monoid r) => Text -> [Param] -> ([[SqlValue]] -> m r) -> m r + -- | Like `runStmt` but instead of collecting all results at once in a + -- list, collects chunks of results and passes them one by one to a given + -- "callback" action. The callbacks may collect and return any `Monoid`. + , runStmtStreaming + :: forall m r . (MonadIO m, MonadMask m, Monoid r) + => Text -> [Param] -> ([[SqlValue]] -> m r) -> m r -- | Execute an SQL statement and return the last inserted primary key, -- where the primary key is auto-incrementing. @@ -232,6 +237,12 @@ data SeldaBackend b -- | Execute a prepared statement. , runPrepared :: Dynamic -> [Param] -> IO (Int, [[SqlValue]]) + -- | Like `runPrepared` but with streaming support, analogously to + -- `runStmt` and `runStmtStreaming`. + , runPreparedStreaming + :: forall m r . (MonadIO m, MonadMask m, Monoid r) + => Dynamic -> [Param] -> ([[SqlValue]] -> m r) -> m r + -- | Get a list of all columns in the given table, with the type and any -- modifiers for each column. -- Return an empty list if the given table does not exist. From 86069aa891926d7eef8206fe0ffb786f1a13c5fb Mon Sep 17 00:00:00 2001 From: Mirek Kratochvil Date: Fri, 24 Oct 2025 11:31:17 +0200 Subject: [PATCH 36/39] psql prep-stmt streaming: *clicks in place* --- selda-postgresql/src/Database/Selda/PostgreSQL.hs | 2 ++ 1 file changed, 2 insertions(+) diff --git a/selda-postgresql/src/Database/Selda/PostgreSQL.hs b/selda-postgresql/src/Database/Selda/PostgreSQL.hs index 0715bfb..119d46a 100644 --- a/selda-postgresql/src/Database/Selda/PostgreSQL.hs +++ b/selda-postgresql/src/Database/Selda/PostgreSQL.hs @@ -205,6 +205,8 @@ pgBackend c = SeldaBackend , runStmtWithPK = \q ps -> left <$> pgQueryRunner c True q ps , prepareStmt = pgPrepare c , runPrepared = pgRun c + , runPreparedStreaming + = pgRunStream c , getTableInfo = pgGetTableInfo c . rawTableName , backendId = PostgreSQL , ppConfig = pgPPConfig From 92066ee37a90ccb18f03efb191eabb3dcf9cab70 Mon Sep 17 00:00:00 2001 From: Mirek Kratochvil Date: Fri, 24 Oct 2025 11:46:49 +0200 Subject: [PATCH 37/39] add prep-stmt streaming to sqlite backend --- selda-sqlite/src/Database/Selda/SQLite.hs | 29 +++++++++++++++++++---- 1 file changed, 24 insertions(+), 5 deletions(-) diff --git a/selda-sqlite/src/Database/Selda/SQLite.hs b/selda-sqlite/src/Database/Selda/SQLite.hs index 5db79a4..862cf89 100644 --- a/selda-sqlite/src/Database/Selda/SQLite.hs +++ b/selda-sqlite/src/Database/Selda/SQLite.hs @@ -76,6 +76,8 @@ sqliteBackend db = SeldaBackend , runStmtWithPK = \q ps -> fst <$> sqliteQueryRunner db q ps , prepareStmt = \_ _ -> sqlitePrepare db , runPrepared = sqliteRunPrepared db + , runPreparedStreaming + = sqliteRunPreparedStreaming db , getTableInfo = sqliteGetTableInfo db . fromTableName , ppConfig = defPPConfig {ppMaxInsertParams = Just 999} , backendId = SQLite @@ -186,6 +188,24 @@ sqliteRunPrepared db hdl params = do Left e@(SQLError{}) -> throwM (SqlError (show e)) Right res -> return (snd res) +sqliteRunPreparedStreaming :: + (MonadIO m, MonadMask m, Monoid r) + => Database + -> Dynamic + -> [Param] + -> ([[SqlValue]] -> m r) + -> m r +sqliteRunPreparedStreaming _ hdl params k = do + -- note for self: this does not need the `db` handle because it does not + -- return the last-inserted row ID. + eres <- try $ do + let Just stm = fromDynamic hdl + sqliteRunStmtStreaming stm params k `finally` do + liftIO $ clearBindings stm >> reset stm + case eres of + Left e@(SQLError{}) -> throwM (SqlError (show e)) + Right res -> return res + sqliteQueryRunner :: Database -> QueryRunner (Int64, (Int, [[SqlValue]])) sqliteQueryRunner db qry params = do eres <- try $ do @@ -224,20 +244,19 @@ sqliteQueryRunnerStreaming db qry params k = do eres <- try $ do stm <- liftIO $ prepare db qry - sqliteRunStmtStreaming db stm params k `finally` do + sqliteRunStmtStreaming stm params k `finally` do liftIO $ finalize stm case eres of Left e@(SQLError {}) -> throwM (SqlError (show e)) Right res -> return res sqliteRunStmtStreaming :: - (MonadIO m, MonadMask m, MonadIO m, Monoid r) - => Database - -> Statement + (MonadIO m, MonadMask m, Monoid r) + => Statement -> [Param] -> ([[SqlValue]] -> m r) -> m r -sqliteRunStmtStreaming db stm params k = do +sqliteRunStmtStreaming stm params k = do liftIO $ bind stm [toSqlData p | Param p <- params] streamRows stm 1024 $ k . map (map fromSqlData) From 21eddf60f8b22d1f38b68ad415059078781bcb61 Mon Sep 17 00:00:00 2001 From: Mirek Kratochvil Date: Fri, 24 Oct 2025 11:56:08 +0200 Subject: [PATCH 38/39] bump versions of everything to mark a new feature, add changelog --- ChangeLog.hs | 12 ++++++++++-- selda-postgresql/selda-postgresql.cabal | 22 +++++++++++----------- selda-sqlite/selda-sqlite.cabal | 8 ++++---- selda/selda.cabal | 2 +- 4 files changed, 26 insertions(+), 18 deletions(-) diff --git a/ChangeLog.hs b/ChangeLog.hs index 425eaf7..7bb2f80 100644 --- a/ChangeLog.hs +++ b/ChangeLog.hs @@ -10,13 +10,21 @@ import Text.Read changeLog :: ChangeLog changeLog = - [ Version "0.5.2.0" "2022-09-18" + [ Version "0.5.3.0" "2025-10-24" + "Streaming support and maintenance" + [ "Add support streaming of large results (#200)" + , "Publish the website via GitHub pages (#201)" + , "Properly use cascading deletes with foreign keys (#191)" + , "Schema query validation fix for PostgreSQL (#183)" + , "Many small version compatibility fixes" + ] + , Version "0.5.2.0" "2022-09-18" "Quality of life improvements" [ "Add support for GHC versions 8.10-9.2" , "Add typed UUIDs" , "Allow literal rows in update queries. (#139)" , "Add support for UNION/UNION ALL. (#140)" - , "Suppor raw PostgreSQL connetion strings. (#136)" + , "Support raw PostgreSQL connetion strings. (#136)" , "Move to Docker and GitHub actions for testing." , "Drop support for GHC versions <8.8." , "Various bugfixes." diff --git a/selda-postgresql/selda-postgresql.cabal b/selda-postgresql/selda-postgresql.cabal index 3505bda..a9ac4fb 100644 --- a/selda-postgresql/selda-postgresql.cabal +++ b/selda-postgresql/selda-postgresql.cabal @@ -1,5 +1,5 @@ name: selda-postgresql -version: 0.1.8.2 +version: 0.1.9.0 synopsis: PostgreSQL backend for the Selda database EDSL. description: PostgreSQL backend for the Selda database EDSL. Requires the PostgreSQL @libpq@ development libraries to be @@ -23,16 +23,16 @@ library OverloadedStrings CPP build-depends: - base >=4.9 && <5 - , bytestring >=0.9 && <0.13 - , exceptions >=0.8 && <0.11 - , selda >=0.5 && <0.6 - , selda-json >=0.1 && <0.2 - , text >=1.0 && <2.2 - , postgresql-binary >=0.12 && <0.15 - , postgresql-libpq >=0.9 && <0.12 - , time >=1.5 && <1.16 - , uuid-types >=1.0 && <1.1 + base >=4.9 && <5 + , bytestring >=0.9 && <0.13 + , exceptions >=0.8 && <0.11 + , selda >=0.5.3 && <0.6 + , selda-json >=0.1 && <0.2 + , text >=1.0 && <2.2 + , postgresql-binary >=0.12 && <0.15 + , postgresql-libpq >=0.9 && <0.12 + , time >=1.5 && <1.16 + , uuid-types >=1.0 && <1.1 hs-source-dirs: src default-language: diff --git a/selda-sqlite/selda-sqlite.cabal b/selda-sqlite/selda-sqlite.cabal index 122c904..feb1825 100644 --- a/selda-sqlite/selda-sqlite.cabal +++ b/selda-sqlite/selda-sqlite.cabal @@ -1,5 +1,5 @@ name: selda-sqlite -version: 0.1.7.2 +version: 0.1.8.0 synopsis: SQLite backend for the Selda database EDSL. description: Allows the Selda database EDSL to be used with SQLite databases. @@ -20,9 +20,9 @@ library GADTs CPP build-depends: - base >=4.9 && <5 - , selda >=0.5 && <0.6 - , text >=1.0 && <2.2 + base >=4.9 && <5 + , selda >=0.5.3 && <0.6 + , text >=1.0 && <2.2 , bytestring >=0.10 && <0.13 , direct-sqlite >=2.2 && <2.4 , directory >=1.2.2 && <1.4 diff --git a/selda/selda.cabal b/selda/selda.cabal index e9860da..a978fc9 100644 --- a/selda/selda.cabal +++ b/selda/selda.cabal @@ -1,5 +1,5 @@ name: selda -version: 0.5.2.1 +version: 0.5.3.0 synopsis: Multi-backend, high-level EDSL for interacting with SQL databases. description: This package provides an EDSL for writing portable, type-safe, high-level database code. Its feature set includes querying and modifying databases, From 80b28174f67a95a43337ccd64de0d91ffea6d5be Mon Sep 17 00:00:00 2001 From: Mirek Kratochvil Date: Fri, 24 Oct 2025 18:24:13 +0200 Subject: [PATCH 39/39] fmt, add docs --- selda/src/Database/Selda/Frontend.hs | 16 ++++++++++++---- 1 file changed, 12 insertions(+), 4 deletions(-) diff --git a/selda/src/Database/Selda/Frontend.hs b/selda/src/Database/Selda/Frontend.hs index 04a330a..9269ac4 100644 --- a/selda/src/Database/Selda/Frontend.hs +++ b/selda/src/Database/Selda/Frontend.hs @@ -58,10 +58,18 @@ import Control.Monad.IO.Class ( MonadIO(..) ) query :: (MonadSelda m, Result a) => Query (Backend m) a -> m [Res a] query q = withBackend (flip queryWith q . runStmt) --- | Run a query within a Selda monad like `query` and stream the results. -forQuery :: forall m a r . (MonadSelda m, MonadMask m, Result a, Monoid r) => Query (Backend m) a -> (Res a -> m r) -> m r -forQuery q k = withBackend $ \b -> - uncurry (runStmtStreaming b) (ppConfig b `compileWith` q) $ fmap mconcat . traverse k . mkResults (Proxy :: Proxy a) +-- | Run a query within a Selda monad like `query` and stream the results to a +-- callback function. The callbacks may produce a Monoid which will be +-- concatenated into the final result. +forQuery :: + forall m a r. (MonadSelda m, MonadMask m, Result a, Monoid r) + => Query (Backend m) a + -> (Res a -> m r) + -> m r +forQuery q k = + withBackend $ \b -> + uncurry (runStmtStreaming b) (ppConfig b `compileWith` q) + $ fmap mconcat . traverse k . mkResults (Proxy :: Proxy a) -- | Perform the given query, and insert the result into the given table. -- Returns the number of inserted rows.