Address review comments; Add documentation; Shared buffering mechanism for input and output streams

This commit is contained in:
Zahary Karadjov 2020-04-29 21:19:18 +03:00
commit b24300bd3f
No known key found for this signature in database
GPG key ID: C8936F8A3073D609
24 changed files with 2396 additions and 1060 deletions

405
README.md
View file

@ -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. <br />
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`. <br />
It can represent any Chronos `Transport` as an input stream.
* `asyncSocketInput` (async)
Enabled by importing `faststreams/std_adapters`. <br />
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. <br />
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`. <br />
It can represent any Chronos `Transport` as an input stream.
* `asyncSocketOutput` (async)
Enabled by importing `faststreams/std_adapters`. <br />
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

View file

@ -1,6 +1,6 @@
import
faststreams/[input_stream, output_stream]
faststreams/[inputs, outputs]
export
input_stream, output_stream
inputs, outputs

View file

@ -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"

View file

@ -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

226
faststreams/buffers.nim Normal file
View file

@ -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)

View file

@ -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),

View file

@ -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)

577
faststreams/inputs.nim Normal file
View file

@ -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)

View file

@ -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,

View file

@ -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

685
faststreams/outputs.nim Normal file
View file

@ -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))

View file

@ -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

View file

@ -0,0 +1,2 @@
import
async_backend

5
faststreams/stdin.nim Normal file
View file

@ -0,0 +1,5 @@
import
inputs
let fsStdIn* {.threadvar.} = fileInput(system.stdin)

5
faststreams/stdout.nim Normal file
View file

@ -0,0 +1,5 @@
import
outputs
let fsStdOut* {.threadvar.} = fileOutput(system.stdout)

92
faststreams/textio.nim Normal file
View file

@ -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)

View file

@ -1,5 +1,6 @@
import
test_input_stream,
test_output_stream,
test_pipelines
test_inputs,
test_outputs,
test_pipelines,
test_readme_examples

View file

@ -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):

View file

@ -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))

36
tests/test_inputs.nim Normal file
View file

@ -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))

View file

@ -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

240
tests/test_outputs.nim Normal file
View file

@ -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

View file

@ -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,

View file

@ -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\")"