diff --git a/README.md b/README.md index a2fc533..459e9b2 100644 --- a/README.md +++ b/README.md @@ -6,7 +6,410 @@ [![License: MIT](https://img.shields.io/badge/License-MIT-blue.svg)](https://opensource.org/licenses/MIT) ![Stability: experimental](https://img.shields.io/badge/stability-experimental-orange.svg) -Nearly zero-overhead input/output streams for Nim +FastStreams is a highly efficient library for all your I/O needs. + +It offers nearly zero-overhead synchronous and asynchronous streams +for handling inputs and outputs of various types: + +* Memory inputs and outputs for serialization frameworks and parsers +* File inputs and outputs +* Pipes and Process I/O +* Networking + +The library aims to provide a common interface between all stream types +that allows the application code to be easily portable to different back-end +event loops. In particular, [Chronos](https://github.com/status-im/nim-chronos) +and [AsyncDispatch](https://nim-lang.org/docs/asyncdispatch.html) +are already supported. It's envisioned that the library will also +gain support for the Nginx event loop to allow the creation of web +applications running as Nginx run-time modules and the [SeaStar event loop](http://seastar.io/) +for the development of extremely low-latency services taking advantage +of [kernel-bypass networking](https://blog.cloudflare.com/kernel-bypass/). + +## What does zero-overhead mean? + +Even though FastStreams support multiple stream types, the API is designed +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 +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: + +* Parsers for data formats and protocols employing formal grammars +* Block ciphers +* Compressors and decompressors +* Stream multiplexers + +The zero-copy behavior and low-memory usage is maintained even when multiple +streams are layered on top of each other while back-pressure is properly +accounted for. This makes FastStreams ideal for implementing highly-flexible +networking stacks such as [LibP2P](https://github.com/status-im/nim-libp2p). + +## The key ideas in the FastStreams design + +FastStreams is heavily inspired by the `System.IO.Pipelines` API which was +developed and released by Microsoft in 2018 and is considered the result of +multiple years of evolution over similar APIs shipped in previous SDKs. + +We highly recommend reading the following two articles which provide an in-depth +explanation for the benefits of the design: + +* https://blog.marcgravell.com/2018/07/pipe-dreams-part-1.html +* https://blog.marcgravell.com/2018/07/pipe-dreams-part-2.html + +Here, we'll only summarize the main insights: + +### Obtaining data from the input device is not the same as consuming it. + +When protocols and formats are layered on top of each other, it's highly +inconvenient to handle a read operation that can return an arbitrary amount +of data. If not enough data was returned, you may need to copy the available +bytes into a local buffer and then repeat the reading operation until enough +data is gathered and the local buffer can be processed. On the other hand, +if more data was received, you need to complete the current stage of processing +and then somehow feed the remaining bytes into the next stage of processing +(e.g this might be a nested format or a different parsing branch in the formal + grammar of the protocol). Both of these scenarios require logic that is +difficult to write correctly and results in unnecessary copying of the input +bytes. + +A major difference in the FastStreams design is that the arbitrary-length +data obtained from the input device is managed by the stream itself and you +are provided with an API allowing you to control precisely how much data +is consumed from the stream. Consuming the buffered content does not invoke +costly asynchronous calls and you are allowed to peek at the stream contents +before deciding which step to take next (something crucial for handling formal +grammars). Thus, using the FastStreams API results in code that is both highly +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 +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 +handles the buffer switching logic automatically for the user. + +Nevertheless, the buffering logic can be configured for unbuffered reads +and writes and it supports efficiently various common real-world patterns +such as: + +* Length prefixes + + To handle protocols with length prefixes without any memory overhead, + the output streams support "delayed writes" where a portion of the + stream content is specified only after the prefixed content is written + to the stream. + +* 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. + +* Content with known length + + Some streams have 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: + +### `InputStream` + +An `InputStream` manages a particular input device. The library offers out +of the box the following input stream types: + +* `fileInput` + + For reading files through the familiar `fread` API from the C run-time. + +* `memFileInput` + + For reading memory mapped files which provides best performance. + +* `unsafeMemoryInput` + + For handling strings, sequences and openarrays as an input stream.
+ You are responsible for ensuring that the backing buffer won't be invalidated + while the stream is being used. + +* `chronosInput` (async) + + Enabled by importing `faststreams/chronos_adapters`.
+ It can represent any Chronos `Transport` as an input stream. + +* `asyncSocketInput` (async) + + Enabled by importing `faststreams/std_adapters`.
+ Allows using Nim's standard library `AsyncSocket` type as an input stream. + +You can extend the library with new `InputStream` types without modifying it. +Please see the inline code documentation of `InputStreamVTable` for more details. + +All of the above APIs are possible constructors for creating an `InputStream`. +The stream instances will manage their resources through destructors, but you +might want to `close` them explicitly in async context or when you need to +handle the possible errors from the closing operation. + +Here is an example usage: + +```nim +var + jsonString = "[1, 2, 3]" + jsonNodes = parseJson(unsafeMemoryInput(jsonString)) + moreNodes = parseJson(fileInput("data.json")) +``` + +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 = + result = newStringToken() + + advance stream # skip the opening quote + + while stream.readable: + let nextChar = stream.read.char + case nextChar + of '\'': + if stream.readable: + let escaped = stream.read.char + case escaped + of 'n': result.add '\n' + of 't': result.add '\t' + else: result.add escaped + else: + error(UnexpectedEndOfFile) + of '"' + return + else: + result.add nextChar + + error(UnexpectedEndOfFile) + +proc nextToken(stream: InputStream): JsonToken = + while stream.readable: + case stream.peek.char + of '"': + result = scanString(stream) + of '0'..'9': + result = scanNumber(stream) + of 'a'..'z', 'A'..'Z', '_': + result = scanIdentifier(stream) + of '{': + advance stream # skip the character + result = objectStartToken + ... + + return eofToken + +proc parseJson(stream: InputStream): JsonNode = + while (let token = nextToken(stream); token != eofToken): + case token + of numberToken: + result = newJsonNumber(token.num) + of stringToken: + result = newJsonString(token.str) + of objectStartToken: + result = parseObject(stream) + ... +``` + +The above example is nothing but a toy program, but we can already see many +usage patterns of the `InputStream` type. For a more sophisticated and complete +implementation of a JSON parser, please see the [nim-json-serialization](https://github.com/status-im/nim-json-serialization) +package. + +As we can see from the example above, calling `stream.read` should always be +preceded by a call to `stream.readable`. When the stream is in the readable +state, we can also `peek` at the next character before we decide how to +proceed. Besides calling `read`, we can also mark the data as consumed by +calling `stream.advance`. + +The above APIs demonstrate how you can consume the data one byte at the time. +Common wisdom might tell you that this should be inefficient, but that's not +the case with FastStreams. The loop `while stream.readable: stream.read` will +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). +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. + +Sometimes, you need to check whether the stream contains at least a specific +number of bytes. You can use the `stream.readable(N)` API to achieve this. + +Reading multiple bytes at once is then possible with `stream.read(N)`, but +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. + +### `OutputStream` + +An `OutputStream` manages a particular output device. The library offers out +of the box the following output stream types: + +* `writeFileOutput` + + For writing files through the familiar `fwrite` API from the C run-time. + +* `memoryOutput` + + For building a `string` or a `seq[byte]` result. + +* `unsafeMemoryOutput` + + For writing to an arbitrary existing buffer.
+ You are responsible for ensuring that the backing buffer won't be invalidated + while the stream is being used. + +* `chronosOutput` (async) + + Enabled by importing `faststreams/chronos_adapters`.
+ It can represent any Chronos `Transport` as an input stream. + +* `asyncSocketOutput` (async) + + Enabled by importing `faststreams/std_adapters`.
+ Allows using Nim's standard library `AsyncSocket` type as an output stream. + +You can extend the library with new `OutputStream` types without modifying it. +Please see the inline code documentation of `OutputStreamVTable` for more details. + +All of the above APIs are possible constructors for creating an `OutputStream`. +The stream instances will manage their resources through destructors, but you +might want to `close` them explicitly in async context or when you need to +handle the possible errors from the closing operation. + +Here is an example usage: + +```nim +type + ABC = object + a: int + b: char + c: string + +var stream = memoryOutput() +stream.writeNimRepr(ABC(a: 1, b: 'b', c: "str")) +var repr = stream.getOutput(string) +``` + +The `writeNimRepr` in the above example is not part of the library, but +let's see how it can be implemented: + +```nim +import + typetraits, faststreams + +proc writeNimRepr*(stream: OutputStream, str: string) = + stream.write '"' + + for c in str: + if c == '"': + stream.write ['\'', '"'] + else: + stream.write c + + stream.write '"' + +proc writeNimRepr*(stream: OutputStream, x: char) = + stream.write ['\'', x, '\''] + +proc writeNimRepr*(stream: OutputStream, x: int) = + stream.write $x # Making this more optimal has been left + # as an exercise for the reader + +proc writeNimRepr*[T](stream: OutputStream, obj: T) = + stream.write typetraits.name(T) + stream.write '(' + + var firstField = true + for name, val in fieldPairs(obj): + if not firstField: + stream.write ", " + + stream.write name + stream.write ": " + stream.writeNimRepr val + + firstField = false + + stream.write ')' +``` + +When the stream is created, its output buffers will be initialized with a +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. + +Please note that even in async context, `write` will complete immediately. +To handle back-pressure properly, use `stream.flush` or `stream.waitForConsumer` +which will ensure that the buffered data is drained to a specified number of +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`. + +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. + +### `Pipeline` + +(This section is a stub and it will be expanded with more details in the future) + +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`. + +Each transformation step is a function of the kind: + +```nim +type PipelineStep* = proc (i: InputStream, o: OutputStream) + {.gcsafe, raises: [Defect, CatchableError].} +``` + +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. + +The pipeline transformation steps are usually employing the `fsMultiSync` +pragma to make them usable in both synchronous and asynchronous scenarios. + +Please note that the above higher-level APIs are just about simplifying the +instantiation of multiple `Pipe` objects that can be used to hook input and +output streams in arbitrary ways. + +A stream multiplexer for example is likely to rely on the lower-level `Pipe` +objects and the underlying `PageBuffers` directly. ## License diff --git a/faststreams.nim b/faststreams.nim index 8a8d538..6fcda7d 100644 --- a/faststreams.nim +++ b/faststreams.nim @@ -1,6 +1,6 @@ import - faststreams/[input_stream, output_stream] + faststreams/[inputs, outputs] export - input_stream, output_stream + inputs, outputs diff --git a/faststreams.nimble b/faststreams.nimble index 7fa698d..628719a 100644 --- a/faststreams.nimble +++ b/faststreams.nimble @@ -1,7 +1,7 @@ mode = ScriptMode.Verbose packageName = "faststreams" -version = "0.1.0" +version = "0.2.0" author = "Status Research & Development GmbH" description = "Nearly zero-overhead input/output streams for Nim" license = "Apache License 2.0" @@ -12,6 +12,7 @@ requires "nim >= 1.2.0", "chronos" task test, "Run all tests": - exec "nim c -r --threads:off tests/all_tests" - exec "nim c -r --threads:on tests/all_tests" + exec "nim c -r -d:debug --threads:on tests/all_tests" + exec "nim c -r -d:release --threads:on tests/all_tests" + exec "nim c -r -d:danger --threads:on tests/all_tests" diff --git a/faststreams/async_backend.nim b/faststreams/async_backend.nim index 6d39875..30de7b5 100644 --- a/faststreams/async_backend.nim +++ b/faststreams/async_backend.nim @@ -1,18 +1,29 @@ const faststreams_async_backend {.strdefine.} = "chronos" -when faststreams_async_backend == "chronos": - import chronos # import chronos/[asyncfutures2, asyncmacro2] - export chronos # export asyncfutures2, asyncmacro2 +type + CloseBehavior* = enum + waitAsyncClose + dontWaitAsyncClose - template faststreamsAwait*(f: Future): untyped = +when faststreams_async_backend == "chronos": + import + chronos + + export + chronos + + template fsAwait*(f: Future): untyped = await f elif faststreams_async_backend in ["std", "asyncdispatch"]: - import std/[asyncfutures, asyncmacro] - export asyncfutures, asyncmacro + import + std/[asyncfutures, asyncmacro] + + export + asyncfutures, asyncmacro - template faststreamsAwait*(awaited: Future[T]): untyped = + template fsAwait*(awaited: Future[T]): untyped = # TODO revisit after https://github.com/nim-lang/Nim/pull/12085/ is merged let f = awaited yield f @@ -23,9 +34,15 @@ elif faststreams_async_backend in ["std", "asyncdispatch"]: else: {.fatal: "Unrecognized network backend: " & faststreams_async_backend.} -template raiseFaststreamsError*(errMsg: string, body: untyped) = +template fsTranslateErrors*(errMsg: string, body: untyped) = try: body - except CatchableError as err: - raise newException(IOError, errMsg, err) + except Exception as err: + if err[] of Defect: + raise (ref Defect)(err) + else: + raise newException(IOError, errMsg, err) + +template noAwait*(expr: untyped): untyped = + expr diff --git a/faststreams/buffers.nim b/faststreams/buffers.nim new file mode 100644 index 0000000..9c91a8c --- /dev/null +++ b/faststreams/buffers.nim @@ -0,0 +1,226 @@ +import + deques, + stew/[ptrops, ranges/ptr_arith], + async_backend + +type + PageKind* = enum + userPage + stringPage + mallocPage + + PageSpan* = object + startAddr*, endAddr*: ptr byte + + Page* = object + startOffset*: Natural + endOffset*: Natural + case kind*: PageKind + of userPage, mallocPage: + bufferStart, bufferEnd: ptr byte + of stringPage: + data*: ref string + + PageRef* = ref Page + + PageBuffers* = ref object + pageSize*: Natural + maxWriteSize*: Natural + backPressureLimit*: Natural + + queue*: Deque[PageRef] + getters: seq[Future[void]] + putters: seq[Future[void]] + + eofReached: bool + + totalBytesRead*: Natural + totalBytesWritten*: Natural + +const + nimPageSize* = 4096 + pageMetadataSize* = offsetof(Page, data) + nimAllocatorMetadataSize* = 32 + # TODO: Get this legally from the Nim allocator. + # The goal is to make perfect page-aligned allocations + # that get fast O(0) treatment. + defaultPageSize* = 4096 - (pageMetadataSize + nimAllocatorMetadataSize) + maxStackUsage* = 16384 + +func pageBaseAddr*(page: PageRef): ptr byte = + if page.kind == stringPage: + cast[ptr byte](addr page.data[][0]) + else: + page.bufferStart + +func pageStartAddr*(page: PageRef): ptr byte = + if page.kind == stringPage: + offset(cast[ptr byte](addr page.data[][0]), page.startOffset) + else: + offset(page.bufferStart, page.startOffset) + +func pageEndAddr*(page: PageRef): ptr byte = + if page.kind == stringPage: + offset(cast[ptr byte](addr page.data[][0]), page.endOffset) + else: + offset(page.bufferStart, page.endOffset) + +template pageChars*(page: PageRef): untyped = + let baseAddr = cast[ptr UncheckedArray[char]](pageBaseAddr(page)) + toOpenArray(baseAddr, page.startOffset, page.endOffset - 1) + +func span*(page: PageRef, writable: static[bool] = false): PageSpan = + if page.kind == stringPage: + let baseAddr = cast[ptr byte](addr page.data[][0]) + PageSpan(startAddr: offset(baseAddr, page.startOffset), + endAddr: offset(baseAddr, when writable: page.data[].len + else: page.endOffset)) + else: + PageSpan(startAddr: offset(page.bufferStart, page.startOffset), + endAddr: when writable: page.bufferEnd + else: offset(page.bufferStart, page.endOffset)) + +template writableSpan*(page: PageRef): PageSpan = + span(page, writable = true) + +func initPageBuffers*(pageSize: Natural, + maxWriteSize = high(int)): PageBuffers = + if pageSize > 0: + return PageBuffers(pageSize: pageSize, + maxWriteSize: maxWriteSize) + +template allocRef[T: not ref](x: T): ref T = + let res = new type(x) + res[] = x + res + +func getWritablePage*(buffers: PageBuffers): PageRef = + # TODO: The semantics of this func are quite unusual + # I should find a more appropriate name + if buffers.queue.len == 0: + result = PageRef(kind: stringPage, + data: allocRef newString(buffers.pageSize), + endOffset: buffers.pageSize) + buffers.queue.addLast result + else: + result = buffers.queue[0] + +func addWritablePage*(buffers: PageBuffers, pageSize: Natural): PageRef = + result = PageRef(kind: stringPage, + data: allocRef newString(pageSize), + endOffset: pageSize) + buffers.queue.addLast result + +func addWritablePage*(buffers: PageBuffers): PageRef = + buffers.addWritablePage(buffers.pageSize) + +template getWritableSpan*(buffers: PageBuffers): PageSpan = + getWritablePage(buffers).span(writable = true) + +func ensureRunway*(buffers: PageBuffers, neededRunway: Natural): PageSpan = + doAssert buffers.queue.len == 0 + buffers.pageSize = neededRunway + getWritableSpan(buffers) + +template len*(buffers: PageBuffers): int = + buffers.queue.len + +template popFirst*(buffers: PageBuffers): PageRef = + buffers.queue.popFirst + +template `[]`*(buffers: PageBuffers, idx: Natural): PageRef = + buffers.queue[idx] + +func splitLastPageAt*(buffers: PageBuffers, address: ptr byte) = + var + topPage = buffers.queue.peekLast + newPage = PageRef() + splitPosition = distance(topPage.pageBaseAddr, address) + + newPage[] = topPage[] + topPage.endOffset = splitPosition + newPage.startOffset = splitPosition + + buffers.queue.addLast newPage + +func endLastPageAt*(buffers: PageBuffers, address: ptr byte) = + if buffers != nil and buffers.queue.len > 0: + var topPage = buffers.queue.peekLast + topPage.endOffset = distance(topPage.pageBaseAddr, address) + +func trackPageWrite*(page: PageRef, bytesWritten: Natural) {.inline.} = + page.endOffset = page.startOffset + bytesWritten + +template writeToSpan*(buffersParam: PageBuffers, + spanVarName, writeExpr: untyped) = + var + buffers = buffersParam + page = buffers.getWritablePage + spanVarName = page.writableSpan + + # TODO: what if we exit with an exception here? + # Are the side-effects of `getWritablePage` above OK to keep? + + let bytesWritten = writeExpr + trackPageWrite(page, bytesWritten) + + if bytesWritten == 0: + buffers.eofReached = true + +func nextAlignedSize*(minSize, pageSize: Natural): Natural = + # TODO: This is not perfectly accurate. Revisit later + ((minSize div pageSize) + 1) * pageSize + +template consumeAllPages*(buffersParam: PageBuffers, + pageAddrVar, pageLenVar, body: untyped) = + let buffers = buffersParam + doAssert buffers != nil + + var recycledPage: PageRef + for page in buffers.queue: + let + pageAddrVar = page.pageStartAddr + pageLenVar = page.endOffset - page.startOffset + + if page.kind == stringPage and page.data[].len == buffers.pageSize: + recycledPage = page + + # TODO: what if the body throws an exception? + # Should we do anything with the remaining pages? + body + + buffers.queue.clear() + + if recycledPage != nil: + recycledPage.startOffset = 0 + recycledPage.endOffset = 0 + buffers.queue.addLast recycledPage + +template wasEofReached*(buffers: PageBuffers): bool = + buffers.eofReached + +# BEWARE! These templates violate the double evaluation +# safety measures in order to produce better inlined +# code. We are using a `var` type to make it harder +# to accidentally misuse them. +template len*(span: var PageSpan): Natural = + distance(span.startAddr, span.endAddr) + +template atEnd*(span: var PageSpan): bool = + span.startAddr == span.endAddr + +template hasRunway*(span: var PageSpan): bool = + span.startAddr != span.endAddr + +template bumpPointer*(span: var PageSpan, numberOfBytes: Natural = 1) = + span.startAddr = offset(span.startAddr, numberOfBytes) + +template writeByte*(span: var PageSpan, val: byte) = + span.startAddr[] = val + span.startAddr = offset(span.startAddr, 1) + +template charsToBytes*(chars: openArray[char]): untyped = + bind makeOpenArray + var charsStart = unsafeAddr chars[0] + makeOpenArray(cast[ptr byte](charsStart), chars.len) + diff --git a/faststreams/chronos_adapters.nim b/faststreams/chronos_adapters.nim index 39a10e4..450bb76 100644 --- a/faststreams/chronos_adapters.nim +++ b/faststreams/chronos_adapters.nim @@ -1,6 +1,6 @@ import chronos, - input_stream, output_stream, multisync + inputs, outputs, buffers, multisync export chronos, fsMultiSync @@ -21,39 +21,43 @@ const writeIncompleteErrMsg = "Failed to write all bytes to Chronos transport" proc fsCloseWait(t: StreamTransport) {.async, raises: [Defect, IOError].} = - raiseFaststreamsError closingErrMsg: + fsTranslateErrors closingErrMsg: await t.closeWait() proc fsReadOnce(t: StreamTransport, - buffer: ptr byte, bufSize: int): Future[int] {.async, raises: [Defect, IOError].} = - raiseFaststreamsError readingErrMsg: - return t.readOnce(pointer(buffer), bufSize) + buffer: ptr byte, bufSize: int) + {.raises: [Defect, IOError], async.} = + fsTranslateErrors readingErrMsg: + buffers.writeToSpan(span): + await t.readOnce(span.startAddr, span.len) # TODO: Use the Raising type here let ChronosInputStreamVTable = InputStreamVTable( - readSync: proc (s: InputStream, buffer: ptr byte, bufSize: int): int + readSync: proc (s: InputStream, buffers: PageBuffers) {.nimcall, gcsafe, raises: [IOError, Defect].} = var cs = ChronosInputStream(s) doAssert cs.allowWaitFor - raiseFaststreamsError readingErrMsg: - return waitFor cs.transport.readOnce(pointer(buffer), bufSize) + + fsTranslateErrors readingErrMsg: + buffers.writeToSpan(span): + waitFor cs.transport.readOnce(span.startAddr, span.len) , - readAsync: proc (s: InputStream, buffer: ptr byte, bufSize: int): Future[int] + readAsync: proc (s: InputStream, buffers: PageBuffers): Future[Natural] {.nimcall, gcsafe, raises: [IOError, Defect].} = - ChronosInputStream(s).transport.fsReadOnce(buffer, bufSize) + ChronosInputStream(s).transport.fsReadOnce(buffers) , closeSync: proc (s: InputStream) {.nimcall, gcsafe, raises: [IOError, Defect].} = - raiseFaststreamsError closingErrMsg: + fsTranslateErrors closingErrMsg: ChronosInputStream(s).transport.close() , - closeAsync: proc (s: InputStream, cb: CloseAsyncCallback): Future[void] + closeAsync: proc (s: InputStream): Future[void] {.nimcall, gcsafe, raises: [IOError, Defect].} = ChronosInputStream(s).transport.fsCloseWait() ) func chronosInput*(s: StreamTransport, - pageSize = output_stream.defaultPageSize, + pageSize = buffers.defaultPageSize, allowWaitFor = false): InputStreamHandle = InputStreamHandle(s: ChronosInputStream( vtable: vtableAddr ChronosInputStreamVTable, @@ -65,7 +69,7 @@ let ChronosOutputStreamVTable = OutputStreamVTable( {.nimcall, gcsafe, raises: [IOError, Defect].} = var cs = ChronosOutputStream(s) doAssert cs.allowWaitFor - let bytesWritten = raiseFaststreamsError writingErrMsg: + let bytesWritten = fsTranslateErrors writingErrMsg: waitFor cs.transport.write(unsafeAddr page[0], page.len) if bytesWritten != page.len: raise newException(IOError, writeIncompleteErrMsg) @@ -107,7 +111,7 @@ let ChronosOutputStreamVTable = OutputStreamVTable( ) func chronosOutput*(s: StreamTransport, - pageSize = output_stream.defaultPageSize, + pageSize = buffers.defaultPageSize, allowWaitFor = false): OutputStreamHandle = var stream = ChronosOutputStream( vtable: vtableAddr(SnappyStreamVTable), diff --git a/faststreams/input_stream.nim b/faststreams/input_stream.nim deleted file mode 100644 index 0c1894d..0000000 --- a/faststreams/input_stream.nim +++ /dev/null @@ -1,305 +0,0 @@ -import - memfiles, options, - stew/[ptrops, ranges/ptr_arith], - async_backend - -type - InputStream* = ref object of RootObj - vtable*: ptr InputStreamVTable - head*: ptr byte - pageSize*: int - bufferSize: int - bufferStart, bufferEnd: ptr byte - bufferEndPos: int - - LayeredInputStream* = ref object of InputStream - subStream*: InputStream - - InputStreamHandle* = object - s*: InputStream - - AsyncInputStream* {.borrow: `.`.} = distinct InputStream - - ReadSyncProc* = proc (s: InputStream, buffer: ptr byte, bufSize: int): int - {.nimcall, gcsafe, raises: [IOError, Defect].} - - ReadAsyncProc* = proc (s: InputStream, buffer: ptr byte, bufSize: int): Future[int] - {.nimcall, gcsafe, raises: [IOError, Defect].} - - CloseSyncProc* = proc (s: InputStream) - {.nimcall, gcsafe, raises: [IOError, Defect].} - - CloseAsyncProc* = proc (s: InputStream): Future[void] - {.nimcall, gcsafe, raises: [IOError, Defect].} - - GetLenSyncProc* = proc (s: InputStream): int - {.nimcall, gcsafe, raises: [IOError, Defect].} - - InputStreamVTable* = object - readSync*: ReadSyncProc - readAsync*: ReadAsyncProc - closeSync*: CloseSyncProc - closeAsync*: CloseAsyncProc - getLenSync*: GetLenSyncProc - - FileInputStream = ref object of InputStream - file: MemFile - -const - lengthUnknown* = -1 - debugHelpers = false - nimAllocatorMetadataSize* = 0 - # TODO: Get this from Nim's allocator. - # The goal is to make perfect page-aligned allocations - # defaultPageSize = 4096 - nimAllocatorMetadataSize - -proc preventFurtherReading(s: InputStream) = - s.vtable = nil - s.head = nil - s.bufferEnd = nil - -proc close*(s: InputStream) {.raises: [IOError, Defect].} = - if s != nil: - if s.vtable != nil and s.vtable.closeSync != nil: - s.vtable.closeSync(s) - - s.preventFurtherReading() - -# TODO -# The destructors are currently disabled because they seem to cause -# mysterious segmentation faults related to corrupted GC internal -# data structures. -#[ -proc `=destroy`*(h: var InputStreamHandle) {.raises: [Defect].} = - if h.s != nil: - if h.s.vtable != nil and h.s.vtable.closeSync != nil: - try: - h.s.vtable.closeSync(h.s) - except IOError: - # Since this is a destructor, there is not much we can do here. - # If the user wanted to handle the error, they would have called - # `close` manually. - discard # TODO - # TODO ATTENTION! - # Uncommenting the following line will lead to a GC heap corruption. - # Most likely this leads to Nim collecting some object prematurely. - # h.s = nil - # We work-around the problem through more indirect incapacitatation - # of the stream object: - h.s.preventFurtherReading() -]# - -converter implicitDeref*(h: InputStreamHandle): InputStream = - h.s - -let FileStreamVTable = InputStreamVTable( - closeSync: proc (s: InputStream) - {.nimcall, gcsafe, raises: [IOError, Defect].} = - try: - close FileInputStream(s).file - except OSError as err: - raise newException(IOError, "Failed to close file", err) - , - getLenSync: proc (s: InputStream): int - {.nimcall, gcsafe, raises: [IOError, Defect].} = - distance(s.head, s.bufferEnd) -) - -template vtableAddr*(vtable: InputStreamVTable): ptr InputStreamVTable = - ## This is a simple work-around for the somewhat broken side - ## effects analysis of Nim - reading from global let variables - ## is considered a side-effect. - {.noSideEffect.}: - unsafeAddr vtable - -proc fileInput*(filename: string): InputStreamHandle = - let - memFile = memfiles.open(filename) - head = cast[ptr byte](memFile.mem) - fileSize = memFile.size - - var stream = FileInputStream( - vtable: vtableAddr FileStreamVTable, - head: head, - bufferEnd: offset(head, fileSize), - bufferEndPos: fileSize, - file: memFile) - - when debugHelpers: - stream.bufferStart = head - - InputStreamHandle(s: stream) - -proc memoryInput*(mem: openarray[byte]): InputStreamHandle = - let head = unsafeAddr mem[0] - InputStreamHandle(s: InputStream( - head: head, - bufferEnd: offset(head, mem.len), - bufferEndPos: mem.len)) - -proc memoryInput*(str: string): InputStreamHandle = - memoryInput str.toOpenArrayByte(0, str.len - 1) - -# TODO: Is this used, should we deprecate it? -proc endPos*(s: InputStream): int = - doAssert s.vtable == nil or s.vtable.getLenSync != nil - return s.bufferEndPos - -# TODO The return type here could be Option[Natural] if Nim had -# the Option[range] optimisation that will make it equvalent to `int`. -proc len*(s: InputStream): int {.raises: [Defect, IOError].} = - if s.vtable == nil: - distance(s.head, s.bufferEnd) - elif s.vtable.getLenSync != nil: - s.vtable.getLenSync(s) - else: - lengthUnknown - -template len*(s: AsyncInputStream): int = - len InputStream(s) - -proc bufferMoreDataSync(s: InputStream): bool = - # Returns true if more data was successfully buffered - if s.vtable == nil or s.vtable.readSync == nil: - return false - - let bytesRead = s.vtable.readSync(s, s.bufferStart, s.bufferSize) - if bytesRead == 0: - # TODO close the input device - s.vtable = nil - return false - else: - s.bufferEnd = offset(s.bufferStart, bytesRead) - s.bufferEndPos += bytesRead - return true - -proc bufferMoreDataAsync(s: AsyncInputStream): Future[bool] {.async.} = - # Returns true if more data was successfully buffered - return false - -proc readable*(s: InputStream): bool = - if s.head != s.bufferEnd: - true - else: - s.bufferMoreDataSync() - -template readable*(sp: AsyncInputStream): bool = - let s = sp - if s.head != s.bufferEnd: - true - else: - faststreamsAwait bufferMoreDataAsync(s) - -proc readable*(s: InputStream, n: int): bool = - if distance(s.head, s.bufferEnd) >= n: - return true - - if s.vtable == nil or s.vtable.readSync == nil: - return false - - # TODO - doAssert false, "Multi-buffer reading will be implemented later" - -template readable*(sp: AsyncInputStream, n: int): bool = - let s = sp - - if distance(s.head, s.bufferEnd) >= n: - return true - - if s.vtable == nil: - return false - - # TODO - doAssert false, "Multi-buffer reading will be implemented later" - -template close*(s: AsyncInputStream) = - close InputStream(s) - -proc peek*(s: InputStream): byte {.inline.} = - doAssert s.head != s.bufferEnd - return s.head[] - -template peek*(s: AsyncInputStream): byte = - peek InputStream(s) - -proc peekAt*(s: InputStream, pos: int): byte {.inline.} = - # TODO implement page flipping - let peekHead = offset(s.head, pos) - doAssert cast[uint](peekHead) < cast[uint](s.bufferEnd) - return peekHead[] - -template peekAt*(s: AsyncInputStream, pos: int): byte = - peekAt InputStream(s) - -when debugHelpers: - proc showPosition*(s: InputStream) = - echo "head at ", distance(s.bufferStart, s.head), "/", - distance(s.bufferStart, s.bufferEnd) - -proc advance*(s: InputStream) = - if s.head != s.bufferEnd: - s.head = offset(s.head, 1) - else: - discard s.bufferMoreDataSync() - -template advance*(sp: AsyncInputStream) = - let s = sp - if s.head != s.bufferEnd: - s.head = offset(s.head, 1) - else: - discard faststreamsAwait(bufferMoreDataAsync(s)) - -proc read*(s: InputStream): byte = - result = s.peek() - advance s - -template read*(sp: AsyncInputStream): byte = - let s = sp - let res = s.peek() - advance(s) - res - -proc checkReadAhead(s: InputStream, n: int): ptr byte = - result = s.head - doAssert distance(s.head, s.bufferEnd) >= n - s.head = offset(s.head, n) - -template read*(s: InputStream, n: int): auto = - makeOpenArray(checkReadAhead(s, n), n) - -proc next*(s: InputStream): Option[byte] = - if readable(s): - result = some read(s) - -template next*(sp: AsyncInputStream): Option[byte] = - let s = sp - if readable(s): - some read(s) - else: - none byte - -proc bufferPos(s: InputStream, pos: int): ptr byte = - let offsetFromEnd = pos - s.bufferEndPos - doAssert offsetFromEnd < 0 - result = offset(s.bufferEnd, offsetFromEnd) - doAssert result >= s.bufferStart - -proc pos*(s: InputStream): int {.inline.} = - s.bufferEndPos - distance(s.head, s.bufferEnd) - -template pos*(s: AsyncInputStream): int = - pos InputStream(s) - -proc firstAccessiblePos*(s: InputStream): int {.inline.} = - s.bufferEndPos - distance(s.bufferStart, s.bufferEnd) - -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 - -proc rewindTo*(s: InputStream, pos: int) {.inline.} = - s.head = s.bufferPos(pos) - diff --git a/faststreams/inputs.nim b/faststreams/inputs.nim new file mode 100644 index 0000000..9809b83 --- /dev/null +++ b/faststreams/inputs.nim @@ -0,0 +1,577 @@ +import + os, memfiles, options, + stew/[ptrops, ranges/ptr_arith], + async_backend, buffers + +export + options, CloseBehavior + +type + InputStream* = ref object of RootObj + vtable: ptr InputStreamVTable # This is nil for unsafe memory inputs + buffers: PageBuffers # This is nil for unsafe memory inputs + span: PageSpan + spanEndPos: Natural + closeFut: Future[void] # This is nil before `close` is called + + LayeredInputStream* = ref object of InputStream + subStream*: InputStream + + InputStreamHandle* = object + s*: InputStream + + AsyncInputStream* {.borrow: `.`.} = distinct InputStream + + ReadSyncProc* = proc (s: InputStream) + {.nimcall, gcsafe, raises: [IOError, Defect].} + + ReadAsyncProc* = proc (s: InputStream): Future[void] + {.nimcall, gcsafe, raises: [IOError, Defect].} + + CloseSyncProc* = proc (s: InputStream) + {.nimcall, gcsafe, raises: [IOError, Defect].} + + CloseAsyncProc* = proc (s: InputStream): Future[void] + {.nimcall, gcsafe, raises: [IOError, Defect].} + + GetLenSyncProc* = proc (s: InputStream): Natural + {.nimcall, gcsafe, raises: [IOError, Defect].} + + InputStreamVTable* = object + readSync*: ReadSyncProc + readAsync*: ReadAsyncProc + closeSync*: CloseSyncProc + closeAsync*: CloseAsyncProc + getLenSync*: GetLenSyncProc + + MemFileInputStream = ref object of InputStream + file: MemFile + + FileInputStream = ref object of InputStream + file: File + +proc disconnectInputDevice(s: InputStream) = + # TODO + # Document the behavior that closeAsync is preferred + if s.vtable != nil: + if s.vtable.closeAsync != nil: + s.closeFut = s.vtable.closeAsync(s) + elif s.vtable.closeSync != nil: + s.vtable.closeSync(s) + s.vtable = nil + +template disconnectInputDevice(s: AsyncInputStream) = + disconnectInputDevice InputStream(s) + +proc preventFurtherReading(s: InputStream) = + s.vtable = nil + s.span = default(PageSpan) + +template preventFurtherReading(s: AsyncInputStream) = + preventFurtherReading InputStream(s) + +template makeHandle*(sp: InputStream): InputStreamHandle = + let s = sp + InputStreamHandle(s: s) + +proc close*(s: InputStream, + behavior = dontWaitAsyncClose) + {.raises: [IOError, Defect].} = + ## Closes the stream. Any resources associated with the stream + ## will be released and no further reading will be possible. + ## + ## If the underlying input device requires asynchronous closing + ## and `behavior` is set to `waitAsyncClose`, this proc will use + ## `waitFor` to block until the async operation completes. + s.disconnectInputDevice() + s.preventFurtherReading() + if s.closeFut != nil: + fsTranslateErrors "Stream closing failed": + if behavior == waitAsyncClose: + waitFor s.closeFut + else: + asyncCheck s.closeFut + +proc close*(s: AsyncInputStream): Future[void] + {.raises: [IOError, Defect].} = + ## Starts the asychronous closing of the stream and returns a future that + ## tracks the closing operation. + s.disconnectInputDevice() + s.preventFurtherReading() + result = InputStream(s).closeFut + doAssert result != nil + +template closeNoWait*(sp: AsyncInputStream|InputStream) = + ## Close the stream without waiting even if's async. + ## This operation will use `asyncCheck` internally to detect unhandled + ## errors from the closing operation. + close(InputStream(s), dontWaitAsyncClose) + +# TODO +# The destructors are currently disabled because they seem to cause +# mysterious segmentation faults related to corrupted GC internal +# data structures. +#[ +proc `=destroy`*(h: var InputStreamHandle) {.raises: [Defect].} = + if h.s != nil: + if h.s.vtable != nil and h.s.vtable.closeSync != nil: + try: + h.s.vtable.closeSync(h.s) + except IOError: + # Since this is a destructor, there is not much we can do here. + # If the user wanted to handle the error, they would have called + # `close` manually. + discard # TODO + # TODO ATTENTION! + # Uncommenting the following line will lead to a GC heap corruption. + # Most likely this leads to Nim collecting some object prematurely. + # h.s = nil + # We work-around the problem through more indirect incapacitatation + # of the stream object: + h.s.preventFurtherReading() +]# + +converter implicitDeref*(h: InputStreamHandle): InputStream = + ## Any `InputStreamHandle` value can be implicitly converted to an + ## `InputStream` or an `AsyncInputStream` value. + h.s + +template vtableAddr*(vtable: InputStreamVTable): ptr InputStreamVTable = + # This is a simple work-around for the somewhat broken side + # effects analysis of Nim - reading from global let variables + # is considered a side-effect. + {.noSideEffect.}: + unsafeAddr vtable + +let MemFileInputVTable = InputStreamVTable( + closeSync: proc (s: InputStream) + {.nimcall, gcsafe, raises: [IOError, Defect].} = + try: + close MemFileInputStream(s).file + except OSError as err: + raise newException(IOError, "Failed to close file", err) + , + getLenSync: proc (s: InputStream): Natural + {.nimcall, gcsafe, raises: [IOError, Defect].} = + s.span.len +) + +proc memFileInput*(filename: string, mappedSize = -1, offset = 0): InputStreamHandle + {.raises: [IOError, OSError].} = + ## Creates an input stream for reading the contents of a memory-mapped file. + ## + ## Using this API will provide better performance than `fileInput`, + ## but this comes at a cost of higher address space usage which may + ## be problematic when working with extremely large files. + ## + ## All parameters are forwarded to Nim's memfiles.open function: + ## + ## ``filename`` + ## The name of the file to read. + ## + ## ``mappedSize`` and ``offset`` + ## can be used to map only a slice of the file. + ## + ## ``offset`` must be multiples of the PAGE SIZE of your OS + ## (usually 4K or 8K, but is unique to your OS) + + # Nim's memfiles module will fail to map an empty file, + # but we don't consider this a problem. The stream will + # be in non-readable state from the start. + let fileSize = getFileSize(filename) + if fileSize == 0: + return makeHandle InputStream() + + let + memFile = memfiles.open(filename, + mode = fmRead, + mappedSize = mappedSize, + offset = offset) + head = cast[ptr byte](memFile.mem) + mappedSize = memFile.size + + makeHandle MemFileInputStream( + vtable: vtableAddr MemFileInputVTable, + span: PageSpan( + startAddr: head, + endAddr: offset(head, mappedSize)), + file: memFile) + +proc readableNow*(s: InputStream): bool = + (not s.span.atEnd) or (s.buffers != nil and s.buffers.len > 1) + +template readableNow*(s: AsyncInputStream): bool = + readableNow InputStream(s) + +func totalUnconsumedBytes*(s: InputStream): Natural = + ## Returns the number of bytes that are currently sitting within the stream + ## buffers and that can be consumed with `read` or `advance`. + result = s.span.len + if s.buffers != nil: + result += s.buffers.totalBytesRead - s.spanEndPos + +template totalUnconsumedBytes*(s: AsyncInputStream): Natural = + ## Alias for InputStream.totalUnconsumedBytes + totalUnconsumedBytes InputStream(s) + +let FileInputVTable = InputStreamVTable( + readSync: proc (s: InputStream) + {.nimcall, gcsafe, raises: [IOError, Defect].} = + let file = FileInputStream(s).file + s.buffers.writeToSpan(span): + file.readBuffer(span.startAddr, span.len) + , + getLenSync: proc (s: InputStream): Natural + {.nimcall, gcsafe, raises: [IOError, Defect].} = + let + s = FileInputStream(s) + runway = s.totalUnconsumedBytes + + let preservedPos = getFilePos(s.file) + setFilePos(s.file, 0, fspEnd) + let endPos = getFilePos(s.file) + setFilePos(s.file, preservedPos) + + endPos - preservedPos + runway + , + closeSync: proc (s: InputStream) + {.nimcall, gcsafe, raises: [IOError, Defect].} = + try: + close FileInputStream(s).file + except OSError as err: + raise newException(IOError, "Failed to close file", err) +) + +proc fileInput*(filename: string, + offset = 0, + pageSize = defaultPageSize): InputStreamHandle + {.raises: [IOError, OSError].} = + ## Creates an input stream for reading the contents of a file + ## through Nim's `io` module. + ## + ## Parameters: + ## + ## ``filename`` + ## The name of the file to read. + ## + ## ``offset`` + ## Initial position in the file where reading should start. + ## + let file = system.open(filename, fmRead) + + if offset != 0: + setFilePos(file, offset) + + makeHandle FileInputStream( + vtable: vtableAddr FileInputVTable, + buffers: initPageBuffers(pageSize), + file: file) + +proc unsafeMemoryInput*(mem: openarray[byte]): InputStreamHandle = + let head = unsafeAddr mem[0] + + makeHandle InputStream( + span: PageSpan( + startAddr: head, + endAddr: offset(head, mem.len)), + spanEndPos: mem.len) + +proc unsafeMemoryInput*(str: string): InputStreamHandle = + unsafeMemoryInput str.toOpenArrayByte(0, str.len - 1) + +proc len*(s: InputStream): Option[Natural] {.raises: [Defect, IOError].} = + if s.vtable == nil: + some s.span.len + elif s.vtable.getLenSync != nil: + some s.vtable.getLenSync(s) + else: + none Natural + +template len*(s: AsyncInputStream): int = + len InputStream(s) + +proc flipPage(s: InputStream) = + doAssert s.buffers.len > 1 + discard s.buffers.popFirst + s.span = s.buffers[0].span + s.spanEndPos += s.span.len + +proc continueAfterRead(s: InputStream): bool = + # Please note that this is extracted into a proc only to reduce the code + # that ends up inlined into async procs by `bufferMoreDataImpl`. + # The inlining itself is required to support the await-free operation of + # the `readable` APIs. + let firstReadPage = s.buffers[0] + + s.span = firstReadPage.span + let bytesRead = s.span.len + s.spanEndPos += bytesRead + + # The read might have been incomplete which signals the EOF of the stream. + # If this is the case, we disconnect the input device which prevents any + # further attempts to read from it: + if wasEofReached(s.buffers): + s.disconnectInputDevice() + + # If we read some bytes anyway, we tell the user code that our buffers + # contain some unconsumed data: + bytesRead > 0 + +template bufferMoreDataImpl(s, awaiter, readOp: untyped): bool = + # This template is always called when the current page has been + # completely exhausted. It should produce `true` if more data was + # successfully buffered, so reading can continue. + # + # The vtable will be `nil` for a memory stream and `vtable.readOp` + # will be `nil` for a memFile. If we've reached here, this is the + # end of the memory buffer, so we can signal EOF: + if s.buffers == nil or s.vtable == nil or s.vtable.readOp == nil: + false + else: + # There might be additional pages in our buffer queue. If so, we + # just jump to the next one: + if s.buffers.len > 1: + flipPage s + true + else: + # We ask our input device to populate our page queue with newly + # read pages. The state of the queue afterwards will tell us if + # the read was successful. In `continueAfterRead`, we examine if + # EOF was reached, but please note that some data might have been + # read anyway: + awaiter s.vtable.readOp(s) + continueAfterRead(s) + +proc bufferMoreDataSync(s: InputStream): bool = + # This proc exists only to avoid inlining of the code of + # `bufferMoreDataImpl` into `readable` (which in turn is + # a template inlined in the user code). + bufferMoreDataImpl(s, noAwait, readSync) + +template readable*(sp: InputStream): bool = + ## Checks whether reading more data from the stream is possible. + ## + ## If there is any unconsumed data in the stream buffers, the + ## operation returns `true` immediately. You can call `read` + ## or `peek` afterwards to consume or examine the next byte + ## in the stream. + ## + ## If the stream buffers are empty, the operation may block + ## until more data becomes available. The end of the stream + ## may be reached at this point, which will be indicated by + ## a `false` return value. Any attempt to call `read` or + ## `peek` afterwards is considered a `Defect`. + ## + ## Please note that this API is intended for stream consumers + ## who need to consume the data one byte at a time. A typical + ## usage will be the following: + ## + ## ```nim + ## while stream.readable: + ## case stream.peek.char + ## of '"': + ## parseString(stream) + ## of '0'..'9': + ## parseNumber(stream) + ## of '\': + ## discard stream.read # skip the slash + ## let escapedChar = stream.read + ## ``` + ## + ## Even though the user code consumes the data one byte at a time, + ## in the majority of cases this consist of simply incrementing a + ## pointer within the stream buffers. Only when the stream buffers + ## are exhausted, a new read operation will be executed throught + ## the stream input device which may repopulate the buffers with + ## fresh data. See `Stream Pages` for futher discussion of this. + + # This is a template, because we want the pointer check to be + # inlined at the call sites. Only if it fails, we call into the + # larger non-inlined proc: + let s = sp + hasRunway(s.span) or bufferMoreDataSync(s) + +template readable*(sp: AsyncInputStream): bool = + ## Async version of `readable`. + ## The intended API usage is the same. Instead of blocking, an async + ## stream will use `await` while waiting for more data. + let s = sp + if hasRunway(s.span): + true + else: + bufferMoreDataImpl(s, fsAsync, readAsync) + +template readableNImpl(s, n, awaiter, readOp: untyped): bool = + let runway = s.totalUnconsumedBytes + + if runway >= n: + true + elif s.buffers == nil or s.vtable == nil or s.vtable.readOp == nil: + false + else: + var + bytesDeficit = n - runway + targetBytesRead = s.buffers.totalBytesRead + bytesDeficit + res = false + + while true: + awaiter s.vtable.readOp(s) + + if wasEofReached(s.buffers): + s.disconnectInputDevice() + res = s.buffers.totalBytesRead >= targetBytesRead + break + + if s.buffers.totalBytesRead >= targetBytesRead: + res = true + break + + res + +proc readable*(s: InputStream, n: int): bool = + ## Checks whether reading `n` bytes from the input stream is possible. + ## + ## If there is enough unconsumed data in the stream buffers, the + ## operation will return `true` immediately. You can use `read`, + ## `peek`, `read(n)` or `peek(n)` afterwards to consume up to the + ## number of verified bytes. Please note that consuming more bytes + ## will be considered a `Defect`. + ## + ## If the stream buffers do not contain enough data, the operation + ## may block until more data becomes available. The end of the stream + ## may be reached at this point, which will be indicated by a `false` + ## return value. Please note that the stream might still contain some + ## unconsumed bytes after `readable(n)` returned false. You can use + ## `totalUnconsumedBytes` or a combination of `readable` and `read` + ## to consume the remaining bytes if desired. + ## + ## If possible, prefer consuming the data one byte at a time. This + ## ensures the most optimal usage of the stream buffers. Even after + ## calling `readable(n)`, it's still preferrable to continue with + ## `read` instead of `read(n)` because the later may require the + ## resulting bytes to be copied to a freshly allocated sequence. + ## + ## In the situation where the consumed bytes need to be copied to + ## an existing external buffer, `readInto` will provide the best + ## performance instead. + ## + ## Just like `readable`, this operation will invoke reads on the + ## stream input device only when necessary. See `Stream Pages` + ## for futher discussion of this. + readableNImpl(s, n, noAwait, readSync) + +template readable*(sp: AsyncInputStream, np: int): bool = + ## Async version of `readable(n)`. + ## The intended API usage is the same. Instead of blocking, an async + ## stream will use `await` while waiting for more data. + let + s = sp + n = np + + readableNImpl(s, n, fsAwait, readAsync) + +proc peek*(s: InputStream): byte {.inline.} = + doAssert hasRunway(s.span) + return s.span.startAddr[] + +template peek*(s: AsyncInputStream): byte = + peek 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) + return peekHead[] + +template peekAt*(s: AsyncInputStream, pos: int): byte = + peekAt InputStream(s) + +proc advance*(s: InputStream) = + if hasRunway(s.span): + bumpPointer s.span + elif s.buffers != nil and s.buffers.len > 1: + 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 readIntoEx*(s: InputStream, target: var openarray[byte]): int = + ## Read data into the destination buffer. + ## + ## Returns the number of bytes that were successfully + ## written to the buffer. The function will return a + ## number smaller than the buffer length only if EOF + ## was reached before the buffer was fully populated. + discard + +proc readInto*(s: InputStream, target: var openarray[byte]): bool = + ## Read data into the destination buffer. + ## + ## Returns `false` if EOF was reached before the buffer + ## was fully populated. if you need precise information + ## regarding the number of bytes read, see `readIntoEx`. + s.readIntoEx(target) == target.len + +template readInto*(s: AsyncInputStream, target: var openarray[byte]): bool = + ## Asynchronously read data into the destination buffer. + ## + ## Returns `false` if EOF was reached before the buffer + ## was fully populated. if you need precise information + ## regarding the number of bytes read, see `readIntoEx`. + ## + ## If there are enough bytes already buffered by the stream, + ## the expression will complete immediately. + ## Otherwise, it will await more bytes to become available. + discard + +proc checkReadAhead(s: InputStream, n: Natural): ptr byte = + # TODO: handle multi-page + result = s.span.startAddr + doAssert s.span.len >= n + bumpPointer s.span, n + +template read*(s: InputStream, n: Natural): auto = + makeOpenArray(checkReadAhead(s, n), n) + +proc next*(s: InputStream): Option[byte] = + if readable(s): + result = some read(s) + +template next*(sp: AsyncInputStream): Option[byte] = + let s = sp + if readable(s): + some read(s) + else: + none byte + +proc pos*(s: InputStream): int {.inline.} = + s.spanEndPos - s.span.len + +template pos*(s: AsyncInputStream): int = + pos InputStream(s) + +when false: + # Obsolete APIs for removal + proc bufferPos(s: InputStream, pos: int): ptr byte = + let offsetFromEnd = pos - s.spanEndPos + doAssert offsetFromEnd < 0 + result = offset(s.span.endAddr, offsetFromEnd) + doAssert 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 + + proc rewindTo*(s: InputStream, pos: int) {.inline.} = + s.head = s.bufferPos(pos) + diff --git a/faststreams/multisync.nim b/faststreams/multisync.nim index c718cac..6bf1a01 100644 --- a/faststreams/multisync.nim +++ b/faststreams/multisync.nim @@ -1,6 +1,6 @@ import stew/shims/macros, - async_backend, input_stream, output_stream + async_backend, inputs, outputs macro fsMultiSync*(body: untyped) = # We will produce an identical copy of the annotated proc, diff --git a/faststreams/output_stream.nim b/faststreams/output_stream.nim deleted file mode 100644 index e49b512..0000000 --- a/faststreams/output_stream.nim +++ /dev/null @@ -1,540 +0,0 @@ -import - deques, typetraits, - stew/[ptrops, strings, ranges/ptr_arith], - async_backend - -type - OutputPage = object - buffer: string - startOffset: int - - OutputStream* = ref object of RootObj - vtable*: ptr OutputStreamVTable - cursor*: WriteCursor - pages: Deque[OutputPage] - endPos: int - extCursorsCount: int - pageSize*: int - maxWriteSize*: int - minWriteSize*: int - - LayeredOutputStream* = ref object of OutputStream - subStream*: OutputStream - - OutputStreamHandle* = object - s*: OutputStream - - AsyncOutputStream* {.borrow: `.`.} = distinct OutputStream - - WritePageSyncProc* = proc (s: OutputStream, page: openarray[byte]) - {.nimcall, gcsafe, raises: [IOError, Defect].} - - WritePageAsyncProc* = proc (s: OutputStream, buf: pointer, bufLen: int): Future[void] - {.nimcall, gcsafe, raises: [IOError, Defect].} - - FlushSyncProc* = proc (s: OutputStream) - {.nimcall, gcsafe, raises: [IOError, Defect].} - - FlushAsyncProc* = proc (s: OutputStream): Future[void] - {.nimcall, gcsafe, raises: [IOError, Defect].} - - CloseSyncProc* = proc (s: OutputStream) - {.nimcall, gcsafe, raises: [IOError, Defect].} - - CloseAsyncProc* = proc (s: OutputStream): Future[void] - {.nimcall, gcsafe, raises: [IOError, Defect].} - - OutputStreamVTable* = object - writePageSync*: WritePageSyncProc - writePageAsync*: WritePageAsyncProc - flushSync*: FlushSyncProc - flushAsync*: FlushAsyncProc - closeSync*: CloseSyncProc - closeAsyncProc*: CloseAsyncProc - - WriteCursor* = object - head, bufferEnd: ptr byte - stream: OutputStream - - VarSizeWriteCursor* = distinct WriteCursor - - FileOutputStream = ref object of OutputStream - file: File - -const - nimAllocatorMetadataSize* = 0 - # TODO: Get this from Nim's allocator. - # The goal is to make perfect page-aligned allocations - defaultPageSize = 4096 - nimAllocatorMetadataSize - 1 # 1 byte for the null terminator - -proc close*(s: OutputStream) {.raises: [IOError, Defect].} = - if s != nil: - if s.vtable != nil and s.vtable.closeSync != nil: - s.vtable.closeSync(s) - -# TODO -# The destructors are currently disabled because they seem to cause -# mysterious segmentation faults related to corrupted GC internal -# data structures. -#[ -proc `=destroy`*(h: var OutputStreamHandle) {.raises: [Defect].} = - if h.s != nil: - if h.s.vtable != nil and h.s.vtable.closeSync != nil: - try: - h.s.vtable.closeSync(h.s) - except IOError: - # Since this is a destructor, there is not much we can do here. - # If the user wanted to handle the error, they would have called - # `close` manually. - discard # TODO - # h.s = nil -]# - -converter implicitDeref*(h: OutputStreamHandle): OutputStream = - h.s - -template canExtendOutput(s: OutputStream): bool = - # Streams writing to pre-allocated existing buffers cannot be grown - s != nil and s.pageSize > 0 - -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) - -func runway*(c: var WriteCursor): int {.inline.} = - distance(c.head, c.bufferEnd) - -proc prepareRunway*(s: OutputStream, length: int) = - # TODO implement this - discard - -template prepareRunway*(s: AsyncOutputStream, length: int) = - prepareRunway OutputStream(s) - -proc flipPage(s: OutputStream) = - s.cursor.head = cast[ptr byte](addr s.pages[s.pages.len - 1].buffer[0]) - # TODO: There is an assumption here and elsewhere that `s.pages[^1]` has - # a length equal to `s.pageSize` - s.cursor.bufferEnd = cast[ptr byte](offset(s.cursor.head, s.pageSize)) - s.endPos += s.pageSize - -proc addPage(s: OutputStream) = - doAssert s.pageSize > 0 - s.pages.addLast OutputPage(buffer: newString(s.pageSize), - startOffset: 0) - s.flipPage - -proc initWithSinglePage*(s: OutputStream) = - s.pages = initDeque[OutputPage]() - s.addPage() - s.cursor.stream = s - -proc memoryOutput*(pageSize = defaultPageSize): OutputStreamHandle = - var stream = OutputStream( - pageSize: pageSize, - minWriteSize: 1, - maxWriteSize: high(int)) - - stream.initWithSinglePage() - - OutputStreamHandle(s: stream) - -proc memoryOutput*(buffer: pointer, len: int): OutputStreamHandle = - let buffer = cast[ptr byte](buffer) - - var stream = OutputStream() - stream.cursor.head = buffer - stream.cursor.bufferEnd = offset(buffer, len) - stream.cursor.stream = stream - stream.endPos = len - - OutputStreamHandle(s: stream) - -let FileStreamVTable = OutputStreamVTable( - writePageSync: proc (s: OutputStream, data: openarray[byte]) - {.nimcall, gcsafe, raises: [IOError, Defect].} = - var file = FileOutputStream(s).file - var written = file.writeBuffer(unsafeAddr data[0], data.len) - if written != data.len: - raise newException(IOError, "Failed to write OutputStream page.") - , - flushSync: proc (s: OutputStream) - {.nimcall, gcsafe, raises: [IOError, Defect].} = - flushFile FileOutputStream(s).file - , - closeSync: proc (s: OutputStream) - {.nimcall, gcsafe, raises: [IOError, Defect].} = - close FileOutputStream(s).file -) - -template vtableAddr*(vtable: OutputStreamVTable): ptr OutputStreamVTable = - ## This is a simple work-around for the somewhat broken side - ## effects analysis of Nim - reading from global let variables - ## is considered a side-effect. - {.noSideEffect.}: - unsafeAddr vtable - -proc fileOutput*(filename: string, - fileMode: FileMode = fmWrite, - pageSize = defaultPageSize): OutputStreamHandle {. - raises: [IOError, Defect] -.} = - let f = open(filename, fileMode) - - var stream = FileOutputStream( - vtable: vtableAddr FileStreamVTable, - pageSize: pageSize, - minWriteSize: 1, - maxWriteSize: high(int), - file: f) - - stream.initWithSinglePage() - - OutputStreamHandle(s: stream) - -proc pos*(s: OutputStream): int = - s.endPos - s.cursor.runway - -proc safeWritePage(s: OutputStream, data: openarray[byte]) {.inline.} = - if data.len > 0: s.vtable.writePageSync(s, data) - -proc writePages(s: OutputStream, skipLast = 0) = - assert s.vtable != nil - for i in 0 ..< s.pages.len - skipLast: - s.safeWritePage s.pages[i].buffer.toOpenArrayByte(0, s.pages[i].buffer.len - 1) - -proc writePartialPage(s: OutputStream, page: var OutputPage) = - assert s.vtable != nil - let - unwrittenBytes = s.cursor.runway - pageEndPos = s.pageSize - unwrittenBytes - 1 - pageStartPos = page.startOffset - - s.safeWritePage page.buffer.toOpenArrayByte(pageStartPos, pageEndPos) - s.endPos -= unwrittenBytes - - page.startOffset = 0 - s.flipPage - -proc flush*(s: OutputStream) = - doAssert s.extCursorsCount == 0 - if s.vtable != nil: - # We write all pages except the last one - s.writePages(skipLast = 1) - # Then we erase them from the list - s.pages.shrink(fromFirst = s.pages.len - 1) - # Then we write the current page, which is probably incomplete - s.writePartialPage s.pages[0] - # Finally, we flush - s.vtable.flushSync(s) - -proc writePendingPagesAndLeaveOne(s: OutputStream) {.inline.} = - s.writePages - s.pages.shrink(fromFirst = s.pages.len - 1) - s.pages[0].startOffset = 0 - s.flipPage - -proc tryFlushing(s: OutputStream) {.inline.} = - # Pre-conditions: - # * The cursor has reached the current buffer end - # - # Post-conditions: - # * All completed pages are written - # * There is a fresh page ready for writing at the top - # (we can reuse a previously existing page for this) - # * The head and bufferEnd pointers point to the new top page - if s.vtable != nil and s.extCursorsCount == 0: - s.writePendingPagesAndLeaveOne - else: - s.addPage - -func endAddr(s: string): ptr byte {.inline.} = - let a = unsafeAddr s[0] - offset(cast[ptr byte](a), s.len) - -template startAddr(s: string): ptr byte = - cast[ptr byte](unsafeAddr s[0]) - -func boundingAddrs(s: string): (ptr byte, ptr byte) {.inline.} = - (startAddr s, endAddr s) - -proc findNextPage(c: var WriteCursor): int = - let cursorBufferEnd = c.bufferEnd - for i in 0 .. c.stream.pages.len - 2: - let pageEnd = endAddr c.stream.pages[i].buffer - if cursorBufferEnd == pageEnd: - return i + 1 - - doAssert false # There is no next page the cursor can move to - -proc moveToPage(c: var WriteCursor, p: var OutputPage) = - doAssert p.startOffset > 0 - c.head = cast[ptr byte](unsafeAddr p.buffer[0]) - c.bufferEnd = offset(c.head, p.startOffset) - p.startOffset = 0 - -proc moveToNextPage(c: var WriteCursor) = - c.moveToPage c.stream.pages[c.findNextPage()] - -proc append*(c: var WriteCursor, b: byte) = - if c.head == c.bufferEnd: - doAssert c.stream.canExtendOutput - if c.isExternalCursor: - c.moveToNextPage() - else: - c.stream.tryFlushing() - - c.head[] = b - c.head = offset(c.head, 1) - -template append*(c: var WriteCursor, x: char) = - bind append - c.append byte(x) - -proc writeDataAsPages(s: OutputStream, data: ptr byte, dataLen: int) = - var - data = data - dataLen = dataLen - - if dataLen > s.pageSize: - if dataLen < s.maxWriteSize: - s.vtable.writePageSync(s, makeOpenArray(data, dataLen)) - s.endPos += dataLen - return - - while dataLen > s.pageSize: - s.vtable.writePageSync(s, makeOpenArray(data, s.pageSize)) - data = offset(data, s.pageSize) - dec dataLen, s.pageSize - s.endPos += s.pageSize - - copyMem(s.cursor.head, data, dataLen) - s.cursor.head = offset(s.cursor.head, dataLen) - -proc newStringFromBytes(input: ptr byte, inputLen: int): string = - assert inputLen > 0 - result = newString(inputLen) - copyMem(addr result[0], input, inputLen) - -proc handleLongAppend*(c: var WriteCursor, bytes: openarray[byte]) = - var - pageRemaining = c.runway - inputPos = unsafeAddr bytes[0] - inputLen = bytes.len - stream = c.stream - - template reduceInput(delta: int) = - inputPos = offset(inputPos, delta) - inputLen -= delta - - # Since the input is longer, we first make sure that the top-most - # page is filled to the top: - doAssert c.stream.canExtendOutput - copyMem(c.head, inputPos, pageRemaining) - reduceInput pageRemaining - - if c.isExternalCursor: - var - totalPages = stream.pages.len - nextPageIdx = c.findNextPage - - while nextPageIdx < totalPages: - let - pageStart = startAddr stream.pages[nextPageIdx].buffer - pageRunway = stream.pages[nextPageIdx].startOffset - pageLen = stream.pageSize - - doAssert pageRunway > 0 - stream.pages[nextPageIdx].startOffset = 0 - - if pageRunway < pageLen: - doAssert inputLen <= pageRunway - copyMem(pageStart, inputPos, inputLen) - c.head = offset(pageStart, inputLen) - c.bufferEnd = offset(pageStart, pageRunway) - return - else: - if inputLen <= pageLen: - copyMem(pageStart, inputPos, inputLen) - c.head = offset(pageStart, inputLen) - c.bufferEnd = offset(pageStart, pageLen) - return - else: - copyMem(pageStart, inputPos, pageLen) - reduceInput pageLen - inc nextPageIdx - - doAssert false # If we reached here, this means that we've ran out - # of pages, so this is a write past the end of the - # pre-allocated space for the delayed write. - - elif c.stream.vtable != nil and c.stream.extCursorsCount == 0: - # This stream has an output device and we are ready to flush - # all the pending pages. One fresh page will be left on top. - # The input is yet to be written: - stream.writePendingPagesAndLeaveOne - # This will directly send our input to the output device. - # Since the output device has a preference for pageSize and - # maxWriteSize, we'll send some full pages to it and then - # some bytes will be written to the fresh page created above: - stream.writeDataAsPages(inputPos, inputLen) - else: - # We are not ready to flush, so we must create pending pages. - # We'll try to create them as large as possible: - let maxPageSize = c.stream.maxWriteSize - - # We know how much the endPos will advance, but please note that - # it may be corrected later if we end up writing a portion of the - # input to an incomplete page: - stream.endPos += inputLen - - # Try to create big pages until we have more data: - while inputLen > maxPageSize: - stream.pages.addLast OutputPage( - buffer: newStringFromBytes(inputPos, maxPageSize), - startOffset: 0) - reduceInput maxPageSize - - # Here the remaining input is smaller than a max page, but it may be - # still larger than a regular page. If this is the case, we just create - # one final oversized page and then we leave one empty fresh page where - # the writing will continue: - if inputLen > c.stream.pageSize: - stream.pages.addLast OutputPage( - buffer: newStringFromBytes(inputPos, inputLen), - startOffset: 0) - stream.addPage - else: - # We don't have enough remaining bytes for a full page, so we'll just - # allocate a new empty page and we'll write the remaining input there. - # This will also reset the cursor to the start of the page: - stream.addPage - copyMem(c.head, inputPos, inputLen) - c.head = offset(c.head, inputLen) - # We must correct endPos, because it must mark the end of the top-most - # page. Since `addPage` advances the endPos as well and our remaining - # input was written to the newly created page, our initial increase of - # endPos was overestimated: - stream.endPos -= inputLen - -proc append*(c: var WriteCursor, bytes: openarray[byte]) {.inline.} = - # We have a short inlinable function handling the case when the input is - # short enough to fit in the current page. We'll keep buffering until the - # page is full: - let - pageRemaining = c.runway - inputLen = bytes.len - - if inputLen == 0: return - if inputLen <= pageRemaining: - copyMem(c.head, unsafeAddr bytes[0], inputLen) - c.head = offset(c.head, inputLen) - else: - handleLongAppend(c, bytes) - -proc append*(c: var WriteCursor, chars: openarray[char]) {.inline.} = - var charsStart = unsafeAddr chars[0] - c.append makeOpenArray(cast[ptr byte](charsStart), chars.len) - -template appendMemCopy*[T](c: var WriteCursor, value: T) = - bind append - static: - type TT = T # TODO This deals with a Nim bug - assert supportsCopyMem(TT) - let valueAddr = unsafeAddr value - c.append makeOpenArray(cast[ptr byte](valueAddr), sizeof(value)) - -template append*(c: var WriteCursor, str: string) = - bind append - append c, str.toOpenArrayByte(0, str.len - 1) - -template append*(s: OutputStream, value: auto) = - bind append - append s.cursor, value - -template appendMemCopy*(s: OutputStream, value: auto) = - bind append - append s.cursor, value - -proc getOutput*(s: OutputStream, T: type string): string = - doAssert s.vtable == nil and s.extCursorsCount == 0 and s.pageSize > 0 - - s.pages[s.pages.len - 1].buffer.setLen(s.pageSize - s.cursor.runway) - - if s.pages.len == 1 and s.pages[0].startOffset == 0: - result.swap s.pages[0].buffer - else: - result = newStringOfCap(s.pos) - for page in items(s.pages): - result.add page.buffer.toOpenArray(page.startOffset.int, - page.buffer.len - 1) - -template getOutput*(s: OutputStream, T: type seq[byte]): seq[byte] = - cast[seq[byte]](s.getOutput(string)) - -template getOutput*(s: OutputStream): seq[byte] = - cast[seq[byte]](s.getOutput(string)) - -proc finishPageEarly(s: OutputStream, unwrittenBytes: int) {.inline.} = - s.pages[s.pages.len - 1].buffer.setLen(s.pageSize - unwrittenBytes) - s.endPos -= unwrittenBytes - s.tryFlushing() - -proc createCursor(s: OutputStream, size: int): WriteCursor = - inc s.extCursorsCount - - result = WriteCursor(head: s.cursor.head, - bufferEnd: offset(s.cursor.head, size), - stream: s) - - s.cursor.head = result.bufferEnd - -proc delayFixedSizeWrite*(s: OutputStream, size: Natural): WriteCursor = - let remainingBytesInPage = s.cursor.runway - if size <= remainingBytesInPage: - result = createCursor(s, size) - else: - result = createCursor(s, remainingBytesInPage) - var size = size - remainingBytesInPage - s.endPos += size - while size > s.pageSize: - s.pages.addLast OutputPage(buffer: newString(s.pageSize), - startOffset: s.pageSize) - size -= s.pageSize - - s.pages.addLast OutputPage(buffer: newString(s.pageSize), - startOffset: size) - - let (pageStart, pageEnd) = boundingAddrs s.pages[s.pages.len - 1].buffer - s.cursor.head = offset(pageStart, size) - s.cursor.bufferEnd = pageEnd - s.endPos += (s.pageSize - size) - -proc delayVarSizeWrite*(s: OutputStream, maxSize: Natural): VarSizeWriteCursor = - doAssert maxSize < s.pageSize - s.finishPageEarly s.cursor.runway - VarSizeWriteCursor createCursor(s, maxSize) - -proc finalize*(cursor: var WriteCursor) = - doAssert cursor.stream.extCursorsCount > 0 - dec cursor.stream.extCursorsCount - -proc writeAndFinalize*(cursor: var WriteCursor, data: openarray[byte]) = - doAssert data.len == cursor.runway - copyMem(cursor.head, unsafeAddr data[0], data.len) - finalize cursor - -proc writeAndFinalize*(c: var VarSizeWriteCursor, data: openarray[byte]) = - template cursor: auto = WriteCursor(c) - - for page in mitems(cursor.stream.pages): - if unsafeAddr(page.buffer[0]) == cursor.head: - let overestimatedBytes = cursor.runway - data.len - doAssert overestimatedBytes >= 0 - page.startOffset = overestimatedBytes - copyMem(offset(cursor.head, overestimatedBytes), unsafeAddr data[0], data.len) - finalize cursor - return - - doAssert false - diff --git a/faststreams/outputs.nim b/faststreams/outputs.nim new file mode 100644 index 0000000..213ef1f --- /dev/null +++ b/faststreams/outputs.nim @@ -0,0 +1,685 @@ +## Please note that the use of unbuffered streams comes with a number +## of restrictions: +## +## * Delayed writes are not supported. +## * Output consuming operations such as `getOutput`, `consumeOutputs` and +## `consumeContiguousOutput` should not be used with them. +## * They cannot participate as intermediate steps in pipelines. + +import + deques, typetraits, + stew/[ptrops, strings, ranges/ptr_arith], + buffers, async_backend + +export + CloseBehavior + +type + OutputStream* = ref object of RootObj + vtable*: ptr OutputStreamVTable # This is nil for any memory output + buffers*: PageBuffers # This is nil for unsafe memory outputs + span: PageSpan + spanEndPos: Natural + extCursorsCount: int + closeFut: Future[void] + + WriteCursor* = object + span: PageSpan + stream: OutputStream + + LayeredOutputStream* = ref object of OutputStream + subStream*: OutputStream + + OutputStreamHandle* = object + s*: OutputStream + + AsyncOutputStream* {.borrow: `.`.} = distinct OutputStream + + WriteSyncProc* = proc (s: OutputStream, buf: pointer, bufLen: Natural) + {.nimcall, gcsafe, raises: [IOError, Defect].} + + WriteAsyncProc* = proc (s: OutputStream, buf: pointer, bufLen: Natural): Future[void] + {.nimcall, gcsafe, raises: [IOError, Defect].} + + FlushSyncProc* = proc (s: OutputStream) + {.nimcall, gcsafe, raises: [IOError, Defect].} + + FlushAsyncProc* = proc (s: OutputStream): Future[void] + {.nimcall, gcsafe, raises: [IOError, Defect].} + + CloseSyncProc* = proc (s: OutputStream) + {.nimcall, gcsafe, raises: [IOError, Defect].} + + CloseAsyncProc* = proc (s: OutputStream): Future[void] + {.nimcall, gcsafe, raises: [IOError, Defect].} + + OutputStreamVTable* = object + writeSync*: WriteSyncProc + writeAsync*: WriteAsyncProc + flushSync*: FlushSyncProc + flushAsync*: FlushAsyncProc + closeSync*: CloseSyncProc + closeAsync*: CloseAsyncProc + + VarSizeWriteCursor* = distinct WriteCursor + + FileOutputStream = ref object of OutputStream + file: File + +const + nimAllocatorMetadataSize* = 0 + # TODO: Get this from Nim's allocator. + # The goal is to make perfect page-aligned allocations + +proc disconnectOutputDevice(s: OutputStream) = + if s.vtable != nil: + if s.vtable.closeAsync != nil: + s.closeFut = s.vtable.closeAsync(s) + elif s.vtable.closeSync != nil: + s.vtable.closeSync(s) + s.vtable = nil + +template disconnectOutputDevice(s: AsyncOutputStream) = + disconnectOutputDevice OutputStream(s) + +proc close*(s: OutputStream, + behavior = dontWaitAsyncClose) + {.raises: [IOError, Defect].} = + disconnectOutputDevice(s) + if s.closeFut != nil: + fsTranslateErrors "Stream closing failed": + if behavior == waitAsyncClose: + waitFor s.closeFut + else: + asyncCheck s.closeFut + +proc close*(s: AsyncOutputStream): Future[void] + {.raises: [IOError, Defect].} = + disconnectOutputDevice(s) + result = OutputStream(s).closeFut + doAssert result != nil + +template closeNoWait*(sp: AsyncOutputStream|OutputStream) = + ## Close the stream without waiting even if's async. + ## This operation will use `asyncCheck` internally to detect unhandled + ## errors from the closing operation. + close(InputStream(s), dontWaitAsyncClose) + +# TODO +# The destructors are currently disabled because they seem to cause +# mysterious segmentation faults related to corrupted GC internal +# data structures. +#[ +proc `=destroy`*(h: var OutputStreamHandle) {.raises: [Defect].} = + if h.s != nil: + if h.s.vtable != nil and h.s.vtable.closeSync != nil: + try: + h.s.vtable.closeSync(h.s) + except IOError: + # Since this is a destructor, there is not much we can do here. + # If the user wanted to handle the error, they would have called + # `close` manually. + discard # TODO + # h.s = nil +]# + +converter implicitDeref*(h: OutputStreamHandle): OutputStream = + h.s + +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) = + s.span = s.buffers.addWritablePage().writableSpan + s.spanEndPos += s.span.len + +template makeHandle*(sp: OutputStream): OutputStreamHandle = + let s = sp + OutputStreamHandle(s: s) + +proc memoryOutput*(pageSize = defaultPageSize): OutputStreamHandle = + doAssert pageSize > 0 + makeHandle OutputStream(buffers: initPageBuffers(pageSize)) + +proc unsafeMemoryOutput*(buffer: pointer, len: Natural): OutputStreamHandle = + let buffer = cast[ptr byte](buffer) + + makeHandle OutputStream( + span: PageSpan(startAddr: buffer, endAddr: offset(buffer, len)), + spanEndPos: len) + +proc ensureRunway*(s: OutputStream, neededRunway: Natural) = + ## The hint provided in `ensureRunway` overrides any previous + ## hint specified at stream creation with `pageSize`. + let runway = s.span.len + + # This is a temporary requirement. + # ensureRunway should be called immediately after creating the OutputStream + # In the future, we'll relax this by implementing more logic in buffers.nim + doAssert runway == 0, "call ensureRunway immediately after stream creation" + + if neededRunway > runway: + # 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" + s.span = s.buffers.ensureRunway(neededRunway - runway) + +template ensureRunway*(s: AsyncOutputStream, neededRunway: Natural) = + ensureRunway OutputStream(s, neededRunway) + +let FileOutputVTable = OutputStreamVTable( + writeSync: proc (s: OutputStream, buf: pointer, bufLen: Natural) + {.nimcall, gcsafe, raises: [IOError, Defect].} = + var file = FileOutputStream(s).file + + template fail = + raise newException(IOError, "Failed to write OutputStream page.") + + if s.buffers != nil: + s.buffers.consumeAllPages(pageAddr, pageLen): + let written = file.writeBuffer(pageAddr, pageLen) + if written != pageLen: fail() + + if bufLen > 0: + doAssert buf != nil + var written = file.writeBuffer(buf, bufLen) + if written != bufLen: fail() + , + flushSync: proc (s: OutputStream) + {.nimcall, gcsafe, raises: [IOError, Defect].} = + flushFile FileOutputStream(s).file + , + closeSync: proc (s: OutputStream) + {.nimcall, gcsafe, raises: [IOError, Defect].} = + close FileOutputStream(s).file +) + +template vtableAddr*(vtable: OutputStreamVTable): ptr OutputStreamVTable = + ## This is a simple work-around for the somewhat broken side + ## effects analysis of Nim - reading from global let variables + ## is considered a side-effect. + {.noSideEffect.}: + unsafeAddr vtable + +proc fileOutput*(filename: string, + fileMode: FileMode = fmWrite, + pageSize = defaultPageSize): OutputStreamHandle + {.raises: [IOError, Defect].} = + let f = open(filename, fileMode) + + makeHandle FileOutputStream( + vtable: vtableAddr FileOutputVTable, + buffers: initPageBuffers(pageSize), + file: f) + +proc pos*(s: OutputStream): int = + s.spanEndPos - s.span.len + +template pos*(s: AsyncOutputStream): int = + pos OutputStream(s) + +# +# Pre-conditions for `drainAllBuffers(Sync/Async)` +# * The cursor has reached the current span end +# * We are working with a vtable-enabled stream +# +# Post-conditions: +# * All completed pages are written +# * There is a fresh page ready for writing at the top +# (we can reuse a previously existing page for this) +# * The stream cursor is re-initialized at the start of the top page +# +proc drainAllBuffersSync(s: OutputStream, buf: pointer, bufSize: Natural) = + s.vtable.writeSync(s, buf, bufSize) + if s.buffers != nil: + s.span = s.buffers.getWritableSpan() + s.spanEndPos += s.span.len + +proc drainAllBuffersAsync(s: OutputStream, buf: pointer, bufSize: Natural) {.async.} = + await s.vtable.writeAsync(s, buf, bufSize) + s.span = s.buffers.getWritableSpan() + s.spanEndPos += s.span.len + +proc createCursor(s: OutputStream, size: int): WriteCursor = + inc s.extCursorsCount + + let + # The start address matches the current stream main cursor location + startAddr = s.span.startAddr + endAddr = offset(startAddr, size) + + result = WriteCursor( + stream: s, + span: PageSpan(startAddr: startAddr, endAddr: endAddr)) + + # Adjust the stream main cursor to point past the end + # of the newly created cursor: + s.span.startAddr = endAddr + +proc delayFixedSizeWrite*(s: OutputStream, size: Natural): WriteCursor = + let runway = s.span.len + if size <= runway: + result = createCursor(s, size) + else: + result = createCursor(s, runway) + + let + runwayDeficit = size - runway + nextPageSize = nextAlignedSize(runwayDeficit, s.buffers.pageSize) + nextPage = s.buffers.addWritablePage(nextPageSize) + nextPageSpan = nextPage.writableSpan + + s.span = PageSpan(startAddr: offset(nextPageSpan.startAddr, runwayDeficit), + endAddr: nextPageSpan.endAddr) + + # See the explanation about split cursors above + nextPage.startOffset = -runwayDeficit + + s.spanEndPos += nextPageSize + +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 + + let runway = s.span.len + if maxSize <= runway: + let + startAddr = s.span.startAddr + endAddr = offset(startAddr, maxSize) + + result = VarSizeWriteCursor WriteCursor( + stream: s, + span: PageSpan(startAddr: startAddr, endAddr: endAddr)) + + s.buffers.splitLastPageAt(endAddr) + s.span.startAddr = endAddr + + else: + s.buffers.endLastPageAt(s.span.startAddr) + let + nextPageSize = nextAlignedSize(maxSize, s.buffers.pageSize) + nextPageSpan = s.buffers.addWritablePage(nextPageSize).writableSpan + cursorEndAddr = offset(nextPageSpan.startAddr, maxSize) + + result = VarSizeWriteCursor WriteCursor( + stream: s, + span: PageSpan(startAddr: nextPageSpan.startAddr, + endAddr: cursorEndAddr)) + + s.span = PageSpan(startAddr: cursorEndAddr, + endAddr: nextPageSpan.endAddr) + s.spanEndPos += nextPageSize + +proc finalize*(cursor: var WriteCursor) = + doAssert cursor.stream.extCursorsCount > 0 + dec cursor.stream.extCursorsCount + +proc finalWrite*(cursor: var WriteCursor, data: openArray[byte]) = + doAssert data.len == cursor.span.len + copyMem(cursor.span.startAddr, unsafeAddr data[0], data.len) + finalize cursor + +proc finalWrite*(c: var VarSizeWriteCursor, data: openArray[byte]) = + template cursor: auto = WriteCursor(c) + + let overestimatedBytes = cursor.span.len - data.len + doAssert overestimatedBytes >= 0 + + for page in items(cursor.stream.buffers.queue): + let baseAddr = page.pageBaseAddr + if page.pageEndAddr == cursor.span.endAddr: + # This is a page ending cursor + page.endOffset = distance(baseAddr, cursor.span.startAddr) + data.len + copyMem(cursor.span.startAddr, unsafeAddr data[0], data.len) + finalize cursor + return + + if cursor.span.startAddr == baseAddr: + # This is page starting cursor + page.startOffset = overestimatedBytes + copyMem(offset(baseAddr, overestimatedBytes), unsafeAddr data[0], data.len) + finalize cursor + return + + doAssert false + +proc tryMovingToNextPage(c: var WriteCursor) = + # A split cursor is a fixed-size cursor that ended up on page boundary. + # + # Part of the cursor used the last few bytes of the first page and we've + # left some empty space at the beginning of the second page. + # + # Even if the cursor size was very large, we've made sure the next + # page is big enough to hold all the data. When we created the cursor, + # we've taken a note regarding the number of bytes on the second page + # that are reserved by writing them as a negative value for the page + # `startOffset`. + # + # All we need to do here is update the cursor span to point to the next + # page and set the now final `endAddr`. The page `startOffset` is updated + # to 0 to indicate that the cursor has made the flip. + # + # If you are wondering, var-sized cursors cannot be split, because our + # strategy is to always place them at the beggining or end of pages. + # + # When we try to create a var-sized cursor, we check if there are enough + # bytes on the current page to contain the worst case scenario (the var + # sized cursor has an upper size limit). If there are enough bytes, we + # end the page prematurely (it will end up with an `endOffset`). We can + # then recycle the same memory for the next page that will use an adjusted + # `startOffset`. The `endOffset` of the first page will be written when + # the cursor is finalized and its final size becomes known. + # + # If there weren't enough bytes (a much more rare event), we allocate a + # new page. We adjust the `endOffset` of the current page to mark it's + # premature end and we mark the cursor as special by writing a + + # The split cursor is definetely not on the last page, so we can iterate + # only over the preceeding pages to find where it was: + var prevPage = c.stream.buffers.queue[0] + for i in 1 ..< c.stream.buffers.queue.len: + let page = c.stream.buffers.queue[i] + if c.span.endAddr == prevPage.pageEndAddr and page.startOffset < 0: + # We found what we need, so let's get to business: + c.span.startAddr = page.pageBaseAddr + c.span.endAddr = offset(c.span.startAddr, -page.startOffset) + page.startOffset = 0 + return + prevPage = page + + # 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" + +template flushImpl(s: OutputStream, awaiter, writeOp, flushOp: untyped) = + doAssert s.extCursorsCount == 0 + if s.vtable != nil: + if s.buffers != nil: + s.buffers.endLastPageAt s.span.startAddr + awaiter s.vtable.writeOp(s, nil, 0) + s.span = s.buffers.getWritableSpan() + s.spanEndPos += s.span.len + + if s.vtable.flushOp != nil: + awaiter s.vtable.flushOp(s) + +proc flush*(s: OutputStream) = + flushImpl(s, noAwait, writeSync, flushSync) + +template flush*(s: AsyncOutputStream) = + let s = sp + flushImpl(s, fsAwait, writeAsync, flushAsync) + +template writeByteImpl(s: OutputStream, b: byte, awaiter, writeOp, drainOp: untyped) = + if s.span.atEnd: + # Unsafe memory outputs don't use pages at all, so if our cursor + # reached here, this is a range violation defect: + doAssert canExtendOutput(s) + + if s.vtable == nil or s.extCursorsCount > 0: + # This is the main cursor of a stream, but we are either not + # ready to flush due to outstanding delayed writes or this is + # just a memory output stream. In both cases, we just need to + # allocate more memory and continue writing: + addPage(s) + elif s.buffers == nil: + awaiter s.vtable.writeOp(nil, unsafeAddr b, 1) + else: + awaiter drainOp(s, nil, 0) + + writeByte(s.span, b) + +proc write*(c: var WriteCursor, b: byte) = + if c.span.atEnd: + # The cursor has reached the end of its buffer, but it may be a + # split cursor. If that's the case, the following function will + # succeed. If that's not a split cursor, we'll raise a Defect. + tryMovingToNextPage(c) + + writeByte(c.span, b) + +proc write*(s: OutputStream, b: byte) = + writeByteImpl(s, b, noAwait, writeSync, drainAllBuffersSync) + +template write*(s: AsyncOutputStream, b: byte) = + # TODO: I should do something with the write async Futures + bind write + write OutputStream(s) + +template writeAndWait*(sp: AsyncOutputStream, b: byte) = + let s = sp + writeByteImpl(s, b, fsAwait, writeAsync, drainAllBuffersAsync) + +template write*(s: OutputStream|AsyncOutputStream|var WriteCursor, x: char) = + bind write + write s, byte(x) + +proc writeToANewPage(s: OutputStream, bytes: openArray[byte]) = + var + runway = s.span.len + inputPos = unsafeAddr bytes[0] + inputLen = bytes.len + + template reduceInput(delta: int) = + inputPos = offset(inputPos, delta) + inputLen -= delta + + if runway > 0: + copyMem(s.span.startAddr, inputPos, runway) + reduceInput runway + + doAssert s.buffers != nil + + let nextPageSize = nextAlignedSize(inputLen, s.buffers.pageSize) + let nextPage = s.buffers.addWritablePage(nextPageSize) + + s.span = nextPage.writableSpan + s.spanEndPos += s.span.len + + copyMem(s.span.startAddr, inputPos, inputLen) + s.span.startAddr = offset(s.span.startAddr, inputLen) + +template writeBytesImpl(s: OutputStream, + bytes: openArray[byte], + drainOp: untyped) = + let inputLen = bytes.len + if inputLen == 0: return + + # We have a short inlinable function handling the case when the input is + # short enough to fit in the current page. We'll keep buffering until the + # page is full: + let runway = s.span.len + if inputLen <= runway: + copyMem(s.span.startAddr, unsafeAddr bytes[0], inputLen) + s.span.startAddr = offset(s.span.startAddr, inputLen) + elif s.vtable == nil or s.extCursorsCount > 0: + # We are not ready to flush, so we must create pending pages. + # We'll try to create them as large as possible: + s.writeToANewPage(bytes) + else: + s.buffers.endLastPageAt(s.span.startAddr) + drainOp + +proc write*(s: OutputStream, bytes: openArray[byte]) = + writeBytesImpl(s, bytes): + drainAllBuffersSync(s, unsafeAddr bytes[0], bytes.len) + +proc write*(s: OutputStream, chars: openArray[char]) = + write s, charsToBytes(chars) + +proc write*(s: OutputStream, value: string) {.inline.} = + write s, value.toOpenArrayByte(0, value.len - 1) + +template memCopyToBytes(value: auto): untyped = + type T = type(value) + static: assert supportsCopyMem(T) + let valueAddr = unsafeAddr value + makeOpenArray(cast[ptr byte](valueAddr), sizeof(T)) + +proc writeMemCopy*(s: OutputStream, value: auto) = + bind write + write s, memCopyToBytes(value) + +proc writeBytesAsyncImpl(sp: AsyncOutputStream, + bytes: openarray[byte]): Future[void] = + let s = OutputStream(sp) + writeBytesImpl(s, bytes): + return s.vtable.writeAsync(s, unsafeAddr bytes[0], bytes.len) + +proc writeBytesAsyncImpl(s: AsyncOutputStream, + chars: openarray[char]): Future[void] = + writeBytesAsyncImpl s, charsToBytes(chars) + +proc writeBytesAsyncImpl(s: AsyncOutputStream, + str: string): Future[void] = + writeBytesAsyncImpl s, toOpenArray(str, 0, str.len - 1) + +template writeAndWait*(sp: AsyncOutputStream, value: auto) = + bind writeBytesAsyncImpl + + let + s = sp + f = writeBytesAsyncImpl(s, value) + + if f != nil: + fsAwait(f) + s.span = s.buffers.getWritableSpan() + s.spanEndPos += s.span.len + +template writeMemCopyAndWait*(sp: AsyncOutputStream, value: auto) = + writeAndWait(sp, memCopyToBytes(value)) + +proc writeBytesToCursor(c: var WriteCursor, bytes: openarray[byte]) = + var + runway = c.span.len + inputPos = unsafeAddr bytes[0] + inputLen = bytes.len + + template reduceInput(delta: int) = + inputPos = offset(inputPos, delta) + inputLen -= delta + + if inputLen <= runway: + copyMem(c.span.startAddr, inputPos, inputLen) + c.span.startAddr = offset(c.span.startAddr, inputLen) + else: + # This must be a split cursor. We need to complete its first page first, + # then switch to the second and continue the write there. + copyMem(c.span.startAddr, unsafeAddr bytes[0], runway) + reduceInput runway + # If this really is a split cursor, the following operation will succeed. + # Otherwise, it will Defect and the conclusion is that this was a write + # past the cursor end. + c.tryMovingToNextPage() + # 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 + copyMem(c.span.startAddr, inputPos, inputLen) + c.span.startAddr = offset(c.span.startAddr, inputLen) + +template write*(c: var WriteCursor, bytes: openarray[byte]) = + bind writeBytesToCursor + writeBytesToCursor(c, bytes) + +proc write*(c: var WriteCursor, chars: openarray[char]) {.inline.} = + var charsStart = unsafeAddr chars[0] + writeBytesToCursor(c, makeOpenArray(cast[ptr byte](charsStart), chars.len)) + +proc writeMemCopy*[T](c: var WriteCursor, value: T) = + bind writeBytesToCursor + writeBytesToCursor(c, memCopyToBytes(value)) + +proc write*(c: var WriteCursor, str: string) = + writeBytesToCursor(c, str.toOpenArrayByte(0, str.len - 1)) + +template consumeOutputs*(sp: OutputStream, bytesVar, body: untyped) = + ## Please note that calling `consumeOutputs` on an unbuffered stream + ## or an unsafe memory stream is considered a Defect. + ## + ## Before consuming the outputs, all outstanding delayed writes must be finalized. + let s = sp + doAssert s.extCursorsCount == 0 and s.buffers != nil + + consumeAllPages(s.buffers, pageStartAddr, pageLen): + template bytesVar: untyped = + makeOpenArray(pageStartAddr, pageLen) + + body + +template consumeContiguousOutput*(sp: OutputStream, bytesVar, body: untyped) = + ## Please note that calling `consumeContiguousOutput` on an unbuffered stream + ## or an unsafe memory stream is considered a Defect. + ## + ## Before consuming the output, all outstanding delayed writes must be finalized. + ## + + # TODO: This code is a bit too much to be inlined. Maybe this should be a proc + # with a callback, but this will restrict the types of variables it can write to. + # OTOH, perhaps only `consumeAllPages` is the offending part. + var + s = sp + contigiousBytes: string # this may remain null + bytesPtr: ptr byte + bytesLen: int + + doAssert s.extCursorsCount == 0 and s.buffers != nil + + if s.buffers.queue.len == 1: + let page = s.buffers.queue[0] + bytesPtr = page.pageStartAddr + bytesLen = page.endOffset - pageStartOffset + # We need to reset the page to an empty state, so it can be reused + page.startOffset = 0 + page.endOffset = 0 + else: + contigiousBytes = newStringOfCap(s.pos) + + consumeAllPages(s.buffers, pageStartAddr, pageLen): + contigiousBytes.add makeOpenArray(cast[ptr char](pageStartAddr), pageLen) + + bytesPtr = addr contigiousBytes[0] + bytesLen = contigiousBytes.len + + template bytesVar: untyped = + makeOpenArray(bytesPtr, bytesLen) + + body + +proc getOutput*(s: OutputStream, T: type string): string = + ## Please note that calling `getOutput` on an unbuffered stream + ## or an unsafe memory stream is considered a Defect. + ## + ## Before consuming the output, all outstanding delayed writes must be finalized. + ## + doAssert s.extCursorsCount == 0 and s.buffers != nil + s.buffers.endLastPageAt s.span.startAddr + + if s.buffers.queue.len == 1: + let page = s.buffers.queue[0] + if page.kind == stringPage and page.startOffset == 0: + result.swap page.data[] + result.setLen page.endOffset + # We clear the buffers, so the stream will be in pristine state. + # The next write is going to create a fresh new starting page. + s.buffers.queue.clear() + return + + result = newStringOfCap(s.pos) + for page in items(s.buffers.queue): + result.add page.pageChars + +template getOutput*(s: OutputStream, T: type seq[byte]): seq[byte] = + cast[seq[byte]](s.getOutput(string)) + +template getOutput*(s: OutputStream): seq[byte] = + cast[seq[byte]](s.getOutput(string)) + diff --git a/faststreams/pipelines.nim b/faststreams/pipelines.nim index f662d9c..1d7e231 100644 --- a/faststreams/pipelines.nim +++ b/faststreams/pipelines.nim @@ -1,9 +1,9 @@ import macros, - input_stream, output_stream + inputs, outputs export - input_stream, output_stream + inputs, outputs macro executePipeline*(start: InputStream, steps: varargs[untyped]) = var input = start @@ -21,7 +21,7 @@ macro executePipeline*(start: InputStream, steps: varargs[untyped]) = `step`(`input`, `outputVar`) input = quote do: - memoryInput(getOutput(`outputVar`)) + unsafeMemoryInput(getOutput(`outputVar`)) if defined(debugMacros) or defined(debugPipelines): echo result.repr diff --git a/faststreams/std_adapters.nim b/faststreams/std_adapters.nim new file mode 100644 index 0000000..a3feb3f --- /dev/null +++ b/faststreams/std_adapters.nim @@ -0,0 +1,2 @@ +import + async_backend diff --git a/faststreams/stdin.nim b/faststreams/stdin.nim new file mode 100644 index 0000000..954a68b --- /dev/null +++ b/faststreams/stdin.nim @@ -0,0 +1,5 @@ +import + inputs + +let fsStdIn* {.threadvar.} = fileInput(system.stdin) + diff --git a/faststreams/stdout.nim b/faststreams/stdout.nim new file mode 100644 index 0000000..beaccdd --- /dev/null +++ b/faststreams/stdout.nim @@ -0,0 +1,5 @@ +import + outputs + +let fsStdOut* {.threadvar.} = fileOutput(system.stdout) + diff --git a/faststreams/textio.nim b/faststreams/textio.nim new file mode 100644 index 0000000..900e3cc --- /dev/null +++ b/faststreams/textio.nim @@ -0,0 +1,92 @@ +import + stew/ptrops, + inputs, outputs, buffers + +# The following code implements writing numbers to a stream without going +# through Nim's `$` operator which will allocate memory. +# It's based on some speed comparisons of different methods presented here: +# http://www.zverovich.net/2013/09/07/integer-to-string-conversion-in-cplusplus.html + +# TODO Maybe the `writeText` proc shouldn't be instantiated for every integer +# type, but only for the largest "native" one. We can promote the rest with +# a template. + +const + digitsTable = block: + var s = "" + for i in 0..99: + if i < 10: s.add '0' + s.add $i + s + + maxLen = ($BiggestInt.high).len + 4 # null terminator, sign + +proc writeText*(s: OutputStream, x: SomeUnsignedInt) = + var + num: array[maxLen, char] + pos = num.len + + template writeByteInReverse(c: char) = + dec pos + num[pos] = c + + var val = x + while val > 99: + # Integer division is slow so do it for a group of two digits instead + # of for every digit. The idea comes from the talk by Alexandrescu + # "Three Optimization Tips for C++". + let base100digitIdx = (val mod 100) * 2 + val = val div 100 + + writeByteInReverse digitsTable[base100digitIdx + 1] + writeByteInReverse digitsTable[base100digitIdx] + + when true: + if val < 10: + writeByteInReverse char(ord('0') + val) + else: + let base100digitIdx = val * 2 + writeByteInReverse digitsTable[base100digitIdx + 1] + writeByteInReverse digitsTable[base100digitIdx] + else: + # Alternative idea: + # We now know enough to write digits directly to the stream. + if val < 10: + write s, byte(ord('\0') + val) + else: + let base100digitIdx = val * 2 + write s, digitsTable[base100digitIdx] + write s, digitsTable[base100digitIdx + 1] + + write s, num.toOpenArray(pos, static(num.len - 1)) + +proc writeText*(s: OutputStream, x: SomeSignedInt) = + # TODO: Determine this accurately + type MatchingUInt = BiggestUInt + + if x < 0: + s.write '-' + # The `0 - x` trick below takes care of one corner case: + # How do we get the abs value of low(int)? + # The naive `-x` triggers an overflow, because low(int8) + # is -128, while high(int8) is 127. + writeText(s, MatchingUInt(0) - MatchingUInt(x)) + else: + writeText(s, MatchingUInt(x)) + +template writeText*(s: OutputStream, str: string) = + write s, str + +template writeText*(s: OutputStream, val: auto) = + write s, $val + +proc writeHex*(s: OutputStream, bytes: openarray[byte]) = + const hexChars = "0123456789abcdef" + + for b in bytes: + s.write hexChars[int b shr 4 and 0xF] + s.write hexChars[int b and 0xF] + +proc writeHex*(s: OutputStream, chars: openarray[char]) = + writeHex s, charsToBytes(chars) + diff --git a/tests/all_tests.nim b/tests/all_tests.nim index f5b0e65..1cf4289 100644 --- a/tests/all_tests.nim +++ b/tests/all_tests.nim @@ -1,5 +1,6 @@ import - test_input_stream, - test_output_stream, - test_pipelines + test_inputs, + test_outputs, + test_pipelines, + test_readme_examples diff --git a/tests/base64.nim b/tests/base64.nim index 341cb1d..fc3369c 100644 --- a/tests/base64.nim +++ b/tests/base64.nim @@ -13,8 +13,6 @@ const invalidChar = 255 paddingByte = byte('=') -template encodeSize(size: int): int = (size * 4 div 3) + 6 - proc base64encode*(i: InputStream, o: OutputStream) = var n: uint32 @@ -25,7 +23,7 @@ proc base64encode*(i: InputStream, o: OutputStream) = n = exp template outputChar(x: typed) = - o.append cb64[x and 63] + o.write cb64[x and 63] while i.readable(3): inputByte(b shl 16) @@ -43,12 +41,12 @@ proc base64encode*(i: InputStream, o: OutputStream) = outputChar(n shr 18) outputChar(n shr 12) outputChar(n shr 6) - o.append paddingByte + o.write paddingByte else: outputChar(n shr 18) outputChar(n shr 12) - o.append paddingByte - o.append paddingByte + o.write paddingByte + o.write paddingByte proc initDecodeTable*(): array[256, char] = # computes a decode table at compile time @@ -80,11 +78,11 @@ proc base64decode*(i: InputStream, o: OutputStream) = raiseInvalidChar(c, i.pos - 1) template outputChar(x: untyped) = - o.append char(x and 255) + o.write char(x and 255) let inputLen = i.len - if inputLen != lengthUnknown: - o.prepareRunway decodeSize(inputLen) + if inputLen.isSome: + o.ensureRunway decodeSize(inputLen.get) # hot loop: read 4 characters at at time while i.readable(8): diff --git a/tests/test_input_stream.nim b/tests/test_input_stream.nim deleted file mode 100644 index 40b5d1f..0000000 --- a/tests/test_input_stream.nim +++ /dev/null @@ -1,12 +0,0 @@ -import - unittest, strutils, stew/ranges/ptr_arith, - ../faststreams - -suite "input stream": - test "string input": - var input = repeat("1234 5678 90AB CDEF\n", 1000) - var stream = memoryInput(input) - - check: - (stream.read(4) == "1234".toOpenArrayByte(0, 3)) - diff --git a/tests/test_inputs.nim b/tests/test_inputs.nim new file mode 100644 index 0000000..ceb1110 --- /dev/null +++ b/tests/test_inputs.nim @@ -0,0 +1,36 @@ +{.used.} + +import + os, unittest, strutils, stew/ranges/ptr_arith, + ../faststreams + +suite "input stream": + test "empty input": + var str = "" + var i = unsafeMemoryInput(str) + + check: + i.readable == false + i.next.isNone + + expect Defect: + echo i.read + + test "missing file input": + const fileName = "there-is-no-such-faststreams-file-1" + + check not fileExists(fileName) + expect CatchableError: discard fileInput(fileName) + + check not fileExists(fileName) + expect CatchableError: discard memFileInput(fileName) + + check not fileExists(fileName) + + test "simple": + var input = repeat("1234 5678 90AB CDEF\n", 1000) + var stream = unsafeMemoryInput(input) + + check: + (stream.read(4) == "1234".toOpenArrayByte(0, 3)) + diff --git a/tests/test_output_stream.nim b/tests/test_output_stream.nim deleted file mode 100644 index 7aacd68..0000000 --- a/tests/test_output_stream.nim +++ /dev/null @@ -1,154 +0,0 @@ -import - os, unittest, random, - stew/ranges/ptr_arith, - ../faststreams - -proc bytes(s: string): seq[byte] = - result = newSeqOfCap[byte](s.len) - for c in s: result.add byte(c) - -template bytes(c: char): byte = byte(c) -template bytes(b: seq[byte]): seq[byte] = b - -proc repeat(b: byte, count: int): seq[byte] = - result = newSeq[byte](count) - for i in 0 ..< count: result[i] = b - -proc randomBytes(n: int): seq[byte] = - result.newSeq n - for i in 0 ..< n: - result[i] = byte(rand(255)) - -suite "output stream": - setup: - var memStream = memoryOutput() - var altOutput: seq[byte] = @[] - var tempFilePath = getTempDir() / "faststreams_testfile" - var fileStream = fileOutput(tempFilePath) - - const bufferSize = 1000000 - var buffer = alloc(bufferSize) - var existingBufferStream = memoryOutput(buffer, bufferSize) - - teardown: - removeFile tempFilePath - - template output(val: auto) {.dirty.} = - altOutput.add bytes(val) - - memStream.append val - fileStream.append val - existingBufferStream.append val - - template checkOutputsMatch = - fileStream.flush - - let - fileContents = readFile(tempFilePath).string.bytes - memStreamContents = memStream.getOutput - - let outputsMatch = - altOutput == memStreamContents and - altOutput == makeOpenArray(cast[ptr byte](buffer), - existingBufferStream.pos) - - check outputsMatch - - test "no appends produce an empty output": - checkOutputsMatch() - - test "append zero length slice": - output "" - checkOutputsMatch() - - test "string output": - for i in 0 .. 1: - output $i - output " bottles on the wall" - output '\n' - - checkOutputsMatch() - - test "delayed write": - output "initial output\n" - const delayedWriteContent = bytes "delayed write\n" - - var cursor = memStream.delayFixedSizeWrite(delayedWriteContent.len) - let cursorStart = memStream.pos - altOutput.add delayedWriteContent - - fileStream.append delayedWriteContent - existingBufferStream.append delayedWriteContent - - var totalBytesWritten = 0 - for i, count in [12, 342, 2121, 23, 1, 34012, 932]: - output repeat(byte(i), count) - totalBytesWritten += count - check memStream.pos - cursorStart == totalBytesWritten - - cursor.writeAndFinalize delayedWriteContent - - checkOutputsMatch() - - test "multi-page delayed writes": - randomize(1000) - - type - DelayedWrite = object - cursor: WriteCursor - content: seq[byte] - written: int - - var delayedWrites = newSeq[DelayedWrite]() - - for i in 0..50: - let - size = rand(8000) + 2000 - randomBytes = randomBytes(size) - decision = rand(100) - - if decision < 70: - # Write at some random cursor - if delayedWrites.len == 0: - continue - - let - i = rand(delayedWrites.len - 1) - written = delayedWrites[i].written - remaining = delayedWrites[i].content.len - written - toWrite = min(rand(remaining) + 10, remaining) - - delayedWrites[i].cursor.append delayedWrites[i].content[written ..< written + toWrite] - delayedWrites[i].written += toWrite - - if remaining - toWrite == 0: - finalize delayedWrites[i].cursor - if i != delayedWrites.len - 1: - swap(delayedWrites[i], delayedWrites[^1]) - delayedWrites.setLen(delayedWrites.len - 1) - - elif decision < 90: - # Normal write - memStream.append randomBytes - altOutput.add randomBytes - - else: - # Create cursor - altOutput.add randomBytes - delayedWrites.add DelayedWrite( - cursor: memStream.delayFixedSizeWrite(randomBytes.len), - content: randomBytes, - written: 0) - - # Check that the stream position is consistently tracked at every step - check altOutput.len == memStream.pos - - # Write all unwritten data to all outstanding cursors - for dw in mitems(delayedWrites): - let remaining = dw.content.len - dw.written - dw.cursor.append dw.content[dw.written ..< dw.written + remaining] - finalize dw.cursor - - # The final outputs are the same - check altOutput == memStream.getOutput - diff --git a/tests/test_outputs.nim b/tests/test_outputs.nim new file mode 100644 index 0000000..8d842cf --- /dev/null +++ b/tests/test_outputs.nim @@ -0,0 +1,240 @@ +{.used.} + +import + os, unittest, random, + stew/ranges/ptr_arith, + ../faststreams, ../faststreams/textio + +proc bytes(s: string): seq[byte] = + result = newSeqOfCap[byte](s.len) + for c in s: result.add byte(c) + +template bytes(c: char): byte = byte(c) +template bytes(b: seq[byte]): seq[byte] = b +template bytes[N, T](b: array[N, T]): seq[byte] = @b + +proc repeat(b: byte, count: int): seq[byte] = + result = newSeq[byte](count) + for i in 0 ..< count: result[i] = b + +proc randomBytes(n: int): seq[byte] = + result.newSeq n + for i in 0 ..< n: + result[i] = byte(rand(255)) + +proc readAllAndClose(s: InputStream): seq[byte] = + while s.readable: + result.add s.read + + close(s) + +import memfiles + +suite "output stream": + setup: + var + nimSeq: seq[byte] = @[] + + memStream = memoryOutput() + smallPageSizeStream = memoryOutput(pageSize = 10) + largePageSizeStream = memoryOutput(pageSize = 1000000) + + fileOutputPath = getTempDir() / "faststreams_testfile" + unbufferedFileOutputPath = getTempDir() / "faststreams_testfile_unbuffered" + + fileStream = fileOutput(fileOutputPath) + unbufferedFileStream = fileOutput(unbufferedFileOutputPath, pageSize = 0) + + bufferSize = 1000000 + buffer = alloc(bufferSize) + streamWritingToExistingBuffer = unsafeMemoryOutput(buffer, bufferSize) + + teardown: + removeFile fileOutputPath + removeFile unbufferedFileOutputPath + dealloc buffer + + template output(val: auto) {.dirty.} = + nimSeq.add bytes(val) + + memStream.write val + smallPageSizeStream.write val + largePageSizeStream.write val + + fileStream.write val + unbufferedFileStream.write val + + streamWritingToExistingBuffer.write val + + template outputText(val: auto) = + let valAsStr = $val + nimSeq.add valAsStr.toOpenArrayByte(0, valAsStr.len - 1) + + memStream.writeText val + smallPageSizeStream.writeText val + largePageSizeStream.writeText val + + fileStream.writeText val + unbufferedFileStream.writeText val + + streamWritingToExistingBuffer.writeText val + + template checkOutputsMatch(showResults = false, + skipUnbufferedFile = false) = + flush fileStream + close fileStream + + flush unbufferedFileStream + close unbufferedFileStream + + check fileExists(fileOutputPath) and + fileExists(unbufferedFileOutputPath) + + let + memStreamRes = memStream.getOutput + readFileRes = readFile(fileOutputPath).string.bytes + fileInputRes = fileInput(fileOutputPath).readAllAndClose + memFileInputRes = memFileInput(fileOutputPath).readAllAndClose + fileInputWithSmallPagesRes = fileInput(fileOutputPath, pageSize = 10).readAllAndClose + + when showResults: + checkpoint "Nim seq result" + checkpoint $nimSeq + + checkpoint "Writes to existing buffer result" + checkpoint $makeOpenArray(cast[ptr byte](buffer), + streamWritingToExistingBuffer.pos) + + checkpoint "mem stream result" + checkpoint $memStreamRes + + checkpoint "readFile result" + checkpoint $readFileRes + + checkpoint "fileInput result" + checkpoint $fileInputRes + + checkpoint "memFileInput result" + checkpoint $memFileInputRes + + checkpoint "fileInput with small pageSize result" + checkpoint $fileInputWithSmallPagesRes + + let outputsMatch = + nimSeq == makeOpenArray(cast[ptr byte](buffer), + streamWritingToExistingBuffer.pos) and + nimSeq == memStreamRes and + nimSeq == readFileRes and + nimSeq == fileInputRes and + nimSeq == memFileInputRes and + nimSeq == fileInputWithSmallPagesRes + + check outputsMatch + + when not skipUnbufferedFile: + let unbufferedFileRes = readFile(unbufferedFileOutputPath).string.bytes + check nimSeq == unbufferedFileRes + + test "no appends produce an empty output": + checkOutputsMatch() + + test "write zero length slices": + output "" + output newSeq[byte]() + var arr: array[0, byte] + output arr + + check nimSeq.len == 0 + checkOutputsMatch() + + test "text output": + for i in 1 .. 100: + outputText i + outputText " bottles on the wall" + outputText '\n' + + checkOutputsMatch() + + test "delayed write": + output "initial output\n" + const delayedWriteContent = bytes "delayed write\n" + + var cursor = memStream.delayFixedSizeWrite(delayedWriteContent.len) + let cursorStart = memStream.pos + + nimSeq.add delayedWriteContent + fileStream.write delayedWriteContent + streamWritingToExistingBuffer.write delayedWriteContent + + var totalBytesWritten = 0 + for i, count in [2]: # 12, 342, 2121, 23, 1, 34012, 932]: + output repeat(byte(i), count) + totalBytesWritten += count + check memStream.pos - cursorStart == totalBytesWritten + + cursor.finalWrite delayedWriteContent + + checkOutputsMatch(skipUnbufferedFile = true) + + test "multi-page delayed writes": + randomize(1000) + + type + DelayedWrite = object + cursor: WriteCursor + content: seq[byte] + written: int + + var delayedWrites = newSeq[DelayedWrite]() + + for i in 0..50: + let + size = rand(8000) + 2000 + randomBytes = randomBytes(size) + decision = rand(100) + + if decision < 70: + # Write at some random cursor + if delayedWrites.len == 0: + continue + + let + i = rand(delayedWrites.len - 1) + written = delayedWrites[i].written + remaining = delayedWrites[i].content.len - written + toWrite = min(rand(remaining) + 10, remaining) + + delayedWrites[i].cursor.write delayedWrites[i].content[written ..< written + toWrite] + delayedWrites[i].written += toWrite + + if remaining - toWrite == 0: + finalize delayedWrites[i].cursor + if i != delayedWrites.len - 1: + swap(delayedWrites[i], delayedWrites[^1]) + delayedWrites.setLen(delayedWrites.len - 1) + + elif decision < 90: + # Normal write + memStream.write randomBytes + nimSeq.add randomBytes + + else: + # Create cursor + nimSeq.add randomBytes + delayedWrites.add DelayedWrite( + cursor: memStream.delayFixedSizeWrite(randomBytes.len), + content: randomBytes, + written: 0) + + # Check that the stream position is consistently tracked at every step + check nimSeq.len == memStream.pos + + # Write all unwritten data to all outstanding cursors + for dw in mitems(delayedWrites): + let remaining = dw.content.len - dw.written + dw.cursor.write dw.content[dw.written ..< dw.written + remaining] + finalize dw.cursor + + # The final outputs are the same + check nimSeq == memStream.getOutput + diff --git a/tests/test_pipelines.nim b/tests/test_pipelines.nim index e4b6f49..481ea20 100644 --- a/tests/test_pipelines.nim +++ b/tests/test_pipelines.nim @@ -1,3 +1,5 @@ +{.used.} + import std/[unittest, strutils, base64], ../faststreams/pipelines, @@ -13,7 +15,7 @@ type proc upcaseAllCharacters(i: InputStream, o: OutputStream) = while i.readable: - o.append toUpperAscii(char i.read()) + o.write toUpperAscii(char i.read()) template timeit(timerVar: var Nanos, code: untyped) = let t0 = getTicks() @@ -40,7 +42,7 @@ suite "pipelines": timeIt times.fsPipeline: var memOut = memoryOutput() - executePipeline(memoryInput(loremIpsum), + executePipeline(unsafeMemoryInput(loremIpsum), upcaseAllCharacters, base64encode, base64decode, diff --git a/tests/test_readme_examples.nim b/tests/test_readme_examples.nim new file mode 100644 index 0000000..6e5ff68 --- /dev/null +++ b/tests/test_readme_examples.nim @@ -0,0 +1,53 @@ +{.used.} + +import + typetraits, ../faststreams + +proc writeNimRepr*(stream: OutputStream, str: string) = + stream.write '"' + + for c in str: + if c == '"': + stream.write ['\'', '"'] + else: + stream.write c + + stream.write '"' + +proc writeNimRepr*(stream: OutputStream, x: char) = + stream.write ['\'', x, '\''] + +proc writeNimRepr*(stream: OutputStream, x: int) = + stream.write $x # Making this more optimal has been left + # as an exercise for the reader + +proc writeNimRepr*[T](stream: OutputStream, obj: T) = + stream.write typetraits.name(T) + stream.write '(' + + var firstField = true + for name, val in fieldPairs(obj): + if not firstField: + stream.write ", " + + stream.write name + stream.write ": " + stream.writeNimRepr val + + firstField = false + + stream.write ')' + +type + ABC = object + a: int + b: char + c: string + +block: + var stream = memoryOutput() + stream.writeNimRepr(ABC(a: 1, b: 'b', c: "str")) + var repr = stream.getOutput(string) + + doAssert repr == "ABC(a: 1, b: 'b', c: \"str\")" +