Implement the full behavior of withReadableRange; close #12

This commit is contained in:
Zahary Karadjov 2020-05-22 09:44:25 +03:00
commit 442ad53b00
No known key found for this signature in database
GPG key ID: C8936F8A3073D609
4 changed files with 160 additions and 121 deletions

View file

@ -13,6 +13,7 @@ type
Page* = object Page* = object
consumedTo*: Natural consumedTo*: Natural
writtenTo*: Natural writtenTo*: Natural
fauxEofAt*: Natural
data*: ref string data*: ref string
PageRef* = ref Page PageRef* = ref Page
@ -26,6 +27,7 @@ type
waitingWriter*: Future[void] waitingWriter*: Future[void]
eofReached*: bool eofReached*: bool
fauxEofPos*: Natural
const const
nimPageSize* = 4096 nimPageSize* = 4096
@ -36,6 +38,8 @@ const
defaultPageSize* = 4096 - nimAllocatorMetadataSize defaultPageSize* = 4096 - nimAllocatorMetadataSize
maxStackUsage* = 16384 maxStackUsage* = 16384
noEofChange* = high(Natural)
when debugHelpers: when debugHelpers:
proc describeBuffers*(context: static string, buffers: PageBuffers) = proc describeBuffers*(context: static string, buffers: PageBuffers) =
debugEcho context, " :: buffers" debugEcho context, " :: buffers"
@ -78,11 +82,39 @@ template pageChars*(page: PageRef): untyped =
let baseAddr = cast[ptr UncheckedArray[char]](allocationStart(page)) let baseAddr = cast[ptr UncheckedArray[char]](allocationStart(page))
toOpenArray(baseAddr, page.consumedTo, page.writtenTo - 1) toOpenArray(baseAddr, page.consumedTo, page.writtenTo - 1)
func obtainReadableSpan*(page: PageRef, writable: static[bool] = false): PageSpan = func obtainReadableSpan*(buffers: PageBuffers,
let baseAddr = page.allocationStart currentSpanEndPos: var Natural): PageSpan =
result = PageSpan(startAddr: offset(baseAddr, page.consumedTo), if buffers.queue.len == 0:
endAddr: offset(baseAddr, page.writtenTo)) return default(PageSpan)
page.consumedTo = page.writtenTo
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 = func writableSpan*(page: PageRef): PageSpan =
let baseAddr = allocationStart(page) let baseAddr = allocationStart(page)
@ -115,6 +147,13 @@ func trackWrittenTo*(buffers: PageBuffers, spanHeadPos: ptr byte) =
var topPage = buffers.queue.peekLast var topPage = buffers.queue.peekLast
topPage.writtenTo = distance(topPage.allocationStart, spanHeadPos) 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 = func addWritablePage*(buffers: PageBuffers, pageSize: Natural): PageRef =
trackWrittenToEnd(buffers) trackWrittenToEnd(buffers)
result = PageRef(data: allocRef newString(pageSize)) result = PageRef(data: allocRef newString(pageSize))
@ -136,25 +175,6 @@ template getWritableSpan*(buffers: PageBuffers): PageSpan =
let page = getWritablePage(buffers, buffers.pageSize) let page = getWritablePage(buffers, buffers.pageSize)
writableSpan(page) 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 = func stringFromBytes(src: pointer, srcLen: Natural): string =
result = newString(srcLen) result = newString(srcLen)
copyMem(addr result[0], src, srcLen) copyMem(addr result[0], src, srcLen)

View file

@ -206,8 +206,19 @@ proc memFileInput*(filename: string, mappedSize = -1, offset = 0): InputStreamHa
endAddr: offset(head, mappedSize)), endAddr: offset(head, mappedSize)),
file: memFile) 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 = 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 = template readableNow*(s: AsyncInputStream): bool =
readableNow InputStream(s) readableNow InputStream(s)
@ -229,8 +240,7 @@ proc readOnce*(sp: AsyncInputStream): Future[Natural] {.async.} =
disconnectInputDevice(s) disconnectInputDevice(s)
if result > 0 and s.span.len == 0: if result > 0 and s.span.len == 0:
s.buffers.nextReadableSpan(s.span) getNewSpan s
s.spanEndPos += s.span.len
proc timeoutToNextByteImpl(s: AsyncInputStream, proc timeoutToNextByteImpl(s: AsyncInputStream,
deadline: Future): Future[bool] {.async.} = deadline: Future): Future[bool] {.async.} =
@ -261,32 +271,11 @@ proc closeAsync*(s: AsyncInputStream) {.async.} =
# TODO: End of purely async interface # 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 = func getBestContiguousRunway(s: InputStream): Natural =
result = s.span.len result = s.span.len
if result == 0: if result == 0 and s.buffers != nil:
if s.buffers != nil and s.buffers.len > 1: getNewSpan s
flipPage s result = s.span.len
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
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
@ -297,7 +286,7 @@ func totalUnconsumedBytes*(s: InputStream): Natural =
else: s.buffers.totalBufferedBytes else: s.buffers.totalBufferedBytes
if localRunway == 0 and runwayInBuffers > 0: if localRunway == 0 and runwayInBuffers > 0:
flipPage s getNewSpan s
localRunway + runwayInBuffers localRunway + runwayInBuffers
@ -305,6 +294,38 @@ template totalUnconsumedBytes*(s: AsyncInputStream): Natural =
## Alias for InputStream.totalUnconsumedBytes ## Alias for InputStream.totalUnconsumedBytes
totalUnconsumedBytes InputStream(s) 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( let fileInputVTable = InputStreamVTable(
readSync: proc (s: InputStream, dst: pointer, dstLen: Natural): Natural readSync: proc (s: InputStream, dst: pointer, dstLen: Natural): Natural
{.nimcall, gcsafe, raises: [IOError, Defect].} = {.nimcall, gcsafe, raises: [IOError, Defect].} =
@ -384,24 +405,30 @@ template len*(s: AsyncInputStream): Option[Natural] =
len InputStream(s) len InputStream(s)
func memoryInput*(buffers: PageBuffers): InputStreamHandle = func memoryInput*(buffers: PageBuffers): InputStreamHandle =
var spanEndPos = Natural 0
var span = if buffers.len == 0: default(PageSpan) var span = if buffers.len == 0: default(PageSpan)
else: obtainReadableSpan buffers.queue[0] else: buffers.obtainReadableSpan(spanEndPos)
makeHandle InputStream(buffers: buffers, makeHandle InputStream(buffers: buffers,
span: span, span: span,
spanEndPos: span.len) spanEndPos: spanEndPos)
func memoryInput*(data: openarray[byte]): InputStreamHandle = func memoryInput*(data: openarray[byte]): InputStreamHandle =
let let stream = if data.len > 0:
buffers = initPageBuffers(data.len) let
page = buffers.addWritablePage(data.len) buffers = initPageBuffers(data.len)
pageSpan = page.fullSpan 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, InputStream(buffers: buffers,
span: pageSpan, span: pageSpan,
spanEndPos: data.len) spanEndPos: data.len)
else:
InputStream()
makeHandle stream
func memoryInput*(data: openarray[char]): InputStreamHandle = func memoryInput*(data: openarray[char]): InputStreamHandle =
memoryInput charsToBytes(data) memoryInput charsToBytes(data)
@ -409,9 +436,9 @@ func memoryInput*(data: openarray[char]): InputStreamHandle =
proc resetBuffers*(s: InputStream, buffers: PageBuffers) = proc resetBuffers*(s: InputStream, buffers: PageBuffers) =
# This should be used only on safe memory input streams # This should be used only on safe memory input streams
fsAssert s.vtable == nil and s.buffers != nil and buffers.len > 0 fsAssert s.vtable == nil and s.buffers != nil and buffers.len > 0
s.spanEndPos = 0
s.buffers = buffers s.buffers = buffers
s.span = obtainReadableSpan buffers.queue[0] getNewSpan s
s.spanEndPos = s.span.len
proc continueAfterRead(s: InputStream, bytesRead: Natural): bool = proc continueAfterRead(s: InputStream, bytesRead: Natural): bool =
# Please note that this is extracted into a proc only to reduce the code # 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) disconnectInputDevice(s)
if bytesRead > 0: if bytesRead > 0:
s.buffers.nextReadableSpan(s.span) getNewSpan s
s.spanEndPos += s.span.len
return true return true
else: else:
return false 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` # 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 # will be `nil` for a memFile. If we've reached here, this is the
# end of the memory buffer, so we can signal EOF: # 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 false
else: else:
# There might be additional pages in our buffer queue. If so, we # There might be additional pages in our buffer queue. If so, we
# just jump to the next one: # just jump to the next one:
if s.buffers.len > 1: getNewSpan s
flipPage s if hasRunway(s.span):
true true
else: elif s.vtable != nil and s.vtable.readOp != nil:
# We ask our input device to populate our page queue with newly # We ask our input device to populate our page queue with newly
# read pages. The state of the queue afterwards will tell us if # read pages. The state of the queue afterwards will tell us if
# the read was successful. In `continueAfterRead`, we examine if # the read was successful. In `continueAfterRead`, we examine if
# EOF was reached, but please note that some data might have been # EOF was reached, but please note that some data might have been
# read anyway: # read anyway:
continueAfterRead(s, awaiter s.vtable.readOp(s, nil, 0)) continueAfterRead(s, awaiter s.vtable.readOp(s, nil, 0))
else:
false
proc bufferMoreDataSync(s: InputStream): bool = proc bufferMoreDataSync(s: InputStream): bool =
# This proc exists only to avoid inlining of the code of # This proc exists only to avoid inlining of the code of
@ -518,8 +546,7 @@ template readable*(sp: AsyncInputStream): bool =
func continueAfterReadN(s: InputStream, func continueAfterReadN(s: InputStream,
runwayBeforeRead, bytesRead: Natural) = runwayBeforeRead, bytesRead: Natural) =
if runwayBeforeRead == 0 and bytesRead > 0: if runwayBeforeRead == 0 and bytesRead > 0:
s.buffers.nextReadableSpan(s.span) getNewSpan s
s.spanEndPos += s.span.len
template readableNImpl(s, n, awaiter, readOp: untyped): bool = template readableNImpl(s, n, awaiter, readOp: untyped): bool =
let runway = totalUnconsumedBytes(s) let runway = totalUnconsumedBytes(s)
@ -590,27 +617,22 @@ template readable*(sp: AsyncInputStream, np: int): bool =
readableNImpl(s, n, fsAwait, readAsync) 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 = template peek*(sp: InputStream): byte =
let s = sp let s = sp
if hasRunway(s.span): if hasRunway(s.span):
s.span.startAddr[] s.span.startAddr[]
else: else:
flipPage s getNewSpanOrDieTrying s
s.span.startAddr[] s.span.startAddr[]
template peek*(s: AsyncInputStream): byte = template peek*(s: AsyncInputStream): byte =
peek InputStream(s) peek InputStream(s)
func readFromNewSpan(s: InputStream): byte =
getNewSpanOrDieTrying s
result = s.span.startAddr[]
bumpPointer s.span
template read*(sp: InputStream): byte = template read*(sp: InputStream): byte =
let s = sp let s = sp
if hasRunway(s.span): if hasRunway(s.span):
@ -618,7 +640,7 @@ template read*(sp: InputStream): byte =
bumpPointer(s.span) bumpPointer(s.span)
res res
else: else:
flipPageRead s readFromNewSpan s
template read*(s: AsyncInputStream): byte = template read*(s: AsyncInputStream): byte =
read InputStream(s) read InputStream(s)
@ -636,7 +658,7 @@ proc advance*(s: InputStream) =
if hasRunway(s.span): if hasRunway(s.span):
bumpPointer s.span bumpPointer s.span
else: else:
flipPage s getNewSpan s
proc advance*(s: InputStream, n: Natural) = proc advance*(s: InputStream, n: Natural) =
# TODO This is silly, implement it properly # 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: if s.buffers != nil:
# Since we reached the end of the current page, # 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? # TODO: what if the page was extended?
if s.buffers.len > 0: if s.buffers.len > 0:
@ -717,21 +739,28 @@ proc drainBuffersInto*(s: InputStream, dstAddr: ptr byte, dstLen: Natural): Natu
template readIntoExImpl(s: InputStream, template readIntoExImpl(s: InputStream,
dst: ptr byte, dstLen: Natural, dst: ptr byte, dstLen: Natural,
awaiter, readOp: untyped): Natural = awaiter, readOp: untyped): Natural =
var bytesRead = drainBuffersInto(s, dst, dstLen) let totalBytesDrained = drainBuffersInto(s, dst, dstLen)
var bytesDeficit = (dstLen - totalBytesDrained)
while bytesRead < dstLen: if bytesDeficit > 0:
let var adjustedDst = offset(dst, totalBytesDrained)
bytesDeficit = dstLen - bytesRead
adjustedDst = offset(dst, bytesRead)
bytesRead += awaiter s.vtable.readOp(s, adjustedDst, bytesDeficit) while true:
let newBytesRead = awaiter s.vtable.readOp(s, adjustedDst, bytesDeficit)
if s.buffers.eofReached: s.spanEndPos += newBytesRead
disconnectInputDevice(s) bytesDeficit -= newBytesRead
break
s.spanEndPos += bytesRead if s.buffers.eofReached:
bytesRead disconnectInputDevice(s)
break
if bytesDeficit == 0:
break
adjustedDst = offset(dst, newBytesRead)
dstLen - bytesDeficit
proc readIntoEx*(s: InputStream, dst: var openarray[byte]): int = proc readIntoEx*(s: InputStream, dst: var openarray[byte]): int =
## Read data into the destination buffer. ## Read data into the destination buffer.
@ -854,21 +883,3 @@ proc pos*(s: InputStream): int {.inline.} =
template pos*(s: AsyncInputStream): int = template pos*(s: AsyncInputStream): int =
pos InputStream(s) 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)

View file

@ -172,8 +172,9 @@ func pipeInput*(source: InputStream,
func pipeInput*(buffers: PageBuffers, func pipeInput*(buffers: PageBuffers,
allowWaitFor = false, allowWaitFor = false,
source: InputStream = nil): AsyncInputStream = source: InputStream = nil): AsyncInputStream =
var spanEndPos = Natural 0
var span = if buffers.len == 0: default(PageSpan) var span = if buffers.len == 0: default(PageSpan)
else: obtainReadableSpan buffers.queue[0] else: buffers.obtainReadableSpan(spanEndPos)
AsyncInputStream LayeredInputStream( AsyncInputStream LayeredInputStream(
vtable: vtableAddr pipeInputVTable, vtable: vtableAddr pipeInputVTable,

View file

@ -53,10 +53,10 @@ procSuite "input stream":
var input = unsafeMemoryInput(str) var input = unsafeMemoryInput(str)
emptyInputTests "fileInput": emptyInputTests "fileInput":
var input = memFileInput("files" / "empty_file") var input = fileInput("files" / "empty_file")
emptyInputTests "memFileInput": emptyInputTests "memFileInput":
var input = fileInput("files" / "empty_file") var input = memFileInput("files" / "empty_file")
template asciiTableFileTest(name: string, body: untyped) = template asciiTableFileTest(name: string, body: untyped) =
test name & " of ascii table with regular pageSize": test name & " of ascii table with regular pageSize":
@ -137,11 +137,18 @@ procSuite "input stream":
check not fileExists(fileName) check not fileExists(fileName)
test "non-blocking reads": test "non-blocking reads":
let s = fileInput(asciiTableFile, pageSize = 20) let s = fileInput(asciiTableFile, pageSize = 100)
if s.readable: if s.readable(20):
s.withReadableRange(20, r): s.withReadableRange(20, r):
check r.readAll.len == 20 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": test "simple":
var input = repeat("1234 5678 90AB CDEF\n", 1000) var input = repeat("1234 5678 90AB CDEF\n", 1000)