The previous commit shipped Firestore rules that reference an org_id claim
nothing issues yet, and an org_id filter nothing writes yet - this is the
commit that makes both real. Backend half of SAAS_PLAN.md B2/B2b/B2c.
Data model: organizations/{org_id} and org_members/{uid} are new
collections (models.py OrganizationRecord/OrgMember). org_id is now an
Optional field on NodeRecord, SystemRecord, CallRecord, IncidentRecord,
AlertRule, and AlertEvent - optional because every existing document
predates it; scripts/backfill_org_id.py (written, not run - it touches
production Firestore and Firebase Auth claims) is what closes that gap
later. plan_id/subscription_status/stripe_* on OrganizationRecord are
deliberately None: no billing or pricing model has been decided, so this is
a seam, not a promise. app/internal/tenancy.py holds FOUNDING_ORG_ID, the
org every pre-tenancy document and every legacy enrollment path resolves
into.
Where org_id comes from, end to end: a customer's node enrolls with a
per-org token (new enrollment_tokens/{token_hash} collection, minted via
POST /org/enrollment-tokens - new routers/org.py) instead of the old
fleet-wide ENROLLMENT_TOKEN, which still works as a fallback that resolves
to FOUNDING_ORG_ID so an already-deployed node's .env doesn't start failing
today. The node's org_id then flows onto every call it produces
(mqtt_handler.py's call_start/call_end, upload.py's /upload handler all
resolve it from the node doc), and onto every incident correlated from
those calls (incident_correlator.py's _create_incident/_create_master_incident).
That last one is the part that isn't just a read filter: _build_context's
`all_active = collection_list("incidents", status="active")` fed every
correlation candidate - fast-path talkgroup match, unit-continuity,
disambiguation - from the entire incidents collection, unscoped. Without
scoping it to the call's own org_id, a call from org A could link into an
incident org B already owns, which is a cross-tenant data merge at
correlation time, not just an over-broad read. Same shape of bug in
alerter.py: rule matching pulled every enabled alert_rule regardless of
org, so org A's keyword rule could fire (and POST org A's Discord webhook)
on org B's radio traffic. Both now resolve org_id from the call doc itself
rather than threading a new parameter through every caller.
Every list/get route gained org scoping via a new resolve_caller_org_id()
helper in internal/auth.py, which handles the three credential shapes those
routes accept (service key, node api_key, Firebase user) uniformly and
returns None (unrestricted) for the service key and platform admins -
preserving today's single-org behaviour exactly while closing the leak for
everyone else: GET /nodes, /systems, /calls, /incidents, /alerts,
/alert-rules. Write routes for nodes/systems (approve, create, delete, etc.)
deliberately stay platform-admin-only for now rather than being loosened to
org-owner/operator - that's a real gap called out in SAAS_PLAN.md 2.4's
"should be" column, but it's a separate authorization redesign the 12-item
build order doesn't actually enumerate, and doing it half-considered here
risked being exactly the "half-applied filter is worse than none" failure
mode the plan warns about. Today's founding org keeps working unchanged;
loosening node/system management to org owners is follow-up work, flagged
rather than guessed at.
Also closed the four spend/access-attack routes SAAS_PLAN.md B2c called out
by file and line: POST /calls/{id}/reprocess is now admin-only (was any
signed-in viewer looping the Whisper+Gemini pipeline for free - DEFERRED.md
had this as a live, independent-of-SaaS exploit) plus a per-call rate
limiter as a second guard; POST /alerts/{id}/acknowledge now checks the
alert's org_id; GET /admin/features moved from require_firebase_token to
require_admin_token; and trips.py's four unauthenticated mutation routes
(create_trip, update_trip_tags, create_event, update_event) are now
restricted to the founding org (or the bot's service key, or a platform
admin) - trips has no org_id of its own and isn't getting one, since
[[trips-feature-intentional]] says it's an internal utility riding along on
this stack, not a tenant-scoped product surface.
New public-but-scoped seam: POST /auth/signup (routers/links.py, alongside
the existing /auth/link* routes) provisions an organizations doc and an
owner org_members doc for a just-created Firebase user, then sets their
org_id/org_role claims - idempotent, so a double-submit doesn't create two
orgs. This is the only route that turns "has a Firebase account" into "can
read anything," which is what the frontend AuthProvider no-claim guard
(next commit) is built around.
Also new: GET/PATCH /org for the organization profile (closes the disabled
"Save changes" button noted in DEFERRED.md - there was no organizations
concept to save into before this), and POST /waitlist (public, source-IP
rate-limited, not coupled to any plan or tier - the commercial model is
still an open decision per SAAS_PLAN.md section 6).
Verified: all touched files py_compile clean; c2-core pytest is 69
passed / 10 failed, matching the documented pre-existing baseline exactly
(DEFERRED.md - mqtt_handler/node_sweeper test-vs-code drift, unrelated to
this change) - no new failures. flake8 --max-line-length=120 shows no new
violations in any touched file (checked each new E501/E221/E30x against
`git diff` to confirm it predates this commit); c2-core has no CI lint gate
regardless (CLAUDE.md - flake8 only runs in Client CI).
No new environment variables. Firestore composite indexes for the queries
this introduces were already shipped in the previous commit
(infra/firestore/firestore.indexes.json).
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
371 lines
16 KiB
Python
371 lines
16 KiB
Python
import asyncio
|
|
import json
|
|
from datetime import datetime, timezone, timedelta
|
|
from typing import Optional
|
|
import paho.mqtt.client as mqtt
|
|
from app.config import settings
|
|
from app.internal.logger import logger
|
|
from app.internal import firestore as fstore
|
|
from app.internal.tenancy import FOUNDING_ORG_ID
|
|
|
|
|
|
class MQTTHandler:
|
|
def __init__(self):
|
|
self._client: Optional[mqtt.Client] = None
|
|
self._loop: Optional[asyncio.AbstractEventLoop] = None
|
|
self._connected = False
|
|
|
|
def _build_client(self) -> mqtt.Client:
|
|
client = mqtt.Client(
|
|
callback_api_version=mqtt.CallbackAPIVersion.VERSION2,
|
|
client_id="drb-c2-core",
|
|
)
|
|
if settings.mqtt_user:
|
|
client.username_pw_set(settings.mqtt_user, settings.mqtt_pass)
|
|
|
|
client.on_connect = self._on_connect
|
|
client.on_disconnect = self._on_disconnect
|
|
client.on_message = self._on_message
|
|
return client
|
|
|
|
def _on_connect(self, client, userdata, flags, reason_code, properties):
|
|
if reason_code == 0:
|
|
self._connected = True
|
|
client.subscribe("nodes/+/checkin", qos=1)
|
|
client.subscribe("nodes/+/status", qos=1)
|
|
client.subscribe("nodes/+/metadata", qos=1)
|
|
# TODO(mqtt-cutover): drop this subscribe once the enrollment/HTTP
|
|
# credentials flow (routers/enrollment.py) is stable in prod and
|
|
# node-26 (the one live node) has been migrated. See
|
|
# MQTT-PUBLIC-AUTH-PLAN.md "Rollout order" step 6.
|
|
client.subscribe("nodes/+/key_request", qos=1)
|
|
logger.info("MQTT connected — subscribed to node topics.")
|
|
else:
|
|
logger.error(f"MQTT connect refused: {reason_code}")
|
|
|
|
def _on_disconnect(self, client, userdata, disconnect_flags, reason_code, properties):
|
|
self._connected = False
|
|
logger.warning(f"MQTT disconnected: {reason_code}")
|
|
|
|
def _on_message(self, client, userdata, msg):
|
|
try:
|
|
payload = json.loads(msg.payload.decode())
|
|
except Exception:
|
|
logger.warning(f"Non-JSON MQTT message on {msg.topic}")
|
|
return
|
|
|
|
asyncio.run_coroutine_threadsafe(
|
|
self._dispatch(msg.topic, payload), self._loop
|
|
)
|
|
|
|
async def _dispatch(self, topic: str, payload: dict):
|
|
parts = topic.split("/")
|
|
# Expected: nodes/{node_id}/{type}
|
|
if len(parts) != 3 or parts[0] != "nodes":
|
|
return
|
|
|
|
node_id = parts[1]
|
|
msg_type = parts[2]
|
|
|
|
try:
|
|
if msg_type == "checkin":
|
|
await self._handle_checkin(node_id, payload)
|
|
elif msg_type == "status":
|
|
await self._handle_status(node_id, payload)
|
|
elif msg_type == "metadata":
|
|
await self._handle_metadata(node_id, payload)
|
|
elif msg_type == "key_request":
|
|
await self._handle_key_request(node_id)
|
|
except Exception as e:
|
|
logger.error(f"MQTT dispatch error [{msg_type}] from {node_id}: {e}")
|
|
|
|
# ------------------------------------------------------------------
|
|
# Checkin — upsert node; flag new unconfigured nodes
|
|
# ------------------------------------------------------------------
|
|
|
|
async def _handle_checkin(self, node_id: str, payload: dict):
|
|
existing = await fstore.doc_get("nodes", node_id)
|
|
now = datetime.now(timezone.utc)
|
|
|
|
if not existing:
|
|
# First time we've seen this node — create it as unconfigured, pending approval.
|
|
# This branch only fires for a node_id that has never gone through
|
|
# POST /nodes/enroll (routers/enrollment.py) — a properly-enrolled
|
|
# node already has a Firestore doc, with its real org_id, by the time
|
|
# its first checkin arrives, so `existing` would be truthy and this
|
|
# branch wouldn't run. What's left is the legacy shared-MQTT-password
|
|
# path (node-26 — see the TODO(mqtt-cutover) notes in this file),
|
|
# which has no enrollment token to resolve org_id from at all.
|
|
# Default it to FOUNDING_ORG_ID, same as enrollment.py's own
|
|
# legacy-token fallback.
|
|
doc = {
|
|
"node_id": node_id,
|
|
"org_id": FOUNDING_ORG_ID,
|
|
"name": payload.get("name", node_id),
|
|
"lat": payload.get("lat", 0.0),
|
|
"lon": payload.get("lon", 0.0),
|
|
"status": "unconfigured",
|
|
"configured": False,
|
|
"last_seen": now.isoformat(),
|
|
"assigned_system_id": None,
|
|
"approval_status": "pending",
|
|
"node_type": payload.get("node_type", "fixed"),
|
|
"enforce_override_timeout": payload.get("enforce_override_timeout", True),
|
|
"is_overridden": False,
|
|
"override_system_id": None,
|
|
"override_timeout_at": None,
|
|
}
|
|
await fstore.doc_set("nodes", node_id, doc, merge=False)
|
|
logger.info(f"New node registered: {node_id} — pending admin approval.")
|
|
else:
|
|
updates = {
|
|
"last_seen": now.isoformat(),
|
|
"name": payload.get("name", existing.get("name", node_id)),
|
|
"lat": payload.get("lat", existing.get("lat", 0.0)),
|
|
"lon": payload.get("lon", existing.get("lon", 0.0)),
|
|
}
|
|
# Update status on checkin (don't clobber an active recording)
|
|
if existing.get("status") not in ("recording",):
|
|
if existing.get("configured"):
|
|
updates["status"] = "online"
|
|
elif existing.get("approval_status") == "approved":
|
|
# Approved but not yet configured — restore reachable status after reboot
|
|
updates["status"] = "unconfigured"
|
|
|
|
node_type = payload.get("node_type", existing.get("node_type", "fixed"))
|
|
enforce_timeout = payload.get("enforce_override_timeout", existing.get("enforce_override_timeout", True))
|
|
is_overridden = payload.get("is_overridden", False)
|
|
override_system_id = payload.get("override_system_id")
|
|
|
|
updates["node_type"] = node_type
|
|
updates["enforce_override_timeout"] = enforce_timeout
|
|
|
|
if node_type == "portable":
|
|
updates["is_overridden"] = False
|
|
updates["override_system_id"] = None
|
|
updates["override_timeout_at"] = None
|
|
else:
|
|
updates["is_overridden"] = is_overridden
|
|
updates["override_system_id"] = override_system_id
|
|
|
|
if is_overridden:
|
|
existing_timeout = existing.get("override_timeout_at")
|
|
existing_override_id = existing.get("override_system_id")
|
|
if enforce_timeout:
|
|
if not existing_timeout or existing_override_id != override_system_id:
|
|
updates["override_timeout_at"] = (now + timedelta(hours=24)).isoformat()
|
|
else:
|
|
updates["override_timeout_at"] = None
|
|
else:
|
|
updates["override_timeout_at"] = None
|
|
|
|
await fstore.doc_update("nodes", node_id, updates)
|
|
|
|
# NOTE: discord_connected in checkins is informational only — do NOT release the
|
|
# token here. The bot watchdog reconnects on transient Discord drops, so a single
|
|
# checkin with discord_connected=False during a brief reconnect window would
|
|
# incorrectly free the token while the bot is still active. Token release is
|
|
# handled exclusively by the discord_leave command and the node offline sweeper.
|
|
|
|
# ------------------------------------------------------------------
|
|
# Status update
|
|
# ------------------------------------------------------------------
|
|
|
|
async def _handle_status(self, node_id: str, payload: dict):
|
|
status = payload.get("status")
|
|
if not status:
|
|
return
|
|
try:
|
|
await fstore.doc_update("nodes", node_id, {
|
|
"status": status,
|
|
"last_seen": datetime.now(timezone.utc).isoformat(),
|
|
})
|
|
except Exception as e:
|
|
if "No document to update" in str(e):
|
|
logger.info(f"Status from deleted/unknown node {node_id} — ignoring (no Firestore doc)")
|
|
else:
|
|
raise
|
|
|
|
# ------------------------------------------------------------------
|
|
# Metadata — call_start / call_end events
|
|
# ------------------------------------------------------------------
|
|
|
|
async def _handle_metadata(self, node_id: str, payload: dict):
|
|
event = payload.get("event")
|
|
if event == "call_start":
|
|
await self._on_call_start(node_id, payload)
|
|
elif event == "call_end":
|
|
await self._on_call_end(node_id, payload)
|
|
|
|
async def _on_call_start(self, node_id: str, payload: dict):
|
|
call_id = payload.get("call_id")
|
|
if not call_id:
|
|
return
|
|
|
|
# Look up assigned system for this node (cached — assignment rarely changes)
|
|
node = await fstore.doc_get_cached("nodes", node_id)
|
|
system_id = node.get("assigned_system_id") if node else None
|
|
# org_id is inherited from the node, not carried in the MQTT payload —
|
|
# this is the load-bearing tenancy stamp (SAAS_PLAN.md B2b): every call
|
|
# and, downstream, every incident correlated from it, traces back to
|
|
# this. None only for a call from a node that predates tenancy and
|
|
# hasn't been through scripts/backfill_org_id.py yet.
|
|
org_id = node.get("org_id") if node else None
|
|
|
|
started_at_raw = payload.get("started_at")
|
|
started_at = (
|
|
datetime.fromisoformat(started_at_raw)
|
|
if started_at_raw
|
|
else datetime.now(timezone.utc)
|
|
)
|
|
|
|
# Prefer the name from OP25 metadata; fall back to the system config
|
|
tgid_name = payload.get("tgid_name") or ""
|
|
if not tgid_name and system_id and payload.get("tgid"):
|
|
system_doc = await fstore.doc_get_cached("systems", system_id)
|
|
if system_doc:
|
|
tgid_int = int(payload["tgid"])
|
|
for tg in system_doc.get("config", {}).get("talkgroups", []):
|
|
if int(tg.get("id", -1)) == tgid_int:
|
|
tgid_name = tg.get("name", "")
|
|
break
|
|
|
|
doc = {
|
|
"call_id": call_id,
|
|
"node_id": node_id,
|
|
"org_id": org_id,
|
|
"system_id": system_id,
|
|
"talkgroup_id": payload.get("tgid"),
|
|
"talkgroup_name": tgid_name,
|
|
"freq": payload.get("freq"),
|
|
"srcaddr": payload.get("srcaddr"),
|
|
"started_at": started_at,
|
|
"ended_at": None,
|
|
"audio_url": None,
|
|
"transcript": None,
|
|
"incident_id": None,
|
|
"location": None,
|
|
"tags": [],
|
|
"status": "active",
|
|
}
|
|
await fstore.doc_set("calls", call_id, doc, merge=False)
|
|
logger.info(f"Call start: {call_id} (node={node_id}, tgid={payload.get('tgid')})")
|
|
|
|
async def _on_call_end(self, node_id: str, payload: dict):
|
|
call_id = payload.get("call_id")
|
|
if not call_id:
|
|
return
|
|
|
|
ended_at_raw = payload.get("ended_at")
|
|
ended_at = (
|
|
datetime.fromisoformat(ended_at_raw)
|
|
if ended_at_raw
|
|
else datetime.now(timezone.utc)
|
|
)
|
|
|
|
updates = {
|
|
"ended_at": ended_at,
|
|
"status": "ended",
|
|
}
|
|
# doc_set below is a merge, so if call_start already wrote org_id this
|
|
# is a no-op write of the same value. But DEFERRED.md notes call_end
|
|
# can in principle arrive before call_start (ordering relies on MQTT
|
|
# preserving per-topic order, which holds in practice but isn't
|
|
# guaranteed) — in that case doc_set would CREATE the calls doc here
|
|
# with no org_id at all unless it's resolved independently.
|
|
node = await fstore.doc_get_cached("nodes", node_id)
|
|
if node and node.get("org_id"):
|
|
updates["org_id"] = node["org_id"]
|
|
if payload.get("audio_url"):
|
|
updates["audio_url"] = payload["audio_url"]
|
|
|
|
await fstore.doc_set("calls", call_id, updates)
|
|
logger.info(f"Call end: {call_id}")
|
|
|
|
# ------------------------------------------------------------------
|
|
# Key request — re-deliver an existing approved key to a node that
|
|
# lost its credentials (e.g. after a directory move / fresh volume)
|
|
# TODO(mqtt-cutover): remove this handler + publish_node_key() below,
|
|
# and the key_request subscribe above, in the separate post-cutover
|
|
# pass called out in MQTT-PUBLIC-AUTH-PLAN.md. Left in place for now so
|
|
# node-26 (currently live, using the shared-password MQTT path) keeps
|
|
# working until the enrollment flow has replaced it in prod.
|
|
# ------------------------------------------------------------------
|
|
|
|
async def _handle_key_request(self, node_id: str):
|
|
key_doc = await fstore.doc_get("node_keys", node_id)
|
|
if not key_doc or not key_doc.get("api_key"):
|
|
logger.warning(f"Key request from {node_id} but no key found in Firestore — node may not be approved yet.")
|
|
return
|
|
self.publish_node_key(node_id, key_doc["api_key"])
|
|
logger.info(f"Re-delivered API key to {node_id} on request.")
|
|
|
|
# ------------------------------------------------------------------
|
|
# Outbound — send a command to a specific node
|
|
# ------------------------------------------------------------------
|
|
|
|
def send_command(self, node_id: str, payload: dict) -> bool:
|
|
topic = f"nodes/{node_id}/commands"
|
|
if self._client and self._connected:
|
|
self._client.publish(topic, json.dumps(payload), qos=1)
|
|
logger.info(f"Command sent to {node_id}: {payload.get('action')}")
|
|
return True
|
|
logger.warning(f"MQTT not connected — could not send command to {node_id}")
|
|
return False
|
|
|
|
def push_config(self, node_id: str, system_config: dict):
|
|
topic = f"nodes/{node_id}/config"
|
|
if self._client and self._connected:
|
|
self._client.publish(topic, json.dumps(system_config), qos=1)
|
|
logger.info(f"Config pushed to {node_id}")
|
|
else:
|
|
logger.warning(f"MQTT not connected — could not push config to {node_id}")
|
|
|
|
def publish_node_key(self, node_id: str, api_key: str):
|
|
"""Publish the provisioned API key to the node (retained so it survives reconnects).
|
|
TODO(mqtt-cutover): dead once nodes.py's callers switch to the HTTP
|
|
credentials poll (routers/enrollment.py) exclusively. See note above
|
|
_handle_key_request."""
|
|
topic = f"nodes/{node_id}/api_key"
|
|
if self._client and self._connected:
|
|
self._client.publish(topic, json.dumps({"api_key": api_key}), qos=2, retain=True)
|
|
logger.info(f"API key provisioned to {node_id}")
|
|
else:
|
|
logger.warning(f"MQTT not connected — could not provision key to {node_id}")
|
|
|
|
# ------------------------------------------------------------------
|
|
# Lifecycle
|
|
# ------------------------------------------------------------------
|
|
|
|
async def connect(self):
|
|
self._loop = asyncio.get_event_loop()
|
|
self._client = self._build_client()
|
|
# Start the paho network loop first so it drives reconnects automatically,
|
|
# then keep attempting the initial TCP connect until it succeeds.
|
|
self._client.loop_start()
|
|
asyncio.create_task(self._connect_with_retry())
|
|
|
|
async def _connect_with_retry(self):
|
|
delay = 5
|
|
logger.info(f"MQTT connecting to {settings.mqtt_broker}:{settings.mqtt_port}")
|
|
while True:
|
|
try:
|
|
self._client.connect(settings.mqtt_broker, settings.mqtt_port, keepalive=60)
|
|
return # paho loop_start + reconnect_delay_set handles the rest
|
|
except Exception as e:
|
|
logger.warning(f"MQTT connect failed ({e}) — retrying in {delay}s")
|
|
await asyncio.sleep(delay)
|
|
delay = min(delay * 2, 60)
|
|
|
|
async def disconnect(self):
|
|
if self._client:
|
|
self._client.loop_stop()
|
|
self._client.disconnect()
|
|
|
|
@property
|
|
def is_connected(self) -> bool:
|
|
return self._connected
|
|
|
|
|
|
mqtt_handler = MQTTHandler()
|