Compare commits
11
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
cea094d66b | ||
|
|
01c146e21e | ||
|
|
d60fef67ad | ||
|
|
fe643924c7 | ||
|
|
bccb3e0316 | ||
|
|
1a631d65d0 | ||
|
|
3a944f35c1 | ||
|
|
0712e7a437 | ||
|
|
a739fa64f0 | ||
|
|
7189ba03e4 | ||
|
|
ef1e3d7f9d |
@@ -63,6 +63,7 @@ jobs:
|
||||
NEXT_PUBLIC_FIREBASE_MESSAGING_SENDER_ID=${{ secrets.FIREBASE_MESSAGING_SENDER_ID }}
|
||||
NEXT_PUBLIC_FIREBASE_APP_ID=${{ secrets.FIREBASE_APP_ID }}
|
||||
NEXT_PUBLIC_FIRESTORE_DATABASE=${{ secrets.FIRESTORE_DATABASE }}
|
||||
NEXT_PUBLIC_MAP_TILE_URL=https://tile.openstreetmap.org/{z}/{x}/{y}.png
|
||||
|
||||
deploy:
|
||||
name: Deploy to VM
|
||||
|
||||
@@ -33,6 +33,13 @@ SUMMARY_INTERVAL_MINUTES=15
|
||||
CORRELATION_WINDOW_HOURS=4
|
||||
EMBEDDING_SIMILARITY_THRESHOLD=0.82
|
||||
|
||||
# Browser origins allowed to call this API cross-origin (JSON list). The only
|
||||
# browser caller is the frontend's Archive page (GET /calls/search). Set this
|
||||
# to the exact origin the frontend is served from — scheme + host, no path.
|
||||
# Defaults to https://drb.cusano.net. A "*" entry works for local dev but is
|
||||
# logged as a probable misconfiguration and never gets a credentialed response.
|
||||
CORS_ORIGINS=["https://drb.cusano.net"]
|
||||
|
||||
# Fleet-wide token edge nodes present as X-Enrollment-Token on first boot
|
||||
# (POST /nodes/enroll). Shared across every node — NOT a per-node secret.
|
||||
# Generate with: openssl rand -hex 32
|
||||
|
||||
@@ -180,16 +180,18 @@ class Settings(BaseSettings):
|
||||
# between genuinely separate transmissions on a busy dispatch channel.
|
||||
duplicate_window_seconds: int = 10
|
||||
|
||||
# CORS — set to your frontend origin(s) in production, e.g. ["https://app.example.com"]
|
||||
# Defaults to "*" for local development only.
|
||||
# Browser origins allowed to call this API cross-origin. The only browser
|
||||
# caller is the frontend's Archive page (GET /calls/search) — every other
|
||||
# page reads Firestore directly. The frontend is served on the BARE domain
|
||||
# (see infra Caddyfile.j2 — only drb. and api. have DNS records), so the
|
||||
# default is that origin, not app.<domain>. Override via CORS_ORIGINS (JSON
|
||||
# list) if the frontend ever moves; keep infra/.../c2-core.env.j2 in sync.
|
||||
#
|
||||
# Leaving this as "*" is not merely permissive: main.py turns OFF
|
||||
# allow_credentials when it sees a wildcard, because Starlette would
|
||||
# otherwise reflect each caller's origin back WITH
|
||||
# Access-Control-Allow-Credentials. So a production deployment that
|
||||
# forgets to set this gets a loud ERROR at startup and loses credentialed
|
||||
# cross-origin requests, rather than silently accepting every origin.
|
||||
cors_origins: list[str] = ["*"]
|
||||
# A "*" entry here still works for local dev but is refused a credentialed
|
||||
# response: main.py never enables allow_credentials (auth is a Bearer
|
||||
# header, not a cookie), and it logs a loud ERROR when it sees a wildcard
|
||||
# in a deployment so a forgotten override is visible.
|
||||
cors_origins: list[str] = ["https://drb.cusano.net"]
|
||||
|
||||
# Discord webhook URL that app/internal/ai_health.py posts to when an AI
|
||||
# tier (transcription/correlation) transitions into or out of degraded
|
||||
|
||||
@@ -108,16 +108,31 @@ _ROAD_RE = re.compile(
|
||||
)
|
||||
|
||||
|
||||
# Street-type synonyms collapsed to one token so "Mohegan Park Avenue" and
|
||||
# "Mohegan Park Ave" produce the same road id (server-26#115 — that one
|
||||
# difference was splitting a car-alarm incident into two).
|
||||
_ROAD_SUFFIX_CANON = {
|
||||
"avenue": "ave", "street": "st", "road": "rd", "drive": "dr",
|
||||
"boulevard": "blvd", "lane": "ln", "court": "ct", "place": "pl",
|
||||
"highway": "hwy", "parkway": "pkwy",
|
||||
}
|
||||
|
||||
|
||||
def _extract_road_ids(text: str) -> set[str]:
|
||||
"""
|
||||
Extract normalised road/route identifiers from a location string.
|
||||
e.g. "suspect east on Route 202" → {"route 202"}
|
||||
"at Main Street and Oak Ave" → {"main street", "oak ave"}
|
||||
"at Main Street and Oak Ave" → {"main st", "oak ave"}
|
||||
"""
|
||||
return {
|
||||
re.sub(r"[\s.\-]+", " ", m.group().lower()).strip()
|
||||
for m in _ROAD_RE.finditer(text)
|
||||
}
|
||||
ids: set[str] = set()
|
||||
for m in _ROAD_RE.finditer(text):
|
||||
key = re.sub(r"[\s.\-]+", " ", m.group().lower()).strip()
|
||||
parts = key.split()
|
||||
if parts and parts[-1] in _ROAD_SUFFIX_CANON:
|
||||
parts[-1] = _ROAD_SUFFIX_CANON[parts[-1]]
|
||||
key = " ".join(parts)
|
||||
ids.add(key)
|
||||
return ids
|
||||
|
||||
|
||||
def _location_mentions_road_overlap(new_location: str, inc_mentions: list[str]) -> bool:
|
||||
@@ -670,6 +685,7 @@ async def correlate_call(
|
||||
reassignment: bool = False,
|
||||
embedding: Optional[list] = None,
|
||||
severity: Optional[str] = None,
|
||||
transcript: Optional[str] = None,
|
||||
) -> Optional[str]:
|
||||
"""
|
||||
Link call_id to an existing incident or create a new one.
|
||||
@@ -686,7 +702,7 @@ async def correlate_call(
|
||||
system_id=system_id, talkgroup_id=talkgroup_id, talkgroup_name=talkgroup_name,
|
||||
tags=tags, incident_type=incident_type, location=location,
|
||||
reassignment=reassignment, create_if_new=create_if_new,
|
||||
embedding=embedding, severity=severity,
|
||||
embedding=embedding, severity=severity, transcript=transcript,
|
||||
)
|
||||
decision = _run_decision(ctx)
|
||||
return await _apply_and_log(decision, ctx)
|
||||
@@ -710,6 +726,7 @@ async def preview_correlation(
|
||||
reassignment: bool = False,
|
||||
embedding: Optional[list] = None,
|
||||
severity: Optional[str] = None,
|
||||
transcript: Optional[str] = None,
|
||||
) -> dict:
|
||||
"""
|
||||
Run the rules engine and return the decision WITHOUT committing to Firestore.
|
||||
@@ -730,7 +747,7 @@ async def preview_correlation(
|
||||
system_id=system_id, talkgroup_id=talkgroup_id, talkgroup_name=talkgroup_name,
|
||||
tags=tags, incident_type=incident_type, location=location,
|
||||
reassignment=reassignment, create_if_new=create_if_new,
|
||||
embedding=embedding, severity=severity,
|
||||
embedding=embedding, severity=severity, transcript=transcript,
|
||||
)
|
||||
decision = _run_decision(ctx)
|
||||
return {"decision": decision, "ctx": ctx}
|
||||
@@ -765,6 +782,7 @@ async def _build_context(
|
||||
create_if_new: bool,
|
||||
embedding: Optional[list] = None,
|
||||
severity: Optional[str] = None,
|
||||
transcript: Optional[str] = None,
|
||||
) -> dict:
|
||||
now = reference_time or datetime.now(timezone.utc)
|
||||
window = timedelta(hours=settings.correlation_window_hours)
|
||||
@@ -804,6 +822,13 @@ async def _build_context(
|
||||
call_vehicles = vehicles if vehicles is not None else (call_doc.get("vehicles") or [])
|
||||
call_cleared = cleared_units if cleared_units is not None else (call_doc.get("cleared_units") or [])
|
||||
call_severity = severity or "routine"
|
||||
# The transcript the LLM correlation tier reasons over. Prefer the SCENE's
|
||||
# own words (server-26#102) — passed by upload.py's scene loop — and fall
|
||||
# back to the call doc only when no scene text was supplied (the
|
||||
# recorrelation sweep, and single-scene calls where the two are identical).
|
||||
# Without this, every non-primary scene of a multi-scene call was judged by
|
||||
# the LLM against a transcript containing the OTHER scenes.
|
||||
scene_transcript = transcript or call_doc.get("transcript_corrected") or call_doc.get("transcript")
|
||||
# A string that is not a place is not a location anywhere downstream — not
|
||||
# in the fit tests, not in the thin-call test, not in the LLM prompt, and
|
||||
# not on the incident. Its coordinates go with it: coords are geocoded
|
||||
@@ -826,6 +851,7 @@ async def _build_context(
|
||||
return {
|
||||
"call_id": call_id, "org_id": org_id, "all_active": all_active, "recent": recent,
|
||||
"call_doc": call_doc, "call_embedding": call_embedding,
|
||||
"scene_transcript": scene_transcript,
|
||||
"call_units": call_units, "call_vehicles": call_vehicles,
|
||||
"call_cleared": call_cleared, "call_severity": call_severity,
|
||||
"coords": coords, "is_thin_call": is_thin_call, "now": now,
|
||||
|
||||
@@ -172,7 +172,7 @@ async def extract_scenes(
|
||||
|
||||
Each scene dict contains:
|
||||
tags, incident_type, location, location_coords, resolved,
|
||||
severity, vehicles, units, transcript_corrected,
|
||||
severity, vehicles, units, transcript, transcript_corrected,
|
||||
segment_indices, embedding
|
||||
|
||||
Side-effect: updates calls/{call_id} in Firestore with merged tags,
|
||||
@@ -337,6 +337,10 @@ async def extract_scenes(
|
||||
)
|
||||
embedding = await asyncio.to_thread(_sync_embed, scene_text)
|
||||
|
||||
scene_transcript = _scene_transcript_text(
|
||||
transcript, segments, segment_indices, transcript_corrected
|
||||
)
|
||||
|
||||
processed.append({
|
||||
"tags": tags,
|
||||
"incident_type": incident_type,
|
||||
@@ -348,6 +352,7 @@ async def extract_scenes(
|
||||
"severity": severity,
|
||||
"resolved": resolved,
|
||||
"reassignment": reassignment,
|
||||
"transcript": scene_transcript,
|
||||
"transcript_corrected": transcript_corrected,
|
||||
"segment_indices": segment_indices,
|
||||
"embedding": embedding,
|
||||
@@ -571,11 +576,49 @@ def _municipality_from_tg(tg_name: Optional[str]) -> Optional[str]:
|
||||
def _build_transcript_block(transcript: str, segments: Optional[list[dict]]) -> str:
|
||||
"""Format transcript as numbered transmissions if segments are available."""
|
||||
if segments and len(segments) > 1:
|
||||
lines = [f"{i+1}. [{s['start']}s] {s['text']}" for i, s in enumerate(segments)]
|
||||
# 0-based labels, matching the prompt's "0-based indices into the
|
||||
# numbered transmissions" — the model echoes these back as
|
||||
# `segment_indices`, which _build_scene_embed_text and the per-scene
|
||||
# `transcript` (server-26#102) then slice with directly.
|
||||
lines = [f"{i}. [{s['start']}s] {s['text']}" for i, s in enumerate(segments)]
|
||||
return f"Transmissions ({len(segments)}):\n" + "\n".join(lines)
|
||||
return f"Transcript:\n{transcript}"
|
||||
|
||||
|
||||
def _scene_transcript_text(
|
||||
transcript: str,
|
||||
segments: Optional[list[dict]],
|
||||
segment_indices: Optional[list[int]],
|
||||
transcript_corrected: Optional[str],
|
||||
) -> str:
|
||||
"""
|
||||
This scene's own words, unprefixed — the segments it owns, joined.
|
||||
|
||||
server-26#102: the correlator's LLM tier reads this per scene instead of
|
||||
the call doc's whole-call transcript, so on a multi-scene call scene N is
|
||||
no longer judged against scenes 1..N-1's text.
|
||||
|
||||
Never returns "". Anything that would leave the slice empty — no
|
||||
`segment_indices` (a single-segment call is never numbered by
|
||||
`_build_transcript_block`), or indices that are out of range / not ints —
|
||||
falls back to the whole-call transcript, which for a single-scene call is
|
||||
the same text and for a mis-sliced multi-scene call is at least this
|
||||
call's own words. `_sync_extract`'s prompt documents 0-based indices and
|
||||
`_build_transcript_block` numbers to match, so no base normalisation here.
|
||||
"""
|
||||
if transcript_corrected:
|
||||
return transcript_corrected
|
||||
if segments and segment_indices:
|
||||
joined = " ".join(
|
||||
segments[i]["text"]
|
||||
for i in segment_indices
|
||||
if isinstance(i, int) and 0 <= i < len(segments)
|
||||
)
|
||||
if joined:
|
||||
return joined
|
||||
return transcript
|
||||
|
||||
|
||||
def _build_scene_embed_text(
|
||||
transcript: str,
|
||||
segments: Optional[list[dict]],
|
||||
|
||||
@@ -45,7 +45,18 @@ def _fmt_idle(inc: dict, now: datetime) -> str:
|
||||
|
||||
|
||||
def _inc_summary(inc: dict, now: datetime) -> str:
|
||||
# server-26#115: the model was given no title and no talkgroup, so it
|
||||
# could not tell that "car alarms, Mohegan Park Ave" and "car alarms,
|
||||
# Mohegan Park Avenue" on the same channel were one incident — it defaulted
|
||||
# to "new". Title is the single strongest human-readable signal for "is
|
||||
# this the same event"; talkgroup is what makes same-channel continuation
|
||||
# obvious.
|
||||
parts = [f"id:{inc['incident_id']}", f"type:{inc.get('type') or '?'}"]
|
||||
tgs = inc.get("talkgroup_ids") or []
|
||||
if tgs:
|
||||
parts.append(f"tg:[{', '.join(str(t) for t in tgs[:3])}]")
|
||||
if inc.get("title"):
|
||||
parts.append(f"title:{inc['title']!r}")
|
||||
if inc.get("location"):
|
||||
parts.append(f"loc:{inc['location']}")
|
||||
units = inc.get("units") or []
|
||||
@@ -61,7 +72,13 @@ def _inc_summary(inc: dict, now: datetime) -> str:
|
||||
def _call_block(ctx: dict) -> str:
|
||||
lines = []
|
||||
call_doc = ctx["call_doc"]
|
||||
transcript = call_doc.get("transcript_corrected") or call_doc.get("transcript")
|
||||
# The SCENE's own transcript, resolved in _build_context (server-26#102).
|
||||
# Falls back to the call doc for a ctx built without a scene (tests, sweep).
|
||||
transcript = (
|
||||
ctx.get("scene_transcript")
|
||||
or call_doc.get("transcript_corrected")
|
||||
or call_doc.get("transcript")
|
||||
)
|
||||
if transcript:
|
||||
lines.append(f"Transcript: {transcript[:700]}")
|
||||
if ctx["tags"]:
|
||||
@@ -74,19 +91,50 @@ def _call_block(ctx: dict) -> str:
|
||||
lines.append(f"Units: {ctx['call_units']}")
|
||||
if ctx["call_vehicles"]:
|
||||
lines.append(f"Vehicles: {ctx['call_vehicles']}")
|
||||
if ctx["talkgroup_name"]:
|
||||
lines.append(f"Talkgroup: {ctx['talkgroup_name']}")
|
||||
if ctx["talkgroup_name"] or ctx.get("talkgroup_id") is not None:
|
||||
# Both the name and the id — _inc_summary emits numeric tg ids, so the
|
||||
# id is what makes the "same talkgroup" rule in _RULES evaluable
|
||||
# (server-26#115 review).
|
||||
tgid = ctx.get("talkgroup_id")
|
||||
name = ctx["talkgroup_name"] or "?"
|
||||
lines.append(f"Talkgroup: {name}" + (f" (id {tgid})" if tgid is not None else ""))
|
||||
return "\n".join(lines) if lines else "(no details)"
|
||||
|
||||
|
||||
def _prompt_incidents(recent: list[dict]) -> list[dict]:
|
||||
"""The ≤20 candidates shown to the model, most-recently-active first.
|
||||
|
||||
`ctx["recent"]` is an unordered slice of a Firestore result with no
|
||||
order_by, so a busy 2h window (~40 active incidents) meant the model saw
|
||||
an arbitrary half of the candidates (server-26#115 review). Sorting by
|
||||
updated_at desc also makes each row's `idle:` field monotonic.
|
||||
"""
|
||||
def _key(inc: dict):
|
||||
return str(inc.get("updated_at") or inc.get("started_at") or "")
|
||||
return sorted(recent, key=_key, reverse=True)[:20]
|
||||
|
||||
|
||||
_SCHEMA = '{"action": "link" | "new" | "orphan", "incident_id": "<id_string or null>", "reasoning": "<one sentence>"}'
|
||||
|
||||
_RULES = """
|
||||
Rules:
|
||||
- "link" only with clear positive evidence: same units, same geocoded location, or semantically identical scene on the same talkgroup within the last few minutes.
|
||||
- A call on a DIFFERENT talkgroup than an incident requires unit overlap or geocoded location match — topic similarity alone is not enough.
|
||||
- "new" only if the call has a clear incident_type AND describes a distinct, identifiable scene.
|
||||
- "orphan" when in doubt — conservative is always correct.
|
||||
Rules (this system OVER-SPLITS — a real incident routinely gets shattered into
|
||||
5-10 duplicates. A wrong link is cheap; a duplicate incident is the failure
|
||||
mode. Bias accordingly.):
|
||||
- Prefer "link" when the call plausibly continues a recent incident ON THE SAME
|
||||
TALKGROUP: same or overlapping units, the same or an adjacent location (treat
|
||||
"Ave"/"Avenue", "St"/"Street", "Rd"/"Road" as identical; a house number plus
|
||||
the same street is the same place), the same subject/vehicle/case number, or a
|
||||
follow-up beat ("units clearing", "negative contact", "tow en route", "event
|
||||
number 214-201", a status update) to an incident that is only a few minutes
|
||||
idle. The bar for "link" on the same talkgroup is LOW.
|
||||
- Reserve "new" for a call that clearly describes a DIFFERENT event from every
|
||||
recent incident — a different place, different units, and a different subject,
|
||||
not merely a different transmission about the same job.
|
||||
- "orphan" a call that is not an incident at all: radio checks, roll call,
|
||||
a unit marking on/off duty or 10-8/10-98, mileage/log entries, a bare
|
||||
acknowledgement. Do not open a "new" incident for these.
|
||||
- A call on a DIFFERENT talkgroup than an incident still requires unit overlap
|
||||
or a geocoded/location match — topic similarity alone is not enough there.
|
||||
- Do NOT link just because both calls involve police or both mention a road.
|
||||
"""
|
||||
|
||||
@@ -95,7 +143,7 @@ def _build_decide_prompt(ctx: dict) -> str:
|
||||
now = ctx["now"]
|
||||
recent = ctx["recent"]
|
||||
inc_block = (
|
||||
"\n".join(_inc_summary(inc, now) for inc in recent[:20])
|
||||
"\n".join(_inc_summary(inc, now) for inc in _prompt_incidents(recent))
|
||||
if recent else "(none)"
|
||||
)
|
||||
return (
|
||||
@@ -113,7 +161,7 @@ def _build_tiebreak_prompt(rules_decision: dict, llm_decision: dict, ctx: dict)
|
||||
now = ctx["now"]
|
||||
recent = ctx["recent"]
|
||||
inc_block = (
|
||||
"\n".join(_inc_summary(inc, now) for inc in recent[:20])
|
||||
"\n".join(_inc_summary(inc, now) for inc in _prompt_incidents(recent))
|
||||
if recent else "(none)"
|
||||
)
|
||||
|
||||
|
||||
@@ -108,6 +108,7 @@ async def _recorrelate_orphan(call: dict) -> bool:
|
||||
cleared_units = call.get("cleared_units") or [],
|
||||
embedding = call.get("embedding"),
|
||||
severity = call.get("severity"),
|
||||
transcript = call.get("transcript_corrected") or call.get("transcript"),
|
||||
reference_time = started_at, # anchor window to when the call happened
|
||||
create_if_new = False, # never create — link-only
|
||||
)
|
||||
|
||||
+24
-17
@@ -78,33 +78,40 @@ async def lifespan(app: FastAPI):
|
||||
|
||||
app = FastAPI(title="DRB C2 Core", lifespan=lifespan)
|
||||
|
||||
# "*" plus allow_credentials=True is not the permissive-but-harmless setting it
|
||||
# looks like. Starlette does not refuse the combination -- it reflects the
|
||||
# caller's Origin back and still sends Access-Control-Allow-Credentials: true,
|
||||
# so the effective policy becomes "any origin, with credentials", the opposite
|
||||
# of what a wildcard normally means. Rather than trust every deployment to
|
||||
# remember to override CORS_ORIGINS, make the dangerous pair unrepresentable.
|
||||
# The browser needs CORS to reach this API at all: the frontend's Archive page
|
||||
# calls GET /calls/search with Authorization + Content-Type headers, which
|
||||
# forces a preflight. Without this middleware the OPTIONS gets a bare 405 and
|
||||
# the fetch fails (#110). allow_origins is an explicit list -- never "*" in a
|
||||
# deployment -- so name every host the frontend is served from in CORS_ORIGINS.
|
||||
#
|
||||
# allow_credentials stays False on purpose: auth here is a Bearer header, not a
|
||||
# cookie, so credentialed CORS is never needed, and keeping it False is what
|
||||
# lets an explicit-origin allowlist work without Starlette's "*"-only
|
||||
# restriction. "*" + credentials is the dangerous pair (Starlette reflects the
|
||||
# caller's Origin back WITH Access-Control-Allow-Credentials: true); this code
|
||||
# cannot produce it because credentials are hard-off.
|
||||
def cors_allows_credentials(origins: list[str]) -> bool:
|
||||
"""False when any entry is a wildcard. Extracted so it can be tested
|
||||
without re-importing this module, which drags in every router."""
|
||||
return "*" not in origins
|
||||
"""Always False -- credentialed CORS is never enabled here (Bearer auth,
|
||||
not cookies). Kept as a named predicate so a future edit that wants to
|
||||
turn credentials on has to go through here and confront the "*" case.
|
||||
A wildcard entry would additionally be refused a credentialed response."""
|
||||
return False
|
||||
|
||||
|
||||
_cors_is_wildcard = not cors_allows_credentials(settings.cors_origins)
|
||||
_cors_is_wildcard = "*" in settings.cors_origins
|
||||
if _cors_is_wildcard:
|
||||
logger.error(
|
||||
"CORS_ORIGINS is '*', so credentialed cross-origin requests are being "
|
||||
"DISABLED to avoid reflecting every caller's origin back with "
|
||||
"Access-Control-Allow-Credentials. Set CORS_ORIGINS to your frontend "
|
||||
"origin(s) in production, e.g. [\"https://app.example.com\"]."
|
||||
"CORS_ORIGINS contains '*'. That is fine for local dev but is almost "
|
||||
"certainly a misconfigured deployment -- set CORS_ORIGINS to your "
|
||||
"frontend origin(s), e.g. [\"https://drb.cusano.net\"]."
|
||||
)
|
||||
|
||||
app.add_middleware(
|
||||
CORSMiddleware,
|
||||
allow_origins=settings.cors_origins,
|
||||
allow_methods=["*"],
|
||||
allow_headers=["*"],
|
||||
allow_credentials=not _cors_is_wildcard,
|
||||
allow_methods=["GET", "POST", "PUT", "PATCH", "DELETE", "OPTIONS"],
|
||||
allow_headers=["authorization", "content-type"],
|
||||
allow_credentials=False,
|
||||
)
|
||||
|
||||
app.include_router(nodes.router, dependencies=[Depends(require_service_or_firebase_token)])
|
||||
|
||||
@@ -116,6 +116,7 @@ async def _correlate_with_consensus(
|
||||
reassignment: bool = False,
|
||||
embedding: Optional[list] = None,
|
||||
severity: Optional[str] = None,
|
||||
transcript: Optional[str] = None,
|
||||
) -> Optional[str]:
|
||||
"""
|
||||
Consensus correlator: runs the rules engine and the cheap LLM in sequence.
|
||||
@@ -133,7 +134,7 @@ async def _correlate_with_consensus(
|
||||
tags=tags, incident_type=incident_type, location=location,
|
||||
location_coords=location_coords, units=units, vehicles=vehicles,
|
||||
cleared_units=cleared_units, reassignment=reassignment,
|
||||
embedding=embedding, severity=severity,
|
||||
embedding=embedding, severity=severity, transcript=transcript,
|
||||
)
|
||||
ctx = preview["ctx"]
|
||||
rules_decision = preview["decision"]
|
||||
@@ -226,6 +227,7 @@ async def _run_extraction_pipeline(
|
||||
reassignment=is_reassignment,
|
||||
embedding=scene.get("embedding"),
|
||||
severity=scene.get("severity"),
|
||||
transcript=scene.get("transcript"),
|
||||
)
|
||||
if incident_id and incident_id not in incident_ids:
|
||||
incident_ids.append(incident_id)
|
||||
@@ -343,6 +345,7 @@ async def _run_intelligence_pipeline(
|
||||
reassignment=is_reassignment,
|
||||
embedding=scene.get("embedding"),
|
||||
severity=scene.get("severity"),
|
||||
transcript=scene.get("transcript"),
|
||||
)
|
||||
if incident_id and incident_id not in incident_ids:
|
||||
incident_ids.append(incident_id)
|
||||
|
||||
@@ -0,0 +1,67 @@
|
||||
"""
|
||||
server-26#115 — the tiebreaker manufactured incidents because it was blind to
|
||||
what would tell it two incidents are one.
|
||||
|
||||
Two low-risk supports for the reframed prompt:
|
||||
1. `_extract_road_ids` collapses street-type synonyms, so "Mohegan Park Ave"
|
||||
and "Mohegan Park Avenue" share a road id (they were splitting one
|
||||
car-alarm incident into two).
|
||||
2. `_inc_summary` now carries the incident title and talkgroup, the two
|
||||
signals the model needs to recognise a same-channel continuation.
|
||||
"""
|
||||
from datetime import datetime, timezone
|
||||
|
||||
from app.internal.incident_correlator import (
|
||||
_extract_road_ids, _location_mentions_road_overlap,
|
||||
)
|
||||
from app.internal.llm_correlator import _inc_summary, _prompt_incidents
|
||||
|
||||
NOW = datetime(2026, 9, 7, 8, 0, 0, tzinfo=timezone.utc)
|
||||
|
||||
|
||||
def test_avenue_and_ave_are_the_same_road_id():
|
||||
assert _extract_road_ids("Mohegan Park Avenue") == _extract_road_ids("Mohegan Park Ave")
|
||||
assert _extract_road_ids("191 Broadway Street") == _extract_road_ids("191 Broadway St")
|
||||
assert _extract_road_ids("North State Road") == _extract_road_ids("North State Rd")
|
||||
|
||||
|
||||
def test_road_overlap_matches_across_the_synonym():
|
||||
assert _location_mentions_road_overlap("multiple car alarms Mohegan Park Avenue",
|
||||
["patrol to Mohegan Park Ave"]) is True
|
||||
# still discriminates genuinely different streets
|
||||
assert _location_mentions_road_overlap("Oak Avenue", ["Elm Avenue"]) is False
|
||||
|
||||
|
||||
def test_inc_summary_carries_title_and_talkgroup():
|
||||
s = _inc_summary({
|
||||
"incident_id": "abc123",
|
||||
"type": "police",
|
||||
"talkgroup_ids": [9560],
|
||||
"title": "Nuisance Alarm at Mohegan Park Ave",
|
||||
"location": "Mohegan Park Ave",
|
||||
"units": ["Headquarters"],
|
||||
"tags": ["car-alarm"],
|
||||
"updated_at": NOW.isoformat(),
|
||||
}, NOW)
|
||||
assert "title:'Nuisance Alarm at Mohegan Park Ave'" in s
|
||||
assert "tg:[9560]" in s
|
||||
assert "id:abc123" in s
|
||||
|
||||
|
||||
def test_inc_summary_omits_missing_optional_fields():
|
||||
s = _inc_summary({"incident_id": "x", "updated_at": NOW.isoformat()}, NOW)
|
||||
assert "title:" not in s and "tg:" not in s and "loc:" not in s
|
||||
assert s.startswith("id:x")
|
||||
|
||||
|
||||
def test_prompt_incidents_is_most_recently_active_first_and_capped():
|
||||
recent = [
|
||||
{"incident_id": f"i{n}", "updated_at": f"2026-09-07T0{n}:00:00+00:00"}
|
||||
for n in range(1, 8)
|
||||
]
|
||||
ordered = _prompt_incidents(recent)
|
||||
assert [i["incident_id"] for i in ordered] == ["i7", "i6", "i5", "i4", "i3", "i2", "i1"]
|
||||
assert len(_prompt_incidents(recent * 5)) == 20
|
||||
# falls back to started_at when updated_at is absent, and never raises
|
||||
assert _prompt_incidents([{"incident_id": "a", "started_at": NOW.isoformat()},
|
||||
{"incident_id": "b"}])[0]["incident_id"] == "a"
|
||||
@@ -0,0 +1,66 @@
|
||||
"""
|
||||
End-to-end CORS wiring for the one browser-facing REST surface.
|
||||
|
||||
The frontend's Archive page calls GET /calls/search with Authorization +
|
||||
Content-Type headers, which forces the browser to send a CORS preflight
|
||||
first. Before #110 that OPTIONS got a bare 405 with no Access-Control-*
|
||||
headers and the fetch failed with "TypeError: Failed to fetch". These
|
||||
tests drive the real app through TestClient so a regression in the
|
||||
middleware wiring (not just the helper) is caught.
|
||||
|
||||
TestClient is NOT used as a context manager on purpose: that would run the
|
||||
lifespan (mqtt_handler.connect(), the sweeper loops, dynsec bootstrap),
|
||||
none of which is needed here -- CORSMiddleware answers a preflight before
|
||||
routing or dependencies run.
|
||||
"""
|
||||
from fastapi.testclient import TestClient
|
||||
|
||||
from app.config import settings
|
||||
from app.main import app
|
||||
|
||||
client = TestClient(app)
|
||||
|
||||
ALLOWED_ORIGIN = "https://drb.cusano.net"
|
||||
DISALLOWED_ORIGIN = "https://evil.example.com"
|
||||
|
||||
|
||||
def test_default_allowed_origin_matches_the_deployed_frontend():
|
||||
# The frontend is served on the bare domain (infra Caddyfile.j2), so the
|
||||
# default must allow exactly that origin without any env override.
|
||||
assert ALLOWED_ORIGIN in settings.cors_origins
|
||||
|
||||
|
||||
def test_preflight_for_calls_search_is_allowed():
|
||||
resp = client.options(
|
||||
"/calls/search",
|
||||
headers={
|
||||
"Origin": ALLOWED_ORIGIN,
|
||||
"Access-Control-Request-Method": "GET",
|
||||
"Access-Control-Request-Headers": "authorization,content-type",
|
||||
},
|
||||
)
|
||||
assert resp.status_code == 200
|
||||
assert resp.headers.get("access-control-allow-origin") == ALLOWED_ORIGIN
|
||||
allow_methods = resp.headers.get("access-control-allow-methods", "").upper()
|
||||
assert "GET" in allow_methods
|
||||
# Bearer auth, not cookies -- credentials must never be advertised.
|
||||
assert "access-control-allow-credentials" not in resp.headers
|
||||
|
||||
|
||||
def test_preflight_from_disallowed_origin_gets_no_allow_origin():
|
||||
resp = client.options(
|
||||
"/calls/search",
|
||||
headers={
|
||||
"Origin": DISALLOWED_ORIGIN,
|
||||
"Access-Control-Request-Method": "GET",
|
||||
},
|
||||
)
|
||||
assert resp.headers.get("access-control-allow-origin") is None
|
||||
|
||||
|
||||
def test_simple_get_from_allowed_origin_is_annotated():
|
||||
# Even a non-preflight GET must carry Access-Control-Allow-Origin or the
|
||||
# browser hides the response body from the page.
|
||||
resp = client.get("/health", headers={"Origin": ALLOWED_ORIGIN})
|
||||
assert resp.status_code == 200
|
||||
assert resp.headers.get("access-control-allow-origin") == ALLOWED_ORIGIN
|
||||
@@ -5,8 +5,9 @@ Starlette does not reject `allow_origins=["*"]` combined with
|
||||
`allow_credentials=True`. It reflects the caller's Origin back in
|
||||
Access-Control-Allow-Origin and still sends
|
||||
Access-Control-Allow-Credentials: true, so the effective policy is the
|
||||
opposite of what a wildcard usually means. main.py defuses that by turning
|
||||
credentials off whenever it sees a wildcard; these tests hold it to that.
|
||||
opposite of what a wildcard usually means. main.py never enables
|
||||
credentials at all (auth is a Bearer header, not a cookie), which makes
|
||||
that pair unrepresentable; these tests hold it to that.
|
||||
|
||||
The policy lives in a pure function so it can be exercised directly --
|
||||
reloading app.main to vary settings drags every router back through import
|
||||
@@ -28,11 +29,11 @@ def test_wildcard_among_real_origins_still_disables_credentials():
|
||||
assert cors_allows_credentials(["https://app.example.com", "*"]) is False
|
||||
|
||||
|
||||
def test_named_origins_keep_credentials():
|
||||
# Naming your origins is how you ask for credentialed requests, so a
|
||||
# correctly configured deployment must not be penalised.
|
||||
assert cors_allows_credentials(["https://app.example.com"]) is True
|
||||
assert cors_allows_credentials([]) is True
|
||||
def test_credentials_never_enabled_even_for_named_origins():
|
||||
# Auth here is a Bearer header, not a cookie, so credentialed CORS is
|
||||
# never needed. The predicate is hard-off regardless of the origin list.
|
||||
assert cors_allows_credentials(["https://app.example.com"]) is False
|
||||
assert cors_allows_credentials([]) is False
|
||||
|
||||
|
||||
def test_the_app_actually_mounted_that_policy():
|
||||
@@ -42,6 +43,7 @@ def test_the_app_actually_mounted_that_policy():
|
||||
(mw.kwargs for mw in app.user_middleware if mw.cls is CORSMiddleware), None
|
||||
)
|
||||
assert opts is not None, "CORSMiddleware is not mounted at all"
|
||||
assert opts["allow_credentials"] is False
|
||||
assert opts["allow_credentials"] is cors_allows_credentials(settings.cors_origins)
|
||||
|
||||
|
||||
|
||||
@@ -304,6 +304,39 @@ async def test_a_scene_is_judged_on_its_own_embedding_and_severity():
|
||||
assert ctx["call_severity"] == "major"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_the_llm_tier_reads_the_scene_transcript_not_the_whole_call():
|
||||
"""
|
||||
server-26#102. intelligence.py writes only the primary scene's corrected
|
||||
text to calls/{id}. _call_block (the LLM correlation prompt) must reason
|
||||
over the SCENE being correlated, not a whole-call transcript that also
|
||||
contains the other scenes. _build_context threads the scene's text in;
|
||||
with no scene text it falls back to the call doc (sweep / single-scene).
|
||||
"""
|
||||
with patch("app.internal.incident_correlator.fstore") as mock_fstore:
|
||||
mock_fstore.doc_get = AsyncMock(return_value={
|
||||
"transcript": "scene one about a fire. scene two about a traffic stop.",
|
||||
})
|
||||
mock_fstore.collection_list = AsyncMock(return_value=[])
|
||||
scene = await _build_context(
|
||||
call_id="call-1", units=None, vehicles=None, cleared_units=None,
|
||||
location_coords=None, reference_time=NOW,
|
||||
system_id="sys-1", talkgroup_id=383, talkgroup_name=DISPATCH_TG,
|
||||
tags=[], incident_type="police", location=None,
|
||||
reassignment=False, create_if_new=True,
|
||||
transcript="scene two about a traffic stop.",
|
||||
)
|
||||
fallback = await _build_context(
|
||||
call_id="call-1", units=None, vehicles=None, cleared_units=None,
|
||||
location_coords=None, reference_time=NOW,
|
||||
system_id="sys-1", talkgroup_id=383, talkgroup_name=DISPATCH_TG,
|
||||
tags=[], incident_type="police", location=None,
|
||||
reassignment=False, create_if_new=True,
|
||||
)
|
||||
assert scene["scene_transcript"] == "scene two about a traffic stop."
|
||||
assert fallback["scene_transcript"] == "scene one about a fire. scene two about a traffic stop."
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_a_bare_number_never_becomes_an_incident_location_or_title():
|
||||
inc = await _create(tags=["flames"], location="49", coords=None,
|
||||
|
||||
@@ -0,0 +1,40 @@
|
||||
"""
|
||||
server-26#102 — a scene is correlated on its OWN transcript, not the whole call.
|
||||
|
||||
_scene_transcript_text slices the segments a scene owns. It must never return
|
||||
"" (an empty slice would let incident_correlator._build_context fall back to
|
||||
the call doc's whole-call transcript, re-opening the leak in exactly the case
|
||||
— bad indices — where it matters).
|
||||
"""
|
||||
from app.internal.intelligence import _scene_transcript_text
|
||||
|
||||
SEGS = [
|
||||
{"text": "structure fire, 12 Main"},
|
||||
{"text": "engine 4 responding"},
|
||||
{"text": "traffic stop, plate ABC"},
|
||||
{"text": "one occupant"},
|
||||
]
|
||||
WHOLE = "structure fire, 12 Main engine 4 responding traffic stop, plate ABC one occupant"
|
||||
|
||||
|
||||
def test_scene_owns_a_subset_of_segments():
|
||||
assert _scene_transcript_text(WHOLE, SEGS, [0, 1], None) == "structure fire, 12 Main engine 4 responding"
|
||||
assert _scene_transcript_text(WHOLE, SEGS, [2, 3], None) == "traffic stop, plate ABC one occupant"
|
||||
|
||||
|
||||
def test_corrected_text_wins_when_present():
|
||||
assert _scene_transcript_text(WHOLE, SEGS, [0], "cleaned up text") == "cleaned up text"
|
||||
|
||||
|
||||
def test_no_segment_indices_falls_back_to_whole_call():
|
||||
# single-segment calls are never numbered by _build_transcript_block → null indices
|
||||
assert _scene_transcript_text(WHOLE, SEGS, None, None) == WHOLE
|
||||
assert _scene_transcript_text(WHOLE, None, [0, 1], None) == WHOLE
|
||||
|
||||
|
||||
def test_out_of_range_or_nonint_indices_fall_back_never_empty():
|
||||
assert _scene_transcript_text(WHOLE, SEGS, [9, 10], None) == WHOLE # all out of range
|
||||
assert _scene_transcript_text(WHOLE, SEGS, ["1", "2"], None) == WHOLE # 1-based strings, rejected
|
||||
assert _scene_transcript_text(WHOLE, SEGS, [-1], None) == WHOLE # negative
|
||||
# partial validity: keep what's in range
|
||||
assert _scene_transcript_text(WHOLE, SEGS, [3, 99], None) == "one occupant"
|
||||
Reference in New Issue
Block a user