diff --git a/api/main.py b/api/main.py index ad5afea..261ddfc 100644 --- a/api/main.py +++ b/api/main.py @@ -30,6 +30,7 @@ from api.gdpr import router as gdpr_router from api.memory_state import router as memory_state_router from api.reliability import router as reliability_router +from api.outcomes import router as outcomes_router from api.metrics import record_http_request, record_query, render_metrics from api.decisions import router as decisions_router from api.deps import RolesDep, memory, set_memory_service @@ -111,6 +112,7 @@ async def lifespan(app: FastAPI) -> AsyncIterator[None]: app.include_router(evidence_router) app.include_router(memory_state_router) app.include_router(reliability_router) +app.include_router(outcomes_router) @app.middleware("http") diff --git a/api/outcomes.py b/api/outcomes.py new file mode 100644 index 0000000..b5916b1 --- /dev/null +++ b/api/outcomes.py @@ -0,0 +1,53 @@ +"""POST /outcomes/record — append action ledger entry.""" + +from __future__ import annotations + +from typing import Any + +from fastapi import APIRouter +from pydantic import BaseModel, Field + +from api.deps import RolesDep +from outcomes.ledger import ActionLedgerEntry, record_outcome, update_usefulness + +router = APIRouter(prefix="/outcomes", tags=["outcomes"]) + + +class RecordOutcomeRequest(BaseModel): + workspace_id: str + agent_id: str + task: str + memory_ids: list[str] = Field(default_factory=list) + reliability_decision: str | None = None + action: dict[str, Any] = Field(default_factory=dict) + outcome: str = "unknown" + metrics: dict[str, Any] = Field(default_factory=dict) + prior_usefulness: float = Field(default=0.5, ge=0.0, le=1.0) + trace_id: str | None = None + + +class RecordOutcomeResponse(BaseModel): + id: str + usefulness_score: float + + +@router.post("/record", response_model=RecordOutcomeResponse) +def post_record_outcome(body: RecordOutcomeRequest, _roles: RolesDep) -> RecordOutcomeResponse: + entry = ActionLedgerEntry( + workspace_id=body.workspace_id, + agent_id=body.agent_id, + task=body.task, + memory_ids=body.memory_ids, + reliability_decision=body.reliability_decision, + action=body.action, + outcome=body.outcome, # type: ignore[arg-type] + metrics=body.metrics, + trace_id=body.trace_id, + ) + record_outcome(entry) + success = body.outcome == "success" + usefulness = update_usefulness( + prior_usefulness=body.prior_usefulness, + success=success, + ) + return RecordOutcomeResponse(id=entry.id, usefulness_score=usefulness) diff --git a/outcomes/__init__.py b/outcomes/__init__.py new file mode 100644 index 0000000..f2030cd --- /dev/null +++ b/outcomes/__init__.py @@ -0,0 +1,15 @@ +"""Outcome-driven memory updating package.""" + +from outcomes.ledger import ( + ActionLedgerEntry, + OutcomeRecord, + record_outcome, + update_usefulness, +) + +__all__ = [ + "ActionLedgerEntry", + "OutcomeRecord", + "record_outcome", + "update_usefulness", +] diff --git a/outcomes/ledger.py b/outcomes/ledger.py new file mode 100644 index 0000000..2189029 --- /dev/null +++ b/outcomes/ledger.py @@ -0,0 +1,75 @@ +"""Outcome ledger and usefulness updates (V2 Enhancement 5).""" + +from __future__ import annotations + +import uuid +from datetime import UTC, datetime +from typing import Any, Literal + +from pydantic import BaseModel, Field + +OutcomeResult = Literal["success", "failure", "partial", "unknown"] + + +def _uuid4() -> str: + return str(uuid.uuid4()) + + +def _utcnow() -> datetime: + return datetime.now(UTC) + + +class ActionLedgerEntry(BaseModel): + """Append-oriented audit of agent action + memory used.""" + + id: str = Field(default_factory=_uuid4) + workspace_id: str + agent_id: str + task: str + memory_ids: list[str] = Field(default_factory=list) + reliability_decision: str | None = None + action: dict[str, Any] = Field(default_factory=dict) + outcome: OutcomeResult = "unknown" + metrics: dict[str, Any] = Field(default_factory=dict) + trace_id: str | None = None + created_at: datetime = Field(default_factory=_utcnow) + + +class OutcomeRecord(BaseModel): + """Outcome linked to a decision or procedure.""" + + id: str = Field(default_factory=_uuid4) + workspace_id: str + action_id: str | None = None + procedure_id: str | None = None + decision_id: str | None = None + result: OutcomeResult + success: bool + metrics: dict[str, Any] = Field(default_factory=dict) + observed_at: datetime = Field(default_factory=_utcnow) + evidence_ids: list[str] = Field(default_factory=list) + + +def update_usefulness( + *, + prior_usefulness: float, + success: bool, + learning_rate: float = 0.1, +) -> float: + """Update usefulness score without erasing history (bounded EMA).""" + target = 1.0 if success else 0.0 + value = prior_usefulness + learning_rate * (target - prior_usefulness) + return max(0.0, min(1.0, round(value, 4))) + + +_LEDGER: list[ActionLedgerEntry] = [] + + +def record_outcome(entry: ActionLedgerEntry) -> ActionLedgerEntry: + """Append to in-memory ledger (Neo4j writer can replace later).""" + _LEDGER.append(entry) + return entry + + +def ledger_for_workspace(workspace_id: str) -> list[ActionLedgerEntry]: + return [e for e in _LEDGER if e.workspace_id == workspace_id] diff --git a/tests/outcomes/test_ledger.py b/tests/outcomes/test_ledger.py new file mode 100644 index 0000000..49bf4a3 --- /dev/null +++ b/tests/outcomes/test_ledger.py @@ -0,0 +1,9 @@ +from outcomes.ledger import update_usefulness + + +def test_usefulness_rises_on_success() -> None: + assert update_usefulness(prior_usefulness=0.5, success=True) > 0.5 + + +def test_usefulness_falls_on_failure() -> None: + assert update_usefulness(prior_usefulness=0.5, success=False) < 0.5