[Perl] DAQ system for the FAPG.
refactor isolate backend domain services
Move reading queries, analytics, and series caching out of controllers. Correct weighted rollups, gap markers, current weather selection, and refresh semantics. Harden application secret generation and persistence with focused regression tests.
Changed files
- roles/dashboard/lib/FAPG/DAQ/Dashboard.pm
- roles/dashboard/lib/FAPG/DAQ/Dashboard/Analytics.pm
- roles/dashboard/lib/FAPG/DAQ/Dashboard/Controller/Download.pm
- roles/dashboard/lib/FAPG/DAQ/Dashboard/Controller/Insights.pm
- roles/dashboard/lib/FAPG/DAQ/Dashboard/Controller/Pages.pm
- roles/dashboard/lib/FAPG/DAQ/Dashboard/Controller/Reading.pm
- roles/dashboard/lib/FAPG/DAQ/Dashboard/Controller/Weather.pm
- roles/dashboard/lib/FAPG/DAQ/Dashboard/Model/Reading.pm
- roles/dashboard/lib/FAPG/DAQ/Dashboard/Model/Weather.pm
- roles/dashboard/lib/FAPG/DAQ/Dashboard/Secret.pm
- roles/dashboard/lib/FAPG/DAQ/Dashboard/Service/SeriesCache.pm
- roles/dashboard/lib/FAPG/DAQ/Dashboard/Service/WeatherFetcher.pm
- roles/dashboard/lib/FAPG/DAQ/Dashboard/Timeframe.pm
- roles/dashboard/lib/FAPG/DAQ/Dashboard/Util.pm
- roles/dashboard/t/04-readings-series.t
- roles/dashboard/t/09-weather.t
- roles/dashboard/t/10-secret.t
roles/dashboard/lib/FAPG/DAQ/Dashboard.pm
@@ -2,7 +2,10 @@
2
2
use Mojo::Base 'Mojolicious', -signatures;
3
3
4
4
use Mojo::SQLite;
5
Added:
use FAPG::DAQ::Dashboard::Model::Reading;
5
6
use FAPG::DAQ::Dashboard::Model::Weather;
7
Added:
use FAPG::DAQ::Dashboard::Secret qw(load_or_create_secret read_secret);
8
Added:
use FAPG::DAQ::Dashboard::Service::SeriesCache;
6
9
use FAPG::DAQ::Dashboard::Service::WeatherFetcher;
7
10
8
11
# This method will run once at server start
@@ -14,33 +17,20 @@
14
17
15
18
# Auto-generate and persist a secret (never committed to VCS).
16
19
# Falls back to config secrets if present (backward compatibility).
17
Removed:
my $secret_path = $self->home . '/.mojo-secrets';
18
Removed:
my $secret;
19
Removed:
if ( -f $secret_path && -r $secret_path ) {
20
Removed:
open my $fh, '<', $secret_path;
21
Removed:
$secret = do { local $/; <$fh> };
22
Removed:
close $fh;
23
Removed:
chomp $secret if defined $secret;
20
Added:
my $secret_path = $ENV{FAPG_DAQ_SECRET_FILE}
21
Added:
// $self->home . '/.mojo-secrets';
22
Added:
my $stored_secret = read_secret($secret_path);
23
Added:
24
Added:
if ( defined $stored_secret && length($stored_secret) > 8 ) {
25
Added:
$self->secrets( [$stored_secret] );
24
26
}
25
Removed:
if ( $secret && length $secret > 8 ) {
26
Removed:
$self->secrets( [$secret] );
27
Removed:
}
28
27
elsif ( $config->{secrets} && ref $config->{secrets} eq 'ARRAY' ) {
29
28
$self->secrets( $config->{secrets} );
30
29
}
31
30
else {
32
Removed:
require Digest::SHA;
33
Removed:
$secret = Digest::SHA::sha256_hex( join '', time, $$, rand, rand );
34
Removed:
eval {
35
Removed:
open my $fh, '>', $secret_path or die "open: $!";
36
Removed:
print $fh $secret;
37
Removed:
close $fh;
38
Removed:
chmod 0600, $secret_path;
39
Removed:
$self->log->info("Generated new application secret in $secret_path");
40
Removed:
};
41
Removed:
if ($@) {
42
Removed:
$self->log->warn("Could not persist secret to $secret_path: $@");
43
Removed:
}
31
Added:
my ( $secret, $created ) = load_or_create_secret($secret_path);
32
Added:
$self->log->info("Generated new application secret in $secret_path")
33
Added:
if $created;
44
34
$self->secrets( [$secret] );
45
35
}
46
36
@@ -60,13 +50,21 @@
60
50
61
51
$self->helper( sqlite => sub {$sqlite} );
62
52
53
Added:
my $reading_model
54
Added:
= FAPG::DAQ::Dashboard::Model::Reading->new( sqlite => $sqlite );
55
Added:
my $series_cache = FAPG::DAQ::Dashboard::Service::SeriesCache->new;
56
Added:
57
Added:
$self->helper( reading => sub {$reading_model} );
58
Added:
$self->helper( series_cache => sub {$series_cache} );
59
Added:
63
60
# Weather model and fetcher
64
Removed:
my $weather_model = FAPG::DAQ::Dashboard::Model::Weather->new( sqlite => $sqlite );
61
Added:
my $weather_model
62
Added:
= FAPG::DAQ::Dashboard::Model::Weather->new( sqlite => $sqlite );
65
63
$weather_model->create_table;
66
64
67
65
$self->helper( weather => sub {$weather_model} );
68
66
69
Removed:
my $weather_config = $config->{weather} // {};
67
Added:
my $weather_config = $config->{weather} // {};
70
68
my $weather_fetcher = FAPG::DAQ::Dashboard::Service::WeatherFetcher->new(
71
69
latitude => $weather_config->{latitude} // 46.303,
72
70
longitude => $weather_config->{longitude} // 6.098,
@@ -75,53 +73,52 @@
75
73
76
74
$self->helper( weather_fetcher => sub {$weather_fetcher} );
77
75
78
Removed:
# Schedule periodic weather fetch (skip during tests with in-memory DB)
79
Removed:
if ( $weather_config->{enabled} && ( $ENV{FAPG_DAQ_DB} // '' ) ne ':memory:' ) {
80
Removed:
require Mojo::IOLoop;
81
Removed:
Mojo::IOLoop->timer(
82
Removed:
5 => sub {
83
Removed:
my $rows = $weather_fetcher->fetch;
84
Removed:
$weather_model->store_batch($rows);
85
Removed:
$self->log->info( 'Weather: initial fetch stored ' . scalar(@$rows) . ' points' );
86
Removed:
}
87
Removed:
);
88
Removed:
Mojo::IOLoop->recurring(
89
Removed:
3600 => sub {
90
Removed:
my $rows = $weather_fetcher->fetch;
91
Removed:
$weather_model->store_batch($rows);
92
Removed:
$self->log->info( 'Weather: periodic fetch stored ' . scalar(@$rows) . ' points' );
93
Removed:
}
94
Removed:
);
95
Removed:
}
96
Removed:
97
76
$self->helper(
98
77
probes => sub {
99
78
return [
100
Removed:
{ key => 'ph',
101
Removed:
label => 'pH',
102
Removed:
unit => 'pH',
103
Removed:
node => 'Pi Zero 01',
79
Added:
{ key => 'ph',
80
Added:
label => 'pH',
81
Added:
unit => 'pH',
82
Added:
node => 'Pi Zero 01',
83
Added:
good_min => 6.4,
84
Added:
good_max => 7.2,
104
85
},
105
Removed:
{ key => 'do',
106
Removed:
label => 'Dissolved Oxygen',
107
Removed:
unit => 'mg/L',
108
Removed:
node => 'Pi Zero 03',
86
Added:
{ key => 'do',
87
Added:
label => 'Dissolved Oxygen',
88
Added:
unit => 'mg/L',
89
Added:
node => 'Pi Zero 03',
90
Added:
good_min => 5.5,
91
Added:
good_max => 10,
109
92
},
110
Removed:
{ key => 'orp',
111
Removed:
label => 'ORP',
112
Removed:
unit => 'mV',
113
Removed:
node => 'Pi Zero 04',
93
Added:
{ key => 'orp',
94
Added:
label => 'ORP',
95
Added:
unit => 'mV',
96
Added:
node => 'Pi Zero 04',
97
Added:
good_min => 250,
98
Added:
good_max => 400,
114
99
},
115
Removed:
{ key => 'ec',
116
Removed:
label => 'Electrical Conductivity',
117
Removed:
unit => 'µS/cm',
118
Removed:
node => 'Pi Zero 02',
100
Added:
{ key => 'ec',
101
Added:
label => 'Electrical Conductivity',
102
Added:
unit => 'µS/cm',
103
Added:
node => 'Pi Zero 02',
104
Added:
good_min => 300,
105
Added:
good_max => 1500,
119
106
},
120
107
];
121
108
}
122
109
);
123
110
124
111
$self->helper(
112
Added:
probe_by_key => sub {
113
Added:
my ( $controller, $key ) = @_;
114
Added:
for my $probe ( $controller->probes->@* ) {
115
Added:
return $probe if $probe->{key} eq $key;
116
Added:
}
117
Added:
return;
118
Added:
}
119
Added:
);
120
Added:
121
Added:
$self->helper(
125
122
status_items => sub {
126
123
return [
127
124
{ key => 'hub',
@@ -132,6 +129,16 @@
132
129
}
133
130
);
134
131
132
Added:
$self->helper(
133
Added:
status_item_by_key => sub {
134
Added:
my ( $controller, $key ) = @_;
135
Added:
for my $item ( $controller->status_items->@* ) {
136
Added:
return $item if $item->{key} eq $key;
137
Added:
}
138
Added:
return;
139
Added:
}
140
Added:
);
141
Added:
135
142
my $r = $self->routes;
136
143
137
144
$r->get('/')->to('pages#index');
@@ -157,7 +164,7 @@
157
164
$r->get('/api/weather/hourly')->to('Weather#hourly');
158
165
$r->get('/api/weather/forecast')->to('Weather#forecast');
159
166
$r->get('/api/weather/current')->to('Weather#current');
160
Removed:
$r->get('/api/weather/refresh')->to('Weather#refresh');
167
Added:
$r->post('/api/weather/refresh')->to('Weather#refresh');
161
168
}
162
169
163
170
1;
roles/dashboard/lib/FAPG/DAQ/Dashboard/Analytics.pm
@@ -0,0 +1,188 @@
1
Added:
package FAPG::DAQ::Dashboard::Analytics;
2
Added:
use Mojo::Base -strict;
3
Added:
4
Added:
use Exporter qw(import);
5
Added:
use POSIX qw(strftime);
6
Added:
7
Added:
use FAPG::DAQ::Dashboard::Util qw(utc_timestamp);
8
Added:
9
Added:
our @EXPORT_OK = qw(
10
Added:
aggregate_series
11
Added:
derivative_points
12
Added:
diurnal_summary
13
Added:
stability_points
14
Added:
);
15
Added:
16
Added:
my $HOUR_SECONDS = 3_600;
17
Added:
18
Added:
sub aggregate_series {
19
Added:
my ( $rollups, $window, $gap_threshold ) = @_;
20
Added:
my $points = _series_points($window);
21
Added:
my @weighted_totals;
22
Added:
my $previous_sample_epoch;
23
Added:
24
Added:
for my $rollup (@$rollups) {
25
Added:
my $epoch = 0 + $rollup->{minute_epoch};
26
Added:
my $first = 0 + ( $rollup->{first_sample_epoch} // $epoch );
27
Added:
my $last = 0 + ( $rollup->{last_sample_epoch} // $epoch );
28
Added:
my $has_gap = defined $previous_sample_epoch
29
Added:
&& $first - $previous_sample_epoch > $gap_threshold;
30
Added:
$has_gap ||= ( $rollup->{max_gap_seconds} // 0 ) > $gap_threshold;
31
Added:
$previous_sample_epoch = $last
32
Added:
if !defined $previous_sample_epoch
33
Added:
|| $last > $previous_sample_epoch;
34
Added:
35
Added:
next
36
Added:
if $epoch < $window->{start_epoch}
37
Added:
|| $epoch >= $window->{end_epoch};
38
Added:
39
Added:
my $point_index = int(
40
Added:
( $epoch - $window->{start_epoch} ) / $window->{bucket_stride} );
41
Added:
my $point = $points->[$point_index];
42
Added:
my $count = 0 + $rollup->{count};
43
Added:
44
Added:
$weighted_totals[$point_index]
45
Added:
+= ( 0 + $rollup->{avg_value} ) * $count;
46
Added:
$point->{min} = 0 + $rollup->{value_min}
47
Added:
if !defined $point->{min}
48
Added:
|| $rollup->{value_min} < $point->{min};
49
Added:
$point->{max} = 0 + $rollup->{value_max}
50
Added:
if !defined $point->{max}
51
Added:
|| $rollup->{value_max} > $point->{max};
52
Added:
$point->{count} += $count;
53
Added:
$point->{gap} = 1 if $has_gap;
54
Added:
}
55
Added:
56
Added:
for my $i ( 0 .. $#$points ) {
57
Added:
next if !$points->[$i]{count};
58
Added:
$points->[$i]{value}
59
Added:
= $weighted_totals[$i] / $points->[$i]{count};
60
Added:
}
61
Added:
62
Added:
return $points;
63
Added:
}
64
Added:
65
Added:
sub diurnal_summary {
66
Added:
my ($rollups) = @_;
67
Added:
my %by_day;
68
Added:
69
Added:
for my $row (@$rollups) {
70
Added:
my @time = gmtime $row->{bucket_epoch};
71
Added:
my $date = strftime( '%Y-%m-%d', @time );
72
Added:
$by_day{$date}[ $time[2] ] = 0 + $row->{avg_value};
73
Added:
}
74
Added:
75
Added:
my @traces;
76
Added:
my @sums = (0) x 24;
77
Added:
my @counts = (0) x 24;
78
Added:
79
Added:
for my $date ( sort keys %by_day ) {
80
Added:
my @trace;
81
Added:
for my $hour ( 0 .. 23 ) {
82
Added:
my $value = $by_day{$date}[$hour];
83
Added:
push @trace, $value;
84
Added:
if ( defined $value ) {
85
Added:
$sums[$hour] += $value;
86
Added:
++$counts[$hour];
87
Added:
}
88
Added:
}
89
Added:
push @traces, { date => $date, values => \@trace };
90
Added:
}
91
Added:
92
Added:
my @mean = map { $counts[$_] ? $sums[$_] / $counts[$_] : undef } 0 .. 23;
93
Added:
94
Added:
return ( \@traces, \@mean );
95
Added:
}
96
Added:
97
Added:
sub derivative_points {
98
Added:
my ( $rollups, $probes ) = @_;
99
Added:
my %by_probe;
100
Added:
push $by_probe{ $_->{probe} }->@*, $_ for @$rollups;
101
Added:
102
Added:
my %points;
103
Added:
for my $probe (@$probes) {
104
Added:
my $probe_rollups = $by_probe{$probe} // [];
105
Added:
my @probe_points;
106
Added:
107
Added:
for my $i ( 1 .. $#$probe_rollups ) {
108
Added:
my $previous = $probe_rollups->[ $i - 1 ];
109
Added:
my $current = $probe_rollups->[$i];
110
Added:
my $elapsed
111
Added:
= $current->{bucket_epoch} - $previous->{bucket_epoch};
112
Added:
next if $elapsed <= 0;
113
Added:
114
Added:
my $rate = ( $current->{avg_value} - $previous->{avg_value} )
115
Added:
/ ( $elapsed / $HOUR_SECONDS );
116
Added:
push @probe_points,
117
Added:
{
118
Added:
timestamp => utc_timestamp( $current->{bucket_epoch} ),
119
Added:
value => 0 + sprintf( '%.4f', $rate ),
120
Added:
};
121
Added:
}
122
Added:
123
Added:
$points{$probe} = \@probe_points;
124
Added:
}
125
Added:
126
Added:
return \%points;
127
Added:
}
128
Added:
129
Added:
sub stability_points {
130
Added:
my ( $rollups, $ranges, $window_size ) = @_;
131
Added:
my %by_epoch;
132
Added:
133
Added:
for my $row (@$rollups) {
134
Added:
my $range = $ranges->{ $row->{probe} } or next;
135
Added:
my $span = $range->{max} - $range->{min};
136
Added:
$by_epoch{ $row->{bucket_epoch} }{ $row->{probe} }
137
Added:
= ( $row->{avg_value} - $range->{min} ) / $span;
138
Added:
}
139
Added:
140
Added:
my @epochs = sort { $a <=> $b } keys %by_epoch;
141
Added:
my @points;
142
Added:
for my $i ( $window_size - 1 .. $#epochs ) {
143
Added:
my @values;
144
Added:
for my $epoch ( @epochs[ $i - $window_size + 1 .. $i ] ) {
145
Added:
push @values, values $by_epoch{$epoch}->%*;
146
Added:
}
147
Added:
next if @values < 2;
148
Added:
149
Added:
my $mean = 0;
150
Added:
$mean += $_ for @values;
151
Added:
$mean /= @values;
152
Added:
153
Added:
my $variance = 0;
154
Added:
$variance += ( $_ - $mean )**2 for @values;
155
Added:
$variance /= @values;
156
Added:
157
Added:
push @points,
158
Added:
{
159
Added:
timestamp => utc_timestamp( $epochs[$i] ),
160
Added:
score => 0 + sprintf( '%.4f', sqrt $variance ),
161
Added:
};
162
Added:
}
163
Added:
164
Added:
return \@points;
165
Added:
}
166
Added:
167
Added:
sub _series_points {
168
Added:
my ($window) = @_;
169
Added:
my @points;
170
Added:
171
Added:
for my $index ( 0 .. $window->{bucket_count} - 1 ) {
172
Added:
push @points,
173
Added:
{
174
Added:
timestamp => utc_timestamp(
175
Added:
$window->{start_epoch} + ( $index * $window->{bucket_stride} )
176
Added:
),
177
Added:
value => undef,
178
Added:
min => undef,
179
Added:
max => undef,
180
Added:
count => 0,
181
Added:
gap => 0,
182
Added:
};
183
Added:
}
184
Added:
185
Added:
return \@points;
186
Added:
}
187
Added:
188
Added:
1;
roles/dashboard/lib/FAPG/DAQ/Dashboard/Controller/Download.pm
@@ -3,6 +3,7 @@
3
3
4
4
use Excel::Writer::XLSX;
5
5
use FAPG::DAQ::Dashboard::Timeframe qw(max_series_range series_window);
6
Added:
use FAPG::DAQ::Dashboard::Util qw(utc_timestamp);
6
7
use POSIX qw(strftime);
7
8
use Time::Local qw(timegm);
8
9
@@ -23,7 +24,6 @@
23
24
my @probes = $self->every_param('probe')->@*;
24
25
my $from_epoch = date_epoch($from);
25
26
my $to_epoch = date_epoch($to);
26
Removed:
my %known = map { $_->{key} => 1 } $self->probes->@*;
27
27
28
28
return $self->render(
29
29
status => 400,
@@ -38,22 +38,12 @@
38
38
return $self->render(
39
39
status => 400,
40
40
text => 'Every selected probe must be a known probe type.',
41
Removed:
) if grep { !$known{$_} } @probes;
41
Added:
) if grep { !$self->probe_by_key($_) } @probes;
42
42
43
Removed:
my $probe_filter
44
Removed:
= @probes
45
Removed:
? ' AND probe IN (' . join( ',', ('?') x @probes ) . ')'
46
Removed:
: '';
47
Removed:
my $rows = raw_readings(
48
Removed:
$self,
49
Removed:
'COALESCE(received_at, timestamp) >= ? '
50
Removed:
. 'AND COALESCE(received_at, timestamp) < ?'
51
Removed:
. $probe_filter,
52
Removed:
'COALESCE(received_at, timestamp)',
43
Added:
my $rows
44
Added:
= $self->reading->readings_by_received_range(
53
45
utc_timestamp($from_epoch),
54
Removed:
utc_timestamp( $to_epoch + ( 24 * 60 * 60 ) ),
55
Removed:
@probes,
56
Removed:
);
46
Added:
utc_timestamp( $to_epoch + ( 24 * 60 * 60 ) ), \@probes, );
57
47
my $probe_name = join '-', @probes;
58
48
my $filename
59
49
= !@probes
@@ -67,12 +57,11 @@
67
57
my $probe = $self->param('probe') // '';
68
58
my $timeframe = $self->param('timeframe') // '';
69
59
my $range = $self->param('range') // 1;
70
Removed:
my %known = map { $_->{key} => 1 } $self->probes->@*;
71
60
72
61
return $self->render(
73
62
status => 400,
74
63
text => 'A known probe is required for timeframe downloads.',
75
Removed:
) unless $known{$probe};
64
Added:
) unless $self->probe_by_key($probe);
76
65
77
66
return $self->render(
78
67
status => 400,
@@ -89,40 +78,13 @@
89
78
text => "Unknown timeframe: $timeframe",
90
79
) unless defined $window;
91
80
92
Removed:
my $rows = raw_readings(
93
Removed:
$self,
94
Removed:
'probe = ? AND COALESCE(timestamp, received_at) >= ? '
95
Removed:
. 'AND COALESCE(timestamp, received_at) < ?',
96
Removed:
'COALESCE(timestamp, received_at)',
97
Removed:
$probe,
98
Removed:
$window->{start_iso},
99
Removed:
$window->{end_iso},
100
Removed:
);
81
Added:
my $rows = $self->reading->readings_by_measurement_range( $probe,
82
Added:
$window->{start_iso}, $window->{end_iso}, );
101
83
my $filename = "fapg-daq-$probe-$timeframe-$range.$format";
102
84
103
85
render_download( $self, $rows, $filename, $format, [$probe] );
104
86
}
105
87
106
Removed:
sub raw_readings ( $self, $where, $order_by, @bind ) {
107
Removed:
my $database = $self->sqlite->db;
108
Removed:
my $table
109
Removed:
= $database->query(
110
Removed:
q{SELECT name FROM sqlite_master WHERE type = 'table' AND name = 'readings'}
111
Removed:
)->hash;
112
Removed:
113
Removed:
return [] unless $table;
114
Removed:
115
Removed:
return $database->query(
116
Removed:
qq{
117
Removed:
SELECT timestamp, received_at, node, probe, value, unit
118
Removed:
FROM readings
119
Removed:
WHERE $where
120
Removed:
ORDER BY $order_by, id
121
Removed:
},
122
Removed:
@bind,
123
Removed:
)->hashes->to_array;
124
Removed:
}
125
Removed:
126
88
sub render_download ( $self, $rows, $filename, $format, $selected_probes ) {
127
89
my $pivoted = $selected_probes->@* != 1;
128
90
my ( $headers, $table )
@@ -251,7 +213,7 @@
251
213
}
252
214
253
215
sub spreadsheet_timestamp ($timestamp) {
254
Removed:
return undef if !defined $timestamp;
216
Added:
return if !defined $timestamp;
255
217
256
218
$timestamp =~ s/T/ /;
257
219
$timestamp =~ s/Z\z//;
@@ -268,20 +230,16 @@
268
230
}
269
231
270
232
sub date_epoch ($date) {
271
Removed:
return undef unless $date =~ /\A(\d{4})-(\d{2})-(\d{2})\z/;
233
Added:
return unless $date =~ /\A(\d{4})-(\d{2})-(\d{2})\z/;
272
234
273
235
my ( $year, $month, $day ) = ( $1, $2, $3 );
274
236
my $epoch = eval { timegm( 0, 0, 0, $day, $month - 1, $year ) };
275
237
276
Removed:
return undef if !defined $epoch || $@;
238
Added:
return if !defined $epoch || $@;
277
239
278
Removed:
return undef unless strftime( '%Y-%m-%d', gmtime($epoch) ) eq $date;
240
Added:
return unless strftime( '%Y-%m-%d', gmtime($epoch) ) eq $date;
279
241
280
242
return $epoch;
281
Removed:
}
282
Removed:
283
Removed:
sub utc_timestamp ($epoch) {
284
Removed:
return strftime( '%Y-%m-%dT%H:%M:%SZ', gmtime $epoch );
285
243
}
286
244
287
245
1;
roles/dashboard/lib/FAPG/DAQ/Dashboard/Controller/Insights.pm
@@ -1,9 +1,11 @@
1
1
package FAPG::DAQ::Dashboard::Controller::Insights;
2
2
use Mojo::Base 'Mojolicious::Controller', -signatures;
3
3
4
Removed:
use POSIX qw(strftime);
4
Added:
use FAPG::DAQ::Dashboard::Analytics qw(
5
Added:
derivative_points diurnal_summary stability_points
6
Added:
);
7
Added:
use FAPG::DAQ::Dashboard::Util qw(bounded_integer);
5
8
6
Removed:
my $WEEK_SECONDS = 7 * 24 * 60 * 60;
7
9
my $HOUR_SECONDS = 3_600;
8
10
my $DAY_SECONDS = 24 * 60 * 60;
9
11
@@ -11,75 +13,25 @@
11
13
# Returns 24 hourly slots with per-day traces for the last 7 days.
12
14
sub diurnal ($self) {
13
15
my $probe = $self->param('probe') // '';
14
Removed:
my %known = map { $_->{key} => 1 } $self->probes->@*;
15
16
16
17
return $self->render(
17
18
status => 404,
18
19
json => { error => "Unknown probe: $probe" },
19
Removed:
) unless $known{$probe};
20
Added:
) unless $self->probe_by_key($probe);
20
21
21
Removed:
my $days = $self->param('days') // 7;
22
Removed:
$days = 7 unless $days =~ /\A\d+\z/ && $days >= 1 && $days <= 30;
22
Added:
my $days = bounded_integer( $self->param('days'), 7, 1, 30 );
23
Added:
my $now = time;
24
Added:
my $rollups = $self->reading->rollup_averages( [$probe], $HOUR_SECONDS,
25
Added:
$now - ( $days * $DAY_SECONDS ), $now, );
26
Added:
my ( $traces, $mean ) = diurnal_summary($rollups);
23
27
24
Removed:
my $now = time;
25
Removed:
my $start = $now - ( $days * $DAY_SECONDS );
26
Removed:
my $rollups = $self->sqlite->db->query(
27
Removed:
q{
28
Removed:
SELECT
29
Removed:
bucket_epoch,
30
Removed:
value_total / sample_count AS avg_value
31
Removed:
FROM reading_rollups
32
Removed:
WHERE probe = ?
33
Removed:
AND bucket_seconds = ?
34
Removed:
AND bucket_epoch >= ?
35
Removed:
AND bucket_epoch < ?
36
Removed:
ORDER BY bucket_epoch
37
Removed:
},
38
Removed:
$probe,
39
Removed:
$HOUR_SECONDS,
40
Removed:
$start,
41
Removed:
$now,
42
Removed:
)->hashes->to_array;
43
Removed:
44
Removed:
# Group by day-of-week, with hour-of-day as x-axis
45
Removed:
my %by_day;
46
Removed:
for my $row (@$rollups) {
47
Removed:
my $epoch = $row->{bucket_epoch};
48
Removed:
my @t = gmtime($epoch);
49
Removed:
my $hour = $t[2];
50
Removed:
my $date = strftime( '%Y-%m-%d', @t );
51
Removed:
52
Removed:
$by_day{$date}[$hour] = 0 + $row->{avg_value};
53
Removed:
}
54
Removed:
55
Removed:
# Build traces: one array per day, plus a mean trace
56
Removed:
my @traces;
57
Removed:
my @sums = (0) x 24;
58
Removed:
my @counts = (0) x 24;
59
Removed:
60
Removed:
for my $date ( sort keys %by_day ) {
61
Removed:
my $values = $by_day{$date};
62
Removed:
my @trace;
63
Removed:
for my $h ( 0 .. 23 ) {
64
Removed:
my $v = $values->[$h];
65
Removed:
push @trace, $v;
66
Removed:
if ( defined $v ) {
67
Removed:
$sums[$h] += $v;
68
Removed:
$counts[$h] += 1;
69
Removed:
}
70
Removed:
}
71
Removed:
push @traces, { date => $date, values => \@trace };
72
Removed:
}
73
Removed:
74
Removed:
my @mean = map { $counts[$_] > 0 ? $sums[$_] / $counts[$_] : undef } 0 .. 23;
75
Removed:
76
28
$self->render(
77
29
json => {
78
30
probe => $probe,
79
Removed:
days => 0 + $days,
31
Added:
days => $days,
80
32
hours => [ 0 .. 23 ],
81
Removed:
traces => \@traces,
82
Removed:
mean => \@mean,
33
Added:
traces => $traces,
34
Added:
mean => $mean,
83
35
},
84
36
);
85
37
}
@@ -87,59 +39,17 @@
87
39
# GET /api/insights/derivatives
88
40
# Returns hourly rate-of-change for all probes over the last N hours.
89
41
sub derivatives ($self) {
90
Removed:
my $hours = $self->param('hours') // 48;
91
Removed:
$hours = 48 unless $hours =~ /\A\d+\z/ && $hours >= 1 && $hours <= 720;
42
Added:
my $hours = bounded_integer( $self->param('hours'), 48, 1, 720 );
43
Added:
my $now = time;
44
Added:
my @probes = map { $_->{key} } $self->probes->@*;
45
Added:
my $rollups = $self->reading->rollup_averages( \@probes, $HOUR_SECONDS,
46
Added:
$now - ( $hours * $HOUR_SECONDS ), $now, );
47
Added:
my $probe_data = derivative_points( $rollups, \@probes );
92
48
93
Removed:
my $now = time;
94
Removed:
my $start = $now - ( $hours * $HOUR_SECONDS );
95
Removed:
96
Removed:
my %probe_data;
97
Removed:
for my $probe_info ( $self->probes->@* ) {
98
Removed:
my $probe = $probe_info->{key};
99
Removed:
100
Removed:
my $rollups = $self->sqlite->db->query(
101
Removed:
q{
102
Removed:
SELECT
103
Removed:
bucket_epoch,
104
Removed:
value_total / sample_count AS avg_value
105
Removed:
FROM reading_rollups
106
Removed:
WHERE probe = ?
107
Removed:
AND bucket_seconds = ?
108
Removed:
AND bucket_epoch >= ?
109
Removed:
AND bucket_epoch < ?
110
Removed:
ORDER BY bucket_epoch
111
Removed:
},
112
Removed:
$probe,
113
Removed:
$HOUR_SECONDS,
114
Removed:
$start,
115
Removed:
$now,
116
Removed:
)->hashes->to_array;
117
Removed:
118
Removed:
my @points;
119
Removed:
for my $i ( 1 .. $rollups->$#* ) {
120
Removed:
my $prev = $rollups->[ $i - 1 ];
121
Removed:
my $curr = $rollups->[$i];
122
Removed:
my $dt = $curr->{bucket_epoch} - $prev->{bucket_epoch};
123
Removed:
124
Removed:
next unless $dt > 0;
125
Removed:
126
Removed:
# Rate per hour
127
Removed:
my $rate = ( $curr->{avg_value} - $prev->{avg_value} )
128
Removed:
/ ( $dt / $HOUR_SECONDS );
129
Removed:
130
Removed:
push @points, {
131
Removed:
timestamp => utc_timestamp( $curr->{bucket_epoch} ),
132
Removed:
value => 0 + sprintf( '%.4f', $rate ),
133
Removed:
};
134
Removed:
}
135
Removed:
136
Removed:
$probe_data{$probe} = \@points;
137
Removed:
}
138
Removed:
139
49
$self->render(
140
50
json => {
141
Removed:
hours => 0 + $hours,
142
Removed:
probes => \%probe_data,
51
Added:
hours => $hours,
52
Added:
probes => $probe_data,
143
53
},
144
54
);
145
55
}
@@ -148,94 +58,26 @@
148
58
# Returns a rolling stability score (lower = more stable).
149
59
# Computed as the rolling 6h normalised standard deviation across all probes.
150
60
sub stability ($self) {
151
Removed:
my $hours = $self->param('hours') // 168; # 7 days default
152
Removed:
$hours = 168 unless $hours =~ /\A\d+\z/ && $hours >= 1 && $hours <= 720;
61
Added:
my $hours = bounded_integer( $self->param('hours'), 168, 1, 720 );
62
Added:
my $now = time;
63
Added:
my @probe_info = $self->probes->@*;
64
Added:
my @probes = map { $_->{key} } @probe_info;
65
Added:
my %ranges = map {
66
Added:
$_->{key} => { min => $_->{good_min}, max => $_->{good_max} }
67
Added:
} @probe_info;
68
Added:
my $rollups = $self->reading->rollup_averages( \@probes, $HOUR_SECONDS,
69
Added:
$now - ( $hours * $HOUR_SECONDS ), $now, );
153
70
154
Removed:
my $now = time;
155
Removed:
my $start = $now - ( $hours * $HOUR_SECONDS );
156
Removed:
157
Removed:
# Optimal ranges for normalisation
158
Removed:
my %ranges = (
159
Removed:
ph => { min => 6.4, max => 7.2 },
160
Removed:
do => { min => 5.5, max => 10 },
161
Removed:
orp => { min => 250, max => 400 },
162
Removed:
ec => { min => 300, max => 1500 },
163
Removed:
);
164
Removed:
165
Removed:
# Fetch hourly averages for all probes
166
Removed:
my %by_epoch;
167
Removed:
for my $probe_info ( $self->probes->@* ) {
168
Removed:
my $probe = $probe_info->{key};
169
Removed:
my $range = $ranges{$probe} or next;
170
Removed:
my $span = $range->{max} - $range->{min};
171
Removed:
172
Removed:
my $rollups = $self->sqlite->db->query(
173
Removed:
q{
174
Removed:
SELECT
175
Removed:
bucket_epoch,
176
Removed:
value_total / sample_count AS avg_value
177
Removed:
FROM reading_rollups
178
Removed:
WHERE probe = ?
179
Removed:
AND bucket_seconds = ?
180
Removed:
AND bucket_epoch >= ?
181
Removed:
AND bucket_epoch < ?
182
Removed:
ORDER BY bucket_epoch
183
Removed:
},
184
Removed:
$probe,
185
Removed:
$HOUR_SECONDS,
186
Removed:
$start,
187
Removed:
$now,
188
Removed:
)->hashes->to_array;
189
Removed:
190
Removed:
for my $row (@$rollups) {
191
Removed:
my $normalised = ( $row->{avg_value} - $range->{min} ) / $span;
192
Removed:
$by_epoch{ $row->{bucket_epoch} }{$probe} = $normalised;
193
Removed:
}
194
Removed:
}
195
Removed:
196
Removed:
# Compute rolling stddev over a 6-hour window
197
71
my $window_size = 6;
198
Removed:
my @sorted_epochs = sort { $a <=> $b } keys %by_epoch;
199
Removed:
my @points;
72
Added:
my $points = stability_points( $rollups, \%ranges, $window_size );
200
73
201
Removed:
for my $i ( $window_size - 1 .. $#sorted_epochs ) {
202
Removed:
my @window_epochs
203
Removed:
= @sorted_epochs[ $i - $window_size + 1 .. $i ];
204
Removed:
my @all_values;
205
Removed:
206
Removed:
for my $epoch (@window_epochs) {
207
Removed:
push @all_values, values $by_epoch{$epoch}->%*;
208
Removed:
}
209
Removed:
210
Removed:
next unless @all_values >= 2;
211
Removed:
212
Removed:
my $mean = 0;
213
Removed:
$mean += $_ for @all_values;
214
Removed:
$mean /= @all_values;
215
Removed:
216
Removed:
my $variance = 0;
217
Removed:
$variance += ( $_ - $mean ) ** 2 for @all_values;
218
Removed:
$variance /= @all_values;
219
Removed:
220
Removed:
my $stddev = sqrt($variance);
221
Removed:
222
Removed:
push @points, {
223
Removed:
timestamp => utc_timestamp( $sorted_epochs[$i] ),
224
Removed:
score => 0 + sprintf( '%.4f', $stddev ),
225
Removed:
};
226
Removed:
}
227
Removed:
228
74
$self->render(
229
75
json => {
230
Removed:
hours => 0 + $hours,
76
Added:
hours => $hours,
231
77
window => $window_size,
232
Removed:
points => \@points,
78
Added:
points => $points,
233
79
},
234
80
);
235
Removed:
}
236
Removed:
237
Removed:
sub utc_timestamp ($epoch) {
238
Removed:
return strftime( '%Y-%m-%dT%H:%M:%SZ', gmtime $epoch );
239
81
}
240
82
241
83
1;
roles/dashboard/lib/FAPG/DAQ/Dashboard/Controller/Pages.pm
@@ -26,10 +26,10 @@
26
26
27
27
sub probe_list ($self) {
28
28
$self->render(
29
Removed:
template => 'dashboard/probes',
30
Removed:
probes => $self->probes,
31
Removed:
nav_page => 'probes',
32
Removed:
breadcrumbs =>
29
Added:
template => 'dashboard/probes',
30
Added:
probes => $self->probes,
31
Added:
nav_page => 'probes',
32
Added:
breadcrumbs =>
33
33
[ { label => 'Home', href => '/' }, { label => 'Probes' }, ],
34
34
graph_control_aria_label => 'All probe chart timeframes',
35
35
);
@@ -39,21 +39,9 @@
39
39
my $from = strftime( '%Y-%m-%d', gmtime( time - ( 6 * 24 * 60 * 60 ) ) );
40
40
my $month_from
41
41
= strftime( '%Y-%m-%d', gmtime( time - ( 29 * 24 * 60 * 60 ) ) );
42
Removed:
my $to = strftime( '%Y-%m-%d', gmtime(time) );
43
Removed:
my $database = $self->sqlite->db;
44
Removed:
my $table
45
Removed:
= $database->query(
46
Removed:
q{SELECT name FROM sqlite_master WHERE type = 'table' AND name = 'readings'}
47
Removed:
)->hash;
48
Removed:
my $oldest;
42
Added:
my $to = strftime( '%Y-%m-%d', gmtime(time) );
43
Added:
my $oldest = $self->reading->oldest_received_timestamp;
49
44
50
Removed:
if ($table) {
51
Removed:
$oldest
52
Removed:
= $database->query(
53
Removed:
'SELECT MIN(COALESCE(received_at, timestamp)) AS timestamp FROM readings'
54
Removed:
)->hash->{timestamp};
55
Removed:
}
56
Removed:
57
45
$oldest = $from if !defined $oldest;
58
46
$oldest =~ s/T.*\z//;
59
47
@@ -91,16 +79,16 @@
91
79
92
80
sub graph ($self) {
93
81
my $probe_key = $self->stash('probe');
94
Removed:
my ($probe) = grep { $_->{key} eq $probe_key } $self->probes->@*;
82
Added:
my $probe = $self->probe_by_key($probe_key);
95
83
96
84
return $self->reply->not_found unless defined $probe;
97
85
98
86
$self->render(
99
Removed:
template => 'dashboard/graph',
100
Removed:
probe => $probe,
101
Removed:
nav_page => 'probes',
102
Removed:
single_probe => 1,
103
Removed:
breadcrumbs => [
87
Added:
template => 'dashboard/graph',
88
Added:
probe => $probe,
89
Added:
nav_page => 'probes',
90
Added:
single_probe => 1,
91
Added:
breadcrumbs => [
104
92
{ label => 'Home', href => '/' },
105
93
{ label => 'Probes', href => '/probes' },
106
94
{ label => $probe->{label} },
roles/dashboard/lib/FAPG/DAQ/Dashboard/Controller/Reading.pm
@@ -1,57 +1,34 @@
1
1
package FAPG::DAQ::Dashboard::Controller::Reading;
2
2
use Mojo::Base 'Mojolicious::Controller', -signatures;
3
3
4
Added:
use FAPG::DAQ::Dashboard::Analytics qw(aggregate_series);
4
5
use FAPG::DAQ::Dashboard::Timeframe qw(max_series_range series_window);
5
Removed:
use POSIX qw(strftime);
6
Added:
use FAPG::DAQ::Dashboard::Util qw(utc_timestamp);
6
7
7
8
my $STATUS_REACHABLE_SECONDS = 120;
8
9
my $MAX_SERIES_RANGE = max_series_range();
9
10
my $MINUTE_SECONDS = 60;
10
Removed:
my $MAX_SERIES_CACHE_ENTRIES = 256;
11
Removed:
my %SERIES_CACHE;
12
11
13
12
sub list ($self) {
14
13
my $probe = $self->param('probe') // '';
15
14
16
Removed:
my %known = map { $_->{key} => 1 } $self->probes->@*;
17
Removed:
18
15
return $self->render(
19
16
status => 404,
20
17
json => { error => "Unknown probe type: $probe", },
21
Removed:
) unless $known{$probe};
18
Added:
) unless $self->probe_by_key($probe);
22
19
23
20
my $limit = $self->param('limit') // 300;
24
21
25
22
$limit = 300 unless $limit =~ /^\d+$/;
26
23
$limit = 2_000 if $limit > 2_000;
27
24
28
Removed:
my @where = ('probe = ?');
29
Removed:
my @bind = ($probe);
30
25
my $since = $self->param('since');
26
Added:
$since = undef
27
Added:
if !defined $since
28
Added:
|| $since !~ /\A\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}Z\z/;
31
29
32
Removed:
if ( defined $since
33
Removed:
&& $since =~ /\A\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}Z\z/ )
34
Removed:
{
35
Removed:
push @where, 'COALESCE(timestamp, received_at) >= ?';
36
Removed:
push @bind, $since;
37
Removed:
}
30
Added:
my $rows = $self->reading->list( $probe, $limit, $since );
38
31
39
Removed:
my $rows = $self->sqlite->db->query(
40
Removed:
q{
41
Removed:
SELECT timestamp, received_at, node, probe, value, unit
42
Removed:
FROM readings
43
Removed:
WHERE
44
Removed:
}
45
Removed:
. join( "\n AND ", @where ) . q{
46
Removed:
ORDER BY COALESCE(timestamp, received_at) DESC, id DESC
47
Removed:
LIMIT ?
48
Removed:
},
49
Removed:
@bind,
50
Removed:
$limit,
51
Removed:
)->hashes->to_array;
52
Removed:
53
Removed:
$rows = [ reverse $rows->@* ];
54
Removed:
55
32
$self->render(
56
33
json => {
57
34
probe => $probe,
@@ -63,12 +40,10 @@
63
40
sub series ($self) {
64
41
my $probe = $self->param('probe') // '';
65
42
66
Removed:
my %known = map { $_->{key} => 1 } $self->probes->@*;
67
Removed:
68
43
return $self->render(
69
44
status => 404,
70
45
json => { error => "Unknown probe type: $probe", },
71
Removed:
) unless $known{$probe};
46
Added:
) unless $self->probe_by_key($probe);
72
47
73
48
my $timeframe = $self->param('timeframe') // 'day';
74
49
my $range = $self->param('range') // 1;
@@ -92,36 +67,19 @@
92
67
= int( $now / $rollup_seconds ) * $rollup_seconds + $rollup_seconds;
93
68
my $cache_key = join q{:}, $probe, $timeframe, $range,
94
69
$window->{start_epoch}, $window->{end_epoch};
95
Removed:
my $cached = $SERIES_CACHE{$cache_key};
70
Added:
my $cached = $self->series_cache->get( $cache_key, $now );
96
71
97
Removed:
if ( $cached && $cached->{expires_epoch} > $now ) {
72
Added:
if ($cached) {
98
73
return render_series_payload( $self, $cached->{payload}, 'HIT',
99
74
$cached->{expires_epoch} );
100
75
}
101
76
102
Removed:
my $rollups = $self->sqlite->db->query(
103
Removed:
q{
104
Removed:
SELECT
105
Removed:
bucket_epoch AS minute_epoch,
106
Removed:
sample_count AS count,
107
Removed:
value_total / sample_count AS avg_value,
108
Removed:
value_min,
109
Removed:
value_max
110
Removed:
FROM reading_rollups
111
Removed:
WHERE probe = ?
112
Removed:
AND bucket_seconds = ?
113
Removed:
AND bucket_epoch >= ?
114
Removed:
AND bucket_epoch < ?
115
Removed:
ORDER BY bucket_epoch
116
Removed:
},
117
Removed:
$probe,
118
Removed:
$rollup_seconds,
77
Added:
my $rollups
78
Added:
= $self->reading->series_rollups( $probe, $rollup_seconds,
119
79
$window->{start_epoch} - ( 2 * $rollup_seconds ),
120
Removed:
$window->{end_epoch},
121
Removed:
)->hashes->to_array;
80
Added:
$window->{end_epoch}, );
81
Added:
my $points = aggregate_series( $rollups, $window, 2 * $rollup_seconds );
122
82
123
Removed:
my $points = aggregate_series( $rollups, $window );
124
Removed:
125
83
my $payload = {
126
84
probe => $probe,
127
85
timeframe => $timeframe,
@@ -129,32 +87,11 @@
129
87
points => $points,
130
88
};
131
89
132
Removed:
prune_series_cache($now);
133
Removed:
$SERIES_CACHE{$cache_key} = {
134
Removed:
expires_epoch => $expires_epoch,
135
Removed:
payload => $payload,
136
Removed:
};
90
Added:
$self->series_cache->set( $cache_key, $payload, $expires_epoch, $now );
137
91
138
92
return render_series_payload( $self, $payload, 'MISS', $expires_epoch );
139
93
}
140
94
141
Removed:
sub prune_series_cache ($now) {
142
Removed:
for my $key ( keys %SERIES_CACHE ) {
143
Removed:
delete $SERIES_CACHE{$key}
144
Removed:
if $SERIES_CACHE{$key}{expires_epoch} <= $now;
145
Removed:
}
146
Removed:
147
Removed:
if ( keys(%SERIES_CACHE) >= $MAX_SERIES_CACHE_ENTRIES ) {
148
Removed:
my ($oldest) = sort {
149
Removed:
$SERIES_CACHE{$a}{expires_epoch}
150
Removed:
<=> $SERIES_CACHE{$b}{expires_epoch}
151
Removed:
} keys %SERIES_CACHE;
152
Removed:
delete $SERIES_CACHE{$oldest};
153
Removed:
}
154
Removed:
155
Removed:
return;
156
Removed:
}
157
Removed:
158
95
sub render_series_payload ( $self, $payload, $cache_status, $expires_epoch ) {
159
96
my $max_age = int( $expires_epoch - time );
160
97
$max_age = 0 if $max_age < 0;
@@ -167,28 +104,13 @@
167
104
sub status ($self) {
168
105
my $probe = $self->param('probe') // '';
169
106
170
Removed:
my %known = map { $_->{key} => 1 } $self->status_items->@*;
171
Removed:
172
107
return $self->render(
173
108
status => 404,
174
109
json => { error => "Unknown status item: $probe", },
175
Removed:
) unless $known{$probe};
110
Added:
) unless $self->status_item_by_key($probe);
176
111
177
Removed:
my $row = eval {
178
Removed:
$self->sqlite->db->query(
179
Removed:
q{
180
Removed:
SELECT timestamp, received_at, node, probe, status, message, error, device
181
Removed:
FROM node_status
182
Removed:
WHERE probe = ?
183
Removed:
AND valid = 1
184
Removed:
ORDER BY COALESCE(timestamp, received_at) DESC, id DESC
185
Removed:
LIMIT 1
186
Removed:
},
187
Removed:
$probe,
188
Removed:
)->hash;
189
Removed:
};
112
Added:
my $row = $self->reading->latest_status($probe);
190
113
191
Removed:
$row = undef if $@;
192
114
mark_unreachable_if_stale($row) if $row;
193
115
194
116
$self->render(
@@ -200,40 +122,10 @@
200
122
}
201
123
202
124
sub status_list ($self) {
203
Removed:
my @probes = map { $_->{key} } $self->status_items->@*;
204
Removed:
my $placeholders = join ', ', ('?') x @probes;
205
Removed:
my %statuses = map { $_ => undef } @probes;
125
Added:
my @probes = map { $_->{key} } $self->status_items->@*;
126
Added:
my %statuses = map { $_ => undef } @probes;
127
Added:
my $rows = $self->reading->latest_statuses( \@probes );
206
128
207
Removed:
my $rows = eval {
208
Removed:
$self->sqlite->db->query(
209
Removed:
qq{
210
Removed:
WITH ranked AS (
211
Removed:
SELECT
212
Removed:
timestamp,
213
Removed:
received_at,
214
Removed:
node,
215
Removed:
probe,
216
Removed:
status,
217
Removed:
message,
218
Removed:
error,
219
Removed:
device,
220
Removed:
ROW_NUMBER() OVER (
221
Removed:
PARTITION BY probe
222
Removed:
ORDER BY COALESCE(timestamp, received_at) DESC, id DESC
223
Removed:
) AS status_rank
224
Removed:
FROM node_status
225
Removed:
WHERE valid = 1
226
Removed:
AND probe IN ($placeholders)
227
Removed:
)
228
Removed:
SELECT timestamp, received_at, node, probe, status, message, error, device
229
Removed:
FROM ranked
230
Removed:
WHERE status_rank = 1
231
Removed:
},
232
Removed:
@probes,
233
Removed:
)->hashes->to_array;
234
Removed:
};
235
Removed:
236
Removed:
$rows = [] if $@;
237
129
for my $row (@$rows) {
238
130
mark_unreachable_if_stale($row);
239
131
$statuses{ $row->{probe} } = $row;
@@ -252,82 +144,22 @@
252
144
}
253
145
254
146
sub reachable_since {
255
Removed:
return strftime( '%Y-%m-%dT%H:%M:%SZ',
256
Removed:
gmtime( time - $STATUS_REACHABLE_SECONDS ) );
147
Added:
return utc_timestamp( time - $STATUS_REACHABLE_SECONDS );
257
148
}
258
149
259
Removed:
sub series_points ($window) {
260
Removed:
my @points;
261
Removed:
262
Removed:
for my $i ( 0 .. $window->{bucket_count} - 1 ) {
263
Removed:
push @points,
264
Removed:
{
265
Removed:
timestamp => utc_timestamp(
266
Removed:
$window->{start_epoch} + ( $i * $window->{bucket_stride} )
267
Removed:
),
268
Removed:
value => undef,
269
Removed:
min => undef,
270
Removed:
max => undef,
271
Removed:
count => 0,
272
Removed:
gap => 0,
273
Removed:
};
274
Removed:
}
275
Removed:
276
Removed:
return \@points;
277
Removed:
}
278
Removed:
279
Removed:
sub aggregate_series ( $rollups, $window ) {
280
Removed:
my $points = series_points($window);
281
Removed:
my @values;
282
Removed:
283
Removed:
for my $rollup ( @$rollups ) {
284
Removed:
my $epoch = 0 + $rollup->{minute_epoch};
285
Removed:
next
286
Removed:
if $epoch < $window->{start_epoch}
287
Removed:
|| $epoch >= $window->{end_epoch};
288
Removed:
289
Removed:
my $point_index = int(
290
Removed:
( $epoch - $window->{start_epoch} ) / $window->{bucket_stride} );
291
Removed:
my $point = $points->[$point_index];
292
Removed:
293
Removed:
push $values[$point_index]->@*, 0 + $rollup->{avg_value};
294
Removed:
$point->{min} = 0 + $rollup->{value_min}
295
Removed:
if !defined $point->{min}
296
Removed:
|| $rollup->{value_min} < $point->{min};
297
Removed:
$point->{max} = 0 + $rollup->{value_max}
298
Removed:
if !defined $point->{max}
299
Removed:
|| $rollup->{value_max} > $point->{max};
300
Removed:
$point->{count} += 0 + $rollup->{count};
301
Removed:
}
302
Removed:
303
Removed:
for my $i ( 0 .. $points->$#* ) {
304
Removed:
my $point_values = $values[$i] // [];
305
Removed:
next if !@$point_values;
306
Removed:
307
Removed:
my $total = 0;
308
Removed:
$total += $_ for @$point_values;
309
Removed:
310
Removed:
$points->[$i]{value} = $total / @$point_values;
311
Removed:
}
312
Removed:
313
Removed:
return $points;
314
Removed:
}
315
Removed:
316
150
sub correlation ($self) {
317
151
my $probe_x = $self->param('probe_x') // '';
318
152
my $probe_y = $self->param('probe_y') // '';
319
153
320
Removed:
my %known = map { $_->{key} => 1 } $self->probes->@*;
321
Removed:
322
154
return $self->render(
323
155
status => 404,
324
156
json => { error => "Unknown probe type: $probe_x", },
325
Removed:
) unless $known{$probe_x};
157
Added:
) unless $self->probe_by_key($probe_x);
326
158
327
159
return $self->render(
328
160
status => 404,
329
161
json => { error => "Unknown probe type: $probe_y", },
330
Removed:
) unless $known{$probe_y};
162
Added:
) unless $self->probe_by_key($probe_y);
331
163
332
164
my $timeframe = $self->param('timeframe') // 'day';
333
165
my $range = $self->param('range') // 1;
@@ -347,33 +179,12 @@
347
179
my $rollup_seconds
348
180
= $window->{bucket_stride} >= 3_600 ? 3_600 : $MINUTE_SECONDS;
349
181
350
Removed:
my $rows = $self->sqlite->db->query(
351
Removed:
q{
352
Removed:
SELECT
353
Removed:
a.bucket_epoch,
354
Removed:
a.value_total / a.sample_count AS x,
355
Removed:
b.value_total / b.sample_count AS y
356
Removed:
FROM reading_rollups a
357
Removed:
JOIN reading_rollups b
358
Removed:
ON b.probe = ?
359
Removed:
AND b.bucket_seconds = a.bucket_seconds
360
Removed:
AND b.bucket_epoch = a.bucket_epoch
361
Removed:
WHERE a.probe = ?
362
Removed:
AND a.bucket_seconds = ?
363
Removed:
AND a.bucket_epoch >= ?
364
Removed:
AND a.bucket_epoch < ?
365
Removed:
ORDER BY a.bucket_epoch
366
Removed:
},
367
Removed:
$probe_y,
368
Removed:
$probe_x,
369
Removed:
$rollup_seconds,
370
Removed:
$window->{start_epoch},
371
Removed:
$window->{end_epoch},
372
Removed:
)->hashes->to_array;
182
Added:
my $rows
183
Added:
= $self->reading->correlation_rollups( $probe_x, $probe_y,
184
Added:
$rollup_seconds, $window->{start_epoch},
185
Added:
$window->{end_epoch}, );
373
186
374
Removed:
my @points = map {
375
Removed:
{ x => 0 + $_->{x}, y => 0 + $_->{y} }
376
Removed:
} @$rows;
187
Added:
my @points = map { { x => 0 + $_->{x}, y => 0 + $_->{y} } } @$rows;
377
188
378
189
$self->render(
379
190
json => {
@@ -384,10 +195,6 @@
384
195
points => \@points,
385
196
},
386
197
);
387
Removed:
}
388
Removed:
389
Removed:
sub utc_timestamp ($epoch) {
390
Removed:
return strftime( '%Y-%m-%dT%H:%M:%SZ', gmtime $epoch );
391
198
}
392
199
393
200
1;
roles/dashboard/lib/FAPG/DAQ/Dashboard/Controller/Weather.pm
@@ -1,12 +1,11 @@
1
1
package FAPG::DAQ::Dashboard::Controller::Weather;
2
2
use Mojo::Base 'Mojolicious::Controller', -signatures;
3
3
4
Removed:
use POSIX qw(strftime);
4
Added:
use FAPG::DAQ::Dashboard::Util qw(bounded_integer utc_timestamp);
5
5
6
6
# GET /api/weather/hourly?hours=168
7
7
sub hourly ($self) {
8
Removed:
my $hours = $self->param('hours') // 168;
9
Removed:
$hours = 168 unless $hours =~ /\A\d+\z/ && $hours >= 1 && $hours <= 720;
8
Added:
my $hours = bounded_integer( $self->param('hours'), 168, 1, 720 );
10
9
11
10
my $now = time;
12
11
my $from = $now - ( $hours * 3600 );
@@ -23,14 +22,14 @@
23
22
24
23
# GET /api/weather/forecast
25
24
sub forecast ($self) {
26
Removed:
my $now = time;
27
Removed:
my $from = $now - ( 24 * 3600 ); # Include past 24h for context
28
Removed:
my $points = $self->weather->get_forecast($from);
25
Added:
my $now = time;
26
Added:
my $from = $now - ( 24 * 3600 ); # Include past 24h for context
27
Added:
my $points = $self->weather->get_forecast($from);
29
28
my @formatted = map { _format_point($_) } @$points;
30
29
31
30
$self->render(
32
31
json => {
33
Removed:
from => _utc_timestamp($from),
32
Added:
from => utc_timestamp($from),
34
33
points => \@formatted,
35
34
},
36
35
);
@@ -38,11 +37,16 @@
38
37
39
38
# GET /api/weather/current
40
39
sub current ($self) {
41
Removed:
my $latest = $self->weather->latest;
40
Added:
my $latest = $self->weather->latest_at(time);
42
41
43
42
unless ($latest) {
44
43
return $self->render(
45
Removed:
json => { temperature => undef, humidity => undef, rain => undef, timestamp => undef },
44
Added:
json => {
45
Added:
temperature => undef,
46
Added:
humidity => undef,
47
Added:
rain => undef,
48
Added:
timestamp => undef
49
Added:
},
46
50
);
47
51
}
48
52
@@ -51,7 +55,7 @@
51
55
temperature => $latest->{temperature},
52
56
humidity => $latest->{humidity},
53
57
rain => $latest->{rain},
54
Removed:
timestamp => _utc_timestamp( $latest->{epoch} ),
58
Added:
timestamp => utc_timestamp( $latest->{epoch} ),
55
59
},
56
60
);
57
61
}
@@ -62,17 +66,36 @@
62
66
63
67
unless ( $config->{enabled} ) {
64
68
return $self->render(
65
Removed:
json => { status => 'disabled', message => 'Weather fetch is disabled in config' },
69
Added:
json => {
70
Added:
status => 'disabled',
71
Added:
message => 'Weather fetch is disabled in config'
72
Added:
},
66
73
);
67
74
}
68
75
69
Removed:
my $fetcher = $self->app->weather_fetcher;
70
Removed:
my $rows = $fetcher->fetch;
71
Removed:
my $stored = $self->weather->store_batch($rows);
76
Added:
my ( $rows, $stored );
77
Added:
my $success = eval {
78
Added:
$rows = $self->app->weather_fetcher->fetch;
79
Added:
$stored = $self->weather->store_batch($rows);
80
Added:
1;
81
Added:
};
72
82
83
Added:
if ( !$success ) {
84
Added:
my $error = $@ || 'Unknown weather refresh failure';
85
Added:
chomp $error;
86
Added:
$self->app->log->error($error);
87
Added:
return $self->render(
88
Added:
status => 502,
89
Added:
json => {
90
Added:
status => 'error',
91
Added:
message => 'Weather data could not be refreshed',
92
Added:
},
93
Added:
);
94
Added:
}
95
Added:
73
96
$self->render(
74
97
json => {
75
Removed:
status => 'ok',
98
Added:
status => 'ok',
76
99
fetched => scalar @$rows,
77
100
stored => $stored // 0,
78
101
},
@@ -81,15 +104,11 @@
81
104
82
105
sub _format_point ($row) {
83
106
return {
84
Removed:
timestamp => _utc_timestamp( $row->{epoch} ),
107
Added:
timestamp => utc_timestamp( $row->{epoch} ),
85
108
temperature => $row->{temperature},
86
109
humidity => $row->{humidity},
87
110
rain => $row->{rain},
88
111
};
89
Removed:
}
90
Removed:
91
Removed:
sub _utc_timestamp ($epoch) {
92
Removed:
return strftime( '%Y-%m-%dT%H:%M:%SZ', gmtime $epoch );
93
112
}
94
113
95
114
1;
roles/dashboard/lib/FAPG/DAQ/Dashboard/Model/Reading.pm
@@ -0,0 +1,209 @@
1
Added:
package FAPG::DAQ::Dashboard::Model::Reading;
2
Added:
use Mojo::Base -base, -signatures;
3
Added:
4
Added:
has 'sqlite';
5
Added:
6
Added:
sub list ( $self, $probe, $limit, $since = undef ) {
7
Added:
my @where = ('probe = ?');
8
Added:
my @bind = ($probe);
9
Added:
10
Added:
if ( defined $since ) {
11
Added:
push @where, 'COALESCE(timestamp, received_at) >= ?';
12
Added:
push @bind, $since;
13
Added:
}
14
Added:
15
Added:
my $rows = $self->sqlite->db->query(
16
Added:
q{
17
Added:
SELECT timestamp, received_at, node, probe, value, unit
18
Added:
FROM readings
19
Added:
WHERE
20
Added:
}
21
Added:
. join( "\n AND ", @where ) . q{
22
Added:
ORDER BY COALESCE(timestamp, received_at) DESC, id DESC
23
Added:
LIMIT ?
24
Added:
},
25
Added:
@bind,
26
Added:
$limit,
27
Added:
)->hashes->to_array;
28
Added:
29
Added:
return [ reverse @$rows ];
30
Added:
}
31
Added:
32
Added:
sub series_rollups ( $self, $probe, $bucket_seconds, $start_epoch,
33
Added:
$end_epoch )
34
Added:
{
35
Added:
return $self->sqlite->db->query(
36
Added:
q{
37
Added:
SELECT
38
Added:
bucket_epoch AS minute_epoch,
39
Added:
sample_count AS count,
40
Added:
value_total / sample_count AS avg_value,
41
Added:
value_min,
42
Added:
value_max,
43
Added:
first_sample_epoch,
44
Added:
last_sample_epoch,
45
Added:
max_gap_seconds
46
Added:
FROM reading_rollups
47
Added:
WHERE probe = ?
48
Added:
AND bucket_seconds = ?
49
Added:
AND bucket_epoch >= ?
50
Added:
AND bucket_epoch < ?
51
Added:
ORDER BY bucket_epoch
52
Added:
},
53
Added:
$probe,
54
Added:
$bucket_seconds,
55
Added:
$start_epoch,
56
Added:
$end_epoch,
57
Added:
)->hashes->to_array;
58
Added:
}
59
Added:
60
Added:
sub correlation_rollups ( $self, $probe_x, $probe_y, $bucket_seconds,
61
Added:
$start_epoch, $end_epoch )
62
Added:
{
63
Added:
return $self->sqlite->db->query(
64
Added:
q{
65
Added:
SELECT
66
Added:
a.bucket_epoch,
67
Added:
a.value_total / a.sample_count AS x,
68
Added:
b.value_total / b.sample_count AS y
69
Added:
FROM reading_rollups a
70
Added:
JOIN reading_rollups b
71
Added:
ON b.probe = ?
72
Added:
AND b.bucket_seconds = a.bucket_seconds
73
Added:
AND b.bucket_epoch = a.bucket_epoch
74
Added:
WHERE a.probe = ?
75
Added:
AND a.bucket_seconds = ?
76
Added:
AND a.bucket_epoch >= ?
77
Added:
AND a.bucket_epoch < ?
78
Added:
ORDER BY a.bucket_epoch
79
Added:
},
80
Added:
$probe_y,
81
Added:
$probe_x,
82
Added:
$bucket_seconds,
83
Added:
$start_epoch,
84
Added:
$end_epoch,
85
Added:
)->hashes->to_array;
86
Added:
}
87
Added:
88
Added:
sub latest_status ( $self, $probe ) {
89
Added:
return if !$self->_table_exists('node_status');
90
Added:
return $self->sqlite->db->query(
91
Added:
q{
92
Added:
SELECT timestamp, received_at, node, probe, status, message, error, device
93
Added:
FROM node_status
94
Added:
WHERE probe = ?
95
Added:
AND valid = 1
96
Added:
ORDER BY COALESCE(timestamp, received_at) DESC, id DESC
97
Added:
LIMIT 1
98
Added:
},
99
Added:
$probe,
100
Added:
)->hash;
101
Added:
}
102
Added:
103
Added:
sub latest_statuses ( $self, $probes ) {
104
Added:
return [] if !@$probes || !$self->_table_exists('node_status');
105
Added:
my $placeholders = join ', ', ('?') x @$probes;
106
Added:
107
Added:
return $self->sqlite->db->query(
108
Added:
qq{
109
Added:
WITH ranked AS (
110
Added:
SELECT
111
Added:
timestamp,
112
Added:
received_at,
113
Added:
node,
114
Added:
probe,
115
Added:
status,
116
Added:
message,
117
Added:
error,
118
Added:
device,
119
Added:
ROW_NUMBER() OVER (
120
Added:
PARTITION BY probe
121
Added:
ORDER BY COALESCE(timestamp, received_at) DESC, id DESC
122
Added:
) AS status_rank
123
Added:
FROM node_status
124
Added:
WHERE valid = 1
125
Added:
AND probe IN ($placeholders)
126
Added:
)
127
Added:
SELECT timestamp, received_at, node, probe, status, message, error, device
128
Added:
FROM ranked
129
Added:
WHERE status_rank = 1
130
Added:
},
131
Added:
@$probes,
132
Added:
)->hashes->to_array;
133
Added:
}
134
Added:
135
Added:
sub rollup_averages ( $self, $probes, $bucket_seconds, $start_epoch,
136
Added:
$end_epoch )
137
Added:
{
138
Added:
return [] if !@$probes;
139
Added:
my $placeholders = join ', ', ('?') x @$probes;
140
Added:
141
Added:
return $self->sqlite->db->query(
142
Added:
qq{
143
Added:
SELECT
144
Added:
probe,
145
Added:
bucket_epoch,
146
Added:
value_total / sample_count AS avg_value
147
Added:
FROM reading_rollups
148
Added:
WHERE probe IN ($placeholders)
149
Added:
AND bucket_seconds = ?
150
Added:
AND bucket_epoch >= ?
151
Added:
AND bucket_epoch < ?
152
Added:
ORDER BY probe, bucket_epoch
153
Added:
},
154
Added:
@$probes,
155
Added:
$bucket_seconds,
156
Added:
$start_epoch,
157
Added:
$end_epoch,
158
Added:
)->hashes->to_array;
159
Added:
}
160
Added:
161
Added:
sub readings_by_received_range ( $self, $from, $to, $probes ) {
162
Added:
my $probe_filter
163
Added:
= @$probes
164
Added:
? ' AND probe IN (' . join( ',', ('?') x @$probes ) . ')'
165
Added:
: '';
166
Added:
return $self->_raw_readings(
167
Added:
'COALESCE(received_at, timestamp) >= ? '
168
Added:
. 'AND COALESCE(received_at, timestamp) < ?'
169
Added:
. $probe_filter,
170
Added:
'COALESCE(received_at, timestamp)',
171
Added:
$from, $to, @$probes,
172
Added:
);
173
Added:
}
174
Added:
175
Added:
sub readings_by_measurement_range ( $self, $probe, $from, $to ) {
176
Added:
return $self->_raw_readings(
177
Added:
'probe = ? AND COALESCE(timestamp, received_at) >= ? '
178
Added:
. 'AND COALESCE(timestamp, received_at) < ?',
179
Added:
'COALESCE(timestamp, received_at)', $probe, $from, $to,
180
Added:
);
181
Added:
}
182
Added:
183
Added:
sub oldest_received_timestamp ($self) {
184
Added:
return if !$self->_table_exists('readings');
185
Added:
return $self->sqlite->db->query(
186
Added:
'SELECT MIN(COALESCE(received_at, timestamp)) AS timestamp FROM readings'
187
Added:
)->hash->{timestamp};
188
Added:
}
189
Added:
190
Added:
sub _raw_readings ( $self, $where, $order_by, @bind ) {
191
Added:
return [] if !$self->_table_exists('readings');
192
Added:
return $self->sqlite->db->query(
193
Added:
qq{
194
Added:
SELECT timestamp, received_at, node, probe, value, unit
195
Added:
FROM readings
196
Added:
WHERE $where
197
Added:
ORDER BY $order_by, id
198
Added:
},
199
Added:
@bind,
200
Added:
)->hashes->to_array;
201
Added:
}
202
Added:
203
Added:
sub _table_exists ( $self, $table ) {
204
Added:
return !!$self->sqlite->db->query(
205
Added:
q{SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = ?},
206
Added:
$table, )->hash;
207
Added:
}
208
Added:
209
Added:
1;
roles/dashboard/lib/FAPG/DAQ/Dashboard/Model/Weather.pm
@@ -1,7 +1,7 @@
1
1
package FAPG::DAQ::Dashboard::Model::Weather;
2
2
use Mojo::Base -base, -signatures;
3
3
4
Removed:
use POSIX qw(strftime);
4
Added:
use FAPG::DAQ::Dashboard::Util qw(utc_timestamp);
5
5
6
6
has 'sqlite';
7
7
@@ -24,7 +24,7 @@
24
24
return unless $rows && @$rows;
25
25
26
26
my $db = $self->sqlite->db;
27
Removed:
my $fetched_at = strftime( '%Y-%m-%dT%H:%M:%SZ', gmtime time );
27
Added:
my $fetched_at = utc_timestamp(time);
28
28
29
29
my $tx = $db->begin;
30
30
for my $row (@$rows) {
@@ -63,6 +63,17 @@
63
63
ORDER BY epoch},
64
64
$from_epoch,
65
65
)->hashes->to_array;
66
Added:
}
67
Added:
68
Added:
sub latest_at ( $self, $epoch ) {
69
Added:
return $self->sqlite->db->query(
70
Added:
q{SELECT epoch, temperature, humidity, rain
71
Added:
FROM weather_hourly
72
Added:
WHERE epoch <= ?
73
Added:
ORDER BY epoch DESC
74
Added:
LIMIT 1},
75
Added:
$epoch,
76
Added:
)->hash;
66
77
}
67
78
68
79
sub latest ($self) {
roles/dashboard/lib/FAPG/DAQ/Dashboard/Secret.pm
@@ -0,0 +1,88 @@
1
Added:
package FAPG::DAQ::Dashboard::Secret;
2
Added:
use Mojo::Base -strict;
3
Added:
4
Added:
use Errno qw(EEXIST);
5
Added:
use Exporter qw(import);
6
Added:
use Fcntl qw(O_CREAT O_EXCL O_WRONLY);
7
Added:
8
Added:
our @EXPORT_OK = qw(load_or_create_secret read_secret);
9
Added:
10
Added:
my $SECRET_BYTES = 32;
11
Added:
my $MIN_SECRET_LENGTH = 8;
12
Added:
13
Added:
sub read_secret {
14
Added:
my ($path) = @_;
15
Added:
return if !-e $path;
16
Added:
17
Added:
open my $fh, '<', $path
18
Added:
or die "Unable to open application secret $path: $!";
19
Added:
my $secret = do { local $/; <$fh> };
20
Added:
close $fh
21
Added:
or die "Unable to close application secret $path: $!";
22
Added:
23
Added:
chomp $secret if defined $secret;
24
Added:
return $secret;
25
Added:
}
26
Added:
27
Added:
sub load_or_create_secret {
28
Added:
my ($path) = @_;
29
Added:
my $existing = read_secret($path);
30
Added:
31
Added:
if ( defined $existing ) {
32
Added:
_validate_secret( $existing, $path );
33
Added:
return ( $existing, 0 );
34
Added:
}
35
Added:
36
Added:
my $secret = _random_secret();
37
Added:
if ( sysopen my $fh, $path, O_WRONLY | O_CREAT | O_EXCL, 0600 ) {
38
Added:
if ( !print {$fh} "$secret\n" ) {
39
Added:
my $error = $!;
40
Added:
close $fh;
41
Added:
unlink $path;
42
Added:
die "Unable to write application secret $path: $error";
43
Added:
}
44
Added:
if ( !close $fh ) {
45
Added:
my $error = $!;
46
Added:
unlink $path;
47
Added:
die "Unable to close application secret $path: $error";
48
Added:
}
49
Added:
50
Added:
return ( $secret, 1 );
51
Added:
}
52
Added:
53
Added:
my $error = 0 + $!;
54
Added:
if ( $error == EEXIST ) {
55
Added:
my $winner = read_secret($path);
56
Added:
_validate_secret( $winner, $path );
57
Added:
return ( $winner, 0 );
58
Added:
}
59
Added:
60
Added:
die "Unable to create application secret $path: $!";
61
Added:
}
62
Added:
63
Added:
sub _random_secret {
64
Added:
open my $fh, '<:raw', '/dev/urandom'
65
Added:
or die "Unable to open /dev/urandom: $!";
66
Added:
67
Added:
my $bytes = '';
68
Added:
while ( length $bytes < $SECRET_BYTES ) {
69
Added:
my $remaining = $SECRET_BYTES - length $bytes;
70
Added:
my $chunk;
71
Added:
my $read = read( $fh, $chunk, $remaining );
72
Added:
die "Unable to read /dev/urandom: $!" if !defined $read;
73
Added:
die 'Unexpected end of file while reading /dev/urandom' if !$read;
74
Added:
$bytes .= $chunk;
75
Added:
}
76
Added:
77
Added:
close $fh or die "Unable to close /dev/urandom: $!";
78
Added:
return unpack 'H*', $bytes;
79
Added:
}
80
Added:
81
Added:
sub _validate_secret {
82
Added:
my ( $secret, $path ) = @_;
83
Added:
die "Application secret in $path is missing or too short"
84
Added:
if !defined $secret || length($secret) <= $MIN_SECRET_LENGTH;
85
Added:
return;
86
Added:
}
87
Added:
88
Added:
1;
roles/dashboard/lib/FAPG/DAQ/Dashboard/Service/SeriesCache.pm
@@ -0,0 +1,38 @@
1
Added:
package FAPG::DAQ::Dashboard::Service::SeriesCache;
2
Added:
use Mojo::Base -base, -signatures;
3
Added:
4
Added:
has entries => sub { {} };
5
Added:
has max_entries => 256;
6
Added:
7
Added:
sub get ( $self, $key, $now ) {
8
Added:
my $entry = $self->entries->{$key};
9
Added:
return if !$entry || $entry->{expires_epoch} <= $now;
10
Added:
return $entry;
11
Added:
}
12
Added:
13
Added:
sub set ( $self, $key, $payload, $expires_epoch, $now ) {
14
Added:
$self->_prune($now);
15
Added:
$self->entries->{$key} = {
16
Added:
expires_epoch => $expires_epoch,
17
Added:
payload => $payload,
18
Added:
};
19
Added:
return;
20
Added:
}
21
Added:
22
Added:
sub _prune ( $self, $now ) {
23
Added:
my $entries = $self->entries;
24
Added:
for my $key ( keys %$entries ) {
25
Added:
delete $entries->{$key} if $entries->{$key}{expires_epoch} <= $now;
26
Added:
}
27
Added:
28
Added:
if ( keys(%$entries) >= $self->max_entries ) {
29
Added:
my ($oldest) = sort {
30
Added:
$entries->{$a}{expires_epoch} <=> $entries->{$b}{expires_epoch}
31
Added:
} keys %$entries;
32
Added:
delete $entries->{$oldest};
33
Added:
}
34
Added:
35
Added:
return;
36
Added:
}
37
Added:
38
Added:
1;
roles/dashboard/lib/FAPG/DAQ/Dashboard/Service/WeatherFetcher.pm
@@ -27,11 +27,12 @@
27
27
28
28
if ( my $err = $tx->error ) {
29
29
my $msg = $err->{message} || "HTTP $err->{code}";
30
Removed:
$self->log->warn("Weather fetch failed: $msg");
31
Removed:
return [];
30
Added:
die "Weather fetch failed: $msg";
32
31
}
33
32
34
33
my $json = $tx->result->json;
34
Added:
die 'Weather fetch returned an invalid hourly response'
35
Added:
if !$json || ref $json->{hourly} ne 'HASH';
35
36
return $self->_parse_response($json);
36
37
}
37
38
@@ -39,28 +40,30 @@
39
40
return [] unless $json && $json->{hourly};
40
41
41
42
my $hourly = $json->{hourly};
42
Removed:
my $times = $hourly->{time} || [];
43
Removed:
my $temps = $hourly->{temperature_2m} || [];
43
Added:
my $times = $hourly->{time} || [];
44
Added:
my $temps = $hourly->{temperature_2m} || [];
44
45
my $humids = $hourly->{relative_humidity_2m} || [];
45
Removed:
my $rains = $hourly->{rain} || [];
46
Added:
my $rains = $hourly->{rain} || [];
46
47
47
48
my @rows;
48
49
for my $i ( 0 .. $#$times ) {
49
50
my $epoch = _iso_to_epoch( $times->[$i] );
50
51
next unless defined $epoch;
51
52
52
Removed:
push @rows, {
53
Added:
push @rows,
54
Added:
{
53
55
epoch => $epoch,
54
56
temperature => $temps->[$i],
55
57
humidity => $humids->[$i],
56
58
rain => $rains->[$i],
57
Removed:
};
59
Added:
};
58
60
}
59
61
60
62
return \@rows;
61
63
}
62
64
63
65
sub _iso_to_epoch ($iso) {
66
Added:
64
67
# Format: "2026-08-10T00:00" (UTC, no timezone suffix)
65
68
return unless $iso && $iso =~ /\A(\d{4})-(\d{2})-(\d{2})T(\d{2}):(\d{2})/;
66
69
roles/dashboard/lib/FAPG/DAQ/Dashboard/Timeframe.pm
@@ -2,8 +2,9 @@
2
2
use Mojo::Base -strict;
3
3
4
4
use Exporter qw(import);
5
Removed:
use POSIX qw(strftime);
6
5
6
Added:
use FAPG::DAQ::Dashboard::Util qw(utc_timestamp);
7
Added:
7
8
our @EXPORT_OK = qw(max_series_range series_window);
8
9
9
10
my $MAX_SERIES_RANGE = 24;
@@ -36,7 +37,7 @@
36
37
37
38
sub series_window {
38
39
my ( $timeframe, $range ) = @_;
39
Removed:
my $config = $SERIES_TIMEFRAMES{$timeframe} or return undef;
40
Added:
my $config = $SERIES_TIMEFRAMES{$timeframe} or return;
40
41
my $bucket_stride = $config->{bucket_stride} * $range;
41
42
my $span_seconds = $config->{span_seconds} * $range;
42
43
my $end_epoch
@@ -51,11 +52,6 @@
51
52
start_iso => utc_timestamp($start_epoch),
52
53
end_iso => utc_timestamp($end_epoch),
53
54
};
54
Removed:
}
55
Removed:
56
Removed:
sub utc_timestamp {
57
Removed:
my ($epoch) = @_;
58
Removed:
return strftime( '%Y-%m-%dT%H:%M:%SZ', gmtime $epoch );
59
55
}
60
56
61
57
1;
roles/dashboard/lib/FAPG/DAQ/Dashboard/Util.pm
@@ -0,0 +1,24 @@
1
Added:
package FAPG::DAQ::Dashboard::Util;
2
Added:
use Mojo::Base -strict;
3
Added:
4
Added:
use Exporter qw(import);
5
Added:
use POSIX qw(strftime);
6
Added:
7
Added:
our @EXPORT_OK = qw(bounded_integer utc_timestamp);
8
Added:
9
Added:
sub bounded_integer {
10
Added:
my ( $value, $default, $minimum, $maximum ) = @_;
11
Added:
return $default
12
Added:
if !defined $value
13
Added:
|| $value !~ /\A\d+\z/
14
Added:
|| $value < $minimum
15
Added:
|| $value > $maximum;
16
Added:
return 0 + $value;
17
Added:
}
18
Added:
19
Added:
sub utc_timestamp {
20
Added:
my ($epoch) = @_;
21
Added:
return strftime( '%Y-%m-%dT%H:%M:%SZ', gmtime $epoch );
22
Added:
}
23
Added:
24
Added:
1;
roles/dashboard/t/04-readings-series.t
@@ -11,6 +11,7 @@
11
11
12
12
seed_series_readings($t);
13
13
seed_gap_readings($t);
14
Added:
seed_weighted_rollups($t);
14
15
15
16
$t->get_ok('/api/readings/ec/series?timeframe=hour')
16
17
->status_is(200)
@@ -50,7 +51,12 @@
50
51
$t->get_ok('/api/readings/do/series?timeframe=hour')->status_is(200);
51
52
ok( scalar( grep { !defined $_->{value} } $t->tx->res->json->{points}->@* ),
52
53
'series has null-value points for empty time buckets' );
54
Added:
ok( scalar( grep { $_->{gap} } $t->tx->res->json->{points}->@* ),
55
Added:
'series marks long acquisition gaps' );
53
56
57
Added:
$t->get_ok('/api/readings/orp/series?timeframe=hour')->status_is(200);
58
Added:
assert_series_point( $t, 9, 0, 10, 10 );
59
Added:
54
60
$t->get_ok('/api/readings/ec/series?timeframe=century')->status_is(400);
55
61
$t->get_ok('/api/readings/ec/series?timeframe=week&range=0')->status_is(400);
56
62
$t->get_ok('/api/readings/ec/series?timeframe=week&range=25')->status_is(400);
@@ -66,7 +72,8 @@
66
72
is( $populated[0]{value}, $value, 'series bucket has expected value' );
67
73
is( $populated[0]{min}, $min, 'series bucket has expected min' );
68
74
is( $populated[0]{max}, $max, 'series bucket has expected max' );
69
Removed:
is( $populated[0]{count}, $count, 'series bucket has expected sample count' );
75
Added:
is( $populated[0]{count},
76
Added:
$count, 'series bucket has expected sample count' );
70
77
71
78
return;
72
79
}
@@ -124,6 +131,42 @@
124
131
'do',
125
132
8,
126
133
'mg/L',
134
Added:
);
135
Added:
}
136
Added:
137
Added:
return;
138
Added:
}
139
Added:
140
Added:
sub seed_weighted_rollups {
141
Added:
my ($t) = @_;
142
Added:
my $series_stride = 2 * 60;
143
Added:
my $bucket_start
144
Added:
= int( time / $series_stride ) * $series_stride - $series_stride;
145
Added:
my $db = $t->app->sqlite->db;
146
Added:
147
Added:
for my $rollup (
148
Added:
{ epoch => $bucket_start, count => 1, total => 0, value => 0 },
149
Added:
{ epoch => $bucket_start + 60, count => 9, total => 90, value => 10 },
150
Added:
)
151
Added:
{
152
Added:
$db->query(
153
Added:
q{
154
Added:
INSERT OR REPLACE INTO reading_rollups (
155
Added:
probe, bucket_seconds, bucket_epoch, sample_count,
156
Added:
value_total, value_min, value_max,
157
Added:
first_sample_epoch, last_sample_epoch, max_gap_seconds
158
Added:
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
159
Added:
},
160
Added:
'orp',
161
Added:
60,
162
Added:
$rollup->{epoch},
163
Added:
$rollup->{count},
164
Added:
$rollup->{total},
165
Added:
$rollup->{value},
166
Added:
$rollup->{value},
167
Added:
$rollup->{epoch},
168
Added:
$rollup->{epoch},
169
Added:
0,
127
170
);
128
171
}
129
172
roles/dashboard/t/09-weather.t
@@ -13,7 +13,7 @@
13
13
my $weather = $t->app->weather;
14
14
ok $weather, 'weather helper is available';
15
15
16
Removed:
my $now = time;
16
Added:
my $now = time;
17
17
my @rows = map {
18
18
{ epoch => ( int( $now / 3600 ) - $_ ) * 3600,
19
19
temperature => 20 + $_ * 0.5,
@@ -39,11 +39,20 @@
39
39
my $weather = $t->app->weather;
40
40
my $epoch = int( time / 3600 ) * 3600;
41
41
42
Removed:
$weather->store_batch( [ { epoch => $epoch, temperature => 25, humidity => 70, rain => 0 } ] );
43
Removed:
$weather->store_batch( [ { epoch => $epoch, temperature => 26, humidity => 71, rain => 0.1 } ] );
42
Added:
$weather->store_batch(
43
Added:
[ { epoch => $epoch, temperature => 25, humidity => 70, rain => 0 } ]
44
Added:
);
45
Added:
$weather->store_batch(
46
Added:
[ { epoch => $epoch,
47
Added:
temperature => 26,
48
Added:
humidity => 71,
49
Added:
rain => 0.1
50
Added:
}
51
Added:
]
52
Added:
);
44
53
45
54
my $range = $weather->get_range( $epoch, $epoch );
46
Removed:
is scalar @$range, 1, 'only one row for the same epoch';
55
Added:
is scalar @$range, 1, 'only one row for the same epoch';
47
56
is $range->[0]{temperature}, 26, 'second insert replaces first';
48
57
};
49
58
@@ -51,7 +60,7 @@
51
60
$t->get_ok('/api/weather/hourly?hours=24')
52
61
->status_is(200)
53
62
->json_has('/points')
54
Removed:
->json_is('/hours' => 24);
63
Added:
->json_is( '/hours' => 24 );
55
64
56
65
my $points = $t->tx->res->json->{points};
57
66
ok @$points > 0, 'hourly returns seeded weather points';
@@ -62,12 +71,22 @@
62
71
};
63
72
64
73
subtest 'weather API forecast endpoint' => sub {
74
Added:
65
75
# Seed some future data
66
76
my $future = ( int( time / 3600 ) + 2 ) * 3600;
67
Removed:
$t->app->weather->store_batch( [
68
Removed:
{ epoch => $future, temperature => 28, humidity => 55, rain => 0 },
69
Removed:
{ epoch => $future + 3600, temperature => 27, humidity => 58, rain => 0.5 },
70
Removed:
] );
77
Added:
$t->app->weather->store_batch(
78
Added:
[ { epoch => $future,
79
Added:
temperature => 28,
80
Added:
humidity => 55,
81
Added:
rain => 0
82
Added:
},
83
Added:
{ epoch => $future + 3600,
84
Added:
temperature => 27,
85
Added:
humidity => 58,
86
Added:
rain => 0.5
87
Added:
},
88
Added:
]
89
Added:
);
71
90
72
91
$t->get_ok('/api/weather/forecast')
73
92
->status_is(200)
@@ -78,25 +97,47 @@
78
97
ok @$points > 0, 'forecast returns points';
79
98
};
80
99
81
Removed:
subtest 'weather API current endpoint' => sub {
100
Added:
subtest 'weather API current endpoint excludes future forecasts' => sub {
101
Added:
my $current_hour = int( time / 3600 ) * 3600;
102
Added:
my $future = $current_hour + 3600;
103
Added:
104
Added:
$t->app->weather->store_batch(
105
Added:
[ { epoch => $current_hour,
106
Added:
temperature => 22,
107
Added:
humidity => 65,
108
Added:
rain => 0,
109
Added:
},
110
Added:
{ epoch => $future,
111
Added:
temperature => 99,
112
Added:
humidity => 1,
113
Added:
rain => 10,
114
Added:
},
115
Added:
]
116
Added:
);
117
Added:
118
Added:
my $latest = $t->app->weather->latest_at(time);
119
Added:
is $latest->{epoch}, $current_hour,
120
Added:
'model selects the latest point at or before now';
121
Added:
82
122
$t->get_ok('/api/weather/current')
83
123
->status_is(200)
84
Removed:
->json_has('/temperature')
85
Removed:
->json_has('/humidity')
86
Removed:
->json_has('/rain')
87
Removed:
->json_has('/timestamp');
124
Added:
->json_is( '/temperature' => 22 )
125
Added:
->json_is( '/timestamp' => utc_timestamp($current_hour) );
88
126
};
89
127
90
Removed:
subtest 'weather API refresh endpoint' => sub {
91
Removed:
$t->get_ok('/api/weather/refresh')
128
Added:
subtest 'weather API refresh endpoint is POST-only' => sub {
129
Added:
local $t->app->config->{weather}{enabled} = 0;
130
Added:
131
Added:
$t->get_ok('/api/weather/refresh')->status_is(404);
132
Added:
$t->post_ok('/api/weather/refresh')
92
133
->status_is(200)
93
Removed:
->json_has('/status');
134
Added:
->json_is( '/status' => 'disabled' );
94
135
};
95
136
96
137
subtest 'weather page renders' => sub {
97
138
$t->get_ok('/weather')
98
139
->status_is(200)
99
Removed:
->text_is('title', 'FAPG DAQ Weather')
140
Added:
->text_is( 'title', 'FAPG DAQ Weather' )
100
141
->element_exists('canvas#chart-weather-combined')
101
142
->element_exists('[data-default-timeframe="day"]');
102
143
};
@@ -118,12 +159,17 @@
118
159
my $empty = test_empty_app();
119
160
$empty->get_ok('/api/weather/hourly')
120
161
->status_is(200)
121
Removed:
->json_is('/points' => []);
162
Added:
->json_is( '/points' => [] );
122
163
$empty->get_ok('/api/weather/current')
123
164
->status_is(200)
124
Removed:
->json_is('/temperature' => undef);
125
Removed:
$empty->get_ok('/weather')
126
Removed:
->status_is(200);
165
Added:
->json_is( '/temperature' => undef );
166
Added:
$empty->get_ok('/weather')->status_is(200);
127
167
};
128
168
129
169
done_testing;
170
Added:
171
Added:
sub utc_timestamp {
172
Added:
my ($epoch) = @_;
173
Added:
require POSIX;
174
Added:
return POSIX::strftime( '%Y-%m-%dT%H:%M:%SZ', gmtime $epoch );
175
Added:
}
roles/dashboard/t/10-secret.t
@@ -0,0 +1,43 @@
1
Added:
use Mojo::Base -strict;
2
Added:
3
Added:
use Test2::V0;
4
Added:
5
Added:
use File::Spec;
6
Added:
use File::Temp qw(tempdir);
7
Added:
use FindBin;
8
Added:
use lib "${FindBin::Bin}/../lib/";
9
Added:
10
Added:
use FAPG::DAQ::Dashboard::Secret qw(load_or_create_secret read_secret);
11
Added:
12
Added:
my $directory = tempdir( CLEANUP => 1 );
13
Added:
my $path = File::Spec->catfile( $directory, '.mojo-secrets' );
14
Added:
15
Added:
my ( $generated, $created ) = load_or_create_secret($path);
16
Added:
like $generated, qr/\A[0-9a-f]{64}\z/,
17
Added:
'generated secret has 256 bits encoded as hex';
18
Added:
ok $created, 'missing secret is created';
19
Added:
is read_secret($path), $generated, 'generated secret is persisted';
20
Added:
is( ( stat $path )[2] & 0777, 0600, 'secret file is owner-readable only' );
21
Added:
22
Added:
my ( $reused, $created_again ) = load_or_create_secret($path);
23
Added:
is $reused, $generated, 'existing secret is reused';
24
Added:
ok !$created_again, 'existing secret is not rewritten';
25
Added:
26
Added:
my $invalid_path = File::Spec->catfile( $directory, 'invalid-secret' );
27
Added:
open my $invalid_fh, '>', $invalid_path
28
Added:
or die "Unable to create invalid test secret: $!";
29
Added:
print {$invalid_fh} "short\n";
30
Added:
close $invalid_fh or die "Unable to close invalid test secret: $!";
31
Added:
32
Added:
like dies { load_or_create_secret($invalid_path) },
33
Added:
qr/missing or too short/,
34
Added:
'an existing invalid secret fails closed';
35
Added:
36
Added:
unlink $invalid_path or die "Unable to remove invalid test secret: $!";
37
Added:
symlink 'missing-target', $invalid_path
38
Added:
or die "Unable to create test symlink: $!";
39
Added:
like dies { load_or_create_secret($invalid_path) },
40
Added:
qr/missing or too short/,
41
Added:
'a dangling secret symlink is not replaced';
42
Added:
43
Added:
done_testing;