• 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

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

6
#include "log.h"
7
#include "serial_device.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 MAX_LOW_PRIORITY_LAG = 1s;
17
    const auto SUSPEND_POLL_TIMEOUT = std::chrono::minutes(10);
18

19
    class TDeviceReader
20
    {
21
        PRegisterRange RegisterRange;
22
        milliseconds MaxPollTime;
23
        PPollableDevice Device;
24
        bool ReadAtLeastOneRegister;
25
        const util::TSpentTimeMeter& SessionTime;
26
        TSerialClientDeviceAccessHandler& LastAccessedDevice;
27
        TPort& Port;
28

29
    public:
30
        TDeviceReader(TPort& port,
789✔
31
                      const util::TSpentTimeMeter& sessionTime,
32
                      milliseconds maxPollTime,
33
                      bool readAtLeastOneRegister,
34
                      TSerialClientDeviceAccessHandler& lastAccessedDevice)
35
            : MaxPollTime(maxPollTime),
789✔
36
              ReadAtLeastOneRegister(readAtLeastOneRegister),
37
              SessionTime(sessionTime),
38
              LastAccessedDevice(lastAccessedDevice),
39
              Port(port)
789✔
40
        {}
789✔
41

42
        bool operator()(const PPollableDevice& device, TItemAccumulationPolicy policy, milliseconds pollLimit)
813✔
43
        {
44
            if (Device) {
813✔
45
                return false;
39✔
46
            }
47

48
            pollLimit = std::min(MaxPollTime, pollLimit);
774✔
49
            if (policy != TItemAccumulationPolicy::Force) {
774✔
50
                ReadAtLeastOneRegister = false;
39✔
51
            }
52

53
            RegisterRange =
54
                device->ReadRegisterRange(Port, pollLimit, ReadAtLeastOneRegister, SessionTime, LastAccessedDevice);
774✔
55
            Device = device;
774✔
56
            return !RegisterRange->RegisterList().empty();
774✔
57
        }
58

59
        PRegisterRange GetRegisterRange() const
789✔
60
        {
61
            return RegisterRange;
789✔
62
        }
63

64
        PPollableDevice GetDevice() const
3,414✔
65
        {
66
            return Device;
3,414✔
67
        }
68
    };
69

70
    class TClosedPortDeviceReader
71
    {
72
        std::list<PRegister> Regs;
73
        steady_clock::time_point CurrentTime;
74
        PPollableDevice Device;
75

76
    public:
77
        TClosedPortDeviceReader(steady_clock::time_point currentTime): CurrentTime(currentTime)
4✔
78
        {}
4✔
79

80
        bool operator()(const PPollableDevice& device, TItemAccumulationPolicy policy, milliseconds pollLimit)
4✔
81
        {
82
            if (Device) {
4✔
83
                return false;
×
84
            }
85
            Regs = device->MarkWaitingRegistersAsReadErrorAndReschedule(CurrentTime);
4✔
86
            Device = device;
4✔
87
            return true;
4✔
88
        }
89

90
        std::list<PRegister>& GetRegisters()
16✔
91
        {
92
            return Regs;
16✔
93
        }
94

95
        PPollableDevice GetDevice() const
20✔
96
        {
97
            return Device;
20✔
98
        }
99

100
        void ClearRegisters()
8✔
101
        {
102
            Regs.clear();
8✔
103
            Device.reset();
8✔
104
        }
8✔
105
    };
106
};
107

108
TSerialClientRegisterPoller::TSerialClientRegisterPoller(size_t lowPriorityRateLimit)
95✔
109
    : Scheduler(MAX_LOW_PRIORITY_LAG),
110
      ThrottlingStateLogger(),
111
      LowPriorityRateLimiter(lowPriorityRateLimit)
95✔
112
{}
95✔
113

114
void TSerialClientRegisterPoller::SetDevices(const std::list<PSerialDevice>& devices,
95✔
115
                                             steady_clock::time_point currentTime)
116
{
117
    std::unique_lock lock(Mutex);
190✔
118

119
    for (const auto& dev: devices) {
195✔
120
        auto pollableDevice = std::make_shared<TPollableDevice>(dev, currentTime, TPriority::High);
100✔
121
        if (pollableDevice->HasRegisters()) {
100✔
122
            Scheduler.AddEntry(pollableDevice, currentTime, TPriority::High);
6✔
123
            Devices.insert({dev, pollableDevice});
6✔
124
        }
125
        pollableDevice = std::make_shared<TPollableDevice>(dev, currentTime, TPriority::Low);
100✔
126
        if (pollableDevice->HasRegisters()) {
100✔
127
            Scheduler.AddEntry(pollableDevice, currentTime, TPriority::Low);
98✔
128
            Devices.insert({dev, pollableDevice});
98✔
129
        }
130
        dev->AddOnConnectionStateChangedCallback(
200✔
131
            [this](PSerialDevice device) { OnDeviceConnectionStateChanged(device); });
330✔
132
    }
133
}
95✔
134

135
void TSerialClientRegisterPoller::ScheduleNextPoll(PPollableDevice device)
741✔
136
{
137
    if (device->HasRegisters()) {
741✔
138
        Scheduler.AddEntry(device, device->GetDeadline(), device->GetPriority());
741✔
139
    }
140
}
741✔
141

142
void TSerialClientRegisterPoller::ClosedPortCycle(steady_clock::time_point currentTime, TRegisterCallback callback)
4✔
143
{
144
    Scheduler.ResetLoadBalancing();
4✔
145

146
    RescheduleDisconnectedDevices();
4✔
147
    RescheduleDevicesWithSpendedPoll(currentTime);
4✔
148

149
    std::unique_lock lock(Mutex);
8✔
150

151
    TClosedPortDeviceReader reader(currentTime);
8✔
152
    do {
4✔
153
        reader.ClearRegisters();
8✔
154
        Scheduler.AccumulateNext(currentTime, reader, TItemSelectionPolicy::All);
8✔
155
        for (auto& reg: reader.GetRegisters()) {
16✔
156
            reg->SetError(TRegister::TError::ReadError);
8✔
157
            if (callback) {
8✔
158
                callback(reg);
8✔
159
            }
160
            auto device = reader.GetDevice()->GetDevice();
16✔
161
            device->SetTransferResult(false);
8✔
162
        }
163
        if (reader.GetDevice()) {
8✔
164
            ScheduleNextPoll(reader.GetDevice());
4✔
165
        }
166
    } while (!reader.GetRegisters().empty());
8✔
167
}
4✔
168

169
std::chrono::steady_clock::time_point TSerialClientRegisterPoller::GetDeadline(
752✔
170
    bool lowPriorityRateLimitIsExceeded,
171
    const util::TSpentTimeMeter& spentTime) const
172
{
173
    if (Scheduler.IsEmpty()) {
752✔
174
        return spentTime.GetStartTime() + 1s;
15✔
175
    }
176
    if (lowPriorityRateLimitIsExceeded) {
737✔
177
        auto lowPriorityDeadline = Scheduler.GetLowPriorityDeadline();
×
178
        // There are some low priority items
179
        if (lowPriorityDeadline != std::chrono::steady_clock::time_point::max()) {
×
180
            lowPriorityDeadline = std::max(lowPriorityDeadline, LowPriorityRateLimiter.GetStartTime() + 1s);
×
181
        }
182
        return std::min(Scheduler.GetHighPriorityDeadline(), lowPriorityDeadline);
×
183
    }
184
    return Scheduler.GetDeadline();
737✔
185
}
186

187
TPollResult TSerialClientRegisterPoller::OpenPortCycle(TPort& port,
789✔
188
                                                       const util::TSpentTimeMeter& spentTime,
189
                                                       std::chrono::milliseconds maxPollingTime,
190
                                                       bool readAtLeastOneRegister,
191
                                                       TSerialClientDeviceAccessHandler& lastAccessedDevice,
192
                                                       TRegisterCallback callback)
193
{
194
    RescheduleDisconnectedDevices();
789✔
195
    RescheduleDevicesWithSpendedPoll(spentTime.GetStartTime());
789✔
196

197
    std::unique_lock lock(Mutex);
1,578✔
198

199
    TPollResult res;
789✔
200

201
    TDeviceReader reader(port, spentTime, maxPollingTime, readAtLeastOneRegister, lastAccessedDevice);
1,578✔
202

203
    bool lowPriorityRateLimitIsExceeded = LowPriorityRateLimiter.IsOverLimit(spentTime.GetStartTime());
789✔
204

205
    Scheduler.AccumulateNext(spentTime.GetStartTime(),
789✔
206
                             reader,
207
                             lowPriorityRateLimitIsExceeded ? TItemSelectionPolicy::OnlyHighPriority
208
                                                            : TItemSelectionPolicy::All);
209
    if (lowPriorityRateLimitIsExceeded) {
789✔
210
        auto throttlingMsg = ThrottlingStateLogger.GetMessage();
×
211
        if (!throttlingMsg.empty()) {
×
212
            LOG(Warn) << port.GetDescription() << " " << throttlingMsg;
×
213
        }
214
    }
215
    auto range = reader.GetRegisterRange();
1,578✔
216

217
    if (!range) {
789✔
218
        // Nothing to read
219
        res.Deadline = GetDeadline(lowPriorityRateLimitIsExceeded, spentTime);
15✔
220
        return res;
15✔
221
    }
222

223
    // There are registers waiting read, but they don't fit in allowed poll limit
224
    if (range->RegisterList().empty()) {
774✔
225
        res.NotEnoughTime = true;
37✔
226
        if (reader.GetDevice()->GetPriority() == TPriority::High) {
37✔
227
            // High priority registers are limited by maxPollingTime
228
            res.Deadline = spentTime.GetStartTime() + maxPollingTime;
1✔
229
        } else {
230
            // Low priority registers are limited by high priority and maxPollingTime
231
            res.Deadline = std::min(Scheduler.GetHighPriorityDeadline(), spentTime.GetStartTime() + maxPollingTime);
36✔
232
        }
233
        return res;
37✔
234
    }
235

236
    res.Device = reader.GetDevice()->GetDevice();
737✔
237

238
    for (auto& reg: range->RegisterList()) {
1,903✔
239
        if (callback) {
1,166✔
240
            callback(reg);
1,166✔
241
        }
242
        if (reader.GetDevice()->GetPriority() == TPriority::Low) {
1,166✔
243
            LowPriorityRateLimiter.NewItem(spentTime.GetStartTime());
1,124✔
244
        }
245
    }
246
    ScheduleNextPoll(reader.GetDevice());
737✔
247

248
    Scheduler.UpdateSelectionTime(ceil<milliseconds>(spentTime.GetSpentTime()), reader.GetDevice()->GetPriority());
737✔
249
    res.Deadline = GetDeadline(lowPriorityRateLimitIsExceeded, spentTime);
737✔
250
    return res;
737✔
251
}
252

253
void TSerialClientRegisterPoller::SuspendPoll(PSerialDevice device, std::chrono::steady_clock::time_point currentTime)
3✔
254
{
255
    std::unique_lock lock(Mutex);
6✔
256

257
    if (DevicesWithSpendedPoll.find(device) == DevicesWithSpendedPoll.end()) {
3✔
258
        auto range = Devices.equal_range(device);
3✔
259
        if (range.first == range.second) {
3✔
NEW
260
            throw std::runtime_error("Device " + device->ToString() + " is not included to poll");
×
261
        }
262
        for (auto it = range.first; it != range.second; ++it) {
6✔
263
            Scheduler.Remove(it->second);
3✔
264
        }
265
    }
266

267
    DevicesWithSpendedPoll[device] = currentTime + SUSPEND_POLL_TIMEOUT;
3✔
268
    LOG(Info) << "Device " << device->ToString() << " poll suspended";
3✔
269

270
    if (device->GetConnectionState() != TDeviceConnectionState::DISCONNECTED) {
3✔
271
        device->SetDisconnected();
3✔
272
    }
273
}
3✔
274

275
void TSerialClientRegisterPoller::ResumePoll(PSerialDevice device)
3✔
276
{
277
    std::unique_lock lock(Mutex);
3✔
278

279
    if (DevicesWithSpendedPoll.find(device) == DevicesWithSpendedPoll.end()) {
3✔
NEW
280
        throw std::runtime_error("Device " + device->ToString() + " poll is not suspended");
×
281
    }
282

283
    auto range = Devices.equal_range(device);
3✔
284
    for (auto it = range.first; it != range.second; ++it) {
6✔
285
        it->second->RescheduleAllRegisters();
3✔
286
        Scheduler.AddEntry(it->second, it->second->GetDeadline(), it->second->GetPriority());
3✔
287
    }
288

289
    DevicesWithSpendedPoll.erase(device);
3✔
290
    LOG(Info) << "Device " << device->ToString() << " poll resumed";
3✔
291
}
3✔
292

293
void TSerialClientRegisterPoller::OnDeviceConnectionStateChanged(PSerialDevice device)
130✔
294
{
295
    if (device->GetConnectionState() == TDeviceConnectionState::DISCONNECTED &&
150✔
296
        DevicesWithSpendedPoll.find(device) == DevicesWithSpendedPoll.end())
150✔
297
    {
298
        DisconnectedDevicesWaitingForReschedule.push_back(device);
17✔
299
    }
300
}
130✔
301

302
bool TPollableDeviceComparePredicate::operator()(const PPollableDevice& d1, const PPollableDevice& d2) const
14✔
303
{
304
    if (d1->GetDevice() != d2->GetDevice()) {
14✔
305
        return d1->GetDevice()->DeviceConfig()->SlaveId > d2->GetDevice()->DeviceConfig()->SlaveId;
14✔
306
    }
307
    return false;
×
308
}
309

310
TThrottlingStateLogger::TThrottlingStateLogger(): FirstTime(true)
95✔
311
{}
95✔
312

313
std::string TThrottlingStateLogger::GetMessage()
×
314
{
315
    if (FirstTime) {
×
316
        FirstTime = false;
×
317
        return "Register read rate limit is exceeded";
×
318
    }
319
    return std::string();
×
320
}
321

322
void TSerialClientRegisterPoller::RescheduleDisconnectedDevices()
793✔
323
{
324
    std::unique_lock lock(Mutex);
1,586✔
325
    for (auto& device: DisconnectedDevicesWaitingForReschedule) {
810✔
326
        auto range = Devices.equal_range(device);
17✔
327
        for (auto it = range.first; it != range.second; ++it) {
35✔
328
            Scheduler.Remove(it->second);
18✔
329
            it->second->RescheduleAllRegisters();
18✔
330
            Scheduler.AddEntry(it->second, it->second->GetDeadline(), it->second->GetPriority());
18✔
331
        }
332
    }
333
    DisconnectedDevicesWaitingForReschedule.clear();
793✔
334
}
793✔
335

336
void TSerialClientRegisterPoller::RescheduleDevicesWithSpendedPoll(std::chrono::steady_clock::time_point currentTime)
793✔
337
{
338
    std::list<PSerialDevice> list;
1,586✔
339
    {
340
        std::unique_lock lock(Mutex);
1,586✔
341
        for (auto it = DevicesWithSpendedPoll.begin(); it != DevicesWithSpendedPoll.end(); ++it) {
809✔
342
            if (it->second <= currentTime) {
16✔
343
                list.push_back(it->first);
1✔
344
            }
345
        }
346
    }
347
    for (auto device: list) {
794✔
348
        ResumePoll(device);
1✔
349
    }
350
}
793✔
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