diff --git a/faststreams/buffers.nim b/faststreams/buffers.nim index 30b3e22..2444aa6 100644 --- a/faststreams/buffers.nim +++ b/faststreams/buffers.nim @@ -316,31 +316,6 @@ template charsToBytes*(chars: openArray[char]): untyped = var charsStart = unsafeAddr chars[0] makeOpenArray(cast[ptr byte](charsStart), chars.len) -template implementWrites*(buffersParam: PageBuffers, - srcParam: pointer, - srcLenParam: Natural, - dstDesc: static string, - writeStartVar, writeLenVar, - writeBlock: untyped) = - let - buffers = buffersParam - writeStartVar = srcParam - writeLenVar = srcLenParam - - template raiseError = - raise newException(IOError, "Failed to write all bytes to " & dstDesc) - - if buffers != nil: - for writeStartVar, writeLenVar in consumePageBuffers(s.buffers): - let bytesWritten = writeBlock - # TODO: Can we repair the buffers here? - if bytesWritten != writeLenVar: raiseError() - - if srcLen > 0: - fsAssert src != nil - let bytesWritten = writeBlock - if bytesWritten != writeLenVar: raiseError() - type ReadFlag* = enum partialReadIsEof diff --git a/faststreams/inputs.nim b/faststreams/inputs.nim index 6966c1a..25cad53 100644 --- a/faststreams/inputs.nim +++ b/faststreams/inputs.nim @@ -53,11 +53,11 @@ type FileInputStream = ref object of InputStream file: File -template Async*(s: InputStream): AsyncInputStream = - AsyncInputStream(s) +template Sync*(s: InputStream): InputStream = s +template Async*(s: InputStream): AsyncInputStream = AsyncInputStream(s) -template Sync*(s: AsyncInputStream): InputStream = - InputStream(s) +template Sync*(s: AsyncInputStream): InputStream = InputStream(s) +template Async*(s: AsyncInputStream): AsyncInputStream = s proc disconnectInputDevice(s: InputStream) = # TODO @@ -578,9 +578,17 @@ proc advance*(s: InputStream) = else: flipPage s +proc advance*(s: InputStream, n: Natural) = + # TODO This is silly, implement it properly + for i in 0 ..< n: + advance s + template advance*(s: AsyncInputStream) = advance InputStream(s) +template advance*(s: AsyncInputStream, n: Natural) = + advance InputStream(s), n + proc drainBuffersInto*(s: InputStream, dstAddr: ptr byte, dstLen: Natural): Natural = var dst = dstAddr diff --git a/faststreams/multisync.nim b/faststreams/multisync.nim index 012092b..5f5e644 100644 --- a/faststreams/multisync.nim +++ b/faststreams/multisync.nim @@ -2,6 +2,9 @@ import stew/shims/macros, async_backend, inputs, outputs +export + async_backend + macro fsMultiSync*(body: untyped) = # We will produce an identical copy of the annotated proc, # but taking async parameters and having the async pragma. @@ -13,9 +16,9 @@ macro fsMultiSync*(body: untyped) = # The return types becomes Future[T] if asyncProcParams[0].kind == nnkEmpty: - asyncProcParams[0] = newTree(nnkBracketExpr, ident"Future", ident"void") + asyncProcParams[0] = newTree(nnkBracketExpr, bindSym"Future", ident"void") else: - asyncProcParams[0] = newTree(nnkBracketExpr, ident"Future", asyncProcParams[0]) + asyncProcParams[0] = newTree(nnkBracketExpr, bindSym"Future", asyncProcParams[0]) # We replace all stream inputs with their async counterparts for i in 1 ..< asyncProcParams.len: diff --git a/faststreams/outputs.nim b/faststreams/outputs.nim index d5195fb..7c6cbdf 100644 --- a/faststreams/outputs.nim +++ b/faststreams/outputs.nim @@ -12,7 +12,7 @@ import buffers, async_backend export - CloseBehavior + initPageBuffers, CloseBehavior type OutputStream* = ref object of RootObj @@ -69,11 +69,11 @@ type FileOutputStream = ref object of OutputStream file: File -template Async*(s: OutputStream): AsyncOutputStream = - AsyncOutputStream(s) +template Sync*(s: OutputStream): OutputStream = s +template Async*(s: OutputStream): AsyncOutputStream = AsyncOutputStream(s) -template Sync*(s: AsyncOutputStream): OutputStream = - OuputStream(s) +template Sync*(s: AsyncOutputStream): OutputStream = OutputStream(s) +template Async*(s: AsyncOutputStream): AsyncOutputStream = s proc disconnectOutputDevice(s: OutputStream) = if s.vtable != nil: @@ -199,6 +199,31 @@ proc ensureRunway*(s: OutputStream, neededRunway: Natural) = template ensureRunway*(s: AsyncOutputStream, neededRunway: Natural) = ensureRunway OutputStream(s), neededRunway +template implementWrites*(buffersParam: PageBuffers, + srcParam: pointer, + srcLenParam: Natural, + dstDesc: static string, + writeStartVar, writeLenVar, + writeBlock: untyped) = + let + buffers = buffersParam + writeStartVar = srcParam + writeLenVar = srcLenParam + + template raiseError = + raise newException(IOError, "Failed to write all bytes to " & dstDesc) + + if buffers != nil: + for writeStartVar, writeLenVar in consumePageBuffers(s.buffers): + let bytesWritten = writeBlock + # TODO: Can we repair the buffers here? + if bytesWritten != writeLenVar: raiseError() + + if writeLenVar > 0: + fsAssert writeStartVar != nil + let bytesWritten = writeBlock + if bytesWritten != writeLenVar: raiseError() + let fileOutputVTable = OutputStreamVTable( writeSync: proc (s: OutputStream, src: pointer, srcLen: Natural) {.nimcall, gcsafe, raises: [IOError, Defect].} = @@ -583,6 +608,9 @@ proc writeBytesAsyncImpl(s: OutputStream, str: string): Future[void] = writeBytesAsyncImpl s, toOpenArray(str, 0, str.len - 1) +template writeAndWait*(s: OutputStream, value: untyped) = + write s, value + template writeAndWait*(sp: AsyncOutputStream, value: untyped) = bind writeBytesAsyncImpl