Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
76 changes: 70 additions & 6 deletions lib/data/model/server/status_history.dart
Original file line number Diff line number Diff line change
Expand Up @@ -116,13 +116,26 @@ class StatusHistory {
_netTx.add(netTxs, time.length);
}

/// Replaces the buffer with [samples], oldest first. Used to prefill from
/// monitor's stored history before live sampling takes over; a no-op once
/// live samples exist, so a late-arriving history response can't rewind
/// what has already been charted.
/// Puts the agent's stored [samples] (oldest first) before the live ones,
/// so a page opened on a monitor server draws the last hour at once.
///
/// Only the part older than the first live sample goes in: the live samples
/// are what has been drawn, and a late response must not rewrite them. They
/// were often there already — the home page samples every server as it
/// lists it — and when that first one has no rate yet (it is the first),
/// keeping the history out until the buffer was empty left the chart with
/// nothing to draw until the agent's next cycle. The capacity keeps the
/// newest: past it, the oldest stored samples go, never live ones.
void seed(List<StatusHistorySample> samples) {
if (!isEmpty || samples.isEmpty) return;
for (final s in samples) {
final firstLive = isEmpty ? null : time.first;
final older = [
for (final s in samples)
if (firstLive == null || s.timeMs < firstLive) s,
];
if (older.isEmpty) return;
final live = [for (var i = 0; i < length; i++) _sampleAt(i)];
_clear();
for (final s in older) {
add(
timeMs: s.timeMs,
cpu: s.cpu,
Expand All @@ -137,6 +150,47 @@ class StatusHistory {
battery: s.battery,
);
}
for (final replay in live) {
replay();
}
}

/// The sample at [i] as a call that appends it again, devices included.
void Function() _sampleAt(int i) {
final (t, c, m, sw, d, rx, tx, dr, dw, g, te, b) = (
time[i], cpu[i], mem[i], swap[i], disk[i], netRx[i], netTx[i],
diskRead[i], diskWrite[i], gpu[i], temp[i], battery[i],
);
final temps = _temps.at(i), reads = _diskReads.at(i), writes = _diskWrites.at(i);
final rxs = _netRx.at(i), txs = _netTx.at(i);
return () => add(
timeMs: t,
cpu: c,
mem: m,
swap: sw,
disk: d,
netRx: rx,
netTx: tx,
diskRead: dr,
diskWrite: dw,
gpu: g,
temp: te,
temps: temps,
diskReads: reads,
diskWrites: writes,
netRxs: rxs,
netTxs: txs,
battery: b,
);
}

void _clear() {
for (final f in [time, cpu, mem, swap, disk, netRx, netTx, diskRead, diskWrite, gpu, temp, battery]) {
f.clear();
}
for (final d in [_temps, _diskReads, _diskWrites, _netRx, _netTx]) {
d.byDevice.clear();
}
}
}

Expand All @@ -160,8 +214,18 @@ class _DeviceHistory {
return f;
});
series.add(values[device]);
// Gone for a whole buffer: nothing of it is left to draw, and kept, a
// host cycling through short-lived interfaces grows a series for each.
if (!values.containsKey(device) && series.every((e) => e == null)) {
byDevice.remove(device);
}
}
}

/// What sample [i] had per device, the devices without a reading left out.
Map<String, double> at(int i) => {
for (final MapEntry(:key, :value) in byDevice.entries) key: ?value[i],
};
}

/// One point handed to [StatusHistory.seed].
Expand Down
3 changes: 3 additions & 0 deletions lib/data/model/server/time_seq.dart
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,9 @@ class Fifo<T> extends ListBase<T> {
@override
void operator []=(int index, T value) => _list[index] = value;

/// Empties it. [ListBase.clear] would go through the [length] setter.
@override
void clear() => _list.clear();
}

/// A two-sample window over a list of counters that gets re-collected every
Expand Down
21 changes: 19 additions & 2 deletions lib/data/provider/server/single.dart
Original file line number Diff line number Diff line change
Expand Up @@ -144,6 +144,11 @@ class ServerNotifier extends _$ServerNotifier {
bool _usePersistentShellForStatus = true;
int _operationGeneration = 0;

/// The generation [seedHistory] last asked the source for, so each
/// connection asks once however many places want the trend. Cleared when
/// the asking fails, for the next poll to try again.
int? _seededOperation;

/// Last [ShellFunc.statusExt] output, appended to every status refresh's raw
/// output so its segments don't blank out on the polls in between — see
/// [_extendedStatusInterval]
Expand Down Expand Up @@ -177,6 +182,10 @@ class ServerNotifier extends _$ServerNotifier {
@override
ServerState build(String serverId) {
ref.onDispose(() {
// Supersedes whatever is in flight, so it ends at its next
// [_isRefreshCurrent] — which compares this before reading `state`,
// and reading `state` after disposal throws.
_operationGeneration++;
unawaited(_disposePersistentShell());
_source?.close();
try {
Expand Down Expand Up @@ -458,6 +467,10 @@ class ServerNotifier extends _$ServerNotifier {
if (_refreshingOperation == operation) return;

_refreshingOperation = operation;
// Alongside the first poll rather than after it: the card in the list
// draws the same trend as the page, and built up from live samples alone
// its first point waited for two agent cycles — rates need a pair.
unawaited(seedHistory());
try {
// Somebody pressed Retry, and the first thing they are owed is that it
// registered. Both paths raise a connecting/loading state of their own,
Expand Down Expand Up @@ -812,8 +825,9 @@ class ServerNotifier extends _$ServerNotifier {
/// Prefills [ServerStatus.history] from whatever trend data the source
/// already holds, so a freshly opened detail page shows a trend instead of
/// building one up from scratch. A no-op for sources without
/// [ServerCapabilities.storedHistory], and once live samples exist — see
/// [StatusHistory.seed].
/// [ServerCapabilities.storedHistory]; what is older than the first live
/// sample goes before it — see [StatusHistory.seed]. Once per connection:
/// [refresh] asks with every poll, and the page again when it opens.
Future<void> seedHistory({int minutes = 60}) async {
final generation = _operationGeneration;
final spi = state.spi;
Expand All @@ -822,8 +836,10 @@ class ServerNotifier extends _$ServerNotifier {
// leading one meant `capabilities.storedHistory` advertised a trend the
// page then never seeded, and the chart built up from empty exactly where
// an agent had months of it.
if (_seededOperation == generation) return;
final credential = _historyCredential(spi);
if (credential == null) return;
_seededOperation = generation;
try {
// Asking for exactly what the buffer holds. Any more is averaged down
// on the agent's side instead of being carried here and dropped by
Expand All @@ -840,6 +856,7 @@ class ServerNotifier extends _$ServerNotifier {
updateStatus(_copyStatus(state.status));
} catch (e, s) {
if (!_isRefreshCurrent(generation, spi)) return;
_seededOperation = null;
Loggers.app.warning('Seed history for ${spi.name}', e, s);
}
}
Expand Down
48 changes: 48 additions & 0 deletions test/unit/server/status_history_test.dart
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,19 @@ void main() {
expect(h.netTxByDevice['tun0'], [2, null]);
});

test('a device gone for a whole buffer is dropped', () {
final h = StatusHistory();

h.add(timeMs: 0, netTxs: const {'eth0': 1, 'tun0': 2});
for (var t = 1; t < StatusHistory.capacity; t++) {
h.add(timeMs: t, netTxs: const {'eth0': 1});
}
expect(h.netTxByDevice, contains('tun0'));

h.add(timeMs: StatusHistory.capacity, netTxs: const {'eth0': 1});
expect(h.netTxByDevice.keys, ['eth0']);
});

test('each metric keeps its own devices', () {
final h = StatusHistory();

Expand Down Expand Up @@ -88,4 +101,39 @@ void main() {

expect(h.gpu, [41, null]);
});

group('seed', () {
test('fills an empty buffer', () {
final h = StatusHistory();
h.seed(const [StatusHistorySample(timeMs: 1, cpu: 10), StatusHistorySample(timeMs: 2, cpu: 20)]);
expect(h.time.toList(), [1, 2]);
expect(h.cpu.toList(), [10, 20]);
});

test('goes before the live samples, which keep everything they had', () {
final h = StatusHistory();
// The first live sample has no CPU rate yet.
h.add(timeMs: 100, cpu: null, mem: 50, netRxs: const {'eth0': 7});
h.seed(const [
StatusHistorySample(timeMs: 10, cpu: 1),
StatusHistorySample(timeMs: 20, cpu: 2),
// Newer than what is live: the live sample stands.
StatusHistorySample(timeMs: 100, cpu: 99),
StatusHistorySample(timeMs: 150, cpu: 99),
]);
expect(h.time.toList(), [10, 20, 100]);
expect(h.cpu.toList(), [1, 2, null]);
expect(h.mem.toList(), [null, null, 50]);
expect(h.netRxByDevice['eth0']?.toList(), [null, null, 7]);
});

test('past the capacity the oldest stored samples go, never live ones', () {
final h = StatusHistory();
h.add(timeMs: 1000000, cpu: 5);
h.seed([for (var i = 0; i < StatusHistory.capacity; i++) StatusHistorySample(timeMs: i, cpu: 1)]);
expect(h.length, StatusHistory.capacity);
expect(h.time.last, 1000000);
expect(h.time.first, 1);
});
});
}
Loading