diff --git a/faststreams/buffers.nim b/faststreams/buffers.nim index 7017a02..75111fb 100644 --- a/faststreams/buffers.nim +++ b/faststreams/buffers.nim @@ -151,9 +151,13 @@ proc setFauxEof*(buffers: PageBuffers, pos: Natural): Natural = proc restoreEof*(buffers: PageBuffers, pos: Natural) = buffers.fauxEofPos = pos +template allocWritablePage*(pageSize: Natural, writtenToParam: Natural = 0): auto = + PageRef(data: allocRef newString(pageSize), + writtenTo: writtenToParam) + func addWritablePage*(buffers: PageBuffers, pageSize: Natural): PageRef = trackWrittenToEnd(buffers) - result = PageRef(data: allocRef newString(pageSize)) + result = allocWritablePage(pageSize) buffers.queue.addLast result func getWritablePage*(buffers: PageBuffers, diff --git a/faststreams/outputs.nim b/faststreams/outputs.nim index 75f4113..53adc68 100644 --- a/faststreams/outputs.nim +++ b/faststreams/outputs.nim @@ -379,6 +379,7 @@ proc delayVarSizeWrite*(s: OutputStream, maxSize: Natural): VarSizeWriteCursor = startAddr = s.span.startAddr endAddr = offset(startAddr, maxSize) + inc s.extCursorsCount result = VarSizeWriteCursor WriteCursor( stream: s, span: PageSpan(startAddr: startAddr, endAddr: endAddr)) @@ -387,14 +388,17 @@ proc delayVarSizeWrite*(s: OutputStream, maxSize: Natural): VarSizeWriteCursor = s.span.startAddr = endAddr else: + trackWrittenTo(s.buffers, s.span.startAddr) + let nextPageSize = nextAlignedSize(maxSize, s.buffers.pageSize) - nextPage = s.buffers.addWritablePage(nextPageSize) + nextPage = allocWritablePage(nextPageSize, maxSize) nextPageSpan = nextPage.fullSpan cursorEndAddr = offset(nextPageSpan.startAddr, maxSize) - nextPage.consumedTo = -maxSize + s.buffers.queue.addLast nextPage + inc s.extCursorsCount result = VarSizeWriteCursor WriteCursor( stream: s, span: PageSpan(startAddr: nextPageSpan.startAddr, @@ -402,7 +406,7 @@ proc delayVarSizeWrite*(s: OutputStream, maxSize: Natural): VarSizeWriteCursor = s.span = PageSpan(startAddr: cursorEndAddr, endAddr: nextPageSpan.endAddr) - s.spanEndPos += nextPageSize + s.spanEndPos += nextPageSize - runway proc finalize*(cursor: var WriteCursor) = fsAssert cursor.stream.extCursorsCount > 0 @@ -419,15 +423,10 @@ proc finalWrite*(c: var VarSizeWriteCursor, data: openArray[byte]) = let overestimatedBytes = cursor.span.len - data.len fsAssert overestimatedBytes >= 0 + WriteCursor(c).stream.spanEndPos -= overestimatedBytes + for page in items(cursor.stream.buffers.queue): let baseAddr = page.allocationStart - if page.allocationEnd == cursor.span.endAddr: - # This is a page ending cursor - page.writtenTo = distance(baseAddr, cursor.span.startAddr) + data.len - copyMem(cursor.span.startAddr, unsafeAddr data[0], data.len) - finalize cursor - return - if cursor.span.startAddr == baseAddr: # This is page starting cursor page.consumedTo = overestimatedBytes @@ -435,6 +434,13 @@ proc finalWrite*(c: var VarSizeWriteCursor, data: openArray[byte]) = finalize cursor return + if page.readableEnd == cursor.span.endAddr: + # This is a page ending cursor + page.writtenTo = distance(baseAddr, cursor.span.startAddr) + data.len + copyMem(cursor.span.startAddr, unsafeAddr data[0], data.len) + finalize cursor + return + fsAssert false proc tryMovingToNextPage(c: var WriteCursor) = diff --git a/tests/test_outputs.nim b/tests/test_outputs.nim index dc461bb..ed1d079 100644 --- a/tests/test_outputs.nim +++ b/tests/test_outputs.nim @@ -1,7 +1,7 @@ {.used.} import - os, unittest, random, + os, unittest, random, strformat, stew/ranges/ptr_arith, ../faststreams, ../faststreams/textio @@ -17,10 +17,12 @@ proc repeat(b: byte, count: int): seq[byte] = result = newSeq[byte](count) for i in 0 ..< count: result[i] = b +const line = "123456789123456789123456789123456789\n\n\n\n\n" + proc randomBytes(n: int): seq[byte] = result.newSeq n for i in 0 ..< n: - result[i] = byte(rand(255)) + result[i] = byte(rand(line)) proc readAllAndClose(s: InputStream): seq[byte] = while s.readable: @@ -152,90 +154,32 @@ suite "output stream": checkOutputsMatch() + template undelayedOutput(content: seq[byte]) {.dirty.} = + nimSeq.add content + streamWritingToExistingBuffer.write content + test "delayed write": output "initial output\n" const delayedWriteContent = bytes "delayed write\n" - var cursor = memStream.delayFixedSizeWrite(delayedWriteContent.len) + var memCursor = memStream.delayFixedSizeWrite(delayedWriteContent.len) + var fileCursor = fileStream.delayVarSizeWrite(delayedWriteContent.len + 50) + let cursorStart = memStream.pos - nimSeq.add delayedWriteContent - fileStream.write delayedWriteContent - streamWritingToExistingBuffer.write delayedWriteContent + undelayedOutput delayedWriteContent var bytesWritten = 0 - for i, count in [2]: # 12, 342, 2121, 23, 1, 34012, 932]: + for i, count in [2, 12, 342, 2121, 23, 1, 34012, 932]: output repeat(byte(i), count) bytesWritten += count check memStream.pos - cursorStart == bytesWritten - cursor.finalWrite delayedWriteContent + memCursor.finalWrite delayedWriteContent + fileCursor.finalWrite delayedWriteContent checkOutputsMatch(skipUnbufferedFile = true) - 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.write delayedWrites[i].content[written ..< written + toWrite] - delayedWrites[i].written += toWrite - - if remaining - toWrite == 0: - finalize delayedWrites[i].cursor - if i != delayedWrites.len - 1: - swap(delayedWrites[i], delayedWrites[^1]) - delayedWrites.setLen(delayedWrites.len - 1) - - elif decision < 90: - # Normal write - memStream.write randomBytes - nimSeq.add randomBytes - - else: - # Create cursor - nimSeq.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 nimSeq.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.write dw.content[dw.written ..< dw.written + remaining] - finalize dw.cursor - - # The final outputs are the same - let resultsAreEqual = nimSeq == memStream.getOutput - check resultsAreEqual - test "float output": let basic: float64 = 12345.125 let small: float32 = 12345.125 @@ -248,3 +192,136 @@ suite "output stream": outputText tiny checkOutputsMatch() + +suite "randomized tests": + type + WriteTypes = enum + FixedSize + VarSize + Mixed + + DelayedWrite = object + isFixedSize: bool + fixedSizeCursor: WriteCursor + varSizeCursor: VarSizeWriteCursor + content: seq[byte] + written: int + + proc randomizedCursorsTestImpl(stream: OutputStream, + seed = 1000, + iterations = 1000, + minWriteSize = 500, + maxWriteSize = 1000, + writeTypes = Mixed, + varSizeVariance = 50): seq[byte] = + randomize seed + + var delayedWrites = newSeq[DelayedWrite]() + let writeSizeSpread = maxWriteSize - minWriteSize + + for i in 0 ..< iterations: + let decision = rand(100) + + if decision < 20: + # Write at some random cursor + if delayedWrites.len > 0: + let + i = rand(delayedWrites.len - 1) + written = delayedWrites[i].written + remaining = delayedWrites[i].content.len - written + toWrite = min(rand(remaining) + 10, remaining) + + if delayedWrites[i].isFixedSize: + delayedWrites[i].fixedSizeCursor.write delayedWrites[i].content[written ..< written + toWrite] + + delayedWrites[i].written += toWrite + + if remaining - toWrite == 0: + if delayedWrites[i].isFixedSize: + finalize delayedWrites[i].fixedSizeCursor + else: + finalWrite delayedWrites[i].varSizeCursor, delayedWrites[i].content + + if i != delayedWrites.len - 1: + swap(delayedWrites[i], delayedWrites[^1]) + delayedWrites.setLen(delayedWrites.len - 1) + + continue + + let + size = rand(writeSizeSpread) + minWriteSize + randomBytes = randomBytes(size) + + if decision < 90: + # Normal write + result.add randomBytes + stream.write randomBytes + + else: + # Create cursor + result.add randomBytes + + let isFixedSize = case writeTypes + of FixedSize: true + of VarSize: false + of Mixed: rand(10) > 3 + + if isFixedSize: + let cursor = stream.delayFixedSizeWrite(randomBytes.len) + + delayedWrites.add DelayedWrite( + fixedSizeCursor: cursor, + content: randomBytes, + written: 0, + isFixedSize: true) + else: + let + overestimatedBytes = rand(varSizeVariance) + cursorSize = randomBytes.len + overestimatedBytes + cursor = stream.delayVarSizeWrite(cursorSize) + + delayedWrites.add DelayedWrite( + varSizeCursor: cursor, + content: randomBytes, + written: 0, + isFixedSize: false) + + # Write all unwritten data to all outstanding cursors + if stream != nil: + for dw in mitems(delayedWrites): + if dw.isFixedSize: + let remaining = dw.content.len - dw.written + dw.fixedSizeCursor.write dw.content[dw.written ..< dw.written + remaining] + finalize dw.fixedSizeCursor + else: + dw.varSizeCursor.finalWrite dw.content + + template randomizedCursorsTest(streamExpr: OutputStreamHandle, + writeTypesExpr: WriteTypes, + varSizeVarianceExpr: int, + customChecks: untyped = nil) = + const testName = "randomized cursor test [" & astToStr(streamExpr) & + ";writes=" & $writeTypesExpr & ",variance=" & $varSizeVarianceExpr & "]" + test testName: + let s = streamExpr + var referenceResult = randomizedCursorsTestImpl(stream = s, + writeTypes = writeTypesExpr, + varSizeVariance = varSizeVarianceExpr) + + when astToStr(customChecks) == "nil": + let streamResult = s.getOutput() + let resultsMatch = streamResult == referenceResult + when false: + if not resultsMatch: + writeFile("reference-result.txt", referenceResult) + writeFile("stream-result.txt", streamResult) + check resultsMatch + else: + customChecks + + check referenceResult.len == s.pos + + randomizedCursorsTest(memoryOutput(), FixedSize, 0) + randomizedCursorsTest(memoryOutput(), VarSize, 100) + randomizedCursorsTest(memoryOutput(pageSize = 10), Mixed, 10) +