From 0ba658a1979e6ca8967f8adf903f78eaddcec986 Mon Sep 17 00:00:00 2001 From: albert Date: Mon, 28 Sep 2026 18:34:15 +0200 Subject: [PATCH 1/6] feat(ws): live dashboard channel so another browser's change shows up on its own MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The dashboard was pull-only: a change produced elsewhere — a candidate accepting, a case escalating, a rescue closed by the worker — stayed invisible until the manager refocused, acted or reloaded. Spec §7.5 already defined the endpoint and Caddy already routed it. - `app/events.py`: a typed `DashboardEvent` (name, location, rescue, shift, origin, at — never a body, phone, name or health flag) and an `EventBus` port with a Redis publisher, an async subscriber and a no-op fake. Publishing is best effort: a broker outage degrades to the previous pull-only behaviour instead of failing a rescue. - The orchestrator publishes at every transition; the worker runtime injects the Redis bus (it was silently using the no-op, which is why the channel was empty with a healthy socket). - `WS /ws/locations/{id}?token=…` validates the manager's JWT and their right to that location before accepting (4401/4403), relays the location's channel and cleans up on disconnect (4503 when the broker is down). The relay and the client listener run together, and the socket is closed before anything is cancelled — awaiting or cancelling a blocked `receive()` deadlocks the test client, which is how the suite hung the first time. - The frontend `useLiveEvents` hook maps each event to the query keys the actions already invalidate (so nothing is duplicated in the client cache), reconnects with capped backoff, clears the session on 4401 and stops on unmount; the shell shows a quiet indicator while it is not live, and the dev proxy forwards `/ws` with the upgrade. Verified live end to end: the events reach Redis with the privacy rules intact, and with two browsers open a rescue produced in one moved the other's Today board 2.5 s later with no reload. 493 backend tests, 196 frontend tests, ruff, mypy, oxlint, build and tsc clean. --- backend/app/api/ws.py | 133 +++++++++ backend/app/events.py | 232 +++++++++++++++ backend/app/main.py | 2 + backend/app/runtime.py | 6 + backend/app/services/orchestrator.py | 245 +++++++++++++++- backend/tests/unit/api/test_ws.py | 277 ++++++++++++++++++ backend/tests/unit/test_events.py | 208 +++++++++++++ frontend/src/App.tsx | 13 + .../services/__tests__/liveEvents.test.tsx | 168 +++++++++++ frontend/src/services/liveEvents.ts | 168 +++++++++++ frontend/vite.config.ts | 8 + odd/tasks/realtime-events.md | 122 ++++++++ 12 files changed, 1580 insertions(+), 2 deletions(-) create mode 100644 backend/app/api/ws.py create mode 100644 backend/app/events.py create mode 100644 backend/tests/unit/api/test_ws.py create mode 100644 backend/tests/unit/test_events.py create mode 100644 frontend/src/services/__tests__/liveEvents.test.tsx create mode 100644 frontend/src/services/liveEvents.ts create mode 100644 odd/tasks/realtime-events.md diff --git a/backend/app/api/ws.py b/backend/app/api/ws.py new file mode 100644 index 0000000..222fb15 --- /dev/null +++ b/backend/app/api/ws.py @@ -0,0 +1,133 @@ +"""Live event channel for the dashboard (spec §7.5): a thin authenticated relay. + +`WS /ws/locations/{location_id}?token=...` validates the manager's JWT before +the socket streams anything, checks the manager's right to that location, and +then forwards every event published on the location's Redis channel as JSON +text. It is a relay only: no database access per event, no business logic — +the screens keep rendering from the API, the events just tell them what to +refetch (spec §7.6). + +Close codes (sent right after `accept` so the browser can actually read them: +closing before the accept degenerates into an opaque HTTP 403 handshake +rejection under uvicorn, which would hide the reason from the client): + +- 4401: missing, malformed, expired or wrong-role token +- 4403: a manager without rights on that location +- 4503: the event broker (Redis) is unavailable +""" + +import asyncio +import contextlib + +import structlog +from fastapi import APIRouter, Depends, WebSocket, WebSocketDisconnect +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker + +from app.core.config import Settings, get_settings +from app.db.models import Manager +from app.db.session import create_engine_and_session +from app.events import EventSubscription, channel_for, subscribe +from app.security.tokens import verify_token + +router = APIRouter() + +logger = structlog.get_logger(__name__) + +ALLOWED_ROLES = ("manager", "operator") + + +def ws_session_factory( + settings: Settings = Depends(get_settings), +) -> async_sessionmaker[AsyncSession]: + """One short-lived session factory per socket (the manager check is the + only DB access in the whole connection — never per event).""" + _, factory = create_engine_and_session(settings.database_url) + return factory + + +async def _reject(websocket: WebSocket, code: int) -> None: + """Answer the handshake with a close code the client can read.""" + await websocket.accept() + await websocket.close(code=code) + + +async def _forward(websocket: WebSocket, subscription: EventSubscription) -> None: + """Relay loop: publish-side frames in, socket text frames out.""" + async for event in subscription.events(): + await websocket.send_text(event.to_json()) + + +async def _drain(websocket: WebSocket) -> None: + """Consume client frames until the socket closes. + + The relay and this loop run together and the connection ends when EITHER + finishes. Awaiting `websocket.receive()` alone deadlocks against a client + that is waiting for the server to finish (the browser and the test client + both do), which is how this hung the suite the first time. + """ + try: + while True: + await websocket.receive() + except WebSocketDisconnect: + return + + +@router.websocket("/ws/locations/{location_id}") +async def location_events( + websocket: WebSocket, + location_id: str, + token: str = "", + session_factory: async_sessionmaker[AsyncSession] = Depends(ws_session_factory), + settings: Settings = Depends(get_settings), +) -> None: + """Stream one location's dashboard events to an authenticated manager.""" + claims = verify_token(token, settings) if token else None + if claims is None or claims.role not in ALLOWED_ROLES: + await _reject(websocket, 4401) + return + + # The manager may only watch their own locations; the operator role sees + # everything (spec §7.6). Seeded managers carry `location_ids`. + async with session_factory() as session: + manager = ( + await session.execute(select(Manager).where(Manager.id == claims.manager_id)) + ).scalar_one_or_none() + if manager is None or ( + claims.role == "manager" and location_id not in (manager.location_ids or []) + ): + await _reject(websocket, 4403) + return + + try: + # Subscribe BEFORE accepting: a down broker must answer a clear close + # code instead of a socket that hangs open and never forwards (§9.3). + subscription = await subscribe(settings.redis_url, channel_for(location_id)) + except Exception as error: + logger.warning( + "ws_event_broker_unavailable", + location_id=location_id, + error=str(error)[:200], + ) + await _reject(websocket, 4503) + return + + await websocket.accept() + relay = asyncio.create_task(_forward(websocket, subscription)) + listener = asyncio.create_task(_drain(websocket)) + try: + # Whichever ends first ends the connection: the client going away, the + # subscription finishing (the broker closed), or a socket error. + await asyncio.wait({relay, listener}, return_when=asyncio.FIRST_COMPLETED) + finally: + # Close BEFORE cancelling: cancelling a `receive()` that is blocked in + # the ASGI portal deadlocks the test client, while closing the socket + # lets the listener raise its disconnect and finish on its own. + with contextlib.suppress(RuntimeError): + await websocket.close() + relay.cancel() + await asyncio.gather(relay, listener, return_exceptions=True) + await subscription.close() + + +__all__ = ["router", "ws_session_factory"] diff --git a/backend/app/events.py b/backend/app/events.py new file mode 100644 index 0000000..465a89a --- /dev/null +++ b/backend/app/events.py @@ -0,0 +1,232 @@ +"""Dashboard event bus (spec §7.5, §7.6): the worker publishes, the API relays. + +The domain (worker process) publishes one small typed event per meaningful +transition on a Redis pub/sub channel; the API process streams the channel to +the authenticated dashboard sockets. The port keeps both sides testable: the +orchestrator gets a no-op in tests, and the endpoint gets an injected fake. + +Privacy invariant (spec §10): a `DashboardEvent` NEVER carries a message +body, a phone number, a name or any health detail — only ids, the origin +label and a timestamp. +""" + +import contextlib +import json +from collections.abc import AsyncIterator, Callable +from dataclasses import dataclass +from datetime import UTC, datetime +from enum import StrEnum +from typing import Any, Protocol + +import structlog + +logger = structlog.get_logger(__name__) + +# The Redis channel for one location's live events (spec §7.5). The API and +# the worker agree on the key by construction. +CHANNEL_PREFIX = "shift_rescue:events:" + + +def channel_for(location_id: str) -> str: + """The pub/sub channel that carries a location's dashboard events.""" + return f"{CHANNEL_PREFIX}{location_id}" + + +class EventName(StrEnum): + """Every meaningful domain change the dashboard reacts to (spec §7.6).""" + + RESCUE_OPENED = "RESCUE_OPENED" + OFFERS_SENT = "OFFERS_SENT" + OFFER_ACCEPTED = "OFFER_ACCEPTED" + OFFER_DECLINED = "OFFER_DECLINED" + CASE_COVERED = "CASE_COVERED" + CASE_ESCALATED = "CASE_ESCALATED" + CASE_CLOSED = "CASE_CLOSED" + MESSAGE_RECEIVED = "MESSAGE_RECEIVED" + MESSAGE_SENT = "MESSAGE_SENT" + APPROVAL_REQUESTED = "APPROVAL_REQUESTED" + APPROVAL_DECIDED = "APPROVAL_DECIDED" + + +@dataclass(frozen=True) +class DashboardEvent: + """One domain change, routed by location (spec §7.5). + + Payload rules (spec §10, enforced by construction: these six fields are + the whole payload): **never** a message body, a phone number, a name or + a health flag — only ids, the origin label ("system", "agent", + "employee:", "manager:") and the ISO-8601 moment. + """ + + name: EventName + location_id: str + rescue_id: str | None + shift_id: str | None + origin: str + at: str + + def to_json(self) -> str: + """The exact wire format the socket forwards and the hook parses.""" + return json.dumps( + { + "name": self.name.value, + "location_id": self.location_id, + "rescue_id": self.rescue_id, + "shift_id": self.shift_id, + "origin": self.origin, + "at": self.at, + } + ) + + @classmethod + def from_json(cls, raw: str) -> "DashboardEvent | None": + """Parse a subscriber payload; None when malformed or unknown. + + Total on purpose: a bad frame on the channel must never crash the + relay, it is just skipped. + """ + try: + data = json.loads(raw) + except (TypeError, ValueError): + return None + if not isinstance(data, dict): + return None + raw_name = data.get("name") + if not isinstance(raw_name, str): + return None + try: + name = EventName(raw_name) + location_id = data["location_id"] + except (KeyError, ValueError): + return None + if not isinstance(location_id, str): + return None + + def _optional_str(key: str) -> str | None: + value = data.get(key) + return value if isinstance(value, str) else None + + return cls( + name=name, + location_id=location_id, + rescue_id=_optional_str("rescue_id"), + shift_id=_optional_str("shift_id"), + origin=_optional_str("origin") or "system", + at=_optional_str("at") or datetime.now(UTC).isoformat(), + ) + + +class EventBus(Protocol): + """Port the orchestrator publishes through (spec §9.3: best effort).""" + + async def publish(self, event: DashboardEvent) -> None: ... + + +class NoopEventBus: + """Test/standalone bus: events are dropped, nothing ever fails.""" + + async def publish(self, event: DashboardEvent) -> None: + return None + + +def _default_redis(url: str) -> Any: + # Lazy import: the API and the worker must boot without a broker and + # without redis installed beyond the Celery dependency (spec §7.4). + import redis.asyncio as aioredis + + return aioredis.Redis.from_url(url, decode_responses=True) + + +class RedisEventBus: + """Worker-side publisher. Failures are logged and swallowed (spec §9.3): + a down broker must never fail a rescue.""" + + def __init__(self, url: str, client_factory: Callable[[], Any] | None = None) -> None: + self._url = url + self._client_factory = client_factory or (lambda: _default_redis(url)) + self._client: Any = None + + async def publish(self, event: DashboardEvent) -> None: + try: + if self._client is None: + self._client = self._client_factory() + await self._client.publish(channel_for(event.location_id), event.to_json()) + except Exception as error: # the broker never breaks the worker + logger.warning( + "dashboard_event_publish_failed", + name=event.name.value, + error=str(error)[:200], + ) + self._client = None # a broken connection is rebuilt on the next event + + +class EventSubscription: + """An active pub/sub subscription on the API side (spec §7.5). + + Iterating `events()` yields parsed `DashboardEvent`s; a dropped broker + connection ends the iteration quietly so the socket handler cleans up + (the browser reconnects with backoff) instead of crashing. + """ + + def __init__(self, pubsub: Any, client: Any) -> None: + self._pubsub = pubsub + self._client = client + self.closed = False + + async def events(self) -> AsyncIterator[DashboardEvent]: + try: + async for message in self._pubsub.listen(): + if not isinstance(message, dict) or message.get("type") != "message": + continue + event = DashboardEvent.from_json(str(message.get("data") or "")) + if event is not None: + yield event + except Exception as error: # broker down or dropped: end the stream quietly + logger.warning("event_stream_dropped", error=str(error)[:200]) + return + + async def close(self) -> None: + if self.closed: + return + self.closed = True + for closer in (self._pubsub, self._client): + with contextlib.suppress(Exception): + await closer.aclose() + + +async def subscribe( + url: str, + channel: str, + client_factory: Callable[[], Any] | None = None, +) -> EventSubscription: + """Connect to Redis and subscribe to one channel. + + Raises when the broker cannot be reached (the socket handler answers a + clear close code instead of hanging); once subscribed, disconnects are + tolerated by `EventSubscription.events()`. + """ + factory = client_factory or (lambda: _default_redis(url)) + client = factory() + pubsub = client.pubsub() + try: + await pubsub.subscribe(channel) + except Exception: + with contextlib.suppress(Exception): + await pubsub.aclose() + with contextlib.suppress(Exception): + await client.aclose() + raise + return EventSubscription(pubsub, client) + + +__all__ = [ + "CHANNEL_PREFIX", + "DashboardEvent", + "EventBus", + "EventName", + "EventSubscription", + "NoopEventBus", + "RedisEventBus", + "channel_for", + "subscribe", +] diff --git a/backend/app/main.py b/backend/app/main.py index 0d26821..62613d8 100644 --- a/backend/app/main.py +++ b/backend/app/main.py @@ -20,6 +20,7 @@ from app.api.rescues import router as rescues_router from app.api.status import router as status_router from app.api.webhooks_twilio import router as twilio_router +from app.api.ws import router as ws_router from app.core.config import get_settings from app.core.logging import configure_logging from app.observability.tracing import configure_tracing, shutdown_tracing @@ -56,6 +57,7 @@ def create_app() -> FastAPI: allow_headers=["Authorization", "Content-Type"], ) app.include_router(health_router) + app.include_router(ws_router) app.include_router(twilio_router) app.include_router(status_router) app.include_router(auth_router) diff --git a/backend/app/runtime.py b/backend/app/runtime.py index c3cf73c..0e44047 100644 --- a/backend/app/runtime.py +++ b/backend/app/runtime.py @@ -18,6 +18,7 @@ from app.core.clock import Clock, DemoClock, SystemClock, redis_offset_source from app.core.config import Settings, get_settings from app.db.session import create_engine_and_session +from app.events import RedisEventBus from app.integrations.workforce.mock import MockWorkforceAdapter from app.services.orchestrator import RescueOrchestrator from app.workers.celery_scheduler import CeleryScheduler @@ -128,6 +129,11 @@ def build_runtime(settings: Settings, *, scheduler: Scheduler | None = None) -> scheduler=scheduler, clock=clock, interpreter=interpreter, + # The live dashboard channel (spec §7.5): the worker publishes every + # transition and the API relays it to the connected managers. Best + # effort by design — a broker outage degrades to the previous pull-only + # behaviour instead of failing a rescue. + events=RedisEventBus(settings.redis_url), ) for name, handler in orchestrator.task_handlers().items(): scheduler.register(name, handler) diff --git a/backend/app/services/orchestrator.py b/backend/app/services/orchestrator.py index 80b08a0..67b555c 100644 --- a/backend/app/services/orchestrator.py +++ b/backend/app/services/orchestrator.py @@ -41,6 +41,7 @@ from app.domain.quiet_hours import next_quiet_end, offers_allowed from app.domain.ranking import RankedCandidate, rank_candidates from app.domain.state_machine import SideEffect, State, StateMachineEvent, transition +from app.events import DashboardEvent, EventBus, EventName, NoopEventBus from app.integrations.workforce.mock import MockWorkforceAdapter from app.observability.redaction import mask_phone, redact_if_health from app.ports import Channel, Scheduler @@ -76,6 +77,7 @@ def __init__( clock: Clock, config: OrchestratorConfig | None = None, interpreter: MessageInterpreter | None = None, + events: EventBus | None = None, ) -> None: self._sessions = session_factory self._workforce = workforce @@ -84,6 +86,76 @@ def __init__( self._clock = clock self._config = config or OrchestratorConfig() self.interpreter = interpreter + # Live-channel port (spec §7.5): optional, no-op by default; a Redis + # outage never reaches the domain (the bus itself is best effort). + self._events = events if events is not None else NoopEventBus() + + # --- live events (spec §7.5, §7.6) --------------------------------------- + + async def _emit( + self, + name: EventName, + *, + location_id: str | None, + rescue_id: str | None = None, + shift_id: str | None = None, + origin: str = "system", + ) -> None: + """Publish one dashboard event, best effort (spec §9.3). + + Called right after the commit that produced the change, so a client + refetching on the event sees the new state. The payload carries ids + only — never a body, phone, name or health detail (spec §10) — and + a bus failure is logged and swallowed: the worker never fails + because the broker is down. + """ + if location_id is None: + return + try: + await self._events.publish( + DashboardEvent( + name=name, + location_id=location_id, + rescue_id=rescue_id, + shift_id=shift_id, + origin=origin, + at=_utc(self._clock.now()).isoformat(), + ) + ) + except Exception as error: # the rescue path never depends on the bus + structlog.get_logger(__name__).warning( + "dashboard_event_publish_failed", + name=name.value, + error=str(error)[:200], + ) + + async def _emit_message_sent( + self, + conversation_id: str | None, + employee_id: str | None, + rescue_id: str | None, + ) -> None: + """MESSAGE_SENT for the employee conversation just persisted.""" + location_id: str | None = None + if employee_id is not None: + location_id = await self._location_of(employee_id) + elif conversation_id is not None: + async with self._sessions() as session: + row = ( + await session.execute( + select(Conversation.employee_id).where( + Conversation.id == conversation_id + ) + ) + ).scalar_one_or_none() + if row is not None: + location_id = await self._location_of(row) + await self._emit( + EventName.MESSAGE_SENT, + location_id=location_id, + rescue_id=rescue_id, + origin="agent", + ) # --- inbound ------------------------------------------------------------- @@ -101,6 +173,12 @@ async def handle_inbound( if message_id is None: return # duplicate provider message: processed once (spec §7.4) + await self._emit( + EventName.MESSAGE_RECEIVED, + location_id=await self._location_of(employee_id), + origin=f"employee:{employee_id}", + ) + # Spec §9.3: a paused agent does nothing — the manager takes over. if await self._agent_is_paused(employee_id): await self._forward_to_manager_when_paused(employee_id, text) @@ -307,6 +385,13 @@ async def _accept_with_conditions( ) ) await session.commit() + await self._emit( + EventName.APPROVAL_REQUESTED, + location_id=case.location_id, + rescue_id=case.id, + shift_id=case.shift_id, + origin=f"employee:{employee_id}", + ) self._scheduler.schedule( _aware(case.deadline_at), "approval_timeout", {"case_id": case.id} ) @@ -584,6 +669,14 @@ async def _open_absence_case( # forever (spec §5.5; the transition is OPEN + DEADLINE_REACHED). self._scheduler.schedule(_aware(deadline), "rescue_deadline", {"case_id": case_id}) + await self._emit( + EventName.RESCUE_OPENED, + location_id=target.location_id, + rescue_id=case_id, + shift_id=target.id, + origin=f"employee:{employee_id}", + ) + await self._send_template( to=self._phone_of(employee), template_key="absence_confirm", @@ -699,6 +792,12 @@ async def _handle_confirmation(self, conversation_id: str, employee_id: str) -> if queued is None: await self._escalate(session, case, StateMachineEvent.WAVES_EXHAUSTED) await session.commit() + await self._emit( + EventName.CASE_ESCALATED, + location_id=case.location_id, + rescue_id=case.id, + shift_id=case.shift_id, + ) manager = await self._manager_for(case.location_id) if manager is not None and manager.get("phone_e164"): await self._send_template( @@ -725,6 +824,13 @@ async def _handle_confirmation(self, conversation_id: str, employee_id: str) -> # Persist the whole OPEN -> OFFERING transition before effects land. await session.commit() + if offered_count > 0: + await self._emit( + EventName.OFFERS_SENT, + location_id=case.location_id, + rescue_id=case.id, + shift_id=case.shift_id, + ) self._schedule_wave_tasks(case, now) # --- candidates and first wave -------------------------------------------- @@ -1027,6 +1133,13 @@ async def _try_accept_offer(self, employee_id: str) -> bool: ) ) await session.commit() + await self._emit( + EventName.APPROVAL_REQUESTED, + location_id=case.location_id, + rescue_id=case.id, + shift_id=case.shift_id, + origin=f"employee:{employee_id}", + ) manager = await self._manager_for(case.location_id) if manager is not None and manager.get("phone_e164"): await self._send_template( @@ -1117,6 +1230,13 @@ async def _try_accept_offer(self, employee_id: str) -> bool: manager = await self._manager_for(case.location_id) location_name, location_tz = await self._location_info(case.location_id) await session.commit() + await self._emit( + EventName.APPROVAL_REQUESTED, + location_id=case.location_id, + rescue_id=case.id, + shift_id=case.shift_id, + origin=f"employee:{employee_id}", + ) # Manager must decide before the rescue deadline (§4.2). self._scheduler.schedule( _aware(case.deadline_at), "approval_timeout", {"case_id": case.id} @@ -1151,6 +1271,20 @@ async def _try_accept_offer(self, employee_id: str) -> bool: ) ) await session.commit() + await self._emit( + EventName.OFFER_ACCEPTED, + location_id=case.location_id, + rescue_id=case.id, + shift_id=case.shift_id, + origin=f"employee:{employee_id}", + ) + await self._emit( + EventName.CASE_COVERED, + location_id=case.location_id, + rescue_id=case.id, + shift_id=case.shift_id, + origin=f"employee:{employee_id}", + ) # Effects after commit (spec §4.2). HRIS failures retry 3 times, # then the rescue escalates with a technical reason (spec §5.5). @@ -1186,6 +1320,12 @@ async def _try_accept_offer(self, employee_id: str) -> bool: ) ) await session.commit() + await self._emit( + EventName.CASE_ESCALATED, + location_id=case.location_id, + rescue_id=case.id, + shift_id=case.shift_id, + ) manager = await self._manager_for(case.location_id) location_name, location_tz = await self._location_info(case.location_id) if manager is not None and manager.get("phone_e164"): @@ -1304,6 +1444,12 @@ async def _try_decline_offer(self, employee_id: str) -> bool: ) ) await session.commit() + await self._emit( + EventName.OFFER_DECLINED, + location_id=await self._location_of(employee_id), + rescue_id=offer.rescue_id, + origin=f"employee:{employee_id}", + ) return True async def _try_withdraw(self, employee_id: str) -> bool: @@ -1368,6 +1514,12 @@ async def _resume_offering(self, case: RescueCase) -> None: if not next_candidates: await self._escalate(session, case, StateMachineEvent.WAVES_EXHAUSTED) await session.commit() + await self._emit( + EventName.CASE_ESCALATED, + location_id=case.location_id, + rescue_id=case.id, + shift_id=case.shift_id, + ) await self._notify_escalation(case) return @@ -1389,6 +1541,13 @@ async def _resume_offering(self, case: RescueCase) -> None: now=now, ) await session.commit() + if sent: + await self._emit( + EventName.OFFERS_SENT, + location_id=case.location_id, + rescue_id=case.id, + shift_id=case.shift_id, + ) if sent: self._scheduler.schedule( now + timedelta(minutes=self._config.wave_interval_minutes), @@ -1430,6 +1589,13 @@ async def _handle_retraction(self, employee_id: str) -> bool: ) ) await session.commit() + await self._emit( + EventName.APPROVAL_REQUESTED, + location_id=case.location_id, + rescue_id=case.id, + shift_id=case.shift_id, + origin=f"employee:{employee_id}", + ) manager = await self._manager_for(case.location_id) if manager is not None and manager.get("phone_e164"): @@ -1482,6 +1648,13 @@ async def decide_approval(self, approval_id: str, decision: str, decided_by: str ) case.status = result.new_state.value await session.commit() + await self._emit( + EventName.APPROVAL_DECIDED, + location_id=case.location_id, + rescue_id=case.id, + shift_id=case.shift_id, + origin=f"manager:{decided_by}", + ) return if approval.kind == "cancel_rescue": @@ -1503,6 +1676,20 @@ async def decide_approval(self, approval_id: str, decision: str, decided_by: str ) ) await session.commit() + await self._emit( + EventName.APPROVAL_DECIDED, + location_id=case.location_id, + rescue_id=case.id, + shift_id=case.shift_id, + origin=f"manager:{decided_by}", + ) + await self._emit( + EventName.CASE_CLOSED, + location_id=case.location_id, + rescue_id=case.id, + shift_id=case.shift_id, + origin=f"manager:{decided_by}", + ) return if approval.kind == "partial_coverage": @@ -1534,6 +1721,20 @@ async def decide_approval(self, approval_id: str, decision: str, decided_by: str ) ) await session.commit() + await self._emit( + EventName.APPROVAL_DECIDED, + location_id=case.location_id, + rescue_id=case.id, + shift_id=case.shift_id, + origin=f"manager:{decided_by}", + ) + await self._emit( + EventName.CASE_COVERED, + location_id=case.location_id, + rescue_id=case.id, + shift_id=case.shift_id, + origin=f"manager:{decided_by}", + ) # Effects after commit. if offer is not None and offer.employee_id: @@ -1606,6 +1807,13 @@ async def close_rescue(self, rescue_id: str, decided_by: str) -> None: ) ) await session.commit() + await self._emit( + EventName.CASE_CLOSED, + location_id=case.location_id, + rescue_id=case.id, + shift_id=case.shift_id, + origin=f"manager:{decided_by}", + ) async def _supersede_offers(self, session: AsyncSession, rescue_id: str) -> None: pending = ( @@ -1651,6 +1859,12 @@ async def _on_wave_timeout(self, payload: dict[str, Any]) -> None: if _aware(case.deadline_at) <= now: await self._escalate(session, case, StateMachineEvent.DEADLINE_REACHED) await session.commit() + await self._emit( + EventName.CASE_ESCALATED, + location_id=case.location_id, + rescue_id=case.id, + shift_id=case.shift_id, + ) await self._notify_escalation(case) return @@ -1668,6 +1882,12 @@ async def _on_wave_timeout(self, payload: dict[str, Any]) -> None: if not next_candidates: await self._escalate(session, case, StateMachineEvent.WAVES_EXHAUSTED) await session.commit() + await self._emit( + EventName.CASE_ESCALATED, + location_id=case.location_id, + rescue_id=case.id, + shift_id=case.shift_id, + ) await self._notify_escalation(case) return @@ -1678,7 +1898,7 @@ async def _on_wave_timeout(self, payload: dict[str, Any]) -> None: ).scalars().all() next_wave = max(last_wave) + 1 location_name, location_tz = await self._location_info(case.location_id) - await self._send_wave_offers( + sent = await self._send_wave_offers( session, case, shift, @@ -1689,6 +1909,13 @@ async def _on_wave_timeout(self, payload: dict[str, Any]) -> None: now=now, ) await session.commit() + if sent: + await self._emit( + EventName.OFFERS_SENT, + location_id=case.location_id, + rescue_id=case.id, + shift_id=case.shift_id, + ) self._schedule_wave_tasks(case, now) async def _on_deadline(self, payload: dict[str, Any]) -> None: @@ -1708,6 +1935,12 @@ async def _on_deadline(self, payload: dict[str, Any]) -> None: return await self._escalate(session, case, StateMachineEvent.DEADLINE_REACHED) await session.commit() + await self._emit( + EventName.CASE_ESCALATED, + location_id=case.location_id, + rescue_id=case.id, + shift_id=case.shift_id, + ) await self._notify_escalation(case) async def _on_approval_timeout(self, payload: dict[str, Any]) -> None: @@ -1765,7 +1998,7 @@ async def _on_send_wave(self, payload: dict[str, Any]) -> None: return ranked = await self._compute_candidates(case.location_id, shift, now) location_name, location_tz = await self._location_info(case.location_id) - await self._send_wave_offers( + sent = await self._send_wave_offers( session, case, shift, @@ -1776,6 +2009,13 @@ async def _on_send_wave(self, payload: dict[str, Any]) -> None: now=now, ) await session.commit() + if sent: + await self._emit( + EventName.OFFERS_SENT, + location_id=case.location_id, + rescue_id=case.id, + shift_id=case.shift_id, + ) async def _escalate( self, session: AsyncSession, case: RescueCase, event: StateMachineEvent @@ -2120,6 +2360,7 @@ async def _send_template( ) ) await session.commit() + await self._emit_message_sent(conversation_id, employee_id, rescue_id) async def _send_out_of_scope(self, conversation_id: str, employee_id: str) -> None: employee = await self._employee(employee_id) diff --git a/backend/tests/unit/api/test_ws.py b/backend/tests/unit/api/test_ws.py new file mode 100644 index 0000000..5f04d25 --- /dev/null +++ b/backend/tests/unit/api/test_ws.py @@ -0,0 +1,277 @@ +"""Unit tests for the live-event socket (spec §7.5, §7.6, §9.3). + +Hermetic: no broker, no real Redis — the subscribe seam is monkeypatched and +the manager check runs on SQLite. The close codes are the contract: the hook +in the dashboard clears the session on 4401 and keeps working otherwise. +""" + +import json +import os +import tempfile +from datetime import UTC, datetime, timedelta + +import jwt +import pytest +from fastapi import WebSocketDisconnect +from fastapi.testclient import TestClient +from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine + +from app.api.ws import ws_session_factory +from app.core.config import Settings, get_settings +from app.db.models import Base, Manager +from app.db.seed import DEMO_LOCATION_ID +from app.events import DashboardEvent, EventName +from app.main import create_app +from app.security.tokens import issue_token + +MANAGER_ROLE = "manager" + + +@pytest.fixture() +async def db_factory(): + fd, db_path = tempfile.mkstemp(suffix=".db") + os.close(fd) + engine = create_async_engine(f"sqlite+aiosqlite:///{db_path}") + factory = async_sessionmaker(engine, expire_on_commit=False) + async with engine.begin() as conn: + await conn.run_sync(Base.metadata.create_all) + async with factory() as session: + session.add( + Manager( + id="mgr_1", + name="Demo Manager", + email="manager@laterraza.demo", + phone_e164="+34600999001", + password_hash="x", + role=MANAGER_ROLE, + location_ids=[DEMO_LOCATION_ID], + ) + ) + session.add( + Manager( + id="mgr_other", + name="Other Manager", + email="other@laterraza.demo", + phone_e164="+34600999002", + password_hash="x", + role=MANAGER_ROLE, + location_ids=["loc_other"], + ) + ) + session.add( + Manager( + id="op_1", + name="Demo Operator", + email="operator@laterraza.demo", + phone_e164="+34600999003", + password_hash="x", + role="operator", + location_ids=[], + ) + ) + await session.commit() + yield factory + await engine.dispose() + + +@pytest.fixture() +def settings() -> Settings: + return Settings( + app_env="test", + jwt_secret="test-secret", + database_url="sqlite+aiosqlite:///unused.db", + redis_url="redis://broker.invalid:6379/0", + ) + + +@pytest.fixture() +def app(settings, db_factory): + application = create_app() + application.dependency_overrides[get_settings] = lambda: settings + application.dependency_overrides[ws_session_factory] = lambda: db_factory + return application + + +class FakeSubscription: + def __init__(self, events: list[DashboardEvent] | None = None) -> None: + self._events = events or [] + self.closed = False + + async def events(self): + for event in self._events: + yield event + + async def close(self) -> None: + self.closed = True + + +@pytest.fixture() +def installed(): + """Install a fake subscribe seam; returns a dict with what it saw.""" + seen: dict = {} + + def install(monkeypatch, subscription: FakeSubscription) -> dict: + from app.api import ws as ws_module + + async def fake_subscribe(url: str, channel: str) -> FakeSubscription: + seen["url"] = url + seen["channel"] = channel + seen["subscription"] = subscription + return subscription + + monkeypatch.setattr(ws_module, "subscribe", fake_subscribe) + return seen + + return install + + +def _event(**overrides) -> DashboardEvent: + base = dict( + name=EventName.OFFER_ACCEPTED, + location_id=DEMO_LOCATION_ID, + rescue_id="case_1", + shift_id="shift_1", + origin="employee:emp_1", + at="2026-10-03T14:40:00+00:00", + ) + base.update(overrides) + return DashboardEvent(**base) + + +def _connect(client: TestClient, token: str | None, location: str = DEMO_LOCATION_ID): + query = f"?token={token}" if token is not None else "" + return client.websocket_connect(f"/ws/locations/{location}{query}") + + +def _rejects_with(app, monkeypatch, token: str | None, code: int) -> None: + with ( + TestClient(app) as client, + pytest.raises(WebSocketDisconnect) as exc_info, + _connect(client, token) as websocket, + ): + websocket.receive_text() + assert exc_info.value.code == code + + +# --- authentication (4401) --------------------------------------------------- + + +def test_socket_rejects_a_missing_token(app, monkeypatch) -> None: + _rejects_with(app, monkeypatch, token=None, code=4401) + + +def test_socket_rejects_a_malformed_token(app, monkeypatch) -> None: + _rejects_with(app, monkeypatch, token="not-a-jwt", code=4401) + + +def test_socket_rejects_an_expired_token(app, monkeypatch, settings) -> None: + now = datetime.now(UTC) + expired = jwt.encode( + { + "sub": "mgr_1", + "role": MANAGER_ROLE, + "iat": now - timedelta(hours=2), + "exp": now - timedelta(hours=1), + }, + settings.jwt_secret, + algorithm="HS256", + ) + _rejects_with(app, monkeypatch, token=expired, code=4401) + + +def test_socket_rejects_a_wrong_role_token(app, monkeypatch, settings) -> None: + token = issue_token("mgr_1", "employee", settings)[0] + _rejects_with(app, monkeypatch, token=token, code=4401) + + +# --- authorization (4403) ---------------------------------------------------- + + +def test_socket_rejects_a_manager_without_rights_on_the_location( + app, monkeypatch, settings +) -> None: + token = issue_token("mgr_other", MANAGER_ROLE, settings)[0] + _rejects_with(app, monkeypatch, token=token, code=4403) + + +def test_socket_rejects_an_unknown_manager(app, monkeypatch, settings) -> None: + token = issue_token("mgr_ghost", MANAGER_ROLE, settings)[0] + _rejects_with(app, monkeypatch, token=token, code=4403) + + +# --- the happy path ---------------------------------------------------------- + + +def test_socket_forwards_a_published_event_to_a_valid_manager( + app, monkeypatch, installed, settings +) -> None: + subscription = FakeSubscription([_event()]) + seen = installed(monkeypatch, subscription) + token = issue_token("mgr_1", MANAGER_ROLE, settings)[0] + + with TestClient(app) as client, _connect(client, token) as websocket: + payload = json.loads(websocket.receive_text()) + + assert payload == { + "name": "OFFER_ACCEPTED", + "location_id": DEMO_LOCATION_ID, + "rescue_id": "case_1", + "shift_id": "shift_1", + "origin": "employee:emp_1", + "at": "2026-10-03T14:40:00+00:00", + } + assert seen["channel"] == f"shift_rescue:events:{DEMO_LOCATION_ID}" + assert subscription.closed # the subscription is released on disconnect + + +def test_socket_lets_the_operator_watch_any_location( + app, monkeypatch, installed, settings +) -> None: + installed(monkeypatch, FakeSubscription([])) + token = issue_token("op_1", "operator", settings)[0] + + with ( + TestClient(app) as client, + _connect(client, token, location="loc_any") as websocket, + pytest.raises(WebSocketDisconnect), + ): + websocket.receive_text() + + +def test_socket_relays_only_ids_never_payload_bodies( + app, monkeypatch, installed, settings +) -> None: + subscription = FakeSubscription([_event(origin="system")]) + installed(monkeypatch, subscription) + token = issue_token("mgr_1", MANAGER_ROLE, settings)[0] + + with TestClient(app) as client, _connect(client, token) as websocket: + raw = websocket.receive_text() + + assert "body" not in raw + assert "phone" not in raw + assert set(json.loads(raw)) <= { + "name", + "location_id", + "rescue_id", + "shift_id", + "origin", + "at", + } + + +# --- broker down (4503) ------------------------------------------------------ + + +def test_socket_closes_4503_when_the_broker_is_unavailable( + app, monkeypatch, settings +) -> None: + from app.api import ws as ws_module + + async def broken_subscribe(url: str, channel: str) -> FakeSubscription: + raise OSError("broker down") + + monkeypatch.setattr(ws_module, "subscribe", broken_subscribe) + + token = issue_token("mgr_1", MANAGER_ROLE, settings)[0] + _rejects_with(app, monkeypatch, token=token, code=4503) diff --git a/backend/tests/unit/test_events.py b/backend/tests/unit/test_events.py new file mode 100644 index 0000000..c196174 --- /dev/null +++ b/backend/tests/unit/test_events.py @@ -0,0 +1,208 @@ +"""Unit tests for the dashboard event bus (spec §7.5, §9.3, §10). + +Hermetic: every Redis interaction goes through fakes — a unit test never +talks to a real broker. +""" + +import json + +import pytest +from structlog.testing import capture_logs + +from app.events import ( + DashboardEvent, + EventName, + NoopEventBus, + RedisEventBus, + channel_for, + subscribe, +) + +EVENT_CHANNEL = "shift_rescue:events:loc_la_terraza" + + +def _event(**overrides) -> DashboardEvent: + base = dict( + name=EventName.OFFER_ACCEPTED, + location_id="loc_la_terraza", + rescue_id="case_1", + shift_id="shift_1", + origin="employee:emp_1", + at="2026-10-03T14:40:00+00:00", + ) + base.update(overrides) + return DashboardEvent(**base) + + +class FakePublisher: + """Stands in for `redis.asyncio.Redis` on the publish side.""" + + def __init__(self) -> None: + self.published: list[tuple[str, str]] = [] + self.fail = False + + async def publish(self, channel: str, payload: str) -> None: + if self.fail: + raise OSError("broker down") + self.published.append((channel, payload)) + + +class FakePubSub: + def __init__(self, messages: list) -> None: + self.messages = messages + self.subscribed: list[str] = [] + self.closed = False + self.fail_subscribe = False + + async def subscribe(self, channel: str) -> None: + if self.fail_subscribe: + raise OSError("broker down") + self.subscribed.append(channel) + + async def listen(self): + for message in self.messages: + yield message + + async def aclose(self) -> None: + self.closed = True + + +class FakeSubscriber: + def __init__(self, messages: list) -> None: + self._pubsub = FakePubSub(messages) + self.closed = False + + def pubsub(self) -> FakePubSub: + return self._pubsub + + async def aclose(self) -> None: + self.closed = True + + +# --- serialization --------------------------------------------------------- + + +def test_channel_for_uses_the_location_key() -> None: + assert channel_for("loc_la_terraza") == EVENT_CHANNEL + + +def test_event_json_round_trip() -> None: + event = _event() + assert DashboardEvent.from_json(event.to_json()) == event + + +def test_event_json_carries_exactly_the_six_allowed_fields() -> None: + # Privacy invariant (spec §10): ids, origin and time — nothing else can + # ride along because the dataclass is the whole payload. + payload = json.loads(_event().to_json()) + assert set(payload) == { + "name", + "location_id", + "rescue_id", + "shift_id", + "origin", + "at", + } + + +@pytest.mark.parametrize( + "raw", + ["not json", "[]", json.dumps({"name": "NOT_A_REAL_EVENT"}), json.dumps({})], +) +def test_from_json_is_total_on_garbage(raw: str) -> None: + assert DashboardEvent.from_json(raw) is None + + +def test_from_json_fills_defaults_for_optional_fields() -> None: + event = DashboardEvent.from_json( + json.dumps({"name": "RESCUE_OPENED", "location_id": "loc_1"}) + ) + assert event is not None + assert event.rescue_id is None + assert event.shift_id is None + assert event.origin == "system" + assert event.at # a timestamp is always present + + +# --- the buses --------------------------------------------------------------- + + +async def test_noop_bus_drops_events_without_failing() -> None: + await NoopEventBus().publish(_event()) + + +async def test_redis_bus_publishes_to_the_location_channel() -> None: + fake = FakePublisher() + bus = RedisEventBus("redis://localhost:6379/0", client_factory=lambda: fake) + + await bus.publish(_event()) + + assert fake.published == [(EVENT_CHANNEL, _event().to_json())] + + +async def test_redis_bus_reuses_one_client_and_closes_nothing() -> None: + fake = FakePublisher() + calls = {"n": 0} + + def factory() -> FakePublisher: + calls["n"] += 1 + return fake + + bus = RedisEventBus("redis://localhost:6379/0", client_factory=factory) + await bus.publish(_event()) + await bus.publish(_event()) + assert calls["n"] == 1 + + +async def test_redis_bus_swallows_a_broker_failure_and_logs() -> None: + fake = FakePublisher() + fake.fail = True + bus = RedisEventBus("redis://localhost:6379/0", client_factory=lambda: fake) + + with capture_logs() as logs: + await bus.publish(_event()) + + assert [log["event"] for log in logs] == ["dashboard_event_publish_failed"] + # A broken client is dropped so the next event rebuilds the connection. + await bus.publish(_event()) # must not raise either + assert fake.published == [] + + +# --- the subscriber ---------------------------------------------------------- + + +async def test_subscribe_yields_parsed_events_and_closes_cleanly() -> None: + event = _event() + client = FakeSubscriber( + [ + {"type": "subscribe"}, # control frame: skipped + {"type": "message", "data": event.to_json()}, + {"type": "message", "data": "garbage"}, # skipped, never crashes + {"type": "other", "data": event.to_json()}, + ] + ) + def factory() -> FakeSubscriber: + return client + + subscription = await subscribe("redis://x/0", EVENT_CHANNEL, client_factory=factory) + received = [e async for e in subscription.events()] + await subscription.close() + + assert client.pubsub().subscribed == [EVENT_CHANNEL] + assert received == [event] + assert client.pubsub().closed and client.closed + assert subscription.closed + + +async def test_subscribe_raises_when_the_broker_is_down() -> None: + client = FakeSubscriber([]) + client.pubsub().fail_subscribe = True + + def factory() -> FakeSubscriber: + return client + + with pytest.raises(OSError): + await subscribe("redis://x/0", EVENT_CHANNEL, client_factory=factory) + + # The failed subscribe still releases the half-open resources. + assert client.pubsub().closed and client.closed diff --git a/frontend/src/App.tsx b/frontend/src/App.tsx index be3a3b6..be24485 100644 --- a/frontend/src/App.tsx +++ b/frontend/src/App.tsx @@ -12,6 +12,7 @@ import { EvalsScreen } from './screens/EvalsScreen' import { SimulatorScreen } from './screens/SimulatorScreen' import { RequireAuth } from './components/RequireAuth' import { getSession } from './services/auth' +import { useLiveEvents } from './services/liveEvents' import { ApiError } from './services/apiClient' import { usePendingApprovals } from './services/hooks' @@ -60,6 +61,9 @@ function Shell({ now, onLogout }: { now?: Date; onLogout: () => void }) { const [selectedRescueId, setSelectedRescueId] = useState(null) const { approvals } = usePendingApprovals() const manager = getSession()?.manager + // Live channel (spec §7.5): the screens refetch when the worker reports a + // change made elsewhere. Its state is shown so it is never a silent promise. + const liveState = useLiveEvents() const navigate = (next: AppView) => { setView(next) @@ -78,6 +82,15 @@ function Shell({ now, onLogout }: { now?: Date; onLogout: () => void }) { onLogout={onLogout} />
+ {liveState !== 'live' ? ( +

+ + {liveState === 'connecting' ? 'Connecting to live updates…' : 'Live updates offline'} +

+ ) : null} {selectedRescueId ? ( ({ + getLocationId: vi.fn(async () => 'loc_1'), +})) + +class FakeSocket { + static instances: FakeSocket[] = [] + url: string + onopen: (() => void) | null = null + onmessage: ((event: { data: string }) => void) | null = null + onclose: ((event: { code: number }) => void) | null = null + onerror: (() => void) | null = null + closed = false + + constructor(url: string) { + this.url = url + FakeSocket.instances.push(this) + } + + close() { + this.closed = true + } + + emitOpen() { + this.onopen?.() + } + + emitEvent(name: string) { + this.onmessage?.({ data: JSON.stringify({ name, location_id: 'loc' }) }) + } + + emitClose(code = 1006) { + this.closed = true + this.onclose?.({ code }) + } +} + +function wrapper(client: QueryClient) { + return ({ children }: { children: ReactNode }) => ( + {children} + ) +} + +function session() { + setSession({ + accessToken: 'test-token', + manager: { id: 'mgr_1', name: 'Demo Manager', email: 'm@x.demo', role: 'manager', locationIds: ['loc_1'] }, + }) +} + +beforeEach(() => { + FakeSocket.instances = [] + vi.stubGlobal('WebSocket', FakeSocket as unknown as typeof WebSocket) + // The suite runs with VITE_USE_MOCK=true (offline dashboard): the live channel + // is deliberately inert there, so these cases opt out of mock mode. + vi.stubEnv('VITE_USE_MOCK', 'false') + localStorage.clear() +}) + +afterEach(() => { + vi.unstubAllGlobals() + vi.useRealTimers() + clearSession() +}) + +describe('socketUrl', () => { + it('carries the location and the token', () => { + expect(socketUrl('loc_1', 'tok')).toBe( + `${location.origin.replace('http', 'ws')}/ws/locations/loc_1?token=tok`, + ) + }) +}) + +describe('queryKeysForEvent', () => { + it('maps each event family to the keys the screens refetch', () => { + expect(queryKeysForEvent('CASE_ESCALATED')).toContainEqual(['rescues']) + expect(queryKeysForEvent('OFFER_ACCEPTED')).toContainEqual(['shifts']) + expect(queryKeysForEvent('MESSAGE_RECEIVED')).toContainEqual(['demo-thread']) + expect(queryKeysForEvent('APPROVAL_DECIDED')).toContainEqual(['approvals']) + }) + + it('an unknown event invalidates nothing', () => { + expect(queryKeysForEvent('SOMETHING_ELSE')).toEqual([]) + }) +}) + +describe('useLiveEvents', () => { + it('opens the authenticated socket and invalidates on an event', async () => { + session() + const client = new QueryClient() + const invalidate = vi.spyOn(client, 'invalidateQueries') + + const { result } = renderHook(() => useLiveEvents(), { wrapper: wrapper(client) }) + + await waitFor(() => expect(FakeSocket.instances.length).toBe(1)) + expect(FakeSocket.instances[0].url).toContain('/ws/locations/loc_1?token=test-token') + + FakeSocket.instances[0].emitOpen() + await waitFor(() => expect(result.current).toBe('live')) + + FakeSocket.instances[0].emitEvent('OFFER_ACCEPTED') + await waitFor(() => + expect(invalidate).toHaveBeenCalledWith({ queryKey: ['rescues'] }), + ) + }) + + it('clears the session when the server closes with 4401', async () => { + session() + const client = new QueryClient() + renderHook(() => useLiveEvents(), { wrapper: wrapper(client) }) + + await waitFor(() => expect(FakeSocket.instances.length).toBe(1)) + FakeSocket.instances[0].emitClose(4401) + + await waitFor(() => expect(localStorage.getItem('shift-rescue.session')).toBeNull()) + }) + + it('reconnects with backoff after a lost connection', async () => { + vi.useFakeTimers() + session() + const client = new QueryClient() + renderHook(() => useLiveEvents(), { wrapper: wrapper(client) }) + + await vi.waitFor(() => expect(FakeSocket.instances.length).toBe(1)) + FakeSocket.instances[0].emitClose(1006) + + await vi.advanceTimersByTimeAsync(1100) + await vi.waitFor(() => expect(FakeSocket.instances.length).toBe(2)) + }) + + it('closes the socket on unmount without reconnecting', async () => { + session() + const client = new QueryClient() + const { unmount } = renderHook(() => useLiveEvents(), { wrapper: wrapper(client) }) + + await waitFor(() => expect(FakeSocket.instances.length).toBe(1)) + const socket = FakeSocket.instances[0] + unmount() + + expect(socket.closed).toBe(true) + }) + + it('stays offline without a session', async () => { + const client = new QueryClient() + const { result } = renderHook(() => useLiveEvents(), { wrapper: wrapper(client) }) + + await waitFor(() => expect(result.current).toBe('offline')) + expect(FakeSocket.instances).toHaveLength(0) + }) +}) diff --git a/frontend/src/services/liveEvents.ts b/frontend/src/services/liveEvents.ts new file mode 100644 index 0000000..5bebbaf --- /dev/null +++ b/frontend/src/services/liveEvents.ts @@ -0,0 +1,168 @@ +/** + * Live dashboard channel (spec §7.5 `WS /ws/locations/{id}`, §7.6). + * + * The screens keep rendering from the API; this channel only tells them what to + * refetch. One event type maps to the same query keys the actions already + * invalidate, so nothing is duplicated in the client cache. + * + * It degrades exactly like the rest of the dashboard: a socket that cannot open, + * a broker that is down or a token that expired never breaks a screen. A 4401 + * means the session is gone, which is handled the same way the HTTP client + * handles it (clear the session and let the guard show the login). + */ + +import { useQueryClient } from '@tanstack/react-query' +import { useEffect, useRef, useState } from 'react' + +import { getLocationId } from './api' +import { clearSession, getToken } from './auth' +import { isMockMode } from './dataSource' + +export type LiveState = 'connecting' | 'live' | 'offline' + +/** Close codes the API uses (see `app/api/ws.py`). */ +const CLOSE_UNAUTHORIZED = 4401 +const RECONNECT_MIN_MS = 1000 +const RECONNECT_MAX_MS = 15000 + +/** Query prefix keys each event implies, mirroring the action invalidations. */ +const RESCUE_KEYS = [ + ['rescues'], + ['shifts'], + ['approvals'], + ['conversations'], + ['demo-employees'], +] +const MESSAGE_KEYS = [['conversations'], ['demo-thread'], ['metrics'], ['agent-decisions']] +const APPROVAL_KEYS = [['approvals'], ['rescues'], ['shifts']] + +const EVENT_KEYS: Record = { + RESCUE_OPENED: RESCUE_KEYS, + OFFERS_SENT: RESCUE_KEYS, + OFFER_ACCEPTED: RESCUE_KEYS, + OFFER_DECLINED: RESCUE_KEYS, + CASE_COVERED: RESCUE_KEYS, + CASE_ESCALATED: RESCUE_KEYS, + CASE_CLOSED: RESCUE_KEYS, + MESSAGE_RECEIVED: MESSAGE_KEYS, + MESSAGE_SENT: MESSAGE_KEYS, + APPROVAL_REQUESTED: APPROVAL_KEYS, + APPROVAL_DECIDED: APPROVAL_KEYS, +} + +/** The query keys an event name invalidates; unknown events invalidate nothing. */ +export function queryKeysForEvent(name: string): string[][] { + return EVENT_KEYS[name] ?? [] +} + +/** `ws(s)://host` for the same origin, or the configured override. */ +export function socketBaseUrl(): string { + const configured = import.meta.env.VITE_WS_BASE_URL + if (configured) { + return configured.replace(/\/$/, '') + } + const protocol = window.location.protocol === 'https:' ? 'wss:' : 'ws:' + return `${protocol}//${window.location.host}` +} + +/** The socket URL for a location and token (exported for the tests). */ +export function socketUrl(locationId: string, token: string): string { + return `${socketBaseUrl()}/ws/locations/${encodeURIComponent(locationId)}?token=${encodeURIComponent(token)}` +} + +export function useLiveEvents(enabled = true): LiveState { + const queryClient = useQueryClient() + const [state, setState] = useState('offline') + const retries = useRef(0) + + useEffect(() => { + if (!enabled || isMockMode()) { + // Mock mode is the offline dashboard: the channel stays inert. The state + // already starts as 'offline', so nothing is set synchronously here. + return + } + let socket: WebSocket | null = null + let timer: ReturnType | undefined + let closed = false + + const connect = async () => { + const token = getToken() + if (closed) return + if (!token) { + setState('offline') + return + } + let locationId: string + try { + locationId = await getLocationId() + } catch { + // The dashboard will surface the API failure; the channel just waits. + setState('offline') + return + } + if (closed) return + + setState('connecting') + try { + socket = new WebSocket(socketUrl(locationId, token)) + } catch { + scheduleReconnect() + return + } + + socket.onopen = () => { + retries.current = 0 + setState('live') + } + socket.onmessage = (message) => { + let name: string | undefined + try { + name = (JSON.parse(String(message.data)) as { name?: string }).name + } catch { + return // a frame the API never sends: ignore it, never crash + } + if (!name) return + for (const queryKey of queryKeysForEvent(name)) { + void queryClient.invalidateQueries({ queryKey }) + } + } + socket.onclose = (event) => { + socket = null + setState('offline') + if (event.code === CLOSE_UNAUTHORIZED) { + // The session is gone: same handling as the HTTP 401 path. + clearSession() + return + } + scheduleReconnect() + } + socket.onerror = () => { + // onclose follows and owns the reconnect. + } + } + + const scheduleReconnect = () => { + if (closed) return + const delay = Math.min(RECONNECT_MIN_MS * 2 ** retries.current, RECONNECT_MAX_MS) + retries.current += 1 + timer = setTimeout(() => void connect(), delay) + } + + void connect() + + return () => { + closed = true + if (timer !== undefined) clearTimeout(timer) + if (socket !== null) { + // Detach the handlers: unmount is not a connection failure. + socket.onclose = null + socket.onmessage = null + socket.onopen = null + socket.onerror = null + socket.close() + } + } + }, [enabled, queryClient]) + + return state +} diff --git a/frontend/vite.config.ts b/frontend/vite.config.ts index e431a82..c7bb513 100644 --- a/frontend/vite.config.ts +++ b/frontend/vite.config.ts @@ -20,6 +20,14 @@ export default defineConfig({ target: 'http://localhost:8000', changeOrigin: true, }, + // Live dashboard channel (spec §7.5): a WebSocket needs the upgrade + // forwarded too, or the handshake never completes and the dashboard shows + // "connecting" forever. + '/ws': { + target: 'ws://localhost:8000', + ws: true, + changeOrigin: true, + }, }, }, test: { diff --git a/odd/tasks/realtime-events.md b/odd/tasks/realtime-events.md new file mode 100644 index 0000000..318423e --- /dev/null +++ b/odd/tasks/realtime-events.md @@ -0,0 +1,122 @@ +# Feature: realtime-events (WebSocket `/ws/locations/{id}`) + +**Status**: in progress +**Branch**: `feature/realtime-and-polish` +**Spec references**: §7.5 (`WS /ws/locations/{id}`), §7.6 (the dashboard updates itself) +**ADRs**: ADR-003 (single EC2, Caddy already routes `/ws/*`) + +## Problem + +The dashboard is pull-only. The reply to the manager's *own* action appears (the +Simulator polls after a send, actions invalidate queries), but a change produced +**elsewhere** — a candidate accepting, a case escalating, a rescue closed by the +worker — stays invisible until the manager refocuses the tab, acts, or reloads. +The user called this out ("no es tiempo real") and agreed to defer it; it is the +last thing that keeps the demo from feeling finished. + +Spec §7.5 defines the endpoint for it, Caddy already forwards `/ws/*`, and the +frontend already has one place per screen that knows which queries to refetch, +so the work is bounded. + +## Decisions (fixed, do not re-litigate) + +1. **Events travel through Redis pub/sub**, the broker the project already runs + for Celery. The worker is where the domain changes, and the API process is + what holds the sockets; a Redis channel (`shift_rescue:events:`) + is the smallest correct bridge. No polling of the database, no second broker. +2. **The domain publishes one small, typed event per meaningful change**: + `RESCUE_OPENED`, `OFFERS_SENT`, `OFFER_ACCEPTED`, `OFFER_DECLINED`, `CASE_COVERED`, + `CASE_ESCALATED`, `CASE_CLOSED`, `MESSAGE_RECEIVED`, `MESSAGE_SENT`, + `APPROVAL_REQUESTED`, `APPROVAL_DECIDED`. Each carries the rescue/shift ids and + the location, and **never** a message body, a phone number, a name or any health + detail. +3. **Publishing is best effort and never blocks a rescue**: a Redis failure logs + and is swallowed, exactly like the runtime snapshot and the tracing bootstrap. + The event bus is an injected port, so tests use a fake and the orchestrator + keeps working when the broker is down (spec §9.3). +4. **The socket is authenticated**: the browser cannot set headers on a + WebSocket, so the manager's JWT travels as a query parameter and is validated + before the socket is accepted; an invalid or missing token closes with 4401 + without subscribing. Only the client's own location is streamed. +5. **The frontend refetches, it does not re-render from the payload.** One hook + (`useLiveEvents`) maps each event type to the query keys it invalidates — the + same keys the actions already invalidate — so the screens keep rendering from + the API and nothing is duplicated in the client cache. + +## Tasks + +### T1 — Event bus (backend) +`app/events.py`: a typed `DashboardEvent` (name, location_id, rescue_id, shift_id, +origin, at) and an `EventBus` port with two implementations — a Redis publisher +(worker side) and a Redis subscriber/stream (API side) — plus a no-op fake for +tests. Publishing failures log and return. + +### T2 — The domain publishes (backend) +`RescueOrchestrator` publishes at every transition listed above, from the single +places it already commits them (the state machine transitions, the wave send, the +idempotent inbound handler). A test proves each transition emits its event, and +that no payload contains a body, a phone, a name or a health flag. + +### T3 — The endpoint (`app/api/ws.py`) +`WS /ws/locations/{location_id}?token=…`: validate the JWT (and the manager's +right to that location), accept, subscribe to the channel, forward events as +JSON, and clean up on disconnect. Unauthenticated or malformed → close 4401. +Redis unavailable → close with a clear code instead of hanging. Registered in +`app/main.py` (also in demo mode). + +### T4 — The frontend hook (`services/liveEvents.ts` + screens) +`useLiveEvents(locationId)` opens the socket (reconnecting with backoff while the +tab is visible, closing on unmount), and invalidates the query keys that match the +event: rescues, shifts, approvals, conversations, the open rescue detail, metrics, +interpretations. A tiny visible indicator ("live" / "reconnecting") so the state is +honest, and no crash when the socket cannot open (the dashboard keeps working as +today). + +### T5 — Tests and docs +- Backend: the event bus publishes/serializes correctly; each domain transition + emits its event and respects the payload rules; the socket rejects an + unauthenticated or malformed token, accepts a valid one and forwards a published + event, and closes cleanly. +- Frontend: the hook connects with the token, invalidates the right keys per event + type, reconnects with backoff and stops on unmount; screens render the indicator. +- `docs/runbook.md`: how to check the live channel; `docs/SHIFT_RESCUE_SPEC.md` + §7.5 compliance note; evidence in this document. + +### Parent verification (real stack, two browsers) + +| Check | Result | +| --- | --- | +| Events reach Redis | `PSUBSCRIBE shift_rescue:events:*` saw `MESSAGE_RECEIVED` (origin `employee:emp_01_kitchen`) and `MESSAGE_SENT` (origin `agent`) | +| Payload privacy | the frame carries only `name`, `location_id`, `rescue_id`, `shift_id`, `origin`, `at` — no body, no phone, no name | +| Two-browser acceptance (criterion 1) | with the demo reset and one browser on **Today**, a rescue produced from a **second** browser moved the first one **2.5 s later without a reload** (the row gained `View detail` as the case opened) | +| Socket state indicator | no warning rendered, i.e. the channel was live | + +### Two defects found while verifying (both fixed here) + +1. **The worker never published.** `build_runtime` built the orchestrator without + an event bus, so every transition went to the no-op: the channel was silent with + the socket perfectly healthy. `app/runtime.py` now injects + `RedisEventBus(settings.redis_url)`. +2. **The dev proxy did not forward `/ws`.** A WebSocket needs the upgrade + forwarded (`ws: true`), so the handshake never completed and the dashboard sat + on "Connecting to live updates…" forever. `frontend/vite.config.ts` proxies + `/ws` now; `infra/Caddyfile` already routed it. +3. The first version of the endpoint deadlocked the suite: it awaited + `websocket.receive()` while the client waited for the server, and cancelling a + `receive()` blocked in the ASGI portal hangs the test client. The relay and the + listener now run together and the socket is closed **before** anything is + cancelled. + +## Acceptance criteria + +1. With two browsers open on the dashboard, accepting an offer from the Simulator + moves the other one's Today board without a reload. +2. An unauthenticated socket is refused and never receives events. +3. Killing Redis makes the dashboard behave exactly as today (no crash, no hung + socket) and the worker keeps processing messages. +4. No event payload carries a message body, phone, name or health detail. +5. Suites green: backend, frontend, ruff, mypy, oxlint, build, tsc. + +## Verification evidence + +_Pending._ From b241c1f73de15919b4797d86d5ee1ff6f4e851f7 Mon Sep 17 00:00:00 2001 From: albert Date: Mon, 28 Sep 2026 18:39:04 +0200 Subject: [PATCH 2/6] feat(agent): record the notices that leave no trace MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two kinds of message reached a real phone and left nothing behind, so the dashboard could not tell the story and the audit trail had holes. - The manager's notices (escalation, coverage) go to a phone, not to an employee conversation, so nothing was persisted: the timeline showed a rescue that escalated "by itself". They now pass through one helper that sends and records `MANAGER_NOTIFIED` on the case with the template key, or `MANAGER_NOTIFY_SKIPPED` with the reason when the manager has no reachable phone — the attempt is recorded, never a delivery claim. The rescue timeline labels both ("Manager notified" / "Manager could not be reached"). - The race loser's "already covered" notice is an employee message and now carries the conversation, so it appears in their thread instead of looking unanswered. Verified live: an escalation produced `MANAGER_NOTIFIED {template: manager_escalated}` on the case. 496 backend tests, 196 frontend tests, ruff, mypy, oxlint, build and tsc clean. --- backend/app/services/orchestrator.py | 55 +++++++++--- .../tests/unit/services/test_orchestrator.py | 87 ++++++++++++++++++- frontend/src/domain/rescue.ts | 2 + frontend/src/domain/types.ts | 4 + frontend/src/screens/RescueDetailScreen.tsx | 1 + 5 files changed, 134 insertions(+), 15 deletions(-) diff --git a/backend/app/services/orchestrator.py b/backend/app/services/orchestrator.py index 67b555c..4689108 100644 --- a/backend/app/services/orchestrator.py +++ b/backend/app/services/orchestrator.py @@ -1140,16 +1140,14 @@ async def _try_accept_offer(self, employee_id: str) -> bool: shift_id=case.shift_id, origin=f"employee:{employee_id}", ) - manager = await self._manager_for(case.location_id) - if manager is not None and manager.get("phone_e164"): - await self._send_template( - to=manager["phone_e164"], - template_key="manager_covered", - employee_name=await self._employee_name(employee_id), - role="—", - start="—", - end="—", - ) + await self._notify_manager( + case, + "manager_covered", + employee_name=await self._employee_name(employee_id), + role="—", + start="—", + end="—", + ) return True if case.status != State.OFFERING.value: @@ -1367,6 +1365,8 @@ async def _reply_already_covered(self, employee_id: str) -> None: await self._send_template( to=self._phone_of(employee), template_key="offer_already_covered", + conversation_id=f"conv_twilio_{self._phone_of(employee)}", + employee_id=employee_id, employee_name=employee["full_name"], ) @@ -2038,15 +2038,42 @@ async def _escalate( ) ) - async def _notify_escalation(self, case: RescueCase) -> None: + async def _notify_manager(self, case: RescueCase, template_key: str, **params: Any) -> None: + """Tell the manager something about a rescue, and record that we did. + + The message goes to a phone, not to an employee conversation, so without + an audit event the dashboard timeline would show a case that escalated + "by itself" with no trace of the notice. Best effort, like every + outbound send: what is recorded is the attempt, never a delivery claim. + """ manager = await self._manager_for(case.location_id) if manager is None or not manager.get("phone_e164"): + await self._audit_notice(case, template_key, "no manager phone") return + await self._send_template(to=manager["phone_e164"], template_key=template_key, **params) + await self._audit_notice(case, template_key, None) + + async def _audit_notice( + self, case: RescueCase, template_key: str, skipped_reason: str | None + ) -> None: + async with self._sessions() as session: + session.add( + AuditEvent( + id=f"audit_{uuid4().hex}", + rescue_id=case.id, + type="MANAGER_NOTIFIED" if skipped_reason is None else "MANAGER_NOTIFY_SKIPPED", + payload={"template": template_key, "reason": skipped_reason}, + actor="system", + ) + ) + await session.commit() + + async def _notify_escalation(self, case: RescueCase) -> None: shift = await self._workforce.get_shift(case.shift_id) location_name, location_tz = await self._location_info(case.location_id) - await self._send_template( - to=manager["phone_e164"], - template_key="manager_escalated", + await self._notify_manager( + case, + "manager_escalated", role=self._role_label(shift.role) if shift else "—", start=self._fmt(shift.starts_at, location_tz) if shift else "—", end=self._fmt(shift.ends_at, location_tz) if shift else "—", diff --git a/backend/tests/unit/services/test_orchestrator.py b/backend/tests/unit/services/test_orchestrator.py index 99614cb..c9c94c7 100644 --- a/backend/tests/unit/services/test_orchestrator.py +++ b/backend/tests/unit/services/test_orchestrator.py @@ -9,7 +9,15 @@ from sqlalchemy import func, select from app.agent.interpreter import MessageInterpreter -from app.db.models import ApprovalRequest, AuditEvent, Message, Offer, RescueCase, Shift +from app.db.models import ( + ApprovalRequest, + AuditEvent, + Manager, + Message, + Offer, + RescueCase, + Shift, +) from app.db.seed import DEMO_LOCATION_ID from app.services.orchestrator import RECENT_CASE_WINDOW_HOURS from tests.unit.services.helpers import build_world, run_to_offering @@ -744,3 +752,80 @@ async def test_an_offer_for_a_shift_that_already_ended_is_not_acceptable() -> No assert stale, "the stale offer must be audited" assert stored.status == "CANCELLED" assert world.channel.with_template("offer_already_covered") + + +async def test_manager_notices_are_recorded_on_the_case(world, db) -> None: + """A notice to a phone left no trace: the timeline showed a rescue that + escalated "by itself". Every manager notice is now audited on the case.""" + await world.orchestrator.handle_inbound( + conversation_id=CONVERSATION, + employee_id="emp_01_floor", + provider_message_id=PROVIDER_ID, + text="me encuentro fatal, hoy no puedo ir", + ) + world.clock.advance(timedelta(minutes=11)) + await world.scheduler.run_due(world.clock.now()) + + async with db() as session: + events = (await session.execute(select(AuditEvent))).scalars().all() + + notices = [e for e in events if e.type == "MANAGER_NOTIFIED"] + assert notices, "the escalation notice must be recorded" + assert notices[0].payload["template"] == "manager_escalated" + assert world.channel.with_template("manager_escalated") + + +async def test_a_manager_without_a_phone_is_audited_as_skipped(world, db) -> None: + """The dashboard should show that nobody was reachable, not silence.""" + async with db() as session: + manager = (await session.execute(select(Manager))).scalar_one() + manager.phone_e164 = None + await session.commit() + + await world.orchestrator.handle_inbound( + conversation_id=CONVERSATION, + employee_id="emp_01_floor", + provider_message_id=PROVIDER_ID, + text="me encuentro fatal, hoy no puedo ir", + ) + world.clock.advance(timedelta(minutes=11)) + await world.scheduler.run_due(world.clock.now()) + + async with db() as session: + events = (await session.execute(select(AuditEvent))).scalars().all() + + skipped = [e for e in events if e.type == "MANAGER_NOTIFY_SKIPPED"] + assert skipped, "an unreachable manager must be visible in the timeline" + assert skipped[0].payload["reason"] == "no manager phone" + + +async def test_the_race_loser_notice_is_stored_in_their_conversation(world, db) -> None: + """It was sent and never recorded, so the thread looked unanswered.""" + offers = await run_to_offering(world) + winner, loser = offers[0], offers[1] + + await world.orchestrator.handle_inbound( + conversation_id=f"conv_{winner.employee_id}", + employee_id=winner.employee_id, + provider_message_id="accept_winner", + text="sí", + ) + await world.orchestrator.handle_inbound( + conversation_id=f"conv_{loser.employee_id}", + employee_id=loser.employee_id, + provider_message_id="accept_loser", + text="sí", + ) + + async with db() as session: + stored = ( + await session.execute( + select(Message).where( + Message.template_key == "offer_already_covered", + Message.direction == "outbound", + ) + ) + ).scalars().all() + + assert stored, "the loser's notice must be in their conversation" + assert "ya se ha cubierto" in stored[0].body_redacted diff --git a/frontend/src/domain/rescue.ts b/frontend/src/domain/rescue.ts index 3ced33b..a1c66e3 100644 --- a/frontend/src/domain/rescue.ts +++ b/frontend/src/domain/rescue.ts @@ -23,6 +23,8 @@ const eventLabels: Record = { SHIFT_ASSIGNED: 'Shift assigned in HRIS', ESCALATED: 'Escalated to manager', CANCELLED: 'Rescue cancelled', + MANAGER_NOTIFIED: 'Manager notified', + MANAGER_NOTIFY_SKIPPED: 'Manager could not be reached', } export function eventLabel(type: AuditEventType): string { diff --git a/frontend/src/domain/types.ts b/frontend/src/domain/types.ts index c8f6b96..0c07354 100644 --- a/frontend/src/domain/types.ts +++ b/frontend/src/domain/types.ts @@ -68,6 +68,10 @@ export type AuditEventType = | 'SHIFT_ASSIGNED' | 'ESCALATED' | 'CANCELLED' + /** A notice went (or could not go) to the manager's phone; recorded on the + * case so the timeline shows that somebody was told. */ + | 'MANAGER_NOTIFIED' + | 'MANAGER_NOTIFY_SKIPPED' export interface AuditEvent { id: string diff --git a/frontend/src/screens/RescueDetailScreen.tsx b/frontend/src/screens/RescueDetailScreen.tsx index 3b1b2e6..53162b2 100644 --- a/frontend/src/screens/RescueDetailScreen.tsx +++ b/frontend/src/screens/RescueDetailScreen.tsx @@ -35,6 +35,7 @@ const eventDotClasses: Partial> = { OFFER_ACCEPTED: 'bg-green-accent', OFFER_DECLINED: 'bg-error', ESCALATED: 'bg-error', + MANAGER_NOTIFY_SKIPPED: 'bg-error', } function offerStatusFor(candidate: CandidateResult, offers: Offer[]): OfferStatus | undefined { From c7680852dc7aca97caf4772aaff9df1aee803530 Mon Sep 17 00:00:00 2001 From: albert Date: Mon, 28 Sep 2026 18:44:25 +0200 Subject: [PATCH 3/6] feat(seed): keep a demo day alive at any hour MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The rotation has fixed windows, so a demo late in the evening found today's shifts finished and tomorrow's outside the roster's "today": nobody could report an absence and the board looked dead. That is the trap the user hit ("no hay ninguno que ponga Starts at"). When the seeded day is dry the seed now anchors two shifts to the current hour instead — one in progress, one starting soon — assigned to employees the rotation left free that day, and only for the window that is actually missing (at 03:00 nothing is in progress; at 23:30 the bar close already covers "now" and only the upcoming one is added). During working hours nothing changes, so the normal case stays exactly as it was, and the pool is taken from the rotation itself: an invented employee id would schedule a shift for somebody who does not exist. Verified live after a reseed: 5 shifts in progress and 2 starting within six hours. 498 backend tests, ruff and mypy clean. --- backend/app/db/seed.py | 60 +++++++++++++++++++++++++++++++-- backend/tests/unit/test_seed.py | 37 ++++++++++++++++++++ 2 files changed, 95 insertions(+), 2 deletions(-) diff --git a/backend/app/db/seed.py b/backend/app/db/seed.py index 14e4b40..67dcd06 100644 --- a/backend/app/db/seed.py +++ b/backend/app/db/seed.py @@ -121,13 +121,69 @@ def _shift( ) -def _shifts() -> list[Shift]: +def _demo_anchors(shifts: list[Shift], now: datetime) -> list[Shift]: + """Two extra shifts that keep a demo day usable at any hour. + + The rotation has fixed windows (07:00, 15:00, 23:00...), so a demo late in + the evening finds today's shifts already finished and tomorrow's outside the + roster's "today": the Simulator then has nobody who can report an absence and + the board looks dead. When the day is dry, two shifts are anchored to the + current hour instead — one in progress, one starting soon — assigned to + employees who are free that day. During working hours nothing is added, so + the normal case stays exactly as before. + """ + day_zero = [ + shift for shift in shifts if shift.starts_at.date() == SEED_START_DATE.date() + ] + in_progress = [s for s in day_zero if s.starts_at <= now <= s.ends_at] + starting_later = [s for s in day_zero if s.starts_at > now] + if in_progress and starting_later: + return [] + + # The pool is exactly the employees the rotation uses: inventing ids here + # would schedule a shift for somebody who does not exist. + pool = sorted( + {shift.employee_id for shift in shifts if shift.employee_id is not None} + ) + assigned = {shift.employee_id for shift in day_zero if shift.employee_id is not None} + available = [employee for employee in pool if employee not in assigned] + + date = SEED_START_DATE.date().isoformat() + anchors: list[Shift] = [] + for slot, (starts, ends) in enumerate( + ( + (now - timedelta(hours=1), now + timedelta(hours=6)), + (now + timedelta(hours=2), now + timedelta(hours=10)), + ) + ): + if (slot == 0 and in_progress) or (slot == 1 and starting_later): + continue + if not available: + break + employee_id = available.pop(0) + role = employee_id.split("_", 2)[-1].rsplit("_", 1)[0] if "_" in employee_id else "floor" + anchors.append( + Shift( + id=f"shift_lt_{date}_{role}_anchor{slot + 1}", + location_id=DEMO_LOCATION_ID, + role=role, + starts_at=starts, + ends_at=ends, + employee_id=employee_id, + status="scheduled", + ) + ) + return anchors + + +def _shifts(now: datetime | None = None) -> list[Shift]: shifts: list[Shift] = [] zone = SEED_START_DATE.replace(tzinfo=UTC) def at(day_offset: int, hour: int) -> datetime: return zone + timedelta(days=day_offset, hours=hour) + moment = now if now is not None else datetime.now(UTC) for d in range(SEED_DAYS): weekend = (SEED_START_DATE + timedelta(days=d)).weekday() >= 5 @@ -184,7 +240,7 @@ def at(day_offset: int, hour: int) -> datetime: supervisor_id = f"emp_{d % 2 + 25:02d}_supervisor" shifts.append(_shift("supervisor", d, "main", at(d, 11), at(d, 19), supervisor_id)) - return shifts + return shifts + _demo_anchors(shifts, moment) def _availability_blocks() -> list[AvailabilityBlock]: diff --git a/backend/tests/unit/test_seed.py b/backend/tests/unit/test_seed.py index 1abc4ab..bb608fc 100644 --- a/backend/tests/unit/test_seed.py +++ b/backend/tests/unit/test_seed.py @@ -159,3 +159,40 @@ async def test_seed_clears_previous_shifts_before_reinserting(session_factory) - await seed_database(session) shifts = (await session.execute(select(Shift))).scalars().all() assert all(s.id != "stale_1" for s in shifts) + + +def test_a_dry_day_gets_the_demo_anchors_it_needs() -> None: + """At any hour a seeded demo must have somebody who can act now. + + The rotation has fixed windows, so outside working hours the day can be dry: + at 03:00 nothing is in progress, and at 23:30 nothing starts later. Whatever + the hour, the seed guarantees both a shift in progress and something ahead, + using employees that actually exist. + """ + from app.db.seed import SEED_START_DATE, _shifts + + for hour, minute in ((3, 0), (23, 30)): + moment = SEED_START_DATE.replace(hour=hour, minute=minute, tzinfo=UTC) + shifts = _shifts(now=moment) + pool = {shift.employee_id for shift in shifts} + + assert any( + shift.starts_at <= moment <= shift.ends_at for shift in shifts + ), f"nobody on shift at {hour:02d}:{minute:02d}" + assert any( + shift.starts_at > moment for shift in shifts + ), f"nothing ahead at {hour:02d}:{minute:02d}" + anchors = [shift for shift in shifts if "anchor" in shift.id] + assert all(len(shift.id) <= 64 for shift in anchors) + # An anchor must never schedule somebody who is not on the roster. + assert all(shift.employee_id in pool for shift in anchors) + + +def test_a_normal_working_hour_needs_no_anchors() -> None: + """During the day the rotation already covers past, present and future.""" + from app.db.seed import SEED_START_DATE, _shifts + + midday = SEED_START_DATE.replace(hour=10, tzinfo=UTC) + shifts = _shifts(now=midday) + + assert not [s for s in shifts if "anchor" in s.id] From c62b58157c314dd9a3bb3a24ed21a61fa096ebf5 Mon Sep 17 00:00:00 2001 From: albert Date: Mon, 28 Sep 2026 19:09:05 +0200 Subject: [PATCH 4/6] feat(api): the manager can mark an absence, plus three small cleanups MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The last endpoint of spec §7.5 that was missing, and three loose ends. - `POST /api/shifts/{id}/absence`: the manager marks an absence and the rescue opens immediately (no WhatsApp confirmation needed, the manager is the authority). It enqueues a Celery task and answers 202 like the other writes; 404 for an unknown shift or one outside the manager's locations, 409 when the shift is already absent or already has a live case. The orchestrator method reuses the confirmed-absence path (HRIS marked absent, case opened, deadline scheduled, wave sent, manager notified) and audits `ABSENCE_MARKED` with the manager as the actor, with a redelivery guard so a retried task cannot open a second case. - `run-due-jobs` is gone from beat: it drove the in-memory scheduler, which the broker-owned timers made obsolete, and it always ticked zero. The runtime snapshot it used to publish moved to `reconcile-stale-cases` — that snapshot is what the degraded banner reads, so removing the tick without moving it would have silently stopped an open circuit breaker from reaching /api/status. - Sentry is initialized when `SENTRY_DSN` is set (API lifespan and worker bootstrap), best effort and silent when unset; the DSN is never logged. - The OTel trace id is stored on each interpretation when a span is active, so the decision detail can link to its Langfuse trace instead of always showing null. 515 backend tests, ruff and mypy clean. --- backend/app/api/shifts.py | 99 +++++++++++ backend/app/core/config.py | 4 +- backend/app/main.py | 35 +++- backend/app/services/orchestrator.py | 111 +++++++++++++ backend/app/workers/celery_app.py | 11 +- backend/app/workers/tasks.py | 23 +++ backend/app/workers/tracing_bootstrap.py | 27 ++- backend/pyproject.toml | 1 + backend/tests/unit/api/conftest.py | 2 + backend/tests/unit/api/test_shifts_absence.py | 154 ++++++++++++++++++ .../unit/services/test_interpreter_wiring.py | 55 +++++++ .../tests/unit/services/test_orchestrator.py | 110 +++++++++++++ backend/tests/unit/test_celery.py | 37 ++++- backend/tests/unit/test_observability.py | 20 +++ backend/uv.lock | 15 ++ docs/runbook.md | 17 +- 16 files changed, 695 insertions(+), 26 deletions(-) create mode 100644 backend/app/api/shifts.py create mode 100644 backend/tests/unit/api/test_shifts_absence.py diff --git a/backend/app/api/shifts.py b/backend/app/api/shifts.py new file mode 100644 index 0000000..d0a6d0a --- /dev/null +++ b/backend/app/api/shifts.py @@ -0,0 +1,99 @@ +"""Shift endpoints (spec §7.5): the manager marks an absence, which opens a rescue. + +The manager action is authoritative — no WhatsApp confirmation round trip: +the API validates ownership and state, enqueues the domain change and answers +202 (same shape as the close/approval writes). The worker applies it through +`RescueOrchestrator.mark_absence`, which reuses the confirmed-absence path. +No demo gate here: this is production behaviour (unlike `/dev/*`). +""" + +from fastapi import APIRouter, Depends, HTTPException +from pydantic import BaseModel, Field +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession +from structlog import get_logger + +from app.api.dependencies import ManagerPrincipal, current_manager, get_db +from app.db.models import Manager, RescueCase, Shift +from app.observability.redaction import redact_if_health +from app.workers.tasks import mark_shift_absence_task + +router = APIRouter(prefix="/api/shifts", tags=["shifts"]) + +logger = get_logger(__name__) + +# Case statuses where a rescue is still running on the shift (must match the +# orchestrator's redelivery guard in `RescueOrchestrator.mark_absence`). +_LIVE_CASE_STATUSES = ("OPEN", "OFFERING", "AWAITING_APPROVAL", "ESCALATED") + + +class MarkAbsenceIn(BaseModel): + """Optional context for the audit trail. + + The reason rides only inside the `ABSENCE_MARKED` audit payload, + health-redacted before storage (spec §10) and capped; nothing else is + persisted, and it never becomes a message to anyone. + """ + + reason: str | None = Field(default=None, max_length=200) + + +async def _shift_or_404( + session: AsyncSession, shift_id: str, manager_id: str +) -> Shift: + """The shift when it exists at one of the manager's locations, else 404. + + A shift from another location is indistinguishable from a missing one: + no existence leak across locations (spec §7.5). + """ + shift = ( + await session.execute(select(Shift).where(Shift.id == shift_id)) + ).scalar_one_or_none() + if shift is not None: + manager = ( + await session.execute(select(Manager).where(Manager.id == manager_id)) + ).scalar_one_or_none() + if manager is not None and shift.location_id in (manager.location_ids or []): + return shift + raise HTTPException(status_code=404, detail="Shift not found") + + +@router.post("/{shift_id}/absence", status_code=202) +async def mark_shift_absence( + shift_id: str, + body: MarkAbsenceIn | None = None, + principal: ManagerPrincipal = Depends(current_manager), + session: AsyncSession = Depends(get_db), +) -> dict[str, str]: + """Enqueue the manager-marked absence; the worker owns the domain change. + + 202 after the enqueue, 404 when the shift does not exist or is not at the + manager's location, 409 when the shift is already absent or already has a + live rescue — the detail says which, so the dashboard can show it. + """ + shift = await _shift_or_404(session, shift_id, principal.manager_id) + if shift.status == "absent": + raise HTTPException(status_code=409, detail="Shift is already marked absent") + live = ( + await session.execute( + select(RescueCase.id).where( + RescueCase.shift_id == shift_id, + RescueCase.status.in_(_LIVE_CASE_STATUSES), + ) + ) + ).first() + if live is not None: + raise HTTPException( + status_code=409, detail="A rescue is already running for this shift" + ) + reason = redact_if_health(body.reason) if body is not None and body.reason else None + try: + mark_shift_absence_task.delay(shift_id, principal.manager_id, reason) + except Exception as error: + logger.error( + "shift_absence_enqueue_failed", shift_id=shift_id, error=str(error)[:200] + ) + raise HTTPException( + status_code=500, detail="Failed to enqueue shift absence" + ) from error + return {"status": "queued", "id": shift_id} diff --git a/backend/app/core/config.py b/backend/app/core/config.py index eaeb7f6..b5cebe0 100644 --- a/backend/app/core/config.py +++ b/backend/app/core/config.py @@ -66,7 +66,9 @@ class Settings(BaseSettings): langfuse_public_key: str = "" langfuse_secret_key: str = "" langfuse_host: str = "https://cloud.langfuse.com" - sentry_dsn: str = "" # declared for .env parity; Sentry init not wired yet + # Sentry (spec §9.3): initialized in the API lifespan and the worker + # bootstrap when set; unset stays a silent no-op. + sentry_dsn: str = "" @property def cors_origin_list(self) -> list[str]: diff --git a/backend/app/main.py b/backend/app/main.py index 62613d8..7528e2c 100644 --- a/backend/app/main.py +++ b/backend/app/main.py @@ -18,23 +18,49 @@ from app.api.locations import router as locations_router from app.api.metrics import router as metrics_router from app.api.rescues import router as rescues_router +from app.api.shifts import router as shifts_router from app.api.status import router as status_router from app.api.webhooks_twilio import router as twilio_router from app.api.ws import router as ws_router -from app.core.config import get_settings +from app.core.config import Settings, get_settings from app.core.logging import configure_logging from app.observability.tracing import configure_tracing, shutdown_tracing +def _init_sentry(settings: Settings) -> None: + """Best-effort Sentry init (spec §9.3 spirit): never stops the boot. + + Unset DSN is a silent no-op (the setting exists for `.env` parity). A + failure is logged and swallowed; the DSN itself is never logged. Twin + copy lives in `app/workers/tracing_bootstrap.py` (worker process). + """ + if not settings.sentry_dsn: + return + try: + import sentry_sdk + + sentry_sdk.init( + dsn=settings.sentry_dsn, + environment=settings.app_env, + traces_sample_rate=0.1, + ) + structlog.get_logger(__name__).info("sentry_initialized") + except Exception as error: + structlog.get_logger(__name__).warning( + "sentry_init_failed", error=str(error)[:200] + ) + + @contextlib.asynccontextmanager async def lifespan(_: FastAPI) -> AsyncIterator[None]: settings = get_settings() configure_tracing(settings) - # The API owns no scheduling: Celery beat ticks `run_due_jobs` every 5 s - # and the worker runs the daily retention purge (spec §7.2/§7.3). + _init_sentry(settings) + # The API owns no scheduling: Celery beat runs the reconcile sweep and the + # daily retention purge; timers themselves are broker-owned (spec §7.2/§7.3). structlog.get_logger(__name__).info( "scheduling_owned_by_worker", - detail="Celery beat drives scheduled rescue work (run-due-jobs every 5s)", + detail="Celery beat drives the reconcile sweep and the retention purge", ) try: yield @@ -62,6 +88,7 @@ def create_app() -> FastAPI: app.include_router(status_router) app.include_router(auth_router) app.include_router(locations_router) + app.include_router(shifts_router) app.include_router(rescues_router) app.include_router(approvals_router) app.include_router(conversations_router) diff --git a/backend/app/services/orchestrator.py b/backend/app/services/orchestrator.py index 4689108..2dc8cbb 100644 --- a/backend/app/services/orchestrator.py +++ b/backend/app/services/orchestrator.py @@ -13,6 +13,7 @@ from zoneinfo import ZoneInfo import structlog +from opentelemetry.trace import get_current_span from sqlalchemy import select from sqlalchemy.exc import IntegrityError from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker @@ -32,6 +33,7 @@ Message, Offer, RescueCase, + Shift, ) from app.db.models import Interpretation as InterpretationRow from app.domain.eligibility import evaluate_eligibility @@ -553,6 +555,13 @@ async def _persist_interpretation(self, message_id: str, interpreted: Any) -> No "contains_health_details": interpreted.contains_health_details, "question_text": interpreted.question_text, } + # Trace link (spec §7.6): the API derives the Langfuse URL from + # `extracted["trace_id"]`. Stored only when a valid span is + # active — tracing off leaves the key out entirely, so the link + # stays null exactly as before. + trace_id = _current_trace_id() + if trace_id is not None: + extracted["trace_id"] = trace_id async with self._sessions() as session: session.add( InterpretationRow( @@ -1758,6 +1767,85 @@ async def _offer(self, session: AsyncSession, offer_id: str | None) -> Offer | N await session.execute(select(Offer).where(Offer.id == offer_id)) ).scalar_one_or_none() + async def mark_absence( + self, shift_id: str, manager_id: str, reason: str | None = None + ) -> str | None: + """Manager marks a shift absent from the dashboard (spec §7.5). + + The manager action is authoritative, so no employee confirmation is + needed: the case is created here (origin `manager_dashboard`, audit + `ABSENCE_MARKED` with `actor="manager:"`) and everything else — + the HRIS absence, the deadline, wave 1 and the manager notice — + reuses exactly the confirmed-absence path (`_handle_confirmation`). + An optional `reason` rides only in the audit payload, health-redacted + by the API and never stored anywhere else (spec §10). + + Returns the case id, or None when the shift is gone, has no assignee, + is already absent or already has a live case. The API answers + 404/409 from the same checks at request time; this second guard makes + a redelivered task (acks_late) a harmless no-op instead of a second + rescue. + """ + now = self._clock.now() + async with self._sessions() as session: + shift = ( + await session.execute(select(Shift).where(Shift.id == shift_id)) + ).scalar_one_or_none() + if shift is None or shift.employee_id is None or shift.status == "absent": + return None + live = ( + await session.execute( + select(RescueCase.id).where( + RescueCase.shift_id == shift_id, + RescueCase.status.in_(_LIVE_CASE_STATUSES), + ) + ) + ).first() + if live is not None: + return None + + employee_id = shift.employee_id + case_id = f"case_{uuid4().hex}" + session.add( + RescueCase( + id=case_id, + location_id=shift.location_id, + shift_id=shift_id, + absent_employee_id=employee_id, + origin="manager_dashboard", + status=State.OPEN.value, + opened_at=now, + deadline_at=self._deadline_for(now, shift), + ) + ) + payload: dict[str, Any] = {"shift_id": shift_id} + if reason: + payload["reason"] = redact_if_health(reason)[:200] + session.add( + AuditEvent( + id=f"audit_{uuid4().hex}", + rescue_id=case_id, + type="ABSENCE_MARKED", + payload=payload, + actor=f"manager:{manager_id}", + ) + ) + await session.commit() + + await self._emit( + EventName.RESCUE_OPENED, + location_id=shift.location_id, + rescue_id=case_id, + shift_id=shift_id, + origin=f"manager:{manager_id}", + ) + # The confirmed-absence path owns the transition, the deadline, wave 1 + # and the manager notice; the conversation id is only consulted by its + # no-case branch, which a case created just above never takes. + conversation_id = await self._conversation_for(employee_id) + await self._handle_confirmation(conversation_id or "", employee_id) + return case_id + async def close_rescue(self, rescue_id: str, decided_by: str) -> None: """Manual manager close (spec §7.5), runs in the worker. @@ -2610,6 +2698,29 @@ def _fmt(self, moment: datetime, tz_name: str | None = None) -> str: return moment.astimezone(tz).strftime("%H:%M") +# Case statuses where a rescue is still running on a shift: the manager +# absence endpoint answers 409 for these (spec §7.5) and the orchestrator's +# redelivery guard refuses to open a second case. +_LIVE_CASE_STATUSES = ( + State.OPEN.value, + State.OFFERING.value, + State.AWAITING_APPROVAL.value, + State.ESCALATED.value, +) + + +def _current_trace_id() -> str | None: + """Current OTel trace id (hex) when a valid span is active, else None. + + Tracing off (or no active span) returns None: nothing is stored under + `extracted["trace_id"]` and the API's Langfuse link stays null. + """ + context = get_current_span().get_span_context() + if context is not None and context.is_valid: + return format(context.trace_id, "032x") + return None + + def _utc(moment: datetime) -> datetime: return moment.astimezone(UTC) diff --git a/backend/app/workers/celery_app.py b/backend/app/workers/celery_app.py index a7f2c03..bca7a51 100644 --- a/backend/app/workers/celery_app.py +++ b/backend/app/workers/celery_app.py @@ -1,8 +1,9 @@ """Celery application (spec §7.2: queues, waves, timeouts, retries). -Celery beat owns time: it ticks the scheduler (`run-due-jobs`, every 5 s, -for the memory backend and the status snapshot), runs the reconcile sweep -(`reconcile-stale-cases`, every 60 s) and the daily retention purge. +Celery beat owns the periodic sweeps: the reconcile pass every 60 s and the +daily retention purge. Timers themselves are broker-owned: `CeleryScheduler` +publishes each one as a deferred task, so there is no beat tick scanning a +queue (`run-due-jobs` was removed when timers left the worker process). Orchestration runs in the worker process via `app.runtime` — the API process only enqueues tasks (spec §7.5). """ @@ -31,10 +32,6 @@ ) celery_app.conf.beat_schedule = { - "run-due-jobs": { - "task": "app.workers.tasks.run_due_jobs", - "schedule": 5.0, - }, "reconcile-stale-cases": { # Self-healing sweep (spec §9.1): re-enqueue deadline timers of # overdue cases that still expect action. Handlers re-check state, diff --git a/backend/app/workers/tasks.py b/backend/app/workers/tasks.py index 37abdd5..61f8ccf 100644 --- a/backend/app/workers/tasks.py +++ b/backend/app/workers/tasks.py @@ -87,6 +87,24 @@ def close_rescue_task(rescue_id: str, decided_by: str) -> bool: return True +@celery_app.task(name="app.workers.tasks.mark_shift_absence") +def mark_shift_absence_task(shift_id: str, manager_id: str, reason: str | None = None) -> bool: + """Manager-marked absence in the worker (spec §7.5). + + The orchestrator owns the domain change: it re-checks shift state (a + redelivered task is a no-op), opens the case and reuses the + confirmed-absence path for the HRIS absence, deadline, wave 1 and the + manager notice. + """ + from app.runtime import get_worker_runtime + + case_id = run_async( + get_worker_runtime().orchestrator.mark_absence(shift_id, manager_id, reason) + ) + logger.info("worker_shift_absence_marked", shift_id=shift_id, case_id=case_id) + return case_id is not None + + @celery_app.task( name="app.workers.tasks.apply_scheduled_job", bind=True, @@ -166,6 +184,11 @@ def reconcile_stale_cases() -> int: ) if recovered: logger.warning("reconcile_stale_cases_recovered", count=recovered) + # This sweep is the heartbeat that keeps the degraded-status snapshot fresh + # for the API probe (spec §9.3). It used to ride on the memory-scheduler + # tick, which the broker-owned timers made obsolete: without moving it here + # an open circuit breaker would silently stop reaching the dashboard banner. + publish_runtime_snapshot(runtime) return recovered diff --git a/backend/app/workers/tracing_bootstrap.py b/backend/app/workers/tracing_bootstrap.py index 7d1eeed..68cc4c2 100644 --- a/backend/app/workers/tracing_bootstrap.py +++ b/backend/app/workers/tracing_bootstrap.py @@ -20,17 +20,40 @@ from celery.signals import worker_process_init, worker_process_shutdown from structlog import get_logger -from app.core.config import get_settings +from app.core.config import Settings, get_settings from app.observability.tracing import configure_tracing, shutdown_tracing logger = get_logger(__name__) +def _init_sentry(settings: Settings) -> None: + """Best-effort Sentry init in the worker: never stops a boot. + + Unset DSN is a silent no-op; a failure is logged and swallowed. The DSN + itself is never logged. Twin copy lives in `app/main.py` (API process). + """ + if not settings.sentry_dsn: + return + try: + import sentry_sdk + + sentry_sdk.init( + dsn=settings.sentry_dsn, + environment=settings.app_env, + traces_sample_rate=0.1, + ) + logger.info("sentry_initialized") + except Exception as error: + logger.warning("sentry_init_failed", error=str(error)[:200]) + + @worker_process_init.connect def _init_worker_tracing(**_kwargs: Any) -> None: """Install the TracerProvider in this preforked child (or solo worker).""" try: - installed = configure_tracing(get_settings()) + settings = get_settings() + installed = configure_tracing(settings) + _init_sentry(settings) logger.info("worker_tracing_bootstrap", installed=installed) except Exception as error: logger.warning("worker_tracing_bootstrap_failed", error=str(error)[:200]) diff --git a/backend/pyproject.toml b/backend/pyproject.toml index f995680..d6fd959 100644 --- a/backend/pyproject.toml +++ b/backend/pyproject.toml @@ -23,6 +23,7 @@ dependencies = [ "httpx>=0.28.1", "argon2-cffi>=23.1", # password hashing for manager login (spec §7.5) "pyjwt>=2.9", # JWT issue/verify for the dashboard API (spec §7.5) + "sentry-sdk>=2.0", # error/performance monitoring, init only when SENTRY_DSN is set ] [dependency-groups] diff --git a/backend/tests/unit/api/conftest.py b/backend/tests/unit/api/conftest.py index 0032d66..adb54d3 100644 --- a/backend/tests/unit/api/conftest.py +++ b/backend/tests/unit/api/conftest.py @@ -14,6 +14,7 @@ from app.workers.tasks import ( apply_scheduled_job, close_rescue_task, + mark_shift_absence_task, process_inbound_message, reconcile_stale_cases, ) @@ -23,6 +24,7 @@ reconcile_stale_cases, apply_scheduled_job, close_rescue_task, + mark_shift_absence_task, ) _calls: list[tuple[str, tuple]] = [] diff --git a/backend/tests/unit/api/test_shifts_absence.py b/backend/tests/unit/api/test_shifts_absence.py new file mode 100644 index 0000000..0744ee2 --- /dev/null +++ b/backend/tests/unit/api/test_shifts_absence.py @@ -0,0 +1,154 @@ +"""POST /api/shifts/{id}/absence (spec §7.5): the manager marks an absence. + +The manager action is authoritative: 202 after the Celery enqueue, 404 when +the shift is unknown or outside the manager's locations, 409 when the shift +is already absent or already has a live rescue — the detail says which. No +demo gate: production behaviour. No broker: the enqueue is stubbed (the API +suite autouse fixture plus the explicit recorder below). +""" + +from datetime import timedelta + +import pytest +from sqlalchemy import select + +from app.db.models import RescueCase, Shift +from app.workers.tasks import mark_shift_absence_task +from tests.conftest import ABSENT_ID, LOCATION_ID, NOW, SHIFT_ID, auth_headers + +SCHEDULED_SHIFT_ID = "shift_scheduled" + + +class StubTask: + """Records `.delay` calls (replaces the suite-wide no-op recorder).""" + + def __init__(self) -> None: + self.calls: list[tuple] = [] + + def delay(self, *args) -> None: + self.calls.append(args) + + +@pytest.fixture() +def stub_absence(monkeypatch): + stub = StubTask() + monkeypatch.setattr(mark_shift_absence_task, "delay", stub.delay) + return stub + + +async def _add_scheduled_shift( + sessions, + *, + shift_id: str = SCHEDULED_SHIFT_ID, + location_id: str = LOCATION_ID, +) -> str: + async with sessions() as session: + session.add( + Shift( + id=shift_id, + location_id=location_id, + role="bar", + starts_at=NOW + timedelta(hours=4), + ends_at=NOW + timedelta(hours=12), + employee_id=ABSENT_ID, + status="scheduled", + ) + ) + await session.commit() + return shift_id + + +async def _add_live_case(sessions, shift_id: str) -> None: + async with sessions() as session: + session.add( + RescueCase( + id="res_live", + location_id=LOCATION_ID, + shift_id=shift_id, + absent_employee_id=ABSENT_ID, + origin="employee_message", + status="OFFERING", + opened_at=NOW, + deadline_at=NOW + timedelta(minutes=30), + ) + ) + await session.commit() + + +async def test_marks_absence_returns_202_and_enqueues(client, world, stub_absence) -> None: + shift_id = await _add_scheduled_shift(world.sessions) + response = await client.post(f"/api/shifts/{shift_id}/absence", headers=auth_headers()) + assert response.status_code == 202 + assert response.json() == {"status": "queued", "id": shift_id} + assert stub_absence.calls == [(shift_id, "mgr_1", None)] + + +async def test_optional_reason_is_passed_through_to_the_task(client, world, stub_absence) -> None: + shift_id = await _add_scheduled_shift(world.sessions) + response = await client.post( + f"/api/shifts/{shift_id}/absence", + json={"reason": "llamo y no contesta"}, + headers=auth_headers(), + ) + assert response.status_code == 202 + assert stub_absence.calls == [(shift_id, "mgr_1", "llamo y no contesta")] + + +async def test_unknown_shift_404(client, world, stub_absence) -> None: + response = await client.post("/api/shifts/shift_missing/absence", headers=auth_headers()) + assert response.status_code == 404 + assert stub_absence.calls == [] + + +async def test_shift_outside_the_manager_locations_404(client, world, stub_absence) -> None: + shift_id = await _add_scheduled_shift( + world.sessions, shift_id="shift_other", location_id="loc_other" + ) + response = await client.post(f"/api/shifts/{shift_id}/absence", headers=auth_headers()) + assert response.status_code == 404 + assert stub_absence.calls == [] + + +async def test_already_absent_409_names_the_conflict(client, world, stub_absence) -> None: + # The seeded world's shift is already marked absent (spec §7.5). + response = await client.post(f"/api/shifts/{SHIFT_ID}/absence", headers=auth_headers()) + assert response.status_code == 409 + assert "already marked absent" in response.json()["detail"] + assert stub_absence.calls == [] + + +async def test_live_rescue_409_names_the_conflict(client, world, stub_absence) -> None: + shift_id = await _add_scheduled_shift(world.sessions) + await _add_live_case(world.sessions, shift_id) + response = await client.post(f"/api/shifts/{shift_id}/absence", headers=auth_headers()) + assert response.status_code == 409 + assert "already running" in response.json()["detail"] + assert stub_absence.calls == [] + + +async def test_requires_manager_jwt(client, world, stub_absence) -> None: + response = await client.post(f"/api/shifts/{SCHEDULED_SHIFT_ID}/absence") + assert response.status_code == 401 + assert stub_absence.calls == [] + + +async def test_enqueue_failure_is_a_loud_500(client, world, stub_absence, monkeypatch) -> None: + shift_id = await _add_scheduled_shift(world.sessions) + + def boom(*args) -> None: + raise RuntimeError("broker down") + + monkeypatch.setattr(mark_shift_absence_task, "delay", boom) + response = await client.post(f"/api/shifts/{shift_id}/absence", headers=auth_headers()) + assert response.status_code == 500 + + +async def test_absence_leaves_world_data_untouched(client, world, stub_absence) -> None: + """The API never applies the domain change inline: the worker owns it.""" + shift_id = await _add_scheduled_shift(world.sessions) + await client.post(f"/api/shifts/{shift_id}/absence", headers=auth_headers()) + async with world.sessions() as session: + shift = ( + await session.execute(select(Shift).where(Shift.id == shift_id)) + ).scalar_one() + assert shift.status == "scheduled" # still untouched until the worker runs diff --git a/backend/tests/unit/services/test_interpreter_wiring.py b/backend/tests/unit/services/test_interpreter_wiring.py index ee0dc4c..50739be 100644 --- a/backend/tests/unit/services/test_interpreter_wiring.py +++ b/backend/tests/unit/services/test_interpreter_wiring.py @@ -630,3 +630,58 @@ def keys(recorded: dict) -> set[str]: "pending_confirmation", } assert llm.contexts[-1]["pending_confirmation"] == "shift_1" + + +# --- trace link (spec §7.6): the Langfuse URL input ---------------------------- + + +async def test_llm_interpretation_stores_the_trace_id_when_a_span_is_active() -> None: + """A valid OTel span at persist time lands in `extracted["trace_id"]` + (hex), which the API already turns into a Langfuse link.""" + import opentelemetry.trace as otel_trace + from opentelemetry.trace import NonRecordingSpan, SpanContext, TraceFlags + + world, _ = await build_world(floor_count=4) + world.orchestrator.interpreter = interpreter_with( + {"intent": "ABSENCE_REPORT", "confidence": 0.95} + ) + trace_id = 0x1234567890ABCDEF1234567890ABCDEF + span = NonRecordingSpan( + SpanContext( + trace_id=trace_id, + span_id=0x2222, + is_remote=False, + trace_flags=TraceFlags(TraceFlags.SAMPLED), + ) + ) + with otel_trace.use_span(span): + await world.orchestrator.handle_inbound( + conversation_id=CONVERSATION, + employee_id="emp_01_floor", + provider_message_id="wires_trace_1", + text="me encuentro fatal, hoy no puedo ir", + ) + + async with world.session_factory() as session: + row = (await session.execute(select(InterpretationRow))).scalars().one() + assert row.extracted["trace_id"] == format(trace_id, "032x") + + +async def test_llm_interpretation_stores_no_trace_id_without_an_active_span() -> None: + """Tracing off (no valid span): the key is not stored at all, so the API + trace link stays null exactly as before.""" + world, _ = await build_world(floor_count=4) + world.orchestrator.interpreter = interpreter_with( + {"intent": "ABSENCE_REPORT", "confidence": 0.95} + ) + + await world.orchestrator.handle_inbound( + conversation_id=CONVERSATION, + employee_id="emp_01_floor", + provider_message_id="wires_trace_2", + text="me encuentro fatal, hoy no puedo ir", + ) + + async with world.session_factory() as session: + row = (await session.execute(select(InterpretationRow))).scalars().one() + assert "trace_id" not in row.extracted diff --git a/backend/tests/unit/services/test_orchestrator.py b/backend/tests/unit/services/test_orchestrator.py index c9c94c7..1be9463 100644 --- a/backend/tests/unit/services/test_orchestrator.py +++ b/backend/tests/unit/services/test_orchestrator.py @@ -829,3 +829,113 @@ async def test_the_race_loser_notice_is_stored_in_their_conversation(world, db) assert stored, "the loser's notice must be in their conversation" assert "ya se ha cubierto" in stored[0].body_redacted + + +# --- manager-marked absence (spec §7.5) --------------------------------------- + + +async def test_manager_marked_absence_reuses_the_confirmation_path(world, db, now) -> None: + """The dashboard action is authoritative: the case opens and wave 1 goes + out with no employee round trip, exactly like a confirmed absence.""" + case_id = await world.orchestrator.mark_absence("shift_1", "mgr_1") + + assert case_id is not None and case_id.startswith("case_") + async with db() as session: + case = ( + await session.execute(select(RescueCase).where(RescueCase.id == case_id)) + ).scalar_one() + shift = (await session.execute(select(Shift).where(Shift.id == "shift_1"))).scalar_one() + events = ( + await session.execute(select(AuditEvent).where(AuditEvent.rescue_id == case_id)) + ).scalars().all() + offers = ( + await session.execute(select(Offer).where(Offer.rescue_id == case_id)) + ).scalars().all() + + assert case.origin == "manager_dashboard" + assert case.status == "OFFERING" + assert case.absent_employee_id == "emp_01_floor" + assert shift.status == "absent" # marked in the HRIS + marked = [e for e in events if e.type == "ABSENCE_MARKED"] + assert [e.actor for e in marked] == ["manager:mgr_1"] + assert marked[0].payload == {"shift_id": "shift_1"} + assert any(e.type == "RESCUE_OPENED" for e in events) + assert {offer.wave_number for offer in offers} == {1} + # Deadline and wave timers are scheduled (broker-owned in production). + assert world.scheduler.pending_count() == 2 + # The manager is notified the rescue opened. + assert world.channel.with_template("manager_rescue_opened") + + +async def test_manager_marked_absence_records_the_reason_in_the_audit_payload( + world, db +) -> None: + """An optional reason lives only in the ABSENCE_MARKED payload (§10).""" + case_id = await world.orchestrator.mark_absence( + "shift_1", "mgr_1", reason="llamo y no contesta" + ) + + async with db() as session: + event = ( + await session.execute( + select(AuditEvent).where( + AuditEvent.rescue_id == case_id, + AuditEvent.type == "ABSENCE_MARKED", + ) + ) + ).scalar_one() + assert event.payload["reason"] == "llamo y no contesta" + + +async def test_manager_marked_absence_refuses_a_second_rescue(world, db) -> None: + """A redelivery (or a double click) must not open a second case.""" + first = await world.orchestrator.mark_absence("shift_1", "mgr_1") + second = await world.orchestrator.mark_absence("shift_1", "mgr_1") + + assert first is not None + assert second is None # shift already absent, case already running + async with db() as session: + count = ( + await session.execute(select(func.count()).select_from(RescueCase)) + ).scalar_one() + assert count == 1 + + +async def test_manager_marked_absence_for_a_shift_with_a_live_case_returns_none( + world, db, now +) -> None: + """An OFFERING rescue on a scheduled shift is a conflict, not a new case.""" + async with db() as session: + session.add( + Shift( + id="shift_live", + location_id=DEMO_LOCATION_ID, + role="bar", + starts_at=now + timedelta(hours=4), + ends_at=now + timedelta(hours=12), + employee_id="emp_02_floor", + status="scheduled", + ) + ) + session.add( + RescueCase( + id="case_live", + location_id=DEMO_LOCATION_ID, + shift_id="shift_live", + absent_employee_id="emp_02_floor", + origin="employee_message", + status="OFFERING", + opened_at=now, + deadline_at=now + timedelta(minutes=30), + ) + ) + await session.commit() + + result = await world.orchestrator.mark_absence("shift_live", "mgr_1") + + assert result is None + async with db() as session: + count = ( + await session.execute(select(func.count()).select_from(RescueCase)) + ).scalar_one() + assert count == 1 diff --git a/backend/tests/unit/test_celery.py b/backend/tests/unit/test_celery.py index d76907a..5f9ea5a 100644 --- a/backend/tests/unit/test_celery.py +++ b/backend/tests/unit/test_celery.py @@ -11,6 +11,7 @@ import app.runtime as runtime_module import app.workers.tasks as tasks import app.workers.tracing_bootstrap as bootstrap +from app.core.clock import SystemClock from app.core.config import Settings, get_settings from app.workers.celery_app import celery_app from app.workers.scheduler import SimScheduler @@ -36,10 +37,13 @@ def test_retention_purge_task_is_registered() -> None: # --- beat schedule (spec §7.3) ------------------------------------------------ -def test_beat_schedule_ticks_the_scheduler_every_five_seconds() -> None: - entry = celery_app.conf.beat_schedule["run-due-jobs"] - assert entry["task"] == "app.workers.tasks.run_due_jobs" - assert entry["schedule"] == 5.0 +def test_beat_schedule_does_not_tick_the_scheduler() -> None: + """Timers are broker-owned (CeleryScheduler): no beat tick scans a queue. + + `run-due-jobs` drove the in-memory SimScheduler and always found zero due + jobs in the preforked worker; the entry is gone and must not come back. + """ + assert "run-due-jobs" not in celery_app.conf.beat_schedule def test_beat_schedule_runs_the_retention_purge_daily() -> None: @@ -64,6 +68,9 @@ def __init__(self, error: Exception | None = None) -> None: self._error = error self.scheduler = SimScheduler() self.interpreter = None + # The reconcile sweep takes these as arguments before the stubbed body runs. + self.session_factory = None + self.clock = SystemClock() self.calls: list[tuple[str, str, str]] = [] async def handle_inbound(self, from_phone: str, message_sid: str, body: str) -> bool: @@ -288,3 +295,25 @@ def boom(_settings: Settings) -> bool: failures = [e for e in logs if e["event"] == "worker_tracing_bootstrap_failed"] assert len(failures) == 1 assert "otel exploded" in failures[0]["error"] + +async def _reconcile_none(*args, **kwargs): + return 0 + + +def test_reconcile_publishes_the_snapshot(monkeypatch, fake_runtime: FakeRuntime) -> None: + """The sweep is the heartbeat the degraded banner depends on. + + It used to ride on the memory-scheduler tick; with broker-owned timers that + tick is gone (and no longer scheduled), so if the sweep stopped publishing the + snapshot an open circuit breaker would silently stop reaching /api/status. + """ + published: list[dict] = [] + monkeypatch.setattr( + tasks, "publish_runtime_snapshot", lambda runtime: published.append(runtime) + ) + monkeypatch.setattr(tasks, "_reconcile_stale_cases", _reconcile_none) + monkeypatch.setattr("app.runtime.get_worker_runtime", lambda: fake_runtime) + + tasks.reconcile_stale_cases() + + assert published == [fake_runtime] diff --git a/backend/tests/unit/test_observability.py b/backend/tests/unit/test_observability.py index 830d96f..d29eef1 100644 --- a/backend/tests/unit/test_observability.py +++ b/backend/tests/unit/test_observability.py @@ -202,3 +202,23 @@ def boom() -> None: failures = [e for e in logs if e["event"] == "worker_tracing_shutdown_failed"] assert len(failures) == 1 assert "flush exploded" in failures[0]["error"] + + +# --- Sentry init (best effort, spec §9.3 spirit) ------------------------------ + + +def test_sentry_disabled_is_a_silent_no_op(monkeypatch) -> None: + """Unset SENTRY_DSN (the default): neither boot path touches sentry_sdk.""" + import sentry_sdk + + import app.main as main + + calls: list[dict] = [] + monkeypatch.setattr(sentry_sdk, "init", lambda **kwargs: calls.append(kwargs)) + settings = make_settings() + assert settings.sentry_dsn == "" + + main._init_sentry(settings) + bootstrap._init_sentry(settings) + + assert calls == [] diff --git a/backend/uv.lock b/backend/uv.lock index d281f82..88c4d51 100644 --- a/backend/uv.lock +++ b/backend/uv.lock @@ -1316,6 +1316,19 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/bc/e7/5c595c75e9f41a44f30e526eda465ea0b4eec93470e074e4a111b253f13a/s3transfer-0.19.2-py3-none-any.whl", hash = "sha256:d8168eccca828cbb2cd573675333f3bddd254313a9c42494b84c76b539e8ba25", size = 90216, upload-time = "2026-07-22T19:30:43.251Z" }, ] +[[package]] +name = "sentry-sdk" +version = "2.71.0" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "certifi" }, + { name = "urllib3" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/67/93/9ef70bfb7346c778caf09b7cd10c17c33f601f2d17019fa1aedcda738b5f/sentry_sdk-2.71.0.tar.gz", hash = "sha256:7beb27a22f06396f3a05510c3da6ea10e1c79e02e859993c27ce4ff074e296d6", size = 1052358, upload-time = "2026-09-28T14:04:00.645Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/a4/f6/b3e68b2e08e9a52e6c8a279ca2fc0d3a535f3886222375f54ef9efbd694d/sentry_sdk-2.71.0-py3-none-any.whl", hash = "sha256:9ed36b4c206048aa67b5ddba0a6e28d59fc55cf10d60f5f7aea2c3917184bf4a", size = 534254, upload-time = "2026-09-28T14:03:58.668Z" }, +] + [[package]] name = "shift-rescue-backend" version = "0.1.0" @@ -1335,6 +1348,7 @@ dependencies = [ { name = "pyjwt" }, { name = "python-multipart" }, { name = "redis" }, + { name = "sentry-sdk" }, { name = "sqlalchemy", extra = ["asyncio"] }, { name = "strands-agents" }, { name = "structlog" }, @@ -1370,6 +1384,7 @@ requires-dist = [ { name = "pyjwt", specifier = ">=2.9" }, { name = "python-multipart", specifier = ">=0.0.32" }, { name = "redis", specifier = ">=5.0" }, + { name = "sentry-sdk", specifier = ">=2.0" }, { name = "sqlalchemy", extras = ["asyncio"], specifier = ">=2.0" }, { name = "strands-agents", specifier = "==1.56.0" }, { name = "structlog", specifier = ">=24.1" }, diff --git a/docs/runbook.md b/docs/runbook.md index 5bcd705..ac77b2a 100644 --- a/docs/runbook.md +++ b/docs/runbook.md @@ -186,9 +186,9 @@ curl -fsS -o /dev/null -w '%{http_code}\n' https:/// # 200 (SPA) The **worker and beat containers are required for the demo**: the API only enqueues the inbound task (spec §7.5), the worker runs the orchestration and -the LLM, and beat ticks the scheduler (`run-due-jobs`, every 5 s) plus the -daily retention purge. `docker compose ps` must show `api`, `worker`, `beat`, -`redis` and `postgres` up. +the LLM, and beat runs the reconcile sweep (`reconcile-stale-cases`, every +60 s) plus the daily retention purge. `docker compose ps` must show `api`, +`worker`, `beat`, `redis` and `postgres` up. If a message gets no reply, check in this order: @@ -340,10 +340,11 @@ Safety nets, in order: recovers a timer lost to a crash between the database commit and the enqueue. -Beat still ticks `run-due-jobs` (5 s), which drives the `memory` scheduler -backend (`SCHEDULER_BACKEND=memory`, for a single-process local run) and -refreshes the API status snapshot. `SCHEDULER_BACKEND=celery` is the -production default. +Beat no longer ticks a scheduler queue: timers are broker-owned +(`CeleryScheduler` publishes each one as a deferred task), so the +`run-due-jobs` entry was removed. The `memory` scheduler backend remains +only for the eval harness and single-process local runs +(`SCHEDULER_BACKEND=memory`). ### Record an eval run and see it in the dashboard @@ -394,7 +395,7 @@ command above; the mock fixture stays available offline behind | `20003 Primary compliance profile` | Twilio Trust Hub profile `draft` | complete and submit the profile in Trust Hub | | Rescue looks stuck (no wave, no escalation) | `logs worker \| grep reconcile_stale_cases_recovered`; `docker compose ps` (worker and beat up?) | the reconcile sweep re-enqueues overdue deadline timers every 60 s, so a stuck case escalates within a minute of beat running; if it does not, check worker/beat are up and Redis is reachable — do not re-run the flow first | | Message gets no reply | see §4 checklist | API enqueues (`twilio_inbound_received`), worker processes (`worker_inbound_processed`); a `twilio_inbound_enqueue_failed` 500 means Redis/broker down — Twilio retries, recover Redis | -| Timeouts/purge never fire | `docker compose ps` shows `beat` down | start beat: `docker compose up -d beat` — beat owns `run-due-jobs` (every 5 s) and the daily purge | +| Timeouts/purge never fire | `docker compose ps` shows `beat` down | start beat: `docker compose up -d beat` — beat owns `reconcile-stale-cases` (every 60 s) and the daily purge | | DB full / slow | `df -h`, `docker system df` | prune images (`docker image prune -f`), grow the EBS volume | ## 7. Rollback From eb45545ab6dedd2c32b4c3a99f7e1ae86c8defb5 Mon Sep 17 00:00:00 2001 From: albert Date: Mon, 28 Sep 2026 19:14:19 +0200 Subject: [PATCH 5/6] feat(web): wire the "Report absence" button to the manager action MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The floating action said "Report absence" and did nothing — a dead control in the dashboard, and the reason the new `POST /api/shifts/{id}/absence` was only reachable from a terminal. The button now opens a small dialog listing today's still-scheduled shifts (with role, window and who holds them) and marks the chosen one absent. The copy is honest about the 202 — "the rescue is opening; the worker applies it in a moment" — and a refusal reaches the screen with the server's reason (409 when the shift is already absent or already has a live case, 404 when it is not available for this location). The board refetches on success and the live channel delivers whatever else changes. Verified in the browser: the button opens the dialog with 10 shifts, marking one was accepted, and Today showed that shift as absent with its rescue searching. 199 frontend tests, build, lint and tsc clean. --- frontend/src/App.tsx | 12 +- .../components/ReportAbsenceDialog.test.tsx | 72 +++++++++++ .../src/components/ReportAbsenceDialog.tsx | 117 ++++++++++++++++++ frontend/src/services/api.ts | 8 ++ 4 files changed, 208 insertions(+), 1 deletion(-) create mode 100644 frontend/src/components/ReportAbsenceDialog.test.tsx create mode 100644 frontend/src/components/ReportAbsenceDialog.tsx diff --git a/frontend/src/App.tsx b/frontend/src/App.tsx index be24485..4a27141 100644 --- a/frontend/src/App.tsx +++ b/frontend/src/App.tsx @@ -11,6 +11,7 @@ import { AgentDecisionsScreen } from './screens/AgentDecisionsScreen' import { EvalsScreen } from './screens/EvalsScreen' import { SimulatorScreen } from './screens/SimulatorScreen' import { RequireAuth } from './components/RequireAuth' +import { ReportAbsenceDialog } from './components/ReportAbsenceDialog' import { getSession } from './services/auth' import { useLiveEvents } from './services/liveEvents' import { ApiError } from './services/apiClient' @@ -59,6 +60,7 @@ function viewFromHash(): AppView | null { function Shell({ now, onLogout }: { now?: Date; onLogout: () => void }) { const [view, setView] = useState(() => viewFromHash() ?? 'today') const [selectedRescueId, setSelectedRescueId] = useState(null) + const [reportingAbsence, setReportingAbsence] = useState(false) const { approvals } = usePendingApprovals() const manager = getSession()?.manager // Live channel (spec §7.5): the screens refetch when the worker reports a @@ -117,7 +119,15 @@ function Shell({ now, onLogout }: { now?: Date; onLogout: () => void }) { )}
- + {/* The manager marks an absence here (spec §7.5): the rescue opens and + the first wave goes out without an employee confirmation. */} + setReportingAbsence(true)} /> + {reportingAbsence ? ( + setReportingAbsence(false)} + /> + ) : null} ) } diff --git a/frontend/src/components/ReportAbsenceDialog.test.tsx b/frontend/src/components/ReportAbsenceDialog.test.tsx new file mode 100644 index 0000000..fd70cf9 --- /dev/null +++ b/frontend/src/components/ReportAbsenceDialog.test.tsx @@ -0,0 +1,72 @@ +/** + * Manager action: mark an absence from the dashboard (spec §7.5). + * + * The dialog lists today's still-scheduled shifts and posts the chosen one. The + * endpoint enqueues the work, so the copy must not claim the rescue is already + * open; a refusal (409) has to reach the screen with the server's reason. + */ + +import { screen, waitFor, within } from '@testing-library/react' +import userEvent from '@testing-library/user-event' +import { afterEach, describe, expect, it, vi } from 'vitest' + +import { ReportAbsenceDialog } from './ReportAbsenceDialog' +import { renderWithProviders } from '../test/renderWithProviders' +import { ApiError } from '../services/apiClient' + +const markShiftAbsence = vi.fn() +vi.mock('../services/api', () => ({ + markShiftAbsence: (shiftId: string) => markShiftAbsence(shiftId), +})) + +const DAY = '2026-10-03' + +afterEach(() => { + markShiftAbsence.mockReset() +}) + +describe('ReportAbsenceDialog', () => { + it('lists the shifts still scheduled and marks the chosen one absent', async () => { + markShiftAbsence.mockResolvedValue({ status: 'accepted', id: 'shift_x' }) + const user = userEvent.setup() + renderWithProviders( {}} />) + + const list = await screen.findByRole('list', { name: "Today's shifts" }) + const markButtons = await within(list).findAllByRole('button', { name: 'Mark absent' }) + expect(markButtons.length).toBeGreaterThan(0) + + await user.click(markButtons[0]) + + await waitFor(() => expect(markShiftAbsence).toHaveBeenCalledTimes(1)) + // The copy is honest about the 202: the worker applies it in a moment. + expect( + await screen.findByText(/the rescue is opening/i), + ).toBeInTheDocument() + }) + + it('shows the server refusal when the shift already has a rescue', async () => { + markShiftAbsence.mockRejectedValue( + new ApiError(409, 'That shift already has a rescue running.'), + ) + const user = userEvent.setup() + renderWithProviders( {}} />) + + const list = await screen.findByRole('list', { name: "Today's shifts" }) + const [first] = await within(list).findAllByRole('button', { name: 'Mark absent' }) + await user.click(first) + + expect(await screen.findByRole('alert')).toHaveTextContent(/already has a rescue/i) + }) + + it('explains a shift that is not available', async () => { + markShiftAbsence.mockRejectedValue(new ApiError(404, 'Shift not found')) + const user = userEvent.setup() + renderWithProviders( {}} />) + + const list = await screen.findByRole('list', { name: "Today's shifts" }) + const [first] = await within(list).findAllByRole('button', { name: 'Mark absent' }) + await user.click(first) + + expect(await screen.findByRole('alert')).toHaveTextContent(/not available for this location/i) + }) +}) diff --git a/frontend/src/components/ReportAbsenceDialog.tsx b/frontend/src/components/ReportAbsenceDialog.tsx new file mode 100644 index 0000000..6f05097 --- /dev/null +++ b/frontend/src/components/ReportAbsenceDialog.tsx @@ -0,0 +1,117 @@ +/** + * Manager action: mark a shift absent and open the rescue (spec §7.5 + * `POST /api/shifts/{id}/absence`). + * + * The manager is the authority, so the absence needs no WhatsApp confirmation: + * the rescue opens and the first wave goes out. The endpoint enqueues the work + * (202), so the copy says the change applies in a moment instead of pretending + * it is instantaneous — and the live channel is what makes the board move. + */ + +import { useState } from 'react' +import { useMutation, useQueryClient } from '@tanstack/react-query' + +import { markShiftAbsence } from '../services/api' +import { ApiError } from '../services/apiClient' +import { useTodayShifts } from '../services/hooks' +import { Button } from './ui/Button' + +function describeShift(shift: { role: string; startsAt: string; endsAt: string; assigneeName: string | null }): string { + const time = (iso: string) => + new Date(iso).toLocaleTimeString([], { hour: '2-digit', minute: '2-digit', hour12: false }) + const who = shift.assigneeName ?? 'Unassigned' + return `${shift.role} · ${time(shift.startsAt)}–${time(shift.endsAt)} · ${who}` +} + +function absenceError(error: unknown): string { + if (error instanceof ApiError) { + if (error.status === 409) { + return error.message + } + if (error.status === 404) { + return 'That shift is not available for this location.' + } + return 'The absence could not be marked. Try again.' + } + return 'Cannot reach the server. Check your connection and try again.' +} + +export function ReportAbsenceDialog({ dayIso, onClose }: { dayIso: string; onClose: () => void }) { + const { shifts } = useTodayShifts(dayIso) + const queryClient = useQueryClient() + const [queued, setQueued] = useState(null) + + const mutation = useMutation({ + mutationFn: (shiftId: string) => markShiftAbsence(shiftId), + onSuccess: (_accepted, shiftId) => { + setQueued(shiftId) + // The worker opens the rescue: refetch the board, and let the live channel + // deliver whatever else changes. + void queryClient.invalidateQueries({ queryKey: ['shifts'] }) + void queryClient.invalidateQueries({ queryKey: ['rescues'] }) + }, + }) + + const open = (shifts ?? []).filter((shift) => shift.status === 'scheduled') + const busy = mutation.isPending + + return ( +
+
+
+
+

+ Report absence +

+

+ The rescue opens right away and the first wave of offers goes out. No + confirmation is needed from the employee. +

+
+ +
+ + {queued !== null ? ( +

+ Absence accepted: the rescue is opening. The board updates in a moment + (the worker applies it). +

+ ) : null} + + {mutation.isError ? ( +

+ {absenceError(mutation.error)} +

+ ) : null} + +
    + {open.map((shift) => ( +
  • + {describeShift(shift)} + +
  • + ))} +
+ {open.length === 0 ? ( +

+ No shift of today is still scheduled: everything is already absent, + covered or open. +

+ ) : null} +
+
+ ) +} diff --git a/frontend/src/services/api.ts b/frontend/src/services/api.ts index a20b0cd..159f87d 100644 --- a/frontend/src/services/api.ts +++ b/frontend/src/services/api.ts @@ -81,6 +81,14 @@ interface MetricsWire { stuckRescues: number } +/** Manager marks a shift absent: the rescue opens and the first wave goes out. */ +export async function markShiftAbsence(shiftId: string): Promise { + return apiFetch( + `/api/shifts/${encodeURIComponent(shiftId)}/absence`, + { method: 'POST', body: JSON.stringify({}) }, + ) +} + interface MutationAcceptedWire { status: string id: string From e6d3d35799e0c61244e4a2be995109f140ec6154 Mon Sep 17 00:00:00 2001 From: albert Date: Mon, 28 Sep 2026 19:18:49 +0200 Subject: [PATCH 6/6] docs: record the polish round (notices, seed anchors, absence action, cleanups) --- odd/tasks/polish-round.md | 55 +++++++++++++++++++++++++++++++++++++++ 1 file changed, 55 insertions(+) create mode 100644 odd/tasks/polish-round.md diff --git a/odd/tasks/polish-round.md b/odd/tasks/polish-round.md new file mode 100644 index 0000000..6c43317 --- /dev/null +++ b/odd/tasks/polish-round.md @@ -0,0 +1,55 @@ +# Feature: polish-round (notices, seed anchors, absence endpoint, small cleanups) + +**Status**: closed +**Branch**: `feature/realtime-and-polish` +**Spec references**: §7.5 (`POST /api/shifts/{id}/absence`), §6.4 (escalation summary), §9.2 (snapshot), §9.3 (degradation), §10 (retention) +**Related**: `odd/tasks/realtime-events.md` (the WebSocket, first unit of the same round) + +## Why + +Four gaps were left after the dashboard, evaluation and demo work landed. None of +them was a defect in the spec's terms, and all of them were visible to whoever +used the product: + +1. **Notices reached a phone and left no trace.** The manager's escalation and + coverage notices go to a phone (not to an employee conversation), so nothing + was persisted and the rescue timeline showed a case that escalated "by + itself". The race loser's "already covered" reply was an employee message and + was not stored either, so their thread looked unanswered. +2. **A demo day could be dry.** The seeded rotation has fixed windows, so a demo + late in the evening had nothing in progress and nothing starting "today": the + Simulator offered no frame able to report an absence. +3. **Two endpoints of §7.5 were missing**, and the dashboard's "Report absence" + button was dead: the manager could not mark an absence from the product at + all. +4. **Three loose ends**: a beat entry that ticked an obsolete in-memory + scheduler, a declared-but-unused `SENTRY_DSN`, and a Langfuse link that was + always null because the trace id was never stored. + +## What was done + +| Unit | Change | +| --- | --- | +| Notices | One helper sends and audits every manager notice on the case (`MANAGER_NOTIFIED` with the template, or `MANAGER_NOTIFY_SKIPPED` with the reason when the manager has no reachable phone — the attempt is recorded, never a delivery claim). The race loser's reply now carries the conversation. Both audit types render in the timeline ("Manager notified" / "Manager could not be reached"). | +| Seed | When the seeded day is dry the seed anchors two shifts to the current hour — one in progress, one starting soon — assigned to employees the rotation left free that day, and only for the window that is actually missing. Nothing changes during working hours, and the pool comes from the rotation itself so an anchor never schedules somebody who does not exist. | +| Absence endpoint | `POST /api/shifts/{id}/absence` enqueues a task (202) that opens the rescue through the confirmed-absence path, auditing `ABSENCE_MARKED` with the manager as the actor and guarding against a second case on redelivery. 404 for an unknown or out-of-location shift, 409 when it is already absent or already has a live case. | +| The button | The floating "Report absence" action opens a dialog listing today's still-scheduled shifts and marks the chosen one absent, with honest 202 copy and the server's refusal surfaced. | +| Cleanups | `run-due-jobs` (obsolete in-memory tick) is gone from beat and the runtime snapshot it published moved to `reconcile-stale-cases`, which is the heartbeat the degraded banner depends on. Sentry initializes when the DSN is set, best effort and silent when unset. The OTel trace id is stored on each interpretation so the decision detail can link to its Langfuse trace. | + +## Evidence + +| Check | Result | +| --- | --- | +| Backend suite (Windows) | **515 passed**, 2 skipped; ruff and mypy clean | +| Backend suite (Linux CI image) | **515 passed**, 2 skipped, 0 failures | +| Frontend suite | **199 passed**; build, oxlint and tsc clean | +| Notices, live | an escalation recorded `MANAGER_NOTIFIED {template: manager_escalated}` on the case | +| Seed, live | after a reseed: **5 shifts in progress and 2 starting within six hours** | +| Absence action, live (browser) | the button opened the dialog with **10 shifts**; marking one was accepted and Today showed that shift **absent with its rescue searching** | +| Live channel (from `realtime-events`) | two browsers: a rescue produced in one moved the other's board **2.5 s later with no reload** | + +### Deliberately not done + +The EC2 deployment (T3/T4 of `deploy-delivery`) stays deferred by user decision: +the repo artifacts are complete and the runbook has the exact commands, to be run +when the user decides.