diff --git a/Server/src/overseer/__init__.py b/Server/src/overseer/__init__.py new file mode 100644 index 000000000..23bd2e390 --- /dev/null +++ b/Server/src/overseer/__init__.py @@ -0,0 +1,22 @@ +"""Provider-neutral control-plane primitives for supervising commercial AI agents.""" + +from .ledger import EventLedger, LedgerEvent, TeamSummary +from .simulation import run_support_ticket_simulation +from .teams import TEAM_SPECS + + +def create_overseer_app(*args, **kwargs): + """Lazily create the optional FastAPI dashboard adapter.""" + from .api import create_overseer_app as _create_overseer_app + + return _create_overseer_app(*args, **kwargs) + + +__all__ = [ + "EventLedger", + "LedgerEvent", + "TeamSummary", + "create_overseer_app", + "run_support_ticket_simulation", + "TEAM_SPECS", +] diff --git a/Server/src/overseer/api.py b/Server/src/overseer/api.py new file mode 100644 index 000000000..d32ec6fa7 --- /dev/null +++ b/Server/src/overseer/api.py @@ -0,0 +1,80 @@ +"""Read-only HTTP view of the overseer ledger.""" + +from __future__ import annotations + +from dataclasses import asdict + +from hmac import compare_digest +from decimal import Decimal +from typing import Collection + +from fastapi import FastAPI, Header, HTTPException, Query + +from .ledger import EventLedger + + +def create_overseer_app( + ledger: EventLedger, + *, + api_key: str, + authorized_team_ids: Collection[str] | None = None, +) -> FastAPI: + """Create an authenticated, read-only API backed by an existing ledger.""" + if not api_key: + raise ValueError("api_key is required") + app = FastAPI(title="Overseer Control Plane", version="1.0") + scoped_team_ids = ( + frozenset(authorized_team_ids) if authorized_team_ids is not None else None + ) + + def serialize(value: object) -> object: + """Convert ledger values into JSON-compatible response values.""" + if isinstance(value, Decimal): + return float(value) + if isinstance(value, dict): + return {key: serialize(item) for key, item in value.items()} + if isinstance(value, list): + return [serialize(item) for item in value] + return value + + def authorize(requested_team_id: str | None, presented_key: str | None) -> None: + """Validate the API key and optional team scope for an endpoint request.""" + if presented_key is None or not compare_digest(presented_key, api_key): + raise HTTPException(status_code=401, detail="authentication required") + if scoped_team_ids is not None: + if requested_team_id is None or requested_team_id not in scoped_team_ids: + raise HTTPException(status_code=403, detail="team access denied") + + @app.get("/events") + def events( + team_id: str | None = Query(default=None), + limit: int = Query(default=100, ge=1, le=500), + x_overseer_api_key: str | None = Header(default=None), + ) -> list[dict]: + """Return recent ledger events visible to the authenticated caller.""" + authorize(team_id, x_overseer_api_key) + return [serialize(asdict(event)) for event in ledger.list_events(team_id, limit)] + + @app.get("/summaries") + def summaries( + team_id: str | None = Query(default=None), + x_overseer_api_key: str | None = Header(default=None), + ) -> list[dict]: + """Return financial summaries visible to the authenticated caller.""" + authorize(team_id, x_overseer_api_key) + return [serialize(asdict(summary)) for summary in ledger.summarize(team_id)] + + @app.get("/approvals") + def pending_approvals( + team_id: str | None = Query(default=None), + limit: int = Query(default=100, ge=1, le=500), + x_overseer_api_key: str | None = Header(default=None), + ) -> list[dict]: + """Return pending approval events visible to the authenticated caller.""" + authorize(team_id, x_overseer_api_key) + return [ + serialize(asdict(event)) + for event in ledger.list_pending_approvals(team_id, limit) + ] + + return app diff --git a/Server/src/overseer/ledger.py b/Server/src/overseer/ledger.py new file mode 100644 index 000000000..aa7c6db4c --- /dev/null +++ b/Server/src/overseer/ledger.py @@ -0,0 +1,327 @@ +"""Durable activity, approval, and economics ledger for the overseer dashboard. + +The ledger deliberately stores facts rather than provider-specific objects. Connectors +for a CRM, billing provider, or model runtime can translate their callbacks into events. +""" + +from __future__ import annotations + +from dataclasses import asdict, dataclass +from datetime import datetime, timezone +from decimal import Decimal, InvalidOperation +import json +import sqlite3 +from threading import RLock +from typing import Any +from uuid import uuid4 + + +@dataclass(frozen=True) +class LedgerEvent: + """Immutable representation of one activity, revenue, or cost event.""" + + id: str + category: str + event_type: str + team_id: str + agent_id: str + task_id: str | None + amount: Decimal + currency: str + requires_approval: bool + approved_by: str | None + metadata: dict[str, Any] + created_at: str + + +@dataclass(frozen=True) +class TeamSummary: + """Aggregated financial and approval state for one team and currency.""" + + team_id: str + currency: str + revenue: Decimal + costs: Decimal + profit: Decimal + activity_count: int + pending_approvals: int + + +class EventLedger: + """SQLite-backed ledger suitable for simulation and a later API adapter.""" + + def __init__(self, connection: sqlite3.Connection): + """Initialize the ledger and migrate an existing database if needed.""" + self._lock = RLock() + database_path = connection.execute("PRAGMA database_list").fetchone()[2] + if database_path: + self._connection = sqlite3.connect(database_path, check_same_thread=False) + else: + self._connection = sqlite3.connect(":memory:", check_same_thread=False) + connection.backup(self._connection) + self._connection.row_factory = sqlite3.Row + self._connection.execute("PRAGMA foreign_keys = ON") + self._connection.executescript( + """ + CREATE TABLE IF NOT EXISTS ledger_events ( + id TEXT PRIMARY KEY, + category TEXT NOT NULL CHECK (category IN ('activity', 'revenue', 'cost')), + event_type TEXT NOT NULL, + team_id TEXT NOT NULL, + agent_id TEXT NOT NULL, + task_id TEXT, + amount TEXT NOT NULL DEFAULT '0', + currency TEXT NOT NULL, + requires_approval INTEGER NOT NULL DEFAULT 0, + approved_by TEXT, + metadata_json TEXT NOT NULL, + created_at TEXT NOT NULL + ); + """ + ) + self._migrate_legacy_amount_column() + self._connection.executescript( + """ + CREATE INDEX IF NOT EXISTS idx_ledger_events_team + ON ledger_events(team_id, created_at); + CREATE INDEX IF NOT EXISTS idx_ledger_events_approval + ON ledger_events(requires_approval, approved_by); + """ + ) + self._connection.commit() + + def _migrate_legacy_amount_column(self) -> None: + """Rewrite legacy REAL amounts as text to preserve exact decimals.""" + columns = self._connection.execute("PRAGMA table_info(ledger_events)").fetchall() + amount_column = next((column for column in columns if column["name"] == "amount"), None) + if amount_column is None or amount_column["type"].upper() != "REAL": + return + + self._connection.execute("BEGIN") + try: + self._connection.execute("ALTER TABLE ledger_events RENAME TO ledger_events_legacy") + self._connection.execute( + """ + CREATE TABLE ledger_events ( + id TEXT PRIMARY KEY, + category TEXT NOT NULL CHECK (category IN ('activity', 'revenue', 'cost')), + event_type TEXT NOT NULL, + team_id TEXT NOT NULL, + agent_id TEXT NOT NULL, + task_id TEXT, + amount TEXT NOT NULL DEFAULT '0', + currency TEXT NOT NULL, + requires_approval INTEGER NOT NULL DEFAULT 0, + approved_by TEXT, + metadata_json TEXT NOT NULL, + created_at TEXT NOT NULL + ) + """ + ) + self._connection.execute( + """ + INSERT INTO ledger_events + SELECT id, category, event_type, team_id, agent_id, task_id, + CAST(amount AS TEXT), currency, requires_approval, approved_by, + metadata_json, created_at + FROM ledger_events_legacy + """ + ) + self._connection.execute("DROP TABLE ledger_events_legacy") + self._connection.commit() + except Exception: + self._connection.rollback() + raise + + def record( + self, + *, + category: str, + event_type: str, + team_id: str, + agent_id: str, + task_id: str | None = None, + amount: Decimal | int | float | str = 0, + currency: str = "USD", + requires_approval: bool = False, + approved_by: str | None = None, + metadata: dict[str, Any] | None = None, + created_at: datetime | None = None, + ) -> LedgerEvent: + """Validate and persist one ledger event, returning its immutable record.""" + if category not in {"activity", "revenue", "cost"}: + raise ValueError("category must be activity, revenue, or cost") + if not team_id or not agent_id or not event_type: + raise ValueError("team_id, agent_id, and event_type are required") + if requires_approval and approved_by is not None: + raise ValueError("an event cannot require approval and already be approved") + try: + exact_amount = Decimal(str(amount)) + except (InvalidOperation, ValueError): + raise ValueError("amount must be a finite non-negative number") from None + if not exact_amount.is_finite(): + raise ValueError("amount must be finite and non-negative") + if exact_amount < 0: + raise ValueError("amount must be non-negative; use category to distinguish revenue and cost") + if created_at is not None: + if created_at.tzinfo is None or created_at.utcoffset() is None: + raise ValueError("created_at must be timezone-aware") + normalized_created_at = created_at.astimezone(timezone.utc) + else: + normalized_created_at = datetime.now(timezone.utc) + + event = LedgerEvent( + id=str(uuid4()), + category=category, + event_type=event_type, + team_id=team_id, + agent_id=agent_id, + task_id=task_id, + amount=exact_amount, + currency=currency.upper(), + requires_approval=requires_approval, + approved_by=approved_by, + metadata=metadata or {}, + created_at=normalized_created_at.isoformat(), + ) + with self._lock: + self._connection.execute( + """ + INSERT INTO ledger_events + (id, category, event_type, team_id, agent_id, task_id, amount, currency, + requires_approval, approved_by, metadata_json, created_at) + VALUES (:id, :category, :event_type, :team_id, :agent_id, :task_id, :amount, + :currency, :requires_approval, :approved_by, :metadata_json, :created_at) + """, + {**asdict(event), "requires_approval": int(event.requires_approval), + "amount": str(event.amount), + "metadata_json": json.dumps(event.metadata, sort_keys=True)}, + ) + self._connection.commit() + return event + + def approve(self, event_id: str, approver_id: str) -> LedgerEvent: + """Approve a pending event and return the updated record.""" + if not approver_id: + raise ValueError("approver_id is required") + with self._lock: + cursor = self._connection.execute( + """ + UPDATE ledger_events + SET requires_approval = 0, approved_by = ? + WHERE id = ? AND requires_approval = 1 AND approved_by IS NULL + """, + (approver_id, event_id), + ) + if cursor.rowcount != 1: + raise LookupError("pending approval not found") + self._connection.commit() + return self.get(event_id) + + def get(self, event_id: str) -> LedgerEvent: + """Load one event by ID or raise when it does not exist.""" + with self._lock: + row = self._connection.execute( + "SELECT * FROM ledger_events WHERE id = ?", (event_id,) + ).fetchone() + if row is None: + raise LookupError("ledger event not found") + return self._row_to_event(row) + + def list_events(self, team_id: str | None = None, limit: int | None = None) -> list[LedgerEvent]: + """Return newest events, optionally filtered by team and limited in count.""" + limit_sql = "" if limit is None else " LIMIT ?" + limit_params: tuple[Any, ...] = () if limit is None else (limit,) + with self._lock: + if team_id is None: + rows = self._connection.execute( + "SELECT * FROM ledger_events ORDER BY created_at DESC" + limit_sql, + limit_params, + ).fetchall() + else: + rows = self._connection.execute( + "SELECT * FROM ledger_events WHERE team_id = ? ORDER BY created_at DESC" + limit_sql, + (team_id, *limit_params), + ).fetchall() + return [self._row_to_event(row) for row in rows] + + def list_pending_approvals( + self, team_id: str | None = None, limit: int | None = None + ) -> list[LedgerEvent]: + """Return newest events that still require human approval.""" + query = "SELECT * FROM ledger_events WHERE requires_approval = 1" + params: tuple[Any, ...] = () + if team_id is not None: + query += " AND team_id = ?" + params += (team_id,) + query += " ORDER BY created_at DESC" + if limit is not None: + query += " LIMIT ?" + params += (limit,) + with self._lock: + rows = self._connection.execute(query, params).fetchall() + return [self._row_to_event(row) for row in rows] + + def summarize(self, team_id: str | None = None) -> list[TeamSummary]: + """Aggregate revenue, costs, activity, and approvals by team and currency.""" + where = "" if team_id is None else "WHERE team_id = ?" + params: tuple[Any, ...] = () if team_id is None else (team_id,) + with self._lock: + rows = self._connection.execute( + """ + SELECT team_id, currency, category, requires_approval, amount + FROM ledger_events + """ + where + """ + ORDER BY team_id, currency + """, + params, + ).fetchall() + totals: dict[tuple[str, str], dict[str, Any]] = {} + for row in rows: + key = (row["team_id"], row["currency"]) + total = totals.setdefault( + key, + { + "revenue": Decimal("0"), + "costs": Decimal("0"), + "activity_count": 0, + "pending_approvals": 0, + }, + ) + total["activity_count"] += 1 + if row["requires_approval"]: + total["pending_approvals"] += 1 + elif row["category"] == "revenue": + total["revenue"] += Decimal(row["amount"]) + elif row["category"] == "cost": + total["costs"] += Decimal(row["amount"]) + return [ + TeamSummary( + team_id=team_id, + currency=currency, + revenue=total["revenue"], + costs=total["costs"], + profit=total["revenue"] - total["costs"], + activity_count=total["activity_count"], + pending_approvals=total["pending_approvals"], + ) + for (team_id, currency), total in sorted(totals.items()) + ] + + @staticmethod + def _row_to_event(row: sqlite3.Row) -> LedgerEvent: + """Convert a SQLite row into the public immutable event model.""" + return LedgerEvent( + id=row["id"], + category=row["category"], + event_type=row["event_type"], + team_id=row["team_id"], + agent_id=row["agent_id"], + task_id=row["task_id"], + amount=Decimal(str(row["amount"])), + currency=row["currency"], + requires_approval=bool(row["requires_approval"]), + approved_by=row["approved_by"], + metadata=json.loads(row["metadata_json"]), + created_at=row["created_at"], + ) diff --git a/Server/src/overseer/simulation.py b/Server/src/overseer/simulation.py new file mode 100644 index 000000000..d4552c799 --- /dev/null +++ b/Server/src/overseer/simulation.py @@ -0,0 +1,67 @@ +"""Deterministic support-agent simulation for exercising the overseer ledger.""" + +from __future__ import annotations + +from decimal import Decimal, InvalidOperation + +from .ledger import EventLedger, LedgerEvent + + +def run_support_ticket_simulation( + ledger: EventLedger, + *, + ticket_id: str = "T-100", + subscription_amount: float = 499, + model_cost: float = 31.25, +) -> list[LedgerEvent]: + """Record one complete support-ticket lifecycle without external side effects.""" + try: + exact_subscription_amount = Decimal(str(subscription_amount)) + exact_model_cost = Decimal(str(model_cost)) + except (InvalidOperation, TypeError, ValueError): + raise ValueError( + "subscription_amount and model_cost must be finite and non-negative" + ) from None + if any( + not amount.is_finite() or amount < 0 + for amount in (exact_subscription_amount, exact_model_cost) + ): + raise ValueError("subscription_amount and model_cost must be finite and non-negative") + task_id = f"ticket:{ticket_id}" + events = [ + ledger.record( + category="activity", + event_type="ticket_received", + team_id="support", + agent_id="triage-agent", + task_id=task_id, + metadata={"channel": "simulation"}, + ), + ledger.record( + category="activity", + event_type="ticket_resolved", + team_id="support", + agent_id="support-agent", + task_id=task_id, + metadata={"confidence": 0.94, "source": "approved-product-docs"}, + ), + ledger.record( + category="cost", + event_type="model_usage", + team_id="support", + agent_id="support-agent", + task_id=task_id, + amount=exact_model_cost, + currency="USD", + ), + ledger.record( + category="revenue", + event_type="subscription_paid", + team_id="support", + agent_id="billing-agent", + task_id=task_id, + amount=exact_subscription_amount, + currency="USD", + ), + ] + return events diff --git a/Server/src/overseer/teams/__init__.py b/Server/src/overseer/teams/__init__.py new file mode 100644 index 000000000..261e2bf4b --- /dev/null +++ b/Server/src/overseer/teams/__init__.py @@ -0,0 +1,24 @@ +"""Provider-neutral definitions for the five commercial agent servers.""" + +from .content_operations import TEAM_SPEC as CONTENT_OPERATIONS +from .customer_support import TEAM_SPEC as CUSTOMER_SUPPORT +from .inventory_operations import TEAM_SPEC as INVENTORY_OPERATIONS +from .lead_generation import TEAM_SPEC as LEAD_GENERATION +from .market_intelligence import TEAM_SPEC as MARKET_INTELLIGENCE + +TEAM_SPECS = ( + LEAD_GENERATION, + CUSTOMER_SUPPORT, + CONTENT_OPERATIONS, + MARKET_INTELLIGENCE, + INVENTORY_OPERATIONS, +) + +__all__ = [ + "CONTENT_OPERATIONS", + "CUSTOMER_SUPPORT", + "INVENTORY_OPERATIONS", + "LEAD_GENERATION", + "MARKET_INTELLIGENCE", + "TEAM_SPECS", +] diff --git a/Server/src/overseer/teams/content_operations.py b/Server/src/overseer/teams/content_operations.py new file mode 100644 index 000000000..0171130ca --- /dev/null +++ b/Server/src/overseer/teams/content_operations.py @@ -0,0 +1,16 @@ +"""Automated content operations server definition.""" + +TEAM_SPEC = { + "team_id": "content-operations", + "name": "Automated Content and Repurposing", + "purpose": "Turn approved audio or video into edited clips, transcripts, and channel-ready drafts.", + "agents": [ + "media-intake-agent", + "transcription-agent", + "highlight-editor-agent", + "social-copy-agent", + "content-quality-agent", + ], + "requires_human_approval_for": ["publication", "brand_claim", "copyrighted_asset_use"], + "planned_revenue_model": "monthly content-operations package", +} diff --git a/Server/src/overseer/teams/customer_support.py b/Server/src/overseer/teams/customer_support.py new file mode 100644 index 000000000..506c5bd73 --- /dev/null +++ b/Server/src/overseer/teams/customer_support.py @@ -0,0 +1,16 @@ +"""Specialized customer support server definition.""" + +TEAM_SPEC = { + "team_id": "customer-support", + "name": "Specialized Customer Support", + "purpose": "Resolve approved tier-one support requests using customer documentation and systems.", + "agents": [ + "ticket-triage-agent", + "knowledge-retrieval-agent", + "response-drafting-agent", + "escalation-agent", + "support-analyst", + ], + "requires_human_approval_for": ["refund", "account_change", "production_action"], + "planned_revenue_model": "subscription priced by support volume", +} diff --git a/Server/src/overseer/teams/inventory_operations.py b/Server/src/overseer/teams/inventory_operations.py new file mode 100644 index 000000000..b63851f79 --- /dev/null +++ b/Server/src/overseer/teams/inventory_operations.py @@ -0,0 +1,16 @@ +"""Inventory and operations server definition.""" + +TEAM_SPEC = { + "team_id": "inventory-operations", + "name": "Inventory and Operations Management", + "purpose": "Monitor stock, identify demand signals, and prepare purchasing recommendations.", + "agents": [ + "inventory-monitor-agent", + "demand-forecast-agent", + "reorder-planning-agent", + "purchase-order-agent", + "operations-analyst", + ], + "requires_human_approval_for": ["purchase_order", "supplier_change", "stock_adjustment"], + "planned_revenue_model": "setup fee plus monthly maintenance", +} diff --git a/Server/src/overseer/teams/lead_generation.py b/Server/src/overseer/teams/lead_generation.py new file mode 100644 index 000000000..0571e435a --- /dev/null +++ b/Server/src/overseer/teams/lead_generation.py @@ -0,0 +1,16 @@ +"""Lead generation and sales qualification server definition.""" + +TEAM_SPEC = { + "team_id": "lead-generation", + "name": "Lead Generation and Sales Qualification", + "purpose": "Find suitable prospects, prepare personalized drafts, and qualify inbound leads.", + "agents": [ + "prospect-researcher", + "data-enrichment-agent", + "personalization-agent", + "qualification-agent", + "pipeline-analyst", + ], + "requires_human_approval_for": ["outbound_message", "contact_import", "pricing_offer"], + "planned_revenue_model": "monthly retainer or qualified-meeting fee", +} diff --git a/Server/src/overseer/teams/market_intelligence.py b/Server/src/overseer/teams/market_intelligence.py new file mode 100644 index 000000000..d9b3155ee --- /dev/null +++ b/Server/src/overseer/teams/market_intelligence.py @@ -0,0 +1,16 @@ +"""Financial and market intelligence server definition.""" + +TEAM_SPEC = { + "team_id": "market-intelligence", + "name": "Financial and Market Intelligence", + "purpose": "Collect public business signals and produce traceable research reports and metrics.", + "agents": [ + "source-monitor-agent", + "filings-extraction-agent", + "competitor-pricing-agent", + "report-builder-agent", + "evidence-review-agent", + ], + "requires_human_approval_for": ["client_report", "investment_claim", "restricted_source_access"], + "planned_revenue_model": "custom research packs or recurring intelligence subscription", +} diff --git a/Server/tests/test_overseer_api.py b/Server/tests/test_overseer_api.py new file mode 100644 index 000000000..ae809efeb --- /dev/null +++ b/Server/tests/test_overseer_api.py @@ -0,0 +1,127 @@ +import sqlite3 +from decimal import Decimal + +import pytest +from fastapi.testclient import TestClient + +from overseer import EventLedger, create_overseer_app, run_support_ticket_simulation + + +@pytest.mark.parametrize( + ("subscription_amount", "model_cost"), + [(float("nan"), 31.25), (499, float("inf")), (-1, 31.25), (499, -1)], +) +def test_simulation_validates_amounts_before_recording_events( + subscription_amount, model_cost +): + ledger = EventLedger(sqlite3.connect(":memory:")) + + with pytest.raises(ValueError, match="finite and non-negative"): + run_support_ticket_simulation( + ledger, + subscription_amount=subscription_amount, + model_cost=model_cost, + ) + + assert ledger.list_events() == [] + + +def test_simulation_records_full_lifecycle_for_valid_amounts(): + ledger = EventLedger(sqlite3.connect(":memory:")) + + events = run_support_ticket_simulation( + ledger, + ticket_id="T-42", + subscription_amount=Decimal("499.00"), + model_cost=Decimal("31.25"), + ) + + assert [event.event_type for event in events] == [ + "ticket_received", + "ticket_resolved", + "model_usage", + "subscription_paid", + ] + assert len(ledger.list_events()) == 4 + assert events[2].amount == Decimal("31.25") + assert events[3].amount == Decimal("499.00") + + +def test_api_exposes_simulation_events_and_financial_summary(): + ledger = EventLedger(sqlite3.connect(":memory:")) + run_support_ticket_simulation(ledger, ticket_id="T-42") + client = TestClient(create_overseer_app(ledger, api_key="test-key")) + + headers = {"X-Overseer-Api-Key": "test-key"} + events = client.get("/events?team_id=support", headers=headers).json() + summary = client.get("/summaries?team_id=support", headers=headers).json() + + assert len(events) == 4 + assert summary == [ + { + "team_id": "support", + "currency": "USD", + "revenue": 499, + "costs": 31.25, + "profit": 467.75, + "activity_count": 4, + "pending_approvals": 0, + } + ] + + +def test_api_only_returns_pending_approvals(): + ledger = EventLedger(sqlite3.connect(":memory:")) + pending = ledger.record( + category="cost", + event_type="refund_requested", + team_id="support", + agent_id="billing-agent", + amount=50, + requires_approval=True, + ) + ledger.record( + category="activity", + event_type="ticket_resolved", + team_id="support", + agent_id="support-agent", + ) + client = TestClient(create_overseer_app(ledger, api_key="test-key")) + + approvals = client.get("/approvals", headers={"X-Overseer-Api-Key": "test-key"}).json() + + assert len(approvals) == 1 + assert approvals[0]["id"] == pending.id + + +def test_api_requires_authentication_and_enforces_team_scope(): + ledger = EventLedger(sqlite3.connect(":memory:")) + client = TestClient( + create_overseer_app( + ledger, + api_key="test-key", + authorized_team_ids={"support"}, + ) + ) + + assert client.get("/events").status_code == 401 + assert ( + client.get( + "/events", + headers={"X-Overseer-Api-Key": "test-key"}, + ).status_code + == 403 + ) + assert ( + client.get( + "/events?team_id=lead-generation", + headers={"X-Overseer-Api-Key": "test-key"}, + ).status_code + == 403 + ) + + for path in ("/summaries", "/approvals"): + assert ( + client.get(path, headers={"X-Overseer-Api-Key": "test-key"}).status_code + == 403 + ) diff --git a/Server/tests/test_overseer_ledger.py b/Server/tests/test_overseer_ledger.py new file mode 100644 index 000000000..a520f9fec --- /dev/null +++ b/Server/tests/test_overseer_ledger.py @@ -0,0 +1,251 @@ +import sqlite3 +from concurrent.futures import ThreadPoolExecutor +from datetime import datetime, timezone +from decimal import Decimal + +import pytest + +from overseer import EventLedger + + +@pytest.fixture +def ledger(): + return EventLedger(sqlite3.connect(":memory:")) + + +def test_records_activity_and_summarizes_revenue_cost_and_profit(ledger): + ledger.record( + category="revenue", + event_type="subscription_paid", + team_id="support", + agent_id="billing-agent", + amount=500, + currency="usd", + ) + ledger.record( + category="cost", + event_type="model_usage", + team_id="support", + agent_id="support-agent", + amount=125, + currency="USD", + ) + ledger.record( + category="activity", + event_type="ticket_resolved", + team_id="support", + agent_id="support-agent", + metadata={"ticket_id": "T-1"}, + ) + + summary = ledger.summarize()[0] + assert summary.team_id == "support" + assert summary.revenue == 500 + assert summary.costs == 125 + assert summary.profit == 375 + assert summary.activity_count == 3 + + +def test_exact_money_is_currency_aware_and_pending_amounts_are_excluded(ledger): + ledger.record( + category="revenue", + event_type="paid", + team_id="support", + agent_id="billing", + amount="0.10", + currency="usd", + ) + ledger.record( + category="revenue", + event_type="paid", + team_id="support", + agent_id="billing", + amount="0.20", + currency="usd", + ) + pending_revenue = ledger.record( + category="revenue", + event_type="pending", + team_id="support", + agent_id="billing", + amount="100", + currency="eur", + requires_approval=True, + ) + pending_cost = ledger.record( + category="cost", + event_type="refund_requested", + team_id="support", + agent_id="billing", + amount="0.05", + currency="usd", + requires_approval=True, + ) + + summaries = ledger.summarize() + by_currency = {summary.currency: summary for summary in summaries} + assert by_currency["USD"].revenue == Decimal("0.30") + assert by_currency["USD"].costs == Decimal("0") + assert by_currency["USD"].profit == Decimal("0.30") + assert by_currency["EUR"].revenue == Decimal("0") + assert by_currency["EUR"].pending_approvals == 1 + assert by_currency["EUR"].costs == Decimal("0") + assert by_currency["EUR"].profit == Decimal("0") + ledger.approve(pending_revenue.id, "operator-1") + ledger.approve(pending_cost.id, "operator-1") + by_currency = {summary.currency: summary for summary in ledger.summarize()} + assert by_currency["EUR"].revenue == Decimal("100") + assert by_currency["EUR"].profit == Decimal("100") + assert by_currency["USD"].costs == Decimal("0.05") + assert by_currency["USD"].profit == Decimal("0.25") + + +def test_migrates_legacy_real_amount_column_before_new_writes(tmp_path): + connection = sqlite3.connect(tmp_path / "legacy.db") + connection.execute( + """ + CREATE TABLE ledger_events ( + id TEXT PRIMARY KEY, + category TEXT NOT NULL, + event_type TEXT NOT NULL, + team_id TEXT NOT NULL, + agent_id TEXT NOT NULL, + task_id TEXT, + amount REAL NOT NULL DEFAULT 0, + currency TEXT NOT NULL, + requires_approval INTEGER NOT NULL DEFAULT 0, + approved_by TEXT, + metadata_json TEXT NOT NULL, + created_at TEXT NOT NULL + ) + """ + ) + + ledger = EventLedger(connection) + amount_type = connection.execute( + "PRAGMA table_info(ledger_events)" + ).fetchall()[6][2] + ledger.record( + category="revenue", + event_type="paid", + team_id="support", + agent_id="billing", + amount="123.456789", + currency="USD", + ) + + assert amount_type == "TEXT" + assert ledger.list_events(limit=1)[0].amount == Decimal("123.456789") + + +def test_rejects_non_finite_or_naive_timestamps(ledger): + with pytest.raises(ValueError, match="finite"): + ledger.record(category="cost", event_type="x", team_id="t", agent_id="a", amount=float("nan")) + with pytest.raises(ValueError, match="finite"): + ledger.record(category="cost", event_type="x", team_id="t", agent_id="a", amount=float("inf")) + with pytest.raises(ValueError, match="finite"): + ledger.record(category="cost", event_type="x", team_id="t", agent_id="a", amount=float("-inf")) + with pytest.raises(ValueError, match="timezone-aware"): + ledger.record( + category="activity", + event_type="x", + team_id="t", + agent_id="a", + created_at=datetime(2025, 1, 1), + ) + + event = ledger.record( + category="activity", + event_type="x", + team_id="t", + agent_id="a", + created_at=datetime(2025, 1, 1, tzinfo=timezone.utc), + ) + assert event.created_at.endswith("+00:00") + + offset_event = ledger.record( + category="activity", + event_type="x", + team_id="t", + agent_id="a", + created_at=datetime.fromisoformat("2025-01-01T01:00:00+01:00"), + ) + assert offset_event.created_at == "2025-01-01T00:00:00+00:00" + + +def test_orders_events_by_utc_time_when_offsets_differ(ledger): + ledger.record( + category="activity", + event_type="older", + team_id="t", + agent_id="a", + created_at=datetime.fromisoformat("2025-01-01T00:00:00+14:00"), + ) + ledger.record( + category="activity", + event_type="newer", + team_id="t", + agent_id="a", + created_at=datetime.fromisoformat("2024-12-31T23:00:00+00:00"), + ) + + assert [event.event_type for event in ledger.list_events()] == ["newer", "older"] + + +def test_requires_approval_is_visible_and_can_be_approved(ledger): + event = ledger.record( + category="cost", + event_type="refund_requested", + team_id="support", + agent_id="support-agent", + amount=50, + requires_approval=True, + ) + + assert ledger.summarize()[0].pending_approvals == 1 + approved = ledger.approve(event.id, "operator-1") + assert approved.approved_by == "operator-1" + assert not approved.requires_approval + assert ledger.summarize()[0].pending_approvals == 0 + + +def test_rejects_unsafe_or_ambiguous_events(ledger): + with pytest.raises(ValueError, match="category"): + ledger.record( + category="unknown", + event_type="x", + team_id="support", + agent_id="agent", + ) + with pytest.raises(ValueError, match="non-negative"): + ledger.record( + category="cost", + event_type="x", + team_id="support", + agent_id="agent", + amount=-1, + ) + with pytest.raises(ValueError, match="already be approved"): + ledger.record( + category="activity", + event_type="x", + team_id="support", + agent_id="agent", + requires_approval=True, + approved_by="operator", + ) + + +def test_serializes_concurrent_worker_thread_access(ledger): + def record_event(index): + ledger.record( + category="activity", + event_type="worker_event", + team_id="support", + agent_id=f"agent-{index}", + ) + + with ThreadPoolExecutor(max_workers=8) as executor: + list(executor.map(record_event, range(40))) + + assert len(ledger.list_events(team_id="support")) == 40 diff --git a/docs/overseer-control-plane.md b/docs/overseer-control-plane.md new file mode 100644 index 000000000..ddeb212b0 --- /dev/null +++ b/docs/overseer-control-plane.md @@ -0,0 +1,96 @@ +# Overseer control plane + +The first commercial-agent capability is an event ledger, not an autonomous +outreach or payment connector. `Server/src/overseer/ledger.py` records provider- +neutral facts that a future dashboard can display: + +- agent activity and task history; +- revenue and operating costs; +- pending and completed human approvals; and +- per-team profit summaries. + +## Simulation + +```python +import sqlite3 +from overseer import EventLedger + +ledger = EventLedger(sqlite3.connect("overseer.db", check_same_thread=False)) +ledger.record( + category="activity", + event_type="ticket_resolved", + team_id="support", + agent_id="support-agent", + task_id="ticket:T-100", + metadata={"source": "approved-product-docs", "confidence": 0.94}, +) +ledger.record( + category="revenue", + event_type="subscription_paid", + team_id="support", + agent_id="billing-agent", + amount=499, + currency="USD", +) +ledger.record( + category="cost", + event_type="model_usage", + team_id="support", + agent_id="support-agent", + amount=31.25, + currency="USD", +) +print(ledger.summarize()) +``` + +## Read-only dashboard API + +The API adapter exposes the same facts without granting dashboard clients +permission to mutate them: + +```python +from overseer import create_overseer_app + +app = create_overseer_app(ledger, api_key="set-this-from-secret-storage") +``` + +It provides authenticated `GET /events`, `GET /summaries`, and `GET /approvals` +using the `X-Overseer-Api-Key` header. Pass `authorized_team_ids` when the +operator should be restricted to a subset of teams; scoped requests must +include a team ID from that set. The +`run_support_ticket_simulation()` helper can populate a local demo ledger so +the dashboard can be built and reviewed before any provider credentials exist. +When the API is used, `EventLedger` copies the SQLite connection to one that is +safe for API handlers. File- +backed databases remain persistent, while in-memory databases are copied into +an isolated database owned by the ledger, with access serialized for FastAPI +worker threads. + +Amounts are stored as exact decimals and summaries are separated by currency. +Existing ledgers with the earlier SQLite `REAL` amount column are migrated to +the exact `TEXT` representation before new events are written. +Pending approval events remain visible in the activity and approval feeds but +are excluded from realized revenue, cost, and profit totals until approved. + +The ledger intentionally does not send emails, place calls, issue refunds, or +connect to a payment provider. Those actions must be implemented as connectors +that emit ledger events and use `requires_approval=True` for irreversible work. +The dashboard should read the same event stream, so the operator can inspect +what happened before enabling production credentials. + +## Five agent servers + +The planned commercial servers are defined as provider-neutral configuration +under `Server/src/overseer/teams/`. Each definition contains five specialist +agents, the server purpose, its planned revenue model, and actions that remain +human-approval gated: + +- `lead_generation.py` +- `customer_support.py` +- `content_operations.py` +- `market_intelligence.py` +- `inventory_operations.py` + +These files do not connect to email, CRM, voice, payment, customer, or +production systems. They are the supervised server boundaries that connectors +and workflows can be added to later.