Update main script and modules for DAQ node role.

Commit
c34e790f5dd48b8891bc97e6dd162700eb2c2ccd
Author
Marius Peter <dev@marius-peter.com>
Author date
Committer
Marius Peter <dev@marius-peter.com>
Committer date
Changed files
roles/daq-node/bin/fapg-daq-node
index 55a6d599..ece75007 100755..100755
@@ -8,28 +8,40 @@
8 8 use FindBin;
9 9 use lib "${FindBin::Bin}/../lib/perl5";
10 10
11 Removed: use FAPG::DAQ::EZO;
11 Added: use FAPG::DAQ::EZO::USB qw(
12 Added: open_ezo_usb
13 Added: read_probe_info
14 Added: read_once
15 Added: );
12 16
17 Added: use FAPG::DAQ::MQTT qw(publish_mqtt);
18 Added:
13 19 =head1 NAME
14 20
15 21 fapg-daq - read sensor data and send over MQTT
16 22
17 23 =cut
18 24
19 Removed: my $probe = FAPG::DAQ::EZO::identify_probe();
25 Added: my $device = $ENV{EZO_SERIAL_DEVICE} // '/dev/ttyUSB0';
26 Added: my $serial = open_ezo_usb( device => $device );
27 Added: my $info = read_probe_info($serial);
20 28
21 Removed: print "Found probe "
22 Removed: . $probe->type
23 Removed: . " on "
24 Removed: . $probe->device
25 Removed: . ".\n";
29 Added: print <<"EOF";
30 Added: Device $device
31 Added: Probe $info->{probe}
32 Added: Firmware $info->{firmware}
33 Added: Unit $info->{unit}
34 Added: EOF
26 35
27 36 while (1) {
28 Removed: my $reading = $probe->get_reading();
29 Removed:
30 Removed: print "Reading: " . $reading->value . " " . $reading->unit . "\n";
31 Removed:
32 Removed: $reading->publish_mqtt();
33 Removed:
37 Added: try {
38 Added: my $reading = read_once( $serial, probe => $info->{probe} );
39 Added: say "$reading->{probe}: $reading->{value} $reading->{unit}";
40 Added: publish_mqtt($reading);
41 Added: } catch ($e) {
42 Added: warn "DAQ loop error: $e";
43 Added: sleep 30;
44 Added: next
45 Added: };
34 46 sleep 5;
35 47 }
roles/daq-node/lib/perl5/FAPG/DAQ/EZO.pm
index 83466d01..00000000 100644..000000
@@ -1,362 +0,0 @@
1 Removed: package FAPG::DAQ::EZO;
2 Removed: # -*- mode: perl-ts; -*-
3 Removed:
4 Removed: use 5.40.1;
5 Removed: use strict;
6 Removed: use warnings;
7 Removed:
8 Removed: use Device::SerialPort;
9 Removed: use Time::HiRes qw(time sleep);
10 Removed:
11 Removed: my %SUPPORTED_PROBES = (
12 Removed: ph => { unit => 'pH' },
13 Removed: do => { unit => 'mg/L' },
14 Removed: orp => { unit => 'mV' },
15 Removed: ec => { unit => 'uS/cm' },
16 Removed: );
17 Removed:
18 Removed: sub identify_probe {
19 Removed: my $device = _find_usb_device();
20 Removed: my $port = _open_serial_port($device);
21 Removed:
22 Removed: # Stop continuous mode if enabled.
23 Removed: # This makes later "R" commands deterministic.
24 Removed: _drain_serial($port, 0.25);
25 Removed: _write_ezo_command($port, 'C,0');
26 Removed: sleep 0.4;
27 Removed: _drain_serial($port, 0.25);
28 Removed:
29 Removed: my @lines = _ezo_command($port, 'I', 2.0);
30 Removed:
31 Removed: my ($raw_type, $firmware, $raw_reply);
32 Removed:
33 Removed: for my $line (@lines) {
34 Removed: if ($line =~ /^\?I,([^,\r\n]+),([^,\r\n]+)/i) {
35 Removed: $raw_type = $1;
36 Removed: $firmware = $2;
37 Removed: $raw_reply = $line;
38 Removed: last;
39 Removed: }
40 Removed: }
41 Removed:
42 Removed: die "Could not identify EZO probe.\n"
43 Removed: unless defined $raw_type;
44 Removed:
45 Removed: my $type = _normalize_probe_type($raw_type);
46 Removed:
47 Removed: die "Unsupported EZO probe type '$raw_type'.\n"
48 Removed: unless exists $SUPPORTED_PROBES{$type};
49 Removed:
50 Removed: return FAPG::DAQ::EZO::Probe->new(
51 Removed: device => $device,
52 Removed: port => $port,
53 Removed: type => $type,
54 Removed: raw_type => $raw_type,
55 Removed: firmware => $firmware,
56 Removed: unit => $SUPPORTED_PROBES{$type}->{unit},
57 Removed: raw_reply => $raw_reply,
58 Removed: );
59 Removed: }
60 Removed:
61 Removed: sub _find_usb_device {
62 Removed: return $ENV{EZO_SERIAL_DEVICE}
63 Removed: if defined $ENV{EZO_SERIAL_DEVICE} && length $ENV{EZO_SERIAL_DEVICE};
64 Removed:
65 Removed: my @devices = grep { /UART|FTDI|Atlas/i } glob '/dev/serial/by-id/*';
66 Removed:
67 Removed: die "Found no EZO USB serial device under /dev/serial/by-id/.\n"
68 Removed: unless @devices;
69 Removed:
70 Removed: die "Found multiple USB serial devices. Set EZO_SERIAL_DEVICE explicitly.\n"
71 Removed: if @devices > 1;
72 Removed:
73 Removed: return $devices[0];
74 Removed: }
75 Removed:
76 Removed: sub _open_serial_port {
77 Removed: my ($device) = @_;
78 Removed:
79 Removed: my $port = Device::SerialPort->new($device)
80 Removed: or die "Cannot open serial port $device: $!";
81 Removed:
82 Removed: $port->baudrate(9600);
83 Removed: $port->databits(8);
84 Removed: $port->parity('none');
85 Removed: $port->stopbits(1);
86 Removed: $port->handshake('none');
87 Removed:
88 Removed: $port->read_const_time(100);
89 Removed: $port->read_char_time(0);
90 Removed:
91 Removed: $port->write_settings
92 Removed: or die "Cannot apply serial settings to $device.\n";
93 Removed:
94 Removed: return $port;
95 Removed: }
96 Removed:
97 Removed: sub _ezo_command {
98 Removed: my ($port, $command, $timeout) = @_;
99 Removed:
100 Removed: _drain_serial($port, 0.05);
101 Removed: _write_ezo_command($port, $command);
102 Removed:
103 Removed: return _read_ezo_lines($port, $timeout);
104 Removed: }
105 Removed:
106 Removed: sub _write_ezo_command {
107 Removed: my ($port, $command) = @_;
108 Removed:
109 Removed: my $payload = "$command\r";
110 Removed: my $written = $port->write($payload);
111 Removed:
112 Removed: die "Could not write EZO command '$command'.\n"
113 Removed: unless defined $written && $written == length($payload);
114 Removed:
115 Removed: return 1;
116 Removed: }
117 Removed:
118 Removed: sub _read_ezo_lines {
119 Removed: my ($port, $timeout) = @_;
120 Removed:
121 Removed: my $deadline = time + $timeout;
122 Removed: my $buffer = '';
123 Removed: my @lines;
124 Removed:
125 Removed: while (time < $deadline) {
126 Removed: my ($count, $chunk) = $port->read(128);
127 Removed:
128 Removed: if ($count && defined $chunk) {
129 Removed: $buffer .= $chunk;
130 Removed:
131 Removed: while ($buffer =~ s/^([^\r]*)\r//) {
132 Removed: my $line = $1;
133 Removed: $line =~ s/^\s+//;
134 Removed: $line =~ s/\s+$//;
135 Removed:
136 Removed: push @lines, $line if length $line;
137 Removed: }
138 Removed: }
139 Removed:
140 Removed: sleep 0.02;
141 Removed: }
142 Removed:
143 Removed: return @lines;
144 Removed: }
145 Removed:
146 Removed: sub _drain_serial {
147 Removed: my ($port, $seconds) = @_;
148 Removed:
149 Removed: my $deadline = time + $seconds;
150 Removed:
151 Removed: while (time < $deadline) {
152 Removed: my ($count, undef) = $port->read(255);
153 Removed: sleep($count ? 0.01 : 0.05);
154 Removed: }
155 Removed: }
156 Removed:
157 Removed: sub _normalize_probe_type {
158 Removed: my ($raw_type) = @_;
159 Removed:
160 Removed: my $type = lc $raw_type;
161 Removed: $type =~ s/[^a-z0-9]//g;
162 Removed:
163 Removed: return $type;
164 Removed: }
165 Removed:
166 Removed: sub _parse_first_number {
167 Removed: my ($line) = @_;
168 Removed:
169 Removed: my ($first_field) = split /,/, $line;
170 Removed:
171 Removed: return undef unless defined $first_field;
172 Removed:
173 Removed: $first_field =~ s/^\s+//;
174 Removed: $first_field =~ s/\s+$//;
175 Removed:
176 Removed: return undef
177 Removed: unless $first_field =~ /^([+-]?(?:\d+(?:\.\d*)?|\.\d+))$/;
178 Removed:
179 Removed: return 0 + $1;
180 Removed: }
181 Removed:
182 Removed: 1;
183 Removed:
184 Removed: package FAPG::DAQ::EZO::Probe;
185 Removed: # -*- mode: perl-ts; -*-
186 Removed:
187 Removed: use 5.40.1;
188 Removed: use strict;
189 Removed: use warnings;
190 Removed:
191 Removed: sub new {
192 Removed: my ($class, %args) = @_;
193 Removed: return bless \%args, $class;
194 Removed: }
195 Removed:
196 Removed: sub device {
197 Removed: my ($self) = @_;
198 Removed: return $self->{device};
199 Removed: }
200 Removed:
201 Removed: sub port {
202 Removed: my ($self) = @_;
203 Removed: return $self->{port};
204 Removed: }
205 Removed:
206 Removed: sub type {
207 Removed: my ($self) = @_;
208 Removed: return $self->{type};
209 Removed: }
210 Removed:
211 Removed: sub raw_type {
212 Removed: my ($self) = @_;
213 Removed: return $self->{raw_type};
214 Removed: }
215 Removed:
216 Removed: sub firmware {
217 Removed: my ($self) = @_;
218 Removed: return $self->{firmware};
219 Removed: }
220 Removed:
221 Removed: sub unit {
222 Removed: my ($self) = @_;
223 Removed: return $self->{unit};
224 Removed: }
225 Removed:
226 Removed: sub get_reading {
227 Removed: my ($self) = @_;
228 Removed:
229 Removed: my @lines = FAPG::DAQ::EZO::_ezo_command($self->port, 'R', 3.0);
230 Removed:
231 Removed: for my $line (@lines) {
232 Removed: die "EZO probe returned error while reading: $line\n"
233 Removed: if $line =~ /^\*ER/i;
234 Removed:
235 Removed: next if $line =~ /^\*/;
236 Removed: next if $line =~ /^\?/;
237 Removed:
238 Removed: my $value = FAPG::DAQ::EZO::_parse_first_number($line);
239 Removed:
240 Removed: next unless defined $value;
241 Removed:
242 Removed: return FAPG::DAQ::EZO::Reading->new(
243 Removed: probe => $self->type,
244 Removed: value => $value,
245 Removed: unit => $self->unit,
246 Removed: device => $self->device,
247 Removed: raw_reply => $line,
248 Removed: );
249 Removed: }
250 Removed:
251 Removed: die "Could not parse EZO reading from response: " . join(' | ', @lines) . "\n";
252 Removed: }
253 Removed:
254 Removed: 1;
255 Removed:
256 Removed: package FAPG::DAQ::EZO::Reading;
257 Removed: # -*- mode: perl-ts; -*-
258 Removed:
259 Removed: use 5.40.1;
260 Removed: use strict;
261 Removed: use warnings;
262 Removed:
263 Removed: use JSON::PP qw(encode_json);
264 Removed: use Net::MQTT::Simple;
265 Removed: use POSIX qw(strftime);
266 Removed: use Sys::Hostname qw(hostname);
267 Removed:
268 Removed: sub new {
269 Removed: my ($class, %args) = @_;
270 Removed:
271 Removed: $args{timestamp} //= _utc_timestamp();
272 Removed: $args{node} //= $ENV{DAQ_NODE} || hostname();
273 Removed:
274 Removed: return bless \%args, $class;
275 Removed: }
276 Removed:
277 Removed: sub probe {
278 Removed: my ($self) = @_;
279 Removed: return $self->{probe};
280 Removed: }
281 Removed:
282 Removed: sub value {
283 Removed: my ($self) = @_;
284 Removed: return $self->{value};
285 Removed: }
286 Removed:
287 Removed: sub unit {
288 Removed: my ($self) = @_;
289 Removed: return $self->{unit};
290 Removed: }
291 Removed:
292 Removed: sub node {
293 Removed: my ($self) = @_;
294 Removed: return $self->{node};
295 Removed: }
296 Removed:
297 Removed: sub timestamp {
298 Removed: my ($self) = @_;
299 Removed: return $self->{timestamp};
300 Removed: }
301 Removed:
302 Removed: sub topic {
303 Removed: my ($self) = @_;
304 Removed:
305 Removed: return join '/',
306 Removed: 'fapg',
307 Removed: 'daq',
308 Removed: $self->probe,
309 Removed: $self->node,
310 Removed: 'reading';
311 Removed: }
312 Removed:
313 Removed: sub as_hash {
314 Removed: my ($self) = @_;
315 Removed:
316 Removed: return {
317 Removed: timestamp => $self->timestamp,
318 Removed: probe => $self->probe,
319 Removed: value => $self->value,
320 Removed: unit => $self->unit,
321 Removed: node => $self->node,
322 Removed: };
323 Removed: }
324 Removed:
325 Removed: sub as_json {
326 Removed: my ($self) = @_;
327 Removed: return encode_json($self->as_hash);
328 Removed: }
329 Removed:
330 Removed: sub publish_mqtt {
331 Removed: my ($self) = @_;
332 Removed:
333 Removed: my $host = $ENV{MQTT_HOST} // 'fapg-daq-five-01';
334 Removed: my $port = $ENV{MQTT_PORT} // 1883;
335 Removed:
336 Removed: my $server = "$host:$port";
337 Removed:
338 Removed: my $mqtt = Net::MQTT::Simple->new($server);
339 Removed:
340 Removed: my $username = exists $ENV{MQTT_USERNAME}
341 Removed: ? $ENV{MQTT_USERNAME}
342 Removed: : 'fapg_zero';
343 Removed:
344 Removed: if (defined $username && length $username) {
345 Removed: my $password = $ENV{MQTT_PASSWORD};
346 Removed:
347 Removed: die "MQTT_USERNAME is set but MQTT_PASSWORD is missing.\n"
348 Removed: unless defined $password && length $password;
349 Removed:
350 Removed: $mqtt->login($username, $password);
351 Removed: }
352 Removed:
353 Removed: $mqtt->publish($self->topic => $self->as_json);
354 Removed:
355 Removed: return 1;
356 Removed: }
357 Removed:
358 Removed: sub _utc_timestamp {
359 Removed: return strftime('%Y-%m-%dT%H:%M:%SZ', gmtime);
360 Removed: }
361 Removed:
362 Removed: 1;
roles/daq-node/lib/perl5/FAPG/DAQ/EZO/USB.pm
index 00000000..e2b8e92a 000000..100644
@@ -0,0 +1,200 @@
1 Added: # -*- mode: cperl; -*-
2 Added:
3 Added: package FAPG::DAQ::EZO::USB;
4 Added:
5 Added: use v5.40.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: open_ezo_usb
15 Added: drain_serial
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 Added: );
24 Added:
25 Added: my $DEFAULT_BAUDRATE = 9_600;
26 Added: my $DEFAULT_TIMEOUT_S = 1.5;
27 Added: my $EOL = "\r";
28 Added:
29 Added: sub open_ezo_usb (%opt) {
30 Added: my $device = $opt{device} // croak 'open_ezo_usb requires device => "/dev/tty..."';
31 Added: my $baudrate = $opt{baudrate} // $DEFAULT_BAUDRATE;
32 Added: my $timeout_ms = $opt{timeout_ms} // 250;
33 Added:
34 Added: require Device::SerialPort;
35 Added:
36 Added: my $port = Device::SerialPort->new($device)
37 Added: or croak "cannot open serial device $device";
38 Added:
39 Added: $port->baudrate($baudrate) or croak "cannot set baudrate on $device";
40 Added: $port->databits(8) or croak "cannot set databits on $device";
41 Added: $port->parity('none') or croak "cannot set parity on $device";
42 Added: $port->stopbits(1) or croak "cannot set stopbits on $device";
43 Added: $port->handshake('none') or croak "cannot set handshake on $device";
44 Added: $port->read_char_time(0);
45 Added: $port->read_const_time($timeout_ms);
46 Added: $port->write_settings or croak "cannot apply serial settings on $device";
47 Added:
48 Added: return $port;
49 Added: }
50 Added:
51 Added: sub drain_serial ( $io, %opt ) {
52 Added: my $max_s = $opt{max_s} // 0.20;
53 Added: my $chunk_sz = $opt{chunk_sz} // 255;
54 Added: my $deadline = time + $max_s;
55 Added: my $drained = q{};
56 Added:
57 Added: while ( time < $deadline ) {
58 Added: my ( $count, $chunk ) = $io->read($chunk_sz);
59 Added: last if not defined $count || $count == 0;
60 Added: $drained .= $chunk // q{};
61 Added: }
62 Added:
63 Added: return $drained;
64 Added: }
65 Added:
66 Added: sub ezo_command ( $io, $command, %opt ) {
67 Added: croak 'command must not contain carriage returns or newlines'
68 Added: if $command =~ /[\r\n]/;
69 Added:
70 Added: drain_serial( $io, max_s => ( $opt{drain_s} // 0.05 ) )
71 Added: if $opt{drain} // 1;
72 Added:
73 Added: my $bytes = $command . $EOL;
74 Added: my $wrote = $io->write($bytes);
75 Added:
76 Added: croak "short write to EZO device: wrote $wrote of " . length($bytes) . ' bytes'
77 Added: if not defined $wrote || $wrote != length($bytes);
78 Added:
79 Added: my $response = _read_ezo_line( $io, timeout_s => ( $opt{timeout_s} // $DEFAULT_TIMEOUT_S ), );
80 Added:
81 Added: croak "EZO command '$command' returned an error: $response"
82 Added: if $response =~ /\A\?(?:ERROR|ER)\z/i;
83 Added:
84 Added: return $response;
85 Added: }
86 Added:
87 Added: sub read_probe_info ( $io, %opt ) {
88 Added: my $raw = ezo_command( $io, 'I', %opt );
89 Added:
90 Added: my ( undef, $device, $firmware ) = split qr{,}, $raw, 3;
91 Added:
92 Added: croak "unexpected EZO info response: $raw"
93 Added: if not defined $device || $raw !~ /\A\?I,/;
94 Added:
95 Added: my $probe = normalize_probe_type($device);
96 Added:
97 Added: return {
98 Added: raw => $raw,
99 Added: device => $device,
100 Added: firmware => $firmware,
101 Added: probe => $probe,
102 Added: unit => unit_for_probe_type($probe),
103 Added: };
104 Added: }
105 Added:
106 Added: sub detect_probe_type ( $io, %opt ) {
107 Added: return read_probe_info( $io, %opt )->{probe};
108 Added: }
109 Added:
110 Added: sub read_once ( $io, %opt ) {
111 Added: my $probe = $opt{probe};
112 Added:
113 Added: $probe = detect_probe_type( $io, %opt ) if not defined $probe;
114 Added: $probe = normalize_probe_type($probe);
115 Added:
116 Added: my $raw = ezo_command( $io, 'R', %opt );
117 Added: my $parsed = parse_reading( $raw, probe => $probe );
118 Added:
119 Added: return {
120 Added: probe => $probe,
121 Added: unit => unit_for_probe_type($probe),
122 Added: raw => $raw,
123 Added: value => $parsed->{value},
124 Added: values => $parsed->{values},
125 Added: };
126 Added: }
127 Added:
128 Added: sub parse_reading ( $raw, %opt ) {
129 Added: croak 'empty EZO reading'
130 Added: if not defined $raw || $raw eq q{};
131 Added:
132 Added: my @fields = split /,/, $raw;
133 Added: my @values;
134 Added:
135 Added: for my $field (@fields) {
136 Added: croak "non-numeric EZO reading field '$field' in '$raw'"
137 Added: if $field !~ /\A[+-]?(?:\d+(?:\.\d*)?|\.\d+)\z/;
138 Added:
139 Added: push @values, 0 + $field;
140 Added: }
141 Added:
142 Added: return {
143 Added: raw => $raw,
144 Added: value => $values[0],
145 Added: values => \@values,
146 Added: };
147 Added: }
148 Added:
149 Added: sub normalize_probe_type ($probe) {
150 Added: croak 'probe type is required'
151 Added: if not defined $probe || $probe eq q{};
152 Added:
153 Added: my $p = lc $probe;
154 Added: $p =~ s/[\s_-]+//g;
155 Added:
156 Added: return
157 Added: ( $p eq 'ph' ) ? 'ph'
158 Added: : ( $p eq 'do' || $p eq 'dissolvedoxygen' ) ? 'do'
159 Added: : ( $p eq 'ec' || $p eq 'conductivity' ) ? 'ec'
160 Added: : ( $p eq 'orp' || $p eq 'redox' ) ? 'redox'
161 Added: : croak "unknown EZO probe type '$probe'";
162 Added: }
163 Added:
164 Added: sub unit_for_probe_type ($probe) {
165 Added: my $p = normalize_probe_type($probe);
166 Added:
167 Added: return 'pH' if $p eq 'ph';
168 Added: return 'mg/L' if $p eq 'do';
169 Added: return 'mV' if $p eq 'redox';
170 Added: return 'uS/cm' if $p eq 'ec';
171 Added:
172 Added: croak "unknown normalized probe type '$p'";
173 Added: }
174 Added:
175 Added: sub _read_ezo_line ( $io, %opt ) {
176 Added: my $timeout_s = $opt{timeout_s} // $DEFAULT_TIMEOUT_S;
177 Added: my $deadline = time + $timeout_s;
178 Added: my $buf = q{};
179 Added:
180 Added: while ( time < $deadline ) {
181 Added: my ( $count, $chunk ) = $io->read(1);
182 Added:
183 Added: if ( !defined $count || $count == 0 ) {
184 Added: sleep 0.01;
185 Added: next;
186 Added: }
187 Added:
188 Added: $buf .= $chunk;
189 Added: last if $buf =~ /\r\z/;
190 Added: }
191 Added:
192 Added: croak "timeout waiting for EZO response after ${timeout_s}s"
193 Added: if $buf !~ /\r\z/;
194 Added:
195 Added: $buf =~ s/[\r\n]+\z//;
196 Added:
197 Added: return $buf;
198 Added: }
199 Added:
200 Added: 1;
roles/daq-node/lib/perl5/FAPG/DAQ/MQTT.pm
index 00000000..98e72602 000000..100644
@@ -0,0 +1,161 @@
1 Added: # -*- mode: cperl; -*-
2 Added:
3 Added: package FAPG::DAQ::MQTT;
4 Added:
5 Added: use v5.40.1;
6 Added: use warnings;
7 Added:
8 Added: use Carp qw(croak);
9 Added: use Exporter 'import';
10 Added: use JSON::PP qw(encode_json);
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: mqtt_payload_for_reading
17 Added: mqtt_topic_for_reading
18 Added: mqtt_broker_endpoint
19 Added: );
20 Added:
21 Added: my $DEFAULT_HOST = 'fapg-daq-five-01';
22 Added: my $DEFAULT_PORT = 1883;
23 Added: my $DEFAULT_SCHEMA = 'fapg.daq.reading.v1';
24 Added: my $TOPIC_PREFIX = 'fapg/daq';
25 Added:
26 Added: sub publish_mqtt ($reading, %opt) {
27 Added: my $client = $opt{client} // _mqtt_client(%opt);
28 Added: my $topic = mqtt_topic_for_reading($reading, %opt);
29 Added: my $payload = mqtt_payload_for_reading($reading, %opt);
30 Added:
31 Added: $client->publish($topic => $payload);
32 Added:
33 Added: return {
34 Added: topic => $topic,
35 Added: payload => $payload,
36 Added: };
37 Added: }
38 Added:
39 Added: sub mqtt_payload_for_reading ($reading, %opt) {
40 Added: _assert_reading($reading);
41 Added:
42 Added: my $probe = lc $reading->{probe};
43 Added: my $node = _node_name($reading, %opt);
44 Added:
45 Added: my %payload = (
46 Added: schema => $opt{schema} // $DEFAULT_SCHEMA,
47 Added: timestamp => $reading->{timestamp} // _utc_timestamp(),
48 Added: node => $node,
49 Added: probe => $probe,
50 Added: value => 0 + $reading->{value},
51 Added: unit => $reading->{unit} // 'n/a',
52 Added: source => $opt{source} // 'fapg-daq-node',
53 Added: );
54 Added:
55 Added: $payload{raw} = $reading->{raw}
56 Added: if exists $reading->{raw} && defined $reading->{raw};
57 Added:
58 Added: $payload{values} = [ map { 0 + $_ } $reading->{values}->@* ]
59 Added: if ref($reading->{values}) eq 'ARRAY';
60 Added:
61 Added: return JSON::PP->new->canonical(1)->encode(\%payload);
62 Added: }
63 Added:
64 Added: sub mqtt_topic_for_reading ($reading, %opt) {
65 Added: _assert_reading($reading);
66 Added:
67 Added: my $prefix = $opt{topic_prefix} // $TOPIC_PREFIX;
68 Added: my $probe = _topic_level('probe', lc $reading->{probe});
69 Added: my $node = _topic_level('node', _node_name($reading, %opt));
70 Added:
71 Added: $prefix =~ s{/+\z}{};
72 Added:
73 Added: return join '/', $prefix, $probe, $node, 'reading';
74 Added: }
75 Added:
76 Added: sub mqtt_broker_endpoint (%opt) {
77 Added: my $host = $opt{host}
78 Added: // $ENV{FAPG_MQTT_HOST}
79 Added: // $ENV{MQTT_HOST}
80 Added: // $DEFAULT_HOST;
81 Added:
82 Added: my $port = $opt{port}
83 Added: // $ENV{FAPG_MQTT_PORT}
84 Added: // $ENV{MQTT_PORT}
85 Added: // $DEFAULT_PORT;
86 Added:
87 Added: return $host if $host =~ /:\d+\z/;
88 Added: return "$host:$port";
89 Added: }
90 Added:
91 Added: sub _mqtt_client (%opt) {
92 Added: my $endpoint = mqtt_broker_endpoint(%opt);
93 Added:
94 Added: my $username = $opt{username}
95 Added: // $ENV{FAPG_MQTT_USERNAME}
96 Added: // $ENV{MQTT_USERNAME}
97 Added: // 'fapg_zero';
98 Added:
99 Added: my $password = $opt{password}
100 Added: // $ENV{FAPG_MQTT_PASSWORD}
101 Added: // $ENV{MQTT_PASSWORD};
102 Added:
103 Added: state %clients;
104 Added:
105 Added: my $cache_key = join "\0", $endpoint, ($username // q{}), ($password // q{});
106 Added: return $clients{$cache_key} if exists $clients{$cache_key};
107 Added:
108 Added: require Net::MQTT::Simple;
109 Added:
110 Added: my $client = Net::MQTT::Simple->new($endpoint)
111 Added: or croak "cannot connect to MQTT broker $endpoint";
112 Added:
113 Added: if (defined $username && $username ne q{}) {
114 Added: croak 'MQTT password is required when MQTT username is set'
115 Added: if !defined $password;
116 Added:
117 Added: $client->login($username, $password);
118 Added: }
119 Added:
120 Added: return $clients{$cache_key} = $client;
121 Added: }
122 Added:
123 Added: sub _assert_reading ($reading) {
124 Added: croak 'reading must be a HASH reference'
125 Added: if ref($reading) ne 'HASH';
126 Added:
127 Added: croak 'reading requires a probe field'
128 Added: if !defined $reading->{probe} || $reading->{probe} eq q{};
129 Added:
130 Added: croak 'reading requires a value field'
131 Added: if !exists $reading->{value} || !defined $reading->{value};
132 Added:
133 Added: croak "reading value is not numeric: $reading->{value}"
134 Added: if $reading->{value} !~ /\A[+-]?(?:\d+(?:\.\d*)?|\.\d+)\z/;
135 Added:
136 Added: return;
137 Added: }
138 Added:
139 Added: sub _node_name ($reading, %opt) {
140 Added: return $reading->{node}
141 Added: // $opt{node}
142 Added: // $ENV{FAPG_DAQ_NODE}
143 Added: // $ENV{DAQ_NODE}
144 Added: // hostname();
145 Added: }
146 Added:
147 Added: sub _topic_level ($name, $value) {
148 Added: croak "$name topic level is required"
149 Added: if !defined $value || $value eq q{};
150 Added:
151 Added: croak "$name topic level must not contain '/', '+', '#', or NUL: $value"
152 Added: if $value =~ m{[\0/+\#]};
153 Added:
154 Added: return $value;
155 Added: }
156 Added:
157 Added: sub _utc_timestamp () {
158 Added: return strftime '%Y-%m-%dT%H:%M:%SZ', gmtime;
159 Added: }
160 Added:
161 Added: 1;