View raw

1 #!/usr/bin/env perl 2 3 use v5.32.1; 4 use strict; 5 use warnings; 6 7 BEGIN { 8 # Net::MQTT::Simple intentionally requires this when using username/password 9 # over plain MQTT. The broker is expected to be reachable only over the VPN/LAN. 10 $ENV{MQTT_SIMPLE_ALLOW_INSECURE_LOGIN} //= 1; 11 } 12 13 use DBI; 14 use JSON::PP qw(decode_json); 15 use Net::MQTT::Simple; 16 use POSIX qw(strftime); 17 use Scalar::Util qw(looks_like_number); 18 19 $| = 1; 20 21 my $db_path = env( 'FAPG_DAQ_DB', '/var/lib/fapg-daq/readings.sqlite3' ); 22 my $mqtt_host = required_env('MQTT_HOST'); 23 my $mqtt_port = env( 'MQTT_PORT', '1883' ); 24 my $mqtt_username = env( 'MQTT_USERNAME', 'fapg_vps' ); 25 my $mqtt_password = required_env('MQTT_PASSWORD'); 26 my $mqtt_topic = env( 'MQTT_TOPIC', 'fapg/daq/#' ); 27 28 my $dbh = connect_db($db_path); 29 init_schema($dbh); 30 31 my $mqtt = Net::MQTT::Simple->new("$mqtt_host:$mqtt_port"); 32 $mqtt->login( $mqtt_username, $mqtt_password ); 33 34 $SIG{INT} = $SIG{TERM} = sub { 35 warn "Shutting down fapg-daq-mqtt-sqlite\n"; 36 eval { $mqtt->disconnect; }; 37 eval { $dbh->disconnect; }; 38 exit 0; 39 }; 40 41 warn "Subscribing to mqtt://$mqtt_host:$mqtt_port/$mqtt_topic\n"; 42 warn "Writing MQTT messages to $db_path\n"; 43 44 $mqtt->run( 45 $mqtt_topic => sub { 46 my ( $topic, $message ) = @_; 47 48 eval { 49 store_message( $dbh, $topic, $message ); 50 1; 51 } or do { 52 my $error = $@ || 'unknown error'; 53 chomp $error; 54 warn "Failed to store MQTT message from $topic: $error\n"; 55 }; 56 }, 57 ); 58 59 sub env { 60 my ( $name, $default ) = @_; 61 return exists $ENV{$name} ? $ENV{$name} : $default; 62 } 63 64 sub required_env { 65 my ($name) = @_; 66 die "Missing required environment variable: $name\n" 67 if !exists $ENV{$name} || $ENV{$name} eq ''; 68 return $ENV{$name}; 69 } 70 71 sub connect_db { 72 my ($path) = @_; 73 74 my $dbh = DBI->connect( 75 "dbi:SQLite:dbname=$path", 76 '', '', 77 { RaiseError => 1, 78 PrintError => 0, 79 AutoCommit => 1, 80 sqlite_unicode => 1, 81 }, 82 ); 83 84 $dbh->do('PRAGMA journal_mode = WAL'); 85 $dbh->do('PRAGMA synchronous = NORMAL'); 86 $dbh->do('PRAGMA busy_timeout = 5000'); 87 $dbh->do('PRAGMA foreign_keys = ON'); 88 89 return $dbh; 90 } 91 92 sub init_schema { 93 my ($dbh) = @_; 94 95 $dbh->do( 96 q{ 97 CREATE TABLE IF NOT EXISTS readings ( 98 id INTEGER PRIMARY KEY AUTOINCREMENT, 99 received_at TEXT NOT NULL, 100 topic TEXT NOT NULL, 101 kind TEXT, 102 schema TEXT, 103 timestamp TEXT, 104 probe TEXT, 105 node TEXT, 106 value REAL, 107 unit TEXT, 108 source TEXT, 109 payload TEXT NOT NULL, 110 valid INTEGER NOT NULL DEFAULT 1, 111 error TEXT 112 ) 113 } 114 ); 115 116 $dbh->do( 117 q{ 118 CREATE INDEX IF NOT EXISTS readings_timestamp_idx 119 ON readings(timestamp) 120 } 121 ); 122 123 $dbh->do( 124 q{ 125 CREATE INDEX IF NOT EXISTS readings_probe_node_timestamp_idx 126 ON readings(probe, node, timestamp) 127 } 128 ); 129 130 $dbh->do( 131 q{ 132 CREATE INDEX IF NOT EXISTS readings_probe_observed_idx 133 ON readings(probe, COALESCE(timestamp, received_at), id) 134 } 135 ); 136 137 $dbh->do( 138 q{ 139 CREATE INDEX IF NOT EXISTS readings_topic_idx 140 ON readings(topic) 141 } 142 ); 143 144 init_reading_rollups($dbh); 145 146 $dbh->do( 147 q{ 148 CREATE TABLE IF NOT EXISTS node_status ( 149 id INTEGER PRIMARY KEY AUTOINCREMENT, 150 received_at TEXT NOT NULL, 151 topic TEXT NOT NULL, 152 schema TEXT, 153 timestamp TEXT, 154 probe TEXT, 155 node TEXT, 156 status TEXT, 157 message TEXT, 158 error TEXT, 159 device TEXT, 160 source TEXT, 161 payload TEXT NOT NULL, 162 valid INTEGER NOT NULL DEFAULT 1 163 ) 164 } 165 ); 166 167 $dbh->do( 168 q{ 169 CREATE INDEX IF NOT EXISTS node_status_probe_node_timestamp_idx 170 ON node_status(probe, node, timestamp) 171 } 172 ); 173 174 $dbh->do( 175 q{ 176 CREATE INDEX IF NOT EXISTS node_status_valid_probe_observed_idx 177 ON node_status(valid, probe, COALESCE(timestamp, received_at), id) 178 } 179 ); 180 181 $dbh->do( 182 q{ 183 CREATE INDEX IF NOT EXISTS node_status_topic_idx 184 ON node_status(topic) 185 } 186 ); 187 } 188 189 sub init_reading_rollups { 190 my ($dbh) = @_; 191 192 $dbh->do( 193 q{ 194 CREATE TABLE IF NOT EXISTS reading_rollups ( 195 probe TEXT NOT NULL, 196 bucket_seconds INTEGER NOT NULL, 197 bucket_epoch INTEGER NOT NULL, 198 sample_count INTEGER NOT NULL, 199 value_total REAL NOT NULL, 200 value_min REAL NOT NULL, 201 value_max REAL NOT NULL, 202 first_sample_epoch INTEGER NOT NULL, 203 last_sample_epoch INTEGER NOT NULL, 204 max_gap_seconds INTEGER NOT NULL DEFAULT 0, 205 PRIMARY KEY (probe, bucket_seconds, bucket_epoch) 206 ) WITHOUT ROWID 207 } 208 ); 209 210 for my $bucket_seconds ( 60, 3_600 ) { 211 $dbh->do( 212 qq{ 213 WITH samples AS ( 214 SELECT 215 id, 216 probe, 217 CAST( 218 CAST(strftime('%s', COALESCE(timestamp, received_at)) AS INTEGER) 219 / $bucket_seconds AS INTEGER 220 ) * $bucket_seconds AS bucket_epoch, 221 CAST(strftime('%s', COALESCE(timestamp, received_at)) AS INTEGER) 222 AS sample_epoch, 223 CAST(value AS REAL) AS value 224 FROM readings 225 WHERE valid = 1 226 AND probe IS NOT NULL 227 AND value IS NOT NULL 228 AND COALESCE(timestamp, received_at) IS NOT NULL 229 ), 230 ordered AS ( 231 SELECT 232 *, 233 sample_epoch - LAG(sample_epoch) OVER ( 234 PARTITION BY probe, bucket_epoch 235 ORDER BY sample_epoch, id 236 ) AS gap_seconds 237 FROM samples 238 ) 239 INSERT INTO reading_rollups ( 240 probe, bucket_seconds, bucket_epoch, sample_count, 241 value_total, value_min, value_max, 242 first_sample_epoch, last_sample_epoch, max_gap_seconds 243 ) 244 SELECT 245 probe, 246 $bucket_seconds, 247 bucket_epoch, 248 COUNT(*), 249 SUM(value), 250 MIN(value), 251 MAX(value), 252 MIN(sample_epoch), 253 MAX(sample_epoch), 254 MAX(COALESCE(gap_seconds, 0)) 255 FROM ordered 256 WHERE NOT EXISTS ( 257 SELECT 1 FROM reading_rollups 258 WHERE bucket_seconds = $bucket_seconds 259 LIMIT 1 260 ) 261 GROUP BY probe, bucket_epoch 262 } 263 ); 264 } 265 266 $dbh->do( 267 q{ 268 CREATE TRIGGER IF NOT EXISTS readings_rollup_insert 269 AFTER INSERT ON readings 270 WHEN NEW.valid = 1 271 AND NEW.probe IS NOT NULL 272 AND NEW.value IS NOT NULL 273 AND COALESCE(NEW.timestamp, NEW.received_at) IS NOT NULL 274 BEGIN 275 INSERT INTO reading_rollups ( 276 probe, bucket_seconds, bucket_epoch, sample_count, 277 value_total, value_min, value_max, 278 first_sample_epoch, last_sample_epoch, max_gap_seconds 279 ) 280 VALUES ( 281 NEW.probe, 282 60, 283 CAST( 284 CAST(strftime('%s', COALESCE(NEW.timestamp, NEW.received_at)) AS INTEGER) 285 / 60 AS INTEGER 286 ) * 60, 287 1, 288 CAST(NEW.value AS REAL), 289 CAST(NEW.value AS REAL), 290 CAST(NEW.value AS REAL), 291 CAST(strftime('%s', COALESCE(NEW.timestamp, NEW.received_at)) AS INTEGER), 292 CAST(strftime('%s', COALESCE(NEW.timestamp, NEW.received_at)) AS INTEGER), 293 0 294 ) 295 ON CONFLICT (probe, bucket_seconds, bucket_epoch) DO UPDATE SET 296 sample_count = sample_count + 1, 297 value_total = value_total + excluded.value_total, 298 value_min = MIN(value_min, excluded.value_min), 299 value_max = MAX(value_max, excluded.value_max), 300 max_gap_seconds = MAX( 301 max_gap_seconds, 302 CASE 303 WHEN excluded.first_sample_epoch > last_sample_epoch 304 THEN excluded.first_sample_epoch - last_sample_epoch 305 ELSE 0 306 END 307 ), 308 first_sample_epoch = MIN(first_sample_epoch, excluded.first_sample_epoch), 309 last_sample_epoch = MAX(last_sample_epoch, excluded.last_sample_epoch); 310 311 INSERT INTO reading_rollups ( 312 probe, bucket_seconds, bucket_epoch, sample_count, 313 value_total, value_min, value_max, 314 first_sample_epoch, last_sample_epoch, max_gap_seconds 315 ) 316 VALUES ( 317 NEW.probe, 318 3600, 319 CAST( 320 CAST(strftime('%s', COALESCE(NEW.timestamp, NEW.received_at)) AS INTEGER) 321 / 3600 AS INTEGER 322 ) * 3600, 323 1, 324 CAST(NEW.value AS REAL), 325 CAST(NEW.value AS REAL), 326 CAST(NEW.value AS REAL), 327 CAST(strftime('%s', COALESCE(NEW.timestamp, NEW.received_at)) AS INTEGER), 328 CAST(strftime('%s', COALESCE(NEW.timestamp, NEW.received_at)) AS INTEGER), 329 0 330 ) 331 ON CONFLICT (probe, bucket_seconds, bucket_epoch) DO UPDATE SET 332 sample_count = sample_count + 1, 333 value_total = value_total + excluded.value_total, 334 value_min = MIN(value_min, excluded.value_min), 335 value_max = MAX(value_max, excluded.value_max), 336 max_gap_seconds = MAX( 337 max_gap_seconds, 338 CASE 339 WHEN excluded.first_sample_epoch > last_sample_epoch 340 THEN excluded.first_sample_epoch - last_sample_epoch 341 ELSE 0 342 END 343 ), 344 first_sample_epoch = MIN(first_sample_epoch, excluded.first_sample_epoch), 345 last_sample_epoch = MAX(last_sample_epoch, excluded.last_sample_epoch); 346 END 347 } 348 ); 349 } 350 351 sub store_message { 352 my ( $dbh, $topic, $message ) = @_; 353 354 my ( $topic_probe, $topic_node, $topic_kind ) 355 = $topic =~ m{\Afapg/daq/([^/]+)/([^/]+)/([^/]+)\z}; 356 die "invalid MQTT topic: $topic\n" 357 if !defined $topic_kind || $topic_kind !~ /\A(?:reading|status)\z/; 358 359 my $valid = 1; 360 my $error; 361 my $data = eval { decode_json($message) }; 362 363 if ($@) { 364 $valid = 0; 365 $error = $@; 366 chomp $error; 367 $data = {}; 368 } 369 elsif ( ref $data ne 'HASH' ) { 370 $valid = 0; 371 $error = 'JSON payload is not an object'; 372 $data = {}; 373 } 374 375 if ( $topic_kind eq 'status' ) { 376 store_status_message( $dbh, $topic, $message, $topic_probe, 377 $topic_node, $data, $valid, $error ); 378 return; 379 } 380 381 my $value = $data->{value}; 382 $value = undef if defined $value && !looks_like_number($value); 383 384 my $sth = $dbh->prepare_cached( 385 q{ 386 INSERT INTO readings ( 387 received_at, 388 topic, 389 kind, 390 schema, 391 timestamp, 392 probe, 393 node, 394 value, 395 unit, 396 source, 397 payload, 398 valid, 399 error 400 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) 401 } 402 ); 403 404 $sth->execute( 405 utc_now(), $topic, 406 $topic_kind, $data->{schema}, 407 $data->{timestamp}, $data->{probe} // $topic_probe, 408 $data->{node} // $topic_node, $value, 409 $data->{unit}, $data->{source}, 410 $message, $valid, 411 $error, 412 ); 413 } 414 415 sub store_status_message { 416 my ( $dbh, $topic, $message, $topic_probe, $topic_node, $data, $valid, 417 $error ) 418 = @_; 419 420 my $status = $data->{status}; 421 if ( defined $status ) { 422 $status = lc $status; 423 if ( $status !~ /\A(?:ok|probe_error)\z/ ) { 424 $valid = 0; 425 $error 426 = defined $error 427 ? "$error; unsupported status: $status" 428 : "unsupported status: $status"; 429 } 430 } 431 else { 432 $valid = 0; 433 $error = defined $error ? "$error; missing status" : 'missing status'; 434 } 435 436 my $sth = $dbh->prepare_cached( 437 q{ 438 INSERT INTO node_status ( 439 received_at, 440 topic, 441 schema, 442 timestamp, 443 probe, 444 node, 445 status, 446 message, 447 error, 448 device, 449 source, 450 payload, 451 valid 452 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) 453 } 454 ); 455 456 $sth->execute( 457 utc_now(), $topic, 458 $data->{schema}, $data->{timestamp}, 459 $data->{probe} // $topic_probe, $data->{node} // $topic_node, 460 $status, $data->{message}, 461 $data->{error} // $error, $data->{device}, 462 $data->{source}, $message, 463 $valid, 464 ); 465 } 466 467 sub utc_now { 468 return strftime( '%Y-%m-%dT%H:%M:%SZ', gmtime ); 469 } 470