• 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

91.53
/src/serial_client.cpp
1
#include "serial_client.h"
2
#include <chrono>
3
#include <iostream>
4
#include <unistd.h>
5

6
#include "modbus_ext_common.h"
7
#include "write_channel_serial_client_task.h"
8

9
using namespace std::chrono_literals;
10
using namespace std::chrono;
11

12
#define LOG(logger) logger.Log() << "[serial client] "
13

14
namespace
15
{
16
    const auto PORT_OPEN_ERROR_NOTIFICATION_INTERVAL = 5min;
17
    const auto CLOSED_PORT_CYCLE_TIME = 500ms;
18
    const auto MAX_POLL_TIME = 100ms;
19
    // const auto MAX_FLUSHES_WHEN_POLL_IS_DUE = 20;
20
    const auto BALANCING_THRESHOLD = 500ms;
21
    // const auto MIN_READ_EVENTS_TIME = 25ms;
22
    const size_t MAX_EVENT_READ_ERRORS = 10;
23

24
    std::chrono::milliseconds GetReadEventsPeriod(const TPort& port)
79✔
25
    {
26
        auto sendByteTime = port.GetSendTimeBytes(1);
79✔
27
        // >= 115200
28
        if (sendByteTime < 100us) {
79✔
29
            return 50ms;
2✔
30
        }
31
        // >= 38400
32
        if (sendByteTime < 300us) {
77✔
33
            return 100ms;
×
34
        }
35
        // < 38400
36
        return 200ms;
77✔
37
    }
38
};
39

40
TSerialClient::TSerialClient(PPort port,
99✔
41
                             const TPortOpenCloseLogic::TSettings& openCloseSettings,
42
                             util::TGetNowFn nowFn,
43
                             size_t lowPriorityRateLimit)
99✔
44
    : Port(port),
45
      OpenCloseLogic(openCloseSettings, nowFn),
46
      ConnectLogger(PORT_OPEN_ERROR_NOTIFICATION_INTERVAL, "[serial client] "),
47
      NowFn(nowFn),
48
      LowPriorityRateLimit(lowPriorityRateLimit)
99✔
49
{}
99✔
50

51
TSerialClient::~TSerialClient()
99✔
52
{
53
    if (Port->IsOpen()) {
99✔
54
        Port->Close();
52✔
55
    }
56
}
99✔
57

58
void TSerialClient::AddDevice(PSerialDevice device)
85✔
59
{
60
    if (RegReader)
85✔
61
        throw TSerialDeviceException("can't add registers to the active client");
×
62
    for (const auto& reg: device->GetRegisters()) {
619✔
63
        if (Handlers.find(reg) != Handlers.end())
534✔
64
            throw TSerialDeviceException("duplicate register");
×
65
        auto handler = Handlers[reg] = std::make_shared<TRegisterHandler>(reg);
534✔
66
        RegList.push_back(reg);
534✔
67
        LOG(Debug) << "AddRegister: " << reg;
534✔
68
    }
69
    Devices.push_back(device);
85✔
70
}
85✔
71

72
void TSerialClient::Activate()
620✔
73
{
74
    if (!RegReader) {
620✔
75
        RegReader = std::make_unique<TSerialClientRegisterAndEventsReader>(Devices,
79✔
76
                                                                           GetReadEventsPeriod(*Port),
79✔
77
                                                                           NowFn,
79✔
78
                                                                           LowPriorityRateLimit);
79✔
79
        LastAccessedDevice = std::make_unique<TSerialClientDeviceAccessHandler>(RegReader->GetEventsReader());
79✔
80
    }
81
}
620✔
82

83
void TSerialClient::Connect()
×
84
{
85
    OpenCloseLogic.OpenIfAllowed(Port);
×
86
}
87

88
void TSerialClient::WaitForPollAndFlush(steady_clock::time_point currentTime, steady_clock::time_point waitUntil)
620✔
89
{
90
    if (currentTime > waitUntil) {
620✔
91
        waitUntil = currentTime;
616✔
92
    }
93

94
    if (Debug.IsEnabled()) {
620✔
95
        LOG(Debug) << Port->GetDescription() << duration_cast<milliseconds>(currentTime.time_since_epoch()).count()
×
96
                   << ": Wait until " << duration_cast<milliseconds>(waitUntil.time_since_epoch()).count();
×
97
    }
98

99
    std::vector<PSerialClientTask> retryTasks;
1,240✔
100
    {
101
        std::unique_lock<std::mutex> lock(TasksMutex);
1,240✔
102

103
        while (TasksCv.wait_until(lock, waitUntil, [this]() { return !Tasks.empty(); })) {
1,978✔
104
            std::vector<PSerialClientTask> tasks;
118✔
105
            Tasks.swap(tasks);
59✔
106
            lock.unlock();
59✔
107
            for (auto& task: tasks) {
184✔
108
                if (task->Run(Port, *LastAccessedDevice, Devices) == ISerialClientTask::TRunResult::RETRY) {
125✔
109
                    retryTasks.push_back(task);
8✔
110
                }
111
            }
112
            lock.lock();
59✔
113
        }
114
    }
115
    for (auto& task: retryTasks) {
628✔
116
        AddTask(task);
8✔
117
    }
118
}
620✔
119

120
void TSerialClient::ProcessPolledRegister(PRegister reg)
1,053✔
121
{
122
    if (reg->GetErrorState().test(TRegister::ReadError) || reg->GetErrorState().test(TRegister::WriteError)) {
1,053✔
123
        if (RegisterErrorCallback) {
155✔
124
            RegisterErrorCallback(reg);
155✔
125
        }
126
    } else {
127
        if (RegisterReadCallback) {
898✔
128
            RegisterReadCallback(reg);
898✔
129
        }
130
    }
131
}
1,053✔
132

133
void TSerialClient::Cycle()
620✔
134
{
135
    Activate();
620✔
136

137
    try {
138
        OpenCloseLogic.OpenIfAllowed(Port);
623✔
139
    } catch (const std::exception& e) {
6✔
140
        ConnectLogger.Log(e.what(), Debug, Error);
3✔
141
    }
142

143
    if (Port->IsOpen()) {
620✔
144
        ConnectLogger.DropTimeout();
616✔
145
        OpenPortCycle();
616✔
146
    } else {
147
        ClosedPortCycle();
4✔
148
    }
149
}
620✔
150

151
void TSerialClient::ClosedPortCycle()
4✔
152
{
153
    auto currentTime = NowFn();
4✔
154
    auto waitUntil = currentTime + CLOSED_PORT_CYCLE_TIME;
4✔
155
    WaitForPollAndFlush(currentTime, waitUntil);
4✔
156

157
    RegReader->ClosedPortCycle(waitUntil, [this](PRegister reg) { ProcessPolledRegister(reg); });
12✔
158
}
4✔
159

160
void TSerialClient::SetTextValue(PRegister reg, const std::string& value)
115✔
161
{
162
    auto handler = GetHandler(reg);
230✔
163
    handler->SetTextValue(value);
115✔
164
    auto serialClientTask =
165
        std::make_shared<TWriteChannelSerialClientTask>(handler, RegisterReadCallback, RegisterErrorCallback);
115✔
166
    AddTask(serialClientTask);
115✔
167
}
115✔
168

169
void TSerialClient::SetReadCallback(const TSerialClient::TRegisterCallback& callback)
99✔
170
{
171
    RegisterReadCallback = callback;
99✔
172
}
99✔
173

174
void TSerialClient::SetErrorCallback(const TSerialClient::TRegisterCallback& callback)
99✔
175
{
176
    RegisterErrorCallback = callback;
99✔
177
}
99✔
178

179
PRegisterHandler TSerialClient::GetHandler(PRegister reg) const
115✔
180
{
181
    auto it = Handlers.find(reg);
115✔
182
    if (it == Handlers.end())
115✔
183
        throw TSerialDeviceException("register not found");
×
184
    return it->second;
230✔
185
}
186

187
void TSerialClient::OpenPortCycle()
616✔
188
{
189
    auto currentTime = NowFn();
616✔
190
    auto waitUntil = RegReader->GetDeadline(currentTime);
616✔
191
    // Limit waiting time to be responsive
192
    waitUntil = std::min(waitUntil, currentTime + MAX_POLL_TIME);
616✔
193
    WaitForPollAndFlush(currentTime, waitUntil);
616✔
194

195
    auto device = RegReader->OpenPortCycle(
196
        *Port,
616✔
197
        [this](PRegister reg) { ProcessPolledRegister(reg); },
1,045✔
198
        *LastAccessedDevice);
1,848✔
199

200
    if (device) {
616✔
201
        OpenCloseLogic.CloseIfNeeded(Port, device->GetConnectionState() == TDeviceConnectionState::DISCONNECTED);
616✔
202
    }
203
}
616✔
204

205
PPort TSerialClient::GetPort()
×
206
{
207
    return Port;
×
208
}
209

210
std::list<PSerialDevice> TSerialClient::GetDevices()
1✔
211
{
212
    return Devices;
1✔
213
}
214

215
void TSerialClient::AddTask(PSerialClientTask task)
126✔
216
{
217
    {
218
        std::unique_lock<std::mutex> lock(TasksMutex);
252✔
219
        Tasks.push_back(task);
126✔
220
    }
221
    TasksCv.notify_all();
126✔
222
}
126✔
223

NEW
224
void TSerialClient::SuspendPoll(PSerialDevice device, std::chrono::steady_clock::time_point currentTime)
×
225
{
NEW
226
    RegReader->SuspendPoll(device, currentTime);
×
227
}
228

NEW
229
void TSerialClient::ResumePoll(PSerialDevice device)
×
230
{
NEW
231
    RegReader->ResumePoll(device);
×
232
}
233

234
TSerialClientRegisterAndEventsReader::TSerialClientRegisterAndEventsReader(const std::list<PSerialDevice>& devices,
95✔
235
                                                                           std::chrono::milliseconds readEventsPeriod,
236
                                                                           util::TGetNowFn nowFn,
237
                                                                           size_t lowPriorityRateLimit)
95✔
238
    : EventsReader(std::make_shared<TSerialClientEventsReader>(MAX_EVENT_READ_ERRORS)),
239
      RegisterPoller(lowPriorityRateLimit),
240
      TimeBalancer(BALANCING_THRESHOLD),
241
      ReadEventsPeriod(readEventsPeriod),
242
      SpentTime(nowFn),
243
      LastCycleWasTooSmallToPoll(false),
244
      NowFn(nowFn)
95✔
245
{
246
    auto currentTime = NowFn();
95✔
247
    RegisterPoller.SetDevices(devices, currentTime);
95✔
248
    EventsReader->SetDevices(devices);
95✔
249
    TimeBalancer.AddEntry(TClientTaskType::POLLING, currentTime, TPriority::Low);
95✔
250
}
95✔
251

252
void TSerialClientRegisterAndEventsReader::ClosedPortCycle(std::chrono::steady_clock::time_point currentTime,
4✔
253
                                                           TRegisterCallback regCallback)
254
{
255
    EventsReader->SetReadErrors(regCallback);
4✔
256
    RegisterPoller.ClosedPortCycle(currentTime, regCallback);
4✔
257
}
4✔
258

259
class TSerialClientTaskHandler
260
{
261
public:
262
    TClientTaskType TaskType;
263
    TItemAccumulationPolicy Policy;
264
    milliseconds PollLimit;
265
    bool NotReady = true;
266

267
    bool operator()(TClientTaskType task, TItemAccumulationPolicy policy, milliseconds pollLimit)
873✔
268
    {
269
        PollLimit = pollLimit;
873✔
270
        TaskType = task;
873✔
271
        Policy = policy;
873✔
272
        NotReady = false;
873✔
273
        return true;
873✔
274
    }
275
};
276

277
PSerialDevice TSerialClientRegisterAndEventsReader::OpenPortCycle(TPort& port,
873✔
278
                                                                  TRegisterCallback regCallback,
279
                                                                  TSerialClientDeviceAccessHandler& lastAccessedDevice)
280
{
281
    // Count idle time as high priority task time to faster reach time balancing threshold
282
    if (LastCycleWasTooSmallToPoll) {
873✔
283
        TimeBalancer.UpdateSelectionTime(ceil<milliseconds>(SpentTime.GetSpentTime()), TPriority::High);
58✔
284
    }
285

286
    SpentTime.Start();
873✔
287
    TSerialClientTaskHandler handler;
873✔
288
    TimeBalancer.AccumulateNext(SpentTime.GetStartTime(), handler, TItemSelectionPolicy::All);
873✔
289
    if (handler.NotReady) {
873✔
290
        return nullptr;
×
291
    }
292

293
    if (handler.TaskType == TClientTaskType::EVENTS) {
873✔
294
        if (EventsReader && EventsReader->HasDevicesWithEnabledEvents()) {
84✔
295
            lastAccessedDevice.PrepareToAccess(port, nullptr);
81✔
296
            EventsReader->ReadEvents(port, MAX_POLL_TIME, regCallback, NowFn);
81✔
297
            TimeBalancer.UpdateSelectionTime(ceil<milliseconds>(SpentTime.GetSpentTime()), TPriority::High);
81✔
298
            TimeBalancer.AddEntry(TClientTaskType::EVENTS,
81✔
299
                                  SpentTime.GetStartTime() + ReadEventsPeriod,
162✔
300
                                  TPriority::High);
301
        }
302
        SpentTime.Start();
84✔
303

304
        // TODO: Need to notify port open/close logic about errors
305
        return nullptr;
84✔
306
    }
307

308
    // Some registers can have theoretical read time more than poll limit.
309
    // Define special cases when reading can exceed poll limit to read the registers:
310
    // 1. TimeBalancer can force reading of such registers.
311
    // 2. If there are not devices with enabled events, the only limiting timeout is MAX_POLL_TIME.
312
    //    We can miss it and read at least one register.
313
    const bool readAtLeastOneRegister =
314
        (handler.Policy == TItemAccumulationPolicy::Force) || !EventsReader->HasDevicesWithEnabledEvents();
789✔
315

316
    auto res = RegisterPoller.OpenPortCycle(port,
317
                                            SpentTime,
789✔
318
                                            std::min(handler.PollLimit, MAX_POLL_TIME),
1,578✔
319
                                            readAtLeastOneRegister,
320
                                            lastAccessedDevice,
321
                                            regCallback);
1,578✔
322

323
    TimeBalancer.AddEntry(TClientTaskType::POLLING, res.Deadline, TPriority::Low);
789✔
324
    if (res.NotEnoughTime) {
789✔
325
        LastCycleWasTooSmallToPoll = true;
37✔
326
    } else {
327
        LastCycleWasTooSmallToPoll = false;
752✔
328
        TimeBalancer.UpdateSelectionTime(ceil<milliseconds>(SpentTime.GetSpentTime()), TPriority::Low);
752✔
329
    }
330

331
    if (EventsReader->HasDevicesWithEnabledEvents() && !TimeBalancer.Contains(TClientTaskType::EVENTS)) {
789✔
332
        TimeBalancer.AddEntry(TClientTaskType::EVENTS, SpentTime.GetStartTime() + ReadEventsPeriod, TPriority::High);
12✔
333
    }
334

335
    SpentTime.Start();
789✔
336
    return res.Device;
789✔
337
}
338

339
std::chrono::steady_clock::time_point TSerialClientRegisterAndEventsReader::GetDeadline(
873✔
340
    std::chrono::steady_clock::time_point currentTime) const
341
{
342
    return TimeBalancer.GetDeadline();
873✔
343
}
344

345
PSerialClientEventsReader TSerialClientRegisterAndEventsReader::GetEventsReader() const
95✔
346
{
347
    return EventsReader;
95✔
348
}
349

350
void TSerialClientRegisterAndEventsReader::SuspendPoll(PSerialDevice device,
3✔
351
                                                       std::chrono::steady_clock::time_point currentTime)
352
{
353
    RegisterPoller.SuspendPoll(device, currentTime);
3✔
354
}
3✔
355

356
void TSerialClientRegisterAndEventsReader::ResumePoll(PSerialDevice device)
2✔
357
{
358
    RegisterPoller.ResumePoll(device);
2✔
359
}
2✔
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