Harmonize the APIs and the coding practices in input_stream and output_stream

This commit is contained in:
Zahary Karadjov 2020-04-10 15:30:58 +03:00
commit 3efad6f7f2
No known key found for this signature in database
GPG key ID: C8936F8A3073D609
4 changed files with 186 additions and 155 deletions

View file

@ -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())

View file

@ -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)