#!/usr/bin/env perl use v5.32.1; use strict; use warnings; BEGIN { # Net::MQTT::Simple intentionally requires this when using username/password # over plain MQTT. The broker is expected to be reachable only over the VPN/LAN. $ENV{MQTT_SIMPLE_ALLOW_INSECURE_LOGIN} //= 1; } use DBI; use JSON::PP qw(decode_json); use Net::MQTT::Simple; use POSIX qw(strftime); use Scalar::Util qw(looks_like_number); $| = 1; my $db_path = env( 'FAPG_DAQ_DB', '/var/lib/fapg-daq/readings.sqlite3' ); my $mqtt_host = required_env('MQTT_HOST'); my $mqtt_port = env( 'MQTT_PORT', '1883' ); my $mqtt_username = env( 'MQTT_USERNAME', 'fapg_vps' ); my $mqtt_password = required_env('MQTT_PASSWORD'); my $mqtt_topic = env( 'MQTT_TOPIC', 'fapg/daq/#' ); my $dbh = connect_db($db_path); init_schema($dbh); my $mqtt = Net::MQTT::Simple->new("$mqtt_host:$mqtt_port"); $mqtt->login( $mqtt_username, $mqtt_password ); $SIG{INT} = $SIG{TERM} = sub { warn "Shutting down fapg-daq-mqtt-sqlite\n"; eval { $mqtt->disconnect; }; eval { $dbh->disconnect; }; exit 0; }; warn "Subscribing to mqtt://$mqtt_host:$mqtt_port/$mqtt_topic\n"; warn "Writing MQTT messages to $db_path\n"; $mqtt->run( $mqtt_topic => sub { my ( $topic, $message ) = @_; eval { store_message( $dbh, $topic, $message ); 1; } or do { my $error = $@ || 'unknown error'; chomp $error; warn "Failed to store MQTT message from $topic: $error\n"; }; }, ); sub env { my ( $name, $default ) = @_; return exists $ENV{$name} ? $ENV{$name} : $default; } sub required_env { my ($name) = @_; die "Missing required environment variable: $name\n" if !exists $ENV{$name} || $ENV{$name} eq ''; return $ENV{$name}; } sub connect_db { my ($path) = @_; my $dbh = DBI->connect( "dbi:SQLite:dbname=$path", '', '', { RaiseError => 1, PrintError => 0, AutoCommit => 1, sqlite_unicode => 1, }, ); $dbh->do('PRAGMA journal_mode = WAL'); $dbh->do('PRAGMA synchronous = NORMAL'); $dbh->do('PRAGMA busy_timeout = 5000'); $dbh->do('PRAGMA foreign_keys = ON'); return $dbh; } sub init_schema { my ($dbh) = @_; $dbh->do( q{ CREATE TABLE IF NOT EXISTS readings ( id INTEGER PRIMARY KEY AUTOINCREMENT, received_at TEXT NOT NULL, topic TEXT NOT NULL, kind TEXT, schema TEXT, timestamp TEXT, probe TEXT, node TEXT, value REAL, unit TEXT, source TEXT, payload TEXT NOT NULL, valid INTEGER NOT NULL DEFAULT 1, error TEXT ) } ); $dbh->do( q{ CREATE INDEX IF NOT EXISTS readings_timestamp_idx ON readings(timestamp) } ); $dbh->do( q{ CREATE INDEX IF NOT EXISTS readings_probe_node_timestamp_idx ON readings(probe, node, timestamp) } ); $dbh->do( q{ CREATE INDEX IF NOT EXISTS readings_probe_observed_idx ON readings(probe, COALESCE(timestamp, received_at), id) } ); $dbh->do( q{ CREATE INDEX IF NOT EXISTS readings_topic_idx ON readings(topic) } ); init_reading_rollups($dbh); $dbh->do( q{ CREATE TABLE IF NOT EXISTS node_status ( id INTEGER PRIMARY KEY AUTOINCREMENT, received_at TEXT NOT NULL, topic TEXT NOT NULL, schema TEXT, timestamp TEXT, probe TEXT, node TEXT, status TEXT, message TEXT, error TEXT, device TEXT, source TEXT, payload TEXT NOT NULL, valid INTEGER NOT NULL DEFAULT 1 ) } ); $dbh->do( q{ CREATE INDEX IF NOT EXISTS node_status_probe_node_timestamp_idx ON node_status(probe, node, timestamp) } ); $dbh->do( q{ CREATE INDEX IF NOT EXISTS node_status_valid_probe_observed_idx ON node_status(valid, probe, COALESCE(timestamp, received_at), id) } ); $dbh->do( q{ CREATE INDEX IF NOT EXISTS node_status_topic_idx ON node_status(topic) } ); } sub init_reading_rollups { my ($dbh) = @_; $dbh->do( q{ CREATE TABLE IF NOT EXISTS reading_rollups ( probe TEXT NOT NULL, bucket_seconds INTEGER NOT NULL, bucket_epoch INTEGER NOT NULL, sample_count INTEGER NOT NULL, value_total REAL NOT NULL, value_min REAL NOT NULL, value_max REAL NOT NULL, first_sample_epoch INTEGER NOT NULL, last_sample_epoch INTEGER NOT NULL, max_gap_seconds INTEGER NOT NULL DEFAULT 0, PRIMARY KEY (probe, bucket_seconds, bucket_epoch) ) WITHOUT ROWID } ); for my $bucket_seconds ( 60, 3_600 ) { $dbh->do( qq{ WITH samples AS ( SELECT id, probe, CAST( CAST(strftime('%s', COALESCE(timestamp, received_at)) AS INTEGER) / $bucket_seconds AS INTEGER ) * $bucket_seconds AS bucket_epoch, CAST(strftime('%s', COALESCE(timestamp, received_at)) AS INTEGER) AS sample_epoch, CAST(value AS REAL) AS value FROM readings WHERE valid = 1 AND probe IS NOT NULL AND value IS NOT NULL AND COALESCE(timestamp, received_at) IS NOT NULL ), ordered AS ( SELECT *, sample_epoch - LAG(sample_epoch) OVER ( PARTITION BY probe, bucket_epoch ORDER BY sample_epoch, id ) AS gap_seconds FROM samples ) INSERT INTO reading_rollups ( probe, bucket_seconds, bucket_epoch, sample_count, value_total, value_min, value_max, first_sample_epoch, last_sample_epoch, max_gap_seconds ) SELECT probe, $bucket_seconds, bucket_epoch, COUNT(*), SUM(value), MIN(value), MAX(value), MIN(sample_epoch), MAX(sample_epoch), MAX(COALESCE(gap_seconds, 0)) FROM ordered WHERE NOT EXISTS ( SELECT 1 FROM reading_rollups WHERE bucket_seconds = $bucket_seconds LIMIT 1 ) GROUP BY probe, bucket_epoch } ); } $dbh->do( q{ CREATE TRIGGER IF NOT EXISTS readings_rollup_insert AFTER INSERT ON readings WHEN NEW.valid = 1 AND NEW.probe IS NOT NULL AND NEW.value IS NOT NULL AND COALESCE(NEW.timestamp, NEW.received_at) IS NOT NULL BEGIN INSERT INTO reading_rollups ( probe, bucket_seconds, bucket_epoch, sample_count, value_total, value_min, value_max, first_sample_epoch, last_sample_epoch, max_gap_seconds ) VALUES ( NEW.probe, 60, CAST( CAST(strftime('%s', COALESCE(NEW.timestamp, NEW.received_at)) AS INTEGER) / 60 AS INTEGER ) * 60, 1, CAST(NEW.value AS REAL), CAST(NEW.value AS REAL), CAST(NEW.value AS REAL), CAST(strftime('%s', COALESCE(NEW.timestamp, NEW.received_at)) AS INTEGER), CAST(strftime('%s', COALESCE(NEW.timestamp, NEW.received_at)) AS INTEGER), 0 ) ON CONFLICT (probe, bucket_seconds, bucket_epoch) DO UPDATE SET sample_count = sample_count + 1, value_total = value_total + excluded.value_total, value_min = MIN(value_min, excluded.value_min), value_max = MAX(value_max, excluded.value_max), max_gap_seconds = MAX( max_gap_seconds, CASE WHEN excluded.first_sample_epoch > last_sample_epoch THEN excluded.first_sample_epoch - last_sample_epoch ELSE 0 END ), first_sample_epoch = MIN(first_sample_epoch, excluded.first_sample_epoch), last_sample_epoch = MAX(last_sample_epoch, excluded.last_sample_epoch); INSERT INTO reading_rollups ( probe, bucket_seconds, bucket_epoch, sample_count, value_total, value_min, value_max, first_sample_epoch, last_sample_epoch, max_gap_seconds ) VALUES ( NEW.probe, 3600, CAST( CAST(strftime('%s', COALESCE(NEW.timestamp, NEW.received_at)) AS INTEGER) / 3600 AS INTEGER ) * 3600, 1, CAST(NEW.value AS REAL), CAST(NEW.value AS REAL), CAST(NEW.value AS REAL), CAST(strftime('%s', COALESCE(NEW.timestamp, NEW.received_at)) AS INTEGER), CAST(strftime('%s', COALESCE(NEW.timestamp, NEW.received_at)) AS INTEGER), 0 ) ON CONFLICT (probe, bucket_seconds, bucket_epoch) DO UPDATE SET sample_count = sample_count + 1, value_total = value_total + excluded.value_total, value_min = MIN(value_min, excluded.value_min), value_max = MAX(value_max, excluded.value_max), max_gap_seconds = MAX( max_gap_seconds, CASE WHEN excluded.first_sample_epoch > last_sample_epoch THEN excluded.first_sample_epoch - last_sample_epoch ELSE 0 END ), first_sample_epoch = MIN(first_sample_epoch, excluded.first_sample_epoch), last_sample_epoch = MAX(last_sample_epoch, excluded.last_sample_epoch); END } ); } sub store_message { my ( $dbh, $topic, $message ) = @_; my ( $topic_probe, $topic_node, $topic_kind ) = $topic =~ m{\Afapg/daq/([^/]+)/([^/]+)/([^/]+)\z}; die "invalid MQTT topic: $topic\n" if !defined $topic_kind || $topic_kind !~ /\A(?:reading|status)\z/; my $valid = 1; my $error; my $data = eval { decode_json($message) }; if ($@) { $valid = 0; $error = $@; chomp $error; $data = {}; } elsif ( ref $data ne 'HASH' ) { $valid = 0; $error = 'JSON payload is not an object'; $data = {}; } if ( $topic_kind eq 'status' ) { store_status_message( $dbh, $topic, $message, $topic_probe, $topic_node, $data, $valid, $error ); return; } my $value = $data->{value}; $value = undef if defined $value && !looks_like_number($value); my $sth = $dbh->prepare_cached( q{ INSERT INTO readings ( received_at, topic, kind, schema, timestamp, probe, node, value, unit, source, payload, valid, error ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) } ); $sth->execute( utc_now(), $topic, $topic_kind, $data->{schema}, $data->{timestamp}, $data->{probe} // $topic_probe, $data->{node} // $topic_node, $value, $data->{unit}, $data->{source}, $message, $valid, $error, ); } sub store_status_message { my ( $dbh, $topic, $message, $topic_probe, $topic_node, $data, $valid, $error ) = @_; my $status = $data->{status}; if ( defined $status ) { $status = lc $status; if ( $status !~ /\A(?:ok|probe_error)\z/ ) { $valid = 0; $error = defined $error ? "$error; unsupported status: $status" : "unsupported status: $status"; } } else { $valid = 0; $error = defined $error ? "$error; missing status" : 'missing status'; } my $sth = $dbh->prepare_cached( q{ INSERT INTO node_status ( received_at, topic, schema, timestamp, probe, node, status, message, error, device, source, payload, valid ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) } ); $sth->execute( utc_now(), $topic, $data->{schema}, $data->{timestamp}, $data->{probe} // $topic_probe, $data->{node} // $topic_node, $status, $data->{message}, $data->{error} // $error, $data->{device}, $data->{source}, $message, $valid, ); } sub utc_now { return strftime( '%Y-%m-%dT%H:%M:%SZ', gmtime ); }