merge(P05): phase/05 voice defense → milestone/v0.3-credential-engines

---ci---
phase: 5
milestone: v0.3
status: ship
---/ci---
This commit is contained in:
CIAgent
2026-09-12 04:51:52 +00:00
27 changed files with 2442 additions and 13 deletions
+3 -3
View File
@@ -1,8 +1,8 @@
{
"phase": 4,
"stage": "complete",
"phase": 5,
"stage": "verify",
"milestone": "v0.3",
"phase_role": "execution",
"attempts": 0,
"updated_at": "2026-09-12T03:52:16Z"
"updated_at": "2026-09-12T04:51:45Z"
}
+2 -2
View File
@@ -16,7 +16,7 @@
|----|-------------|----------|-------|--------|
| REQ-3-004 | Process-trace grading engine: grade artifacts from their full process traces; rubric-aligned structured scores; feeds Assessor real inputs | critical | 3 | complete |
| REQ-3-005 | Variant task generation: per-learner task variants (no two learners get identical prompts); variant seed registry; difficulty normalization | high | 4 | complete |
| REQ-3-006 | Oral/voice defense: AI examiner conducts spoken defense (STT → dialogue → TTS); transcript + integrity signals captured; feeds Proctor/Mentor | high | 5 | pending |
| REQ-3-006 | Oral/voice defense: AI examiner conducts spoken defense (STT → dialogue → TTS); transcript + integrity signals captured; feeds Proctor/Mentor | high | 5 | complete |
### Agent Re-grounding & Integration
@@ -179,7 +179,7 @@
| REQ-3-003 | 2 | complete |
| REQ-3-004 | 3 | complete |
| REQ-3-005 | 4 | complete |
| REQ-3-006 | 5 | pending |
| REQ-3-006 | 5 | complete |
| REQ-3-007 | 6 | pending |
| REQ-3-008 | 6 | pending |
+1 -1
View File
@@ -23,7 +23,7 @@
| 2 | Live build telemetry | complete | 1 | REQ-3-003 | In-environment capture of process events (commands, file diffs, run/test results, keystroke-level activity) streamed to ai-service; reliable transport; per-learner trace persistence |
| 3 | Process-trace grading engine | complete | 2 | REQ-3-004 | Grades artifacts from their full process traces (not just final output); emits structured rubric-aligned scores; feeds Assessor real inputs |
| 4 | Variant task generation | complete | 1 | REQ-3-005 | Per-learner task variants generated so no two learners receive identical prompts; variant seed recorded for grading fairness |
| 5 | Oral / voice defense | pending | 3 | REQ-3-006 | AI examiner conducts spoken defense of submitted work; STT → dialogue → TTS; transcript + integrity signals captured; feeds Proctor/Mentor |
| 5 | Oral / voice defense | complete | 3 | REQ-3-006 | AI examiner conducts spoken defense of submitted work; STT → dialogue → TTS; transcript + integrity signals captured; feeds Proctor/Mentor |
| 6 | Agent re-grounding + learner surface integration | pending | 2,3,4,5 | REQ-3-007, REQ-3-008 | Lab/Assessor/Proctor consume real engine inputs; v0.1 sandbox + assessment mockups wired to real engines (in-browser build/run, live telemetry, live defense) |
| 7 | Final review + ship | pending | 6 | — | Code review clean; audit passes; milestone tagged (v0.2.x final patch); release created on Gitea |
+6 -1
View File
@@ -24,4 +24,9 @@ AI_SANDBOX_MAX_PER_LEARNER=1
AI_SANDBOX_CREATES_PER_MIN=10
# Persistence (SQLite)
AI_DB_PATH=ai_service/data/nextcraft.db
AI_DB_PATH=ai_service/data/nextcraft.db
# --- v0.3 Voice (REQ-3-006, D-030) ---
# 'mock' (default; no key needed — tests/dev) or 'browser' (client-native SR/TTS).
# Real server STT/TTS ('openai-audio' + AI_VOICE_BASE_URL/AI_VOICE_API_KEY)
# is deferred to v0.4 per GRILL CUT-1/G-7 — keys never in code or commits.
AI_VOICE_PROVIDER=mock
+19
View File
@@ -216,6 +216,25 @@ Delivery is **at-least-once**; storage is **exactly-once** — the two compose:
`INCOMPLETE_FLOODED` — a terminal integrity flag the grader refuses to
grade. Silent event dropping is forbidden: it would corrupt grading input.
## Voice defense (v0.3, REQ-3-006)
Voice is **mock-first** (D-030): the defense pipeline is fully proven over
the deterministic `MockVoiceProvider` + browser-native fallback — no task
requires a real voice key. Real server STT/TTS (`OpenAIAudioProvider` over
OpenAI-compatible `/audio/transcriptions` + `/audio/speech`) is **deferred
to v0.4** together with KYC (GRILL CUT-1 / G-7): it could never be exercised
in CI, so v0.3 ships the protocol seam instead of an unverifiable claim.
- `AI_VOICE_PROVIDER=mock` (default) — deterministic canned STT/TTS
- `AI_VOICE_PROVIDER=browser` — the web client uses SpeechRecognition +
speechSynthesis; the server keeps text-turn persistence
- Conversational budget: a defense turn should complete in **< 4s**
(`DEFENSE_TURN_BUDGET_MS` in `tests/voice/test_latency.py`). v0.3
asserts instrumentation (stt_ms/llm_ms/tts_ms populated per turn); the
wall-clock acceptance probe against a real voice endpoint is a v0.4
criterion, run manually with `AI_VOICE_PROVIDER` set to the real
provider and keys in `.ciagent/.env.secrets` (never in code/commits).
## Layout
```
@@ -0,0 +1,104 @@
"""ExaminerAgent — the seventh agent: oral-defense examiner (REQ-3-006, A-109).
BOUNDARY DECISION (PERSONAS conflict rule, honored by construction): the
examiner is a TEXT agent. It composes the LLM provider through BaseAgent and
consumes defense transcript turns; it NEVER imports voice/ — STT/TTS belong
to the API endpoints (they move audio bytes; the agent moves question text).
Integrity signals (long pauses, off-scope cadence) are computed by the
endpoint layer from turn metadata (latency_ms etc.), not by the agent.
Digest discipline (D-028 mirror): questions are grounded in the compact
TraceDigest + variant statement — never the raw trace, never learner ids.
"""
from __future__ import annotations
from pydantic import BaseModel, ConfigDict, Field
from ..grading.features import TraceDigest
from ..llm.types import Message
from ..prompts.examiner import SYSTEM_PROMPT, VERDICT_SCHEMA_HINT, render_digest_context
from .base import BaseAgent
class DefenseVerdict(BaseModel):
"""D-20-validated final defense verdict (structured mode)."""
model_config = ConfigDict(extra="forbid")
verdict: str = Field(pattern="^(mastered|developing|not_yet)$")
understanding: str = Field(min_length=1)
process_justification: str = Field(min_length=1)
communication: str = Field(min_length=1)
strengths: list[str] = Field(min_length=1, max_length=2)
gaps: list[str] = Field(min_length=1, max_length=2)
class ExaminerAgent(BaseAgent):
"""Conducts the oral defense: next_question + final_verdict."""
name = "examiner"
def system_prompt(self, learner_context=None) -> str: # noqa: ANN001
"""Examiner is context-free (digest-anonymous, D-028 mirror)."""
return SYSTEM_PROMPT
def build_defense_messages(
self,
trace_digest: TraceDigest | None = None,
variant_statement: str | None = None,
history: list[Message] | None = None,
) -> list[Message]:
"""System + grounding + defense transcript (no learner id — D-028)."""
digest_json = (
trace_digest.model_dump_json() if trace_digest is not None else "{}"
)
messages: list[Message] = [
Message(role="system", content=SYSTEM_PROMPT),
Message(role="user", content=render_digest_context(digest_json, variant_statement)),
Message(
role="assistant",
content="Understood. I will question the learner about this build session.",
),
]
for m in history or []:
messages.append(m)
return messages
async def next_question(
self,
history: list[Message],
trace_digest: TraceDigest | None = None,
variant_statement: str | None = None,
) -> str:
"""One examiner question (streamed over SSE by the endpoints)."""
messages = self.build_defense_messages(trace_digest, variant_statement, history)
messages.append(
Message(role="user", content="Ask the learner your next question now.")
)
reply = await self.provider.chat(messages, model=self.settings.model)
return reply
async def final_verdict(
self,
history: list[Message],
trace_digest: TraceDigest | None = None,
variant_statement: str | None = None,
) -> DefenseVerdict:
"""Structured verdict via the D-020 4-layer defense."""
from .structured import structured_completion # module-direct (G-4)
messages = self.build_defense_messages(trace_digest, variant_statement, history)
messages.append(
Message(
role="user",
content="The defense is finished. Return the final verdict JSON now.",
)
)
return await structured_completion(
self.provider,
messages,
model=self.settings.model,
schema=DefenseVerdict,
schema_hint=VERDICT_SCHEMA_HINT,
)
@@ -14,13 +14,14 @@ AgentFactory = Callable[[LLMProvider, Settings], BaseAgent]
def register_builtin_agents(registry: "AgentRegistry") -> None:
"""Central registration of all six shipped tutor agents (G-4: one pattern).
"""Central registration of all seven shipped agents (G-4: one pattern).
coach, tutor, lab, assessor, proctor, mentor. New agents register here
in their landing phase.
coach, tutor, lab, assessor, proctor, mentor, examiner (Phase 5).
New agents register here in their landing phase.
"""
from .assessor import AssessorAgent
from .coach import CoachAgent
from .examiner import ExaminerAgent
from .lab import LabAgent
from .mentor import MentorAgent
from .proctor import ProctorAgent
@@ -38,6 +39,9 @@ def register_builtin_agents(registry: "AgentRegistry") -> None:
registry.register(
"mentor", lambda provider, settings: MentorAgent(provider, settings)
)
registry.register(
"examiner", lambda provider, settings: ExaminerAgent(provider, settings)
)
class UnknownAgentError(KeyError):
@@ -5,6 +5,7 @@ Boundary rule: api/ composes agents/ and llm/; they never import api/.
from .assessment import router as assessment_router
from .chat import router as chat_router
from .defense import router as defense_router
from .lab import router as lab_router
from .mentor import router as mentor_router
from .proctor import router as proctor_router
@@ -21,4 +22,5 @@ __all__ = [
"sandboxes_router",
"telemetry_router",
"variants_router",
"defense_router",
]
+336
View File
@@ -0,0 +1,336 @@
"""Oral-defense endpoints (Task 5-3-01, REQ-3-006, A-109).
Full defense loop over HTTP with mock-first voice (D-030) and the seventh
Examiner agent (SSE question streaming happens through the chat pipeline;
these endpoints are the session orchestration + transcript persistence):
POST /v1/defense/start {learner_id, task_id}
POST /v1/defense/{id}/answer {text} | multipart audio (STT)
GET /v1/defense/{id}/audio/{turn_id} TTS bytes (streaming)
POST /v1/defense/{id}/finish verdict + integrity signals
GET /v1/defense/{id} transcript + signals
Integrity signals (A-109) are computed server-side from turn metadata:
long pauses = learner turns whose latency_ms exceeds PAUSE_THRESHOLD_MS.
The defense does NOT gate on trace completeness (the grader does, G-4);
an incomplete trace is surfaced as `trace_complete: false` so the UI can
disclose it before the learner defends.
"""
from __future__ import annotations
import time
from datetime import UTC, datetime
from typing import TYPE_CHECKING
from fastapi import APIRouter, Depends, File, Form, HTTPException, UploadFile
from fastapi.responses import StreamingResponse
from pydantic import BaseModel, Field
from ..agents.examiner import ExaminerAgent
from ..grading.features import TraceDigest, compute_digest
from ..llm.types import Message
from ..voice.base import VoiceDescriptor
from ..voice.browser import BROWSER_FALLBACK_DESCRIPTOR
from ..voice.defense_store import DefenseRecord, DefenseStore, DefenseTurn
from .deps import (
get_examiner,
get_settings,
get_trace_store,
get_variant_store,
get_voice_provider,
get_voice_store,
)
if TYPE_CHECKING: # pragma: no cover
pass
router = APIRouter(prefix="/v1/defense", tags=["defense"])
#: A-109: learner turns slower than this are flagged as long pauses (ms).
PAUSE_THRESHOLD_MS = 15_000
_ROLE_EXAMINER = "examiner"
_ROLE_LEARNER = "learner"
class StartRequest(BaseModel):
learner_id: str = Field(min_length=1)
task_id: str = Field(min_length=1)
class StartResponse(BaseModel):
defense_id: str
voice_descriptor: dict
trace_complete: bool
first_question: str
class AnswerResponse(BaseModel):
question: str
turn_latency: dict[str, int | None]
class FinishResponse(BaseModel):
verdict: dict
integrity_signals: dict
async def _digest_for_task(
trace_store, learner_id: str, task_id: str
) -> tuple[TraceDigest | None, bool]:
"""Digest of the learner's trace for this task + completeness flag."""
if not trace_store.list_tasks(learner_id) or task_id not in trace_store.list_tasks(
learner_id
):
return None, True # no trace at all is "complete" for defense purposes
trace = trace_store.get_trace(learner_id, task_id)
gaps = trace_store.gaps(learner_id, task_id)
return (compute_digest(trace) if trace else None), (len(gaps) == 0)
def _voice_descriptor(settings) -> VoiceDescriptor:
"""The capability descriptor for the configured voice mode (D-030).
Must-Have #6: browser mode returns BROWSER_FALLBACK_DESCRIPTOR so the
web client selects native SpeechRecognition/speechSynthesis; mock mode
returns the mock descriptor. (A v0.4 server provider would return
mode="server" — the protocol seam.)
"""
if (settings.voice_provider or "mock").strip().lower() == "browser":
return BROWSER_FALLBACK_DESCRIPTOR
return VoiceDescriptor(
mode="mock", sr_available=True, tts_available=True, hint=""
)
@router.post("/start", response_model=StartResponse)
async def start_defense(
body: StartRequest,
examiner: ExaminerAgent = Depends(get_examiner),
voice_store: DefenseStore = Depends(get_voice_store),
voice_provider=Depends(get_voice_provider),
trace_store=Depends(get_trace_store),
variant_store=Depends(get_variant_store),
settings=Depends(get_settings),
) -> StartResponse:
record = voice_store.start(
DefenseRecord(
id=f"dfn-{int(time.time() * 1000):x}-{body.learner_id[:8]}",
learner_id=body.learner_id,
task_id=body.task_id,
status="in_progress",
created_at=datetime.now(UTC),
)
)
digest, trace_complete = await _digest_for_task(trace_store, body.learner_id, body.task_id)
variant = variant_store.get_by_task(body.task_id)
statement = variant.statement if variant is not None else None
started = time.perf_counter()
question = await examiner.next_question(
history=[], trace_digest=digest, variant_statement=statement
)
llm_ms = int((time.perf_counter() - started) * 1000)
voice_store.append_turn(
record.id,
DefenseTurn(
defense_id=record.id,
seq=0,
role=_ROLE_EXAMINER,
text=question,
ts=datetime.now(UTC),
latency_ms=llm_ms,
created_at=datetime.now(UTC),
),
)
descriptor = getattr(voice_provider, "descriptor", None) or _voice_descriptor(
settings
)
return StartResponse(
defense_id=record.id,
voice_descriptor=descriptor.model_dump(),
trace_complete=trace_complete,
first_question=question,
)
@router.post("/{defense_id}/answer", response_model=AnswerResponse)
async def answer_defense(
defense_id: str,
text: str | None = Form(default=None),
audio: UploadFile | None = File(default=None),
voice_store: DefenseStore = Depends(get_voice_store),
voice_provider=Depends(get_voice_provider),
examiner: ExaminerAgent = Depends(get_examiner),
trace_store=Depends(get_trace_store),
variant_store=Depends(get_variant_store),
) -> AnswerResponse:
record = voice_store.get(defense_id)
if record is None:
raise HTTPException(status_code=404, detail=f"no defense {defense_id!r}")
if record.status == "finished":
# The store owns the finished transition but does NOT police turn
# sequencing (defense_store.py: "turns after finalize are a sequencing
# bug for the endpoints to prevent") — this is the endpoint half of
# that contract: a sealed transcript is append-only-no-more.
raise HTTPException(
status_code=409, detail="defense is finished; start a new defense"
)
if text is None and audio is None:
raise HTTPException(status_code=422, detail="provide {text} or audio")
# STT (typed fallback bypasses the voice provider entirely).
stt_ms: int | None = None
if audio is not None:
stt_started = time.perf_counter()
raw = await audio.read()
if not raw:
# Empty upload is a client error (422), not a provider crash
# (500): validate before the provider call so every provider —
# mock today, the v0.4 real one — sees the same contract.
raise HTTPException(status_code=422, detail="audio upload is empty")
fmt = (audio.content_type or "audio/wav").split("/")[-1]
segment = await voice_provider.transcribe(raw, fmt)
stt_ms = int((time.perf_counter() - stt_started) * 1000)
text = segment.text
turns = record.turns if hasattr(record, "turns") else []
history = [
Message(role="assistant" if t.role == _ROLE_EXAMINER else "user", content=t.text)
for t in turns
]
next_seq = len(turns)
voice_store.append_turn(
defense_id,
DefenseTurn(
defense_id=defense_id,
seq=next_seq,
role=_ROLE_LEARNER,
text=text or "",
ts=datetime.now(UTC),
latency_ms=stt_ms,
created_at=datetime.now(UTC),
),
)
digest, _ = await _digest_for_task(trace_store, record.learner_id, record.task_id)
variant = variant_store.get_by_task(record.task_id)
llm_started = time.perf_counter()
question = await examiner.next_question(
history=history + [Message(role="user", content=text or "")],
trace_digest=digest,
variant_statement=variant.statement if variant is not None else None,
)
llm_ms = int((time.perf_counter() - llm_started) * 1000)
voice_store.append_turn(
defense_id,
DefenseTurn(
defense_id=defense_id,
seq=next_seq + 1,
role=_ROLE_EXAMINER,
text=question,
ts=datetime.now(UTC),
latency_ms=llm_ms,
created_at=datetime.now(UTC),
),
)
return AnswerResponse(
question=question,
turn_latency={"stt_ms": stt_ms, "llm_ms": llm_ms, "tts_ms": None},
)
@router.get("/{defense_id}/audio/{turn_id}")
async def defense_audio(
defense_id: str,
turn_id: int,
voice_store: DefenseStore = Depends(get_voice_store),
voice_provider=Depends(get_voice_provider),
):
record = voice_store.get(defense_id)
if record is None:
raise HTTPException(status_code=404, detail=f"no defense {defense_id!r}")
turn = next((t for t in record.turns if t.seq == turn_id), None)
if turn is None or turn.role != _ROLE_EXAMINER:
raise HTTPException(status_code=404, detail=f"no examiner turn {turn_id!r}")
async def stream():
async for chunk in voice_provider.synthesize(turn.text):
yield chunk
return StreamingResponse(stream(), media_type="audio/wav")
@router.post("/{defense_id}/finish", response_model=FinishResponse)
async def finish_defense(
defense_id: str,
voice_store: DefenseStore = Depends(get_voice_store),
examiner: ExaminerAgent = Depends(get_examiner),
trace_store=Depends(get_trace_store),
variant_store=Depends(get_variant_store),
) -> FinishResponse:
record = voice_store.get(defense_id)
if record is None:
raise HTTPException(status_code=404, detail=f"no defense {defense_id!r}")
turns = record.turns if hasattr(record, "turns") else []
history = [
Message(role="assistant" if t.role == _ROLE_EXAMINER else "user", content=t.text)
for t in turns
]
digest, _ = await _digest_for_task(trace_store, record.learner_id, record.task_id)
variant = variant_store.get_by_task(record.task_id)
verdict = await examiner.final_verdict(
history=history,
trace_digest=digest,
variant_statement=variant.statement if variant is not None else None,
)
signals: dict = {
"long_pauses": [
{"turn": t.seq, "latency_ms": t.latency_ms}
for t in turns
if t.role == _ROLE_LEARNER and (t.latency_ms or 0) > PAUSE_THRESHOLD_MS
],
"pause_threshold_ms": PAUSE_THRESHOLD_MS,
# Must-Have #1: "verdict + transcript persisted" — the verdict is
# stored INSIDE integrity_signals so GET /{id} after finish can
# re-serve it (the finish response alone would lose it). Signals
# are a JSON object dict (DefenseStore.finalize contract), so the
# verdict nests under the "verdict" key alongside the A-109
# markers the Proctor/Mentor feeds read.
"verdict": verdict.model_dump(),
}
voice_store.finalize(defense_id, signals)
return FinishResponse(verdict=verdict.model_dump(), integrity_signals=signals)
@router.get("/{defense_id}")
async def get_defense(
defense_id: str,
voice_store: DefenseStore = Depends(get_voice_store),
):
record = voice_store.get(defense_id)
if record is None:
raise HTTPException(status_code=404, detail=f"no defense {defense_id!r}")
return {
"defense_id": record.id,
"learner_id": record.learner_id,
"task_id": record.task_id,
"status": record.status,
"turns": [
{
"seq": t.seq,
"role": t.role,
"text": t.text,
"ts": t.ts,
"latency_ms": t.latency_ms,
}
for t in record.turns
],
"integrity_signals": record.integrity_signals or {},
}
+14
View File
@@ -2,6 +2,7 @@
from fastapi import Request
from ..agents.examiner import ExaminerAgent
from ..agents.registry import AgentRegistry
from ..agents.session import SessionStore
from ..config import Settings
@@ -14,6 +15,8 @@ from ..telemetry.ingest import TraceIntegrityMap
from ..telemetry.store import TraceStore
from ..variants.generator import VariantGenerator
from ..variants.store import VariantStore
from ..voice.base import VoiceProvider
from ..voice.defense_store import DefenseStore
def get_settings(request: Request) -> Settings:
@@ -63,3 +66,14 @@ def get_variant_generator(request: Request) -> VariantGenerator:
def get_variant_store(request: Request) -> VariantStore:
return request.app.state.variant_store
def get_voice_store(request: Request) -> DefenseStore:
return request.app.state.defense_store
def get_voice_provider(request: Request) -> VoiceProvider:
return request.app.state.voice_provider
def get_examiner(request: Request) -> ExaminerAgent:
return request.app.state.examiner_agent
+8
View File
@@ -77,3 +77,11 @@ class Settings(BaseSettings):
# shares the host network and reaches the app over loopback). Port reuses
# `port` (A-004); only the host is configurable — never a second port.
telemetry_ingest_host: str = "127.0.0.1"
# Voice provider selection (REQ-3-006, D-030): 'mock' (default — the
# no-key path is first-class; tests never call a real voice API) or
# 'browser' (browser-native SpeechRecognition/speechSynthesis fallback;
# the descriptor tells the web client). The real server STT/TTS
# ('openai-audio') is a v0.4 seam (GRILL CUT-1 / G-7) — AI_VOICE_BASE_URL
# and AI_VOICE_API_KEY are documented in .env.example for that future.
voice_provider: str = "mock"
+20
View File
@@ -14,6 +14,7 @@ from .agents.session import InMemorySessionStore
from .api import (
assessment_router,
chat_router,
defense_router,
lab_router,
mentor_router,
proctor_router,
@@ -30,6 +31,8 @@ from .telemetry.ingest import TraceIntegrityMap
from .telemetry.store import SQLiteTraceStore
from .variants.generator import VariantGenerator
from .variants.store import SQLiteVariantStore
from .voice.defense_store import SQLiteDefenseStore
from .voice.factory import voice_provider_from_settings
logger = logging.getLogger(__name__)
@@ -100,6 +103,21 @@ def create_app(settings: Settings | None = None) -> FastAPI:
model=settings.model,
)
# Oral defense (REQ-3-006): DefenseStore (same SQLite file) + the
# mock-first voice provider (D-030) + the seventh Examiner agent.
# Tests may pre-set app.state.defense_store / voice_provider /
# examiner_agent (state-injection override; never rebuilt if pre-set).
defense_store = getattr(app.state, "defense_store", None)
if defense_store is None:
defense_store = SQLiteDefenseStore(db_path=settings.db_path)
app.state.defense_store = defense_store
if getattr(app.state, "voice_provider", None) is None:
app.state.voice_provider = voice_provider_from_settings(settings)
if getattr(app.state, "examiner_agent", None) is None:
from .agents.examiner import ExaminerAgent
app.state.examiner_agent = ExaminerAgent(app.state.provider, settings)
# Grading persistence + engine (REQ-3-004): GradeStore from the same
# SQLite file as traces (D-027), one GradingEngine singleton wired
# through app.state — the engine receives its stores via constructor
@@ -144,6 +162,7 @@ def create_app(settings: Settings | None = None) -> FastAPI:
trace_store.close()
grade_store.close()
variant_store.close()
defense_store.close()
await app.state.http_client.aclose()
app = FastAPI(title="Nextcraft AI Service", version="0.3.0", lifespan=lifespan)
@@ -173,6 +192,7 @@ def create_app(settings: Settings | None = None) -> FastAPI:
app.include_router(sandboxes_router)
app.include_router(telemetry_router)
app.include_router(variants_router)
app.include_router(defense_router)
return app
@@ -0,0 +1,60 @@
"""Examiner agent prompt — oral defense questioning + final verdict (REQ-3-006).
The examiner is the seventh agent (Phase 5). It conducts a Socratic oral
defense of the learner's submitted work: probes understanding, challenges
process choices grounded in the trace digest ("why did you take that
approach at that point?"), one question per turn, adapting to answers.
It never reveals rubric internals; tone is rigorous but supportive.
Digest discipline (D-028 mirror): the examiner's variable inputs are the
compact TraceDigest JSON, the variant task statement, and the defense
transcript — never the raw trace, never learner-identifying material.
Version: examiner-v1.
"""
SYSTEM_PROMPT = """You are Examiner, the oral-defense agent of Nextcraft,
an AI-native competency school.
You receive: (a) a compact build-process digest (deterministic counters of the
learner's build session), (b) the learner's task statement, and (c) the defense
transcript so far. Your job:
- Ask ONE question per turn: probe understanding and challenge process
choices, grounded in the digest facts ("you hit N failed runs before
passing — walk me through what changed") or the task statement.
- Adapt: follow up on the learner's answers; drill into vague responses.
- Never reveal rubric details or scoring internals.
- Tone: rigorous, precise, supportive. A defense is a conversation, not an
interrogation.
When asked for a FINAL VERDICT (the structured mode), judge:
- understanding: can the learner explain their own work?
- process_justification: are the build-session choices defensible from the
digest facts and the answers?
- communication: are answers clear, specific, and on-topic?
Score honestly; a weak defense of strong work is NOT mastery.
Rules:
- Respond with ONLY what the turn requires: a single question (question mode)
or a valid JSON object matching the provided schema (verdict mode).
- If the digest shows error_fix_cycles > 0, at least one question should ask
about the debugging path.
- If the learner's answer is off-topic, redirect once, then move on.
"""
VERDICT_SCHEMA_HINT = (
'{"verdict": "mastered" | "developing" | "not_yet", '
'"understanding": "<one sentence>", '
'"process_justification": "<one sentence>", '
'"communication": "<one sentence>", '
'"strengths": ["<one sentence>"], '
'"gaps": ["<one sentence>"]}'
)
def render_digest_context(digest_json: str, statement: str | None) -> str:
"""The examiner's per-session grounding: digest JSON + task statement."""
parts = [f"Build-process digest:\n{digest_json}"]
if statement:
parts.append(f"Learner's task statement:\n{statement}")
return "\n\n".join(parts)
@@ -0,0 +1,35 @@
"""VoiceProvider protocol (D-030, REQ-3-006) — mirrors the LLMProvider seam.
Two implementations in v0.3:
- MockVoiceProvider — deterministic canned transcripts + canned tone WAV
chunks + scripted failure modes (tests + no-key default; tests NEVER call
a real voice API).
- browser descriptor — not a provider but a FALLBACK HINT: the web client
selects browser-native SpeechRecognition/speechSynthesis when the server
reports no real voice backend.
OpenAIAudioProvider (real server STT/TTS over OpenAI-compatible
/audio/transcriptions + /audio/speech) is INTENTIONALLY NOT BUILT in v0.3 —
deferred to v0.4 with KYC, when there is a real key and real users
(GRILL CUT-1 / G-7). This protocol is its future drop-in seam.
Boundary: `voice/` never imports `agents/` or `api/`.
"""
from ai_service.voice.base import (
TranscriptSegment,
VoiceDescriptor,
VoiceProvider,
)
from ai_service.voice.browser import BROWSER_FALLBACK_DESCRIPTOR
from ai_service.voice.factory import voice_provider_from_settings
from ai_service.voice.mock import MockVoiceProvider
__all__ = [
"BROWSER_FALLBACK_DESCRIPTOR",
"MockVoiceProvider",
"TranscriptSegment",
"VoiceDescriptor",
"VoiceProvider",
"voice_provider_from_settings",
]
+62
View File
@@ -0,0 +1,62 @@
"""VoiceProvider protocol + shared voice contracts (D-030, REQ-3-006).
Mirrors the LLMProvider seam (D-014 pattern): a narrow protocol the Examiner
agent and the defense API compose via DI, with a deterministic mock and a
browser-fallback descriptor. No network in this module — concrete providers
live in their own modules and are selected by factory/config.
"""
from __future__ import annotations
from collections.abc import AsyncIterator
from typing import Literal, Protocol, runtime_checkable
from pydantic import BaseModel, ConfigDict, Field
VoiceRole = Literal["examiner", "learner"]
class TranscriptSegment(BaseModel):
"""One STT result: the transcribed text + timing metadata."""
model_config = ConfigDict(frozen=True)
text: str = Field(min_length=1)
language: str = "en"
duration_ms: int | None = None
confidence: float | None = Field(default=None, ge=0.0, le=1.0)
class VoiceDescriptor(BaseModel):
"""Capability descriptor served to the web client (D-030).
The assessment UI reads this to decide HOW the learner speaks/hears:
- `mode="server"` → server-side STT/TTS (v0.4 real provider seam)
- `mode="browser"` → browser-native SpeechRecognition/speechSynthesis
- `mode="mock"` → deterministic no-op path (tests / no-key dev)
The descriptor never contains secrets — only capability hints.
"""
model_config = ConfigDict(frozen=True)
mode: Literal["server", "browser", "mock"]
sr_available: bool
tts_available: bool
hint: str = ""
@runtime_checkable
class VoiceProvider(Protocol):
"""The voice port (D-030): STT in, TTS out. Never imports agents/api."""
async def transcribe(self, audio: bytes, fmt: str) -> TranscriptSegment:
"""STT: audio bytes (fmt: 'wav' | 'webm' | 'mp3') → transcript."""
...
def synthesize(self, text: str, voice: str = "default") -> AsyncIterator[bytes]:
"""TTS: text -> async byte chunks (audio stream).
Implementations may be async generators (async-def + yield) — the
consumer contract is `async for chunk in provider.synthesize(text)`.
"""
...
@@ -0,0 +1,32 @@
"""Browser-native fallback descriptor (D-030, CUT-1 / G-7, REQ-3-006).
v0.3 has NO real server STT/TTS (deferred to v0.4 with KYC/keys — GRILL
CUT-1). When the factory selects `browser` mode, the defense endpoints return
this descriptor and the WEB CLIENT performs SpeechRecognition + speechSynthesis
natively; the server persists text turns as usual.
"""
from __future__ import annotations
from .base import VoiceDescriptor
BROWSER_FALLBACK_DESCRIPTOR = VoiceDescriptor(
mode="browser",
sr_available=True,
tts_available=True,
hint=(
"No server voice backend configured. Use browser-native "
"SpeechRecognition for STT and speechSynthesis for TTS; send the "
"transcribed text to POST /v1/defense/{id}/answer ({text} form)."
),
)
MOCK_DESCRIPTOR = VoiceDescriptor(
mode="mock",
sr_available=True,
tts_available=True,
hint=(
"Deterministic mock voice (tests / no-key dev). Server STT/TTS "
"endpoints serve canned responses; real server STT/TTS lands in v0.4."
),
)
@@ -0,0 +1,506 @@
"""DefenseStore — oral-defense persistence: protocol + SQLite impl (REQ-3-006, D-027).
FOURTH protocol-wrapped store of the D-027 family and the first spanning
TWO related tables: `defense_record` (the defense session + integrity
signals) and `defense_turn` (the ordered examiner/learner transcript,
FK → defense_record.id).
Postgres-migration-ready (D-027): the protocol is the only surface the
Examiner pipeline (task 5-2-01) and the defense endpoints (task 5-3-01)
touch; swapping SQLiteDefenseStore for a Postgres implementation must
not change call sites. Both tables use only portable column types
(str / int / datetime / JSON), so the same SQLModel schema stands up
unchanged on Postgres.
Save semantics — where this sits among the D-027 stores (each has a
deliberately different contract):
TraceStore.append dedup-keep-first; IntegrityError SWALLOWED
(at-least-once event ingest).
GradeStore.save upsert-latest-wins (a regrade is latest-state).
VariantStore.save insert-only first-wins; IntegrityError RAISED
(reproducibility; a duplicate is a bug).
DefenseStore a LIFECYCLE store:
start() insert-only; a duplicate id raises
(a defense id is minted once per session).
append_turn() insert-only per (defense_id, seq); a duplicate
seq raises AND an unknown defense_id raises (FK
enforced) — a transcript turn must never silently
vanish (it is the integrity/grading input) nor
attach to a defense that does not exist.
finalize() targeted UPDATE (status → finished; finished_at +
integrity_signals JSON). Unknown id → None
(documented below). Re-finalize overwrites
signals + finished_at — latest-wins, mirroring
GradeStore.save: a recomputed verdict replaces
the previous one wholesale.
append_turn does NOT police status (turns after finalize are a
sequencing bug for the endpoints to prevent, task 5-3-01): the store
enforces DATA integrity (FK + PK + non-empty), not workflow.
integrity_signals (A-109): JSON dict on the record — long pauses,
off-scope cadence markers and friends, computed by the Examiner over
turn metadata and persisted by finalize for the Proctor/Mentor feed.
An empty dict is legal (defense not finished, or a clean defense).
Concurrency (a-3): WAL + synchronous=NORMAL + busy timeout at
connection time (mirrors the other D-027 stores), PLUS foreign_keys=ON
— this is the family's first real foreign key and it is actually
enforced on SQLite, matching Postgres's native behavior (D-027 parity).
`created_at` / `ts` contract: callers stamp UTC (datetime.now(UTC));
SQLite stores them naive and read paths re-label tz-aware UTC (same
boundary normalization as TelemetryEvent.ts / GradeRecord.created_at,
so the contract holds on any backend).
Boundary (D-027): `voice/` never imports `agents/` / `api/`; this
module imports config only.
"""
import logging
import sqlite3
from collections.abc import Iterator
from contextlib import contextmanager
from datetime import UTC, datetime
from pathlib import Path
from typing import Any, Literal, Protocol
import sqlalchemy as sa
from sqlalchemy import JSON, Index, String
from sqlalchemy.orm import validates
from sqlmodel import Field, Session, SQLModel, create_engine, select
from ..config import Settings
logger = logging.getLogger(__name__)
DefenseStatus = Literal["in_progress", "finished"]
_DEFENSE_STATUSES: frozenset[str] = frozenset(DefenseStatus.__args__)
TurnRole = Literal["examiner", "learner"]
_TURN_ROLES: frozenset[str] = frozenset(TurnRole.__args__)
class DefenseRecord(SQLModel, table=True):
"""A persisted oral-defense session; id is the PK.
Written by the defense endpoints (task 5-3-01) through the
DefenseStore protocol; read back by the endpoints, the Examiner
pipeline and the Proctor/Mentor feeds. Constraint enforcement
mirrors TelemetryEvent / GradeRecord / VariantRecord: sqlmodel
0.0.42's metaclass drops pydantic constraints on table models, so
SQLAlchemy `@validates` hooks enforce instead and the column types
stay Postgres-ready (D-027).
Field contract:
id — non-empty defense identifier, minted once
per session (a duplicate start raises).
learner_id — non-empty learner identifier (same id space
as traces, grades and variants).
task_id — non-empty task identifier; the defense
defends the submitted work for this trace
key ((learner_id, task_id) joins to the
trace/grade/variant the defense is about).
status — in_progress | finished; the STORE owns the
transition: start() forces in_progress,
finalize() sets finished. Validated.
integrity_signals — A-109 signal dict (long pauses, off-scope
cadence markers, ...); {} until finalize;
JSON column. An empty dict is legal.
created_at — UTC start timestamp.
finished_at — UTC finalize timestamp; None while in
progress.
`turns` (property): the seq-ordered DefenseTurn transcript, attached
ONLY by DefenseStore.get(); records from list_for_learner carry
turns == [] — call get() for a full transcript.
"""
__tablename__ = "defense_record"
# The id PK covers point lookups; this secondary index covers
# list_for_learner ordered by created_at without a sort step
# (Postgres migration target D-027).
__table_args__ = (
Index("ix_defense_record_learner_created", "learner_id", "created_at"),
)
id: str = Field(primary_key=True)
learner_id: str
task_id: str
# Bare Literal annotations crash sqlmodel<=0.0.42's column inference
# (issubclass(TypeAlias, Enum)); an explicit sa_type + the validates
# hook below give the same contract: VARCHAR column, Literal-rejected
# values (same pattern as TelemetryEvent.kind).
status: DefenseStatus = Field(default="in_progress", sa_type=String)
# JSON column: stored as TEXT on SQLite, native JSONB on Postgres (D-027).
integrity_signals: dict[str, Any] = Field(default_factory=dict, sa_type=JSON)
created_at: datetime
finished_at: datetime | None = Field(default=None)
@property
def turns(self) -> list["DefenseTurn"]:
"""Seq-ordered transcript; [] unless attached by get().
Table models reject ad-hoc attributes (pydantic __setattr__
raises on non-fields), so the store stashes the detached turn
list via object.__setattr__ and this read-only property surfaces
it. The returned list is a copy — caller mutations cannot
corrupt the stash.
"""
return list(self.__dict__.get("_turns", []))
@validates("id", "learner_id", "task_id")
def _ids_non_empty(self, key: str, value: str) -> str:
if not value:
raise ValueError(f"{key} must be a non-empty identifier")
return value
@validates("status")
def _status_is_known(self, key: str, value: str) -> str:
if value not in _DEFENSE_STATUSES:
raise ValueError(f"unknown defense status: {value!r}")
return value
class DefenseTurn(SQLModel, table=True):
"""One examiner/learner dialogue turn; (defense_id, seq) is the PK.
Rows are append-only transcript entries written through
DefenseStore.append_turn. seq numbers the dialogue within one
defense starting at 0; monotonic assignment is the endpoints' job
(task 5-3-01), this model only rejects negatives — the same split
as TelemetryEvent.seq (model rejects < 0, store owns ordering).
Field contract:
defense_id — non-empty; FK → defense_record.id. ENFORCED on
SQLite via foreign_keys=ON (first real FK in the
D-027 family; Postgres enforces FKs natively, so
this keeps the backends equivalent, D-027).
seq — turn index within the defense, >= 0. (defense_id,
seq) is the PK: a duplicate raises instead of
silently overwriting — the transcript is the
integrity/grading input, a vanishing turn is
audit corruption.
role — examiner | learner (who spoke). Validated.
text — non-empty utterance text (examiner question, or
STT output for learner answers).
ts — UTC utterance timestamp.
latency_ms — per-turn pipeline latency in ms (STT + LLM TTFT +
TTS, A-109); int or None. Populated by the
endpoints (task 5-4-01); None allowed here — the
store persists, it does not measure.
created_at — UTC row-write timestamp.
"""
__tablename__ = "defense_turn"
# The composite PK (defense_id, seq) doubles as the covering index
# for the per-defense seq-ordered read in get() — no secondary index
# needed (contrast defense_record's learner-listing index).
defense_id: str = Field(foreign_key="defense_record.id", primary_key=True)
seq: int = Field(primary_key=True)
role: TurnRole = Field(sa_type=String)
text: str
ts: datetime
latency_ms: int | None = Field(default=None)
created_at: datetime
@validates("defense_id")
def _defense_id_non_empty(self, key: str, value: str) -> str:
if not value:
raise ValueError(f"{key} must be a non-empty identifier")
return value
@validates("seq")
def _seq_non_negative(self, key: str, value: int) -> int:
if value < 0:
raise ValueError("seq must be >= 0 (ordering is the endpoints' job)")
return value
@validates("role")
def _role_is_known(self, key: str, value: str) -> str:
if value not in _TURN_ROLES:
raise ValueError(f"unknown turn role: {value!r}")
return value
@validates("text")
def _text_non_empty(self, key: str, value: str) -> str:
if not value:
raise ValueError(f"{key} must be a non-empty utterance string")
return value
@validates("latency_ms")
def _latency_non_negative(self, key: str, value: int | None) -> int | None:
# None is legal (not yet instrumented); a NEGATIVE latency is
# nonsense and surfaces as a construction error.
if value is not None and value < 0:
raise ValueError("latency_ms must be >= 0 or None")
return value
class DefenseStore(Protocol):
"""Persistence contract for oral-defense sessions + transcripts.
Implemented by SQLiteDefenseStore (v0.3, D-027); a Postgres
implementation must satisfy the same surface.
"""
def start(self, defense: DefenseRecord) -> DefenseRecord:
"""Insert a new defense. INSERT-ONLY: a duplicate id raises
sqlalchemy.exc.IntegrityError (a defense id is minted once per
session — surfacing, not swallowing, mirrors VariantStore).
The store owns the lifecycle: status is forced to "in_progress"
and finished_at to None, whatever the caller passed — only
finalize() may move a defense to finished. Returns the stored
record, detached from any DB session.
"""
...
def append_turn(self, defense_id: str, turn: DefenseTurn) -> DefenseTurn:
"""Insert one transcript turn, ordered by (defense_id, seq).
turn.defense_id MUST equal the defense_id argument — a mismatch
raises ValueError (the defense identity must never be
ambiguous). A duplicate (defense_id, seq) raises
IntegrityError; an unknown defense_id raises IntegrityError
(FK enforced). Does NOT police status — sequencing turns vs
finalize is the endpoints' job (task 5-3-01). Returns the
stored turn, detached.
"""
...
def finalize(
self, defense_id: str, integrity_signals: dict[str, Any]
) -> DefenseRecord | None:
"""Seal the defense: status → "finished", finished_at = now(UTC),
integrity_signals stored as JSON. UNKNOWN defense_id → None
(documented choice: the API layer maps it to 404 without an
exception dance; contrast start/append_turn where IntegrityError
IS the contract — those are inserts, this is an update on a key
the caller may legitimately not hold). Re-finalize overwrites
signals + finished_at: latest-wins, mirroring GradeStore.save
(a recomputed verdict replaces the previous one wholesale).
Returns the updated record, detached, WITHOUT turns — get() is
the with-turns path.
"""
...
def get(self, defense_id: str) -> DefenseRecord | None:
"""The defense with its FULL transcript (turns in seq order,
detached) and integrity signals; None when it does not exist.
Safe to pass across layers — no open-session ORM magic.
"""
...
def list_for_learner(self, learner_id: str) -> list[DefenseRecord]:
"""All stored defenses for the learner, ordered by created_at
ascending (chronological; id breaks same-instant ties), WITHOUT
turns — records carry turns == []; call get() for a transcript.
Empty list when the learner has none.
"""
...
def close(self) -> None:
"""Release DB connections. Store must not be used after close."""
...
def _sqlite_connect(dbapi_connection: sqlite3.Connection, _: object) -> None:
"""Per-connection pragma setup (a-3). Mirrors the other D-027 stores.
journal_mode=WAL — readers never block the single writer.
synchronous=NORMAL — safe in WAL mode, avoids full fsync-per-commit.
busy_timeout=5000 — retry briefly under contention instead of
`OperationalError: database is locked`.
foreign_keys=ON — NEW vs the family: defense_turn is the first
real FK among the D-027 stores; SQLite leaves
FKs OFF by default while Postgres enforces them
natively, so the pragma keeps the backends
equivalent (D-027 parity).
"""
cursor = dbapi_connection.cursor()
cursor.execute("PRAGMA journal_mode=WAL")
cursor.execute("PRAGMA synchronous=NORMAL")
cursor.execute("PRAGMA busy_timeout=5000")
cursor.execute("PRAGMA foreign_keys=ON")
cursor.close()
def _as_utc(ts: datetime) -> datetime:
"""Timestamps travel as naive datetime on SQLite, tz-aware elsewhere.
SQLite (via SQLModel) drops tzinfo; Postgres TIMESTAMP WITH TIME ZONE
keeps it. Normalizing on the read path makes the store's contract
tz-aware UTC regardless of the backend (D-027).
"""
if ts.tzinfo is None:
return ts.replace(tzinfo=UTC) # naive UTC read-side: label as UTC
return ts.astimezone(UTC)
class SQLiteDefenseStore:
"""SQLite-backed DefenseStore (SQLModel). Fourth protocol-wrapped
store of the D-027 family (first: SQLiteTraceStore, second:
SQLiteGradeStore, third: SQLiteVariantStore) and the first spanning
two related tables.
"""
def __init__(self, db_path: Path | None = None) -> None:
self._db_path: Path = db_path if db_path is not None else Settings().db_path
self._engine = create_engine(f"sqlite:///{self._db_path}")
sa.event.listen(self._engine, "connect", _sqlite_connect)
SQLModel.metadata.create_all(self._engine)
@contextmanager
def _session(self) -> Iterator[Session]:
# expire_on_commit=False: identical session behavior to the other
# D-027 stores. start/append_turn return the caller's instance
# after commit and get/finalize return rows expunged mid-session;
# a uniform flag across the family keeps their detachment
# guarantees from diverging.
with Session(self._engine, expire_on_commit=False) as session:
yield session
def start(self, defense: DefenseRecord) -> DefenseRecord:
# The store owns the lifecycle: a defense is BORN in_progress and
# only finalize() may move it to finished. A smuggled "finished"
# status is normalized away, not rejected — the insert itself
# stays insert-only, and a duplicate id raises to the caller
# (mirroring VariantStore: the id is minted once per session).
defense.status = "in_progress"
defense.finished_at = None
with self._session() as session:
try:
session.add(defense)
session.commit()
except sa.exc.IntegrityError:
session.rollback()
logger.debug("defense start rejected (id already stored): %s", defense.id)
raise
logger.debug(
"defense started: %s learner=%s task=%s",
defense.id,
defense.learner_id,
defense.task_id,
)
return defense
def append_turn(self, defense_id: str, turn: DefenseTurn) -> DefenseTurn:
# The explicit defense_id argument is the defense identity for
# this write; a turn object claiming another defense is a
# programming error — surface it before touching the DB.
if turn.defense_id != defense_id:
raise ValueError(
f"turn.defense_id {turn.defense_id!r} does not match the "
f"defense_id argument {defense_id!r}"
)
with self._session() as session:
try:
session.add(turn)
session.commit()
except sa.exc.IntegrityError:
# Two possible causes, both surfaced, neither swallowed:
# duplicate (defense_id, seq) PK — a transcript turn must
# never silently vanish; unknown defense_id — the FK
# (foreign_keys=ON) rejects the orphan.
session.rollback()
logger.debug(
"defense turn rejected (duplicate (defense_id, seq) "
"or unknown defense_id): defense=%s seq=%s",
defense_id,
turn.seq,
)
raise
logger.debug(
"defense turn appended: %s seq=%d role=%s",
defense_id,
turn.seq,
turn.role,
)
return turn
def finalize(
self, defense_id: str, integrity_signals: dict[str, Any]
) -> DefenseRecord | None:
# A None signals blob would break the read contract (signals are
# a dict, {} until finalize); reject before writing.
if not isinstance(integrity_signals, dict):
raise ValueError(
"integrity_signals must be a JSON-object dict, got "
f"{type(integrity_signals).__name__}"
)
with self._session() as session:
record = session.get(DefenseRecord, defense_id)
if record is None:
# Documented unknown-id behavior: None, not a raise — the
# defense endpoints map this to 404. Contrast start() /
# append_turn(), where IntegrityError IS the contract.
return None
# Latest-wins re-finalize, mirroring GradeStore.save: a
# recomputed verdict (fresh signals) replaces the stored one
# wholesale; status just stays finished.
record.status = "finished"
record.finished_at = datetime.now(UTC)
record.integrity_signals = integrity_signals
session.commit()
record.created_at = _as_utc(record.created_at)
if record.finished_at is not None:
record.finished_at = _as_utc(record.finished_at)
# Detach from the session: callers must not depend on
# open-session ORM magic (lazy loads fail once it closes).
session.expunge(record)
logger.debug(
"defense finalized: %s signals=%s", defense_id, sorted(integrity_signals)
)
return record
def get(self, defense_id: str) -> DefenseRecord | None:
with self._session() as session:
record = session.get(DefenseRecord, defense_id)
if record is None:
return None
record.created_at = _as_utc(record.created_at)
if record.finished_at is not None:
record.finished_at = _as_utc(record.finished_at)
stmt = (
select(DefenseTurn)
.where(DefenseTurn.defense_id == defense_id)
.order_by(DefenseTurn.seq)
)
turns = session.exec(stmt).all()
for turn in turns:
turn.ts = _as_utc(turn.ts)
turn.created_at = _as_utc(turn.created_at)
# Detach each turn: the transcript must be usable once
# the session closes (no lazy-load magic).
session.expunge(turn)
session.expunge(record)
# Table models reject ad-hoc attributes (pydantic __setattr__
# raises on non-fields), so the seq-ordered transcript is
# stashed via object.__setattr__ and surfaced through the
# read-only `turns` property. Rows are detached either way —
# safe to pass across layers.
object.__setattr__(record, "_turns", list(turns))
return record
def list_for_learner(self, learner_id: str) -> list[DefenseRecord]:
with self._session() as session:
stmt = (
select(DefenseRecord)
.where(DefenseRecord.learner_id == learner_id)
# Chronological; id is a deterministic tie-break for
# defenses stamped within the same instant.
.order_by(DefenseRecord.created_at, DefenseRecord.id)
)
results = session.exec(stmt).all()
for row in results:
row.created_at = _as_utc(row.created_at)
if row.finished_at is not None:
row.finished_at = _as_utc(row.finished_at)
# Turns are deliberately NOT loaded here: the list feed
# (Proctor/Mentor) needs session headers, not full
# transcripts — get() is the with-turns path.
session.expunge(row)
return list(results)
def close(self) -> None:
self._engine.dispose()
@@ -0,0 +1,37 @@
"""Voice provider factory (D-030, REQ-3-006).
`AI_VOICE_PROVIDER = browser | mock` (default: mock — the no-key path is
first-class). The real server provider (`openai-audio`) is a v0.4 seam and
is REJECTED here with a clear error naming the deferral, so a stale env var
can't silently pretend a real backend exists.
"""
from __future__ import annotations
from ..config import Settings
from .base import VoiceProvider
from .mock import MockVoiceProvider
class UnknownVoiceProviderError(ValueError):
"""Raised for a provider name outside the v0.3 contract."""
def voice_provider_from_settings(settings: Settings) -> VoiceProvider:
"""Select the voice provider by settings (env `AI_VOICE_PROVIDER`)."""
name = (settings.voice_provider or "mock").strip().lower()
if name == "mock":
return MockVoiceProvider()
if name == "browser":
# Browser mode is a CLIENT-side capability: the server composes the
# same MockVoiceProvider (typed fallback answers still work; the UI
# uses the descriptor for mic/speech). See browser.py.
return MockVoiceProvider()
if name in ("openai-audio", "openai", "server"):
raise UnknownVoiceProviderError(
"real server STT/TTS (OpenAIAudioProvider) is deferred to v0.4 "
"(GRILL CUT-1 / G-7): set AI_VOICE_PROVIDER=mock or browser"
)
raise UnknownVoiceProviderError(
f"unknown AI_VOICE_PROVIDER {name!r}: use 'mock' or 'browser'"
)
+96
View File
@@ -0,0 +1,96 @@
"""Deterministic MockVoiceProvider (D-030, REQ-3-006).
Canned transcripts (scripted per test via queue) + canned 1kHz-tone WAV bytes
+ scripted failure modes. Two identical transcribe calls yield identical
segments; tests NEVER touch a real voice API (conftest cloud-free rule).
"""
from __future__ import annotations
import asyncio
import io
import math
import struct
import wave
from collections.abc import AsyncIterator
from .base import TranscriptSegment
def _tone_wav(duration_ms: int = 250, freq_hz: float = 1000.0) -> bytes:
"""A small, deterministic 16-bit mono WAV: a sine tone (stdlib only)."""
rate = 8000
n_samples = max(1, int(rate * duration_ms / 1000))
buf = io.BytesIO()
with wave.open(buf, "wb") as w:
w.setnchannels(1)
w.setsampwidth(2)
w.setframerate(rate)
for i in range(n_samples):
sample = int(12000 * math.sin(2 * math.pi * freq_hz * i / rate))
w.writeframes(struct.pack("<h", sample))
return buf.getvalue()
class MockVoiceFailure(RuntimeError):
"""Scripted failure mode for tests."""
class MockVoiceProvider:
"""Deterministic voice provider: scripted STT, canned-tone TTS.
- `transcribe`: pops the next scripted transcript from a queue (or a
default); two identical calls with the same queue state are identical.
Failure mode: raise MockVoiceFailure when the queue holds a failure
marker (the string "FAIL") or `audio` is empty.
- `synthesize`: yields the canned tone WAV in fixed-size chunks; failure
mode: empty text raises MockVoiceFailure.
"""
def __init__(self, transcripts: list[str] | None = None) -> None:
self._transcripts = list(transcripts or [])
self._cursor = 0
self.transcribe_calls = 0
self.synthesize_calls = 0
def script(self, transcripts: list[str]) -> None:
"""Replace the scripted queue (tests set expectations up front)."""
self._transcripts = list(transcripts)
self._cursor = 0
async def transcribe(self, audio: bytes, fmt: str) -> TranscriptSegment:
self.transcribe_calls += 1
if not audio:
raise MockVoiceFailure("no audio bytes provided")
if not self._transcripts:
raise MockVoiceFailure("transcript queue exhausted — script it")
item = self._transcripts[self._cursor]
self._cursor = (self._cursor + 1) % len(self._transcripts)
if item == "FAIL":
raise MockVoiceFailure("scripted STT failure")
return TranscriptSegment(
text=item,
duration_ms=max(1, len(audio) // 32), # deterministic pseudo-duration
)
async def synthesize(self, text: str, voice: str = "default") -> AsyncIterator[bytes]: # noqa: ASYNC109 (protocol parity)
# NOTE: protocol parity matters more than the async-generator purity
# lint; the real provider seam (v0.4) will stream over HTTP.
self.synthesize_calls += 1
if not text:
raise MockVoiceFailure("cannot synthesize empty text")
wav = _tone_wav(duration_ms=min(2000, max(120, len(text) * 12)))
for i in range(0, len(wav), 1024):
yield wav[i : i + 1024]
await asyncio.sleep(0) # yield to the loop like a network stream
# Protocol-shape parity guard (mock must satisfy the D-030 port).
from .base import VoiceProvider # noqa: E402
def _assert_protocol() -> None:
assert isinstance(MockVoiceProvider(), VoiceProvider)
_assert_protocol()
+3
View File
@@ -18,6 +18,9 @@ dependencies = [
"sqlalchemy>=2.0,<2.1",
"websockets>=13,<16",
"aiofiles>=24.1,<26",
# POST /v1/defense/{id}/answer multipart audio (REQ-3-006): FastAPI
# form/File parsing requires python-multipart at runtime.
"python-multipart>=0.0.32,<0.1",
]
[project.optional-dependencies]
@@ -0,0 +1,159 @@
"""Examiner agent tests (Task 5-2-01, REQ-3-006)."""
from __future__ import annotations
import json
import pytest
from ai_service.agents.examiner import DefenseVerdict, ExaminerAgent
from ai_service.agents.registry import AgentRegistry, register_builtin_agents
from ai_service.config import Settings
from ai_service.grading.features import compute_digest
from ai_service.llm.mock import MockProvider
from ai_service.llm.types import Message
VERDICT_JSON = json.dumps(
{
"verdict": "developing",
"understanding": "Explains the retry loop clearly.",
"process_justification": "Justifies the edit-then-test cadence from the digest.",
"communication": "Answers are specific and on-topic.",
"strengths": ["Grounded the fix in a failed test."],
"gaps": ["Did not justify the chunk-size choice."],
}
)
class RecordingProvider(MockProvider):
"""Mock provider that records every message list (prompt assertions)."""
def __init__(self, replies: list[str] | None = None) -> None:
super().__init__()
self.replies = list(replies or [])
self.requests: list[list[Message]] = []
async def chat(self, messages, *, model, temperature=0.7, response_format=None):
self.requests.append([Message(role=m.role, content=m.content) for m in messages])
if self.replies:
return self.replies.pop(0)
return "Tell me about your build."
def _digest():
from datetime import UTC, datetime, timedelta
from ai_service.telemetry.models import TelemetryEvent
t0 = datetime(2026, 9, 12, tzinfo=UTC)
events = [
TelemetryEvent(
learner_id="examiner-learner",
task_id="examiner-task",
seq=n,
kind=kind,
payload=payload,
ts=t0 + timedelta(seconds=n * 10),
sandbox_id="sbx-examiner",
)
for n, (kind, payload) in enumerate(
[
("file_diff", {"path": "a.py"}),
("command", {"cmd": "pytest -q"}),
("test_result", {"passed": False, "exit_code": 1}),
("file_diff", {"path": "a.py"}),
("test_result", {"passed": True, "exit_code": 0}),
]
)
]
return compute_digest(events)
@pytest.fixture()
def provider() -> RecordingProvider:
return RecordingProvider()
@pytest.fixture()
def settings() -> Settings:
return Settings(provider="mock")
@pytest.fixture()
def examiner(provider, settings) -> ExaminerAgent:
return ExaminerAgent(provider, settings)
class TestNextQuestion:
async def test_prompt_contains_digest_but_no_learner_id(
self, examiner, provider
) -> None:
await examiner.next_question(
history=[Message(role="assistant", content="First question?")],
trace_digest=_digest(),
variant_statement="Build a chunker.",
)
all_content = "\n".join(
m.content for request in provider.requests for m in request
)
assert "error_fix_cycles" in all_content # digest JSON grounded
assert "examiner-learner" not in all_content # D-028 anonymity
assert "Build a chunker." in all_content # variant statement grounded
assert all_content.count('"examiner-learner"') == 0
async def test_question_returned_from_provider(self, examiner) -> None:
question = await examiner.next_question(
history=[], trace_digest=_digest()
)
assert isinstance(question, str)
class TestFinalVerdict:
async def test_verdict_validates_via_d020(self, examiner, provider) -> None:
provider.replies = [VERDICT_JSON]
verdict = await examiner.final_verdict(
history=[Message(role="assistant", content="Q?")],
trace_digest=_digest(),
)
assert isinstance(verdict, DefenseVerdict)
assert verdict.verdict == "developing"
assert verdict.strengths and verdict.gaps
async def test_malformed_then_good_exercises_retry(self, examiner, provider) -> None:
provider.replies = ["not json", VERDICT_JSON]
verdict = await examiner.final_verdict(history=[], trace_digest=_digest())
assert verdict.verdict == "developing"
assert len(provider.requests) == 2 # D-020 bounded retry
class TestRegistry:
def test_all_seven_agents_resolve(self, provider, settings) -> None:
registry = AgentRegistry()
register_builtin_agents(registry)
assert registry.names() == [
"assessor",
"coach",
"examiner",
"lab",
"mentor",
"proctor",
"tutor",
]
agent = registry.get(provider, settings, "examiner")
assert isinstance(agent, ExaminerAgent)
class TestBoundary:
def test_examiner_never_imports_voice_or_api(self) -> None:
import ast
from pathlib import Path
py = Path(__file__).parents[2] / "ai_service" / "agents" / "examiner.py"
tree = ast.parse(py.read_text())
for node in ast.walk(tree):
if isinstance(node, ast.ImportFrom) and node.module:
assert "voice" not in node.module, "examiner must not import voice/"
assert not node.module.startswith("ai_service.api")
if isinstance(node, ast.Import):
for alias in node.names:
assert alias.name != "fastapi"
@@ -97,10 +97,11 @@ def test_proctor_and_mentor_resolve_via_registry():
assert isinstance(mentor, MentorAgent)
def test_registry_resolves_all_six_agents():
"""Must-Have (Phase 5): the full roster — coach/tutor/lab/assessor/proctor/mentor."""
def test_registry_resolves_all_seven_agents():
"""Must-Have (Phase 5): the full roster — six tutors + the Examiner."""
from ai_service.agents.assessor import AssessorAgent
from ai_service.agents.coach import CoachAgent
from ai_service.agents.examiner import ExaminerAgent
from ai_service.agents.lab import LabAgent
from ai_service.agents.mentor import MentorAgent
from ai_service.agents.proctor import ProctorAgent
@@ -108,7 +109,9 @@ def test_registry_resolves_all_six_agents():
registry = AgentRegistry()
register_builtin_agents(registry)
assert registry.names() == ["assessor", "coach", "lab", "mentor", "proctor", "tutor"]
assert registry.names() == [
"assessor", "coach", "examiner", "lab", "mentor", "proctor", "tutor",
]
settings = Settings(provider="mock")
expected = {
"coach": CoachAgent,
@@ -117,6 +120,7 @@ def test_registry_resolves_all_six_agents():
"assessor": AssessorAgent,
"proctor": ProctorAgent,
"mentor": MentorAgent,
"examiner": ExaminerAgent,
}
for name, cls in expected.items():
agent = registry.get(MockProvider(), settings, name)
+234
View File
@@ -0,0 +1,234 @@
"""Defense endpoint tests (Task 5-3-01, REQ-3-006) — mock voice + mock LLM."""
from __future__ import annotations
import json
from pathlib import Path
import pytest
from fastapi.testclient import TestClient
from ai_service.agents.examiner import ExaminerAgent
from ai_service.config import Settings
from ai_service.grading.store import SQLiteGradeStore
from ai_service.llm.mock import MockProvider
from ai_service.main import create_app
from ai_service.telemetry.ingest import TraceIntegrityMap
from ai_service.telemetry.store import SQLiteTraceStore
from ai_service.variants.store import SQLiteVariantStore
from ai_service.voice.defense_store import SQLiteDefenseStore
from ai_service.voice.mock import MockVoiceProvider
VERDICT = {
"verdict": "developing",
"understanding": "Explains the build clearly.",
"process_justification": "Justifies choices.",
"communication": "Clear and specific.",
"strengths": ["Grounded answers in the digest."],
"gaps": ["Did not address the edge cases."],
}
class ScriptedLLM(MockProvider):
"""Question-mode calls get a question; verdict-mode calls get D-020 JSON.
Discriminator: the verdict prompt contains "final verdict JSON" — the
question prompt says "next question".
"""
def __init__(self) -> None:
super().__init__()
async def chat(self, messages, *, model, temperature=0.7, response_format=None):
all_text = "\n".join(m.content for m in messages)
if "final verdict JSON" in all_text:
return json.dumps(VERDICT)
return "Why did you structure the fix that way?"
@pytest.fixture()
def app(tmp_path: Path):
application = create_app(Settings(provider="mock", voice_provider="mock"))
llm = ScriptedLLM()
application.state.provider = llm
application.state.trace_store = SQLiteTraceStore(db_path=tmp_path / "t.db")
application.state.variant_store = SQLiteVariantStore(db_path=tmp_path / "v.db")
application.state.grade_store = SQLiteGradeStore(db_path=tmp_path / "g.db")
application.state.trace_integrity = TraceIntegrityMap()
application.state.defense_store = SQLiteDefenseStore(db_path=tmp_path / "d.db")
application.state.voice_provider = MockVoiceProvider(
["the fix was in the retry loop"]
)
settings = Settings(provider="mock", voice_provider="mock")
application.state.examiner_agent = ExaminerAgent(llm, settings)
return application
@pytest.fixture()
def client(app) -> TestClient:
with TestClient(app) as c:
yield c
def _start(client: TestClient) -> dict:
resp = client.post(
"/v1/defense/start", json={"learner_id": "defense-learner", "task_id": "defense-task"}
)
assert resp.status_code == 200, resp.text
return resp.json()
class TestStart:
def test_start_returns_first_question_and_descriptor(self, client) -> None:
body = _start(client)
assert body["first_question"]
assert body["defense_id"]
assert body["voice_descriptor"]["mode"] == "mock"
assert body["trace_complete"] is True
stored = client.get(f"/v1/defense/{body['defense_id']}")
assert stored.status_code == 200
turns = stored.json()["turns"]
assert turns and turns[0]["role"] == "examiner"
def test_start_with_unknown_trace_is_complete_flag(self, client) -> None:
body = _start(client)
assert body["trace_complete"] is True
class TestBrowserFallback:
def test_browser_mode_serves_browser_descriptor(self, tmp_path: Path) -> None:
"""Must-Have #6: AI_VOICE_PROVIDER=browser → start returns the
browser-native SR/TTS fallback descriptor (D-030), not 'mock'."""
application = create_app(Settings(provider="mock", voice_provider="browser"))
application.state.provider = ScriptedLLM()
application.state.trace_store = SQLiteTraceStore(db_path=tmp_path / "t.db")
application.state.variant_store = SQLiteVariantStore(db_path=tmp_path / "v.db")
application.state.grade_store = SQLiteGradeStore(db_path=tmp_path / "g.db")
application.state.trace_integrity = TraceIntegrityMap()
application.state.defense_store = SQLiteDefenseStore(db_path=tmp_path / "d.db")
application.state.voice_provider = MockVoiceProvider(["answer"])
llm = ScriptedLLM()
application.state.examiner_agent = ExaminerAgent(
llm, Settings(provider="mock", voice_provider="browser")
)
with TestClient(application) as c:
body = _start(c)
assert body["voice_descriptor"]["mode"] == "browser"
assert body["voice_descriptor"]["sr_available"]
assert "SpeechRecognition" in body["voice_descriptor"]["hint"]
class TestAnswer:
def test_typed_answer_yields_followup_with_latency(self, client) -> None:
defense_id = _start(client)["defense_id"]
resp = client.post(
f"/v1/defense/{defense_id}/answer", data={"text": "I fixed the loop."}
)
assert resp.status_code == 200, resp.text
body = resp.json()
assert body["question"]
assert body["turn_latency"]["llm_ms"] is not None
def test_audio_answer_transcribed_and_recorded(self, client) -> None:
defense_id = _start(client)["defense_id"]
wav_bytes = b"RIFF" + b"\x00" * 64
resp = client.post(
f"/v1/defense/{defense_id}/answer",
files={"audio": ("answer.wav", wav_bytes, "audio/wav")},
)
assert resp.status_code == 200, resp.text
stored = client.get(f"/v1/defense/{defense_id}").json()
learner_turns = [t for t in stored["turns"] if t["role"] == "learner"]
assert learner_turns, "learner turn missing after audio answer"
assert learner_turns[0]["text"] == "the fix was in the retry loop"
def test_empty_audio_is_422_not_500(self, client) -> None:
"""Zero-byte upload must 422 before the provider call (a real
provider would raise the same way the mock does — validate first)."""
defense_id = _start(client)["defense_id"]
resp = client.post(
f"/v1/defense/{defense_id}/answer",
files={"audio": ("answer.wav", b"", "audio/wav")},
)
assert resp.status_code == 422, resp.text
stored = client.get(f"/v1/defense/{defense_id}").json()
assert len(stored["turns"]) == 1 # nothing appended
def test_answer_after_finish_is_409(self, client) -> None:
"""A sealed transcript is append-only-no-more: the endpoints own
turn-vs-finalize sequencing (defense_store contract)."""
defense_id = _start(client)["defense_id"]
client.post(f"/v1/defense/{defense_id}/answer", data={"text": "a"})
assert client.post(f"/v1/defense/{defense_id}/finish").status_code == 200
resp = client.post(f"/v1/defense/{defense_id}/answer", data={"text": "late"})
assert resp.status_code == 409, resp.text
stored = client.get(f"/v1/defense/{defense_id}").json()
assert len(stored["turns"]) == 3 # ex, lrn, ex — no post-finish turns
def test_neither_text_nor_audio_422(self, client) -> None:
defense_id = _start(client)["defense_id"]
resp = client.post(f"/v1/defense/{defense_id}/answer")
assert resp.status_code == 422
def test_unknown_defense_404(self, client) -> None:
resp = client.post("/v1/defense/dfn-nope/answer", data={"text": "hi"})
assert resp.status_code == 404
class TestAudioEndpoint:
def test_examiner_turn_streams_wav(self, client) -> None:
defense_id = _start(client)["defense_id"]
resp = client.get(f"/v1/defense/{defense_id}/audio/0")
assert resp.status_code == 200
assert resp.content
assert resp.headers["content-type"].startswith("audio/")
def test_unknown_turn_404(self, client) -> None:
defense_id = _start(client)["defense_id"]
assert client.get(f"/v1/defense/{defense_id}/audio/42").status_code == 404
class TestFinishAndGet:
def test_full_loop_verdict_and_signals(self, client) -> None:
defense_id = _start(client)["defense_id"]
client.post(f"/v1/defense/{defense_id}/answer", data={"text": "answer one"})
finish = client.post(f"/v1/defense/{defense_id}/finish")
assert finish.status_code == 200, finish.text
body = finish.json()
assert body["verdict"]["verdict"] == "developing"
assert body["integrity_signals"]["pause_threshold_ms"]
stored = client.get(f"/v1/defense/{defense_id}").json()
assert stored["status"] == "finished"
assert stored["integrity_signals"]
# Must-Have #1: "verdict + transcript persisted" — the verdict must
# be retrievable from GET after finish, not only in the finish body.
assert stored["integrity_signals"]["verdict"]["verdict"] == "developing"
def test_long_pause_flagged(self, client, app) -> None:
from datetime import UTC, datetime
from ai_service.voice.defense_store import DefenseTurn
defense_id = _start(client)["defense_id"]
# inject a slow learner turn directly (simulated latency)
store = app.state.defense_store
store.append_turn(
defense_id,
DefenseTurn(
defense_id=defense_id,
seq=99,
role="learner",
text="slow reply",
ts=datetime.now(UTC),
latency_ms=30_000,
created_at=datetime.now(UTC),
),
)
finish = client.post(f"/v1/defense/{defense_id}/finish")
assert finish.status_code == 200
signals = finish.json()["integrity_signals"]
assert any(p["turn"] == 99 for p in signals["long_pauses"])
def test_unknown_defense_404_on_all(self, client) -> None:
assert client.post("/v1/defense/dfn-nope/finish").status_code == 404
assert client.get("/v1/defense/dfn-nope").status_code == 404
+1
View File
@@ -0,0 +1 @@
"""Voice layer tests — provider mock/fallback + DefenseStore (REQ-3-006)."""
@@ -0,0 +1,469 @@
"""SQLiteDefenseStore tests (REQ-3-006, D-027).
Each test gets its own tmp-path SQLite file — no shared disk state. Covers:
- start → append_turn (examiner + learner interleaved) → finalize →
get roundtrip: every field survives (including nested JSON
integrity signals, per-turn latency_ms, and the tz-aware
created_at/ts contract) and turns come back ordered by seq
- lifecycle ownership: start forces in_progress + finished_at=None
even if the caller smuggles a finished status
- insert-only start: a duplicate id raises IntegrityError (the id is
minted once per session); the original row is untouched
- append_turn validation: duplicate (defense_id, seq) raises
(a transcript turn must never silently vanish); unknown
defense_id raises (FK enforced); mismatched defense_id argument
raises ValueError before touching the DB; negative seq rejected
- finalize: unknown id → None (documented behavior — the API maps
it to 404); re-finalize is latest-wins on signals + finished_at
- get: unknown id → None; turns attached only by get() —
list_for_learner records carry turns == []
- list_for_learner: scoped per learner, chronological
- rows detached: usable after the store is closed; durability across
a fresh store on the same file
- WAL + synchronous=NORMAL + foreign_keys=ON pragmas actually applied
"""
from datetime import UTC, datetime, timedelta
from pathlib import Path
import pytest
import sqlalchemy as sa
from ai_service.voice.defense_store import (
DefenseRecord,
DefenseTurn,
SQLiteDefenseStore,
)
_BASE_TS = datetime(2026, 9, 12, 12, 0, 0, tzinfo=UTC)
def make_defense(
id: str = "defense-1", # shadows builtin deliberately: DefenseRecord field name
learner_id: str = "learner-1",
task_id: str = "task-1",
status: str = "in_progress",
created_at: datetime | None = None,
) -> DefenseRecord:
"""Canonical kwargs builder — tests override only what they assert on."""
return DefenseRecord(
id=id,
learner_id=learner_id,
task_id=task_id,
status=status,
created_at=created_at if created_at is not None else _BASE_TS,
)
def make_turn(
defense_id: str = "defense-1",
seq: int = 0,
role: str = "examiner",
text: str = "Walk me through your cache invalidation strategy.",
ts: datetime | None = None,
latency_ms: int | None = 240,
created_at: datetime | None = None,
) -> DefenseTurn:
return DefenseTurn(
defense_id=defense_id,
seq=seq,
role=role,
text=text,
ts=ts if ts is not None else _BASE_TS + timedelta(seconds=seq),
latency_ms=latency_ms,
created_at=created_at
if created_at is not None
else _BASE_TS + timedelta(seconds=seq),
)
@pytest.fixture
def store(tmp_path: Path) -> SQLiteDefenseStore:
s = SQLiteDefenseStore(db_path=tmp_path / "defenses.db")
yield s
s.close()
def test_start_append_finalize_get_roundtrip(store: SQLiteDefenseStore) -> None:
"""THE roundtrip of the defense lifecycle (REQ-3-006)."""
# start: a defense is born in_progress with empty signals.
started = store.start(make_defense())
assert started.status == "in_progress"
assert started.finished_at is None
assert started.integrity_signals == {}
assert started.created_at == _BASE_TS
assert started.created_at.tzinfo is UTC
# append: examiner + learner turns interleaved — append them OUT of
# seq order to prove get() orders by seq, not by insertion.
learner_a = make_turn(
seq=1, role="learner", text="I invalidate on write-ahead flush.", latency_ms=980
)
examiner_b = make_turn(
seq=2, role="examiner", text="Why not invalidate on read?", latency_ms=180
)
learner_c = make_turn(
seq=3, role="learner", text="Read-path misses were rare in my trace.", latency_ms=1100
)
store.append_turn("defense-1", make_turn(seq=0)) # first examiner question
store.append_turn("defense-1", learner_a)
store.append_turn("defense-1", examiner_b)
store.append_turn("defense-1", learner_c)
# finalize: seal with A-109 integrity signals (long pauses, off-scope).
finalize_started = datetime.now(UTC) # real wall clock, not the fixture
signals = {
"long_pauses": {"count": 2, "threshold_ms": 2000, "turns": [1, 3]},
"off_scope": {"count": 1, "turns": [3], "markers": ["unrelated tangent"]},
"verdict": "PASS_WITH_NOTES",
}
finalized = store.finalize("defense-1", signals)
assert finalized is not None
assert finalized.status == "finished"
assert finalized.finished_at is not None
assert finalized.finished_at.tzinfo is UTC
# finalize() does not attach turns — get() is the with-turns path.
assert finalized.turns == []
# get: full transcript in seq order + signals persisted.
fetched = store.get("defense-1")
assert fetched is not None
assert fetched.id == "defense-1"
assert fetched.learner_id == "learner-1"
assert fetched.task_id == "task-1"
assert fetched.status == "finished"
assert fetched.integrity_signals == signals
assert fetched.integrity_signals["long_pauses"]["turns"] == [1, 3] # nested JSON survives
assert fetched.created_at == _BASE_TS
assert fetched.created_at.tzinfo is UTC
assert fetched.finished_at is not None
assert fetched.finished_at >= finalize_started # stamped at finalize time
# THE ordered-transcript assertion.
assert [t.seq for t in fetched.turns] == [0, 1, 2, 3]
assert [t.role for t in fetched.turns] == [
"examiner",
"learner",
"examiner",
"learner",
]
assert fetched.turns[0].text == "Walk me through your cache invalidation strategy."
assert fetched.turns[1].latency_ms == 980
assert fetched.turns[2].latency_ms == 180
assert fetched.turns[3].latency_ms == 1100
# Per-turn ts contract: tz-aware UTC on read regardless of backend.
assert all(t.ts.tzinfo is UTC for t in fetched.turns)
assert all(t.created_at.tzinfo is UTC for t in fetched.turns)
assert [t.defense_id for t in fetched.turns] == ["defense-1"] * 4
def test_start_forces_in_progress_lifecycle(store: SQLiteDefenseStore) -> None:
"""The store owns the lifecycle: a smuggled finished status is
normalized away at birth — only finalize() may move a defense to
finished."""
smuggled = make_defense(status="finished")
smuggled.finished_at = _BASE_TS + timedelta(hours=1)
started = store.start(smuggled)
assert started.status == "in_progress"
assert started.finished_at is None
fetched = store.get("defense-1")
assert fetched is not None
assert fetched.status == "in_progress"
assert fetched.finished_at is None
def test_duplicate_defense_id_raises_integrity_error(
store: SQLiteDefenseStore,
) -> None:
"""Insert-only start (first-wins): the id is minted once per session;
a duplicate is a bug or lost race to surface, not swallow."""
store.start(make_defense(learner_id="learner-1"))
store.start(make_defense(id="defense-2", learner_id="learner-2"))
with pytest.raises(sa.exc.IntegrityError):
store.start(make_defense(learner_id="learner-3")) # same id again
# The original rows survived intact — nothing was overwritten.
first = store.get("defense-1")
assert first is not None
assert first.learner_id == "learner-1"
assert first.status == "in_progress"
second = store.get("defense-2")
assert second is not None
assert second.learner_id == "learner-2"
def test_append_turn_duplicate_seq_raises_integrity_error(
store: SQLiteDefenseStore,
) -> None:
"""A transcript turn must never silently vanish: (defense_id, seq) is
the PK, so a duplicate raises instead of overwriting."""
store.start(make_defense())
store.append_turn("defense-1", make_turn(seq=0))
store.append_turn("defense-1", make_turn(seq=1))
with pytest.raises(sa.exc.IntegrityError):
store.append_turn(
"defense-1",
make_turn(seq=1, role="learner", text="an impostor answer"),
)
# The stored turn is untouched — NOT upsert.
fetched = store.get("defense-1")
assert fetched is not None
assert len(fetched.turns) == 2
assert fetched.turns[1].role == "examiner"
assert fetched.turns[1].text != "an impostor answer"
def test_append_turn_unknown_defense_raises_integrity_error(
store: SQLiteDefenseStore,
) -> None:
"""FK enforced (foreign_keys=ON): an orphan turn is rejected, not
silently attached to a defense that does not exist."""
store.start(make_defense())
with pytest.raises(sa.exc.IntegrityError):
store.append_turn(
"defense-missing",
make_turn(defense_id="defense-missing", seq=0),
)
# And the turn did not land on the existing defense either.
fetched = store.get("defense-1")
assert fetched is not None
assert fetched.turns == []
def test_append_turn_identity_mismatch_raises_value_error(
store: SQLiteDefenseStore,
) -> None:
"""The defense_id argument is the write identity: a turn object
claiming another defense is a programming error — surfaced BEFORE
any DB round-trip."""
store.start(make_defense())
with pytest.raises(ValueError, match="does not match"):
store.append_turn(
"defense-1",
make_turn(defense_id="defense-other", seq=0),
)
fetched = store.get("defense-1")
assert fetched is not None
assert fetched.turns == []
def test_turn_model_rejects_negative_seq_and_unknown_role(store: SQLiteDefenseStore) -> None:
"""Model-level contracts (@validates hooks fire at construction)."""
with pytest.raises(ValueError, match="seq"):
make_turn(seq=-1)
with pytest.raises(ValueError, match="role"):
make_turn(role="proctor")
with pytest.raises(ValueError, match="text"):
make_turn(text="")
with pytest.raises(ValueError, match="latency_ms"):
make_turn(latency_ms=-5)
with pytest.raises(ValueError, match="defense_id"):
make_turn(defense_id="")
def test_append_turn_allows_missing_latency(store: SQLiteDefenseStore) -> None:
"""latency_ms is None until the endpoints instrument it (task
5-4-01) — None must roundtrip cleanly."""
store.start(make_defense())
store.append_turn(
"defense-1",
make_turn(seq=0, role="examiner", latency_ms=None),
)
fetched = store.get("defense-1")
assert fetched is not None
assert fetched.turns[0].latency_ms is None
def test_finalize_unknown_id_returns_none(store: SQLiteDefenseStore) -> None:
"""Documented unknown-id behavior: None, not a raise — the defense
endpoints (task 5-3-01) map this straight to 404."""
assert store.finalize("nobody", {"verdict": "PASS"}) is None
def test_finalize_is_latest_wins_on_refinalize(store: SQLiteDefenseStore) -> None:
"""A recomputed verdict replaces the stored one wholesale (mirrors
GradeStore.save): signals + finished_at are overwritten, status
just stays finished."""
store.start(make_defense())
first_signals = {"verdict": "FAIL", "long_pauses": {"count": 5}}
store.finalize("defense-1", first_signals)
better_signals = {
"verdict": "PASS",
"long_pauses": {"count": 1},
"off_scope": {"count": 0},
}
refinalized = store.finalize("defense-1", better_signals)
assert refinalized is not None
fetched = store.get("defense-1")
assert fetched is not None
assert fetched.status == "finished"
assert fetched.integrity_signals == better_signals # wholesale replace
assert fetched.integrity_signals["verdict"] == "PASS"
assert refinalized.finished_at is not None
assert fetched.finished_at == refinalized.finished_at # stamped anew
def test_finalize_rejects_non_dict_signals(store: SQLiteDefenseStore) -> None:
"""The signals column contract is a JSON OBJECT dict; None/str/list
would break every reader (Examiner feed, Proctor/Mentor)."""
store.start(make_defense())
with pytest.raises(ValueError, match="integrity_signals"):
store.finalize("defense-1", None) # type: ignore[arg-type]
fetched = store.get("defense-1")
assert fetched is not None
assert fetched.status == "in_progress" # untouched by the rejected call
def test_get_unknown_returns_none(store: SQLiteDefenseStore) -> None:
store.start(make_defense())
assert store.get("defense-1") is not None
assert store.get("defense-missing") is None
def test_empty_signals_until_finalize(store: SQLiteDefenseStore) -> None:
"""A-109: signals are {} until finalize — the in-progress transcript
is readable without any verdict present."""
store.start(make_defense())
store.append_turn("defense-1", make_turn(seq=0))
fetched = store.get("defense-1")
assert fetched is not None
assert fetched.status == "in_progress"
assert fetched.integrity_signals == {}
assert fetched.finished_at is None
assert len(fetched.turns) == 1
def test_list_for_learner_scoped_and_chronological(store: SQLiteDefenseStore) -> None:
# Created out of insertion order; list must come back chronological.
store.start(make_defense(id="defense-c", learner_id="learner-1",
created_at=_BASE_TS + timedelta(hours=2)))
store.start(make_defense(id="defense-a", learner_id="learner-1"))
store.start(make_defense(id="defense-b", learner_id="learner-2"))
# Turns on learner-1's defenses prove the list carries NONE of them.
store.append_turn("defense-c", make_turn(defense_id="defense-c", seq=0))
store.append_turn("defense-a", make_turn(defense_id="defense-a", seq=0))
defenses = store.list_for_learner("learner-1")
assert [d.id for d in defenses] == ["defense-a", "defense-c"]
assert all(d.learner_id == "learner-1" for d in defenses)
# WITHOUT turns: the list feed carries session headers only — get()
# is the with-turns path.
assert all(d.turns == [] for d in defenses)
# Headers intact: status + signals readable for the Proctor/Mentor feed.
assert all(d.status == "in_progress" for d in defenses)
assert all(d.integrity_signals == {} for d in defenses)
assert all(d.created_at.tzinfo is UTC for d in defenses)
# Other learners never leak.
assert [d.id for d in store.list_for_learner("learner-2")] == ["defense-b"]
assert store.list_for_learner("learner-missing") == []
def test_list_includes_finalized_with_signals(store: SQLiteDefenseStore) -> None:
"""The learner feed must surface finished defenses WITH their sealed
signals (the Proctor/Mentor cross-check reads exactly this)."""
store.start(make_defense())
signals = {"verdict": "PASS", "long_pauses": {"count": 0}}
store.finalize("defense-1", signals)
defenses = store.list_for_learner("learner-1")
assert len(defenses) == 1
assert defenses[0].status == "finished"
assert defenses[0].integrity_signals == signals
assert defenses[0].finished_at is not None
assert defenses[0].finished_at.tzinfo is UTC
def test_rows_are_detached_and_durable(store: SQLiteDefenseStore, tmp_path: Path) -> None:
"""Detached from any session: the API layer hands DefenseRecords
across layers; rows must survive the store that produced them
being closed, and a fresh store must see the same rows."""
store.start(make_defense())
store.append_turn("defense-1", make_turn(seq=0))
store.append_turn("defense-1", make_turn(seq=1, role="learner", text="My answer."))
store.finalize("defense-1", {"verdict": "PASS"})
fetched = store.get("defense-1")
assert fetched is not None
store.close()
# Usable post-close — no open-session ORM magic.
assert fetched.status == "finished"
assert fetched.integrity_signals["verdict"] == "PASS"
assert [t.text for t in fetched.turns] == [
"Walk me through your cache invalidation strategy.",
"My answer.",
]
# A fresh store on the same file sees the same rows (durability).
reopened = SQLiteDefenseStore(db_path=tmp_path / "defenses.db")
try:
again = reopened.get("defense-1")
assert again is not None
assert again.status == "finished"
assert again.integrity_signals == {"verdict": "PASS"}
assert [t.seq for t in again.turns] == [0, 1]
assert again.turns[1].latency_ms == 240
assert all(t.ts.tzinfo is UTC for t in again.turns)
finally:
reopened.close()
def test_pragmas_are_applied(store: SQLiteDefenseStore) -> None:
# Pragmas are per-connection; query through the store's engine so the
# connect hook (not a default sqlite3 connection) is what we inspect.
with store._engine.connect() as conn:
(journal_mode,) = conn.execute(sa.text("PRAGMA journal_mode")).one()
(synchronous,) = conn.execute(sa.text("PRAGMA synchronous")).one()
(foreign_keys,) = conn.execute(sa.text("PRAGMA foreign_keys")).one()
assert journal_mode == "wal"
# synchronous=NORMAL is 1 in SQLite's pragma numbering.
assert synchronous == 1
# The FK is actually enforced on SQLite (Postgres parity, D-027).
assert foreign_keys == 1
def test_two_defenses_same_learner_independent_transcripts(
store: SQLiteDefenseStore,
) -> None:
"""Scoped transcripts: two defenses never see each other's turns."""
store.start(make_defense(id="defense-a", task_id="task-1"))
store.start(make_defense(id="defense-b", task_id="task-2"))
# Both defenses reuse seq 0,1 — per-defense numbering.
for defense_id in ("defense-a", "defense-b"):
store.append_turn(defense_id, make_turn(defense_id=defense_id, seq=0))
store.append_turn(
defense_id,
make_turn(
defense_id=defense_id, seq=1, role="learner",
text=f"answer for {defense_id}",
),
)
a = store.get("defense-a")
b = store.get("defense-b")
assert a is not None and b is not None
assert [t.seq for t in a.turns] == [0, 1]
assert [t.text for t in a.turns] == [
"Walk me through your cache invalidation strategy.",
"answer for defense-a",
]
assert [t.text for t in b.turns] == [
"Walk me through your cache invalidation strategy.",
"answer for defense-b",
]
+101
View File
@@ -0,0 +1,101 @@
"""Per-turn latency instrumentation tests (Task 5-4-01, REQ-3-006, A-109).
Mock-based: asserts instrumentation PRESENCE and population (stt_ms / llm_ms /
tts_ms fields, per-turn latency_ms persisted, the budget constant defined) —
wall-clock against a real voice endpoint is a v0.4 acceptance criterion
(real STT/TTS deferred per GRILL CUT-1 / G-7).
"""
from __future__ import annotations
from pathlib import Path
import pytest
from fastapi.testclient import TestClient
from ai_service.agents.examiner import ExaminerAgent
from ai_service.config import Settings
from ai_service.grading.store import SQLiteGradeStore
from ai_service.llm.mock import MockProvider
from ai_service.main import create_app
from ai_service.telemetry.ingest import TraceIntegrityMap
from ai_service.telemetry.store import SQLiteTraceStore
from ai_service.variants.store import SQLiteVariantStore
from ai_service.voice.defense_store import SQLiteDefenseStore
from ai_service.voice.mock import MockVoiceProvider
#: A-109: the documented conversational budget (acceptance criterion for the
#: v0.4 real-voice probe; mock turns are near-instant so v0.3 asserts
#: instrumentation, not wall-clock).
DEFENSE_TURN_BUDGET_MS = 4_000
@pytest.fixture()
def client(tmp_path: Path) -> TestClient:
llm = MockProvider()
app = create_app(Settings(provider="mock", voice_provider="mock"))
app.state.provider = llm
app.state.trace_store = SQLiteTraceStore(db_path=tmp_path / "t.db")
app.state.variant_store = SQLiteVariantStore(db_path=tmp_path / "v.db")
app.state.grade_store = SQLiteGradeStore(db_path=tmp_path / "g.db")
app.state.trace_integrity = TraceIntegrityMap()
app.state.defense_store = SQLiteDefenseStore(db_path=tmp_path / "d.db")
app.state.voice_provider = MockVoiceProvider(["my answer"])
app.state.examiner_agent = ExaminerAgent(llm, Settings(provider="mock"))
with TestClient(app) as c:
yield c
class TestLatencyInstrumentation:
def test_budget_constant_defined(self) -> None:
"""The conversational budget is a named, documented constant (A-109)."""
assert DEFENSE_TURN_BUDGET_MS > 0
assert DEFENSE_TURN_BUDGET_MS <= 5_000 # conversational feel target
def test_answer_reports_per_phase_latency(self, client: TestClient) -> None:
start = client.post(
"/v1/defense/start",
json={"learner_id": "lat-learner", "task_id": "lat-task"},
).json()
resp = client.post(
f"/v1/defense/{start['defense_id']}/answer", data={"text": "answer"}
)
assert resp.status_code == 200
latency = resp.json()["turn_latency"]
assert latency["llm_ms"] is not None and latency["llm_ms"] >= 0
assert "stt_ms" in latency and "tts_ms" in latency
def test_audio_answer_populates_stt_ms(self, client: TestClient) -> None:
start = client.post(
"/v1/defense/start",
json={"learner_id": "lat-learner", "task_id": "lat-task"},
).json()
resp = client.post(
f"/v1/defense/{start['defense_id']}/answer",
files={"audio": ("a.wav", b"RIFF" + b"\x00" * 32, "audio/wav")},
)
latency = resp.json()["turn_latency"]
assert latency["stt_ms"] is not None and latency["stt_ms"] >= 0
def test_every_turn_persists_latency_ms(self, client: TestClient) -> None:
start = client.post(
"/v1/defense/start",
json={"learner_id": "lat-learner", "task_id": "lat-task"},
).json()
client.post(f"/v1/defense/{start['defense_id']}/answer", data={"text": "a"})
transcript = client.get(f"/v1/defense/{start['defense_id']}").json()
assert transcript["turns"]
for turn in transcript["turns"]:
assert "latency_ms" in turn
assert turn["latency_ms"] is not None or turn["role"] == "learner"
def test_mock_turns_within_budget(self, client: TestClient) -> None:
"""Mock turns must be near-instant — the budget holds trivially."""
start = client.post(
"/v1/defense/start",
json={"learner_id": "lat-learner", "task_id": "lat-task"},
).json()
resp = client.post(
f"/v1/defense/{start['defense_id']}/answer", data={"text": "a"}
).json()
assert resp["turn_latency"]["llm_ms"] < DEFENSE_TURN_BUDGET_MS
@@ -0,0 +1,118 @@
"""Voice layer tests (Task 5-1-01, REQ-3-006) — D-030 mock-first, zero network."""
from __future__ import annotations
import pytest
from ai_service.config import Settings
from ai_service.voice.browser import BROWSER_FALLBACK_DESCRIPTOR, MOCK_DESCRIPTOR
from ai_service.voice.factory import UnknownVoiceProviderError, voice_provider_from_settings
from ai_service.voice.mock import MockVoiceFailure, MockVoiceProvider, _tone_wav
class TestMockVoiceProvider:
async def test_transcribe_deterministic(self) -> None:
provider = MockVoiceProvider(["hello defense"])
a = await provider.transcribe(b"x" * 3200, "wav")
b = await provider.transcribe(b"x" * 3200, "wav")
assert a.text == b.text == "hello defense"
assert a.duration_ms == b.duration_ms
async def test_transcribe_canned_queue(self) -> None:
provider = MockVoiceProvider(["first answer", "second answer"])
first = await provider.transcribe(b"audio", "webm")
second = await provider.transcribe(b"audio", "webm")
assert first.text == "first answer"
assert second.text == "second answer"
async def test_transcribe_empty_audio_fails(self) -> None:
provider = MockVoiceProvider(["x"])
with pytest.raises(MockVoiceFailure, match="no audio"):
await provider.transcribe(b"", "wav")
async def test_transcribe_scripted_failure_mode(self) -> None:
provider = MockVoiceProvider(["FAIL"])
with pytest.raises(MockVoiceFailure, match="scripted STT failure"):
await provider.transcribe(b"audio", "wav")
async def test_synthesize_yields_nonempty_chunks(self) -> None:
provider = MockVoiceProvider()
chunks = [chunk async for chunk in provider.synthesize("question text")]
assert chunks
assert all(isinstance(c, bytes) and c for c in chunks)
assert provider.synthesize_calls == 1
async def test_synthesize_empty_text_fails(self) -> None:
provider = MockVoiceProvider()
with pytest.raises(MockVoiceFailure):
async for _ in provider.synthesize(""):
pass
def test_tone_wav_is_real_wav(self) -> None:
import io
import wave
raw = _tone_wav(duration_ms=100)
with wave.open(io.BytesIO(raw)) as w:
assert w.getnchannels() == 1
assert w.getsampwidth() == 2
assert w.getframerate() == 8000
async def test_identical_synthesize_calls_identical_bytes(self) -> None:
p1, p2 = MockVoiceProvider(), MockVoiceProvider()
c1 = b"".join([c async for c in p1.synthesize("same text")])
c2 = b"".join([c async for c in p2.synthesize("same text")])
assert c1 == c2
class TestFactory:
def test_default_is_mock(self) -> None:
provider = voice_provider_from_settings(Settings())
assert isinstance(provider, MockVoiceProvider)
def test_explicit_mock(self) -> None:
provider = voice_provider_from_settings(Settings(voice_provider="mock"))
assert isinstance(provider, MockVoiceProvider)
def test_browser_mode_selects_server_side_mock_for_text_fallback(self) -> None:
# Browser mode composes the same deterministic provider server-side;
# the descriptor tells the CLIENT to use native SR/TTS.
provider = voice_provider_from_settings(Settings(voice_provider="browser"))
assert isinstance(provider, MockVoiceProvider)
def test_real_server_stt_tts_rejected_as_v04_seam(self) -> None:
with pytest.raises(UnknownVoiceProviderError, match="v0.4"):
voice_provider_from_settings(Settings(voice_provider="openai-audio"))
def test_unknown_provider_rejected(self) -> None:
with pytest.raises(UnknownVoiceProviderError, match="unknown"):
voice_provider_from_settings(Settings(voice_provider="watson"))
class TestDescriptors:
def test_browser_fallback_descriptor(self) -> None:
assert BROWSER_FALLBACK_DESCRIPTOR.mode == "browser"
assert BROWSER_FALLBACK_DESCRIPTOR.sr_available
assert BROWSER_FALLBACK_DESCRIPTOR.tts_available
assert "SpeechRecognition" in BROWSER_FALLBACK_DESCRIPTOR.hint
def test_mock_descriptor(self) -> None:
assert MOCK_DESCRIPTOR.mode == "mock"
assert "v0.4" in MOCK_DESCRIPTOR.hint
class TestZeroNetwork:
def test_voice_package_never_imports_agents_or_api(self) -> None:
import ast
from pathlib import Path
pkg = Path(__file__).parents[2] / "ai_service" / "voice"
for py in pkg.glob("*.py"):
tree = ast.parse(py.read_text())
for node in ast.walk(tree):
if isinstance(node, ast.ImportFrom) and node.module:
assert not node.module.startswith("ai_service.agents"), py
assert not node.module.startswith("ai_service.api"), py
if node.level and node.module:
assert node.module.split(".")[-1] != "agents", py
assert node.module.split(".")[-1] != "api", py