[Perl] DAQ system for the FAPG.
Run perltidy.
Changed files
- roles/daq-node/bin/fapg-daq-node
- roles/daq-node/lib/perl5/FAPG/DAQ/MQTT.pm
- roles/daq-node/t/03-mqtt.t
- roles/daq-recorder/bin/fapg-daq-mqtt-sqlite
- roles/dashboard/lib/dashboard.pm
- roles/dashboard/lib/dashboard/Controller/Dashboard.pm
- roles/dashboard/lib/dashboard/Controller/Reading.pm
- roles/dashboard/t/basic.t
- t/01-ssh-from-dev.t
roles/daq-node/bin/fapg-daq-node
@@ -41,7 +41,7 @@
41
41
} or do {
42
42
warn "DAQ loop error: $@";
43
43
sleep 30;
44
Removed:
next
44
Added:
next;
45
45
};
46
46
sleep 5;
47
47
}
roles/daq-node/lib/perl5/FAPG/DAQ/MQTT.pm
@@ -7,15 +7,15 @@
7
7
8
8
use Carp qw(croak);
9
9
use Exporter 'import';
10
Removed:
use JSON::PP qw(encode_json);
11
Removed:
use POSIX qw(strftime);
10
Added:
use JSON::PP qw(encode_json);
11
Added:
use POSIX qw(strftime);
12
12
use Sys::Hostname qw(hostname);
13
13
14
14
our @EXPORT_OK = qw(
15
Removed:
publish_mqtt
16
Removed:
mqtt_payload_for_reading
17
Removed:
mqtt_topic_for_reading
18
Removed:
mqtt_broker_endpoint
15
Added:
publish_mqtt
16
Added:
mqtt_payload_for_reading
17
Added:
mqtt_topic_for_reading
18
Added:
mqtt_broker_endpoint
19
19
);
20
20
21
21
my $DEFAULT_HOST = 'fapg-daq-five-01';
@@ -24,146 +24,131 @@
24
24
my $TOPIC_PREFIX = 'fapg/daq';
25
25
26
26
sub publish_mqtt {
27
Removed:
my ($reading, %opt) = @_;
28
Removed:
my $client = $opt{client} // _mqtt_client(%opt);
29
Removed:
my $topic = mqtt_topic_for_reading($reading, %opt);
30
Removed:
my $payload = mqtt_payload_for_reading($reading, %opt);
27
Added:
my ( $reading, %opt ) = @_;
28
Added:
my $client = $opt{client} // _mqtt_client(%opt);
29
Added:
my $topic = mqtt_topic_for_reading( $reading, %opt );
30
Added:
my $payload = mqtt_payload_for_reading( $reading, %opt );
31
31
32
Removed:
$client->publish($topic => $payload);
32
Added:
$client->publish( $topic => $payload );
33
33
34
Removed:
return {
35
Removed:
topic => $topic,
36
Removed:
payload => $payload,
37
Removed:
};
34
Added:
return {
35
Added:
topic => $topic,
36
Added:
payload => $payload,
37
Added:
};
38
38
}
39
39
40
40
sub mqtt_payload_for_reading {
41
Removed:
my ($reading, %opt) = @_;
42
Removed:
_assert_reading($reading);
41
Added:
my ( $reading, %opt ) = @_;
42
Added:
_assert_reading($reading);
43
43
44
Removed:
my $probe = lc $reading->{probe};
45
Removed:
my $node = _node_name($reading, %opt);
44
Added:
my $probe = lc $reading->{probe};
45
Added:
my $node = _node_name( $reading, %opt );
46
46
47
Removed:
my %payload = (
48
Removed:
schema => $opt{schema} // $DEFAULT_SCHEMA,
49
Removed:
timestamp => $reading->{timestamp} // _utc_timestamp(),
50
Removed:
node => $node,
51
Removed:
probe => $probe,
52
Removed:
value => 0 + $reading->{value},
53
Removed:
unit => $reading->{unit} // 'n/a',
54
Removed:
source => $opt{source} // 'fapg-daq-node',
55
Removed:
);
47
Added:
my %payload = (
48
Added:
schema => $opt{schema} // $DEFAULT_SCHEMA,
49
Added:
timestamp => $reading->{timestamp} // _utc_timestamp(),
50
Added:
node => $node,
51
Added:
probe => $probe,
52
Added:
value => 0 + $reading->{value},
53
Added:
unit => $reading->{unit} // 'n/a',
54
Added:
source => $opt{source} // 'fapg-daq-node',
55
Added:
);
56
56
57
Removed:
$payload{raw} = $reading->{raw}
58
Removed:
if exists $reading->{raw} && defined $reading->{raw};
57
Added:
$payload{raw} = $reading->{raw}
58
Added:
if exists $reading->{raw} && defined $reading->{raw};
59
59
60
Removed:
$payload{values} = [ map { 0 + $_ } $reading->{values}->@* ]
61
Removed:
if ref($reading->{values}) eq 'ARRAY';
60
Added:
$payload{values} = [ map { 0 + $_ } $reading->{values}->@* ]
61
Added:
if ref( $reading->{values} ) eq 'ARRAY';
62
62
63
Removed:
return JSON::PP->new->canonical(1)->encode(\%payload);
63
Added:
return JSON::PP->new->canonical(1)->encode( \%payload );
64
64
}
65
65
66
66
sub mqtt_topic_for_reading {
67
Removed:
my ($reading, %opt) = @_;
68
Removed:
_assert_reading($reading);
67
Added:
my ( $reading, %opt ) = @_;
68
Added:
_assert_reading($reading);
69
69
70
Removed:
my $prefix = $opt{topic_prefix} // $TOPIC_PREFIX;
71
Removed:
my $probe = _topic_level('probe', lc $reading->{probe});
72
Removed:
my $node = _topic_level('node', _node_name($reading, %opt));
70
Added:
my $prefix = $opt{topic_prefix} // $TOPIC_PREFIX;
71
Added:
my $probe = _topic_level( 'probe', lc $reading->{probe} );
72
Added:
my $node = _topic_level( 'node', _node_name( $reading, %opt ) );
73
73
74
Removed:
$prefix =~ s{/+\z}{};
74
Added:
$prefix =~ s{/+\z}{};
75
75
76
Removed:
return join '/', $prefix, $probe, $node, 'reading';
76
Added:
return join '/', $prefix, $probe, $node, 'reading';
77
77
}
78
78
79
79
sub mqtt_broker_endpoint {
80
Removed:
my (%opt) = @_;
81
Removed:
my $host = $opt{host}
82
Removed:
// $ENV{FAPG_MQTT_HOST}
83
Removed:
// $ENV{MQTT_HOST}
84
Removed:
// $DEFAULT_HOST;
80
Added:
my (%opt) = @_;
81
Added:
my $host = $opt{host} // $ENV{FAPG_MQTT_HOST} // $ENV{MQTT_HOST} // $DEFAULT_HOST;
85
82
86
Removed:
my $port = $opt{port}
87
Removed:
// $ENV{FAPG_MQTT_PORT}
88
Removed:
// $ENV{MQTT_PORT}
89
Removed:
// $DEFAULT_PORT;
83
Added:
my $port = $opt{port} // $ENV{FAPG_MQTT_PORT} // $ENV{MQTT_PORT} // $DEFAULT_PORT;
90
84
91
Removed:
return $host if $host =~ /:\d+\z/;
92
Removed:
return "$host:$port";
85
Added:
return $host if $host =~ /:\d+\z/;
86
Added:
return "$host:$port";
93
87
}
94
88
95
89
sub _mqtt_client {
96
Removed:
my (%opt) = @_;
97
Removed:
my $endpoint = mqtt_broker_endpoint(%opt);
90
Added:
my (%opt) = @_;
91
Added:
my $endpoint = mqtt_broker_endpoint(%opt);
98
92
99
Removed:
my $username = $opt{username}
100
Removed:
// $ENV{FAPG_MQTT_USERNAME}
101
Removed:
// $ENV{MQTT_USERNAME}
102
Removed:
// 'fapg_zero';
93
Added:
my $username = $opt{username} // $ENV{FAPG_MQTT_USERNAME} // $ENV{MQTT_USERNAME} // 'fapg_zero';
103
94
104
Removed:
my $password = $opt{password}
105
Removed:
// $ENV{FAPG_MQTT_PASSWORD}
106
Removed:
// $ENV{MQTT_PASSWORD};
95
Added:
my $password = $opt{password} // $ENV{FAPG_MQTT_PASSWORD} // $ENV{MQTT_PASSWORD};
107
96
108
Removed:
state %clients;
97
Added:
state %clients;
109
98
110
Removed:
my $cache_key = join "\0", $endpoint, ($username // q{}), ($password // q{});
111
Removed:
return $clients{$cache_key} if exists $clients{$cache_key};
99
Added:
my $cache_key = join "\0", $endpoint, ( $username // q{} ), ( $password // q{} );
100
Added:
return $clients{$cache_key} if exists $clients{$cache_key};
112
101
113
Removed:
require Net::MQTT::Simple;
102
Added:
require Net::MQTT::Simple;
114
103
115
Removed:
my $client = Net::MQTT::Simple->new($endpoint)
116
Removed:
or croak "cannot connect to MQTT broker $endpoint";
104
Added:
my $client = Net::MQTT::Simple->new($endpoint)
105
Added:
or croak "cannot connect to MQTT broker $endpoint";
117
106
118
Removed:
if (defined $username && $username ne q{}) {
119
Removed:
croak 'MQTT password is required when MQTT username is set'
120
Removed:
if !defined $password;
107
Added:
if ( defined $username && $username ne q{} ) {
108
Added:
croak 'MQTT password is required when MQTT username is set'
109
Added:
if !defined $password;
121
110
122
Removed:
$client->login($username, $password);
123
Removed:
}
111
Added:
$client->login( $username, $password );
112
Added:
}
124
113
125
Removed:
return $clients{$cache_key} = $client;
114
Added:
return $clients{$cache_key} = $client;
126
115
}
127
116
128
117
sub _assert_reading {
129
Removed:
my ($reading) = @_;
130
Removed:
croak 'reading must be a HASH reference'
131
Removed:
if ref($reading) ne 'HASH';
118
Added:
my ($reading) = @_;
119
Added:
croak 'reading must be a HASH reference'
120
Added:
if ref($reading) ne 'HASH';
132
121
133
Removed:
croak 'reading requires a probe field'
134
Removed:
if !defined $reading->{probe} || $reading->{probe} eq q{};
122
Added:
croak 'reading requires a probe field'
123
Added:
if !defined $reading->{probe} || $reading->{probe} eq q{};
135
124
136
Removed:
croak 'reading requires a value field'
137
Removed:
if !exists $reading->{value} || !defined $reading->{value};
125
Added:
croak 'reading requires a value field'
126
Added:
if !exists $reading->{value} || !defined $reading->{value};
138
127
139
Removed:
croak "reading value is not numeric: $reading->{value}"
140
Removed:
if $reading->{value} !~ /\A[+-]?(?:\d+(?:\.\d*)?|\.\d+)\z/;
128
Added:
croak "reading value is not numeric: $reading->{value}"
129
Added:
if $reading->{value} !~ /\A[+-]?(?:\d+(?:\.\d*)?|\.\d+)\z/;
141
130
142
Removed:
return;
131
Added:
return;
143
132
}
144
133
145
134
sub _node_name {
146
Removed:
my ($reading, %opt) = @_;
147
Removed:
return $reading->{node}
148
Removed:
// $opt{node}
149
Removed:
// $ENV{FAPG_DAQ_NODE}
150
Removed:
// $ENV{DAQ_NODE}
151
Removed:
// hostname();
135
Added:
my ( $reading, %opt ) = @_;
136
Added:
return $reading->{node} // $opt{node} // $ENV{FAPG_DAQ_NODE} // $ENV{DAQ_NODE} // hostname();
152
137
}
153
138
154
139
sub _topic_level {
155
Removed:
my ($name, $value) = @_;
156
Removed:
croak "$name topic level is required"
157
Removed:
if !defined $value || $value eq q{};
140
Added:
my ( $name, $value ) = @_;
141
Added:
croak "$name topic level is required"
142
Added:
if !defined $value || $value eq q{};
158
143
159
Removed:
croak "$name topic level must not contain '/', '+', '#', or NUL: $value"
160
Removed:
if $value =~ m{[\0/+\#]};
144
Added:
croak "$name topic level must not contain '/', '+', '#', or NUL: $value"
145
Added:
if $value =~ m{[\0/+\#]};
161
146
162
Removed:
return $value;
147
Added:
return $value;
163
148
}
164
149
165
150
sub _utc_timestamp {
166
Removed:
return strftime '%Y-%m-%dT%H:%M:%SZ', gmtime;
151
Added:
return strftime '%Y-%m-%dT%H:%M:%SZ', gmtime;
167
152
}
168
153
169
154
1;
roles/daq-node/t/03-mqtt.t
@@ -11,98 +11,97 @@
11
11
use lib "${FindBin::Bin}/../lib/perl5/";
12
12
13
13
use FAPG::DAQ::MQTT qw(
14
Removed:
publish_mqtt
15
Removed:
mqtt_payload_for_reading
16
Removed:
mqtt_topic_for_reading
17
Removed:
mqtt_broker_endpoint
14
Added:
publish_mqtt
15
Added:
mqtt_payload_for_reading
16
Added:
mqtt_topic_for_reading
17
Added:
mqtt_broker_endpoint
18
18
);
19
19
20
20
{
21
Removed:
package Local::FakeMQTT;
22
21
23
Removed:
use v5.32.1;
24
Removed:
use warnings;
22
Added:
package Local::FakeMQTT;
25
23
26
Removed:
sub new {
27
Removed:
my ($class) = @_;
28
Removed:
return bless { published => [] }, $class;
29
Removed:
}
24
Added:
use v5.32.1;
25
Added:
use warnings;
30
26
31
Removed:
sub publish {
32
Removed:
my ($self, $topic, $payload) = @_;
33
Removed:
push $self->{published}->@*, [$topic, $payload];
34
Removed:
return 1;
35
Removed:
}
27
Added:
sub new {
28
Added:
my ($class) = @_;
29
Added:
return bless { published => [] }, $class;
30
Added:
}
36
31
37
Removed:
sub published {
38
Removed:
my ($self) = @_;
39
Removed:
return $self->{published};
40
Removed:
}
32
Added:
sub publish {
33
Added:
my ( $self, $topic, $payload ) = @_;
34
Added:
push $self->{published}->@*, [ $topic, $payload ];
35
Added:
return 1;
36
Added:
}
37
Added:
38
Added:
sub published {
39
Added:
my ($self) = @_;
40
Added:
return $self->{published};
41
Added:
}
41
42
}
42
43
43
44
subtest 'topic shape' => sub {
44
Removed:
my $topic = mqtt_topic_for_reading(
45
Removed:
{ probe => 'pH', value => 7.12, unit => 'pH' },
46
Removed:
node => 'fapg-daq-zero-01',
47
Removed:
);
45
Added:
my $topic = mqtt_topic_for_reading( { probe => 'pH', value => 7.12, unit => 'pH' },
46
Added:
node => 'fapg-daq-zero-01', );
48
47
49
Removed:
is $topic, 'fapg/daq/ph/fapg-daq-zero-01/reading';
48
Added:
is $topic, 'fapg/daq/ph/fapg-daq-zero-01/reading';
50
49
};
51
50
52
51
subtest 'payload shape' => sub {
53
Removed:
my $payload = mqtt_payload_for_reading(
54
Removed:
{
55
Removed:
probe => 'pH',
56
Removed:
value => 7.12,
57
Removed:
unit => 'pH',
58
Removed:
raw => '7.12',
59
Removed:
timestamp => '2026-07-06T10:00:00Z',
60
Removed:
},
61
Removed:
node => 'fapg-daq-zero-01',
62
Removed:
);
52
Added:
my $payload = mqtt_payload_for_reading(
53
Added:
{
54
Added:
probe => 'pH',
55
Added:
value => 7.12,
56
Added:
unit => 'pH',
57
Added:
raw => '7.12',
58
Added:
timestamp => '2026-07-06T10:00:00Z',
59
Added:
},
60
Added:
node => 'fapg-daq-zero-01',
61
Added:
);
63
62
64
Removed:
my $decoded = decode_json($payload);
63
Added:
my $decoded = decode_json($payload);
65
64
66
Removed:
is $decoded->{schema}, 'fapg.daq.reading.v1';
67
Removed:
is $decoded->{timestamp}, '2026-07-06T10:00:00Z';
68
Removed:
is $decoded->{node}, 'fapg-daq-zero-01';
69
Removed:
is $decoded->{probe}, 'ph';
70
Removed:
is $decoded->{value}, 7.12;
71
Removed:
is $decoded->{unit}, 'pH';
72
Removed:
is $decoded->{raw}, '7.12';
73
Removed:
is $decoded->{source}, 'fapg-daq-node';
65
Added:
is $decoded->{schema}, 'fapg.daq.reading.v1';
66
Added:
is $decoded->{timestamp}, '2026-07-06T10:00:00Z';
67
Added:
is $decoded->{node}, 'fapg-daq-zero-01';
68
Added:
is $decoded->{probe}, 'ph';
69
Added:
is $decoded->{value}, 7.12;
70
Added:
is $decoded->{unit}, 'pH';
71
Added:
is $decoded->{raw}, '7.12';
72
Added:
is $decoded->{source}, 'fapg-daq-node';
74
73
};
75
74
76
75
subtest 'publish uses injected client, so unit tests need no broker' => sub {
77
Removed:
my $fake = Local::FakeMQTT->new;
76
Added:
my $fake = Local::FakeMQTT->new;
78
77
79
Removed:
my $result = publish_mqtt(
80
Removed:
{
81
Removed:
probe => 'DO',
82
Removed:
value => 8.34,
83
Removed:
unit => 'mg/L',
84
Removed:
timestamp => '2026-07-06T10:00:00Z',
85
Removed:
},
86
Removed:
node => 'fapg-daq-zero-02',
87
Removed:
client => $fake,
88
Removed:
);
78
Added:
my $result = publish_mqtt(
79
Added:
{
80
Added:
probe => 'DO',
81
Added:
value => 8.34,
82
Added:
unit => 'mg/L',
83
Added:
timestamp => '2026-07-06T10:00:00Z',
84
Added:
},
85
Added:
node => 'fapg-daq-zero-02',
86
Added:
client => $fake,
87
Added:
);
89
88
90
Removed:
is $result->{topic}, 'fapg/daq/do/fapg-daq-zero-02/reading';
91
Removed:
is scalar $fake->published->@*, 1;
92
Removed:
is $fake->published->[0][0], 'fapg/daq/do/fapg-daq-zero-02/reading';
89
Added:
is $result->{topic}, 'fapg/daq/do/fapg-daq-zero-02/reading';
90
Added:
is scalar $fake->published->@*, 1;
91
Added:
is $fake->published->[0][0], 'fapg/daq/do/fapg-daq-zero-02/reading';
93
92
94
Removed:
my $decoded = decode_json($fake->published->[0][1]);
95
Removed:
is $decoded->{probe}, 'do';
96
Removed:
is $decoded->{value}, 8.34;
93
Added:
my $decoded = decode_json( $fake->published->[0][1] );
94
Added:
is $decoded->{probe}, 'do';
95
Added:
is $decoded->{value}, 8.34;
97
96
};
98
97
99
98
subtest 'broker endpoint defaults to the Pi 5 broker' => sub {
100
Removed:
local $ENV{FAPG_MQTT_HOST};
101
Removed:
local $ENV{MQTT_HOST};
102
Removed:
local $ENV{FAPG_MQTT_PORT};
103
Removed:
local $ENV{MQTT_PORT};
99
Added:
local $ENV{FAPG_MQTT_HOST};
100
Added:
local $ENV{MQTT_HOST};
101
Added:
local $ENV{FAPG_MQTT_PORT};
102
Added:
local $ENV{MQTT_PORT};
104
103
105
Removed:
is mqtt_broker_endpoint(), 'fapg-daq-five-01:1883';
104
Added:
is mqtt_broker_endpoint(), 'fapg-daq-five-01:1883';
106
105
};
107
106
108
107
done_testing;
roles/daq-recorder/bin/fapg-daq-mqtt-sqlite
@@ -5,96 +5,96 @@
5
5
use warnings;
6
6
7
7
BEGIN {
8
Removed:
# Net::MQTT::Simple intentionally requires this when using username/password
9
Removed:
# over plain MQTT. The broker is expected to be reachable only over the VPN/LAN.
10
Removed:
$ENV{MQTT_SIMPLE_ALLOW_INSECURE_LOGIN} //= 1;
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
11
}
12
12
13
13
use DBI;
14
14
use JSON::PP qw(decode_json);
15
15
use Net::MQTT::Simple;
16
Removed:
use POSIX qw(strftime);
16
Added:
use POSIX qw(strftime);
17
17
use Scalar::Util qw(looks_like_number);
18
18
19
19
$| = 1;
20
20
21
Removed:
my $db_path = env('FAPG_DAQ_DB', '/var/lib/fapg-daq/readings.sqlite3');
21
Added:
my $db_path = env( 'FAPG_DAQ_DB', '/var/lib/fapg-daq/readings.sqlite3' );
22
22
my $mqtt_host = required_env('MQTT_HOST');
23
Removed:
my $mqtt_port = env('MQTT_PORT', '1883');
24
Removed:
my $mqtt_username = env('MQTT_USERNAME', 'fapg_vps');
23
Added:
my $mqtt_port = env( 'MQTT_PORT', '1883' );
24
Added:
my $mqtt_username = env( 'MQTT_USERNAME', 'fapg_vps' );
25
25
my $mqtt_password = required_env('MQTT_PASSWORD');
26
Removed:
my $mqtt_topic = env('MQTT_TOPIC', 'fapg/daq/#');
26
Added:
my $mqtt_topic = env( 'MQTT_TOPIC', 'fapg/daq/#' );
27
27
28
28
my $dbh = connect_db($db_path);
29
29
init_schema($dbh);
30
30
31
31
my $mqtt = Net::MQTT::Simple->new("$mqtt_host:$mqtt_port");
32
Removed:
$mqtt->login($mqtt_username, $mqtt_password);
32
Added:
$mqtt->login( $mqtt_username, $mqtt_password );
33
33
34
34
$SIG{INT} = $SIG{TERM} = sub {
35
Removed:
warn "Shutting down fapg-daq-mqtt-sqlite\n";
36
Removed:
eval { $mqtt->disconnect; };
37
Removed:
eval { $dbh->disconnect; };
38
Removed:
exit 0;
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
39
};
40
40
41
41
warn "Subscribing to mqtt://$mqtt_host:$mqtt_port/$mqtt_topic\n";
42
42
warn "Writing MQTT messages to $db_path\n";
43
43
44
44
$mqtt->run(
45
Removed:
$mqtt_topic => sub {
46
Removed:
my ($topic, $message) = @_;
45
Added:
$mqtt_topic => sub {
46
Added:
my ( $topic, $message ) = @_;
47
47
48
Removed:
eval {
49
Removed:
store_message($dbh, $topic, $message);
50
Removed:
1;
51
Removed:
} or do {
52
Removed:
my $error = $@ || 'unknown error';
53
Removed:
chomp $error;
54
Removed:
warn "Failed to store MQTT message from $topic: $error\n";
55
Removed:
};
56
Removed:
},
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
57
);
58
58
59
59
sub env {
60
Removed:
my ($name, $default) = @_;
61
Removed:
return exists $ENV{$name} ? $ENV{$name} : $default;
60
Added:
my ( $name, $default ) = @_;
61
Added:
return exists $ENV{$name} ? $ENV{$name} : $default;
62
62
}
63
63
64
64
sub required_env {
65
Removed:
my ($name) = @_;
66
Removed:
die "Missing required environment variable: $name\n"
67
Removed:
if !exists $ENV{$name} || $ENV{$name} eq '';
68
Removed:
return $ENV{$name};
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
69
}
70
70
71
71
sub connect_db {
72
Removed:
my ($path) = @_;
72
Added:
my ($path) = @_;
73
73
74
Removed:
my $dbh = DBI->connect(
75
Removed:
"dbi:SQLite:dbname=$path",
76
Removed:
'',
77
Removed:
'',
78
Removed:
{
79
Removed:
RaiseError => 1,
80
Removed:
PrintError => 0,
81
Removed:
AutoCommit => 1,
82
Removed:
sqlite_unicode => 1,
83
Removed:
},
84
Removed:
);
74
Added:
my $dbh = DBI->connect(
75
Added:
"dbi:SQLite:dbname=$path",
76
Added:
'', '',
77
Added:
{
78
Added:
RaiseError => 1,
79
Added:
PrintError => 0,
80
Added:
AutoCommit => 1,
81
Added:
sqlite_unicode => 1,
82
Added:
},
83
Added:
);
85
84
86
Removed:
$dbh->do('PRAGMA journal_mode = WAL');
87
Removed:
$dbh->do('PRAGMA synchronous = NORMAL');
88
Removed:
$dbh->do('PRAGMA busy_timeout = 5000');
89
Removed:
$dbh->do('PRAGMA foreign_keys = ON');
85
Added:
$dbh->do('PRAGMA journal_mode = WAL');
86
Added:
$dbh->do('PRAGMA synchronous = NORMAL');
87
Added:
$dbh->do('PRAGMA busy_timeout = 5000');
88
Added:
$dbh->do('PRAGMA foreign_keys = ON');
90
89
91
Removed:
return $dbh;
90
Added:
return $dbh;
92
91
}
93
92
94
93
sub init_schema {
95
Removed:
my ($dbh) = @_;
94
Added:
my ($dbh) = @_;
96
95
97
Removed:
$dbh->do(q{
96
Added:
$dbh->do(
97
Added:
q{
98
98
CREATE TABLE IF NOT EXISTS readings (
99
99
id INTEGER PRIMARY KEY AUTOINCREMENT,
100
100
received_at TEXT NOT NULL,
@@ -111,50 +111,56 @@
111
111
valid INTEGER NOT NULL DEFAULT 1,
112
112
error TEXT
113
113
)
114
Removed:
});
114
Added:
}
115
Added:
);
115
116
116
Removed:
$dbh->do(q{
117
Added:
$dbh->do(
118
Added:
q{
117
119
CREATE INDEX IF NOT EXISTS readings_timestamp_idx
118
120
ON readings(timestamp)
119
Removed:
});
121
Added:
}
122
Added:
);
120
123
121
Removed:
$dbh->do(q{
124
Added:
$dbh->do(
125
Added:
q{
122
126
CREATE INDEX IF NOT EXISTS readings_probe_node_timestamp_idx
123
127
ON readings(probe, node, timestamp)
124
Removed:
});
128
Added:
}
129
Added:
);
125
130
126
Removed:
$dbh->do(q{
131
Added:
$dbh->do(
132
Added:
q{
127
133
CREATE INDEX IF NOT EXISTS readings_topic_idx
128
134
ON readings(topic)
129
Removed:
});
135
Added:
}
136
Added:
);
130
137
}
131
138
132
139
sub store_message {
133
Removed:
my ($dbh, $topic, $message) = @_;
140
Added:
my ( $dbh, $topic, $message ) = @_;
134
141
135
Removed:
my ($topic_probe, $topic_node, $topic_kind) =
136
Removed:
$topic =~ m{\Afapg/daq/([^/]+)/([^/]+)/([^/]+)\z};
142
Added:
my ( $topic_probe, $topic_node, $topic_kind ) = $topic =~ m{\Afapg/daq/([^/]+)/([^/]+)/([^/]+)\z};
137
143
138
Removed:
my $valid = 1;
139
Removed:
my $error;
140
Removed:
my $data = eval { decode_json($message) };
144
Added:
my $valid = 1;
145
Added:
my $error;
146
Added:
my $data = eval { decode_json($message) };
141
147
142
Removed:
if ($@) {
143
Removed:
$valid = 0;
144
Removed:
$error = $@;
145
Removed:
chomp $error;
146
Removed:
$data = {};
147
Removed:
}
148
Removed:
elsif (ref $data ne 'HASH') {
149
Removed:
$valid = 0;
150
Removed:
$error = 'JSON payload is not an object';
151
Removed:
$data = {};
152
Removed:
}
148
Added:
if ($@) {
149
Added:
$valid = 0;
150
Added:
$error = $@;
151
Added:
chomp $error;
152
Added:
$data = {};
153
Added:
} elsif ( ref $data ne 'HASH' ) {
154
Added:
$valid = 0;
155
Added:
$error = 'JSON payload is not an object';
156
Added:
$data = {};
157
Added:
}
153
158
154
Removed:
my $value = $data->{value};
155
Removed:
$value = undef if defined $value && !looks_like_number($value);
159
Added:
my $value = $data->{value};
160
Added:
$value = undef if defined $value && !looks_like_number($value);
156
161
157
Removed:
my $sth = $dbh->prepare_cached(q{
162
Added:
my $sth = $dbh->prepare_cached(
163
Added:
q{
158
164
INSERT INTO readings (
159
165
received_at,
160
166
topic,
@@ -170,25 +176,18 @@
170
176
valid,
171
177
error
172
178
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
173
Removed:
});
179
Added:
}
180
Added:
);
174
181
175
Removed:
$sth->execute(
176
Removed:
utc_now(),
177
Removed:
$topic,
178
Removed:
$topic_kind,
179
Removed:
$data->{schema},
180
Removed:
$data->{timestamp},
181
Removed:
$data->{probe} // $topic_probe,
182
Removed:
$data->{node} // $topic_node,
183
Removed:
$value,
184
Removed:
$data->{unit},
185
Removed:
$data->{source},
186
Removed:
$message,
187
Removed:
$valid,
188
Removed:
$error,
189
Removed:
);
182
Added:
$sth->execute(
183
Added:
utc_now(), $topic, $topic_kind,
184
Added:
$data->{schema}, $data->{timestamp}, $data->{probe} // $topic_probe,
185
Added:
$data->{node} // $topic_node, $value, $data->{unit},
186
Added:
$data->{source}, $message, $valid,
187
Added:
$error,
188
Added:
);
190
189
}
191
190
192
191
sub utc_now {
193
Removed:
return strftime('%Y-%m-%dT%H:%M:%SZ', gmtime);
192
Added:
return strftime( '%Y-%m-%dT%H:%M:%SZ', gmtime );
194
193
}
roles/dashboard/lib/dashboard.pm
@@ -8,60 +8,59 @@
8
8
# This method will run once at server start
9
9
sub startup ($self) {
10
10
11
Removed:
# Load configuration from config file
12
Removed:
my $config = $self->plugin('NotYAMLConfig');
11
Added:
# Load configuration from config file
12
Added:
my $config = $self->plugin('NotYAMLConfig');
13
13
14
Removed:
# Configure the application
15
Removed:
$self->secrets( $config->{secrets} );
14
Added:
# Configure the application
15
Added:
$self->secrets( $config->{secrets} );
16
16
17
Removed:
$self->config(
18
Removed:
hypnotoad => {
19
Removed:
listen => ['http://127.0.0.1:3000'],
20
Removed:
pid_file => '/run/fapg-daq-dashboard/hypnotoad.pid',
21
Removed:
workers => 2,
22
Removed:
proxy => 1,
23
Removed:
},
24
Removed:
);
17
Added:
$self->config(
18
Added:
hypnotoad => {
19
Added:
listen => ['http://127.0.0.1:3000'],
20
Added:
pid_file => '/run/fapg-daq-dashboard/hypnotoad.pid',
21
Added:
workers => 2,
22
Added:
proxy => 1,
23
Added:
},
24
Added:
);
25
25
26
Removed:
my $db_path = $ENV{FAPG_DAQ_DB}
27
Removed:
// $self->home->rel_file('fapg-daq.sqlite3');
26
Added:
my $db_path = $ENV{FAPG_DAQ_DB} // $self->home->rel_file('fapg-daq.sqlite3');
28
27
29
Removed:
my $sqlite = Mojo::SQLite->new("sqlite:$db_path");
28
Added:
my $sqlite = Mojo::SQLite->new("sqlite:$db_path");
30
29
31
Removed:
$self->helper( sqlite => sub { $sqlite } );
30
Added:
$self->helper( sqlite => sub { $sqlite } );
32
31
33
Removed:
$self->helper(
34
Removed:
probes => sub {
35
Removed:
return [
36
Removed:
{
37
Removed:
key => 'ph',
38
Removed:
label => 'pH',
39
Removed:
unit => 'pH',
40
Removed:
},
41
Removed:
{
42
Removed:
key => 'do',
43
Removed:
label => 'Dissolved Oxygen',
44
Removed:
unit => 'mg/L',
45
Removed:
},
46
Removed:
{
47
Removed:
key => 'orp',
48
Removed:
label => 'Redox / ORP',
49
Removed:
unit => 'mV',
50
Removed:
},
51
Removed:
{
52
Removed:
key => 'ec',
53
Removed:
label => 'Electrical Conductivity',
54
Removed:
unit => 'µS/cm',
55
Removed:
},
56
Removed:
];
57
Removed:
}
58
Removed:
);
32
Added:
$self->helper(
33
Added:
probes => sub {
34
Added:
return [
35
Added:
{
36
Added:
key => 'ph',
37
Added:
label => 'pH',
38
Added:
unit => 'pH',
39
Added:
},
40
Added:
{
41
Added:
key => 'do',
42
Added:
label => 'Dissolved Oxygen',
43
Added:
unit => 'mg/L',
44
Added:
},
45
Added:
{
46
Added:
key => 'orp',
47
Added:
label => 'Redox / ORP',
48
Added:
unit => 'mV',
49
Added:
},
50
Added:
{
51
Added:
key => 'ec',
52
Added:
label => 'Electrical Conductivity',
53
Added:
unit => 'µS/cm',
54
Added:
},
55
Added:
];
56
Added:
}
57
Added:
);
59
58
60
Removed:
my $r = $self->routes;
59
Added:
my $r = $self->routes;
61
60
62
Removed:
$r->get('/')->to('dashboard#index');
61
Added:
$r->get('/')->to('dashboard#index');
63
62
64
Removed:
$r->get('/api/readings/:probe')->to('Reading#list');
63
Added:
$r->get('/api/readings/:probe')->to('Reading#list');
65
64
}
66
65
67
66
1;
roles/dashboard/lib/dashboard/Controller/Dashboard.pm
@@ -4,10 +4,10 @@
4
4
use Mojo::Base 'Mojolicious::Controller', -signatures;
5
5
6
6
sub index ($self) {
7
Removed:
$self->render(
8
Removed:
template => 'dashboard/index',
9
Removed:
probes => $self->probes,
10
Removed:
);
7
Added:
$self->render(
8
Added:
template => 'dashboard/index',
9
Added:
probes => $self->probes,
10
Added:
);
11
11
}
12
12
13
13
1;
roles/dashboard/lib/dashboard/Controller/Reading.pm
@@ -4,42 +4,42 @@
4
4
use Mojo::Base 'Mojolicious::Controller', -signatures;
5
5
6
6
sub list ($self) {
7
Removed:
my $probe = $self->param('probe') // '';
7
Added:
my $probe = $self->param('probe') // '';
8
8
9
Removed:
my %known = map { $_->{key} => 1 } $self->probes->@*;
9
Added:
my %known = map { $_->{key} => 1 } $self->probes->@*;
10
10
11
Removed:
return $self->render(
12
Removed:
status => 404,
13
Removed:
json => {
14
Removed:
error => "Unknown probe type: $probe",
15
Removed:
},
16
Removed:
) unless $known{$probe};
11
Added:
return $self->render(
12
Added:
status => 404,
13
Added:
json => {
14
Added:
error => "Unknown probe type: $probe",
15
Added:
},
16
Added:
) unless $known{$probe};
17
17
18
Removed:
my $limit = $self->param('limit') // 300;
18
Added:
my $limit = $self->param('limit') // 300;
19
19
20
Removed:
$limit = 300 unless $limit =~ /^\d+$/;
21
Removed:
$limit = 2_000 if $limit > 2_000;
20
Added:
$limit = 300 unless $limit =~ /^\d+$/;
21
Added:
$limit = 2_000 if $limit > 2_000;
22
22
23
Removed:
my $rows = $self->sqlite->db->query(
24
Removed:
q{
23
Added:
my $rows = $self->sqlite->db->query(
24
Added:
q{
25
25
SELECT timestamp, node, probe, value, unit
26
26
FROM readings
27
27
WHERE probe = ?
28
28
ORDER BY timestamp DESC
29
29
LIMIT ?
30
30
},
31
Removed:
$probe,
32
Removed:
$limit,
33
Removed:
)->hashes->to_array;
31
Added:
$probe,
32
Added:
$limit,
33
Added:
)->hashes->to_array;
34
34
35
Removed:
$rows = [ reverse $rows->@* ];
35
Added:
$rows = [ reverse $rows->@* ];
36
36
37
Removed:
$self->render(
38
Removed:
json => {
39
Removed:
probe => $probe,
40
Removed:
readings => $rows,
41
Removed:
},
42
Removed:
);
37
Added:
$self->render(
38
Added:
json => {
39
Added:
probe => $probe,
40
Added:
readings => $rows,
41
Added:
},
42
Added:
);
43
43
}
44
44
45
45
1;
roles/dashboard/t/basic.t
@@ -13,7 +13,8 @@
13
13
14
14
my $t = Test::Mojo->new('dashboard');
15
15
16
Removed:
$t->app->sqlite->db->query(q{
16
Added:
$t->app->sqlite->db->query(
17
Added:
q{
17
18
CREATE TABLE readings (
18
19
id INTEGER PRIMARY KEY AUTOINCREMENT,
19
20
timestamp TEXT NOT NULL,
@@ -23,12 +24,15 @@
23
24
unit TEXT,
24
25
raw_json TEXT
25
26
)
26
Removed:
});
27
Added:
}
28
Added:
);
27
29
28
Removed:
$t->app->sqlite->db->query(q{
30
Added:
$t->app->sqlite->db->query(
31
Added:
q{
29
32
INSERT INTO readings (timestamp, node, probe, value, unit)
30
33
VALUES (?, ?, ?, ?, ?)
31
Removed:
}, '2026-07-05T12:00:00Z', 'fapg-daq-zero-ph-01', 'ph', 7.12, 'pH');
34
Added:
}, '2026-07-05T12:00:00Z', 'fapg-daq-zero-ph-01', 'ph', 7.12, 'pH'
35
Added:
);
32
36
33
37
$t->get_ok('/')
34
38
->status_is(200)
@@ -37,10 +41,9 @@
37
41
38
42
$t->get_ok('/api/readings/ph')
39
43
->status_is(200)
40
Removed:
->json_is('/probe' => 'ph')
41
Removed:
->json_is('/readings/0/value' => 7.12);
44
Added:
->json_is( '/probe' => 'ph' )
45
Added:
->json_is( '/readings/0/value' => 7.12 );
42
46
43
Removed:
$t->get_ok('/api/readings/nope')
44
Removed:
->status_is(404);
47
Added:
$t->get_ok('/api/readings/nope')->status_is(404);
45
48
46
49
done_testing;
t/01-ssh-from-dev.t
@@ -19,7 +19,7 @@
19
19
local $SIG{ALRM} = sub { die "SSH timeout after $timeout seconds.\n" };
20
20
alarm $timeout;
21
21
22
Removed:
my $ssh = Net::SSH::Perl->new( $h->{hostname});
22
Added:
my $ssh = Net::SSH::Perl->new( $h->{hostname} );
23
23
$ssh->login( $h->{user}, $h->{password} );
24
24
( $out, $err, $exit ) = $ssh->cmd('true');
25
25