Files
praxis/server/session_recorder.py
Praxis CI ec397f2c65 docs(milestone): complete v0.5-live-assist — v0.1.13 tagged, milestone release, merged to main
v0.5 (Live Assist — on-the-job voice companion) milestone complete.
4 phases: P0 (pre-execution, v0.1.10) → P1 (assist core + guardrail,
v0.1.11) → P2 (integration + tech-debt + NFR, v0.1.12) → P3 (final
review + ship, v0.1.13 = milestone release).

16/16 REQs covered (3 ASSIST + 4 NFR + 9 IDEATE). 4 v0.6 backlog.
469 tests passed, 0 failed. 1 P0 fixed (guardrail processor safety).
8 P1+ flagged for v0.6. 8 v0.4 P1+ tech-debt addressed.
G-049 + G-067 grill MUSTs resolved. ESCALATION-01 (PIPEDA) OPEN for
human legal review before assist surface go-live.

---ci---
project: praxis
phase: 3
milestone: v0.5
status: complete
requirements:
  covered: [REQ-ASSIST-01, REQ-ASSIST-02, REQ-ASSIST-03, REQ-NFR-ASSIST-01, REQ-NFR-ASSIST-02, REQ-NFR-ASSIST-03, REQ-NFR-ASSIST-04, REQ-IDEATE-01, REQ-IDEATE-02, REQ-IDEATE-03, REQ-IDEATE-04, REQ-IDEATE-05, REQ-IDEATE-06, REQ-IDEATE-07, REQ-IDEATE-08, REQ-IDEATE-09]
  partial: []
---/ci---
2026-08-04 22:35:56 +00:00

399 lines
16 KiB
Python

"""Session recorder — wires the SQLite store into the pipeline lifecycle (TASK-04-03).
On session start: create a sessions row.
Per turn: log a turns row with ASR/TTS text + latency.
On branch decision: update branch_path.
On session end: set outcome + update progress + store cost + debrief.
After end(): the caller may invoke `run_mastery_flow()` to run the off-voice-path
mastery scoring pipeline (SLICE-07 TASK-07-01): evidence extraction → rubric
scoring → scenario score → IRT theta update → path gate check + week advance →
SQLite gate-event audit → optional VC issuance (SLICE-09, lazy import).
No auth — learner_id is the hardcoded 'learner-1' (D-007).
"""
from __future__ import annotations
import asyncio
import datetime as _dt
import json
import logging
import uuid
from typing import Any, Awaitable, Callable
from db.store import PraxisStore, HARDCODED_LEARNER_ID
from server.cost import CostBreakdown, derive_cost
log = logging.getLogger(__name__)
def _now_iso() -> str:
return _dt.datetime.now(_dt.timezone.utc).isoformat()
class SessionRecorder:
"""Records a voice session to SQLite (TASK-04-03)."""
def __init__(
self,
store: PraxisStore,
learner_id: str = HARDCODED_LEARNER_ID,
scenario_id: str = "cs_refund_ca_v01",
pg_store: Any = None,
session_type: str = "practice",
) -> None:
self.store = store
self.learner_id = learner_id
self.scenario_id = scenario_id
self.pg_store = pg_store
self.session_type = session_type
self.session_id: str | None = None
self._turn_seq = 0
# Cost inputs accumulated over the session.
self._llm_input_tokens = 0
self._llm_output_tokens = 0
self._deepgram_minutes = 0.0
self._tts_chars = 0
self._debrief_input_tokens = 0
self._debrief_output_tokens = 0
self._branch_path: list[str] = []
# Transcribed turns captured for the post-session mastery flow.
# Each entry: {"role": "learner"|"customer"|"assistant", "content": str}.
self._mastery_turns: list[dict[str, str]] = []
# Populated by run_mastery_flow(); surfaced to the debrief caller.
self.mastery_result: dict[str, Any] | None = None
async def start(self) -> str:
"""Create the session row; return the session id."""
self.session_id = await self.store.start_session(self.learner_id, self.scenario_id)
return self.session_id
async def log_turn(
self,
role: str,
asr_text: str | None = None,
tts_text: str | None = None,
latency_ms: float | None = None,
) -> None:
"""Log one turn to the turns table."""
if self.session_id is None:
return
await self.store.log_turn(
self.session_id, self._turn_seq, role, asr_text, tts_text, latency_ms
)
self._turn_seq += 1
# Accumulate cost inputs.
if asr_text:
# Rough: 1 token ≈ 4 chars.
self._llm_input_tokens += len(asr_text) // 4
self._mastery_turns.append({"role": role, "content": asr_text})
if tts_text:
self._tts_chars += len(tts_text)
self._llm_output_tokens += len(tts_text) // 4
if role == "assistant" and not asr_text:
self._mastery_turns.append({"role": role, "content": tts_text})
if latency_ms and role == "assistant":
# Rough audio-minutes estimate from latency (placeholder for real metering).
pass
def add_audio_minutes(self, minutes: float) -> None:
self._deepgram_minutes += minutes
def add_debrief_tokens(self, input_tokens: int, output_tokens: int) -> None:
self._debrief_input_tokens += input_tokens
self._debrief_output_tokens += output_tokens
def set_branch_path(self, branch_path: list[str]) -> None:
self._branch_path = branch_path
def set_mastery_turns(self, turns: list[dict[str, str]]) -> None:
"""Override the captured transcript turns used by run_mastery_flow()."""
self._mastery_turns = list(turns)
async def end(
self,
outcome: str,
tts_provider: str = "cartesia",
debrief_text: str | None = None,
schedule_mastery: bool = False,
mastery_deps: "MasteryFlowDeps | None" = None,
) -> CostBreakdown:
"""End the session: derive cost, write the session row, update progress.
If `schedule_mastery=True` and `mastery_deps` is provided, the mastery
flow is scheduled as a fire-and-forget asyncio task (off the voice
path). The task result lands in `self.mastery_result` once it completes.
"""
if self.session_id is None:
raise RuntimeError("SessionRecorder.end() called before start()")
breakdown = derive_cost(
llm_input_tokens=self._llm_input_tokens,
llm_output_tokens=self._llm_output_tokens,
deepgram_audio_minutes=self._deepgram_minutes,
tts_characters=self._tts_chars,
debrief_input_tokens=self._debrief_input_tokens,
debrief_output_tokens=self._debrief_output_tokens,
tts_provider=tts_provider,
)
await self.store.end_session(
self.session_id,
branch_path=self._branch_path,
outcome=outcome,
cost_cents=breakdown.derived_cents,
cost_breakdown=breakdown.as_dict(),
debrief_text=debrief_text,
)
await self.store.update_progress(self.learner_id, self.scenario_id, outcome)
if schedule_mastery and mastery_deps is not None:
asyncio.create_task(
self._run_mastery_flow_guarded(mastery_deps)
)
# v0.4 P2 (D-054): fire-and-forget cohort aggregation hook. Runs in
# parallel with the mastery flow — aggregation only needs the session
# outcome (available after session end), not the mastery scoring
# result. Rubric-dependent metrics are reconciled by the nightly job.
# Off the voice path (C-8, D-054). No-op if pg_store is None.
if self.pg_store is not None:
session_outcome = self._build_session_outcome(outcome)
asyncio.create_task(self._run_cohort_aggregation(session_outcome))
return breakdown
def _build_session_outcome(self, outcome: str) -> dict[str, Any]:
"""Construct the session_outcome dict for the aggregation hook.
v0.5 (D-062): includes session_type ('practice' | 'assist') so the
aggregator can branch. Assist shifts set session_type='assist' via
AssistSession (which reuses this pattern); practice sessions default
to 'practice'.
"""
rubric_scores: list[dict[str, Any]] = []
if self.mastery_result and isinstance(self.mastery_result, dict):
rubric_scores = list(self.mastery_result.get("rubric_scores") or [])
return {
"learner_ref": self.learner_id,
"path": self._path_slug(),
"scenario_id": self.scenario_id,
"outcome": outcome,
"session_type": self.session_type,
"rubric_scores": rubric_scores,
"failure_mode": self._failure_mode(),
"branch_path": list(self._branch_path),
"timestamp": _now_iso(),
}
def _path_slug(self) -> str:
# The scenario_id encodes the path loosely; default to customer_service.
if self.scenario_id and self.scenario_id.startswith("cs_"):
return "customer_service"
return "default"
def _failure_mode(self) -> str | None:
if self.mastery_result and isinstance(self.mastery_result, dict):
return self.mastery_result.get("failure_mode")
return None
async def _run_cohort_aggregation(self, session_outcome: dict[str, Any]) -> None:
"""Fire-and-forget wrapper around the cohort aggregation hook (D-054)."""
try:
from server.cohort.hook import on_session_end
await on_session_end(self.pg_store, session_outcome)
except Exception:
log.exception("cohort aggregation dispatch failed for session %s", self.session_id)
async def _run_mastery_flow_guarded(self, deps: "MasteryFlowDeps") -> None:
try:
await self.run_mastery_flow(deps)
except Exception:
log.exception("mastery flow failed for session %s", self.session_id)
async def run_mastery_flow(self, deps: "MasteryFlowDeps") -> dict[str, Any]:
"""Run the off-voice-path mastery scoring pipeline (SLICE-07 TASK-07-01).
Steps:
1. evidence_extractor.extract_evidence(turns, rubric_criteria, llm)
2. if ExtractionResult.scoring_inconclusive → return inconclusive
status (no score, no gate event, no progress change). The caller
surfaces a retry in the debrief (grill Axis 4 MUST #3).
3. rubric_scorer.score(evidence, rubric)
4. mastery_score.compute_scenario_score(criterion_scores, rubric)
5. irt.update_theta + persist via store.upsert_ability
6. path_engine.check_gate + advance_week + persist via store.upsert_progress
7. record mastery_gate_event in SQLite (audit, REQ-NFR-MAST-02)
8. if week-final gate open → vc_issuer.issue_credential (lazy import;
SLICE-09 may not be present yet → ImportError is swallowed)
Returns a dict describing the result (status, scenario_score, theta,
week, gate_open, ...). Stored on `self.mastery_result`.
"""
from server.mastery import evidence_extractor as _ev
from server.mastery import mastery_score as _ms
from server.mastery import rubric_scorer as _rs
rubric = deps.load_rubric()
scenario = deps.load_scenario()
criterion_ids = [m.criterion_id for m in scenario.rubric_criteria] or rubric.criterion_ids()
path_slug = scenario.path
extraction = await _ev.extract_evidence(
self._mastery_turns, criterion_ids, deps.llm
)
if extraction.scoring_inconclusive:
self.mastery_result = {
"status": "scoring_inconclusive",
"attempts": extraction.attempts,
"rejected_quotes": extraction.rejected_quotes,
"retry_advised": True,
}
return self.mastery_result
criterion_scores = _rs.score(extraction.evidence, rubric)
scenario_score = _ms.compute_scenario_score(criterion_scores, rubric)
progress_row = await self.store.get_progress(self.learner_id, path_slug)
if progress_row is not None:
progress = dict(progress_row)
scenarios_passed: list[str] = list(
json.loads(progress.get("scenarios_passed_json") or "[]")
)
else:
progress = {}
scenarios_passed = []
if scenario_score.passed and self.scenario_id not in scenarios_passed:
scenarios_passed.append(self.scenario_id)
# Recompute the path score over the passing set we know about.
path_score = _ms.compute_path_score(
[scenario_score] if scenario_score.passed else []
)
# If prior passing scenario scores are tracked elsewhere, they'd be
# folded in here; the mastery_progress row stores the cumulative mean.
path = deps.load_path()
week = deps.path_engine.current_week(progress) if progress else 1
gate_open = deps.path_engine.check_gate(
{"distinct_passed": len(scenarios_passed), "mastery_score": path_score},
week,
path,
)
# IRT theta update (uses scenario difficulty as the item parameter b).
ability_row = await self.store.get_ability(self.learner_id, path_slug)
if ability_row is not None:
theta = float(ability_row["theta"])
sigma_sq = float(ability_row["sigma_sq"])
observations = int(ability_row["observations"])
else:
theta = 0.0
sigma_sq = 1.0
observations = 0
outcome = 1.0 if scenario_score.passed else 0.0
b = float(scenario.difficulty)
new_theta, new_sigma_sq = deps.irt.update_theta(theta, sigma_sq, outcome, b)
new_observations = observations + 1
await self.store.upsert_ability(
self.learner_id, path_slug, new_theta, new_sigma_sq, new_observations
)
# Advance the week only if the gate is open (D-048).
new_progress = progress
if gate_open:
new_progress = deps.path_engine.advance_week(progress or {"current_week": week})
new_progress["distinct_passed"] = len(scenarios_passed)
new_progress["mastery_score"] = path_score
else:
new_progress = dict(progress or {"current_week": week})
new_progress["distinct_passed"] = len(scenarios_passed)
new_progress["mastery_score"] = path_score
new_week = int(new_progress.get("current_week", week))
await self.store.upsert_progress(
self.learner_id,
path_slug,
new_week,
scenarios_passed,
path_score,
gate_open,
)
# Audit log (REQ-NFR-MAST-02). scoring_inconclusive never reaches here.
rubric_scores_json = [cs.model_dump() for cs in criterion_scores]
await self.store.record_gate_event(
self.learner_id,
path_slug,
week,
scenarios_passed,
rubric_scores_json,
path_score,
gate_open,
)
# VC issuance — week-final gate open (grill Axis 8 MUST). SLICE-09 may
# not exist yet; the lazy import is wrapped so P1 ships independently.
vc_credential_id: str | None = None
path_complete = gate_open and new_week >= 6
if path_complete:
try:
from server.vc.issuer import issue_credential as _issue_credential # type: ignore
vc_credential_id = await _issue_credential(
store=self.store,
learner_id=self.learner_id,
path=path_slug,
scenarios_passed=scenarios_passed,
rubric_score=path_score,
completed_weeks=new_week,
evidence=rubric_scores_json,
)
except ImportError:
log.info("vc_issuer not available (SLICE-09 pending); skipping issuance")
except Exception:
log.exception("vc issuance failed for learner %s", self.learner_id)
self.mastery_result = {
"status": "scored",
"scenario_id": self.scenario_id,
"weighted_mean": scenario_score.weighted_mean,
"passed": scenario_score.passed,
"fail_reason": scenario_score.fail_reason,
"theta": new_theta,
"sigma_sq": new_sigma_sq,
"observations": new_observations,
"week": week,
"new_week": new_week,
"gate_open": gate_open,
"path_complete": path_complete,
"vc_credential_id": vc_credential_id,
"attempts": extraction.attempts,
}
return self.mastery_result
class MasteryFlowDeps:
"""Dependency bundle for SessionRecorder.run_mastery_flow().
Injected by the caller (DI): keeps session_recorder.py decoupled from the
concrete rubric/scenario/path loaders and the LLM provider.
"""
def __init__(
self,
llm: Any,
irt: Any,
path_engine: Any,
load_rubric: Callable[[], Any],
load_scenario: Callable[[], Any],
load_path: Callable[[], Any],
) -> None:
self.llm = llm
self.irt = irt
self.path_engine = path_engine
self.load_rubric = load_rubric
self.load_scenario = load_scenario
self.load_path = load_path
__all__ = ["SessionRecorder", "MasteryFlowDeps"]