diff --git a/README.md b/README.md
index 459e9b2..ffcb903 100644
--- a/README.md
+++ b/README.md
@@ -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.
In particular, reading from a `memoryInput` or writing to a `memoryOutput`
-will have the equivalent performance to a loop iterating over an `openarray`
-or another loop populating a pre-allocated `string`. `memFileInput` offers
-similar performance characteristics when working with files. The idiomatic
+will have similar performance to a loop iterating over an `openarray` or
+another loop populating a pre-allocated `string`. `memFileInput` offers
+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
efficient memory allocation patterns and zero-copy performance in a great
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.
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
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
@@ -109,19 +109,19 @@ such as:
* Block compressors and Block ciphers
- These can benefit significantly from a more precise control of the size
- of the buffered pages which can be configured to match the block size
- of the encoder.
+ These can benefit significantly from a more precise control over
+ the stride of the buffered pages which can be configured to match
+ the block size of the encoder.
* 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
make sure such cases are handled as optimally as possible.
## Basic API usage
-The FastStreams API consists of 3 major object types:
+The FastStreams API consists of ony few major object types:
### `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
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)
Enabled by importing `faststreams/chronos_adapters`.
@@ -173,7 +183,7 @@ The example above assumes we might have a `parseJson` function accepting an
`InputStream`. Here how this function could be defined:
```nim
-proc scanString(stream: InputStream): JsonToken =
+proc scanString(stream: InputStream): JsonToken {.fsMultiSync.} =
result = newStringToken()
advance stream # skip the opening quote
@@ -197,7 +207,7 @@ proc scanString(stream: InputStream): JsonToken =
error(UnexpectedEndOfFile)
-proc nextToken(stream: InputStream): JsonToken =
+proc nextToken(stream: InputStream): JsonToken {.fsMultiSync.} =
while stream.readable:
case stream.peek.char
of '"':
@@ -213,7 +223,7 @@ proc nextToken(stream: InputStream): JsonToken =
return eofToken
-proc parseJson(stream: InputStream): JsonNode =
+proc parseJson(stream: InputStream): JsonNode {.fsMultiSync.} =
while (let token = nextToken(stream); token != eofToken):
case token
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
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
will invoke a new read operation on the backing input device and this may
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
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
-situations where your communicating party is failing to send data in time.
+#### `AsyncInputStream` and `fsMultiSync`
-### `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
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
while the stream is being used.
+* `pipeOutput` (async)
+
+ For arbitrary conmmunication between a produced and a consumer.
+
* `chronosOutput` (async)
Enabled by importing `faststreams/chronos_adapters`.
@@ -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
it would be sent to the output device.
-Writes larger than a page will be sent to the output device immediately,
-so setting the `pageSize` to zero enables unbuffered mode of operation.
+As the example demonstrates, a `memoryOutput` will continue buffering
+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.
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
can use the `stream.writeAndWait` API.
-Fixed-size and variable-size length prefixes can be handled without
-additional memory allocations through the `stream.delayFixedSizeWrite`
-and `stream.delayVarSizeWrite` APIs which return a `WriteCursor` object
-that must be `finalized` after the length-prefix is written. You can do
-this in one step with `cursor.finalWrite`.
+Many protocols and formats employ fixed-size and variable-size length prefixes
+that have been tradionally difficult to handle because they require you to
+either measure the size of the content before writing it to the stream, or
+even worse, serialize it to a memory buffer in order to determine its size.
-As the example demonstrates, a `memoryOutput` will continue buffering
-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 if a maximum memory usage policy is specified.
+FastStreams supports handling such length prefixes with a zero-copy mechanism
+that doesn't require additional memory allocations. `stream.delayFixedSizeWrite`
+and `stream.delayVarSizeWrite` are APIs that return a `WriteCursor` object that
+can be used to implement a delayed write to the stream. After obtaining the
+write cursor you can take a note of the current `pos` in the stream and then
+continue issuing `stream.write` operations normally. After all of the content
+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`
@@ -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
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:
@@ -397,6 +467,15 @@ type PipelineStep* = proc (i: InputStream, o: OutputStream)
{.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
`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.
diff --git a/faststreams/async_backend.nim b/faststreams/async_backend.nim
index c542f80..784544c 100644
--- a/faststreams/async_backend.nim
+++ b/faststreams/async_backend.nim
@@ -36,6 +36,13 @@ elif faststreams_async_backend in ["std", "asyncdispatch"]:
else:
{.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) =
try:
body
diff --git a/faststreams/buffers.nim b/faststreams/buffers.nim
index 144455c..30b3e22 100644
--- a/faststreams/buffers.nim
+++ b/faststreams/buffers.nim
@@ -142,16 +142,16 @@ func nextReadableSpan*(buffers: PageBuffers, span: var PageSpan) =
pageReadableEnd = firstPage.readableEnd
if span.endAddr == nil:
- doAssert buffers.queue.len > 0
+ fsAssert buffers.queue.len > 0
span = obtainReadableSpan buffers.queue[0]
elif span.endAddr != pageReadableEnd:
# 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
span.endAddr = pageReadableEnd
firstPage.consumedTo = firstPage.writtenTo
else:
- doAssert buffers.queue.len > 1
+ fsAssert buffers.queue.len > 1
discard buffers.queue.popFirst
span = obtainReadableSpan buffers.queue[0]
@@ -209,7 +209,7 @@ func ensureRunway*(buffers: PageBuffers,
# This is a more complicated path that should almost never
# trigger in practice in a typically implemented code that
# calls `ensureRunway` at the beggining of a transformation.
- doAssert buffers.queue.len > 0
+ fsAssert buffers.queue.len > 0
let currPage = buffers.queue.peekLast
if currPage.hasDelayedWritesAtPageStart:
@@ -266,7 +266,7 @@ func splitLastPageAt*(buffers: PageBuffers, address: ptr byte) =
buffers.queue.addLast newPage
iterator consumePages*(buffers: PageBuffers): PageRef =
- doAssert buffers != nil
+ fsAssert buffers != nil
var recycledPage: PageRef
while buffers.queue.len > 0:
@@ -337,7 +337,7 @@ template implementWrites*(buffersParam: PageBuffers,
if bytesWritten != writeLenVar: raiseError()
if srcLen > 0:
- doAssert src != nil
+ fsAssert src != nil
let bytesWritten = writeBlock
if bytesWritten != writeLenVar: raiseError()
diff --git a/faststreams/chronos_adapters.nim b/faststreams/chronos_adapters.nim
index 4d94d47..4dcc503 100644
--- a/faststreams/chronos_adapters.nim
+++ b/faststreams/chronos_adapters.nim
@@ -45,7 +45,7 @@ let chronosInputVTable = InputStreamVTable(
readSync: proc (s: InputStream, dst: pointer, dstLen: Natural): Natural
{.nimcall, gcsafe, raises: [IOError, Defect].} =
var cs = ChronosInputStream(s)
- doAssert cs.allowWaitFor
+ fsAssert cs.allowWaitFor
waitFor chronosReadOnce(cs, dst, dstLen)
,
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)
{.nimcall, gcsafe, raises: [IOError, Defect].} =
var cs = ChronosOutputStream(s)
- doAssert cs.allowWaitFor
+ fsAssert cs.allowWaitFor
waitFor chronosWrites(cs, src, srcLen)
,
writeAsync: proc (s: OutputStream, src: pointer, srcLen: Natural): Future[void]
diff --git a/faststreams/inputs.nim b/faststreams/inputs.nim
index c69c924..6966c1a 100644
--- a/faststreams/inputs.nim
+++ b/faststreams/inputs.nim
@@ -216,7 +216,7 @@ template readableNow*(s: AsyncInputStream): bool =
readableNow InputStream(s)
func flipPage(s: InputStream) =
- doAssert s.buffers.len > 1
+ fsAssert s.buffers != nil and s.buffers.len > 1
discard s.buffers.popFirst
s.span = obtainReadableSpan s.buffers[0]
s.spanEndPos += s.span.len
@@ -348,7 +348,7 @@ func memoryInput*(data: openarray[char]): InputStreamHandle =
proc resetBuffers*(s: InputStream, buffers: PageBuffers) =
# 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.span = obtainReadableSpan buffers.queue[0]
s.spanEndPos = s.span.len
@@ -530,17 +530,43 @@ template readable*(sp: AsyncInputStream, np: int): bool =
readableNImpl(s, n, fsAwait, readAsync)
-proc peek*(s: InputStream): byte {.inline.} =
- doAssert hasRunway(s.span)
- return s.span.startAddr[]
+when false:
+ func flipPagePeek(s: InputStream): byte =
+ 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 =
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.} =
# TODO implement page flipping
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[]
template peekAt*(s: AsyncInputStream, pos: int): byte =
@@ -549,19 +575,12 @@ template peekAt*(s: AsyncInputStream, pos: int): byte =
proc advance*(s: InputStream) =
if hasRunway(s.span):
bumpPointer s.span
- elif s.buffers != nil and s.buffers.len > 1:
+ else:
flipPage s
template advance*(s: AsyncInputStream) =
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 =
var
dst = dstAddr
@@ -691,7 +710,7 @@ template readInto*(sp: AsyncInputStream, dst: var openarray[byte]): bool =
proc readOnce*(sp: AsyncInputStream): Future[Natural] =
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)
when defined(windows):
@@ -724,7 +743,8 @@ template readNImpl(sp: InputStream,
if n > runway:
startAddr = allocMem(tmpSeq, n, np)
- doAssert drainBuffersInto(s, startAddr, n) == n
+ let drained {.used.} = drainBuffersInto(s, startAddr, n)
+ fsAssert drained == n
else:
startAddr = s.span.startAddr
bumpPointer s.span, n
@@ -775,16 +795,16 @@ when false:
# Obsolete APIs for removal
proc bufferPos(s: InputStream, pos: int): ptr byte =
let offsetFromEnd = pos - s.spanEndPos
- doAssert offsetFromEnd < 0
+ fsAssert offsetFromEnd < 0
result = offset(s.span.endAddr, offsetFromEnd)
- doAssert result >= s.bufferStart
+ fsAssert result >= s.bufferStart
proc `[]`*(s: InputStream, pos: int): byte {.inline.} =
s.bufferPos(pos)[]
proc rewind*(s: InputStream, delta: int) =
s.head = offset(s.head, -delta)
- doAssert s.head >= s.bufferStart
+ fsAssert s.head >= s.bufferStart
proc rewindTo*(s: InputStream, pos: int) {.inline.} =
s.head = s.bufferPos(pos)
diff --git a/faststreams/outputs.nim b/faststreams/outputs.nim
index 27b7cf8..22d3ac9 100644
--- a/faststreams/outputs.nim
+++ b/faststreams/outputs.nim
@@ -87,7 +87,7 @@ template disconnectOutputDevice(s: AsyncOutputStream) =
disconnectOutputDevice OutputStream(s)
template flushImpl(s: OutputStream, awaiter, writeOp, flushOp: untyped) =
- doAssert s.extCursorsCount == 0
+ fsAssert s.extCursorsCount == 0
if s.vtable != nil:
if s.buffers != nil:
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
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) =
let
nextPageSize = s.buffers.pageSize
@@ -175,7 +171,7 @@ template makeHandle*(sp: OutputStream): OutputStreamHandle =
OutputStreamHandle(s: s)
proc memoryOutput*(pageSize = defaultPageSize): OutputStreamHandle =
- doAssert pageSize > 0
+ fsAssert pageSize > 0
# We are not creating an initial output page, because `ensureRunway`
# can determine the most appropriate size.
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
# it will have a large enough size to hold the data you are
# 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.spanEndPos += (s.span.len - runway)
@@ -246,7 +242,7 @@ template pos*(s: AsyncOutputStream): int =
pos OutputStream(s)
proc getBuffers*(s: OutputStream): PageBuffers =
- doAssert s.buffers != nil
+ fsAssert s.buffers != nil
s.buffers.trackWrittenTo s.span.startAddr
return s.buffers
@@ -334,7 +330,7 @@ proc delayFixedSizeWrite*(s: OutputStream, size: Natural): WriteCursor =
proc delayVarSizeWrite*(s: OutputStream, maxSize: Natural): VarSizeWriteCursor =
## Please note that using variable sized writes are not supported
## for unbuffered streams and unsafe memory inputs.
- doAssert s.buffers != nil
+ fsAssert s.buffers != nil
let runway = s.span.len
if maxSize <= runway:
@@ -368,11 +364,11 @@ proc delayVarSizeWrite*(s: OutputStream, maxSize: Natural): VarSizeWriteCursor =
s.spanEndPos += nextPageSize
proc finalize*(cursor: var WriteCursor) =
- doAssert cursor.stream.extCursorsCount > 0
+ fsAssert cursor.stream.extCursorsCount > 0
dec cursor.stream.extCursorsCount
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)
finalize cursor
@@ -380,7 +376,7 @@ proc finalWrite*(c: var VarSizeWriteCursor, data: openArray[byte]) =
template cursor: auto = WriteCursor(c)
let overestimatedBytes = cursor.span.len - data.len
- doAssert overestimatedBytes >= 0
+ fsAssert overestimatedBytes >= 0
for page in items(cursor.stream.buffers.queue):
let baseAddr = page.allocationStart
@@ -398,7 +394,7 @@ proc finalWrite*(c: var VarSizeWriteCursor, data: openArray[byte]) =
finalize cursor
return
- doAssert false
+ fsAssert false
proc tryMovingToNextPage(c: var WriteCursor) =
# 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
# 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)
- 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) =
if atEnd(s.span):
# Unsafe memory outputs don't use pages at all, so if our cursor
# reached here, this is a range violation defect:
- doAssert canExtendOutput(s)
+ fsAssert canExtendOutput(s)
if s.vtable == nil or s.extCursorsCount > 0:
# 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)
reduceInput runway
- doAssert s.buffers != nil
+ fsAssert s.buffers != nil
let nextPageSize = nextAlignedSize(inputLen, s.buffers.pageSize)
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
runway = c.span.len
# The write shouldn't go past the end of the new runway
- doAssert inputLen <= runway
+ fsAssert inputLen <= runway
copyMem(c.span.startAddr, inputPos, 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
## be finalized.
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):
template bytesVar: untyped =
@@ -673,7 +669,7 @@ template consumeContiguousOutput*(sp: OutputStream, bytesVar, body: untyped) =
bytesPtr: ptr byte
bytesLen: int
- doAssert s.extCursorsCount == 0 and s.buffers != nil
+ fsAssert s.extCursorsCount == 0 and s.buffers != nil
if s.buffers.queue.len == 1:
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.
##
- doAssert s.extCursorsCount == 0 and s.buffers != nil
+ fsAssert s.extCursorsCount == 0 and s.buffers != nil
s.buffers.trackWrittenTo s.span.startAddr
if s.buffers.queue.len == 1:
diff --git a/faststreams/pipelines.nim b/faststreams/pipelines.nim
index 19aeee3..734a5b7 100644
--- a/faststreams/pipelines.nim
+++ b/faststreams/pipelines.nim
@@ -5,11 +5,6 @@ import
export
inputs, outputs, async_backend
-template clearAndWait(ep: AsyncEvent) =
- let e = ep
- clear e
- await e.wait()
-
type
FsAsyncPipe* = ref object
# TODO: Make these stream handles
@@ -38,22 +33,17 @@ proc pipeRead(s: LayeredInputStream,
minBytesExpected = max(1, dstLen)
bytesInBuffersNow = bytesInBuffersAtStart
- describeBuffers "at start", buffers
-
while bytesInBuffersNow < minBytesExpected:
awake buffers.waitingWriter
- echo "About to wait for writer"
buffers.waitingReader.enterWait "waiting for writer to buffer more data"
- echo "Awaken from wait"
bytesInBuffersNow = buffers.totalBufferedBytes
if buffers.eofReached:
- echo "read bytes ", bytesInBuffersNow - bytesInBuffersAtStart
- describeBuffers "at end", buffers
return bytesInBuffersNow - bytesInBuffersAtStart
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
@@ -61,7 +51,6 @@ proc pipeRead(s: LayeredInputStream,
proc pipeWrite(s: LayeredOutputStream, src: pointer, srcLen: Natural) {.async.} =
let buffers = s.buffers
- echo "pipe write"
while buffers.canAcceptWrite(srcLen) == false:
buffers.waitingWriter.enterWait "waiting for reader to drain the buffers"
@@ -81,7 +70,7 @@ let pipeInputVTable = InputStreamVTable(
{.nimcall, gcsafe, raises: [IOError, Defect].} =
fsTranslateErrors "Failed to read from pipe":
let ls = LayeredInputStream(s)
- doAssert ls.allowWaitFor
+ fsAssert ls.allowWaitFor
return waitFor pipeRead(ls, dst, dstLen)
,
readAsync: proc (s: InputStream, dst: pointer, dstLen: Natural): Future[Natural]
@@ -117,7 +106,7 @@ let pipeOutputVTable = OutputStreamVTable(
{.nimcall, gcsafe, raises: [IOError, Defect].} =
fsTranslateErrors "Failed to write all bytes to pipe":
var ls = LayeredOutputStream(s)
- doAssert ls.allowWaitFor
+ fsAssert ls.allowWaitFor
waitFor pipeWrite(ls, src, srcLen)
,
writeAsync: proc (s: OutputStream, src: pointer, srcLen: Natural): Future[void]
@@ -146,7 +135,6 @@ let pipeOutputVTable = OutputStreamVTable(
{.nimcall, gcsafe, raises: [IOError, Defect].} =
s.buffers.eofReached = true
- echo "writer closes the stream"
fsTranslateErrors "Unexpected error from Future.complete":
awake s.buffers.waitingReader
@@ -173,7 +161,7 @@ let pipeOutputVTable = OutputStreamVTable(
func pipeInput*(source: InputStream,
pageSize = defaultPageSize,
allowWaitFor = false): AsyncInputStream =
- doAssert pageSize > 0
+ fsAssert pageSize > 0
AsyncInputStream LayeredInputStream(
vtable: vtableAddr pipeInputVTable,
@@ -199,7 +187,7 @@ proc pipeOutput*(destination: OutputStream,
pageSize = defaultPageSize,
maxBufferedBytes = defaultPageSize * 4,
allowWaitFor = false): AsyncOutputStream =
- doAssert pageSize > 0
+ fsAssert pageSize > 0
var
buffers = initPageBuffers pageSize
@@ -231,7 +219,7 @@ proc pipeOutput*(buffers: PageBuffers,
func asyncPipe*(pageSize = defaultPageSize,
maxBufferedBytes = defaultPageSize * 4): FsAsyncPipe =
- doAssert pageSize > 0
+ fsAssert pageSize > 0
FsAsyncPipe(buffers: initPageBuffers(pageSize, maxBufferedBytes))
func initReader*(pipe: FsAsyncPipe): AsyncInputStream =
diff --git a/faststreams/textio.nim b/faststreams/textio.nim
index 6e201fe..9b7ba09 100644
--- a/faststreams/textio.nim
+++ b/faststreams/textio.nim
@@ -1,6 +1,6 @@
import
stew/ptrops,
- inputs, outputs, buffers, multisync
+ inputs, outputs, buffers, async_backend, multisync
template matchingIntType(T: type int64): type = uint64
template matchingIntType(T: type int32): type = uint32
@@ -112,7 +112,7 @@ const
Digits* = {'0'..'9'}
proc readLine*(s: InputStream, keepEol = false): TaintedString =
- doAssert readableNow(s)
+ fsAssert readableNow(s)
while s.readable:
let c = s.peek.char
@@ -131,7 +131,7 @@ proc readLine*(s: InputStream, keepEol = false): TaintedString =
proc readUntil*(s: InputStream,
sep: openarray[char]): Option[TaintedString] =
- doAssert readableNow(s)
+ fsAssert readableNow(s)
var res = ""
while s.readable(sep.len):
if s.lookAheadMatch(charsToBytes(sep)):
@@ -150,7 +150,7 @@ iterator lines*(s: InputStream, keepEol = false): TaintedString =
yield readLine(s, keepEol)
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 =
advance s
diff --git a/tests/test_inputs.nim b/tests/test_inputs.nim
index 958865f..2d8d7b5 100644
--- a/tests/test_inputs.nim
+++ b/tests/test_inputs.nim
@@ -13,7 +13,7 @@ proc bytes(s: string): seq[byte] =
proc str(bytes: openarray[byte]): string =
result = newStringOfCap(bytes.len)
- for b in bytes:
+ for b in items(bytes):
result.add b.char
proc countLines(s: InputStream): Natural =
@@ -31,13 +31,15 @@ procSuite "input stream":
test "input is not readable with read":
check not input.readable
- expect Defect:
- echo "This read should not complete: ", input.read
+ when not defined(danger):
+ expect Defect:
+ echo "This read should not complete: ", input.read
test "input is not readable with read(n)":
check not input.readable(10)
- expect Defect:
- echo "This read should not complete: ", input.read(10)
+ when not defined(danger):
+ expect Defect:
+ echo "This read should not complete: ", input.read(10)
test "next returns none":
check input.next.isNone
diff --git a/tests/test_pipelines.nim b/tests/test_pipelines.nim
index 68cd266..4405a23 100644
--- a/tests/test_pipelines.nim
+++ b/tests/test_pipelines.nim
@@ -26,7 +26,6 @@ proc upcaseAllCharacters(i: InputStream, o: OutputStream) {.fsMultiSync.} =
while i.readable:
o.write toUpperAscii(i.read.char)
- echo "closing upcase"
close o
proc printTimes(t: TestTimes) =
@@ -51,7 +50,6 @@ procSuite "pipelines":
"""
- #[
test "upper-case/base64 pipeline benchmark":
var
times: TestTimes
@@ -61,9 +59,8 @@ procSuite "pipelines":
let inputText = loremIpsum.repeat(5000)
- when debugHelpers:
- echo "Input len: ", inputText.len
- echo "Base 64 len: ", base64.encode(inputText).len
+ timeIt times.stdFunctionCalls:
+ stdRes = base64.decode(base64.encode(toUpperAscii(inputText)))
timeIt times.fsPipeline:
fsRes = executePipeline(unsafeMemoryInput(inputText),
@@ -79,21 +76,14 @@ procSuite "pipelines":
base64decode,
getOutput string)
- timeIt times.stdFunctionCalls:
- stdRes = base64.decode(base64.encode(toUpperAscii(inputText)))
-
check fsAsyncRes == stdRes
check fsRes == stdRes
printTimes times
- ]#
asyncTest "upper-case/base64 async pipeline":
let pipe = asyncPipe()
- let inputText = repeat(loremIpsum, 8)
-
- when debugHelpers:
- echo "Input len: ", inputText.len
+ let inputText = repeat(loremIpsum, 100)
proc pipeFeeder(s: AsyncOutputStream) {.gcsafe, async.} =
randomize 1234
@@ -112,7 +102,6 @@ procSuite "pipelines":
let sleep = rand(50) - 45
if sleep > 0:
- echo "written ", pos
await sleepAsync(sleep.milliseconds)
close s
diff --git a/tests/test_pipelines.nim.cfg b/tests/test_pipelines.nim.cfg
new file mode 100644
index 0000000..9734574
--- /dev/null
+++ b/tests/test_pipelines.nim.cfg
@@ -0,0 +1,4 @@
+--linedir:off
+--linetrace:off
+--stacktrace:off
+