• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

wirenboard / wb-mqtt-serial / 681

08 Aug 2025 02:21PM UTC coverage: 73.114% (+0.1%) from 73.001%
681

push

github

web-flow
Add device/SetPoll RPC for device poll suspending and resuming

6595 of 9378 branches covered (70.32%)

121 of 144 new or added lines in 5 files covered. (84.03%)

12531 of 17139 relevant lines covered (73.11%)

373.28 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

13.44
/src/rpc/rpc_device_handler.cpp
1
#include "rpc_device_handler.h"
2
#include "rpc_device_load_config_task.h"
3
#include "rpc_device_load_task.h"
4
#include "rpc_device_probe_task.h"
5
#include "rpc_device_set_task.h"
6
#include "rpc_helpers.h"
7

8
#define LOG(logger) ::logger.Log() << "[RPC] "
9

10
void TRPCDeviceParametersCache::RegisterCallbacks(PHandlerConfig handlerConfig)
×
11
{
12
    for (const auto& portConfig: handlerConfig->PortConfigs) {
×
13
        for (const auto& device: portConfig->Devices) {
×
14
            std::string id = GetId(*portConfig->Port, device->Device->DeviceConfig()->SlaveId);
×
15
            device->Device->AddOnConnectionStateChangedCallback([this, id](PSerialDevice device) {
×
16
                if (device->GetConnectionState() == TDeviceConnectionState::DISCONNECTED) {
×
17
                    Remove(id);
×
18
                }
19
            });
×
20
        }
21
    }
22
}
23

24
std::string TRPCDeviceParametersCache::GetId(const TPort& port, const std::string& slaveId) const
×
25
{
26
    return port.GetDescription(false) + ":" + slaveId;
×
27
}
28

29
void TRPCDeviceParametersCache::Add(const std::string& id, const Json::Value& value)
×
30
{
31
    std::unique_lock lock(Mutex);
×
32
    DeviceParameters[id] = value;
×
33
}
34

35
void TRPCDeviceParametersCache::Remove(const std::string& id)
×
36
{
37
    std::unique_lock lock(Mutex);
×
38
    DeviceParameters.erase(id);
×
39
}
40

41
bool TRPCDeviceParametersCache::Contains(const std::string& id) const
×
42
{
43
    std::unique_lock lock(Mutex);
×
44
    return DeviceParameters.find(id) != DeviceParameters.end();
×
45
}
46

47
const Json::Value& TRPCDeviceParametersCache::Get(const std::string& id, const Json::Value& defaultValue) const
×
48
{
49
    std::unique_lock lock(Mutex);
×
50
    auto it = DeviceParameters.find(id);
×
51
    return it != DeviceParameters.end() ? it->second : defaultValue;
×
52
};
53

54
TRPCDeviceHelper::TRPCDeviceHelper(const Json::Value& request,
×
55
                                   const TSerialDeviceFactory& deviceFactory,
56
                                   PTemplateMap templates,
57
                                   TSerialClientTaskRunner& serialClientTaskRunner)
×
58
{
59
    auto params = serialClientTaskRunner.GetSerialClientParams(request);
×
60
    if (params.Device == nullptr) {
×
61
        DeviceTemplate = templates->GetTemplate(request["device_type"].asString());
×
62
        auto config = std::make_shared<TDeviceConfig>("RPC Device",
63
                                                      request["slave_id"].asString(),
×
64
                                                      DeviceTemplate->GetProtocol());
×
65
        if (DeviceTemplate->GetProtocol() == "modbus") {
×
66
            config->MaxRegHole = Modbus::MAX_HOLE_CONTINUOUS_16_BIT_REGISTERS;
×
67
            config->MaxBitHole = Modbus::MAX_HOLE_CONTINUOUS_1_BIT_REGISTERS;
×
68
            config->MaxReadRegisters = Modbus::MAX_READ_REGISTERS;
×
69
        }
70
        ProtocolParams = deviceFactory.GetProtocolParams(DeviceTemplate->GetProtocol());
×
71
        Device = ProtocolParams.factory->CreateDevice(DeviceTemplate->GetTemplate(), config, ProtocolParams.protocol);
×
72
    } else {
73
        Device = params.Device;
×
74
        DeviceTemplate = templates->GetTemplate(Device->DeviceConfig()->DeviceType);
×
75
        ProtocolParams = deviceFactory.GetProtocolParams(DeviceTemplate->GetProtocol());
×
76
        DeviceFromConfig = true;
×
77
    }
78
    if (DeviceTemplate->WithSubdevices()) {
×
79
        throw TRPCException("Device \"" + DeviceTemplate->Type + "\" is not supported by this RPC",
×
80
                            TRPCResultCode::RPC_WRONG_PARAM_VALUE);
×
81
    }
82
}
83

84
TRPCDeviceRequest::TRPCDeviceRequest(const TDeviceProtocolParams& protocolParams,
×
85
                                     PSerialDevice device,
86
                                     PDeviceTemplate deviceTemplate,
87
                                     bool deviceFromConfig)
×
88
    : ProtocolParams(protocolParams),
89
      Device(device),
90
      DeviceTemplate(deviceTemplate),
91
      DeviceFromConfig(deviceFromConfig)
×
92
{
93
    Json::Value responseTimeout = DeviceTemplate->GetTemplate()["response_timeout_ms"];
×
94
    if (responseTimeout.isInt()) {
×
95
        ResponseTimeout = std::chrono::milliseconds(responseTimeout.asInt());
×
96
    }
97

98
    Json::Value frameTimeout = DeviceTemplate->GetTemplate()["frame_timeout_ms"];
×
99
    if (frameTimeout.isInt()) {
×
100
        FrameTimeout = std::chrono::milliseconds(frameTimeout.asInt());
×
101
    }
102
}
103

104
void TRPCDeviceRequest::ParseSettings(const Json::Value& request,
×
105
                                      WBMQTT::TMqttRpcServer::TResultCallback onResult,
106
                                      WBMQTT::TMqttRpcServer::TErrorCallback onError)
107
{
108
    SerialPortSettings = ParseRPCSerialPortSettings(request);
×
109
    WBMQTT::JSON::Get(request, "response_timeout", ResponseTimeout);
×
110
    WBMQTT::JSON::Get(request, "frame_timeout", FrameTimeout);
×
111
    WBMQTT::JSON::Get(request, "total_timeout", TotalTimeout);
×
112
    OnResult = onResult;
×
113
    OnError = onError;
×
114
}
115

116
TRPCDeviceHandler::TRPCDeviceHandler(const std::string& requestDeviceLoadConfigSchemaFilePath,
×
117
                                     const std::string& requestDeviceLoadSchemaFilePath,
118
                                     const std::string& requestDeviceSetSchemaFilePath,
119
                                     const std::string& requestDeviceProbeSchemaFilePath,
120
                                     const std::string& requestDeviceSetPollSchemaFilePath,
121
                                     const TSerialDeviceFactory& deviceFactory,
122
                                     PTemplateMap templates,
123
                                     TSerialClientTaskRunner& serialClientTaskRunner,
124
                                     TRPCDeviceParametersCache& parametersCache,
125
                                     WBMQTT::PMqttRpcServer rpcServer)
×
126
    : DeviceFactory(deviceFactory),
127
      RequestDeviceLoadConfigSchema(LoadRPCRequestSchema(requestDeviceLoadConfigSchemaFilePath, "device/LoadConfig")),
×
128
      RequestDeviceLoadSchema(LoadRPCRequestSchema(requestDeviceLoadSchemaFilePath, "device/Load")),
×
129
      RequestDeviceSetSchema(LoadRPCRequestSchema(requestDeviceSetSchemaFilePath, "device/Set")),
×
130
      RequestDeviceProbeSchema(LoadRPCRequestSchema(requestDeviceProbeSchemaFilePath, "device/Probe")),
×
NEW
131
      RequestDeviceSetPollSchema(LoadRPCRequestSchema(requestDeviceSetPollSchemaFilePath, "device/SetPoll")),
×
132
      Templates(templates),
133
      SerialClientTaskRunner(serialClientTaskRunner),
134
      ParametersCache(parametersCache)
×
135
{
136
    rpcServer->RegisterAsyncMethod("device",
×
137
                                   "LoadConfig",
138
                                   std::bind(&TRPCDeviceHandler::LoadConfig,
×
139
                                             this,
×
140
                                             std::placeholders::_1,
141
                                             std::placeholders::_2,
142
                                             std::placeholders::_3));
×
143
    rpcServer->RegisterAsyncMethod("device",
×
144
                                   "Load",
145
                                   std::bind(&TRPCDeviceHandler::Load, //
×
146
                                             this,
×
147
                                             std::placeholders::_1,
148
                                             std::placeholders::_2,
149
                                             std::placeholders::_3));
×
150
    rpcServer->RegisterAsyncMethod("device",
×
151
                                   "Set",
152
                                   std::bind(&TRPCDeviceHandler::Set, //
×
153
                                             this,
×
154
                                             std::placeholders::_1,
155
                                             std::placeholders::_2,
156
                                             std::placeholders::_3));
×
157
    rpcServer->RegisterAsyncMethod("device",
×
158
                                   "Probe",
159
                                   std::bind(&TRPCDeviceHandler::Probe,
×
160
                                             this,
×
161
                                             std::placeholders::_1,
162
                                             std::placeholders::_2,
163
                                             std::placeholders::_3));
×
164

NEW
165
    rpcServer->RegisterMethod("device", "SetPoll", std::bind(&TRPCDeviceHandler::SetPoll, this, std::placeholders::_1));
×
166
}
167

168
void TRPCDeviceHandler::LoadConfig(const Json::Value& request,
×
169
                                   WBMQTT::TMqttRpcServer::TResultCallback onResult,
170
                                   WBMQTT::TMqttRpcServer::TErrorCallback onError)
171
{
172
    ValidateRPCRequest(request, RequestDeviceLoadConfigSchema);
×
173
    try {
174
        auto helper = TRPCDeviceHelper(request, DeviceFactory, Templates, SerialClientTaskRunner);
×
175
        auto rpcRequest = ParseRPCDeviceLoadConfigRequest(request,
176
                                                          helper.ProtocolParams,
177
                                                          helper.Device,
178
                                                          helper.DeviceTemplate,
179
                                                          helper.DeviceFromConfig,
180
                                                          ParametersCache,
181
                                                          onResult,
182
                                                          onError);
×
183
        SerialClientTaskRunner.RunTask(request, std::make_shared<TRPCDeviceLoadConfigSerialClientTask>(rpcRequest));
×
184
    } catch (const TRPCException& e) {
×
185
        ProcessException(e, onError);
×
186
    }
187
}
188

189
void TRPCDeviceHandler::Load(const Json::Value& request,
×
190
                             WBMQTT::TMqttRpcServer::TResultCallback onResult,
191
                             WBMQTT::TMqttRpcServer::TErrorCallback onError)
192
{
193
    ValidateRPCRequest(request, RequestDeviceLoadSchema);
×
194
    try {
195
        auto helper = TRPCDeviceHelper(request, DeviceFactory, Templates, SerialClientTaskRunner);
×
196
        auto rpcRequest = ParseRPCDeviceLoadRequest(request,
197
                                                    helper.ProtocolParams,
198
                                                    helper.Device,
199
                                                    helper.DeviceTemplate,
200
                                                    helper.DeviceFromConfig,
201
                                                    onResult,
202
                                                    onError);
×
203
        SerialClientTaskRunner.RunTask(request, std::make_shared<TRPCDeviceLoadSerialClientTask>(rpcRequest));
×
204
    } catch (const TRPCException& e) {
×
205
        ProcessException(e, onError);
×
206
    }
207
}
208

209
void TRPCDeviceHandler::Set(const Json::Value& request,
×
210
                            WBMQTT::TMqttRpcServer::TResultCallback onResult,
211
                            WBMQTT::TMqttRpcServer::TErrorCallback onError)
212
{
213
    ValidateRPCRequest(request, RequestDeviceSetSchema);
×
214
    try {
215
        auto helper = TRPCDeviceHelper(request, DeviceFactory, Templates, SerialClientTaskRunner);
×
216
        auto rpcRequest = ParseRPCDeviceSetRequest(request,
217
                                                   helper.ProtocolParams,
218
                                                   helper.Device,
219
                                                   helper.DeviceTemplate,
220
                                                   helper.DeviceFromConfig,
221
                                                   onResult,
222
                                                   onError);
×
223
        SerialClientTaskRunner.RunTask(request, std::make_shared<TRPCDeviceSetSerialClientTask>(rpcRequest));
×
224
    } catch (const TRPCException& e) {
×
225
        ProcessException(e, onError);
×
226
    }
227
}
228

229
void TRPCDeviceHandler::Probe(const Json::Value& request,
×
230
                              WBMQTT::TMqttRpcServer::TResultCallback onResult,
231
                              WBMQTT::TMqttRpcServer::TErrorCallback onError)
232
{
233
    ValidateRPCRequest(request, RequestDeviceProbeSchema);
×
234
    try {
235
        SerialClientTaskRunner.RunTask(request,
×
236
                                       std::make_shared<TRPCDeviceProbeSerialClientTask>(request, onResult, onError));
×
237
    } catch (const TRPCException& e) {
×
238
        ProcessException(e, onError);
×
239
    }
240
}
241

NEW
242
Json::Value TRPCDeviceHandler::SetPoll(const Json::Value& request)
×
243
{
NEW
244
    ValidateRPCRequest(request, RequestDeviceSetPollSchema);
×
NEW
245
    auto params = SerialClientTaskRunner.GetSerialClientParams(request);
×
NEW
246
    if (!params.SerialClient || !params.Device) {
×
NEW
247
        throw TRPCException("Port or device not found", TRPCResultCode::RPC_WRONG_PARAM_VALUE);
×
248
    }
249
    try {
NEW
250
        if (!request["poll"].asBool()) {
×
NEW
251
            params.SerialClient->SuspendPoll(params.Device, std::chrono::steady_clock::now());
×
252
        } else {
NEW
253
            params.SerialClient->ResumePoll(params.Device);
×
254
        }
NEW
255
    } catch (const std::runtime_error& e) {
×
NEW
256
        LOG(Warn) << e.what();
×
NEW
257
        throw TRPCException(e.what(), TRPCResultCode::RPC_WRONG_PARAM_VALUE);
×
258
    }
NEW
259
    return Json::Value();
×
260
}
261

262
TRPCRegisterList CreateRegisterList(const TDeviceProtocolParams& protocolParams,
4✔
263
                                    const PSerialDevice& device,
264
                                    const Json::Value& templateItems,
265
                                    const Json::Value& knownItems,
266
                                    const std::string& fwVersion)
267
{
268
    TRPCRegisterList registerList;
4✔
269
    for (auto it = templateItems.begin(); it != templateItems.end(); ++it) {
30✔
270
        const auto& item = *it;
26✔
271
        auto id = templateItems.isObject() ? it.key().asString() : item["id"].asString();
26✔
272
        bool duplicate = false;
26✔
273
        for (const auto& item: registerList) {
67✔
274
            if (item.first == id) {
42✔
275
                duplicate = true;
1✔
276
                break;
1✔
277
            }
278
        }
279
        if (duplicate || item["address"].isNull() || item["readonly"].asBool() || !knownItems[id].isNull()) {
26✔
280
            continue;
5✔
281
        }
282
        if (!fwVersion.empty()) {
21✔
283
            std::string fw = item["fw"].asString();
12✔
284
            if (!fw.empty() && util::CompareVersionStrings(fw, fwVersion) > 0) {
12✔
285
                continue;
4✔
286
            }
287
        }
288
        auto config = LoadRegisterConfig(item,
289
                                         *protocolParams.protocol->GetRegTypes(),
17✔
290
                                         std::string(),
34✔
291
                                         *protocolParams.factory,
17✔
292
                                         protocolParams.factory->GetRegisterAddressFactory().GetBaseRegisterAddress(),
17✔
293
                                         0);
34✔
294
        auto reg = std::make_shared<TRegister>(device, config.RegisterConfig);
17✔
295
        reg->SetAvailable(TRegisterAvailability::AVAILABLE);
17✔
296
        registerList.push_back(std::make_pair(id, reg));
17✔
297
    }
298
    return registerList;
4✔
299
}
300

301
void ReadRegisterList(TPort& port,
×
302
                      PSerialDevice device,
303
                      TRPCRegisterList& registerList,
304
                      Json::Value& result,
305
                      int maxRetries)
306
{
307
    if (registerList.size() == 0) {
×
308
        return;
×
309
    }
310
    TRegisterComparePredicate compare;
311
    std::sort(registerList.begin(),
×
312
              registerList.end(),
313
              [compare](std::pair<std::string, PRegister>& a, std::pair<std::string, PRegister>& b) {
×
314
                  return compare(b.second, a.second);
×
315
              });
316

317
    std::string error;
×
318
    for (int i = 0; i <= maxRetries; i++) {
×
319
        try {
320
            device->Prepare(port, TDevicePrepareMode::WITHOUT_SETUP);
×
321
            break;
×
322
        } catch (const TSerialDeviceException& e) {
×
323
            if (i == maxRetries) {
×
324
                error = std::string("Failed to prepare session: ") + e.what();
×
325
                LOG(Warn) << port.GetDescription() << " " << device->ToString() << ": " << error;
×
326
                throw TRPCException(error, TRPCResultCode::RPC_WRONG_PARAM_VALUE);
×
327
            }
328
        }
329
    }
330

331
    size_t index = 0;
×
332
    while (index < registerList.size() && error.empty()) {
×
333
        auto first = registerList[index].second;
×
334
        auto range = device->CreateRegisterRange();
×
335
        while (index < registerList.size() &&
×
336
               range->Add(port, registerList[index].second, std::chrono::milliseconds::max()))
×
337
        {
338
            ++index;
×
339
        }
340
        for (int i = 0; i <= maxRetries; ++i) {
×
341
            try {
342
                device->ReadRegisterRange(port, range, true);
×
343
                break;
×
344
            } catch (const TSerialDeviceException& e) {
×
345
                if (i == maxRetries) {
×
346
                    error = "Failed to read " + std::to_string(range->RegisterList().size()) +
×
347
                            " registers starting from <" + first->GetConfig()->ToString() + ">: " + e.what();
×
348
                }
349
            }
350
        }
351
    }
352

353
    try {
354
        device->EndSession(port);
×
355
    } catch (const TSerialDeviceException& e) {
×
356
        LOG(Warn) << port.GetDescription() << " " << device->ToString() << " unable to end session: " << e.what();
×
357
    }
358

359
    if (!error.empty()) {
×
360
        LOG(Warn) << port.GetDescription() << " " << device->ToString() << ": " << error;
×
361
        throw TRPCException(error, TRPCResultCode::RPC_WRONG_PARAM_VALUE);
×
362
    }
363

364
    for (size_t i = 0; i < registerList.size(); ++i) {
×
365
        auto& reg = registerList[i];
×
366
        result[reg.first] = RawValueToJSON(*reg.second->GetConfig(), reg.second->GetValue());
×
367
    }
368
}
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc