packages feed

streamly-lmdb 0.6.0 → 0.7.0

raw patch · 6 files changed

+174/−126 lines, 6 filesPVP ok

version bump matches the API change (PVP)

API changes (from Hackage documentation)

- Streamly.External.LMDB: [overwriteOptions] :: WriteOptions -> !OverwriteOptions
+ Streamly.External.LMDB: [readUnsafeFFI] :: ReadOptions -> !Bool
+ Streamly.External.LMDB: [writeOverwriteOptions] :: WriteOptions -> !OverwriteOptions
+ Streamly.External.LMDB: [writeUnsafeFFI] :: WriteOptions -> !Bool
+ Streamly.External.LMDB.Internal.Foreign: c_mdb_cursor_close_unsafe :: Ptr MDB_cursor -> IO ()
+ Streamly.External.LMDB.Internal.Foreign: c_mdb_cursor_get_unsafe :: Ptr MDB_cursor -> Ptr MDB_val -> Ptr MDB_val -> MDB_cursor_op_t -> IO CInt
+ Streamly.External.LMDB.Internal.Foreign: c_mdb_cursor_open_unsafe :: Ptr MDB_txn -> MDB_dbi_t -> Ptr (Ptr MDB_cursor) -> IO CInt
+ Streamly.External.LMDB.Internal.Foreign: c_mdb_get_unsafe :: Ptr MDB_txn -> MDB_dbi_t -> Ptr MDB_val -> Ptr MDB_val -> IO CInt
+ Streamly.External.LMDB.Internal.Foreign: c_mdb_put_unsafe_ :: Ptr MDB_txn -> MDB_dbi_t -> Ptr CChar -> CSize -> Ptr CChar -> CSize -> CUInt -> IO CInt
+ Streamly.External.LMDB.Internal.Foreign: c_mdb_txn_abort_unsafe :: Ptr MDB_txn -> IO ()
+ Streamly.External.LMDB.Internal.Foreign: c_mdb_txn_begin_unsafe :: Ptr MDB_env -> Ptr MDB_txn -> CUInt -> Ptr (Ptr MDB_txn) -> IO CInt
+ Streamly.External.LMDB.Internal.Foreign: c_mdb_txn_commit_unsafe :: Ptr MDB_txn -> IO CInt
+ Streamly.External.LMDB.Internal.Foreign: mdb_cursor_open_unsafe :: Ptr MDB_txn -> MDB_dbi_t -> IO (Ptr MDB_cursor)
+ Streamly.External.LMDB.Internal.Foreign: mdb_put_unsafe_ :: Ptr MDB_txn -> MDB_dbi_t -> Ptr CChar -> CSize -> Ptr CChar -> CSize -> CUInt -> IO ()
+ Streamly.External.LMDB.Internal.Foreign: mdb_txn_begin_unsafe :: Ptr MDB_env -> Ptr MDB_txn -> CUInt -> IO (Ptr MDB_txn)
+ Streamly.External.LMDB.Internal.Foreign: mdb_txn_commit_unsafe :: Ptr MDB_txn -> IO ()
- Streamly.External.LMDB: ReadOptions :: !ReadDirection -> !Maybe ByteString -> ReadOptions
+ Streamly.External.LMDB: ReadOptions :: !ReadDirection -> !Maybe ByteString -> !Bool -> ReadOptions
- Streamly.External.LMDB: WriteOptions :: !Int -> !OverwriteOptions -> !Bool -> WriteOptions
+ Streamly.External.LMDB: WriteOptions :: !Int -> !OverwriteOptions -> !Bool -> !Bool -> WriteOptions

Files

ChangeLog.md view
@@ -1,3 +1,7 @@+## 0.7.0++* Added `readUnsafeFFI` and `writeUnsafeFFI` options.+ ## 0.6.0  * Updated for Streamly 0.9.0.
README.md view
@@ -67,11 +67,17 @@  ## Benchmarks -See `bench/README.md`. Summary (with rough figures from our machine<sup>†</sup> using GHC 9.2.7 [GHC 8.10.7]):+See `bench/README.md`. Summary (with rough figures from our machine<sup>†</sup>): -* **Reading.** For reading a fully cached LMDB database, this library (when `unsafeReadLMDB` is used instead of `readLMDB`) has roughly a 5 ns/pair [15 ns/pair] overhead compared to plain Haskell `IO` code, which has roughly another 10 ns/pair overhead compared to C. (The first two being similar fulfills the promise of [streamly](https://hackage.haskell.org/package/streamly) and stream fusion.) We deduce that if your total workload per pair takes longer than around 15 ns [25 ns], your bottleneck will not be your usage of this library as opposed to C.-* **Writing**. Writing with plain Haskell `IO` code and with this library is, respectively, roughly 15% [30%] and 50% slower than writing with C. We have not dug further into these differences because this write performance is currently good enough for our purposes.+* Reading:+  - For iterating through a fully cached LMDB database, this library has roughly a 110 ns/pair overhead compared to C. (Plain Haskell `IO` code has roughly a 70 ns/pair overhead compared to C. The two preceding figures being similar fulfills the promise of [streamly](https://hackage.haskell.org/package/streamly) and stream fusion.)+  - By using `unsafeReadLMDB` instead of `readLMDB`, we can get the overhead down to roughly 100 ns/pair.+  - By additionally using the `readUnsafeFFI` option (to use `unsafe` FFI calls under the hood), we can get the overhead down to roughly 40 ns/pair. -(There have apparently been some performance improvements between GHC 8.10.7 and 9.2.7.)+* Writing:+  - For writing to an LMDB database, this library has roughly a 210 ns/pair overhead compared to C. (Plain Haskell `IO` code has roughly a 100 ns/pair overhead compared to C. The two preceding figures being similar fulfills the promise of [streamly](https://hackage.haskell.org/package/streamly) and stream fusion.)+  - By using the `writeUnsafeFFI` option (to use `unsafe` FFI calls under the hood), we can get the overhead down to roughly 140 ns/pair. -<sup>†</sup> April 2023; [Linode](https://linode.com); Debian 11, Dedicated 32GB: 16 CPU, 640GB Storage, 32GB RAM.+* For most Haskell programs, these differences will not cause problems. (For instance, note that merely opening and reading 1 byte from a file with C already takes us tens of *microseconds*.)++<sup>†</sup> May 2023; [Linode](https://linode.com); Debian 11, Dedicated 32GB: 16 CPU, 640GB SSD storage, 32GB RAM.
src/Streamly/External/LMDB.hs view
@@ -80,44 +80,6 @@ import Foreign.Storable (poke) import Streamly.External.LMDB.Internal (Database (..), Mode (..), ReadOnly, ReadWrite) import Streamly.External.LMDB.Internal.Foreign-  ( LMDB_Error (..),-    MDB_ErrCode (..),-    MDB_cursor,-    MDB_env,-    MDB_txn,-    MDB_val (MDB_val, mv_data, mv_size),-    c_mdb_cursor_close,-    c_mdb_cursor_get,-    c_mdb_dbi_close,-    c_mdb_env_close,-    c_mdb_get,-    c_mdb_txn_abort,-    combineOptions,-    mdb_append,-    mdb_clear,-    mdb_create,-    mdb_cursor_open,-    mdb_dbi_open,-    mdb_env_create,-    mdb_env_open,-    mdb_env_set_mapsize,-    mdb_env_set_maxdbs,-    mdb_env_set_maxreaders,-    mdb_first,-    mdb_last,-    mdb_next,-    mdb_nooverwrite,-    mdb_nosubdir,-    mdb_notfound,-    mdb_notls,-    mdb_prev,-    mdb_put_,-    mdb_rdonly,-    mdb_set_range,-    mdb_txn_begin,-    mdb_txn_commit,-    throwLMDBErrNum,-  ) import Streamly.Internal.Data.Fold (Fold (Fold), Step (Partial)) import Streamly.Internal.Data.IOFinalizer (newIOFinalizer, runIOFinalizer) import Streamly.Internal.Data.Stream.StreamD.Type (Step (Stop, Yield))@@ -309,6 +271,23 @@         (Forward, Just _) -> (mdb_set_range, mdb_next)         (Backward, Nothing) -> (mdb_last, mdb_prev)         (Backward, Just _) -> (mdb_set_range, mdb_prev)+      (txn_begin, cursor_open, cursor_get, cursor_close, txn_abort) =+        if readUnsafeFFI ropts+          then+            ( mdb_txn_begin_unsafe,+              mdb_cursor_open_unsafe,+              c_mdb_cursor_get_unsafe,+              c_mdb_cursor_close_unsafe,+              c_mdb_txn_abort_unsafe+            )+          else+            ( mdb_txn_begin,+              mdb_cursor_open,+              c_mdb_cursor_get,+              c_mdb_cursor_close,+              c_mdb_txn_abort+            )+       supply = lmap . const    in supply firstOp $         Unfold@@ -317,22 +296,23 @@                 liftIO $                   if op == mdb_set_range && subsequentOp == mdb_prev                     then do-                      -- Reverse MDB_SET_RANGE.+                      -- A “reverse MDB_SET_RANGE” (i.e., a “less than or equal to”) is not+                      -- available in LMDB, so we simulate it ourselves.                       kfst' <- peek pk                       kfst <- packCStringLen (mv_data kfst', fromIntegral $ mv_size kfst')-                      rc <- c_mdb_cursor_get pcurs pk pv op+                      rc <- cursor_get pcurs pk pv op                       if rc /= 0 && rc == mdb_notfound-                        then c_mdb_cursor_get pcurs pk pv mdb_last+                        then cursor_get pcurs pk pv mdb_last                         else                           if rc == 0                             then do                               k' <- peek pk                               k <- unsafePackCStringLen (mv_data k', fromIntegral $ mv_size k')                               if k /= kfst-                                then c_mdb_cursor_get pcurs pk pv mdb_prev+                                then cursor_get pcurs pk pv mdb_prev                                 else return rc                             else return rc-                    else c_mdb_cursor_get pcurs pk pv op+                    else cursor_get pcurs pk pv op                found <-                 liftIO $@@ -355,25 +335,28 @@               (pcurs, pk, pv, ref) <- liftIO $                 mask_ $ do                   (ptxn, pcurs) <- case mtxncurs of-                    Nothing -> liftIO $ do-                      ptxn <- mdb_txn_begin penv nullPtr mdb_rdonly-                      pcurs <- mdb_cursor_open ptxn dbi+                    Nothing -> do+                      ptxn <- txn_begin penv nullPtr mdb_rdonly+                      pcurs <- cursor_open ptxn dbi                       return (ptxn, pcurs)                     Just (ReadOnlyTxn ptxn, Cursor pcurs) ->                       return (ptxn, pcurs)-                  pk <- liftIO malloc-                  pv <- liftIO malloc+                  pk <- malloc+                  pv <- malloc                    _ <- case readStart ropts of                     Nothing -> return ()                     Just k -> unsafeUseAsCStringLen k $ \(kp, kl) ->                       poke pk (MDB_val (fromIntegral kl) kp) -                  ref <- liftIO . newIOFinalizer $ do+                  ref <- newIOFinalizer $ do                     free pv >> free pk                     when (isNothing mtxncurs) $-                      -- There is no need to commit this read-only transaction.-                      c_mdb_cursor_close pcurs >> c_mdb_txn_abort ptxn+                      -- With LMDB, there is ordinarily no need to commit read-only transactions.+                      -- (The exception is when we want to make databases that were opened during+                      -- the transaction available later, but that’s not applicable here.) We can+                      -- therefore abort ptxn, both for failure (exceptions) and success.+                      cursor_close pcurs >> txn_abort ptxn                   return (pcurs, pk, pv, ref)               return (op, pcurs, pk, pv, ref)           )@@ -412,35 +395,60 @@   { readDirection :: !ReadDirection,     -- | If 'Nothing', a forward [backward] iteration starts at the beginning [end] of the database.     -- Otherwise, it starts at the first key that is greater [less] than or equal to the 'Just' key.-    readStart :: !(Maybe ByteString)+    readStart :: !(Maybe ByteString),+    -- | Use @unsafe@ FFI calls under the hood. This can increase iteration speed, but one should+    -- bear in mind that @unsafe@ FFI calls can have an adverse impact on the performance of the+    -- rest of the program (e.g., its ability to effectively spawn green threads).+    readUnsafeFFI :: !Bool   }   deriving (Show) --- | By default, we start reading from the beginning of the database (i.e., from the smallest key).+-- | By default, we start reading from the beginning of the database (i.e., from the smallest key),+-- and we don’t use unsafe FFI calls. defaultReadOptions :: ReadOptions defaultReadOptions =   ReadOptions     { readDirection = Forward,-      readStart = Nothing+      readStart = Nothing,+      readUnsafeFFI = False     }  -- | Direction of key iteration. data ReadDirection = Forward | Backward deriving (Show) -data OverwriteOptions = OverwriteAllow | OverwriteAllowSame | OverwriteDisallow deriving (Eq)+data OverwriteOptions+  = -- | When a key reoccurs, overwrite the value.+    OverwriteAllow+  | -- | When a key reoccurs, throw an exception except when the value is the same.+    OverwriteAllowSame+  | -- | When a key reoccurs, throw an exception.+    OverwriteDisallow+  deriving (Eq)  data WriteOptions = WriteOptions-  { writeTransactionSize :: !Int,-    overwriteOptions :: !OverwriteOptions,-    writeAppend :: !Bool+  { -- | The number of key-value pairs per write transaction.+    writeTransactionSize :: !Int,+    writeOverwriteOptions :: !OverwriteOptions,+    -- | Assume the input data is already ordered. This allows the use of @MDB_APPEND@ under the+    -- hood and substantially improves write performance. An exception will be thrown if the+    -- assumption about the ordering is not true.+    writeAppend :: !Bool,+    -- | Use @unsafe@ FFI calls under the hood. This can increase write performance, but one should+    -- bear in mind that @unsafe@ FFI calls can have an adverse impact on the performance of the+    -- rest of the program (e.g., its ability to effectively spawn green threads).+    writeUnsafeFFI :: !Bool   } +-- | By default, we use a write transaction size of 1 (one write transaction for each key-value+-- pair), allow overwriting, don’t assume that the input data is already ordered, and don’t use+-- unsafe FFI calls. defaultWriteOptions :: WriteOptions defaultWriteOptions =   WriteOptions     { writeTransactionSize = 1,-      overwriteOptions = OverwriteAllow,-      writeAppend = False+      writeOverwriteOptions = OverwriteAllow,+      writeAppend = False,+      writeUnsafeFFI = False     }  newtype ExceptionString = ExceptionString String deriving (Show)@@ -459,13 +467,17 @@ -- KB chunks and benchmark from there. {-# INLINE writeLMDB #-} writeLMDB :: (MonadIO m) => Database ReadWrite -> WriteOptions -> Fold m (ByteString, ByteString) ()-writeLMDB (Database penv dbi) options =-  let txnSize = max 1 (writeTransactionSize options)-      overwriteOpt = overwriteOptions options+writeLMDB (Database penv dbi) wopts =+  let txnSize = max 1 (writeTransactionSize wopts)+      overwriteOpt = writeOverwriteOptions wopts       flags =         combineOptions $           [mdb_nooverwrite | overwriteOpt `elem` [OverwriteAllowSame, OverwriteDisallow]]-            ++ [mdb_append | writeAppend options]+            ++ [mdb_append | writeAppend wopts]+      (txn_begin, txn_commit, put_, get) =+        if writeUnsafeFFI wopts+          then (mdb_txn_begin_unsafe, mdb_txn_commit_unsafe, mdb_put_unsafe_, c_mdb_get_unsafe)+          else (mdb_txn_begin, mdb_txn_commit, mdb_put_, c_mdb_get)    in Fold         ( \(threadId, iter, currChunkSz, mtxn) (k, v) -> do             -- In the first few iterations, ascertain that we are still on the same (bound) thread.@@ -493,21 +505,21 @@               if currChunkSz' == 0                 then liftIO $                   mask_ $ do-                    ptxn <- mdb_txn_begin penv nullPtr 0-                    ref <- newIOFinalizer $ mdb_txn_commit ptxn+                    ptxn <- txn_begin penv nullPtr 0+                    ref <- newIOFinalizer $ txn_commit ptxn                     return (ptxn, ref)                 else return $ fromJust mtxn              liftIO $               unsafeUseAsCStringLen k $ \(kp, kl) -> unsafeUseAsCStringLen v $ \(vp, vl) ->                 catch-                  (mdb_put_ ptxn dbi kp (fromIntegral kl) vp (fromIntegral vl) flags)+                  (put_ ptxn dbi kp (fromIntegral kl) vp (fromIntegral vl) flags)                   ( \(e :: LMDB_Error) -> do                       -- Discard error if OverwriteAllowSame was specified and the error from LMDB                       -- was due to the exact same key-value pair already existing in the database.                       ok <- with (MDB_val (fromIntegral kl) kp) $ \pk ->                         alloca $ \pv -> do-                          rc <- c_mdb_get ptxn dbi pk pv+                          rc <- get ptxn dbi pk pv                           if rc == 0                             then do                               v' <- peek pv
src/Streamly/External/LMDB/Internal/Foreign.hsc view
@@ -40,61 +40,85 @@         #{poke MDB_val, mv_data} ptr pd     {-# INLINE poke #-} -foreign import ccall unsafe "lmdb.h mdb_strerror"+foreign import ccall safe "lmdb.h mdb_strerror"     c_mdb_strerror :: CInt -> IO CString -foreign import ccall unsafe "lmdb.h mdb_env_create"+foreign import ccall safe "lmdb.h mdb_env_create"     c_mdb_env_create :: Ptr (Ptr MDB_env) -> IO CInt -foreign import ccall unsafe "lmdb.h mdb_env_set_mapsize"+foreign import ccall safe "lmdb.h mdb_env_set_mapsize"     c_mdb_env_set_mapsize :: Ptr MDB_env -> CSize -> IO CInt -foreign import ccall unsafe "lmdb.h mdb_env_set_maxreaders"+foreign import ccall safe "lmdb.h mdb_env_set_maxreaders"     c_mdb_env_set_maxreaders :: Ptr MDB_env -> CUInt -> IO CInt -foreign import ccall unsafe "lmdb.h mdb_env_set_maxdbs"+foreign import ccall safe "lmdb.h mdb_env_set_maxdbs"     c_mdb_env_set_maxdbs :: Ptr MDB_env -> MDB_dbi_t -> IO CInt -foreign import ccall unsafe "lmdb.h mdb_env_open"+foreign import ccall safe "lmdb.h mdb_env_open"     c_mdb_env_open :: Ptr MDB_env -> CString -> CUInt -> MDB_mode_t -> IO CInt -foreign import ccall unsafe "lmdb.h mdb_txn_begin"+foreign import ccall safe "lmdb.h mdb_txn_begin"     c_mdb_txn_begin :: Ptr MDB_env -> Ptr MDB_txn -> CUInt -> Ptr (Ptr MDB_txn) -> IO CInt -foreign import ccall unsafe "lmdb.h mdb_dbi_open"+foreign import ccall unsafe "lmdb.h mdb_txn_begin"+    c_mdb_txn_begin_unsafe :: Ptr MDB_env -> Ptr MDB_txn -> CUInt -> Ptr (Ptr MDB_txn) -> IO CInt++foreign import ccall safe "lmdb.h mdb_dbi_open"     c_mdb_dbi_open :: Ptr MDB_txn -> CString -> CUInt -> Ptr MDB_dbi_t -> IO CInt -foreign import ccall unsafe "lmdb.h mdb_txn_commit"+foreign import ccall safe "lmdb.h mdb_txn_commit"     c_mdb_txn_commit :: Ptr MDB_txn -> IO CInt -foreign import ccall unsafe "lmdb.h mdb_txn_abort"+foreign import ccall unsafe "lmdb.h mdb_txn_commit"+    c_mdb_txn_commit_unsafe :: Ptr MDB_txn -> IO CInt++foreign import ccall safe "lmdb.h mdb_txn_abort"     c_mdb_txn_abort :: Ptr MDB_txn -> IO () -foreign import ccall unsafe "lmdb.h mdb_cursor_open"+foreign import ccall unsafe "lmdb.h mdb_txn_abort"+    c_mdb_txn_abort_unsafe :: Ptr MDB_txn -> IO ()++foreign import ccall safe "lmdb.h mdb_cursor_open"     c_mdb_cursor_open :: Ptr MDB_txn -> MDB_dbi_t -> Ptr (Ptr MDB_cursor) -> IO CInt -foreign import ccall unsafe "lmdb.h mdb_cursor_get"+foreign import ccall unsafe "lmdb.h mdb_cursor_open"+    c_mdb_cursor_open_unsafe :: Ptr MDB_txn -> MDB_dbi_t -> Ptr (Ptr MDB_cursor) -> IO CInt++foreign import ccall safe "lmdb.h mdb_cursor_get"     c_mdb_cursor_get :: Ptr MDB_cursor -> Ptr MDB_val -> Ptr MDB_val -> MDB_cursor_op_t -> IO CInt -foreign import ccall unsafe "lmdb.h mdb_cursor_close"+foreign import ccall unsafe "lmdb.h mdb_cursor_get"+    c_mdb_cursor_get_unsafe :: Ptr MDB_cursor -> Ptr MDB_val -> Ptr MDB_val -> MDB_cursor_op_t -> IO CInt++foreign import ccall safe "lmdb.h mdb_cursor_close"     c_mdb_cursor_close :: Ptr MDB_cursor -> IO () -foreign import ccall unsafe "lmdb.h mdb_dbi_close"+foreign import ccall unsafe "lmdb.h mdb_cursor_close"+    c_mdb_cursor_close_unsafe :: Ptr MDB_cursor -> IO ()++foreign import ccall safe "lmdb.h mdb_dbi_close"     c_mdb_dbi_close :: Ptr MDB_env -> MDB_dbi_t -> IO () -foreign import ccall unsafe "lmdb.h mdb_env_close"+foreign import ccall safe "lmdb.h mdb_env_close"     c_mdb_env_close :: Ptr MDB_env -> IO () -foreign import ccall unsafe "lmdb.h mdb_get"+foreign import ccall safe "lmdb.h mdb_get"     c_mdb_get :: Ptr MDB_txn -> MDB_dbi_t -> Ptr MDB_val -> Ptr MDB_val -> IO CInt -foreign import ccall unsafe "lmdb.h mdb_put"+foreign import ccall unsafe "lmdb.h mdb_get"+    c_mdb_get_unsafe :: Ptr MDB_txn -> MDB_dbi_t -> Ptr MDB_val -> Ptr MDB_val -> IO CInt++foreign import ccall safe "lmdb.h mdb_put"     c_mdb_put :: Ptr MDB_txn -> MDB_dbi_t -> Ptr MDB_val -> Ptr MDB_val -> CUInt -> IO CInt -foreign import ccall unsafe "streamly_lmdb_foreign.h mdb_put_"+foreign import ccall safe "streamly_lmdb_foreign.h mdb_put_"     c_mdb_put_ :: Ptr MDB_txn -> MDB_dbi_t -> Ptr CChar -> CSize -> Ptr CChar -> CSize -> CUInt -> IO CInt -foreign import ccall unsafe "lmdb.h mdb_drop"+foreign import ccall unsafe "streamly_lmdb_foreign.h mdb_put_"+    c_mdb_put_unsafe_ :: Ptr MDB_txn -> MDB_dbi_t -> Ptr CChar -> CSize -> Ptr CChar -> CSize -> CUInt -> IO CInt++foreign import ccall safe "lmdb.h mdb_drop"     c_mdb_drop :: Ptr MDB_txn -> MDB_dbi_t -> CInt -> IO CInt  data LMDB_Error = LMDB_Error@@ -237,17 +261,33 @@     alloca $ \pptxn -> c_mdb_txn_begin penv parent flags pptxn >>= \rc ->         if rc /= 0 then throwLMDBErrNum "mdb_txn_begin" rc else peek pptxn +mdb_txn_begin_unsafe :: Ptr MDB_env -> Ptr MDB_txn -> CUInt -> IO (Ptr MDB_txn)+mdb_txn_begin_unsafe penv parent flags =+    alloca $ \pptxn -> c_mdb_txn_begin_unsafe penv parent flags pptxn >>= \rc ->+        if rc /= 0 then throwLMDBErrNum "mdb_txn_begin" rc else peek pptxn+ -- If the commit fails, aborts the transaction. mdb_txn_commit :: Ptr MDB_txn -> IO () mdb_txn_commit ptxn =     c_mdb_txn_commit ptxn >>= \rc ->         when (rc /= 0) $ c_mdb_txn_abort ptxn >> throwLMDBErrNum "mdb_txn_commit" rc +-- If the commit fails, aborts the transaction.+mdb_txn_commit_unsafe :: Ptr MDB_txn -> IO ()+mdb_txn_commit_unsafe ptxn =+    c_mdb_txn_commit_unsafe ptxn >>= \rc ->+        when (rc /= 0) $ c_mdb_txn_abort_unsafe ptxn >> throwLMDBErrNum "mdb_txn_commit" rc+ mdb_cursor_open :: Ptr MDB_txn -> MDB_dbi_t -> IO (Ptr MDB_cursor) mdb_cursor_open ptxn dbi =     alloca $ \ppcurs -> c_mdb_cursor_open ptxn dbi ppcurs >>= \rc ->         if rc /= 0 then c_mdb_txn_abort ptxn >> throwLMDBErrNum "mdb_cursor_open" rc else peek ppcurs +mdb_cursor_open_unsafe :: Ptr MDB_txn -> MDB_dbi_t -> IO (Ptr MDB_cursor)+mdb_cursor_open_unsafe ptxn dbi =+    alloca $ \ppcurs -> c_mdb_cursor_open_unsafe ptxn dbi ppcurs >>= \rc ->+        if rc /= 0 then c_mdb_txn_abort ptxn >> throwLMDBErrNum "mdb_cursor_open" rc else peek ppcurs+ mdb_dbi_open :: Ptr MDB_txn -> Maybe String -> CUInt -> IO MDB_dbi_t mdb_dbi_open ptxn name flags = do     withCStringMaybe name $ \cname ->@@ -264,6 +304,12 @@ mdb_put_ :: Ptr MDB_txn -> MDB_dbi_t -> Ptr CChar -> CSize -> Ptr CChar -> CSize -> CUInt -> IO () mdb_put_ ptxn dbi pk ks pv vs flags =     c_mdb_put_ ptxn dbi pk ks pv vs flags >>= \rc ->+        when (rc /= 0) $ throwLMDBErrNum "mdb_put_" rc++{-# INLINE mdb_put_unsafe_ #-}+mdb_put_unsafe_ :: Ptr MDB_txn -> MDB_dbi_t -> Ptr CChar -> CSize -> Ptr CChar -> CSize -> CUInt -> IO ()+mdb_put_unsafe_ ptxn dbi pk ks pv vs flags =+    c_mdb_put_unsafe_ ptxn dbi pk ks pv vs flags >>= \rc ->         when (rc /= 0) $ throwLMDBErrNum "mdb_put_" rc  mdb_clear :: Ptr MDB_txn -> MDB_dbi_t -> IO ()
streamly-lmdb.cabal view
@@ -1,6 +1,6 @@ cabal-version:  3.0 name:           streamly-lmdb-version:        0.6.0+version:        0.7.0 synopsis:       Stream data to or from LMDB databases using the streamly library. description:    Please see the README on GitHub at <https://github.com/shlok/streamly-lmdb#readme> category:       Database, Streaming, Streamly
test/Streamly/External/LMDB/Tests.hs view
@@ -11,38 +11,11 @@ import Data.List (find, foldl', nubBy, sort) import Data.Word (Word8) import Foreign (castPtr, nullPtr, with)+import Streamly.Data.Stream.Prelude (fromList, toList, unfold)+import qualified Streamly.Data.Stream.Prelude as S import Streamly.External.LMDB-  ( Cursor,-    Environment,-    Mode,-    OverwriteOptions (..),-    ReadDirection (..),-    ReadOnlyTxn,-    ReadOptions (..),-    ReadWrite,-    WriteOptions (..),-    abortReadOnlyTxn,-    beginReadOnlyTxn,-    clearDatabase,-    closeCursor,-    defaultReadOptions,-    defaultWriteOptions,-    openCursor,-    readLMDB,-    unsafeReadLMDB,-    writeLMDB,-  ) import Streamly.External.LMDB.Internal (Database (..)) import Streamly.External.LMDB.Internal.Foreign-  ( MDB_val (..),-    combineOptions,-    mdb_nooverwrite,-    mdb_put,-    mdb_txn_begin,-    mdb_txn_commit,-  )-import Streamly.Data.Stream.Prelude (fromList, toList, unfold)-import qualified Streamly.Data.Stream.Prelude as S import Test.QuickCheck (Gen, NonEmptyList (..), choose, elements, frequency) import Test.QuickCheck.Monadic (PropertyM, monadicIO, pick, run) import Test.Tasty (TestTree)@@ -117,12 +90,14 @@   run $ clearDatabase db    chunkSz <- pick arbitrary+  unsafeFFI <- pick arbitrary    let fol' =         writeLMDB db $           defaultWriteOptions             { writeTransactionSize = chunkSz,-              overwriteOptions = OverwriteAllow+              writeOverwriteOptions = OverwriteAllow,+              writeUnsafeFFI = unsafeFFI             }    -- TODO: Run with new "bound" functionality in streamly.@@ -143,13 +118,15 @@   run $ clearDatabase db    chunkSz <- pick arbitrary+  unsafeFFI <- pick arbitrary    -- TODO: Run with new "bound" functionality in streamly.   let fol' =         writeLMDB db $           defaultWriteOptions             { writeTransactionSize = chunkSz,-              overwriteOptions = OverwriteDisallow+              writeOverwriteOptions = OverwriteDisallow,+              writeUnsafeFFI = unsafeFFI             }   e <- run $ try @SomeException $ (asyncBound (S.fold fol' (fromList keyValuePairs)) >>= wait)   exceptionAsExpected <-@@ -174,13 +151,15 @@   run $ clearDatabase db    chunkSz <- pick arbitrary+  unsafeFFI <- pick arbitrary    -- TODO: Run with new "bound" functionality in streamly.   let fol' =         writeLMDB db $           defaultWriteOptions             { writeTransactionSize = chunkSz,-              overwriteOptions = OverwriteAllowSame+              writeOverwriteOptions = OverwriteAllowSame,+              writeUnsafeFFI = unsafeFFI             }   e <- run $ try @SomeException $ (asyncBound (S.fold fol' (fromList keyValuePairs)) >>= wait)   exceptionAsExpected <-@@ -310,9 +289,10 @@ readOptionsAndResults pairsInDb = do   forw <- arbitrary   let dir = if forw then Forward else Backward+  unsafeFFI <- arbitrary   let len = length pairsInDb   readAll <- frequency [(1, return True), (3, return False)]-  let ropts = defaultReadOptions {readDirection = dir}+  let ropts = defaultReadOptions {readDirection = dir, readUnsafeFFI = unsafeFFI}   if readAll     then return (ropts {readStart = Nothing}, (if forw then id else reverse) pairsInDb)     else