[Perl] DAQ system for the FAPG.
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