From 2bc3941e7626e3b4e9c05bc7c2c371a7aee273f0 Mon Sep 17 00:00:00 2001 From: Zahary Karadjov Date: Fri, 26 Jul 2019 23:52:16 +0300 Subject: [PATCH] Support for arbitrarily large fixed delayed writes --- faststreams/output_stream.nim | 196 +++++++++++++++++++++++++--------- tests/test_output_stream.nim | 81 ++++++++++++-- 2 files changed, 218 insertions(+), 59 deletions(-) diff --git a/faststreams/output_stream.nim b/faststreams/output_stream.nim index 2f11af0..0c62f15 100644 --- a/faststreams/output_stream.nim +++ b/faststreams/output_stream.nim @@ -6,7 +6,7 @@ type buffer: string startOffset: int - OutputStream* = object of RootObj + OutputStream* = object cursor*: WriteCursor pages: Deque[OutputPage] endPos: int @@ -15,6 +15,7 @@ type extCursorsCount: int pageSize: int maxWriteSize*: int + minWriteSize*: int WriteCursor* = object head, bufferEnd: ptr byte @@ -25,38 +26,43 @@ type OutputStreamVar* = ref OutputStream + MemOutputStream* = distinct OutputStreamVar + # Writing to a MemOutputStream produces no side-effects + # Keep this temporary for backward-compatibility DelayedWriteCursor* = WriteCursor VarSizeWriteCursor* = distinct WriteCursor OutputStreamVTable* = object + # TODO - the noSideEffect is temporary here until we switch to Nim 0.20.2 + # where noSideEffects overrides may make it possible to implement + # the MemOutputStream handling. writePage*: proc (s: OutputStreamVar, page: openarray[byte]) - {.nimcall, gcsafe, raises: [IOError, Defect] .} + {.noSideEffect, nimcall, gcsafe, raises: [IOError, Defect] .} flush*: proc (s: OutputStreamVar) - {.nimcall, gcsafe, raises: [IOError, Defect].} + {.noSideEffect, nimcall, gcsafe, raises: [IOError, Defect].} const allocatorMetadata = 0 # TODO: Get this from Nim's allocator. # The goal is to make perfect page-aligned allocations defaultPageSize = 4096 - allocatorMetadata - 1 # 1 byte for the null terminator -func canExtendOutput(c: var WriteCursor): bool {.inline.} = - # ATTENTION! - # The `var` modifier is intentional here. Nim tries to optimise small values - # such as WriteCursors by copying them instead of passing them by reference. - # This will break the address checks below. - # - # Only the original stream cursor is allowed to write past its buffer end by - # allocating new memory pages. Streams writing to pre-allocated existing - # buffers cannot be grown as well: - unsafeAddr(c) == unsafeAddr(c.stream.cursor) and c.stream.pageSize > 0 +template canExtendOutput(s: OutputStreamVar): bool = + # Streams writing to pre-allocated existing buffers cannot be grown + s.pageSize > 0 + +template isExternalCursor(c: var WriteCursor): bool = + # Is this the original stream cursor or is it one created by a "delayed write" + addr(c) != addr(c.stream.cursor) func remainingBytesToWrite*(c: var WriteCursor): int {.inline.} = distance(c.head, c.bufferEnd) proc flipPage(s: OutputStreamVar) = s.cursor.head = cast[ptr byte](addr s.pages[s.pages.len - 1].buffer[0]) + # TODO: There is an assumption here and elsewhere that `s.pages[^1]` has + # a length equal to `s.pageSize` s.cursor.bufferEnd = cast[ptr byte](shift(s.cursor.head, s.pageSize)) s.endPos += s.pageSize @@ -65,9 +71,13 @@ proc addPage(s: OutputStreamVar) = startOffset: 0) s.flipPage -proc initWithSinglePage*(s: OutputStreamVar, pageSize, maxWriteSize: int) = +proc initWithSinglePage*(s: OutputStreamVar, + pageSize: int, + maxWriteSize: int, + minWriteSize = 1) = s.pageSize = pageSize s.maxWriteSize = maxWriteSize + s.minWriteSize = minWriteSize s.pages = initDeque[OutputPage]() s.addPage s.cursor.stream = s @@ -77,25 +87,28 @@ proc init*(T: type OutputStream, new result result.initWithSinglePage pageSize, high(int) -let FileStreamVTable = OutputStreamVTable( - writePage: proc (s: OutputStreamVar, data: openarray[byte]) {.nimcall, gcsafe.} = - var output = FileOutput(s.outputDevice) - var written = output.file.writeBuffer(unsafeAddr data[0], data.len) - if written != data.len: - raise newException(IOError, "Failed to write OutputStream page.") - , - flush: proc (s: OutputStreamVar) {.nimcall, gcsafe.} = - var output = FileOutput(s.outputDevice) - flushFile output.file -) +when false: + # TODO: revisit this when we switch to Nim 0.20.2 and we have working + # noSideEffect overrides. + let FileStreamVTable = OutputStreamVTable( + writePage: proc (s: OutputStreamVar, data: openarray[byte]) {.nimcall, gcsafe.} = + var output = FileOutput(s.outputDevice) + var written = output.file.writeBuffer(unsafeAddr data[0], data.len) + if written != data.len: + raise newException(IOError, "Failed to write OutputStream page.") + , + flush: proc (s: OutputStreamVar) {.nimcall, gcsafe.} = + var output = FileOutput(s.outputDevice) + flushFile output.file + ) -proc init*(T: type OutputStream, - filename: string, - pageSize = defaultPageSize): ref OutputStream = - new result - result.outputDevice = FileOutput(file: open(filename, fmWrite)) - result.vtable = unsafeAddr FileStreamVTable - result.initWithSinglePage pageSize, high(int) + proc init*(T: type OutputStream, + filename: string, + pageSize = defaultPageSize): ref OutputStream = + new result + result.outputDevice = FileOutput(file: open(filename, fmWrite)) + result.vtable = unsafeAddr FileStreamVTable + result.initWithSinglePage pageSize, high(int) proc init*(T: type OutputStream, buffer: pointer, len: int): ref OutputStream = @@ -162,10 +175,41 @@ proc tryFlushing(s: OutputStreamVar) {.inline.} = else: s.addPage +func endAddr(s: string): ptr byte {.inline.} = + let a = unsafeAddr s[0] + shift(cast[ptr byte](a), s.len) + +template startAddr(s: string): ptr byte = + cast[ptr byte](unsafeAddr s[0]) + +func boundingAddrs(s: string): (ptr byte, ptr byte) {.inline.} = + (startAddr s, endAddr s) + +proc findNextPage(c: var WriteCursor): int = + let cursorBufferEnd = c.bufferEnd + for i in 0 .. c.stream.pages.len - 2: + let pageEnd = endAddr c.stream.pages[i].buffer + if cursorBufferEnd == pageEnd: + return i + 1 + + doAssert false # There is no next page the cursor can move to + +proc moveToPage(c: var WriteCursor, p: var OutputPage) = + doAssert p.startOffset > 0 + c.head = cast[ptr byte](unsafeAddr p.buffer[0]) + c.bufferEnd = shift(c.head, p.startOffset) + p.startOffset = 0 + +proc moveToNextPage(c: var WriteCursor) = + c.moveToPage c.stream.pages[c.findNextPage()] + proc append*(c: var WriteCursor, b: byte) = if c.head == c.bufferEnd: - doAssert c.canExtendOutput - c.stream.tryFlushing() + doAssert c.stream.canExtendOutput + if c.isExternalCursor: + c.moveToNextPage() + else: + c.stream.tryFlushing() c.head[] = b c.head = shift(c.head, 1) @@ -203,6 +247,7 @@ proc handleLongAppend*(c: var WriteCursor, bytes: openarray[byte]) = pageRemaining = c.remainingBytesToWrite inputPos = unsafeAddr bytes[0] inputLen = bytes.len + stream = c.stream template reduceInput(delta: int) = inputPos = shift(inputPos, delta) @@ -210,20 +255,55 @@ proc handleLongAppend*(c: var WriteCursor, bytes: openarray[byte]) = # Since the input is longer, we first make sure that the top-most # page is filled to the top: - doAssert c.canExtendOutput + doAssert c.stream.canExtendOutput copyMem(c.head, inputPos, pageRemaining) reduceInput pageRemaining - if c.stream.vtable != nil and c.stream.extCursorsCount == 0: + if c.isExternalCursor: + var + totalPages = stream.pages.len + nextPageIdx = c.findNextPage + + while nextPageIdx < totalPages: + let + pageStart = startAddr stream.pages[nextPageIdx].buffer + pageRunway = stream.pages[nextPageIdx].startOffset + pageLen = stream.pageSize + + doAssert pageRunway > 0 + stream.pages[nextPageIdx].startOffset = 0 + + if pageRunway < pageLen: + doAssert inputLen <= pageRunway + copyMem(pageStart, inputPos, inputLen) + c.head = shift(pageStart, inputLen) + c.bufferEnd = shift(pageStart, pageRunway) + return + else: + if inputLen <= pageLen: + copyMem(pageStart, inputPos, inputLen) + c.head = shift(pageStart, inputLen) + c.bufferEnd = shift(pageStart, pageLen) + return + else: + copyMem(pageStart, inputPos, pageLen) + reduceInput pageLen + inc nextPageIdx + + doAssert false # If we reached here, this means that we've ran out + # of pages, so this is a write past the end of the + # pre-allocated space for the delayed write. + + elif c.stream.vtable != nil and c.stream.extCursorsCount == 0: # This stream has an output device and we are ready to flush # all the pending pages. One fresh page will be left on top. # The input is yet to be written: - c.stream.writePendingPagesAndLeaveOne + stream.writePendingPagesAndLeaveOne # This will directly send our input to the output device. # Since the output device has a preference for pageSize and # maxWriteSize, we'll send some full pages to it and then # some bytes will be written to the fresh page created above: - c.stream.writeDataAsPages(inputPos, inputLen) + stream.writeDataAsPages(inputPos, inputLen) else: # We are not ready to flush, so we must create pending pages. # We'll try to create them as large as possible: @@ -232,11 +312,11 @@ proc handleLongAppend*(c: var WriteCursor, bytes: openarray[byte]) = # We know how much the endPos will advance, but please note that # it may be corrected later if we end up writing a portion of the # input to an incomplete page: - c.stream.endPos += inputLen + stream.endPos += inputLen # Try to create big pages until we have more data: while inputLen > maxPageSize: - c.stream.pages.addLast OutputPage( + stream.pages.addLast OutputPage( buffer: newStringFromBytes(inputPos, maxPageSize), startOffset: 0) reduceInput maxPageSize @@ -246,22 +326,22 @@ proc handleLongAppend*(c: var WriteCursor, bytes: openarray[byte]) = # one final oversized page and then we leave one empty fresh page where # the writing will continue: if inputLen > c.stream.pageSize: - c.stream.pages.addLast OutputPage( + stream.pages.addLast OutputPage( buffer: newStringFromBytes(inputPos, inputLen), startOffset: 0) - c.stream.addPage + stream.addPage else: # We don't have enough remaining bytes for a full page, so we'll just # allocate a new empty page and we'll write the remaining input there. # This will also reset the cursor to the start of the page: - c.stream.addPage + stream.addPage copyMem(c.head, inputPos, inputLen) c.head = shift(c.head, inputLen) # We must correct endPos, because it must mark the end of the top-most # page. Since `addPage` advances the endPos as well and our remaining # input was written to the newly created page, our initial increase of # endPos was overestimated: - c.stream.endPos -= inputLen + stream.endPos -= inputLen proc append*(c: var WriteCursor, bytes: openarray[byte]) {.inline.} = # We have a short inlinable function handling the case when the input is @@ -269,11 +349,10 @@ proc append*(c: var WriteCursor, bytes: openarray[byte]) {.inline.} = # page is full: let pageRemaining = c.remainingBytesToWrite - inputPos = unsafeAddr bytes[0] inputLen = bytes.len if inputLen <= pageRemaining: - copyMem(c.head, inputPos, inputLen) + copyMem(c.head, unsafeAddr bytes[0], inputLen) c.head = shift(c.head, inputLen) else: handleLongAppend(c, bytes) @@ -335,15 +414,28 @@ proc createCursor(s: OutputStreamVar, size: int): WriteCursor = s.cursor.head = result.bufferEnd -proc delayFixedSizeWrite*(s: OutputStreamVar, size: int): WriteCursor = +proc delayFixedSizeWrite*(s: OutputStreamVar, size: Natural): WriteCursor = let remainingBytesInPage = s.cursor.remainingBytesToWrite - if size > remainingBytesInPage: - doAssert size < s.pageSize - s.finishPageEarly remainingBytesInPage + if size <= remainingBytesInPage: + result = s.createCursor(size) + else: + result = s.createCursor(remainingBytesInPage) + var size = size - remainingBytesInPage + s.endPos += size + while size > s.pageSize: + s.pages.addLast OutputPage(buffer: newString(s.pageSize), + startOffset: s.pageSize) + size -= s.pageSize - s.createCursor(size) + s.pages.addLast OutputPage(buffer: newString(s.pageSize), + startOffset: size) -proc delayVarSizeWrite*(s: OutputStreamVar, maxSize: int): VarSizeWriteCursor = + let (pageStart, pageEnd) = boundingAddrs s.pages[s.pages.len - 1].buffer + s.cursor.head = shift(pageStart, size) + s.cursor.bufferEnd = pageEnd + s.endPos += (s.pageSize - size) + +proc delayVarSizeWrite*(s: OutputStreamVar, maxSize: Natural): VarSizeWriteCursor = doAssert maxSize < s.pageSize s.finishPageEarly s.cursor.remainingBytesToWrite VarSizeWriteCursor s.createCursor(maxSize) diff --git a/tests/test_output_stream.nim b/tests/test_output_stream.nim index 6495bc8..b7c4170 100644 --- a/tests/test_output_stream.nim +++ b/tests/test_output_stream.nim @@ -1,5 +1,5 @@ import - os, unittest, + os, unittest, random, stew/ranges/ptr_arith, ../faststreams @@ -14,12 +14,17 @@ proc repeat(b: byte, count: int): seq[byte] = result = newSeq[byte](count) for i in 0 ..< count: result[i] = b +proc randomBytes(n: int): seq[byte] = + result.newSeq n + for i in 0 ..< n: + result[i] = byte(rand(255)) + suite "output stream": setup: var memStream = OutputStream.init var altOutput: seq[byte] = @[] var tempFilePath = getTempDir() / "faststreams_testfile" - var fileStream = OutputStream.init tempFilePath + # var fileStream = OutputStream.init tempFilePath const bufferSize = 1000000 var buffer = alloc(bufferSize) @@ -32,18 +37,18 @@ suite "output stream": altOutput.add bytes(val) memStream.append val - fileStream.append val + # fileStream.append val existingBufferStream.append val template checkOutputsMatch = - fileStream.flush + # fileStream.flush let - fileContents = readFile(tempFilePath).string.bytes + # fileContents = readFile(tempFilePath).string.bytes memStreamContents = memStream.getOutput check altOutput == memStreamContents - check altOutput == fileContents + # check altOutput == fileContents check altOutput == makeOpenArray(cast[ptr byte](buffer), existingBufferStream.pos) @@ -66,7 +71,7 @@ suite "output stream": let cursorStart = memStream.pos altOutput.add delayedWriteContent - fileStream.append delayedWriteContent + # fileStream.append delayedWriteContent existingBufferStream.append delayedWriteContent var totalBytesWritten = 0 @@ -79,3 +84,65 @@ suite "output stream": checkOutputsMatch() + test "multi-page delayed writes": + randomize(1000) + + type + DelayedWrite = object + cursor: WriteCursor + content: seq[byte] + written: int + + var delayedWrites = newSeq[DelayedWrite]() + + for i in 0..50: + let + size = rand(8000) + 2000 + randomBytes = randomBytes(size) + decision = rand(100) + + if decision < 70: + # Write at some random cursor + if delayedWrites.len == 0: + continue + + let + i = rand(delayedWrites.len - 1) + written = delayedWrites[i].written + remaining = delayedWrites[i].content.len - written + toWrite = min(rand(remaining) + 10, remaining) + + delayedWrites[i].cursor.append delayedWrites[i].content[written ..< written + toWrite] + delayedWrites[i].written += toWrite + + if remaining - toWrite == 0: + dispose delayedWrites[i].cursor + if i != delayedWrites.len - 1: + swap(delayedWrites[i], delayedWrites[^1]) + delayedWrites.setLen(delayedWrites.len - 1) + + elif decision < 90: + # Normal write + memStream.append randomBytes + altOutput.add randomBytes + + else: + # Create cursor + altOutput.add randomBytes + delayedWrites.add DelayedWrite( + cursor: memStream.delayFixedSizeWrite(randomBytes.len), + content: randomBytes, + written: 0) + + # Check that the stream position is consistently tracked at every step + check altOutput.len == memStream.pos + + # Write all unwritten data to all outstanding cursors + for dw in mitems(delayedWrites): + let remaining = dw.content.len - dw.written + dw.cursor.append dw.content[dw.written ..< dw.written + remaining] + dispose dw.cursor + + # The final outputs are the same + check altOutput == memStream.getOutput +