[Perl] DAQ system for the FAPG.
Flatten all Perl library paths.
Changed files
AGENTS.org
@@ -69,7 +69,7 @@
69
69
* Repository Map
70
70
71
71
- =README.org= :: describes the overall architecture.
72
Removed:
- =lib/perl5/FAPG/DAQ/= :: contains shared Perl modules.
72
Added:
- =lib/FAPG/DAQ/= :: contains shared Perl modules.
73
73
- =roles/daq-node/= :: contains the probe reader script and node
74
74
service file. The Pi Zeros handle this role.
75
75
- =roles/daq-hub/= :: contains MQTT broker configuration. The Pi 5
@@ -96,9 +96,9 @@
96
96
97
97
** Structure
98
98
99
Removed:
- Shared DAQ behavior SHALL appear under =lib/perl5/FAPG/DAQ/= when it
99
Added:
- Shared DAQ behavior SHALL appear under =lib/FAPG/DAQ/= when it
100
100
is used by more than one role.
101
Removed:
- Role-local code SHALL stay under that role's =lib/perl5=, =bin=,
101
Added:
- Role-local code SHALL stay under that role's =lib=, =bin=,
102
102
=etc=, or =t= tree.
103
103
104
104
** Perl
lib/FAPG/DAQ/MQTT.pm
@@ -0,0 +1,226 @@
1
Added:
# -*- mode: cperl; -*-
2
Added:
3
Added:
package FAPG::DAQ::MQTT;
4
Added:
5
Added:
use v5.32.1;
6
Added:
use warnings;
7
Added:
8
Added:
use Carp qw(croak);
9
Added:
use Exporter 'import';
10
Added:
use JSON::PP;
11
Added:
use POSIX qw(strftime);
12
Added:
use Sys::Hostname qw(hostname);
13
Added:
14
Added:
our @EXPORT_OK = qw(
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
Added:
);
23
Added:
24
Added:
my $DEFAULT_HOST = 'fapg-daq-five-01';
25
Added:
my $DEFAULT_PORT = 1883;
26
Added:
my $DEFAULT_SCHEMA = 'fapg.daq.reading.v1';
27
Added:
my $DEFAULT_STATUS_SCHEMA = 'fapg.daq.status.v1';
28
Added:
my $TOPIC_PREFIX = 'fapg/daq';
29
Added:
30
Added:
sub publish_mqtt {
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
Added:
36
Added:
$client->publish( $topic => $payload );
37
Added:
38
Added:
return {
39
Added:
topic => $topic,
40
Added:
payload => $payload,
41
Added:
};
42
Added:
}
43
Added:
44
Added:
sub publish_status_mqtt {
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
Added:
50
Added:
$client->publish( $topic => $payload );
51
Added:
52
Added:
return {
53
Added:
topic => $topic,
54
Added:
payload => $payload,
55
Added:
};
56
Added:
}
57
Added:
58
Added:
sub mqtt_payload_for_reading {
59
Added:
my ( $reading, %opt ) = @_;
60
Added:
_assert_reading($reading);
61
Added:
62
Added:
my $probe = lc $reading->{probe};
63
Added:
my $node = _node_name( $reading, %opt );
64
Added:
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
Added:
75
Added:
$payload{raw} = $reading->{raw}
76
Added:
if exists $reading->{raw} && defined $reading->{raw};
77
Added:
78
Added:
$payload{values} = [ map { 0 + $_ } $reading->{values}->@* ]
79
Added:
if ref( $reading->{values} ) eq 'ARRAY';
80
Added:
81
Added:
return JSON::PP->new->canonical(1)->encode( \%payload );
82
Added:
}
83
Added:
84
Added:
sub mqtt_payload_for_status {
85
Added:
my ( $status, %opt ) = @_;
86
Added:
_assert_status($status);
87
Added:
88
Added:
my $probe = lc $status->{probe};
89
Added:
my $node = _node_name( $status, %opt );
90
Added:
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
Added:
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
Added:
105
Added:
return JSON::PP->new->canonical(1)->encode( \%payload );
106
Added:
}
107
Added:
108
Added:
sub mqtt_topic_for_reading {
109
Added:
my ( $reading, %opt ) = @_;
110
Added:
_assert_reading($reading);
111
Added:
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
Added:
116
Added:
$prefix =~ s{/+\z}{};
117
Added:
118
Added:
return join '/', $prefix, $probe, $node, 'reading';
119
Added:
}
120
Added:
121
Added:
sub mqtt_topic_for_status {
122
Added:
my ( $status, %opt ) = @_;
123
Added:
_assert_status($status);
124
Added:
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
Added:
129
Added:
$prefix =~ s{/+\z}{};
130
Added:
131
Added:
return join '/', $prefix, $probe, $node, 'status';
132
Added:
}
133
Added:
134
Added:
sub mqtt_broker_endpoint {
135
Added:
my (%opt) = @_;
136
Added:
my $host = $opt{host} // $ENV{FAPG_MQTT_HOST} // $ENV{MQTT_HOST} // $DEFAULT_HOST;
137
Added:
138
Added:
my $port = $opt{port} // $ENV{FAPG_MQTT_PORT} // $ENV{MQTT_PORT} // $DEFAULT_PORT;
139
Added:
140
Added:
return $host if $host =~ /:\d+\z/;
141
Added:
return "$host:$port";
142
Added:
}
143
Added:
144
Added:
sub _mqtt_client {
145
Added:
my (%opt) = @_;
146
Added:
my $endpoint = mqtt_broker_endpoint(%opt);
147
Added:
148
Added:
my $username = $opt{username} // $ENV{FAPG_MQTT_USERNAME} // $ENV{MQTT_USERNAME} // 'fapg_zero';
149
Added:
150
Added:
my $password = $opt{password} // $ENV{FAPG_MQTT_PASSWORD} // $ENV{MQTT_PASSWORD};
151
Added:
152
Added:
state %clients;
153
Added:
154
Added:
my $cache_key = join "\0", $endpoint, ( $username // q{} ), ( $password // q{} );
155
Added:
return $clients{$cache_key} if exists $clients{$cache_key};
156
Added:
157
Added:
require Net::MQTT::Simple;
158
Added:
159
Added:
my $client = Net::MQTT::Simple->new($endpoint)
160
Added:
or croak "cannot connect to MQTT broker $endpoint";
161
Added:
162
Added:
if ( defined $username && $username ne q{} ) {
163
Added:
croak 'MQTT password is required when MQTT username is set'
164
Added:
if !defined $password;
165
Added:
166
Added:
$client->login( $username, $password );
167
Added:
}
168
Added:
169
Added:
return $clients{$cache_key} = $client;
170
Added:
}
171
Added:
172
Added:
sub _assert_reading {
173
Added:
my ($reading) = @_;
174
Added:
croak 'reading must be a HASH reference'
175
Added:
if ref($reading) ne 'HASH';
176
Added:
177
Added:
croak 'reading requires a probe field'
178
Added:
if !defined $reading->{probe} || $reading->{probe} eq q{};
179
Added:
180
Added:
croak 'reading requires a value field'
181
Added:
if !exists $reading->{value} || !defined $reading->{value};
182
Added:
183
Added:
croak "reading value is not numeric: $reading->{value}"
184
Added:
if $reading->{value} !~ /\A[+-]?(?:\d+(?:\.\d*)?|\.\d+)\z/;
185
Added:
186
Added:
return;
187
Added:
}
188
Added:
189
Added:
sub _assert_status {
190
Added:
my ($status) = @_;
191
Added:
croak 'status must be a HASH reference'
192
Added:
if ref($status) ne 'HASH';
193
Added:
194
Added:
croak 'status requires a probe field'
195
Added:
if !defined $status->{probe} || $status->{probe} eq q{};
196
Added:
197
Added:
croak 'status requires a status field'
198
Added:
if !defined $status->{status} || $status->{status} eq q{};
199
Added:
200
Added:
croak "unsupported node status '$status->{status}'"
201
Added:
if lc( $status->{status} ) !~ /\A(?:ok|probe_error)\z/;
202
Added:
203
Added:
return;
204
Added:
}
205
Added:
206
Added:
sub _node_name {
207
Added:
my ( $reading, %opt ) = @_;
208
Added:
return $reading->{node} // $opt{node} // $ENV{FAPG_DAQ_NODE} // $ENV{DAQ_NODE} // hostname();
209
Added:
}
210
Added:
211
Added:
sub _topic_level {
212
Added:
my ( $name, $value ) = @_;
213
Added:
croak "$name topic level is required"
214
Added:
if !defined $value || $value eq q{};
215
Added:
216
Added:
croak "$name topic level must not contain '/', '+', '#', or NUL: $value"
217
Added:
if $value =~ m{[\0/+\#]};
218
Added:
219
Added:
return $value;
220
Added:
}
221
Added:
222
Added:
sub _utc_timestamp {
223
Added:
return strftime '%Y-%m-%dT%H:%M:%SZ', gmtime;
224
Added:
}
225
Added:
226
Added:
1;
lib/perl5/FAPG/DAQ/MQTT.pm
@@ -1,226 +0,0 @@
1
Removed:
# -*- mode: cperl; -*-
2
Removed:
3
Removed:
package FAPG::DAQ::MQTT;
4
Removed:
5
Removed:
use v5.32.1;
6
Removed:
use warnings;
7
Removed:
8
Removed:
use Carp qw(croak);
9
Removed:
use Exporter 'import';
10
Removed:
use JSON::PP;
11
Removed:
use POSIX qw(strftime);
12
Removed:
use Sys::Hostname qw(hostname);
13
Removed:
14
Removed:
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
22
Removed:
);
23
Removed:
24
Removed:
my $DEFAULT_HOST = 'fapg-daq-five-01';
25
Removed:
my $DEFAULT_PORT = 1883;
26
Removed:
my $DEFAULT_SCHEMA = 'fapg.daq.reading.v1';
27
Removed:
my $DEFAULT_STATUS_SCHEMA = 'fapg.daq.status.v1';
28
Removed:
my $TOPIC_PREFIX = 'fapg/daq';
29
Removed:
30
Removed:
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 );
35
Removed:
36
Removed:
$client->publish( $topic => $payload );
37
Removed:
38
Removed:
return {
39
Removed:
topic => $topic,
40
Removed:
payload => $payload,
41
Removed:
};
42
Removed:
}
43
Removed:
44
Removed:
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 );
49
Removed:
50
Removed:
$client->publish( $topic => $payload );
51
Removed:
52
Removed:
return {
53
Removed:
topic => $topic,
54
Removed:
payload => $payload,
55
Removed:
};
56
Removed:
}
57
Removed:
58
Removed:
sub mqtt_payload_for_reading {
59
Removed:
my ( $reading, %opt ) = @_;
60
Removed:
_assert_reading($reading);
61
Removed:
62
Removed:
my $probe = lc $reading->{probe};
63
Removed:
my $node = _node_name( $reading, %opt );
64
Removed:
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:
);
74
Removed:
75
Removed:
$payload{raw} = $reading->{raw}
76
Removed:
if exists $reading->{raw} && defined $reading->{raw};
77
Removed:
78
Removed:
$payload{values} = [ map { 0 + $_ } $reading->{values}->@* ]
79
Removed:
if ref( $reading->{values} ) eq 'ARRAY';
80
Removed:
81
Removed:
return JSON::PP->new->canonical(1)->encode( \%payload );
82
Removed:
}
83
Removed:
84
Removed:
sub mqtt_payload_for_status {
85
Removed:
my ( $status, %opt ) = @_;
86
Removed:
_assert_status($status);
87
Removed:
88
Removed:
my $probe = lc $status->{probe};
89
Removed:
my $node = _node_name( $status, %opt );
90
Removed:
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:
);
99
Removed:
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:
}
104
Removed:
105
Removed:
return JSON::PP->new->canonical(1)->encode( \%payload );
106
Removed:
}
107
Removed:
108
Removed:
sub mqtt_topic_for_reading {
109
Removed:
my ( $reading, %opt ) = @_;
110
Removed:
_assert_reading($reading);
111
Removed:
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 ) );
115
Removed:
116
Removed:
$prefix =~ s{/+\z}{};
117
Removed:
118
Removed:
return join '/', $prefix, $probe, $node, 'reading';
119
Removed:
}
120
Removed:
121
Removed:
sub mqtt_topic_for_status {
122
Removed:
my ( $status, %opt ) = @_;
123
Removed:
_assert_status($status);
124
Removed:
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 ) );
128
Removed:
129
Removed:
$prefix =~ s{/+\z}{};
130
Removed:
131
Removed:
return join '/', $prefix, $probe, $node, 'status';
132
Removed:
}
133
Removed:
134
Removed:
sub mqtt_broker_endpoint {
135
Removed:
my (%opt) = @_;
136
Removed:
my $host = $opt{host} // $ENV{FAPG_MQTT_HOST} // $ENV{MQTT_HOST} // $DEFAULT_HOST;
137
Removed:
138
Removed:
my $port = $opt{port} // $ENV{FAPG_MQTT_PORT} // $ENV{MQTT_PORT} // $DEFAULT_PORT;
139
Removed:
140
Removed:
return $host if $host =~ /:\d+\z/;
141
Removed:
return "$host:$port";
142
Removed:
}
143
Removed:
144
Removed:
sub _mqtt_client {
145
Removed:
my (%opt) = @_;
146
Removed:
my $endpoint = mqtt_broker_endpoint(%opt);
147
Removed:
148
Removed:
my $username = $opt{username} // $ENV{FAPG_MQTT_USERNAME} // $ENV{MQTT_USERNAME} // 'fapg_zero';
149
Removed:
150
Removed:
my $password = $opt{password} // $ENV{FAPG_MQTT_PASSWORD} // $ENV{MQTT_PASSWORD};
151
Removed:
152
Removed:
state %clients;
153
Removed:
154
Removed:
my $cache_key = join "\0", $endpoint, ( $username // q{} ), ( $password // q{} );
155
Removed:
return $clients{$cache_key} if exists $clients{$cache_key};
156
Removed:
157
Removed:
require Net::MQTT::Simple;
158
Removed:
159
Removed:
my $client = Net::MQTT::Simple->new($endpoint)
160
Removed:
or croak "cannot connect to MQTT broker $endpoint";
161
Removed:
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;
165
Removed:
166
Removed:
$client->login( $username, $password );
167
Removed:
}
168
Removed:
169
Removed:
return $clients{$cache_key} = $client;
170
Removed:
}
171
Removed:
172
Removed:
sub _assert_reading {
173
Removed:
my ($reading) = @_;
174
Removed:
croak 'reading must be a HASH reference'
175
Removed:
if ref($reading) ne 'HASH';
176
Removed:
177
Removed:
croak 'reading requires a probe field'
178
Removed:
if !defined $reading->{probe} || $reading->{probe} eq q{};
179
Removed:
180
Removed:
croak 'reading requires a value field'
181
Removed:
if !exists $reading->{value} || !defined $reading->{value};
182
Removed:
183
Removed:
croak "reading value is not numeric: $reading->{value}"
184
Removed:
if $reading->{value} !~ /\A[+-]?(?:\d+(?:\.\d*)?|\.\d+)\z/;
185
Removed:
186
Removed:
return;
187
Removed:
}
188
Removed:
189
Removed:
sub _assert_status {
190
Removed:
my ($status) = @_;
191
Removed:
croak 'status must be a HASH reference'
192
Removed:
if ref($status) ne 'HASH';
193
Removed:
194
Removed:
croak 'status requires a probe field'
195
Removed:
if !defined $status->{probe} || $status->{probe} eq q{};
196
Removed:
197
Removed:
croak 'status requires a status field'
198
Removed:
if !defined $status->{status} || $status->{status} eq q{};
199
Removed:
200
Removed:
croak "unsupported node status '$status->{status}'"
201
Removed:
if lc( $status->{status} ) !~ /\A(?:ok|probe_error)\z/;
202
Removed:
203
Removed:
return;
204
Removed:
}
205
Removed:
206
Removed:
sub _node_name {
207
Removed:
my ( $reading, %opt ) = @_;
208
Removed:
return $reading->{node} // $opt{node} // $ENV{FAPG_DAQ_NODE} // $ENV{DAQ_NODE} // hostname();
209
Removed:
}
210
Removed:
211
Removed:
sub _topic_level {
212
Removed:
my ( $name, $value ) = @_;
213
Removed:
croak "$name topic level is required"
214
Removed:
if !defined $value || $value eq q{};
215
Removed:
216
Removed:
croak "$name topic level must not contain '/', '+', '#', or NUL: $value"
217
Removed:
if $value =~ m{[\0/+\#]};
218
Removed:
219
Removed:
return $value;
220
Removed:
}
221
Removed:
222
Removed:
sub _utc_timestamp {
223
Removed:
return strftime '%Y-%m-%dT%H:%M:%SZ', gmtime;
224
Removed:
}
225
Removed:
226
Removed:
1;
roles/daq-hub/bin/fapg-daq-hub-status
@@ -6,7 +6,7 @@
6
6
use warnings;
7
7
8
8
use FindBin qw($Bin);
9
Removed:
use lib "$Bin/../../../lib/perl5";
9
Added:
use lib "$Bin/../../../lib";
10
10
11
11
use FAPG::DAQ::MQTT qw(publish_status_mqtt);
12
12
use Sys::Hostname qw(hostname);
roles/daq-node/bin/fapg-daq-node
@@ -6,8 +6,8 @@
6
6
use warnings;
7
7
8
8
use FindBin qw($Bin);
9
Removed:
use lib "$Bin/../../../lib/perl5";
10
Removed:
use lib "$Bin/../lib/perl5";
9
Added:
use lib "$Bin/../../../lib";
10
Added:
use lib "$Bin/../lib";
11
11
12
12
use FAPG::DAQ::EZO::USB qw(
13
13
discover_ezo_usb_device
@@ -131,7 +131,7 @@
131
131
}
132
132
133
133
my $hostname = hostname();
134
Removed:
return $1 if $hostname =~ /(?:\A|-)(ph|do|ec|orp)(?:-|\z)/;
134
Added:
return $1 if $hostname =~ /(?:\A|-)(ph|do|ec|orp)(?:-|\z)/;
135
135
136
136
return undef;
137
137
}
roles/daq-node/lib/FAPG/DAQ/EZO/USB.pm
@@ -0,0 +1,243 @@
1
Added:
# -*- mode: cperl; -*-
2
Added:
3
Added:
package FAPG::DAQ::EZO::USB;
4
Added:
5
Added:
use v5.32.1;
6
Added:
use strict;
7
Added:
use warnings;
8
Added:
9
Added:
use Exporter 'import';
10
Added:
use Carp qw(croak);
11
Added:
use Time::HiRes qw(time sleep);
12
Added:
13
Added:
our @EXPORT_OK = qw(
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
Added:
);
25
Added:
26
Added:
my $DEFAULT_BAUDRATE = 9_600;
27
Added:
my $DEFAULT_TIMEOUT_S = 1.5;
28
Added:
my $EOL = "\r";
29
Added:
30
Added:
sub discover_ezo_usb_device {
31
Added:
my %opt = @_;
32
Added:
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
Added:
40
Added:
my @candidates;
41
Added:
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
Added:
46
Added:
push @candidates, grep { /$preferred_pattern/ } @by_id;
47
Added:
push @candidates, grep { $_ !~ /$preferred_pattern/ } @by_id;
48
Added:
}
49
Added:
50
Added:
for my $pattern (@fallback_patterns) {
51
Added:
push @candidates, sort grep { -e $_ && !-d $_ } glob $pattern;
52
Added:
}
53
Added:
54
Added:
my %seen;
55
Added:
@candidates = grep { !$seen{$_}++ } @candidates;
56
Added:
57
Added:
return $candidates[0] if @candidates;
58
Added:
59
Added:
croak "cannot discover EZO USB serial device under $by_id_dir, /dev/ttyUSB*, or /dev/ttyACM*";
60
Added:
}
61
Added:
62
Added:
sub open_ezo_usb {
63
Added:
my %opt = @_;
64
Added:
my $device = $opt{device} // croak 'open_ezo_usb requires device => "/dev/tty..."';
65
Added:
my $baudrate = $opt{baudrate} // $DEFAULT_BAUDRATE;
66
Added:
my $timeout_ms = $opt{timeout_ms} // 250;
67
Added:
68
Added:
require Device::SerialPort;
69
Added:
70
Added:
my $port = Device::SerialPort->new($device)
71
Added:
or croak "cannot open serial device $device";
72
Added:
73
Added:
$port->baudrate($baudrate) or croak "cannot set baudrate on $device";
74
Added:
$port->databits(8) or croak "cannot set databits on $device";
75
Added:
$port->parity('none') or croak "cannot set parity on $device";
76
Added:
$port->stopbits(1) or croak "cannot set stopbits on $device";
77
Added:
$port->handshake('none') or croak "cannot set handshake on $device";
78
Added:
$port->read_char_time(0);
79
Added:
$port->read_const_time($timeout_ms);
80
Added:
$port->write_settings or croak "cannot apply serial settings on $device";
81
Added:
82
Added:
return $port;
83
Added:
}
84
Added:
85
Added:
sub drain_serial {
86
Added:
my ( $io, %opt ) = @_;
87
Added:
my $max_s = $opt{max_s} // 0.20;
88
Added:
my $chunk_sz = $opt{chunk_sz} // 255;
89
Added:
my $deadline = time + $max_s;
90
Added:
my $drained = q{};
91
Added:
92
Added:
while ( time < $deadline ) {
93
Added:
my ( $count, $chunk ) = $io->read($chunk_sz);
94
Added:
last if not defined $count || $count == 0;
95
Added:
$drained .= $chunk // q{};
96
Added:
}
97
Added:
98
Added:
return $drained;
99
Added:
}
100
Added:
101
Added:
sub ezo_command {
102
Added:
my ( $io, $command, %opt ) = @_;
103
Added:
croak 'command must not contain carriage returns or newlines'
104
Added:
if $command =~ /[\r\n]/;
105
Added:
106
Added:
drain_serial( $io, max_s => ( $opt{drain_s} // 0.05 ) )
107
Added:
if $opt{drain} // 1;
108
Added:
109
Added:
my $bytes = $command . $EOL;
110
Added:
my $wrote = $io->write($bytes);
111
Added:
112
Added:
croak "short write to EZO device: wrote $wrote of " . length($bytes) . ' bytes'
113
Added:
if not defined $wrote || $wrote != length($bytes);
114
Added:
115
Added:
my $response = _read_ezo_line( $io, timeout_s => ( $opt{timeout_s} // $DEFAULT_TIMEOUT_S ), );
116
Added:
117
Added:
croak "EZO command '$command' returned an error: $response"
118
Added:
if $response =~ /\A\?(?:ERROR|ER)\z/i;
119
Added:
120
Added:
return $response;
121
Added:
}
122
Added:
123
Added:
sub read_probe_info {
124
Added:
my ( $io, %opt ) = @_;
125
Added:
my $raw = ezo_command( $io, 'I', %opt );
126
Added:
127
Added:
my ( undef, $device, $firmware ) = split qr{,}, $raw, 3;
128
Added:
129
Added:
croak "unexpected EZO info response: $raw"
130
Added:
if not defined $device || $raw !~ /\A\?I,/;
131
Added:
132
Added:
my $probe = normalize_probe_type($device);
133
Added:
134
Added:
return {
135
Added:
raw => $raw,
136
Added:
device => $device,
137
Added:
firmware => $firmware,
138
Added:
probe => $probe,
139
Added:
unit => unit_for_probe_type($probe),
140
Added:
};
141
Added:
}
142
Added:
143
Added:
sub detect_probe_type {
144
Added:
my ( $io, %opt ) = @_;
145
Added:
return read_probe_info( $io, %opt )->{probe};
146
Added:
}
147
Added:
148
Added:
sub read_once {
149
Added:
my ( $io, %opt ) = @_;
150
Added:
my $probe = $opt{probe};
151
Added:
152
Added:
$probe = detect_probe_type( $io, %opt ) if not defined $probe;
153
Added:
$probe = normalize_probe_type($probe);
154
Added:
155
Added:
my $raw = ezo_command( $io, 'R', %opt );
156
Added:
my $parsed = parse_reading( $raw, probe => $probe );
157
Added:
158
Added:
return {
159
Added:
probe => $probe,
160
Added:
unit => unit_for_probe_type($probe),
161
Added:
raw => $raw,
162
Added:
value => $parsed->{value},
163
Added:
values => $parsed->{values},
164
Added:
};
165
Added:
}
166
Added:
167
Added:
sub parse_reading {
168
Added:
my ( $raw, %opt ) = @_;
169
Added:
croak 'empty EZO reading'
170
Added:
if not defined $raw || $raw eq q{};
171
Added:
172
Added:
my @fields = split /,/, $raw;
173
Added:
my @values;
174
Added:
175
Added:
for my $field (@fields) {
176
Added:
croak "non-numeric EZO reading field '$field' in '$raw'"
177
Added:
if $field !~ /\A[+-]?(?:\d+(?:\.\d*)?|\.\d+)\z/;
178
Added:
179
Added:
push @values, 0 + $field;
180
Added:
}
181
Added:
182
Added:
return {
183
Added:
raw => $raw,
184
Added:
value => $values[0],
185
Added:
values => \@values,
186
Added:
};
187
Added:
}
188
Added:
189
Added:
sub normalize_probe_type {
190
Added:
my ($probe) = @_;
191
Added:
croak 'probe type is required'
192
Added:
if not defined $probe || $probe eq q{};
193
Added:
194
Added:
my $p = lc $probe;
195
Added:
$p =~ s/[\s_-]+//g;
196
Added:
197
Added:
return
198
Added:
( $p eq 'ph' ) ? 'ph'
199
Added:
: ( $p eq 'do' || $p eq 'dissolvedoxygen' ) ? 'do'
200
Added:
: ( $p eq 'ec' || $p eq 'conductivity' ) ? 'ec'
201
Added:
: ( $p eq 'orp' ) ? 'orp'
202
Added:
: croak "unknown EZO probe type '$probe'";
203
Added:
}
204
Added:
205
Added:
sub unit_for_probe_type {
206
Added:
my ($probe) = @_;
207
Added:
my $p = normalize_probe_type($probe);
208
Added:
209
Added:
return 'pH' if $p eq 'ph';
210
Added:
return 'mg/L' if $p eq 'do';
211
Added:
return 'mV' if $p eq 'orp';
212
Added:
return 'uS/cm' if $p eq 'ec';
213
Added:
214
Added:
croak "unknown normalized probe type '$p'";
215
Added:
}
216
Added:
217
Added:
sub _read_ezo_line {
218
Added:
my ( $io, %opt ) = @_;
219
Added:
my $timeout_s = $opt{timeout_s} // $DEFAULT_TIMEOUT_S;
220
Added:
my $deadline = time + $timeout_s;
221
Added:
my $buf = q{};
222
Added:
223
Added:
while ( time < $deadline ) {
224
Added:
my ( $count, $chunk ) = $io->read(1);
225
Added:
226
Added:
if ( !defined $count || $count == 0 ) {
227
Added:
sleep 0.01;
228
Added:
next;
229
Added:
}
230
Added:
231
Added:
$buf .= $chunk;
232
Added:
last if $buf =~ /\r\z/;
233
Added:
}
234
Added:
235
Added:
croak "timeout waiting for EZO response after ${timeout_s}s"
236
Added:
if $buf !~ /\r\z/;
237
Added:
238
Added:
$buf =~ s/[\r\n]+\z//;
239
Added:
240
Added:
return $buf;
241
Added:
}
242
Added:
243
Added:
1;
roles/daq-node/lib/perl5/FAPG/DAQ/EZO/USB.pm
@@ -1,243 +0,0 @@
1
Removed:
# -*- mode: cperl; -*-
2
Removed:
3
Removed:
package FAPG::DAQ::EZO::USB;
4
Removed:
5
Removed:
use v5.32.1;
6
Removed:
use strict;
7
Removed:
use warnings;
8
Removed:
9
Removed:
use Exporter 'import';
10
Removed:
use Carp qw(croak);
11
Removed:
use Time::HiRes qw(time sleep);
12
Removed:
13
Removed:
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
24
Removed:
);
25
Removed:
26
Removed:
my $DEFAULT_BAUDRATE = 9_600;
27
Removed:
my $DEFAULT_TIMEOUT_S = 1.5;
28
Removed:
my $EOL = "\r";
29
Removed:
30
Removed:
sub discover_ezo_usb_device {
31
Removed:
my %opt = @_;
32
Removed:
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*);
39
Removed:
40
Removed:
my @candidates;
41
Removed:
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;
45
Removed:
46
Removed:
push @candidates, grep { /$preferred_pattern/ } @by_id;
47
Removed:
push @candidates, grep { $_ !~ /$preferred_pattern/ } @by_id;
48
Removed:
}
49
Removed:
50
Removed:
for my $pattern (@fallback_patterns) {
51
Removed:
push @candidates, sort grep { -e $_ && !-d $_ } glob $pattern;
52
Removed:
}
53
Removed:
54
Removed:
my %seen;
55
Removed:
@candidates = grep { !$seen{$_}++ } @candidates;
56
Removed:
57
Removed:
return $candidates[0] if @candidates;
58
Removed:
59
Removed:
croak "cannot discover EZO USB serial device under $by_id_dir, /dev/ttyUSB*, or /dev/ttyACM*";
60
Removed:
}
61
Removed:
62
Removed:
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;
67
Removed:
68
Removed:
require Device::SerialPort;
69
Removed:
70
Removed:
my $port = Device::SerialPort->new($device)
71
Removed:
or croak "cannot open serial device $device";
72
Removed:
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";
81
Removed:
82
Removed:
return $port;
83
Removed:
}
84
Removed:
85
Removed:
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{};
91
Removed:
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:
}
97
Removed:
98
Removed:
return $drained;
99
Removed:
}
100
Removed:
101
Removed:
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]/;
105
Removed:
106
Removed:
drain_serial( $io, max_s => ( $opt{drain_s} // 0.05 ) )
107
Removed:
if $opt{drain} // 1;
108
Removed:
109
Removed:
my $bytes = $command . $EOL;
110
Removed:
my $wrote = $io->write($bytes);
111
Removed:
112
Removed:
croak "short write to EZO device: wrote $wrote of " . length($bytes) . ' bytes'
113
Removed:
if not defined $wrote || $wrote != length($bytes);
114
Removed:
115
Removed:
my $response = _read_ezo_line( $io, timeout_s => ( $opt{timeout_s} // $DEFAULT_TIMEOUT_S ), );
116
Removed:
117
Removed:
croak "EZO command '$command' returned an error: $response"
118
Removed:
if $response =~ /\A\?(?:ERROR|ER)\z/i;
119
Removed:
120
Removed:
return $response;
121
Removed:
}
122
Removed:
123
Removed:
sub read_probe_info {
124
Removed:
my ( $io, %opt ) = @_;
125
Removed:
my $raw = ezo_command( $io, 'I', %opt );
126
Removed:
127
Removed:
my ( undef, $device, $firmware ) = split qr{,}, $raw, 3;
128
Removed:
129
Removed:
croak "unexpected EZO info response: $raw"
130
Removed:
if not defined $device || $raw !~ /\A\?I,/;
131
Removed:
132
Removed:
my $probe = normalize_probe_type($device);
133
Removed:
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:
};
141
Removed:
}
142
Removed:
143
Removed:
sub detect_probe_type {
144
Removed:
my ( $io, %opt ) = @_;
145
Removed:
return read_probe_info( $io, %opt )->{probe};
146
Removed:
}
147
Removed:
148
Removed:
sub read_once {
149
Removed:
my ( $io, %opt ) = @_;
150
Removed:
my $probe = $opt{probe};
151
Removed:
152
Removed:
$probe = detect_probe_type( $io, %opt ) if not defined $probe;
153
Removed:
$probe = normalize_probe_type($probe);
154
Removed:
155
Removed:
my $raw = ezo_command( $io, 'R', %opt );
156
Removed:
my $parsed = parse_reading( $raw, probe => $probe );
157
Removed:
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:
};
165
Removed:
}
166
Removed:
167
Removed:
sub parse_reading {
168
Removed:
my ( $raw, %opt ) = @_;
169
Removed:
croak 'empty EZO reading'
170
Removed:
if not defined $raw || $raw eq q{};
171
Removed:
172
Removed:
my @fields = split /,/, $raw;
173
Removed:
my @values;
174
Removed:
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/;
178
Removed:
179
Removed:
push @values, 0 + $field;
180
Removed:
}
181
Removed:
182
Removed:
return {
183
Removed:
raw => $raw,
184
Removed:
value => $values[0],
185
Removed:
values => \@values,
186
Removed:
};
187
Removed:
}
188
Removed:
189
Removed:
sub normalize_probe_type {
190
Removed:
my ($probe) = @_;
191
Removed:
croak 'probe type is required'
192
Removed:
if not defined $probe || $probe eq q{};
193
Removed:
194
Removed:
my $p = lc $probe;
195
Removed:
$p =~ s/[\s_-]+//g;
196
Removed:
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'";
203
Removed:
}
204
Removed:
205
Removed:
sub unit_for_probe_type {
206
Removed:
my ($probe) = @_;
207
Removed:
my $p = normalize_probe_type($probe);
208
Removed:
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';
213
Removed:
214
Removed:
croak "unknown normalized probe type '$p'";
215
Removed:
}
216
Removed:
217
Removed:
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{};
222
Removed:
223
Removed:
while ( time < $deadline ) {
224
Removed:
my ( $count, $chunk ) = $io->read(1);
225
Removed:
226
Removed:
if ( !defined $count || $count == 0 ) {
227
Removed:
sleep 0.01;
228
Removed:
next;
229
Removed:
}
230
Removed:
231
Removed:
$buf .= $chunk;
232
Removed:
last if $buf =~ /\r\z/;
233
Removed:
}
234
Removed:
235
Removed:
croak "timeout waiting for EZO response after ${timeout_s}s"
236
Removed:
if $buf !~ /\r\z/;
237
Removed:
238
Removed:
$buf =~ s/[\r\n]+\z//;
239
Removed:
240
Removed:
return $buf;
241
Removed:
}
242
Removed:
243
Removed:
1;
roles/daq-node/t/01-ezo-usb.t
@@ -9,7 +9,7 @@
9
9
use File::Temp qw(tempdir);
10
10
11
11
use FindBin;
12
Removed:
use lib "${FindBin::Bin}/../lib/perl5/";
12
Added:
use lib "${FindBin::Bin}/../lib/";
13
13
14
14
use FAPG::DAQ::EZO::USB qw(
15
15
discover_ezo_usb_device
roles/daq-node/t/03-mqtt.t
@@ -8,8 +8,8 @@
8
8
use JSON::PP qw(decode_json);
9
9
10
10
use FindBin;
11
Removed:
use lib "${FindBin::Bin}/../../../lib/perl5/";
12
Removed:
use lib "${FindBin::Bin}/../lib/perl5/";
11
Added:
use lib "${FindBin::Bin}/../../../lib/";
12
Added:
use lib "${FindBin::Bin}/../lib/";
13
13
14
14
use FAPG::DAQ::MQTT qw(
15
15
publish_mqtt