Implement InputStream.nonBlockingReads and document it
This commit is contained in:
parent
fc76250d02
commit
ec5c3af313
3 changed files with 56 additions and 0 deletions
29
README.md
29
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
|
It is assumed that in traditional async code, timeouts will be managed more
|
||||||
explicitly with `sleepAsync` and the `or` operator defined over futures.
|
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`
|
### `OutputStream` and `AsyncOutputStream`
|
||||||
|
|
||||||
An `OutputStream` manages a particular output device. The library offers out
|
An `OutputStream` manages a particular output device. The library offers out
|
||||||
|
|
|
||||||
|
|
@ -274,6 +274,22 @@ func getBestContiguousRunway(s: InputStream): Natural =
|
||||||
flipPage s
|
flipPage s
|
||||||
result = s.span.len
|
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 =
|
func totalUnconsumedBytes*(s: InputStream): Natural =
|
||||||
## Returns the number of bytes that are currently sitting within the stream
|
## Returns the number of bytes that are currently sitting within the stream
|
||||||
## buffers and that can be consumed with `read` or `advance`.
|
## buffers and that can be consumed with `read` or `advance`.
|
||||||
|
|
|
||||||
|
|
@ -20,6 +20,10 @@ proc countLines(s: InputStream): Natural =
|
||||||
for s in lines(s):
|
for s in lines(s):
|
||||||
inc result
|
inc result
|
||||||
|
|
||||||
|
proc readAll(s: InputStream): seq[byte] =
|
||||||
|
while s.readable:
|
||||||
|
result.add s.read
|
||||||
|
|
||||||
const
|
const
|
||||||
asciiTableFile = "files" / "ascii_table.txt"
|
asciiTableFile = "files" / "ascii_table.txt"
|
||||||
asciiTableContents = slurp(asciiTableFile)
|
asciiTableContents = slurp(asciiTableFile)
|
||||||
|
|
@ -132,6 +136,13 @@ procSuite "input stream":
|
||||||
|
|
||||||
check not fileExists(fileName)
|
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":
|
test "simple":
|
||||||
var input = repeat("1234 5678 90AB CDEF\n", 1000)
|
var input = repeat("1234 5678 90AB CDEF\n", 1000)
|
||||||
var stream = unsafeMemoryInput(input)
|
var stream = unsafeMemoryInput(input)
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue