diff --git a/faststreams/buffers.nim b/faststreams/buffers.nim index 2444aa6..ba421c2 100644 --- a/faststreams/buffers.nim +++ b/faststreams/buffers.nim @@ -13,6 +13,7 @@ type Page* = object consumedTo*: Natural writtenTo*: Natural + fauxEofAt*: Natural data*: ref string PageRef* = ref Page @@ -26,6 +27,7 @@ type waitingWriter*: Future[void] eofReached*: bool + fauxEofPos*: Natural const nimPageSize* = 4096 @@ -36,6 +38,8 @@ const defaultPageSize* = 4096 - nimAllocatorMetadataSize maxStackUsage* = 16384 + noEofChange* = high(Natural) + when debugHelpers: proc describeBuffers*(context: static string, buffers: PageBuffers) = debugEcho context, " :: buffers" @@ -78,11 +82,39 @@ template pageChars*(page: PageRef): untyped = let baseAddr = cast[ptr UncheckedArray[char]](allocationStart(page)) toOpenArray(baseAddr, page.consumedTo, page.writtenTo - 1) -func obtainReadableSpan*(page: PageRef, writable: static[bool] = false): PageSpan = - let baseAddr = page.allocationStart - result = PageSpan(startAddr: offset(baseAddr, page.consumedTo), - endAddr: offset(baseAddr, page.writtenTo)) - page.consumedTo = page.writtenTo +func obtainReadableSpan*(buffers: PageBuffers, + currentSpanEndPos: var Natural): PageSpan = + if buffers.queue.len == 0: + return default(PageSpan) + + var page = buffers.queue[0] + var unconsumedLen = page.writtenTo - page.consumedTo + if unconsumedLen == 0: + if buffers.queue.len > 1: + discard buffers.queue.popFirst + page = buffers.queue[0] + unconsumedLen = page.writtenTo - page.consumedTo + else: + return default(PageSpan) + + let + baseAddr = page.allocationStart + startAddr = offset(baseAddr, page.consumedTo) + + let usableLen = if buffers.fauxEofPos != 0: + let maxSize = buffers.fauxEofPos - currentSpanEndPos + if maxSize < unconsumedLen: + maxSize + else: + unconsumedLen + else: + unconsumedLen + + page.consumedTo += usableLen + currentSpanEndPos += usableLen + + PageSpan(startAddr: startAddr, + endAddr: offset(startAddr, usableLen)) func writableSpan*(page: PageRef): PageSpan = let baseAddr = allocationStart(page) @@ -115,6 +147,13 @@ func trackWrittenTo*(buffers: PageBuffers, spanHeadPos: ptr byte) = var topPage = buffers.queue.peekLast topPage.writtenTo = distance(topPage.allocationStart, spanHeadPos) +proc setFauxEof*(buffers: PageBuffers, pos: Natural): Natural = + result = buffers.fauxEofPos + buffers.fauxEofPos = pos + +proc restoreEof*(buffers: PageBuffers, pos: Natural) = + buffers.fauxEofPos = pos + func addWritablePage*(buffers: PageBuffers, pageSize: Natural): PageRef = trackWrittenToEnd(buffers) result = PageRef(data: allocRef newString(pageSize)) @@ -136,25 +175,6 @@ template getWritableSpan*(buffers: PageBuffers): PageSpan = let page = getWritablePage(buffers, buffers.pageSize) writableSpan(page) -func nextReadableSpan*(buffers: PageBuffers, span: var PageSpan) = - let - firstPage = buffers.queue.peekFirst - pageReadableEnd = firstPage.readableEnd - - if span.endAddr == nil: - fsAssert buffers.queue.len > 0 - span = obtainReadableSpan buffers.queue[0] - elif span.endAddr != pageReadableEnd: - # Check whether the span points within the current page: - fsAssert distance(firstPage.allocationStart, span.endAddr) >= 0 and - distance(span.endAddr, pageReadableEnd) >= 0 - span.endAddr = pageReadableEnd - firstPage.consumedTo = firstPage.writtenTo - else: - fsAssert buffers.queue.len > 1 - discard buffers.queue.popFirst - span = obtainReadableSpan buffers.queue[0] - func stringFromBytes(src: pointer, srcLen: Natural): string = result = newString(srcLen) copyMem(addr result[0], src, srcLen) diff --git a/faststreams/inputs.nim b/faststreams/inputs.nim index 60d4aef..a38420a 100644 --- a/faststreams/inputs.nim +++ b/faststreams/inputs.nim @@ -206,8 +206,19 @@ proc memFileInput*(filename: string, mappedSize = -1, offset = 0): InputStreamHa endAddr: offset(head, mappedSize)), file: memFile) +func getNewSpan(s: InputStream) = + let buffers = s.buffers + fsAssert buffers != nil + s.span = buffers.obtainReadableSpan(s.spanEndPos) + +func getNewSpanOrDieTrying(s: InputStream) = + getNewSpan s + fsAssert s.span.hasRunway + proc readableNow*(s: InputStream): bool = - (not s.span.atEnd) or (s.buffers != nil and s.buffers.len > 1) + if s.span.hasRunway: return true + getNewSpan s + s.span.hasRunway template readableNow*(s: AsyncInputStream): bool = readableNow InputStream(s) @@ -229,8 +240,7 @@ proc readOnce*(sp: AsyncInputStream): Future[Natural] {.async.} = disconnectInputDevice(s) if result > 0 and s.span.len == 0: - s.buffers.nextReadableSpan(s.span) - s.spanEndPos += s.span.len + getNewSpan s proc timeoutToNextByteImpl(s: AsyncInputStream, deadline: Future): Future[bool] {.async.} = @@ -261,32 +271,11 @@ proc closeAsync*(s: AsyncInputStream) {.async.} = # TODO: End of purely async interface -func flipPage(s: InputStream) = - fsAssert s.buffers != nil and s.buffers.len > 1 - discard s.buffers.popFirst - s.span = obtainReadableSpan s.buffers[0] - s.spanEndPos += s.span.len - func getBestContiguousRunway(s: InputStream): Natural = result = s.span.len - if result == 0: - if s.buffers != nil and s.buffers.len > 1: - flipPage s - result = s.span.len - -template withReadableRange*(sp: InputStream|AsyncInputStream, - rangeLen: Natural, - rangeStreamVarName, blk: untyped) = - let s = InputStream sp - - let vtable = s.vtable - s.vtable = nil - - try: - let `rangeStreamVarName` {.inject.} = s - blk - finally: - s.vtable = vtable + if result == 0 and s.buffers != nil: + getNewSpan s + result = s.span.len func totalUnconsumedBytes*(s: InputStream): Natural = ## Returns the number of bytes that are currently sitting within the stream @@ -297,7 +286,7 @@ func totalUnconsumedBytes*(s: InputStream): Natural = else: s.buffers.totalBufferedBytes if localRunway == 0 and runwayInBuffers > 0: - flipPage s + getNewSpan s localRunway + runwayInBuffers @@ -305,6 +294,38 @@ template totalUnconsumedBytes*(s: AsyncInputStream): Natural = ## Alias for InputStream.totalUnconsumedBytes totalUnconsumedBytes InputStream(s) +proc limitReadableRange(s: InputStream, rangeLen: Natural): Natural = + fsAssert rangeLen > 0 + s.vtable = nil + + let runway = s.span.len + if rangeLen > runway: + return s.buffers.setFauxEof(rangeLen - runway + s.spanEndPos) + else: + s.span.endAddr = offset(s.span.startAddr, rangeLen) + let bytesToUnconsume = runway - rangeLen + s.spanEndPos -= bytesToUnconsume + if s.buffers != nil: + s.buffers.queue.peekFirst.consumedTo -= bytesToUnconsume + return s.buffers.setFauxEof(s.spanEndPos) + +template withReadableRange*(sp: InputStream|AsyncInputStream, + rangeLen: Natural, + rangeStreamVarName, blk: untyped) = + let + s = InputStream sp + vtable = s.vtable + origEndAddr = s.span.endAddr + origEof = limitReadableRange(s, rangeLen) + + try: + let rangeStreamVarName {.inject.} = s + blk + finally: + s.vtable = vtable + s.span.endAddr = origEndAddr + restoreEof(s.buffers, origEof) + let fileInputVTable = InputStreamVTable( readSync: proc (s: InputStream, dst: pointer, dstLen: Natural): Natural {.nimcall, gcsafe, raises: [IOError, Defect].} = @@ -384,24 +405,30 @@ template len*(s: AsyncInputStream): Option[Natural] = len InputStream(s) func memoryInput*(buffers: PageBuffers): InputStreamHandle = + var spanEndPos = Natural 0 var span = if buffers.len == 0: default(PageSpan) - else: obtainReadableSpan buffers.queue[0] + else: buffers.obtainReadableSpan(spanEndPos) makeHandle InputStream(buffers: buffers, span: span, - spanEndPos: span.len) + spanEndPos: spanEndPos) func memoryInput*(data: openarray[byte]): InputStreamHandle = - let - buffers = initPageBuffers(data.len) - page = buffers.addWritablePage(data.len) - pageSpan = page.fullSpan + let stream = if data.len > 0: + let + buffers = initPageBuffers(data.len) + page = buffers.addWritablePage(data.len) + pageSpan = page.fullSpan - copyMem(pageSpan.startAddr, unsafeAddr data[0], data.len) + copyMem(pageSpan.startAddr, unsafeAddr data[0], data.len) - makeHandle InputStream(buffers: buffers, - span: pageSpan, - spanEndPos: data.len) + InputStream(buffers: buffers, + span: pageSpan, + spanEndPos: data.len) + else: + InputStream() + + makeHandle stream func memoryInput*(data: openarray[char]): InputStreamHandle = memoryInput charsToBytes(data) @@ -409,9 +436,9 @@ func memoryInput*(data: openarray[char]): InputStreamHandle = proc resetBuffers*(s: InputStream, buffers: PageBuffers) = # This should be used only on safe memory input streams fsAssert s.vtable == nil and s.buffers != nil and buffers.len > 0 + s.spanEndPos = 0 s.buffers = buffers - s.span = obtainReadableSpan buffers.queue[0] - s.spanEndPos = s.span.len + getNewSpan s proc continueAfterRead(s: InputStream, bytesRead: Natural): bool = # Please note that this is extracted into a proc only to reduce the code @@ -426,8 +453,7 @@ proc continueAfterRead(s: InputStream, bytesRead: Natural): bool = disconnectInputDevice(s) if bytesRead > 0: - s.buffers.nextReadableSpan(s.span) - s.spanEndPos += s.span.len + getNewSpan s return true else: return false @@ -440,21 +466,23 @@ template bufferMoreDataImpl(s, awaiter, readOp: untyped): bool = # 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: + if s.buffers == 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 + getNewSpan s + if hasRunway(s.span): true - else: + elif s.vtable != nil and s.vtable.readOp != nil: # 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: continueAfterRead(s, awaiter s.vtable.readOp(s, nil, 0)) + else: + false proc bufferMoreDataSync(s: InputStream): bool = # This proc exists only to avoid inlining of the code of @@ -518,8 +546,7 @@ template readable*(sp: AsyncInputStream): bool = func continueAfterReadN(s: InputStream, runwayBeforeRead, bytesRead: Natural) = if runwayBeforeRead == 0 and bytesRead > 0: - s.buffers.nextReadableSpan(s.span) - s.spanEndPos += s.span.len + getNewSpan s template readableNImpl(s, n, awaiter, readOp: untyped): bool = let runway = totalUnconsumedBytes(s) @@ -590,27 +617,22 @@ template readable*(sp: AsyncInputStream, np: int): bool = readableNImpl(s, n, fsAwait, readAsync) -when false: - func flipPagePeek(s: InputStream): byte = - flipPage s - result = s.span.startAddr[] - -func flipPageRead(s: InputStream): byte = - flipPage s - result = s.span.startAddr[] - bumpPointer s.span - template peek*(sp: InputStream): byte = let s = sp if hasRunway(s.span): s.span.startAddr[] else: - flipPage s + getNewSpanOrDieTrying s s.span.startAddr[] template peek*(s: AsyncInputStream): byte = peek InputStream(s) +func readFromNewSpan(s: InputStream): byte = + getNewSpanOrDieTrying s + result = s.span.startAddr[] + bumpPointer s.span + template read*(sp: InputStream): byte = let s = sp if hasRunway(s.span): @@ -618,7 +640,7 @@ template read*(sp: InputStream): byte = bumpPointer(s.span) res else: - flipPageRead s + readFromNewSpan s template read*(s: AsyncInputStream): byte = read InputStream(s) @@ -636,7 +658,7 @@ proc advance*(s: InputStream) = if hasRunway(s.span): bumpPointer s.span else: - flipPage s + getNewSpan s proc advance*(s: InputStream, n: Natural) = # TODO This is silly, implement it properly @@ -666,7 +688,7 @@ proc drainBuffersInto*(s: InputStream, dstAddr: ptr byte, dstLen: Natural): Natu if s.buffers != nil: # Since we reached the end of the current page, - # we have to do the equivalent of `flipPage`: + # we have to do the equivalent of `getNewSpan`: # TODO: what if the page was extended? if s.buffers.len > 0: @@ -717,21 +739,28 @@ proc drainBuffersInto*(s: InputStream, dstAddr: ptr byte, dstLen: Natural): Natu template readIntoExImpl(s: InputStream, dst: ptr byte, dstLen: Natural, awaiter, readOp: untyped): Natural = - var bytesRead = drainBuffersInto(s, dst, dstLen) + let totalBytesDrained = drainBuffersInto(s, dst, dstLen) + var bytesDeficit = (dstLen - totalBytesDrained) - while bytesRead < dstLen: - let - bytesDeficit = dstLen - bytesRead - adjustedDst = offset(dst, bytesRead) + if bytesDeficit > 0: + var adjustedDst = offset(dst, totalBytesDrained) - bytesRead += awaiter s.vtable.readOp(s, adjustedDst, bytesDeficit) + while true: + let newBytesRead = awaiter s.vtable.readOp(s, adjustedDst, bytesDeficit) - if s.buffers.eofReached: - disconnectInputDevice(s) - break + s.spanEndPos += newBytesRead + bytesDeficit -= newBytesRead - s.spanEndPos += bytesRead - bytesRead + if s.buffers.eofReached: + disconnectInputDevice(s) + break + + if bytesDeficit == 0: + break + + adjustedDst = offset(dst, newBytesRead) + + dstLen - bytesDeficit proc readIntoEx*(s: InputStream, dst: var openarray[byte]): int = ## Read data into the destination buffer. @@ -854,21 +883,3 @@ proc pos*(s: InputStream): int {.inline.} = 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 - fsAssert offsetFromEnd < 0 - result = offset(s.span.endAddr, offsetFromEnd) - fsAssert 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) - fsAssert s.head >= s.bufferStart - - proc rewindTo*(s: InputStream, pos: int) {.inline.} = - s.head = s.bufferPos(pos) - diff --git a/faststreams/pipelines.nim b/faststreams/pipelines.nim index 9f5a805..6e5c3ac 100644 --- a/faststreams/pipelines.nim +++ b/faststreams/pipelines.nim @@ -172,8 +172,9 @@ func pipeInput*(source: InputStream, func pipeInput*(buffers: PageBuffers, allowWaitFor = false, source: InputStream = nil): AsyncInputStream = + var spanEndPos = Natural 0 var span = if buffers.len == 0: default(PageSpan) - else: obtainReadableSpan buffers.queue[0] + else: buffers.obtainReadableSpan(spanEndPos) AsyncInputStream LayeredInputStream( vtable: vtableAddr pipeInputVTable, diff --git a/tests/test_inputs.nim b/tests/test_inputs.nim index 54107d8..74f2b9e 100644 --- a/tests/test_inputs.nim +++ b/tests/test_inputs.nim @@ -53,10 +53,10 @@ procSuite "input stream": var input = unsafeMemoryInput(str) emptyInputTests "fileInput": - var input = memFileInput("files" / "empty_file") + var input = fileInput("files" / "empty_file") emptyInputTests "memFileInput": - var input = fileInput("files" / "empty_file") + var input = memFileInput("files" / "empty_file") template asciiTableFileTest(name: string, body: untyped) = test name & " of ascii table with regular pageSize": @@ -137,11 +137,18 @@ procSuite "input stream": check not fileExists(fileName) test "non-blocking reads": - let s = fileInput(asciiTableFile, pageSize = 20) - if s.readable: + let s = fileInput(asciiTableFile, pageSize = 100) + if s.readable(20): s.withReadableRange(20, r): check r.readAll.len == 20 - check s.readable + + check s.readable + + if s.readable(200): + s.withReadableRange(200, r): + check r.readAll.len == 200 + + check s.readable test "simple": var input = repeat("1234 5678 90AB CDEF\n", 1000)