-
Notifications
You must be signed in to change notification settings - Fork 47
Expand file tree
/
Copy pathendpoint.h
More file actions
121 lines (88 loc) · 3.08 KB
/
Copy pathendpoint.h
File metadata and controls
121 lines (88 loc) · 3.08 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
#pragma once
#include "public.h"
#include "app_context.h"
#include "executor.h"
#include "server.h"
#include "vhost.h"
#include <cloud/blockstore/libs/diagnostics/incomplete_requests.h>
#include <cloud/blockstore/libs/diagnostics/server_stats.h>
#include <cloud/blockstore/libs/service/context.h>
#include <cloud/blockstore/libs/service/public.h>
#include <library/cpp/threading/future/future.h>
#include <util/generic/intrlist.h>
#include <util/generic/string.h>
#include <util/system/spinlock.h>
#include <atomic>
#include <memory>
namespace NCloud::NBlockStore::NVhost {
////////////////////////////////////////////////////////////////////////////////
struct TRequest
: public TIntrusiveListItem<TRequest>
, TAtomicRefCount<TRequest>
{
const TVhostRequestPtr VhostRequest;
const TCallContextPtr CallContext;
TMetricRequest MetricRequest;
std::atomic_flag Completed = 0;
TRequest(ui64 requestId, TVhostRequestPtr vhostRequest)
: VhostRequest(std::move(vhostRequest))
, CallContext(MakeIntrusive<TCallContext>(requestId))
, MetricRequest(VhostRequest->Type)
{}
};
using TRequestPtr = TIntrusivePtr<TRequest>;
////////////////////////////////////////////////////////////////////////////////
// A single vhost block device exposed to the guest. Translates the vhost
// requests dispatched to it by its executor into the IDeviceHandler API and
// keeps track of the requests in flight.
class TEndpoint final
: public IRequestProcessor
, public std::enable_shared_from_this<TEndpoint>
{
private:
TAppContext& AppCtx;
const IDeviceHandlerPtr DeviceHandler;
const TString SocketPath;
const TStorageOptions Options;
const ui32 SocketAccessMode;
TExecutor* const Executor;
IVhostDevicePtr VhostDevice;
TIntrusiveList<TRequest> RequestsInFlight;
TAdaptiveLock RequestsLock;
std::atomic_flag Stopped = false;
public:
TEndpoint(
TAppContext& appCtx,
IDeviceHandlerPtr deviceHandler,
TString socketPath,
TStorageOptions options,
ui32 socketAccessMode,
TExecutor* executor);
void* GetCookie()
{
return static_cast<IRequestProcessor*>(this);
}
TExecutor* GetExecutor() const
{
return Executor;
}
void SetVhostDevice(IVhostDevicePtr vhostDevice);
NProto::TError Start();
NThreading::TFuture<NProto::TError> Stop(bool deleteSocket);
void Update(ui64 blocksCount);
ui32 GetVhostQueuesCount() const
{
return Options.VhostQueuesCount;
}
size_t CollectRequests(const TIncompleteRequestsCollector& collector);
void ProcessRequest(TVhostRequestPtr vhostRequest) override;
private:
template <typename TMethod>
void ProcessRequest(TRequestPtr request);
TRequestPtr RegisterRequest(TVhostRequestPtr vhostRequest);
void CompleteRequest(TRequest& request, const NProto::TError& error);
void UnregisterRequest(TRequest& request);
TVhostRequest::EResult GetResult(NProto::TError& error);
};
using TEndpointPtr = std::shared_ptr<TEndpoint>;
} // namespace NCloud::NBlockStore::NVhost