Fixes #268
This commit is contained in:
parent
40b611cc2f
commit
d6632ad973
3 changed files with 65 additions and 18 deletions
|
|
@ -126,6 +126,7 @@ type
|
||||||
handleTask*: proc (s: PAsyncSocket) {.closure.}
|
handleTask*: proc (s: PAsyncSocket) {.closure.}
|
||||||
|
|
||||||
lineBuffer: TaintedString ## Temporary storage for ``recvLine``
|
lineBuffer: TaintedString ## Temporary storage for ``recvLine``
|
||||||
|
sendBuffer: string ## Temporary storage for ``send``
|
||||||
sslNeedAccept: bool
|
sslNeedAccept: bool
|
||||||
proto: TProtocol
|
proto: TProtocol
|
||||||
deleg: PDelegate
|
deleg: PDelegate
|
||||||
|
|
@ -155,6 +156,7 @@ proc newAsyncSocket(): PAsyncSocket =
|
||||||
result.handleTask = (proc (s: PAsyncSocket) = nil)
|
result.handleTask = (proc (s: PAsyncSocket) = nil)
|
||||||
|
|
||||||
result.lineBuffer = "".TaintedString
|
result.lineBuffer = "".TaintedString
|
||||||
|
result.sendBuffer = ""
|
||||||
|
|
||||||
proc AsyncSocket*(domain: TDomain = AF_INET, typ: TType = SOCK_STREAM,
|
proc AsyncSocket*(domain: TDomain = AF_INET, typ: TType = SOCK_STREAM,
|
||||||
protocol: TProtocol = IPPROTO_TCP,
|
protocol: TProtocol = IPPROTO_TCP,
|
||||||
|
|
@ -225,10 +227,22 @@ proc asyncSockHandleWrite(h: PObject) =
|
||||||
else:
|
else:
|
||||||
PAsyncSocket(h).deleg.mode = fmReadWrite
|
PAsyncSocket(h).deleg.mode = fmReadWrite
|
||||||
else:
|
else:
|
||||||
if PAsyncSocket(h).handleWrite != nil:
|
if PAsyncSocket(h).sendBuffer != "":
|
||||||
PAsyncSocket(h).handleWrite(PAsyncSocket(h))
|
let sock = PAsyncSocket(h)
|
||||||
|
let bytesSent = sock.socket.sendAsync(sock.sendBuffer)
|
||||||
|
assert bytesSent > 0
|
||||||
|
if bytesSent != sock.sendBuffer.len:
|
||||||
|
sock.sendBuffer = sock.sendBuffer[bytesSent .. -1]
|
||||||
|
elif bytesSent == sock.sendBuffer.len:
|
||||||
|
sock.sendBuffer = ""
|
||||||
|
|
||||||
|
if PAsyncSocket(h).handleWrite != nil:
|
||||||
|
PAsyncSocket(h).handleWrite(PAsyncSocket(h))
|
||||||
else:
|
else:
|
||||||
PAsyncSocket(h).deleg.mode = fmRead
|
if PAsyncSocket(h).handleWrite != nil:
|
||||||
|
PAsyncSocket(h).handleWrite(PAsyncSocket(h))
|
||||||
|
else:
|
||||||
|
PAsyncSocket(h).deleg.mode = fmRead
|
||||||
|
|
||||||
when defined(ssl):
|
when defined(ssl):
|
||||||
proc asyncSockDoHandshake(h: PObject) =
|
proc asyncSockDoHandshake(h: PObject) =
|
||||||
|
|
@ -340,7 +354,8 @@ proc acceptAddr*(server: PAsyncSocket, client: var PAsyncSocket,
|
||||||
# deleg.open is set in ``toDelegate``.
|
# deleg.open is set in ``toDelegate``.
|
||||||
|
|
||||||
client.socket = c
|
client.socket = c
|
||||||
client.lineBuffer = ""
|
client.lineBuffer = "".TaintedString
|
||||||
|
client.sendBuffer = ""
|
||||||
client.info = SockConnected
|
client.info = SockConnected
|
||||||
|
|
||||||
proc accept*(server: PAsyncSocket, client: var PAsyncSocket) =
|
proc accept*(server: PAsyncSocket, client: var PAsyncSocket) =
|
||||||
|
|
@ -445,6 +460,26 @@ proc recvLine*(s: PAsyncSocket, line: var TaintedString): bool =
|
||||||
of RecvFail:
|
of RecvFail:
|
||||||
result = false
|
result = false
|
||||||
|
|
||||||
|
proc send*(sock: PAsyncSocket, data: string) =
|
||||||
|
## Sends ``data`` to socket ``sock``. This is basically a nicer implementation
|
||||||
|
## of ``sockets.sendAsync``.
|
||||||
|
##
|
||||||
|
## If ``data`` cannot be sent immediately it will be buffered and sent
|
||||||
|
## when ``sock`` becomes writeable (during the ``handleWrite`` event).
|
||||||
|
## It's possible that only a part of ``data`` will be sent immediately, while
|
||||||
|
## the rest of it will be buffered and sent later.
|
||||||
|
if sock.sendBuffer.len != 0:
|
||||||
|
sock.sendBuffer.add(data)
|
||||||
|
return
|
||||||
|
let bytesSent = sock.socket.sendAsync(data)
|
||||||
|
assert bytesSent >= 0
|
||||||
|
if bytesSent == 0:
|
||||||
|
sock.sendBuffer.add(data)
|
||||||
|
sock.deleg.mode = fmReadWrite
|
||||||
|
elif bytesSent != data.len:
|
||||||
|
sock.sendBuffer.add(data[bytesSent .. -1])
|
||||||
|
sock.deleg.mode = fmReadWrite
|
||||||
|
|
||||||
proc timeValFromMilliseconds(timeout = 500): TTimeVal =
|
proc timeValFromMilliseconds(timeout = 500): TTimeVal =
|
||||||
if timeout != -1:
|
if timeout != -1:
|
||||||
var seconds = timeout div 1000
|
var seconds = timeout div 1000
|
||||||
|
|
|
||||||
|
|
@ -454,11 +454,13 @@ proc doUpload(ftp: PFTPClient, async = false): bool =
|
||||||
if ftp.dsockConnected:
|
if ftp.dsockConnected:
|
||||||
if ftp.job.toStore.len() > 0:
|
if ftp.job.toStore.len() > 0:
|
||||||
assert(async)
|
assert(async)
|
||||||
if ftp.asyncDSock.sendAsync(ftp.job.toStore):
|
let bytesSent = ftp.asyncDSock.sendAsync(ftp.job.toStore)
|
||||||
|
if bytesSent == ftp.job.toStore.len:
|
||||||
ftp.job.toStore = ""
|
ftp.job.toStore = ""
|
||||||
ftp.job.progress.inc(ftp.job.toStore.len)
|
elif bytesSent != ftp.job.toStore.len and bytesSent != 0:
|
||||||
ftp.job.oneSecond.inc(ftp.job.toStore.len)
|
ftp.job.toStore = ftp.job.toStore[bytesSent .. -1]
|
||||||
|
ftp.job.progress.inc(bytesSent)
|
||||||
|
ftp.job.oneSecond.inc(bytesSent)
|
||||||
else:
|
else:
|
||||||
var s = newStringOfCap(4000)
|
var s = newStringOfCap(4000)
|
||||||
var len = ftp.job.file.readBuffer(addr(s[0]), 4000)
|
var len = ftp.job.file.readBuffer(addr(s[0]), 4000)
|
||||||
|
|
@ -476,8 +478,12 @@ proc doUpload(ftp: PFTPClient, async = false): bool =
|
||||||
if not async:
|
if not async:
|
||||||
getDSock(ftp).send(s)
|
getDSock(ftp).send(s)
|
||||||
else:
|
else:
|
||||||
if not ftp.asyncDSock.sendAsync(s):
|
let bytesSent = ftp.asyncDSock.sendAsync(s)
|
||||||
ftp.job.toStore = s
|
if bytesSent == 0:
|
||||||
|
ftp.job.toStore.add(s)
|
||||||
|
elif bytesSent != s.len:
|
||||||
|
ftp.job.toStore.add(s[bytesSent .. -1])
|
||||||
|
len = bytesSent
|
||||||
|
|
||||||
ftp.job.progress.inc(len)
|
ftp.job.progress.inc(len)
|
||||||
ftp.job.oneSecond.inc(len)
|
ftp.job.oneSecond.inc(len)
|
||||||
|
|
|
||||||
|
|
@ -1317,13 +1317,18 @@ proc send*(socket: TSocket, data: string) {.tags: [FWriteIO].} =
|
||||||
|
|
||||||
OSError()
|
OSError()
|
||||||
|
|
||||||
proc sendAsync*(socket: TSocket, data: string): bool {.tags: [FWriteIO].} =
|
proc sendAsync*(socket: TSocket, data: string): int {.tags: [FWriteIO].} =
|
||||||
## sends data to a non-blocking socket. Returns whether ``data`` was sent.
|
## sends data to a non-blocking socket.
|
||||||
result = true
|
## Returns ``0`` if no data could be sent, if data has been sent
|
||||||
var bytesSent = send(socket, cstring(data), data.len)
|
## returns the amount of bytes of ``data`` that was successfully sent. This
|
||||||
|
## number may not always be the length of ``data`` but typically is.
|
||||||
|
##
|
||||||
|
## An EOS (or ESSL if socket is an SSL socket) exception is raised if an error
|
||||||
|
## occurs.
|
||||||
|
result = send(socket, cstring(data), data.len)
|
||||||
when defined(ssl):
|
when defined(ssl):
|
||||||
if socket.isSSL:
|
if socket.isSSL:
|
||||||
if bytesSent <= 0:
|
if result <= 0:
|
||||||
let ret = SSLGetError(socket.sslHandle, bytesSent.cint)
|
let ret = SSLGetError(socket.sslHandle, bytesSent.cint)
|
||||||
case ret
|
case ret
|
||||||
of SSL_ERROR_ZERO_RETURN:
|
of SSL_ERROR_ZERO_RETURN:
|
||||||
|
|
@ -1339,18 +1344,19 @@ proc sendAsync*(socket: TSocket, data: string): bool {.tags: [FWriteIO].} =
|
||||||
else: SSLError("Unknown Error")
|
else: SSLError("Unknown Error")
|
||||||
else:
|
else:
|
||||||
return
|
return
|
||||||
if bytesSent == -1:
|
if result == -1:
|
||||||
when defined(windows):
|
when defined(windows):
|
||||||
var err = WSAGetLastError()
|
var err = WSAGetLastError()
|
||||||
# TODO: Test on windows.
|
# TODO: Test on windows.
|
||||||
if err == WSAEINPROGRESS:
|
if err == WSAEINPROGRESS:
|
||||||
return false
|
return 0
|
||||||
else: OSError()
|
else: OSError()
|
||||||
else:
|
else:
|
||||||
if errno == EAGAIN or errno == EWOULDBLOCK:
|
if errno == EAGAIN or errno == EWOULDBLOCK:
|
||||||
return false
|
return 0
|
||||||
else: OSError()
|
else: OSError()
|
||||||
|
|
||||||
|
|
||||||
proc trySend*(socket: TSocket, data: string): bool {.tags: [FWriteIO].} =
|
proc trySend*(socket: TSocket, data: string): bool {.tags: [FWriteIO].} =
|
||||||
## safe alternative to ``send``. Does not raise an EOS when an error occurs,
|
## safe alternative to ``send``. Does not raise an EOS when an error occurs,
|
||||||
## and instead returns ``false`` on failure.
|
## and instead returns ``false`` on failure.
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue