From 6223b27cb1b035496a5ad259396f5a9938900f30 Mon Sep 17 00:00:00 2001 From: Ryan Zhou Date: Mon, 24 Aug 2026 11:40:12 -0500 Subject: [PATCH 1/5] feat: define acknowledged decision trace evidence --- src/sleeper_manager/domain/planning.py | 8 +- .../persistence/acknowledgements.py | 250 +++++++++ tests/unit/test_acknowledgement_queries.py | 487 ++++++++++++++++++ tests/unit/test_planning_models.py | 29 ++ 4 files changed, 771 insertions(+), 3 deletions(-) create mode 100644 src/sleeper_manager/persistence/acknowledgements.py create mode 100644 tests/unit/test_acknowledgement_queries.py diff --git a/src/sleeper_manager/domain/planning.py b/src/sleeper_manager/domain/planning.py index f4797e8..23962aa 100644 --- a/src/sleeper_manager/domain/planning.py +++ b/src/sleeper_manager/domain/planning.py @@ -255,11 +255,13 @@ def __post_init__(self) -> None: _require_text(self.provenance, "Acknowledged provenance") if self.action is AcknowledgedAction.LOCK: if self.slot_index is None or self.slot_position is None: - raise PlanningStateError("Acknowledged locks require slot evidence") - if self.slot_index < 0: + if self.reconciled: + raise PlanningStateError("Acknowledged locks require slot evidence") + elif self.slot_index < 0: raise PlanningStateError("Acknowledged lock slot indices must be non-negative") if self.accepted_fantasy_score is None or not isfinite(self.accepted_fantasy_score): - raise PlanningStateError("Acknowledged locks require a finite accepted score") + if self.reconciled: + raise PlanningStateError("Acknowledged locks require a finite accepted score") return if ( self.slot_index is not None diff --git a/src/sleeper_manager/persistence/acknowledgements.py b/src/sleeper_manager/persistence/acknowledgements.py new file mode 100644 index 0000000..6089c82 --- /dev/null +++ b/src/sleeper_manager/persistence/acknowledgements.py @@ -0,0 +1,250 @@ +from __future__ import annotations + +import json +from collections import defaultdict +from collections.abc import Mapping, Sequence +from dataclasses import dataclass, replace +from datetime import datetime +from math import isfinite +from typing import Any + +from sleeper_manager.domain.planning import AcknowledgedAction, AcknowledgedDecisionEvidence + +LOCK_IN_DECISION_TYPE = "lock_in" +ACKNOWLEDGEMENT_TRACE_SCHEMA_VERSION = 1 +ACKNOWLEDGEMENT_PROVENANCE = "repository_acknowledgement_v1" + +ACKNOWLEDGED_DECISIONS_QUERY = """ +SELECT + r.recommendation_id AS recommendation_id, + r.player_id AS player_id, + r.game_id AS game_id, + r.status AS recommendation_status, + r.acknowledged_action AS recommendation_action, + r.acknowledged_at AS recommendation_acknowledged_at, + r.trace_json AS trace_json, + a.action AS acknowledgement_action, + a.acknowledged_at AS acknowledged_at +FROM recommendations r +JOIN acknowledgements a ON a.recommendation_id = r.recommendation_id +WHERE r.league_id = ? + AND r.fantasy_week = ? + AND r.decision_type = ? +""" + +ACKNOWLEDGED_DECISIONS_INDEX_SQL = """ +CREATE INDEX IF NOT EXISTS recommendations_league_week_decision_status_idx +ON recommendations (league_id, fantasy_week, decision_type, status) +""" + +_LOCK_FRAGMENT_FIELDS = ("slot_index", "slot_position", "accepted_fantasy_score") +_ACTION_BY_STORED_VALUE = { + "locked": AcknowledgedAction.LOCK, + "passed": AcknowledgedAction.PASS, +} + + +class AcknowledgementQueryError(RuntimeError): + """Raised when stored acknowledgement rows cannot form identity-bearing evidence.""" + + +@dataclass(frozen=True, slots=True) +class AcknowledgementRawRow: + recommendation_id: object + player_id: object + game_id: object + recommendation_status: object + recommendation_action: object + recommendation_acknowledged_at: object + trace_json: object + acknowledgement_action: object + acknowledged_at: object + + +def raw_row_from_mapping(row: Mapping[str, Any]) -> AcknowledgementRawRow: + return AcknowledgementRawRow( + recommendation_id=row.get("recommendation_id"), + player_id=row.get("player_id"), + game_id=row.get("game_id"), + recommendation_status=row.get("recommendation_status"), + recommendation_action=row.get("recommendation_action"), + recommendation_acknowledged_at=row.get("recommendation_acknowledged_at"), + trace_json=row.get("trace_json"), + acknowledgement_action=row.get("acknowledgement_action"), + acknowledged_at=row.get("acknowledged_at"), + ) + + +def raw_row_from_sequence(row: Sequence[object]) -> AcknowledgementRawRow: + return AcknowledgementRawRow( + recommendation_id=row[0], + player_id=row[1], + game_id=row[2], + recommendation_status=row[3], + recommendation_action=row[4], + recommendation_acknowledged_at=row[5], + trace_json=row[6], + acknowledgement_action=row[7], + acknowledged_at=row[8], + ) + + +def decode_acknowledged_decisions( + rows: Sequence[AcknowledgementRawRow], + *, + as_of: datetime, +) -> tuple[AcknowledgedDecisionEvidence, ...]: + if as_of.tzinfo is None: + raise AcknowledgementQueryError("as_of must be timezone-aware") + decoded: list[AcknowledgedDecisionEvidence] = [] + for row in rows: + evidence = _decode_row(row) + if evidence.decided_at > as_of: + continue + decoded.append(evidence) + return _reconcile(_order(decoded)) + + +def _decode_row(row: AcknowledgementRawRow) -> AcknowledgedDecisionEvidence: + decision_id = _require_identity(row.recommendation_id, "recommendation ID") + player_id = _require_identity(row.player_id, "player ID") + game_id = _require_identity(row.game_id, "game ID") + action = _ACTION_BY_STORED_VALUE.get(_text_or_none(row.acknowledgement_action) or "") + if action is None: + raise AcknowledgementQueryError("unknown acknowledgement action") + decided_at = _parse_aware(row.acknowledged_at, label="acknowledgement timestamp") + slot_index, slot_position, accepted_score, trace_ok = _parse_trace(row.trace_json, action) + duplicates_ok = _duplicates_match(row, action, decided_at) + return AcknowledgedDecisionEvidence( + decision_id=decision_id, + player_id=player_id, + game_id=game_id, + action=action, + decided_at=decided_at, + provenance=ACKNOWLEDGEMENT_PROVENANCE, + slot_index=slot_index, + slot_position=slot_position, + accepted_fantasy_score=accepted_score, + reconciled=trace_ok and duplicates_ok, + ) + + +def _duplicates_match( + row: AcknowledgementRawRow, + action: AcknowledgedAction, + decided_at: datetime, +) -> bool: + status = _text_or_none(row.recommendation_status) + stored_action = _ACTION_BY_STORED_VALUE.get(_text_or_none(row.recommendation_action) or "") + try: + stored_at = _parse_aware( + row.recommendation_acknowledged_at, + label="recommendation acknowledgement timestamp", + ) + except AcknowledgementQueryError: + return False + return status == "acknowledged" and stored_action is action and stored_at == decided_at + + +def _parse_trace( + trace_json: object, + action: AcknowledgedAction, +) -> tuple[int | None, str | None, float | None, bool]: + if not isinstance(trace_json, str): + return None, None, None, False + try: + payload = json.loads(trace_json) + except json.JSONDecodeError: + return None, None, None, False + if not isinstance(payload, dict): + return None, None, None, False + fragment = payload.get("acknowledgement") + if not isinstance(fragment, dict): + return None, None, None, False + if fragment.get("schema_version") != ACKNOWLEDGEMENT_TRACE_SCHEMA_VERSION or isinstance( + fragment.get("schema_version"), bool + ): + return None, None, None, False + slot_index = _parse_slot_index(fragment.get("slot_index")) + slot_position = _parse_slot_position(fragment.get("slot_position")) + accepted_score = _parse_accepted_score(fragment.get("accepted_fantasy_score")) + has_lock_member = any(key in fragment for key in _LOCK_FRAGMENT_FIELDS) + if action is AcknowledgedAction.PASS: + return None, None, None, not has_lock_member + complete = slot_index is not None and slot_position is not None and accepted_score is not None + return slot_index, slot_position, accepted_score, complete + + +def _parse_slot_index(value: object) -> int | None: + if isinstance(value, bool) or not isinstance(value, int) or value < 0: + return None + return value + + +def _parse_slot_position(value: object) -> str | None: + if not isinstance(value, str) or not value.strip(): + return None + return value.strip().upper() + + +def _parse_accepted_score(value: object) -> float | None: + if isinstance(value, bool) or not isinstance(value, int | float): + return None + score = float(value) + if not isfinite(score): + return None + return score + + +def _require_identity(value: object, label: str) -> str: + text = _text_or_none(value) + if text is None: + raise AcknowledgementQueryError(f"acknowledgement identity is missing {label}") + return text + + +def _text_or_none(value: object) -> str | None: + if value is None or isinstance(value, bool): + return None + text = str(value).strip() + return text or None + + +def _parse_aware(value: object, *, label: str) -> datetime: + if not isinstance(value, str) or not value.strip(): + raise AcknowledgementQueryError(f"{label} is missing") + try: + parsed = datetime.fromisoformat(value) + except ValueError as error: + raise AcknowledgementQueryError(f"{label} is invalid") from error + if parsed.tzinfo is None: + raise AcknowledgementQueryError(f"{label} must be timezone-aware") + return parsed + + +def _order( + evidence: Sequence[AcknowledgedDecisionEvidence], +) -> tuple[AcknowledgedDecisionEvidence, ...]: + return tuple(sorted(evidence, key=lambda item: (item.decided_at, item.decision_id))) + + +def _reconcile( + evidence: Sequence[AcknowledgedDecisionEvidence], +) -> tuple[AcknowledgedDecisionEvidence, ...]: + player_games: dict[tuple[str, str], set[str]] = defaultdict(set) + locked_slots: dict[int, set[str]] = defaultdict(set) + locked_players: dict[str, set[str]] = defaultdict(set) + for item in evidence: + player_games[item.player_id, item.game_id].add(item.decision_id) + if item.action is AcknowledgedAction.LOCK: + locked_players[item.player_id].add(item.decision_id) + if item.slot_index is not None: + locked_slots[item.slot_index].add(item.decision_id) + conflicted: set[str] = set() + for ids in (*player_games.values(), *locked_slots.values(), *locked_players.values()): + if len(ids) > 1: + conflicted.update(ids) + return tuple( + replace(item, reconciled=False) if item.decision_id in conflicted else item + for item in evidence + ) diff --git a/tests/unit/test_acknowledgement_queries.py b/tests/unit/test_acknowledgement_queries.py new file mode 100644 index 0000000..deddd4e --- /dev/null +++ b/tests/unit/test_acknowledgement_queries.py @@ -0,0 +1,487 @@ +from __future__ import annotations + +import json +from datetime import UTC, datetime, timedelta, timezone +from math import inf, nan + +import pytest + +from sleeper_manager.domain.planning import AcknowledgedAction, AcknowledgedDecisionEvidence +from sleeper_manager.persistence.acknowledgements import ( + ACKNOWLEDGEMENT_PROVENANCE, + ACKNOWLEDGEMENT_TRACE_SCHEMA_VERSION, + LOCK_IN_DECISION_TYPE, + AcknowledgementQueryError, + AcknowledgementRawRow, + decode_acknowledged_decisions, +) + +AS_OF = datetime(2026, 1, 7, 18, tzinfo=UTC) +EASTERN = timezone(timedelta(hours=-5)) + + +def _row( + *, + recommendation_id: object = "rec-1", + player_id: object = "p1", + game_id: object = "g1", + recommendation_status: object = "acknowledged", + recommendation_action: object = "locked", + recommendation_acknowledged_at: object = AS_OF.isoformat(), + trace_json: object | None = None, + acknowledgement_action: object = "locked", + acknowledged_at: object = AS_OF.isoformat(), +) -> AcknowledgementRawRow: + if trace_json is None: + trace_json = _lock_trace() + return AcknowledgementRawRow( + recommendation_id=recommendation_id, + player_id=player_id, + game_id=game_id, + recommendation_status=recommendation_status, + recommendation_action=recommendation_action, + recommendation_acknowledged_at=recommendation_acknowledged_at, + trace_json=trace_json, + acknowledgement_action=acknowledgement_action, + acknowledged_at=acknowledged_at, + ) + + +def _lock_trace( + *, + slot_index: object = 2, + slot_position: object = "UTIL", + accepted_fantasy_score: object = 34.7, + schema_version: object = 1, + omit: tuple[str, ...] = (), + extra_top_level: dict[str, object] | None = None, + extra_fragment: dict[str, object] | None = None, +) -> str: + fragment: dict[str, object] = { + "schema_version": schema_version, + "slot_index": slot_index, + "slot_position": slot_position, + "accepted_fantasy_score": accepted_fantasy_score, + } + for key in omit: + fragment.pop(key, None) + if extra_fragment: + fragment.update(extra_fragment) + payload: dict[str, object] = {"acknowledgement": fragment} + if extra_top_level: + payload.update(extra_top_level) + return json.dumps(payload) + + +def _pass_trace( + *, + schema_version: object = 1, + extra_fragment: dict[str, object] | None = None, + extra_top_level: dict[str, object] | None = None, +) -> str: + fragment: dict[str, object] = {"schema_version": schema_version} + if extra_fragment: + fragment.update(extra_fragment) + payload: dict[str, object] = {"acknowledgement": fragment} + if extra_top_level: + payload.update(extra_top_level) + return json.dumps(payload) + + +def test_canonical_constants_are_stable() -> None: + assert LOCK_IN_DECISION_TYPE == "lock_in" + assert ACKNOWLEDGEMENT_TRACE_SCHEMA_VERSION == 1 + assert ACKNOWLEDGEMENT_PROVENANCE == "repository_acknowledgement_v1" + + +def test_naive_as_of_is_rejected_before_rows_are_processed() -> None: + with pytest.raises(AcknowledgementQueryError, match="timezone-aware"): + decode_acknowledged_decisions((_row(),), as_of=datetime(2026, 1, 7, 18)) + + +@pytest.mark.parametrize( + "field,value", + [ + ("recommendation_id", ""), + ("recommendation_id", " "), + ("recommendation_id", None), + ("player_id", ""), + ("player_id", None), + ("game_id", ""), + ("game_id", " "), + ("game_id", None), + ], +) +def test_missing_identity_fields_are_storage_integrity_failures( + field: str, + value: object, +) -> None: + with pytest.raises(AcknowledgementQueryError, match="identity"): + decode_acknowledged_decisions((_row(**{field: value}),), as_of=AS_OF) + + +def test_unknown_authoritative_action_is_a_storage_integrity_failure() -> None: + with pytest.raises(AcknowledgementQueryError, match="action"): + decode_acknowledged_decisions((_row(acknowledgement_action="lock"),), as_of=AS_OF) + + +def test_naive_authoritative_timestamp_is_a_storage_integrity_failure() -> None: + naive = "2026-01-07T18:00:00" + with pytest.raises(AcknowledgementQueryError, match="timezone-aware"): + decode_acknowledged_decisions((_row(acknowledged_at=naive),), as_of=AS_OF) + + +def test_invalid_authoritative_timestamp_is_a_storage_integrity_failure() -> None: + with pytest.raises(AcknowledgementQueryError, match="timestamp"): + decode_acknowledged_decisions((_row(acknowledged_at="not-a-time"),), as_of=AS_OF) + + +def test_point_in_time_includes_equivalent_offsets_and_excludes_newer() -> None: + included_eastern = AS_OF.astimezone(EASTERN).isoformat() + later = (AS_OF + timedelta(seconds=1)).isoformat() + earlier = (AS_OF - timedelta(minutes=1)).isoformat() + rows = ( + _row( + recommendation_id="rec-later", + acknowledged_at=later, + recommendation_acknowledged_at=later, + ), + _row( + recommendation_id="rec-boundary", + acknowledged_at=included_eastern, + recommendation_acknowledged_at=included_eastern, + ), + _row( + recommendation_id="rec-earlier", + acknowledged_at=earlier, + recommendation_acknowledged_at=earlier, + ), + ) + + result = decode_acknowledged_decisions(rows, as_of=AS_OF) + + assert [item.decision_id for item in result] == ["rec-earlier", "rec-boundary"] + + +def test_result_order_uses_datetime_comparison_not_iso_strings() -> None: + later_utc = datetime(2026, 1, 7, 19, tzinfo=UTC) + earlier_offset = datetime(2026, 1, 7, 12, tzinfo=timezone(timedelta(hours=-8))) + rows = ( + _row( + recommendation_id="rec-z", + acknowledged_at=later_utc.isoformat(), + recommendation_acknowledged_at=later_utc.isoformat(), + ), + _row( + recommendation_id="rec-a", + acknowledged_at=earlier_offset.isoformat(), + recommendation_acknowledged_at=earlier_offset.isoformat(), + ), + ) + + result = decode_acknowledged_decisions(rows, as_of=datetime(2026, 1, 8, tzinfo=UTC)) + + assert [item.decision_id for item in result] == ["rec-z", "rec-a"] + assert result[0].decided_at == later_utc + assert result[1].decided_at == earlier_offset + + +def test_canonical_lock_trace_decodes_to_reconciled_evidence() -> None: + result = decode_acknowledged_decisions((_row(trace_json=_lock_trace()),), as_of=AS_OF) + + assert result == ( + AcknowledgedDecisionEvidence( + decision_id="rec-1", + player_id="p1", + game_id="g1", + action=AcknowledgedAction.LOCK, + decided_at=AS_OF, + provenance=ACKNOWLEDGEMENT_PROVENANCE, + slot_index=2, + slot_position="UTIL", + accepted_fantasy_score=34.7, + reconciled=True, + ), + ) + + +def test_canonical_pass_trace_decodes_to_reconciled_evidence() -> None: + row = _row( + acknowledgement_action="passed", + recommendation_action="passed", + trace_json=_pass_trace(extra_top_level={"policy": {"ignored": True}}), + ) + + result = decode_acknowledged_decisions((row,), as_of=AS_OF) + + assert result == ( + AcknowledgedDecisionEvidence( + decision_id="rec-1", + player_id="p1", + game_id="g1", + action=AcknowledgedAction.PASS, + decided_at=AS_OF, + provenance=ACKNOWLEDGEMENT_PROVENANCE, + reconciled=True, + ), + ) + + +@pytest.mark.parametrize( + "trace_json", + ["{", "[1]", json.dumps({"acknowledgement": []}), json.dumps({"other": {}}), "null"], +) +def test_malformed_trace_yields_unreconciled_identity_bearing_evidence(trace_json: str) -> None: + result = decode_acknowledged_decisions((_row(trace_json=trace_json),), as_of=AS_OF) + + assert len(result) == 1 + assert result[0].decision_id == "rec-1" + assert result[0].action is AcknowledgedAction.LOCK + assert result[0].reconciled is False + assert result[0].slot_index is None + assert result[0].slot_position is None + assert result[0].accepted_fantasy_score is None + + +def test_unsupported_trace_version_is_unreconciled() -> None: + result = decode_acknowledged_decisions( + (_row(trace_json=_lock_trace(schema_version=2)),), + as_of=AS_OF, + ) + + assert result[0].reconciled is False + + +@pytest.mark.parametrize("omit", [("slot_index",), ("slot_position",), ("accepted_fantasy_score",)]) +def test_lock_missing_a_required_field_is_unreconciled(omit: tuple[str, ...]) -> None: + result = decode_acknowledged_decisions( + (_row(trace_json=_lock_trace(omit=omit)),), + as_of=AS_OF, + ) + + assert result[0].reconciled is False + if "slot_index" in omit: + assert result[0].slot_index is None + if "slot_position" in omit: + assert result[0].slot_position is None + if "accepted_fantasy_score" in omit: + assert result[0].accepted_fantasy_score is None + + +@pytest.mark.parametrize("slot_index", [True, False, -1, 1.5, "2", None]) +def test_invalid_slot_index_is_unreconciled(slot_index: object) -> None: + result = decode_acknowledged_decisions( + (_row(trace_json=_lock_trace(slot_index=slot_index)),), + as_of=AS_OF, + ) + + assert result[0].reconciled is False + assert result[0].slot_index is None + + +def test_non_negative_integer_slot_index_is_valid() -> None: + result = decode_acknowledged_decisions( + (_row(trace_json=_lock_trace(slot_index=0)),), + as_of=AS_OF, + ) + + assert result[0].reconciled is True + assert result[0].slot_index == 0 + + +@pytest.mark.parametrize("slot_position", ["", " "]) +def test_blank_slot_position_is_unreconciled(slot_position: str) -> None: + result = decode_acknowledged_decisions( + (_row(trace_json=_lock_trace(slot_position=slot_position)),), + as_of=AS_OF, + ) + + assert result[0].reconciled is False + assert result[0].slot_position is None + + +def test_lowercase_slot_position_is_normalized() -> None: + result = decode_acknowledged_decisions( + (_row(trace_json=_lock_trace(slot_position="util")),), + as_of=AS_OF, + ) + + assert result[0].reconciled is True + assert result[0].slot_position == "UTIL" + + +@pytest.mark.parametrize("score", [nan, inf, -inf, True, False, "34.7", None]) +def test_invalid_accepted_score_is_unreconciled(score: object) -> None: + result = decode_acknowledged_decisions( + (_row(trace_json=_lock_trace(accepted_fantasy_score=score)),), + as_of=AS_OF, + ) + + assert result[0].reconciled is False + assert result[0].accepted_fantasy_score is None + + +@pytest.mark.parametrize("score", [0, 0.0, -3.5, 12]) +def test_finite_accepted_scores_including_negatives_are_valid(score: float | int) -> None: + result = decode_acknowledged_decisions( + (_row(trace_json=_lock_trace(accepted_fantasy_score=score)),), + as_of=AS_OF, + ) + + assert result[0].reconciled is True + assert result[0].accepted_fantasy_score == float(score) + + +@pytest.mark.parametrize( + "extra_fragment", + [ + {"slot_index": 1}, + {"slot_position": "G"}, + {"accepted_fantasy_score": 12.0}, + ], +) +def test_pass_with_lock_only_members_is_unreconciled(extra_fragment: dict[str, object]) -> None: + row = _row( + acknowledgement_action="passed", + recommendation_action="passed", + trace_json=_pass_trace(extra_fragment=extra_fragment), + ) + + result = decode_acknowledged_decisions((row,), as_of=AS_OF) + + assert result[0].action is AcknowledgedAction.PASS + assert result[0].reconciled is False + assert result[0].slot_index is None + assert result[0].slot_position is None + assert result[0].accepted_fantasy_score is None + + +def test_duplicated_recommendation_fields_that_disagree_are_unreconciled() -> None: + mismatched = _row( + recommendation_status="pending", + recommendation_action="passed", + recommendation_acknowledged_at=(AS_OF + timedelta(minutes=5)).isoformat(), + ) + + result = decode_acknowledged_decisions((mismatched,), as_of=AS_OF) + + assert result[0].reconciled is False + assert result[0].action is AcknowledgedAction.LOCK + assert result[0].decided_at == AS_OF + assert result[0].slot_index == 2 + + +def test_decoder_never_returns_token_hashes() -> None: + row = _row( + trace_json=_lock_trace(extra_top_level={"token_hash": "secret"}), + ) + + dumped = repr(decode_acknowledged_decisions((row,), as_of=AS_OF)) + + assert "secret" not in dumped + assert "token_hash" not in dumped + + +def test_distinct_recommendation_ids_for_the_same_player_game_conflict() -> None: + rows = ( + _row(recommendation_id="rec-a", trace_json=_lock_trace()), + _row( + recommendation_id="rec-b", + acknowledgement_action="locked", + recommendation_action="locked", + trace_json=_lock_trace(slot_index=3, slot_position="G"), + ), + ) + + result = decode_acknowledged_decisions(rows, as_of=AS_OF) + + assert {item.decision_id for item in result} == {"rec-a", "rec-b"} + assert all(item.reconciled is False for item in result) + + +def test_lock_and_pass_for_the_same_player_game_conflict() -> None: + rows = ( + _row(recommendation_id="rec-lock"), + _row( + recommendation_id="rec-pass", + acknowledgement_action="passed", + recommendation_action="passed", + trace_json=_pass_trace(), + ), + ) + + result = decode_acknowledged_decisions(rows, as_of=AS_OF) + + assert all(item.reconciled is False for item in result) + + +def test_two_locks_for_the_same_slot_conflict() -> None: + rows = ( + _row(recommendation_id="rec-a", player_id="p1", game_id="g1", trace_json=_lock_trace()), + _row( + recommendation_id="rec-b", + player_id="p2", + game_id="g2", + trace_json=_lock_trace(), + ), + ) + + result = decode_acknowledged_decisions(rows, as_of=AS_OF) + + assert all(item.reconciled is False for item in result) + + +def test_two_locks_for_the_same_player_conflict() -> None: + rows = ( + _row(recommendation_id="rec-a", game_id="g1", trace_json=_lock_trace(slot_index=0)), + _row( + recommendation_id="rec-b", + game_id="g2", + trace_json=_lock_trace(slot_index=1, slot_position="G"), + ), + ) + + result = decode_acknowledged_decisions(rows, as_of=AS_OF) + + assert all(item.reconciled is False for item in result) + + +def test_passes_for_different_games_and_later_lock_are_valid() -> None: + earlier = (AS_OF - timedelta(hours=2)).isoformat() + later = (AS_OF - timedelta(hours=1)).isoformat() + rows = ( + _row( + recommendation_id="rec-pass-g1", + game_id="g1", + acknowledgement_action="passed", + recommendation_action="passed", + acknowledged_at=earlier, + recommendation_acknowledged_at=earlier, + trace_json=_pass_trace(), + ), + _row( + recommendation_id="rec-pass-g2", + game_id="g2", + acknowledgement_action="passed", + recommendation_action="passed", + acknowledged_at=earlier, + recommendation_acknowledged_at=earlier, + trace_json=_pass_trace(), + ), + _row( + recommendation_id="rec-lock-g3", + game_id="g3", + acknowledged_at=later, + recommendation_acknowledged_at=later, + trace_json=_lock_trace(), + ), + ) + + result = decode_acknowledged_decisions(rows, as_of=AS_OF) + + assert [item.reconciled for item in result] == [True, True, True] + assert [item.action for item in result] == [ + AcknowledgedAction.PASS, + AcknowledgedAction.PASS, + AcknowledgedAction.LOCK, + ] diff --git a/tests/unit/test_planning_models.py b/tests/unit/test_planning_models.py index b879d78..bb24de1 100644 --- a/tests/unit/test_planning_models.py +++ b/tests/unit/test_planning_models.py @@ -279,3 +279,32 @@ def test_acknowledged_decision_preserves_reconciliation_state() -> None: reconciled=False, ) assert evidence.reconciled is False + + +def test_unreconciled_lock_may_omit_incomplete_slot_evidence() -> None: + evidence = AcknowledgedDecisionEvidence( + decision_id="rec-5", + player_id="p1", + game_id="g1", + action=AcknowledgedAction.LOCK, + decided_at=NOW - timedelta(hours=1), + provenance="repository-query", + reconciled=False, + ) + assert evidence.slot_index is None + assert evidence.slot_position is None + assert evidence.accepted_fantasy_score is None + assert evidence.reconciled is False + + +def test_reconciled_lock_still_requires_complete_slot_evidence() -> None: + with pytest.raises(PlanningStateError, match="slot evidence"): + AcknowledgedDecisionEvidence( + decision_id="rec-6", + player_id="p1", + game_id="g1", + action=AcknowledgedAction.LOCK, + decided_at=NOW - timedelta(hours=1), + provenance="repository-query", + reconciled=True, + ) From e6d7db2fb7c31e05a7a281a83f2ef3d853aff6eb Mon Sep 17 00:00:00 2001 From: Ryan Zhou Date: Mon, 24 Aug 2026 11:44:38 -0500 Subject: [PATCH 2/5] feat: load acknowledged team-week decisions --- .../0002_acknowledged_team_week_decisions.sql | 2 + .../persistence/async_sqlite.py | 14 + src/sleeper_manager/persistence/base.py | 17 + src/sleeper_manager/persistence/d1.py | 44 ++ src/sleeper_manager/persistence/sqlite.py | 26 + tests/unit/test_acknowledgement_queries.py | 480 ++++++++++++++++++ tests/unit/test_d1.py | 9 +- 7 files changed, 591 insertions(+), 1 deletion(-) create mode 100644 infra/cloudflare/migrations/0002_acknowledged_team_week_decisions.sql diff --git a/infra/cloudflare/migrations/0002_acknowledged_team_week_decisions.sql b/infra/cloudflare/migrations/0002_acknowledged_team_week_decisions.sql new file mode 100644 index 0000000..ac9f267 --- /dev/null +++ b/infra/cloudflare/migrations/0002_acknowledged_team_week_decisions.sql @@ -0,0 +1,2 @@ +CREATE INDEX IF NOT EXISTS recommendations_league_week_decision_status_idx +ON recommendations (league_id, fantasy_week, decision_type, status); diff --git a/src/sleeper_manager/persistence/async_sqlite.py b/src/sleeper_manager/persistence/async_sqlite.py index c024696..de5e974 100644 --- a/src/sleeper_manager/persistence/async_sqlite.py +++ b/src/sleeper_manager/persistence/async_sqlite.py @@ -1,6 +1,7 @@ from datetime import datetime from pathlib import Path +from sleeper_manager.domain.planning import AcknowledgedDecisionEvidence from sleeper_manager.persistence.base import ( AcknowledgementAction, AcknowledgementResult, @@ -79,3 +80,16 @@ async def record_lock_acknowledgement( async def is_locked(self, recommendation_id: str) -> bool: return self._repository.is_locked(recommendation_id) + + async def load_acknowledged_decisions( + self, + league_id: str, + fantasy_week: int, + *, + as_of: datetime, + ) -> tuple[AcknowledgedDecisionEvidence, ...]: + return self._repository.load_acknowledged_decisions( + league_id, + fantasy_week, + as_of=as_of, + ) diff --git a/src/sleeper_manager/persistence/base.py b/src/sleeper_manager/persistence/base.py index b978c58..3c027f7 100644 --- a/src/sleeper_manager/persistence/base.py +++ b/src/sleeper_manager/persistence/base.py @@ -4,6 +4,7 @@ from typing import Protocol from sleeper_manager.domain.nba import DataQualityState +from sleeper_manager.domain.planning import AcknowledgedDecisionEvidence @dataclass(frozen=True, slots=True) @@ -181,6 +182,14 @@ def record_lock_acknowledgement( def is_locked(self, recommendation_id: str) -> bool: ... + def load_acknowledged_decisions( + self, + league_id: str, + fantasy_week: int, + *, + as_of: datetime, + ) -> tuple[AcknowledgedDecisionEvidence, ...]: ... + class AsyncStateRepository(Protocol): async def initialize(self) -> None: ... @@ -224,3 +233,11 @@ async def record_lock_acknowledgement( ) -> None: ... async def is_locked(self, recommendation_id: str) -> bool: ... + + async def load_acknowledged_decisions( + self, + league_id: str, + fantasy_week: int, + *, + as_of: datetime, + ) -> tuple[AcknowledgedDecisionEvidence, ...]: ... diff --git a/src/sleeper_manager/persistence/d1.py b/src/sleeper_manager/persistence/d1.py index cd2f867..35604c8 100644 --- a/src/sleeper_manager/persistence/d1.py +++ b/src/sleeper_manager/persistence/d1.py @@ -5,6 +5,15 @@ from typing import Any from sleeper_manager.domain.nba import DataQualityState +from sleeper_manager.domain.planning import AcknowledgedDecisionEvidence +from sleeper_manager.persistence.acknowledgements import ( + ACKNOWLEDGED_DECISIONS_QUERY, + LOCK_IN_DECISION_TYPE, + AcknowledgementQueryError, + decode_acknowledged_decisions, + raw_row_from_mapping, + raw_row_from_sequence, +) from sleeper_manager.persistence.base import ( AcknowledgementAction, AcknowledgementOutcome, @@ -103,6 +112,9 @@ player_id TEXT NOT NULL, acknowledged_at TEXT NOT NULL ); + +CREATE INDEX IF NOT EXISTS recommendations_league_week_decision_status_idx +ON recommendations (league_id, fantasy_week, decision_type, status); """ @@ -123,6 +135,15 @@ async def _first(self, query: str, *params: object) -> dict[str, Any] | None: row = await self._statement(query, params).first() return dict(row) if isinstance(row, Mapping) else None + async def _all(self, query: str, *params: object) -> Sequence[object]: + result = await self._statement(query, params).all() + if not isinstance(result, Mapping): + raise AcknowledgementQueryError("unexpected D1 result envelope") + rows = result.get("results") + if isinstance(rows, str | bytes) or not isinstance(rows, Sequence): + raise AcknowledgementQueryError("unexpected D1 result envelope") + return rows + async def _run(self, query: str, *params: object) -> Any: return await self._statement(query, params).run() @@ -493,3 +514,26 @@ async def is_locked(self, recommendation_id: str) -> bool: ) is not None ) + + async def load_acknowledged_decisions( + self, + league_id: str, + fantasy_week: int, + *, + as_of: datetime, + ) -> tuple[AcknowledgedDecisionEvidence, ...]: + rows = await self._all( + ACKNOWLEDGED_DECISIONS_QUERY, + league_id, + fantasy_week, + LOCK_IN_DECISION_TYPE, + ) + decoded_rows = [] + for row in rows: + if isinstance(row, Mapping): + decoded_rows.append(raw_row_from_mapping(row)) + elif isinstance(row, Sequence) and not isinstance(row, str | bytes): + decoded_rows.append(raw_row_from_sequence(row)) + else: + raise AcknowledgementQueryError("unexpected D1 result envelope") + return decode_acknowledged_decisions(tuple(decoded_rows), as_of=as_of) diff --git a/src/sleeper_manager/persistence/sqlite.py b/src/sleeper_manager/persistence/sqlite.py index fa9d333..7bf5d83 100644 --- a/src/sleeper_manager/persistence/sqlite.py +++ b/src/sleeper_manager/persistence/sqlite.py @@ -5,6 +5,14 @@ from pathlib import Path from sleeper_manager.domain.nba import DataQualityState +from sleeper_manager.domain.planning import AcknowledgedDecisionEvidence +from sleeper_manager.persistence.acknowledgements import ( + ACKNOWLEDGED_DECISIONS_INDEX_SQL, + ACKNOWLEDGED_DECISIONS_QUERY, + LOCK_IN_DECISION_TYPE, + decode_acknowledged_decisions, + raw_row_from_sequence, +) from sleeper_manager.persistence.base import ( AcknowledgementAction, AcknowledgementOutcome, @@ -134,6 +142,24 @@ def initialize(self) -> None: ) """ ) + connection.execute(ACKNOWLEDGED_DECISIONS_INDEX_SQL) + + def load_acknowledged_decisions( + self, + league_id: str, + fantasy_week: int, + *, + as_of: datetime, + ) -> tuple[AcknowledgedDecisionEvidence, ...]: + with self._connect() as connection: + rows = connection.execute( + ACKNOWLEDGED_DECISIONS_QUERY, + (league_id, fantasy_week, LOCK_IN_DECISION_TYPE), + ).fetchall() + return decode_acknowledged_decisions( + tuple(raw_row_from_sequence(row) for row in rows), + as_of=as_of, + ) @staticmethod def _recommendation(row: tuple[object, ...]) -> RecommendationRecord: diff --git a/tests/unit/test_acknowledgement_queries.py b/tests/unit/test_acknowledgement_queries.py index deddd4e..1c97ac5 100644 --- a/tests/unit/test_acknowledgement_queries.py +++ b/tests/unit/test_acknowledgement_queries.py @@ -1,10 +1,14 @@ from __future__ import annotations +import asyncio import json +import sqlite3 from datetime import UTC, datetime, timedelta, timezone from math import inf, nan +from pathlib import Path import pytest +from test_d1 import FakeD1 from sleeper_manager.domain.planning import AcknowledgedAction, AcknowledgedDecisionEvidence from sleeper_manager.persistence.acknowledgements import ( @@ -15,6 +19,17 @@ AcknowledgementRawRow, decode_acknowledged_decisions, ) +from sleeper_manager.persistence.async_sqlite import AsyncSQLiteStateRepository +from sleeper_manager.persistence.base import ( + AcknowledgementAction, + AcknowledgementOutcome, + AcknowledgementResult, + ActionTokenRecord, + RecommendationRecord, +) +from sleeper_manager.persistence.d1 import D1_SCHEMA, D1StateRepository +from sleeper_manager.persistence.sqlite import SQLiteStateRepository +from sleeper_manager.persistence.tokens import hash_action_token AS_OF = datetime(2026, 1, 7, 18, tzinfo=UTC) EASTERN = timezone(timedelta(hours=-5)) @@ -485,3 +500,468 @@ def test_passes_for_different_games_and_later_lock_are_valid() -> None: AcknowledgedAction.PASS, AcknowledgedAction.LOCK, ] + + +class _Backend: + def __init__(self, kind: str, tmp_path: Path) -> None: + self.kind = kind + if kind == "sqlite": + self.path = tmp_path / "state.db" + self.repo: SQLiteStateRepository | D1StateRepository = SQLiteStateRepository(self.path) + self.repo.initialize() + self.database: FakeD1 | None = None + else: + self.path = tmp_path / "d1.db" + self.database = FakeD1() + asyncio.run(self.database.exec(D1_SCHEMA)) + self.repo = D1StateRepository(self.database) + + def call(self, method: str, *args: object, **kwargs: object) -> object: + result = getattr(self.repo, method)(*args, **kwargs) + if asyncio.iscoroutine(result): + return asyncio.run(result) + return result + + def load( + self, league_id: str = "league-1", week: int = 1, *, as_of: datetime = AS_OF + ) -> tuple[AcknowledgedDecisionEvidence, ...]: + loaded = self.call("load_acknowledged_decisions", league_id, week, as_of=as_of) + assert isinstance(loaded, tuple) + return loaded + + def execute(self, sql: str, params: tuple[object, ...] = ()) -> sqlite3.Cursor: + if self.kind == "sqlite": + assert isinstance(self.repo, SQLiteStateRepository) + with self.repo._connect() as connection: + cursor = connection.execute(sql, params) + connection.commit() + return cursor + assert self.database is not None + cursor = self.database.connection.execute(sql, params) + self.database.connection.commit() + return cursor + + +@pytest.fixture(params=["sqlite", "d1"]) +def backend(request: pytest.FixtureRequest, tmp_path: Path) -> _Backend: + return _Backend(request.param, tmp_path) + + +def _recommendation(**overrides: object) -> RecommendationRecord: + values: dict[str, object] = { + "recommendation_id": "rec-1", + "idempotency_key": "idemp-1", + "league_id": "league-1", + "fantasy_week": 1, + "player_id": "p1", + "game_id": "g1", + "decision_type": LOCK_IN_DECISION_TYPE, + "title": "Lock", + "message": "Lock the player", + "deadline": AS_OF + timedelta(hours=1), + "policy_version": "v1", + "created_at": AS_OF - timedelta(hours=1), + "trace_json": _lock_trace(), + } + values.update(overrides) + return RecommendationRecord(**values) # type: ignore[arg-type] + + +def _acknowledge( + backend: _Backend, + record: RecommendationRecord, + action: AcknowledgementAction, + *, + token: str, + at: datetime, +) -> AcknowledgementOutcome: + backend.call("create_recommendation", record) + backend.call( + "create_action_token", + ActionTokenRecord( + token_hash=hash_action_token(token), + recommendation_id=record.recommendation_id, + action=action, + created_at=record.created_at, + expires_at=record.deadline or at, + ), + ) + result = backend.call("consume_action_token", hash_action_token(token), action, at) + assert isinstance(result, AcknowledgementResult) + return result.outcome + + +def test_empty_league_week_returns_no_acknowledgements(backend: _Backend) -> None: + assert backend.load() == () + + +def test_league_and_week_are_isolated(backend: _Backend) -> None: + _acknowledge( + backend, + _recommendation(), + AcknowledgementAction.LOCKED, + token="lock-token", + at=AS_OF - timedelta(minutes=1), + ) + _acknowledge( + backend, + _recommendation( + recommendation_id="rec-other-league", + idempotency_key="idemp-other-league", + league_id="league-2", + ), + AcknowledgementAction.LOCKED, + token="other-league", + at=AS_OF - timedelta(minutes=1), + ) + _acknowledge( + backend, + _recommendation( + recommendation_id="rec-other-week", + idempotency_key="idemp-other-week", + fantasy_week=2, + game_id="g2", + ), + AcknowledgementAction.LOCKED, + token="other-week", + at=AS_OF - timedelta(minutes=1), + ) + + loaded = backend.load("league-1", 1) + assert [item.decision_id for item in loaded] == ["rec-1"] + + +def test_as_of_includes_boundary_and_excludes_future(backend: _Backend) -> None: + _acknowledge( + backend, + _recommendation(), + AcknowledgementAction.LOCKED, + token="boundary", + at=AS_OF, + ) + _acknowledge( + backend, + _recommendation( + recommendation_id="rec-future", + idempotency_key="idemp-future", + game_id="g2", + ), + AcknowledgementAction.LOCKED, + token="future", + at=AS_OF + timedelta(seconds=1), + ) + + loaded = backend.load(as_of=AS_OF) + assert [item.decision_id for item in loaded] == ["rec-1"] + assert loaded[0].decided_at == AS_OF + + +def test_canonical_lock_and_pass_round_trip(backend: _Backend) -> None: + decided_lock = AS_OF - timedelta(hours=2) + decided_pass = AS_OF - timedelta(hours=1) + _acknowledge( + backend, + _recommendation(trace_json=_lock_trace()), + AcknowledgementAction.LOCKED, + token="lock", + at=decided_lock, + ) + _acknowledge( + backend, + _recommendation( + recommendation_id="rec-pass", + idempotency_key="idemp-pass", + player_id="p2", + game_id="g2", + trace_json=_pass_trace(), + ), + AcknowledgementAction.PASSED, + token="pass", + at=decided_pass, + ) + + loaded = backend.load() + assert loaded == ( + AcknowledgedDecisionEvidence( + decision_id="rec-1", + player_id="p1", + game_id="g1", + action=AcknowledgedAction.LOCK, + decided_at=decided_lock, + provenance=ACKNOWLEDGEMENT_PROVENANCE, + slot_index=2, + slot_position="UTIL", + accepted_fantasy_score=34.7, + ), + AcknowledgedDecisionEvidence( + decision_id="rec-pass", + player_id="p2", + game_id="g2", + action=AcknowledgedAction.PASS, + decided_at=decided_pass, + provenance=ACKNOWLEDGEMENT_PROVENANCE, + ), + ) + + +def test_placeholder_and_weekly_lineup_records_are_excluded(backend: _Backend) -> None: + _acknowledge( + backend, + _recommendation( + recommendation_id="rec-placeholder", + idempotency_key="idemp-placeholder", + decision_type="placeholder_lock_in", + ), + AcknowledgementAction.LOCKED, + token="placeholder", + at=AS_OF - timedelta(minutes=1), + ) + _acknowledge( + backend, + _recommendation( + recommendation_id="rec-lineup", + idempotency_key="idemp-lineup", + player_id="p2", + decision_type="weekly_lineup", + trace_json=_pass_trace(), + ), + AcknowledgementAction.PASSED, + token="lineup", + at=AS_OF - timedelta(minutes=1), + ) + + assert backend.load() == () + + +def test_duplicate_token_replay_returns_one_constraint(backend: _Backend) -> None: + record = _recommendation() + first = _acknowledge( + backend, + record, + AcknowledgementAction.LOCKED, + token="once", + at=AS_OF - timedelta(minutes=1), + ) + replay = backend.call( + "consume_action_token", + hash_action_token("once"), + AcknowledgementAction.LOCKED, + AS_OF, + ) + + assert first is AcknowledgementOutcome.APPLIED + assert replay.outcome is AcknowledgementOutcome.ALREADY_USED # type: ignore[union-attr] + assert [item.decision_id for item in backend.load()] == ["rec-1"] + + +def test_malformed_trace_loads_as_unreconciled_evidence(backend: _Backend) -> None: + _acknowledge( + backend, + _recommendation(trace_json="{"), + AcknowledgementAction.LOCKED, + token="bad-trace", + at=AS_OF - timedelta(minutes=1), + ) + + loaded = backend.load() + assert len(loaded) == 1 + assert loaded[0].reconciled is False + assert loaded[0].decision_id == "rec-1" + + +def test_direct_sql_duplicate_column_mismatch_is_unreconciled(backend: _Backend) -> None: + # Public writes refuse this inconsistency; mutate the duplicated recommendation copy. + _acknowledge( + backend, + _recommendation(), + AcknowledgementAction.LOCKED, + token="mismatch", + at=AS_OF - timedelta(minutes=1), + ) + backend.execute( + """ + UPDATE recommendations + SET status = ?, acknowledged_action = ? + WHERE recommendation_id = ? + """, + ("pending", "passed", "rec-1"), + ) + + loaded = backend.load() + assert loaded[0].reconciled is False + assert loaded[0].action is AcknowledgedAction.LOCK + + +def test_sqlite_index_leading_columns(tmp_path: Path) -> None: + repository = SQLiteStateRepository(tmp_path / "state.db") + repository.initialize() + with repository._connect() as connection: + columns = [ + row[2] + for row in connection.execute( + "PRAGMA index_info(recommendations_league_week_decision_status_idx)" + ) + ] + assert columns[:3] == ["league_id", "fantasy_week", "decision_type"] + assert columns[-1] == "status" + + +def test_async_sqlite_matches_sync_after_reopen(tmp_path: Path) -> None: + path = tmp_path / "state.db" + sync = SQLiteStateRepository(path) + sync.initialize() + backend = _Backend("sqlite", tmp_path) + backend.repo = sync + backend.path = path + _acknowledge( + backend, + _recommendation(), + AcknowledgementAction.LOCKED, + token="durable", + at=AS_OF - timedelta(minutes=1), + ) + expected = sync.load_acknowledged_decisions("league-1", 1, as_of=AS_OF) + + reopened = AsyncSQLiteStateRepository(path) + actual = asyncio.run(reopened.load_acknowledged_decisions("league-1", 1, as_of=AS_OF)) + assert actual == expected + + +def test_d1_all_returns_every_joined_row() -> None: + database = FakeD1() + asyncio.run(database.exec(D1_SCHEMA)) + repository = D1StateRepository(database) + backend = _Backend("d1", Path("/tmp")) + backend.database = database + backend.repo = repository + _acknowledge( + backend, + _recommendation(), + AcknowledgementAction.LOCKED, + token="one", + at=AS_OF - timedelta(hours=1), + ) + _acknowledge( + backend, + _recommendation( + recommendation_id="rec-2", + idempotency_key="idemp-2", + player_id="p2", + game_id="g2", + trace_json=_pass_trace(), + ), + AcknowledgementAction.PASSED, + token="two", + at=AS_OF - timedelta(minutes=1), + ) + + loaded = asyncio.run(repository.load_acknowledged_decisions("league-1", 1, as_of=AS_OF)) + assert [item.decision_id for item in loaded] == ["rec-1", "rec-2"] + + +def test_d1_binds_league_week_and_decision_type() -> None: + database = FakeD1() + asyncio.run(database.exec(D1_SCHEMA)) + repository = D1StateRepository(database) + asyncio.run(repository.load_acknowledged_decisions("league-9", 4, as_of=AS_OF)) + assert database.last_bound == ("league-9", 4, LOCK_IN_DECISION_TYPE) + + +def test_d1_unexpected_result_envelope_raises() -> None: + class _Broken: + def prepare(self, query: str) -> object: + class _Statement: + def bind(self, *params: object) -> _Statement: + return self + + async def all(self) -> object: + return {"results": "not-rows"} + + return _Statement() + + repository = D1StateRepository(_Broken()) + with pytest.raises(AcknowledgementQueryError, match="envelope"): + asyncio.run(repository.load_acknowledged_decisions("league-1", 1, as_of=AS_OF)) + + +def test_migration_adds_index_without_rewriting_phase3_rows(tmp_path: Path) -> None: + path = tmp_path / "migrated.db" + connection = sqlite3.connect(path) + root = Path(__file__).resolve().parents[2] + connection.executescript( + (root / "infra/cloudflare/migrations/0001_phase3.sql").read_text(encoding="utf-8") + ) + connection.execute( + """ + INSERT INTO recommendations ( + recommendation_id, idempotency_key, league_id, fantasy_week, player_id, game_id, + decision_type, title, message, deadline, policy_version, created_at, status, + acknowledged_action, acknowledged_at, trace_json + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + """, + ( + "rec-phase3", + "idemp-phase3", + "league-1", + 1, + "p1", + "g1", + "placeholder_lock_in", + "Lock", + "placeholder", + AS_OF.isoformat(), + "v1", + AS_OF.isoformat(), + "acknowledged", + "locked", + AS_OF.isoformat(), + _lock_trace(), + ), + ) + connection.execute( + """ + INSERT INTO action_tokens ( + token_hash, recommendation_id, action, created_at, expires_at, used_at + ) VALUES (?, ?, ?, ?, ?, ?) + """, + ( + "token-hash", + "rec-phase3", + "locked", + AS_OF.isoformat(), + AS_OF.isoformat(), + AS_OF.isoformat(), + ), + ) + connection.execute( + """ + INSERT INTO acknowledgements ( + acknowledgement_id, recommendation_id, action, acknowledged_at, token_hash + ) VALUES (?, ?, ?, ?, ?) + """, + ("ack-phase3", "rec-phase3", "locked", AS_OF.isoformat(), "token-hash"), + ) + connection.commit() + connection.executescript( + (root / "infra/cloudflare/migrations/0002_acknowledged_team_week_decisions.sql").read_text( + encoding="utf-8" + ) + ) + remaining = connection.execute( + "SELECT decision_type, trace_json FROM recommendations WHERE recommendation_id = ?", + ("rec-phase3",), + ).fetchone() + columns = [ + row[2] + for row in connection.execute( + "PRAGMA index_info(recommendations_league_week_decision_status_idx)" + ) + ] + connection.close() + + assert remaining is not None + assert remaining[0] == "placeholder_lock_in" + assert columns[:3] == ["league_id", "fantasy_week", "decision_type"] + + repository = SQLiteStateRepository(path) + assert repository.load_acknowledged_decisions("league-1", 1, as_of=AS_OF) == () diff --git a/tests/unit/test_d1.py b/tests/unit/test_d1.py index 688119f..2a452a3 100644 --- a/tests/unit/test_d1.py +++ b/tests/unit/test_d1.py @@ -25,7 +25,9 @@ def __init__(self, database: "FakeD1", query: str, params: tuple[Any, ...] = ()) self.params = params def bind(self, *params: object) -> "FakeStatement": - return FakeStatement(self.database, self.query, params) + bound = FakeStatement(self.database, self.query, params) + self.database.last_bound = params + return bound def execute(self) -> sqlite3.Cursor: return self.database.connection.execute(self.query, self.params) @@ -40,11 +42,16 @@ async def first(self) -> dict[str, Any] | None: row = cursor.fetchone() return dict(row) if row is not None else None + async def all(self) -> dict[str, Any]: + cursor = self.execute() + return {"results": [dict(row) for row in cursor.fetchall()]} + class FakeD1: def __init__(self) -> None: self.connection = sqlite3.connect(":memory:") self.connection.row_factory = sqlite3.Row + self.last_bound: tuple[object, ...] = () def prepare(self, query: str) -> FakeStatement: return FakeStatement(self, query) From fe9c20424a9d384201832f8e8f6ed0fd1f0ce898 Mon Sep 17 00:00:00 2001 From: Ryan Zhou Date: Mon, 24 Aug 2026 11:47:30 -0500 Subject: [PATCH 3/5] feat: source live acknowledgement constraints --- .../workflows/planning_collection.py | 10 +- tests/unit/test_live_planning_smoke.py | 6 +- tests/unit/test_planning_collection.py | 171 +++++++++++++++++- 3 files changed, 178 insertions(+), 9 deletions(-) diff --git a/src/sleeper_manager/workflows/planning_collection.py b/src/sleeper_manager/workflows/planning_collection.py index 3ab9c35..5be4761 100644 --- a/src/sleeper_manager/workflows/planning_collection.py +++ b/src/sleeper_manager/workflows/planning_collection.py @@ -56,10 +56,10 @@ async def players(self) -> Mapping[str, Mapping[str, Any]]: ... class AcknowledgementSource(Protocol): - """Supplies durable manager decisions; the real repository query lands in Task 5.2.""" + """Supplies durable manager decisions from the acknowledgement repository.""" - async def load( - self, league_id: str, week: int, *, as_of: datetime + async def load_acknowledged_decisions( + self, league_id: str, fantasy_week: int, *, as_of: datetime ) -> tuple[AcknowledgedDecisionEvidence, ...]: ... @@ -141,7 +141,9 @@ def tick() -> datetime: decision_time=decision_time, ) acknowledgements = ( - await acknowledgement_source.load(profile.league_id, week_window.week, as_of=decision_time) + await acknowledgement_source.load_acknowledged_decisions( + profile.league_id, week_window.week, as_of=decision_time + ) if acknowledgement_source is not None else () ) diff --git a/tests/unit/test_live_planning_smoke.py b/tests/unit/test_live_planning_smoke.py index db16f7b..fc7ad04 100644 --- a/tests/unit/test_live_planning_smoke.py +++ b/tests/unit/test_live_planning_smoke.py @@ -63,9 +63,9 @@ def load(self, target, *, before: datetime) -> HistoricalFeatureSlice: # noqa: class _NoAcknowledgements: - """Replaced by the repository query in Task 5.2.""" - - async def load(self, league_id: str, week: int, *, as_of: datetime) -> tuple: # noqa: ANN401 + async def load_acknowledged_decisions( + self, league_id: str, fantasy_week: int, *, as_of: datetime + ) -> tuple: # noqa: ANN401 return () diff --git a/tests/unit/test_planning_collection.py b/tests/unit/test_planning_collection.py index 58bc037..edc515f 100644 --- a/tests/unit/test_planning_collection.py +++ b/tests/unit/test_planning_collection.py @@ -1,4 +1,5 @@ from datetime import UTC, datetime, timedelta +from pathlib import Path import pytest @@ -27,6 +28,19 @@ PlanningReasonCode, ) from sleeper_manager.domain.scoring import ScoringPolicy +from sleeper_manager.persistence.acknowledgements import ( + ACKNOWLEDGEMENT_PROVENANCE, + LOCK_IN_DECISION_TYPE, + AcknowledgementQueryError, +) +from sleeper_manager.persistence.async_sqlite import AsyncSQLiteStateRepository +from sleeper_manager.persistence.base import ( + AcknowledgementAction, + AcknowledgementOutcome, + ActionTokenRecord, + RecommendationRecord, +) +from sleeper_manager.persistence.tokens import hash_action_token from sleeper_manager.projections.live_baseline import LiveProjectionTarget from sleeper_manager.workflows.planning_collection import ( PlanningCollectionError, @@ -437,8 +451,10 @@ def __init__(self, records: tuple[AcknowledgedDecisionEvidence, ...]) -> None: self.records = records self.requested: tuple[str, int] | None = None - async def load(self, league_id: str, week: int, *, as_of: datetime): - self.requested = (league_id, week) + async def load_acknowledged_decisions( + self, league_id: str, fantasy_week: int, *, as_of: datetime + ): + self.requested = (league_id, fantasy_week) return self.records @@ -465,6 +481,157 @@ async def run() -> None: } +def _lock_in_recommendation(**overrides: object) -> RecommendationRecord: + values: dict[str, object] = { + "recommendation_id": "rec-live-1", + "idempotency_key": "idemp-live-1", + "league_id": "league-1", + "fantasy_week": 1, + "player_id": "p1", + "game_id": "g1", + "decision_type": LOCK_IN_DECISION_TYPE, + "title": "Decision", + "message": "Decide", + "deadline": NOW + timedelta(hours=1), + "policy_version": "v1", + "created_at": NOW - timedelta(hours=3), + "trace_json": ( + '{"acknowledgement":{"schema_version":1,' + '"slot_index":0,"slot_position":"PG","accepted_fantasy_score":24}}' + ), + } + values.update(overrides) + return RecommendationRecord(**values) # type: ignore[arg-type] + + +async def _seed_acknowledgement( + repository: AsyncSQLiteStateRepository, + record: RecommendationRecord, + action: AcknowledgementAction, + *, + token: str, + at: datetime, +) -> None: + await repository.create_recommendation(record) + await repository.create_action_token( + ActionTokenRecord( + token_hash=hash_action_token(token), + recommendation_id=record.recommendation_id, + action=action, + created_at=record.created_at, + expires_at=record.deadline or at, + ) + ) + result = await repository.consume_action_token(hash_action_token(token), action, at) + assert result.outcome is AcknowledgementOutcome.APPLIED + + +def test_async_sqlite_pass_survives_reopen_into_live_state(tmp_path: Path) -> None: + path = tmp_path / "state.db" + writer = AsyncSQLiteStateRepository(path) + asyncio_run(writer.initialize()) + asyncio_run( + _seed_acknowledgement( + writer, + _lock_in_recommendation( + trace_json='{"acknowledgement":{"schema_version":1}}', + ), + AcknowledgementAction.PASSED, + token="pass-live", + at=NOW - timedelta(hours=2), + ) + ) + source = AsyncSQLiteStateRepository(path) + + async def run() -> object: + return await _collect(_nba(), _RecordingProjections(), acknowledgement_source=source) + + evidence = asyncio_run(run()) + state = build_live_team_week_state(evidence.inputs, decision_time=evidence.decision_time) + remaining = {(item.sleeper_player_id, item.game_id) for item in state.unpassed_opportunities} + assert ("p1", "g1") not in remaining + assert ("p1", "g2") in remaining + assert state.passed_opportunities[0].decision_time == NOW - timedelta(hours=2) + assert not state.is_blocked + + +def test_async_sqlite_lock_round_trips_into_fixed_slot(tmp_path: Path) -> None: + path = tmp_path / "state.db" + writer = AsyncSQLiteStateRepository(path) + asyncio_run(writer.initialize()) + asyncio_run( + _seed_acknowledgement( + writer, + _lock_in_recommendation(), + AcknowledgementAction.LOCKED, + token="lock-live", + at=NOW - timedelta(hours=2), + ) + ) + source = AsyncSQLiteStateRepository(path) + + async def run() -> object: + return await _collect(_nba(), _RecordingProjections(), acknowledgement_source=source) + + evidence = asyncio_run(run()) + state = build_live_team_week_state(evidence.inputs, decision_time=evidence.decision_time) + assert len(state.fixed_slots) == 1 + fixed = state.fixed_slots[0] + assert fixed.slot_index == 0 + assert fixed.slot_position == "PG" + assert fixed.accepted_fantasy_score == 24 + assert fixed.decision_time == NOW - timedelta(hours=2) + assert fixed.provenance == ACKNOWLEDGEMENT_PROVENANCE + assert not state.is_blocked + + +def test_malformed_repository_trace_blocks_without_becoming_a_constraint( + tmp_path: Path, +) -> None: + path = tmp_path / "state.db" + writer = AsyncSQLiteStateRepository(path) + asyncio_run(writer.initialize()) + asyncio_run( + _seed_acknowledgement( + writer, + _lock_in_recommendation(trace_json="{"), + AcknowledgementAction.LOCKED, + token="bad-live", + at=NOW - timedelta(hours=2), + ) + ) + source = AsyncSQLiteStateRepository(path) + + async def run() -> object: + return await _collect(_nba(), _RecordingProjections(), acknowledgement_source=source) + + evidence = asyncio_run(run()) + assert len(evidence.inputs.acknowledgements) == 1 + assert evidence.inputs.acknowledgements[0].reconciled is False + state = build_live_team_week_state(evidence.inputs, decision_time=evidence.decision_time) + assert PlanningReasonCode.ACKNOWLEDGEMENT_CONFLICT in state.blocking_reasons + assert state.fixed_slots == () + assert state.passed_opportunities == () + + +def test_repository_acknowledgement_errors_propagate() -> None: + class _BrokenSource: + async def load_acknowledged_decisions( + self, league_id: str, fantasy_week: int, *, as_of: datetime + ) -> tuple[AcknowledgedDecisionEvidence, ...]: + raise AcknowledgementQueryError("unexpected D1 result envelope") + + async def run() -> object: + return await _collect( + _nba(), + _RecordingProjections(), + acknowledgement_source=_BrokenSource(), + ) + + with pytest.raises(AcknowledgementQueryError, match="envelope"): + asyncio_run(run()) + + def asyncio_run(awaitable): # noqa: ANN001 import asyncio From 25b626fa76c7f5f1e20fd19b3b6f2886e02a6a15 Mon Sep 17 00:00:00 2001 From: Ryan Zhou Date: Mon, 24 Aug 2026 12:07:33 -0500 Subject: [PATCH 4/5] chore: forbid commit co-author trailers --- .cursor/rules/no-cursor-coauthor.mdc | 13 +++++++++++++ AGENTS.md | 1 + 2 files changed, 14 insertions(+) create mode 100644 .cursor/rules/no-cursor-coauthor.mdc diff --git a/.cursor/rules/no-cursor-coauthor.mdc b/.cursor/rules/no-cursor-coauthor.mdc new file mode 100644 index 0000000..7711798 --- /dev/null +++ b/.cursor/rules/no-cursor-coauthor.mdc @@ -0,0 +1,13 @@ +--- +description: Never add Cursor or any other Co-authored-by trailer to git commits +alwaysApply: true +--- + +# No commit co-authors + +Never add a `Co-authored-by` trailer to a git commit, including +`Co-authored-by: Cursor `. + +Commit messages in this repo are a single subject line with no body unless the +user explicitly requests a body. Do not append attribution, AI co-author lines, +or any other trailer. diff --git a/AGENTS.md b/AGENTS.md index c989d93..28fdaac 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -19,6 +19,7 @@ - Prefer several focused commits over one large commit that mixes unrelated behavior, refactoring, tests, or infrastructure changes. - Do not create noisy checkpoint or broken WIP commits merely for frequency. Each commit should represent a meaningful completed step and pass the relevant formatting, linting, and tests. - Use short, one-line commit messages with no body unless the user explicitly requests otherwise. +- Never add a `Co-authored-by` trailer, including Cursor (`Co-authored-by: Cursor `). Do not append any commit attribution trailers. - Follow the concise `: ` style, for example `chore: initialize sleeper manager` or `feat: sync league profile`. - Do not mention local roadmap, design, implementation, or validation phases—or untracked planning documents—in commit messages; those are not durable repository context for reviewers. - Make the subject describe the outcome of the commit, not the implementation process. Avoid vague messages such as `updates`, `changes`, or `fix stuff`. From b8214bae0e81445b9b9eb43d368cd2e97171e845 Mon Sep 17 00:00:00 2001 From: Ryan Zhou Date: Mon, 24 Aug 2026 12:07:33 -0500 Subject: [PATCH 5/5] fix: fail closed on live acknowledgement queries --- .../persistence/acknowledgements.py | 12 +-- src/sleeper_manager/persistence/d1.py | 29 +++++- tests/unit/test_acknowledgement_queries.py | 88 +++++++++++++++++++ tests/unit/test_d1.py | 2 +- 4 files changed, 121 insertions(+), 10 deletions(-) diff --git a/src/sleeper_manager/persistence/acknowledgements.py b/src/sleeper_manager/persistence/acknowledgements.py index 6089c82..de6d5cf 100644 --- a/src/sleeper_manager/persistence/acknowledgements.py +++ b/src/sleeper_manager/persistence/acknowledgements.py @@ -98,21 +98,23 @@ def decode_acknowledged_decisions( raise AcknowledgementQueryError("as_of must be timezone-aware") decoded: list[AcknowledgedDecisionEvidence] = [] for row in rows: - evidence = _decode_row(row) - if evidence.decided_at > as_of: + decided_at = _parse_aware(row.acknowledged_at, label="acknowledgement timestamp") + if decided_at > as_of: continue - decoded.append(evidence) + decoded.append(_decode_row(row, decided_at)) return _reconcile(_order(decoded)) -def _decode_row(row: AcknowledgementRawRow) -> AcknowledgedDecisionEvidence: +def _decode_row( + row: AcknowledgementRawRow, + decided_at: datetime, +) -> AcknowledgedDecisionEvidence: decision_id = _require_identity(row.recommendation_id, "recommendation ID") player_id = _require_identity(row.player_id, "player ID") game_id = _require_identity(row.game_id, "game ID") action = _ACTION_BY_STORED_VALUE.get(_text_or_none(row.acknowledgement_action) or "") if action is None: raise AcknowledgementQueryError("unknown acknowledgement action") - decided_at = _parse_aware(row.acknowledged_at, label="acknowledgement timestamp") slot_index, slot_position, accepted_score, trace_ok = _parse_trace(row.trace_json, action) duplicates_ok = _duplicates_match(row, action, decided_at) return AcknowledgedDecisionEvidence( diff --git a/src/sleeper_manager/persistence/d1.py b/src/sleeper_manager/persistence/d1.py index 35604c8..aa68f08 100644 --- a/src/sleeper_manager/persistence/d1.py +++ b/src/sleeper_manager/persistence/d1.py @@ -117,6 +117,23 @@ ON recommendations (league_id, fantasy_week, decision_type, status); """ +_MISSING = object() + + +def _js_to_python(value: object) -> object: + converter = getattr(value, "to_py", None) + if callable(converter): + return converter() + return value + + +def _d1_field(original: object, payload: object, name: str) -> object: + if isinstance(payload, Mapping) and name in payload: + return _js_to_python(payload[name]) + if hasattr(original, name): + return _js_to_python(getattr(original, name)) + return _MISSING + class D1StateRepository(AsyncStateRepository): """Async repository backed by a Cloudflare D1 binding.""" @@ -137,12 +154,16 @@ async def _first(self, query: str, *params: object) -> dict[str, Any] | None: async def _all(self, query: str, *params: object) -> Sequence[object]: result = await self._statement(query, params).all() - if not isinstance(result, Mapping): + payload = _js_to_python(result) + rows = _d1_field(result, payload, "results") + if rows is _MISSING or isinstance(rows, str | bytes) or not isinstance(rows, Sequence): raise AcknowledgementQueryError("unexpected D1 result envelope") - rows = result.get("results") - if isinstance(rows, str | bytes) or not isinstance(rows, Sequence): + success = _d1_field(result, payload, "success") + if success is False: + raise AcknowledgementQueryError("unsuccessful D1 query") + if success is not True: raise AcknowledgementQueryError("unexpected D1 result envelope") - return rows + return tuple(_js_to_python(row) for row in rows) async def _run(self, query: str, *params: object) -> Any: return await self._statement(query, params).run() diff --git a/tests/unit/test_acknowledgement_queries.py b/tests/unit/test_acknowledgement_queries.py index 1c97ac5..6cc74f1 100644 --- a/tests/unit/test_acknowledgement_queries.py +++ b/tests/unit/test_acknowledgement_queries.py @@ -178,6 +178,25 @@ def test_point_in_time_includes_equivalent_offsets_and_excludes_newer() -> None: assert [item.decision_id for item in result] == ["rec-earlier", "rec-boundary"] +def test_future_rows_are_skipped_before_identity_decoding() -> None: + future = (AS_OF + timedelta(minutes=1)).isoformat() + rows = ( + _row(recommendation_id="rec-current"), + _row( + recommendation_id=None, + player_id=None, + game_id=None, + acknowledgement_action="lock", + acknowledged_at=future, + recommendation_acknowledged_at=future, + ), + ) + + result = decode_acknowledged_decisions(rows, as_of=AS_OF) + + assert [item.decision_id for item in result] == ["rec-current"] + + def test_result_order_uses_datetime_comparison_not_iso_strings() -> None: later_utc = datetime(2026, 1, 7, 19, tzinfo=UTC) earlier_offset = datetime(2026, 1, 7, 12, tzinfo=timezone(timedelta(hours=-8))) @@ -884,6 +903,75 @@ async def all(self) -> object: asyncio.run(repository.load_acknowledged_decisions("league-1", 1, as_of=AS_OF)) +class _JsProxy: + """Minimal Python Worker D1 proxy: not a Mapping, convertible via to_py().""" + + def __init__(self, value: object) -> None: + self._value = value + + def to_py(self) -> object: + return self._value + + +class _AttributeD1Result: + def __init__(self, results: object, success: object = True) -> None: + self.results = results + self.success = success + + +def _joined_lock_row() -> dict[str, object]: + return { + "recommendation_id": "rec-1", + "player_id": "p1", + "game_id": "g1", + "recommendation_status": "acknowledged", + "recommendation_action": "locked", + "recommendation_acknowledged_at": AS_OF.isoformat(), + "trace_json": _lock_trace(), + "acknowledgement_action": "locked", + "acknowledged_at": AS_OF.isoformat(), + } + + +def _repository_returning(result: object) -> D1StateRepository: + class _Statement: + def bind(self, *params: object) -> _Statement: + return self + + async def all(self) -> object: + return result + + class _Database: + def prepare(self, query: str) -> object: + return _Statement() + + return D1StateRepository(_Database()) + + +def test_d1_python_worker_proxy_results_decode() -> None: + row = _joined_lock_row() + proxy_result = _JsProxy({"success": True, "results": [_JsProxy(row)]}) + loaded = asyncio.run( + _repository_returning(proxy_result).load_acknowledged_decisions("league-1", 1, as_of=AS_OF) + ) + assert [item.decision_id for item in loaded] == ["rec-1"] + assert loaded[0].reconciled is True + + attribute_result = _AttributeD1Result(results=[_JsProxy(row)], success=True) + loaded = asyncio.run( + _repository_returning(attribute_result).load_acknowledged_decisions( + "league-1", 1, as_of=AS_OF + ) + ) + assert [item.decision_id for item in loaded] == ["rec-1"] + + +def test_d1_unsuccessful_query_raises_instead_of_empty_constraints() -> None: + repository = _repository_returning({"success": False, "results": []}) + with pytest.raises(AcknowledgementQueryError, match="unsuccessful"): + asyncio.run(repository.load_acknowledged_decisions("league-1", 1, as_of=AS_OF)) + + def test_migration_adds_index_without_rewriting_phase3_rows(tmp_path: Path) -> None: path = tmp_path / "migrated.db" connection = sqlite3.connect(path) diff --git a/tests/unit/test_d1.py b/tests/unit/test_d1.py index 2a452a3..caec41f 100644 --- a/tests/unit/test_d1.py +++ b/tests/unit/test_d1.py @@ -44,7 +44,7 @@ async def first(self) -> dict[str, Any] | None: async def all(self) -> dict[str, Any]: cursor = self.execute() - return {"results": [dict(row) for row in cursor.fetchall()]} + return {"success": True, "results": [dict(row) for row in cursor.fetchall()]} class FakeD1: