Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d3ccd6f94e | ||
|
|
ea062af86a |
@@ -29,3 +29,52 @@ See [`CONTRIBUTING.md`](CONTRIBUTING.md) for the development workflow and TDD re
|
|||||||
- TypeScript: strict mode, functional components and hooks, no extra UI libraries.
|
- TypeScript: strict mode, functional components and hooks, no extra UI libraries.
|
||||||
- YAML: 2-space indent.
|
- YAML: 2-space indent.
|
||||||
- Diagrams: Mermaid inside Markdown.
|
- Diagrams: Mermaid inside Markdown.
|
||||||
|
|
||||||
|
## Decision records
|
||||||
|
|
||||||
|
Significant architectural decisions are recorded chronologically as **ADRs** in
|
||||||
|
[`docs/adr/`](docs/adr/) — one immutable file per decision, in the
|
||||||
|
Status / Context / Options / Decision / Consequences format. Start from
|
||||||
|
[`docs/adr/README.md`](docs/adr/README.md).
|
||||||
|
|
||||||
|
- **ADR** = a point-in-time decision for *this* repo → `docs/adr/`.
|
||||||
|
- **ASR** = a standing, ecosystem-wide standard (Go-first, secrets, layout) →
|
||||||
|
lives in the `inventar` repo, **not** here.
|
||||||
|
|
||||||
|
When a change alters a boundary, contract, or trade-off, add an ADR (copy
|
||||||
|
[`docs/adr/TEMPLATE.md`](docs/adr/TEMPLATE.md), take the next number, update the
|
||||||
|
index) and cross-link it from the affected design doc. A later reversal gets a
|
||||||
|
new ADR that supersedes the old one — never edit an Accepted record.
|
||||||
|
|
||||||
|
## Runtime & operations (sim environment)
|
||||||
|
|
||||||
|
The k3d ground miniature is `swarm-sim`. Services are exposed as NodePorts:
|
||||||
|
|
||||||
|
| Service | URL |
|
||||||
|
| --- | --- |
|
||||||
|
| Explorer (T1 live lake) | `http://localhost:30088` |
|
||||||
|
| Warehouse explorer (T3) | `http://localhost:30089` |
|
||||||
|
| Grafana (dashboards `swarm-fleet`, `swarm-platform`) | `http://localhost:30300` |
|
||||||
|
| Prometheus | `http://localhost:30990` |
|
||||||
|
| MinIO console | `http://localhost:30901` |
|
||||||
|
| Prototype dev server (`npm run dev`) | `http://localhost:5173` |
|
||||||
|
|
||||||
|
Rebuild the simulator image and roll it out:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
cd simulator && docker build -t swarm-house/simulator:dev .
|
||||||
|
k3d image import swarm-house/simulator:dev -c swarm-sim
|
||||||
|
kubectl rollout restart -n swarm statefulset/drone deployment/explorer deployment/exporter
|
||||||
|
cd infra/terraform/sim-env && terraform apply
|
||||||
|
```
|
||||||
|
|
||||||
|
**CPU budget knobs** (see [ADR-0009](docs/adr/ADR-0009-bounded-lake-scans.md)):
|
||||||
|
`METRIC_FLIGHT_WINDOW` (scan only recent flights), `SCAN_INTERVAL_S`,
|
||||||
|
`VIEW_REFRESH_S`, `TREE_CACHE_S`, `KEEP_ALIVE=1` (drones idle after seal instead
|
||||||
|
of restarting), `IMU_HZ`. Never scan the full lake on a hot path — it grows with
|
||||||
|
pod churn. Long-horizon analysis is a T3 warehouse query, not an exporter job.
|
||||||
|
|
||||||
|
Runtime data lives under `var/t1` / `var/t3`, owned by uid **10001**
|
||||||
|
([ADR-0007](docs/adr/ADR-0007-non-root-image.md),
|
||||||
|
[ADR-0008](docs/adr/ADR-0008-runtime-data-paths.md)); both are gitignored.
|
||||||
|
Terraform state is local and never committed.
|
||||||
|
|||||||
@@ -81,9 +81,18 @@ identifying names of people, companies, or locations.
|
|||||||
- [ ] Tests added or updated; `pytest tests/ -q` passes locally
|
- [ ] Tests added or updated; `pytest tests/ -q` passes locally
|
||||||
- [ ] `docker build` in `simulator/` succeeds (test stage included)
|
- [ ] `docker build` in `simulator/` succeeds (test stage included)
|
||||||
- [ ] Documentation updated if behaviour or boundaries changed
|
- [ ] Documentation updated if behaviour or boundaries changed
|
||||||
|
- [ ] **ADR added** if the change alters an architectural boundary, contract, or trade-off (see [`docs/adr/`](docs/adr/))
|
||||||
- [ ] English only; no identifying references
|
- [ ] English only; no identifying references
|
||||||
- [ ] Conventional commit message
|
- [ ] Conventional commit message
|
||||||
|
|
||||||
|
## Decisions
|
||||||
|
|
||||||
|
Architectural decisions are recorded chronologically as ADRs in
|
||||||
|
[`docs/adr/`](docs/adr/). Read [`docs/adr/README.md`](docs/adr/README.md) before
|
||||||
|
changing a boundary or contract; add a new ADR (from
|
||||||
|
[`docs/adr/TEMPLATE.md`](docs/adr/TEMPLATE.md)) rather than editing an Accepted
|
||||||
|
one when you reverse course.
|
||||||
|
|
||||||
## Questions
|
## Questions
|
||||||
|
|
||||||
Open design questions live in [`docs/09-open-questions.md`](docs/09-open-questions.md).
|
Open design questions live in [`docs/09-open-questions.md`](docs/09-open-questions.md).
|
||||||
|
|||||||
@@ -33,6 +33,9 @@ This repository describes how to build, deliver, test, and operate that platform
|
|||||||
| [09 — Open questions](docs/09-open-questions.md) | Known unknowns and proposed answers |
|
| [09 — Open questions](docs/09-open-questions.md) | Known unknowns and proposed answers |
|
||||||
| [10 — Domain context](docs/10-domain-context.md) | Swarm autonomy principles this design builds on |
|
| [10 — Domain context](docs/10-domain-context.md) | Swarm autonomy principles this design builds on |
|
||||||
| [11 — CI/CD & delivery](docs/11-cicd-delivery.md) | Pipelines, registry, dev containers, fleet releases |
|
| [11 — CI/CD & delivery](docs/11-cicd-delivery.md) | Pipelines, registry, dev containers, fleet releases |
|
||||||
|
| [ADRs](docs/adr/) | Architecture Decision Records — the decisions behind the above, in the order they were made |
|
||||||
|
|
||||||
|
The design principles below are the *what*; the [ADRs](docs/adr/) are the *why and when* — each principle traces to a dated, immutable decision record.
|
||||||
|
|
||||||
## Runnable parts
|
## Runnable parts
|
||||||
|
|
||||||
|
|||||||
@@ -2,6 +2,8 @@
|
|||||||
|
|
||||||
Zero trust, fully air-gapped, everything provisioned before wheels-up.
|
Zero trust, fully air-gapped, everything provisioned before wheels-up.
|
||||||
|
|
||||||
|
> **Decision:** [ADR-0005 — WireGuard beneath SSH](adr/ADR-0005-wireguard-beneath-ssh.md).
|
||||||
|
|
||||||
## Trust model
|
## Trust model
|
||||||
|
|
||||||
- **Nothing joins the swarm at runtime.** Every device receives its identity (key pair + certificate) during ground provisioning, signed by the fleet's offline CA. There is no trust-on-first-use, no runtime enrollment endpoint, no exception path.
|
- **Nothing joins the swarm at runtime.** Every device receives its identity (key pair + certificate) during ground provisioning, signed by the fleet's offline CA. There is no trust-on-first-use, no runtime enrollment endpoint, no exception path.
|
||||||
|
|||||||
@@ -2,6 +2,8 @@
|
|||||||
|
|
||||||
Three infrastructures, one data contract. The layout, schemas, and pipelines are identical everywhere; only scale and lifetime differ.
|
Three infrastructures, one data contract. The layout, schemas, and pipelines are identical everywhere; only scale and lifetime differ.
|
||||||
|
|
||||||
|
> **Decisions:** [ADR-0006 — IaC boundaries](adr/ADR-0006-iac-boundaries.md), [ADR-0008 — runtime data paths](adr/ADR-0008-runtime-data-paths.md).
|
||||||
|
|
||||||
```mermaid
|
```mermaid
|
||||||
graph LR
|
graph LR
|
||||||
subgraph dev [3 — Dev and simulation]
|
subgraph dev [3 — Dev and simulation]
|
||||||
|
|||||||
@@ -2,6 +2,8 @@
|
|||||||
|
|
||||||
The platform watches itself with the same discipline it applies to sensor data — and largely through the same pipeline.
|
The platform watches itself with the same discipline it applies to sensor data — and largely through the same pipeline.
|
||||||
|
|
||||||
|
> **Decision:** [ADR-0009 — bound observability lake scans to a flight window and cache](adr/ADR-0009-bounded-lake-scans.md).
|
||||||
|
|
||||||
## Two kinds of signals, one storage
|
## Two kinds of signals, one storage
|
||||||
|
|
||||||
| Signal | Examples | Where it goes |
|
| Signal | Examples | Where it goes |
|
||||||
@@ -53,6 +55,13 @@ graph LR
|
|||||||
| Which link degraded first, and was it distance or interference? | RSSI telemetry joined with pose distance |
|
| Which link degraded first, and was it distance or interference? | RSSI telemetry joined with pose distance |
|
||||||
| Is the new writer version flushing slower on ARM64? | CI simulation runs emit the same metrics; diff across fleet releases |
|
| Is the new writer version flushing slower on ARM64? | CI simulation runs emit the same metrics; diff across fleet releases |
|
||||||
|
|
||||||
A working miniature of this doctrine ships in [`simulator/`](../simulator/): the `monitoring` Compose profile runs a DuckDB-based exporter over the generated Parquet plus Prometheus and a provisioned Grafana dashboard — fleet statistics derived from the data platform itself, with no agent on the (virtual) drones.
|
A working miniature of this doctrine ships in [`simulator/`](../simulator/): the `monitoring` Compose profile runs a DuckDB-based exporter over the generated Parquet plus Prometheus and two provisioned Grafana dashboards:
|
||||||
|
|
||||||
|
| Dashboard | Focus |
|
||||||
|
| --- | --- |
|
||||||
|
| **Swarm Fleet — generated data** | Row counts, detections, pose frames, Parquet bytes, battery, RSSI |
|
||||||
|
| **Swarm Platform — CPU & scan health** | Node CPU/memory (node-exporter), exporter/explorer scan durations, flight-window gauge |
|
||||||
|
|
||||||
|
The exporter and explorer **scan only the most recent flight partitions** (`METRIC_FLIGHT_WINDOW`, default 5) and cache tree/views — otherwise CPU climbs as every pod restart appends a new `flight=` tree to the shared lake. Drones in k3d set `KEEP_ALIVE=1` after sealing so they idle instead of exiting and spawning another flight.
|
||||||
|
|
||||||
Alerting on the ground follows standard practice (Prometheus alert rules for infrastructure, CI gates for regression in simulated staleness/throughput budgets). In flight there is nobody to page — the platform's job is to degrade in the documented order and record everything for the post-mortem.
|
Alerting on the ground follows standard practice (Prometheus alert rules for infrastructure, CI gates for regression in simulated staleness/throughput budgets). In flight there is nobody to page — the platform's job is to degrade in the documented order and record everything for the post-mortem.
|
||||||
|
|||||||
@@ -61,3 +61,9 @@ Who owns the device registry (drone ids, key issuance, revocation) organizationa
|
|||||||
Sub-250 g class units cannot run the full stack (no GPU, minimal CPU/storage).
|
Sub-250 g class units cannot run the full stack (no GPU, minimal CPU/storage).
|
||||||
|
|
||||||
*Proposal:* define a minimal profile early — state publisher + UDP-multicast pose broadcast, no local Parquet store, no video — so the swarm protocol never assumes full-stack peers.
|
*Proposal:* define a minimal profile early — state publisher + UDP-multicast pose broadcast, no local Parquet store, no video — so the swarm protocol never assumes full-stack peers.
|
||||||
|
|
||||||
|
## 11 — Lake retention and partition pruning
|
||||||
|
|
||||||
|
The observability path now scans only the most recent flights ([ADR-0009](adr/ADR-0009-bounded-lake-scans.md)), but old `flight=` partitions still accumulate on T1 disk after offload to T3.
|
||||||
|
|
||||||
|
*Proposal:* a scheduled prune of T1 partitions whose flights are confirmed present in the T3 warehouse (offload as the retention gate), configurable by age and free-space watermark. On the drone the same policy is bounded by the NVMe quota ladder ([03](03-data-platform.md)); in the sim environment it is a ground CronJob alongside the offload job.
|
||||||
|
|||||||
@@ -2,6 +2,8 @@
|
|||||||
|
|
||||||
How software, models, and configuration reach the fleet — reproducibly, scanned, versioned, and atomically rollback-able. Everything lives inside the air gap.
|
How software, models, and configuration reach the fleet — reproducibly, scanned, versioned, and atomically rollback-able. Everything lives inside the air gap.
|
||||||
|
|
||||||
|
> **Decisions:** [ADR-0006 — IaC boundaries & fleet manifest](adr/ADR-0006-iac-boundaries.md), [ADR-0007 — non-root image](adr/ADR-0007-non-root-image.md).
|
||||||
|
|
||||||
## Source and pipelines: GitLab
|
## Source and pipelines: GitLab
|
||||||
|
|
||||||
GitLab (self-managed, on-prem) is the backbone: repositories, CI/CD, and the container registry in one system.
|
GitLab (self-managed, on-prem) is the backbone: repositories, CI/CD, and the container registry in one system.
|
||||||
|
|||||||
@@ -0,0 +1,38 @@
|
|||||||
|
# ADR-0001: No swarm-wide orchestrator; coordinate through data
|
||||||
|
|
||||||
|
## Status
|
||||||
|
|
||||||
|
Accepted (2026-07-08)
|
||||||
|
|
||||||
|
## Context
|
||||||
|
|
||||||
|
The fleet flies fully autonomously with no internet uplink and only
|
||||||
|
intermittent mesh connectivity between drones. A conventional control plane
|
||||||
|
(Kubernetes across the swarm, a leader election, a central scheduler) assumes
|
||||||
|
stable membership and low-latency links — exactly what a swarm does not have.
|
||||||
|
Partitions are the normal case, not the exception.
|
||||||
|
|
||||||
|
## Options
|
||||||
|
|
||||||
|
| Option | Pros | Cons |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| A — Cluster orchestrator spanning the swarm | Familiar tooling | Assumes stable membership; a partition stalls coordination; single point of failure in the air |
|
||||||
|
| B — Each drone autonomous, coordinating through exchanged data | Survives partitions; no air-side control plane to fail | No global view; behaviour must be derivable from local state |
|
||||||
|
|
||||||
|
## Decision
|
||||||
|
|
||||||
|
Every drone is an autonomous node. There is **no orchestrator spanning the
|
||||||
|
swarm**. Coordination happens by exchanging derived state (see
|
||||||
|
[ADR-0003](ADR-0003-sync-derived-state-only.md)), and each drone makes its own
|
||||||
|
in-flight decisions against its local data window. Orchestration tooling
|
||||||
|
(Kubernetes, Flux) is confined to the **ground** segment
|
||||||
|
([ADR-0006](ADR-0006-iac-boundaries.md)).
|
||||||
|
|
||||||
|
## Consequences
|
||||||
|
|
||||||
|
- The platform degrades gracefully under partition: a drone that loses peers
|
||||||
|
keeps flying and recording.
|
||||||
|
- There is no single "fleet state" to query in the air; health questions are
|
||||||
|
answered from each drone's own store and reconciled on the ground.
|
||||||
|
- Design principle 1 in [`../../README.md`](../../README.md) and the on-board
|
||||||
|
model in [`../02-architecture.md`](../02-architecture.md) follow from this.
|
||||||
@@ -0,0 +1,32 @@
|
|||||||
|
# ADR-0002: One storage format on every tier — Parquet + DuckDB, Hive layout
|
||||||
|
|
||||||
|
## Status
|
||||||
|
|
||||||
|
Accepted (2026-07-08)
|
||||||
|
|
||||||
|
## Context
|
||||||
|
|
||||||
|
Data lives in three places: on the drone's NVMe in flight (T1), in the ground
|
||||||
|
warehouse after landing (T3), and in the dev/simulation environment. If each
|
||||||
|
tier used a different format or layout, offload would need transformation code,
|
||||||
|
schemas would drift, and a query proven on the bench would not run unchanged on
|
||||||
|
flight data.
|
||||||
|
|
||||||
|
## Decision
|
||||||
|
|
||||||
|
Use **one storage format on every floor**: columnar **Parquet** files in an
|
||||||
|
identical **Hive-partitioned** layout (`dataset=…/flight=…/drone=…/…`), queried
|
||||||
|
with **DuckDB**. Flight offload T1 → T3 is therefore a plain **mirror**
|
||||||
|
(copy/rsync of partition directories), not an ETL step.
|
||||||
|
|
||||||
|
## Consequences
|
||||||
|
|
||||||
|
- Offload is a file operation, testable with a checksum; no format converters
|
||||||
|
to maintain or version.
|
||||||
|
- The same SQL runs against T1, T3, and simulated data — the simulator
|
||||||
|
([`../../simulator/`](../../simulator/)) generates the *real* layout, so tests
|
||||||
|
exercise the production format.
|
||||||
|
- Partitioning choices are load-bearing; changing them is a schema migration.
|
||||||
|
Detailed in [`../03-data-platform.md`](../03-data-platform.md).
|
||||||
|
- DuckDB's single-file, no-server model fits an air-gapped drone with no room
|
||||||
|
for a database daemon.
|
||||||
@@ -0,0 +1,33 @@
|
|||||||
|
# ADR-0003: Sync derived state only; raw telemetry stays local
|
||||||
|
|
||||||
|
## Status
|
||||||
|
|
||||||
|
Accepted (2026-07-08)
|
||||||
|
|
||||||
|
## Context
|
||||||
|
|
||||||
|
Each drone produces high-rate sensor telemetry and on-board video detections —
|
||||||
|
far more than the mesh can carry. But peers only need enough to coordinate
|
||||||
|
flight: where everyone is, where they are heading, what was detected. Pushing
|
||||||
|
raw feeds across the air would saturate the radio and drain batteries for data
|
||||||
|
nobody reads in flight.
|
||||||
|
|
||||||
|
## Decision
|
||||||
|
|
||||||
|
Only **derived state** crosses the air — position, attitude, detections — as
|
||||||
|
small pub/sub broadcasts (~5 Hz). **Raw telemetry stays on local NVMe** until
|
||||||
|
the drone lands, then travels to the warehouse during offload
|
||||||
|
([ADR-0002](ADR-0002-one-storage-format.md)). The design is **event-driven**:
|
||||||
|
new derived data triggers downstream action through hooks; nothing polls asking
|
||||||
|
"anything new yet?".
|
||||||
|
|
||||||
|
## Consequences
|
||||||
|
|
||||||
|
- Air bandwidth scales with the number of peers and the broadcast rate, not
|
||||||
|
with sensor resolution.
|
||||||
|
- The warehouse is the only place with the full picture; in-flight decisions use
|
||||||
|
the local window plus peer broadcasts.
|
||||||
|
- Broadcast staleness per peer becomes a first-class health signal
|
||||||
|
([`../07-observability.md`](../07-observability.md)).
|
||||||
|
- What syncs and what does not is specified in
|
||||||
|
[`../04-swarm-sync.md`](../04-swarm-sync.md).
|
||||||
@@ -0,0 +1,41 @@
|
|||||||
|
# ADR-0004: Read-only SQL-over-SSH as the peer query contract; no debug API
|
||||||
|
|
||||||
|
## Status
|
||||||
|
|
||||||
|
Accepted (2026-07-08). Supersedes the earlier sketch of a bespoke query service.
|
||||||
|
|
||||||
|
## Context
|
||||||
|
|
||||||
|
Beyond the 5 Hz broadcast, a peer (or an engineer on the bench) sometimes needs
|
||||||
|
to *ask* a drone a richer question: "show me your detections in the last minute
|
||||||
|
near this position". Inventing a service for this means a new port, a new
|
||||||
|
protocol, an auth layer, and a second code path that only runs during
|
||||||
|
debugging — the classic "works on the bench, untested in flight" trap.
|
||||||
|
|
||||||
|
## Options
|
||||||
|
|
||||||
|
| Option | Pros | Cons |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| A — Custom query microservice + REST API | Flexible | New port/protocol/auth to secure; a debug-only path that never flies |
|
||||||
|
| B — Read-only DuckDB SQL over an SSH forced command | Reuses SSH trust + the query language already on both ends; identical path on bench and in flight | SQL must be gated to read-only |
|
||||||
|
|
||||||
|
## Decision
|
||||||
|
|
||||||
|
Peer data access is **read-only DuckDB SQL executed through an SSH forced
|
||||||
|
command**. SQL is already the query language on both ends
|
||||||
|
([ADR-0002](ADR-0002-one-storage-format.md)), and SSH already carries the trust
|
||||||
|
model ([ADR-0005](ADR-0005-wireguard-beneath-ssh.md)), so no new service, port,
|
||||||
|
or protocol is invented. A **SQL gate** rejects anything that writes
|
||||||
|
(`COPY`, `INSERT`, `ATTACH`, pragmas). There is **no separate debug API**: an
|
||||||
|
engineer debugging on the ground runs the identical query through the identical
|
||||||
|
wrapper, permissions, and output format a peer drone would use.
|
||||||
|
|
||||||
|
## Consequences
|
||||||
|
|
||||||
|
- One access path is built, secured, and tested — "what you test is what
|
||||||
|
flies".
|
||||||
|
- The SQL gate is safety-critical and is covered by unit tests
|
||||||
|
(`simulator/tests/test_sql_gate.py`).
|
||||||
|
- The explorer and the prototype's live mode are *just another read-only
|
||||||
|
consumer* of this same contract ([`../04-swarm-sync.md`](../04-swarm-sync.md)).
|
||||||
|
- Design principles 7 and 8 in [`../../README.md`](../../README.md) restate this.
|
||||||
@@ -0,0 +1,40 @@
|
|||||||
|
# ADR-0005: WireGuard beneath SSH for the mesh transport
|
||||||
|
|
||||||
|
## Status
|
||||||
|
|
||||||
|
Accepted (2026-07-08)
|
||||||
|
|
||||||
|
## Context
|
||||||
|
|
||||||
|
Peer bulk sync and SQL-over-SSH ([ADR-0004](ADR-0004-sql-over-ssh-contract.md))
|
||||||
|
need an encrypted, authenticated channel across an ad-hoc radio mesh. SSH alone
|
||||||
|
would work at the application layer, but it leaves the drones' listening
|
||||||
|
surface (sshd, any other service port) exposed on the raw mesh network, and it
|
||||||
|
gives no cheap, uniform way to authenticate *the peer* before the application
|
||||||
|
handshake.
|
||||||
|
|
||||||
|
## Options
|
||||||
|
|
||||||
|
| Option | Pros | Cons |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| A — SSH only | One daemon, familiar | Service ports exposed on the raw mesh; per-service auth; no network-layer peer identity |
|
||||||
|
| B — WireGuard tunnel, SSH inside it | Network-layer peer auth via static keys; only the WG port is exposed; SSH speaks over a private, encrypted subnet | Two layers to provision |
|
||||||
|
|
||||||
|
## Decision
|
||||||
|
|
||||||
|
Run **WireGuard as the network layer** and **SSH on top of it**. WireGuard
|
||||||
|
authenticates peers by pre-shared static public keys and presents a private
|
||||||
|
encrypted subnet; SSH forced commands then carry the SQL and bulk-sync contract
|
||||||
|
over that subnet. Keys are issued **per device before deployment** — nothing
|
||||||
|
joins the mesh at runtime (design principle 5).
|
||||||
|
|
||||||
|
## Consequences
|
||||||
|
|
||||||
|
- The only thing exposed on the raw radio is the WireGuard port; sshd and the
|
||||||
|
data services listen only on the WG interface.
|
||||||
|
- Peer identity is established once, at the network layer, and reused by every
|
||||||
|
application above it.
|
||||||
|
- Provisioning must place both WG and SSH keys during host prep — handled by
|
||||||
|
`infra/ansible/drone-provision.yml` ([ADR-0006](ADR-0006-iac-boundaries.md)).
|
||||||
|
- Rationale and the "why not SSH alone" argument live in
|
||||||
|
[`../05-network-security.md`](../05-network-security.md).
|
||||||
@@ -0,0 +1,41 @@
|
|||||||
|
# ADR-0006: IaC split — Ansible hosts, Terraform ground, Flux ground-only, fleet manifest for drones
|
||||||
|
|
||||||
|
## Status
|
||||||
|
|
||||||
|
Accepted (2026-07-08)
|
||||||
|
|
||||||
|
## Context
|
||||||
|
|
||||||
|
Three very different things need provisioning: (1) drone hosts, which are
|
||||||
|
air-gapped and never reachable by a controller in flight; (2) the ground
|
||||||
|
segment (warehouse, observability, offload), a normal k3s cluster; (3) the
|
||||||
|
simulation environment that mirrors the ground segment locally. Using one tool
|
||||||
|
for all three would force a GitOps controller or a Terraform apply loop onto the
|
||||||
|
drone — impossible for an autonomous, disconnected node
|
||||||
|
([ADR-0001](ADR-0001-no-swarm-orchestrator.md)).
|
||||||
|
|
||||||
|
## Decision
|
||||||
|
|
||||||
|
Split infrastructure ownership by lifecycle:
|
||||||
|
|
||||||
|
| Layer | Tool | Scope |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| Host preparation | **Ansible** | Drone bench prep (keys, forced commands, WireGuard, Compose bundle) and the local k3d sim cluster |
|
||||||
|
| Ground workloads | **Terraform** | k3s/k3d resources: virtual fleet + observability (`sim-env`), warehouse + T1→T3 offload (`ground`) |
|
||||||
|
| Ground continuous delivery | **Flux (GitOps)** | Ground segment **only** — policy labels, dashboard bundles |
|
||||||
|
| Drone delivery | **Fleet release manifest** | A versioned, digest-pinned artifact list; no controller pulls to the drone |
|
||||||
|
|
||||||
|
Drones are delivered by the **fleet release manifest** (one version for the
|
||||||
|
whole fleet, atomic rollback), never by a controller reaching into the air.
|
||||||
|
|
||||||
|
## Consequences
|
||||||
|
|
||||||
|
- No GitOps agent or `terraform apply` ever targets a drone; the air side stays
|
||||||
|
controller-free.
|
||||||
|
- Terraform state is environment-local and kept out of Git
|
||||||
|
([ADR-0008](ADR-0008-runtime-data-paths.md)).
|
||||||
|
- The sim cluster is a faithful miniature: the same Terraform describes it and
|
||||||
|
the ground segment.
|
||||||
|
- Boundaries are documented in
|
||||||
|
[`../06-environments.md`](../06-environments.md#iac-boundaries); delivery in
|
||||||
|
[`../11-cicd-delivery.md`](../11-cicd-delivery.md).
|
||||||
@@ -0,0 +1,27 @@
|
|||||||
|
# ADR-0007: Run the simulator image as a non-root user
|
||||||
|
|
||||||
|
## Status
|
||||||
|
|
||||||
|
Accepted (2026-07-08)
|
||||||
|
|
||||||
|
## Context
|
||||||
|
|
||||||
|
The Trivy scan in CI flagged **DS-0002**: the simulator image ran as `root`
|
||||||
|
because no `USER` was set. A container writing the data lake as root also
|
||||||
|
produces root-owned Parquet on bind mounts, which then blocks a later non-root
|
||||||
|
process from the same files.
|
||||||
|
|
||||||
|
## Decision
|
||||||
|
|
||||||
|
Create a dedicated unprivileged user **`swarm` (uid/gid 10001)** in the
|
||||||
|
`simulator/Dockerfile` and switch to it with `USER swarm` before the entrypoint.
|
||||||
|
The uid is fixed (not auto-assigned) so host-side ownership of the shared data
|
||||||
|
paths ([ADR-0008](ADR-0008-runtime-data-paths.md)) is deterministic.
|
||||||
|
|
||||||
|
## Consequences
|
||||||
|
|
||||||
|
- Trivy DS-0002 clears; the image follows least-privilege.
|
||||||
|
- Bind-mounted data is owned by `10001:10001`; host prep (`sim-cluster.yml`) and
|
||||||
|
Ansible `chown` the `var/t1` / `var/t3` paths to this uid so k3d pods can write.
|
||||||
|
- Any future service image reusing this data must run as the same uid or share
|
||||||
|
the group.
|
||||||
@@ -0,0 +1,40 @@
|
|||||||
|
# ADR-0008: Shared `var/t1` / `var/t3` runtime paths; state out of Git
|
||||||
|
|
||||||
|
## Status
|
||||||
|
|
||||||
|
Accepted (2026-07-08)
|
||||||
|
|
||||||
|
## Context
|
||||||
|
|
||||||
|
Compose, the k3d fleet, and the ground offload job all need to read and write
|
||||||
|
the same data lake, so that "go live" in the prototype shows the data the
|
||||||
|
simulator just generated. Early on, Terraform state and the provider cache were
|
||||||
|
accidentally committed (~6 MB of `*.tfstate` and a vendored provider binary),
|
||||||
|
polluting history and risking stale/secret data in Git.
|
||||||
|
|
||||||
|
## Decision
|
||||||
|
|
||||||
|
Standardise two **gitignored** runtime directories at the repo root:
|
||||||
|
|
||||||
|
| Path | Tier | Contents |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| `var/t1/` | T1 live lake | In-flight Parquet from Compose or the k3d fleet |
|
||||||
|
| `var/t3/` | T3 warehouse | Offloaded historical partitions |
|
||||||
|
|
||||||
|
They are created automatically and overridable via
|
||||||
|
`SWARM_T1_DIR` / `SWARM_WAREHOUSE_DATA` / `SWARM_SIM_DATA`. Treat `var/` as
|
||||||
|
**runtime-only** per the ecosystem canonical layout. `.gitignore` excludes
|
||||||
|
`*.tfstate`, `*.tfstate.backup`, and `.terraform/` while **keeping**
|
||||||
|
`.terraform.lock.hcl` (the lock is source, the state is not).
|
||||||
|
|
||||||
|
## Consequences
|
||||||
|
|
||||||
|
- One lake feeds every component; the live prototype reflects real generated
|
||||||
|
data.
|
||||||
|
- Git history carries no machine state or large binaries; `terraform` treats
|
||||||
|
its state as local/remote-backend concern, not versioned source.
|
||||||
|
- Ownership of these paths is fixed to uid 10001
|
||||||
|
([ADR-0007](ADR-0007-non-root-image.md)).
|
||||||
|
- Recovery from earlier state drift is documented in
|
||||||
|
[`../../infra/terraform/README.md`](../../infra/terraform/README.md)
|
||||||
|
(`terraform import`).
|
||||||
@@ -0,0 +1,52 @@
|
|||||||
|
# ADR-0009: Bound observability lake scans to a flight window and cache
|
||||||
|
|
||||||
|
## Status
|
||||||
|
|
||||||
|
Accepted (2026-07-08)
|
||||||
|
|
||||||
|
## Context
|
||||||
|
|
||||||
|
The k3d sim node was running at ~245% CPU. Three causes compounded:
|
||||||
|
|
||||||
|
1. **Explorer** rebuilt a DuckDB view over the *entire* lake on every HTTP
|
||||||
|
request (~2000 Parquet files, 200 MB+) — roughly 1.5 cores.
|
||||||
|
2. **Exporter** scanned the full lake every 5 s — roughly 0.7 cores.
|
||||||
|
3. **Drones** exited after sealing a flight; Kubernetes restarted each pod,
|
||||||
|
and every restart appended a fresh `flight=` partition tree, so the lake grew
|
||||||
|
without bound and every scan got more expensive.
|
||||||
|
|
||||||
|
The cost is inherent to scanning an append-only lake whose size grows with pod
|
||||||
|
churn, not to any single slow query. Metrics and the live view only ever need
|
||||||
|
recent flights.
|
||||||
|
|
||||||
|
## Decision
|
||||||
|
|
||||||
|
Bound and cache the scans instead of reading the whole lake:
|
||||||
|
|
||||||
|
- **Flight window.** The exporter and explorer scan only the most recent
|
||||||
|
`METRIC_FLIGHT_WINDOW` flights (default 5) via a shared helper
|
||||||
|
(`simulator/monitoring/lake.py`). Older partitions stay on disk but out of the
|
||||||
|
hot path.
|
||||||
|
- **Caching.** Explorer caches DuckDB views (`VIEW_REFRESH_S=30`) and the
|
||||||
|
partition tree (`TREE_CACHE_S=15`); the exporter scans every
|
||||||
|
`SCAN_INTERVAL_S=30` instead of 5 s.
|
||||||
|
- **Stop restart churn.** Drones set `KEEP_ALIVE=1` and idle after sealing
|
||||||
|
instead of exiting, so no new flight tree is created per restart. In k3d,
|
||||||
|
`IMU_HZ=20` and resource limits cap per-drone load; the live prototype polls
|
||||||
|
every 2 s.
|
||||||
|
- **Self-observability.** Add node-exporter, a
|
||||||
|
`swarm_exporter_scan_duration_seconds` metric, and a **Swarm Platform**
|
||||||
|
Grafana dashboard (node CPU/memory + scan health) alongside the existing fleet
|
||||||
|
dashboard.
|
||||||
|
|
||||||
|
## Consequences
|
||||||
|
|
||||||
|
- Steady state dropped to explorer ~29 m, exporter ~16 m CPU; the k3d node to
|
||||||
|
~56% (from ~245%).
|
||||||
|
- Metrics reflect recent flights only; long-horizon analysis is a warehouse
|
||||||
|
(T3) query, not an exporter concern — consistent with
|
||||||
|
[ADR-0003](ADR-0003-sync-derived-state-only.md).
|
||||||
|
- Old flight partitions accumulate on disk; pruning them is a follow-up
|
||||||
|
(tracked in [`../09-open-questions.md`](../09-open-questions.md)).
|
||||||
|
- Observability design is documented in
|
||||||
|
[`../07-observability.md`](../07-observability.md).
|
||||||
@@ -0,0 +1,40 @@
|
|||||||
|
# Architecture Decision Records
|
||||||
|
|
||||||
|
This directory records the significant decisions behind Swarm House, in the
|
||||||
|
order they were made. Each record is immutable once **Accepted** — a later
|
||||||
|
decision that changes course gets a new number and supersedes the old one
|
||||||
|
rather than editing it.
|
||||||
|
|
||||||
|
## ADR vs ASR
|
||||||
|
|
||||||
|
| Kind | What it captures | Where it lives |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| **ADR** (Architecture Decision Record) | A point-in-time decision for *this* repository, with the options considered and their consequences | Here, in `docs/adr/` |
|
||||||
|
| **ASR** (Architecture Standard Record) | A standing, cross-repository policy every project must obey (Go-first, secrets, layout) | The ecosystem `inventar` repo, not here |
|
||||||
|
|
||||||
|
The rules in [`../../AGENTS.md`](../../AGENTS.md) are this repo's local standing
|
||||||
|
constraints; anything ecosystem-wide belongs in `inventar` as an ASR. Design
|
||||||
|
decisions specific to the platform are ADRs and belong here.
|
||||||
|
|
||||||
|
## Index
|
||||||
|
|
||||||
|
| ADR | Decision | Status |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| [ADR-0001](ADR-0001-no-swarm-orchestrator.md) | No swarm-wide orchestrator; coordinate through data | Accepted |
|
||||||
|
| [ADR-0002](ADR-0002-one-storage-format.md) | One storage format on every tier — Parquet + DuckDB, Hive layout | Accepted |
|
||||||
|
| [ADR-0003](ADR-0003-sync-derived-state-only.md) | Sync derived state only; raw telemetry stays local | Accepted |
|
||||||
|
| [ADR-0004](ADR-0004-sql-over-ssh-contract.md) | Read-only SQL-over-SSH as the peer query contract; no debug API | Accepted |
|
||||||
|
| [ADR-0005](ADR-0005-wireguard-beneath-ssh.md) | WireGuard beneath SSH for the mesh transport | Accepted |
|
||||||
|
| [ADR-0006](ADR-0006-iac-boundaries.md) | IaC split: Ansible hosts, Terraform ground, Flux ground-only, fleet manifest for drones | Accepted |
|
||||||
|
| [ADR-0007](ADR-0007-non-root-image.md) | Run the simulator image as a non-root user | Accepted |
|
||||||
|
| [ADR-0008](ADR-0008-runtime-data-paths.md) | Shared `var/t1` / `var/t3` runtime paths; state out of Git | Accepted |
|
||||||
|
| [ADR-0009](ADR-0009-bounded-lake-scans.md) | Bound observability lake scans to a flight window and cache | Accepted |
|
||||||
|
|
||||||
|
## Writing a new ADR
|
||||||
|
|
||||||
|
1. Copy [`TEMPLATE.md`](TEMPLATE.md) to `ADR-NNNN-short-slug.md` (next free number).
|
||||||
|
2. Fill in Context, Options (if more than one was real), Decision, Consequences.
|
||||||
|
3. Add a row to the index above.
|
||||||
|
4. Cross-link the ADR from the design doc it affects (and vice versa).
|
||||||
|
|
||||||
|
Keep the tone laconic: state the decision and why, skip the narrative.
|
||||||
@@ -0,0 +1,29 @@
|
|||||||
|
# ADR-NNNN: Short decision title
|
||||||
|
|
||||||
|
## Status
|
||||||
|
|
||||||
|
Proposed | Accepted | Superseded by ADR-XXXX
|
||||||
|
|
||||||
|
## Context
|
||||||
|
|
||||||
|
What forces are in play? What problem or constraint triggered the decision?
|
||||||
|
Keep it to the facts that actually shaped the choice.
|
||||||
|
|
||||||
|
## Options
|
||||||
|
|
||||||
|
| Option | Pros | Cons |
|
||||||
|
| --- | --- | --- |
|
||||||
|
| A — … | … | … |
|
||||||
|
| B — … | … | … |
|
||||||
|
|
||||||
|
(Drop this section when only one option was ever real.)
|
||||||
|
|
||||||
|
## Decision
|
||||||
|
|
||||||
|
The choice, stated plainly. One paragraph.
|
||||||
|
|
||||||
|
## Consequences
|
||||||
|
|
||||||
|
- What becomes easier.
|
||||||
|
- What becomes harder or is now forbidden.
|
||||||
|
- Follow-ups, links to affected docs.
|
||||||
@@ -63,6 +63,14 @@ resource "kubernetes_stateful_set" "drone" {
|
|||||||
name = "DATA_DIR"
|
name = "DATA_DIR"
|
||||||
value = "/data"
|
value = "/data"
|
||||||
}
|
}
|
||||||
|
env {
|
||||||
|
name = "KEEP_ALIVE"
|
||||||
|
value = "1"
|
||||||
|
}
|
||||||
|
env {
|
||||||
|
name = "IMU_HZ"
|
||||||
|
value = "20"
|
||||||
|
}
|
||||||
|
|
||||||
volume_mount {
|
volume_mount {
|
||||||
name = "lake"
|
name = "lake"
|
||||||
@@ -118,9 +126,21 @@ resource "kubernetes_deployment" "exporter" {
|
|||||||
name = "EXPORTER_PORT"
|
name = "EXPORTER_PORT"
|
||||||
value = "9105"
|
value = "9105"
|
||||||
}
|
}
|
||||||
|
env {
|
||||||
|
name = "SCAN_INTERVAL_S"
|
||||||
|
value = "30"
|
||||||
|
}
|
||||||
|
env {
|
||||||
|
name = "METRIC_FLIGHT_WINDOW"
|
||||||
|
value = "5"
|
||||||
|
}
|
||||||
port {
|
port {
|
||||||
container_port = 9105
|
container_port = 9105
|
||||||
}
|
}
|
||||||
|
resources {
|
||||||
|
requests = { cpu = "50m", memory = "64Mi" }
|
||||||
|
limits = { cpu = "500m", memory = "256Mi" }
|
||||||
|
}
|
||||||
volume_mount {
|
volume_mount {
|
||||||
name = "lake"
|
name = "lake"
|
||||||
mount_path = "/data"
|
mount_path = "/data"
|
||||||
@@ -182,9 +202,25 @@ resource "kubernetes_deployment" "explorer" {
|
|||||||
name = "EXPLORER_PORT"
|
name = "EXPLORER_PORT"
|
||||||
value = "8088"
|
value = "8088"
|
||||||
}
|
}
|
||||||
|
env {
|
||||||
|
name = "VIEW_REFRESH_S"
|
||||||
|
value = "30"
|
||||||
|
}
|
||||||
|
env {
|
||||||
|
name = "TREE_CACHE_S"
|
||||||
|
value = "15"
|
||||||
|
}
|
||||||
|
env {
|
||||||
|
name = "METRIC_FLIGHT_WINDOW"
|
||||||
|
value = "5"
|
||||||
|
}
|
||||||
port {
|
port {
|
||||||
container_port = 8088
|
container_port = 8088
|
||||||
}
|
}
|
||||||
|
resources {
|
||||||
|
requests = { cpu = "50m", memory = "64Mi" }
|
||||||
|
limits = { cpu = "750m", memory = "256Mi" }
|
||||||
|
}
|
||||||
volume_mount {
|
volume_mount {
|
||||||
name = "lake"
|
name = "lake"
|
||||||
mount_path = "/data"
|
mount_path = "/data"
|
||||||
@@ -226,11 +262,17 @@ resource "kubernetes_config_map" "prometheus" {
|
|||||||
data = {
|
data = {
|
||||||
"prometheus.yml" = <<-EOT
|
"prometheus.yml" = <<-EOT
|
||||||
global:
|
global:
|
||||||
scrape_interval: 5s
|
scrape_interval: 15s
|
||||||
scrape_configs:
|
scrape_configs:
|
||||||
- job_name: swarm
|
- job_name: swarm
|
||||||
static_configs:
|
static_configs:
|
||||||
- targets: ["exporter:9105"]
|
- targets: ["exporter:9105"]
|
||||||
|
- job_name: explorer
|
||||||
|
static_configs:
|
||||||
|
- targets: ["explorer:8088"]
|
||||||
|
- job_name: node
|
||||||
|
static_configs:
|
||||||
|
- targets: ["node-exporter:9100"]
|
||||||
EOT
|
EOT
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -327,7 +369,91 @@ resource "kubernetes_config_map" "grafana_dashboard" {
|
|||||||
namespace = kubernetes_namespace.swarm.metadata[0].name
|
namespace = kubernetes_namespace.swarm.metadata[0].name
|
||||||
}
|
}
|
||||||
data = {
|
data = {
|
||||||
"swarm.json" = file("${path.module}/../../../simulator/monitoring/grafana/dashboards/swarm.json")
|
"swarm.json" = file("${path.module}/../../../simulator/monitoring/grafana/dashboards/swarm.json")
|
||||||
|
"swarm-platform.json" = file("${path.module}/../../../simulator/monitoring/grafana/dashboards/swarm-platform.json")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
# Host metrics for the k3d node (CPU / memory on the platform dashboard)
|
||||||
|
resource "kubernetes_deployment" "node_exporter" {
|
||||||
|
metadata {
|
||||||
|
name = "node-exporter"
|
||||||
|
namespace = kubernetes_namespace.swarm.metadata[0].name
|
||||||
|
labels = local.labels
|
||||||
|
}
|
||||||
|
|
||||||
|
spec {
|
||||||
|
replicas = 1
|
||||||
|
selector {
|
||||||
|
match_labels = { app = "node-exporter" }
|
||||||
|
}
|
||||||
|
template {
|
||||||
|
metadata {
|
||||||
|
labels = merge(local.labels, { app = "node-exporter" })
|
||||||
|
}
|
||||||
|
spec {
|
||||||
|
host_network = true
|
||||||
|
host_pid = true
|
||||||
|
container {
|
||||||
|
name = "node-exporter"
|
||||||
|
image = "prom/node-exporter:v1.8.2"
|
||||||
|
args = [
|
||||||
|
"--path.procfs=/host/proc",
|
||||||
|
"--path.sysfs=/host/sys",
|
||||||
|
"--path.rootfs=/host/root",
|
||||||
|
"--collector.filesystem.mount-points-exclude=^/(sys|proc|dev|host|etc)($$|/)",
|
||||||
|
]
|
||||||
|
port {
|
||||||
|
container_port = 9100
|
||||||
|
}
|
||||||
|
volume_mount {
|
||||||
|
name = "proc"
|
||||||
|
mount_path = "/host/proc"
|
||||||
|
read_only = true
|
||||||
|
}
|
||||||
|
volume_mount {
|
||||||
|
name = "sys"
|
||||||
|
mount_path = "/host/sys"
|
||||||
|
read_only = true
|
||||||
|
}
|
||||||
|
volume_mount {
|
||||||
|
name = "root"
|
||||||
|
mount_path = "/host/root"
|
||||||
|
read_only = true
|
||||||
|
}
|
||||||
|
resources {
|
||||||
|
requests = { cpu = "20m", memory = "32Mi" }
|
||||||
|
limits = { cpu = "200m", memory = "64Mi" }
|
||||||
|
}
|
||||||
|
}
|
||||||
|
volume {
|
||||||
|
name = "proc"
|
||||||
|
host_path { path = "/proc" }
|
||||||
|
}
|
||||||
|
volume {
|
||||||
|
name = "sys"
|
||||||
|
host_path { path = "/sys" }
|
||||||
|
}
|
||||||
|
volume {
|
||||||
|
name = "root"
|
||||||
|
host_path { path = "/" }
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
resource "kubernetes_service" "node_exporter" {
|
||||||
|
metadata {
|
||||||
|
name = "node-exporter"
|
||||||
|
namespace = kubernetes_namespace.swarm.metadata[0].name
|
||||||
|
}
|
||||||
|
spec {
|
||||||
|
selector = { app = "node-exporter" }
|
||||||
|
port {
|
||||||
|
port = 9100
|
||||||
|
target_port = 9100
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -25,7 +25,7 @@ function displayAspect(): number {
|
|||||||
return Math.max(1, Math.min(MAX_ASPECT, q));
|
return Math.max(1, Math.min(MAX_ASPECT, q));
|
||||||
}
|
}
|
||||||
|
|
||||||
const LIVE_POLL_MS = 1000;
|
const LIVE_POLL_MS = 2000;
|
||||||
|
|
||||||
export default function App(): JSX.Element {
|
export default function App(): JSX.Element {
|
||||||
const [droneCount, setDroneCount] = useState(8);
|
const [droneCount, setDroneCount] = useState(8);
|
||||||
|
|||||||
@@ -24,7 +24,8 @@ services:
|
|||||||
command: ["python", "monitoring/exporter.py"]
|
command: ["python", "monitoring/exporter.py"]
|
||||||
environment:
|
environment:
|
||||||
DATA_DIR: /data
|
DATA_DIR: /data
|
||||||
SCAN_INTERVAL_S: "5"
|
SCAN_INTERVAL_S: "30"
|
||||||
|
METRIC_FLIGHT_WINDOW: "5"
|
||||||
volumes:
|
volumes:
|
||||||
- ${SWARM_T1_DIR:-../var/t1}:/data:ro
|
- ${SWARM_T1_DIR:-../var/t1}:/data:ro
|
||||||
- ./monitoring:/app/monitoring:ro
|
- ./monitoring:/app/monitoring:ro
|
||||||
@@ -37,6 +38,9 @@ services:
|
|||||||
command: ["python", "explorer/server.py"]
|
command: ["python", "explorer/server.py"]
|
||||||
environment:
|
environment:
|
||||||
DATA_DIR: /data
|
DATA_DIR: /data
|
||||||
|
VIEW_REFRESH_S: "30"
|
||||||
|
TREE_CACHE_S: "15"
|
||||||
|
METRIC_FLIGHT_WINDOW: "5"
|
||||||
volumes:
|
volumes:
|
||||||
- ${SWARM_T1_DIR:-../var/t1}:/data:ro
|
- ${SWARM_T1_DIR:-../var/t1}:/data:ro
|
||||||
- ./explorer:/app/explorer:ro
|
- ./explorer:/app/explorer:ro
|
||||||
|
|||||||
@@ -1,33 +1,28 @@
|
|||||||
"""Data-plane explorer: a web view over the Hive-partitioned Parquet lake.
|
"""Data-plane explorer: a web view over the Hive-partitioned Parquet lake."""
|
||||||
|
|
||||||
Serves three things:
|
|
||||||
/ single-page UI (partition tree + read-only SQL console)
|
|
||||||
/api/tree partition hierarchy with file counts and bytes, live
|
|
||||||
/api/query gated read-only DuckDB SQL, same statement rules as the
|
|
||||||
peer query channel (SELECT/WITH only, single statement)
|
|
||||||
|
|
||||||
The point is doctrinal, not just convenient: the explorer reuses the exact
|
|
||||||
read-only SQL contract that drones expose to each other, so "looking at the
|
|
||||||
data plane" on the bench exercises the same path a peer would use in flight.
|
|
||||||
"""
|
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import json
|
import json
|
||||||
import os
|
import os
|
||||||
import re
|
import re
|
||||||
|
import sys
|
||||||
|
import threading
|
||||||
|
import time
|
||||||
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
|
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
import duckdb
|
import duckdb
|
||||||
|
|
||||||
|
sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "monitoring"))
|
||||||
|
from lake import flight_window, iter_parquet_files, parquet_reader # noqa: E402
|
||||||
|
|
||||||
DATA_DIR = Path(os.environ.get("DATA_DIR", "./data"))
|
DATA_DIR = Path(os.environ.get("DATA_DIR", "./data"))
|
||||||
PORT = int(os.environ.get("EXPLORER_PORT", "8088"))
|
PORT = int(os.environ.get("EXPLORER_PORT", "8088"))
|
||||||
ROW_LIMIT = int(os.environ.get("ROW_LIMIT", "500"))
|
ROW_LIMIT = int(os.environ.get("ROW_LIMIT", "500"))
|
||||||
|
VIEW_REFRESH_S = float(os.environ.get("VIEW_REFRESH_S", "30"))
|
||||||
|
TREE_CACHE_S = float(os.environ.get("TREE_CACHE_S", "15"))
|
||||||
STATIC_DIR = Path(__file__).parent
|
STATIC_DIR = Path(__file__).parent
|
||||||
|
|
||||||
# Same spirit as the forced-command gate on a real drone: one statement,
|
|
||||||
# must be a read, no statement that could write, configure, or reach out.
|
|
||||||
_ALLOWED_START = re.compile(r"^\s*(SELECT|WITH|DESCRIBE|SUMMARIZE|SHOW)\b", re.IGNORECASE)
|
_ALLOWED_START = re.compile(r"^\s*(SELECT|WITH|DESCRIBE|SUMMARIZE|SHOW)\b", re.IGNORECASE)
|
||||||
_FORBIDDEN = re.compile(
|
_FORBIDDEN = re.compile(
|
||||||
r"\b(INSERT|UPDATE|DELETE|CREATE|DROP|ALTER|ATTACH|DETACH|COPY|EXPORT"
|
r"\b(INSERT|UPDATE|DELETE|CREATE|DROP|ALTER|ATTACH|DETACH|COPY|EXPORT"
|
||||||
@@ -35,9 +30,17 @@ _FORBIDDEN = re.compile(
|
|||||||
re.IGNORECASE,
|
re.IGNORECASE,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
_lock = threading.Lock()
|
||||||
|
_con: duckdb.DuckDBPyConnection | None = None
|
||||||
|
_views_at = 0.0
|
||||||
|
_tree_cache: dict | None = None
|
||||||
|
_tree_at = 0.0
|
||||||
|
_last_tree_s = 0.0
|
||||||
|
_last_query_s = 0.0
|
||||||
|
_query_count = 0
|
||||||
|
|
||||||
|
|
||||||
def gate(sql: str) -> str | None:
|
def gate(sql: str) -> str | None:
|
||||||
"""Return a rejection reason, or None if the statement passes."""
|
|
||||||
stripped = re.sub(r"--[^\n]*|/\*.*?\*/", " ", sql, flags=re.DOTALL).strip().rstrip(";")
|
stripped = re.sub(r"--[^\n]*|/\*.*?\*/", " ", sql, flags=re.DOTALL).strip().rstrip(";")
|
||||||
if not stripped:
|
if not stripped:
|
||||||
return "empty statement"
|
return "empty statement"
|
||||||
@@ -50,27 +53,39 @@ def gate(sql: str) -> str | None:
|
|||||||
return None
|
return None
|
||||||
|
|
||||||
|
|
||||||
def connect() -> duckdb.DuckDBPyConnection:
|
def _refresh_views() -> None:
|
||||||
"""Fresh connection with the three datasets pre-registered as views."""
|
global _con, _views_at
|
||||||
con = duckdb.connect()
|
con = duckdb.connect()
|
||||||
for ds in ("telemetry", "detections", "state"):
|
for ds in ("telemetry", "detections", "state"):
|
||||||
pattern = f"{DATA_DIR}/dataset={ds}/**/*.parquet"
|
reader = parquet_reader(DATA_DIR, ds)
|
||||||
try:
|
try:
|
||||||
con.execute(
|
con.execute(f"CREATE OR REPLACE VIEW {ds} AS SELECT * FROM read_parquet({reader})")
|
||||||
f"CREATE VIEW {ds} AS SELECT * FROM read_parquet("
|
|
||||||
f"'{pattern}', hive_partitioning=true, union_by_name=true)"
|
|
||||||
)
|
|
||||||
except duckdb.Error:
|
except duckdb.Error:
|
||||||
pass # dataset not written yet; view simply won't exist
|
pass
|
||||||
return con
|
with _lock:
|
||||||
|
if _con is not None:
|
||||||
|
_con.close()
|
||||||
|
_con = con
|
||||||
|
_views_at = time.monotonic()
|
||||||
|
|
||||||
|
|
||||||
|
def connect() -> duckdb.DuckDBPyConnection:
|
||||||
|
if _con is None or time.monotonic() - _views_at > VIEW_REFRESH_S:
|
||||||
|
_refresh_views()
|
||||||
|
assert _con is not None
|
||||||
|
return _con
|
||||||
|
|
||||||
|
|
||||||
def tree() -> dict:
|
def tree() -> dict:
|
||||||
"""Partition hierarchy: dataset -> flight -> drone -> leafs, with sizes."""
|
global _tree_cache, _tree_at, _last_tree_s
|
||||||
|
now = time.monotonic()
|
||||||
|
if _tree_cache is not None and now - _tree_at < TREE_CACHE_S:
|
||||||
|
return _tree_cache
|
||||||
|
started = now
|
||||||
root: dict = {}
|
root: dict = {}
|
||||||
total_bytes = 0
|
total_bytes = 0
|
||||||
total_files = 0
|
total_files = 0
|
||||||
for f in sorted(DATA_DIR.rglob("*.parquet")):
|
for f in sorted(iter_parquet_files(DATA_DIR)):
|
||||||
rel = f.relative_to(DATA_DIR)
|
rel = f.relative_to(DATA_DIR)
|
||||||
size = f.stat().st_size
|
size = f.stat().st_size
|
||||||
total_bytes += size
|
total_bytes += size
|
||||||
@@ -80,7 +95,6 @@ def tree() -> dict:
|
|||||||
node = node.setdefault("children", {}).setdefault(part, {})
|
node = node.setdefault("children", {}).setdefault(part, {})
|
||||||
leaf = node.setdefault("children", {}).setdefault(rel.parts[-1], {})
|
leaf = node.setdefault("children", {}).setdefault(rel.parts[-1], {})
|
||||||
leaf["bytes"] = size
|
leaf["bytes"] = size
|
||||||
# roll sizes up the tree
|
|
||||||
node = root
|
node = root
|
||||||
node["bytes"] = node.get("bytes", 0) + size
|
node["bytes"] = node.get("bytes", 0) + size
|
||||||
node["files"] = node.get("files", 0) + 1
|
node["files"] = node.get("files", 0) + 1
|
||||||
@@ -88,27 +102,50 @@ def tree() -> dict:
|
|||||||
node = node["children"][part]
|
node = node["children"][part]
|
||||||
node["bytes"] = node.get("bytes", 0) + size
|
node["bytes"] = node.get("bytes", 0) + size
|
||||||
node["files"] = node.get("files", 0) + 1
|
node["files"] = node.get("files", 0) + 1
|
||||||
return {"tree": root, "total_bytes": total_bytes, "total_files": total_files}
|
_last_tree_s = time.monotonic() - started
|
||||||
|
_tree_cache = {"tree": root, "total_bytes": total_bytes, "total_files": total_files}
|
||||||
|
_tree_at = now
|
||||||
|
return _tree_cache
|
||||||
|
|
||||||
|
|
||||||
def run_query(sql: str) -> dict:
|
def run_query(sql: str) -> dict:
|
||||||
|
global _last_query_s, _query_count
|
||||||
reason = gate(sql)
|
reason = gate(sql)
|
||||||
if reason:
|
if reason:
|
||||||
return {"error": f"rejected by read-only gate: {reason}"}
|
return {"error": f"rejected by read-only gate: {reason}"}
|
||||||
con = connect()
|
started = time.monotonic()
|
||||||
try:
|
with _lock:
|
||||||
cur = con.sql(sql)
|
con = connect()
|
||||||
columns = cur.columns
|
try:
|
||||||
rows = cur.fetchmany(ROW_LIMIT)
|
cur = con.sql(sql)
|
||||||
return {
|
columns = cur.columns
|
||||||
"columns": columns,
|
rows = cur.fetchmany(ROW_LIMIT)
|
||||||
"rows": [[repr(v) if isinstance(v, bytes) else v for v in row] for row in rows],
|
_query_count += 1
|
||||||
"truncated": len(rows) == ROW_LIMIT,
|
_last_query_s = time.monotonic() - started
|
||||||
}
|
return {
|
||||||
except duckdb.Error as exc:
|
"columns": columns,
|
||||||
return {"error": str(exc)}
|
"rows": [[repr(v) if isinstance(v, bytes) else v for v in row] for row in rows],
|
||||||
finally:
|
"truncated": len(rows) == ROW_LIMIT,
|
||||||
con.close()
|
}
|
||||||
|
except duckdb.Error as exc:
|
||||||
|
return {"error": str(exc)}
|
||||||
|
|
||||||
|
|
||||||
|
def metrics_text() -> str:
|
||||||
|
return "\n".join([
|
||||||
|
"# HELP swarm_explorer_query_duration_seconds Wall time of the last SQL query",
|
||||||
|
"# TYPE swarm_explorer_query_duration_seconds gauge",
|
||||||
|
f"swarm_explorer_query_duration_seconds {_last_query_s:.4f}",
|
||||||
|
"# HELP swarm_explorer_tree_duration_seconds Wall time of the last partition tree build",
|
||||||
|
"# TYPE swarm_explorer_tree_duration_seconds gauge",
|
||||||
|
f"swarm_explorer_tree_duration_seconds {_last_tree_s:.4f}",
|
||||||
|
"# HELP swarm_explorer_queries_total Read-only queries served",
|
||||||
|
"# TYPE swarm_explorer_queries_total counter",
|
||||||
|
f"swarm_explorer_queries_total {_query_count}",
|
||||||
|
"# HELP swarm_metric_flight_window Flight partitions in DuckDB views",
|
||||||
|
"# TYPE swarm_metric_flight_window gauge",
|
||||||
|
f"swarm_metric_flight_window {flight_window()}",
|
||||||
|
]) + "\n"
|
||||||
|
|
||||||
|
|
||||||
class Handler(BaseHTTPRequestHandler):
|
class Handler(BaseHTTPRequestHandler):
|
||||||
@@ -116,13 +153,11 @@ class Handler(BaseHTTPRequestHandler):
|
|||||||
self.send_response(code)
|
self.send_response(code)
|
||||||
self.send_header("Content-Type", ctype)
|
self.send_header("Content-Type", ctype)
|
||||||
self.send_header("Content-Length", str(len(body)))
|
self.send_header("Content-Length", str(len(body)))
|
||||||
# Dev CORS: lets the prototype's live mode poll the API from another
|
|
||||||
# origin. Everything behind this is read-only by construction.
|
|
||||||
self.send_header("Access-Control-Allow-Origin", "*")
|
self.send_header("Access-Control-Allow-Origin", "*")
|
||||||
self.end_headers()
|
self.end_headers()
|
||||||
self.wfile.write(body)
|
self.wfile.write(body)
|
||||||
|
|
||||||
def do_OPTIONS(self) -> None: # noqa: N802 — http.server API
|
def do_OPTIONS(self) -> None: # noqa: N802
|
||||||
self.send_response(204)
|
self.send_response(204)
|
||||||
self.send_header("Access-Control-Allow-Origin", "*")
|
self.send_header("Access-Control-Allow-Origin", "*")
|
||||||
self.send_header("Access-Control-Allow-Methods", "GET, POST, OPTIONS")
|
self.send_header("Access-Control-Allow-Methods", "GET, POST, OPTIONS")
|
||||||
@@ -132,15 +167,17 @@ class Handler(BaseHTTPRequestHandler):
|
|||||||
def _json(self, payload: dict, code: int = 200) -> None:
|
def _json(self, payload: dict, code: int = 200) -> None:
|
||||||
self._send(code, json.dumps(payload, default=str).encode(), "application/json")
|
self._send(code, json.dumps(payload, default=str).encode(), "application/json")
|
||||||
|
|
||||||
def do_GET(self) -> None: # noqa: N802 — http.server API
|
def do_GET(self) -> None: # noqa: N802
|
||||||
if self.path in ("/", "/index.html"):
|
if self.path in ("/", "/index.html"):
|
||||||
self._send(200, (STATIC_DIR / "index.html").read_bytes(), "text/html; charset=utf-8")
|
self._send(200, (STATIC_DIR / "index.html").read_bytes(), "text/html; charset=utf-8")
|
||||||
elif self.path == "/api/tree":
|
elif self.path == "/api/tree":
|
||||||
self._json(tree())
|
self._json(tree())
|
||||||
|
elif self.path == "/metrics":
|
||||||
|
self._send(200, metrics_text().encode(), "text/plain; version=0.0.4")
|
||||||
else:
|
else:
|
||||||
self._send(404, b"not found", "text/plain")
|
self._send(404, b"not found", "text/plain")
|
||||||
|
|
||||||
def do_POST(self) -> None: # noqa: N802 — http.server API
|
def do_POST(self) -> None: # noqa: N802
|
||||||
if self.path != "/api/query":
|
if self.path != "/api/query":
|
||||||
self._send(404, b"not found", "text/plain")
|
self._send(404, b"not found", "text/plain")
|
||||||
return
|
return
|
||||||
@@ -156,6 +193,20 @@ class Handler(BaseHTTPRequestHandler):
|
|||||||
pass
|
pass
|
||||||
|
|
||||||
|
|
||||||
|
def _view_loop() -> None:
|
||||||
|
while True:
|
||||||
|
try:
|
||||||
|
_refresh_views()
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
time.sleep(VIEW_REFRESH_S)
|
||||||
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
if __name__ == "__main__":
|
||||||
print(f"data-plane explorer on :{PORT}, reading {DATA_DIR}")
|
_refresh_views()
|
||||||
|
threading.Thread(target=_view_loop, daemon=True).start()
|
||||||
|
print(
|
||||||
|
f"data-plane explorer on :{PORT}, views refresh every {VIEW_REFRESH_S}s, "
|
||||||
|
f"last {flight_window()} flights, reading {DATA_DIR}"
|
||||||
|
)
|
||||||
ThreadingHTTPServer(("", PORT), Handler).serve_forever()
|
ThreadingHTTPServer(("", PORT), Handler).serve_forever()
|
||||||
|
|||||||
@@ -1,13 +1,14 @@
|
|||||||
"""Prometheus exporter over the simulator's Parquet output.
|
"""Prometheus exporter over the simulator's Parquet output.
|
||||||
|
|
||||||
Periodically scans DATA_DIR with DuckDB and exposes fleet statistics as
|
Periodically scans DATA_DIR with DuckDB and exposes fleet statistics as
|
||||||
/metrics. Zero dependencies beyond duckdb: the exposition format is plain
|
/metrics. Scans only the most recent flight partitions by default so CPU
|
||||||
text, served with the standard library HTTP server.
|
stays bounded as the lake grows across pod restarts.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import os
|
import os
|
||||||
|
import sys
|
||||||
import threading
|
import threading
|
||||||
import time
|
import time
|
||||||
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
|
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
|
||||||
@@ -15,27 +16,29 @@ from pathlib import Path
|
|||||||
|
|
||||||
import duckdb
|
import duckdb
|
||||||
|
|
||||||
|
sys.path.insert(0, str(Path(__file__).resolve().parent))
|
||||||
|
from lake import flight_window, iter_parquet_files, parquet_reader # noqa: E402
|
||||||
|
|
||||||
DATA_DIR = Path(os.environ.get("DATA_DIR", "./data"))
|
DATA_DIR = Path(os.environ.get("DATA_DIR", "./data"))
|
||||||
PORT = int(os.environ.get("EXPORTER_PORT", "9105"))
|
PORT = int(os.environ.get("EXPORTER_PORT", "9105"))
|
||||||
SCAN_INTERVAL_S = float(os.environ.get("SCAN_INTERVAL_S", "5"))
|
SCAN_INTERVAL_S = float(os.environ.get("SCAN_INTERVAL_S", "15"))
|
||||||
|
|
||||||
_lock = threading.Lock()
|
_lock = threading.Lock()
|
||||||
_payload = "# swarm exporter starting\n"
|
_payload = "# swarm exporter starting\n"
|
||||||
|
_last_scan_s = 0.0
|
||||||
|
|
||||||
|
|
||||||
def _q(con: duckdb.DuckDBPyConnection, sql: str) -> list[tuple]:
|
def _q(con: duckdb.DuckDBPyConnection, sql: str) -> list[tuple]:
|
||||||
try:
|
try:
|
||||||
return con.sql(sql).fetchall()
|
return con.sql(sql).fetchall()
|
||||||
except duckdb.Error:
|
except duckdb.Error:
|
||||||
return [] # partitions may not exist yet while the swarm warms up
|
return []
|
||||||
|
|
||||||
|
|
||||||
def collect() -> str:
|
def collect() -> str:
|
||||||
|
global _last_scan_s
|
||||||
|
started = time.monotonic()
|
||||||
con = duckdb.connect()
|
con = duckdb.connect()
|
||||||
# union_by_name: sensors have different schemas under one dataset glob
|
|
||||||
glob = lambda ds: ( # noqa: E731
|
|
||||||
f"'{DATA_DIR}/dataset={ds}/**/*.parquet', hive_partitioning=true, union_by_name=true"
|
|
||||||
)
|
|
||||||
lines: list[str] = []
|
lines: list[str] = []
|
||||||
|
|
||||||
def metric(name: str, help_text: str, mtype: str, rows: list[str]) -> None:
|
def metric(name: str, help_text: str, mtype: str, rows: list[str]) -> None:
|
||||||
@@ -43,27 +46,38 @@ def collect() -> str:
|
|||||||
lines.append(f"# TYPE {name} {mtype}")
|
lines.append(f"# TYPE {name} {mtype}")
|
||||||
lines.extend(rows)
|
lines.extend(rows)
|
||||||
|
|
||||||
metric(
|
for ds, name in (
|
||||||
"swarm_rows_total", "Telemetry rows written per drone and sensor", "gauge",
|
("telemetry", "swarm_rows_total"),
|
||||||
[f'swarm_rows_total{{drone="{d}",sensor="{s}"}} {n}'
|
("detections", "swarm_detections_total"),
|
||||||
for d, s, n in _q(con, f"SELECT drone, sensor, count(*) FROM read_parquet({glob('telemetry')}) GROUP BY 1,2")],
|
("state", "swarm_state_frames_total"),
|
||||||
)
|
):
|
||||||
metric(
|
reader = parquet_reader(DATA_DIR, ds)
|
||||||
"swarm_detections_total", "Detection events per drone and class", "gauge",
|
if ds == "telemetry":
|
||||||
[f'swarm_detections_total{{drone="{d}",cls="{c}"}} {n}'
|
metric(
|
||||||
for d, c, n in _q(con, f"SELECT drone, cls, count(*) FROM read_parquet({glob('detections')}) GROUP BY 1,2")],
|
name, f"Rows in dataset={ds} (recent {flight_window()} flights)", "gauge",
|
||||||
)
|
[f'{name}{{drone="{d}",sensor="{s}"}} {n}'
|
||||||
metric(
|
for d, s, n in _q(con, f"SELECT drone, sensor, count(*) FROM read_parquet({reader}) GROUP BY 1,2")],
|
||||||
"swarm_state_frames_total", "Pose broadcast frames per drone and direction", "gauge",
|
)
|
||||||
[f'swarm_state_frames_total{{drone="{d}",direction="{dr}"}} {n}'
|
elif ds == "detections":
|
||||||
for d, dr, n in _q(con, f"SELECT drone, direction, count(*) FROM read_parquet({glob('state')}) GROUP BY 1,2")],
|
metric(
|
||||||
)
|
name, "Detection events per drone and class", "gauge",
|
||||||
|
[f'{name}{{drone="{d}",cls="{c}"}} {n}'
|
||||||
|
for d, c, n in _q(con, f"SELECT drone, cls, count(*) FROM read_parquet({reader}) GROUP BY 1,2")],
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
metric(
|
||||||
|
name, "Pose broadcast frames per drone and direction", "gauge",
|
||||||
|
[f'{name}{{drone="{d}",direction="{dr}"}} {n}'
|
||||||
|
for d, dr, n in _q(con, f"SELECT drone, direction, count(*) FROM read_parquet({reader}) GROUP BY 1,2")],
|
||||||
|
)
|
||||||
|
|
||||||
|
telem = parquet_reader(DATA_DIR, "telemetry")
|
||||||
metric(
|
metric(
|
||||||
"swarm_battery_pct", "Latest battery level per drone", "gauge",
|
"swarm_battery_pct", "Latest battery level per drone", "gauge",
|
||||||
[f'swarm_battery_pct{{drone="{d}"}} {v}'
|
[f'swarm_battery_pct{{drone="{d}"}} {v}'
|
||||||
for d, v in _q(con, f"""
|
for d, v in _q(con, f"""
|
||||||
SELECT drone, arg_max(level_pct, ts_ns)
|
SELECT drone, arg_max(level_pct, ts_ns)
|
||||||
FROM read_parquet({glob('telemetry')})
|
FROM read_parquet({telem})
|
||||||
WHERE sensor='battery' GROUP BY drone""")],
|
WHERE sensor='battery' GROUP BY drone""")],
|
||||||
)
|
)
|
||||||
metric(
|
metric(
|
||||||
@@ -71,26 +85,24 @@ def collect() -> str:
|
|||||||
[f'swarm_rssi_dbm{{drone="{d}",peer="{p}"}} {v}'
|
[f'swarm_rssi_dbm{{drone="{d}",peer="{p}"}} {v}'
|
||||||
for d, p, v in _q(con, f"""
|
for d, p, v in _q(con, f"""
|
||||||
SELECT drone, peer_id, arg_max(rssi_dbm, ts_ns)
|
SELECT drone, peer_id, arg_max(rssi_dbm, ts_ns)
|
||||||
FROM read_parquet({glob('telemetry')})
|
FROM read_parquet({telem})
|
||||||
WHERE sensor='rssi' GROUP BY drone, peer_id""")],
|
|
||||||
)
|
|
||||||
metric(
|
|
||||||
"swarm_peer_distance_m", "Latest inter-drone distance estimate", "gauge",
|
|
||||||
[f'swarm_peer_distance_m{{drone="{d}",peer="{p}"}} {v}'
|
|
||||||
for d, p, v in _q(con, f"""
|
|
||||||
SELECT drone, peer_id, arg_max(distance_m, ts_ns)
|
|
||||||
FROM read_parquet({glob('telemetry')})
|
|
||||||
WHERE sensor='rssi' GROUP BY drone, peer_id""")],
|
WHERE sensor='rssi' GROUP BY drone, peer_id""")],
|
||||||
)
|
)
|
||||||
|
|
||||||
files = list(DATA_DIR.rglob("*.parquet"))
|
files = iter_parquet_files(DATA_DIR)
|
||||||
metric(
|
metric(
|
||||||
"swarm_parquet_bytes", "Bytes on disk per dataset", "gauge",
|
"swarm_parquet_bytes", "Bytes on disk per dataset (recent flights)", "gauge",
|
||||||
[f'swarm_parquet_bytes{{dataset="{ds}"}} {sum(f.stat().st_size for f in files if f"dataset={ds}" in str(f))}'
|
[f'swarm_parquet_bytes{{dataset="{ds}"}} {sum(f.stat().st_size for f in files if f"dataset={ds}" in str(f))}'
|
||||||
for ds in ("telemetry", "detections", "state")],
|
for ds in ("telemetry", "detections", "state")],
|
||||||
)
|
)
|
||||||
metric("swarm_parquet_files", "Parquet files on disk", "gauge",
|
metric("swarm_parquet_files", "Parquet files scanned (recent flights)", "gauge",
|
||||||
[f"swarm_parquet_files {len(files)}"])
|
[f"swarm_parquet_files {len(files)}"])
|
||||||
|
metric("swarm_metric_flight_window", "Flight partitions included per dataset", "gauge",
|
||||||
|
[f"swarm_metric_flight_window {flight_window()}"])
|
||||||
|
|
||||||
|
_last_scan_s = time.monotonic() - started
|
||||||
|
metric("swarm_exporter_scan_duration_seconds", "Wall time of the last metrics scan", "gauge",
|
||||||
|
[f"swarm_exporter_scan_duration_seconds {_last_scan_s:.4f}"])
|
||||||
con.close()
|
con.close()
|
||||||
return "\n".join(lines) + "\n"
|
return "\n".join(lines) + "\n"
|
||||||
|
|
||||||
@@ -101,11 +113,11 @@ def scanner() -> None:
|
|||||||
started = time.monotonic()
|
started = time.monotonic()
|
||||||
try:
|
try:
|
||||||
payload = collect()
|
payload = collect()
|
||||||
except Exception as exc: # keep serving stale metrics over dying
|
except Exception as exc:
|
||||||
payload = f"# collect error: {exc}\n"
|
payload = f"# collect error: {exc}\n"
|
||||||
with _lock:
|
with _lock:
|
||||||
_payload = payload
|
_payload = payload
|
||||||
time.sleep(max(0.5, SCAN_INTERVAL_S - (time.monotonic() - started)))
|
time.sleep(max(1.0, SCAN_INTERVAL_S - (time.monotonic() - started)))
|
||||||
|
|
||||||
|
|
||||||
class Handler(BaseHTTPRequestHandler):
|
class Handler(BaseHTTPRequestHandler):
|
||||||
@@ -123,10 +135,13 @@ class Handler(BaseHTTPRequestHandler):
|
|||||||
self.wfile.write(body)
|
self.wfile.write(body)
|
||||||
|
|
||||||
def log_message(self, *_args: object) -> None:
|
def log_message(self, *_args: object) -> None:
|
||||||
pass # scrapes every few seconds; keep the log quiet
|
pass
|
||||||
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
if __name__ == "__main__":
|
||||||
threading.Thread(target=scanner, daemon=True).start()
|
threading.Thread(target=scanner, daemon=True).start()
|
||||||
print(f"swarm exporter on :{PORT}/metrics, scanning {DATA_DIR} every {SCAN_INTERVAL_S}s")
|
print(
|
||||||
|
f"swarm exporter on :{PORT}/metrics, scanning last {flight_window()} flights "
|
||||||
|
f"every {SCAN_INTERVAL_S}s under {DATA_DIR}"
|
||||||
|
)
|
||||||
ThreadingHTTPServer(("", PORT), Handler).serve_forever()
|
ThreadingHTTPServer(("", PORT), Handler).serve_forever()
|
||||||
|
|||||||
@@ -0,0 +1,90 @@
|
|||||||
|
{
|
||||||
|
"uid": "swarm-platform",
|
||||||
|
"title": "Swarm Platform — CPU & scan health",
|
||||||
|
"tags": ["swarm", "platform"],
|
||||||
|
"timezone": "utc",
|
||||||
|
"schemaVersion": 39,
|
||||||
|
"version": 1,
|
||||||
|
"refresh": "10s",
|
||||||
|
"time": { "from": "now-30m", "to": "now" },
|
||||||
|
"panels": [
|
||||||
|
{
|
||||||
|
"id": 1, "type": "timeseries", "title": "Node CPU %",
|
||||||
|
"gridPos": { "h": 8, "w": 12, "x": 0, "y": 0 },
|
||||||
|
"datasource": { "type": "prometheus", "uid": "swarm-prom" },
|
||||||
|
"targets": [{
|
||||||
|
"expr": "100 - (avg by (instance) (rate(node_cpu_seconds_total{mode=\"idle\"}[1m])) * 100)",
|
||||||
|
"legendFormat": "{{instance}}", "refId": "A"
|
||||||
|
}],
|
||||||
|
"fieldConfig": { "defaults": { "unit": "percent", "min": 0, "max": 100 }, "overrides": [] }
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"id": 2, "type": "timeseries", "title": "Node memory used %",
|
||||||
|
"gridPos": { "h": 8, "w": 12, "x": 12, "y": 0 },
|
||||||
|
"datasource": { "type": "prometheus", "uid": "swarm-prom" },
|
||||||
|
"targets": [{
|
||||||
|
"expr": "(1 - node_memory_MemAvailable_bytes / node_memory_MemTotal_bytes) * 100",
|
||||||
|
"legendFormat": "used", "refId": "A"
|
||||||
|
}],
|
||||||
|
"fieldConfig": { "defaults": { "unit": "percent", "min": 0, "max": 100 }, "overrides": [] }
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"id": 3, "type": "timeseries", "title": "Exporter scan duration",
|
||||||
|
"gridPos": { "h": 7, "w": 8, "x": 0, "y": 8 },
|
||||||
|
"datasource": { "type": "prometheus", "uid": "swarm-prom" },
|
||||||
|
"targets": [{
|
||||||
|
"expr": "swarm_exporter_scan_duration_seconds",
|
||||||
|
"legendFormat": "scan seconds", "refId": "A"
|
||||||
|
}],
|
||||||
|
"fieldConfig": { "defaults": { "unit": "s" }, "overrides": [] }
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"id": 4, "type": "timeseries", "title": "Explorer query duration",
|
||||||
|
"gridPos": { "h": 7, "w": 8, "x": 8, "y": 8 },
|
||||||
|
"datasource": { "type": "prometheus", "uid": "swarm-prom" },
|
||||||
|
"targets": [{
|
||||||
|
"expr": "swarm_explorer_query_duration_seconds",
|
||||||
|
"legendFormat": "last query", "refId": "A"
|
||||||
|
}],
|
||||||
|
"fieldConfig": { "defaults": { "unit": "s" }, "overrides": [] }
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"id": 5, "type": "timeseries", "title": "Explorer tree build duration",
|
||||||
|
"gridPos": { "h": 7, "w": 8, "x": 16, "y": 8 },
|
||||||
|
"datasource": { "type": "prometheus", "uid": "swarm-prom" },
|
||||||
|
"targets": [{
|
||||||
|
"expr": "swarm_explorer_tree_duration_seconds",
|
||||||
|
"legendFormat": "tree seconds", "refId": "A"
|
||||||
|
}],
|
||||||
|
"fieldConfig": { "defaults": { "unit": "s" }, "overrides": [] }
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"id": 6, "type": "stat", "title": "Parquet files scanned",
|
||||||
|
"gridPos": { "h": 5, "w": 6, "x": 0, "y": 15 },
|
||||||
|
"datasource": { "type": "prometheus", "uid": "swarm-prom" },
|
||||||
|
"targets": [{ "expr": "swarm_parquet_files", "instant": true, "refId": "A" }],
|
||||||
|
"options": { "reduceOptions": { "calcs": ["lastNotNull"] } }
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"id": 7, "type": "stat", "title": "Flight window",
|
||||||
|
"gridPos": { "h": 5, "w": 6, "x": 6, "y": 15 },
|
||||||
|
"datasource": { "type": "prometheus", "uid": "swarm-prom" },
|
||||||
|
"targets": [{ "expr": "swarm_metric_flight_window", "instant": true, "refId": "A" }],
|
||||||
|
"options": { "reduceOptions": { "calcs": ["lastNotNull"] } }
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"id": 8, "type": "stat", "title": "Telemetry rows (window)",
|
||||||
|
"gridPos": { "h": 5, "w": 6, "x": 12, "y": 15 },
|
||||||
|
"datasource": { "type": "prometheus", "uid": "swarm-prom" },
|
||||||
|
"targets": [{ "expr": "sum(swarm_rows_total)", "instant": true, "refId": "A" }],
|
||||||
|
"options": { "reduceOptions": { "calcs": ["lastNotNull"] } }
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"id": 9, "type": "stat", "title": "Min fleet battery %",
|
||||||
|
"gridPos": { "h": 5, "w": 6, "x": 18, "y": 15 },
|
||||||
|
"datasource": { "type": "prometheus", "uid": "swarm-prom" },
|
||||||
|
"targets": [{ "expr": "min(swarm_battery_pct)", "instant": true, "refId": "A" }],
|
||||||
|
"fieldConfig": { "defaults": { "unit": "percent", "min": 0, "max": 100 }, "overrides": [] }
|
||||||
|
}
|
||||||
|
]
|
||||||
|
}
|
||||||
@@ -0,0 +1,43 @@
|
|||||||
|
"""Scan the Hive-partitioned lake without re-reading every historical flight."""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import os
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
|
||||||
|
def flight_window() -> int:
|
||||||
|
return max(1, int(os.environ.get("METRIC_FLIGHT_WINDOW", "5")))
|
||||||
|
|
||||||
|
|
||||||
|
def recent_flights(data_dir: Path, dataset: str, limit: int | None = None) -> list[Path]:
|
||||||
|
"""Most recently modified flight= partitions for a dataset."""
|
||||||
|
limit = limit or flight_window()
|
||||||
|
base = data_dir / f"dataset={dataset}"
|
||||||
|
if not base.is_dir():
|
||||||
|
return []
|
||||||
|
flights = [p for p in base.iterdir() if p.is_dir() and p.name.startswith("flight=")]
|
||||||
|
flights.sort(key=lambda p: p.stat().st_mtime, reverse=True)
|
||||||
|
return flights[:limit]
|
||||||
|
|
||||||
|
|
||||||
|
def parquet_reader(data_dir: Path, dataset: str, *, limit: int | None = None) -> str:
|
||||||
|
"""DuckDB read_parquet() source limited to recent flights."""
|
||||||
|
flights = recent_flights(data_dir, dataset, limit)
|
||||||
|
if not flights:
|
||||||
|
path = data_dir / f"dataset={dataset}" / "**" / "*.parquet"
|
||||||
|
return f"'{path}', hive_partitioning=true, union_by_name=true"
|
||||||
|
if len(flights) == 1:
|
||||||
|
return f"'{flights[0]}/**/*.parquet', hive_partitioning=true, union_by_name=true"
|
||||||
|
inner = ", ".join(f"'{f}/**/*.parquet'" for f in flights)
|
||||||
|
return f"[{inner}], hive_partitioning=true, union_by_name=true"
|
||||||
|
|
||||||
|
|
||||||
|
def iter_parquet_files(data_dir: Path, *, flight_limit: int | None = None) -> list[Path]:
|
||||||
|
"""Parquet paths under recent flights only — avoids full-lake rglob."""
|
||||||
|
limit = flight_limit or flight_window()
|
||||||
|
out: list[Path] = []
|
||||||
|
for ds in ("telemetry", "detections", "state"):
|
||||||
|
for flight in recent_flights(data_dir, ds, limit):
|
||||||
|
out.extend(flight.rglob("*.parquet"))
|
||||||
|
return out
|
||||||
@@ -1,7 +1,10 @@
|
|||||||
global:
|
global:
|
||||||
scrape_interval: 5s
|
scrape_interval: 15s
|
||||||
|
|
||||||
scrape_configs:
|
scrape_configs:
|
||||||
- job_name: swarm
|
- job_name: swarm
|
||||||
static_configs:
|
static_configs:
|
||||||
- targets: ["exporter:9105"]
|
- targets: ["exporter:9105"]
|
||||||
|
- job_name: explorer
|
||||||
|
static_configs:
|
||||||
|
- targets: ["explorer:8088"]
|
||||||
|
|||||||
@@ -0,0 +1,40 @@
|
|||||||
|
"""Lake scan helpers — bounded to recent flight partitions."""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
from monitoring.lake import flight_window, iter_parquet_files, parquet_reader, recent_flights
|
||||||
|
|
||||||
|
|
||||||
|
def test_parquet_reader_limits_to_recent_flights(tmp_path: Path) -> None:
|
||||||
|
for i, name in enumerate(("flight=aaa", "flight=bbb", "flight=ccc")):
|
||||||
|
hour = tmp_path / f"dataset=telemetry/{name}/drone=dr-01/sensor=imu/year=2026/month=07/day=08/hour=10"
|
||||||
|
hour.mkdir(parents=True)
|
||||||
|
f = hour / "data.parquet"
|
||||||
|
f.write_bytes(b"x" * (i + 1))
|
||||||
|
# Make later names newer
|
||||||
|
import os
|
||||||
|
import time
|
||||||
|
os.utime(f, (time.time() + i, time.time() + i))
|
||||||
|
|
||||||
|
reader = parquet_reader(tmp_path, "telemetry", limit=1)
|
||||||
|
assert "flight=ccc" in reader
|
||||||
|
assert "flight=bbb" not in reader
|
||||||
|
|
||||||
|
|
||||||
|
def test_iter_parquet_files_skips_old_flights(tmp_path: Path) -> None:
|
||||||
|
old = tmp_path / "dataset=state/flight=old/drone=dr-01/year=2026/month=07/day=08/hour=09"
|
||||||
|
old.mkdir(parents=True)
|
||||||
|
(old / "data.parquet").write_bytes(b"old")
|
||||||
|
new = tmp_path / "dataset=state/flight=new/drone=dr-01/year=2026/month=07/day=08/hour=10"
|
||||||
|
new.mkdir(parents=True)
|
||||||
|
new_file = new / "data.parquet"
|
||||||
|
new_file.write_bytes(b"new")
|
||||||
|
import os
|
||||||
|
import time
|
||||||
|
os.utime(new_file, (time.time() + 10, time.time() + 10))
|
||||||
|
|
||||||
|
files = iter_parquet_files(tmp_path, flight_limit=1)
|
||||||
|
assert len(files) == 1
|
||||||
|
assert "flight=new" in str(files[0])
|
||||||
@@ -2,6 +2,7 @@
|
|||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import os
|
||||||
import random
|
import random
|
||||||
import select
|
import select
|
||||||
import time
|
import time
|
||||||
@@ -108,6 +109,10 @@ def run(cfg: Config) -> None:
|
|||||||
writer.seal()
|
writer.seal()
|
||||||
print(f"[{cfg.drone_id}] done: {frames_sent} state frames sent, "
|
print(f"[{cfg.drone_id}] done: {frames_sent} state frames sent, "
|
||||||
f"{len(peers_seen)} peers seen {sorted(peers_seen)}; sealed to {root}")
|
f"{len(peers_seen)} peers seen {sorted(peers_seen)}; sealed to {root}")
|
||||||
|
if os.environ.get("KEEP_ALIVE", "0") == "1":
|
||||||
|
print(f"[{cfg.drone_id}] KEEP_ALIVE=1 — idle after seal (no pod restart churn)")
|
||||||
|
while True:
|
||||||
|
time.sleep(3600)
|
||||||
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
if __name__ == "__main__":
|
||||||
|
|||||||
Reference in New Issue
Block a user