Add tests for delayed var-size writes
This commit is contained in:
parent
4c5464a7e0
commit
5d7cad792f
3 changed files with 169 additions and 82 deletions
|
|
@ -151,9 +151,13 @@ proc setFauxEof*(buffers: PageBuffers, pos: Natural): Natural =
|
||||||
proc restoreEof*(buffers: PageBuffers, pos: Natural) =
|
proc restoreEof*(buffers: PageBuffers, pos: Natural) =
|
||||||
buffers.fauxEofPos = pos
|
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 =
|
func addWritablePage*(buffers: PageBuffers, pageSize: Natural): PageRef =
|
||||||
trackWrittenToEnd(buffers)
|
trackWrittenToEnd(buffers)
|
||||||
result = PageRef(data: allocRef newString(pageSize))
|
result = allocWritablePage(pageSize)
|
||||||
buffers.queue.addLast result
|
buffers.queue.addLast result
|
||||||
|
|
||||||
func getWritablePage*(buffers: PageBuffers,
|
func getWritablePage*(buffers: PageBuffers,
|
||||||
|
|
|
||||||
|
|
@ -379,6 +379,7 @@ proc delayVarSizeWrite*(s: OutputStream, maxSize: Natural): VarSizeWriteCursor =
|
||||||
startAddr = s.span.startAddr
|
startAddr = s.span.startAddr
|
||||||
endAddr = offset(startAddr, maxSize)
|
endAddr = offset(startAddr, maxSize)
|
||||||
|
|
||||||
|
inc s.extCursorsCount
|
||||||
result = VarSizeWriteCursor WriteCursor(
|
result = VarSizeWriteCursor WriteCursor(
|
||||||
stream: s,
|
stream: s,
|
||||||
span: PageSpan(startAddr: startAddr, endAddr: endAddr))
|
span: PageSpan(startAddr: startAddr, endAddr: endAddr))
|
||||||
|
|
@ -387,14 +388,17 @@ proc delayVarSizeWrite*(s: OutputStream, maxSize: Natural): VarSizeWriteCursor =
|
||||||
s.span.startAddr = endAddr
|
s.span.startAddr = endAddr
|
||||||
|
|
||||||
else:
|
else:
|
||||||
|
trackWrittenTo(s.buffers, s.span.startAddr)
|
||||||
|
|
||||||
let
|
let
|
||||||
nextPageSize = nextAlignedSize(maxSize, s.buffers.pageSize)
|
nextPageSize = nextAlignedSize(maxSize, s.buffers.pageSize)
|
||||||
nextPage = s.buffers.addWritablePage(nextPageSize)
|
nextPage = allocWritablePage(nextPageSize, maxSize)
|
||||||
nextPageSpan = nextPage.fullSpan
|
nextPageSpan = nextPage.fullSpan
|
||||||
cursorEndAddr = offset(nextPageSpan.startAddr, maxSize)
|
cursorEndAddr = offset(nextPageSpan.startAddr, maxSize)
|
||||||
|
|
||||||
nextPage.consumedTo = -maxSize
|
s.buffers.queue.addLast nextPage
|
||||||
|
|
||||||
|
inc s.extCursorsCount
|
||||||
result = VarSizeWriteCursor WriteCursor(
|
result = VarSizeWriteCursor WriteCursor(
|
||||||
stream: s,
|
stream: s,
|
||||||
span: PageSpan(startAddr: nextPageSpan.startAddr,
|
span: PageSpan(startAddr: nextPageSpan.startAddr,
|
||||||
|
|
@ -402,7 +406,7 @@ proc delayVarSizeWrite*(s: OutputStream, maxSize: Natural): VarSizeWriteCursor =
|
||||||
|
|
||||||
s.span = PageSpan(startAddr: cursorEndAddr,
|
s.span = PageSpan(startAddr: cursorEndAddr,
|
||||||
endAddr: nextPageSpan.endAddr)
|
endAddr: nextPageSpan.endAddr)
|
||||||
s.spanEndPos += nextPageSize
|
s.spanEndPos += nextPageSize - runway
|
||||||
|
|
||||||
proc finalize*(cursor: var WriteCursor) =
|
proc finalize*(cursor: var WriteCursor) =
|
||||||
fsAssert cursor.stream.extCursorsCount > 0
|
fsAssert cursor.stream.extCursorsCount > 0
|
||||||
|
|
@ -419,15 +423,10 @@ proc finalWrite*(c: var VarSizeWriteCursor, data: openArray[byte]) =
|
||||||
let overestimatedBytes = cursor.span.len - data.len
|
let overestimatedBytes = cursor.span.len - data.len
|
||||||
fsAssert overestimatedBytes >= 0
|
fsAssert overestimatedBytes >= 0
|
||||||
|
|
||||||
|
WriteCursor(c).stream.spanEndPos -= overestimatedBytes
|
||||||
|
|
||||||
for page in items(cursor.stream.buffers.queue):
|
for page in items(cursor.stream.buffers.queue):
|
||||||
let baseAddr = page.allocationStart
|
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:
|
if cursor.span.startAddr == baseAddr:
|
||||||
# This is page starting cursor
|
# This is page starting cursor
|
||||||
page.consumedTo = overestimatedBytes
|
page.consumedTo = overestimatedBytes
|
||||||
|
|
@ -435,6 +434,13 @@ proc finalWrite*(c: var VarSizeWriteCursor, data: openArray[byte]) =
|
||||||
finalize cursor
|
finalize cursor
|
||||||
return
|
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
|
fsAssert false
|
||||||
|
|
||||||
proc tryMovingToNextPage(c: var WriteCursor) =
|
proc tryMovingToNextPage(c: var WriteCursor) =
|
||||||
|
|
|
||||||
|
|
@ -1,7 +1,7 @@
|
||||||
{.used.}
|
{.used.}
|
||||||
|
|
||||||
import
|
import
|
||||||
os, unittest, random,
|
os, unittest, random, strformat,
|
||||||
stew/ranges/ptr_arith,
|
stew/ranges/ptr_arith,
|
||||||
../faststreams, ../faststreams/textio
|
../faststreams, ../faststreams/textio
|
||||||
|
|
||||||
|
|
@ -17,10 +17,12 @@ proc repeat(b: byte, count: int): seq[byte] =
|
||||||
result = newSeq[byte](count)
|
result = newSeq[byte](count)
|
||||||
for i in 0 ..< count: result[i] = b
|
for i in 0 ..< count: result[i] = b
|
||||||
|
|
||||||
|
const line = "123456789123456789123456789123456789\n\n\n\n\n"
|
||||||
|
|
||||||
proc randomBytes(n: int): seq[byte] =
|
proc randomBytes(n: int): seq[byte] =
|
||||||
result.newSeq n
|
result.newSeq n
|
||||||
for i in 0 ..< n:
|
for i in 0 ..< n:
|
||||||
result[i] = byte(rand(255))
|
result[i] = byte(rand(line))
|
||||||
|
|
||||||
proc readAllAndClose(s: InputStream): seq[byte] =
|
proc readAllAndClose(s: InputStream): seq[byte] =
|
||||||
while s.readable:
|
while s.readable:
|
||||||
|
|
@ -152,90 +154,32 @@ suite "output stream":
|
||||||
|
|
||||||
checkOutputsMatch()
|
checkOutputsMatch()
|
||||||
|
|
||||||
|
template undelayedOutput(content: seq[byte]) {.dirty.} =
|
||||||
|
nimSeq.add content
|
||||||
|
streamWritingToExistingBuffer.write content
|
||||||
|
|
||||||
test "delayed write":
|
test "delayed write":
|
||||||
output "initial output\n"
|
output "initial output\n"
|
||||||
const delayedWriteContent = bytes "delayed write\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
|
let cursorStart = memStream.pos
|
||||||
|
|
||||||
nimSeq.add delayedWriteContent
|
undelayedOutput delayedWriteContent
|
||||||
fileStream.write delayedWriteContent
|
|
||||||
streamWritingToExistingBuffer.write delayedWriteContent
|
|
||||||
|
|
||||||
var bytesWritten = 0
|
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)
|
output repeat(byte(i), count)
|
||||||
bytesWritten += count
|
bytesWritten += count
|
||||||
check memStream.pos - cursorStart == bytesWritten
|
check memStream.pos - cursorStart == bytesWritten
|
||||||
|
|
||||||
cursor.finalWrite delayedWriteContent
|
memCursor.finalWrite delayedWriteContent
|
||||||
|
fileCursor.finalWrite delayedWriteContent
|
||||||
|
|
||||||
checkOutputsMatch(skipUnbufferedFile = true)
|
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":
|
test "float output":
|
||||||
let basic: float64 = 12345.125
|
let basic: float64 = 12345.125
|
||||||
let small: float32 = 12345.125
|
let small: float32 = 12345.125
|
||||||
|
|
@ -248,3 +192,136 @@ suite "output stream":
|
||||||
outputText tiny
|
outputText tiny
|
||||||
|
|
||||||
checkOutputsMatch()
|
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)
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue