mirror of
https://github.com/turnstonelabs/turnstone.git
synced 2026-08-12 23:12:23 -06:00
36419a9809
Workstreams attached to a private project were visible -- including their conversation content -- to holders of admin.cluster.inspect / admin.coordinator (both default builtin-admin permissions), defeating the project's confidentiality boundary. Enforce that a private project's resources are visible only to people IN the project (owner, workstream creator, or an explicit member), even for admins. Surfaces closed: - WorkstreamProjectVisibility bypass narrowed to service scope only (node->console machine plumbing, re-filtered per-user at the console edge). No human principal bypasses; admin.cluster.inspect gates the inspect surfaces, not tenancy. This flows to /dashboard, session listings, the attachment row-gate, cluster_workstreams, cluster_node_detail, and cluster_snapshot/SSE. - cluster_ws_detail 404-masks a workstream in a private project the caller can't see; cluster_ws_live_bulk routes such ids to the denied list (no private-project oracle). - Coordinator operator verbs (history/export/detail/send/approve/set_title/open/ children/tasks/attachments) now enforce project tenancy: _coordinator_tenant_check on coord_endpoint_config, the gate in _resolve_coordinator_or_404 (children/tasks), the tenant_check now run in make_open_handler before rehydrate, and a project-visibility check in _coord_attachment_owner. admin.coordinator gates the surface cluster-wide, but a non-member is 404-masked. The tenant-check mirrors the manager-first + coordinator-kind ladder so kind-isolation is preserved. - service scope is no longer user-assignable: admin_create_token and both turnstone-admin CLI mint paths reject it via reject_unassignable_scopes, so an admin.users holder cannot self-mint a service token and restore the bypass. Service scope is minted only by ServiceTokenManager / the JWT secret. - The events/global node proxy (service-elevated cross-tenant firehose) is gated on admin.cluster.inspect so a plain authenticated user cannot reach it through the console proxy. Updates the OpenAPI description, the row-gate/tenancy-filter docstrings, and adds tests for every surface (visibility predicate + cluster detail/bulk + coordinator history/export/children/open/attachments + events/global proxy + scope-mint rejection); inverts the tests that pinned the old admin-bypass contract.
555 lines
19 KiB
Python
555 lines
19 KiB
Python
"""Tests for the phase-6 polish endpoints (#q-1).
|
|
|
|
Covers:
|
|
|
|
- GET /v1/api/cluster/ws/live — bulk live-block fetch (admin.cluster.inspect).
|
|
- GET /v1/api/workstreams/{ws_id}/metrics — per-coordinator health snapshot.
|
|
|
|
Both endpoints ride on the same test harness as
|
|
``test_coordinator_endpoints.py`` — a minimal Starlette app with an
|
|
auth-injecting middleware, TestClient + MockTransport for the
|
|
upstream node fetches.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import pytest
|
|
from starlette.applications import Starlette
|
|
from starlette.middleware import Middleware
|
|
from starlette.routing import Route
|
|
from starlette.testclient import TestClient
|
|
|
|
from tests._coord_test_helpers import (
|
|
_AuthMiddleware,
|
|
_build_mgr,
|
|
_fake_registry,
|
|
_FakeConfigStore,
|
|
)
|
|
from turnstone.console.server import (
|
|
cluster_ws_live_bulk,
|
|
coordinator_metrics,
|
|
)
|
|
from turnstone.core.storage._sqlite import SQLiteBackend
|
|
|
|
|
|
@pytest.fixture
|
|
def storage(tmp_path):
|
|
return SQLiteBackend(str(tmp_path / "phase6.db"))
|
|
|
|
|
|
def _make_client(storage, *, coord_mgr=None) -> TestClient:
|
|
app = Starlette(
|
|
routes=[
|
|
Route("/v1/api/cluster/ws/live", cluster_ws_live_bulk, methods=["GET"]),
|
|
Route(
|
|
"/v1/api/workstreams/{ws_id}/metrics",
|
|
coordinator_metrics,
|
|
methods=["GET"],
|
|
),
|
|
],
|
|
middleware=[Middleware(_AuthMiddleware)],
|
|
)
|
|
app.state.coord_mgr = coord_mgr
|
|
app.state.coord_adapter = coord_mgr._adapter if coord_mgr is not None else None
|
|
app.state.config_store = _FakeConfigStore({"coordinator.model_alias": "gpt-4"})
|
|
app.state.coord_registry = _fake_registry() if coord_mgr is not None else None
|
|
app.state.coord_registry_error = "" if coord_mgr else "registry missing"
|
|
app.state.auth_storage = storage
|
|
app.state.jwt_secret = "x" * 64
|
|
return TestClient(app)
|
|
|
|
|
|
def _seed_workstream(
|
|
storage: SQLiteBackend,
|
|
*,
|
|
ws_id: str,
|
|
node_id: str,
|
|
user_id: str = "user-1",
|
|
kind: str = "interactive",
|
|
state: str = "idle",
|
|
parent_ws_id: str | None = None,
|
|
created: str | None = None,
|
|
) -> None:
|
|
storage.register_workstream(
|
|
ws_id,
|
|
node_id=node_id,
|
|
user_id=user_id,
|
|
name=f"ws-{ws_id[:4]}",
|
|
state=state,
|
|
kind=kind,
|
|
parent_ws_id=parent_ws_id,
|
|
)
|
|
if created is not None:
|
|
# Override the created timestamp directly — register_workstream
|
|
# stamps "now", so we need a second write to test the
|
|
# spawns_last_hour boundary.
|
|
import sqlalchemy as sa
|
|
|
|
from turnstone.core.storage._sqlite import workstreams
|
|
|
|
with storage._conn() as conn:
|
|
conn.execute(
|
|
sa.update(workstreams).where(workstreams.c.ws_id == ws_id).values(created=created)
|
|
)
|
|
conn.commit()
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# GET /v1/api/cluster/ws/live — bulk live-block fetch
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
_ADMIN_HEADERS = {"X-Test-User": "user-1", "X-Test-Perms": "admin.cluster.inspect"}
|
|
_OWNER_HEADERS = _ADMIN_HEADERS # same caller; permission grants inspect
|
|
|
|
|
|
def test_bulk_live_requires_permission(storage):
|
|
client = _make_client(storage, coord_mgr=_build_mgr(storage))
|
|
resp = client.get(
|
|
"/v1/api/cluster/ws/live?ids=" + "a" * 32,
|
|
headers={"X-Test-User": "u", "X-Test-Perms": "read"},
|
|
)
|
|
assert resp.status_code == 403
|
|
|
|
|
|
def test_bulk_live_empty_ids_returns_empty_body(storage):
|
|
client = _make_client(storage, coord_mgr=_build_mgr(storage))
|
|
resp = client.get("/v1/api/cluster/ws/live?ids=", headers=_ADMIN_HEADERS)
|
|
assert resp.status_code == 200
|
|
body = resp.json()
|
|
assert body == {"results": {}, "denied": [], "truncated": False}
|
|
|
|
|
|
def test_bulk_live_strips_invalid_ids(storage):
|
|
"""IDs failing the hex-regex are silently dropped; duplicates
|
|
collapse."""
|
|
client = _make_client(storage, coord_mgr=_build_mgr(storage))
|
|
# NOT-HEX is invalid; the valid id is 32 chars hex but unknown to
|
|
# storage → shows up as denied.
|
|
resp = client.get(
|
|
"/v1/api/cluster/ws/live?ids=NOT-HEX,NOT-HEX,," + ("a" * 32) + "," + ("a" * 32),
|
|
headers=_ADMIN_HEADERS,
|
|
)
|
|
assert resp.status_code == 200
|
|
body = resp.json()
|
|
# Invalid / empty / duplicate ids trimmed; only the one valid-but-
|
|
# missing id is reported as denied.
|
|
assert body["denied"] == ["a" * 32]
|
|
assert body["results"] == {}
|
|
|
|
|
|
def test_bulk_live_caps_ids_at_50(storage):
|
|
"""Ids past the server-side cap truncate with truncated=true."""
|
|
client = _make_client(storage, coord_mgr=_build_mgr(storage))
|
|
# 60 fake ids → cap=50 keeps the first 50 (dedup preserves order).
|
|
ids = ",".join(f"{i:064x}" for i in range(60))
|
|
resp = client.get(
|
|
"/v1/api/cluster/ws/live?ids=" + ids,
|
|
headers=_ADMIN_HEADERS,
|
|
)
|
|
assert resp.status_code == 200
|
|
body = resp.json()
|
|
assert body["truncated"] is True
|
|
# All 50 kept ids resolve to 'denied' (no storage rows) — their
|
|
# inclusion in the response proves the cap took the head 50.
|
|
assert len(body["denied"]) == 50
|
|
|
|
|
|
def test_bulk_live_admin_bypass_returns_live(storage):
|
|
"""An admin user (holds admin.users or admin.roles, not just
|
|
admin.cluster.inspect) bypasses tenancy and sees non-owned rows'
|
|
live blocks. Coordinator live-block synthesis is in-process, so
|
|
results is populated without any upstream node fetch."""
|
|
mgr = _build_mgr(storage)
|
|
ws = mgr.create(user_id="other-user")
|
|
client = _make_client(storage, coord_mgr=mgr)
|
|
resp = client.get(
|
|
"/v1/api/cluster/ws/live?ids=" + ws.id,
|
|
headers={
|
|
"X-Test-User": "user-1",
|
|
# admin.users grants the _is_admin bypass in addition to
|
|
# admin.cluster.inspect for the endpoint itself.
|
|
"X-Test-Perms": "admin.cluster.inspect,admin.users",
|
|
},
|
|
)
|
|
assert resp.status_code == 200
|
|
body = resp.json()
|
|
assert ws.id in body["results"]
|
|
assert body["denied"] == []
|
|
|
|
|
|
def test_bulk_live_cluster_wide_visibility(storage):
|
|
"""A project-less workstream has no tenancy to enforce, so any
|
|
``admin.cluster.inspect`` caller sees it in ``results``. ``denied``
|
|
is reserved for ids that don't correspond to a persisted workstream
|
|
(no existence oracle for unknown ids)."""
|
|
ws_id = "b" * 32
|
|
_seed_workstream(storage, ws_id=ws_id, node_id="node-a", user_id="stranger")
|
|
client = _make_client(storage, coord_mgr=_build_mgr(storage))
|
|
resp = client.get(
|
|
f"/v1/api/cluster/ws/live?ids={ws_id}",
|
|
headers={"X-Test-User": "user-1", "X-Test-Perms": "admin.cluster.inspect"},
|
|
)
|
|
assert resp.status_code == 200
|
|
body = resp.json()
|
|
assert ws_id in body["results"]
|
|
assert body["denied"] == []
|
|
|
|
|
|
def test_bulk_live_private_project_row_routes_to_denied(storage):
|
|
"""A workstream in a private project the caller isn't a member of
|
|
routes to ``denied``, not ``results`` — a cluster admin gets no
|
|
private-project oracle from the bulk surface either."""
|
|
storage.create_project("proj-secret", "Secret", "alice")
|
|
ws_id = "c" * 32
|
|
storage.register_workstream(ws_id, node_id="node-a", user_id="alice", project_id="proj-secret")
|
|
client = _make_client(storage, coord_mgr=_build_mgr(storage))
|
|
resp = client.get(
|
|
f"/v1/api/cluster/ws/live?ids={ws_id}",
|
|
headers={"X-Test-User": "stranger", "X-Test-Perms": "admin.cluster.inspect"},
|
|
)
|
|
assert resp.status_code == 200
|
|
body = resp.json()
|
|
assert body["results"] == {}
|
|
assert body["denied"] == [ws_id]
|
|
|
|
|
|
def test_bulk_live_private_project_row_visible_to_member(storage):
|
|
"""A project member sees the row (routes to ``results``); the live
|
|
block is null only because the coordinator row isn't loaded."""
|
|
storage.create_project("proj-secret", "Secret", "alice")
|
|
storage.add_project_member("proj-secret", "member-bob")
|
|
ws_id = "c" * 32
|
|
storage.register_workstream(
|
|
ws_id,
|
|
node_id="console",
|
|
user_id="alice",
|
|
kind="coordinator",
|
|
project_id="proj-secret",
|
|
)
|
|
client = _make_client(storage, coord_mgr=_build_mgr(storage))
|
|
resp = client.get(
|
|
f"/v1/api/cluster/ws/live?ids={ws_id}",
|
|
headers={"X-Test-User": "member-bob", "X-Test-Perms": "admin.cluster.inspect"},
|
|
)
|
|
assert resp.status_code == 200
|
|
body = resp.json()
|
|
assert ws_id in body["results"]
|
|
assert body["denied"] == []
|
|
|
|
|
|
def test_bulk_live_unknown_ids_route_to_denied(storage):
|
|
"""Unknown ids (not in storage) land in ``denied`` so the endpoint
|
|
can't be used as an existence oracle."""
|
|
ws_id = "c" * 32 # not seeded
|
|
client = _make_client(storage, coord_mgr=_build_mgr(storage))
|
|
resp = client.get(
|
|
f"/v1/api/cluster/ws/live?ids={ws_id}",
|
|
headers={"X-Test-User": "user-1", "X-Test-Perms": "admin.cluster.inspect"},
|
|
)
|
|
assert resp.status_code == 200
|
|
body = resp.json()
|
|
assert body["denied"] == [ws_id]
|
|
assert body["results"] == {}
|
|
|
|
|
|
def test_bulk_live_coordinator_row_uses_manager_snapshot(storage):
|
|
"""A coordinator ws_id routes through _fetch_live_block's
|
|
coordinator branch — live is populated from the in-process manager
|
|
even though the pseudo-node has no /dashboard endpoint."""
|
|
mgr = _build_mgr(storage)
|
|
ws = mgr.create(user_id="user-1")
|
|
client = _make_client(storage, coord_mgr=mgr)
|
|
resp = client.get(
|
|
f"/v1/api/cluster/ws/live?ids={ws.id}",
|
|
headers=_OWNER_HEADERS,
|
|
)
|
|
assert resp.status_code == 200
|
|
body = resp.json()
|
|
assert ws.id in body["results"]
|
|
live = body["results"][ws.id]
|
|
assert live is not None
|
|
assert "pending_approval" in live
|
|
# The details list is always present on the wire — empty when no
|
|
# approval is pending so the JS can `key in row` without surprise.
|
|
# Replaces 1.6's singular ``pending_approval_detail`` null
|
|
# (breaking, 1.7).
|
|
assert "pending_approval_details" in live
|
|
assert live["pending_approval_details"] == []
|
|
|
|
|
|
def test_bulk_live_coordinator_row_includes_pending_approval_details(storage):
|
|
"""When an approval cycle is live on a coord UI, the live block
|
|
surfaces one detail entry per cycle with merged items +
|
|
judge_verdict through the coord-pseudo-node path. End-to-end
|
|
equivalent of the dashboard test in test_server_authz, but for
|
|
the console live-bulk endpoint that the coord tree UI actually
|
|
consumes."""
|
|
from turnstone.core.session_ui_base import ApprovalCycle
|
|
|
|
mgr = _build_mgr(storage)
|
|
ws = mgr.create(user_id="user-1")
|
|
items = [
|
|
{
|
|
"call_id": "c-99",
|
|
"header": "spawn_workstream",
|
|
"preview": "{...}",
|
|
"func_name": "spawn_workstream",
|
|
"approval_label": "spawn_workstream",
|
|
"needs_approval": True,
|
|
}
|
|
]
|
|
card = {
|
|
"type": "approve_request",
|
|
"cycle_id": "cyc-99",
|
|
"items": ws.ui._serialize_approval_items(items),
|
|
"judge_pending": False,
|
|
}
|
|
ws.ui._register_approval_cycle(ApprovalCycle(items, card, None))
|
|
ws.ui._llm_verdicts["c-99"] = {
|
|
"recommendation": "approve",
|
|
"risk_level": "low",
|
|
"tier": "llm",
|
|
}
|
|
client = _make_client(storage, coord_mgr=mgr)
|
|
resp = client.get(
|
|
f"/v1/api/cluster/ws/live?ids={ws.id}",
|
|
headers=_OWNER_HEADERS,
|
|
)
|
|
assert resp.status_code == 200
|
|
live = resp.json()["results"][ws.id]
|
|
assert live["pending_approval"] is True # boolean derived flag
|
|
details = live["pending_approval_details"]
|
|
assert len(details) == 1
|
|
detail = details[0]
|
|
assert detail["cycle_id"] == "cyc-99"
|
|
assert detail["call_id"] == "c-99"
|
|
assert detail["items"][0]["func_name"] == "spawn_workstream"
|
|
assert detail["items"][0]["judge_verdict"]["recommendation"] == "approve"
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# GET /v1/api/workstreams/{ws_id}/metrics — per-coordinator health snapshot
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
_METRICS_HEADERS = {"X-Test-User": "user-1", "X-Test-Perms": "admin.coordinator"}
|
|
|
|
|
|
def test_metrics_requires_permission(storage):
|
|
mgr = _build_mgr(storage)
|
|
ws = mgr.create(user_id="user-1")
|
|
client = _make_client(storage, coord_mgr=mgr)
|
|
resp = client.get(
|
|
f"/v1/api/workstreams/{ws.id}/metrics",
|
|
headers={"X-Test-User": "user-1", "X-Test-Perms": "read"},
|
|
)
|
|
assert resp.status_code == 403
|
|
|
|
|
|
def test_metrics_invalid_ws_id_400(storage):
|
|
mgr = _build_mgr(storage)
|
|
client = _make_client(storage, coord_mgr=mgr)
|
|
resp = client.get(
|
|
"/v1/api/workstreams/NOT-HEX/metrics",
|
|
headers=_METRICS_HEADERS,
|
|
)
|
|
assert resp.status_code == 400
|
|
|
|
|
|
def test_metrics_any_admin_coordinator_caller_can_read(storage):
|
|
"""Trusted-team visibility: metrics are readable by any caller
|
|
with ``admin.coordinator`` regardless of the coordinator owner."""
|
|
mgr = _build_mgr(storage)
|
|
ws = mgr.create(user_id="stranger")
|
|
client = _make_client(storage, coord_mgr=mgr)
|
|
resp = client.get(
|
|
f"/v1/api/workstreams/{ws.id}/metrics",
|
|
headers=_METRICS_HEADERS,
|
|
)
|
|
assert resp.status_code == 200
|
|
assert resp.json()["ws_id"] == ws.id
|
|
|
|
|
|
def test_metrics_empty_coordinator_defaults(storage):
|
|
"""A freshly created coordinator with no spawns / no verdicts
|
|
returns zero / empty defaults."""
|
|
mgr = _build_mgr(storage)
|
|
ws = mgr.create(user_id="user-1")
|
|
client = _make_client(storage, coord_mgr=mgr)
|
|
resp = client.get(
|
|
f"/v1/api/workstreams/{ws.id}/metrics",
|
|
headers=_METRICS_HEADERS,
|
|
)
|
|
assert resp.status_code == 200
|
|
body = resp.json()
|
|
assert body["ws_id"] == ws.id
|
|
assert body["spawns_total"] == 0
|
|
assert body["spawns_last_hour"] == 0
|
|
assert body["child_state_counts"] == {}
|
|
assert body["judge_fallback_rate"] == 0.0
|
|
assert body["wait_completions"] == 0
|
|
assert body["wait_timeouts"] == 0
|
|
assert body["wait_avg_elapsed"] == 0.0
|
|
|
|
|
|
def test_metrics_spawns_and_state_counts(storage):
|
|
"""spawns_total counts ALL children (including closed); state
|
|
histogram groups by current state. All children share the
|
|
coordinator's owner so the non-admin tenant filter on the
|
|
aggregate queries counts them all (see next test for the
|
|
cross-tenant filter behaviour)."""
|
|
mgr = _build_mgr(storage)
|
|
ws = mgr.create(user_id="user-1")
|
|
_seed_workstream(
|
|
storage,
|
|
ws_id="aa" * 16,
|
|
node_id="node-a",
|
|
user_id="user-1",
|
|
parent_ws_id=ws.id,
|
|
state="idle",
|
|
)
|
|
_seed_workstream(
|
|
storage,
|
|
ws_id="bb" * 16,
|
|
node_id="node-a",
|
|
user_id="user-1",
|
|
parent_ws_id=ws.id,
|
|
state="running",
|
|
)
|
|
_seed_workstream(
|
|
storage,
|
|
ws_id="cc" * 16,
|
|
node_id="node-a",
|
|
user_id="user-1",
|
|
parent_ws_id=ws.id,
|
|
state="closed",
|
|
)
|
|
client = _make_client(storage, coord_mgr=mgr)
|
|
resp = client.get(
|
|
f"/v1/api/workstreams/{ws.id}/metrics",
|
|
headers=_METRICS_HEADERS,
|
|
)
|
|
assert resp.status_code == 200
|
|
body = resp.json()
|
|
assert body["spawns_total"] == 3
|
|
assert body["child_state_counts"] == {"idle": 1, "running": 1, "closed": 1}
|
|
|
|
|
|
def test_metrics_cluster_wide_aggregates(storage):
|
|
"""Trusted-team model: aggregates are cluster-wide across every
|
|
caller with ``admin.coordinator``. Every child under the
|
|
coordinator counts, regardless of the ``user_id`` on the row.
|
|
"""
|
|
mgr = _build_mgr(storage)
|
|
ws = mgr.create(user_id="alice")
|
|
_seed_workstream(
|
|
storage,
|
|
ws_id="aa" * 16,
|
|
node_id="node-a",
|
|
user_id="alice",
|
|
parent_ws_id=ws.id,
|
|
state="idle",
|
|
)
|
|
_seed_workstream(
|
|
storage,
|
|
ws_id="bb" * 16,
|
|
node_id="node-a",
|
|
user_id="bob",
|
|
parent_ws_id=ws.id,
|
|
state="running",
|
|
)
|
|
client = _make_client(storage, coord_mgr=mgr)
|
|
|
|
# Every admin.coordinator caller sees both children.
|
|
for caller in ("alice", "bob", "admin-1"):
|
|
resp = client.get(
|
|
f"/v1/api/workstreams/{ws.id}/metrics",
|
|
headers={"X-Test-User": caller, "X-Test-Perms": "admin.coordinator"},
|
|
)
|
|
assert resp.status_code == 200, caller
|
|
body = resp.json()
|
|
assert body["spawns_total"] == 2, caller
|
|
assert body["child_state_counts"] == {"idle": 1, "running": 1}, caller
|
|
|
|
|
|
def test_metrics_judge_fallback_rate_substring_match(storage):
|
|
"""judge_fallback_rate is computed from any verdict whose ``tier``
|
|
field contains 'fallback' (case-insensitive). Supports tiers like
|
|
'llm_fallback', 'LLM_FALLBACK', 'fallback_deterministic'."""
|
|
import uuid
|
|
|
|
mgr = _build_mgr(storage)
|
|
ws = mgr.create(user_id="user-1")
|
|
# Three verdicts, two marked fallback (one LLM_FALLBACK, one
|
|
# llm_fallback → both match case-insensitive substring).
|
|
for tier in ("llm_primary", "LLM_FALLBACK", "llm_fallback"):
|
|
storage.create_intent_verdict(
|
|
verdict_id=uuid.uuid4().hex,
|
|
ws_id=ws.id,
|
|
call_id="c-" + tier,
|
|
func_name="f",
|
|
func_args="{}",
|
|
intent_summary="",
|
|
risk_level="low",
|
|
confidence=0.9,
|
|
recommendation="allow",
|
|
reasoning="",
|
|
evidence="",
|
|
tier=tier,
|
|
judge_model="j",
|
|
latency_ms=1,
|
|
)
|
|
client = _make_client(storage, coord_mgr=mgr)
|
|
resp = client.get(
|
|
f"/v1/api/workstreams/{ws.id}/metrics",
|
|
headers=_METRICS_HEADERS,
|
|
)
|
|
assert resp.status_code == 200
|
|
body = resp.json()
|
|
# 2 / 3 verdicts matched → 0.667 (rounded to 3 places).
|
|
assert body["judge_fallback_rate"] == pytest.approx(0.667, abs=1e-3)
|
|
assert body["intent_verdicts_sample"] == 3
|
|
|
|
|
|
def test_metrics_spawns_last_hour_boundary(storage):
|
|
"""Only children whose created timestamp is within the last 3600s
|
|
count toward spawns_last_hour; older children count toward
|
|
spawns_total but not the hour bucket."""
|
|
import time
|
|
from datetime import UTC, datetime, timedelta
|
|
|
|
mgr = _build_mgr(storage)
|
|
ws = mgr.create(user_id="user-1")
|
|
# One child created "now" (within the window); one created 2
|
|
# hours ago (outside the window).
|
|
recent_iso = datetime.fromtimestamp(time.time(), tz=UTC).strftime("%Y-%m-%dT%H:%M:%S")
|
|
old_iso = (datetime.fromtimestamp(time.time(), tz=UTC) - timedelta(hours=2)).strftime(
|
|
"%Y-%m-%dT%H:%M:%S"
|
|
)
|
|
_seed_workstream(
|
|
storage,
|
|
ws_id="aa" * 16,
|
|
node_id="node-a",
|
|
parent_ws_id=ws.id,
|
|
state="idle",
|
|
created=recent_iso,
|
|
)
|
|
_seed_workstream(
|
|
storage,
|
|
ws_id="bb" * 16,
|
|
node_id="node-a",
|
|
parent_ws_id=ws.id,
|
|
state="closed",
|
|
created=old_iso,
|
|
)
|
|
client = _make_client(storage, coord_mgr=mgr)
|
|
resp = client.get(
|
|
f"/v1/api/workstreams/{ws.id}/metrics",
|
|
headers=_METRICS_HEADERS,
|
|
)
|
|
assert resp.status_code == 200
|
|
body = resp.json()
|
|
assert body["spawns_total"] == 2
|
|
assert body["spawns_last_hour"] == 1
|