More comprehensive OutputStream; Basic support for stream cursors
This commit is contained in:
parent
31590a79ec
commit
7903f1680f
5 changed files with 164 additions and 49 deletions
|
|
@ -8,7 +8,8 @@ license = "Apache License 2.0"
|
||||||
skipDirs = @["tests"]
|
skipDirs = @["tests"]
|
||||||
|
|
||||||
requires "nim >= 0.17.0",
|
requires "nim >= 0.17.0",
|
||||||
"ranges"
|
"ranges",
|
||||||
|
"std_shims"
|
||||||
|
|
||||||
task test, "run tests":
|
task test, "run tests":
|
||||||
exec "nim c -r tests/test_output_stream.nim"
|
exec "nim c -r tests/test_output_stream.nim"
|
||||||
|
|
|
||||||
|
|
@ -12,14 +12,19 @@ type
|
||||||
CloseStreamProc = proc ()
|
CloseStreamProc = proc ()
|
||||||
|
|
||||||
ByteStream* = object of RootObj
|
ByteStream* = object of RootObj
|
||||||
head: ptr byte
|
head*: ptr byte
|
||||||
bufferSize: int
|
bufferSize: int
|
||||||
bufferStart, bufferEnd: ptr byte
|
bufferStart, bufferEnd: ptr byte
|
||||||
|
cursorsList: ptr StreamCursor
|
||||||
reader: StreamReader
|
reader: StreamReader
|
||||||
asyncReader: AsyncStreamReader
|
asyncReader: AsyncStreamReader
|
||||||
closeStream: CloseStreamProc
|
closeStream: CloseStreamProc
|
||||||
bufferEndPos: int
|
bufferEndPos: int
|
||||||
|
|
||||||
|
StreamCursor* = object
|
||||||
|
head, bufferEnd: ptr byte
|
||||||
|
nextCursor: ptr StreamCursor
|
||||||
|
|
||||||
BufferedStream*[size: static int] = object of ByteStream
|
BufferedStream*[size: static int] = object of ByteStream
|
||||||
buffer: array[size, byte]
|
buffer: array[size, byte]
|
||||||
|
|
||||||
|
|
@ -27,7 +32,12 @@ type
|
||||||
UnicodeStream* = distinct ByteStream
|
UnicodeStream* = distinct ByteStream
|
||||||
ObjectStream*[T] = distinct ByteStream
|
ObjectStream*[T] = distinct ByteStream
|
||||||
|
|
||||||
proc openFile*(filename: string): ByteStream =
|
# TODO: ByteStreamVar should become a `var` short-lived pointer once this is supported
|
||||||
|
ByteStreamVar* = ref ByteStream
|
||||||
|
AsciiStreamVar* = ref AsciiStream
|
||||||
|
|
||||||
|
proc openFile*(filename: string): ByteStreamVar =
|
||||||
|
new result
|
||||||
var memFile = memfiles.open(filename)
|
var memFile = memfiles.open(filename)
|
||||||
result.head = cast[ptr byte](memFile.mem)
|
result.head = cast[ptr byte](memFile.mem)
|
||||||
when debugHelpers:
|
when debugHelpers:
|
||||||
|
|
@ -46,12 +56,13 @@ proc init*(T: type ByteStream,
|
||||||
result.reader = reader
|
result.reader = reader
|
||||||
result.asyncReader = asyncReader
|
result.asyncReader = asyncReader
|
||||||
|
|
||||||
template memoryStream*(mem: openarray[byte]): ByteStream =
|
proc memoryStream*(mem: openarray[byte]): ByteStreamVar = # TODO: mark the result with `from mem`
|
||||||
ByteStream.init(mem)
|
new result
|
||||||
|
result[] = ByteStream.init(mem)
|
||||||
|
|
||||||
template memoryStream*(str: string): ByteStream =
|
proc memoryStream*(str: string): ByteStreamVar = # TODO: mark the result with `from str`
|
||||||
bind init
|
new result
|
||||||
init(ByteStream, str.toOpenArrayByte(0, str.len - 1))
|
result[] = init(ByteStream, str.toOpenArrayByte(0, str.len - 1))
|
||||||
|
|
||||||
proc init*(T: type BufferedStream,
|
proc init*(T: type BufferedStream,
|
||||||
reader = StreamReader(nil),
|
reader = StreamReader(nil),
|
||||||
|
|
@ -68,6 +79,15 @@ proc syncRead(s: var ByteStream): bool =
|
||||||
s.bufferEndPos += bytesRead
|
s.bufferEndPos += bytesRead
|
||||||
return false
|
return false
|
||||||
|
|
||||||
|
proc ensureBytes*(s: var ByteStream, n: int): bool =
|
||||||
|
if distance(s.head, s.bufferEnd) >= n:
|
||||||
|
return true
|
||||||
|
|
||||||
|
if s.reader == nil:
|
||||||
|
return false
|
||||||
|
|
||||||
|
doAssert false, "Multi-buffer reading will be implemented later"
|
||||||
|
|
||||||
proc eof*(s: var ByteStream): bool =
|
proc eof*(s: var ByteStream): bool =
|
||||||
if s.head != s.bufferEnd:
|
if s.head != s.bufferEnd:
|
||||||
return false
|
return false
|
||||||
|
|
@ -99,6 +119,14 @@ proc read*(s: var ByteStream): byte =
|
||||||
result = s.peek()
|
result = s.peek()
|
||||||
advance s
|
advance s
|
||||||
|
|
||||||
|
proc checkReadAhead(s: ByteStreamVar, n: int): ptr byte =
|
||||||
|
result = s.head
|
||||||
|
assert distance(s.head, s.bufferEnd) >= n
|
||||||
|
s.head = s.head.shift(n)
|
||||||
|
|
||||||
|
template readBytes*(s: ByteStreamVar, n: int): auto =
|
||||||
|
makeOpenArray(checkReadAhead(s, n), n)
|
||||||
|
|
||||||
proc next*(s: var ByteStream): Option[byte] =
|
proc next*(s: var ByteStream): Option[byte] =
|
||||||
if not s.eof:
|
if not s.eof:
|
||||||
result = some s.read()
|
result = some s.read()
|
||||||
|
|
|
||||||
|
|
@ -1,81 +1,155 @@
|
||||||
import
|
import
|
||||||
ranges/ptr_arith
|
deques, ranges/ptr_arith, std_shims/strings
|
||||||
|
|
||||||
type
|
type
|
||||||
OutputStreamVTable = tuple
|
OutputStreamVTable = tuple
|
||||||
prepareOutput: proc (s: ptr OutputStream, size: int) {.nimcall.}
|
prepareOutput: proc (s: ptr OutputStream, size: int) {.nimcall.}
|
||||||
finish: proc (s: ptr OutputStream) {.nimcall.}
|
finish: proc (s: ptr OutputStream) {.nimcall.}
|
||||||
|
|
||||||
OutputStream* = object {.inheritable.}
|
OutputPage = object
|
||||||
|
buffer: string
|
||||||
|
startOffset: int16
|
||||||
|
delayedWrites: int16
|
||||||
|
|
||||||
|
OutputStream* = object
|
||||||
head, bufferEnd: ptr byte
|
head, bufferEnd: ptr byte
|
||||||
|
pages: Deque[OutputPage]
|
||||||
|
firstPage: int
|
||||||
|
endPos: int
|
||||||
|
outputResource: pointer
|
||||||
vtable: ptr OutputStreamVTable
|
vtable: ptr OutputStreamVTable
|
||||||
|
|
||||||
MemoryOutputStream*[T] = object of OutputStream
|
DelayedWriteCursor* = object
|
||||||
output: T
|
head, bufferEnd: ptr byte
|
||||||
|
stream: OutputStreamVar
|
||||||
|
page: int
|
||||||
|
|
||||||
StringOutputStream* = MemoryOutputStream[string]
|
OutputStreamVar* = ref OutputStream
|
||||||
BytesOutputStream* = MemoryOutputStream[seq[byte]]
|
|
||||||
|
|
||||||
FileOutputStream* = object of OutputStream
|
|
||||||
file: File
|
|
||||||
|
|
||||||
const
|
const
|
||||||
pageSize = 4096
|
allocatorMetadata = 0 # TODO: Get this from Nim's allocator.
|
||||||
|
# The goal is to make perfect page-aligned allocations
|
||||||
|
pageSize = 4096 - allocatorMetadata - 1 # 1 byte for the null terminator
|
||||||
|
|
||||||
proc prepareOutputImpl[T](s: ptr OutputStream, size: int) =
|
proc addPage(s: OutputStreamVar) =
|
||||||
let s = cast[ptr T](s)
|
s.pages.addLast OutputPage(buffer: newString(pageSize),
|
||||||
let currentSize = s.output.len
|
delayedWrites: 0,
|
||||||
s.output.setLen(s.output.len + pageSize)
|
startOffset: 0)
|
||||||
s.head = cast[ptr byte](addr s.output[currentSize])
|
s.head = cast[ptr byte](addr s.pages[s.pages.len - 1].buffer[0])
|
||||||
s.bufferEnd = cast[ptr byte](shift(addr s.output[0], s.output.len))
|
s.bufferEnd = cast[ptr byte](shift(s.head, pageSize))
|
||||||
|
s.endPos += pageSize
|
||||||
|
|
||||||
proc finishImpl[T](s: ptr OutputStream) =
|
proc init*(T: type OutputStream): ref OutputStream =
|
||||||
let s = cast[ptr T](s)
|
new result
|
||||||
s.output.setLen(s.output.len - distance(s.head, s.bufferEnd))
|
result.vtable = nil
|
||||||
|
result.pages = initDeque[OutputPage]()
|
||||||
|
result.addPage()
|
||||||
|
|
||||||
proc init*(T: type MemoryOutputStream): T =
|
proc pos*(s: OutputStreamVar): int =
|
||||||
var vtable {.global.} = (prepareOutputImpl[T], finishImpl[T])
|
s.endPos - distance(s.head, s.bufferEnd)
|
||||||
result.vtable = addr vtable
|
|
||||||
result.vtable.prepareOutput(addr result, pageSize)
|
|
||||||
|
|
||||||
proc append*(s: var OutputStream, b: byte) =
|
proc tryFlushing(s: OutputStreamVar) =
|
||||||
|
# TODO This is relevant when writing to files and layered streams (e.g. zip)
|
||||||
|
# Post-conditions:
|
||||||
|
# * All completed pages are written
|
||||||
|
# * There is a fresh page ready for writing at the top
|
||||||
|
# (we can reuse a previously existing page for this)
|
||||||
|
# * The head and bufferEnd pointers point to the new top page
|
||||||
|
s.addPage()
|
||||||
|
|
||||||
|
proc append*(s: OutputStreamVar, b: byte) =
|
||||||
if s.head == s.bufferEnd:
|
if s.head == s.bufferEnd:
|
||||||
s.vtable.prepareOutput(addr s, pageSize)
|
s.tryFlushing()
|
||||||
|
|
||||||
s.head[] = b
|
s.head[] = b
|
||||||
s.head = shift(s.head, 1)
|
s.head = shift(s.head, 1)
|
||||||
|
|
||||||
template append*(s: var OutputStream, c: char) =
|
template append*(s: OutputStreamVar, c: char) =
|
||||||
s.append byte(c)
|
s.append byte(c)
|
||||||
|
|
||||||
proc flush*(s: var OutputStream) =
|
proc append*(s: OutputStreamVar, bytes: openarray[byte]) =
|
||||||
s.vtable.finish(addr s)
|
|
||||||
|
|
||||||
proc getOutput*(s: var MemoryOutputStream): auto =
|
|
||||||
flush s
|
|
||||||
shallow s.output
|
|
||||||
return s.output
|
|
||||||
|
|
||||||
proc append*(s: var OutputStream, bytes: openarray[byte]) =
|
|
||||||
# TODO: this can use copyMem
|
# TODO: this can use copyMem
|
||||||
for b in bytes:
|
for b in bytes:
|
||||||
s.append b
|
s.append b
|
||||||
|
|
||||||
proc append*(s: var OutputStream, chars: openarray[char]) =
|
proc append*(s: OutputStreamVar, chars: openarray[char]) =
|
||||||
# TODO: this can use copyMem
|
# TODO: this can use copyMem
|
||||||
for c in chars:
|
for c in chars:
|
||||||
s.append byte(c)
|
s.append byte(c)
|
||||||
|
|
||||||
template append*(s: var OutputStream, str: string) =
|
template append*(s: OutputStreamVar, str: string) =
|
||||||
s.append str.toOpenArrayByte(0, str.len - 1)
|
s.append str.toOpenArrayByte(0, str.len - 1)
|
||||||
|
|
||||||
|
proc flush*(s: OutputStreamVar) =
|
||||||
|
s.vtable.finish(addr s[])
|
||||||
|
|
||||||
|
proc getOutput*(s: OutputStreamVar, T: type string): string =
|
||||||
|
s.pages[s.pages.len - 1].buffer.setLen(pageSize - distance(s.head, s.bufferEnd))
|
||||||
|
|
||||||
|
if s.pages.len == 1 and s.pages[0].startOffset == 0:
|
||||||
|
result.swap s.pages[0].buffer
|
||||||
|
else:
|
||||||
|
result = newStringOfCap(s.pos)
|
||||||
|
for page in s.pages:
|
||||||
|
assert page.delayedWrites == 0
|
||||||
|
result.add page.buffer.toOpenArray(page.startOffset.int,
|
||||||
|
page.buffer.len - 1)
|
||||||
|
|
||||||
|
template getOutput*(s: OutputStreamVar, T: type seq[byte]): seq[byte] =
|
||||||
|
cast[seq[byte]](s.getOutput(string))
|
||||||
|
|
||||||
|
proc getOutput*(s: OutputStreamVar): seq[byte] =
|
||||||
|
# TODO: is the extra copy here optimized away?
|
||||||
|
# Turning this proc into a template creates problems at the moment.
|
||||||
|
s.getOutput(seq[byte])
|
||||||
|
|
||||||
|
proc flushDelayedPages*(s: OutputStreamVar) =
|
||||||
|
for i in 0 .. s.pages.len - 2:
|
||||||
|
if s.pages[i].delayedWrites > 0: return
|
||||||
|
# TODO:
|
||||||
|
# Send to output
|
||||||
|
|
||||||
|
proc delayFixedSizeWrite*(s: OutputStreamVar, size: int): DelayedWriteCursor =
|
||||||
|
assert size < pageSize
|
||||||
|
|
||||||
|
let remainingBytesInPage = distance(s.head, s.bufferEnd)
|
||||||
|
if size > remainingBytesInPage:
|
||||||
|
s.pages[s.pages.len - 1].buffer.setLen(pageSize - remainingBytesInPage)
|
||||||
|
s.endPos -= remainingBytesInPage
|
||||||
|
s.tryFlushing()
|
||||||
|
|
||||||
|
let curPageIdx = s.pages.len - 1
|
||||||
|
inc s.pages[curPageIdx].delayedWrites
|
||||||
|
|
||||||
|
result = DelayedWriteCursor(head: s.head, bufferEnd: s.head.shift(size),
|
||||||
|
page: s.firstPage + curPageIdx, stream: s)
|
||||||
|
|
||||||
|
s.head = result.bufferEnd
|
||||||
|
|
||||||
|
proc delayVarSizeWrite*(s: OutputStreamVar, maxSize: int): DelayedWriteCursor =
|
||||||
|
# TODO
|
||||||
|
discard
|
||||||
|
|
||||||
|
proc decRef(x: var int16): int16 =
|
||||||
|
result = x - 1
|
||||||
|
assert result >= 0
|
||||||
|
x = result
|
||||||
|
|
||||||
|
proc endWrite*(cursor: DelayedWriteCursor, data: openarray[byte]) =
|
||||||
|
if data.len != distance(cursor.head, cursor.bufferEnd):
|
||||||
|
doAssert false
|
||||||
|
|
||||||
|
copyMem(cursor.head, unsafeAddr data[0], data.len)
|
||||||
|
if cursor.stream.pages[cursor.page - cursor.stream.firstPage].delayedWrites.decRef <= 0:
|
||||||
|
cursor.stream.flushDelayedPages()
|
||||||
|
|
||||||
# Any stream
|
# Any stream
|
||||||
|
|
||||||
proc appendNumberImpl(s: var OutputStream, number: BiggestInt) =
|
proc appendNumberImpl(s: OutputStreamVar, number: BiggestInt) =
|
||||||
# TODO: don't allocate
|
# TODO: don't allocate
|
||||||
s.append $number
|
s.append $number
|
||||||
|
|
||||||
proc appendNumberImpl(s: var OutputStream, number: BiggestUInt) =
|
proc appendNumberImpl(s: OutputStreamVar, number: BiggestUInt) =
|
||||||
# TODO: don't allocate
|
# TODO: don't allocate
|
||||||
s.append $number
|
s.append $number
|
||||||
|
|
||||||
|
|
@ -85,7 +159,7 @@ template toBiggestRepr(i: SomeUnsignedInt): BiggestUInt =
|
||||||
template toBiggestRepr(i: SomeSignedInt): BiggestInt =
|
template toBiggestRepr(i: SomeSignedInt): BiggestInt =
|
||||||
BiggestInt(i)
|
BiggestInt(i)
|
||||||
|
|
||||||
template appendNumber*(s: var OutputStream, i: SomeInteger) =
|
template appendNumber*(s: OutputStreamVar, i: SomeInteger) =
|
||||||
# TODO: specify radix/base
|
# TODO: specify radix/base
|
||||||
appendNumberImpl(s, toBiggestRepr(i))
|
appendNumberImpl(s, toBiggestRepr(i))
|
||||||
|
|
||||||
|
|
|
||||||
12
tests/test_input_stream.nim
Normal file
12
tests/test_input_stream.nim
Normal file
|
|
@ -0,0 +1,12 @@
|
||||||
|
import
|
||||||
|
unittest, strutils, ranges/ptr_arith,
|
||||||
|
../faststreams
|
||||||
|
|
||||||
|
suite "input stream":
|
||||||
|
test "string output":
|
||||||
|
var input = repeat("1234 5678 90AB CDEF\n", 1000)
|
||||||
|
var stream = memoryStream(input)
|
||||||
|
|
||||||
|
check:
|
||||||
|
(stream.readBytes(4) == "1234".toOpenArrayByte(0, 3))
|
||||||
|
|
||||||
|
|
@ -4,7 +4,7 @@ import
|
||||||
|
|
||||||
suite "output stream":
|
suite "output stream":
|
||||||
test "string output":
|
test "string output":
|
||||||
var s = init StringOutputStream
|
var s = init OutputStream
|
||||||
var altOutput = ""
|
var altOutput = ""
|
||||||
|
|
||||||
for i in 0 .. 1000:
|
for i in 0 .. 1000:
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue