Files
asepharyana 1cdb82a76f Sync config from arch
- hypr/apps.lua
- hypr/autostart.lua
- hypr/envs.lua
- hypr/hyprland.lua
- hypr/hyprsunset.conf
- hypr/input.lua
- hypr/looknfeel.lua
- hypr/omasettings.lua
- hypr/xdph.conf
- omarchy/branding/about.txt
- omarchy/branding/screensaver.txt
- omarchy/extensions/omarchy-menu.jsonc
- omarchy/hooks/battery-low.d/play-warning-sound.sample
- omarchy/hooks/font-set.d/show-font-notification.sample
- omarchy/hooks/post-boot.d/weather.sample
- omarchy/hooks/post-update.d/install-voxtype.hook
- omarchy/hooks/post-update.d/setup-agent.hook
- omarchy/hooks/post-update.d/setup-fingerprint.hook
- omarchy/hooks/post-update.d/show-update-notification.sample
- omarchy/hooks/pre-refresh-pacman.d/add-custom-repo.sample
- omarchy/hooks/theme-set.d/show-theme-notification.sample
- omarchy/shell.json
- omarchy/shell.toml
- omarchy/theme.name
- omarchy/themes/azure-glow/README.md
- omarchy/themes/azure-glow/alacritty.toml
- omarchy/themes/azure-glow/btop.theme
- omarchy/themes/azure-glow/hyprland.conf
- omarchy/themes/azure-glow/hyprlock.conf
- omarchy/themes/azure-glow/icons.theme
- … 269 more
2026-09-23 15:19:12 +07:00

2057 lines
94 KiB
Python

"""The nexthopd main loop.
Owns the probes, folds their samples into live.json / recent.json /
apps.json / history.db, detects outages, and runs the scheduled content
check. The QML side never talks to this process — the files are the whole
contract, so either side can restart without the other noticing.
"""
import fcntl
import json
import os
import re
import shutil
import signal
import stat as stat_module
import subprocess
import sys
import threading
import time
from typing import NamedTuple
from collections import deque
from . import __version__, apps, linkevents, net, score, speedtest
from .paths import (ensure_state_dir, ensure_runtime_dir, runtime_dir,
live_path, recent_path, db_path, lock_path, apps_path)
from .instruments import Bench, MergedSeries, merged_stats
from .probes import Series, PingProbe, TcpProbe
from .state import write_atomic, retire_legacy_snapshots
from .store import Store
from .update import UpdateWatch
# Unbroken silence on a leg — no probe of any kind answering — before we
# call it down. Measured on the probe stream itself, from the timestamp of
# the first lost sample after the last reply, so it means four seconds
# whatever the probe cadence or the loop's tick. Long enough to skip a
# Wi-Fi roam, short enough that the alarm still feels immediate.
#
# History, because it cost a year: until 0.2.15 this was a COUNT of loop
# ticks (eight) in which a 3 s any-reply window came back empty, so a
# "loss" was three seconds wide, an outage took ~6.5 s while every comment
# said four, and an interruption shorter than 3 s could not register at
# all. Reading the stream is what makes the number mean what it says.
OUTAGE_AFTER_S = 4.0
# A run of silence shorter than an outage but longer than noise: an
# interruption the user may well have felt — a call breaking up, a stream
# rebuffering. Recorded at recovery, when its length is known, and charged
# to Reliability at half weight (score.reliability). Three lost probes at
# the default 500 ms.
DISRUPTION_AFTER_S = 1.5
# One lost probe is never an event, at any cadence. At a slow probe interval
# a single loss is followed by seconds with no sample at all, and measured
# in seconds alone that would read as an outage.
MIN_LOST_SAMPLES = 2
# A content check measures the line, and a Wi-Fi link that has just
# associated is not the line yet: it may still be on the band it landed on
# rather than the one it will roam to, and its transmit rate may still be
# climbing. Measured on this laptop: a check 55 s after associating read
# 63 Mbps on a 380 Mbps line, from 2.4 GHz at a 16 Mbps tx rate, nine
# minutes before the link moved itself to 5 GHz.
CHECK_SETTLE_S = 60.0
# But a link that is simply slow must still be measured eventually, or Speed
# never scores at all. Defer this long at most, then take what is there.
CHECK_DEFER_MAX_S = 600.0
# How recently a peak test must have run, on this network, to be allowed to
# contradict the everyday basis.
PEAK_FRESH_S = 3600.0
def check_ready(now: float, assoc_since, rate_low_since, waiting_since,
settle_s: float = CHECK_SETTLE_S,
max_defer_s: float = CHECK_DEFER_MAX_S) -> bool:
"""Is the link in a fit state to be measured? Pure, so it is testable.
Two reasons to wait: the association is younger than the settle window,
or the transmit rate is currently down (`LinkWatch.low_since`, the same
signal that opens a rate-drop event). Either way the cap wins in the
end — a permanently poor link gets an honest low number rather than no
number, and the median guard is what protects the score from one bad
sample.
"""
if waiting_since is not None and now - waiting_since >= max_defer_s:
return True
if assoc_since is not None and now - assoc_since < settle_s:
return False
if rate_low_since is not None:
return False
return True
# A stream whose newest sample is older than this has stopped talking — a
# dead ping process, a stopped probe — which is not an outage. `ping -O`
# and the TCP probe keep emitting losses through a real one, so a stale
# stream is never mistaken for a dead line: it is unknown, and not counted.
LEG_STALE_S = 6.0
# How far back to read a leg's stream for the start of the current run. Past
# OUTAGE_AFTER_S with margin; the watch remembers an older start itself.
LEG_STREAM_WINDOW_S = 12.0
# Probes needed on each side of the idle/loaded split before their ratio is
# reported. Below this the comparison is sampling noise.
MIN_LOAD_SPLIT_SAMPLES = 10
# What counts as a busy link, for the loaded/idle latency split only.
#
# This used to borrow `LinkWatch.TRAFFIC_FLOOR_BPS`, whose own comment says it
# exists so Wi-Fi power save does not fire spurious rate-drop events. That is
# a question about the radio; this is a question about the line, and one
# constant cannot answer both. At 25 kB/s it answered neither: on this laptop
# the median minute carries 18 kB/s and p75 is 42 kB/s, so the floor sat
# inside the IDLE distribution and tagged 36.6% of all minutes "loaded".
#
# The proof it measured nothing is in the stored history: across 8,971
# minutes carrying both figures, the loaded half was FASTER than the idle
# half 57% of the time, with medians 18.7 and 18.8 ms. A link cannot answer
# faster while busy; a coin flip is what two buckets holding the same thing
# look like. Selecting minutes by how much they actually carried recovers the
# signal, and only well up the range: at 250 kB/s inversions are 51%, at
# 1 MB/s 50%, at 2.5 MB/s 47%, and only at 5 MB/s do they fall to 29% with
# loaded 19.9 ms against idle 16.6 — the direction physics requires.
#
# That selection is weaker than it looks and the fraction below rests mostly
# on the 57% above, not on it. A minute's stored `rx_bps` is `self.rates`,
# the 3-second sliding window, so it describes the END of a minute rather
# than the minute: of 151 content checks, the median stored rate in the
# check's own minute is 29 kB/s, below even the whole-minute average of
# 233 kB/s, because a sub-second check rarely lands in the stored window.
# The probes' own load flags — which is what the 57% is built from — are
# unaffected, since each probe carries the tag it was measured under.
#
# 5 MB/s is a tenth of what this line carries, and a tenth is the number
# worth keeping rather than the 5, because a fixed rate cannot serve a
# 10 Mbps line and a gigabit one at once — the same lesson the content check
# learned about fixed transfer sizes. Below the floor nothing is called busy,
# so a line whose capacity is unknown does not tag its own background chatter.
LOAD_FRACTION_OF_LINE = 0.10
LOAD_FLOOR_BPS = 125_000
# Queueing can only ADD delay, so a loaded/idle ratio below 1 says the link
# answered faster while busy, which is not a measurement. The sample floor
# above does not catch it: 0.87 was published live on 716 samples per side.
# Within a few percent of 1 the two populations are simply indistinguishable
# and the honest reading is "no inflation"; further below, the split itself
# is untrustworthy — the loaded samples likely landed in a quiet moment — so
# withhold rather than report. A plausibility floor, distinct from a sample
# floor, and the guard the socket metric already has.
MIN_PLAUSIBLE_INFLATION = 0.95
# The TCP instruments: a handshake to port 443, once a second per target.
# Slower than the ICMP cadence on purpose — each one opens a real connection
# to someone else's server. Since 0.2.0 these are seated instruments in the
# bench, so a TCP series feeds the scored internet leg and the outage watch
# whenever it holds a seat (instruments.py).
TCP_PROBE_INTERVAL_S = 1.0
TCP_PROBE_PORT = 443
# The rest of the instrument pool (see instruments.py). Cloudflare edge
# is a host the daemon already fetches from; dns.google is the one
# probe target outside Cloudflare, so a Cloudflare incident cannot
# silence the whole pool. TCP handshakes only — no payload.
CF_EDGE_HOST = "speed.cloudflare.com"
DIVERSITY_HOST = "dns.google"
# A benched instrument idles at a tenth of its seated cadence: enough
# to stay rankable, cheap enough to keep around.
STANDBY_FACTOR = 10.0
BENCH_EVAL_EVERY_S = 60.0
def proc_start_ticks(pid):
"""The process start time in clock ticks, from /proc/<pid>/stat.
Together with the pid it forms a start identity: pids are recycled,
but a recycled pid never reproduces the same start time. The shell
service checks this before it will signal anything.
"""
try:
with open("/proc/%d/stat" % pid, "rb") as f:
data = f.read(4096)
# Field 22, counted after the parenthesised comm (which may itself
# contain spaces and parentheses).
rest = data[data.rindex(b")") + 2:].split()
return int(rest[19])
except (OSError, ValueError, IndexError):
return None
# An interruption that self-heals in under this is recorded but not
# notified. The constant was written with 0.1.0 and never read: outages
# alarmed the instant they were declared, so a six-second blip on a flaky
# link fired a desktop notification the user could do nothing about. The
# event is always logged; only the interruption goes quiet.
NOTIFY_AFTER_S = 5.0
# A content check that failed outright (curl error, endpoint down) used to
# wait the full interval before trying again — an hour of stale Speed for
# a transient fault. One retry after this long; a second failure waits the
# interval, so a blocked endpoint is not hammered.
CONTENT_RETRY_S = 300.0
class Config:
"""Settings, read from the file the QML side writes.
Every value is validated against the same ranges the manifest schema
promises, and the file itself has a size cap — a config the daemon
cannot trust in full is a config it ignores in full. Nothing read
here can grow retained state beyond its documented bounds.
"""
MAX_BYTES = 64 * 1024
# key -> (default, validator). Ranges mirror manifest.json's schema.
# A hostname or address never begins with a dash, and the anchor is
# the last argument to `ping` and `ip route get`, where a leading dash
# would be read as an option. Refused at the setting, so no call site
# has to remember an option terminator.
ANCHOR_RE = re.compile(r"^[A-Za-z0-9:][A-Za-z0-9.:\-]{0,252}$")
@staticmethod
def _int(lo, hi):
def check(v):
if isinstance(v, bool) or not isinstance(v, (int, float)):
return None
n = int(v)
return n if lo <= n <= hi else None
return check
@staticmethod
def _bool(v):
return v if isinstance(v, bool) else None
SCHEMA = {
"internetAnchor": ("1.1.1.1",
lambda v: v if isinstance(v, str)
and Config.ANCHOR_RE.match(v) else None),
"probeIntervalMs": (500, None), # filled below
"contentSpeed": (True, None),
"contentSpeedIntervalMin": (60, None),
"peakEngine": ("Auto",
lambda v: v if v in ("Auto", "Ookla", "Cloudflare",
"fast.com") else None),
"planDownMbps": (0, None),
"planUpMbps": (0, None),
"notifyOutage": (True, None),
"updateCheck": (True, None),
"meteredCare": (True, None),
"historyDays": (7, None),
"throughputWindowS": (3, None),
}
DEFAULTS = {k: v[0] for k, v in SCHEMA.items()}
def __init__(self, state_dir):
self.path = state_dir / "config.json"
self.values = dict(self.DEFAULTS)
self._mtime = 0
def refresh(self):
# Everything is checked on the file descriptor actually read — a
# stat followed by a separate open is a race an attacker wins by
# swapping the file in between. O_NOFOLLOW refuses symlinks, fstat
# types and dates the very fd we read, and the size bound is
# enforced by the bounded read itself, not by a prior check.
try:
fd = os.open(self.path, os.O_RDONLY | os.O_NOFOLLOW | os.O_CLOEXEC)
except OSError:
return
try:
st = os.fstat(fd)
if not stat_module.S_ISREG(st.st_mode):
return
if st.st_mtime == self._mtime:
return
self._mtime = st.st_mtime
chunks = []
remaining = self.MAX_BYTES + 1
while remaining > 0:
chunk = os.read(fd, remaining)
if not chunk:
break
chunks.append(chunk)
remaining -= len(chunk)
data = b"".join(chunks)
except OSError:
return
finally:
os.close(fd)
if len(data) > self.MAX_BYTES:
return
try:
loaded = json.loads(data)
except ValueError:
return
if not isinstance(loaded, dict):
return
merged = dict(self.DEFAULTS)
for key, (default, validate) in self.SCHEMA.items():
if key not in loaded or loaded[key] is None:
continue
checked = validate(loaded[key]) if validate else None
if checked is not None:
merged[key] = checked
self.values = merged
def __getitem__(self, key):
return self.values[key]
# Range validators mirror the manifest schema exactly; a value outside its
# documented range is discarded, never clamped — silence over surprise.
Config.SCHEMA["probeIntervalMs"] = (500, Config._int(250, 5000))
Config.SCHEMA["contentSpeed"] = (True, Config._bool)
Config.SCHEMA["contentSpeedIntervalMin"] = (60, Config._int(15, 1440))
Config.SCHEMA["planDownMbps"] = (0, Config._int(0, 10000))
Config.SCHEMA["planUpMbps"] = (0, Config._int(0, 10000))
Config.SCHEMA["notifyOutage"] = (True, Config._bool)
Config.SCHEMA["updateCheck"] = (True, Config._bool)
Config.SCHEMA["meteredCare"] = (True, Config._bool)
Config.SCHEMA["historyDays"] = (7, Config._int(1, 90))
Config.SCHEMA["throughputWindowS"] = (3, Config._int(1, 30))
class LinkWatch:
"""Watches the Wi-Fi link state and writes events worth remembering:
roams, kicks, drops, associations, sustained rate drops. Instant events
are stored closed; a rate drop stays open until the rate recovers, so
its row carries a duration.
A BSSID change is attributed when `events` (an NlEvents) knows who
ended the previous association: the AP (a kick, with its 802.11
reason), this machine with the next authentication already under way
(a roam), or this machine after a scan (a drop — the link was lost).
Without that knowledge every change is a roam, as it always was.
"""
# A drop only counts when the rate stays below this fraction of the
# recent ceiling for a sustained stretch — rate control flaps all the
# time and a log that records every flap teaches people to ignore it.
LOW_FRACTION = 0.4
RECOVER_FRACTION = 0.6
SUSTAIN_S = 10.0
# Rate drops are only meaningful under traffic: Wi-Fi power save
# renegotiates a low bitrate the moment the link idles, and logging
# that teaches people to ignore the log. Below this many bytes/sec of
# combined throughput the link counts as idle.
TRAFFIC_FLOOR_BPS = 25_000
# Consecutive empty link reads before the link counts as genuinely
# gone. A single failed `iw` call is a hiccup, not a disassociation —
# and its recovery must not be logged as a fresh association.
GAP_SAMPLES = 5
# A local deauth followed by a new authentication within this long is
# the client roaming: mac80211 emits that deauth from inside the call
# that starts the new authentication, so a roam's gap is milliseconds.
# A lost link is followed by a scan first, and its gap is seconds. The
# threshold sits between the two by orders of magnitude.
ROAM_FOLLOW_S = 1.0
def __init__(self, store, events=None):
self.store = store
self.events = events # NlEvents, or anything with cause_for()
self.prev = None # last non-empty link, None until first seen
self.last_link_t = None # when the link was last seen up
self.gap_count = 0
self.disassociated = False
self.rate_ceiling = 0.0
self.low_since = None
self.low_floor = None
self.rate_event_id = None
self.last_sample = 0.0
# When the current association began, so a measurement can wait for
# a link that has only just come up — see check_ready.
self.assoc_since = None
def _instant(self, ts, kind, detail):
# Severity travels WITH the event so the panel can colour a kind it
# has never heard of. A kick or a drop is a fault; a roam, an
# association or a channel change is information.
severity = "warn" if kind in ("kick", "drop") else "info"
eid = self.store.open_event(int(ts), kind, severity, "local", detail)
self.store.close_event(eid, int(ts))
def _cause(self, bssid, since, now):
"""Who ended our association with `bssid`, if anything was seen
since we last saw that link up."""
if not self.events or not bssid or since is None:
return None
try:
return self.events.cause_for(bssid, now, now - since + 2.0)
except Exception:
return None # an attribution failure must not cost the event
def _blame(self, cause, old, new, gap_s=None):
"""(kind, text) for a change away from `old` with a known cause.
gap_s is how long the link was down when that is known (a confirmed
gap); otherwise the deauth-to-reauth delay stands in for it.
"""
why = linkevents.reason_text(cause["reason"], cause["by_ap"])
follow = cause.get("gap_s")
if cause["by_ap"]:
kind, lead = "kick", "Kicked by AP %s (%s)" % (old, why)
elif gap_s is None and follow is not None and follow < self.ROAM_FOLLOW_S:
return "roam", "Roamed to " + new
else:
kind, lead = "drop", "Dropped by this machine (%s)" % why
down = gap_s if gap_s is not None else follow
text = lead + ", rejoined" + ("" if new == old else " via " + new)
if down is not None and down >= 0.5:
text += " after " + _short_duration(down)
return kind, text
def sample(self, now, link, traffic_bps=0.0):
# The caller runs twice a second; once a second is plenty here.
if now - self.last_sample < 1.0:
return
self.last_sample = now
if not link or not link.get("bssid"):
# An empty read is a hiccup until it persists: `iw` times out
# now and then, and treating each blink as a disassociation
# spammed the log with fake re-associations.
self.gap_count += 1
if self.gap_count == self.GAP_SAMPLES:
self.disassociated = True
self._close_rate_event(now)
self.rate_ceiling = 0.0
return
self.gap_count = 0
since, self.last_link_t = self.last_link_t, now
prev, self.prev = self.prev, dict(link)
bssid = link.get("bssid", "")
prev_bssid = prev.get("bssid", "") if prev else ""
ssid = link.get("ssid", "")
if prev is None:
# The daemon's first sighting of an existing link is not an
# association — logging it stamped every daemon restart into
# the event log. It is still the moment we learned of this one,
# so the settle window starts here.
self.disassociated = False
self.assoc_since = now
return
if self.disassociated:
self.disassociated = False
self.assoc_since = now
cause = self._cause(prev_bssid, since, now)
if cause:
# One row for the whole incident: who ended it, how long it
# took to come back, and where. The plain association is
# for gaps nobody claimed — suspend, or no `iw event`.
kind, text = self._blame(cause, prev_bssid, bssid, gap_s=now - since)
self._instant(now, kind, text)
else:
self._instant(now, "associate", "Associated with " + (ssid or bssid))
elif bssid and prev_bssid and bssid != prev_bssid:
cause = self._cause(prev_bssid, since, now)
kind, lead = self._blame(cause, prev_bssid, bssid) if cause \
else ("roam", "Roamed to " + bssid)
parts = [lead]
if prev.get("channel") and link.get("channel") \
and prev["channel"] != link["channel"]:
parts.append("channel %s \u2192 %s" % (prev["channel"], link["channel"]))
if prev.get("signal_dbm") is not None and link.get("signal_dbm") is not None:
parts.append("%s \u2192 %s dBm" % (prev["signal_dbm"], link["signal_dbm"]))
self._instant(now, kind, ", ".join(parts))
# A different AP has a different honest ceiling, and a different
# band: this is a fresh association as far as measuring goes.
self.rate_ceiling = 0.0
self.assoc_since = now
self._close_rate_event(now)
elif bssid == prev_bssid and prev.get("channel") and link.get("channel") \
and prev["channel"] != link["channel"]:
self._instant(now, "channel-change",
"Channel changed %s \u2192 %s on the same AP"
% (prev["channel"], link["channel"]))
tx = link.get("tx_mbps")
if tx is None or tx <= 0:
return
if (traffic_bps or 0) < self.TRAFFIC_FLOOR_BPS:
# Idle link: whatever bitrate power save negotiated is
# unobservable to the user. Freeze the tracker — and close an
# open drop event, since its duration would otherwise count
# idle time as suffering.
self._close_rate_event(now)
return
# A slowly decaying ceiling: the best rate seen lately, with a
# half-life of a few minutes so an old burst does not set the bar
# forever. Only rates seen under traffic feed it.
self.rate_ceiling = max(tx, self.rate_ceiling * 0.998)
if self.rate_ceiling < 100:
return # too slow a link for a drop to mean anything
if tx < self.rate_ceiling * self.LOW_FRACTION:
self.low_floor = tx if self.low_floor is None else min(self.low_floor, tx)
if self.low_since is None:
self.low_since = now
elif self.rate_event_id is None and now - self.low_since >= self.SUSTAIN_S:
self.rate_event_id = self.store.open_event(
int(self.low_since), "rate-drop", "warn", "local",
"Tx rate dropped to %d Mbps" % round(self.low_floor))
elif tx >= self.rate_ceiling * self.RECOVER_FRACTION:
self._close_rate_event(now)
def _close_rate_event(self, now):
if self.rate_event_id is not None:
detail = "Tx rate dropped to %d Mbps" % round(self.low_floor or 0)
self.store.close_event(self.rate_event_id, int(now), detail)
self.rate_event_id = None
self.low_since = None
self.low_floor = None
def _short_duration(seconds):
s = int(round(seconds))
if s < 60:
return "%d s" % max(1, s)
if s < 3600:
return "%d min" % (s // 60)
return "%d h %d min" % (s // 3600, (s % 3600) // 60)
class LegState(NamedTuple):
"""What a leg's probe stream says right now — see leg_state()."""
ok: bool # the newest sample is a reply
ts: float # timestamp of the newest sample
run_since: object # first lost sample of the trailing run; None when ok
lost: int # lost samples in that run; 0 when ok
def leg_state(samples, now: float):
"""Read a leg's recent samples into a LegState, or None when the stream
is empty or stale — the probe stopped talking, which is not an outage.
`samples` are (ts, rtt_or_None, ...) in time order: one Series, or the
MergedSeries of the seated instruments. The trailing run of losses is
the whole question — how long since anything answered, counted from the
first sample that did not. This replaces a rolling "did anything reply
in the last 3 s" window, whose width was silently added to every
threshold built on it.
"""
if not samples:
return None
ts = samples[-1][0]
if now - ts > LEG_STALE_S:
return None
lost, run_since = 0, None
for smp in reversed(samples):
if smp[1] is not None:
break
lost += 1
run_since = smp[0]
if lost == 0:
return LegState(True, ts, None, 0)
return LegState(False, ts, run_since, lost)
class LegWatch:
"""Outage state for one leg: how long its stream has been silent, when
that silence became an outage, and what a recovered run looked like."""
def __init__(self):
self.down_since = None
self.run_since = None # first lost sample of the current run
self.lost = 0 # lost samples seen in that run
self.blip = None # (from, to) of a run that just recovered
def sample(self, state: LegState, now: float):
"""Returns 'down' / 'up' / 'disruption' on a transition, else None.
`disruption` is a run that recovered before reaching the outage
threshold. It is reported at recovery rather than at onset because
that is the first moment its length is known, and its length is what
Reliability charges. Both thresholds are seconds of silence on the
stream, from the first lost sample to the first reply after it.
"""
if state.ok:
began, lost = self.run_since, self.lost
self.run_since, self.lost = None, 0
if self.down_since is not None:
self.down_since = None
return "up"
if (began is not None and lost >= MIN_LOST_SAMPLES
and state.ts - began >= DISRUPTION_AFTER_S):
self.blip = (began, state.ts)
return "disruption"
return None
# The run began when the packets started going missing, not when we
# noticed: keep the earliest start seen for this run.
if self.run_since is None or state.run_since < self.run_since:
self.run_since = state.run_since
self.lost = max(self.lost, state.lost)
if (self.down_since is None and self.lost >= MIN_LOST_SAMPLES
and now - self.run_since >= OUTAGE_AFTER_S):
self.down_since = self.run_since
return "down"
return None
class LegArbiter:
"""What a leg going quiet MEANS, given whether anything beyond it answered.
A run of losses on one leg is not the same event as the internet being
gone. The evidence that separates them is already in hand: if anything
past this leg is still answering, packets are crossing it, so it is not
unreachable — it is merely refusing our probes. That opens a warn-toned
"quiet" event (`QUIET_KIND`) instead of an outage: logged, excluded from
outage_stats, no notification, the bar stays calm. Only every path
falling silent opens a real `outage`, which alarms once it has lasted
NOTIFY_AFTER_S.
Escalation is one-way. If the far side goes quiet too during a quiet
spell, the quiet event closes and a real outage opens, because an outage
that begins mid-spell must still alarm. Nothing walks back the other
way: flapping between verdicts would teach people to ignore both.
The two legs differ only in wording — which kind, which detail, which
notification — so those are class attributes and the mechanism is shared.
Until 0.2.16 this was two classes that had been copied and string-edited
apart, which is exactly how 0.2.4's gateway-quiet reached the store
correctly and the Events tab not at all: a fix to one copy that missed
the other.
"""
LEG = None # "wan" | "local"
QUIET_KIND = None # "icmp-quiet" | "gateway-quiet" (names predate the
QUIET_DETAIL = None # bench/arbiter; not worth a stored-value migration)
OUTAGE_DETAIL = None
ALARM = None # (summary, body) when a real outage is declared
RECOVERED = None # (summary, body) when a real outage recovers
def __init__(self, store, notify):
self.store = store
self.notify = notify
self.event_id = None
self.kind = None # QUIET_KIND | "outage" while down
self._notify_at = None # when the alarm becomes due
self._notified = False # whether it actually fired
@property
def real_outage(self) -> bool:
return self.kind == "outage"
def down(self, now, beyond_ok: bool, since=None):
# `since` is when the silence began; the row carries the onset, not
# the tick that crossed the threshold.
began = int(since if since is not None else now)
if beyond_ok:
self.kind = self.QUIET_KIND
self.event_id = self.store.open_event(
began, self.QUIET_KIND, "warn", self.LEG, self.QUIET_DETAIL)
return
self.kind = "outage"
self.event_id = self.store.open_event(
began, "outage", "critical", self.LEG, self.OUTAGE_DETAIL)
# Logged now, alarmed only if it lasts — see NOTIFY_AFTER_S.
self._notify_at = now + NOTIFY_AFTER_S
self._notified = False
def tick(self, now, beyond_ok: bool):
if self.kind == self.QUIET_KIND and not beyond_ok:
self.store.close_event(self.event_id, int(now))
self.down(now, False)
return
if (self.kind == "outage" and not self._notified
and self._notify_at is not None and now >= self._notify_at):
self._notified = True
self.notify(self.ALARM[0], self.ALARM[1], True)
def up(self, now):
if self.event_id is not None:
self.store.close_event(self.event_id, int(now))
# Only say it came back if we said it went away. A recovery
# notice with no matching alarm is a message about nothing.
if self.kind == "outage" and self._notified:
self.notify(self.RECOVERED[0], self.RECOVERED[1])
self.event_id = None
self.kind = None
self._notify_at = None
self._notified = False
class WanEventArbiter(LegArbiter):
"""The internet leg. `beyond_ok` here means some instrument still
answered: an ISP or middlebox that stops answering one probe while
others still flow used to be recorded — notified, charged to Reliability
— as an outage the user never experienced. See LegArbiter."""
LEG = "wan"
QUIET_KIND = "icmp-quiet"
QUIET_DETAIL = ("Scored probes went quiet; another instrument on "
"the same path kept answering")
OUTAGE_DETAIL = "router answers, nothing past it does"
ALARM = ("No internet",
"The router answers but nothing past it does — "
"the fault is on the ISP side.")
RECOVERED = ("Internet recovered", "Replies from the internet again.")
class LocalEventArbiter(LegArbiter):
"""The local leg. `beyond_ok` here means something past the gateway
answered, so packets are crossing it and it is not unreachable — just
refusing pings, which hotel and captive networks routinely do. This is
0.1.17's wan arbitration applied to the leg it had never covered. See
LegArbiter."""
LEG = "local"
QUIET_KIND = "gateway-quiet"
QUIET_DETAIL = ("Router stopped answering pings; traffic through it "
"kept working")
OUTAGE_DETAIL = "router unreachable"
ALARM = ("Router unreachable",
"Nothing on the local network is answering.")
RECOVERED = ("Local network recovered",
"The router is answering again.")
class CaptiveWatch:
"""Are we behind a sign-in page rather than on the internet?
A probe reply proves a packet came back; it does not prove what sent it.
So this asks for two things at once and only claims interception when it
has both: something IS answering our probes, and the reachability check
cannot prove the real internet answered. Packets going somewhere, but
not to the internet, is what a captive portal looks like from here.
Neither half is enough alone. Probes answering with no reachability check
is the state we were in before, and it read as a healthy internet. A
failed check with nothing answering is simply no internet, which the
wan arbiter already handles — calling that "captive" would put a sign-in
prompt in front of a user whose line is down.
Confirmation takes two consecutive checks, because one failed fetch is a
failed fetch. The decision itself is pure so it can be argued with and
tested; only the fetching is not.
The fetch runs OFF the loop and on suspicion, not on a clock. `tick()`
starts a check when one is due and collects the result on a later tick;
it never waits. Until 0.2.13 it ran the curl inline every 30 s for as
long as anything answered — i.e. always — which stalled the bar for up
to 8 s on exactly the slow networks this exists for, and cost 2,880
fetches a day where the WAN-address fetch it duplicated cost 24. Now: a
check on every new network and whenever the internet comes back; every
30 s only while the last answer was not proof of the internet; hourly
once it was. That hourly check is also where the WAN address comes
from, so the separate fetch is gone.
"""
# While a sign-in page is suspected, re-check soon.
CHECK_EVERY_S = 30.0
# Once the real internet has answered, once an hour keeps the address
# fresh and would notice a portal that appears mid-session.
RECHECK_OPEN_S = 3600.0
CONFIRM_AFTER = 2
def __init__(self, check, spawn=None):
self._check = check # injected: () -> {"verdict", "proof"}
# How a check is run. Off the loop by default; tests pass a
# synchronous spawn so the result lands within the same tick.
self._spawn = spawn or self._in_thread
self._lock = threading.Lock()
self._pending = None # (generation, result) awaiting a tick
self._inflight = False
self._gen = 0 # bumped by request(): stale results drop
self.verdict = "unknown"
self.proof = None
self.checked_ts = None
self.strikes = 0
self._next = 0.0
self._confirmed = False
@staticmethod
def _in_thread(fn):
threading.Thread(target=fn, name="reach", daemon=True).start()
@staticmethod
def captive(verdict: str, probes_answering: bool, strikes: int) -> bool:
"""The whole claim, in one place: replies but no proof of internet."""
return (probes_answering
and verdict in ("intercepted", "silent")
and strikes >= CaptiveWatch.CONFIRM_AFTER)
@property
def confirmed(self) -> bool:
return self._confirmed
def request(self):
"""A new network, or the internet just came back: start over.
The old verdict described a different situation, and a check still
in flight for it must not be mistaken for an answer about this one.
"""
self._gen += 1
self._inflight = False
self._next = 0.0
self.verdict = "unknown"
self.proof = None
self.strikes = 0
self._confirmed = False
def tick(self, now: float, probes_answering: bool):
self._collect(now)
if not probes_answering:
# Nothing is answering at all: not our verdict to make. Drop the
# suspicion rather than carrying it into a real outage.
self.strikes = 0
self._confirmed = False
return
if self._inflight or now < self._next:
self._confirmed = self.captive(
self.verdict, probes_answering, self.strikes)
return
self._inflight = True
gen = self._gen
self._spawn(lambda: self._run(gen))
# A synchronous spawn has already delivered; a thread delivers to a
# later tick.
self._collect(now)
def _run(self, gen: int):
try:
result = self._check() or {"verdict": "silent", "proof": None}
except Exception:
result = {"verdict": "silent", "proof": None}
with self._lock:
self._pending = (gen, result)
self._inflight = False
def _collect(self, now: float):
with self._lock:
pending, self._pending = self._pending, None
if pending is None:
return
gen, result = pending
if gen != self._gen:
return # answered a question we stopped asking
self.verdict = result.get("verdict") or "silent"
self.proof = result.get("proof")
if self.verdict == "open":
self.strikes = 0
else:
self.strikes += 1
self.checked_ts = round(now)
self._next = now + (self.RECHECK_OPEN_S if self.verdict == "open"
else self.CHECK_EVERY_S)
self._confirmed = self.captive(self.verdict, True, self.strikes)
def snapshot(self) -> dict:
if self.verdict == "unknown":
return None
return {"verdict": self.verdict, "captive": self._confirmed,
"checked_ts": self.checked_ts}
class LinkCollector(threading.Thread):
"""Reads the local end — route, interface, Wi-Fi link and station — on
its own thread, and keeps the latest snapshot for the loop to read.
net.snapshot() is three subprocesses (ip, iw link, iw station): about
10 ms when all is well, and up to their 2 s timeouts EACH when it is
not — a laptop coming out of suspend fails all of them at once, and
the loop that owns the outage watch used to wait for every one, twice
a second. Now it reads a dict. The first snapshot is taken
synchronously in start(), so there is always one to read.
A snapshot is never "too old" to hand out: if this thread is stuck
behind a wedged `iw`, the loop keeps the last known link, which is
what a timed-out read produced before as well — only now the loop
does not stop for it.
"""
INTERVAL_S = 0.5
def __init__(self, anchor_fn, snapshot_fn=None, interval_s=INTERVAL_S):
super().__init__(name="link", daemon=True)
self._anchor_fn = anchor_fn
self._snapshot = snapshot_fn or net.snapshot
self.interval = interval_s
self._stop = threading.Event()
self._lock = threading.Lock()
self._latest = {}
self.taken_at = 0.0
def _take(self):
try:
snap = self._snapshot(self._anchor_fn())
except Exception: # noqa: BLE001 — a collector never raises into the daemon
return
with self._lock:
self._latest = snap if isinstance(snap, dict) else {}
self.taken_at = time.time()
def start(self):
self._take()
super().start()
def run(self):
while not self._stop.wait(self.interval):
self._take()
def stop(self):
self._stop.set()
@property
def latest(self) -> dict:
"""A copy: callers annotate it (retry_pct) and must not share."""
with self._lock:
snap = dict(self._latest)
if isinstance(snap.get("station"), dict):
snap["station"] = dict(snap["station"])
return snap
class Daemon:
def __init__(self):
self.state_dir = ensure_state_dir()
ensure_runtime_dir()
self.config = Config(self.state_dir)
self.config.refresh()
self.store = Store(db_path())
self.local = Series()
# The internet leg is measured by a bench of instruments — see
# instruments.py. Each instrument feeds its own Series; the scored
# series (self.total) is a merged view over whichever two hold the
# seats, so every consumer downstream keeps reading one "internet
# leg". The anchor's ICMP series keeps a name of its own too: the
# minute rows record what ICMP alone would have scored (`lag_icmp`).
self.bench = Bench(self._instrument_pool())
self._instrument_series = {}
self._instrument_probes = {}
self._new_instrument_series()
self.total = MergedSeries(self._active_series)
self._last_bench_eval = 0.0
self.probes = []
self._local_probe = None
# (anchor, interval) the running probes were built from, so a
# settings change is noticed on the next tick — see
# restart_probes_if_settings_changed.
self._probe_settings = None
self.running = True
self.route = {}
# Sliding window of (t, rx, tx) counter samples. Rates are computed
# across the whole window, not tick-to-tick — a half-second sample is
# instantaneous chatter, and displaying it twice a second reads as
# flicker rather than as a number.
self.counter_samples = []
self.rates = (None, None) # bytes/sec
# Whether the link is carrying real traffic right now. Probes read
# this as each sample lands, which is what separates idle latency
# from latency under load — the gap between them IS bufferbloat.
# Same floor the link-event logic uses, for the same reason: below
# it, Wi-Fi power save makes the link look busy when nobody is.
self.link_loaded = False
# 5-second aux samples riding along in recent.json: throughput and
# signal, so the panel's charts have history the moment they open.
self.aux_ring = deque(maxlen=400)
# Elapsed-time gate for the 5 s flush. It used to be `int(now) % 5
# == 0`, which is true on BOTH half-second ticks of a qualifying
# second, so the ring filled twice as fast and the Wi-Fi and
# throughput charts held ~17 min under a 30-minute label.
self.last_recent_flush = 0.0
self.last_signal = None
self.watch_local = LegWatch()
self.watch_wan = LegWatch()
self.wan_events = WanEventArbiter(self.store, self.notify)
self.local_events = LocalEventArbiter(self.store, self.notify)
self.captive = CaptiveWatch(net.reachability)
# Is this connection someone's phone sharing its data? Recomputed
# whenever the route changes, which is the only thing that can
# change the answer.
self.metered = None
# The address this connection appears from — live.json only, never
# recent.json or history: shown, not archived.
self.wan_ip = None
# checked_ts of the reachability result the address came from, so
# one proof is not re-adopted every tick.
self._wan_ip_at = 0
# Who ended each Wi-Fi association — read from nl80211 via `iw
# event`, unprivileged. Without it the link log still works; it
# just cannot tell a kick from a roam.
self.nl_events = linkevents.NlEvents()
self.link_watch = LinkWatch(self.store, self.nl_events)
# Notify-only: asks origin whether this checkout is behind and
# never touches it. Off when the user turns updateCheck off.
self.update_watch = UpdateWatch(enabled=bool(self.config["updateCheck"]))
self.app_traffic = apps.AppTraffic()
# The local end, read off the loop — see LinkCollector.
self.link = LinkCollector(lambda: self.config["internetAnchor"])
self.last_apps_poll = 0.0
self.last_content_test = 0.0
# Set while a due content check is waiting for the link to settle.
self._check_waiting_since = None
self.last_minute_flush = 0.0
self.last_rollup = 0.0
self.peak_requested = threading.Event()
self.peak_running = False
self.content_running = False
self._content_retry_used = False
self._lock_fh = None
# ------------------------------------------------------------- lifecycle
def acquire_lock(self) -> bool:
"""One daemon per user. The shell service and a systemd unit can both
try to start us; whoever loses the lock just exits quietly.
The lock file sits at a predictable path, so it is opened without
truncation and without following symlinks — a planted symlink must
fail the open, never redirect a truncation somewhere else — and it
is only ever truncated after this process holds the flock.
"""
try:
fd = os.open(lock_path(),
os.O_RDWR | os.O_CREAT | os.O_NOFOLLOW | os.O_CLOEXEC,
0o600)
except OSError:
return False
try:
if not stat_module.S_ISREG(os.fstat(fd).st_mode):
os.close(fd)
return False
# The creation mode only applies to new files; a lock file left
# by an older version keeps its old permissions until this.
os.fchmod(fd, 0o600)
fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB)
except OSError:
os.close(fd)
return False
os.ftruncate(fd, 0)
os.write(fd, str(os.getpid()).encode())
self._lock_fh = os.fdopen(fd, "r+")
return True
def _instrument_pool(self):
anchor = self.config["internetAnchor"]
return [("icmp-anchor", "icmp", anchor),
("tcp-anchor", "tcp", "%s:443" % anchor),
("tcp-cf", "tcp", CF_EDGE_HOST + ":443"),
("tcp-google", "tcp", DIVERSITY_HOST + ":443")]
def _new_instrument_series(self):
self._instrument_series = {
key: Series() for key, _, _ in self._instrument_pool()}
# The name the rest of the daemon has always read.
self.icmp_anchor = self._instrument_series["icmp-anchor"]
def _active_series(self):
return [self._instrument_series[i.key] for i in self.bench.actives()
if i.key in self._instrument_series]
def _instrument_stats(self, window_s: float = Bench.WINDOW_S):
return {k: Series.stats(v.since(window_s))
for k, v in self._instrument_series.items()}
def _apply_seats(self, changes):
interval = int(self.config["probeIntervalMs"]) / 1000.0
for key, active in changes:
probe = self._instrument_probes.get(key)
if not probe:
continue
kind = self.bench.instruments[key].kind
base = interval if kind == "icmp" else TCP_PROBE_INTERVAL_S
probe.set_interval(base if active else base * STANDBY_FACTOR)
def start_probes(self):
anchor = self.config["internetAnchor"]
self.route = net.route_to(anchor)
# Answer this before the first content check can be due, not only
# when the route later changes — a daemon started on a hotspot must
# not spend a check to find out it is on one.
self.refresh_metered()
interval = int(self.config["probeIntervalMs"])
self._probe_settings = (anchor, interval)
gw = self.route.get("gateway", "")
self._local_probe = None
if gw:
p = PingProbe(gw, self.local, interval, "local",
loaded_fn=lambda: self.link_loaded)
p.start()
self.probes.append(p)
self._local_probe = p
for key, kind, target in self._instrument_pool():
series = self._instrument_series[key]
host = target.rsplit(":", 1)[0] if kind == "tcp" else target
if kind == "icmp":
p = PingProbe(host, series, interval, key,
loaded_fn=lambda: self.link_loaded)
base = interval / 1000.0
else:
p = TcpProbe(host, series, TCP_PROBE_INTERVAL_S, key,
loaded_fn=lambda: self.link_loaded,
port=TCP_PROBE_PORT)
base = TCP_PROBE_INTERVAL_S
self.bench.instruments[key].target = target
if not self.bench.instruments[key].active:
p.set_interval(base * STANDBY_FACTOR)
p.start()
self.probes.append(p)
self._instrument_probes[key] = p
def restart_probes_if_route_changed(self):
"""New default route (roamed networks, docked, VPN up) — new targets."""
anchor = self.config["internetAnchor"]
fresh = net.route_to(anchor)
if not fresh.get("gateway"):
# No route at all is an outage, not a different network, and
# resetting on it threw away the one window a user wants
# afterwards — the run-up to the drop. It also fired twice per
# disconnect, once on the way down and once on the way back.
# Nothing new can contaminate the distributions while there is
# no network, so keep them, and keep the probes running: their
# losses are what the outage watch is reading.
return
if fresh.get("gateway") == self.route.get("gateway") and \
fresh.get("iface") == self.route.get("iface"):
return
self._rebuild_probes(fresh)
def restart_probes_if_settings_changed(self):
"""The anchor or the probe interval changed under us.
Both used to be read only when probes were built, so a change sat
unapplied until the next network change or restart while the panel
said settings apply live. An interval change is applied in place —
the distributions stay valid, only the cadence moves. A new anchor
rebuilds the probes: half the instruments now point somewhere else,
and their old samples describe a host nobody is measuring any more.
"""
if self._probe_settings is None:
return # probes not started yet
anchor = self.config["internetAnchor"]
interval = int(self.config["probeIntervalMs"])
if (anchor, interval) == self._probe_settings:
return
if anchor != self._probe_settings[0]:
self._rebuild_probes(net.route_to(anchor))
return
self._probe_settings = (anchor, interval)
if self._local_probe is not None:
self._local_probe.set_interval(interval / 1000.0)
self._apply_seats([(k, i.active)
for k, i in self.bench.instruments.items()])
def _rebuild_probes(self, route: dict):
"""Stop every probe and start over against `route`: fresh network,
fresh distributions. The bench keeps its seats — continuity until
the new windows hold enough samples to argue about."""
for p in self.probes:
p.stop()
self.probes.clear()
self._instrument_probes = {}
self.local = Series()
self._new_instrument_series()
self.route = route
self.counter_samples = []
# A new route means a new apparent address; drop the stale one
# rather than display it wrong until the next check — and ask for
# that check now: a new network is exactly where a sign-in page is
# likeliest.
self.wan_ip = None
self._wan_ip_at = 0
self.captive.request()
self.start_probes()
def stop(self, *_):
self.running = False
# ------------------------------------------------------------ measuring
def load_floor_bps(self, now: float) -> float:
"""Bytes per second above which this line counts as busy.
A tenth of what the line has been measured to carry, floored. Read
from the same baseline the Speed score uses — this network's own p90
download — and cached for a minute, because it moves at content-check
cadence and this is asked twice a second.
"""
cache = getattr(self, "_load_floor_cache", None)
if not cache or now - cache[0] > 60:
snap = self.link.latest if self.link else {}
network = snap.get("ssid") or snap.get("name") or ""
try:
baseline = self.store.baseline_speed(network=network, now=now,
fallback=False)
except Exception:
baseline = None
floor = LOAD_FLOOR_BPS
if baseline:
floor = max(floor, baseline * 1e6 / 8 * LOAD_FRACTION_OF_LINE)
cache = (now, floor)
self._load_floor_cache = cache
return cache[1]
def throughput(self, now: float, iface: str):
c = net.counters(iface)
if c is None:
self.counter_samples = []
self.rates = (None, None)
self.link_loaded = False
return
window = max(1, min(30, int(self.config["throughputWindowS"])))
self.counter_samples.append((now, c[0], c[1]))
cutoff = now - window - 0.25
while len(self.counter_samples) > 2 and self.counter_samples[0][0] < cutoff:
self.counter_samples.pop(0)
# Hard cap independent of config: the window can never retain more
# than a minute of half-second samples, whatever the file says.
if len(self.counter_samples) > 128:
del self.counter_samples[:len(self.counter_samples) - 128]
if len(self.counter_samples) >= 2:
t0, rx0, tx0 = self.counter_samples[0]
t1, rx1, tx1 = self.counter_samples[-1]
dt = t1 - t0
if dt > 0 and rx1 >= rx0 and tx1 >= tx0:
self.rates = ((rx1 - rx0) / dt, (tx1 - tx0) / dt)
self.link_loaded = ((self.rates[0] or 0.0) + (self.rates[1] or 0.0)
>= self.load_floor_bps(now))
else:
# Counter reset (interface bounced) — start the window over.
self.counter_samples = [self.counter_samples[-1]]
def notify(self, summary: str, body: str, urgent: bool = False):
if not self.config["notifyOutage"]:
return
cmd = None
if shutil.which("omarchy-notification-send"):
cmd = ["omarchy-notification-send", summary, body]
elif shutil.which("notify-send"):
cmd = ["notify-send", "-a", "Nexthop"]
if urgent:
cmd += ["-u", "critical"]
cmd += [summary, body]
if cmd:
try:
subprocess.Popen(cmd, stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL)
except OSError:
pass
def watch_outages(self, now: float):
"""Outage logic on each leg's probe stream (leg_state / LegWatch).
The wan watch counts silence only when the local leg answered, or
when the gateway is merely quiet (LocalEventArbiter): if the router
is confirmed unreachable, the internet probes' losses say nothing
about the ISP.
Arbitration is judged on the run itself — did anything beyond the
leg answer DURING the silence — not on a trailing window: a reply
from just before a 1.5 s blip must not vouch for the blip.
"""
local = leg_state(self.local.since(LEG_STREAM_WINDOW_S), now)
total = leg_state(self.total.since(LEG_STREAM_WINDOW_S), now)
local_ok = None if local is None else local.ok
if local is not None:
move = self.watch_local.sample(local, now)
if move == "down":
# Is anything past the gateway answering? If so the gateway
# is forwarding and merely refuses pings — see
# LocalEventArbiter.
since = self.watch_local.down_since
self.local_events.down(
now, self._any_instrument_replied_between(since, now), since)
elif move == "up":
self.local_events.up(now)
elif move == "disruption":
self.record_disruption("local", self.watch_local)
elif self.watch_local.down_since is not None:
self.local_events.tick(
now, self._any_instrument_alive(OUTAGE_AFTER_S))
# The wan watch normally ignores any window where the local leg lost
# packets: if the router is unreachable, the internet probe's losses
# say nothing about the ISP. But a gateway that merely refuses pings
# is "lost" forever, and skipping the wan watch on such a link would
# mean never noticing a real internet outage there. So the skip
# applies to a confirmed local outage, not to a quiet gateway.
gateway_quiet = (self.watch_local.down_since is not None
and not self.local_events.real_outage)
if total is not None and (local_ok is not False or gateway_quiet):
move = self.watch_wan.sample(total, now)
if move == "down":
# Did another instrument keep answering while the seated
# pair fell silent? Samples are stamped at send time, so a
# handshake that merely straddled the moment the line died
# cannot vouch for the window after it.
since = self.watch_wan.down_since
self.wan_events.down(
now, self._any_instrument_replied_between(since, now), since)
elif move == "up":
self.wan_events.up(now)
# The internet is back — or something answering for it is.
# Ask the reachability check now rather than wait its hour.
self.captive.request()
elif move == "disruption":
self.record_disruption("wan", self.watch_wan)
elif self.watch_wan.down_since is not None:
self.wan_events.tick(
now, self._any_instrument_alive(OUTAGE_AFTER_S))
def record_disruption(self, leg: str, watch, beyond_ok=None):
"""A run that recovered before it became an outage.
Arbitrated exactly like an outage: if something past this leg kept
answering DURING the run, the leg did not interrupt anything — a
gateway dropping three pings while traffic crosses it is not a
disruption, and logging it would fill the log with noise the user
never felt. `beyond_ok` is computed from the run's own interval
unless a caller supplies it.
"""
if not watch.blip:
return
began, ended = watch.blip
watch.blip = None
if beyond_ok is None:
beyond_ok = self._any_instrument_replied_between(began, ended)
if beyond_ok:
return
# Stored closed, with its duration, because Reliability charges
# interruptions in time. Integer seconds are what the table holds,
# so guarantee a non-zero span: outage_stats drops any row whose
# end is not after its start, and a blip that vanished from the
# score would be worse than one rounded up by a second.
eid = self.store.open_event(
int(began), "disruption", "warn", leg,
"brief interruption, recovered on its own")
self.store.close_event(eid, max(int(ended), int(began) + 1))
def _any_instrument_alive(self, window_s: float) -> bool:
"""Some instrument — seated or benched — heard the internet this
recently. While a leg is down, the arbiters read this each tick to
decide whether a quiet spell has become an outage."""
for series in self._instrument_series.values():
if any(s[1] is not None for s in series.since(window_s)):
return True
return False
def _any_instrument_replied_between(self, a: float, b: float) -> bool:
"""Did some instrument hear the internet during [a, b]? Judged on the
interval itself, so a reply from just before a run of silence cannot
vouch for it — which a trailing window longer than the run would."""
span = max(0.0, time.time() - a) + 1.0
for series in self._instrument_series.values():
for smp in series.since(span):
if smp[1] is not None and a <= smp[0] <= b:
return True
return False
def adopt_wan_ip(self):
"""The address this connection appears from, taken from the
reachability check's proof — the same fetch, so no second request
and no second host learns the address. A check that could not prove
the internet leaves the last answer standing: the route is what
invalidates it, and the route path clears it."""
proof, checked = self.captive.proof, self.captive.checked_ts
if not proof or not checked or checked == self._wan_ip_at:
return
self._wan_ip_at = checked
self.wan_ip = dict(proof, checked_ts=checked)
def refresh_metered(self):
"""Tethering, or a connection the user has marked metered.
Two independent signals, and neither is a guess: the gateway sitting
in a range only tethering hands out, and NetworkManager being told
so explicitly. Published so the panel can name the phone instead of
drawing a router, and consulted before anything spends data.
"""
gw = self.route.get("gateway", "")
iface = self.route.get("iface", "")
tether = net.tether_from_gateway(gw)
explicit = net.nm_metered(iface)
if not tether and not explicit:
self.metered = None
return
self.metered = {
"tethered": tether is not None,
"kind": tether["kind"] if tether else "declared",
# What to call the middle node of the path.
"label": tether["label"] if tether else "Metered link",
"explicit": explicit,
}
def maybe_content_test(self, now: float):
if not self.config["contentSpeed"]:
return
interval = max(15, int(self.config["contentSpeedIntervalMin"])) * 60
boost_at = getattr(self, "_content_boost_at", None)
due = (now - self.last_content_test >= interval) or \
(boost_at is not None and now >= boost_at)
if not due:
return
# Skip while down — a failed transfer during an outage is not a
# speed measurement, and skip while another test owns the line.
if self.watch_wan.down_since or self.watch_local.down_since \
or not self.tests_idle():
return
# Someone's phone is paying for this. The check is ~14 MB and runs
# hourly, which is around 336 MB a day of a data plan the user did
# not offer. Skipped rather than shrunk: a smaller sample would
# still cost money and would measure worse. Speed then scores None
# on this network, and the index already skips a component it does
# not have rather than inventing one.
if self.metered and self.config["meteredCare"]:
return
# Wait for a link worth measuring, but not forever — see check_ready.
if not check_ready(now, self.link_watch.assoc_since,
self.link_watch.low_since,
self._check_waiting_since):
if self._check_waiting_since is None:
self._check_waiting_since = now
return
self._check_waiting_since = None
self._content_boost_at = None
self.last_content_test = now
snap = self.link.latest
network = snap.get("ssid") or snap.get("name") or ""
self.content_running = True
down_hint, up_hint = self._content_hint(network)
def run():
try:
r = speedtest.content_test(down_hint_mbps=down_hint,
up_hint_mbps=up_hint)
after = self.link.latest
if (after.get("ssid") or after.get("name") or "") != network:
# The network changed under the transfer, so the sample
# belongs to neither. The change has already scheduled
# a fresh check of its own.
return
if r["ok"]:
self.store.put_test(int(r["started"]), "content", r["engine"],
down_mbps=r["down_mbps"], up_mbps=r["up_mbps"],
bytes=r["bytes"], ok=True, network=network)
# A fresh result should reprice the baseline promptly.
self._baseline_cache = None
self._content_retry_used = False
elif not self._content_retry_used:
self._content_retry_used = True
self._content_boost_at = time.time() + CONTENT_RETRY_S
finally:
self.content_running = False
threading.Thread(target=run, daemon=True, name="content-test").start()
def _content_hint(self, network: str):
"""What this network has shown, so the next check can size itself.
The best of the recent checks rather than the last. Sizing from a
reading that happened to come in low would make the next transfer
shorter, which reads lower again — a ratchet the floor alone would
stop only at the bottom. The best recent reading is also the honest
answer to "what can this line do", which is the question the size is
being chosen against.
None for a network with no history: the first check sends the cap and
produces the hint that every check after it uses.
"""
downs, ups = [], []
for t in self.store.tests(limit=8, kind="content"):
if not t["ok"] or (t["network"] or "") != network:
continue
if t["down_mbps"]:
downs.append(t["down_mbps"])
if t["up_mbps"]:
ups.append(t["up_mbps"])
return (max(downs) if downs else None), (max(ups) if ups else None)
def tests_idle(self) -> bool:
"""May a bandwidth test start? One at a time.
Two saturating transfers invalidate each other's rate and share the
probes' loaded-latency window. The scheduled check always yielded
to a running peak; until 0.2.21 a peak did not yield to a running
check, because nothing recorded that one was running.
"""
return not (self.peak_running or self.content_running)
def run_peak_test(self):
"""On demand, in its own thread; loaded latency comes from the probes."""
if not self.tests_idle():
return
self.peak_running = True
snap = self.link.latest
network = snap.get("ssid") or snap.get("name") or ""
def run():
try:
idle = score.lag_ms(self.total.stats(60))
started = time.time()
r = speedtest.peak_test(self.config["peakEngine"])
loaded_st = merged_stats([[s for s in lst if s[0] >= started]
for lst in self.total.each()])
loaded = (round(loaded_st["p50"], 1)
if loaded_st.get("p50") is not None else None)
if r["ok"]:
self.store.put_test(
int(r["started"]), "peak", r["engine"],
down_mbps=r.get("down_mbps"), up_mbps=r.get("up_mbps"),
ping_idle=r.get("ping_idle") or idle, ping_loaded=loaded,
jitter=r.get("jitter"), bytes=r.get("bytes"),
server=r.get("server"), ok=True,
detail=r.get("url", ""), network=network)
else:
self.store.put_test(int(r["started"]), "peak", r["engine"],
ok=False, network=network)
finally:
self.peak_running = False
threading.Thread(target=run, daemon=True, name="peak-test").start()
# -------------------------------------------------------------- writing
def speed_score(self, now: float, network: str):
"""(score, ctx) for the Speed component.
Plan configured -> scored against it. Otherwise the absolute
experience curve, with a degradation penalty against this network's
own recent p90. The baseline is cached for a minute — it moves at
content-check cadence, not at probe cadence.
"""
tests = [t for t in self.store.tests(limit=12, kind="content")
if t["ok"] and t["down_mbps"] is not None]
plan_d = self.config["planDownMbps"]
plan_u = self.config["planUpMbps"]
if plan_d:
if not tests:
return None, {"basis": "plan", "plan_down": plan_d,
"last_down": None, "last_up": None}
last = tests[0]
spd = score.speed(last["down_mbps"], last["up_mbps"],
plan_d, plan_u or 0)
return spd, {"basis": "plan", "plan_down": plan_d,
"plan_up": plan_u,
"last_down": last["down_mbps"],
"last_up": last["up_mbps"]}
# Checks describe the network they ran on. A result from another
# network says nothing about this one, so on a network with no
# checks yet the component is honestly unknown (and the changed
# network has already scheduled a prompt check).
mine = [t for t in tests
if (t.get("network") or "") == network] if network else tests
if not mine:
return None, {"basis": "auto", "baseline_down": None,
"last_down": None, "last_up": None,
"pending": True}
# Median of the last few checks here, so one bad sample — a check
# that ran mid-roam or during someone's upload — cannot pin the
# score until the next hourly run.
recent = mine[:3]
downs = sorted(t["down_mbps"] for t in recent)
down = downs[len(downs) // 2]
ups = sorted(t["up_mbps"] for t in recent
if t["up_mbps"] is not None)
up = ups[len(ups) // 2] if ups else None
cache = getattr(self, "_baseline_cache", None)
if not cache or now - cache[0] > 60 or cache[2] != network:
baseline = self.store.baseline_speed(network=network, now=now,
fallback=False)
cache = (now, baseline, network)
self._baseline_cache = cache
baseline = cache[1]
spd = score.speed(down, up, baseline_down=baseline)
# A saturating test of the same line, run recently and by hand, is
# better evidence of what the line can do than a 12 MB sample. It
# still does not become the score — a manual test must not flatter
# it — but it can withdraw a figure it contradicts.
peak_down = None
for t in self.store.tests(limit=6, kind="peak"):
if not t["ok"] or t["down_mbps"] is None:
continue
if network and (t.get("network") or "") != network:
continue
if now - t["ts"] > PEAK_FRESH_S:
break
peak_down = t["down_mbps"]
break
scored = score.speed_scored(down, len(recent), peak_down)
return spd, {"basis": "auto", "baseline_down": baseline,
"last_down": down, "last_up": up,
"samples": len(recent), "scored": scored,
"peak_down": peak_down}
def bufferbloat(self, window_s: float = 300.0) -> dict:
"""Lag while the link was idle vs while it was carrying traffic.
The gap between them is bufferbloat, and it is the failure a plain
latency number misses entirely: a line can answer in 15 ms at rest,
sit at 300 ms whenever anyone downloads anything, and still look
excellent on every idle measurement anyone takes of it.
Both figures come from the same probe stream — no extra traffic is
generated to produce them. That is the whole point of tagging each
sample as it lands: the user's own usage supplies the load.
"""
# Split each seated instrument's own stream by load, then merge the
# idle halves and the loaded halves with the instruments counting
# equally — see merged_stats for why the pooled stream may not be
# fed to Series.stats.
lists = self.total.each(window_s)
splits = [Series.split_by_load(lst) for lst in lists]
idle_st = merged_stats([sp[0] for sp in splits])
loaded_st = merged_stats([sp[1] for sp in splits])
n_idle, n_loaded = idle_st["count"], loaded_st["count"]
idle_lag = score.lag_ms(idle_st) if n_idle else None
loaded_lag = score.lag_ms(loaded_st) if n_loaded else None
# A handful of samples on either side produces noise, not a ratio —
# observed live, a five-sample loaded window read as 0.59, i.e. the
# link answering *faster* under load. Both sides need enough
# samples before the comparison means anything.
inflation = None
if (idle_lag and loaded_lag and idle_lag > 0
and n_idle >= MIN_LOAD_SPLIT_SAMPLES
and n_loaded >= MIN_LOAD_SPLIT_SAMPLES):
ratio = loaded_lag / idle_lag
if ratio >= MIN_PLAUSIBLE_INFLATION:
# Clamped at 1: a ratio a hair under it means the two are
# indistinguishable, not that load made the link quicker.
inflation = round(max(1.0, ratio), 2)
# Percentiles over the loaded samples ALONE. The headline stats span
# a fixed 30 s window, so a ten-second burst is averaged with twenty
# seconds of quiet and reads far milder than it was: measured against
# another tool on the same event, 107 ms against its 246. Scoping the
# percentile to the samples that were actually taken under load is
# the same idea as their per-phase percentile, using the tagging
# 0.1.11 already put on every probe.
return {"idle": idle_lag, "loaded": loaded_lag,
"inflation": inflation, "loaded_samples": n_loaded,
"idle_samples": n_idle,
"loaded_p50": loaded_st.get("p50"),
"loaded_p95": loaded_st.get("p95"),
"idle_p50": idle_st.get("p50"),
# How fast the queue emptied once traffic stopped. Depth is
# what everyone reports; duration is what a user feels after
# the download finishes.
"drain": self._drain(lists, splits, self._active_keys())}
def _active_keys(self):
"""Seated instrument keys, in the order `_active_series` yields them."""
return [i.key for i in self.bench.actives()
if i.key in self._instrument_series]
@staticmethod
def _drain(lists, splits, keys=None) -> dict:
"""drain_after_load per instrument, each against its own idle
floor, the slowest one reported. Two instruments with different
base round trips cannot share a baseline: measured against the
lower one's floor the higher one never settles, and the lower one
settles the moment its first post-load sample lands. The queue
they drained through is the same, so the pessimistic view is the
honest one.
That last sentence is under review and the numbers beside `ms` are
why. Each instrument's value is the gap to its next observation, so
it is an UPPER BOUND floored at that instrument's own cadence —
which means taking the largest reliably selects whichever instrument
looks least often, and publishes it as the line being slow to drain.
`min_ms` is the tightest bound the same window offers and `src` names
the instrument the published value came from. Both are recorded per
minute so the choice can be settled from stored history rather than
argued from first principles, which is how it has been argued so far.
"""
out = {"ms": None, "settled": None, "min_ms": None, "src": None}
keys = keys or []
for i, (samples, (idle, _)) in enumerate(zip(lists, splits)):
base = Series.stats(idle).get("p50") if idle else None
d = score.drain_after_load(samples, base)
ms = d.get("ms")
if ms is None:
continue
if out["ms"] is None or ms > out["ms"]:
out["ms"] = ms
out["settled"] = d.get("settled")
out["src"] = keys[i] if i < len(keys) else None
if out["min_ms"] is None or ms < out["min_ms"]:
out["min_ms"] = ms
return out
def compose_live(self, now: float) -> dict:
"""live.json, twice a second.
Some keys here have no panel reader — `reach`, `lag.inflation`,
`pressure.source`, `metered.kind`, `update.state`, `sockets.rejected`
and a few more. They stay on purpose: `nexthop live` is how every
measurement in this file was verified, and a payload trimmed to what
the panel draws would leave the operator blind. Trim deliberately,
never because a grep found no reader.
"""
ls = Series.stats(self.local.since(30))
ts = self.total.stats(30)
ws = score.wan_from(ts, ls)
lag = score.lag_ms(ts)
resp = score.responsiveness(lag) if ts["count"] else None
# Idle vs loaded over a longer window than the headline: bufferbloat
# only shows when the link has actually been used, and 30 s of an
# idle laptop would almost never contain a loaded sample. Reported,
# not yet scored — the number has to be trusted before it can move
# anyone's index.
bloat = self.bufferbloat(300.0)
out_frac, disruptions, disrupt_frac = self.store.outage_stats(24 * 3600, now)
rel = score.reliability(out_frac, disruptions,
disruption_fraction=disrupt_frac)
snap = self.link.latest
network = snap.get("ssid") or snap.get("name") or ""
self.last_signal = snap.get("signal_dbm")
prev = getattr(self, "_content_network", None)
if network and prev is not None and network != prev:
# New network: the hourly cadence would leave Speed unknown or
# stale for up to an hour here. Measure soon — after a settle
# delay, so a roam in progress is not sampled as the network's
# capability.
self._content_boost_at = now + 90
if network:
self._content_network = network
if snap.get("kind") == "wifi":
self.link_watch.sample(now, snap,
(self.rates[0] or 0) + (self.rates[1] or 0))
st = snap.get("station") or {}
if st.get("tx_packets"):
st["retry_pct"] = round(
100.0 * (st.get("tx_retries") or 0) / st["tx_packets"], 2)
spd, speed_ctx = self.speed_score(now, network)
band = score.lag_band(ts)
# An under-sampled or contradicted Speed figure is published and
# left out of the headline — see score.speed_scored.
idx = score.index(resp, rel, spd if speed_ctx.get("scored") else None)
state = "online"
if self.captive.confirmed:
# Outranks both leg verdicts because it explains them: on a
# portal the gateway often refuses pings and something answers
# for the anchor, so "router unreachable" and "internet fine"
# are both artefacts of the same interception.
state = "captive"
elif self.watch_local.down_since and self.local_events.real_outage:
# A silent gateway is not an unreachable one; the arbiter decides.
state = "local-down"
elif self.watch_wan.down_since and self.wan_events.real_outage:
# Pings alone cannot declare this; see WanEventArbiter. During
# an icmp-quiet spell the bar stays its ordinary colour — the
# user's internet is working, and the log holds the anomaly.
state = "wan-down"
elif idx is not None and idx < 70:
state = "degraded"
# An index computed while a leg is confirmed down scores a
# connection that is not there — see score.scored_now. Withheld,
# not lowered; the state is the headline and the panel draws "--".
headline = idx if score.scored_now(state) else None
return {
"v": 1,
"t": round(now, 3),
"state": state,
"index": headline,
"band": score.band(headline),
"scores": {"responsiveness": resp, "reliability": rel, "speed": spd},
"speed_ctx": speed_ctx,
# best/typical/worst all come from the same fold — see
# score.lag_band. `now` stays the scored p75-based figure.
"lag": {"now": lag,
"best": band.get("best"), "worst": band.get("worst"),
"typical": band.get("typical"),
"idle": bloat["idle"], "loaded": bloat["loaded"],
"inflation": bloat["inflation"],
"loaded_samples": bloat["loaded_samples"],
"idle_samples": bloat["idle_samples"],
"loaded_p50": bloat["loaded_p50"],
"loaded_p95": bloat["loaded_p95"],
"idle_p50": bloat["idle_p50"],
"drain_ms": bloat["drain"]["ms"],
"drain_settled": bloat["drain"]["settled"]},
"local": ls, "total": ts, "wan": ws,
"wan_ip": self.wan_ip,
# A phone sharing its data, or a link the user marked metered.
# `care` rides along so the panel can say whether anything is
# actually being held back, and is read fresh each time rather
# than frozen when the link was detected.
"metered": (dict(self.metered, care=bool(self.config["meteredCare"]))
if self.metered else None),
# Proof the real internet answered, or why it did not.
"reach": self.captive.snapshot(),
# What the user's own TCP connections are experiencing, straight
# from the kernel: their real traffic to their real destinations.
"sockets": self.app_traffic.latency,
# What the connection is doing right now, as opposed to lately.
# The index answers the second question and cannot answer the
# first — see score.pressure.
"pressure": score.pressure(
socket_queue_ms=(self.app_traffic.latency or {}).get("queue_p50"),
loaded_ms=bloat["loaded"], idle_ms=bloat["idle"]),
# Whether a newer version is published. A notice, not an
# action: nothing here updates anything.
"update": self.update_watch.snapshot(),
"instruments": self.bench.snapshot(now, self._instrument_stats()),
"rates": {"rx_bps": self.rates[0], "tx_bps": self.rates[1],
"rx_total": self.counter_samples[-1][1] if self.counter_samples else None,
"tx_total": self.counter_samples[-1][2] if self.counter_samples else None},
"link": snap,
"down_since": self.watch_local.down_since or self.watch_wan.down_since,
"peak_running": self.peak_running,
"content_running": self.content_running,
"pid": os.getpid(),
"pid_start": proc_start_ticks(os.getpid()),
"daemon_version": __version__,
}
def flush_recent(self, now: float):
"""recent.json: last 30 min at 5-second resolution, ~360 points."""
self.last_recent_flush = now
self.aux_ring.append((now, self.rates[0], self.rates[1],
self.last_signal))
points = []
bucket = 5.0
start = now - 1800
locs = self.local.all()
insts = self.total.each()
def fold(samples):
out = {}
for smp in samples:
t, r = smp[0], smp[1]
if t < start:
continue
b = int((t - start) / bucket)
out.setdefault(b, []).append(r)
return out
lb, tbs = fold(locs), [fold(lst) for lst in insts]
aux_b = {}
for at, rx, tx, sig in self.aux_ring:
if at >= start:
aux_b[int((at - start) / bucket)] = (rx, tx, sig)
for b in range(int(1800 / bucket)):
l = lb.get(b, [])
lr = [x for x in l if x is not None]
# Each seated instrument's bucket mean, then the mean of those:
# a seat that probes twice as often must not count twice. Loss
# stays a pooled count — on the chart it is a tick, present or
# not, and the readout's percentage is of everything sent.
means, n_t, lost_t = [], 0, 0
for t in (tb.get(b, []) for tb in tbs):
tr = [x for x in t if x is not None]
n_t += len(t)
lost_t += len(t) - len(tr)
if tr:
means.append(sum(tr) / len(tr))
a = aux_b.get(b)
points.append({
"t": round(start + b * bucket, 1),
"local": round(sum(lr) / len(lr), 2) if lr else None,
"total": round(sum(means) / len(means), 2) if means else None,
# The ISP leg per point, derived here rather than in QML so
# the inversion guard has one implementation. A panel that
# subtracted these itself would be a second copy of a rule
# whose whole purpose is refusing to answer, and the copy
# that forgets to refuse is the one that ships.
"wan": score.wan_point_ms(
round(sum(means) / len(means), 2) if means else None,
round(sum(lr) / len(lr), 2) if lr else None),
# None means no probe was sent in this bucket, which is a gap.
# A figure with no `total` means probes went out and nothing
# came back, which is down. The two must stay distinguishable:
# a gap is drawn as nothing, a down as an outage.
"loss": round((len(l) - len(lr) + lost_t) /
max(1, len(l) + n_t), 3) if (l or n_t) else None,
"rx": round(a[0], 1) if a and a[0] is not None else None,
"tx": round(a[1], 1) if a and a[1] is not None else None,
"sig": a[2] if a else None,
})
write_atomic(recent_path(), {"v": 1, "t": now, "bucket_s": bucket,
"points": points})
def flush_minute(self, now: float):
ls = Series.stats(self.local.since(60))
ts = self.total.stats(60)
ws = score.wan_from(ts, ls)
lag = score.lag_ms(ts)
resp = score.responsiveness(lag) if ts["count"] else None
# The old basis kept beside the new: the 0.2.0 switch to
# instrument-scored lag must stay auditable against what ICMP
# alone would have said — the 0.1.10 discipline, applied to
# ourselves.
icmp_stats = Series.stats(self.icmp_anchor.since(60))
icmp_lag = score.lag_ms(icmp_stats) if icmp_stats["count"] else None
out_frac, disruptions, disrupt_frac = self.store.outage_stats(24 * 3600, now)
rel = score.reliability(out_frac, disruptions,
disruption_fraction=disrupt_frac)
bloat = self.bufferbloat(300.0)
snap_link = self.link.latest
if snap_link.get("kind") != "wifi":
snap_link = {}
spd, spd_ctx = self.speed_score(now, snap_link.get("ssid", ""))
idx = score.index(resp, rel, spd if spd_ctx.get("scored") else None)
self.store.put_minute(
int(now // 60) * 60,
{
"local_p50": ls.get("p50"), "local_p95": ls.get("p95"),
"local_jitter": ls.get("jitter"), "local_loss": ls.get("loss"),
"wan_p50": ws.get("p50"), "wan_p95": ws.get("p95"),
"wan_jitter": ws.get("jitter"), "wan_loss": ws.get("loss"),
# Recorded, not scored — see SAMPLE_COLUMNS.
"local_p75": ls.get("p75"), "local_max": ls.get("max"),
"wan_p75": ws.get("p75"), "wan_max": ws.get("max"),
"lag": lag,
"rx_bps": self.rates[0], "tx_bps": self.rates[1],
"signal_dbm": snap_link.get("signal_dbm"),
"resp": resp, "rel": rel, "spd": spd, "idx": idx,
"lag_idle": bloat["idle"], "lag_loaded": bloat["loaded"],
"lag_icmp": icmp_lag,
# Stored so a published figure can be checked afterwards.
# `settled` as 1/0 rather than a bool: the column is REAL like
# its neighbours, and None stays None so "never measured" and
# "measured, did not settle" remain different answers.
"drain_ms": (bloat.get("drain") or {}).get("ms"),
"drain_min_ms": (bloat.get("drain") or {}).get("min_ms"),
"drain_settled": (
None if (bloat.get("drain") or {}).get("settled") is None
else float(bool((bloat["drain"])["settled"]))),
"drain_src": (bloat.get("drain") or {}).get("src"),
},
iface=self.route.get("iface", ""),
network=snap_link.get("ssid", ""),
probes="+".join(sorted(i.key for i in self.bench.actives())),
)
def flush_apps(self, now: float):
"""apps.json: top apps by TCP traffic, plus the honest remainder.
The interface moves bytes that no unprivileged tool can attribute —
UDP and with it QUIC, protocol overhead, other users' processes.
That remainder is published as its own bucket instead of being
left to look like the top apps account for everything.
"""
if not self.app_traffic.poll():
return
tcp_rx = sum(a["rx_bps"] for a in self.app_traffic.rates)
tcp_tx = sum(a["tx_bps"] for a in self.app_traffic.rates)
iface_rx = self.rates[0] or 0.0
iface_tx = self.rates[1] or 0.0
write_atomic(apps_path(), {
"v": 1,
"t": round(now, 1),
"apps": self.app_traffic.top(8),
"other": {
"rx_bps": round(max(0.0, iface_rx - tcp_rx), 1),
"tx_bps": round(max(0.0, iface_tx - tcp_tx), 1),
},
})
# ----------------------------------------------------------------- main
def loop(self):
tick = 0.5
while self.running:
now = time.time()
self.config.refresh()
self.restart_probes_if_settings_changed()
self.watch_outages(now)
self.throughput(now, self.route.get("iface", ""))
write_atomic(live_path(), self.compose_live(now))
if now - self.last_apps_poll >= 3.0:
self.last_apps_poll = now
self.flush_apps(now)
if now - self.last_minute_flush >= 60:
self.last_minute_flush = now
self.flush_minute(now)
self.flush_recent(now)
self.restart_probes_if_route_changed()
elif now - self.last_recent_flush >= 5.0:
self.flush_recent(now)
if now - self.last_rollup >= 3600:
self.last_rollup = now
self.store.rollup_hours(now)
self.store.prune(minute_days=int(self.config["historyDays"]),
now=now)
# Off the loop: the check is a curl, and the bar must keep
# updating at 2 Hz while it runs. tick() only starts and
# collects it; the address it proves is adopted here.
self.captive.tick(now, self._any_instrument_alive(6.0))
self.adopt_wan_ip()
self.update_watch.enabled = bool(self.config["updateCheck"])
self.update_watch.tick(now)
self.maybe_content_test(now)
if now - self._last_bench_eval >= BENCH_EVAL_EVERY_S:
self._last_bench_eval = now
self._apply_seats(
self.bench.evaluate(now, self._instrument_stats()))
# A peak asked for during the scheduled check waits for it to
# finish rather than being dropped or run on top of it.
if self.peak_requested.is_set() and not self.content_running:
self.peak_requested.clear()
self.run_peak_test()
time.sleep(max(0.1, tick - (time.time() - now)))
# Exit code contract with the shell service: LOCK_HELD means another
# instance owns the measurement and the service must not respawn us.
# Every other exit — including a clean 0 from SIGTERM — deserves a
# respawn, because a daemon that was asked to stop is still a daemon
# that is no longer measuring.
EXIT_LOCK_HELD = 3
def run(self):
if not self.acquire_lock():
print("nexthopd: another instance holds the lock, exiting",
file=sys.stderr)
return self.EXIT_LOCK_HELD
# Only now that the lock is ours: whatever a previous daemon left
# open, it will never close. See Store.close_orphans.
self.store.close_orphans(time.time())
retire_legacy_snapshots(self.state_dir, runtime_dir(), time.time())
signal.signal(signal.SIGTERM, self.stop)
signal.signal(signal.SIGINT, self.stop)
# SIGUSR1 is the "run a peak test" doorbell — file-free, and safe to
# send from a QML Process one-liner.
signal.signal(signal.SIGUSR1, lambda *_: self.peak_requested.set())
self.nl_events.start()
self.link.start()
self.start_probes()
try:
self.loop()
finally:
for p in self.probes:
p.stop()
self.link.stop()
self.nl_events.stop()
self.store.close()
return 0
def main():
return Daemon().run()
if __name__ == "__main__":
sys.exit(main())