• Home
  • Line#
  • Scopes#
  • Navigate#
  • Raw
  • Download
1 /*
2  * Copyright (c) 2022 Huawei Device Co., Ltd.
3  * Licensed under the Apache License, Version 2.0 (the "License");
4  * you may not use this file except in compliance with the License.
5  * You may obtain a copy of the License at
6  *
7  *     http://www.apache.org/licenses/LICENSE-2.0
8  *
9  * Unless required by applicable law or agreed to in writing, software
10  * distributed under the License is distributed on an "AS IS" BASIS,
11  * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12  * See the License for the specific language governing permissions and
13  * limitations under the License.
14  */
15 
16 #ifndef REMOTE_EXECUTOR_H
17 #define REMOTE_EXECUTOR_H
18 
19 #include <deque>
20 #include <queue>
21 
22 #include "db_types.h"
23 #include "distributeddb/result_set.h"
24 #include "icommunicator.h"
25 #include "isync_interface.h"
26 #include "message.h"
27 #include "relational_db_sync_interface.h"
28 #include "relational_result_set_impl.h"
29 #include "remote_executor_packet.h"
30 #include "runtime_context.h"
31 
32 namespace DistributedDB {
33 class RemoteExecutor : public RefObject {
34 public:
35     enum class Status {
36         WAITING = 0,
37         WORKING
38     };
39 
40     using OnFinished = std::function<void(int32_t, std::shared_ptr<ResultSet>)>;
41 
42     struct Task {
43         Status status = Status::WAITING;
44         uint32_t taskId = 0u;
45         uint64_t timeout = 0u;
46         uint32_t targetCount = 0;
47         uint32_t currentCount = 0;
48         uint64_t connectionId = 0u;
49         std::string target;
50         RemoteCondition condition;
51         OnFinished onFinished = nullptr;
52         std::shared_ptr<RelationalResultSetImpl> result;
53     };
54 
55     RemoteExecutor();
56 
57     ~RemoteExecutor() = default;
58 
59     int Initialize(ISyncInterface *syncInterface, ICommunicator *communicator);
60 
61     int RemoteQuery(const std::string &device, const RemoteCondition &condition,
62         uint64_t timeout, uint64_t connectionId, std::shared_ptr<ResultSet> &result);
63 
64     // receive request and ack, and process in another thread
65     int ReceiveMessage(const std::string &targetDev, Message *inMsg);
66 
67     void NotifyDeviceOffline(const std::string &device);
68 
69     void NotifyUserChange();
70 
71     void Close();
72 
73     void NotifyConnectionClosed(uint64_t connectionId);
74 
75     int32_t GetTaskCount() const;
76 protected:
77     virtual void ParseOneRequestMessage(const std::string &device, Message *inMsg);
78 
79     virtual bool IsPacketValid(uint32_t sessionId);
80 
81     int ResponseFailed(int errCode, uint32_t sessionId, uint32_t sequenceId, const std::string &device);
82 
83 private:
84     struct SendMessage {
85         uint32_t sessionId = 0u;
86         uint32_t sequenceId = 0u;
87         bool isLast = true;
88         SecurityOption option;
89     };
90 
91     void ReceiveMessageInner(const std::string &targetDev, Message *inMsg);
92 
93     int ReceiveRemoteExecutorRequest(const std::string &targetDev, Message *inMsg);
94 
95     int ReceiveRemoteExecutorAck(const std::string &targetDev, Message *inMsg);
96 
97     int CheckPermissions(const std::string &device, Message *inMsg);
98 
99     int SendRemoteExecutorData(const std::string &device, const Message *inMsg);
100 
101     bool CheckParamValid(const std::string &device, uint64_t timeout) const;
102 
103     bool CheckTaskExeStatus(const std::string &device);
104 
105     uint32_t GenerateSessionId();
106     uint32_t GenerateTaskId();
107 
108     int RemoteQueryInner(const Task &task);
109     void TryExecuteTaskInLock(const std::string &device);
110     void DoRollBack(uint32_t sessionId);
111 
112     int RequestStart(uint32_t sessionId);
113     int SendRequestMessage(const std::string &target, Message *message, uint32_t sessionId);
114 
115     int ResponseData(RelationalRowDataSet &&dataSet, const SendMessage &sendMessage, const std::string &device);
116     int ResponseStart(RemoteExecutorAckPacket *packet, uint32_t sessionId, uint32_t sequenceId,
117         const std::string &device);
118 
119     void StartTimer(uint64_t timeout, uint32_t sessionId);
120     void RemoveTimer(uint32_t sessionId);
121     int TimeoutCallBack(TimerId timerId);
122     void DoTimeout(TimerId timerId);
123 
124     void DoSendFailed(uint32_t sessionId, int errCode);
125     void DoFinished(uint32_t sessionId, int errCode);
126 
127     int ClearTaskInfo(uint32_t sessionId, Task &task);
128     void ClearInnerSource();
129 
130     int FillRequestPacket(RemoteExecutorRequestPacket *packet, uint32_t sessionId, std::string &target);
131 
132     void ReceiveDataWithValidSession(const std::string &targetDev, uint32_t sessionId, uint32_t sequenceId,
133         const RemoteExecutorAckPacket *packet);
134 
135     void RemoveTaskByDevice(const std::string &device, std::vector<uint32_t> &removeList);
136     void RemoveAllTask(int errCode);
137     void RemoveTaskByConnection(uint64_t connectionId, std::vector<uint32_t> &removeList);
138 
139     int GetPacketSize(const std::string &device, size_t &packetSize) const;
140     int CheckSecurityOption(ISyncInterface *storage, ICommunicator *communicator, const SecurityOption &remoteOption);
141     bool CheckRemoteSecurityOption(const std::string &device, const SecurityOption &remoteOption,
142         const SecurityOption &localOption);
143     int ResponseRemoteQueryRequest(RelationalDBSyncInterface *storage, const PreparedStmt &stmt,
144         const std::string &device, uint32_t sessionId);
145 
146     ICommunicator *GetAndIncCommunicator() const;
147     ISyncInterface *GetAndIncSyncInterface() const;
148     static int CheckRemoteRecvData(const std::string &device, SyncGenericInterface *storage, int32_t remoteSecLabel,
149         uint32_t remoteVersion);
150 
151     mutable std::mutex taskLock_;
152     std::map<std::string, std::deque<uint32_t>> searchTaskQueue_; // key is device, value is sessionId queue
153     std::map<std::string, std::set<uint32_t>> deviceWorkingSet_; // key is device, value is sessionId set
154     std::map<uint32_t, Task> taskMap_; // key is sessionId
155 
156     std::mutex timeoutLock_;
157     std::map<TimerId, uint32_t> timeoutMap_; // use to abort task when timeout
158     std::map<uint32_t, TimerId> taskFinishMap_; // use to wake up timer when task finished
159 
160     std::mutex msgQueueLock_;
161     std::queue<std::pair<std::string, Message *>> searchMessageQueue_;
162     std::atomic<uint32_t> workingThreadsCount_;
163 
164     mutable std::mutex innerSourceLock_;
165     ISyncInterface *syncInterface_;
166     ICommunicator *communicator_;
167 
168     uint32_t lastSessionId_;
169     uint32_t lastTaskId_;
170     std::atomic<bool> closed_;
171 
172     std::condition_variable clearCV_;  // msgQueueLock_
173 };
174 }
175 #endif