Support for arbitrarily large fixed delayed writes
This commit is contained in:
parent
26e4f63c57
commit
2bc3941e76
2 changed files with 218 additions and 59 deletions
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue