Files
Logan CusanoandClaude Opus 5 ee633cbe46
Build & Deploy / Build & push images (push) Failing after 42s
Build & Deploy / Deploy to VM (push) Has been skipped
Secure the broker for public exposure: TLS and per-node credentials
Edge nodes are deployed to arbitrary locations by arbitrary people, so the
broker has to be reachable from the internet and secured on its own merits
rather than by a VPN.

Three defects made that impossible. The broker only had a plaintext 1883
listener; every node shared one drb-node password; and the ACL pattern used
%c, the client-supplied client id, so any holder of that shared password
could set client_id to another node and take over its namespace. The comment
claiming this cryptographically prevented cross-node access was wrong and is
gone.

Authentication now uses mosquitto 2.x's built-in dynamic-security plugin on
the stock eclipse-mosquitto image. c2-core administers it over the control
topic, creating each node's client on approval with username=<node_id> and
password=<its node_keys api_key>, attached to a role whose ACL is nodes/%u/#
against the authenticated username. One credential, one revocation point.
An HTTP-callback plugin was implemented first and rejected: that project is
archived upstream, which is not an acceptable dependency on an
internet-facing broker.

Because dynsec state is a second source of truth alongside Firestore,
approve/reissue/delete now write to the broker first and surface a 502
rather than drifting, and c2-core reconciles every approved node into dynsec
on startup.

Adds node self-enrollment (POST /nodes/enroll, GET /nodes/{id}/credentials)
so a new node can obtain its key over HTTPS without an operator handling
secrets by hand. Enrolling an already-approved node_id is refused on the
fleet token alone — otherwise a leaked token plus a guessable id would let
an attacker steal a live node's key before the real node asked for it.
Pickup secrets are stored hashed and returned once, and the endpoint is rate
limited per source IP.

Infrastructure: an 8883 TLS listener fed by Caddy's certificate via a
systemd path unit, a firewall rule for it, and Caddy now 404s /internal/*
so the api vhost cannot proxy internal routes.

Also fixes CORS, which allowed https://app.<domain> while the frontend is
served on the bare domain — every call from the portal would have failed —
and widens the vault gitignore to a glob, since ansible-vault leaves
backup siblings that the exact-name rule left committable.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-16 09:34:44 -04:00

342 lines
16 KiB
Python

"""
Client for mosquitto's built-in dynamic-security plugin.
WHY THIS EXISTS: MQTT-PUBLIC-AUTH-PLAN.md originally specced the
mosquitto-go-auth plugin (HTTP backend). That project was archived by its
maintainer 2025-08-06 ("no more changes") — unacceptable for a broker that's
about to be reachable from the public internet, no way to get a CVE fix.
Replaced with mosquitto 2.x's own `dynamic-security` plugin, which ships in
and is maintained alongside the official eclipse-mosquitto image itself.
HOW DYNSEC WORKS (verified against plugin source on
github.com/eclipse-mosquitto/mosquitto, 2026-08-16 — see citations inline;
NOT verified by running anything, per instruction not to execute/deploy
anything from this machine):
- The broker persists clients/roles/ACLs in a JSON file at
`plugin_opt_config_file` (we point this at /mosquitto/data/, the same
volume `persistence_location` already uses — one durable volume for all
broker state, see docker-compose.yml).
- Admin commands are plain MQTT publishes: JSON `{"commands": [...]}` to
`$CONTROL/dynamic-security/v1` (source: plugin.c,
`mosquitto_callback_register(plg_id, MOSQ_EVT_CONTROL,
dynsec_control_callback, "$CONTROL/dynamic-security/v1", ...)`).
Replies come back on `$CONTROL/dynamic-security/v1/response`
(source: control.c, `#define RESPONSE_TOPIC
"$CONTROL/dynamic-security/v1/response"`).
- Per-command JSON fields (verified against clients.c / roles.c handlers
and the plugin README):
createClient: username, password, clientid, textname, textdescription,
roles: [{rolename, priority}], groups: [...]
modifyClient: same fields, username identifies the existing client
deleteClient: username
createRole: rolename, textname, textdescription,
acls: [{acltype, topic, priority, allow}]
acltype values: publishClientSend, publishClientReceive,
subscribeLiteral, subscribePattern, unsubscribeLiteral,
unsubscribePattern. %u (username) and %c (clientid) are valid
substitutions in `topic` for every type except the two *Literal ones.
- On first boot, if `plugin_opt_config_file` doesn't exist, the plugin
bootstraps itself (config_init.c): reads env var
`MOSQUITTO_DYNSEC_PASSWORD` (or `plugin_opt_password_init_file`) and
creates a client literally named "admin" (hardcoded string, NOT
configurable — verified in config_init.c's `client_add_admin()`) with
three roles: `super-admin` (full pub/sub on `$CONTROL/#` — i.e. this is
what makes a client capable of issuing further dynsec commands, and
it's an ordinary role, nothing hardcoded beyond the initial grant),
`sys-observe` ($SYS/# read-only), `topic-observe` (# read-only, NOT
read-write). If MOSQUITTO_DYNSEC_PASSWORD is set (we always set it),
no `democlient` demo account gets created — that only happens in the
"no password provided, generate one randomly" path.
- This is a genuine backend swap, not just config: `allow_anonymous
false` plus the *absence* of `password_file`/`acl_file` directives
means dynsec is the only auth backend registered — nothing else is
there to conflict with it. (Inferred from plugin architecture — every
mosquitto auth backend, built-in or plugin, registers the same
basic-auth/ACL callback hooks; there's no "layering" mechanism, so
without password_file/acl_file directives there is nothing else to
check credentials or topics.)
WHAT COULD NOT BE VERIFIED (see also the plan doc + final report):
- The exact JSON envelope of a *response* message (only individual
command outcomes were confirmed: `mosquitto_control_command_reply(cmd,
NULL)` for success, `mosquitto_control_command_reply(cmd, "error
string")` for failure — the wrapping object shape, e.g. whether it's
`{"responses": [{"command": ..., "error": ...}]}`, was not directly
read from source). This client parses defensively: it treats ANY
dict containing a non-null "error"/"Error" key anywhere in the
top-level response payload as failure, presence of "already exists" in
that string as an idempotent success, and a response with no such key
within the timeout as success. A response timeout is always a hard
failure (never assumed to mean success).
- "Client already exists" was confirmed verbatim as createClient's
error string; "already exists" for createRole is assumed analogous,
not directly confirmed.
TWO-SOURCES-OF-TRUTH: Firestore's `node_keys` collection is the source of
truth for node credentials (nothing changes there); dynamic-security.json
is a derived cache the broker uses to authenticate. `reconcile_all()`
rebuilds every approved node's dynsec client from Firestore and is called
on every c2-core startup — so a lost/corrupted dynamic-security.json (e.g.
volume wiped) self-heals on the next restart instead of silently locking
out every node. `upsert_node_client()`/`delete_node_client()` are also
called synchronously from routers/nodes.py's approve/reissue/delete
handlers and raise on failure — those endpoints now fail loudly (502)
instead of updating Firestore while dynsec silently didn't get the memo.
"""
import asyncio
import json
import time
import uuid
import paho.mqtt.client as mqtt
from app.config import settings
from app.internal.logger import logger
from app.internal import firestore as fstore
CONTROL_TOPIC = "$CONTROL/dynamic-security/v1"
RESPONSE_TOPIC = "$CONTROL/dynamic-security/v1/response"
_RESPONSE_TIMEOUT_SECONDS = 10
# Role every approved node's dynsec client is attached to. %u = the
# authenticated username (the node_id) — this is the fixed version of the
# old `pattern readwrite nodes/%c/#`, where %c was the client-supplied,
# spoofable client ID.
NODE_ROLE = "node"
# Role c2-core's own login gets, in addition to being handed the plugin's
# built-in `super-admin` role (see grant_c2core_admin()). Mirrors the old
# `topic readwrite #` superuser line.
C2CORE_ROLE = "c2core"
class DynsecError(Exception):
"""A dynsec command was rejected, or no response arrived in time."""
def _run_commands_sync(commands: list[dict], username: str, password: str) -> list[dict]:
"""
Blocking: open a short-lived MQTT connection, publish one or more dynsec
commands, wait for the matching replies, disconnect. Always called via
asyncio.to_thread — see the async wrappers below. A fresh connection per
call (rather than reusing mqtt_handler's long-lived client) keeps this
request/response exchange simple and isolated from that client's
async-callback-driven subscribe state.
"""
responses: list[dict] = []
done = {"got": False, "error": None}
def _on_connect(client, userdata, flags, reason_code, properties):
if reason_code != 0:
done["error"] = f"connect refused: {reason_code}"
done["got"] = True
return
client.subscribe(RESPONSE_TOPIC, qos=1)
client.publish(CONTROL_TOPIC, json.dumps({"commands": commands}), qos=1)
def _on_message(client, userdata, msg):
try:
payload = json.loads(msg.payload.decode())
except Exception:
return
responses.append(payload)
done["got"] = True
client = mqtt.Client(
callback_api_version=mqtt.CallbackAPIVersion.VERSION2,
client_id=f"drb-c2-core-dynsec-{uuid.uuid4().hex[:8]}",
)
client.username_pw_set(username, password)
client.on_connect = _on_connect
client.on_message = _on_message
try:
client.connect(settings.mqtt_broker, settings.mqtt_port, keepalive=30)
except Exception as e:
raise DynsecError(f"could not connect to mosquitto for dynsec command: {e}")
client.loop_start()
deadline = time.monotonic() + _RESPONSE_TIMEOUT_SECONDS
try:
while not done["got"] and time.monotonic() < deadline:
time.sleep(0.05)
finally:
client.loop_stop()
client.disconnect()
if done["error"]:
raise DynsecError(str(done["error"]))
if not responses:
raise DynsecError(
f"no response on {RESPONSE_TOPIC} within {_RESPONSE_TIMEOUT_SECONDS}s for commands: "
f"{[c.get('command') for c in commands]}"
)
return responses
def _check_responses_ok(responses: list[dict], tolerate_already_exists: bool = False) -> None:
"""Raise DynsecError unless every response payload is error-free (or,
when tolerate_already_exists, only contains an 'already exists'-style
error — see the module docstring's "could not verify" note on why this
is a substring match rather than a structured error code check)."""
for payload in responses:
# Defensive: walk the payload looking for any *-cased "error" key
# with a non-empty value, since the exact envelope shape wasn't
# confirmed from source. Covers both a flat {"error": "..."} and a
# {"responses": [{"error": "..."}]}-style wrapper.
errors = _find_error_strings(payload)
for err in errors:
if tolerate_already_exists and "already exist" in err.lower():
continue
raise DynsecError(f"dynsec command failed: {err}")
def _find_error_strings(obj) -> list[str]:
found = []
if isinstance(obj, dict):
for k, v in obj.items():
if k.lower() == "error" and v:
found.append(str(v))
else:
found.extend(_find_error_strings(v))
elif isinstance(obj, list):
for item in obj:
found.extend(_find_error_strings(item))
return found
# ---------------------------------------------------------------------------
# Async wrappers (all real work happens in the thread pool)
# ---------------------------------------------------------------------------
async def _admin_publish(commands: list[dict], tolerate_already_exists: bool = False) -> list[dict]:
if not settings.mqtt_dynsec_admin_pass:
raise DynsecError("MQTT_DYNSEC_ADMIN_PASS / mqtt_dynsec_admin_pass is not configured")
responses = await asyncio.to_thread(
_run_commands_sync, commands, settings.mqtt_dynsec_admin_user, settings.mqtt_dynsec_admin_pass
)
_check_responses_ok(responses, tolerate_already_exists=tolerate_already_exists)
return responses
async def ensure_roles_and_c2core_grant() -> None:
"""
Idempotent, safe to run on every startup:
1. createRole "node" — nodes/%u/# publish+subscribe (both directions)
2. createRole "c2core" — full "#" publish+subscribe, same reach the
old `topic readwrite #` superuser line gave c2-core
3. createClient/modifyClient drb-c2-core (settings.mqtt_user) with
BOTH roles above AND the plugin's built-in "super-admin" role —
i.e. c2-core's existing login is handed the actual dynsec admin
role, not a separate identity, per the design decision.
Runs over the dedicated "admin" bootstrap login (step 3 assigns
super-admin to c2-core's own login for the record / future use, but
THIS module still authenticates its own ongoing calls as "admin" — see
the module docstring for why: it's the one identity guaranteed by
mosquitto's own source to hold super-admin, so control-plane calls
don't depend on step 3's grant having actually landed).
"""
node_acl_types = ["publishClientSend", "publishClientReceive", "subscribePattern", "unsubscribePattern"]
await _admin_publish([{
"command": "createRole",
"rolename": NODE_ROLE,
"textname": "DRB edge node — own namespace only",
"acls": [{"acltype": t, "topic": "nodes/%u/#", "priority": 0, "allow": True} for t in node_acl_types],
}], tolerate_already_exists=True)
c2core_acl_types = ["publishClientSend", "publishClientReceive", "subscribePattern", "unsubscribePattern"]
await _admin_publish([{
"command": "createRole",
"rolename": C2CORE_ROLE,
"textname": "DRB c2-core — full broker access",
"acls": [{"acltype": t, "topic": "#", "priority": 0, "allow": True} for t in c2core_acl_types],
}], tolerate_already_exists=True)
if not settings.mqtt_user or not settings.mqtt_pass:
logger.warning("dynsec: MQTT_USER/MQTT_PASS not configured — skipping c2-core client grant")
return
roles = [{"rolename": C2CORE_ROLE, "priority": 1}, {"rolename": "super-admin", "priority": 2}]
try:
await _admin_publish([{
"command": "createClient",
"username": settings.mqtt_user,
"password": settings.mqtt_pass,
"roles": roles,
}])
logger.info(f"dynsec: created client {settings.mqtt_user!r} with roles {C2CORE_ROLE}, super-admin")
except DynsecError as e:
if "already exist" in str(e).lower():
await _admin_publish([{
"command": "modifyClient",
"username": settings.mqtt_user,
"password": settings.mqtt_pass,
"roles": roles,
}])
logger.info(f"dynsec: updated existing client {settings.mqtt_user!r} with roles {C2CORE_ROLE}, super-admin")
else:
raise
async def upsert_node_client(node_id: str, api_key: str) -> None:
"""Create or update a node's dynsec client — called from
routers/nodes.py approve_node()/reissue_node_key(), and from
reconcile_all() on startup. Raises DynsecError on failure; callers
must not write Firestore as if this succeeded when it didn't."""
try:
await _admin_publish([{
"command": "createClient",
"username": node_id,
"password": api_key,
"roles": [{"rolename": NODE_ROLE, "priority": 1}],
}])
except DynsecError as e:
if "already exist" not in str(e).lower():
raise
await _admin_publish([{
"command": "modifyClient",
"username": node_id,
"password": api_key,
"roles": [{"rolename": NODE_ROLE, "priority": 1}],
}])
async def delete_node_client(node_id: str) -> None:
"""Best-effort: a node that was never enrolled in dynsec (or already
removed) is treated as already-deleted, not an error."""
try:
await _admin_publish([{"command": "deleteClient", "username": node_id}])
except DynsecError as e:
if "not found" not in str(e).lower():
raise
async def reconcile_all() -> None:
"""
Rebuild dynsec state for every approved node from Firestore
(node_keys is the source of truth). Called once at c2-core startup,
after ensure_roles_and_c2core_grant(). Self-heals a lost/corrupted
dynamic-security.json (e.g. volume wiped, or a prior approve/reissue's
dynsec publish silently failed to persist for some other reason) —
without this, a broker restart with an intact Firestore but an empty
dynsec store would lock out every previously-approved node until
someone noticed and manually re-approved each one.
"""
nodes = await fstore.collection_list("nodes", approval_status="approved")
if not nodes:
return
ok, failed = 0, 0
for node in nodes:
node_id = node.get("node_id")
if not node_id:
continue
key_doc = await fstore.doc_get("node_keys", node_id)
if not key_doc or not key_doc.get("api_key"):
logger.warning(f"dynsec reconcile: node {node_id!r} is approved but has no node_keys entry — skipping")
continue
try:
await upsert_node_client(node_id, key_doc["api_key"])
ok += 1
except DynsecError as e:
failed += 1
logger.error(f"dynsec reconcile: failed to sync node {node_id!r}: {e}")
logger.info(f"dynsec reconcile: {ok} node(s) synced, {failed} failed")