Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
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
|
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
|
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
|
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
|
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
|
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
|
tg_fast_path_idle_minutes: int = 90 # fast path: max minutes since incident last updated
|
||||||
|
|||||||
@@ -618,6 +618,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
|
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.
|
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 [])
|
call_count = len(inc.get("call_ids") or [])
|
||||||
if call_count >= settings.incident_max_calls:
|
if call_count >= settings.incident_max_calls:
|
||||||
return f"call_cap:{call_count}"
|
return f"call_cap:{call_count}"
|
||||||
@@ -851,8 +858,16 @@ async def _build_context(
|
|||||||
# the whole collection rather than being unable to correlate at all.
|
# the whole collection rather than being unable to correlate at all.
|
||||||
if org_id is not None:
|
if org_id is not None:
|
||||||
all_active = await fstore.collection_list("incidents", status="active", org_id=org_id)
|
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:
|
else:
|
||||||
all_active = await fstore.collection_list("incidents", status="active")
|
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
|
# 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"] /
|
# neither the rules engine nor the LLM tier (which reads ctx["recent"] /
|
||||||
# ctx["all_active"]) can propose linking into one.
|
# ctx["all_active"]) can propose linking into one.
|
||||||
@@ -2086,6 +2101,10 @@ async def _update_incident(
|
|||||||
# thin traffic rides along without extending its life.
|
# thin traffic rides along without extending its life.
|
||||||
if refresh_activity:
|
if refresh_activity:
|
||||||
updates["updated_at"] = _floor_at_started_at(inc, now).isoformat()
|
updates["updated_at"] = _floor_at_started_at(inc, now).isoformat()
|
||||||
|
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:
|
else:
|
||||||
updates["last_thin_at"] = now.isoformat()
|
updates["last_thin_at"] = now.isoformat()
|
||||||
# Update incident type when a re-classified call provides a concrete type.
|
# Update incident type when a re-classified call provides a concrete type.
|
||||||
@@ -2103,6 +2122,12 @@ async def _update_incident(
|
|||||||
# Signal-based auto-resolve: every tracked unit has cleared, none still active.
|
# 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
|
# 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).
|
# don't fire on incidents where units were never tracked (no unit mentions at all).
|
||||||
|
if inc.get("status") == "resolved":
|
||||||
|
# A timer close was provisional and a related call just arrived.
|
||||||
|
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:
|
if units_cleared and not units_active:
|
||||||
updates["status"] = "resolved"
|
updates["status"] = "resolved"
|
||||||
updates["resolved_at"] = now.isoformat()
|
updates["resolved_at"] = now.isoformat()
|
||||||
@@ -2171,6 +2196,7 @@ async def _create_incident(
|
|||||||
**location_fields,
|
**location_fields,
|
||||||
"location_mentions": [location] if location else [],
|
"location_mentions": [location] if location else [],
|
||||||
"call_ids": [call_id],
|
"call_ids": [call_id],
|
||||||
|
"substantive_call_count": 1,
|
||||||
"talkgroup_ids": [str(talkgroup_id)] if talkgroup_id is not None else [],
|
"talkgroup_ids": [str(talkgroup_id)] if talkgroup_id is not None else [],
|
||||||
"system_ids": [system_id] if system_id else [],
|
"system_ids": [system_id] if system_id else [],
|
||||||
"tags": tags + ["auto-generated"],
|
"tags": tags + ["auto-generated"],
|
||||||
|
|||||||
@@ -142,15 +142,41 @@ async def _summarize_incident(inc: dict) -> None:
|
|||||||
await fstore.doc_set("incidents", incident_id, updates)
|
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:
|
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")
|
all_active = await fstore.collection_list("incidents", status="active")
|
||||||
if not all_active:
|
if not all_active:
|
||||||
return
|
return
|
||||||
|
|
||||||
from app.internal import clock
|
from app.internal import clock
|
||||||
now = clock.now()
|
now = clock.now()
|
||||||
cutoff = timedelta(minutes=settings.incident_auto_resolve_minutes)
|
|
||||||
count = 0
|
count = 0
|
||||||
|
|
||||||
for inc in all_active:
|
for inc in all_active:
|
||||||
@@ -164,11 +190,12 @@ async def _resolve_stale_incidents() -> None:
|
|||||||
if updated_dt.tzinfo is None:
|
if updated_dt.tzinfo is None:
|
||||||
updated_dt = updated_dt.replace(tzinfo=timezone.utc)
|
updated_dt = updated_dt.replace(tzinfo=timezone.utc)
|
||||||
idle_minutes = (now - updated_dt).total_seconds() / 60
|
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, {
|
await fstore.doc_set("incidents", incident_id, {
|
||||||
"status": "resolved",
|
"status": "resolved",
|
||||||
"resolved_at": now.isoformat(),
|
"resolved_at": now.isoformat(),
|
||||||
"resolved_via": "idle_timeout",
|
"resolved_via": "idle_timeout",
|
||||||
|
"reopenable": True,
|
||||||
})
|
})
|
||||||
from app.internal.incident_correlator import maybe_resolve_parent
|
from app.internal.incident_correlator import maybe_resolve_parent
|
||||||
await maybe_resolve_parent(incident_id)
|
await maybe_resolve_parent(incident_id)
|
||||||
|
|||||||
@@ -66,3 +66,30 @@ def test_only_numbered_units_hold_an_incident_open():
|
|||||||
assert not ic._is_trackable_unit(junk), junk
|
assert not ic._is_trackable_unit(junk), junk
|
||||||
for real in ("45-9", "11-Adam", "Whitestone 1", "E-14", "Highway 3-4", "7"):
|
for real in ("45-9", "11-Adam", "Whitestone 1", "E-14", "Highway 3-4", "7"):
|
||||||
assert ic._is_trackable_unit(real), real
|
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")
|
||||||
|
|||||||
@@ -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.
|
# 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())
|
groups = sorted(sorted(i["call_ids"]) for i in sb_incidents.values())
|
||||||
assert groups == [["call-1", "call-2"], ["call-3"]]
|
assert groups == [["call-1", "call-2"], ["call-3"]]
|
||||||
# Each aged out on the replayed clock the way it would have live —
|
# Each aged out on the replayed clock the way it would have live — its
|
||||||
# incident_auto_resolve_minutes after its last activity, not "now".
|
# severity's quiet timer after its last activity, not "now".
|
||||||
assert run["metrics"]["resolved_via"] == {"idle_timeout": 2}
|
assert run["metrics"]["resolved_via"] == {"idle_timeout": 2}
|
||||||
first = next(i for i in sb_incidents.values() if "call-1" in i["call_ids"])
|
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"])
|
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 run["metrics"]["calls"] == 3
|
||||||
assert set(store.data[f"{root}/scenes"]) == {"call-1", "call-2", "call-3"}
|
assert set(store.data[f"{root}/scenes"]) == {"call-1", "call-2", "call-3"}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user