[Perl] DAQ system for the FAPG.
Introduce DAQ recorder role.
roles/daq-recorder/bin/fapg-daq-mqtt-sqlite
@@ -0,0 +1,194 @@
1
Added:
#!/usr/bin/env perl
2
Added:
3
Added:
use v5.32.1;
4
Added:
use strict;
5
Added:
use warnings;
6
Added:
7
Added:
BEGIN {
8
Added:
# Net::MQTT::Simple intentionally requires this when using username/password
9
Added:
# over plain MQTT. The broker is expected to be reachable only over the VPN/LAN.
10
Added:
$ENV{MQTT_SIMPLE_ALLOW_INSECURE_LOGIN} //= 1;
11
Added:
}
12
Added:
13
Added:
use DBI;
14
Added:
use JSON::PP qw(decode_json);
15
Added:
use Net::MQTT::Simple;
16
Added:
use POSIX qw(strftime);
17
Added:
use Scalar::Util qw(looks_like_number);
18
Added:
19
Added:
$| = 1;
20
Added:
21
Added:
my $db_path = env('FAPG_DAQ_DB', '/var/lib/fapg-daq/readings.sqlite3');
22
Added:
my $mqtt_host = required_env('MQTT_HOST');
23
Added:
my $mqtt_port = env('MQTT_PORT', '1883');
24
Added:
my $mqtt_username = env('MQTT_USERNAME', 'fapg_vps');
25
Added:
my $mqtt_password = required_env('MQTT_PASSWORD');
26
Added:
my $mqtt_topic = env('MQTT_TOPIC', 'fapg/daq/#');
27
Added:
28
Added:
my $dbh = connect_db($db_path);
29
Added:
init_schema($dbh);
30
Added:
31
Added:
my $mqtt = Net::MQTT::Simple->new("$mqtt_host:$mqtt_port");
32
Added:
$mqtt->login($mqtt_username, $mqtt_password);
33
Added:
34
Added:
$SIG{INT} = $SIG{TERM} = sub {
35
Added:
warn "Shutting down fapg-daq-mqtt-sqlite\n";
36
Added:
eval { $mqtt->disconnect; };
37
Added:
eval { $dbh->disconnect; };
38
Added:
exit 0;
39
Added:
};
40
Added:
41
Added:
warn "Subscribing to mqtt://$mqtt_host:$mqtt_port/$mqtt_topic\n";
42
Added:
warn "Writing MQTT messages to $db_path\n";
43
Added:
44
Added:
$mqtt->run(
45
Added:
$mqtt_topic => sub {
46
Added:
my ($topic, $message) = @_;
47
Added:
48
Added:
eval {
49
Added:
store_message($dbh, $topic, $message);
50
Added:
1;
51
Added:
} or do {
52
Added:
my $error = $@ || 'unknown error';
53
Added:
chomp $error;
54
Added:
warn "Failed to store MQTT message from $topic: $error\n";
55
Added:
};
56
Added:
},
57
Added:
);
58
Added:
59
Added:
sub env {
60
Added:
my ($name, $default) = @_;
61
Added:
return exists $ENV{$name} ? $ENV{$name} : $default;
62
Added:
}
63
Added:
64
Added:
sub required_env {
65
Added:
my ($name) = @_;
66
Added:
die "Missing required environment variable: $name\n"
67
Added:
if !exists $ENV{$name} || $ENV{$name} eq '';
68
Added:
return $ENV{$name};
69
Added:
}
70
Added:
71
Added:
sub connect_db {
72
Added:
my ($path) = @_;
73
Added:
74
Added:
my $dbh = DBI->connect(
75
Added:
"dbi:SQLite:dbname=$path",
76
Added:
'',
77
Added:
'',
78
Added:
{
79
Added:
RaiseError => 1,
80
Added:
PrintError => 0,
81
Added:
AutoCommit => 1,
82
Added:
sqlite_unicode => 1,
83
Added:
},
84
Added:
);
85
Added:
86
Added:
$dbh->do('PRAGMA journal_mode = WAL');
87
Added:
$dbh->do('PRAGMA synchronous = NORMAL');
88
Added:
$dbh->do('PRAGMA busy_timeout = 5000');
89
Added:
$dbh->do('PRAGMA foreign_keys = ON');
90
Added:
91
Added:
return $dbh;
92
Added:
}
93
Added:
94
Added:
sub init_schema {
95
Added:
my ($dbh) = @_;
96
Added:
97
Added:
$dbh->do(q{
98
Added:
CREATE TABLE IF NOT EXISTS readings (
99
Added:
id INTEGER PRIMARY KEY AUTOINCREMENT,
100
Added:
received_at TEXT NOT NULL,
101
Added:
topic TEXT NOT NULL,
102
Added:
kind TEXT,
103
Added:
schema TEXT,
104
Added:
timestamp TEXT,
105
Added:
probe TEXT,
106
Added:
node TEXT,
107
Added:
value REAL,
108
Added:
unit TEXT,
109
Added:
source TEXT,
110
Added:
payload TEXT NOT NULL,
111
Added:
valid INTEGER NOT NULL DEFAULT 1,
112
Added:
error TEXT
113
Added:
)
114
Added:
});
115
Added:
116
Added:
$dbh->do(q{
117
Added:
CREATE INDEX IF NOT EXISTS readings_timestamp_idx
118
Added:
ON readings(timestamp)
119
Added:
});
120
Added:
121
Added:
$dbh->do(q{
122
Added:
CREATE INDEX IF NOT EXISTS readings_probe_node_timestamp_idx
123
Added:
ON readings(probe, node, timestamp)
124
Added:
});
125
Added:
126
Added:
$dbh->do(q{
127
Added:
CREATE INDEX IF NOT EXISTS readings_topic_idx
128
Added:
ON readings(topic)
129
Added:
});
130
Added:
}
131
Added:
132
Added:
sub store_message {
133
Added:
my ($dbh, $topic, $message) = @_;
134
Added:
135
Added:
my ($topic_probe, $topic_node, $topic_kind) =
136
Added:
$topic =~ m{\Afapg/daq/([^/]+)/([^/]+)/([^/]+)\z};
137
Added:
138
Added:
my $valid = 1;
139
Added:
my $error;
140
Added:
my $data = eval { decode_json($message) };
141
Added:
142
Added:
if ($@) {
143
Added:
$valid = 0;
144
Added:
$error = $@;
145
Added:
chomp $error;
146
Added:
$data = {};
147
Added:
}
148
Added:
elsif (ref $data ne 'HASH') {
149
Added:
$valid = 0;
150
Added:
$error = 'JSON payload is not an object';
151
Added:
$data = {};
152
Added:
}
153
Added:
154
Added:
my $value = $data->{value};
155
Added:
$value = undef if defined $value && !looks_like_number($value);
156
Added:
157
Added:
my $sth = $dbh->prepare_cached(q{
158
Added:
INSERT INTO readings (
159
Added:
received_at,
160
Added:
topic,
161
Added:
kind,
162
Added:
schema,
163
Added:
timestamp,
164
Added:
probe,
165
Added:
node,
166
Added:
value,
167
Added:
unit,
168
Added:
source,
169
Added:
payload,
170
Added:
valid,
171
Added:
error
172
Added:
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
173
Added:
});
174
Added:
175
Added:
$sth->execute(
176
Added:
utc_now(),
177
Added:
$topic,
178
Added:
$topic_kind,
179
Added:
$data->{schema},
180
Added:
$data->{timestamp},
181
Added:
$data->{probe} // $topic_probe,
182
Added:
$data->{node} // $topic_node,
183
Added:
$value,
184
Added:
$data->{unit},
185
Added:
$data->{source},
186
Added:
$message,
187
Added:
$valid,
188
Added:
$error,
189
Added:
);
190
Added:
}
191
Added:
192
Added:
sub utc_now {
193
Added:
return strftime('%Y-%m-%dT%H:%M:%SZ', gmtime);
194
Added:
}
roles/daq-recorder/deploy
@@ -0,0 +1,14 @@
1
Added:
#!/usr/bin/env bash
2
Added:
3
Added:
sudo apt install \
4
Added:
sqlite3 \
5
Added:
libdbi-perl \
6
Added:
libdbd-sqlite3-perl \
7
Added:
libnet-mqtt-simple-perl
8
Added:
9
Added:
sudo useradd \
10
Added:
--system \
11
Added:
--home-dir /var/lib/fapg-daq \
12
Added:
--create-home \
13
Added:
--shell /usr/sbin/nologin \
14
Added:
fapg-daq