Datalogger: A Semi-Realistic Project (Events Version)#
Polling-Thread To Consume Both Message Queues#
Two message queues: sensor data, snapshot requests
Collapse both threads into one
Rip out the two thread while-loop bodies
Create respective lambdas/callback for event notification
Important
No need to protect CSV writing against snapshot rotation ⟶ remove lock
#include <why.h>
#include <why-sensor.h>
#include <why-irq.h>
#include <why-messagequeue.h>
#include <why-time.h>
#include <why-watchdog.h>
#include <why-thread.h>
#include <why-file.h>
#include <why-mqtt.h>
#include <why-mutex.h>
#include <why-eventloop.h>
#include <string>
#include <format>
struct sensor_data
{
Why::TimeSpec timestamp;
std::string_view sensorname;
uint64_t value;
};
int main()
{
Why::init();
Why::RandomSensor s1(0, 100);
Why::RandomSensor s2(100, 200);
auto sensor_data_queue = Why::MessageQueue<sensor_data>::create();
assert(sensor_data_queue);
auto snapshot_queue = Why::MessageQueue<Why::TimeSpec/*timestamp*/>::create();
assert(snapshot_queue);
auto csv_file = Why::File::open("data.csv", Why::File::WriteOnly|Why::File::Append|Why::File::Create);
assert(csv_file);
auto mqtt = Why::MQTTPublisher::create("why-topic", "127.0.0.1");
assert(mqtt);
auto poller = Why::Thread::create(
[&sensor_data_queue, &snapshot_queue, &csv_file, &mqtt](){
Why::Eventloop loop;
// incoming data
loop.watch(
*sensor_data_queue,
[&sensor_data_queue, &csv_file, &mqtt](){
auto sample = sensor_data_queue->get(); // <-- does not block
assert(sample);
auto [timestamp, sensorname, value] = *sample;
// write to CSV
std::string line = std::format("{};{};{}\n", timestamp.to_seconds(), sensorname, value);
auto written = csv_file->write((const uint8_t*)line.c_str(), line.size());
assert(written);
// publish MQTT
std::string json = std::format("{{\"timestamp\": {}, \"sensorname\": {}, \"value\": {} }}",
timestamp.to_seconds(), sensorname, value);
auto ok = mqtt->publish(json);
assert(ok);
});
// snapshot button
loop.watch(
*snapshot_queue,
[&snapshot_queue, &csv_file](){
auto timestamp = snapshot_queue->get(); // <-- does not block
assert(timestamp);
std::string snapshot_filename = std::format("snapshot-{}.csv", timestamp->to_seconds());
Why::File::rename("data.csv", snapshot_filename);
csv_file = Why::File::open("data.csv", Why::File::WriteOnly|Why::File::Append|Why::File::Create);
});
loop(); // <-- here we block
});
assert(poller);
// activate snapshot button
Why::IRQ::connect(666, [&snapshot_queue](int /*irqnum*/){
auto ok = snapshot_queue->put_noblock(Why::now_monotonic());
assert(ok);
});
// activate sensor interrupts
{
auto isr = [&sensor_data_queue, &s1, &s2](int irqnum){
sensor_data data;
if (irqnum == 7)
data = {.timestamp=Why::now_monotonic(), .sensorname="s1", .value=s1.get_value()};
else if (irqnum == 42)
data = {.timestamp=Why::now_monotonic(), .sensorname="s2", .value=s2.get_value()};
else
assert(!"unexpected irq");
auto ok = sensor_data_queue->put_noblock(data);
assert(ok);
};
Why::IRQ::connect(7, isr);
Why::IRQ::connect(42, isr);
}
// activate and feed watchdog
const Why::TimeSpec watchdog_timeout(2, 0);
Why::Watchdog::activate(watchdog_timeout);
while (true) {
Why::Watchdog::feed();
Why::sleep(watchdog_timeout - Why::TimeSpec(0, 500'000'000));
}
return 0;
}
Watchdog Timer Into Polling-Thread (Semaphore Version)#
Currently, we do a (blocking) sleep in main loop
Plan:
Start a timer
In the timer callback function - in ISR context -, give a semaphore
In event loop, watch semaphore for takeability
Feed watchdog in event callback
#include <why.h>
#include <why-sensor.h>
#include <why-irq.h>
#include <why-messagequeue.h>
#include <why-semaphore.h>
#include <why-time.h>
#include <why-watchdog.h>
#include <why-thread.h>
#include <why-file.h>
#include <why-mqtt.h>
#include <why-mutex.h>
#include <why-eventloop.h>
#include <why-timer.h>
#include <string>
#include <format>
struct sensor_data
{
Why::TimeSpec timestamp;
std::string_view sensorname;
uint64_t value;
};
int main()
{
Why::init();
Why::RandomSensor s1(0, 100);
Why::RandomSensor s2(100, 200);
auto sensor_data_queue = Why::MessageQueue<sensor_data>::create();
assert(sensor_data_queue);
auto snapshot_queue = Why::MessageQueue<Why::TimeSpec/*timestamp*/>::create();
assert(snapshot_queue);
auto watchdog_notify_sem = Why::Semaphore::create(0);
auto csv_file = Why::File::open("data.csv", Why::File::WriteOnly|Why::File::Append|Why::File::Create);
assert(csv_file);
auto mqtt = Why::MQTTPublisher::create("why-topic", "127.0.0.1");
assert(mqtt);
auto poller = Why::Thread::create(
[&sensor_data_queue, &snapshot_queue, &watchdog_notify_sem, &csv_file, &mqtt](){
Why::Eventloop loop;
// incoming data
loop.watch(
*sensor_data_queue,
[&sensor_data_queue, &csv_file, &mqtt](){
auto sample = sensor_data_queue->get();
assert(sample);
auto [timestamp, sensorname, value] = *sample;
// write to CSV
std::string line = std::format("{};{};{}\n", timestamp.to_seconds(), sensorname, value);
auto written = csv_file->write((const uint8_t*)line.c_str(), line.size());
assert(written);
// publish MQTT
std::string json = std::format("{{\"timestamp\": {}, \"sensorname\": {}, \"value\": {} }}",
timestamp.to_seconds(), sensorname, value);
auto ok = mqtt->publish(json);
assert(ok);
});
// snapshot button
loop.watch(
*snapshot_queue,
[&snapshot_queue, &csv_file](){
auto timestamp = snapshot_queue->get();
assert(timestamp);
std::string snapshot_filename = std::format("snapshot-{}.csv", timestamp->to_seconds());
Why::File::rename("data.csv", snapshot_filename);
csv_file = Why::File::open("data.csv", Why::File::WriteOnly|Why::File::Append|Why::File::Create);
});
// watchdog feeding
loop.watch(
*watchdog_notify_sem,
[&watchdog_notify_sem](){
watchdog_notify_sem->take(); // <-- does not block
Why::Watchdog::feed();
});
loop();
});
assert(poller);
// activate snapshot button
Why::IRQ::connect(666, [&snapshot_queue](int /*irqnum*/){
auto ok = snapshot_queue->put_noblock(Why::now_monotonic());
assert(ok);
});
// activate sensor interrupts
{
auto isr = [&sensor_data_queue, &s1, &s2](int irqnum){
sensor_data data;
if (irqnum == 7)
data = {.timestamp=Why::now_monotonic(), .sensorname="s1", .value=s1.get_value()};
else if (irqnum == 42)
data = {.timestamp=Why::now_monotonic(), .sensorname="s2", .value=s2.get_value()};
else
assert(!"unexpected irq");
auto ok = sensor_data_queue->put_noblock(data);
assert(ok);
};
Why::IRQ::connect(7, isr);
Why::IRQ::connect(42, isr);
}
// activate watchdog, and start feed timer
const Why::TimeSpec watchdog_timeout(2, 0);
Why::Watchdog::activate(watchdog_timeout);
Why::Timer watchdog_timer([&watchdog_notify_sem](){watchdog_notify_sem->give();});
watchdog_timer.start(
Why::TimeSpec(1,0), // initial expiration in a second
watchdog_timeout - Why::TimeSpec(0, 500'000'000) // period
);
Why::pause(); // <-- main thread is useless
return 0;
}
Watchdog Timer Into Polling-Thread (Poll-Signal Version)#
Semaphore is heavyweight: can be waited for which is extra logic
Poll signal is rather lightweight
Specific to event driven applications
#include <why.h>
#include <why-sensor.h>
#include <why-irq.h>
#include <why-messagequeue.h>
#include <why-pollsignal.h>
#include <why-time.h>
#include <why-watchdog.h>
#include <why-thread.h>
#include <why-file.h>
#include <why-mqtt.h>
#include <why-mutex.h>
#include <why-eventloop.h>
#include <why-timer.h>
#include <string>
#include <format>
struct sensor_data
{
Why::TimeSpec timestamp;
std::string_view sensorname;
uint64_t value;
};
int main()
{
Why::init();
Why::RandomSensor s1(0, 100);
Why::RandomSensor s2(100, 200);
auto sensor_data_queue = Why::MessageQueue<sensor_data>::create();
assert(sensor_data_queue);
auto snapshot_queue = Why::MessageQueue<Why::TimeSpec/*timestamp*/>::create();
assert(snapshot_queue);
auto watchdog_pollsignal = Why::PollSignal::create();
assert(watchdog_pollsignal);
auto csv_file = Why::File::open("data.csv", Why::File::WriteOnly|Why::File::Append|Why::File::Create);
assert(csv_file);
auto mqtt = Why::MQTTPublisher::create("why-topic", "127.0.0.1");
assert(mqtt);
auto poller = Why::Thread::create(
[&sensor_data_queue, &snapshot_queue, &watchdog_pollsignal, &csv_file, &mqtt](){
Why::Eventloop loop;
// incoming data
loop.watch(
*sensor_data_queue,
[&sensor_data_queue, &csv_file, &mqtt](){
auto sample = sensor_data_queue->get();
assert(sample);
auto [timestamp, sensorname, value] = *sample;
// write to CSV
std::string line = std::format("{};{};{}\n", timestamp.to_seconds(), sensorname, value);
auto written = csv_file->write((const uint8_t*)line.c_str(), line.size());
assert(written);
// publish MQTT
std::string json = std::format("{{\"timestamp\": {}, \"sensorname\": {}, \"value\": {} }}",
timestamp.to_seconds(), sensorname, value);
auto ok = mqtt->publish(json);
assert(ok);
});
// snapshot button
loop.watch(
*snapshot_queue,
[&snapshot_queue, &csv_file](){
auto timestamp = snapshot_queue->get();
assert(timestamp);
std::string snapshot_filename = std::format("snapshot-{}.csv", timestamp->to_seconds());
Why::File::rename("data.csv", snapshot_filename);
csv_file = Why::File::open("data.csv", Why::File::WriteOnly|Why::File::Append|Why::File::Create);
});
// watchdog feeding
loop.watch(
*watchdog_pollsignal,
[&watchdog_pollsignal](){
watchdog_pollsignal->reset();
Why::Watchdog::feed();
});
loop();
});
assert(poller);
// activate snapshot button
Why::IRQ::connect(666, [&snapshot_queue](int /*irqnum*/){
auto ok = snapshot_queue->put_noblock(Why::now_monotonic());
assert(ok);
});
// activate sensor interrupts
{
auto isr = [&sensor_data_queue, &s1, &s2](int irqnum){
sensor_data data;
if (irqnum == 7)
data = {.timestamp=Why::now_monotonic(), .sensorname="s1", .value=s1.get_value()};
else if (irqnum == 42)
data = {.timestamp=Why::now_monotonic(), .sensorname="s2", .value=s2.get_value()};
else
assert(!"unexpected irq");
auto ok = sensor_data_queue->put_noblock(data);
assert(ok);
};
Why::IRQ::connect(7, isr);
Why::IRQ::connect(42, isr);
}
// activate watchdog, and start feed timer
const Why::TimeSpec watchdog_timeout(2, 0);
Why::Watchdog::activate(watchdog_timeout);
Why::Timer watchdog_timer([&watchdog_pollsignal](){watchdog_pollsignal->raise(42);});
watchdog_timer.start(
Why::TimeSpec(1,0), // initial expiration in a second
watchdog_timeout - Why::TimeSpec(0, 500'000'000) // period
);
Why::pause(); // <-- main thread is useless
return 0;
}
And Single-Threaded Programs?#
Main thread is useless now
All logic in the poller thread
Why not eliminate that thread altogether?
#include <why.h>
#include <why-sensor.h>
#include <why-irq.h>
#include <why-messagequeue.h>
#include <why-pollsignal.h>
#include <why-time.h>
#include <why-watchdog.h>
#include <why-thread.h>
#include <why-file.h>
#include <why-mqtt.h>
#include <why-mutex.h>
#include <why-eventloop.h>
#include <why-timer.h>
#include <string>
#include <format>
struct sensor_data
{
Why::TimeSpec timestamp;
std::string_view sensorname;
uint64_t value;
};
int main()
{
Why::init();
Why::RandomSensor s1(0, 100);
Why::RandomSensor s2(100, 200);
auto sensor_data_queue = Why::MessageQueue<sensor_data>::create();
assert(sensor_data_queue);
auto snapshot_queue = Why::MessageQueue<Why::TimeSpec/*timestamp*/>::create();
assert(snapshot_queue);
auto watchdog_pollsignal = Why::PollSignal::create();
assert(watchdog_pollsignal);
auto csv_file = Why::File::open("data.csv", Why::File::WriteOnly|Why::File::Append|Why::File::Create);
assert(csv_file);
auto mqtt = Why::MQTTPublisher::create("why-topic", "127.0.0.1");
assert(mqtt);
// activate snapshot button
Why::IRQ::connect(666, [&snapshot_queue](int /*irqnum*/){
auto ok = snapshot_queue->put_noblock(Why::now_monotonic());
assert(ok);
});
// activate sensor interrupts
{
auto isr = [&sensor_data_queue, &s1, &s2](int irqnum){
sensor_data data;
if (irqnum == 7)
data = {.timestamp=Why::now_monotonic(), .sensorname="s1", .value=s1.get_value()};
else if (irqnum == 42)
data = {.timestamp=Why::now_monotonic(), .sensorname="s2", .value=s2.get_value()};
else
assert(!"unexpected irq");
auto ok = sensor_data_queue->put_noblock(data);
assert(ok);
};
Why::IRQ::connect(7, isr);
Why::IRQ::connect(42, isr);
}
// activate watchdog, and start feed timer
const Why::TimeSpec watchdog_timeout(2, 0);
Why::Watchdog::activate(watchdog_timeout);
Why::Timer watchdog_timer([&watchdog_pollsignal](){watchdog_pollsignal->raise(42);});
watchdog_timer.start(
Why::TimeSpec(1,0), // initial expiration in a second
watchdog_timeout - Why::TimeSpec(0, 500'000'000) // period
);
// start main work: watch out for events
Why::Eventloop loop;
// incoming data
loop.watch(
*sensor_data_queue,
[&sensor_data_queue, &csv_file, &mqtt](){
auto sample = sensor_data_queue->get();
assert(sample);
auto [timestamp, sensorname, value] = *sample;
// write to CSV
std::string line = std::format("{};{};{}\n", timestamp.to_seconds(), sensorname, value);
auto written = csv_file->write((const uint8_t*)line.c_str(), line.size());
assert(written);
// publish MQTT
std::string json = std::format("{{\"timestamp\": {}, \"sensorname\": {}, \"value\": {} }}",
timestamp.to_seconds(), sensorname, value);
auto ok = mqtt->publish(json);
assert(ok);
});
// snapshot button
loop.watch(
*snapshot_queue,
[&snapshot_queue, &csv_file](){
auto timestamp = snapshot_queue->get();
assert(timestamp);
std::string snapshot_filename = std::format("snapshot-{}.csv", timestamp->to_seconds());
Why::File::rename("data.csv", snapshot_filename);
csv_file = Why::File::open("data.csv", Why::File::WriteOnly|Why::File::Append|Why::File::Create);
});
// watchdog feeding
loop.watch(
*watchdog_pollsignal,
[&watchdog_pollsignal](){
watchdog_pollsignal->reset();
Why::Watchdog::feed();
});
loop(); // <-- innocent-looking but cool
return 0;
}