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;
}