Writer cursors that precisely track to which page they belong
This commit is contained in:
parent
5135173f60
commit
f1c376cd91
1 changed files with 77 additions and 37 deletions
|
|
@ -12,31 +12,47 @@ type
|
||||||
delayedWrites: int16
|
delayedWrites: int16
|
||||||
|
|
||||||
OutputStream* = object
|
OutputStream* = object
|
||||||
head, bufferEnd: ptr byte
|
cursor: WriteCursor
|
||||||
pages: Deque[OutputPage]
|
pages: Deque[OutputPage]
|
||||||
firstPage: int
|
|
||||||
endPos: int
|
endPos: int
|
||||||
outputResource: pointer
|
|
||||||
vtable: ptr OutputStreamVTable
|
vtable: ptr OutputStreamVTable
|
||||||
|
outputDevice: RootRef
|
||||||
|
|
||||||
DelayedWriteCursor* = object
|
WriteCursor* = object
|
||||||
head, bufferEnd: ptr byte
|
head, bufferEnd: ptr byte
|
||||||
stream: OutputStreamVar
|
stream: OutputStreamVar
|
||||||
page: int
|
absPageIdx: int
|
||||||
|
|
||||||
OutputStreamVar* = ref OutputStream
|
OutputStreamVar* = ref OutputStream
|
||||||
|
|
||||||
|
# Keep this temporary for backward-compatibility
|
||||||
|
DelayedWriteCursor* = WriteCursor
|
||||||
|
|
||||||
const
|
const
|
||||||
allocatorMetadata = 0 # TODO: Get this from Nim's allocator.
|
allocatorMetadata = 0 # TODO: Get this from Nim's allocator.
|
||||||
# The goal is to make perfect page-aligned allocations
|
# The goal is to make perfect page-aligned allocations
|
||||||
pageSize = 4096 - allocatorMetadata - 1 # 1 byte for the null terminator
|
pageSize = 4096 - allocatorMetadata - 1 # 1 byte for the null terminator
|
||||||
|
|
||||||
|
func remainingBytesToWrite*(c: WriteCursor): int {.inline.} =
|
||||||
|
distance(c.head, c.bufferEnd)
|
||||||
|
|
||||||
|
template relToAbsPageIdx(s: OutputStreamVar, idx: int): int =
|
||||||
|
# original code: s.firstPage + idx
|
||||||
|
idx - s.cursor.absPageIdx - 1
|
||||||
|
|
||||||
|
template absToRelPageIdx(s: OutputStreamVar, idx: int): int =
|
||||||
|
# original code: idx - s.firstPage
|
||||||
|
idx + s.cursor.absPageIdx + 1
|
||||||
|
|
||||||
|
func relPage(c: WriteCursor): int {.inline.} =
|
||||||
|
c.stream.absToRelPageIdx c.absPageIdx
|
||||||
|
|
||||||
proc addPage(s: OutputStreamVar) =
|
proc addPage(s: OutputStreamVar) =
|
||||||
s.pages.addLast OutputPage(buffer: newString(pageSize),
|
s.pages.addLast OutputPage(buffer: newString(pageSize),
|
||||||
delayedWrites: 0,
|
delayedWrites: 0,
|
||||||
startOffset: 0)
|
startOffset: 0)
|
||||||
s.head = cast[ptr byte](addr s.pages[s.pages.len - 1].buffer[0])
|
s.cursor.head = cast[ptr byte](addr s.pages[s.pages.len - 1].buffer[0])
|
||||||
s.bufferEnd = cast[ptr byte](shift(s.head, pageSize))
|
s.cursor.bufferEnd = cast[ptr byte](shift(s.cursor.head, pageSize))
|
||||||
s.endPos += pageSize
|
s.endPos += pageSize
|
||||||
|
|
||||||
proc init*(T: type OutputStream): ref OutputStream =
|
proc init*(T: type OutputStream): ref OutputStream =
|
||||||
|
|
@ -44,9 +60,11 @@ proc init*(T: type OutputStream): ref OutputStream =
|
||||||
result.vtable = nil
|
result.vtable = nil
|
||||||
result.pages = initDeque[OutputPage]()
|
result.pages = initDeque[OutputPage]()
|
||||||
result.addPage()
|
result.addPage()
|
||||||
|
result.cursor.absPageIdx = -1
|
||||||
|
result.cursor.stream = result
|
||||||
|
|
||||||
proc pos*(s: OutputStreamVar): int =
|
proc pos*(s: OutputStreamVar): int =
|
||||||
s.endPos - distance(s.head, s.bufferEnd)
|
s.endPos - s.cursor.remainingBytesToWrite
|
||||||
|
|
||||||
proc tryFlushing(s: OutputStreamVar) =
|
proc tryFlushing(s: OutputStreamVar) =
|
||||||
# TODO This is relevant when writing to files and layered streams (e.g. zip)
|
# TODO This is relevant when writing to files and layered streams (e.g. zip)
|
||||||
|
|
@ -57,35 +75,55 @@ proc tryFlushing(s: OutputStreamVar) =
|
||||||
# * The head and bufferEnd pointers point to the new top page
|
# * The head and bufferEnd pointers point to the new top page
|
||||||
s.addPage()
|
s.addPage()
|
||||||
|
|
||||||
proc append*(s: OutputStreamVar, b: byte) =
|
template isDelayedWrite(c: WriteCursor): bool =
|
||||||
if s.head == s.bufferEnd:
|
c.absPageIdx >= 0
|
||||||
s.tryFlushing()
|
|
||||||
|
|
||||||
s.head[] = b
|
proc append*(c: var WriteCursor, b: byte) =
|
||||||
s.head = shift(s.head, 1)
|
if c.head == c.bufferEnd:
|
||||||
|
# Delayed write cursors are not allowed to reach
|
||||||
|
# the end of the buffer and allocate new pages:
|
||||||
|
doAssert(not c.isDelayedWrite)
|
||||||
|
c.stream.tryFlushing()
|
||||||
|
|
||||||
template append*(s: OutputStreamVar, c: char) =
|
c.head[] = b
|
||||||
|
c.head = shift(c.head, 1)
|
||||||
|
|
||||||
|
template append*(c: var WriteCursor, x: char) =
|
||||||
bind append
|
bind append
|
||||||
append s, byte(c)
|
c.append byte(x)
|
||||||
|
|
||||||
proc append*(s: OutputStreamVar, bytes: openarray[byte]) =
|
proc append*(c: var WriteCursor, bytes: openarray[byte]) =
|
||||||
# TODO: this can use copyMem
|
# TODO: this can use copyMem
|
||||||
for b in bytes:
|
for b in bytes:
|
||||||
s.append b
|
c.append b
|
||||||
|
|
||||||
proc append*(s: OutputStreamVar, chars: openarray[char]) =
|
proc append*(c: var WriteCursor, chars: openarray[char]) =
|
||||||
# TODO: this can use copyMem
|
# TODO: this can use copyMem
|
||||||
for c in chars:
|
for x in chars:
|
||||||
s.append byte(c)
|
c.append byte(x)
|
||||||
|
|
||||||
template append*(s: OutputStreamVar, str: string) =
|
template appendMemCopy*(c: var WriteCursor, value: auto) =
|
||||||
s.append str.toOpenArrayByte(0, str.len - 1)
|
bind append
|
||||||
|
# TODO: add a check that this is a trivial type
|
||||||
|
c.append makeOpenArray(cast[ptr byte](unsafeAddr(value)), sizeof(value))
|
||||||
|
|
||||||
|
template append*(c: var WriteCursor, str: string) =
|
||||||
|
bind append
|
||||||
|
c.append str.toOpenArrayByte(0, str.len - 1)
|
||||||
|
|
||||||
|
template append*(s: OutputStreamVar, value: auto) =
|
||||||
|
bind append
|
||||||
|
s.cursor.append value
|
||||||
|
|
||||||
|
template appendMemCopy*(s: OutputStreamVar, value: auto) =
|
||||||
|
bind append
|
||||||
|
s.cursor.append value
|
||||||
|
|
||||||
proc flush*(s: OutputStreamVar) =
|
proc flush*(s: OutputStreamVar) =
|
||||||
s.vtable.finish(addr s[])
|
s.vtable.finish(addr s[])
|
||||||
|
|
||||||
proc getOutput*(s: OutputStreamVar, T: type string): string =
|
proc getOutput*(s: OutputStreamVar, T: type string): string =
|
||||||
s.pages[s.pages.len - 1].buffer.setLen(pageSize - distance(s.head, s.bufferEnd))
|
s.pages[s.pages.len - 1].buffer.setLen(pageSize - s.cursor.remainingBytesToWrite)
|
||||||
|
|
||||||
if s.pages.len == 1 and s.pages[0].startOffset == 0:
|
if s.pages.len == 1 and s.pages[0].startOffset == 0:
|
||||||
result.swap s.pages[0].buffer
|
result.swap s.pages[0].buffer
|
||||||
|
|
@ -110,10 +148,10 @@ proc flushDelayedPages*(s: OutputStreamVar) =
|
||||||
# TODO:
|
# TODO:
|
||||||
# Send to output
|
# Send to output
|
||||||
|
|
||||||
proc delayFixedSizeWrite*(s: OutputStreamVar, size: int): DelayedWriteCursor =
|
proc delayFixedSizeWrite*(s: OutputStreamVar, size: int): WriteCursor =
|
||||||
doAssert size < pageSize
|
doAssert size < pageSize
|
||||||
|
|
||||||
let remainingBytesInPage = distance(s.head, s.bufferEnd)
|
let remainingBytesInPage = s.cursor.remainingBytesToWrite
|
||||||
if size > remainingBytesInPage:
|
if size > remainingBytesInPage:
|
||||||
s.pages[s.pages.len - 1].buffer.setLen(pageSize - remainingBytesInPage)
|
s.pages[s.pages.len - 1].buffer.setLen(pageSize - remainingBytesInPage)
|
||||||
s.endPos -= remainingBytesInPage
|
s.endPos -= remainingBytesInPage
|
||||||
|
|
@ -122,12 +160,14 @@ proc delayFixedSizeWrite*(s: OutputStreamVar, size: int): DelayedWriteCursor =
|
||||||
let curPageIdx = s.pages.len - 1
|
let curPageIdx = s.pages.len - 1
|
||||||
inc s.pages[curPageIdx].delayedWrites
|
inc s.pages[curPageIdx].delayedWrites
|
||||||
|
|
||||||
result = DelayedWriteCursor(head: s.head, bufferEnd: s.head.shift(size),
|
result = WriteCursor(head: s.cursor.head,
|
||||||
page: s.firstPage + curPageIdx, stream: s)
|
bufferEnd: s.cursor.head.shift(size),
|
||||||
|
absPageIdx: s.relToAbsPageIdx(curPageIdx),
|
||||||
|
stream: s)
|
||||||
|
|
||||||
s.head = result.bufferEnd
|
s.cursor.head = result.bufferEnd
|
||||||
|
|
||||||
proc delayVarSizeWrite*(s: OutputStreamVar, maxSize: int): DelayedWriteCursor =
|
proc delayVarSizeWrite*(s: OutputStreamVar, maxSize: int): WriteCursor =
|
||||||
# TODO
|
# TODO
|
||||||
discard
|
discard
|
||||||
|
|
||||||
|
|
@ -136,23 +176,23 @@ proc decRef(x: var int16): int16 =
|
||||||
doAssert result >= 0
|
doAssert result >= 0
|
||||||
x = result
|
x = result
|
||||||
|
|
||||||
proc totalBytesWrittenAfterCursor*(cursor: DelayedWriteCursor): int =
|
proc totalBytesWrittenAfterCursor*(cursor: WriteCursor): int =
|
||||||
template s: auto = cursor.stream
|
template s: auto = cursor.stream
|
||||||
|
|
||||||
let
|
let
|
||||||
spanningPagesTotal = (cursor.stream.pages.len - cursor.page) * pageSize
|
relPageIdx = cursor.relPage
|
||||||
deductedFromFirstPage = distance(unsafeAddr s.pages[cursor.page].buffer[0],
|
spanningPagesTotal = (cursor.stream.pages.len - relPageIdx) * pageSize
|
||||||
|
deductedFromFirstPage = distance(unsafeAddr s.pages[relPageIdx].buffer[0],
|
||||||
cursor.bufferEnd)
|
cursor.bufferEnd)
|
||||||
deductedFromLastPage = distance(s.head, s.bufferEnd)
|
deductedFromLastPage = s.cursor.remainingBytesToWrite
|
||||||
|
|
||||||
spanningPagesTotal - deductedFromFirstPage - deductedFromLastPage
|
spanningPagesTotal - deductedFromFirstPage - deductedFromLastPage
|
||||||
|
|
||||||
proc endWrite*(cursor: DelayedWriteCursor, data: openarray[byte]) =
|
proc endWrite*(cursor: WriteCursor, data: openarray[byte]) =
|
||||||
if data.len != distance(cursor.head, cursor.bufferEnd):
|
doAssert data.len == cursor.remainingBytesToWrite
|
||||||
doAssert false
|
|
||||||
|
|
||||||
copyMem(cursor.head, unsafeAddr data[0], data.len)
|
copyMem(cursor.head, unsafeAddr data[0], data.len)
|
||||||
if cursor.stream.pages[cursor.page - cursor.stream.firstPage].delayedWrites.decRef <= 0:
|
if cursor.stream.pages[cursor.relPage].delayedWrites.decRef <= 0:
|
||||||
cursor.stream.flushDelayedPages()
|
cursor.stream.flushDelayedPages()
|
||||||
|
|
||||||
# Any stream
|
# Any stream
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue