Implemented asyncfile for Posix.

This commit is contained in:
Dominik Picheta 2014-09-05 21:14:18 +01:00
commit 52c16a1a79
3 changed files with 95 additions and 16 deletions

View file

@ -143,6 +143,7 @@ proc echoOriginalStackTrace[T](future: Future[T]) =
echo(future.errorStackTrace) echo(future.errorStackTrace)
else: else:
echo("Empty or nil stack trace.") echo("Empty or nil stack trace.")
echo("Continuing...")
proc read*[T](future: Future[T]): T = proc read*[T](future: Future[T]): T =
## Retrieves the value of ``future``. Future must be finished otherwise ## Retrieves the value of ``future``. Future must be finished otherwise
@ -723,7 +724,7 @@ else:
assert sock.SocketHandle in p.selector assert sock.SocketHandle in p.selector
discard p.selector.update(sock.SocketHandle, events) discard p.selector.update(sock.SocketHandle, events)
proc register(sock: TAsyncFD) = proc register*(sock: TAsyncFD) =
let p = getGlobalDispatcher() let p = getGlobalDispatcher()
var data = PData(sock: sock, readCBs: @[], writeCBs: @[]) var data = PData(sock: sock, readCBs: @[], writeCBs: @[])
p.selector.register(sock.SocketHandle, {}, data.PObject) p.selector.register(sock.SocketHandle, {}, data.PObject)
@ -743,14 +744,14 @@ else:
proc unregister*(fd: TAsyncFD) = proc unregister*(fd: TAsyncFD) =
getGlobalDispatcher().selector.unregister(fd.SocketHandle) getGlobalDispatcher().selector.unregister(fd.SocketHandle)
proc addRead(sock: TAsyncFD, cb: TCallback) = proc addRead*(sock: TAsyncFD, cb: TCallback) =
let p = getGlobalDispatcher() let p = getGlobalDispatcher()
if sock.SocketHandle notin p.selector: if sock.SocketHandle notin p.selector:
raise newException(EInvalidValue, "File descriptor not registered.") raise newException(EInvalidValue, "File descriptor not registered.")
p.selector[sock.SocketHandle].data.PData.readCBs.add(cb) p.selector[sock.SocketHandle].data.PData.readCBs.add(cb)
update(sock, p.selector[sock.SocketHandle].events + {EvRead}) update(sock, p.selector[sock.SocketHandle].events + {EvRead})
proc addWrite(sock: TAsyncFD, cb: TCallback) = proc addWrite*(sock: TAsyncFD, cb: TCallback) =
let p = getGlobalDispatcher() let p = getGlobalDispatcher()
if sock.SocketHandle notin p.selector: if sock.SocketHandle notin p.selector:
raise newException(EInvalidValue, "File descriptor not registered.") raise newException(EInvalidValue, "File descriptor not registered.")

View file

@ -28,7 +28,6 @@ when defined(windows):
import winlean import winlean
else: else:
import posix import posix
{.fatal: "Posix not yet supported".}
type type
AsyncFile = ref object AsyncFile = ref object
@ -59,12 +58,15 @@ else:
case mode case mode
of fmRead: of fmRead:
result = O_RDONLY result = O_RDONLY
of fmWrite, fmAppend: of fmWrite:
result = O_WRONLY or O_CREAT result = O_WRONLY or O_CREAT
of fmAppend:
result = O_WRONLY or O_CREAT or O_APPEND
of fmReadWrite: of fmReadWrite:
result = O_RDWR or O_CREAT result = O_RDWR or O_CREAT
of fmReadWriteExisting: of fmReadWriteExisting:
result = O_RDWR result = O_RDWR
result = result or O_NONBLOCK
proc getFileSize*(f: AsyncFile): int64 = proc getFileSize*(f: AsyncFile): int64 =
## Retrieves the specified file's size. ## Retrieves the specified file's size.
@ -102,6 +104,13 @@ proc openAsync*(filename: string, mode = fmRead): AsyncFile =
else: else:
let flags = getPosixFlags(mode) let flags = getPosixFlags(mode)
# RW (Owner), RW (Group), R (Other)
let perm = S_IRUSR or S_IWUSR or S_IRGRP or S_IWGRP or S_IROTH
result.fd = open(filename, flags, perm).TAsyncFD
if result.fd.cint == -1:
raiseOSError()
register(result.fd)
proc read*(f: AsyncFile, size: int): Future[string] = proc read*(f: AsyncFile, size: int): Future[string] =
## Read ``size`` bytes from the specified file asynchronously starting at ## Read ``size`` bytes from the specified file asynchronously starting at
@ -167,6 +176,29 @@ proc read*(f: AsyncFile, size: int): Future[string] =
copyMem(addr data[0], buffer, bytesRead) copyMem(addr data[0], buffer, bytesRead)
f.offset.inc bytesRead f.offset.inc bytesRead
retFuture.complete($data) retFuture.complete($data)
else:
var readBuffer = newString(size)
proc cb(fd: TAsyncFD): bool =
result = true
let res = read(fd.cint, addr readBuffer[0], size.cint)
if res < 0:
let lastError = osLastError()
if lastError.int32 != EAGAIN:
retFuture.fail(newException(EOS, osErrorMsg(lastError)))
else:
result = false # We still want this callback to be called.
elif res == 0:
# EOF
retFuture.complete("")
else:
readBuffer.setLen(res)
f.offset.inc(res)
retFuture.complete(readBuffer)
if not cb(f.fd):
addRead(f.fd, cb)
return retFuture return retFuture
proc getFilePos*(f: AsyncFile): int64 = proc getFilePos*(f: AsyncFile): int64 =
@ -179,6 +211,10 @@ proc setFilePos*(f: AsyncFile, pos: int64) =
## Sets the position of the file pointer that is used for read/write ## Sets the position of the file pointer that is used for read/write
## operations. The file's first byte has the index zero. ## operations. The file's first byte has the index zero.
f.offset = pos f.offset = pos
when not defined(windows):
let ret = lseek(f.fd.cint, pos, SEEK_SET)
if ret == -1:
raiseOSError()
proc readAll*(f: AsyncFile): Future[string] {.async.} = proc readAll*(f: AsyncFile): Future[string] {.async.} =
## Reads all data from the specified file. ## Reads all data from the specified file.
@ -240,6 +276,29 @@ proc write*(f: AsyncFile, data: string): Future[void] =
assert bytesWritten == data.len.int32 assert bytesWritten == data.len.int32
f.offset.inc(data.len) f.offset.inc(data.len)
retFuture.complete() retFuture.complete()
else:
var written = 0
proc cb(fd: TAsyncFD): bool =
result = true
let remainderSize = data.len-written
let res = write(fd.cint, addr copy[written], remainderSize.cint)
if res < 0:
let lastError = osLastError()
if lastError.int32 != EAGAIN:
retFuture.fail(newException(EOS, osErrorMsg(lastError)))
else:
result = false # We still want this callback to be called.
else:
written.inc res
f.offset.inc res
if res != remainderSize:
result = false # We still have data to write.
else:
retFuture.complete()
if not cb(f.fd):
addWrite(f.fd, cb)
return retFuture return retFuture
proc close*(f: AsyncFile) = proc close*(f: AsyncFile) =
@ -247,10 +306,7 @@ proc close*(f: AsyncFile) =
when defined(windows): when defined(windows):
if not closeHandle(f.fd.THandle).bool: if not closeHandle(f.fd.THandle).bool:
raiseOSError() raiseOSError()
else:
if close(f.fd.cint) == -1:
raiseOSError()

View file

@ -5,10 +5,32 @@ discard """
import asyncfile, asyncdispatch, os import asyncfile, asyncdispatch, os
proc main() {.async.} = proc main() {.async.} =
var file = openAsync(getTempDir() / "foobar.txt", fmReadWrite) let fn = getTempDir() / "foobar.txt"
await file.write("test") removeFile(fn)
file.setFilePos(0)
let data = await file.readAll() # Simple write/read test.
doAssert data == "test" block:
var file = openAsync(fn, fmReadWrite)
await file.write("test")
file.setFilePos(0)
await file.write("foo")
file.setFilePos(0)
let data = await file.readAll()
doAssert data == "foot"
file.close()
# Append test
block:
var file = openAsync(fn, fmAppend)
await file.write("\ntest2")
let errorTest = file.readAll()
await errorTest
doAssert errorTest.failed
file.close()
file = openAsync(fn, fmRead)
let data = await file.readAll()
doAssert data == "foot\ntest2"
file.close()
waitFor main() waitFor main()