diff --git a/backend/druks/events/builder.py b/backend/druks/events/builder.py index e31c44e2..604a6930 100644 --- a/backend/druks/events/builder.py +++ b/backend/druks/events/builder.py @@ -1,32 +1,75 @@ -from sqlalchemy import or_, select +from datetime import datetime +from sqlalchemy import select + +from druks.apps.loader import iter_apps from druks.database import db_session +from druks.durable.models import AgentCall, Artifact, Run from druks.events.feed import FeedItem from druks.events.models import Event -_PAGE_LIMIT_DEFAULT = 200 -_FETCH_LIMIT = 500 - async def build_feed( *, app: str | None = None, + q: str | None = None, + kind: str | None = None, + from_at: datetime | None = None, + until: datetime | None = None, before: int | None = None, - limit: int = _PAGE_LIMIT_DEFAULT, + after: int | None = None, + limit: int = 200, ) -> tuple[list[FeedItem], str | None]: - items = [FeedItem.model_validate(event) for event in await _events(app, before)] - items.sort(key=lambda item: item.seq, reverse=True) - page = items[:limit] - next_cursor = str(page[-1].seq) if len(page) == limit and page else None + """Read one page of recorded Activity and its current destination availability.""" + statement = Event.get_history(app=app).order_by(Event.id.desc()) + if q and q.strip(): + pattern = q.strip().replace("/", "//").replace("%", "/%").replace("_", "/_") + statement = statement.where(Event.subject_label.ilike(f"%{pattern}%", escape="/")) + if kind is not None: + statement = statement.where(Event.type == kind) + if from_at: + statement = statement.where(Event.created_at >= from_at) + if until: + statement = statement.where(Event.created_at < until) + if before is not None: + statement = statement.where(Event.id < before) + if after is not None: + statement = statement.where(Event.id > after) + events = list(await db_session().scalars(statement.limit(limit + 1))) + page = [FeedItem.model_validate(event) for event in events[:limit]] + next_cursor = str(page[-1].seq) if len(events) > limit else None + run_ids = {item.run for item in page if item.run} + runs = ( + set(await db_session().scalars(select(Run.id).where(Run.id.in_(run_ids)))) + if run_ids + else set() + ) + artifact_ids = {item.artifact_id for item in page if item.artifact_id} + artifacts = {} + if artifact_ids: + rows = await db_session().execute( + select(Artifact, AgentCall) + .join(AgentCall, AgentCall.id == Artifact.agent_call_id) + .where(Artifact.id.in_(artifact_ids)) + ) + artifacts = { + artifact.id: bool(call.get_file_path(artifact.path)) for artifact, call in rows + } + subjects = { + (owner.name, subject.subject_type): subject + for owner in iter_apps() + if not owner.builtin + for subject in owner.subjects() + } + available_subjects = {} + for item in page: + identity = (item.app, item.subject_type, item.subject_id) + if identity not in available_subjects: + subject = subjects.get((item.app, item.subject_type)) + available_subjects[identity] = bool( + subject and item.subject_id and await subject.get_for_subject_id(item.subject_id) + ) + item.is_subject_available = available_subjects[identity] + item.is_run_available = item.run in runs + item.is_artifact_available = artifacts.get(item.artifact_id, False) return page, next_cursor - - -async def _events(app: str | None, before: int | None) -> list[Event]: - # This app's events plus any unscoped (core) ones. The log stores the app; - # the core never derives it from the subject. - stmt = select(Event).order_by(Event.id.desc()) - if before: - stmt = stmt.where(Event.id < before) - if app: - stmt = stmt.where(or_(Event.app == app, Event.app.is_(None))) - return list((await db_session().scalars(stmt.limit(_FETCH_LIMIT))).all()) diff --git a/backend/druks/events/feed.py b/backend/druks/events/feed.py index 8bd88b41..3a512551 100644 --- a/backend/druks/events/feed.py +++ b/backend/druks/events/feed.py @@ -1,6 +1,7 @@ from datetime import datetime +from typing import Any -from pydantic import AliasPath, ConfigDict, Field, computed_field +from pydantic import AliasChoices, AliasPath, ConfigDict, Field, computed_field from druks.schemas import Schema @@ -21,6 +22,35 @@ class FeedItem(Schema): subject_type: str | None = None subject_id: str | None = None subject_label: str | None = None + run: str | None = Field(default=None, validation_alias=AliasPath("payload", "run")) + gate: str | None = Field(default=None, validation_alias=AliasPath("payload", "gate")) + parked_at: datetime | None = Field( + default=None, validation_alias=AliasPath("payload", "input_requested_at") + ) + input_request: dict[str, Any] | None = Field( + default=None, validation_alias=AliasPath("payload", "input_request") + ) + result: Any = Field(default=None, validation_alias=AliasPath("payload", "result")) + summary: str | None = Field(default=None, validation_alias=AliasPath("payload", "summary")) + reason: str | None = Field( + default=None, + validation_alias=AliasChoices( + AliasPath("payload", "reason"), AliasPath("payload", "failure") + ), + ) + artifact_id: str | None = Field( + default=None, + validation_alias=AliasChoices( + AliasPath("payload", "artifact_id"), + AliasPath("payload", "input_request", "artifact_id"), + ), + ) + agent_call_id: str | None = Field( + default=None, validation_alias=AliasPath("payload", "agent_call_id") + ) + is_subject_available: bool = False + is_run_available: bool = False + is_artifact_available: bool = False @computed_field @property @@ -32,3 +62,4 @@ class FeedResponse(Schema): items: list[FeedItem] # Event sequence cursor for the next (older) page; None at the tail. next_cursor: str | None = None + kinds: list[str] | None = None diff --git a/backend/druks/events/models.py b/backend/druks/events/models.py index c924bf81..7f36fd52 100644 --- a/backend/druks/events/models.py +++ b/backend/druks/events/models.py @@ -1,7 +1,7 @@ from datetime import datetime from typing import Any -from sqlalchemy import Index +from sqlalchemy import Index, Select, and_, or_, select from sqlalchemy.dialects.postgresql import JSONB from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.orm import Mapped, mapped_column @@ -32,6 +32,36 @@ class Event(Base): created_at: Mapped[datetime] = mapped_column(default=Base.utc_now) payload: Mapped[dict[str, Any]] = mapped_column(JSONB, default=dict) + @classmethod + def get_history(cls, *, app: str | None = None) -> Select[tuple["Event"]]: + # App imports Event while the loader imports App. + from druks.apps.loader import iter_apps + + owners = [owner.name for owner in iter_apps() if not owner.builtin] + statement = select(cls).where( + cls.app.in_(owners), + or_( + cls.type.not_like("workflow.%"), + cls.type.in_( + [ + "workflow.scheduled", + "workflow.parked", + "workflow.failed", + "workflow.cancelled", + ] + ), + and_( + cls.type == "workflow.running", + cls.payload["gate"].astext != "", + cls.payload["input_requested_at"].astext != "", + cls.payload["result"].astext.is_not(None), + ), + ), + ) + if app is not None: + statement = statement.where(cls.app == app) + return statement + @classmethod async def emit( cls, diff --git a/backend/druks/events/routes.py b/backend/druks/events/routes.py index c2288d85..86f71393 100644 --- a/backend/druks/events/routes.py +++ b/backend/druks/events/routes.py @@ -1,75 +1,104 @@ import asyncio +from datetime import UTC +from typing import Annotated -from fastapi import APIRouter, HTTPException, Query, Request, status +from fastapi import APIRouter, Depends, HTTPException, Query, Request, status from fastapi.responses import StreamingResponse +from pydantic import AwareDatetime from druks.api.dependencies import EngineDep -from druks.database import session_scope +from druks.database import db_session, session_scope from druks.durable.live import SSE_HEADERS from druks.events.builder import build_feed from druks.events.feed import FeedResponse +from druks.events.models import Event router = APIRouter(prefix="/api/events", tags=["feed"]) -# Per-connection SSE poll cadence. Short enough that the operator's screen feels -# live, long enough that we're not hammering the DB; the cost is one bounded -# read per tick, so cadence is set by perceived latency rather than load. _SSE_POLL_INTERVAL_SECONDS = 2.0 def _parse_cursor(raw: str | None) -> int | None: - # The cursor is a feed sequence (an event's monotonic pk), opaque to the client — - # it hands back whatever ``next_cursor`` returned. if raw is None: - return None + return try: - return int(raw) - except ValueError as exc: + cursor = int(raw) + if cursor < 0: + raise ValueError + return cursor + except ValueError as error: raise HTTPException( status_code=status.HTTP_400_BAD_REQUEST, - detail=f"Invalid ``before`` cursor: {raw!r}", - ) from exc + detail=f"Invalid event cursor: {raw!r}. Use a returned sequence.", + ) from error -@router.get("", response_model=FeedResponse, response_model_by_alias=True) +def get_filters( + app: str | None = Query(default=None), + q: str | None = Query(default=None), + kind: str | None = Query(default=None), + from_at: Annotated[AwareDatetime | None, Query(alias="from")] = None, + until: AwareDatetime | None = Query(default=None), +) -> dict: + if from_at and until and from_at >= until: + raise HTTPException(status.HTTP_422_UNPROCESSABLE_CONTENT, "until must be after from.") + return { + "app": app, + "q": q, + "kind": kind, + "from_at": from_at.astimezone(UTC) if from_at else None, + "until": until.astimezone(UTC) if until else None, + } + + +@router.get( + "", response_model=FeedResponse, response_model_by_alias=True, response_model_exclude_unset=True +) async def list_feed( + filters: Annotated[dict, Depends(get_filters)], limit: int = Query(default=200, ge=1, le=500), before: str | None = Query(default=None), - app: str | None = Query(default=None), ) -> FeedResponse: cursor = _parse_cursor(before) - items, next_cursor = await build_feed(app=app, before=cursor, limit=limit) - return FeedResponse(items=items, next_cursor=next_cursor) + items, next_cursor = await build_feed(**filters, before=cursor, limit=limit) + response = FeedResponse(items=items, next_cursor=next_cursor) + if before is None: + statement = ( + Event.get_history(app=filters["app"]) + .with_only_columns(Event.type) + .distinct() + .order_by(Event.type) + ) + response.kinds = list(await db_session().scalars(statement)) + return response @router.get("/stream") async def stream_feed( request: Request, engine: EngineDep, - app: str | None = Query(default=None), + filters: Annotated[dict, Depends(get_filters)], + after: str | None = Query(default=None), ) -> StreamingResponse: + last_seq = _parse_cursor(request.headers.get("last-event-id") or after) + async def feed_stream(): - last_seq: int | None = None - first = True + nonlocal last_seq while True: if await request.is_disconnected(): return - # New Session per tick so we don't hold a transaction open across the - # sleep; the open/close cost is irrelevant against the poll cadence. async with session_scope(engine): - items, _next_cursor = await build_feed( - app=app, - before=None, - limit=100 if first else 50, - ) - # Strictly past the last emitted sequence — the monotonic pk never ties, - # so this neither re-sends the boundary event nor drops a same-second one. - fresh = items if last_seq is None else [e for e in items if e.seq > last_seq] - for item in reversed(fresh): # oldest-first within a tick - yield f"data: {item.model_dump_json(by_alias=True)}\n\n" - if fresh: - last_seq = fresh[0].seq # newest just-emitted (page is seq-desc) - first = False + items, cursor = await build_feed(**filters, after=last_seq, limit=100) + # Catch-up can span several pages. Finish the interval before advancing its head. + while last_seq is not None and cursor: + older, cursor = await build_feed( + **filters, after=last_seq, before=int(cursor), limit=100 + ) + items.extend(older) + for item in reversed(items): + yield f"id: {item.seq}\ndata: {item.model_dump_json(by_alias=True)}\n\n" + if items: + last_seq = items[0].seq try: await asyncio.sleep(_SSE_POLL_INTERVAL_SECONDS) except asyncio.CancelledError: diff --git a/backend/tests/test_activity_history.py b/backend/tests/test_activity_history.py new file mode 100644 index 00000000..526b6d93 --- /dev/null +++ b/backend/tests/test_activity_history.py @@ -0,0 +1,255 @@ +from contextlib import asynccontextmanager +from datetime import UTC, datetime, timedelta +from types import SimpleNamespace +from unittest.mock import AsyncMock + +import pytest +from conftest import installation_key +from druks.accounts.dependencies import current_account +from druks.api.server import app +from druks.database import db_session +from druks.durable.models import AgentCall, Artifact +from druks.events import routes +from druks.events.builder import build_feed +from druks.events.models import Event +from druks.testing import seed_run +from druks_field_notes.models import Note +from druks_field_notes.workflows import Summarize +from fastapi import HTTPException +from sqlalchemy import event as sqlalchemy_event + + +@pytest.fixture +async def history(druks_db): + db_session.registry.set(druks_db) + start = datetime(2026, 9, 9, tzinfo=UTC) + rows = [ + Event( + type="summary.ready", app="field_notes", subject_label="ACME%_ One", created_at=start + ), + Event( + type="summary.ready", + app="field_notes", + subject_label="acme%_ Two", + created_at=start + timedelta(hours=1), + ), + Event( + type="plan.prepared", + app="software_factory", + subject_label="ACMEZZ Two", + created_at=start + timedelta(hours=2), + ), + Event( + type="later.kind", + app="field_notes", + subject_label="Unrelated", + created_at=start + timedelta(days=1), + ), + Event(type="hidden", app="core", subject_label="ACME%_ Core"), + Event(type="hidden", app="usage", subject_label="ACME%_ Usage"), + Event(type="hidden", app="not_installed", subject_label="ACME%_ Other"), + Event(type="hidden", subject_label="ACME%_ Unowned"), + ] + druks_db.add_all(rows) + await druks_db.flush() + return rows + + +async def test_literal_search_and_half_open_dates(druks_client, history): + response = await druks_client.get( + "/api/events", + params={ + "q": " aCmE%_ ", + "app": "field_notes", + "kind": "summary.ready", + "from": "2026-09-09T00:00:00Z", + "until": "2026-09-09T01:00:00Z", + }, + ) + assert response.status_code == 200 + data = response.json() + assert [item["seq"] for item in data["items"]] == [history[0].id] + assert data["kinds"] == ["later.kind", "summary.ready"] + blank = (await druks_client.get("/api/events", params={"q": " "})).json() + assert [item["seq"] for item in blank["items"]] == [row.id for row in reversed(history[:4])] + assert blank["kinds"] == ["later.kind", "plan.prepared", "summary.ready"] + + +@pytest.mark.parametrize( + "params", + [ + {"from": "2026-09-10T00:00:00Z", "until": "2026-09-09T00:00:00Z"}, + {"from": "2026-09-09T00:00:00Z", "until": "2026-09-09T00:00:00Z"}, + {"from": "2026-09-09T00:00:00"}, + {"until": "not-a-date"}, + ], +) +async def test_invalid_date_bounds_fail_validation(druks_client, params): + assert (await druks_client.get("/api/events", params=params)).status_code == 422 + + +async def test_kinds_only_query_the_first_page(druks_db, druks_client, history): + statements = [] + + def record(_connection, _cursor, statement, _parameters, _context, _many): + if "SELECT DISTINCT" in statement and "events.type" in statement: + statements.append(statement) + + engine = druks_db.bind.sync_engine + sqlalchemy_event.listen(engine, "before_cursor_execute", record) + try: + first = (await druks_client.get("/api/events", params={"limit": 1})).json() + assert first["kinds"] == ["later.kind", "plan.prepared", "summary.ready"] + assert len(statements) == 1 + second = ( + await druks_client.get( + "/api/events", params={"limit": 1, "before": first["nextCursor"]} + ) + ).json() + assert "kinds" not in second + await build_feed(app="field_notes") + assert len(statements) == 1 + hidden = (await druks_client.get("/api/events", params={"app": "not_installed"})).json() + assert hidden["items"] == hidden["kinds"] == [] + finally: + sqlalchemy_event.remove(engine, "before_cursor_execute", record) + + +async def test_routine_lifecycle_is_not_activity(druks_db): + db_session.registry.set(druks_db) + for kind in ["workflow.running", "workflow.finished", "workflow.step", "workflow.retry"]: + await Event.emit(type=kind, app="field_notes", payload={"run": "gone"}) + await Event.emit( + type="workflow.running", app="field_notes", payload={"gate": "review", "result": {}} + ) + await Event.emit( + type="workflow.running", + app="field_notes", + payload={ + "run": "gone", + "gate": "review", + "input_requested_at": "2026-09-09T01:00:00Z", + "result": {"action": "approve"}, + }, + ) + items, _ = await build_feed() + assert len(items) == 1 + assert items[0].run == "gone" + assert items[0].gate == "review" + assert items[0].parked_at == datetime(2026, 9, 9, 1, tzinfo=UTC) + assert not items[0].is_run_available + + +async def test_exact_artifact_and_recorded_label_survive_later_results_and_deletion( + druks_db, tmp_path, monkeypatch +): + db_session.registry.set(druks_db) + monkeypatch.setenv("DRUKS_DATA_DIR", str(tmp_path)) + note = await Note.create(body="Original work") + run = await seed_run(druks_db, kind=Summarize.kind, subject=note) + key = await installation_key() + calls = [ + AgentCall( + id=f"history-{number}", + run_id=run.id, + agent="field_notes.summarize", + model="test", + sandbox_host_id="test", + api_key_id=key.id, + ) + for number in (1, 2) + ] + druks_db.add_all(calls) + await druks_db.flush() + druks_db.expunge_all() + for call in calls: + await Artifact.record( + call_dir=call.call_dir, + call_id=call.id, + kind="markdown", + title="Result", + content=call.id, + activity={"kind": "summary.ready"}, + ) + items, _ = await build_feed() + assert [item.agent_call_id for item in items] == [calls[1].id, calls[0].id] + assert all( + item.is_artifact_available and item.is_run_available and item.is_subject_available + for item in items + ) + artifact = await Artifact.get_for_call(calls[0].id) + assert items[1].artifact_id == artifact.id + recorded_label = items[1].subject_label + await druks_db.delete(artifact) + await druks_db.delete(await druks_db.merge(note)) + await druks_db.flush() + missing, _ = await build_feed() + assert missing[1].subject_label == recorded_label + assert missing[1].artifact_id == artifact.id + assert not missing[1].is_artifact_available + assert not missing[1].is_subject_available + + +async def test_search_does_not_read_payloads_or_current_subject_text(druks_db): + db_session.registry.set(druks_db) + note = await Note.create(body="Needle") + await Event.emit( + type="build.rejected", + app="field_notes", + subject=note.identity, + label="Recorded work", + payload={"reason": "Needle", "summary": "Needle"}, + ) + assert not (await build_feed(q="needle"))[0] + items, _ = await build_feed(q="recorded") + assert items[0].reason == "Needle" + assert not items[0].run + + +async def test_reconnect_catches_up_all_pages_without_a_kinds_query(druks_db, monkeypatch): + db_session.registry.set(druks_db) + for number in range(205): + await Event.emit(type="summary.ready", app="field_notes", label=str(number)) + first_page, _ = await build_feed(limit=205) + cursor = first_page[-1].seq + queries = [] + + def record(_connection, _cursor, statement, _parameters, _context, _many): + queries.append(statement) + + @asynccontextmanager + async def scope(_engine): + yield druks_db + + monkeypatch.setattr(routes, "session_scope", scope) + monkeypatch.setattr(routes.asyncio, "sleep", AsyncMock()) + request = SimpleNamespace( + headers={"last-event-id": str(cursor)}, is_disconnected=AsyncMock(side_effect=[False, True]) + ) + engine = druks_db.bind.sync_engine + sqlalchemy_event.listen(engine, "before_cursor_execute", record) + try: + response = await routes.stream_feed( + request=request, engine=None, filters={"app": "field_notes"}, after=None + ) + messages = [message async for message in response.body_iterator] + finally: + sqlalchemy_event.remove(engine, "before_cursor_execute", record) + sequences = [int(message.splitlines()[0].removeprefix("id: ")) for message in messages] + assert sequences == sorted(item.seq for item in first_page if item.seq > cursor) + assert len(sequences) == len(set(sequences)) == 204 + assert not any("SELECT DISTINCT" in query for query in queries) + + +async def test_activity_requires_the_existing_account_authorization(druks_client, history): + original = app.dependency_overrides[current_account] + + async def unauthorized(): + raise HTTPException(401, "No operator identity") + + app.dependency_overrides[current_account] = unauthorized + try: + assert (await druks_client.get("/api/events")).status_code == 401 + assert (await druks_client.get("/api/events/stream")).status_code == 401 + finally: + app.dependency_overrides[current_account] = original diff --git a/backend/tests/test_events_feed.py b/backend/tests/test_events_feed.py index 69a2bd01..451d4c04 100644 --- a/backend/tests/test_events_feed.py +++ b/backend/tests/test_events_feed.py @@ -19,7 +19,7 @@ class Pallet(StoredSubject): async def test_feed_carries_what_a_row_is_worded_from(druks_db): note = await Note.create(body="the pump ran hot") await Event.emit( - type="workflow.running", + type="workflow.scheduled", subject=note.identity, label=note.label, app="field_notes", @@ -30,7 +30,7 @@ async def test_feed_carries_what_a_row_is_worded_from(druks_db): by_kind = {row.kind: row for row in (await build_feed())[0]} - started = by_kind["workflow.running"] + started = by_kind["workflow.scheduled"] assert (started.app, started.workflow) == ("field_notes", Summarize.kind) assert (started.subject_type, started.subject_id) == ("note", str(note.id)) # A note declares no label of its own, so it shows itself by identity. @@ -49,12 +49,12 @@ async def test_every_subject_shows_itself(druks_db): assert crate.identity == {"type": "crate", "id": 7} for subject in (crate, pallet): await Event.emit( - type="stocked", subject=subject.identity, label=subject.label, app="faketest" + type="stocked", subject=subject.identity, label=subject.label, app="field_notes" ) await druks_db.delete(crate) await druks_db.flush() - by_type = {row.subject_type: row for row in (await build_feed())[0] if row.app == "faketest"} + by_type = {row.subject_type: row for row in (await build_feed())[0] if row.app == "field_notes"} assert by_type["crate"].subject_label == "CRATE-7" assert by_type["pallet"].subject_label == "pallet 7" @@ -65,7 +65,7 @@ async def test_feed_paginates_same_second_events_without_loss_or_repeat(druks_db # the truncated timestamp used to drop the whole second on the next page; paging on # the monotonic pk covers every event exactly once. for i in range(5): - await Event.emit(type=f"evt-{i}") + await Event.emit(type=f"evt-{i}", app="field_notes") await druks_db.flush() collected = [] diff --git a/frontend/src/api/client.ts b/frontend/src/api/client.ts index a921765f..f512ae21 100644 --- a/frontend/src/api/client.ts +++ b/frontend/src/api/client.ts @@ -8,6 +8,7 @@ import type { Connection, App, FeedResponse, + EventFilters, FileSummary, AppsSettingsResponse, Harness, @@ -222,6 +223,14 @@ async function sendOperation(method: string, path: string, body: unknown): Promi } } +export function eventQuery(params: EventFilters & { limit?: number; before?: string; after?: string }): string { + const query = new URLSearchParams() + for (const [key, value] of Object.entries(params)) { + if (value !== undefined) query.set(key, String(value)) + } + return query.toString() +} + export const api = { dashboardWork: () => getJSON('/api/dashboard/work'), dashboardSchedules: () => getJSON('/api/dashboard/schedules'), @@ -256,13 +265,9 @@ export const api = { postJSON<{ run: string; result: string }>(`/api/runs/${runId}/cancel`, { reason }), retryRun: (runId: string) => postJSON<{ run: string }>(`/api/runs/${runId}/retry`, undefined), - listEvents: (params: { limit?: number; before?: string; app?: string } = {}) => { - const query = new URLSearchParams() - if (params.limit !== undefined) query.set('limit', String(params.limit)) - if (params.before !== undefined) query.set('before', params.before) - if (params.app !== undefined) query.set('app', params.app) - const qs = query.toString() - return getJSON(`/api/events${qs ? `?${qs}` : ''}`) + listEvents: (params: EventFilters & { limit?: number; before?: string } = {}) => { + const query = eventQuery(params) + return getJSON(`/api/events${query ? `?${query}` : ''}`) }, getSettings: () => getJSON('/api/settings'), updateSettings: (body: UpdateSettingsRequest) => diff --git a/frontend/src/api/events.test.ts b/frontend/src/api/events.test.ts new file mode 100644 index 00000000..7d5f9d4a --- /dev/null +++ b/frontend/src/api/events.test.ts @@ -0,0 +1,17 @@ +import { afterEach, expect, it, vi } from 'vitest' +import { api, eventQuery } from './client' + +afterEach(() => vi.unstubAllGlobals()) + +it('uses the same exact filter values for history and stream URLs', async () => { + const filters = { q: 'ACME%_ &?', app: 'field_notes', kind: 'summary.ready', + from: '2026-09-09T00:00:00Z', until: '2026-09-10T00:00:00Z' } + const fetcher = vi.fn(async () => new Response(JSON.stringify({ items: [], nextCursor: null, kinds: [] }))) + vi.stubGlobal('fetch', fetcher) + await api.listEvents(filters) + expect(fetcher).toHaveBeenCalledWith(`/api/events?${eventQuery(filters)}`, expect.anything()) + const query = new URLSearchParams(eventQuery({ ...filters, after: '123' })) + expect(query.get('q')).toBe(filters.q) + expect(query.get('after')).toBe('123') + expect(query.get('kind')).toBe('summary.ready') +}) diff --git a/frontend/src/api/types.ts b/frontend/src/api/types.ts index 2f73a333..87f8f0cf 100644 --- a/frontend/src/api/types.ts +++ b/frontend/src/api/types.ts @@ -749,6 +749,14 @@ export interface UpdateAppsSettingsRequest { } +export interface EventFilters { + q?: string + app?: string + kind?: string + from?: string + until?: string +} + export interface FeedItem { id: string seq: number @@ -764,11 +772,24 @@ export interface FeedItem { // How the subject showed itself ("ENG-767"), snapshotted at write. Absent // exactly when the subject is. subjectLabel?: string | null + run?: string | null + gate?: string | null + parkedAt?: string | null + inputRequest?: InputRequest | null + result?: unknown + summary?: string | null + reason?: string | null + artifactId?: string | null + agentCallId?: string | null + isSubjectAvailable: boolean + isRunAvailable: boolean + isArtifactAvailable: boolean } export interface FeedResponse { items: FeedItem[] nextCursor: string | null + kinds?: string[] } diff --git a/frontend/src/apps/installed.test.tsx b/frontend/src/apps/installed.test.tsx index b0fcbf7b..7b59d9b1 100644 --- a/frontend/src/apps/installed.test.tsx +++ b/frontend/src/apps/installed.test.tsx @@ -1,6 +1,7 @@ import { describe, expect, it } from 'vitest' import type { App } from '../api/types' +import { eventLine } from '../lib/feed' import { registerInstalledApps } from './installed' import { getAppUI } from './registry' @@ -96,4 +97,13 @@ it('uses the declared decision page and preserves encoded subject and request id expect(url.searchParams.get('run')).toBe(target.run) expect(url.searchParams.get('parkedAt')).toBe(target.parkedAt) expect(ui.subjectPath!({ type: 'file', id: '7' }, { run: 'run-one' })).toBe('/decision_app/file/7?run=run-one') + const activity = eventLine({ + id: 'event:1', seq: 1, at: target.parkedAt, kind: 'workflow.parked', app: 'decision_app', + subjectType: 'file', subjectId: '7', run: target.run, parkedAt: target.parkedAt, + isSubjectAvailable: true, isRunAvailable: false, isArtifactAvailable: false, + }) + const destination = new URL(activity.path!, 'https://example.invalid') + expect(destination.pathname).toBe('/decision_app/review/7') + expect(destination.searchParams.get('run')).toBe(target.run) + expect(destination.searchParams.get('parkedAt')).toBe(target.parkedAt) }) diff --git a/frontend/src/apps/registry.tsx b/frontend/src/apps/registry.tsx index a9b1097d..6f801677 100644 --- a/frontend/src/apps/registry.tsx +++ b/frontend/src/apps/registry.tsx @@ -1,5 +1,5 @@ import type { ReactNode } from 'react' -import type { FeedItem, InputRequest } from '../api/types' +import type { FeedItem } from '../api/types' export interface AppRoute { /** A wouter pattern under the router base, such as /notes/:id. */ @@ -7,10 +7,7 @@ export interface AppRoute { render: (params: Record) => ReactNode } -export type ActivityEvent = Pick & { - gate?: string | null - inputRequest?: InputRequest | null -} +export type ActivityEvent = Pick export interface AppUI { name: string diff --git a/frontend/src/apps/software_factory/activity.test.ts b/frontend/src/apps/software_factory/activity.test.ts index 024aa6b1..34478d82 100644 --- a/frontend/src/apps/software_factory/activity.test.ts +++ b/frontend/src/apps/software_factory/activity.test.ts @@ -18,6 +18,7 @@ describe('Factory Activity', () => { ['workflow.cancelled', 'Build stopped'], ])('formats %s through the app registry', (kind, label) => { expect(eventLine({ id: 'event:1', seq: 1, at: '2026-09-09T12:00:00Z', + isSubjectAvailable: true, isRunAvailable: false, isArtifactAvailable: false, kind, app: 'software_factory', workflow: 'software_factory.build' }).label).toBe(label) }) diff --git a/frontend/src/lib/feed.test.ts b/frontend/src/lib/feed.test.ts index 099fa5c2..3303f58a 100644 --- a/frontend/src/lib/feed.test.ts +++ b/frontend/src/lib/feed.test.ts @@ -11,6 +11,9 @@ function event(fields: Partial): FeedItem { seq: 1, at: '2026-07-26T12:00:00Z', kind: 'workflow.running', + isSubjectAvailable: true, + isRunAvailable: false, + isArtifactAvailable: false, ...fields, } } @@ -19,7 +22,7 @@ describe('eventLine', () => { it('names the workflow and what it did', () => { const line = eventLine(event({ kind: 'workflow.running', workflow: 'software_factory.build' })) - expect(line.label).toBe('build started') + expect(line.label).toBe('build response received') expect(line.source).toBe('build') }) @@ -65,7 +68,7 @@ describe('eventLine', () => { }), ) - expect(line.label).toBe('profile started') + expect(line.label).toBe('profile response received') expect(line.subject).toBe('acme/widget') expect(line.path).toBeUndefined() }) @@ -86,3 +89,14 @@ describe('eventLine', () => { expect(line.path).toBeUndefined() }) }) + + +it('retains the recorded run and decision round in Factory links', () => { + const line = eventLine(event({ app: 'software_factory', subjectType: 'work_item', subjectId: '42', + run: 'older-run', parkedAt: '2026-09-09T01:00:00Z', isRunAvailable: true })) + const target = new URL(line.path!, 'https://druks.test') + expect(target.searchParams.get('run')).toBe('older-run') + expect(target.searchParams.get('parkedAt')).toBe('2026-09-09T01:00:00Z') + expect(eventLine(event({ app: 'software_factory', subjectType: 'work_item', subjectId: '42', + isSubjectAvailable: false })).path).toBeUndefined() +}) diff --git a/frontend/src/lib/feed.ts b/frontend/src/lib/feed.ts index 3ccde77b..0ef7911d 100644 --- a/frontend/src/lib/feed.ts +++ b/frontend/src/lib/feed.ts @@ -4,7 +4,8 @@ import { appLabel, getAppUI } from '../apps/registry' // What a workflow doing something is called. The platform owns these words because it // owns the lifecycle; an app's own milestones are already named by their type. const LIFECYCLE_VERBS: Record = { - 'workflow.running': 'started', + 'workflow.scheduled': 'queued', + 'workflow.running': 'response received', 'workflow.parked': 'waiting on you', 'workflow.finished': 'finished', 'workflow.failed': 'failed', @@ -47,9 +48,12 @@ function label(event: FeedItem): string { } function subjectPath(event: FeedItem): string | undefined { - if (event.app && event.subjectType && event.subjectId) { + if (event.isSubjectAvailable && event.app && event.subjectType && event.subjectId) { const ui = getAppUI(event.app) - return ui?.subjectPath?.({ type: event.subjectType, id: event.subjectId }) + const target = event.run + ? { run: event.run, parkedAt: event.parkedAt ?? undefined } + : undefined + return ui?.subjectPath?.({ type: event.subjectType, id: event.subjectId }, target) } return undefined }