diff --git a/faststreams.nimble b/faststreams.nimble index 729526c..f470a6c 100644 --- a/faststreams.nimble +++ b/faststreams.nimble @@ -8,7 +8,8 @@ license = "Apache License 2.0" skipDirs = @["tests"] requires "nim >= 0.17.0", - "ranges" + "ranges", + "std_shims" task test, "run tests": exec "nim c -r tests/test_output_stream.nim" diff --git a/faststreams/input_stream.nim b/faststreams/input_stream.nim index a59fd62..f7753dd 100644 --- a/faststreams/input_stream.nim +++ b/faststreams/input_stream.nim @@ -12,14 +12,19 @@ type CloseStreamProc = proc () ByteStream* = object of RootObj - head: ptr byte + head*: ptr byte bufferSize: int bufferStart, bufferEnd: ptr byte + cursorsList: ptr StreamCursor reader: StreamReader asyncReader: AsyncStreamReader closeStream: CloseStreamProc bufferEndPos: int + StreamCursor* = object + head, bufferEnd: ptr byte + nextCursor: ptr StreamCursor + BufferedStream*[size: static int] = object of ByteStream buffer: array[size, byte] @@ -27,7 +32,12 @@ type UnicodeStream* = 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) result.head = cast[ptr byte](memFile.mem) when debugHelpers: @@ -46,12 +56,13 @@ proc init*(T: type ByteStream, result.reader = reader result.asyncReader = asyncReader -template memoryStream*(mem: openarray[byte]): ByteStream = - ByteStream.init(mem) +proc memoryStream*(mem: openarray[byte]): ByteStreamVar = # TODO: mark the result with `from mem` + new result + result[] = ByteStream.init(mem) -template memoryStream*(str: string): ByteStream = - bind init - init(ByteStream, str.toOpenArrayByte(0, str.len - 1)) +proc memoryStream*(str: string): ByteStreamVar = # TODO: mark the result with `from str` + new result + result[] = init(ByteStream, str.toOpenArrayByte(0, str.len - 1)) proc init*(T: type BufferedStream, reader = StreamReader(nil), @@ -68,6 +79,15 @@ proc syncRead(s: var ByteStream): bool = s.bufferEndPos += bytesRead 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 = if s.head != s.bufferEnd: return false @@ -99,6 +119,14 @@ proc read*(s: var ByteStream): byte = result = s.peek() 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] = if not s.eof: result = some s.read() diff --git a/faststreams/output_stream.nim b/faststreams/output_stream.nim index 69e7006..ff06991 100644 --- a/faststreams/output_stream.nim +++ b/faststreams/output_stream.nim @@ -1,81 +1,155 @@ import - ranges/ptr_arith + deques, ranges/ptr_arith, std_shims/strings type OutputStreamVTable = tuple prepareOutput: proc (s: ptr OutputStream, size: int) {.nimcall.} finish: proc (s: ptr OutputStream) {.nimcall.} - OutputStream* = object {.inheritable.} + OutputPage = object + buffer: string + startOffset: int16 + delayedWrites: int16 + + OutputStream* = object head, bufferEnd: ptr byte + pages: Deque[OutputPage] + firstPage: int + endPos: int + outputResource: pointer vtable: ptr OutputStreamVTable - MemoryOutputStream*[T] = object of OutputStream - output: T + DelayedWriteCursor* = object + head, bufferEnd: ptr byte + stream: OutputStreamVar + page: int - StringOutputStream* = MemoryOutputStream[string] - BytesOutputStream* = MemoryOutputStream[seq[byte]] - - FileOutputStream* = object of OutputStream - file: File + OutputStreamVar* = ref OutputStream 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) = - let s = cast[ptr T](s) - let currentSize = s.output.len - s.output.setLen(s.output.len + pageSize) - s.head = cast[ptr byte](addr s.output[currentSize]) - s.bufferEnd = cast[ptr byte](shift(addr s.output[0], s.output.len)) +proc addPage(s: OutputStreamVar) = + s.pages.addLast OutputPage(buffer: newString(pageSize), + delayedWrites: 0, + startOffset: 0) + s.head = cast[ptr byte](addr s.pages[s.pages.len - 1].buffer[0]) + s.bufferEnd = cast[ptr byte](shift(s.head, pageSize)) + s.endPos += pageSize -proc finishImpl[T](s: ptr OutputStream) = - let s = cast[ptr T](s) - s.output.setLen(s.output.len - distance(s.head, s.bufferEnd)) +proc init*(T: type OutputStream): ref OutputStream = + new result + result.vtable = nil + result.pages = initDeque[OutputPage]() + result.addPage() -proc init*(T: type MemoryOutputStream): T = - var vtable {.global.} = (prepareOutputImpl[T], finishImpl[T]) - result.vtable = addr vtable - result.vtable.prepareOutput(addr result, pageSize) +proc pos*(s: OutputStreamVar): int = + s.endPos - distance(s.head, s.bufferEnd) -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: - s.vtable.prepareOutput(addr s, pageSize) + s.tryFlushing() s.head[] = b s.head = shift(s.head, 1) -template append*(s: var OutputStream, c: char) = +template append*(s: OutputStreamVar, c: char) = s.append byte(c) -proc flush*(s: var OutputStream) = - 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]) = +proc append*(s: OutputStreamVar, bytes: openarray[byte]) = # TODO: this can use copyMem for b in bytes: s.append b -proc append*(s: var OutputStream, chars: openarray[char]) = +proc append*(s: OutputStreamVar, chars: openarray[char]) = # TODO: this can use copyMem for c in chars: 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) +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 -proc appendNumberImpl(s: var OutputStream, number: BiggestInt) = +proc appendNumberImpl(s: OutputStreamVar, number: BiggestInt) = # TODO: don't allocate s.append $number -proc appendNumberImpl(s: var OutputStream, number: BiggestUInt) = +proc appendNumberImpl(s: OutputStreamVar, number: BiggestUInt) = # TODO: don't allocate s.append $number @@ -85,7 +159,7 @@ template toBiggestRepr(i: SomeUnsignedInt): BiggestUInt = template toBiggestRepr(i: SomeSignedInt): BiggestInt = BiggestInt(i) -template appendNumber*(s: var OutputStream, i: SomeInteger) = +template appendNumber*(s: OutputStreamVar, i: SomeInteger) = # TODO: specify radix/base appendNumberImpl(s, toBiggestRepr(i)) diff --git a/tests/test_input_stream.nim b/tests/test_input_stream.nim new file mode 100644 index 0000000..922ab06 --- /dev/null +++ b/tests/test_input_stream.nim @@ -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)) + diff --git a/tests/test_output_stream.nim b/tests/test_output_stream.nim index 32e276f..0d5cc62 100644 --- a/tests/test_output_stream.nim +++ b/tests/test_output_stream.nim @@ -4,7 +4,7 @@ import suite "output stream": test "string output": - var s = init StringOutputStream + var s = init OutputStream var altOutput = "" for i in 0 .. 1000: