feat(P05): voice layer + DefenseStore (Wave 1)
Task 5-1-01: voice/ — VoiceProvider protocol (transcribe/synthesize, D-030 mirroring
LLMProvider), deterministic MockVoiceProvider (scripted STT queue, canned tone-WAV
TTS chunks, failure modes incl. empty audio), browser fallback descriptor (client
native SR/TTS), factory (mock default; browser; openai-audio REJECTED as a v0.4 seam
per CUT-1/G-7), config key AI_VOICE_PROVIDER + .env.example note. voice/ imports no
agents/api (AST-tested).
Task 5-1-03: DefenseStore (4th D-027 store; first FK family) — DefenseRecord +
DefenseTurn (ordered by (defense_id, seq)); start/append_turn/finalize/get/
list_for_learner; PRAGMA foreign_keys=ON for Postgres parity; integrity signals JSON
(A-109); store owns the finished transition.
34 voice tests green; ruff clean.
---ci---
phase: 5
milestone: v0.3
status: execute
requirements: {covered: [REQ-3-006], partial: []}
---/ci---
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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",
|
||||
]
|
||||
@@ -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'"
|
||||
)
|
||||
@@ -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()
|
||||
@@ -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",
|
||||
]
|
||||
@@ -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
|
||||
Reference in New Issue
Block a user