Compare commits
4
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b9e7524817 | ||
|
|
e972cace4a | ||
|
|
969d175a67 | ||
|
|
731b54bed9 |
@@ -89,7 +89,17 @@ class Settings(BaseSettings):
|
||||
embedding_cross_tg_threshold: float = 0.85 # cross-TG path: same dept + 2+ shared units
|
||||
location_proximity_km: float = 0.5 # radius for location-proximity matching
|
||||
geocode_max_km: float = 40.0 # reject geocode results farther than this from the node
|
||||
incident_auto_resolve_minutes: int = 90 # auto-resolve after N minutes with no new calls
|
||||
incident_auto_resolve_minutes: int = 90 # auto-resolve after N minutes with no new calls (major / unknown severity)
|
||||
# Most jobs never say 10-8 on the air (replay of 09-22, server-26#170: ~5 of
|
||||
# ~25 real incidents had an audible clear), so the quiet timer IS the close
|
||||
# for most of them, and one 90-minute timer kept a lockout or a plate check
|
||||
# "active" on the portal an hour after it ended. Scaled by the incident's
|
||||
# severity instead, and made provisional: a timer-closed incident stays
|
||||
# reopenable for incident_reopen_window_minutes, so a long quiet job that
|
||||
# comes back on the air rejoins its own incident rather than splitting.
|
||||
incident_auto_resolve_minutes_routine: int = 30 # routine / minor
|
||||
incident_auto_resolve_minutes_moderate: int = 60
|
||||
incident_reopen_window_minutes: int = 90 # since last substantive call
|
||||
unit_continuity_max_idle_minutes: int = 20 # unit-continuity path: skip if incident idle > this
|
||||
recorrelation_scan_minutes: int = 60 # re-examine orphaned calls ended within this window
|
||||
tg_fast_path_idle_minutes: int = 90 # fast path: max minutes since incident last updated
|
||||
|
||||
@@ -224,6 +224,16 @@ def _normalize_unit(unit: str) -> str:
|
||||
return key or unit.strip().lower()
|
||||
|
||||
|
||||
def _after_close(inc: dict, now: datetime) -> bool:
|
||||
try:
|
||||
closed = datetime.fromisoformat(str(inc.get("resolved_at") or "").replace("Z", "+00:00"))
|
||||
except ValueError:
|
||||
return True
|
||||
if closed.tzinfo is None:
|
||||
closed = closed.replace(tzinfo=timezone.utc)
|
||||
return now > closed
|
||||
|
||||
|
||||
def _is_trackable_unit(unit: str) -> bool:
|
||||
"""
|
||||
Whether a unit is concrete enough to hold an incident open until it clears.
|
||||
@@ -618,6 +628,13 @@ def _incident_at_capacity(inc: dict, now: datetime) -> Optional[str]:
|
||||
it has and still auto-resolves on the normal idle sweep. It just stops
|
||||
being a candidate, so the next call opens a fresh incident.
|
||||
"""
|
||||
# Content-free replies ("10-4", "Cut.", "Copy") are what filled the cap:
|
||||
# the 09-22 bridge MVA hit 40 calls in 32 minutes with ~40% of them thin,
|
||||
# split in half, and its second half took a different job's title. Only
|
||||
# substantive calls count; incidents written before this field existed
|
||||
# fall back to the raw count.
|
||||
call_count = inc.get("substantive_call_count")
|
||||
if call_count is None:
|
||||
call_count = len(inc.get("call_ids") or [])
|
||||
if call_count >= settings.incident_max_calls:
|
||||
return f"call_cap:{call_count}"
|
||||
@@ -851,8 +868,16 @@ async def _build_context(
|
||||
# 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)
|
||||
reopenable = await fstore.collection_list("incidents", status="resolved", reopenable=True, org_id=org_id)
|
||||
else:
|
||||
all_active = await fstore.collection_list("incidents", status="active")
|
||||
reopenable = await fstore.collection_list("incidents", status="resolved", reopenable=True)
|
||||
# A timer-closed incident is provisional (summarizer._resolve_stale_incidents):
|
||||
# inside its reopen window it is still a candidate, and linking a call to
|
||||
# it reopens it (_update_incident). The fast path's own recency gate
|
||||
# (tg_fast_path_idle_minutes) still applies to it like any other candidate.
|
||||
reopen_window = timedelta(minutes=settings.incident_reopen_window_minutes)
|
||||
all_active += [inc for inc in reopenable if _idle_gate_minutes(inc, now) <= reopen_window.total_seconds() / 60]
|
||||
# Incidents past the hard caps are removed from the candidate pool here, so
|
||||
# neither the rules engine nor the LLM tier (which reads ctx["recent"] /
|
||||
# ctx["all_active"]) can propose linking into one.
|
||||
@@ -2014,7 +2039,8 @@ async def _update_incident(
|
||||
incident_id = inc["incident_id"]
|
||||
|
||||
call_ids = list(inc.get("call_ids") or [])
|
||||
if call_id not in call_ids:
|
||||
is_new_call = call_id not in call_ids
|
||||
if is_new_call:
|
||||
call_ids.append(call_id)
|
||||
|
||||
talkgroup_ids = list(inc.get("talkgroup_ids") or [])
|
||||
@@ -2086,6 +2112,11 @@ async def _update_incident(
|
||||
# thin traffic rides along without extending its life.
|
||||
if refresh_activity:
|
||||
updates["updated_at"] = _floor_at_started_at(inc, now).isoformat()
|
||||
if is_new_call: # a second scene of the same call is not a second call
|
||||
updates["substantive_call_count"] = (
|
||||
inc.get("substantive_call_count")
|
||||
if inc.get("substantive_call_count") is not None else len(inc.get("call_ids") or [])
|
||||
) + 1
|
||||
else:
|
||||
updates["last_thin_at"] = now.isoformat()
|
||||
# Update incident type when a re-classified call provides a concrete type.
|
||||
@@ -2103,6 +2134,15 @@ async def _update_incident(
|
||||
# Signal-based auto-resolve: every tracked unit has cleared, none still active.
|
||||
# Requires at least one unit to have explicitly signalled back-in-service so we
|
||||
# don't fire on incidents where units were never tracked (no unit mentions at all).
|
||||
if inc.get("status") == "resolved" and refresh_activity and _after_close(inc, now):
|
||||
# A timer close was provisional and a related, substantive call arrived
|
||||
# after it. A thin "10-4" rides along without reopening (it would not
|
||||
# refresh updated_at, so the next sweep would just close it again),
|
||||
# and neither does a sweep link of a call from before the close.
|
||||
updates.update({"status": "active", "resolved_at": None, "resolved_via": None,
|
||||
"reopenable": False, "reopened_count": (inc.get("reopened_count") or 0) + 1})
|
||||
logger.info(f"Correlator: reopened timer-closed incident {incident_id} (call {call_id})")
|
||||
|
||||
if units_cleared and not units_active:
|
||||
updates["status"] = "resolved"
|
||||
updates["resolved_at"] = now.isoformat()
|
||||
@@ -2171,6 +2211,7 @@ async def _create_incident(
|
||||
**location_fields,
|
||||
"location_mentions": [location] if location else [],
|
||||
"call_ids": [call_id],
|
||||
"substantive_call_count": 1,
|
||||
"talkgroup_ids": [str(talkgroup_id)] if talkgroup_id is not None else [],
|
||||
"system_ids": [system_id] if system_id else [],
|
||||
"tags": tags + ["auto-generated"],
|
||||
@@ -2347,6 +2388,10 @@ async def _find_cross_system_parent(
|
||||
best_score = 0.0
|
||||
|
||||
for inc in recent:
|
||||
# A timer-closed incident is in `recent` only so a related call can
|
||||
# reopen it; it must not be adopted as another agency's parent.
|
||||
if inc.get("status") != "active":
|
||||
continue
|
||||
# Only cross-system candidates
|
||||
if system_id in (inc.get("system_ids") or []):
|
||||
continue
|
||||
|
||||
@@ -142,15 +142,41 @@ async def _summarize_incident(inc: dict) -> None:
|
||||
await fstore.doc_set("incidents", incident_id, updates)
|
||||
|
||||
|
||||
def _auto_resolve_minutes(inc: dict) -> int:
|
||||
"""Quiet time before a timer close, by severity (see config: incident_auto_resolve_minutes_*)."""
|
||||
sev = (inc.get("severity") or "").lower()
|
||||
if sev in ("routine", "minor"):
|
||||
return settings.incident_auto_resolve_minutes_routine
|
||||
if sev == "moderate":
|
||||
return settings.incident_auto_resolve_minutes_moderate
|
||||
return settings.incident_auto_resolve_minutes
|
||||
|
||||
|
||||
async def _expire_reopen_windows(now) -> None:
|
||||
"""A timer-closed incident stops being reopenable once its window passes,
|
||||
so the correlator's reopenable pool stays bounded."""
|
||||
window = timedelta(minutes=settings.incident_reopen_window_minutes)
|
||||
for inc in await fstore.collection_list("incidents", status="resolved", reopenable=True):
|
||||
try:
|
||||
updated = datetime.fromisoformat(str(inc.get("updated_at", "")).replace("Z", "+00:00"))
|
||||
if updated.tzinfo is None:
|
||||
updated = updated.replace(tzinfo=timezone.utc)
|
||||
except ValueError:
|
||||
updated = None
|
||||
if updated is None or now - updated > window:
|
||||
await fstore.doc_set("incidents", inc["incident_id"], {"reopenable": False})
|
||||
|
||||
|
||||
async def _resolve_stale_incidents() -> None:
|
||||
"""Auto-resolve active incidents that have had no new calls for incident_auto_resolve_minutes."""
|
||||
"""Timer-close active incidents that have been quiet longer than their severity allows."""
|
||||
from app.internal import clock
|
||||
await _expire_reopen_windows(clock.now())
|
||||
all_active = await fstore.collection_list("incidents", status="active")
|
||||
if not all_active:
|
||||
return
|
||||
|
||||
from app.internal import clock
|
||||
now = clock.now()
|
||||
cutoff = timedelta(minutes=settings.incident_auto_resolve_minutes)
|
||||
count = 0
|
||||
|
||||
for inc in all_active:
|
||||
@@ -164,11 +190,12 @@ async def _resolve_stale_incidents() -> None:
|
||||
if updated_dt.tzinfo is 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:
|
||||
if idle_minutes > _auto_resolve_minutes(inc):
|
||||
await fstore.doc_set("incidents", incident_id, {
|
||||
"status": "resolved",
|
||||
"resolved_at": now.isoformat(),
|
||||
"resolved_via": "idle_timeout",
|
||||
"reopenable": True,
|
||||
})
|
||||
from app.internal.incident_correlator import maybe_resolve_parent
|
||||
await maybe_resolve_parent(incident_id)
|
||||
|
||||
@@ -66,3 +66,36 @@ def test_only_numbered_units_hold_an_incident_open():
|
||||
assert not ic._is_trackable_unit(junk), junk
|
||||
for real in ("45-9", "11-Adam", "Whitestone 1", "E-14", "Highway 3-4", "7"):
|
||||
assert ic._is_trackable_unit(real), real
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Provisional timer close + reopen, substantive cap (server-26#170 replay)
|
||||
# ---------------------------------------------------------------------------
|
||||
from datetime import datetime, timedelta, timezone # noqa: E402
|
||||
|
||||
from app.config import settings # noqa: E402
|
||||
from app.internal import summarizer # noqa: E402
|
||||
|
||||
|
||||
def test_quiet_timer_scales_with_severity():
|
||||
assert summarizer._auto_resolve_minutes({"severity": "routine"}) == settings.incident_auto_resolve_minutes_routine
|
||||
assert summarizer._auto_resolve_minutes({"severity": "minor"}) == settings.incident_auto_resolve_minutes_routine
|
||||
assert summarizer._auto_resolve_minutes({"severity": "moderate"}) == settings.incident_auto_resolve_minutes_moderate
|
||||
assert summarizer._auto_resolve_minutes({"severity": "major"}) == settings.incident_auto_resolve_minutes
|
||||
assert summarizer._auto_resolve_minutes({}) == settings.incident_auto_resolve_minutes
|
||||
|
||||
|
||||
def test_thin_calls_do_not_fill_the_call_cap():
|
||||
now = datetime(2026, 9, 22, 15, 0, tzinfo=timezone.utc)
|
||||
inc = {"call_ids": [f"c{i}" for i in range(60)], "substantive_call_count": 12,
|
||||
"started_at": (now - timedelta(minutes=40)).isoformat(),
|
||||
"updated_at": now.isoformat()}
|
||||
assert ic._incident_at_capacity(inc, now) is None
|
||||
legacy = {k: v for k, v in inc.items() if k != "substantive_call_count"}
|
||||
assert ic._incident_at_capacity(legacy, now).startswith("call_cap")
|
||||
|
||||
|
||||
def test_reopen_only_for_a_call_after_the_close():
|
||||
closed = {"resolved_at": "2026-09-22T15:00:00+00:00"}
|
||||
assert ic._after_close(closed, datetime(2026, 9, 22, 15, 5, tzinfo=timezone.utc))
|
||||
assert not ic._after_close(closed, datetime(2026, 9, 22, 14, 55, tzinfo=timezone.utc))
|
||||
|
||||
@@ -269,12 +269,14 @@ async def test_run_writes_only_to_its_sandbox_and_pins_the_clock(store):
|
||||
# The two Car 12 calls are one job; the Car 40 call five hours later is another.
|
||||
groups = sorted(sorted(i["call_ids"]) for i in sb_incidents.values())
|
||||
assert groups == [["call-1", "call-2"], ["call-3"]]
|
||||
# Each aged out on the replayed clock the way it would have live —
|
||||
# incident_auto_resolve_minutes after its last activity, not "now".
|
||||
# Each aged out on the replayed clock the way it would have live — its
|
||||
# severity's quiet timer after its last activity, not "now".
|
||||
assert run["metrics"]["resolved_via"] == {"idle_timeout": 2}
|
||||
first = next(i for i in sb_incidents.values() if "call-1" in i["call_ids"])
|
||||
idle = datetime.fromisoformat(first["resolved_at"]) - datetime.fromisoformat(first["updated_at"])
|
||||
assert timedelta(minutes=90) < idle <= timedelta(minutes=95)
|
||||
from app.internal.summarizer import _auto_resolve_minutes
|
||||
limit = timedelta(minutes=_auto_resolve_minutes(first))
|
||||
assert limit < idle <= limit + timedelta(minutes=5)
|
||||
assert run["metrics"]["calls"] == 3
|
||||
assert set(store.data[f"{root}/scenes"]) == {"call-1", "call-2", "call-3"}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user