diff --git a/README.md b/README.md index c6b26ac..06f31c5 100644 --- a/README.md +++ b/README.md @@ -317,6 +317,35 @@ proc performHandshake(c: Connection): bool {.async.} = It is assumed that in traditional async code, timeouts will be managed more explicitly with `sleepAsync` and the `or` operator defined over futures. +#### Non-blocking reads + +Protocols transmitting serialized payloads often provide information regarding +the size of the payload. When you invoke the deserialization routine, it's +preferable if the provided boundaries are treated like an "end of file" marker +for the deserializer. FastStreams provides an easy way to achieve this without +extra copies and memory allocations through the `nonBlockingReads` facility. +Here is a typical usage: + +```nim +proc decodeFrame(s: AsyncInputStream, DecodedType: type): Option[DecodedType] = + if not s.readable(4): + return + + let lengthPrefix = toInt32 s.read(4) + if s.readable(lengthPrefix): + s.nonBlockingReads(lengthPrefix): + s.readValue(Json, DecodedType) +``` + +Please note that the above example uses the [nim-serialization library](https://github.com/status-im/nim-serialization/) + +Simply, inside the `nonBlockingReads` block, `s.readable` will return `false` +as soon as the Json parser has consumed the specified number of bytes. +Furthermore, it's guaranteed that no blocking operations are possible and +thus our `AsyncInputStream` will be treated like a normal `InputStream`. +Depending on the complexity of the stream processors, this will often lead +to much more optimal code. + ### `OutputStream` and `AsyncOutputStream` An `OutputStream` manages a particular output device. The library offers out diff --git a/faststreams/inputs.nim b/faststreams/inputs.nim index f1cdae2..7a511e4 100644 --- a/faststreams/inputs.nim +++ b/faststreams/inputs.nim @@ -274,6 +274,22 @@ func getBestContiguousRunway(s: InputStream): Natural = flipPage s result = s.span.len +template nonBlockingReads*(sp: InputStream|AsyncInputStream, newName, blk: untyped) = + let s = InputStream sp + const sname = astToStr(sp) + + let vtable = s.vtable + s.vtable = nil + + try: + let `newName` {.inject.} = s + blk + finally: + s.vtable = vtable + +template nonBlockingReads*(sp: InputStream|AsyncInputStream, blk: untyped) = + nonBlockingReads(sp, astToStr(sp), blk) + 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`. diff --git a/tests/test_inputs.nim b/tests/test_inputs.nim index 2d8d7b5..688678a 100644 --- a/tests/test_inputs.nim +++ b/tests/test_inputs.nim @@ -20,6 +20,10 @@ proc countLines(s: InputStream): Natural = for s in lines(s): inc result +proc readAll(s: InputStream): seq[byte] = + while s.readable: + result.add s.read + const asciiTableFile = "files" / "ascii_table.txt" asciiTableContents = slurp(asciiTableFile) @@ -132,6 +136,13 @@ procSuite "input stream": check not fileExists(fileName) + test "non-blocking reads": + let s = fileInput(asciiTableFile, pageSize = 20) + if s.readable: + s.nonBlockingReads: + check s.readAll.len == 20 + check s.readable + test "simple": var input = repeat("1234 5678 90AB CDEF\n", 1000) var stream = unsafeMemoryInput(input)