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>
141 lines
4.9 KiB
Python
141 lines
4.9 KiB
Python
"""
|
|
Alert dispatch engine.
|
|
|
|
Loads enabled alert rules from Firestore and checks each one against the call's
|
|
talkgroup ID, tags, and transcript. On a match:
|
|
1. Creates an AlertEvent document in Firestore.
|
|
2. Optionally POSTs a Discord webhook message if the rule has one configured.
|
|
|
|
Never raises — failures are logged as warnings so the pipeline always completes.
|
|
"""
|
|
import uuid
|
|
from datetime import datetime, timezone
|
|
from typing import Optional
|
|
from app.internal.logger import logger
|
|
from app.internal import firestore as fstore
|
|
|
|
|
|
async def check_and_dispatch(
|
|
call_id: str,
|
|
node_id: str,
|
|
talkgroup_id: Optional[int],
|
|
talkgroup_name: Optional[str],
|
|
tags: list[str],
|
|
transcript: Optional[str],
|
|
) -> None:
|
|
"""
|
|
Check all enabled alert rules and fire events for any that match this call.
|
|
"""
|
|
try:
|
|
# 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
|
|
|
|
for rule in rules:
|
|
matched_keywords = _match_rule(rule, talkgroup_id, tags, transcript)
|
|
if not matched_keywords:
|
|
continue
|
|
|
|
alert_id = str(uuid.uuid4())
|
|
snippet = _snippet(transcript)
|
|
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,
|
|
"node_id": node_id,
|
|
"talkgroup_id": talkgroup_id,
|
|
"talkgroup_name": talkgroup_name or "",
|
|
"matched_keywords": matched_keywords,
|
|
"transcript_snippet": snippet,
|
|
"triggered_at": now,
|
|
"acknowledged": False,
|
|
}
|
|
|
|
try:
|
|
await fstore.doc_set("alert_events", alert_id, event, merge=False)
|
|
logger.info(
|
|
f"Alert fired: rule='{rule.get('name')}' call={call_id} "
|
|
f"keywords={matched_keywords}"
|
|
)
|
|
except Exception as e:
|
|
logger.warning(f"Alerter: could not save alert event: {e}")
|
|
continue
|
|
|
|
webhook_url = rule.get("discord_webhook")
|
|
if webhook_url:
|
|
await _post_webhook(webhook_url, rule.get("name", ""), talkgroup_name, matched_keywords, snippet)
|
|
|
|
|
|
def _match_rule(
|
|
rule: dict,
|
|
talkgroup_id: Optional[int],
|
|
tags: list[str],
|
|
transcript: Optional[str],
|
|
) -> list[str]:
|
|
"""Return list of matched keywords/reasons, or empty list if no match."""
|
|
matched: list[str] = []
|
|
|
|
# Talkgroup ID match
|
|
rule_tg_ids = rule.get("talkgroup_ids", [])
|
|
if rule_tg_ids and talkgroup_id is not None and talkgroup_id in rule_tg_ids:
|
|
matched.append(f"talkgroup:{talkgroup_id}")
|
|
|
|
# Keyword match against tags + transcript
|
|
rule_keywords = [kw.lower() for kw in rule.get("keywords", [])]
|
|
for kw in rule_keywords:
|
|
if kw in tags:
|
|
matched.append(kw)
|
|
elif transcript and kw in transcript.lower():
|
|
matched.append(kw)
|
|
|
|
return matched
|
|
|
|
|
|
def _snippet(transcript: Optional[str], max_len: int = 200) -> Optional[str]:
|
|
if not transcript:
|
|
return None
|
|
return transcript[:max_len] + ("…" if len(transcript) > max_len else "")
|
|
|
|
|
|
async def _post_webhook(
|
|
url: str,
|
|
rule_name: str,
|
|
talkgroup_name: Optional[str],
|
|
matched_keywords: list[str],
|
|
snippet: Optional[str],
|
|
) -> None:
|
|
try:
|
|
import httpx
|
|
tg_label = talkgroup_name or "Unknown"
|
|
kw_str = ", ".join(matched_keywords)
|
|
body = (
|
|
f"**Alert: {rule_name}**\n"
|
|
f"Talkgroup: {tg_label}\n"
|
|
f"Matched: {kw_str}"
|
|
)
|
|
if snippet:
|
|
body += f"\n> {snippet}"
|
|
async with httpx.AsyncClient(timeout=5.0) as client:
|
|
await client.post(url, json={"content": body})
|
|
except Exception as e:
|
|
logger.warning(f"Alerter: Discord webhook POST failed: {e}")
|