[Perl] DAQ system for the FAPG.
Reformat Perl files with best practices
Run perltidy --perl-best-practices across tracked Perl modules, scripts, and tests. Use -nst for in-place formatting with this perltidy profile and exclude generated backup files from the commit.
Changed files
- lib/FAPG/DAQ/MQTT.pm
- roles/daq-hub/bin/fapg-daq-hub-status
- roles/daq-node/bin/fapg-daq-node
- roles/daq-node/lib/FAPG/DAQ/EZO/USB.pm
- roles/daq-node/t/01-ezo-usb.t
- roles/daq-node/t/02-probe.t
- roles/daq-node/t/03-mqtt.t
- roles/daq-recorder/bin/fapg-daq-mqtt-sqlite
- roles/dashboard/lib/FAPG/DAQ/Dashboard.pm
- roles/dashboard/lib/FAPG/DAQ/Dashboard/Controller/Dashboard.pm
- roles/dashboard/lib/FAPG/DAQ/Dashboard/Controller/Reading.pm
- roles/dashboard/t/00-homepage.t
- roles/dashboard/t/01-graph-page.t
- roles/dashboard/t/02-readings-api.t
- roles/dashboard/t/03-status-api.t
- roles/dashboard/t/lib/Dashboard/Test.pm
- t/00-ping-from-dev.t
- t/01-ssh-from-dev.t
- t/lib/FAPG/TestSupport.pm
- t/lib/FAPG/TestSupport.pm.example
lib/FAPG/DAQ/MQTT.pm
@@ -12,13 +12,13 @@
12
12
use Sys::Hostname qw(hostname);
13
13
14
14
our @EXPORT_OK = qw(
15
Removed:
publish_mqtt
16
Removed:
publish_status_mqtt
17
Removed:
mqtt_payload_for_reading
18
Removed:
mqtt_payload_for_status
19
Removed:
mqtt_topic_for_reading
20
Removed:
mqtt_topic_for_status
21
Removed:
mqtt_broker_endpoint
15
Added:
publish_mqtt
16
Added:
publish_status_mqtt
17
Added:
mqtt_payload_for_reading
18
Added:
mqtt_payload_for_status
19
Added:
mqtt_topic_for_reading
20
Added:
mqtt_topic_for_status
21
Added:
mqtt_broker_endpoint
22
22
);
23
23
24
24
my $DEFAULT_HOST = 'fapg-daq-five-01';
@@ -28,199 +28,205 @@
28
28
my $TOPIC_PREFIX = 'fapg/daq';
29
29
30
30
sub publish_mqtt {
31
Removed:
my ( $reading, %opt ) = @_;
32
Removed:
my $client = $opt{client} // _mqtt_client(%opt);
33
Removed:
my $topic = mqtt_topic_for_reading( $reading, %opt );
34
Removed:
my $payload = mqtt_payload_for_reading( $reading, %opt );
31
Added:
my ( $reading, %opt ) = @_;
32
Added:
my $client = $opt{client} // _mqtt_client(%opt);
33
Added:
my $topic = mqtt_topic_for_reading( $reading, %opt );
34
Added:
my $payload = mqtt_payload_for_reading( $reading, %opt );
35
35
36
Removed:
$client->publish( $topic => $payload );
36
Added:
$client->publish( $topic => $payload );
37
37
38
Removed:
return {
39
Removed:
topic => $topic,
40
Removed:
payload => $payload,
41
Removed:
};
38
Added:
return {
39
Added:
topic => $topic,
40
Added:
payload => $payload,
41
Added:
};
42
42
}
43
43
44
44
sub publish_status_mqtt {
45
Removed:
my ( $status, %opt ) = @_;
46
Removed:
my $client = $opt{client} // _mqtt_client(%opt);
47
Removed:
my $topic = mqtt_topic_for_status( $status, %opt );
48
Removed:
my $payload = mqtt_payload_for_status( $status, %opt );
45
Added:
my ( $status, %opt ) = @_;
46
Added:
my $client = $opt{client} // _mqtt_client(%opt);
47
Added:
my $topic = mqtt_topic_for_status( $status, %opt );
48
Added:
my $payload = mqtt_payload_for_status( $status, %opt );
49
49
50
Removed:
$client->publish( $topic => $payload );
50
Added:
$client->publish( $topic => $payload );
51
51
52
Removed:
return {
53
Removed:
topic => $topic,
54
Removed:
payload => $payload,
55
Removed:
};
52
Added:
return {
53
Added:
topic => $topic,
54
Added:
payload => $payload,
55
Added:
};
56
56
}
57
57
58
58
sub mqtt_payload_for_reading {
59
Removed:
my ( $reading, %opt ) = @_;
60
Removed:
_assert_reading($reading);
59
Added:
my ( $reading, %opt ) = @_;
60
Added:
_assert_reading($reading);
61
61
62
Removed:
my $probe = lc $reading->{probe};
63
Removed:
my $node = _node_name( $reading, %opt );
62
Added:
my $probe = lc $reading->{probe};
63
Added:
my $node = _node_name( $reading, %opt );
64
64
65
Removed:
my %payload = (
66
Removed:
schema => $opt{schema} // $DEFAULT_SCHEMA,
67
Removed:
timestamp => $reading->{timestamp} // _utc_timestamp(),
68
Removed:
node => $node,
69
Removed:
probe => $probe,
70
Removed:
value => 0 + $reading->{value},
71
Removed:
unit => $reading->{unit} // 'n/a',
72
Removed:
source => $opt{source} // 'fapg-daq-node',
73
Removed:
);
65
Added:
my %payload = (
66
Added:
schema => $opt{schema} // $DEFAULT_SCHEMA,
67
Added:
timestamp => $reading->{timestamp} // _utc_timestamp(),
68
Added:
node => $node,
69
Added:
probe => $probe,
70
Added:
value => 0 + $reading->{value},
71
Added:
unit => $reading->{unit} // 'n/a',
72
Added:
source => $opt{source} // 'fapg-daq-node',
73
Added:
);
74
74
75
Removed:
$payload{raw} = $reading->{raw}
76
Removed:
if exists $reading->{raw} && defined $reading->{raw};
75
Added:
$payload{raw} = $reading->{raw}
76
Added:
if exists $reading->{raw} && defined $reading->{raw};
77
77
78
Removed:
$payload{values} = [ map { 0 + $_ } $reading->{values}->@* ]
79
Removed:
if ref( $reading->{values} ) eq 'ARRAY';
78
Added:
$payload{values} = [ map { 0 + $_ } $reading->{values}->@* ]
79
Added:
if ref( $reading->{values} ) eq 'ARRAY';
80
80
81
Removed:
return JSON::PP->new->canonical(1)->encode( \%payload );
81
Added:
return JSON::PP->new->canonical(1)->encode( \%payload );
82
82
}
83
83
84
84
sub mqtt_payload_for_status {
85
Removed:
my ( $status, %opt ) = @_;
86
Removed:
_assert_status($status);
85
Added:
my ( $status, %opt ) = @_;
86
Added:
_assert_status($status);
87
87
88
Removed:
my $probe = lc $status->{probe};
89
Removed:
my $node = _node_name( $status, %opt );
88
Added:
my $probe = lc $status->{probe};
89
Added:
my $node = _node_name( $status, %opt );
90
90
91
Removed:
my %payload = (
92
Removed:
schema => $opt{schema} // $DEFAULT_STATUS_SCHEMA,
93
Removed:
timestamp => $status->{timestamp} // _utc_timestamp(),
94
Removed:
node => $node,
95
Removed:
probe => $probe,
96
Removed:
status => lc $status->{status},
97
Removed:
source => $opt{source} // 'fapg-daq-node',
98
Removed:
);
91
Added:
my %payload = (
92
Added:
schema => $opt{schema} // $DEFAULT_STATUS_SCHEMA,
93
Added:
timestamp => $status->{timestamp} // _utc_timestamp(),
94
Added:
node => $node,
95
Added:
probe => $probe,
96
Added:
status => lc $status->{status},
97
Added:
source => $opt{source} // 'fapg-daq-node',
98
Added:
);
99
99
100
Removed:
for my $field (qw(message error device)) {
101
Removed:
$payload{$field} = $status->{$field}
102
Removed:
if exists $status->{$field} && defined $status->{$field};
103
Removed:
}
100
Added:
for my $field (qw(message error device)) {
101
Added:
$payload{$field} = $status->{$field}
102
Added:
if exists $status->{$field} && defined $status->{$field};
103
Added:
}
104
104
105
Removed:
return JSON::PP->new->canonical(1)->encode( \%payload );
105
Added:
return JSON::PP->new->canonical(1)->encode( \%payload );
106
106
}
107
107
108
108
sub mqtt_topic_for_reading {
109
Removed:
my ( $reading, %opt ) = @_;
110
Removed:
_assert_reading($reading);
109
Added:
my ( $reading, %opt ) = @_;
110
Added:
_assert_reading($reading);
111
111
112
Removed:
my $prefix = $opt{topic_prefix} // $TOPIC_PREFIX;
113
Removed:
my $probe = _topic_level( 'probe', lc $reading->{probe} );
114
Removed:
my $node = _topic_level( 'node', _node_name( $reading, %opt ) );
112
Added:
my $prefix = $opt{topic_prefix} // $TOPIC_PREFIX;
113
Added:
my $probe = _topic_level( 'probe', lc $reading->{probe} );
114
Added:
my $node = _topic_level( 'node', _node_name( $reading, %opt ) );
115
115
116
Removed:
$prefix =~ s{/+\z}{};
116
Added:
$prefix =~ s{/+\z}{};
117
117
118
Removed:
return join '/', $prefix, $probe, $node, 'reading';
118
Added:
return join '/', $prefix, $probe, $node, 'reading';
119
119
}
120
120
121
121
sub mqtt_topic_for_status {
122
Removed:
my ( $status, %opt ) = @_;
123
Removed:
_assert_status($status);
122
Added:
my ( $status, %opt ) = @_;
123
Added:
_assert_status($status);
124
124
125
Removed:
my $prefix = $opt{topic_prefix} // $TOPIC_PREFIX;
126
Removed:
my $probe = _topic_level( 'probe', lc $status->{probe} );
127
Removed:
my $node = _topic_level( 'node', _node_name( $status, %opt ) );
125
Added:
my $prefix = $opt{topic_prefix} // $TOPIC_PREFIX;
126
Added:
my $probe = _topic_level( 'probe', lc $status->{probe} );
127
Added:
my $node = _topic_level( 'node', _node_name( $status, %opt ) );
128
128
129
Removed:
$prefix =~ s{/+\z}{};
129
Added:
$prefix =~ s{/+\z}{};
130
130
131
Removed:
return join '/', $prefix, $probe, $node, 'status';
131
Added:
return join '/', $prefix, $probe, $node, 'status';
132
132
}
133
133
134
134
sub mqtt_broker_endpoint {
135
Removed:
my (%opt) = @_;
136
Removed:
my $host = $opt{host} // $ENV{FAPG_MQTT_HOST} // $ENV{MQTT_HOST} // $DEFAULT_HOST;
135
Added:
my (%opt) = @_;
136
Added:
my $host = $opt{host} // $ENV{FAPG_MQTT_HOST} // $ENV{MQTT_HOST}
137
Added:
// $DEFAULT_HOST;
137
138
138
Removed:
my $port = $opt{port} // $ENV{FAPG_MQTT_PORT} // $ENV{MQTT_PORT} // $DEFAULT_PORT;
139
Added:
my $port = $opt{port} // $ENV{FAPG_MQTT_PORT} // $ENV{MQTT_PORT}
140
Added:
// $DEFAULT_PORT;
139
141
140
Removed:
return $host if $host =~ /:\d+\z/;
141
Removed:
return "$host:$port";
142
Added:
return $host if $host =~ /:\d+\z/;
143
Added:
return "$host:$port";
142
144
}
143
145
144
146
sub _mqtt_client {
145
Removed:
my (%opt) = @_;
146
Removed:
my $endpoint = mqtt_broker_endpoint(%opt);
147
Added:
my (%opt) = @_;
148
Added:
my $endpoint = mqtt_broker_endpoint(%opt);
147
149
148
Removed:
my $username = $opt{username} // $ENV{FAPG_MQTT_USERNAME} // $ENV{MQTT_USERNAME} // 'fapg_zero';
150
Added:
my $username = $opt{username} // $ENV{FAPG_MQTT_USERNAME}
151
Added:
// $ENV{MQTT_USERNAME} // 'fapg_zero';
149
152
150
Removed:
my $password = $opt{password} // $ENV{FAPG_MQTT_PASSWORD} // $ENV{MQTT_PASSWORD};
153
Added:
my $password = $opt{password} // $ENV{FAPG_MQTT_PASSWORD}
154
Added:
// $ENV{MQTT_PASSWORD};
151
155
152
Removed:
state %clients;
156
Added:
state %clients;
153
157
154
Removed:
my $cache_key = join "\0", $endpoint, ( $username // q{} ), ( $password // q{} );
155
Removed:
return $clients{$cache_key} if exists $clients{$cache_key};
158
Added:
my $cache_key = join "\0", $endpoint, ( $username // q{} ),
159
Added:
( $password // q{} );
160
Added:
return $clients{$cache_key} if exists $clients{$cache_key};
156
161
157
Removed:
require Net::MQTT::Simple;
162
Added:
require Net::MQTT::Simple;
158
163
159
Removed:
my $client = Net::MQTT::Simple->new($endpoint)
160
Removed:
or croak "cannot connect to MQTT broker $endpoint";
164
Added:
my $client = Net::MQTT::Simple->new($endpoint)
165
Added:
or croak "cannot connect to MQTT broker $endpoint";
161
166
162
Removed:
if ( defined $username && $username ne q{} ) {
163
Removed:
croak 'MQTT password is required when MQTT username is set'
164
Removed:
if !defined $password;
167
Added:
if ( defined $username && $username ne q{} ) {
168
Added:
croak 'MQTT password is required when MQTT username is set'
169
Added:
if !defined $password;
165
170
166
Removed:
$client->login( $username, $password );
167
Removed:
}
171
Added:
$client->login( $username, $password );
172
Added:
}
168
173
169
Removed:
return $clients{$cache_key} = $client;
174
Added:
return $clients{$cache_key} = $client;
170
175
}
171
176
172
177
sub _assert_reading {
173
Removed:
my ($reading) = @_;
174
Removed:
croak 'reading must be a HASH reference'
175
Removed:
if ref($reading) ne 'HASH';
178
Added:
my ($reading) = @_;
179
Added:
croak 'reading must be a HASH reference'
180
Added:
if ref($reading) ne 'HASH';
176
181
177
Removed:
croak 'reading requires a probe field'
178
Removed:
if !defined $reading->{probe} || $reading->{probe} eq q{};
182
Added:
croak 'reading requires a probe field'
183
Added:
if !defined $reading->{probe} || $reading->{probe} eq q{};
179
184
180
Removed:
croak 'reading requires a value field'
181
Removed:
if !exists $reading->{value} || !defined $reading->{value};
185
Added:
croak 'reading requires a value field'
186
Added:
if !exists $reading->{value} || !defined $reading->{value};
182
187
183
Removed:
croak "reading value is not numeric: $reading->{value}"
184
Removed:
if $reading->{value} !~ /\A[+-]?(?:\d+(?:\.\d*)?|\.\d+)\z/;
188
Added:
croak "reading value is not numeric: $reading->{value}"
189
Added:
if $reading->{value} !~ /\A[+-]?(?:\d+(?:\.\d*)?|\.\d+)\z/;
185
190
186
Removed:
return;
191
Added:
return;
187
192
}
188
193
189
194
sub _assert_status {
190
Removed:
my ($status) = @_;
191
Removed:
croak 'status must be a HASH reference'
192
Removed:
if ref($status) ne 'HASH';
195
Added:
my ($status) = @_;
196
Added:
croak 'status must be a HASH reference'
197
Added:
if ref($status) ne 'HASH';
193
198
194
Removed:
croak 'status requires a probe field'
195
Removed:
if !defined $status->{probe} || $status->{probe} eq q{};
199
Added:
croak 'status requires a probe field'
200
Added:
if !defined $status->{probe} || $status->{probe} eq q{};
196
201
197
Removed:
croak 'status requires a status field'
198
Removed:
if !defined $status->{status} || $status->{status} eq q{};
202
Added:
croak 'status requires a status field'
203
Added:
if !defined $status->{status} || $status->{status} eq q{};
199
204
200
Removed:
croak "unsupported node status '$status->{status}'"
201
Removed:
if lc( $status->{status} ) !~ /\A(?:ok|probe_error)\z/;
205
Added:
croak "unsupported node status '$status->{status}'"
206
Added:
if lc( $status->{status} ) !~ /\A(?:ok|probe_error)\z/;
202
207
203
Removed:
return;
208
Added:
return;
204
209
}
205
210
206
211
sub _node_name {
207
Removed:
my ( $reading, %opt ) = @_;
208
Removed:
return $reading->{node} // $opt{node} // $ENV{FAPG_DAQ_NODE} // $ENV{DAQ_NODE} // hostname();
212
Added:
my ( $reading, %opt ) = @_;
213
Added:
return $reading->{node} // $opt{node} // $ENV{FAPG_DAQ_NODE}
214
Added:
// $ENV{DAQ_NODE} // hostname();
209
215
}
210
216
211
217
sub _topic_level {
212
Removed:
my ( $name, $value ) = @_;
213
Removed:
croak "$name topic level is required"
214
Removed:
if !defined $value || $value eq q{};
218
Added:
my ( $name, $value ) = @_;
219
Added:
croak "$name topic level is required"
220
Added:
if !defined $value || $value eq q{};
215
221
216
Removed:
croak "$name topic level must not contain '/', '+', '#', or NUL: $value"
217
Removed:
if $value =~ m{[\0/+\#]};
222
Added:
croak "$name topic level must not contain '/', '+', '#', or NUL: $value"
223
Added:
if $value =~ m{[\0/+\#]};
218
224
219
Removed:
return $value;
225
Added:
return $value;
220
226
}
221
227
222
228
sub _utc_timestamp {
223
Removed:
return strftime '%Y-%m-%dT%H:%M:%SZ', gmtime;
229
Added:
return strftime '%Y-%m-%dT%H:%M:%SZ', gmtime;
224
230
}
225
231
226
232
1;
roles/daq-hub/bin/fapg-daq-hub-status
@@ -12,34 +12,36 @@
12
12
use Sys::Hostname qw(hostname);
13
13
14
14
my $publish_interval = positive_number_from_env( HUB_STATUS_INTERVAL => 5 );
15
Removed:
my $node = $ENV{FAPG_DAQ_HUB_NODE} // $ENV{FAPG_DAQ_NODE} // $ENV{DAQ_NODE} // hostname();
15
Added:
my $node = $ENV{FAPG_DAQ_HUB_NODE} // $ENV{FAPG_DAQ_NODE} // $ENV{DAQ_NODE}
16
Added:
// hostname();
16
17
17
18
while (1) {
18
Removed:
eval {
19
Removed:
publish_status_mqtt(
20
Removed:
{
21
Removed:
probe => 'hub',
22
Removed:
node => $node,
23
Removed:
status => 'ok',
24
Removed:
message => 'DAQ hub MQTT broker reachable',
25
Removed:
},
26
Removed:
source => 'fapg-daq-hub',
27
Removed:
);
28
Removed:
1;
29
Removed:
} or do {
30
Removed:
my $error = $@ || 'unknown error';
31
Removed:
warn "MQTT hub status publish error: $error";
32
Removed:
};
19
Added:
eval {
20
Added:
publish_status_mqtt(
21
Added:
{ probe => 'hub',
22
Added:
node => $node,
23
Added:
status => 'ok',
24
Added:
message => 'DAQ hub MQTT broker reachable',
25
Added:
},
26
Added:
source => 'fapg-daq-hub',
27
Added:
);
28
Added:
1;
29
Added:
} or do {
30
Added:
my $error = $@ || 'unknown error';
31
Added:
warn "MQTT hub status publish error: $error";
32
Added:
};
33
33
34
Removed:
sleep $publish_interval;
34
Added:
sleep $publish_interval;
35
35
}
36
36
37
37
sub positive_number_from_env {
38
Removed:
my ( $name, $default ) = @_;
39
Removed:
my $value = $ENV{$name};
38
Added:
my ( $name, $default ) = @_;
39
Added:
my $value = $ENV{$name};
40
40
41
Removed:
return $default
42
Removed:
if !defined $value || $value !~ /\A(?:\d+(?:\.\d*)?|\.\d+)\z/ || $value <= 0;
41
Added:
return $default
42
Added:
if !defined $value
43
Added:
|| $value !~ /\A(?:\d+(?:\.\d*)?|\.\d+)\z/
44
Added:
|| $value <= 0;
43
45
44
Removed:
return 0 + $value;
46
Added:
return 0 + $value;
45
47
}
roles/daq-node/bin/fapg-daq-node
@@ -10,11 +10,11 @@
10
10
use lib "$Bin/../lib";
11
11
12
12
use FAPG::DAQ::EZO::USB qw(
13
Removed:
discover_ezo_usb_device
14
Removed:
normalize_probe_type
15
Removed:
open_ezo_usb
16
Removed:
read_probe_info
17
Removed:
read_once
13
Added:
discover_ezo_usb_device
14
Added:
normalize_probe_type
15
Added:
open_ezo_usb
16
Added:
read_probe_info
17
Added:
read_once
18
18
);
19
19
20
20
use FAPG::DAQ::MQTT qw(publish_mqtt publish_status_mqtt);
@@ -33,136 +33,138 @@
33
33
my ( $serial, $info, $device, $probe );
34
34
35
35
while (1) {
36
Removed:
if ( !defined $serial || !defined $info ) {
37
Removed:
( $serial, $info ) = eval {
38
Removed:
$device = serial_device();
39
Removed:
connect_probe($device);
40
Removed:
};
36
Added:
if ( !defined $serial || !defined $info ) {
37
Added:
( $serial, $info ) = eval {
38
Added:
$device = serial_device();
39
Added:
connect_probe($device);
40
Added:
};
41
41
42
Removed:
if ($@) {
43
Removed:
my $error = $@;
44
Removed:
warn "DAQ probe connection error: $error";
45
Removed:
publish_probe_status(
46
Removed:
probe => $probe // $expected_probe,
47
Removed:
status => 'probe_error',
48
Removed:
error => $error,
49
Removed:
device => $device,
50
Removed:
message => 'Unable to connect to or interrogate EZO probe',
51
Removed:
);
52
Removed:
close_serial($serial);
53
Removed:
( $serial, $info, $device, $probe ) = ();
54
Removed:
sleep $retry_delay;
55
Removed:
next;
42
Added:
if ($@) {
43
Added:
my $error = $@;
44
Added:
warn "DAQ probe connection error: $error";
45
Added:
publish_probe_status(
46
Added:
probe => $probe // $expected_probe,
47
Added:
status => 'probe_error',
48
Added:
error => $error,
49
Added:
device => $device,
50
Added:
message => 'Unable to connect to or interrogate EZO probe',
51
Added:
);
52
Added:
close_serial($serial);
53
Added:
( $serial, $info, $device, $probe ) = ();
54
Added:
sleep $retry_delay;
55
Added:
next;
56
Added:
}
57
Added:
58
Added:
$probe = $info->{probe};
56
59
}
57
60
58
Removed:
$probe = $info->{probe};
59
Removed:
}
61
Added:
my $reading = eval { read_once( $serial, probe => $info->{probe} ) };
60
62
61
Removed:
my $reading = eval { read_once( $serial, probe => $info->{probe} ) };
63
Added:
if ($@) {
64
Added:
my $error = $@;
65
Added:
warn "DAQ probe read error, reconnecting to $device: $error";
66
Added:
publish_probe_status(
67
Added:
probe => $probe // $expected_probe,
68
Added:
status => 'probe_error',
69
Added:
error => $error,
70
Added:
device => $device,
71
Added:
message => 'EZO probe did not return a valid reading',
72
Added:
);
73
Added:
close_serial($serial);
74
Added:
( $serial, $info, $device, $probe ) = ();
75
Added:
sleep $retry_delay;
76
Added:
next;
77
Added:
}
62
78
63
Removed:
if ($@) {
64
Removed:
my $error = $@;
65
Removed:
warn "DAQ probe read error, reconnecting to $device: $error";
79
Added:
say "$reading->{probe}: $reading->{value} $reading->{unit}";
80
Added:
81
Added:
eval { publish_mqtt($reading) } or warn "MQTT publish error: $@";
66
82
publish_probe_status(
67
Removed:
probe => $probe // $expected_probe,
68
Removed:
status => 'probe_error',
69
Removed:
error => $error,
70
Removed:
device => $device,
71
Removed:
message => 'EZO probe did not return a valid reading',
83
Added:
probe => $reading->{probe},
84
Added:
status => 'ok',
85
Added:
device => $device,
86
Added:
message => 'Probe reading valid',
72
87
);
73
Removed:
close_serial($serial);
74
Removed:
( $serial, $info, $device, $probe ) = ();
75
Removed:
sleep $retry_delay;
76
Removed:
next;
77
Removed:
}
78
88
79
Removed:
say "$reading->{probe}: $reading->{value} $reading->{unit}";
80
Removed:
81
Removed:
eval { publish_mqtt($reading) } or warn "MQTT publish error: $@";
82
Removed:
publish_probe_status(
83
Removed:
probe => $reading->{probe},
84
Removed:
status => 'ok',
85
Removed:
device => $device,
86
Removed:
message => 'Probe reading valid',
87
Removed:
);
88
Removed:
89
Removed:
sleep $read_interval;
89
Added:
sleep $read_interval;
90
90
}
91
91
92
92
sub connect_probe {
93
Removed:
my ($device) = @_;
93
Added:
my ($device) = @_;
94
94
95
Removed:
my $serial = open_ezo_usb( device => $device );
96
Removed:
my $info = eval { read_probe_info($serial) };
95
Added:
my $serial = open_ezo_usb( device => $device );
96
Added:
my $info = eval { read_probe_info($serial) };
97
97
98
Removed:
if ($@) {
99
Removed:
my $error = $@;
100
Removed:
close_serial($serial);
101
Removed:
die $error;
102
Removed:
}
98
Added:
if ($@) {
99
Added:
my $error = $@;
100
Added:
close_serial($serial);
101
Added:
die $error;
102
Added:
}
103
103
104
Removed:
print <<"EOF";
104
Added:
print <<"EOF";
105
105
Device $device
106
106
Probe $info->{probe}
107
107
Firmware $info->{firmware}
108
108
Unit $info->{unit}
109
109
EOF
110
110
111
Removed:
return ( $serial, $info );
111
Added:
return ( $serial, $info );
112
112
}
113
113
114
114
sub serial_device {
115
Removed:
return $ENV{EZO_SERIAL_DEVICE}
116
Removed:
if defined $ENV{EZO_SERIAL_DEVICE} && $ENV{EZO_SERIAL_DEVICE} ne q{};
115
Added:
return $ENV{EZO_SERIAL_DEVICE}
116
Added:
if defined $ENV{EZO_SERIAL_DEVICE} && $ENV{EZO_SERIAL_DEVICE} ne q{};
117
117
118
Removed:
if ( defined $ENV{SERIAL_DEVICE} && $ENV{SERIAL_DEVICE} ne q{} ) {
119
Removed:
return $ENV{SERIAL_DEVICE} if -e $ENV{SERIAL_DEVICE};
120
Removed:
warn
121
Removed:
"Configured SERIAL_DEVICE=$ENV{SERIAL_DEVICE} does not exist; attempting USB serial discovery\n";
122
Removed:
}
118
Added:
if ( defined $ENV{SERIAL_DEVICE} && $ENV{SERIAL_DEVICE} ne q{} ) {
119
Added:
return $ENV{SERIAL_DEVICE} if -e $ENV{SERIAL_DEVICE};
120
Added:
warn
121
Added:
"Configured SERIAL_DEVICE=$ENV{SERIAL_DEVICE} does not exist; attempting USB serial discovery\n";
122
Added:
}
123
123
124
Removed:
return discover_ezo_usb_device();
124
Added:
return discover_ezo_usb_device();
125
125
}
126
126
127
127
sub expected_probe {
128
Removed:
for my $name (qw(FAPG_DAQ_PROBE EZO_PROBE PROBE_TYPE)) {
129
Removed:
next if !defined $ENV{$name} || $ENV{$name} eq q{};
130
Removed:
return normalize_probe_type( $ENV{$name} );
131
Removed:
}
128
Added:
for my $name (qw(FAPG_DAQ_PROBE EZO_PROBE PROBE_TYPE)) {
129
Added:
next if !defined $ENV{$name} || $ENV{$name} eq q{};
130
Added:
return normalize_probe_type( $ENV{$name} );
131
Added:
}
132
132
133
Removed:
my $hostname = hostname();
134
Removed:
return $1 if $hostname =~ /(?:\A|-)(ph|do|ec|orp)(?:-|\z)/;
133
Added:
my $hostname = hostname();
134
Added:
return $1 if $hostname =~ /(?:\A|-)(ph|do|ec|orp)(?:-|\z)/;
135
135
136
Removed:
return undef;
136
Added:
return undef;
137
137
}
138
138
139
139
sub publish_probe_status {
140
Removed:
my (%status) = @_;
140
Added:
my (%status) = @_;
141
141
142
Removed:
if ( !defined $status{probe} || $status{probe} eq q{} ) {
143
Removed:
warn "DAQ status not published because probe type is unknown\n";
144
Removed:
return;
145
Removed:
}
142
Added:
if ( !defined $status{probe} || $status{probe} eq q{} ) {
143
Added:
warn "DAQ status not published because probe type is unknown\n";
144
Added:
return;
145
Added:
}
146
146
147
Removed:
eval { publish_status_mqtt( \%status ); 1 }
148
Removed:
or warn "MQTT status publish error: $@";
147
Added:
eval { publish_status_mqtt( \%status ); 1 }
148
Added:
or warn "MQTT status publish error: $@";
149
149
}
150
150
151
151
sub close_serial {
152
Removed:
my ($serial) = @_;
152
Added:
my ($serial) = @_;
153
153
154
Removed:
return if !defined $serial || !$serial->can('close');
154
Added:
return if !defined $serial || !$serial->can('close');
155
155
156
Removed:
eval { $serial->close; 1 }
157
Removed:
or warn "DAQ serial close error: $@";
156
Added:
eval { $serial->close; 1 }
157
Added:
or warn "DAQ serial close error: $@";
158
158
}
159
159
160
160
sub positive_number_from_env {
161
Removed:
my ( $name, $default ) = @_;
162
Removed:
my $value = $ENV{$name};
161
Added:
my ( $name, $default ) = @_;
162
Added:
my $value = $ENV{$name};
163
163
164
Removed:
return $default
165
Removed:
if !defined $value || $value !~ /\A(?:\d+(?:\.\d*)?|\.\d+)\z/ || $value <= 0;
164
Added:
return $default
165
Added:
if !defined $value
166
Added:
|| $value !~ /\A(?:\d+(?:\.\d*)?|\.\d+)\z/
167
Added:
|| $value <= 0;
166
168
167
Removed:
return 0 + $value;
169
Added:
return 0 + $value;
168
170
}
roles/daq-node/lib/FAPG/DAQ/EZO/USB.pm
@@ -11,16 +11,16 @@
11
11
use Time::HiRes qw(time sleep);
12
12
13
13
our @EXPORT_OK = qw(
14
Removed:
discover_ezo_usb_device
15
Removed:
open_ezo_usb
16
Removed:
drain_serial
17
Removed:
ezo_command
18
Removed:
read_probe_info
19
Removed:
detect_probe_type
20
Removed:
read_once
21
Removed:
parse_reading
22
Removed:
normalize_probe_type
23
Removed:
unit_for_probe_type
14
Added:
discover_ezo_usb_device
15
Added:
open_ezo_usb
16
Added:
drain_serial
17
Added:
ezo_command
18
Added:
read_probe_info
19
Added:
detect_probe_type
20
Added:
read_once
21
Added:
parse_reading
22
Added:
normalize_probe_type
23
Added:
unit_for_probe_type
24
24
);
25
25
26
26
my $DEFAULT_BAUDRATE = 9_600;
@@ -28,216 +28,221 @@
28
28
my $EOL = "\r";
29
29
30
30
sub discover_ezo_usb_device {
31
Removed:
my %opt = @_;
31
Added:
my %opt = @_;
32
32
33
Removed:
my $by_id_dir = $opt{by_id_dir} // '/dev/serial/by-id';
34
Removed:
my $preferred_pattern = $opt{preferred_pattern} // qr/UART/i;
35
Removed:
my @fallback_patterns =
36
Removed:
exists $opt{fallback_patterns}
37
Removed:
? $opt{fallback_patterns}->@*
38
Removed:
: qw(/dev/ttyUSB* /dev/ttyACM*);
33
Added:
my $by_id_dir = $opt{by_id_dir} // '/dev/serial/by-id';
34
Added:
my $preferred_pattern = $opt{preferred_pattern} // qr/UART/i;
35
Added:
my @fallback_patterns
36
Added:
= exists $opt{fallback_patterns}
37
Added:
? $opt{fallback_patterns}->@*
38
Added:
: qw(/dev/ttyUSB* /dev/ttyACM*);
39
39
40
Removed:
my @candidates;
40
Added:
my @candidates;
41
41
42
Removed:
if ( opendir my $dh, $by_id_dir ) {
43
Removed:
my @by_id = map { "$by_id_dir/$_" }
44
Removed:
sort grep { $_ !~ /\A\.\.?\z/ && -e "$by_id_dir/$_" } readdir $dh;
42
Added:
if ( opendir my $dh, $by_id_dir ) {
43
Added:
my @by_id = map {"$by_id_dir/$_"}
44
Added:
sort grep { $_ !~ /\A\.\.?\z/ && -e "$by_id_dir/$_" } readdir $dh;
45
45
46
Removed:
push @candidates, grep { /$preferred_pattern/ } @by_id;
47
Removed:
push @candidates, grep { $_ !~ /$preferred_pattern/ } @by_id;
48
Removed:
}
46
Added:
push @candidates, grep {/$preferred_pattern/} @by_id;
47
Added:
push @candidates, grep { $_ !~ /$preferred_pattern/ } @by_id;
48
Added:
}
49
49
50
Removed:
for my $pattern (@fallback_patterns) {
51
Removed:
push @candidates, sort grep { -e $_ && !-d $_ } glob $pattern;
52
Removed:
}
50
Added:
for my $pattern (@fallback_patterns) {
51
Added:
push @candidates, sort grep { -e $_ && !-d $_ } glob $pattern;
52
Added:
}
53
53
54
Removed:
my %seen;
55
Removed:
@candidates = grep { !$seen{$_}++ } @candidates;
54
Added:
my %seen;
55
Added:
@candidates = grep { !$seen{$_}++ } @candidates;
56
56
57
Removed:
return $candidates[0] if @candidates;
57
Added:
return $candidates[0] if @candidates;
58
58
59
Removed:
croak "cannot discover EZO USB serial device under $by_id_dir, /dev/ttyUSB*, or /dev/ttyACM*";
59
Added:
croak
60
Added:
"cannot discover EZO USB serial device under $by_id_dir, /dev/ttyUSB*, or /dev/ttyACM*";
60
61
}
61
62
62
63
sub open_ezo_usb {
63
Removed:
my %opt = @_;
64
Removed:
my $device = $opt{device} // croak 'open_ezo_usb requires device => "/dev/tty..."';
65
Removed:
my $baudrate = $opt{baudrate} // $DEFAULT_BAUDRATE;
66
Removed:
my $timeout_ms = $opt{timeout_ms} // 250;
64
Added:
my %opt = @_;
65
Added:
my $device = $opt{device}
66
Added:
// croak 'open_ezo_usb requires device => "/dev/tty..."';
67
Added:
my $baudrate = $opt{baudrate} // $DEFAULT_BAUDRATE;
68
Added:
my $timeout_ms = $opt{timeout_ms} // 250;
67
69
68
Removed:
require Device::SerialPort;
70
Added:
require Device::SerialPort;
69
71
70
Removed:
my $port = Device::SerialPort->new($device)
71
Removed:
or croak "cannot open serial device $device";
72
Added:
my $port = Device::SerialPort->new($device)
73
Added:
or croak "cannot open serial device $device";
72
74
73
Removed:
$port->baudrate($baudrate) or croak "cannot set baudrate on $device";
74
Removed:
$port->databits(8) or croak "cannot set databits on $device";
75
Removed:
$port->parity('none') or croak "cannot set parity on $device";
76
Removed:
$port->stopbits(1) or croak "cannot set stopbits on $device";
77
Removed:
$port->handshake('none') or croak "cannot set handshake on $device";
78
Removed:
$port->read_char_time(0);
79
Removed:
$port->read_const_time($timeout_ms);
80
Removed:
$port->write_settings or croak "cannot apply serial settings on $device";
75
Added:
$port->baudrate($baudrate) or croak "cannot set baudrate on $device";
76
Added:
$port->databits(8) or croak "cannot set databits on $device";
77
Added:
$port->parity('none') or croak "cannot set parity on $device";
78
Added:
$port->stopbits(1) or croak "cannot set stopbits on $device";
79
Added:
$port->handshake('none') or croak "cannot set handshake on $device";
80
Added:
$port->read_char_time(0);
81
Added:
$port->read_const_time($timeout_ms);
82
Added:
$port->write_settings or croak "cannot apply serial settings on $device";
81
83
82
Removed:
return $port;
84
Added:
return $port;
83
85
}
84
86
85
87
sub drain_serial {
86
Removed:
my ( $io, %opt ) = @_;
87
Removed:
my $max_s = $opt{max_s} // 0.20;
88
Removed:
my $chunk_sz = $opt{chunk_sz} // 255;
89
Removed:
my $deadline = time + $max_s;
90
Removed:
my $drained = q{};
88
Added:
my ( $io, %opt ) = @_;
89
Added:
my $max_s = $opt{max_s} // 0.20;
90
Added:
my $chunk_sz = $opt{chunk_sz} // 255;
91
Added:
my $deadline = time + $max_s;
92
Added:
my $drained = q{};
91
93
92
Removed:
while ( time < $deadline ) {
93
Removed:
my ( $count, $chunk ) = $io->read($chunk_sz);
94
Removed:
last if not defined $count || $count == 0;
95
Removed:
$drained .= $chunk // q{};
96
Removed:
}
94
Added:
while ( time < $deadline ) {
95
Added:
my ( $count, $chunk ) = $io->read($chunk_sz);
96
Added:
last if not defined $count || $count == 0;
97
Added:
$drained .= $chunk // q{};
98
Added:
}
97
99
98
Removed:
return $drained;
100
Added:
return $drained;
99
101
}
100
102
101
103
sub ezo_command {
102
Removed:
my ( $io, $command, %opt ) = @_;
103
Removed:
croak 'command must not contain carriage returns or newlines'
104
Removed:
if $command =~ /[\r\n]/;
104
Added:
my ( $io, $command, %opt ) = @_;
105
Added:
croak 'command must not contain carriage returns or newlines'
106
Added:
if $command =~ /[\r\n]/;
105
107
106
Removed:
drain_serial( $io, max_s => ( $opt{drain_s} // 0.05 ) )
107
Removed:
if $opt{drain} // 1;
108
Added:
drain_serial( $io, max_s => ( $opt{drain_s} // 0.05 ) )
109
Added:
if $opt{drain} // 1;
108
110
109
Removed:
my $bytes = $command . $EOL;
110
Removed:
my $wrote = $io->write($bytes);
111
Added:
my $bytes = $command . $EOL;
112
Added:
my $wrote = $io->write($bytes);
111
113
112
Removed:
croak "short write to EZO device: wrote $wrote of " . length($bytes) . ' bytes'
113
Removed:
if not defined $wrote || $wrote != length($bytes);
114
Added:
croak "short write to EZO device: wrote $wrote of "
115
Added:
. length($bytes)
116
Added:
. ' bytes'
117
Added:
if not defined $wrote || $wrote != length($bytes);
114
118
115
Removed:
my $response = _read_ezo_line( $io, timeout_s => ( $opt{timeout_s} // $DEFAULT_TIMEOUT_S ), );
119
Added:
my $response = _read_ezo_line( $io,
120
Added:
timeout_s => ( $opt{timeout_s} // $DEFAULT_TIMEOUT_S ), );
116
121
117
Removed:
croak "EZO command '$command' returned an error: $response"
118
Removed:
if $response =~ /\A\?(?:ERROR|ER)\z/i;
122
Added:
croak "EZO command '$command' returned an error: $response"
123
Added:
if $response =~ /\A\?(?:ERROR|ER)\z/i;
119
124
120
Removed:
return $response;
125
Added:
return $response;
121
126
}
122
127
123
128
sub read_probe_info {
124
Removed:
my ( $io, %opt ) = @_;
125
Removed:
my $raw = ezo_command( $io, 'I', %opt );
129
Added:
my ( $io, %opt ) = @_;
130
Added:
my $raw = ezo_command( $io, 'I', %opt );
126
131
127
Removed:
my ( undef, $device, $firmware ) = split qr{,}, $raw, 3;
132
Added:
my ( undef, $device, $firmware ) = split qr{,}, $raw, 3;
128
133
129
Removed:
croak "unexpected EZO info response: $raw"
130
Removed:
if not defined $device || $raw !~ /\A\?I,/;
134
Added:
croak "unexpected EZO info response: $raw"
135
Added:
if not defined $device || $raw !~ /\A\?I,/;
131
136
132
Removed:
my $probe = normalize_probe_type($device);
137
Added:
my $probe = normalize_probe_type($device);
133
138
134
Removed:
return {
135
Removed:
raw => $raw,
136
Removed:
device => $device,
137
Removed:
firmware => $firmware,
138
Removed:
probe => $probe,
139
Removed:
unit => unit_for_probe_type($probe),
140
Removed:
};
139
Added:
return {
140
Added:
raw => $raw,
141
Added:
device => $device,
142
Added:
firmware => $firmware,
143
Added:
probe => $probe,
144
Added:
unit => unit_for_probe_type($probe),
145
Added:
};
141
146
}
142
147
143
148
sub detect_probe_type {
144
Removed:
my ( $io, %opt ) = @_;
145
Removed:
return read_probe_info( $io, %opt )->{probe};
149
Added:
my ( $io, %opt ) = @_;
150
Added:
return read_probe_info( $io, %opt )->{probe};
146
151
}
147
152
148
153
sub read_once {
149
Removed:
my ( $io, %opt ) = @_;
150
Removed:
my $probe = $opt{probe};
154
Added:
my ( $io, %opt ) = @_;
155
Added:
my $probe = $opt{probe};
151
156
152
Removed:
$probe = detect_probe_type( $io, %opt ) if not defined $probe;
153
Removed:
$probe = normalize_probe_type($probe);
157
Added:
$probe = detect_probe_type( $io, %opt ) if not defined $probe;
158
Added:
$probe = normalize_probe_type($probe);
154
159
155
Removed:
my $raw = ezo_command( $io, 'R', %opt );
156
Removed:
my $parsed = parse_reading( $raw, probe => $probe );
160
Added:
my $raw = ezo_command( $io, 'R', %opt );
161
Added:
my $parsed = parse_reading( $raw, probe => $probe );
157
162
158
Removed:
return {
159
Removed:
probe => $probe,
160
Removed:
unit => unit_for_probe_type($probe),
161
Removed:
raw => $raw,
162
Removed:
value => $parsed->{value},
163
Removed:
values => $parsed->{values},
164
Removed:
};
163
Added:
return {
164
Added:
probe => $probe,
165
Added:
unit => unit_for_probe_type($probe),
166
Added:
raw => $raw,
167
Added:
value => $parsed->{value},
168
Added:
values => $parsed->{values},
169
Added:
};
165
170
}
166
171
167
172
sub parse_reading {
168
Removed:
my ( $raw, %opt ) = @_;
169
Removed:
croak 'empty EZO reading'
170
Removed:
if not defined $raw || $raw eq q{};
173
Added:
my ( $raw, %opt ) = @_;
174
Added:
croak 'empty EZO reading'
175
Added:
if not defined $raw || $raw eq q{};
171
176
172
Removed:
my @fields = split /,/, $raw;
173
Removed:
my @values;
177
Added:
my @fields = split /,/, $raw;
178
Added:
my @values;
174
179
175
Removed:
for my $field (@fields) {
176
Removed:
croak "non-numeric EZO reading field '$field' in '$raw'"
177
Removed:
if $field !~ /\A[+-]?(?:\d+(?:\.\d*)?|\.\d+)\z/;
180
Added:
for my $field (@fields) {
181
Added:
croak "non-numeric EZO reading field '$field' in '$raw'"
182
Added:
if $field !~ /\A[+-]?(?:\d+(?:\.\d*)?|\.\d+)\z/;
178
183
179
Removed:
push @values, 0 + $field;
180
Removed:
}
184
Added:
push @values, 0 + $field;
185
Added:
}
181
186
182
Removed:
return {
183
Removed:
raw => $raw,
184
Removed:
value => $values[0],
185
Removed:
values => \@values,
186
Removed:
};
187
Added:
return {
188
Added:
raw => $raw,
189
Added:
value => $values[0],
190
Added:
values => \@values,
191
Added:
};
187
192
}
188
193
189
194
sub normalize_probe_type {
190
Removed:
my ($probe) = @_;
191
Removed:
croak 'probe type is required'
192
Removed:
if not defined $probe || $probe eq q{};
195
Added:
my ($probe) = @_;
196
Added:
croak 'probe type is required'
197
Added:
if not defined $probe || $probe eq q{};
193
198
194
Removed:
my $p = lc $probe;
195
Removed:
$p =~ s/[\s_-]+//g;
199
Added:
my $p = lc $probe;
200
Added:
$p =~ s/[\s_-]+//g;
196
201
197
Removed:
return
198
Removed:
( $p eq 'ph' ) ? 'ph'
199
Removed:
: ( $p eq 'do' || $p eq 'dissolvedoxygen' ) ? 'do'
200
Removed:
: ( $p eq 'ec' || $p eq 'conductivity' ) ? 'ec'
201
Removed:
: ( $p eq 'orp' ) ? 'orp'
202
Removed:
: croak "unknown EZO probe type '$probe'";
202
Added:
return
203
Added:
( $p eq 'ph' ) ? 'ph'
204
Added:
: ( $p eq 'do' || $p eq 'dissolvedoxygen' ) ? 'do'
205
Added:
: ( $p eq 'ec' || $p eq 'conductivity' ) ? 'ec'
206
Added:
: ( $p eq 'orp' ) ? 'orp'
207
Added:
: croak "unknown EZO probe type '$probe'";
203
208
}
204
209
205
210
sub unit_for_probe_type {
206
Removed:
my ($probe) = @_;
207
Removed:
my $p = normalize_probe_type($probe);
211
Added:
my ($probe) = @_;
212
Added:
my $p = normalize_probe_type($probe);
208
213
209
Removed:
return 'pH' if $p eq 'ph';
210
Removed:
return 'mg/L' if $p eq 'do';
211
Removed:
return 'mV' if $p eq 'orp';
212
Removed:
return 'uS/cm' if $p eq 'ec';
214
Added:
return 'pH' if $p eq 'ph';
215
Added:
return 'mg/L' if $p eq 'do';
216
Added:
return 'mV' if $p eq 'orp';
217
Added:
return 'uS/cm' if $p eq 'ec';
213
218
214
Removed:
croak "unknown normalized probe type '$p'";
219
Added:
croak "unknown normalized probe type '$p'";
215
220
}
216
221
217
222
sub _read_ezo_line {
218
Removed:
my ( $io, %opt ) = @_;
219
Removed:
my $timeout_s = $opt{timeout_s} // $DEFAULT_TIMEOUT_S;
220
Removed:
my $deadline = time + $timeout_s;
221
Removed:
my $buf = q{};
223
Added:
my ( $io, %opt ) = @_;
224
Added:
my $timeout_s = $opt{timeout_s} // $DEFAULT_TIMEOUT_S;
225
Added:
my $deadline = time + $timeout_s;
226
Added:
my $buf = q{};
222
227
223
Removed:
while ( time < $deadline ) {
224
Removed:
my ( $count, $chunk ) = $io->read(1);
228
Added:
while ( time < $deadline ) {
229
Added:
my ( $count, $chunk ) = $io->read(1);
225
230
226
Removed:
if ( !defined $count || $count == 0 ) {
227
Removed:
sleep 0.01;
228
Removed:
next;
231
Added:
if ( !defined $count || $count == 0 ) {
232
Added:
sleep 0.01;
233
Added:
next;
234
Added:
}
235
Added:
236
Added:
$buf .= $chunk;
237
Added:
last if $buf =~ /\r\z/;
229
238
}
230
239
231
Removed:
$buf .= $chunk;
232
Removed:
last if $buf =~ /\r\z/;
233
Removed:
}
240
Added:
croak "timeout waiting for EZO response after ${timeout_s}s"
241
Added:
if $buf !~ /\r\z/;
234
242
235
Removed:
croak "timeout waiting for EZO response after ${timeout_s}s"
236
Removed:
if $buf !~ /\r\z/;
243
Added:
$buf =~ s/[\r\n]+\z//;
237
244
238
Removed:
$buf =~ s/[\r\n]+\z//;
239
Removed:
240
Removed:
return $buf;
245
Added:
return $buf;
241
246
}
242
247
243
248
1;
roles/daq-node/t/01-ezo-usb.t
@@ -12,118 +12,124 @@
12
12
use lib "${FindBin::Bin}/../lib/";
13
13
14
14
use FAPG::DAQ::EZO::USB qw(
15
Removed:
discover_ezo_usb_device
16
Removed:
ezo_command
17
Removed:
read_probe_info
18
Removed:
detect_probe_type
19
Removed:
read_once
20
Removed:
parse_reading
21
Removed:
normalize_probe_type
22
Removed:
unit_for_probe_type
15
Added:
discover_ezo_usb_device
16
Added:
ezo_command
17
Added:
read_probe_info
18
Added:
detect_probe_type
19
Added:
read_once
20
Added:
parse_reading
21
Added:
normalize_probe_type
22
Added:
unit_for_probe_type
23
23
);
24
24
25
25
{
26
26
27
Removed:
package Local::FakeSerial;
27
Added:
package Local::FakeSerial;
28
28
29
Removed:
use v5.32.1;
30
Removed:
use warnings;
29
Added:
use v5.32.1;
30
Added:
use warnings;
31
31
32
Removed:
sub new {
33
Removed:
my ( $class, @responses ) = @_;
34
Removed:
return bless {
35
Removed:
responses => [ map { $_ =~ /\r\z/ ? $_ : "$_\r" } @responses ],
36
Removed:
rx => q{},
37
Removed:
writes => [],
38
Removed:
}, $class;
39
Removed:
}
32
Added:
sub new {
33
Added:
my ( $class, @responses ) = @_;
34
Added:
return bless {
35
Added:
responses => [ map { $_ =~ /\r\z/ ? $_ : "$_\r" } @responses ],
36
Added:
rx => q{},
37
Added:
writes => [],
38
Added:
}, $class;
39
Added:
}
40
40
41
Removed:
sub write {
42
Removed:
my ( $self, $bytes ) = @_;
43
Removed:
push $self->{writes}->@*, $bytes;
44
Removed:
$self->{rx} .= shift( $self->{responses}->@* ) // q{};
45
Removed:
return length $bytes;
46
Removed:
}
41
Added:
sub write {
42
Added:
my ( $self, $bytes ) = @_;
43
Added:
push $self->{writes}->@*, $bytes;
44
Added:
$self->{rx} .= shift( $self->{responses}->@* ) // q{};
45
Added:
return length $bytes;
46
Added:
}
47
47
48
Removed:
sub read {
49
Removed:
my ( $self, $wanted ) = @_;
50
Removed:
return ( 0, q{} ) if $self->{rx} eq q{};
48
Added:
sub read {
49
Added:
my ( $self, $wanted ) = @_;
50
Added:
return ( 0, q{} ) if $self->{rx} eq q{};
51
51
52
Removed:
my $chunk = substr $self->{rx}, 0, $wanted, q{};
53
Removed:
return ( length($chunk), $chunk );
54
Removed:
}
52
Added:
my $chunk = substr $self->{rx}, 0, $wanted, q{};
53
Added:
return ( length($chunk), $chunk );
54
Added:
}
55
55
56
Removed:
sub writes {
57
Removed:
my ($self) = @_;
58
Removed:
return $self->{writes};
59
Removed:
}
56
Added:
sub writes {
57
Added:
my ($self) = @_;
58
Added:
return $self->{writes};
59
Added:
}
60
60
}
61
61
62
62
subtest 'basic command framing' => sub {
63
Removed:
my $serial = Local::FakeSerial->new('7.12');
63
Added:
my $serial = Local::FakeSerial->new('7.12');
64
64
65
Removed:
is ezo_command( $serial, 'R' ), '7.12', 'returns response without CR';
66
Removed:
is $serial->writes, ["R\r"], 'writes command with CR terminator';
65
Added:
is ezo_command( $serial, 'R' ), '7.12', 'returns response without CR';
66
Added:
is $serial->writes, ["R\r"], 'writes command with CR terminator';
67
67
};
68
68
69
69
subtest 'probe detection from info command' => sub {
70
Removed:
my $serial = Local::FakeSerial->new('?I,pH,2.15');
70
Added:
my $serial = Local::FakeSerial->new('?I,pH,2.15');
71
71
72
Removed:
my $info = read_probe_info($serial);
72
Added:
my $info = read_probe_info($serial);
73
73
74
Removed:
is $info->{raw}, '?I,pH,2.15';
75
Removed:
is $info->{device}, 'pH';
76
Removed:
is $info->{firmware}, '2.15';
77
Removed:
is $info->{probe}, 'ph';
78
Removed:
is $info->{unit}, 'pH';
74
Added:
is $info->{raw}, '?I,pH,2.15';
75
Added:
is $info->{device}, 'pH';
76
Added:
is $info->{firmware}, '2.15';
77
Added:
is $info->{probe}, 'ph';
78
Added:
is $info->{unit}, 'pH';
79
79
80
Removed:
is detect_probe_type( Local::FakeSerial->new('?I,ORP,2.13') ), 'orp';
80
Added:
is detect_probe_type( Local::FakeSerial->new('?I,ORP,2.13') ), 'orp';
81
81
};
82
82
83
83
subtest 'single-value reading' => sub {
84
Removed:
my $serial = Local::FakeSerial->new('8.34');
84
Added:
my $serial = Local::FakeSerial->new('8.34');
85
85
86
Removed:
my $reading = read_once( $serial, probe => 'DO' );
86
Added:
my $reading = read_once( $serial, probe => 'DO' );
87
87
88
Removed:
is $reading->{probe}, 'do';
89
Removed:
is $reading->{unit}, 'mg/L';
90
Removed:
is $reading->{raw}, '8.34';
91
Removed:
is $reading->{value}, 8.34;
92
Removed:
is $reading->{values}, [8.34];
88
Added:
is $reading->{probe}, 'do';
89
Added:
is $reading->{unit}, 'mg/L';
90
Added:
is $reading->{raw}, '8.34';
91
Added:
is $reading->{value}, 8.34;
92
Added:
is $reading->{values}, [8.34];
93
93
};
94
94
95
Removed:
subtest 'EC multi-value reading keeps all values but uses first as primary' => sub {
96
Removed:
my $parsed = parse_reading( '1413,704,0.70,1.001', probe => 'ec' );
95
Added:
subtest 'EC multi-value reading keeps all values but uses first as primary' =>
96
Added:
sub {
97
Added:
my $parsed = parse_reading( '1413,704,0.70,1.001', probe => 'ec' );
97
98
98
Removed:
is $parsed->{value}, 1413;
99
Removed:
is $parsed->{values}, [ 1413, 704, 0.70, 1.001 ];
100
Removed:
};
99
Added:
is $parsed->{value}, 1413;
100
Added:
is $parsed->{values}, [ 1413, 704, 0.70, 1.001 ];
101
Added:
};
101
102
102
103
subtest 'normalization and units' => sub {
103
Removed:
is normalize_probe_type('pH'), 'ph';
104
Removed:
is normalize_probe_type('Dissolved Oxygen'), 'do';
105
Removed:
is normalize_probe_type('ORP'), 'orp';
106
Removed:
is normalize_probe_type('EC'), 'ec';
104
Added:
is normalize_probe_type('pH'), 'ph';
105
Added:
is normalize_probe_type('Dissolved Oxygen'), 'do';
106
Added:
is normalize_probe_type('ORP'), 'orp';
107
Added:
is normalize_probe_type('EC'), 'ec';
107
108
108
Removed:
is unit_for_probe_type('ph'), 'pH';
109
Removed:
is unit_for_probe_type('orp'), 'mV';
110
Removed:
is unit_for_probe_type('ec'), 'uS/cm';
109
Added:
is unit_for_probe_type('ph'), 'pH';
110
Added:
is unit_for_probe_type('orp'), 'mV';
111
Added:
is unit_for_probe_type('ec'), 'uS/cm';
111
112
};
112
113
113
114
subtest 'USB device discovery prefers UART by-id devices' => sub {
114
Removed:
my $tmp = tempdir( CLEANUP => 1 );
115
Removed:
my $by_id = "$tmp/serial/by-id";
115
Added:
my $tmp = tempdir( CLEANUP => 1 );
116
Added:
my $by_id = "$tmp/serial/by-id";
116
117
117
Removed:
make_path($by_id);
118
Added:
make_path($by_id);
118
119
119
Removed:
open my $other, '>', "$by_id/usb-Other_Device" or die "cannot create test device: $!";
120
Removed:
close $other;
120
Added:
open my $other, '>', "$by_id/usb-Other_Device"
121
Added:
or die "cannot create test device: $!";
122
Added:
close $other;
121
123
122
Removed:
open my $uart, '>', "$by_id/usb-FTDI_FT232R_USB_UART_A10" or die "cannot create test device: $!";
123
Removed:
close $uart;
124
Added:
open my $uart, '>', "$by_id/usb-FTDI_FT232R_USB_UART_A10"
125
Added:
or die "cannot create test device: $!";
126
Added:
close $uart;
124
127
125
Removed:
is discover_ezo_usb_device( by_id_dir => $by_id, fallback_patterns => [] ),
126
Removed:
"$by_id/usb-FTDI_FT232R_USB_UART_A10";
128
Added:
is discover_ezo_usb_device(
129
Added:
by_id_dir => $by_id,
130
Added:
fallback_patterns => []
131
Added:
),
132
Added:
"$by_id/usb-FTDI_FT232R_USB_UART_A10";
127
133
};
128
134
129
135
done_testing;
roles/daq-node/t/02-probe.t
@@ -16,7 +16,7 @@
16
16
my $baud = 9600;
17
17
18
18
my $serial = Device::SerialPort->new($port)
19
Removed:
or die "Cannot open serial port: $port";
19
Added:
or die "Cannot open serial port: $port";
20
20
21
21
$serial->baudrate($baud);
22
22
$serial->databits(8);
@@ -40,45 +40,45 @@
40
40
like( $reading, qr/OK/, 'probe returns a single reading' );
41
41
42
42
sub ezo_command {
43
Removed:
my ( $serial, $command, $wait_us ) = @_;
43
Added:
my ( $serial, $command, $wait_us ) = @_;
44
44
45
Removed:
drain_serial($serial);
45
Added:
drain_serial($serial);
46
46
47
Removed:
my $written = $serial->write("$command\r");
48
Removed:
return '' unless defined $written && $written > 0;
47
Added:
my $written = $serial->write("$command\r");
48
Added:
return '' unless defined $written && $written > 0;
49
49
50
Removed:
usleep($wait_us);
50
Added:
usleep($wait_us);
51
51
52
Removed:
my $reply = '';
53
Removed:
while (1) {
54
Removed:
my ( $count, $buffer ) = $serial->read(255);
55
Removed:
last unless $count;
56
Removed:
$reply .= $buffer;
57
Removed:
}
52
Added:
my $reply = '';
53
Added:
while (1) {
54
Added:
my ( $count, $buffer ) = $serial->read(255);
55
Added:
last unless $count;
56
Added:
$reply .= $buffer;
57
Added:
}
58
58
59
Removed:
return clean_reply($reply);
59
Added:
return clean_reply($reply);
60
60
}
61
61
62
62
sub drain_serial {
63
Removed:
my ($serial) = @_;
63
Added:
my ($serial) = @_;
64
64
65
Removed:
while (1) {
66
Removed:
my ( $count, undef ) = $serial->read(255);
67
Removed:
last unless $count;
68
Removed:
}
65
Added:
while (1) {
66
Added:
my ( $count, undef ) = $serial->read(255);
67
Added:
last unless $count;
68
Added:
}
69
69
}
70
70
71
71
sub clean_reply {
72
Removed:
my ($reply) = @_;
72
Added:
my ($reply) = @_;
73
73
74
Removed:
$reply =~ s/\r/\n/g;
75
Removed:
$reply =~ s/\n+/\n/g;
76
Removed:
$reply =~ s/^\n|\n$//g;
74
Added:
$reply =~ s/\r/\n/g;
75
Added:
$reply =~ s/\n+/\n/g;
76
Added:
$reply =~ s/^\n|\n$//g;
77
77
78
Removed:
return $reply;
78
Added:
return $reply;
79
79
}
80
80
81
81
sub printable {
82
Removed:
my ($value) = @_;
83
Removed:
return $value eq '' ? '<no response>' : $value;
82
Added:
my ($value) = @_;
83
Added:
return $value eq '' ? '<no response>' : $value;
84
84
}
roles/daq-node/t/03-mqtt.t
@@ -12,175 +12,172 @@
12
12
use lib "${FindBin::Bin}/../lib/";
13
13
14
14
use FAPG::DAQ::MQTT qw(
15
Removed:
publish_mqtt
16
Removed:
publish_status_mqtt
17
Removed:
mqtt_payload_for_reading
18
Removed:
mqtt_payload_for_status
19
Removed:
mqtt_topic_for_reading
20
Removed:
mqtt_topic_for_status
21
Removed:
mqtt_broker_endpoint
15
Added:
publish_mqtt
16
Added:
publish_status_mqtt
17
Added:
mqtt_payload_for_reading
18
Added:
mqtt_payload_for_status
19
Added:
mqtt_topic_for_reading
20
Added:
mqtt_topic_for_status
21
Added:
mqtt_broker_endpoint
22
22
);
23
23
24
24
{
25
25
26
Removed:
package Local::FakeMQTT;
26
Added:
package Local::FakeMQTT;
27
27
28
Removed:
use v5.32.1;
29
Removed:
use warnings;
28
Added:
use v5.32.1;
29
Added:
use warnings;
30
30
31
Removed:
sub new {
32
Removed:
my ($class) = @_;
33
Removed:
return bless { published => [] }, $class;
34
Removed:
}
31
Added:
sub new {
32
Added:
my ($class) = @_;
33
Added:
return bless { published => [] }, $class;
34
Added:
}
35
35
36
Removed:
sub publish {
37
Removed:
my ( $self, $topic, $payload ) = @_;
38
Removed:
push $self->{published}->@*, [ $topic, $payload ];
39
Removed:
return 1;
40
Removed:
}
36
Added:
sub publish {
37
Added:
my ( $self, $topic, $payload ) = @_;
38
Added:
push $self->{published}->@*, [ $topic, $payload ];
39
Added:
return 1;
40
Added:
}
41
41
42
Removed:
sub published {
43
Removed:
my ($self) = @_;
44
Removed:
return $self->{published};
45
Removed:
}
42
Added:
sub published {
43
Added:
my ($self) = @_;
44
Added:
return $self->{published};
45
Added:
}
46
46
}
47
47
48
48
subtest 'topic shape' => sub {
49
Removed:
my $topic = mqtt_topic_for_reading( { probe => 'pH', value => 7.12, unit => 'pH' },
50
Removed:
node => 'fapg-daq-zero-01', );
49
Added:
my $topic
50
Added:
= mqtt_topic_for_reading(
51
Added:
{ probe => 'pH', value => 7.12, unit => 'pH' },
52
Added:
node => 'fapg-daq-zero-01', );
51
53
52
Removed:
is $topic, 'fapg/daq/ph/fapg-daq-zero-01/reading';
54
Added:
is $topic, 'fapg/daq/ph/fapg-daq-zero-01/reading';
53
55
54
Removed:
my $status_topic =
55
Removed:
mqtt_topic_for_status( { probe => 'pH', status => 'ok' }, node => 'fapg-daq-zero-01', );
56
Added:
my $status_topic
57
Added:
= mqtt_topic_for_status( { probe => 'pH', status => 'ok' },
58
Added:
node => 'fapg-daq-zero-01', );
56
59
57
Removed:
is $status_topic, 'fapg/daq/ph/fapg-daq-zero-01/status';
60
Added:
is $status_topic, 'fapg/daq/ph/fapg-daq-zero-01/status';
58
61
};
59
62
60
63
subtest 'payload shape' => sub {
61
Removed:
my $payload = mqtt_payload_for_reading(
62
Removed:
{
63
Removed:
probe => 'pH',
64
Removed:
value => 7.12,
65
Removed:
unit => 'pH',
66
Removed:
raw => '7.12',
67
Removed:
timestamp => '2026-07-06T10:00:00Z',
68
Removed:
},
69
Removed:
node => 'fapg-daq-zero-01',
70
Removed:
);
64
Added:
my $payload = mqtt_payload_for_reading(
65
Added:
{ probe => 'pH',
66
Added:
value => 7.12,
67
Added:
unit => 'pH',
68
Added:
raw => '7.12',
69
Added:
timestamp => '2026-07-06T10:00:00Z',
70
Added:
},
71
Added:
node => 'fapg-daq-zero-01',
72
Added:
);
71
73
72
Removed:
my $decoded = decode_json($payload);
74
Added:
my $decoded = decode_json($payload);
73
75
74
Removed:
is $decoded->{schema}, 'fapg.daq.reading.v1';
75
Removed:
is $decoded->{timestamp}, '2026-07-06T10:00:00Z';
76
Removed:
is $decoded->{node}, 'fapg-daq-zero-01';
77
Removed:
is $decoded->{probe}, 'ph';
78
Removed:
is $decoded->{value}, 7.12;
79
Removed:
is $decoded->{unit}, 'pH';
80
Removed:
is $decoded->{raw}, '7.12';
81
Removed:
is $decoded->{source}, 'fapg-daq-node';
76
Added:
is $decoded->{schema}, 'fapg.daq.reading.v1';
77
Added:
is $decoded->{timestamp}, '2026-07-06T10:00:00Z';
78
Added:
is $decoded->{node}, 'fapg-daq-zero-01';
79
Added:
is $decoded->{probe}, 'ph';
80
Added:
is $decoded->{value}, 7.12;
81
Added:
is $decoded->{unit}, 'pH';
82
Added:
is $decoded->{raw}, '7.12';
83
Added:
is $decoded->{source}, 'fapg-daq-node';
82
84
};
83
85
84
86
subtest 'status payload shape' => sub {
85
Removed:
my $payload = mqtt_payload_for_status(
86
Removed:
{
87
Removed:
probe => 'pH',
88
Removed:
status => 'probe_error',
89
Removed:
error => 'timeout waiting for EZO response',
90
Removed:
device => '/dev/serial/by-id/usb-FTDI_USB_UART',
91
Removed:
timestamp => '2026-07-06T10:00:05Z',
92
Removed:
},
93
Removed:
node => 'fapg-daq-zero-01',
94
Removed:
);
87
Added:
my $payload = mqtt_payload_for_status(
88
Added:
{ probe => 'pH',
89
Added:
status => 'probe_error',
90
Added:
error => 'timeout waiting for EZO response',
91
Added:
device => '/dev/serial/by-id/usb-FTDI_USB_UART',
92
Added:
timestamp => '2026-07-06T10:00:05Z',
93
Added:
},
94
Added:
node => 'fapg-daq-zero-01',
95
Added:
);
95
96
96
Removed:
my $decoded = decode_json($payload);
97
Added:
my $decoded = decode_json($payload);
97
98
98
Removed:
is $decoded->{schema}, 'fapg.daq.status.v1';
99
Removed:
is $decoded->{timestamp}, '2026-07-06T10:00:05Z';
100
Removed:
is $decoded->{node}, 'fapg-daq-zero-01';
101
Removed:
is $decoded->{probe}, 'ph';
102
Removed:
is $decoded->{status}, 'probe_error';
103
Removed:
is $decoded->{error}, 'timeout waiting for EZO response';
104
Removed:
is $decoded->{device}, '/dev/serial/by-id/usb-FTDI_USB_UART';
105
Removed:
is $decoded->{source}, 'fapg-daq-node';
99
Added:
is $decoded->{schema}, 'fapg.daq.status.v1';
100
Added:
is $decoded->{timestamp}, '2026-07-06T10:00:05Z';
101
Added:
is $decoded->{node}, 'fapg-daq-zero-01';
102
Added:
is $decoded->{probe}, 'ph';
103
Added:
is $decoded->{status}, 'probe_error';
104
Added:
is $decoded->{error}, 'timeout waiting for EZO response';
105
Added:
is $decoded->{device}, '/dev/serial/by-id/usb-FTDI_USB_UART';
106
Added:
is $decoded->{source}, 'fapg-daq-node';
106
107
};
107
108
108
109
subtest 'hub status shape' => sub {
109
Removed:
my $topic = mqtt_topic_for_status(
110
Removed:
{
111
Removed:
probe => 'hub',
112
Removed:
status => 'ok',
113
Removed:
},
114
Removed:
node => 'fapg-daq-five-01',
115
Removed:
);
110
Added:
my $topic = mqtt_topic_for_status(
111
Added:
{ probe => 'hub',
112
Added:
status => 'ok',
113
Added:
},
114
Added:
node => 'fapg-daq-five-01',
115
Added:
);
116
116
117
Removed:
is $topic, 'fapg/daq/hub/fapg-daq-five-01/status';
117
Added:
is $topic, 'fapg/daq/hub/fapg-daq-five-01/status';
118
118
119
Removed:
my $payload = mqtt_payload_for_status(
120
Removed:
{
121
Removed:
probe => 'hub',
122
Removed:
status => 'ok',
123
Removed:
message => 'DAQ hub MQTT broker reachable',
124
Removed:
timestamp => '2026-07-06T10:00:10Z',
125
Removed:
},
126
Removed:
node => 'fapg-daq-five-01',
127
Removed:
source => 'fapg-daq-hub',
128
Removed:
);
119
Added:
my $payload = mqtt_payload_for_status(
120
Added:
{ probe => 'hub',
121
Added:
status => 'ok',
122
Added:
message => 'DAQ hub MQTT broker reachable',
123
Added:
timestamp => '2026-07-06T10:00:10Z',
124
Added:
},
125
Added:
node => 'fapg-daq-five-01',
126
Added:
source => 'fapg-daq-hub',
127
Added:
);
129
128
130
Removed:
my $decoded = decode_json($payload);
129
Added:
my $decoded = decode_json($payload);
131
130
132
Removed:
is $decoded->{schema}, 'fapg.daq.status.v1';
133
Removed:
is $decoded->{timestamp}, '2026-07-06T10:00:10Z';
134
Removed:
is $decoded->{node}, 'fapg-daq-five-01';
135
Removed:
is $decoded->{probe}, 'hub';
136
Removed:
is $decoded->{status}, 'ok';
137
Removed:
is $decoded->{message}, 'DAQ hub MQTT broker reachable';
138
Removed:
is $decoded->{source}, 'fapg-daq-hub';
131
Added:
is $decoded->{schema}, 'fapg.daq.status.v1';
132
Added:
is $decoded->{timestamp}, '2026-07-06T10:00:10Z';
133
Added:
is $decoded->{node}, 'fapg-daq-five-01';
134
Added:
is $decoded->{probe}, 'hub';
135
Added:
is $decoded->{status}, 'ok';
136
Added:
is $decoded->{message}, 'DAQ hub MQTT broker reachable';
137
Added:
is $decoded->{source}, 'fapg-daq-hub';
139
138
};
140
139
141
140
subtest 'publish uses injected client, so unit tests need no broker' => sub {
142
Removed:
my $fake = Local::FakeMQTT->new;
141
Added:
my $fake = Local::FakeMQTT->new;
143
142
144
Removed:
my $result = publish_mqtt(
145
Removed:
{
146
Removed:
probe => 'DO',
147
Removed:
value => 8.34,
148
Removed:
unit => 'mg/L',
149
Removed:
timestamp => '2026-07-06T10:00:00Z',
150
Removed:
},
151
Removed:
node => 'fapg-daq-zero-02',
152
Removed:
client => $fake,
153
Removed:
);
143
Added:
my $result = publish_mqtt(
144
Added:
{ probe => 'DO',
145
Added:
value => 8.34,
146
Added:
unit => 'mg/L',
147
Added:
timestamp => '2026-07-06T10:00:00Z',
148
Added:
},
149
Added:
node => 'fapg-daq-zero-02',
150
Added:
client => $fake,
151
Added:
);
154
152
155
Removed:
is $result->{topic}, 'fapg/daq/do/fapg-daq-zero-02/reading';
156
Removed:
is scalar $fake->published->@*, 1;
157
Removed:
is $fake->published->[0][0], 'fapg/daq/do/fapg-daq-zero-02/reading';
153
Added:
is $result->{topic}, 'fapg/daq/do/fapg-daq-zero-02/reading';
154
Added:
is scalar $fake->published->@*, 1;
155
Added:
is $fake->published->[0][0], 'fapg/daq/do/fapg-daq-zero-02/reading';
158
156
159
Removed:
my $decoded = decode_json( $fake->published->[0][1] );
160
Removed:
is $decoded->{probe}, 'do';
161
Removed:
is $decoded->{value}, 8.34;
157
Added:
my $decoded = decode_json( $fake->published->[0][1] );
158
Added:
is $decoded->{probe}, 'do';
159
Added:
is $decoded->{value}, 8.34;
162
160
163
Removed:
my $status = publish_status_mqtt(
164
Removed:
{
165
Removed:
probe => 'DO',
166
Removed:
status => 'ok',
167
Removed:
},
168
Removed:
node => 'fapg-daq-zero-02',
169
Removed:
client => $fake,
170
Removed:
);
161
Added:
my $status = publish_status_mqtt(
162
Added:
{ probe => 'DO',
163
Added:
status => 'ok',
164
Added:
},
165
Added:
node => 'fapg-daq-zero-02',
166
Added:
client => $fake,
167
Added:
);
171
168
172
Removed:
is $status->{topic}, 'fapg/daq/do/fapg-daq-zero-02/status';
173
Removed:
is scalar $fake->published->@*, 2;
174
Removed:
is $fake->published->[1][0], 'fapg/daq/do/fapg-daq-zero-02/status';
169
Added:
is $status->{topic}, 'fapg/daq/do/fapg-daq-zero-02/status';
170
Added:
is scalar $fake->published->@*, 2;
171
Added:
is $fake->published->[1][0], 'fapg/daq/do/fapg-daq-zero-02/status';
175
172
};
176
173
177
174
subtest 'broker endpoint defaults to the Pi 5 broker' => sub {
178
Removed:
local $ENV{FAPG_MQTT_HOST};
179
Removed:
local $ENV{MQTT_HOST};
180
Removed:
local $ENV{FAPG_MQTT_PORT};
181
Removed:
local $ENV{MQTT_PORT};
175
Added:
local $ENV{FAPG_MQTT_HOST};
176
Added:
local $ENV{MQTT_HOST};
177
Added:
local $ENV{FAPG_MQTT_PORT};
178
Added:
local $ENV{MQTT_PORT};
182
179
183
Removed:
is mqtt_broker_endpoint(), 'fapg-daq-five-01:1883';
180
Added:
is mqtt_broker_endpoint(), 'fapg-daq-five-01:1883';
184
181
};
185
182
186
183
done_testing;
roles/daq-recorder/bin/fapg-daq-mqtt-sqlite
@@ -5,9 +5,9 @@
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;
@@ -18,8 +18,8 @@
18
18
19
19
$| = 1;
20
20
21
Removed:
my $db_path = env( 'FAPG_DAQ_DB', '/var/lib/fapg-daq/readings.sqlite3' );
22
Removed:
my $mqtt_host = required_env('MQTT_HOST');
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
23
my $mqtt_port = env( 'MQTT_PORT', '1883' );
24
24
my $mqtt_username = env( 'MQTT_USERNAME', 'fapg_vps' );
25
25
my $mqtt_password = required_env('MQTT_PASSWORD');
@@ -32,69 +32,68 @@
32
32
$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:
RaiseError => 1,
79
Removed:
PrintError => 0,
80
Removed:
AutoCommit => 1,
81
Removed:
sqlite_unicode => 1,
82
Removed:
},
83
Removed:
);
74
Added:
my $dbh = DBI->connect(
75
Added:
"dbi:SQLite:dbname=$path",
76
Added:
'', '',
77
Added:
{ RaiseError => 1,
78
Added:
PrintError => 0,
79
Added:
AutoCommit => 1,
80
Added:
sqlite_unicode => 1,
81
Added:
},
82
Added:
);
84
83
85
Removed:
$dbh->do('PRAGMA journal_mode = WAL');
86
Removed:
$dbh->do('PRAGMA synchronous = NORMAL');
87
Removed:
$dbh->do('PRAGMA busy_timeout = 5000');
88
Removed:
$dbh->do('PRAGMA foreign_keys = ON');
84
Added:
$dbh->do('PRAGMA journal_mode = WAL');
85
Added:
$dbh->do('PRAGMA synchronous = NORMAL');
86
Added:
$dbh->do('PRAGMA busy_timeout = 5000');
87
Added:
$dbh->do('PRAGMA foreign_keys = ON');
89
88
90
Removed:
return $dbh;
89
Added:
return $dbh;
91
90
}
92
91
93
92
sub init_schema {
94
Removed:
my ($dbh) = @_;
93
Added:
my ($dbh) = @_;
95
94
96
Removed:
$dbh->do(
97
Removed:
q{
95
Added:
$dbh->do(
96
Added:
q{
98
97
CREATE TABLE IF NOT EXISTS readings (
99
98
id INTEGER PRIMARY KEY AUTOINCREMENT,
100
99
received_at TEXT NOT NULL,
@@ -112,31 +111,31 @@
112
111
error TEXT
113
112
)
114
113
}
115
Removed:
);
114
Added:
);
116
115
117
Removed:
$dbh->do(
118
Removed:
q{
116
Added:
$dbh->do(
117
Added:
q{
119
118
CREATE INDEX IF NOT EXISTS readings_timestamp_idx
120
119
ON readings(timestamp)
121
120
}
122
Removed:
);
121
Added:
);
123
122
124
Removed:
$dbh->do(
125
Removed:
q{
123
Added:
$dbh->do(
124
Added:
q{
126
125
CREATE INDEX IF NOT EXISTS readings_probe_node_timestamp_idx
127
126
ON readings(probe, node, timestamp)
128
127
}
129
Removed:
);
128
Added:
);
130
129
131
Removed:
$dbh->do(
132
Removed:
q{
130
Added:
$dbh->do(
131
Added:
q{
133
132
CREATE INDEX IF NOT EXISTS readings_topic_idx
134
133
ON readings(topic)
135
134
}
136
Removed:
);
135
Added:
);
137
136
138
Removed:
$dbh->do(
139
Removed:
q{
137
Added:
$dbh->do(
138
Added:
q{
140
139
CREATE TABLE IF NOT EXISTS node_status (
141
140
id INTEGER PRIMARY KEY AUTOINCREMENT,
142
141
received_at TEXT NOT NULL,
@@ -154,56 +153,57 @@
154
153
valid INTEGER NOT NULL DEFAULT 1
155
154
)
156
155
}
157
Removed:
);
156
Added:
);
158
157
159
Removed:
$dbh->do(
160
Removed:
q{
158
Added:
$dbh->do(
159
Added:
q{
161
160
CREATE INDEX IF NOT EXISTS node_status_probe_node_timestamp_idx
162
161
ON node_status(probe, node, timestamp)
163
162
}
164
Removed:
);
163
Added:
);
165
164
166
Removed:
$dbh->do(
167
Removed:
q{
165
Added:
$dbh->do(
166
Added:
q{
168
167
CREATE INDEX IF NOT EXISTS node_status_topic_idx
169
168
ON node_status(topic)
170
169
}
171
Removed:
);
170
Added:
);
172
171
}
173
172
174
173
sub store_message {
175
Removed:
my ( $dbh, $topic, $message ) = @_;
174
Added:
my ( $dbh, $topic, $message ) = @_;
176
175
177
Removed:
my ( $topic_probe, $topic_node, $topic_kind ) = $topic =~ m{\Afapg/daq/([^/]+)/([^/]+)/([^/]+)\z};
178
Removed:
die "invalid MQTT topic: $topic\n"
179
Removed:
if !defined $topic_kind || $topic_kind !~ /\A(?:reading|status)\z/;
176
Added:
my ( $topic_probe, $topic_node, $topic_kind )
177
Added:
= $topic =~ m{\Afapg/daq/([^/]+)/([^/]+)/([^/]+)\z};
178
Added:
die "invalid MQTT topic: $topic\n"
179
Added:
if !defined $topic_kind || $topic_kind !~ /\A(?:reading|status)\z/;
180
180
181
Removed:
my $valid = 1;
182
Removed:
my $error;
183
Removed:
my $data = eval { decode_json($message) };
181
Added:
my $valid = 1;
182
Added:
my $error;
183
Added:
my $data = eval { decode_json($message) };
184
184
185
Removed:
if ($@) {
186
Removed:
$valid = 0;
187
Removed:
$error = $@;
188
Removed:
chomp $error;
189
Removed:
$data = {};
190
Removed:
} elsif ( ref $data ne 'HASH' ) {
191
Removed:
$valid = 0;
192
Removed:
$error = 'JSON payload is not an object';
193
Removed:
$data = {};
194
Removed:
}
185
Added:
if ($@) {
186
Added:
$valid = 0;
187
Added:
$error = $@;
188
Added:
chomp $error;
189
Added:
$data = {};
190
Added:
} elsif ( ref $data ne 'HASH' ) {
191
Added:
$valid = 0;
192
Added:
$error = 'JSON payload is not an object';
193
Added:
$data = {};
194
Added:
}
195
195
196
Removed:
if ( $topic_kind eq 'status' ) {
197
Removed:
store_status_message( $dbh, $topic, $message, $topic_probe, $topic_node, $data, $valid,
198
Removed:
$error );
199
Removed:
return;
200
Removed:
}
196
Added:
if ( $topic_kind eq 'status' ) {
197
Added:
store_status_message( $dbh, $topic, $message, $topic_probe,
198
Added:
$topic_node, $data, $valid, $error );
199
Added:
return;
200
Added:
}
201
201
202
Removed:
my $value = $data->{value};
203
Removed:
$value = undef if defined $value && !looks_like_number($value);
202
Added:
my $value = $data->{value};
203
Added:
$value = undef if defined $value && !looks_like_number($value);
204
204
205
Removed:
my $sth = $dbh->prepare_cached(
206
Removed:
q{
205
Added:
my $sth = $dbh->prepare_cached(
206
Added:
q{
207
207
INSERT INTO readings (
208
208
received_at,
209
209
topic,
@@ -220,35 +220,41 @@
220
220
error
221
221
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
222
222
}
223
Removed:
);
223
Added:
);
224
224
225
Removed:
$sth->execute(
226
Removed:
utc_now(), $topic, $topic_kind,
227
Removed:
$data->{schema}, $data->{timestamp}, $data->{probe} // $topic_probe,
228
Removed:
$data->{node} // $topic_node, $value, $data->{unit},
229
Removed:
$data->{source}, $message, $valid,
230
Removed:
$error,
231
Removed:
);
225
Added:
$sth->execute(
226
Added:
utc_now(), $topic,
227
Added:
$topic_kind, $data->{schema},
228
Added:
$data->{timestamp}, $data->{probe} // $topic_probe,
229
Added:
$data->{node} // $topic_node, $value,
230
Added:
$data->{unit}, $data->{source},
231
Added:
$message, $valid,
232
Added:
$error,
233
Added:
);
232
234
}
233
235
234
236
sub store_status_message {
235
Removed:
my ( $dbh, $topic, $message, $topic_probe, $topic_node, $data, $valid, $error ) = @_;
237
Added:
my ( $dbh, $topic, $message, $topic_probe, $topic_node, $data, $valid,
238
Added:
$error )
239
Added:
= @_;
236
240
237
Removed:
my $status = $data->{status};
238
Removed:
if ( defined $status ) {
239
Removed:
$status = lc $status;
240
Removed:
if ( $status !~ /\A(?:ok|probe_error)\z/ ) {
241
Removed:
$valid = 0;
242
Removed:
$error =
243
Removed:
defined $error ? "$error; unsupported status: $status" : "unsupported status: $status";
241
Added:
my $status = $data->{status};
242
Added:
if ( defined $status ) {
243
Added:
$status = lc $status;
244
Added:
if ( $status !~ /\A(?:ok|probe_error)\z/ ) {
245
Added:
$valid = 0;
246
Added:
$error
247
Added:
= defined $error
248
Added:
? "$error; unsupported status: $status"
249
Added:
: "unsupported status: $status";
250
Added:
}
251
Added:
} else {
252
Added:
$valid = 0;
253
Added:
$error = defined $error ? "$error; missing status" : 'missing status';
244
254
}
245
Removed:
} else {
246
Removed:
$valid = 0;
247
Removed:
$error = defined $error ? "$error; missing status" : 'missing status';
248
Removed:
}
249
255
250
Removed:
my $sth = $dbh->prepare_cached(
251
Removed:
q{
256
Added:
my $sth = $dbh->prepare_cached(
257
Added:
q{
252
258
INSERT INTO node_status (
253
259
received_at,
254
260
topic,
@@ -265,17 +271,19 @@
265
271
valid
266
272
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
267
273
}
268
Removed:
);
274
Added:
);
269
275
270
Removed:
$sth->execute(
271
Removed:
utc_now(), $topic, $data->{schema},
272
Removed:
$data->{timestamp}, $data->{probe} // $topic_probe, $data->{node} // $topic_node,
273
Removed:
$status, $data->{message}, $data->{error} // $error,
274
Removed:
$data->{device}, $data->{source}, $message,
275
Removed:
$valid,
276
Removed:
);
276
Added:
$sth->execute(
277
Added:
utc_now(), $topic,
278
Added:
$data->{schema}, $data->{timestamp},
279
Added:
$data->{probe} // $topic_probe, $data->{node} // $topic_node,
280
Added:
$status, $data->{message},
281
Added:
$data->{error} // $error, $data->{device},
282
Added:
$data->{source}, $message,
283
Added:
$valid,
284
Added:
);
277
285
}
278
286
279
287
sub utc_now {
280
Removed:
return strftime( '%Y-%m-%dT%H:%M:%SZ', gmtime );
288
Added:
return strftime( '%Y-%m-%dT%H:%M:%SZ', gmtime );
281
289
}
roles/dashboard/lib/FAPG/DAQ/Dashboard.pm
@@ -7,79 +7,75 @@
7
7
8
8
# This method will run once at server start
9
9
sub startup ($self) {
10
Removed:
$self->moniker('dashboard');
10
Added:
$self->moniker('dashboard');
11
11
12
Removed:
# Load configuration from config file
13
Removed:
my $config = $self->plugin('NotYAMLConfig');
12
Added:
# Load configuration from config file
13
Added:
my $config = $self->plugin('NotYAMLConfig');
14
14
15
Removed:
# Configure the application
16
Removed:
$self->secrets( $config->{secrets} );
15
Added:
# Configure the application
16
Added:
$self->secrets( $config->{secrets} );
17
17
18
Removed:
$self->config(
19
Removed:
hypnotoad => {
20
Removed:
listen => ['http://127.0.0.1:3000'],
21
Removed:
pid_file => '/run/fapg-daq-dashboard/hypnotoad.pid',
22
Removed:
workers => 2,
23
Removed:
proxy => 1,
24
Removed:
},
25
Removed:
);
18
Added:
$self->config(
19
Added:
hypnotoad => {
20
Added:
listen => ['http://127.0.0.1:3000'],
21
Added:
pid_file => '/run/fapg-daq-dashboard/hypnotoad.pid',
22
Added:
workers => 2,
23
Added:
proxy => 1,
24
Added:
},
25
Added:
);
26
26
27
Removed:
my $db_path = $ENV{FAPG_DAQ_DB} // $self->home->rel_file('fapg-daq.sqlite3');
27
Added:
my $db_path = $ENV{FAPG_DAQ_DB}
28
Added:
// $self->home->rel_file('fapg-daq.sqlite3');
28
29
29
Removed:
my $sqlite = Mojo::SQLite->new("sqlite:$db_path");
30
Added:
my $sqlite = Mojo::SQLite->new("sqlite:$db_path");
30
31
31
Removed:
$self->helper( sqlite => sub { $sqlite } );
32
Added:
$self->helper( sqlite => sub {$sqlite} );
32
33
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:
node => 'Pi Zero 01',
41
Removed:
},
42
Removed:
{
43
Removed:
key => 'do',
44
Removed:
label => 'Dissolved Oxygen',
45
Removed:
unit => 'mg/L',
46
Removed:
node => 'Pi Zero 03',
47
Removed:
},
48
Removed:
{
49
Removed:
key => 'orp',
50
Removed:
label => 'ORP',
51
Removed:
unit => 'mV',
52
Removed:
node => 'Pi Zero 04',
53
Removed:
},
54
Removed:
{
55
Removed:
key => 'ec',
56
Removed:
label => 'Electrical Conductivity',
57
Removed:
unit => 'µS/cm',
58
Removed:
node => 'Pi Zero 02',
59
Removed:
},
60
Removed:
];
61
Removed:
}
62
Removed:
);
34
Added:
$self->helper(
35
Added:
probes => sub {
36
Added:
return [
37
Added:
{ key => 'ph',
38
Added:
label => 'pH',
39
Added:
unit => 'pH',
40
Added:
node => 'Pi Zero 01',
41
Added:
},
42
Added:
{ key => 'do',
43
Added:
label => 'Dissolved Oxygen',
44
Added:
unit => 'mg/L',
45
Added:
node => 'Pi Zero 03',
46
Added:
},
47
Added:
{ key => 'orp',
48
Added:
label => 'ORP',
49
Added:
unit => 'mV',
50
Added:
node => 'Pi Zero 04',
51
Added:
},
52
Added:
{ key => 'ec',
53
Added:
label => 'Electrical Conductivity',
54
Added:
unit => 'µS/cm',
55
Added:
node => 'Pi Zero 02',
56
Added:
},
57
Added:
];
58
Added:
}
59
Added:
);
63
60
64
Removed:
$self->helper(
65
Removed:
status_items => sub {
66
Removed:
return [
67
Removed:
{
68
Removed:
key => 'hub',
69
Removed:
label => 'DAQ Hub',
70
Removed:
},
71
Removed:
$self->probes->@*,
72
Removed:
];
73
Removed:
}
74
Removed:
);
61
Added:
$self->helper(
62
Added:
status_items => sub {
63
Added:
return [
64
Added:
{ key => 'hub',
65
Added:
label => 'DAQ Hub',
66
Added:
},
67
Added:
$self->probes->@*,
68
Added:
];
69
Added:
}
70
Added:
);
75
71
76
Removed:
my $r = $self->routes;
72
Added:
my $r = $self->routes;
77
73
78
Removed:
$r->get('/')->to('dashboard#index');
79
Removed:
$r->get('/graphs/:probe')->to('dashboard#graph');
74
Added:
$r->get('/')->to('dashboard#index');
75
Added:
$r->get('/graphs/:probe')->to('dashboard#graph');
80
76
81
Removed:
$r->get('/api/readings/:probe')->to('Reading#list');
82
Removed:
$r->get('/api/status/:probe')->to('Reading#status');
77
Added:
$r->get('/api/readings/:probe')->to('Reading#list');
78
Added:
$r->get('/api/status/:probe')->to('Reading#status');
83
79
}
84
80
85
81
1;
roles/dashboard/lib/FAPG/DAQ/Dashboard/Controller/Dashboard.pm
@@ -4,22 +4,22 @@
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
sub graph ($self) {
14
Removed:
my $probe_key = $self->stash('probe');
15
Removed:
my ($probe) = grep { $_->{key} eq $probe_key } $self->probes->@*;
14
Added:
my $probe_key = $self->stash('probe');
15
Added:
my ($probe) = grep { $_->{key} eq $probe_key } $self->probes->@*;
16
16
17
Removed:
return $self->reply->not_found unless defined $probe;
17
Added:
return $self->reply->not_found unless defined $probe;
18
18
19
Removed:
$self->render(
20
Removed:
template => 'dashboard/graph',
21
Removed:
probe => $probe,
22
Removed:
);
19
Added:
$self->render(
20
Added:
template => 'dashboard/graph',
21
Added:
probe => $probe,
22
Added:
);
23
23
}
24
24
25
25
1;
roles/dashboard/lib/FAPG/DAQ/Dashboard/Controller/Reading.pm
@@ -8,70 +8,68 @@
8
8
my $STATUS_REACHABLE_SECONDS = 120;
9
9
10
10
sub list ($self) {
11
Removed:
my $probe = $self->param('probe') // '';
11
Added:
my $probe = $self->param('probe') // '';
12
12
13
Removed:
my %known = map { $_->{key} => 1 } $self->probes->@*;
13
Added:
my %known = map { $_->{key} => 1 } $self->probes->@*;
14
14
15
Removed:
return $self->render(
16
Removed:
status => 404,
17
Removed:
json => {
18
Removed:
error => "Unknown probe type: $probe",
19
Removed:
},
20
Removed:
) unless $known{$probe};
15
Added:
return $self->render(
16
Added:
status => 404,
17
Added:
json => { error => "Unknown probe type: $probe", },
18
Added:
) unless $known{$probe};
21
19
22
Removed:
my $limit = $self->param('limit') // 300;
20
Added:
my $limit = $self->param('limit') // 300;
23
21
24
Removed:
$limit = 300 unless $limit =~ /^\d+$/;
25
Removed:
$limit = 2_000 if $limit > 2_000;
22
Added:
$limit = 300 unless $limit =~ /^\d+$/;
23
Added:
$limit = 2_000 if $limit > 2_000;
26
24
27
Removed:
my @where = ('probe = ?');
28
Removed:
my @bind = ($probe);
29
Removed:
my $since = $self->param('since');
25
Added:
my @where = ('probe = ?');
26
Added:
my @bind = ($probe);
27
Added:
my $since = $self->param('since');
30
28
31
Removed:
if ( defined $since && $since =~ /\A\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}Z\z/ ) {
32
Removed:
push @where, 'COALESCE(received_at, timestamp) >= ?';
33
Removed:
push @bind, $since;
34
Removed:
}
29
Added:
if ( defined $since
30
Added:
&& $since =~ /\A\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}Z\z/ )
31
Added:
{
32
Added:
push @where, 'COALESCE(received_at, timestamp) >= ?';
33
Added:
push @bind, $since;
34
Added:
}
35
35
36
Removed:
my $rows = $self->sqlite->db->query(
37
Removed:
q{
36
Added:
my $rows = $self->sqlite->db->query(
37
Added:
q{
38
38
SELECT timestamp, received_at, node, probe, value, unit
39
39
FROM readings
40
40
WHERE
41
41
}
42
Removed:
. join( "\n AND ", @where ) . q{
42
Added:
. join( "\n AND ", @where ) . q{
43
43
ORDER BY COALESCE(received_at, timestamp) DESC, id DESC
44
44
LIMIT ?
45
45
},
46
Removed:
@bind,
47
Removed:
$limit,
48
Removed:
)->hashes->to_array;
46
Added:
@bind,
47
Added:
$limit,
48
Added:
)->hashes->to_array;
49
49
50
Removed:
$rows = [ reverse $rows->@* ];
50
Added:
$rows = [ reverse $rows->@* ];
51
51
52
Removed:
$self->render(
53
Removed:
json => {
54
Removed:
probe => $probe,
55
Removed:
readings => $rows,
56
Removed:
},
57
Removed:
);
52
Added:
$self->render(
53
Added:
json => {
54
Added:
probe => $probe,
55
Added:
readings => $rows,
56
Added:
},
57
Added:
);
58
58
}
59
59
60
60
sub status ($self) {
61
Removed:
my $probe = $self->param('probe') // '';
61
Added:
my $probe = $self->param('probe') // '';
62
62
63
Removed:
my %known = map { $_->{key} => 1 } $self->status_items->@*;
63
Added:
my %known = map { $_->{key} => 1 } $self->status_items->@*;
64
64
65
Removed:
return $self->render(
66
Removed:
status => 404,
67
Removed:
json => {
68
Removed:
error => "Unknown status item: $probe",
69
Removed:
},
70
Removed:
) unless $known{$probe};
65
Added:
return $self->render(
66
Added:
status => 404,
67
Added:
json => { error => "Unknown status item: $probe", },
68
Added:
) unless $known{$probe};
71
69
72
Removed:
my $row = eval {
73
Removed:
$self->sqlite->db->query(
74
Removed:
q{
70
Added:
my $row = eval {
71
Added:
$self->sqlite->db->query(
72
Added:
q{
75
73
SELECT timestamp, received_at, node, probe, status, message, error, device
76
74
FROM node_status
77
75
WHERE probe = ?
@@ -79,31 +77,33 @@
79
77
ORDER BY COALESCE(timestamp, received_at) DESC, id DESC
80
78
LIMIT 1
81
79
},
82
Removed:
$probe,
83
Removed:
)->hash;
84
Removed:
};
80
Added:
$probe,
81
Added:
)->hash;
82
Added:
};
85
83
86
Removed:
$row = undef if $@;
87
Removed:
mark_unreachable_if_stale($row) if $row;
84
Added:
$row = undef if $@;
85
Added:
mark_unreachable_if_stale($row) if $row;
88
86
89
Removed:
$self->render(
90
Removed:
json => {
91
Removed:
probe => $probe,
92
Removed:
status => $row,
93
Removed:
},
94
Removed:
);
87
Added:
$self->render(
88
Added:
json => {
89
Added:
probe => $probe,
90
Added:
status => $row,
91
Added:
},
92
Added:
);
95
93
}
96
94
97
95
sub mark_unreachable_if_stale ($row) {
98
Removed:
my $received_at = $row->{received_at} // $row->{timestamp};
99
Removed:
return if !defined $received_at || $received_at ge reachable_since();
96
Added:
my $received_at = $row->{received_at} // $row->{timestamp};
97
Added:
return if !defined $received_at || $received_at ge reachable_since();
100
98
101
Removed:
$row->{status} = 'unreachable';
102
Removed:
$row->{message} = "No status received in the last $STATUS_REACHABLE_SECONDS seconds";
99
Added:
$row->{status} = 'unreachable';
100
Added:
$row->{message}
101
Added:
= "No status received in the last $STATUS_REACHABLE_SECONDS seconds";
103
102
}
104
103
105
104
sub reachable_since {
106
Removed:
return strftime( '%Y-%m-%dT%H:%M:%SZ', gmtime( time - $STATUS_REACHABLE_SECONDS ) );
105
Added:
return strftime( '%Y-%m-%dT%H:%M:%SZ',
106
Added:
gmtime( time - $STATUS_REACHABLE_SECONDS ) );
107
107
}
108
108
109
109
1;
roles/dashboard/t/00-homepage.t
@@ -12,51 +12,63 @@
12
12
my $t = test_app();
13
13
14
14
$t->get_ok('/')
15
Removed:
->status_is(200)
16
Removed:
->element_exists('body > header.page-header')
17
Removed:
->element_exists_not('main.site-main > header.page-header')
18
Removed:
->element_exists('input.nav-toggle[type="checkbox"]')
19
Removed:
->text_is( 'a.page-brand[href="/"]', 'FAPG' )
20
Removed:
->element_exists('a.header-status-link[href="/"][data-header-status]')
21
Removed:
->element_exists('a.header-status-link span.header-status-led')
22
Removed:
->element_exists('[data-header-status-item="hub"]')
23
Removed:
->element_exists('[data-header-status-item="ph"]')
24
Removed:
->element_exists('[data-header-status-item="do"]')
25
Removed:
->element_exists('[data-header-status-item="orp"]')
26
Removed:
->element_exists('[data-header-status-item="ec"]')
27
Removed:
->element_exists('nav.probe-nav a[href="/graphs/ph"][data-nav-probe="ph"]')
28
Removed:
->element_exists('nav.probe-nav a[href="/graphs/do"][data-nav-probe="do"]')
29
Removed:
->element_exists('nav.probe-nav a[href="/graphs/orp"][data-nav-probe="orp"]')
30
Removed:
->element_exists('nav.probe-nav a[href="/graphs/ec"][data-nav-probe="ec"]')
31
Removed:
->element_exists('footer.site-footer a[href="https://www.lafermeaquaponique.com/"]')
32
Removed:
->content_like(qr/Copyright © 2026 Marius Peter/)
33
Removed:
->element_exists_not('section#probe-ph[data-panel="ph"]')
34
Removed:
->element_exists_not('section.dashboard-card')
35
Removed:
->content_like(qr/Checking sensors/)
36
Removed:
->content_like(qr/data-overview-status/)
37
Removed:
->content_like(qr/data-sensor-status="ph"/)
38
Removed:
->content_like(qr/Last seen online/)
39
Removed:
->content_like(qr/data-reading-quick-probe="ph"/)
40
Removed:
->content_like(qr/data-reading-quick-probe="do"/)
41
Removed:
->content_like(qr/data-reading-quick-probe="orp"/)
42
Removed:
->content_like(qr/data-reading-quick-probe="ec"/)
43
Removed:
->element_exists('a.reading-quick-link[href="/graphs/ph"] article[data-reading-quick-probe="ph"]')
44
Removed:
->element_exists('a.reading-quick-link[href="/graphs/do"] article[data-reading-quick-probe="do"]')
45
Removed:
->element_exists(
46
Removed:
'a.reading-quick-link[href="/graphs/orp"] article[data-reading-quick-probe="orp"]')
47
Removed:
->element_exists('a.reading-quick-link[href="/graphs/ec"] article[data-reading-quick-probe="ec"]')
48
Removed:
->element_exists_not('article[data-reading-quick-probe="ph"] h3 a')
49
Removed:
->content_like(qr/Dissolved Oxygen/)
50
Removed:
->content_unlike(qr/Hourly/)
51
Removed:
->content_unlike(qr/Daily/)
52
Removed:
->content_unlike(qr/Weekly/)
53
Removed:
->content_unlike(qr/Yearly/)
54
Removed:
->content_unlike(qr/data-smoothing-toggle checked/)
55
Removed:
->content_unlike(qr/smoothing/)
56
Removed:
->content_unlike(qr/Last 10 pH Reading Messages/)
57
Removed:
->content_unlike(qr/Last 10 Dissolved Oxygen Reading Messages/)
58
Removed:
->content_unlike(qr/status-pill-list/)
59
Removed:
->content_unlike(qr/Pi Zero 01/)
60
Removed:
->content_unlike(qr/Connection Status/);
15
Added:
->status_is(200)
16
Added:
->element_exists('body > header.page-header')
17
Added:
->element_exists_not('main.site-main > header.page-header')
18
Added:
->element_exists('input.nav-toggle[type="checkbox"]')
19
Added:
->text_is( 'a.page-brand[href="/"]', 'FAPG' )
20
Added:
->element_exists('a.header-status-link[href="/"][data-header-status]')
21
Added:
->element_exists('a.header-status-link span.header-status-led')
22
Added:
->element_exists('[data-header-status-item="hub"]')
23
Added:
->element_exists('[data-header-status-item="ph"]')
24
Added:
->element_exists('[data-header-status-item="do"]')
25
Added:
->element_exists('[data-header-status-item="orp"]')
26
Added:
->element_exists('[data-header-status-item="ec"]')
27
Added:
->element_exists(
28
Added:
'nav.probe-nav a[href="/graphs/ph"][data-nav-probe="ph"]')
29
Added:
->element_exists(
30
Added:
'nav.probe-nav a[href="/graphs/do"][data-nav-probe="do"]')
31
Added:
->element_exists(
32
Added:
'nav.probe-nav a[href="/graphs/orp"][data-nav-probe="orp"]')
33
Added:
->element_exists(
34
Added:
'nav.probe-nav a[href="/graphs/ec"][data-nav-probe="ec"]')
35
Added:
->element_exists(
36
Added:
'footer.site-footer a[href="https://www.lafermeaquaponique.com/"]')
37
Added:
->content_like(qr/Copyright © 2026 Marius Peter/)
38
Added:
->element_exists_not('section#probe-ph[data-panel="ph"]')
39
Added:
->element_exists_not('section.dashboard-card')
40
Added:
->content_like(qr/Checking sensors/)
41
Added:
->content_like(qr/data-overview-status/)
42
Added:
->content_like(qr/data-sensor-status="ph"/)
43
Added:
->content_like(qr/Last seen online/)
44
Added:
->content_like(qr/data-reading-quick-probe="ph"/)
45
Added:
->content_like(qr/data-reading-quick-probe="do"/)
46
Added:
->content_like(qr/data-reading-quick-probe="orp"/)
47
Added:
->content_like(qr/data-reading-quick-probe="ec"/)
48
Added:
->element_exists(
49
Added:
'a.reading-quick-link[href="/graphs/ph"] article[data-reading-quick-probe="ph"]'
50
Added:
)
51
Added:
->element_exists(
52
Added:
'a.reading-quick-link[href="/graphs/do"] article[data-reading-quick-probe="do"]'
53
Added:
)
54
Added:
->element_exists(
55
Added:
'a.reading-quick-link[href="/graphs/orp"] article[data-reading-quick-probe="orp"]'
56
Added:
)
57
Added:
->element_exists(
58
Added:
'a.reading-quick-link[href="/graphs/ec"] article[data-reading-quick-probe="ec"]'
59
Added:
)
60
Added:
->element_exists_not('article[data-reading-quick-probe="ph"] h3 a')
61
Added:
->content_like(qr/Dissolved Oxygen/)
62
Added:
->content_unlike(qr/Hourly/)
63
Added:
->content_unlike(qr/Daily/)
64
Added:
->content_unlike(qr/Weekly/)
65
Added:
->content_unlike(qr/Yearly/)
66
Added:
->content_unlike(qr/data-smoothing-toggle checked/)
67
Added:
->content_unlike(qr/smoothing/)
68
Added:
->content_unlike(qr/Last 10 pH Reading Messages/)
69
Added:
->content_unlike(qr/Last 10 Dissolved Oxygen Reading Messages/)
70
Added:
->content_unlike(qr/status-pill-list/)
71
Added:
->content_unlike(qr/Pi Zero 01/)
72
Added:
->content_unlike(qr/Connection Status/);
61
73
62
74
done_testing;
roles/dashboard/t/01-graph-page.t
@@ -12,30 +12,35 @@
12
12
my $t = test_app();
13
13
14
14
$t->get_ok('/graphs/ph')
15
Removed:
->status_is(200)
16
Removed:
->text_is( 'a.page-brand[href="/"]', 'FAPG' )
17
Removed:
->element_exists('a.header-status-link[href="/"][data-header-status]')
18
Removed:
->element_exists('a.header-status-link span.header-status-led')
19
Removed:
->element_exists('[data-header-status-item="hub"]')
20
Removed:
->element_exists('[data-header-status-item="ph"]')
21
Removed:
->element_exists('[data-header-status-item="do"]')
22
Removed:
->element_exists('[data-header-status-item="orp"]')
23
Removed:
->element_exists('[data-header-status-item="ec"]')
24
Removed:
->element_exists('nav.probe-nav a[href="/graphs/ph"][data-nav-probe="ph"]')
25
Removed:
->element_exists('nav.probe-nav a[href="/graphs/do"][data-nav-probe="do"]')
26
Removed:
->element_exists('nav.probe-nav a[href="/graphs/orp"][data-nav-probe="orp"]')
27
Removed:
->element_exists('nav.probe-nav a[href="/graphs/ec"][data-nav-probe="ec"]')
28
Removed:
->element_exists('footer.site-footer a[href="https://www.lafermeaquaponique.com/"]')
29
Removed:
->element_exists('section#probe-ph[data-panel="ph"]')
30
Removed:
->element_exists('canvas#chart-ph')
31
Removed:
->element_exists_not('section#probe-do[data-panel="do"]')
32
Removed:
->content_unlike(qr/Sensor Reading Quick Status/)
33
Removed:
->content_like(qr/Hourly/)
34
Removed:
->content_like(qr/Daily/)
35
Removed:
->content_like(qr/Weekly/)
36
Removed:
->content_like(qr/Yearly/)
37
Removed:
->content_like(qr/data-smoothing-toggle checked/)
38
Removed:
->content_like(qr/Latest pH Reading Messages/);
15
Added:
->status_is(200)
16
Added:
->text_is( 'a.page-brand[href="/"]', 'FAPG' )
17
Added:
->element_exists('a.header-status-link[href="/"][data-header-status]')
18
Added:
->element_exists('a.header-status-link span.header-status-led')
19
Added:
->element_exists('[data-header-status-item="hub"]')
20
Added:
->element_exists('[data-header-status-item="ph"]')
21
Added:
->element_exists('[data-header-status-item="do"]')
22
Added:
->element_exists('[data-header-status-item="orp"]')
23
Added:
->element_exists('[data-header-status-item="ec"]')
24
Added:
->element_exists(
25
Added:
'nav.probe-nav a[href="/graphs/ph"][data-nav-probe="ph"]')
26
Added:
->element_exists(
27
Added:
'nav.probe-nav a[href="/graphs/do"][data-nav-probe="do"]')
28
Added:
->element_exists(
29
Added:
'nav.probe-nav a[href="/graphs/orp"][data-nav-probe="orp"]')
30
Added:
->element_exists(
31
Added:
'nav.probe-nav a[href="/graphs/ec"][data-nav-probe="ec"]')
32
Added:
->element_exists(
33
Added:
'footer.site-footer a[href="https://www.lafermeaquaponique.com/"]')
34
Added:
->element_exists('section#probe-ph[data-panel="ph"]')
35
Added:
->element_exists('canvas#chart-ph')
36
Added:
->element_exists_not('section#probe-do[data-panel="do"]')
37
Added:
->content_unlike(qr/Sensor Reading Quick Status/)
38
Added:
->content_like(qr/Hourly/)
39
Added:
->content_like(qr/Daily/)
40
Added:
->content_like(qr/Weekly/)
41
Added:
->content_like(qr/Yearly/)
42
Added:
->content_like(qr/data-smoothing-toggle checked/)
43
Added:
->content_like(qr/Latest pH Reading Messages/);
39
44
40
45
$t->get_ok('/graphs/unknown')->status_is(404);
41
46
roles/dashboard/t/02-readings-api.t
@@ -12,25 +12,25 @@
12
12
my $t = test_app();
13
13
14
14
$t->get_ok('/api/readings/ph')
15
Removed:
->status_is(200)
16
Removed:
->json_is( '/probe' => 'ph' )
17
Removed:
->json_is( '/readings/0/received_at' => '2026-07-05T12:00:10Z' )
18
Removed:
->json_is( '/readings/0/value' => 7.12 )
19
Removed:
->json_hasnt('/readings/1');
15
Added:
->status_is(200)
16
Added:
->json_is( '/probe' => 'ph' )
17
Added:
->json_is( '/readings/0/received_at' => '2026-07-05T12:00:10Z' )
18
Added:
->json_is( '/readings/0/value' => 7.12 )
19
Added:
->json_hasnt('/readings/1');
20
20
21
21
$t->get_ok('/api/readings/ph?since=2026-07-05T12:00:05Z')
22
Removed:
->status_is(200)
23
Removed:
->json_is( '/readings/0/value' => 7.12 )
24
Removed:
->json_hasnt('/readings/1');
22
Added:
->status_is(200)
23
Added:
->json_is( '/readings/0/value' => 7.12 )
24
Added:
->json_hasnt('/readings/1');
25
25
26
26
$t->get_ok('/api/readings/ph?since=2026-07-05T12:00:15Z')
27
Removed:
->status_is(200)
28
Removed:
->json_hasnt('/readings/0');
27
Added:
->status_is(200)
28
Added:
->json_hasnt('/readings/0');
29
29
30
30
$t->get_ok('/api/readings/do')
31
Removed:
->status_is(200)
32
Removed:
->json_is( '/probe' => 'do' )
33
Removed:
->json_is( '/readings/0/value' => 8.34 );
31
Added:
->status_is(200)
32
Added:
->json_is( '/probe' => 'do' )
33
Added:
->json_is( '/readings/0/value' => 8.34 );
34
34
35
35
$t->get_ok('/api/readings/nope')->status_is(404);
36
36
roles/dashboard/t/03-status-api.t
@@ -12,21 +12,22 @@
12
12
my $t = test_app();
13
13
14
14
$t->get_ok('/api/status/ph')
15
Removed:
->status_is(200)
16
Removed:
->json_is( '/probe' => 'ph' )
17
Removed:
->json_is( '/status/status' => 'ok' );
15
Added:
->status_is(200)
16
Added:
->json_is( '/probe' => 'ph' )
17
Added:
->json_is( '/status/status' => 'ok' );
18
18
19
19
$t->get_ok('/api/status/hub')
20
Removed:
->status_is(200)
21
Removed:
->json_is( '/probe' => 'hub' )
22
Removed:
->json_is( '/status/status' => 'ok' )
23
Removed:
->json_is( '/status/node' => 'fapg-daq-five-01' );
20
Added:
->status_is(200)
21
Added:
->json_is( '/probe' => 'hub' )
22
Added:
->json_is( '/status/status' => 'ok' )
23
Added:
->json_is( '/status/node' => 'fapg-daq-five-01' );
24
24
25
25
$t->get_ok('/api/status/do')
26
Removed:
->status_is(200)
27
Removed:
->json_is( '/probe' => 'do' )
28
Removed:
->json_is( '/status/status' => 'unreachable' )
29
Removed:
->json_is( '/status/message' => 'No status received in the last 120 seconds' );
26
Added:
->status_is(200)
27
Added:
->json_is( '/probe' => 'do' )
28
Added:
->json_is( '/status/status' => 'unreachable' )
29
Added:
->json_is(
30
Added:
'/status/message' => 'No status received in the last 120 seconds' );
30
31
31
32
$t->get_ok('/api/status/nope')->status_is(404);
32
33
roles/dashboard/t/lib/Dashboard/Test.pm
@@ -12,24 +12,24 @@
12
12
our @EXPORT_OK = qw(test_app);
13
13
14
14
sub test_app {
15
Removed:
$ENV{FAPG_DAQ_DB} = ':memory:';
15
Added:
$ENV{FAPG_DAQ_DB} = ':memory:';
16
16
17
Removed:
my $t = Test::Mojo->new('FAPG::DAQ::Dashboard');
18
Removed:
my $fresh_status_at = utc_timestamp(time);
19
Removed:
my $stale_status_at = utc_timestamp( time - 121 );
17
Added:
my $t = Test::Mojo->new('FAPG::DAQ::Dashboard');
18
Added:
my $fresh_status_at = utc_timestamp(time);
19
Added:
my $stale_status_at = utc_timestamp( time - 121 );
20
20
21
Removed:
create_schema($t);
22
Removed:
seed_readings($t);
23
Removed:
seed_statuses( $t, $fresh_status_at, $stale_status_at );
21
Added:
create_schema($t);
22
Added:
seed_readings($t);
23
Added:
seed_statuses( $t, $fresh_status_at, $stale_status_at );
24
24
25
Removed:
return $t;
25
Added:
return $t;
26
26
}
27
27
28
28
sub create_schema {
29
Removed:
my ($t) = @_;
29
Added:
my ($t) = @_;
30
30
31
Removed:
$t->app->sqlite->db->query(
32
Removed:
q{
31
Added:
$t->app->sqlite->db->query(
32
Added:
q{
33
33
CREATE TABLE readings (
34
34
id INTEGER PRIMARY KEY AUTOINCREMENT,
35
35
received_at TEXT NOT NULL,
@@ -41,10 +41,10 @@
41
41
raw_json TEXT
42
42
)
43
43
}
44
Removed:
);
44
Added:
);
45
45
46
Removed:
$t->app->sqlite->db->query(
47
Removed:
q{
46
Added:
$t->app->sqlite->db->query(
47
Added:
q{
48
48
CREATE TABLE node_status (
49
49
id INTEGER PRIMARY KEY AUTOINCREMENT,
50
50
received_at TEXT NOT NULL,
@@ -62,88 +62,101 @@
62
62
valid INTEGER NOT NULL DEFAULT 1
63
63
)
64
64
}
65
Removed:
);
65
Added:
);
66
66
67
Removed:
return;
67
Added:
return;
68
68
}
69
69
70
70
sub seed_readings {
71
Removed:
my ($t) = @_;
71
Added:
my ($t) = @_;
72
72
73
Removed:
$t->app->sqlite->db->query(
74
Removed:
q{
73
Added:
$t->app->sqlite->db->query(
74
Added:
q{
75
75
INSERT INTO readings (received_at, timestamp, node, probe, value, unit)
76
76
VALUES (?, ?, ?, ?, ?, ?)
77
77
},
78
Removed:
'2026-07-05T12:00:10Z',
79
Removed:
'2026-07-05T12:00:00Z',
80
Removed:
'fapg-daq-zero-ph-01',
81
Removed:
'ph',
82
Removed:
7.12,
83
Removed:
'pH'
84
Removed:
);
78
Added:
'2026-07-05T12:00:10Z',
79
Added:
'2026-07-05T12:00:00Z',
80
Added:
'fapg-daq-zero-ph-01',
81
Added:
'ph',
82
Added:
7.12,
83
Added:
'pH'
84
Added:
);
85
85
86
Removed:
$t->app->sqlite->db->query(
87
Removed:
q{
86
Added:
$t->app->sqlite->db->query(
87
Added:
q{
88
88
INSERT INTO readings (received_at, timestamp, node, probe, value, unit)
89
89
VALUES (?, ?, ?, ?, ?, ?)
90
90
},
91
Removed:
'2026-07-05T12:00:20Z',
92
Removed:
'2026-07-05T12:00:15Z',
93
Removed:
'fapg-daq-zero-do-01',
94
Removed:
'do',
95
Removed:
8.34,
96
Removed:
'mg/L'
97
Removed:
);
91
Added:
'2026-07-05T12:00:20Z',
92
Added:
'2026-07-05T12:00:15Z',
93
Added:
'fapg-daq-zero-do-01',
94
Added:
'do',
95
Added:
8.34,
96
Added:
'mg/L'
97
Added:
);
98
98
99
Removed:
return;
99
Added:
return;
100
100
}
101
101
102
102
sub seed_statuses {
103
Removed:
my ( $t, $fresh_status_at, $stale_status_at ) = @_;
103
Added:
my ( $t, $fresh_status_at, $stale_status_at ) = @_;
104
104
105
Removed:
insert_status( $t, $fresh_status_at, 'fapg/daq/ph/fapg-daq-zero-ph-01/status',
106
Removed:
'fapg.daq.status.v1', $fresh_status_at, 'ph', 'fapg-daq-zero-ph-01', 'ok',
107
Removed:
'Probe reading valid', '{}', );
105
Added:
insert_status(
106
Added:
$t, $fresh_status_at,
107
Added:
'fapg/daq/ph/fapg-daq-zero-ph-01/status', 'fapg.daq.status.v1',
108
Added:
$fresh_status_at, 'ph',
109
Added:
'fapg-daq-zero-ph-01', 'ok',
110
Added:
'Probe reading valid', '{}',
111
Added:
);
108
112
109
Removed:
insert_status( $t, $stale_status_at, 'fapg/daq/do/fapg-daq-zero-do-01/status',
110
Removed:
'fapg.daq.status.v1', $stale_status_at, 'do', 'fapg-daq-zero-do-01', 'ok',
111
Removed:
'Probe reading valid', '{}', );
113
Added:
insert_status(
114
Added:
$t, $stale_status_at,
115
Added:
'fapg/daq/do/fapg-daq-zero-do-01/status', 'fapg.daq.status.v1',
116
Added:
$stale_status_at, 'do',
117
Added:
'fapg-daq-zero-do-01', 'ok',
118
Added:
'Probe reading valid', '{}',
119
Added:
);
112
120
113
Removed:
insert_status( $t, $fresh_status_at, 'fapg/daq/hub/fapg-daq-five-01/status',
114
Removed:
'fapg.daq.status.v1', $fresh_status_at, 'hub', 'fapg-daq-five-01', 'ok',
115
Removed:
'DAQ hub MQTT broker reachable', '{}', );
121
Added:
insert_status(
122
Added:
$t, $fresh_status_at,
123
Added:
'fapg/daq/hub/fapg-daq-five-01/status', 'fapg.daq.status.v1',
124
Added:
$fresh_status_at, 'hub',
125
Added:
'fapg-daq-five-01', 'ok',
126
Added:
'DAQ hub MQTT broker reachable', '{}',
127
Added:
);
116
128
117
Removed:
return;
129
Added:
return;
118
130
}
119
131
120
132
sub insert_status {
121
Removed:
my ( $t, $received_at, $topic, $schema, $timestamp, $probe, $node, $status, $message, $payload )
122
Removed:
= @_;
133
Added:
my ($t, $received_at, $topic, $schema, $timestamp,
134
Added:
$probe, $node, $status, $message, $payload
135
Added:
) = @_;
123
136
124
Removed:
$t->app->sqlite->db->query(
125
Removed:
q{
137
Added:
$t->app->sqlite->db->query(
138
Added:
q{
126
139
INSERT INTO node_status (
127
140
received_at, topic, schema, timestamp, probe, node, status, message, payload
128
141
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
129
142
},
130
Removed:
$received_at,
131
Removed:
$topic,
132
Removed:
$schema,
133
Removed:
$timestamp,
134
Removed:
$probe,
135
Removed:
$node,
136
Removed:
$status,
137
Removed:
$message,
138
Removed:
$payload,
139
Removed:
);
143
Added:
$received_at,
144
Added:
$topic,
145
Added:
$schema,
146
Added:
$timestamp,
147
Added:
$probe,
148
Added:
$node,
149
Added:
$status,
150
Added:
$message,
151
Added:
$payload,
152
Added:
);
140
153
141
Removed:
return;
154
Added:
return;
142
155
}
143
156
144
157
sub utc_timestamp {
145
Removed:
my ($epoch) = @_;
146
Removed:
return strftime( '%Y-%m-%dT%H:%M:%SZ', gmtime $epoch );
158
Added:
my ($epoch) = @_;
159
Added:
return strftime( '%Y-%m-%dT%H:%M:%SZ', gmtime $epoch );
147
160
}
148
161
149
162
1;
t/00-ping-from-dev.t
@@ -11,8 +11,8 @@
11
11
use FAPG::TestSupport qw(@HOSTS);
12
12
13
13
foreach my $h (@HOSTS) {
14
Removed:
my $cmd = "ping -c 1 $h->{hostname} &> /dev/null";
15
Removed:
ok system($cmd) == 0, "host $h->{hostname} can be pinged";
14
Added:
my $cmd = "ping -c 1 $h->{hostname} &> /dev/null";
15
Added:
ok system($cmd) == 0, "host $h->{hostname} can be pinged";
16
16
}
17
17
18
18
done_testing();
t/01-ssh-from-dev.t
@@ -12,24 +12,25 @@
12
12
use Net::SSH::Perl;
13
13
14
14
foreach my $h (@HOSTS) {
15
Removed:
my ( $out, $err, $exit );
16
Removed:
my $timeout = 5;
15
Added:
my ( $out, $err, $exit );
16
Added:
my $timeout = 5;
17
17
18
Removed:
my $ok = eval {
19
Removed:
local $SIG{ALRM} = sub { die "SSH timeout after $timeout seconds.\n" };
20
Removed:
alarm $timeout;
18
Added:
my $ok = eval {
19
Added:
local $SIG{ALRM}
20
Added:
= sub { die "SSH timeout after $timeout seconds.\n" };
21
Added:
alarm $timeout;
21
22
22
Removed:
my $ssh = Net::SSH::Perl->new( $h->{hostname} );
23
Removed:
$ssh->login( $h->{user}, $h->{password} );
24
Removed:
( $out, $err, $exit ) = $ssh->cmd('true');
23
Added:
my $ssh = Net::SSH::Perl->new( $h->{hostname} );
24
Added:
$ssh->login( $h->{user}, $h->{password} );
25
Added:
( $out, $err, $exit ) = $ssh->cmd('true');
25
26
26
Removed:
1;
27
Removed:
};
27
Added:
1;
28
Added:
};
28
29
29
Removed:
alarm 0;
30
Added:
alarm 0;
30
31
31
Removed:
ok $ok, "host $h->{hostname} can be reached via ssh";
32
Removed:
diag $@ if $@;
32
Added:
ok $ok, "host $h->{hostname} can be reached via ssh";
33
Added:
diag $@ if $@;
33
34
34
35
}
35
36
t/lib/FAPG/TestSupport.pm
@@ -8,48 +8,42 @@
8
8
use Exporter 'import';
9
9
10
10
our @EXPORT_OK = qw(
11
Removed:
@HOSTS
11
Added:
@HOSTS
12
12
);
13
13
14
14
our @HOSTS = (
15
Removed:
{
16
Removed:
user => 'pirou',
17
Removed:
hostname => 'fapg-daq-five-01',
18
Removed:
password => 'pipirou',
19
Removed:
mqtt_from_vps => '10.77.0.2',
20
Removed:
},
21
Removed:
{
22
Removed:
hostname => 'vps',
23
Removed:
user => 'mpeter',
24
Removed:
},
25
Removed:
{
26
Removed:
hostname => 'fapg-daq-zero-ph-01',
27
Removed:
user => 'pirou',
28
Removed:
password => 'pipirou',
29
Removed:
probe => 'ph',
30
Removed:
unit => 'pH',
31
Removed:
},
32
Removed:
{
33
Removed:
hostname => 'fapg-daq-zero-do-01',
34
Removed:
user => 'pirou',
35
Removed:
password => 'pipirou',
36
Removed:
probe => 'do',
37
Removed:
unit => 'mg/L',
38
Removed:
},
39
Removed:
{
40
Removed:
hostname => 'fapg-daq-zero-orp-01',
41
Removed:
user => 'pirou',
42
Removed:
password => 'pipirou',
43
Removed:
probe => 'orp',
44
Removed:
unit => 'mV',
45
Removed:
},
46
Removed:
{
47
Removed:
hostname => 'fapg-daq-zero-ec-01',
48
Removed:
user => 'pirou',
49
Removed:
password => 'pipirou',
50
Removed:
probe => 'ec',
51
Removed:
unit => 'uS/cm',
52
Removed:
},
15
Added:
{ user => 'pirou',
16
Added:
hostname => 'fapg-daq-five-01',
17
Added:
password => 'pipirou',
18
Added:
mqtt_from_vps => '10.77.0.2',
19
Added:
},
20
Added:
{ hostname => 'vps',
21
Added:
user => 'mpeter',
22
Added:
},
23
Added:
{ hostname => 'fapg-daq-zero-ph-01',
24
Added:
user => 'pirou',
25
Added:
password => 'pipirou',
26
Added:
probe => 'ph',
27
Added:
unit => 'pH',
28
Added:
},
29
Added:
{ hostname => 'fapg-daq-zero-do-01',
30
Added:
user => 'pirou',
31
Added:
password => 'pipirou',
32
Added:
probe => 'do',
33
Added:
unit => 'mg/L',
34
Added:
},
35
Added:
{ hostname => 'fapg-daq-zero-orp-01',
36
Added:
user => 'pirou',
37
Added:
password => 'pipirou',
38
Added:
probe => 'orp',
39
Added:
unit => 'mV',
40
Added:
},
41
Added:
{ hostname => 'fapg-daq-zero-ec-01',
42
Added:
user => 'pirou',
43
Added:
password => 'pipirou',
44
Added:
probe => 'ec',
45
Added:
unit => 'uS/cm',
46
Added:
},
53
47
);
54
48
55
49
1;
t/lib/FAPG/TestSupport.pm.example
@@ -8,48 +8,42 @@
8
8
use Exporter 'import';
9
9
10
10
our @EXPORT_OK = qw(
11
Removed:
@HOSTS
11
Added:
@HOSTS
12
12
);
13
13
14
14
our @HOSTS = (
15
Removed:
{
16
Removed:
user => 'pirou',
17
Removed:
hostname => 'fapg-daq-five-01',
18
Removed:
password => 'CHANGE_ME',
19
Removed:
mqtt_from_vps => '10.77.0.2',
20
Removed:
},
21
Removed:
{
22
Removed:
hostname => 'vps',
23
Removed:
user => 'mpeter',
24
Removed:
},
25
Removed:
{
26
Removed:
hostname => 'fapg-daq-zero-ph-01',
27
Removed:
user => 'pirou',
28
Removed:
password => 'CHANGE_ME',
29
Removed:
probe => 'ph',
30
Removed:
unit => 'pH',
31
Removed:
},
32
Removed:
{
33
Removed:
hostname => 'fapg-daq-zero-do-01',
34
Removed:
user => 'pirou',
35
Removed:
password => 'CHANGE_ME',
36
Removed:
probe => 'do',
37
Removed:
unit => 'mg/L',
38
Removed:
},
39
Removed:
{
40
Removed:
hostname => 'fapg-daq-zero-orp-01',
41
Removed:
user => 'pirou',
42
Removed:
password => 'CHANGE_ME',
43
Removed:
probe => 'orp',
44
Removed:
unit => 'mV',
45
Removed:
},
46
Removed:
{
47
Removed:
hostname => 'fapg-daq-zero-ec-01',
48
Removed:
user => 'pirou',
49
Removed:
password => 'CHANGE_ME',
50
Removed:
probe => 'ec',
51
Removed:
unit => 'uS/cm',
52
Removed:
},
15
Added:
{ user => 'pirou',
16
Added:
hostname => 'fapg-daq-five-01',
17
Added:
password => 'CHANGE_ME',
18
Added:
mqtt_from_vps => '10.77.0.2',
19
Added:
},
20
Added:
{ hostname => 'vps',
21
Added:
user => 'mpeter',
22
Added:
},
23
Added:
{ hostname => 'fapg-daq-zero-ph-01',
24
Added:
user => 'pirou',
25
Added:
password => 'CHANGE_ME',
26
Added:
probe => 'ph',
27
Added:
unit => 'pH',
28
Added:
},
29
Added:
{ hostname => 'fapg-daq-zero-do-01',
30
Added:
user => 'pirou',
31
Added:
password => 'CHANGE_ME',
32
Added:
probe => 'do',
33
Added:
unit => 'mg/L',
34
Added:
},
35
Added:
{ hostname => 'fapg-daq-zero-orp-01',
36
Added:
user => 'pirou',
37
Added:
password => 'CHANGE_ME',
38
Added:
probe => 'orp',
39
Added:
unit => 'mV',
40
Added:
},
41
Added:
{ hostname => 'fapg-daq-zero-ec-01',
42
Added:
user => 'pirou',
43
Added:
password => 'CHANGE_ME',
44
Added:
probe => 'ec',
45
Added:
unit => 'uS/cm',
46
Added:
},
53
47
);
54
48
55
49
1;