added Timer and enum return to TCPConnection.waitForDataAsync

This commit is contained in:
Francesco Mecca 2018-02-23 00:27:04 +00:00
parent 0d6ba62f51
commit 20e32cf327

View file

@ -586,29 +586,39 @@ mixin(tracer);
return m_context.readBuffer.length > 0; return m_context.readBuffer.length > 0;
} }
void waitForDataAsync(void delegate(bool) @safe read_callback, Duration timeout = Duration.max) enum WaitForDataAsyncStatus {
noMoreData,
dataAvailable,
waiting
}
WaitForDataAsyncStatus waitForDataAsync(void delegate(bool) @safe read_callback, Duration timeout = Duration.max)
{ {
mixin(tracer); mixin(tracer);
import vibe.core.core : runTask; import vibe.core.core : runTask, setTimer, createTimer;
if (!m_context) { if (!m_context) {
runTask(read_callback, false); runTask(read_callback, false);
return; return WaitForDataAsyncStatus.noMoreData;
} }
if (m_context.readBuffer.length > 0) { if (m_context.readBuffer.length > 0) {
runTask(read_callback, true); runTask(read_callback, true);
return; return WaitForDataAsyncStatus.dataAvailable;
} }
if (timeout <= 0.seconds) {
auto rs = waitForData(0.seconds);
auto mode = rs ? WaitForDataAsyncStatus.dataAvailable : WaitForDataAsyncStatus.noMoreData;
return mode;
}
auto mode = timeout <= 0.seconds ? IOMode.immediate : IOMode.once; auto tm = setTimer(timeout, { eventDriver.sockets.cancelRead(m_socket);
bool cancelled; runTask(read_callback, false); });
IOStatus status; eventDriver.sockets.read(m_socket, m_context.readBuffer.peekDst(), IOMode.once,
size_t nbytes; (sock, st, nb) { tm.stop(); if(st != IOStatus.ok) runTask(read_callback, false);
eventDriver.sockets.waitForData(m_socket,
(sock, st, nb) { if(st != IOStatus.ok) runTask(read_callback, false);
else runTask(read_callback, true); else runTask(read_callback, true);
}); });
return WaitForDataAsyncStatus.waiting;
} }
const(ubyte)[] peek() { return m_context ? m_context.readBuffer.peek() : null; } const(ubyte)[] peek() { return m_context ? m_context.readBuffer.peek() : null; }