DocDoeAI / app /services /adaptive_engine.py
asnannp's picture
deploy: sync backend to Space root (learn-lesson HF cache fix)
d5ee82b
Raw
History Blame Contribute Delete
20.2 kB
"""Adaptive engine orchestration: checkpoints, repair queue, daily planning.
DB-facing layer over the pure ``mastery_engine``. Every mutation is
ownership-scoped, idempotent where the student can double-submit, and stores a
human-readable reason so the product can always answer
"DocDoe chose this task because…".
"""
from __future__ import annotations
import logging
from datetime import datetime, timezone
from fastapi import HTTPException, status
from sqlalchemy import select
from sqlalchemy.exc import IntegrityError
from sqlalchemy.orm import Session
from app.models.learning_state import (
Chapter,
CheckpointEvent,
DailyTask,
RepairItem,
RevisionItem,
StudentProfileState,
TopicMastery,
UsageEvent,
)
from app.schemas.learning_state import (
CheckpointConsequenceResponse,
CheckpointMasteryOut,
CheckpointRequest,
PlanTodayResponse,
RepairItemOut,
)
from app.services import concept_graph
from app.services.learning_events import add_learning_event
from app.services.mastery_engine import (
ERROR_CATEGORY_COPY,
EvidenceEvent,
MasterySnapshot,
apply_evidence,
classify_error,
)
from app.services.today_plan_service import refresh_daily_plan_totals, replan_today
REPAIR_TASK_MINUTES = 10
logger = logging.getLogger(__name__)
def _now() -> datetime:
return datetime.now(timezone.utc)
def _owned_chapter_for_checkpoint(
db: Session,
*,
user_id: str,
payload: CheckpointRequest,
) -> Chapter | None:
if payload.chapter_id:
chapter = db.scalar(
select(Chapter).where(
Chapter.id == payload.chapter_id,
Chapter.user_id == user_id,
)
)
if chapter is not None:
return chapter
return db.scalar(
select(Chapter).where(
Chapter.user_id == user_id,
Chapter.catalog_id == payload.chapter_catalog_id,
)
)
def _sync_revision_item_from_checkpoint(
db: Session,
*,
user_id: str,
payload: CheckpointRequest,
chapter: Chapter | None,
next_review_at: datetime | None,
now: datetime,
) -> None:
if next_review_at is None:
return
client_item_id = f"mastery-review:{payload.concept_key}"[:180]
item = db.scalar(
select(RevisionItem).where(
RevisionItem.user_id == user_id,
RevisionItem.client_item_id == client_item_id,
)
)
due_status = "due" if next_review_at.date() <= now.date() else "pending"
if item is None:
db.add(
RevisionItem(
user_id=user_id,
client_item_id=client_item_id,
subject_id=(chapter.subject_id if chapter else None) or payload.subject_id,
chapter_id=(chapter.id if chapter else None) or payload.chapter_id,
mission_id=payload.mission_id,
topic_key=payload.concept_key,
title=payload.concept_label,
source_kind="checkpoint",
source_ref=payload.question_id,
status=due_status,
due_at=next_review_at,
item_data={
"chapter_catalog_id": payload.chapter_catalog_id,
"estimated_minutes": 10,
},
)
)
return
if item.status in {"completed", "resolved"}:
return
item.due_at = next_review_at
item.status = due_status
item.mission_id = item.mission_id or payload.mission_id
item.chapter_id = item.chapter_id or (chapter.id if chapter else payload.chapter_id)
item.subject_id = item.subject_id or (chapter.subject_id if chapter else payload.subject_id)
def repair_item_out(item: RepairItem) -> RepairItemOut:
return RepairItemOut(
id=item.id,
subject_id=item.subject_id,
chapter_id=item.chapter_id,
concept_key=item.concept_key,
concept_label=item.concept_label,
mission_id=item.mission_id,
error_category=item.error_category,
diagnosis=item.diagnosis,
recommended_activity=item.recommended_activity,
activity_prompt=item.activity_prompt,
priority=item.priority,
estimated_minutes=item.estimated_minutes,
status=item.status,
support_level=item.support_level,
failed_attempts=item.failed_attempts,
retry_result=item.retry_result,
mastery_recovered=item.mastery_recovered,
created_at=item.created_at,
resolved_at=item.resolved_at,
)
def _server_verified_correct(payload: CheckpointRequest) -> bool:
"""Recompute correctness where the server can do so deterministically.
MCQ answers are verifiable by comparison; free-text answers use the client's
deterministic rubric result (also recorded verbatim in the event payload for
audit). The server never silently trusts a claim it can cheaply verify.
"""
if payload.question_type == "mcq":
return payload.student_answer.strip().casefold() == payload.correct_answer.strip().casefold()
return payload.client_correct
def _open_repairs_for_concept(db: Session, user_id: str, concept_key: str) -> list[RepairItem]:
return list(
db.scalars(
select(RepairItem).where(
RepairItem.user_id == user_id,
RepairItem.concept_key == concept_key,
RepairItem.status.in_(("open", "escalated")),
)
)
)
def _apply_checkpoint_task_consequence(
db: Session,
*,
user_id: str,
payload: CheckpointRequest,
correct: bool,
repair_status: str,
now: datetime,
) -> bool:
"""Complete only an exact task whose required evidence this answer satisfies."""
if not payload.daily_task_id:
return False
task = db.scalar(
select(DailyTask).where(
DailyTask.id == payload.daily_task_id,
DailyTask.user_id == user_id,
)
)
if task is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Study task not found.")
if task.mission_id and task.mission_id != payload.mission_id:
raise HTTPException(
status_code=status.HTTP_409_CONFLICT,
detail={"code": "TASK_EVIDENCE_MISMATCH", "message": "This answer belongs to a different lesson task."},
)
evidence_completes = (
correct
and (
(task.task_type == "mistake_repair" and repair_status == "resolved")
or (task.task_type == "revision" and payload.kind in {"revision_recall", "transfer_check", "board_answer"})
or (task.task_type == "board_answer_practice" and payload.kind == "board_answer")
)
)
if not evidence_completes or task.status == "completed":
return False
task.status = "completed"
task.completed_at = now
task.task_metadata = {
**dict(task.task_metadata or {}),
"completion_evidence": {
"kind": "checkpoint_event",
"client_event_id": payload.client_event_id,
"question_id": payload.question_id,
},
}
refresh_daily_plan_totals(db, plan_id=task.daily_plan_id, user_id=user_id)
add_learning_event(
db,
user_id=user_id,
event_type="TASK_COMPLETED",
entity_type="daily_task",
entity_id=task.id,
idempotency_key=f"task-completed:checkpoint:{payload.client_event_id}:{task.id}",
subject_id=task.subject_id,
chapter_id=task.chapter_id,
topic_key=payload.concept_key,
event_data={"daily_plan_id": task.daily_plan_id, "evidence_type": "checkpoint_event"},
)
return True
def record_checkpoint(
db: Session,
*,
user_id: str,
payload: CheckpointRequest,
) -> CheckpointConsequenceResponse:
# Idempotency: a replayed client_event_id returns the stored consequence.
existing = db.scalar(
select(CheckpointEvent).where(
CheckpointEvent.user_id == user_id,
CheckpointEvent.client_event_id == payload.client_event_id,
)
)
if existing is not None:
stored = dict(existing.consequence or {})
stored["replayed"] = True
return CheckpointConsequenceResponse.model_validate(stored)
now = _now()
correct = _server_verified_correct(payload)
error_category: str | None = None
if not correct:
error_category = classify_error(
question_type=payload.question_type,
student_answer=payload.student_answer,
correct_answer=payload.correct_answer,
expected_keywords=payload.expected_keywords,
)
profile = db.scalar(select(StudentProfileState).where(StudentProfileState.user_id == user_id))
exam_date = profile.exam_date if profile is not None else None
chapter = _owned_chapter_for_checkpoint(db, user_id=user_id, payload=payload)
resolved_chapter_id = (chapter.id if chapter else None) or payload.chapter_id
resolved_subject_id = (chapter.subject_id if chapter else None) or payload.subject_id
# --- Repair lifecycle -------------------------------------------------
open_repairs = _open_repairs_for_concept(db, user_id, payload.concept_key)
repair_status: str = "none"
repair_item: RepairItem | None = None
if not correct and error_category is not None:
matching = next((item for item in open_repairs if item.error_category == error_category), None)
node = concept_graph.concept_for(payload.concept_key)
importance = node.exam_importance if node else 0.7
if matching is not None:
matching.failed_attempts += 1
matching.support_level = min(3, matching.support_level + 1)
matching.priority = round(min(2.5, matching.priority + 0.2), 2)
matching.status = "escalated"
matching.chapter_id = matching.chapter_id or resolved_chapter_id
matching.subject_id = matching.subject_id or resolved_subject_id
matching.mission_id = matching.mission_id or payload.mission_id
matching.evidence = {
**(matching.evidence or {}),
"last_failed_question_id": payload.question_id,
"last_failed_at": now.isoformat(),
}
repair_item = matching
repair_status = "escalated"
else:
copy = ERROR_CATEGORY_COPY.get(error_category, ERROR_CATEGORY_COPY["concept_misunderstanding"])
repair_item = RepairItem(
user_id=user_id,
subject_id=resolved_subject_id,
chapter_id=resolved_chapter_id,
concept_key=payload.concept_key,
concept_label=payload.concept_label,
mission_id=payload.mission_id,
error_category=error_category,
diagnosis=copy["diagnosis"],
recommended_activity=copy["activity"],
activity_prompt=copy["activity_prompt"],
priority=round(1.6 + importance * 0.4, 2),
estimated_minutes=REPAIR_TASK_MINUTES,
status="open",
support_level=1,
failed_attempts=1,
source_kind="checkpoint",
source_ref=payload.question_id,
evidence={
"question_id": payload.question_id,
"prompt": payload.prompt[:500],
"student_answer": payload.student_answer[:500],
"at": now.isoformat(),
},
)
db.add(repair_item)
repair_status = "created"
elif correct and open_repairs and (payload.attempt_index > 1 or payload.kind in {"transfer_check", "revision_recall", "board_answer"}):
# A successful retry/transfer closes the concept's open repairs (idempotent —
# already-resolved items are not touched again).
for item in open_repairs:
item.status = "resolved"
item.retry_result = "recovered"
item.resolved_at = now
repair_item = open_repairs[0]
repair_status = "resolved"
# --- Mastery evidence ---------------------------------------------------
mastery_row = db.scalar(
select(TopicMastery)
.where(TopicMastery.user_id == user_id, TopicMastery.topic_key == payload.concept_key)
.with_for_update()
)
snapshot = MasterySnapshot(
score=(mastery_row.score if mastery_row else None) or 0.0,
confidence=(mastery_row.confidence if mastery_row else None) or 0.0,
attempts_count=(mastery_row.attempts_count if mastery_row else None) or 0,
state=(mastery_row.last_result if mastery_row and mastery_row.last_result else "not_started"),
next_review_at=mastery_row.next_review_at if mastery_row else None,
evidence=dict(mastery_row.evidence or {}) if mastery_row else {},
)
# After the lifecycle above, the concept still has an open repair when one
# was just created/escalated, or pre-existing opens were not resolved here.
has_open_repair_after = repair_status in {"created", "escalated"} or (
repair_status != "resolved" and bool(open_repairs)
)
update = apply_evidence(
snapshot,
EvidenceEvent(
kind=payload.kind,
correct=correct,
at=now,
question_id=payload.question_id,
hint_used=payload.hint_used,
error_category=error_category,
time_spent_seconds=payload.time_spent_seconds,
source="guided_class",
),
exam_date=exam_date,
has_open_repair=has_open_repair_after,
)
if mastery_row is None:
mastery_row = TopicMastery(
user_id=user_id,
subject_id=resolved_subject_id,
chapter_id=resolved_chapter_id,
topic_key=payload.concept_key,
topic_label=payload.concept_label,
)
db.add(mastery_row)
mastery_row.subject_id = resolved_subject_id or mastery_row.subject_id
mastery_row.chapter_id = resolved_chapter_id or mastery_row.chapter_id
mastery_row.topic_label = payload.concept_label
mastery_row.score = snapshot.score
mastery_row.confidence = snapshot.confidence
mastery_row.attempts_count = snapshot.attempts_count
mastery_row.last_result = snapshot.state
mastery_row.next_review_at = snapshot.next_review_at
mastery_row.evidence = snapshot.evidence
_sync_revision_item_from_checkpoint(
db,
user_id=user_id,
payload=payload,
chapter=chapter,
next_review_at=snapshot.next_review_at,
now=now,
)
if repair_item is not None and repair_status == "resolved":
repair_item.mastery_recovered = snapshot.score >= 60.0
# --- Message the student can trust ---------------------------------------
if correct and repair_status == "resolved":
message = "Repair recovered — this concept is back on track and your plan will drop the repair task."
elif correct:
message = "Correct. This strengthens the concept's mastery and pushes its next revision further out."
elif repair_status == "escalated":
message = "Still stuck on the same kind of mistake — DocDoe raised the support level for this repair."
else:
message = "Not yet. DocDoe recorded the exact mistake and added a short repair to today's plan."
mastery_out = CheckpointMasteryOut(
before_score=update.before_score,
after_score=update.after_score,
before_state=update.before_state,
after_state=update.after_state,
confidence=update.confidence,
consecutive_success=update.consecutive_success,
next_review_at=update.next_review_at,
)
db.flush() # so repair_item.id exists for the stored consequence
consequence = CheckpointConsequenceResponse(
checkpoint_id="pending",
correct=correct,
replayed=False,
error_category=error_category,
diagnosis=repair_item.diagnosis if repair_item is not None and not correct else None,
recommended_activity=repair_item.recommended_activity if repair_item is not None and not correct else None,
activity_prompt=repair_item.activity_prompt if repair_item is not None and not correct else None,
repair_status=repair_status, # type: ignore[arg-type]
repair_item=repair_item_out(repair_item) if repair_item is not None else None,
mastery=mastery_out,
message=message,
)
event = CheckpointEvent(
user_id=user_id,
client_event_id=payload.client_event_id,
concept_key=payload.concept_key,
question_id=payload.question_id,
attempt_index=payload.attempt_index,
kind=payload.kind,
correct=correct,
error_category=error_category,
hint_used=payload.hint_used,
payload={
"question_type": payload.question_type,
"prompt": payload.prompt[:1000],
"student_answer": payload.student_answer[:1000],
"correct_answer": payload.correct_answer[:1000],
"client_correct": payload.client_correct,
"expected_keywords": payload.expected_keywords,
"time_spent_seconds": payload.time_spent_seconds,
},
consequence={},
)
db.add(event)
try:
db.flush()
except IntegrityError:
# A parallel tab won the unique-constraint race; return its consequence.
db.rollback()
winner = db.scalar(
select(CheckpointEvent).where(
CheckpointEvent.user_id == user_id,
CheckpointEvent.client_event_id == payload.client_event_id,
)
)
if winner is None:
raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail="Checkpoint conflict; retry.")
stored = dict(winner.consequence or {})
stored["replayed"] = True
return CheckpointConsequenceResponse.model_validate(stored)
consequence.checkpoint_id = event.id
event.consequence = consequence.model_dump(mode="json")
db.add(
UsageEvent(
user_id=user_id,
event_type="checkpoint_answered",
resource_type="checkpoint_event",
event_data={"checkpoint_id": event.id, "concept_key": payload.concept_key, "correct": correct},
)
)
add_learning_event(
db,
user_id=user_id,
event_type="ANSWER_CORRECTED" if correct else "ANSWER_SUBMITTED",
entity_type="checkpoint_event",
entity_id=event.id,
idempotency_key=f"checkpoint:{payload.client_event_id}",
subject_id=payload.subject_id,
chapter_id=payload.chapter_id,
topic_key=payload.concept_key,
event_data={"question_id": payload.question_id, "correct": correct, "repair_status": repair_status},
)
_apply_checkpoint_task_consequence(
db,
user_id=user_id,
payload=payload,
correct=correct,
repair_status=repair_status,
now=now,
)
db.commit()
try:
replan_today(
db,
user_id=user_id,
idempotency_key=f"checkpoint:{payload.client_event_id}",
reason="meaningful_evidence",
now=now,
)
except HTTPException as exc:
if exc.status_code != status.HTTP_409_CONFLICT:
raise
logger.info("today_plan_not_replanned user_id=%s reason=study_plan_required", user_id)
return consequence
# ---------------------------------------------------------------------------
# Compatibility
# ---------------------------------------------------------------------------
def plan_today(db: Session, *, user_id: str) -> PlanTodayResponse:
"""Compatibility wrapper for older imports.
The persisted Learning Engine in today_plan_service is the sole
executable planning authority.
"""
from app.services.today_plan_service import ensure_today_plan
return ensure_today_plan(db, user_id=user_id)