From 18f9488beef35fa45b3c6ef657627e1b577fb634 Mon Sep 17 00:00:00 2001 From: Zahary Karadjov Date: Tue, 5 May 2020 19:20:25 +0300 Subject: [PATCH] Expand the documentation; Clean up debugging code; Enable all tests --- README.md | 141 ++++++++++++++++++++++++------- faststreams/async_backend.nim | 7 ++ faststreams/buffers.nim | 12 +-- faststreams/chronos_adapters.nim | 4 +- faststreams/inputs.nim | 58 ++++++++----- faststreams/outputs.nim | 36 ++++---- faststreams/pipelines.nim | 26 ++---- faststreams/textio.nim | 8 +- tests/test_inputs.nim | 12 +-- tests/test_pipelines.nim | 17 +--- tests/test_pipelines.nim.cfg | 4 + 11 files changed, 205 insertions(+), 120 deletions(-) create mode 100644 tests/test_pipelines.nim.cfg 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 +