From 3efad6f7f2f40d3554948bd206c3766a654bb029 Mon Sep 17 00:00:00 2001 From: Zahary Karadjov Date: Fri, 10 Apr 2020 15:30:58 +0300 Subject: [PATCH] Harmonize the APIs and the coding practices in input_stream and output_stream --- faststreams/input_stream.nim | 231 +++++++++++++++++++--------------- faststreams/output_stream.nim | 102 +++++++-------- tests/test_input_stream.nim | 2 +- tests/test_output_stream.nim | 6 +- 4 files changed, 186 insertions(+), 155 deletions(-) diff --git a/faststreams/input_stream.nim b/faststreams/input_stream.nim index e99da95..2a870c6 100644 --- a/faststreams/input_stream.nim +++ b/faststreams/input_stream.nim @@ -1,192 +1,223 @@ import - memfiles, options, stew/[ptrops, ranges/ptr_arith] - -const - # pageSize = 4096 - debugHelpers = false + memfiles, options, + stew/[ptrops, ranges/ptr_arith] type - StreamReader = proc (bufferStart: ptr byte, bufferSize: int): int {.gcsafe.} - # TODO: use openarray once it's supported - AsyncStreamReader = StreamReader # proc (bufferStart: ptr byte, bufferSize: int): Future[int] - CloseStreamProc = proc () - - ByteStream* = object of RootObj + # We inherit from RootObj because in layered streams + # the stream itself is often used as an `outputDevice` + InputStreamObj = object of RootObj head*: ptr byte bufferSize: int bufferStart, bufferEnd: ptr byte - cursorsList: ptr StreamCursor - reader: StreamReader - asyncReader: AsyncStreamReader - closeStream: CloseStreamProc bufferEndPos: int + inputDevice*: RootRef + vtable*: ptr InputStreamVTable - StreamCursor* = object - head, bufferEnd: ptr byte - nextCursor: ptr StreamCursor + InputStream* = ref InputStreamObj - BufferedStream*[size: static int] = object of ByteStream - buffer: array[size, byte] + AsciiInputStream* = distinct InputStream + Utf8InputStream* = distinct InputStream - AsciiStream* = distinct ByteStream - UnicodeStream* = distinct ByteStream - ObjectStream*[T] = distinct ByteStream + ReadSyncProc* = proc (s: InputStream, buffer: ptr byte, bufSize: int): int + {.nimcall, gcsafe, raises: [IOError, Defect].} - # TODO: ByteStreamVar should become a `var` short-lived pointer once this is supported - ByteStreamVar* = ref ByteStream - AsciiStreamVar* = ref AsciiStream + ReadAsyncCallback* = proc (s: InputStream, bytesRead: int) + {.nimcall, gcsafe, raises: [Defect].} - InputStreamVar* = ref ByteStream + ReadAsyncProc* = proc (s: InputStream, + buffer: ptr byte, bufSize: int, + cb: ReadAsyncCallback) + {.nimcall, gcsafe, raises: [IOError, Defect].} + + CloseStreamProc* = proc (s: InputStream) + {.nimcall, gcsafe, raises: [IOError, Defect].} + + InputStreamVTable* = object + readSync*: ReadSyncProc + readAsync*: ReadAsyncProc + closeStream*: CloseStreamProc + hasKnownLen*: bool + + FileInput = ref object of RootObj + file: MemFile + +const + debugHelpers = false + nimAllocatorMetadataSize* = 0 + # TODO: Get this from Nim's allocator. + # The goal is to make perfect page-aligned allocations + defaultPageSize = 4096 - nimAllocatorMetadataSize + +let FileStreamVTable = InputStreamVTable( + readSync: nil, + readAsync: nil, + closeStream: proc (s: InputStream) + {.nimcall, gcsafe, raises: [IOError, Defect].} = + try: + close FileInput(s.inputDevice).file + except OSError as err: + raise newException(IOError, "Failed to close file", err) + , + hasKnownLen: true +) + +template vtableAddr*(vtable: InputStreamVTable): ptr InputStreamVTable = + ## This is a simple work-around for the somewhat broken side + ## effects analysis of Nim - reading from global let variables + ## is considered a side-effect. + {.noSideEffect.}: + unsafeAddr vtable + +proc fileInput*(filename: string): InputStream = + let + inputDevice = FileInput(file: memfiles.open(filename)) + head = cast[ptr byte](inputDevice.file.mem) + fileSize = inputDevice.file.size + + result = InputStream( + head: head, + bufferEnd: offset(head, fileSize), + bufferEndPos: fileSize, + inputDevice: inputDevice, + vtable: vtableAddr FileStreamVTable) -proc openFile*(filename: string): ByteStreamVar = - new result - var memFile = memfiles.open(filename) - result.head = cast[ptr byte](memFile.mem) when debugHelpers: result.bufferStart = result.head - result.bufferEnd = offset(result.head, memFile.size) - result.bufferEndPos = memFile.size - result.closeStream = proc = close memFile -proc init*(T: type ByteStream, - mem: openarray[byte], - reader = StreamReader(nil), - asyncReader = AsyncStreamReader(nil)): ByteStream = - # TODO: the result should use `from mem` once it's supported - result.head = unsafeAddr mem[0] - result.bufferEnd = offset(result.head, mem.len) - result.bufferEndPos = mem.len - result.reader = reader - result.asyncReader = asyncReader +proc implementInputStream*(vtable: ptr InputStreamVTable, + inputDevice: RootRef, + pageSize = defaultPageSize): InputStream = + # TODO: We need to allocate memory here and start reading + InputStream(inputDevice: inputDevice, + vtable: vtable) -proc memoryStream*(mem: openarray[byte]): ByteStreamVar = # TODO: mark the result with `from mem` - new result - result[] = ByteStream.init(mem) +proc memoryInput*(mem: openarray[byte]): InputStream = + let head = unsafeAddr mem[0] + InputStream( + head: head, + bufferEnd: offset(head, mem.len), + bufferEndPos: mem.len) -proc memoryStream*(str: string): ByteStreamVar = # TODO: mark the result with `from str` - new result - result[] = init(ByteStream, str.toOpenArrayByte(0, str.len - 1)) +proc memoryInput*(str: string): InputStream = + memoryInput str.toOpenArrayByte(0, str.len - 1) -proc init*(T: type BufferedStream, - reader = StreamReader(nil), - asyncReader = AsyncStreamReader(nil)): BufferedStream = - result.init result.buffer, reader, asyncReader - -proc endPos*(s: InputStreamVar): int = - # TODO This needs to use a VTable for `system.File` based streams +proc endPos*(s: InputStream): int = + doAssert s.vtable == nil or s.vtable.hasKnownLen return s.bufferEndPos -proc syncRead(s: var ByteStream): bool = - let bytesRead = s.reader(s.bufferStart, s.bufferSize) +proc syncRead(s: InputStream): bool = + let bytesRead = s.vtable.readSync(s, s.bufferStart, s.bufferSize) if bytesRead == 0: - s.reader = nil + # TODO close the input device + s.inputDevice = nil return true else: s.bufferEnd = offset(s.bufferStart, bytesRead) s.bufferEndPos += bytesRead return false -proc ensureBytes*(s: var ByteStream, n: int): bool = +proc ensureBytes*(s: InputStream, n: int): bool = if distance(s.head, s.bufferEnd) >= n: return true - if s.reader == nil: + if s.inputDevice == nil: return false + # TODO doAssert false, "Multi-buffer reading will be implemented later" -proc eof*(s: var ByteStream): bool = +proc eof*(s: InputStream): bool = if s.head != s.bufferEnd: return false - if s.reader == nil: + if s.inputDevice == nil: return true return s.syncRead() -proc eob*(s: ByteStream): bool {.inline.} = - s.head != s.bufferEnd +#proc eob*(s: InputStream): bool {.inline.} = +# s.head != s.bufferEnd -proc peek*(s: ByteStream): byte {.inline.} = +proc peek*(s: InputStream): byte {.inline.} = doAssert s.head != s.bufferEnd return s.head[] when debugHelpers: - proc showPosition*(s: ByteStream) = + proc showPosition*(s: InputStream) = echo "head at ", distance(s.bufferStart, s.head), "/", distance(s.bufferStart, s.bufferEnd) -proc advance*(s: var ByteStream) = +proc advance*(s: InputStream) = if s.head != s.bufferEnd: s.head = offset(s.head, 1) - elif s.reader != nil: + elif s.inputDevice != nil: discard s.syncRead() -proc read*(s: var ByteStream): byte = +proc read*(s: InputStream): byte = result = s.peek() advance s -proc checkReadAhead(s: ByteStreamVar, n: int): ptr byte = +proc checkReadAhead(s: InputStream, n: int): ptr byte = result = s.head doAssert distance(s.head, s.bufferEnd) >= n s.head = offset(s.head, n) -template readBytes*(s: ByteStreamVar, n: int): auto = +template readBytes*(s: InputStream, n: int): auto = makeOpenArray(checkReadAhead(s, n), n) -proc next*(s: var ByteStream): Option[byte] = +proc next*(s: InputStream): Option[byte] = if not s.eof: result = some s.read() -proc bufferPos(s: ByteStream, pos: int): ptr byte = +proc bufferPos(s: InputStream, pos: int): ptr byte = let offsetFromEnd = pos - s.bufferEndPos doAssert offsetFromEnd < 0 result = offset(s.bufferEnd, offsetFromEnd) doAssert result >= s.bufferStart -proc pos*(s: ByteStream): int {.inline.} = +proc pos*(s: InputStream): int {.inline.} = s.bufferEndPos - distance(s.head, s.bufferEnd) -proc firstAccessiblePos*(s: ByteStream): int {.inline.} = +proc firstAccessiblePos*(s: InputStream): int {.inline.} = s.bufferEndPos - distance(s.bufferStart, s.bufferEnd) -proc `[]`*(s: ByteStream, pos: int): byte {.inline.} = +proc `[]`*(s: InputStream, pos: int): byte {.inline.} = s.bufferPos(pos)[] -proc rewind*(s: var ByteStream, delta: int) = +proc rewind*(s: InputStream, delta: int) = s.head = offset(s.head, -delta) doAssert s.head >= s.bufferStart -proc rewindTo*(s: var ByteStream, pos: int) {.inline.} = +proc rewindTo*(s: InputStream, pos: int) {.inline.} = s.head = s.bufferPos(pos) # TODO: use a destructor once we migrate to Nim 0.20 -proc close*(s: var ByteStream) = - if s.closeStream != nil: - s.closeStream() +# TODO: It's not appropriate for this to raise +proc close*(s: InputStream) {.raises: [IOError, Defect].} = + if s.vtable != nil: + s.vtable.closeStream(s) -# TODO: Use `distrinct with` once it's supported -template pos*(s: AsciiStream|UnicodeStream|ObjectStream): int = - ByteStream(s).pos +template pos*(s: AsciiInputStream|Utf8InputStream): int = + InputStream(s).pos -template eof*(s: var (AsciiStream|UnicodeStream|ObjectStream)): bool = - ByteStream(s).eof +template eof*(s: AsciiInputStream|Utf8InputStream): bool = + InputStream(s).eof -template eob*(s: var (AsciiStream|UnicodeStream|ObjectStream)): bool = - ByteStream(s).eob +#template eob*(s: AsciiInputStream|Utf8InputStream): bool = +# InputStream(s).eob -template close*(s: AsciiStream|UnicodeStream|ObjectStream) = - close ByteStream(s) +template close*(s: AsciiInputStream|Utf8InputStream) = + close InputStream(s) -template advance*(s: var AsciiStream) = - advance ByteStream(s) +template advance*(s: AsciiInputStream) = + advance InputStream(s) -template peek*(s: AsciiStream): char = - char ByteStream(s).peek() +template peek*(s: AsciiInputStream): char = + char InputStream(s).peek() -template read*(s: var AsciiStream): char = - char ByteStream(s).read() +template read*(s: var AsciiInputStream): char = + char InputStream(s).read() -template next*(s: var AsciiStream): Option[char] = - cast[Option[char]](ByteStream(s).next()) +template next*(s: var AsciiInputStream): Option[char] = + cast[Option[char]](InputStream(s).next()) diff --git a/faststreams/output_stream.nim b/faststreams/output_stream.nim index 4307094..3b75713 100644 --- a/faststreams/output_stream.nim +++ b/faststreams/output_stream.nim @@ -13,15 +13,14 @@ type cursor*: WriteCursor pages: Deque[OutputPage] endPos: int - vtable*: ptr OutputStreamVTable - outputDevice*: RootRef extCursorsCount: int pageSize: int maxWriteSize*: int minWriteSize*: int + outputDevice*: RootRef + vtable*: ptr OutputStreamVTable OutputStream* = ref OutputStreamObj - OutputStreamVar = OutputStream WritePageProc* = proc (s: OutputStream, page: openarray[byte]) {.nimcall, gcsafe, raises: [IOError, Defect].} @@ -35,7 +34,7 @@ type WriteCursor* = object head, bufferEnd: ptr byte - stream: OutputStreamVar + stream: OutputStream VarSizeWriteCursor* = distinct WriteCursor @@ -43,15 +42,16 @@ type file: File 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 + nimAllocatorMetadataSize* = 0 + # TODO: Get this from Nim's allocator. + # The goal is to make perfect page-aligned allocations + defaultPageSize = 4096 - nimAllocatorMetadataSize - 1 # 1 byte for the null terminator proc createWriteCursor*[R, T](x: var array[R, T]): WriteCursor = let startAddr = cast[ptr byte](addr x[0]) WriteCursor(head: startAddr, bufferEnd: offset(startAddr, sizeof x)) -template canExtendOutput(s: OutputStreamVar): bool = +template canExtendOutput(s: OutputStream): bool = # Streams writing to pre-allocated existing buffers cannot be grown s != nil and s.pageSize > 0 @@ -62,23 +62,23 @@ template isExternalCursor(c: var WriteCursor): bool = func runway*(c: var WriteCursor): int {.inline.} = distance(c.head, c.bufferEnd) -proc flipPage(s: OutputStreamVar) = +proc flipPage(s: OutputStream) = 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](offset(s.cursor.head, s.pageSize)) s.endPos += s.pageSize -proc addPage(s: OutputStreamVar) = +proc addPage(s: OutputStream) = s.pages.addLast OutputPage(buffer: newString(s.pageSize), startOffset: 0) s.flipPage -proc implementOutputStream*(pageSize: int, +proc implementOutputStream*(outputDevice: RootRef, + vtable: ptr OutputStreamVTable, + pageSize = defaultPageSize, maxWriteSize = high(int), - minWriteSize = 1, - vtable: ptr OutputStreamVTable = nil, - outputDevice: RootRef = nil): OutputStream = + minWriteSize = 1): OutputStream = ## This proc is intented for use by module that implement ## their own flavours of OutputStream by providing a custom ## VTable. Such modules should export easier to use high-level @@ -94,18 +94,25 @@ proc implementOutputStream*(pageSize: int, result.addPage result.cursor.stream = result -proc init*(T: type OutputStream, - pageSize = defaultPageSize): OutputStream = - implementOutputStream pageSize +proc memoryOutput*(pageSize = defaultPageSize): OutputStream = + implementOutputStream(outputDevice = nil, vtable = nil, pageSize = pageSize) + +proc memoryOutput*(buffer: pointer, len: int): OutputStream = + result = OutputStream() + let buffer = cast[ptr byte](buffer) + result.cursor.head = buffer + result.cursor.bufferEnd = offset(buffer, len) + result.cursor.stream = result + result.endPos = len let FileStreamVTable = OutputStreamVTable( - writePage: proc (s: OutputStreamVar, data: openarray[byte]) {.nimcall, gcsafe.} = + writePage: proc (s: OutputStream, 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.} = + flush: proc (s: OutputStream) {.nimcall, gcsafe.} = var output = FileOutput(s.outputDevice) flushFile output.file ) @@ -117,34 +124,27 @@ template vtableAddr*(vtable: OutputStreamVTable): ptr OutputStreamVTable = {.noSideEffect.}: unsafeAddr vtable -proc init*(T: type OutputStream, - filename: string, - pageSize = defaultPageSize): OutputStream = - implementOutputStream pageSize, - outputDevice = FileOutput(file: open(filename, fmWrite)), - vtable = vtableAddr FileStreamVTable +proc fileOutput*(filename: string, + fileMode: FileMode = fmWrite, + pageSize = defaultPageSize): OutputStream {. + raises: [IOError, Defect] +.} = + implementOutputStream FileOutput(file: open(filename, fileMode)), + vtable = vtableAddr FileStreamVTable, + pageSize = pageSize -proc init*(T: type OutputStream, - buffer: pointer, len: int): OutputStream = - result = OutputStream() - let buffer = cast[ptr byte](buffer) - result.cursor.head = buffer - result.cursor.bufferEnd = offset(buffer, len) - result.cursor.stream = result - result.endPos = len - -proc pos*(s: OutputStreamVar): int = +proc pos*(s: OutputStream): int = s.endPos - s.cursor.runway -proc safeWritePage(s: OutputStreamVar, data: openarray[byte]) {.inline.} = +proc safeWritePage(s: OutputStream, data: openarray[byte]) {.inline.} = if data.len > 0: s.vtable.writePage(s, data) -proc writePages(s: OutputStreamVar, skipLast = 0) = +proc writePages(s: OutputStream, skipLast = 0) = assert s.vtable != nil for i in 0 ..< s.pages.len - skipLast: s.safeWritePage s.pages[i].buffer.toOpenArrayByte(0, s.pages[i].buffer.len - 1) -proc writePartialPage(s: OutputStreamVar, page: var OutputPage) = +proc writePartialPage(s: OutputStream, page: var OutputPage) = assert s.vtable != nil let unwrittenBytes = s.cursor.runway @@ -157,7 +157,7 @@ proc writePartialPage(s: OutputStreamVar, page: var OutputPage) = page.startOffset = 0 s.flipPage -proc flush*(s: OutputStreamVar) = +proc flush*(s: OutputStream) = doAssert s.extCursorsCount == 0 if s.vtable != nil: # We write all pages except the last one @@ -169,13 +169,13 @@ proc flush*(s: OutputStreamVar) = # Finally, we flush s.vtable.flush(s) -proc writePendingPagesAndLeaveOne(s: OutputStreamVar) {.inline.} = +proc writePendingPagesAndLeaveOne(s: OutputStream) {.inline.} = s.writePages s.pages.shrink(fromFirst = s.pages.len - 1) s.pages[0].startOffset = 0 s.flipPage -proc tryFlushing(s: OutputStreamVar) {.inline.} = +proc tryFlushing(s: OutputStream) {.inline.} = # Pre-conditions: # * The cursor has reached the current buffer end # @@ -232,7 +232,7 @@ template append*(c: var WriteCursor, x: char) = bind append c.append byte(x) -proc writeDataAsPages(s: OutputStreamVar, data: ptr byte, dataLen: int) = +proc writeDataAsPages(s: OutputStream, data: ptr byte, dataLen: int) = var data = data dataLen = dataLen @@ -387,15 +387,15 @@ template append*(c: var WriteCursor, str: string) = bind append c.append str.toOpenArrayByte(0, str.len - 1) -template append*(s: OutputStreamVar, value: auto) = +template append*(s: OutputStream, value: auto) = bind append s.cursor.append value -template appendMemCopy*(s: OutputStreamVar, value: auto) = +template appendMemCopy*(s: OutputStream, value: auto) = bind append s.cursor.append value -proc getOutput*(s: OutputStreamVar, T: type string): string = +proc getOutput*(s: OutputStream, T: type string): string = doAssert s.vtable == nil and s.extCursorsCount == 0 and s.pageSize > 0 s.pages[s.pages.len - 1].buffer.setLen(s.pageSize - s.cursor.runway) @@ -408,18 +408,18 @@ proc getOutput*(s: OutputStreamVar, T: type string): string = result.add page.buffer.toOpenArray(page.startOffset.int, page.buffer.len - 1) -template getOutput*(s: OutputStreamVar, T: type seq[byte]): seq[byte] = +template getOutput*(s: OutputStream, T: type seq[byte]): seq[byte] = cast[seq[byte]](s.getOutput(string)) -template getOutput*(s: OutputStreamVar): seq[byte] = +template getOutput*(s: OutputStream): seq[byte] = cast[seq[byte]](s.getOutput(string)) -proc finishPageEarly(s: OutputStreamVar, unwrittenBytes: int) {.inline.} = +proc finishPageEarly(s: OutputStream, unwrittenBytes: int) {.inline.} = s.pages[s.pages.len - 1].buffer.setLen(s.pageSize - unwrittenBytes) s.endPos -= unwrittenBytes s.tryFlushing() -proc createCursor(s: OutputStreamVar, size: int): WriteCursor = +proc createCursor(s: OutputStream, size: int): WriteCursor = inc s.extCursorsCount result = WriteCursor(head: s.cursor.head, @@ -428,7 +428,7 @@ proc createCursor(s: OutputStreamVar, size: int): WriteCursor = s.cursor.head = result.bufferEnd -proc delayFixedSizeWrite*(s: OutputStreamVar, size: Natural): WriteCursor = +proc delayFixedSizeWrite*(s: OutputStream, size: Natural): WriteCursor = let remainingBytesInPage = s.cursor.runway if size <= remainingBytesInPage: result = s.createCursor(size) @@ -449,7 +449,7 @@ proc delayFixedSizeWrite*(s: OutputStreamVar, size: Natural): WriteCursor = s.cursor.bufferEnd = pageEnd s.endPos += (s.pageSize - size) -proc delayVarSizeWrite*(s: OutputStreamVar, maxSize: Natural): VarSizeWriteCursor = +proc delayVarSizeWrite*(s: OutputStream, maxSize: Natural): VarSizeWriteCursor = doAssert maxSize < s.pageSize s.finishPageEarly s.cursor.runway VarSizeWriteCursor s.createCursor(maxSize) diff --git a/tests/test_input_stream.nim b/tests/test_input_stream.nim index cc9c8d5..6ec37d7 100644 --- a/tests/test_input_stream.nim +++ b/tests/test_input_stream.nim @@ -5,7 +5,7 @@ import suite "input stream": test "string input": var input = repeat("1234 5678 90AB CDEF\n", 1000) - var stream = memoryStream(input) + var stream = memoryInput(input) check: (stream.readBytes(4) == "1234".toOpenArrayByte(0, 3)) diff --git a/tests/test_output_stream.nim b/tests/test_output_stream.nim index 746a989..5415175 100644 --- a/tests/test_output_stream.nim +++ b/tests/test_output_stream.nim @@ -21,14 +21,14 @@ proc randomBytes(n: int): seq[byte] = suite "output stream": setup: - var memStream = OutputStream.init + var memStream = memoryOutput() var altOutput: seq[byte] = @[] var tempFilePath = getTempDir() / "faststreams_testfile" - var fileStream = OutputStream.init tempFilePath + var fileStream = fileOutput(tempFilePath) const bufferSize = 1000000 var buffer = alloc(bufferSize) - var existingBufferStream = OutputStream.init(buffer, bufferSize) + var existingBufferStream = memoryOutput(buffer, bufferSize) teardown: removeFile tempFilePath