From a3681ea6986e8f66f4019934aaa231baacb598fd Mon Sep 17 00:00:00 2001 From: Logan Cusano Date: Tue, 18 Aug 2026 20:28:37 -0400 Subject: [PATCH] Stamp org_id everywhere and gate every route that leaked across tenants 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 --- drb-c2-core/app/internal/alerter.py | 18 +- drb-c2-core/app/internal/auth.py | 97 +++++++++ .../app/internal/incident_correlator.py | 29 ++- drb-c2-core/app/internal/mqtt_handler.py | 29 ++- drb-c2-core/app/internal/tenancy.py | 22 +++ drb-c2-core/app/main.py | 4 +- drb-c2-core/app/models.py | 50 +++++ drb-c2-core/app/routers/admin.py | 11 +- drb-c2-core/app/routers/alerts.py | 34 +++- drb-c2-core/app/routers/calls.py | 33 +++- drb-c2-core/app/routers/enrollment.py | 42 +++- drb-c2-core/app/routers/incidents.py | 21 +- drb-c2-core/app/routers/links.py | 89 ++++++++- drb-c2-core/app/routers/nodes.py | 19 +- drb-c2-core/app/routers/org.py | 121 ++++++++++++ drb-c2-core/app/routers/systems.py | 30 ++- drb-c2-core/app/routers/trips.py | 30 ++- drb-c2-core/app/routers/upload.py | 12 +- drb-c2-core/app/routers/waitlist.py | 57 ++++++ drb-c2-core/scripts/backfill_org_id.py | 186 ++++++++++++++++++ 20 files changed, 889 insertions(+), 45 deletions(-) create mode 100644 drb-c2-core/app/internal/tenancy.py create mode 100644 drb-c2-core/app/routers/org.py create mode 100644 drb-c2-core/app/routers/waitlist.py create mode 100644 drb-c2-core/scripts/backfill_org_id.py diff --git a/drb-c2-core/app/internal/alerter.py b/drb-c2-core/app/internal/alerter.py index ce447dc..cdbc1e3 100644 --- a/drb-c2-core/app/internal/alerter.py +++ b/drb-c2-core/app/internal/alerter.py @@ -27,7 +27,22 @@ async def check_and_dispatch( Check all enabled alert rules and fire events for any that match this call. """ try: - rules = await fstore.collection_list("alert_rules", enabled=True) + # Scoped to the call's own org — an unscoped query here would let an + # alert rule created by one org fire (and POST its Discord webhook) + # on another org's radio traffic. org_id is resolved from the call + # doc rather than threaded through as a new parameter, since every + # caller of check_and_dispatch already has call_id and the call doc + # is the single source of truth for a call's org once mqtt_handler.py + # / upload.py have stamped it. None only for a call from a node that + # predates tenancy and hasn't been through the backfill script yet — + # such calls fall back to the pre-tenancy behaviour of checking + # every rule regardless of org. + call_doc = await fstore.doc_get("calls", call_id) + org_id = (call_doc or {}).get("org_id") + if org_id is not None: + rules = await fstore.collection_list("alert_rules", enabled=True, org_id=org_id) + else: + rules = await fstore.collection_list("alert_rules", enabled=True) except Exception as e: logger.warning(f"Alerter: could not load rules: {e}") return @@ -42,6 +57,7 @@ async def check_and_dispatch( now = datetime.now(timezone.utc).isoformat() event = { "alert_id": alert_id, + "org_id": org_id, "rule_id": rule.get("rule_id", ""), "rule_name": rule.get("name", ""), "call_id": call_id, diff --git a/drb-c2-core/app/internal/auth.py b/drb-c2-core/app/internal/auth.py index 801667c..375eb74 100644 --- a/drb-c2-core/app/internal/auth.py +++ b/drb-c2-core/app/internal/auth.py @@ -84,6 +84,93 @@ def get_role(decoded: dict) -> str: return role if role in ("admin", "operator", "viewer") else "viewer" +# --------------------------------------------------------------------------- +# Tenancy — org_id / org_role claims, set by POST /auth/signup (routers/links.py) +# --------------------------------------------------------------------------- +# `role` above is platform-level (admin/operator/viewer — unrelated to which +# org a user belongs to). `org_role` is the customer-facing one: "owner" or +# "member" of the org named by the `org_id` claim. See SAAS_PLAN.md B2/B4. + +def get_org_role(decoded: dict) -> Optional[str]: + org_role = decoded.get("org_role") + return org_role if org_role in ("owner", "member") else None + + +def require_org(decoded: dict) -> str: + """Return the caller's org_id claim, or 403 if they don't have one. + + A Firebase token with no org_id claim is a real, valid session (the user + signed in) that is nonetheless provisioned into nothing — see + AuthProvider's no-claim guard (SAAS_PLAN.md B3). Every org-scoped route + depends on this rather than trusting a client-supplied org_id, so a + caller can never read/write outside the org their own token names. + """ + org_id = decoded.get("org_id") + if not org_id: + raise HTTPException(403, "This account is not associated with an organization.") + return org_id + + +def resolve_org_scope(decoded: dict, org_id_override: Optional[str] = None) -> str: + """Return the org_id a request should be scoped to. + + Platform admins (role == "admin") may pass ?org_id= to cross into + another org's data for support/debugging — the one exception to "you can + only ever see your own org's data" called out in SAAS_PLAN.md B2. Every + other caller is locked to their own token's org_id claim regardless of + what (if anything) they pass. + """ + if org_id_override and get_role(decoded) == "admin": + return org_id_override + return require_org(decoded) + + +async def resolve_caller_org_id(decoded: dict) -> Optional[str]: + """ + Resolve the org_id a caller should be scoped to, across every credential + shape this file's dependencies can produce (service key, node api_key, + Firebase user) — a single helper so read routes gated by + require_service_or_firebase_token / require_node_service_or_firebase_token + don't each need their own caller-shape switch. + + Returns None for callers that should see across every org: the internal + service key (the Discord bot — a single fleet-wide principal, see + CLAUDE.md's auth section) and platform admins, matching + require_admin_token's existing "admin sees everything" behaviour. A + route that wants admins scoped too should check get_role() itself rather + than relying on this function to do it. + """ + if decoded.get("service"): + return None + if decoded.get("node"): + # Deferred import — same reasoning as require_node_service_or_firebase_token + # above: app.internal.firestore initialises firebase-admin at import + # time, and auth.py is imported from module scope in the routers. + from app.internal import firestore as fstore + node = await fstore.doc_get_cached("nodes", decoded.get("node_id") or "") + return (node or {}).get("org_id") + if get_role(decoded) == "admin": + return None + return require_org(decoded) + + +async def require_org_owner_token( + credentials: Optional[HTTPAuthorizationCredentials] = Security(_bearer), +) -> dict: + """Verify a Firebase ID token AND require org_role == "owner" (or platform admin). + + Used for org-administrative actions a regular member shouldn't be able to + do on their own org — minting/revoking enrollment tokens, renaming the + org. Platform admins pass through regardless of org_role so support can + act on an org that has no reachable owner. + """ + decoded = await require_firebase_token(credentials) + require_org(decoded) + if get_org_role(decoded) != "owner" and get_role(decoded) != "admin": + raise HTTPException(status_code=403, detail="Organization owner access required.") + return decoded + + async def require_admin_token( credentials: Optional[HTTPAuthorizationCredentials] = Security(_bearer), ) -> dict: @@ -165,3 +252,13 @@ trip_chat_limiter = _RateLimiter(max_calls=20, window_seconds=300) summarize_limiter = _RateLimiter(max_calls=5, window_seconds=600) # vocabulary bootstrap: 2 per system per hour bootstrap_limiter = _RateLimiter(max_calls=2, window_seconds=3600) +# per-call reprocess: 3 per call per 10 minutes — reprocess re-runs the full +# Whisper + Gemini pipeline, which is real spend per call; this is now also +# admin-only (see routers/calls.py) but the limiter stays as a second guard +# against a compromised/careless admin session looping it. Keyed by call_id, +# same pattern as summarize_limiter. +reprocess_limiter = _RateLimiter(max_calls=3, window_seconds=600) +# public waitlist submissions: 5 per source IP per hour — POST /waitlist has +# no auth at all by design (SAAS_PLAN.md B6), so this is the only thing +# standing between it and being spammed. +waitlist_limiter = _RateLimiter(max_calls=5, window_seconds=3600) diff --git a/drb-c2-core/app/internal/incident_correlator.py b/drb-c2-core/app/internal/incident_correlator.py index 32f16f1..cc8baad 100644 --- a/drb-c2-core/app/internal/incident_correlator.py +++ b/drb-c2-core/app/internal/incident_correlator.py @@ -335,10 +335,23 @@ async def _build_context( now = reference_time or datetime.now(timezone.utc) window = timedelta(hours=settings.correlation_window_hours) - all_active = await fstore.collection_list("incidents", status="active") + call_doc = await fstore.doc_get("calls", call_id) or {} + org_id = call_doc.get("org_id") + + # Candidate incidents MUST be scoped to the call's own org — without this, + # a call from org A could match/link into an incident belonging to org B + # (fast-path talkgroup match, unit-continuity, disambiguation all pull + # from all_active/recent below), which is a cross-tenant data merge, not + # just an over-broad read. org_id is None for a call from a node that + # predates tenancy and hasn't been through scripts/backfill_org_id.py yet + # — such calls fall back to the pre-tenancy behaviour of matching across + # the whole collection rather than being unable to correlate at all. + if org_id is not None: + all_active = await fstore.collection_list("incidents", status="active", org_id=org_id) + else: + all_active = await fstore.collection_list("incidents", status="active") recent = [inc for inc in all_active if _within_window_of(inc, now, window)] - call_doc = await fstore.doc_get("calls", call_id) or {} call_embedding = call_doc.get("embedding") call_units = units if units is not None else (call_doc.get("units") or []) call_vehicles = vehicles if vehicles is not None else (call_doc.get("vehicles") or []) @@ -348,7 +361,7 @@ async def _build_context( is_thin_call = not call_units and not call_vehicles and not coords return { - "call_id": call_id, "all_active": all_active, "recent": recent, + "call_id": call_id, "org_id": org_id, "all_active": all_active, "recent": recent, "call_doc": call_doc, "call_embedding": call_embedding, "call_units": call_units, "call_vehicles": call_vehicles, "call_cleared": call_cleared, "call_severity": call_severity, @@ -839,6 +852,7 @@ async def _apply_decision(decision: dict, ctx: dict) -> Optional[str]: return None call_id = ctx["call_id"] + org_id = ctx["org_id"] talkgroup_id = ctx["talkgroup_id"] talkgroup_name = ctx["talkgroup_name"] system_id = ctx["system_id"] @@ -891,7 +905,7 @@ async def _apply_decision(decision: dict, ctx: dict) -> Optional[str]: # Create the new agency's child incident first incident_id = await _create_incident( - call_id, incident_type, talkgroup_id, talkgroup_name, system_id, + call_id, org_id, incident_type, talkgroup_id, talkgroup_name, system_id, tags, location, location_coords, call_units, call_vehicles, call_embedding, call_severity, now, ) @@ -910,6 +924,7 @@ async def _apply_decision(decision: dict, ctx: dict) -> Optional[str]: master_id = await _create_master_incident( first_child_id=existing_child_id, second_child_id=incident_id, + org_id=org_id, operational_type=incident_type, location=cross_parent.get("location") or location, location_coords=cross_parent.get("location_coords") or coords, @@ -925,7 +940,7 @@ async def _apply_decision(decision: dict, ctx: dict) -> Optional[str]: else: # Normal single-agency incident creation incident_id = await _create_incident( - call_id, incident_type, talkgroup_id, talkgroup_name, system_id, + call_id, org_id, incident_type, talkgroup_id, talkgroup_name, system_id, tags, location, location_coords, call_units, call_vehicles, call_embedding, call_severity, now, ) @@ -1346,6 +1361,7 @@ async def _update_incident( async def _create_incident( call_id: str, + org_id: Optional[str], incident_type: str, talkgroup_id: Optional[int], talkgroup_name: Optional[str], @@ -1377,6 +1393,7 @@ async def _create_incident( doc = { "incident_id": incident_id, + "org_id": org_id, "title": title, "incident_type": "master", # structural role; "child" set on demotion "type": incident_type, @@ -1424,6 +1441,7 @@ def _merge_embedding_vecs(inc: dict, call_embedding: list[float]) -> dict: async def _create_master_incident( first_child_id: str, second_child_id: str, + org_id: Optional[str], operational_type: str, location: Optional[str], location_coords: Optional[dict], @@ -1437,6 +1455,7 @@ async def _create_master_incident( master_id = str(uuid.uuid4()) doc = { "incident_id": master_id, + "org_id": org_id, "title": f"Multi-agency {operational_type} incident", "incident_type": "master", "type": operational_type, diff --git a/drb-c2-core/app/internal/mqtt_handler.py b/drb-c2-core/app/internal/mqtt_handler.py index ec2e867..1e01bc1 100644 --- a/drb-c2-core/app/internal/mqtt_handler.py +++ b/drb-c2-core/app/internal/mqtt_handler.py @@ -6,6 +6,7 @@ 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: @@ -87,9 +88,19 @@ class MQTTHandler: now = datetime.now(timezone.utc) if not existing: - # First time we've seen this node — create it as unconfigured, pending approval + # 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), @@ -194,6 +205,12 @@ class MQTTHandler: # 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 = ( @@ -216,6 +233,7 @@ class MQTTHandler: 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, @@ -249,6 +267,15 @@ class MQTTHandler: "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"] diff --git a/drb-c2-core/app/internal/tenancy.py b/drb-c2-core/app/internal/tenancy.py new file mode 100644 index 0000000..c331f9f --- /dev/null +++ b/drb-c2-core/app/internal/tenancy.py @@ -0,0 +1,22 @@ +""" +Shared tenancy constants used across auth, enrollment, org provisioning, +trip gating, and scripts/backfill_org_id.py. + +FOUNDING_ORG_ID is the org every pre-tenancy document (nodes, systems, +calls, incidents, alert_rules created before this pass) gets stamped with by +scripts/backfill_org_id.py, and the org the legacy fleet-wide +settings.enrollment_token still resolves to in routers/enrollment.py so an +already-deployed field node doesn't break the day these rules deploy — see +that file's enroll_node() for the fallback path. + +NOTE ON MODEL: per the owner's correction mid-build, DRB's access model is +participation-based (you run a node feeding the network, you get access to +the network's data), not per-seat SaaS — org_role is "owner"/"member" with +no tier axis, and organizations.plan_id/seat_limit/node_limit are inert +placeholders (see models.py OrganizationRecord) until a business-strategy +pass defines what, if anything, gates on them. FOUNDING_ORG_ID plays no +special role in that model beyond being the backfill target and the legacy +token's org — it is not a "free tier" or a privileged org in code. +""" + +FOUNDING_ORG_ID = "founding" diff --git a/drb-c2-core/app/main.py b/drb-c2-core/app/main.py index d6b1e5e..1f0868f 100644 --- a/drb-c2-core/app/main.py +++ b/drb-c2-core/app/main.py @@ -15,7 +15,7 @@ from app.internal.auth import ( require_node_service_or_firebase_token, ) from app.routers import nodes, systems, calls, upload, tokens, incidents, alerts, admin, trips, places, links, users -from app.routers import enrollment, media +from app.routers import enrollment, media, org, waitlist from app.internal import dynsec from app.internal import firestore as fstore @@ -101,6 +101,8 @@ app.include_router(admin.router) # auth is per-endpoint (read: firebase, wri app.include_router(users.router) # auth: admin only app.include_router(links.router) # auth is per-endpoint (generate: firebase, resolve: service key) app.include_router(enrollment.router) # public; auth is the enrollment/pickup-secret tokens, checked inline +app.include_router(org.router) # auth is per-endpoint (read: firebase, write: org owner) +app.include_router(waitlist.router) # public — no auth, source-IP rate limited inline # public by necessity — an