Compare commits
8 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
7e6c6af3dd
|
|||
|
faf49fafd5
|
|||
|
04815f66a4
|
|||
|
fd895de148
|
|||
|
b618ef1819
|
|||
|
407491f055
|
|||
|
2fdf6f0dad
|
|||
|
eb01962553
|
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "purescript-node-stream-pipes",
|
||||
"version": "v1.1.0",
|
||||
"version": "v1.2.3",
|
||||
"type": "module",
|
||||
"dependencies": {
|
||||
"csv-parse": "^5.5.5",
|
||||
|
||||
10
spago.lock
10
spago.lock
@@ -9,8 +9,8 @@ workspace:
|
||||
- either: ">=6.1.0 <7.0.0"
|
||||
- exceptions: ">=6.0.0 <7.0.0"
|
||||
- foldable-traversable: ">=6.0.0 <7.0.0"
|
||||
- foreign-object
|
||||
- lists
|
||||
- foreign-object: ">=4.1.0 <5.0.0"
|
||||
- lists: ">=7.0.0 <8.0.0"
|
||||
- maybe: ">=6.0.0 <7.0.0"
|
||||
- mmorph: ">=7.0.0 <8.0.0"
|
||||
- newtype: ">=5.0.0 <6.0.0"
|
||||
@@ -20,7 +20,7 @@ workspace:
|
||||
- node-path: ">=5.0.0 <6.0.0"
|
||||
- node-streams: ">=9.0.0 <10.0.0"
|
||||
- node-zlib: ">=0.4.0 <0.5.0"
|
||||
- ordered-collections
|
||||
- ordered-collections: ">=3.2.0 <4.0.0"
|
||||
- parallel: ">=6.0.0 <7.0.0"
|
||||
- pipes: ">=8.0.0 <9.0.0"
|
||||
- prelude: ">=6.0.1 <7.0.0"
|
||||
@@ -28,8 +28,8 @@ workspace:
|
||||
- strings: ">=6.0.1 <7.0.0"
|
||||
- tailrec: ">=6.1.0 <7.0.0"
|
||||
- transformers: ">=6.0.0 <7.0.0"
|
||||
- tuples
|
||||
- unordered-collections
|
||||
- tuples: ">=7.0.0 <8.0.0"
|
||||
- unordered-collections: ">=3.1.0 <4.0.0"
|
||||
- unsafe-coerce: ">=6.0.0 <7.0.0"
|
||||
test_dependencies:
|
||||
- console
|
||||
|
||||
12
spago.yaml
12
spago.yaml
@@ -1,7 +1,7 @@
|
||||
package:
|
||||
name: node-stream-pipes
|
||||
publish:
|
||||
version: '1.1.0'
|
||||
version: '1.2.3'
|
||||
license: 'GPL-3.0-or-later'
|
||||
location:
|
||||
githubOwner: 'cakekindel'
|
||||
@@ -10,17 +10,14 @@ package:
|
||||
strict: true
|
||||
pedanticPackages: true
|
||||
dependencies:
|
||||
- foreign-object
|
||||
- lists
|
||||
- ordered-collections
|
||||
- tuples
|
||||
- unordered-collections
|
||||
- aff: ">=7.1.0 <8.0.0"
|
||||
- arrays: ">=7.3.0 <8.0.0"
|
||||
- effect: ">=4.0.0 <5.0.0"
|
||||
- either: ">=6.1.0 <7.0.0"
|
||||
- exceptions: ">=6.0.0 <7.0.0"
|
||||
- foldable-traversable: ">=6.0.0 <7.0.0"
|
||||
- foreign-object: ">=4.1.0 <5.0.0"
|
||||
- lists: ">=7.0.0 <8.0.0"
|
||||
- maybe: ">=6.0.0 <7.0.0"
|
||||
- mmorph: ">=7.0.0 <8.0.0"
|
||||
- newtype: ">=5.0.0 <6.0.0"
|
||||
@@ -30,6 +27,7 @@ package:
|
||||
- node-path: ">=5.0.0 <6.0.0"
|
||||
- node-streams: ">=9.0.0 <10.0.0"
|
||||
- node-zlib: ">=0.4.0 <0.5.0"
|
||||
- ordered-collections: ">=3.2.0 <4.0.0"
|
||||
- parallel: ">=6.0.0 <7.0.0"
|
||||
- pipes: ">=8.0.0 <9.0.0"
|
||||
- prelude: ">=6.0.1 <7.0.0"
|
||||
@@ -37,6 +35,8 @@ package:
|
||||
- strings: ">=6.0.1 <7.0.0"
|
||||
- tailrec: ">=6.1.0 <7.0.0"
|
||||
- transformers: ">=6.0.0 <7.0.0"
|
||||
- tuples: ">=7.0.0 <8.0.0"
|
||||
- unordered-collections: ">=3.1.0 <4.0.0"
|
||||
- unsafe-coerce: ">=6.0.0 <7.0.0"
|
||||
test:
|
||||
main: Test.Main
|
||||
|
||||
1
src/Pipes.Construct.purs
Normal file
1
src/Pipes.Construct.purs
Normal file
@@ -0,0 +1 @@
|
||||
module Pipes.Construct where
|
||||
@@ -2,9 +2,11 @@ module Pipes.Node.FS where
|
||||
|
||||
import Prelude
|
||||
|
||||
import Control.Monad.Error.Class (class MonadThrow)
|
||||
import Data.Maybe (Maybe)
|
||||
import Effect.Aff (Aff)
|
||||
import Effect.Aff.Class (class MonadAff)
|
||||
import Effect.Class (liftEffect)
|
||||
import Effect.Exception (Error)
|
||||
import Node.Buffer (Buffer)
|
||||
import Node.FS.Stream (WriteStreamOptions)
|
||||
import Node.FS.Stream as FS.Stream
|
||||
@@ -22,11 +24,13 @@ import Prim.Row (class Union)
|
||||
-- | See `Pipes.Node.Stream.withEOS` for converting `Producer a`
|
||||
-- | into `Producer (Maybe a)`, emitting `Nothing` before exiting.
|
||||
write
|
||||
:: forall r trash
|
||||
:: forall r trash m
|
||||
. Union r trash WriteStreamOptions
|
||||
=> MonadAff m
|
||||
=> MonadThrow Error m
|
||||
=> Record r
|
||||
-> FilePath
|
||||
-> Consumer (Maybe Buffer) Aff Unit
|
||||
-> Consumer (Maybe Buffer) m Unit
|
||||
write o p = do
|
||||
w <- liftEffect $ FS.Stream.createWriteStream' p o
|
||||
fromWritable $ O.fromBufferWritable w
|
||||
@@ -34,26 +38,26 @@ write o p = do
|
||||
-- | Open a file in write mode, failing if the file already exists.
|
||||
-- |
|
||||
-- | `write {flags: "wx"}`
|
||||
create :: FilePath -> Consumer (Maybe Buffer) Aff Unit
|
||||
create :: forall m. MonadAff m => MonadThrow Error m => FilePath -> Consumer (Maybe Buffer) m Unit
|
||||
create = write { flags: "wx" }
|
||||
|
||||
-- | Open a file in write mode, truncating it if the file already exists.
|
||||
-- |
|
||||
-- | `write {flags: "w"}`
|
||||
truncate :: FilePath -> Consumer (Maybe Buffer) Aff Unit
|
||||
truncate :: forall m. MonadAff m => MonadThrow Error m => FilePath -> Consumer (Maybe Buffer) m Unit
|
||||
truncate = write { flags: "w" }
|
||||
|
||||
-- | Open a file in write mode, appending written contents if the file already exists.
|
||||
-- |
|
||||
-- | `write {flags: "a"}`
|
||||
append :: FilePath -> Consumer (Maybe Buffer) Aff Unit
|
||||
append :: forall m. MonadAff m => MonadThrow Error m => FilePath -> Consumer (Maybe Buffer) m Unit
|
||||
append = write { flags: "a" }
|
||||
|
||||
-- | Creates a `fs.Readable` stream for the file at the given path.
|
||||
-- |
|
||||
-- | Emits `Nothing` before closing. To opt out of this behavior,
|
||||
-- | use `Pipes.Node.Stream.withoutEOS` or `Pipes.Node.Stream.unEOS`.
|
||||
read :: FilePath -> Producer (Maybe Buffer) Aff Unit
|
||||
read :: forall m. MonadAff m => MonadThrow Error m => FilePath -> Producer (Maybe Buffer) m Unit
|
||||
read p = do
|
||||
r <- liftEffect $ FS.Stream.createReadStream p
|
||||
fromReadable $ O.fromBufferReadable r
|
||||
|
||||
@@ -2,28 +2,31 @@ module Pipes.Node.Stream where
|
||||
|
||||
import Prelude
|
||||
|
||||
import Control.Monad.Error.Class (throwError)
|
||||
import Control.Monad.Rec.Class (Step(..), tailRecM, whileJust)
|
||||
import Control.Monad.Error.Class (class MonadThrow, throwError)
|
||||
import Control.Monad.Rec.Class (class MonadRec, Step(..), tailRecM)
|
||||
import Control.Monad.ST.Class (liftST)
|
||||
import Control.Monad.ST.Ref as STRef
|
||||
import Control.Monad.Trans.Class (lift)
|
||||
import Data.Maybe (Maybe(..), maybe)
|
||||
import Data.Maybe (Maybe(..))
|
||||
import Data.Newtype (wrap)
|
||||
import Data.Traversable (for_)
|
||||
import Effect.Aff (Aff, delay)
|
||||
import Effect.Aff.Class (liftAff)
|
||||
import Data.Tuple.Nested ((/\))
|
||||
import Effect.Aff (delay)
|
||||
import Effect.Aff.Class (class MonadAff, liftAff)
|
||||
import Effect.Class (liftEffect)
|
||||
import Effect.Exception (Error)
|
||||
import Node.Stream.Object as O
|
||||
import Pipes (await, yield, (>->))
|
||||
import Pipes (await, yield)
|
||||
import Pipes (for) as P
|
||||
import Pipes.Core (Consumer, Pipe, Producer, Producer_)
|
||||
import Pipes.Prelude (mapFoldable, map) as P
|
||||
import Pipes.Prelude (mapFoldable) as P
|
||||
import Pipes.Util (InvokeResult(..), invoke)
|
||||
|
||||
-- | Convert a `Readable` stream to a `Pipe`.
|
||||
-- |
|
||||
-- | This will yield `Nothing` before exiting, signaling
|
||||
-- | End-of-stream.
|
||||
fromReadable :: forall s a. O.Read s a => s -> Producer_ (Maybe a) Aff Unit
|
||||
fromReadable :: forall s a m. MonadThrow Error m => MonadAff m => O.Read s a => s -> Producer_ (Maybe a) m Unit
|
||||
fromReadable r =
|
||||
let
|
||||
cleanup rmErrorListener = do
|
||||
@@ -38,7 +41,7 @@ fromReadable r =
|
||||
res <- liftEffect $ O.read r
|
||||
case res of
|
||||
O.ReadJust a -> yield (Just a) $> Loop { error, cancel }
|
||||
O.ReadWouldBlock -> lift (O.awaitReadableOrClosed r) $> Loop { error, cancel }
|
||||
O.ReadWouldBlock -> liftAff (O.awaitReadableOrClosed r) $> Loop { error, cancel }
|
||||
O.ReadClosed -> yield Nothing *> cleanup cancel
|
||||
in
|
||||
do
|
||||
@@ -49,7 +52,7 @@ fromReadable r =
|
||||
-- |
|
||||
-- | When `Nothing` is piped to this, the stream will
|
||||
-- | be `end`ed, and the pipe will noop if invoked again.
|
||||
fromWritable :: forall s a. O.Write s a => s -> Consumer (Maybe a) Aff Unit
|
||||
fromWritable :: forall s a m. MonadThrow Error m => MonadAff m => O.Write s a => s -> Consumer (Maybe a) m Unit
|
||||
fromWritable w =
|
||||
let
|
||||
cleanup rmErrorListener = do
|
||||
@@ -82,7 +85,7 @@ fromWritable w =
|
||||
-- |
|
||||
-- | When `Nothing` is piped to this, the `Transform` stream will
|
||||
-- | be `end`ed, and the pipe will noop if invoked again.
|
||||
fromTransform :: forall a b. O.Transform a b -> Pipe (Maybe a) (Maybe b) Aff Unit
|
||||
fromTransform :: forall a b m. MonadThrow Error m => MonadAff m => O.Transform a b -> Pipe (Maybe a) (Maybe b) m Unit
|
||||
fromTransform t =
|
||||
let
|
||||
cleanup removeErrorListener = do
|
||||
@@ -111,7 +114,7 @@ fromTransform t =
|
||||
O.WriteClosed -> cleanup cancel
|
||||
O.WriteOk -> pure $ Loop { error, cancel }
|
||||
O.WriteWouldBlock -> do
|
||||
lift (O.awaitWritableOrClosed t)
|
||||
liftAff $ O.awaitWritableOrClosed t
|
||||
pure $ Loop { error, cancel }
|
||||
in
|
||||
do
|
||||
@@ -121,13 +124,13 @@ fromTransform t =
|
||||
-- | Given a `Producer` of values, wrap them in `Just`.
|
||||
-- |
|
||||
-- | Before the `Producer` exits, emits `Nothing` as an End-of-stream signal.
|
||||
withEOS :: forall a. Producer a Aff Unit -> Producer (Maybe a) Aff Unit
|
||||
withEOS :: forall a m. Monad m => Producer a m Unit -> Producer (Maybe a) m Unit
|
||||
withEOS a = do
|
||||
P.for a (yield <<< Just)
|
||||
yield Nothing
|
||||
|
||||
-- | Strip a pipeline of the EOS signal
|
||||
unEOS :: forall a. Pipe (Maybe a) a Aff Unit
|
||||
unEOS :: forall a m. Monad m => Pipe (Maybe a) a m Unit
|
||||
unEOS = P.mapFoldable identity
|
||||
|
||||
-- | Lift a `Pipe a a` to `Pipe (Maybe a) (Maybe a)`.
|
||||
@@ -140,8 +143,16 @@ unEOS = P.mapFoldable identity
|
||||
-- | `Just` values will be passed to the pipe, and the response(s) will be wrapped in `Just`.
|
||||
-- |
|
||||
-- | `Nothing` will bypass the given pipe entirely, and the pipe will not be invoked again.
|
||||
inEOS :: forall a b. Pipe a b Aff Unit -> Pipe (Maybe a) (Maybe b) Aff Unit
|
||||
inEOS p = whileJust do
|
||||
inEOS :: forall a b m. MonadRec m => Pipe a b m Unit -> Pipe (Maybe a) (Maybe b) m Unit
|
||||
inEOS p = flip tailRecM p \p' -> do
|
||||
ma <- await
|
||||
maybe (yield Nothing) (\a -> yield a >-> p >-> P.map Just) ma
|
||||
pure $ void ma
|
||||
case ma of
|
||||
Just a -> do
|
||||
res <- lift $ invoke p' a
|
||||
case res of
|
||||
Yielded (as /\ p'') -> do
|
||||
for_ (Just <$> as) yield
|
||||
pure $ Loop p''
|
||||
DidNotYield p'' -> pure $ Loop p''
|
||||
Exited -> yield Nothing $> Done unit
|
||||
_ -> yield Nothing $> Done unit
|
||||
|
||||
@@ -2,10 +2,12 @@ module Pipes.Node.Zlib where
|
||||
|
||||
import Prelude
|
||||
|
||||
import Control.Monad.Error.Class (class MonadThrow)
|
||||
import Data.Maybe (Maybe)
|
||||
import Effect (Effect)
|
||||
import Effect.Aff (Aff)
|
||||
import Effect.Aff.Class (class MonadAff)
|
||||
import Effect.Class (liftEffect)
|
||||
import Effect.Exception (Error)
|
||||
import Node.Buffer (Buffer)
|
||||
import Node.Stream.Object as O
|
||||
import Node.Zlib as Zlib
|
||||
@@ -13,28 +15,28 @@ import Node.Zlib.Types (ZlibStream)
|
||||
import Pipes.Core (Pipe)
|
||||
import Pipes.Node.Stream (fromTransform)
|
||||
|
||||
fromZlib :: forall r. Effect (ZlibStream r) -> Pipe (Maybe Buffer) (Maybe Buffer) Aff Unit
|
||||
fromZlib :: forall r m. MonadAff m => MonadThrow Error m => Effect (ZlibStream r) -> Pipe (Maybe Buffer) (Maybe Buffer) m Unit
|
||||
fromZlib z = do
|
||||
raw <- liftEffect $ Zlib.toDuplex <$> z
|
||||
fromTransform $ O.fromBufferTransform raw
|
||||
|
||||
gzip :: Pipe (Maybe Buffer) (Maybe Buffer) Aff Unit
|
||||
gzip :: forall m. MonadAff m => MonadThrow Error m => Pipe (Maybe Buffer) (Maybe Buffer) m Unit
|
||||
gzip = fromZlib Zlib.createGzip
|
||||
|
||||
gunzip :: Pipe (Maybe Buffer) (Maybe Buffer) Aff Unit
|
||||
gunzip :: forall m. MonadAff m => MonadThrow Error m => Pipe (Maybe Buffer) (Maybe Buffer) m Unit
|
||||
gunzip = fromZlib Zlib.createGunzip
|
||||
|
||||
unzip :: Pipe (Maybe Buffer) (Maybe Buffer) Aff Unit
|
||||
unzip :: forall m. MonadAff m => MonadThrow Error m => Pipe (Maybe Buffer) (Maybe Buffer) m Unit
|
||||
unzip = fromZlib Zlib.createUnzip
|
||||
|
||||
inflate :: Pipe (Maybe Buffer) (Maybe Buffer) Aff Unit
|
||||
inflate :: forall m. MonadAff m => MonadThrow Error m => Pipe (Maybe Buffer) (Maybe Buffer) m Unit
|
||||
inflate = fromZlib Zlib.createInflate
|
||||
|
||||
deflate :: Pipe (Maybe Buffer) (Maybe Buffer) Aff Unit
|
||||
deflate :: forall m. MonadAff m => MonadThrow Error m => Pipe (Maybe Buffer) (Maybe Buffer) m Unit
|
||||
deflate = fromZlib Zlib.createDeflate
|
||||
|
||||
brotliCompress :: Pipe (Maybe Buffer) (Maybe Buffer) Aff Unit
|
||||
brotliCompress :: forall m. MonadAff m => MonadThrow Error m => Pipe (Maybe Buffer) (Maybe Buffer) m Unit
|
||||
brotliCompress = fromZlib Zlib.createBrotliCompress
|
||||
|
||||
brotliDecompress :: Pipe (Maybe Buffer) (Maybe Buffer) Aff Unit
|
||||
brotliDecompress :: forall m. MonadAff m => MonadThrow Error m => Pipe (Maybe Buffer) (Maybe Buffer) m Unit
|
||||
brotliDecompress = fromZlib Zlib.createBrotliDecompress
|
||||
|
||||
@@ -3,17 +3,22 @@ module Pipes.Util where
|
||||
import Prelude
|
||||
|
||||
import Control.Monad.Maybe.Trans (MaybeT(..), runMaybeT)
|
||||
import Control.Monad.Rec.Class (whileJust)
|
||||
import Control.Monad.Rec.Class (class MonadRec, forever, whileJust)
|
||||
import Control.Monad.ST.Class (liftST)
|
||||
import Control.Monad.ST.Ref (STRef)
|
||||
import Control.Monad.ST.Ref as STRef
|
||||
import Control.Monad.Trans.Class (lift)
|
||||
import Data.Array.ST (STArray)
|
||||
import Data.Array.ST as Array.ST
|
||||
import Data.HashSet as HashSet
|
||||
import Data.Hashable (class Hashable, hash)
|
||||
import Data.List.NonEmpty (NonEmptyList)
|
||||
import Data.Maybe (Maybe(..))
|
||||
import Data.Tuple.Nested (type (/\), (/\))
|
||||
import Effect.Class (class MonadEffect, liftEffect)
|
||||
import Pipes (await, yield)
|
||||
import Pipes.Core (Pipe)
|
||||
import Pipes.Internal (Proxy(..))
|
||||
|
||||
-- | Yields a separator value `sep` between received values
|
||||
-- |
|
||||
@@ -62,3 +67,57 @@ chunked size = do
|
||||
when (len >= size) $ lift $ yield =<< Just <$> chunkTake
|
||||
yield =<< Just <$> chunkTake
|
||||
yield Nothing
|
||||
|
||||
-- | Equivalent of unix `uniq`, filtering out duplicate values passed to it.
|
||||
-- |
|
||||
-- | Uses a `HashSet` of hashes of `a`; for `n` elements `awaited`, this pipe
|
||||
-- | will occupy O(n) space, and `yield` in O(1) time.
|
||||
uniqHash :: forall a m. Hashable a => MonadEffect m => MonadRec m => Pipe a a m Unit
|
||||
uniqHash = do
|
||||
seenHashesST <- liftEffect $ liftST $ STRef.new HashSet.empty
|
||||
forever do
|
||||
a <- await
|
||||
seenHashes <- liftEffect $ liftST $ STRef.read seenHashesST
|
||||
when (not $ HashSet.member (hash a) seenHashes) do
|
||||
void $ liftEffect $ liftST $ STRef.modify (HashSet.insert $ hash a) seenHashesST
|
||||
yield a
|
||||
|
||||
-- | The result of a single step forward of a pipe.
|
||||
data InvokeResult a b m
|
||||
-- | The pipe `await`ed the value, but did not `yield` a response.
|
||||
= DidNotYield (Pipe a b m Unit)
|
||||
-- | The pipe `await`ed the value, and `yield`ed 1 or more responses.
|
||||
| Yielded (NonEmptyList b /\ Pipe a b m Unit)
|
||||
-- | The pipe `await`ed the value, and exited.
|
||||
| Exited
|
||||
|
||||
data IntermediateInvokeResult a b m
|
||||
= IDidNotYield (Pipe a b m Unit)
|
||||
| IYielded (NonEmptyList b /\ Pipe a b m Unit)
|
||||
| IDidNotAwait (Pipe a b m Unit)
|
||||
|
||||
-- | Pass a single value to a pipe, returning the result of the pipe's invocation.
|
||||
invoke :: forall m a b. Monad m => Pipe a b m Unit -> a -> m (InvokeResult a b m)
|
||||
invoke m a =
|
||||
let
|
||||
go :: IntermediateInvokeResult a b m -> m (InvokeResult a b m)
|
||||
go (IYielded (as /\ n)) =
|
||||
case n of
|
||||
Request _ _ -> pure $ Yielded $ as /\ n
|
||||
Respond rep f -> go (IYielded $ (as <> pure rep) /\ f unit)
|
||||
M o -> go =<< IYielded <$> (as /\ _) <$> o
|
||||
Pure _ -> pure Exited
|
||||
go (IDidNotYield n) =
|
||||
case n of
|
||||
Request _ _ -> pure $ DidNotYield n
|
||||
Respond rep f -> go (IYielded $ pure rep /\ f unit)
|
||||
M o -> go =<< IDidNotYield <$> o
|
||||
Pure _ -> pure Exited
|
||||
go (IDidNotAwait n) =
|
||||
case n of
|
||||
Request _ f -> go (IDidNotYield (f a))
|
||||
Respond rep f -> go (IYielded $ pure rep /\ f unit)
|
||||
M o -> go =<< IDidNotAwait <$> o
|
||||
Pure _ -> pure Exited
|
||||
in
|
||||
go (IDidNotAwait m)
|
||||
|
||||
Reference in New Issue
Block a user