refactor: handle read and write on tcpsocket at the same time
port 94f8336af5
This commit is contained in:
parent
eac59768ea
commit
5430625a7e
1 changed files with 13 additions and 10 deletions
|
|
@ -494,10 +494,11 @@ ISocketMultiplexerJob *TCPSocket::serviceConnected(ISocketMultiplexerJob *job, b
|
||||||
return newJob();
|
return newJob();
|
||||||
}
|
}
|
||||||
|
|
||||||
EJobResult result = kRetry;
|
EJobResult readResult = kRetry;
|
||||||
|
EJobResult writeResult = kRetry;
|
||||||
if (write) {
|
if (write) {
|
||||||
try {
|
try {
|
||||||
result = doWrite();
|
writeResult = doWrite();
|
||||||
} catch (XArchNetworkShutdown &) {
|
} catch (XArchNetworkShutdown &) {
|
||||||
// remote read end of stream hungup. our output side
|
// remote read end of stream hungup. our output side
|
||||||
// has therefore shutdown.
|
// has therefore shutdown.
|
||||||
|
|
@ -507,39 +508,41 @@ ISocketMultiplexerJob *TCPSocket::serviceConnected(ISocketMultiplexerJob *job, b
|
||||||
sendEvent(EventTypes::SocketDisconnected);
|
sendEvent(EventTypes::SocketDisconnected);
|
||||||
m_connected = false;
|
m_connected = false;
|
||||||
}
|
}
|
||||||
result = kNew;
|
writeResult = kNew;
|
||||||
} catch (XArchNetworkDisconnected &) {
|
} catch (XArchNetworkDisconnected &) {
|
||||||
// stream hungup
|
// stream hungup
|
||||||
onDisconnected();
|
onDisconnected();
|
||||||
sendEvent(EventTypes::SocketDisconnected);
|
sendEvent(EventTypes::SocketDisconnected);
|
||||||
result = kNew;
|
writeResult = kNew;
|
||||||
} catch (XArchNetwork &e) {
|
} catch (XArchNetwork &e) {
|
||||||
// other write error
|
// other write error
|
||||||
LOG((CLOG_WARN "error writing socket: %s", e.what()));
|
LOG((CLOG_WARN "error writing socket: %s", e.what()));
|
||||||
onDisconnected();
|
onDisconnected();
|
||||||
sendEvent(EventTypes::StreamOutputError);
|
sendEvent(EventTypes::StreamOutputError);
|
||||||
sendEvent(EventTypes::SocketDisconnected);
|
sendEvent(EventTypes::SocketDisconnected);
|
||||||
result = kNew;
|
writeResult = kNew;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if (read && m_readable) {
|
if (read && m_readable) {
|
||||||
try {
|
try {
|
||||||
result = doRead();
|
readResult = doRead();
|
||||||
} catch (XArchNetworkDisconnected &) {
|
} catch (XArchNetworkDisconnected &) {
|
||||||
// stream hungup
|
// stream hungup
|
||||||
sendEvent(EventTypes::SocketDisconnected);
|
sendEvent(EventTypes::SocketDisconnected);
|
||||||
onDisconnected();
|
onDisconnected();
|
||||||
result = kNew;
|
readResult = kNew;
|
||||||
} catch (XArchNetwork &e) {
|
} catch (XArchNetwork &e) {
|
||||||
// ignore other read error
|
// ignore other read error
|
||||||
LOG((CLOG_WARN "error reading socket: %s", e.what()));
|
LOG((CLOG_WARN "error reading socket: %s", e.what()));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if (result == kBreak) {
|
if (readResult == kBreak || writeResult == kBreak)
|
||||||
return nullptr;
|
return nullptr;
|
||||||
}
|
|
||||||
|
|
||||||
return result == kNew ? newJob() : job;
|
if (writeResult == kNew || readResult == kNew)
|
||||||
|
return newJob();
|
||||||
|
|
||||||
|
return job;
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue