Add DAQ hub status.

Commit
77acf94a14661391c7ce80c5f87aa8d75291ab7f
Author
Marius Peter <dev@marius-peter.com>
Author date
Committer
Marius Peter <dev@marius-peter.com>
Committer date
Changed files
README.org
index 16940a88..8b71b456 100644..100644
@@ -35,6 +35,7 @@
35 35
36 36 Zero -> Broker : MQTT PUBLISH\nfapg/daq/<probe>/<node>/reading\nvalid reading JSON payload
37 37 Zero -> Broker : MQTT PUBLISH\nfapg/daq/<probe>/<node>/status\nnode/probe status JSON payload
38 Added: Broker -> Broker : MQTT PUBLISH\nfapg/daq/hub/<node>/status\nhub status JSON payload
38 39
39 40 note right of Zero
40 41 Example reading payload:
@@ -127,9 +128,15 @@
127 128 fapg/daq/<probe>/<node>/status
128 129 #+end_src
129 130
131 Added: The DAQ hub publishes its own heartbeat status on:
132 Added:
133 Added: #+begin_src text
134 Added: fapg/daq/hub/<node>/status
135 Added: #+end_src
136 Added:
130 137 Status payloads use schema =fapg.daq.status.v1=. A recent =ok=
131 Removed: status means the node is reachable and publishing valid readings. A
132 Removed: recent =probe_error= status means the node is reachable on the farm
133 Removed: network, but the probe or USB serial path is not returning a valid
134 Removed: reading. If the dashboard has no recent status for a node, it treats
135 Removed: that node as unreachable.
138 Added: status means the node is reachable; for probe nodes it also means the
139 Added: node is publishing valid readings. A recent =probe_error= status means
140 Added: the node is reachable on the farm network, but the probe or USB serial
141 Added: path is not returning a valid reading. If the dashboard has no recent
142 Added: status for a node or the DAQ hub, it treats that system as unreachable.
lib/perl5/FAPG/DAQ/MQTT.pm
index 00000000..c100e3ff 000000..100644
@@ -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;
roles/daq-hub/bin/fapg-daq-hub-status
index 00000000..877b67ce 000000..100755
@@ -0,0 +1,45 @@
1 Added: #!/usr/bin/env perl
2 Added: # -*- mode: perl-ts; -*-
3 Added:
4 Added: use v5.32.1;
5 Added: use strict;
6 Added: use warnings;
7 Added:
8 Added: use FindBin qw($Bin);
9 Added: use lib "$Bin/../../../lib/perl5";
10 Added:
11 Added: use FAPG::DAQ::MQTT qw(publish_status_mqtt);
12 Added: use Sys::Hostname qw(hostname);
13 Added:
14 Added: my $publish_interval = positive_number_from_env( HUB_STATUS_INTERVAL => 5 );
15 Added: my $node = $ENV{FAPG_DAQ_HUB_NODE} // $ENV{FAPG_DAQ_NODE} // $ENV{DAQ_NODE} // hostname();
16 Added:
17 Added: while (1) {
18 Added: eval {
19 Added: publish_status_mqtt(
20 Added: {
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 Added:
34 Added: sleep $publish_interval;
35 Added: }
36 Added:
37 Added: sub positive_number_from_env {
38 Added: my ( $name, $default ) = @_;
39 Added: my $value = $ENV{$name};
40 Added:
41 Added: return $default
42 Added: if !defined $value || $value !~ /\A(?:\d+(?:\.\d*)?|\.\d+)\z/ || $value <= 0;
43 Added:
44 Added: return 0 + $value;
45 Added: }
roles/daq-hub/deploy
index 33dca951..9d452124 100755..100755
@@ -5,7 +5,20 @@
5 5 SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
6 6 readonly SCRIPT_DIR
7 7 readonly PASSWD_FILE="/etc/mosquitto/passwd"
8 Removed: readonly SERVICE_NAME="mosquitto.service"
8 Added: readonly BROKER_SERVICE_NAME="mosquitto.service"
9 Added: readonly HUB_STATUS_SERVICE_NAME="fapg-daq-hub-status.service"
10 Added: readonly HUB_STATUS_SERVICE_FILE="/etc/systemd/system/${HUB_STATUS_SERVICE_NAME}"
11 Added: readonly REPO_DIR="/opt/fapg/fapg-daq"
12 Added: readonly HUB_STATUS_BIN="${REPO_DIR}/roles/daq-hub/bin/fapg-daq-hub-status"
13 Added: readonly ENV_DIR="/etc/fapg-daq"
14 Added: readonly HUB_STATUS_ENV_FILE="${ENV_DIR}/hub-status.env"
15 Added: readonly RUN_USER="fapg-daq"
16 Added: readonly RUN_GROUP="fapg-daq"
17 Added: readonly MQTT_HOST="127.0.0.1"
18 Added: readonly MQTT_PORT="1883"
19 Added: readonly MQTT_USERNAME="fapg_hub"
20 Added: readonly HUB_STATUS_INTERVAL="5"
21 Added: HUB_MQTT_PASSWORD=""
9 22
10 23 function die() {
11 24 echo "ERROR: $*" >&2
@@ -41,11 +54,101 @@
41 54 fi
42 55 }
43 56
44 Removed: require_root
57 Added: function read_hub_mqtt_password() {
58 Added: if [[ -n "${HUB_MQTT_PASSWORD:-}" ]]; then
59 Added: return
60 Added: fi
45 61
46 Removed: os_dependencies=(mosquitto mosquitto-clients)
47 Removed: apt update && apt install -y "${os_dependencies[@]}"
62 Added: read -rsp "MQTT password for user '${MQTT_USERNAME}': " HUB_MQTT_PASSWORD
63 Added: echo
48 64
65 Added: [[ -n "${HUB_MQTT_PASSWORD:-}" ]] || die "MQTT password is required."
66 Added: }
67 Added:
68 Added: function ensure_hub_mosquitto_password() {
69 Added: if passwd_has_user "${MQTT_USERNAME}"; then
70 Added: echo "Mosquitto password for ${MQTT_USERNAME} already exists; leaving it unchanged."
71 Added: return
72 Added: fi
73 Added:
74 Added: read_hub_mqtt_password
75 Added:
76 Added: if [[ -f "${PASSWD_FILE}" ]]; then
77 Added: mosquitto_passwd -b "${PASSWD_FILE}" "${MQTT_USERNAME}" "${HUB_MQTT_PASSWORD}"
78 Added: else
79 Added: mosquitto_passwd -b -c "${PASSWD_FILE}" "${MQTT_USERNAME}" "${HUB_MQTT_PASSWORD}"
80 Added: fi
81 Added: }
82 Added:
83 Added: function install_os_dependencies() {
84 Added: apt update && apt install -y \
85 Added: perl \
86 Added: mosquitto \
87 Added: mosquitto-clients \
88 Added: libnet-mqtt-simple-perl
89 Added: }
90 Added:
91 Added: function create_runtime_user() {
92 Added: if ! getent group "${RUN_GROUP}" > /dev/null; then
93 Added: addgroup --system "${RUN_GROUP}"
94 Added: fi
95 Added:
96 Added: if ! id "${RUN_USER}" > /dev/null 2>&1; then
97 Added: adduser \
98 Added: --system \
99 Added: --ingroup "${RUN_GROUP}" \
100 Added: --home /var/lib/fapg-daq \
101 Added: --no-create-home \
102 Added: --disabled-login \
103 Added: "${RUN_USER}"
104 Added: fi
105 Added: }
106 Added:
107 Added: function check_repo_layout() {
108 Added: [[ -d "${REPO_DIR}" ]] || die "Repo not found at ${REPO_DIR}."
109 Added: [[ -x "${HUB_STATUS_BIN}" ]] || die "${HUB_STATUS_BIN} does not exist or is not executable."
110 Added: }
111 Added:
112 Added: function write_hub_status_env_file() {
113 Added: mkdir -p "${ENV_DIR}"
114 Added:
115 Added: if [[ -f "${HUB_STATUS_ENV_FILE}" ]]; then
116 Added: echo "${HUB_STATUS_ENV_FILE} already exists; leaving it unchanged."
117 Added: chown root:"${RUN_GROUP}" "${HUB_STATUS_ENV_FILE}"
118 Added: chmod 0640 "${HUB_STATUS_ENV_FILE}"
119 Added: return
120 Added: fi
121 Added:
122 Added: read_hub_mqtt_password
123 Added:
124 Added: cat > "${HUB_STATUS_ENV_FILE}" <<EOF
125 Added: MQTT_SIMPLE_ALLOW_INSECURE_LOGIN=1
126 Added:
127 Added: MQTT_HOST=${MQTT_HOST}
128 Added: MQTT_PORT=${MQTT_PORT}
129 Added: MQTT_USERNAME=${MQTT_USERNAME}
130 Added: MQTT_PASSWORD=${HUB_MQTT_PASSWORD}
131 Added:
132 Added: HUB_STATUS_INTERVAL=${HUB_STATUS_INTERVAL}
133 Added: EOF
134 Added:
135 Added: chown root:"${RUN_GROUP}" "${HUB_STATUS_ENV_FILE}"
136 Added: chmod 0640 "${HUB_STATUS_ENV_FILE}"
137 Added: }
138 Added:
139 Added: function install_systemd_services() {
140 Added: install -d -m 0755 -o root -g root /etc/systemd/system
141 Added:
142 Added: install -m 0644 -o root -g root \
143 Added: "${SCRIPT_DIR}/etc/systemd/system/fapg-daq-hub-status.service" \
144 Added: "${HUB_STATUS_SERVICE_FILE}"
145 Added: }
146 Added:
147 Added: require_root
148 Added: install_os_dependencies
149 Added: create_runtime_user
150 Added: check_repo_layout
151 Added:
49 152 # Install files
50 153
51 154 install -d -m 0755 -o root -g root /etc/mosquitto/conf.d
@@ -61,6 +164,7 @@
61 164 # Password management
62 165
63 166 ensure_mosquitto_password fapg_zero
167 Added: ensure_hub_mosquitto_password
64 168 ensure_mosquitto_password fapg_vps
65 169
66 170 chown root:mosquitto "${PASSWD_FILE}"
@@ -68,7 +172,13 @@
68 172
69 173 # System service activation
70 174
71 Removed: systemctl enable "${SERVICE_NAME}"
72 Removed: systemctl restart "${SERVICE_NAME}"
175 Added: write_hub_status_env_file
176 Added: install_systemd_services
73 177
74 Removed: systemctl --no-pager --full status "${SERVICE_NAME}"
178 Added: systemctl daemon-reload
179 Added: systemctl enable "${BROKER_SERVICE_NAME}" "${HUB_STATUS_SERVICE_NAME}"
180 Added: systemctl restart "${BROKER_SERVICE_NAME}"
181 Added: systemctl restart "${HUB_STATUS_SERVICE_NAME}"
182 Added:
183 Added: systemctl --no-pager --full status "${BROKER_SERVICE_NAME}"
184 Added: systemctl --no-pager --full status "${HUB_STATUS_SERVICE_NAME}"
roles/daq-hub/etc/mosquitto/acl
index 01644afa..54cb001e 100644..100644
@@ -3,5 +3,8 @@
3 3 topic write fapg/status/#
4 4 topic read fapg/cmd/#
5 5
6 Added: user fapg_hub
7 Added: topic write fapg/daq/hub/+/status
8 Added:
6 9 user fapg_vps
7 10 topic read fapg/#
roles/daq-hub/etc/systemd/system/fapg-daq-hub-status.service
index 00000000..d74a846d 000000..100644
@@ -0,0 +1,24 @@
1 Added: [Unit]
2 Added: Description=FAPG DAQ hub MQTT status heartbeat
3 Added: Wants=network-online.target mosquitto.service
4 Added: After=network-online.target mosquitto.service
5 Added:
6 Added: [Service]
7 Added: Type=simple
8 Added: User=fapg-daq
9 Added: Group=fapg-daq
10 Added:
11 Added: WorkingDirectory=/opt/fapg/fapg-daq
12 Added: EnvironmentFile=/etc/fapg-daq/hub-status.env
13 Added:
14 Added: ExecStart=/opt/fapg/fapg-daq/roles/daq-hub/bin/fapg-daq-hub-status
15 Added:
16 Added: Restart=always
17 Added: RestartSec=5
18 Added:
19 Added: NoNewPrivileges=true
20 Added: ProtectHome=true
21 Added: ProtectSystem=strict
22 Added:
23 Added: [Install]
24 Added: WantedBy=multi-user.target
roles/daq-node/lib/perl5/FAPG/DAQ/MQTT.pm
index f4d2fc74..00000000 100644..000000
@@ -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 qw(encode_json);
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-node/t/03-mqtt.t
index 00d4415a..5fc49c25 100644..100644
@@ -8,6 +8,7 @@
8 8 use JSON::PP qw(decode_json);
9 9
10 10 use FindBin;
11 Added: use lib "${FindBin::Bin}/../../../lib/perl5/";
11 12 use lib "${FindBin::Bin}/../lib/perl5/";
12 13
13 14 use FAPG::DAQ::MQTT qw(
@@ -102,6 +103,39 @@
102 103 is $decoded->{error}, 'timeout waiting for EZO response';
103 104 is $decoded->{device}, '/dev/serial/by-id/usb-FTDI_USB_UART';
104 105 is $decoded->{source}, 'fapg-daq-node';
106 Added: };
107 Added:
108 Added: subtest 'hub status shape' => sub {
109 Added: my $topic = mqtt_topic_for_status(
110 Added: {
111 Added: probe => 'hub',
112 Added: status => 'ok',
113 Added: },
114 Added: node => 'fapg-daq-five-01',
115 Added: );
116 Added:
117 Added: is $topic, 'fapg/daq/hub/fapg-daq-five-01/status';
118 Added:
119 Added: my $payload = mqtt_payload_for_status(
120 Added: {
121 Added: probe => 'hub',
122 Added: status => 'ok',
123 Added: message => 'DAQ hub MQTT broker reachable',
124 Added: timestamp => '2026-07-06T10:00:10Z',
125 Added: },
126 Added: node => 'fapg-daq-five-01',
127 Added: source => 'fapg-daq-hub',
128 Added: );
129 Added:
130 Added: my $decoded = decode_json($payload);
131 Added:
132 Added: is $decoded->{schema}, 'fapg.daq.status.v1';
133 Added: is $decoded->{timestamp}, '2026-07-06T10:00:10Z';
134 Added: is $decoded->{node}, 'fapg-daq-five-01';
135 Added: is $decoded->{probe}, 'hub';
136 Added: is $decoded->{status}, 'ok';
137 Added: is $decoded->{message}, 'DAQ hub MQTT broker reachable';
138 Added: is $decoded->{source}, 'fapg-daq-hub';
105 139 };
106 140
107 141 subtest 'publish uses injected client, so unit tests need no broker' => sub {
roles/dashboard/lib/dashboard.pm
index 4f48167a..aeb8c7fb 100644..100644
@@ -60,6 +60,18 @@
60 60 }
61 61 );
62 62
63 Added: $self->helper(
64 Added: status_items => sub {
65 Added: return [
66 Added: {
67 Added: key => 'hub',
68 Added: label => 'DAQ Hub',
69 Added: },
70 Added: $self->probes->@*,
71 Added: ];
72 Added: }
73 Added: );
74 Added:
63 75 my $r = $self->routes;
64 76
65 77 $r->get('/')->to('dashboard#index');
roles/dashboard/lib/dashboard/Controller/Dashboard.pm
index 89c0bd15..eaa40969 100644..100644
@@ -5,8 +5,9 @@
5 5
6 6 sub index ($self) {
7 7 $self->render(
8 Removed: template => 'dashboard/index',
9 Removed: probes => $self->probes,
8 Added: template => 'dashboard/index',
9 Added: probes => $self->probes,
10 Added: status_items => $self->status_items,
10 11 );
11 12 }
12 13
roles/dashboard/lib/dashboard/Controller/Reading.pm
index 1618a567..700916c2 100644..100644
@@ -60,12 +60,12 @@
60 60 sub status ($self) {
61 61 my $probe = $self->param('probe') // '';
62 62
63 Removed: my %known = map { $_->{key} => 1 } $self->probes->@*;
63 Added: my %known = map { $_->{key} => 1 } $self->status_items->@*;
64 64
65 65 return $self->render(
66 66 status => 404,
67 67 json => {
68 Removed: error => "Unknown probe type: $probe",
68 Added: error => "Unknown status item: $probe",
69 69 },
70 70 ) unless $known{$probe};
71 71
roles/dashboard/public/js/dashboard.js
index caf5bf91..c79ef266 100644..100644
@@ -30,7 +30,7 @@
30 30 let smoothingEnabled = true;
31 31
32 32 function probeRows() {
33 Removed: return Array.from(document.querySelectorAll("[data-latest-probe]"));
33 Added: return Array.from(document.querySelectorAll("[data-latest-status]"));
34 34 }
35 35
36 36 function activeProbe() {
@@ -314,7 +314,7 @@
314 314 async function loadOverview() {
315 315 const rows = probeRows();
316 316 const results = await Promise.allSettled(
317 Removed: rows.map(row => fetchLatestStatus(row.dataset.latestProbe))
317 Added: rows.map(row => fetchLatestStatus(row.dataset.latestStatus))
318 318 );
319 319
320 320 rows.forEach((row, index) => {
roles/dashboard/t/basic.t
index a1f71934..822cb6d4 100644..100644
@@ -112,9 +112,27 @@
112 112 '{}',
113 113 );
114 114
115 Added: $t->app->sqlite->db->query(
116 Added: q{
117 Added: INSERT INTO node_status (
118 Added: received_at, topic, schema, timestamp, probe, node, status, message, payload
119 Added: ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
120 Added: },
121 Added: $fresh_status_at,
122 Added: 'fapg/daq/hub/fapg-daq-five-01/status',
123 Added: 'fapg.daq.status.v1',
124 Added: $fresh_status_at,
125 Added: 'hub',
126 Added: 'fapg-daq-five-01',
127 Added: 'ok',
128 Added: 'DAQ hub MQTT broker reachable',
129 Added: '{}',
130 Added: );
131 Added:
115 132 $t->get_ok('/')
116 133 ->status_is(200)
117 134 ->content_like(qr/FAPG DAQ Dashboard/)
135 Added: ->content_like(qr/DAQ Hub/)
118 136 ->content_like(qr/Dissolved Oxygen/)
119 137 ->content_like(qr/Hourly/)
120 138 ->content_like(qr/Daily/)
@@ -153,6 +171,12 @@
153 171 ->status_is(200)
154 172 ->json_is( '/probe' => 'ph' )
155 173 ->json_is( '/status/status' => 'ok' );
174 Added:
175 Added: $t->get_ok('/api/status/hub')
176 Added: ->status_is(200)
177 Added: ->json_is( '/probe' => 'hub' )
178 Added: ->json_is( '/status/status' => 'ok' )
179 Added: ->json_is( '/status/node' => 'fapg-daq-five-01' );
156 180
157 181 $t->get_ok('/api/status/do')
158 182 ->status_is(200)
roles/dashboard/templates/dashboard/index.html.ep
index a42a205e..c50499b8 100644..100644
@@ -10,12 +10,12 @@
10 10
11 11 <section class="status-band" aria-label="DAQ status overview">
12 12 <ul class="status-pill-list">
13 Removed: % for my $probe ($probes->@*) {
13 Added: % for my $item ($status_items->@*) {
14 14 <li
15 Removed: data-latest-probe="<%= $probe->{key} %>"
16 Removed: data-probe-label="<%= $probe->{label} %>">
15 Added: data-latest-status="<%= $item->{key} %>"
16 Added: data-status-label="<%= $item->{label} %>">
17 17 <span class="status-pill is-unknown" data-latest-state>
18 Removed: <span class="status-pill-probe"><%= $probe->{label} %></span>
18 Added: <span class="status-pill-probe"><%= $item->{label} %></span>
19 19 <span data-status-text>Unknown</span>
20 20 </span>
21 21 </li>