Compare commits
2
Commits
3a786bc227
...
8fdedee25b
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
8fdedee25b | ||
|
|
70d63abeaa |
@@ -51,6 +51,32 @@ _PURSUIT_TAGS = frozenset({
|
||||
"fleeing-vehicle", "suspect-vehicle", "eluding",
|
||||
})
|
||||
|
||||
# Four-level severity ladder, low → high (see intelligence.py EXTRACTION_PROMPT).
|
||||
_SEVERITY_RANK = {"routine": 0, "minor": 1, "moderate": 2, "major": 3}
|
||||
|
||||
|
||||
def _max_severity(current: Optional[str], new: Optional[str]) -> str:
|
||||
"""
|
||||
Highest of two severities on the four-level ladder — the merge rule used
|
||||
every time a call attaches to an existing incident.
|
||||
|
||||
Severity is monotonic by design: it only ever rises, never falls, as more
|
||||
calls link. An incident briefly assessed "major" genuinely was major at
|
||||
that moment; a later call that sounds calmer ("units clear", dispatcher
|
||||
moving on) is evidence the SITUATION is winding down, not that the earlier
|
||||
read was wrong. That's what `status`/`resolved_at` are for — resolution
|
||||
retires an incident, it doesn't retroactively erase how serious it was.
|
||||
Every severity-driven surface (worst-first incident rail, "Major only"
|
||||
filter, map colour) exists to make sure a major event is never missed;
|
||||
downgrading severity mid-incident would silently defeat that on the exact
|
||||
incidents it exists to protect. A malformed/unrecognized value from either
|
||||
side ranks as "routine" so it can never suppress a real escalation.
|
||||
"""
|
||||
current = current if current in _SEVERITY_RANK else "routine"
|
||||
new = new if new in _SEVERITY_RANK else "routine"
|
||||
return current if _SEVERITY_RANK[current] >= _SEVERITY_RANK[new] else new
|
||||
|
||||
|
||||
# Maximum plausible ground speed for a moving incident (pursuit/transport).
|
||||
# ~3 miles/min ≈ 180 mph — well above real pursuit speeds, but generous enough
|
||||
# to tolerate GPS drift and call-timing jitter. Anything faster is a bad geocode
|
||||
@@ -877,6 +903,7 @@ async def _apply_decision(decision: dict, ctx: dict) -> Optional[str]:
|
||||
location, location_coords, call_units, call_vehicles, call_embedding, now,
|
||||
talkgroup_name=talkgroup_name, incident_type=incident_type,
|
||||
cleared_units=call_cleared, refresh_activity=not thin_link,
|
||||
call_severity=call_severity,
|
||||
)
|
||||
return matched_incident["incident_id"]
|
||||
|
||||
@@ -1251,6 +1278,7 @@ async def _update_incident(
|
||||
incident_type: Optional[str] = None,
|
||||
cleared_units: Optional[list[str]] = None,
|
||||
refresh_activity: bool = True,
|
||||
call_severity: Optional[str] = None,
|
||||
) -> None:
|
||||
incident_id = inc["incident_id"]
|
||||
|
||||
@@ -1303,6 +1331,7 @@ async def _update_incident(
|
||||
"units_cleared": units_cleared,
|
||||
"location_mentions": location_mentions,
|
||||
"summary_stale": True,
|
||||
"severity": _max_severity(inc.get("severity"), call_severity),
|
||||
**embedding_updates,
|
||||
}
|
||||
|
||||
@@ -1347,6 +1376,7 @@ async def _update_incident(
|
||||
# don't fire on incidents where units were never tracked (no unit mentions at all).
|
||||
if units_cleared and not units_active:
|
||||
updates["status"] = "resolved"
|
||||
updates["resolved_at"] = now.isoformat()
|
||||
await fstore.doc_set("incidents", incident_id, updates)
|
||||
logger.info(
|
||||
f"Correlator: signal-resolved incident {incident_id} "
|
||||
@@ -1538,7 +1568,10 @@ async def maybe_resolve_parent(incident_id: str) -> None:
|
||||
return # at least one sibling still active
|
||||
|
||||
# All children resolved — close the master
|
||||
await fstore.doc_set("incidents", parent_id, {"status": "resolved"})
|
||||
await fstore.doc_set("incidents", parent_id, {
|
||||
"status": "resolved",
|
||||
"resolved_at": datetime.now(timezone.utc).isoformat(),
|
||||
})
|
||||
logger.info(
|
||||
f"Auto-resolved master incident {parent_id} "
|
||||
f"(all {len(child_ids)} child(ren) resolved)"
|
||||
|
||||
@@ -101,7 +101,10 @@ async def _resolve_stale_incidents() -> None:
|
||||
updated_dt = updated_dt.replace(tzinfo=timezone.utc)
|
||||
idle_minutes = (now - updated_dt).total_seconds() / 60
|
||||
if idle_minutes > settings.incident_auto_resolve_minutes:
|
||||
await fstore.doc_set("incidents", incident_id, {"status": "resolved"})
|
||||
await fstore.doc_set("incidents", incident_id, {
|
||||
"status": "resolved",
|
||||
"resolved_at": now.isoformat(),
|
||||
})
|
||||
from app.internal.incident_correlator import maybe_resolve_parent
|
||||
await maybe_resolve_parent(incident_id)
|
||||
logger.info(
|
||||
|
||||
@@ -173,6 +173,7 @@ async def patch_transcript(
|
||||
await fstore.doc_set("incidents", old_incident_id, {
|
||||
"call_ids": [],
|
||||
"status": "resolved",
|
||||
"resolved_at": datetime.now(timezone.utc).isoformat(),
|
||||
"summary_stale": True,
|
||||
})
|
||||
await fstore.doc_set("calls", call_id, {"incident_ids": [], "incident_id": None})
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
from typing import Optional
|
||||
from datetime import datetime, timezone
|
||||
from fastapi import APIRouter, BackgroundTasks, UploadFile, File, Form, HTTPException, Security
|
||||
from fastapi.security import HTTPBearer, HTTPAuthorizationCredentials
|
||||
from app.internal.storage import upload_audio
|
||||
@@ -202,7 +203,10 @@ async def _run_extraction_pipeline(
|
||||
if incident_id and incident_id not in incident_ids:
|
||||
incident_ids.append(incident_id)
|
||||
if scene["resolved"] and incident_id:
|
||||
await fstore.doc_set("incidents", incident_id, {"status": "resolved"})
|
||||
await fstore.doc_set("incidents", incident_id, {
|
||||
"status": "resolved",
|
||||
"resolved_at": datetime.now(timezone.utc).isoformat(),
|
||||
})
|
||||
await incident_correlator.maybe_resolve_parent(incident_id)
|
||||
logger.info(f"Auto-resolved incident {incident_id} (LLM closure detection)")
|
||||
|
||||
@@ -305,7 +309,10 @@ async def _run_intelligence_pipeline(
|
||||
if incident_id and incident_id not in incident_ids:
|
||||
incident_ids.append(incident_id)
|
||||
if scene["resolved"] and incident_id:
|
||||
await fstore.doc_set("incidents", incident_id, {"status": "resolved"})
|
||||
await fstore.doc_set("incidents", incident_id, {
|
||||
"status": "resolved",
|
||||
"resolved_at": datetime.now(timezone.utc).isoformat(),
|
||||
})
|
||||
await incident_correlator.maybe_resolve_parent(incident_id)
|
||||
logger.info(f"Auto-resolved incident {incident_id} (LLM closure detection)")
|
||||
|
||||
|
||||
@@ -18,6 +18,7 @@ from datetime import datetime, timedelta, timezone
|
||||
from unittest.mock import AsyncMock, patch
|
||||
from app.internal.incident_correlator import (
|
||||
_run_decision, _update_incident, _normalize_unit, _matching_units,
|
||||
_max_severity, maybe_resolve_parent,
|
||||
)
|
||||
|
||||
NOW = datetime(2026, 8, 16, 21, 0, 0, tzinfo=timezone.utc)
|
||||
@@ -247,3 +248,117 @@ def test_normalised_units_link_a_call_that_exact_match_would_orphan():
|
||||
call_units=["K-9A2"], is_thin_call=False, call_severity="routine",
|
||||
))
|
||||
assert decision["action"] == "link"
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Issue #17 — severity re-evaluation as calls attach (monotonic ladder)
|
||||
#
|
||||
# Decision: severity only ever rises, never falls, as more calls link (see
|
||||
# _max_severity's docstring in incident_correlator.py for the full argument).
|
||||
# An incident briefly assessed "major" genuinely was major at that moment;
|
||||
# resolution (status/resolved_at), not a later calmer-sounding call, is what
|
||||
# retires it. These tests lock in both halves of that: escalation raises the
|
||||
# stored severity, and a later lower-severity call does not undo it.
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
@pytest.mark.parametrize("current,new,expected", [
|
||||
("routine", "major", "major"), # escalation — the motivating case
|
||||
("routine", "minor", "minor"),
|
||||
("minor", "moderate", "moderate"),
|
||||
("major", "routine", "major"), # calmer call does NOT downgrade
|
||||
("major", "minor", "major"),
|
||||
("moderate", "moderate", "moderate"), # tie
|
||||
(None, "moderate", "moderate"), # incident with no prior severity
|
||||
("major", None, "major"),
|
||||
("major", "bogus", "major"), # malformed value ranks as routine
|
||||
("bogus", "minor", "minor"),
|
||||
])
|
||||
def test_max_severity_is_monotonic(current, new, expected):
|
||||
assert _max_severity(current, new) == expected
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_escalating_call_raises_stored_incident_severity():
|
||||
"""The #17 motivating case: an incident opened routine, a later call is a
|
||||
working structure fire — the incident's severity must reflect it."""
|
||||
inc = _incident(2.0)
|
||||
inc["severity"] = "routine"
|
||||
with patch("app.internal.incident_correlator.fstore") as mock_fstore:
|
||||
mock_fstore.doc_set = AsyncMock()
|
||||
await _update_incident(
|
||||
inc, "call-2", 9048, "sys-1", [], None, None, [], [], None, NOW,
|
||||
call_severity="major",
|
||||
)
|
||||
updates = mock_fstore.doc_set.await_args.args[2]
|
||||
assert updates["severity"] == "major"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_calmer_followup_call_does_not_downgrade_severity():
|
||||
inc = _incident(2.0)
|
||||
inc["severity"] = "major"
|
||||
with patch("app.internal.incident_correlator.fstore") as mock_fstore:
|
||||
mock_fstore.doc_set = AsyncMock()
|
||||
await _update_incident(
|
||||
inc, "call-2", 9048, "sys-1", [], None, None, [], [], None, NOW,
|
||||
call_severity="routine",
|
||||
)
|
||||
updates = mock_fstore.doc_set.await_args.args[2]
|
||||
assert updates["severity"] == "major"
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Issue #18 — every resolution site stamps resolved_at
|
||||
#
|
||||
# updated_at is not a substitute (thin/ack calls deliberately don't move it,
|
||||
# unrelated field updates do) and existing rows are left null, not backfilled
|
||||
# — null means "resolved before this field existed", not "never resolved".
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_signal_resolve_stamps_resolved_at():
|
||||
"""All tracked units clear -> _update_incident's own auto-resolve path."""
|
||||
inc = _incident(2.0)
|
||||
inc["units_active"] = ["6-Adam"]
|
||||
with patch("app.internal.incident_correlator.fstore") as mock_fstore:
|
||||
mock_fstore.doc_set = AsyncMock()
|
||||
# standalone incident — maybe_resolve_parent's own doc_get short-circuits on None
|
||||
mock_fstore.doc_get = AsyncMock(return_value=None)
|
||||
await _update_incident(
|
||||
inc, "call-2", 9048, "sys-1", [], None, None, [], [], None, NOW,
|
||||
cleared_units=["6-Adam"],
|
||||
)
|
||||
updates = mock_fstore.doc_set.await_args.args[2]
|
||||
assert updates["status"] == "resolved"
|
||||
assert updates["resolved_at"] == NOW.isoformat()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_master_auto_resolve_stamps_resolved_at():
|
||||
"""maybe_resolve_parent closes a master once every child has resolved."""
|
||||
child_a = {"incident_id": "child-a", "parent_incident_id": "master-1"}
|
||||
master = {
|
||||
"incident_id": "master-1",
|
||||
"status": "active",
|
||||
"child_incident_ids": ["child-a", "child-b"],
|
||||
}
|
||||
child_b_resolved = {"incident_id": "child-b", "status": "resolved"}
|
||||
|
||||
async def fake_doc_get(collection, doc_id):
|
||||
return {
|
||||
"child-a": child_a,
|
||||
"master-1": master,
|
||||
"child-b": child_b_resolved,
|
||||
}.get(doc_id)
|
||||
|
||||
with patch("app.internal.incident_correlator.fstore") as mock_fstore:
|
||||
mock_fstore.doc_get = AsyncMock(side_effect=fake_doc_get)
|
||||
mock_fstore.doc_set = AsyncMock()
|
||||
await maybe_resolve_parent("child-a")
|
||||
|
||||
mock_fstore.doc_set.assert_awaited_once()
|
||||
args = mock_fstore.doc_set.await_args.args
|
||||
assert args[0] == "incidents"
|
||||
assert args[1] == "master-1"
|
||||
assert args[2]["status"] == "resolved"
|
||||
assert "resolved_at" in args[2] and args[2]["resolved_at"]
|
||||
|
||||
@@ -109,6 +109,17 @@ export interface CallRecord {
|
||||
/** Legacy field — present on calls recorded before the multi-scene migration. */
|
||||
incident_id?: string | null;
|
||||
location: string | null;
|
||||
/** Per-call geocode, written by intelligence.py. Powers the incident-path polyline on the map. */
|
||||
location_coords?: { lat: number; lng: number } | null;
|
||||
/** Unit callsigns mentioned in this specific transmission (e.g. "15-9", "E-41"). */
|
||||
units?: string[];
|
||||
vehicles?: string[];
|
||||
/** Units this call reported as clearing/back in service. */
|
||||
cleared_units?: string[];
|
||||
/** Set when dedup.py identifies this as a second node's recording of the same transmission — the canonical copy has this null. */
|
||||
duplicate_of?: string | null;
|
||||
/** Radio source address, when the system exposes it. */
|
||||
srcaddr?: string | null;
|
||||
tags: string[];
|
||||
status: "active" | "ended";
|
||||
/** Four-level ladder: routine | minor | moderate | major. Legacy docs may still carry "unknown". */
|
||||
@@ -134,6 +145,14 @@ export interface IncidentRecord {
|
||||
talkgroup_ids: string[];
|
||||
units: string[];
|
||||
vehicles: string[];
|
||||
/** Units currently believed on scene — maintained by incident_correlator.py `_attach`. */
|
||||
units_active?: string[];
|
||||
/** Units that reported clearing/back in service on this incident. */
|
||||
units_cleared?: string[];
|
||||
/** Free-text location mentions accumulated across the incident's calls, beyond the primary `location`. */
|
||||
location_mentions?: string[];
|
||||
/** ISO timestamp of the last thin/status-only call attached (doesn't refresh `updated_at`). */
|
||||
last_thin_at?: string | null;
|
||||
severity: string | null;
|
||||
started_at: string;
|
||||
updated_at: string;
|
||||
|
||||
@@ -52,10 +52,17 @@ export function useCalls(limitCount = 50, dateFrom?: Date, dateTo?: Date) {
|
||||
const toISO = (v: any): string | null =>
|
||||
v?.toDate?.()?.toISOString?.() ?? (typeof v === "string" ? v : null);
|
||||
unsubFirestore = onSnapshot(q, (snap) => {
|
||||
setCalls(snap.docs.map((d) => {
|
||||
const docs = snap.docs
|
||||
.map((d) => {
|
||||
const data = d.data();
|
||||
return { ...data, started_at: toISO(data.started_at) ?? "", ended_at: toISO(data.ended_at) } as CallRecord;
|
||||
}));
|
||||
})
|
||||
// dedup.py flags a second node's recording of the same transmission
|
||||
// with duplicate_of set to the canonical call_id — filtered client-
|
||||
// side (not a where() clause) so this doesn't need a new composite
|
||||
// index alongside the existing org_id/started_at query.
|
||||
.filter((c) => !c.duplicate_of);
|
||||
setCalls(docs);
|
||||
setLoading(false);
|
||||
}, (err: FirestoreError) => { console.error("useCalls:", err); setError(err.message); setLoading(false); });
|
||||
});
|
||||
@@ -92,10 +99,12 @@ export function useCallsByIncident(incidentId: string | null) {
|
||||
where("incident_ids", "array-contains", incidentId)
|
||||
);
|
||||
unsubFirestore = onSnapshot(q, (snap) => {
|
||||
const docs = snap.docs.map((d) => {
|
||||
const docs = snap.docs
|
||||
.map((d) => {
|
||||
const data = d.data();
|
||||
return { ...data, started_at: toISO(data.started_at) ?? "", ended_at: toISO(data.ended_at) } as CallRecord;
|
||||
});
|
||||
})
|
||||
.filter((c) => !c.duplicate_of);
|
||||
docs.sort((a, b) => a.started_at.localeCompare(b.started_at));
|
||||
setCalls(docs);
|
||||
setLoading(false);
|
||||
@@ -131,10 +140,14 @@ export function useActiveCalls() {
|
||||
const toISO = (v: any): string | null =>
|
||||
v?.toDate?.()?.toISOString?.() ?? (typeof v === "string" ? v : null);
|
||||
unsubFirestore = onSnapshot(q, (snap) => {
|
||||
setCalls(snap.docs.map((d) => {
|
||||
setCalls(
|
||||
snap.docs
|
||||
.map((d) => {
|
||||
const data = d.data();
|
||||
return { ...data, started_at: toISO(data.started_at) ?? "", ended_at: toISO(data.ended_at) } as CallRecord;
|
||||
}));
|
||||
})
|
||||
.filter((c) => !c.duplicate_of)
|
||||
);
|
||||
}, (err: FirestoreError) => { console.error("useActiveCalls:", err); });
|
||||
});
|
||||
|
||||
|
||||
Reference in New Issue
Block a user