Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions api/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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")
Expand Down
53 changes: 53 additions & 0 deletions api/outcomes.py
Original file line number Diff line number Diff line change
@@ -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)
15 changes: 15 additions & 0 deletions outcomes/__init__.py
Original file line number Diff line number Diff line change
@@ -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",
]
75 changes: 75 additions & 0 deletions outcomes/ledger.py
Original file line number Diff line number Diff line change
@@ -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]
9 changes: 9 additions & 0 deletions tests/outcomes/test_ledger.py
Original file line number Diff line number Diff line change
@@ -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
Loading