Expand the documentation; Clean up debugging code; Enable all tests

This commit is contained in:
Zahary Karadjov 2020-05-05 19:20:25 +03:00
commit 18f9488bee
No known key found for this signature in database
GPG key ID: C8936F8A3073D609
11 changed files with 205 additions and 120 deletions

141
README.md
View file

@ -33,9 +33,9 @@ in a way that allows the read and write operations to be handled without any
dynamic dispatch in the majority of cases. dynamic dispatch in the majority of cases.
In particular, reading from a `memoryInput` or writing to a `memoryOutput` In particular, reading from a `memoryInput` or writing to a `memoryOutput`
will have the equivalent performance to a loop iterating over an `openarray` will have similar performance to a loop iterating over an `openarray` or
or another loop populating a pre-allocated `string`. `memFileInput` offers another loop populating a pre-allocated `string`. `memFileInput` offers
similar performance characteristics when working with files. The idiomatic the same performance characteristics when working with files. The idiomatic
use of the APIs with the rest of the stream types will result in a highly use of the APIs with the rest of the stream types will result in a highly
efficient memory allocation patterns and zero-copy performance in a great efficient memory allocation patterns and zero-copy performance in a great
variety of real-world use cases such as: variety of real-world use cases such as:
@ -90,7 +90,7 @@ efficient and easy to author.
### Higher efficiency is possible if we say goodbye to the good old single buffer. ### Higher efficiency is possible if we say goodbye to the good old single buffer.
The buffering logic inside the stream divides the data into "pages" which The buffering logic inside the stream divides the data into "pages" which
are allocated with known fast paths in the Nim allocator and which can be are allocated with a known fast path in the Nim allocator and which can be
efficiently transferred between streams and threads in the layered streams efficiently transferred between streams and threads in the layered streams
scenario or in IPC mechanisms such as `AsyncChannel`. The consuming code can scenario or in IPC mechanisms such as `AsyncChannel`. The consuming code can
be aware of this, but doesn't need to. The most idiomatic usage of the API be aware of this, but doesn't need to. The most idiomatic usage of the API
@ -109,19 +109,19 @@ such as:
* Block compressors and Block ciphers * Block compressors and Block ciphers
These can benefit significantly from a more precise control of the size These can benefit significantly from a more precise control over
of the buffered pages which can be configured to match the block size the stride of the buffered pages which can be configured to match
of the encoder. the block size of the encoder.
* Content with known length * Content with known length
Some streams have known length which allows us to accurately estimate Some streams have a known length which allows us to accurately estimate
the size of the transformed content. The `len` and `ensureRunway` APIs the size of the transformed content. The `len` and `ensureRunway` APIs
make sure such cases are handled as optimally as possible. make sure such cases are handled as optimally as possible.
## Basic API usage ## Basic API usage
The FastStreams API consists of 3 major object types: The FastStreams API consists of ony few major object types:
### `InputStream` ### `InputStream`
@ -142,6 +142,16 @@ of the box the following input stream types:
You are responsible for ensuring that the backing buffer won't be invalidated You are responsible for ensuring that the backing buffer won't be invalidated
while the stream is being used. while the stream is being used.
* `memoryInput`
Primarily used to consume the contents written to a previously populated
output stream, but it can also be used to consume the contents of strings
and sequences in a memory-safe way (by creating a copy).
* `pipeInput` (async)
For arbitrary conmmunication between a produced and a consumer.
* `chronosInput` (async) * `chronosInput` (async)
Enabled by importing `faststreams/chronos_adapters`. <br /> Enabled by importing `faststreams/chronos_adapters`. <br />
@ -173,7 +183,7 @@ The example above assumes we might have a `parseJson` function accepting an
`InputStream`. Here how this function could be defined: `InputStream`. Here how this function could be defined:
```nim ```nim
proc scanString(stream: InputStream): JsonToken = proc scanString(stream: InputStream): JsonToken {.fsMultiSync.} =
result = newStringToken() result = newStringToken()
advance stream # skip the opening quote advance stream # skip the opening quote
@ -197,7 +207,7 @@ proc scanString(stream: InputStream): JsonToken =
error(UnexpectedEndOfFile) error(UnexpectedEndOfFile)
proc nextToken(stream: InputStream): JsonToken = proc nextToken(stream: InputStream): JsonToken {.fsMultiSync.} =
while stream.readable: while stream.readable:
case stream.peek.char case stream.peek.char
of '"': of '"':
@ -213,7 +223,7 @@ proc nextToken(stream: InputStream): JsonToken =
return eofToken return eofToken
proc parseJson(stream: InputStream): JsonNode = proc parseJson(stream: InputStream): JsonNode {.fsMultiSync.} =
while (let token = nextToken(stream); token != eofToken): while (let token = nextToken(stream); token != eofToken):
case token case token
of numberToken: of numberToken:
@ -243,7 +253,7 @@ compile to very efficient inlined code that performs nothing more than pointer
increments and comparisons. This will be true even when working with async increments and comparisons. This will be true even when working with async
streams. streams.
The `readable` check is the only place where our code could block (or await). The `readable` check is the only place where our code may block (or await).
Only when all the data in the stream buffers have been consumed, the stream Only when all the data in the stream buffers have been consumed, the stream
will invoke a new read operation on the backing input device and this may will invoke a new read operation on the backing input device and this may
repopulate the buffers with an arbitrary number of new bytes. repopulate the buffers with an arbitrary number of new bytes.
@ -256,10 +266,52 @@ if you need to store the bytes in an object field or another long-term storage
location, consider using `stream.readInto(destination)` which may result in location, consider using `stream.readInto(destination)` which may result in
zero-copy operation. It can also be used to implement unbuffered reading. zero-copy operation. It can also be used to implement unbuffered reading.
In async streams, the `stream.timeoutToNextByte(t)` API can be used to detect #### `AsyncInputStream` and `fsMultiSync`
situations where your communicating party is failing to send data in time.
### `OutputStream` An astute reader might have wondered what is the purpose of the custom pragma
`fsMultiSync` used in the examples above? It is a simple macro generating an
additional `async` copy of our stream processing functions where all the input
types are replaced by their async counterparts (e.g. `AsyncInputStream`) and
the return type is wrapped in a `Future` as usual.
The standard API of `InputStream` and `AsyncInputStream` is exactly the same.
Operations such as `readable` will just invoke `await` behind the scenes, but
there is one key difference - the `await` will be triggered only when there
is not enough data already stored in the stream buffers. Thus, in the great
majority of cases, we avoid the high cost of instantiating a `Future` and
yielding control to the event loop.
We highly recommend implementing most of your stream processing code through
the `fsMultiSync` pragma. This ensures the best possible performance and makes
the code more easily testable (e.g. with inputs stored on disk). FastStreams
ships with a set of fuzzing tools that will help you ensure that your code
behaves correctly with arbitrary data and/or arbitrary interruption points.
Nevertheless, if you need a more traditional async API, please be aware that
all of the functions discussed in this README also have an `*Async` suffix
form that returns a `Future` (e.g. `readableAsync`, `readAsync`, etc).
One exception to the above rule is the helper `stream.timeoutToNextByte(t)`
which can be used to detect situations where your communicating party is
failing to send data in time. It accepts a `Duration` or an existing deadline
`Future` and it's usually used like this:
```nim
proc performHandshake(c: Connection): bool {.async.} =
if c.inputStream.timeoutToNextByte(HANDSHAKE_TIMEOUT):
# The other party didn't send us anything in time,
# We close the connection:
close c
return false
while c.inputStream.readable:
...
```
It is assumed that in traditional async code, timeouts will be managed more
explicitly with `sleepAsync` and the `or` operator defined over futures.
### `OutputStream` and `AsyncOutputStream`
An `OutputStream` manages a particular output device. The library offers out An `OutputStream` manages a particular output device. The library offers out
of the box the following output stream types: of the box the following output stream types:
@ -278,6 +330,10 @@ of the box the following output stream types:
You are responsible for ensuring that the backing buffer won't be invalidated You are responsible for ensuring that the backing buffer won't be invalidated
while the stream is being used. while the stream is being used.
* `pipeOutput` (async)
For arbitrary conmmunication between a produced and a consumer.
* `chronosOutput` (async) * `chronosOutput` (async)
Enabled by importing `faststreams/chronos_adapters`. <br /> Enabled by importing `faststreams/chronos_adapters`. <br />
@ -358,8 +414,18 @@ single page of `pageSize` bytes (specified at stream creation). Calls to
`write` will just populate this page until it becomes full and only then `write` will just populate this page until it becomes full and only then
it would be sent to the output device. it would be sent to the output device.
Writes larger than a page will be sent to the output device immediately, As the example demonstrates, a `memoryOutput` will continue buffering
so setting the `pageSize` to zero enables unbuffered mode of operation. pages until they can be finally concatenated and returned in `stream.getOutput`.
If the output fits within a single page, it will be efficiently moved to
the `getOutput` result. When the output size is known upfront you can ensure
that this optimization is used by calling `stream.ensureRunway` before any
writes, but please note that the library is free to ignore this hint in async
context or if a maximum memory usage policy is specified.
In a non-memory stream, any writes larger than a page or issued through the
`writeNow` API will be sent to the output device immediately.
#### Delayed Writes
Please note that even in async context, `write` will complete immediately. Please note that even in async context, `write` will complete immediately.
To handle back-pressure properly, use `stream.flush` or `stream.waitForConsumer` To handle back-pressure properly, use `stream.flush` or `stream.waitForConsumer`
@ -368,19 +434,23 @@ bytes before continuing. The rationale here is that introducing an interruption
point at every `write` produces less optimal code, but if this is desired you point at every `write` produces less optimal code, but if this is desired you
can use the `stream.writeAndWait` API. can use the `stream.writeAndWait` API.
Fixed-size and variable-size length prefixes can be handled without Many protocols and formats employ fixed-size and variable-size length prefixes
additional memory allocations through the `stream.delayFixedSizeWrite` that have been tradionally difficult to handle because they require you to
and `stream.delayVarSizeWrite` APIs which return a `WriteCursor` object either measure the size of the content before writing it to the stream, or
that must be `finalized` after the length-prefix is written. You can do even worse, serialize it to a memory buffer in order to determine its size.
this in one step with `cursor.finalWrite`.
As the example demonstrates, a `memoryOutput` will continue buffering FastStreams supports handling such length prefixes with a zero-copy mechanism
pages until they can be finally concatenated and returned in `stream.getOutput`. that doesn't require additional memory allocations. `stream.delayFixedSizeWrite`
If the output fits within a single page, it will be efficiently moved to and `stream.delayVarSizeWrite` are APIs that return a `WriteCursor` object that
the `getOutput` result. When the output size is known upfront you can ensure can be used to implement a delayed write to the stream. After obtaining the
that this optimization is used by calling `stream.ensureRunway` before any write cursor you can take a note of the current `pos` in the stream and then
writes, but please note that the library is free to ignore this hint in async continue issuing `stream.write` operations normally. After all of the content
context if a maximum memory usage policy is specified. is written, you obtain `pos` again to determine the final value of the length
prefix. Throughout the whole time, you are free to call `write` on the cursor
to populate the "hole" left in the stream with bytes, but at the end you must
call `finalize` to unlock the stream for flushing. You can also perform the
finalization in one step with `finalWrite` (the one-step approach is manatory
for variable-size prefixes).
### `Pipeline` ### `Pipeline`
@ -388,7 +458,7 @@ context if a maximum memory usage policy is specified.
A `Pipeline` represents a chain of transformations that should be applied to a A `Pipeline` represents a chain of transformations that should be applied to a
stream. It starts with an `InputStream` followed by one or more transformation stream. It starts with an `InputStream` followed by one or more transformation
steps and ending in a `OutputStream`. steps and ending with a result.
Each transformation step is a function of the kind: Each transformation step is a function of the kind:
@ -397,6 +467,15 @@ type PipelineStep* = proc (i: InputStream, o: OutputStream)
{.gcsafe, raises: [Defect, CatchableError].} {.gcsafe, raises: [Defect, CatchableError].}
``` ```
A result obtaining operation is a function of the kind:
```nim
type PipelineResultProc*[T] = proc (i: InputStream): T
{.gcsafe, raises: [Defect, CatchableError].}
```
Please note that `stream.getOutput` is an example of such a function.
Pipelnes can be created with the `cretePipeline` API or executed in place with Pipelnes can be created with the `cretePipeline` API or executed in place with
`executePipeline`. If the first input source is async, then the whole pipeline `executePipeline`. If the first input source is async, then the whole pipeline
with be executing asynchronously which can result in a much lower memory usage. with be executing asynchronously which can result in a much lower memory usage.

View file

@ -36,6 +36,13 @@ elif faststreams_async_backend in ["std", "asyncdispatch"]:
else: else:
{.fatal: "Unrecognized network backend: " & faststreams_async_backend.} {.fatal: "Unrecognized network backend: " & faststreams_async_backend.}
when defined(danger):
template fsAssert*(x) = discard
template fsAssert*(x, msg) = discard
else:
template fsAssert*(x) = doAssert(x)
template fsAssert*(x, msg) = doAssert(x, msg)
template fsTranslateErrors*(errMsg: string, body: untyped) = template fsTranslateErrors*(errMsg: string, body: untyped) =
try: try:
body body

View file

@ -142,16 +142,16 @@ func nextReadableSpan*(buffers: PageBuffers, span: var PageSpan) =
pageReadableEnd = firstPage.readableEnd pageReadableEnd = firstPage.readableEnd
if span.endAddr == nil: if span.endAddr == nil:
doAssert buffers.queue.len > 0 fsAssert buffers.queue.len > 0
span = obtainReadableSpan buffers.queue[0] span = obtainReadableSpan buffers.queue[0]
elif span.endAddr != pageReadableEnd: elif span.endAddr != pageReadableEnd:
# Check whether the span points within the current page: # Check whether the span points within the current page:
doAssert distance(firstPage.allocationStart, span.endAddr) >= 0 and fsAssert distance(firstPage.allocationStart, span.endAddr) >= 0 and
distance(span.endAddr, pageReadableEnd) >= 0 distance(span.endAddr, pageReadableEnd) >= 0
span.endAddr = pageReadableEnd span.endAddr = pageReadableEnd
firstPage.consumedTo = firstPage.writtenTo firstPage.consumedTo = firstPage.writtenTo
else: else:
doAssert buffers.queue.len > 1 fsAssert buffers.queue.len > 1
discard buffers.queue.popFirst discard buffers.queue.popFirst
span = obtainReadableSpan buffers.queue[0] span = obtainReadableSpan buffers.queue[0]
@ -209,7 +209,7 @@ func ensureRunway*(buffers: PageBuffers,
# This is a more complicated path that should almost never # This is a more complicated path that should almost never
# trigger in practice in a typically implemented code that # trigger in practice in a typically implemented code that
# calls `ensureRunway` at the beggining of a transformation. # calls `ensureRunway` at the beggining of a transformation.
doAssert buffers.queue.len > 0 fsAssert buffers.queue.len > 0
let currPage = buffers.queue.peekLast let currPage = buffers.queue.peekLast
if currPage.hasDelayedWritesAtPageStart: if currPage.hasDelayedWritesAtPageStart:
@ -266,7 +266,7 @@ func splitLastPageAt*(buffers: PageBuffers, address: ptr byte) =
buffers.queue.addLast newPage buffers.queue.addLast newPage
iterator consumePages*(buffers: PageBuffers): PageRef = iterator consumePages*(buffers: PageBuffers): PageRef =
doAssert buffers != nil fsAssert buffers != nil
var recycledPage: PageRef var recycledPage: PageRef
while buffers.queue.len > 0: while buffers.queue.len > 0:
@ -337,7 +337,7 @@ template implementWrites*(buffersParam: PageBuffers,
if bytesWritten != writeLenVar: raiseError() if bytesWritten != writeLenVar: raiseError()
if srcLen > 0: if srcLen > 0:
doAssert src != nil fsAssert src != nil
let bytesWritten = writeBlock let bytesWritten = writeBlock
if bytesWritten != writeLenVar: raiseError() if bytesWritten != writeLenVar: raiseError()

View file

@ -45,7 +45,7 @@ let chronosInputVTable = InputStreamVTable(
readSync: proc (s: InputStream, dst: pointer, dstLen: Natural): Natural readSync: proc (s: InputStream, dst: pointer, dstLen: Natural): Natural
{.nimcall, gcsafe, raises: [IOError, Defect].} = {.nimcall, gcsafe, raises: [IOError, Defect].} =
var cs = ChronosInputStream(s) var cs = ChronosInputStream(s)
doAssert cs.allowWaitFor fsAssert cs.allowWaitFor
waitFor chronosReadOnce(cs, dst, dstLen) waitFor chronosReadOnce(cs, dst, dstLen)
, ,
readAsync: proc (s: InputStream, dst: pointer, dstLen: Natural): Future[Natural] readAsync: proc (s: InputStream, dst: pointer, dstLen: Natural): Future[Natural]
@ -74,7 +74,7 @@ let chronosOutputVTable = OutputStreamVTable(
writeSync: proc (s: OutputStream, src: pointer, srcLen: Natural) writeSync: proc (s: OutputStream, src: pointer, srcLen: Natural)
{.nimcall, gcsafe, raises: [IOError, Defect].} = {.nimcall, gcsafe, raises: [IOError, Defect].} =
var cs = ChronosOutputStream(s) var cs = ChronosOutputStream(s)
doAssert cs.allowWaitFor fsAssert cs.allowWaitFor
waitFor chronosWrites(cs, src, srcLen) waitFor chronosWrites(cs, src, srcLen)
, ,
writeAsync: proc (s: OutputStream, src: pointer, srcLen: Natural): Future[void] writeAsync: proc (s: OutputStream, src: pointer, srcLen: Natural): Future[void]

View file

@ -216,7 +216,7 @@ template readableNow*(s: AsyncInputStream): bool =
readableNow InputStream(s) readableNow InputStream(s)
func flipPage(s: InputStream) = func flipPage(s: InputStream) =
doAssert s.buffers.len > 1 fsAssert s.buffers != nil and s.buffers.len > 1
discard s.buffers.popFirst discard s.buffers.popFirst
s.span = obtainReadableSpan s.buffers[0] s.span = obtainReadableSpan s.buffers[0]
s.spanEndPos += s.span.len s.spanEndPos += s.span.len
@ -348,7 +348,7 @@ func memoryInput*(data: openarray[char]): InputStreamHandle =
proc resetBuffers*(s: InputStream, buffers: PageBuffers) = proc resetBuffers*(s: InputStream, buffers: PageBuffers) =
# This should be used only on safe memory input streams # This should be used only on safe memory input streams
doAssert s.vtable == nil and s.buffers != nil and buffers.len > 0 fsAssert s.vtable == nil and s.buffers != nil and buffers.len > 0
s.buffers = buffers s.buffers = buffers
s.span = obtainReadableSpan buffers.queue[0] s.span = obtainReadableSpan buffers.queue[0]
s.spanEndPos = s.span.len s.spanEndPos = s.span.len
@ -530,17 +530,43 @@ template readable*(sp: AsyncInputStream, np: int): bool =
readableNImpl(s, n, fsAwait, readAsync) readableNImpl(s, n, fsAwait, readAsync)
proc peek*(s: InputStream): byte {.inline.} = when false:
doAssert hasRunway(s.span) func flipPagePeek(s: InputStream): byte =
return s.span.startAddr[] flipPage s
result = s.span.startAddr[]
func flipPageRead(s: InputStream): byte =
flipPage s
result = s.span.startAddr[]
bumpPointer s.span
template peek*(sp: InputStream): byte =
let s = sp
if hasRunway(s.span):
s.span.startAddr[]
else:
flipPage s
s.span.startAddr[]
template peek*(s: AsyncInputStream): byte = template peek*(s: AsyncInputStream): byte =
peek InputStream(s) peek InputStream(s)
template read*(sp: InputStream): byte =
let s = sp
if hasRunway(s.span):
let res = s.span.startAddr[]
bumpPointer(s.span)
res
else:
flipPageRead s
template read*(s: AsyncInputStream): byte =
read InputStream(s)
proc peekAt*(s: InputStream, pos: int): byte {.inline.} = proc peekAt*(s: InputStream, pos: int): byte {.inline.} =
# TODO implement page flipping # TODO implement page flipping
let peekHead = offset(s.span.startAddr, pos) let peekHead = offset(s.span.startAddr, pos)
doAssert cast[uint](peekHead) < cast[uint](s.span.endAddr) fsAssert cast[uint](peekHead) < cast[uint](s.span.endAddr)
return peekHead[] return peekHead[]
template peekAt*(s: AsyncInputStream, pos: int): byte = template peekAt*(s: AsyncInputStream, pos: int): byte =
@ -549,19 +575,12 @@ template peekAt*(s: AsyncInputStream, pos: int): byte =
proc advance*(s: InputStream) = proc advance*(s: InputStream) =
if hasRunway(s.span): if hasRunway(s.span):
bumpPointer s.span bumpPointer s.span
elif s.buffers != nil and s.buffers.len > 1: else:
flipPage s flipPage s
template advance*(s: AsyncInputStream) = template advance*(s: AsyncInputStream) =
advance InputStream(s) advance InputStream(s)
proc read*(s: InputStream): byte =
result = s.peek()
advance s
template read*(s: AsyncInputStream): byte =
read InputStream(s)
proc drainBuffersInto*(s: InputStream, dstAddr: ptr byte, dstLen: Natural): Natural = proc drainBuffersInto*(s: InputStream, dstAddr: ptr byte, dstLen: Natural): Natural =
var var
dst = dstAddr dst = dstAddr
@ -691,7 +710,7 @@ template readInto*(sp: AsyncInputStream, dst: var openarray[byte]): bool =
proc readOnce*(sp: AsyncInputStream): Future[Natural] = proc readOnce*(sp: AsyncInputStream): Future[Natural] =
let s = InputStream(sp) let s = InputStream(sp)
doAssert s.buffers != nil and s.vtable != nil fsAssert s.buffers != nil and s.vtable != nil
s.vtable.readAsync(s, nil, 0) s.vtable.readAsync(s, nil, 0)
when defined(windows): when defined(windows):
@ -724,7 +743,8 @@ template readNImpl(sp: InputStream,
if n > runway: if n > runway:
startAddr = allocMem(tmpSeq, n, np) startAddr = allocMem(tmpSeq, n, np)
doAssert drainBuffersInto(s, startAddr, n) == n let drained {.used.} = drainBuffersInto(s, startAddr, n)
fsAssert drained == n
else: else:
startAddr = s.span.startAddr startAddr = s.span.startAddr
bumpPointer s.span, n bumpPointer s.span, n
@ -775,16 +795,16 @@ when false:
# Obsolete APIs for removal # Obsolete APIs for removal
proc bufferPos(s: InputStream, pos: int): ptr byte = proc bufferPos(s: InputStream, pos: int): ptr byte =
let offsetFromEnd = pos - s.spanEndPos let offsetFromEnd = pos - s.spanEndPos
doAssert offsetFromEnd < 0 fsAssert offsetFromEnd < 0
result = offset(s.span.endAddr, offsetFromEnd) result = offset(s.span.endAddr, offsetFromEnd)
doAssert result >= s.bufferStart fsAssert result >= s.bufferStart
proc `[]`*(s: InputStream, pos: int): byte {.inline.} = proc `[]`*(s: InputStream, pos: int): byte {.inline.} =
s.bufferPos(pos)[] s.bufferPos(pos)[]
proc rewind*(s: InputStream, delta: int) = proc rewind*(s: InputStream, delta: int) =
s.head = offset(s.head, -delta) s.head = offset(s.head, -delta)
doAssert s.head >= s.bufferStart fsAssert s.head >= s.bufferStart
proc rewindTo*(s: InputStream, pos: int) {.inline.} = proc rewindTo*(s: InputStream, pos: int) {.inline.} =
s.head = s.bufferPos(pos) s.head = s.bufferPos(pos)

View file

@ -87,7 +87,7 @@ template disconnectOutputDevice(s: AsyncOutputStream) =
disconnectOutputDevice OutputStream(s) disconnectOutputDevice OutputStream(s)
template flushImpl(s: OutputStream, awaiter, writeOp, flushOp: untyped) = template flushImpl(s: OutputStream, awaiter, writeOp, flushOp: untyped) =
doAssert s.extCursorsCount == 0 fsAssert s.extCursorsCount == 0
if s.vtable != nil: if s.vtable != nil:
if s.buffers != nil: if s.buffers != nil:
trackWrittenTo(s.buffers, s.span.startAddr) trackWrittenTo(s.buffers, s.span.startAddr)
@ -159,10 +159,6 @@ template canExtendOutput(s: OutputStream): bool =
# Streams writing to pre-allocated existing buffers cannot be grown # Streams writing to pre-allocated existing buffers cannot be grown
s != nil and s.buffers != nil s != nil and s.buffers != nil
template isExternalCursor(c: var WriteCursor): bool =
# Is this the original stream cursor or is it one created by a "delayed write"
addr(c) != addr(c.stream.cursor)
proc addPage(s: OutputStream) = proc addPage(s: OutputStream) =
let let
nextPageSize = s.buffers.pageSize nextPageSize = s.buffers.pageSize
@ -175,7 +171,7 @@ template makeHandle*(sp: OutputStream): OutputStreamHandle =
OutputStreamHandle(s: s) OutputStreamHandle(s: s)
proc memoryOutput*(pageSize = defaultPageSize): OutputStreamHandle = proc memoryOutput*(pageSize = defaultPageSize): OutputStreamHandle =
doAssert pageSize > 0 fsAssert pageSize > 0
# We are not creating an initial output page, because `ensureRunway` # We are not creating an initial output page, because `ensureRunway`
# can determine the most appropriate size. # can determine the most appropriate size.
makeHandle OutputStream(buffers: initPageBuffers(pageSize)) makeHandle OutputStream(buffers: initPageBuffers(pageSize))
@ -196,7 +192,7 @@ proc ensureRunway*(s: OutputStream, neededRunway: Natural) =
# If you use an unsafe memory output, you must ensure that # If you use an unsafe memory output, you must ensure that
# it will have a large enough size to hold the data you are # it will have a large enough size to hold the data you are
# feeding to it. # feeding to it.
doAssert s.buffers != nil, "Unsafe memory output of insufficient size" fsAssert s.buffers != nil, "Unsafe memory output of insufficient size"
s.buffers.ensureRunway(s.span, neededRunway) s.buffers.ensureRunway(s.span, neededRunway)
s.spanEndPos += (s.span.len - runway) s.spanEndPos += (s.span.len - runway)
@ -246,7 +242,7 @@ template pos*(s: AsyncOutputStream): int =
pos OutputStream(s) pos OutputStream(s)
proc getBuffers*(s: OutputStream): PageBuffers = proc getBuffers*(s: OutputStream): PageBuffers =
doAssert s.buffers != nil fsAssert s.buffers != nil
s.buffers.trackWrittenTo s.span.startAddr s.buffers.trackWrittenTo s.span.startAddr
return s.buffers return s.buffers
@ -334,7 +330,7 @@ proc delayFixedSizeWrite*(s: OutputStream, size: Natural): WriteCursor =
proc delayVarSizeWrite*(s: OutputStream, maxSize: Natural): VarSizeWriteCursor = proc delayVarSizeWrite*(s: OutputStream, maxSize: Natural): VarSizeWriteCursor =
## Please note that using variable sized writes are not supported ## Please note that using variable sized writes are not supported
## for unbuffered streams and unsafe memory inputs. ## for unbuffered streams and unsafe memory inputs.
doAssert s.buffers != nil fsAssert s.buffers != nil
let runway = s.span.len let runway = s.span.len
if maxSize <= runway: if maxSize <= runway:
@ -368,11 +364,11 @@ proc delayVarSizeWrite*(s: OutputStream, maxSize: Natural): VarSizeWriteCursor =
s.spanEndPos += nextPageSize s.spanEndPos += nextPageSize
proc finalize*(cursor: var WriteCursor) = proc finalize*(cursor: var WriteCursor) =
doAssert cursor.stream.extCursorsCount > 0 fsAssert cursor.stream.extCursorsCount > 0
dec cursor.stream.extCursorsCount dec cursor.stream.extCursorsCount
proc finalWrite*(cursor: var WriteCursor, data: openArray[byte]) = proc finalWrite*(cursor: var WriteCursor, data: openArray[byte]) =
doAssert data.len == cursor.span.len fsAssert data.len == cursor.span.len
copyMem(cursor.span.startAddr, unsafeAddr data[0], data.len) copyMem(cursor.span.startAddr, unsafeAddr data[0], data.len)
finalize cursor finalize cursor
@ -380,7 +376,7 @@ proc finalWrite*(c: var VarSizeWriteCursor, data: openArray[byte]) =
template cursor: auto = WriteCursor(c) template cursor: auto = WriteCursor(c)
let overestimatedBytes = cursor.span.len - data.len let overestimatedBytes = cursor.span.len - data.len
doAssert overestimatedBytes >= 0 fsAssert overestimatedBytes >= 0
for page in items(cursor.stream.buffers.queue): for page in items(cursor.stream.buffers.queue):
let baseAddr = page.allocationStart let baseAddr = page.allocationStart
@ -398,7 +394,7 @@ proc finalWrite*(c: var VarSizeWriteCursor, data: openArray[byte]) =
finalize cursor finalize cursor
return return
doAssert false fsAssert false
proc tryMovingToNextPage(c: var WriteCursor) = proc tryMovingToNextPage(c: var WriteCursor) =
# A split cursor is a fixed-size cursor that ended up on page boundary. # A split cursor is a fixed-size cursor that ended up on page boundary.
@ -447,13 +443,13 @@ proc tryMovingToNextPage(c: var WriteCursor) =
# We didn't find any page that this cursor was ending, so this is not # We didn't find any page that this cursor was ending, so this is not
# a split cursor. This means that the user just tried to write past the # a split cursor. This means that the user just tried to write past the
# pre-allocated cursor span, which is considered a Defect (a range error) # pre-allocated cursor span, which is considered a Defect (a range error)
doAssert false, "Attempt to write past the end of a cursor" fsAssert false, "Attempt to write past the end of a cursor"
template writeByteImpl(s: OutputStream, b: byte, awaiter, writeOp, drainOp: untyped) = template writeByteImpl(s: OutputStream, b: byte, awaiter, writeOp, drainOp: untyped) =
if atEnd(s.span): if atEnd(s.span):
# Unsafe memory outputs don't use pages at all, so if our cursor # Unsafe memory outputs don't use pages at all, so if our cursor
# reached here, this is a range violation defect: # reached here, this is a range violation defect:
doAssert canExtendOutput(s) fsAssert canExtendOutput(s)
if s.vtable == nil or s.extCursorsCount > 0: if s.vtable == nil or s.extCursorsCount > 0:
# This is the main cursor of a stream, but we are either not # This is the main cursor of a stream, but we are either not
@ -512,7 +508,7 @@ proc writeToANewPage(s: OutputStream, bytes: openArray[byte]) =
copyMem(s.span.startAddr, inputPos, runway) copyMem(s.span.startAddr, inputPos, runway)
reduceInput runway reduceInput runway
doAssert s.buffers != nil fsAssert s.buffers != nil
let nextPageSize = nextAlignedSize(inputLen, s.buffers.pageSize) let nextPageSize = nextAlignedSize(inputLen, s.buffers.pageSize)
let nextPage = s.buffers.addWritablePage(nextPageSize) let nextPage = s.buffers.addWritablePage(nextPageSize)
@ -618,7 +614,7 @@ proc writeBytesToCursor(c: var WriteCursor, bytes: openarray[byte]) =
# On the next page, we have a new runway # On the next page, we have a new runway
runway = c.span.len runway = c.span.len
# The write shouldn't go past the end of the new runway # The write shouldn't go past the end of the new runway
doAssert inputLen <= runway fsAssert inputLen <= runway
copyMem(c.span.startAddr, inputPos, inputLen) copyMem(c.span.startAddr, inputPos, inputLen)
c.span.startAddr = offset(c.span.startAddr, inputLen) c.span.startAddr = offset(c.span.startAddr, inputLen)
@ -644,7 +640,7 @@ template consumeOutputs*(sp: OutputStream, bytesVar, body: untyped) =
## Before consuming the outputs, all outstanding delayed writes must ## Before consuming the outputs, all outstanding delayed writes must
## be finalized. ## be finalized.
let s = sp let s = sp
doAssert s.extCursorsCount == 0 and s.buffers != nil fsAssert s.extCursorsCount == 0 and s.buffers != nil
for pageReadableStart, pageLen in consumePageBuffers(s.buffers): for pageReadableStart, pageLen in consumePageBuffers(s.buffers):
template bytesVar: untyped = template bytesVar: untyped =
@ -673,7 +669,7 @@ template consumeContiguousOutput*(sp: OutputStream, bytesVar, body: untyped) =
bytesPtr: ptr byte bytesPtr: ptr byte
bytesLen: int bytesLen: int
doAssert s.extCursorsCount == 0 and s.buffers != nil fsAssert s.extCursorsCount == 0 and s.buffers != nil
if s.buffers.queue.len == 1: if s.buffers.queue.len == 1:
let page = s.buffers.queue[0] let page = s.buffers.queue[0]
@ -702,7 +698,7 @@ proc getOutput*(s: OutputStream, T: type string): string =
## ##
## Before consuming the output, all outstanding delayed writes must be finalized. ## Before consuming the output, all outstanding delayed writes must be finalized.
## ##
doAssert s.extCursorsCount == 0 and s.buffers != nil fsAssert s.extCursorsCount == 0 and s.buffers != nil
s.buffers.trackWrittenTo s.span.startAddr s.buffers.trackWrittenTo s.span.startAddr
if s.buffers.queue.len == 1: if s.buffers.queue.len == 1:

View file

@ -5,11 +5,6 @@ import
export export
inputs, outputs, async_backend inputs, outputs, async_backend
template clearAndWait(ep: AsyncEvent) =
let e = ep
clear e
await e.wait()
type type
FsAsyncPipe* = ref object FsAsyncPipe* = ref object
# TODO: Make these stream handles # TODO: Make these stream handles
@ -38,22 +33,17 @@ proc pipeRead(s: LayeredInputStream,
minBytesExpected = max(1, dstLen) minBytesExpected = max(1, dstLen)
bytesInBuffersNow = bytesInBuffersAtStart bytesInBuffersNow = bytesInBuffersAtStart
describeBuffers "at start", buffers
while bytesInBuffersNow < minBytesExpected: while bytesInBuffersNow < minBytesExpected:
awake buffers.waitingWriter awake buffers.waitingWriter
echo "About to wait for writer"
buffers.waitingReader.enterWait "waiting for writer to buffer more data" buffers.waitingReader.enterWait "waiting for writer to buffer more data"
echo "Awaken from wait"
bytesInBuffersNow = buffers.totalBufferedBytes bytesInBuffersNow = buffers.totalBufferedBytes
if buffers.eofReached: if buffers.eofReached:
echo "read bytes ", bytesInBuffersNow - bytesInBuffersAtStart
describeBuffers "at end", buffers
return bytesInBuffersNow - bytesInBuffersAtStart return bytesInBuffersNow - bytesInBuffersAtStart
if dst != nil: if dst != nil:
doAssert drainBuffersInto(s, cast[ptr byte](dst), dstLen) == dstLen let drained {.used.} = drainBuffersInto(s, cast[ptr byte](dst), dstLen)
fsAssert drained == dstLen
awake buffers.waitingWriter awake buffers.waitingWriter
@ -61,7 +51,6 @@ proc pipeRead(s: LayeredInputStream,
proc pipeWrite(s: LayeredOutputStream, src: pointer, srcLen: Natural) {.async.} = proc pipeWrite(s: LayeredOutputStream, src: pointer, srcLen: Natural) {.async.} =
let buffers = s.buffers let buffers = s.buffers
echo "pipe write"
while buffers.canAcceptWrite(srcLen) == false: while buffers.canAcceptWrite(srcLen) == false:
buffers.waitingWriter.enterWait "waiting for reader to drain the buffers" buffers.waitingWriter.enterWait "waiting for reader to drain the buffers"
@ -81,7 +70,7 @@ let pipeInputVTable = InputStreamVTable(
{.nimcall, gcsafe, raises: [IOError, Defect].} = {.nimcall, gcsafe, raises: [IOError, Defect].} =
fsTranslateErrors "Failed to read from pipe": fsTranslateErrors "Failed to read from pipe":
let ls = LayeredInputStream(s) let ls = LayeredInputStream(s)
doAssert ls.allowWaitFor fsAssert ls.allowWaitFor
return waitFor pipeRead(ls, dst, dstLen) return waitFor pipeRead(ls, dst, dstLen)
, ,
readAsync: proc (s: InputStream, dst: pointer, dstLen: Natural): Future[Natural] readAsync: proc (s: InputStream, dst: pointer, dstLen: Natural): Future[Natural]
@ -117,7 +106,7 @@ let pipeOutputVTable = OutputStreamVTable(
{.nimcall, gcsafe, raises: [IOError, Defect].} = {.nimcall, gcsafe, raises: [IOError, Defect].} =
fsTranslateErrors "Failed to write all bytes to pipe": fsTranslateErrors "Failed to write all bytes to pipe":
var ls = LayeredOutputStream(s) var ls = LayeredOutputStream(s)
doAssert ls.allowWaitFor fsAssert ls.allowWaitFor
waitFor pipeWrite(ls, src, srcLen) waitFor pipeWrite(ls, src, srcLen)
, ,
writeAsync: proc (s: OutputStream, src: pointer, srcLen: Natural): Future[void] writeAsync: proc (s: OutputStream, src: pointer, srcLen: Natural): Future[void]
@ -146,7 +135,6 @@ let pipeOutputVTable = OutputStreamVTable(
{.nimcall, gcsafe, raises: [IOError, Defect].} = {.nimcall, gcsafe, raises: [IOError, Defect].} =
s.buffers.eofReached = true s.buffers.eofReached = true
echo "writer closes the stream"
fsTranslateErrors "Unexpected error from Future.complete": fsTranslateErrors "Unexpected error from Future.complete":
awake s.buffers.waitingReader awake s.buffers.waitingReader
@ -173,7 +161,7 @@ let pipeOutputVTable = OutputStreamVTable(
func pipeInput*(source: InputStream, func pipeInput*(source: InputStream,
pageSize = defaultPageSize, pageSize = defaultPageSize,
allowWaitFor = false): AsyncInputStream = allowWaitFor = false): AsyncInputStream =
doAssert pageSize > 0 fsAssert pageSize > 0
AsyncInputStream LayeredInputStream( AsyncInputStream LayeredInputStream(
vtable: vtableAddr pipeInputVTable, vtable: vtableAddr pipeInputVTable,
@ -199,7 +187,7 @@ proc pipeOutput*(destination: OutputStream,
pageSize = defaultPageSize, pageSize = defaultPageSize,
maxBufferedBytes = defaultPageSize * 4, maxBufferedBytes = defaultPageSize * 4,
allowWaitFor = false): AsyncOutputStream = allowWaitFor = false): AsyncOutputStream =
doAssert pageSize > 0 fsAssert pageSize > 0
var var
buffers = initPageBuffers pageSize buffers = initPageBuffers pageSize
@ -231,7 +219,7 @@ proc pipeOutput*(buffers: PageBuffers,
func asyncPipe*(pageSize = defaultPageSize, func asyncPipe*(pageSize = defaultPageSize,
maxBufferedBytes = defaultPageSize * 4): FsAsyncPipe = maxBufferedBytes = defaultPageSize * 4): FsAsyncPipe =
doAssert pageSize > 0 fsAssert pageSize > 0
FsAsyncPipe(buffers: initPageBuffers(pageSize, maxBufferedBytes)) FsAsyncPipe(buffers: initPageBuffers(pageSize, maxBufferedBytes))
func initReader*(pipe: FsAsyncPipe): AsyncInputStream = func initReader*(pipe: FsAsyncPipe): AsyncInputStream =

View file

@ -1,6 +1,6 @@
import import
stew/ptrops, stew/ptrops,
inputs, outputs, buffers, multisync inputs, outputs, buffers, async_backend, multisync
template matchingIntType(T: type int64): type = uint64 template matchingIntType(T: type int64): type = uint64
template matchingIntType(T: type int32): type = uint32 template matchingIntType(T: type int32): type = uint32
@ -112,7 +112,7 @@ const
Digits* = {'0'..'9'} Digits* = {'0'..'9'}
proc readLine*(s: InputStream, keepEol = false): TaintedString = proc readLine*(s: InputStream, keepEol = false): TaintedString =
doAssert readableNow(s) fsAssert readableNow(s)
while s.readable: while s.readable:
let c = s.peek.char let c = s.peek.char
@ -131,7 +131,7 @@ proc readLine*(s: InputStream, keepEol = false): TaintedString =
proc readUntil*(s: InputStream, proc readUntil*(s: InputStream,
sep: openarray[char]): Option[TaintedString] = sep: openarray[char]): Option[TaintedString] =
doAssert readableNow(s) fsAssert readableNow(s)
var res = "" var res = ""
while s.readable(sep.len): while s.readable(sep.len):
if s.lookAheadMatch(charsToBytes(sep)): if s.lookAheadMatch(charsToBytes(sep)):
@ -150,7 +150,7 @@ iterator lines*(s: InputStream, keepEol = false): TaintedString =
yield readLine(s, keepEol) yield readLine(s, keepEol)
proc readUnsignedInt*(s: InputStream, T: type[CompiledUIntTypes]): T = proc readUnsignedInt*(s: InputStream, T: type[CompiledUIntTypes]): T =
doAssert s.readable and s.peek.char in Digits fsAssert s.readable and s.peek.char in Digits
template eatDigitAndPeek: char = template eatDigitAndPeek: char =
advance s advance s

View file

@ -13,7 +13,7 @@ proc bytes(s: string): seq[byte] =
proc str(bytes: openarray[byte]): string = proc str(bytes: openarray[byte]): string =
result = newStringOfCap(bytes.len) result = newStringOfCap(bytes.len)
for b in bytes: for b in items(bytes):
result.add b.char result.add b.char
proc countLines(s: InputStream): Natural = proc countLines(s: InputStream): Natural =
@ -31,13 +31,15 @@ procSuite "input stream":
test "input is not readable with read": test "input is not readable with read":
check not input.readable check not input.readable
expect Defect: when not defined(danger):
echo "This read should not complete: ", input.read expect Defect:
echo "This read should not complete: ", input.read
test "input is not readable with read(n)": test "input is not readable with read(n)":
check not input.readable(10) check not input.readable(10)
expect Defect: when not defined(danger):
echo "This read should not complete: ", input.read(10) expect Defect:
echo "This read should not complete: ", input.read(10)
test "next returns none": test "next returns none":
check input.next.isNone check input.next.isNone

View file

@ -26,7 +26,6 @@ proc upcaseAllCharacters(i: InputStream, o: OutputStream) {.fsMultiSync.} =
while i.readable: while i.readable:
o.write toUpperAscii(i.read.char) o.write toUpperAscii(i.read.char)
echo "closing upcase"
close o close o
proc printTimes(t: TestTimes) = proc printTimes(t: TestTimes) =
@ -51,7 +50,6 @@ procSuite "pipelines":
""" """
#[
test "upper-case/base64 pipeline benchmark": test "upper-case/base64 pipeline benchmark":
var var
times: TestTimes times: TestTimes
@ -61,9 +59,8 @@ procSuite "pipelines":
let inputText = loremIpsum.repeat(5000) let inputText = loremIpsum.repeat(5000)
when debugHelpers: timeIt times.stdFunctionCalls:
echo "Input len: ", inputText.len stdRes = base64.decode(base64.encode(toUpperAscii(inputText)))
echo "Base 64 len: ", base64.encode(inputText).len
timeIt times.fsPipeline: timeIt times.fsPipeline:
fsRes = executePipeline(unsafeMemoryInput(inputText), fsRes = executePipeline(unsafeMemoryInput(inputText),
@ -79,21 +76,14 @@ procSuite "pipelines":
base64decode, base64decode,
getOutput string) getOutput string)
timeIt times.stdFunctionCalls:
stdRes = base64.decode(base64.encode(toUpperAscii(inputText)))
check fsAsyncRes == stdRes check fsAsyncRes == stdRes
check fsRes == stdRes check fsRes == stdRes
printTimes times printTimes times
]#
asyncTest "upper-case/base64 async pipeline": asyncTest "upper-case/base64 async pipeline":
let pipe = asyncPipe() let pipe = asyncPipe()
let inputText = repeat(loremIpsum, 8) let inputText = repeat(loremIpsum, 100)
when debugHelpers:
echo "Input len: ", inputText.len
proc pipeFeeder(s: AsyncOutputStream) {.gcsafe, async.} = proc pipeFeeder(s: AsyncOutputStream) {.gcsafe, async.} =
randomize 1234 randomize 1234
@ -112,7 +102,6 @@ procSuite "pipelines":
let sleep = rand(50) - 45 let sleep = rand(50) - 45
if sleep > 0: if sleep > 0:
echo "written ", pos
await sleepAsync(sleep.milliseconds) await sleepAsync(sleep.milliseconds)
close s close s

View file

@ -0,0 +1,4 @@
--linedir:off
--linetrace:off
--stacktrace:off