nfs4: Make RPC calls asynchronous
This commit is contained in:
@@ -119,8 +119,26 @@ Server::_StartListening()
|
|||||||
status_t
|
status_t
|
||||||
Server::SendCall(Call* call, Reply** reply)
|
Server::SendCall(Call* call, Reply** reply)
|
||||||
{
|
{
|
||||||
status_t result;
|
Request* req;
|
||||||
|
status_t result = SendCallAsync(call, reply, &req);
|
||||||
|
if (result != B_OK)
|
||||||
|
return result;
|
||||||
|
|
||||||
|
result = WaitCall(req);
|
||||||
|
if (result != B_OK) {
|
||||||
|
CancelCall(req);
|
||||||
|
delete req;
|
||||||
|
return result;
|
||||||
|
}
|
||||||
|
|
||||||
|
delete req;
|
||||||
|
return B_OK;
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
status_t
|
||||||
|
Server::SendCallAsync(Call* call, Reply** reply, Request** request)
|
||||||
|
{
|
||||||
if (fThreadError != B_OK)
|
if (fThreadError != B_OK)
|
||||||
return fThreadError;
|
return fThreadError;
|
||||||
|
|
||||||
@@ -137,23 +155,17 @@ Server::SendCall(Call* call, Reply** reply)
|
|||||||
fRequests.AddRequest(req);
|
fRequests.AddRequest(req);
|
||||||
|
|
||||||
XDR::WriteStream& stream = call->Stream();
|
XDR::WriteStream& stream = call->Stream();
|
||||||
result = fConnection->Send(stream.Buffer(), stream.Size());
|
status_t result = fConnection->Send(stream.Buffer(), stream.Size());
|
||||||
if (result != B_OK)
|
if (result != B_OK) {
|
||||||
goto out_cancel;
|
|
||||||
|
|
||||||
result = req->fEvent.Wait(B_RELATIVE_TIMEOUT, kWaitTime);
|
|
||||||
if (result != B_OK)
|
|
||||||
goto out_cancel;
|
|
||||||
|
|
||||||
delete req;
|
|
||||||
return B_OK;
|
|
||||||
|
|
||||||
out_cancel:
|
|
||||||
fRequests.FindRequest(xid);
|
fRequests.FindRequest(xid);
|
||||||
delete req;
|
delete req;
|
||||||
return result;
|
return result;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
*request = req;
|
||||||
|
return B_OK;
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
status_t
|
status_t
|
||||||
Server::Repair()
|
Server::Repair()
|
||||||
|
|||||||
@@ -50,6 +50,12 @@ public:
|
|||||||
|
|
||||||
status_t SendCall(Call* call, Reply** reply);
|
status_t SendCall(Call* call, Reply** reply);
|
||||||
|
|
||||||
|
status_t SendCallAsync(Call* call, Reply** reply,
|
||||||
|
Request** request);
|
||||||
|
inline status_t WaitCall(Request* request,
|
||||||
|
bigtime_t time = kWaitTime);
|
||||||
|
inline status_t CancelCall(Request* request);
|
||||||
|
|
||||||
status_t Repair();
|
status_t Repair();
|
||||||
|
|
||||||
inline const ServerAddress& ID() const;
|
inline const ServerAddress& ID() const;
|
||||||
@@ -75,6 +81,22 @@ private:
|
|||||||
static const bigtime_t kWaitTime = 1000000;
|
static const bigtime_t kWaitTime = 1000000;
|
||||||
};
|
};
|
||||||
|
|
||||||
|
|
||||||
|
inline status_t
|
||||||
|
Server::WaitCall(Request* request, bigtime_t time)
|
||||||
|
{
|
||||||
|
return request->fEvent.Wait(B_RELATIVE_TIMEOUT, time);
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
inline status_t
|
||||||
|
Server::CancelCall(Request* request)
|
||||||
|
{
|
||||||
|
fRequests.FindRequest(request->fXID);
|
||||||
|
return B_OK;
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
inline const ServerAddress&
|
inline const ServerAddress&
|
||||||
Server::ID() const
|
Server::ID() const
|
||||||
{
|
{
|
||||||
|
|||||||
@@ -12,12 +12,28 @@
|
|||||||
|
|
||||||
status_t
|
status_t
|
||||||
Request::Send()
|
Request::Send()
|
||||||
|
{
|
||||||
|
return _TrySend();
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
status_t
|
||||||
|
Request::_TrySend()
|
||||||
{
|
{
|
||||||
RPC::Reply *rpl;
|
RPC::Reply *rpl;
|
||||||
status_t result = fServer->SendCall(fBuilder.Request(), &rpl);
|
RPC::Request *rpc;
|
||||||
|
|
||||||
|
status_t result = fServer->SendCallAsync(fBuilder.Request(), &rpl, &rpc);
|
||||||
if (result != B_OK)
|
if (result != B_OK)
|
||||||
return result;
|
return result;
|
||||||
|
|
||||||
|
result = fServer->WaitCall(rpc);
|
||||||
|
if (result != B_OK) {
|
||||||
|
fServer->CancelCall(rpc);
|
||||||
|
delete rpc;
|
||||||
|
return result;
|
||||||
|
}
|
||||||
|
|
||||||
return fReply.SetTo(rpl);
|
return fReply.SetTo(rpl);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -25,6 +25,8 @@ public:
|
|||||||
void Reset();
|
void Reset();
|
||||||
|
|
||||||
private:
|
private:
|
||||||
|
status_t _TrySend();
|
||||||
|
|
||||||
RPC::Server* fServer;
|
RPC::Server* fServer;
|
||||||
|
|
||||||
RequestBuilder fBuilder;
|
RequestBuilder fBuilder;
|
||||||
|
|||||||
Reference in New Issue
Block a user