"""Explicit SQLite persistence operations for stable Activity identity and revisions."""

from __future__ import annotations

import gzip
import json
from contextlib import contextmanager
from io import BytesIO
from typing import Any, Iterator

from mountain_twin.activity.contracts import ActivitySourceKind, RawProviderRecord, User

from .contracts import (
    ActivityEntity,
    ActivityRevision,
    ActivityRevisionProvenance,
    ActivitySourceLink,
    RawPayloadBlob,
    activity_source_link_id_for,
    canonical_raw_payload_bytes,
    deserialize_activity_snapshot,
    new_activity_entity_id,
    new_activity_source_link,
    raw_payload_sha256,
    raw_provider_record_id_for,
    serialize_activity_snapshot,
)
from .sqlite import ActivitySqliteDatabase, transaction
from .sync import (
    OpaqueProviderState,
    ProviderSyncCheckpoint,
    ProviderSyncRun,
    ProviderSyncSource,
    SyncFrontierPolicy,
    SyncFrontierScope,
    SyncMode,
    SyncRunSourceObservation,
    SyncRunStatus,
    new_provider_sync_run_id,
    provider_sync_source_id_for,
)

_ALLOWED_SYNC_TRANSITIONS = {
    SyncRunStatus.PLANNED: {SyncRunStatus.RUNNING, SyncRunStatus.CANCELLED},
    SyncRunStatus.RUNNING: {
        SyncRunStatus.COMPLETED, SyncRunStatus.PAUSED_RATE_LIMIT, SyncRunStatus.FAILED_RETRYABLE,
        SyncRunStatus.FAILED_ACTIVITY, SyncRunStatus.BLOCKED_AUTH, SyncRunStatus.CANCELLED,
    },
    SyncRunStatus.PAUSED_RATE_LIMIT: {SyncRunStatus.RUNNING, SyncRunStatus.CANCELLED},
    SyncRunStatus.FAILED_RETRYABLE: {SyncRunStatus.RUNNING, SyncRunStatus.CANCELLED},
    SyncRunStatus.FAILED_ACTIVITY: {SyncRunStatus.RUNNING, SyncRunStatus.CANCELLED},
    SyncRunStatus.BLOCKED_AUTH: {SyncRunStatus.RUNNING, SyncRunStatus.CANCELLED},
    SyncRunStatus.COMPLETED: set(),
    SyncRunStatus.CANCELLED: set(),
}


class SqliteActivityIdentityStore:
    """Durable entity/revision and immutable source-truth store; sync and telemetry remain deferred."""

    def __init__(self, database: ActivitySqliteDatabase):
        self.connection = database.connect()

    def close(self) -> None:
        self.connection.close()

    @contextmanager
    def transaction(self) -> Iterator[None]:
        with transaction(self.connection):
            yield

    def ensure_owner(self, user: User) -> User:
        with self.transaction():
            row = self.connection.execute(
                "SELECT created_at FROM activity_owners WHERE user_id = ?", (user.user_id,)
            ).fetchone()
            if row is None:
                self.connection.execute(
                    "INSERT INTO activity_owners(user_id, created_at) VALUES (?, ?)",
                    (user.user_id, user.created_at),
                )
            elif row["created_at"] != user.created_at:
                raise ValueError("activity owner identity already exists with a different creation time")
        return user

    def create_entity(self, user: User, *, activity_entity_id: str | None = None) -> ActivityEntity:
        self.ensure_owner(user)
        entity = ActivityEntity(activity_entity_id or new_activity_entity_id(), user.user_id, user.created_at)
        with self.transaction():
            self.connection.execute(
                "INSERT INTO activity_entities(activity_entity_id, user_id, current_revision_id, created_at) "
                "VALUES (?, ?, NULL, ?)",
                (entity.activity_entity_id, entity.user_id, entity.created_at),
            )
        return entity

    def get_entity(self, activity_entity_id: str) -> ActivityEntity:
        row = self.connection.execute(
            "SELECT activity_entity_id, user_id, created_at, current_revision_id "
            "FROM activity_entities WHERE activity_entity_id = ?", (activity_entity_id,)
        ).fetchone()
        if row is None:
            raise KeyError(activity_entity_id)
        return ActivityEntity(
            row["activity_entity_id"], row["user_id"], row["created_at"], row["current_revision_id"]
        )

    def insert_revision(self, revision: ActivityRevision) -> ActivityRevision:
        entity = self.get_entity(revision.activity_entity_id)
        if entity.user_id != revision.canonical_activity.user_id:
            raise ValueError("activity revision owner differs from its ActivityEntity")
        snapshot = serialize_activity_snapshot(revision.canonical_activity)
        with self.transaction():
            self.connection.execute(
                """INSERT OR IGNORE INTO activity_revisions(
                    activity_revision_id, activity_entity_id, canonical_activity_id,
                    canonical_activity_json, normalizer_version, taxonomy_version, normalized_at,
                    started_at, activity_type, activity_family, distance_m, elapsed_duration_s,
                    moving_duration_s, elevation_gain_m, average_heart_rate_bpm, has_track
                ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""",
                (
                    revision.activity_revision_id, revision.activity_entity_id,
                    revision.canonical_activity.activity_id, snapshot,
                    revision.normalizer_version, revision.taxonomy_version, revision.normalized_at,
                    revision.canonical_activity.started_at,
                    revision.canonical_activity.activity_type.value,
                    revision.canonical_activity.activity_family.value,
                    revision.canonical_activity.distance_m,
                    revision.canonical_activity.elapsed_duration_s,
                    revision.canonical_activity.moving_duration_s,
                    revision.canonical_activity.elevation_gain_m,
                    revision.canonical_activity.average_heart_rate_bpm,
                    int(revision.canonical_activity.track is not None),
                ),
            )
            stored = self._get_revision_row(revision.activity_revision_id)
            stored_revision = self._revision_from_row(stored)
            if (
                stored_revision.activity_entity_id != revision.activity_entity_id
                or stored_revision.normalizer_version != revision.normalizer_version
                or stored_revision.taxonomy_version != revision.taxonomy_version
            ):
                raise ValueError("activity revision identity already exists with different immutable content")
        # A semantically identical repeat may have a later normalized_at audit
        # value.  Keep the first immutable snapshot; never rewrite it.
        return stored_revision

    def get_revision(self, activity_revision_id: str) -> ActivityRevision:
        return self._revision_from_row(self._get_revision_row(activity_revision_id))

    def list_revisions(self, activity_entity_id: str) -> tuple[ActivityRevision, ...]:
        self.get_entity(activity_entity_id)
        rows = self.connection.execute(
            "SELECT * FROM activity_revisions WHERE activity_entity_id = ? "
            "ORDER BY normalized_at, activity_revision_id", (activity_entity_id,)
        ).fetchall()
        return tuple(self._revision_from_row(row) for row in rows)

    def set_current_revision(self, activity_entity_id: str, activity_revision_id: str) -> ActivityEntity:
        with self.transaction():
            revision = self._get_revision_row(activity_revision_id)
            if revision["activity_entity_id"] != activity_entity_id:
                raise ValueError("current revision must belong to the same ActivityEntity")
            updated = self.connection.execute(
                "UPDATE activity_entities SET current_revision_id = ? WHERE activity_entity_id = ?",
                (activity_revision_id, activity_entity_id),
            )
            if updated.rowcount != 1:
                raise KeyError(activity_entity_id)
        return self.get_entity(activity_entity_id)

    def get_current_revision(self, activity_entity_id: str) -> ActivityRevision | None:
        entity = self.get_entity(activity_entity_id)
        return self.get_revision(entity.current_revision_id) if entity.current_revision_id else None

    def resolve_or_create_entity_source_link(
        self,
        user: User,
        source_kind: ActivitySourceKind,
        provider: str | None,
        source_activity_id: str,
        first_seen_at: str,
    ) -> tuple[ActivityEntity, ActivitySourceLink, bool]:
        """Resolve native identity before creating its one durable ActivityEntity."""
        source_link_id = activity_source_link_id_for(
            user.user_id, source_kind, provider, source_activity_id
        )
        with self.transaction():
            self.ensure_owner(user)
            existing = self.connection.execute(
                "SELECT * FROM activity_source_links WHERE source_link_id = ?", (source_link_id,)
            ).fetchone()
            if existing is not None:
                source_link = self._source_link_from_row(existing)
                return self.get_entity(source_link.activity_entity_id), source_link, False

            entity = ActivityEntity(new_activity_entity_id(), user.user_id, first_seen_at)
            self.connection.execute(
                "INSERT INTO activity_entities(activity_entity_id, user_id, current_revision_id, created_at) "
                "VALUES (?, ?, NULL, ?)",
                (entity.activity_entity_id, entity.user_id, entity.created_at),
            )
            source_link = new_activity_source_link(
                entity.activity_entity_id,
                user.user_id,
                source_kind,
                provider,
                source_activity_id,
                first_seen_at,
            )
            self.connection.execute(
                """INSERT INTO activity_source_links(
                    source_link_id, activity_entity_id, user_id, source_kind, provider,
                    source_activity_id, first_seen_at
                ) VALUES (?, ?, ?, ?, ?, ?, ?)""",
                (
                    source_link.source_link_id, source_link.activity_entity_id, source_link.user_id,
                    source_link.source_kind.value, source_link.provider, source_link.source_activity_id,
                    source_link.first_seen_at,
                ),
            )
        return entity, source_link, True

    def resolve_source_link(
        self,
        activity_entity_id: str,
        source_kind: ActivitySourceKind,
        provider: str | None,
        source_activity_id: str,
        first_seen_at: str,
    ) -> ActivitySourceLink:
        entity = self.get_entity(activity_entity_id)
        requested = new_activity_source_link(
            activity_entity_id, entity.user_id, source_kind, provider, source_activity_id, first_seen_at
        )
        with self.transaction():
            self.connection.execute(
                """INSERT OR IGNORE INTO activity_source_links(
                    source_link_id, activity_entity_id, user_id, source_kind, provider,
                    source_activity_id, first_seen_at
                ) VALUES (?, ?, ?, ?, ?, ?, ?)""",
                (
                    requested.source_link_id, requested.activity_entity_id, requested.user_id,
                    requested.source_kind.value, requested.provider, requested.source_activity_id,
                    requested.first_seen_at,
                ),
            )
            stored = self._get_source_link_row(requested.source_link_id)
            source_link = self._source_link_from_row(stored)
            if source_link.activity_entity_id != activity_entity_id:
                raise ValueError("source identity is already linked to another ActivityEntity")
        return source_link

    def get_source_link(self, source_link_id: str) -> ActivitySourceLink:
        return self._source_link_from_row(self._get_source_link_row(source_link_id))

    def put_raw_payload(self, payload: Any) -> RawPayloadBlob:
        canonical_bytes = canonical_raw_payload_bytes(payload)
        blob = RawPayloadBlob(raw_payload_sha256(payload), "application/json", "gzip", len(canonical_bytes))
        compressed_bytes = _deterministic_gzip(canonical_bytes)
        with self.transaction():
            self.connection.execute(
                """INSERT OR IGNORE INTO raw_payload_blobs(
                    payload_sha256, encoding, compression, byte_count, payload_bytes
                ) VALUES (?, ?, ?, ?, ?)""",
                (blob.payload_sha256, blob.encoding, blob.compression, blob.byte_count, compressed_bytes),
            )
            stored = self._get_raw_payload_row(blob.payload_sha256)
            if tuple(stored[name] for name in ("encoding", "compression", "byte_count", "payload_bytes")) != (
                blob.encoding, blob.compression, blob.byte_count, compressed_bytes,
            ):
                raise ValueError("raw payload hash already exists with conflicting immutable content")
        return blob

    def get_raw_payload_blob(self, payload_sha256: str) -> RawPayloadBlob:
        row = self._get_raw_payload_row(payload_sha256)
        return RawPayloadBlob(row["payload_sha256"], row["encoding"], row["compression"], row["byte_count"])

    def get_raw_payload(self, payload_sha256: str) -> Any:
        row = self._get_raw_payload_row(payload_sha256)
        blob = self.get_raw_payload_blob(payload_sha256)
        try:
            canonical_bytes = gzip.decompress(row["payload_bytes"])
        except OSError as error:
            raise ValueError("raw payload blob cannot be decompressed") from error
        if len(canonical_bytes) != blob.byte_count or raw_payload_sha256_from_bytes(canonical_bytes) != blob.payload_sha256:
            raise ValueError("raw payload blob does not match its content hash")
        try:
            payload = json.loads(canonical_bytes.decode("utf-8"))
        except (UnicodeDecodeError, json.JSONDecodeError) as error:
            raise ValueError("raw payload blob is not valid JSON") from error
        if canonical_raw_payload_bytes(payload) != canonical_bytes:
            raise ValueError("raw payload blob is not canonical JSON")
        return payload

    def insert_raw_provider_record(self, source_link_id: str, record: RawProviderRecord) -> RawProviderRecord:
        source_link = self.get_source_link(source_link_id)
        expected_id = raw_provider_record_id_for(
            source_link.user_id, source_link.source_link_id, record.record_kind, record.payload_sha256
        )
        if record.raw_record_id != expected_id:
            raise ValueError("raw provider record identity must match immutable source truth")
        if (
            record.user_id != source_link.user_id
            or record.source_kind is not source_link.source_kind
            or record.provider != source_link.provider
            or record.source_activity_id != source_link.source_activity_id
            or record.payload_reference != f"sqlite:raw-payload:{record.payload_sha256}"
        ):
            raise ValueError("raw provider record differs from its ActivitySourceLink")
        self.get_raw_payload_blob(record.payload_sha256)
        with self.transaction():
            self.connection.execute(
                """INSERT OR IGNORE INTO raw_provider_records(
                    raw_record_id, user_id, source_link_id, record_kind, payload_sha256,
                    source_schema_version, fetched_at
                ) VALUES (?, ?, ?, ?, ?, ?, ?)""",
                (
                    record.raw_record_id, record.user_id, source_link_id, record.record_kind,
                    record.payload_sha256, record.source_schema_version, record.fetched_at,
                ),
            )
            stored = self._get_raw_provider_record_row(record.raw_record_id)
            if tuple(stored[name] for name in (
                "user_id", "source_link_id", "record_kind", "payload_sha256",
            )) != (record.user_id, source_link_id, record.record_kind, record.payload_sha256):
                raise ValueError("raw provider record identity already exists with conflicting immutable content")
        return self._raw_provider_record_from_row(stored)

    def get_raw_provider_record(self, raw_record_id: str) -> RawProviderRecord:
        return self._raw_provider_record_from_row(self._get_raw_provider_record_row(raw_record_id))

    def list_raw_provider_records(self, source_link_id: str) -> tuple[RawProviderRecord, ...]:
        self.get_source_link(source_link_id)
        rows = self.connection.execute(
            "SELECT * FROM raw_provider_records WHERE source_link_id = ? "
            "ORDER BY record_kind, fetched_at, raw_record_id", (source_link_id,)
        ).fetchall()
        return tuple(self._raw_provider_record_from_row(row) for row in rows)

    def attach_revision_provenance(
        self, activity_revision_id: str, raw_record_id: str
    ) -> ActivityRevisionProvenance:
        revision = self.get_revision(activity_revision_id)
        raw_row = self._get_raw_provider_record_row(raw_record_id)
        source_link = self.get_source_link(raw_row["source_link_id"])
        raw_record = self._raw_provider_record_from_row(raw_row)
        if source_link.activity_entity_id != revision.activity_entity_id:
            raise ValueError("revision and source link belong to different ActivityEntities")
        matching_reference = next(
            (
                reference
                for reference in revision.canonical_activity.provenance.source_references
                if (
                    reference.source_kind is raw_record.source_kind
                    and reference.provider == raw_record.provider
                    and reference.source_activity_id == raw_record.source_activity_id
                )
            ),
            None,
        )
        if matching_reference is None:
            raise ValueError("raw provider record source identity is not canonical ActivityRevision provenance")
        provenance = ActivityRevisionProvenance(
            revision.activity_revision_id, revision.activity_entity_id,
            raw_record_id, source_link.source_link_id,
        )
        with self.transaction():
            self.connection.execute(
                """INSERT OR IGNORE INTO activity_revision_provenance(
                    activity_revision_id, activity_entity_id, raw_record_id, source_link_id
                ) VALUES (?, ?, ?, ?)""",
                (
                    provenance.activity_revision_id, provenance.activity_entity_id,
                    provenance.raw_record_id, provenance.source_link_id,
                ),
            )
        return provenance

    def list_revision_provenance(self, activity_revision_id: str) -> tuple[ActivityRevisionProvenance, ...]:
        self.get_revision(activity_revision_id)
        rows = self.connection.execute(
            "SELECT * FROM activity_revision_provenance WHERE activity_revision_id = ? "
            "ORDER BY raw_record_id", (activity_revision_id,)
        ).fetchall()
        return tuple(
            ActivityRevisionProvenance(
                row["activity_revision_id"], row["activity_entity_id"],
                row["raw_record_id"], row["source_link_id"],
            )
            for row in rows
        )

    def resolve_provider_sync_source(
        self, user: User, provider: str, external_subject_id: str, created_at: str
    ) -> ProviderSyncSource:
        source_id = provider_sync_source_id_for(user.user_id, provider, external_subject_id)
        with self.transaction():
            self.ensure_owner(user)
            self.connection.execute(
                """INSERT OR IGNORE INTO provider_sync_sources(
                    sync_source_id, user_id, provider, external_subject_id, created_at
                ) VALUES (?, ?, ?, ?, ?)""",
                (source_id, user.user_id, provider, external_subject_id, created_at),
            )
            source = self._sync_source_from_row(self._get_sync_source_row(source_id))
            if source.user_id != user.user_id:
                raise ValueError("sync source owner differs from requested owner")
        return source

    def get_provider_sync_source(self, sync_source_id: str) -> ProviderSyncSource:
        return self._sync_source_from_row(self._get_sync_source_row(sync_source_id))

    def create_sync_run(
        self,
        sync_source_id: str,
        mode,
        plan: OpaqueProviderState,
        resume: OpaqueProviderState,
        started_at: str,
        *,
        frontier_policy: SyncFrontierPolicy = SyncFrontierPolicy.NONE,
    ) -> ProviderSyncRun:
        source = self.get_provider_sync_source(sync_source_id)
        run = ProviderSyncRun(
            new_provider_sync_run_id(), source.sync_source_id, source.user_id, mode,
            SyncRunStatus.PLANNED, plan, resume, started_at, started_at,
            frontier_policy=frontier_policy,
        )
        with self.transaction():
            self.connection.execute(
                """INSERT INTO provider_sync_runs(
                    sync_run_id, sync_source_id, user_id, mode, status, plan_version, plan_json,
                    resume_version, resume_json, started_at, updated_at, completed_at,
                    pages_completed, activities_persisted, failure_code, frontier_policy,
                    terminal_state_version, terminal_state_json
                ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, NULL, 0, 0, NULL, ?, NULL, NULL)""",
                (run.sync_run_id, run.sync_source_id, run.user_id, run.mode.value, run.status.value,
                 plan.version, plan.serialized, resume.version, resume.serialized, started_at, started_at,
                 frontier_policy.value),
            )
        return run

    def get_sync_run(self, sync_run_id: str) -> ProviderSyncRun:
        return self._sync_run_from_row(self._get_sync_run_row(sync_run_id))

    def transition_sync_run(self, sync_run_id: str, status: SyncRunStatus, updated_at: str,
                            failure_code: str | None = None) -> ProviderSyncRun:
        run = self.get_sync_run(sync_run_id)
        if status is SyncRunStatus.COMPLETED or status not in _ALLOWED_SYNC_TRANSITIONS[run.status]:
            raise ValueError("invalid durable sync run lifecycle transition")
        with self.transaction():
            self.connection.execute(
                "UPDATE provider_sync_runs SET status = ?, updated_at = ?, failure_code = ? WHERE sync_run_id = ?",
                (status.value, updated_at, failure_code, sync_run_id),
            )
        return self.get_sync_run(sync_run_id)

    def advance_sync_run_resume(
        self, sync_run_id: str, resume: OpaqueProviderState, updated_at: str,
        *, pages_completed: int, activities_persisted: int,
    ) -> ProviderSyncRun:
        run = self.get_sync_run(sync_run_id)
        if run.status is not SyncRunStatus.RUNNING:
            raise ValueError("only running sync run may advance durable resume state")
        if pages_completed < run.pages_completed or activities_persisted < run.activities_persisted:
            raise ValueError("sync durable progress cannot move backward")
        with self.transaction():
            self.connection.execute(
                """UPDATE provider_sync_runs SET resume_version = ?, resume_json = ?, updated_at = ?,
                    pages_completed = ?, activities_persisted = ? WHERE sync_run_id = ?""",
                (resume.version, resume.serialized, updated_at, pages_completed, activities_persisted, sync_run_id),
            )
        return self.get_sync_run(sync_run_id)

    def complete_sync_run(
        self, sync_run_id: str, completed_at: str,
        *, pages_completed: int | None = None, activities_persisted: int | None = None,
        terminal_provider_state: OpaqueProviderState | None = None,
        frontier_candidate: OpaqueProviderState | None = None,
    ) -> ProviderSyncCheckpoint | None:
        run = self.get_sync_run(sync_run_id)
        if SyncRunStatus.COMPLETED not in _ALLOWED_SYNC_TRANSITIONS[run.status]:
            raise ValueError("only running sync run may complete")
        eligible = run.frontier_policy is SyncFrontierPolicy.PUBLISH_INCREMENTAL_CONTINUATION
        if eligible != (frontier_candidate is not None):
            raise ValueError("frontier candidate must match the persisted sync run frontier policy")
        final_pages = run.pages_completed if pages_completed is None else pages_completed
        final_activities = run.activities_persisted if activities_persisted is None else activities_persisted
        if final_pages < run.pages_completed or final_activities < run.activities_persisted:
            raise ValueError("completed sync progress cannot move backward")
        result = None
        if frontier_candidate is not None:
            result = ProviderSyncCheckpoint(
                run.sync_source_id, frontier_candidate, run.sync_run_id, completed_at,
                SyncFrontierScope.INCREMENTAL_CONTINUATION,
            )
        with self.transaction():
            self.connection.execute(
                """UPDATE provider_sync_runs SET status = ?, updated_at = ?, completed_at = ?,
                    pages_completed = ?, activities_persisted = ?, failure_code = NULL,
                    terminal_state_version = ?, terminal_state_json = ? WHERE sync_run_id = ?""",
                (SyncRunStatus.COMPLETED.value, completed_at, completed_at, final_pages, final_activities,
                 terminal_provider_state.version if terminal_provider_state else None,
                 terminal_provider_state.serialized if terminal_provider_state else None, sync_run_id),
            )
            if result is not None:
                self.connection.execute(
                    """INSERT INTO provider_sync_checkpoints(
                        sync_source_id, state_version, state_json, completed_sync_run_id, published_at, frontier_scope
                    ) VALUES (?, ?, ?, ?, ?, ?)
                    ON CONFLICT(sync_source_id) DO UPDATE SET
                        state_version = excluded.state_version, state_json = excluded.state_json,
                        completed_sync_run_id = excluded.completed_sync_run_id, published_at = excluded.published_at,
                        frontier_scope = excluded.frontier_scope""",
                    (result.sync_source_id, frontier_candidate.version, frontier_candidate.serialized,
                     result.completed_sync_run_id, result.published_at, result.frontier_scope.value),
                )
        return result

    def get_sync_checkpoint(self, sync_source_id: str) -> ProviderSyncCheckpoint:
        row = self.connection.execute(
            "SELECT * FROM provider_sync_checkpoints WHERE sync_source_id = ?", (sync_source_id,)
        ).fetchone()
        if row is None:
            raise KeyError(sync_source_id)
        return ProviderSyncCheckpoint(
            row["sync_source_id"], OpaqueProviderState.from_serialized(row["state_version"], row["state_json"]),
            row["completed_sync_run_id"], row["published_at"], SyncFrontierScope(row["frontier_scope"]),
        )

    def get_authoritative_frontier(self, sync_source_id: str) -> ProviderSyncCheckpoint:
        row = self.connection.execute(
            """SELECT * FROM provider_sync_checkpoints
               WHERE sync_source_id = ? AND frontier_scope = ?""",
            (sync_source_id, SyncFrontierScope.INCREMENTAL_CONTINUATION.value),
        ).fetchone()
        if row is None:
            raise KeyError(sync_source_id)
        return ProviderSyncCheckpoint(
            row["sync_source_id"], OpaqueProviderState.from_serialized(row["state_version"], row["state_json"]),
            row["completed_sync_run_id"], row["published_at"], SyncFrontierScope(row["frontier_scope"]),
        )

    def record_sync_run_observation(
        self, sync_run_id: str, source_link_id: str, raw_record_id: str, observed_at: str
    ) -> SyncRunSourceObservation:
        run = self.get_sync_run(sync_run_id)
        source_link = self.get_source_link(source_link_id)
        raw = self.get_raw_provider_record(raw_record_id)
        if source_link.user_id != run.user_id or raw.user_id != run.user_id:
            raise ValueError("sync observation owner differs from sync run owner")
        if source_link.provider != self.get_provider_sync_source(run.sync_source_id).provider:
            raise ValueError("sync observation provider differs from sync source")
        if raw.source_activity_id != source_link.source_activity_id:
            raise ValueError("sync observation raw record differs from source link")
        observation = SyncRunSourceObservation(
            run.sync_run_id, run.sync_source_id, run.user_id, source_link_id, raw_record_id, observed_at
        )
        with self.transaction():
            self.connection.execute(
                """INSERT OR IGNORE INTO sync_run_source_observations(
                    sync_run_id, sync_source_id, user_id, source_link_id, raw_record_id, observed_at
                ) VALUES (?, ?, ?, ?, ?, ?)""",
                (observation.sync_run_id, observation.sync_source_id, observation.user_id,
                 observation.source_link_id, observation.raw_record_id, observation.observed_at),
            )
        return observation

    def list_sync_run_observations(self, sync_run_id: str) -> tuple[SyncRunSourceObservation, ...]:
        self.get_sync_run(sync_run_id)
        rows = self.connection.execute(
            "SELECT * FROM sync_run_source_observations WHERE sync_run_id = ? ORDER BY source_link_id, raw_record_id",
            (sync_run_id,),
        ).fetchall()
        return tuple(SyncRunSourceObservation(
            row["sync_run_id"], row["sync_source_id"], row["user_id"], row["source_link_id"],
            row["raw_record_id"], row["observed_at"],
        ) for row in rows)

    def _get_revision_row(self, activity_revision_id: str):
        row = self.connection.execute(
            "SELECT * FROM activity_revisions WHERE activity_revision_id = ?", (activity_revision_id,)
        ).fetchone()
        if row is None:
            raise KeyError(activity_revision_id)
        return row

    def _get_sync_source_row(self, sync_source_id: str):
        row = self.connection.execute(
            "SELECT * FROM provider_sync_sources WHERE sync_source_id = ?", (sync_source_id,)
        ).fetchone()
        if row is None:
            raise KeyError(sync_source_id)
        return row

    def _get_sync_run_row(self, sync_run_id: str):
        row = self.connection.execute("SELECT * FROM provider_sync_runs WHERE sync_run_id = ?", (sync_run_id,)).fetchone()
        if row is None:
            raise KeyError(sync_run_id)
        return row

    def _get_source_link_row(self, source_link_id: str):
        row = self.connection.execute(
            "SELECT * FROM activity_source_links WHERE source_link_id = ?", (source_link_id,)
        ).fetchone()
        if row is None:
            raise KeyError(source_link_id)
        return row

    def _get_raw_payload_row(self, payload_sha256: str):
        row = self.connection.execute(
            "SELECT * FROM raw_payload_blobs WHERE payload_sha256 = ?", (payload_sha256,)
        ).fetchone()
        if row is None:
            raise KeyError(payload_sha256)
        return row

    def _get_raw_provider_record_row(self, raw_record_id: str):
        row = self.connection.execute(
            "SELECT * FROM raw_provider_records WHERE raw_record_id = ?", (raw_record_id,)
        ).fetchone()
        if row is None:
            raise KeyError(raw_record_id)
        return row

    @staticmethod
    def _revision_from_row(row) -> ActivityRevision:
        return ActivityRevision(
            activity_revision_id=row["activity_revision_id"],
            activity_entity_id=row["activity_entity_id"],
            canonical_activity=deserialize_activity_snapshot(row["canonical_activity_json"]),
            normalizer_version=row["normalizer_version"],
            taxonomy_version=row["taxonomy_version"],
            normalized_at=row["normalized_at"],
        )

    @staticmethod
    def _source_link_from_row(row) -> ActivitySourceLink:
        return ActivitySourceLink(
            row["source_link_id"], row["activity_entity_id"], row["user_id"],
            ActivitySourceKind(row["source_kind"]), row["provider"],
            row["source_activity_id"], row["first_seen_at"],
        )

    @staticmethod
    def _sync_source_from_row(row) -> ProviderSyncSource:
        return ProviderSyncSource(
            row["sync_source_id"], row["user_id"], row["provider"], row["external_subject_id"], row["created_at"]
        )

    @staticmethod
    def _sync_run_from_row(row) -> ProviderSyncRun:
        return ProviderSyncRun(
            row["sync_run_id"], row["sync_source_id"], row["user_id"], SyncMode(row["mode"]),
            SyncRunStatus(row["status"]), OpaqueProviderState.from_serialized(row["plan_version"], row["plan_json"]),
            OpaqueProviderState.from_serialized(row["resume_version"], row["resume_json"]), row["started_at"],
            row["updated_at"], row["completed_at"], row["pages_completed"], row["activities_persisted"], row["failure_code"],
            SyncFrontierPolicy(row["frontier_policy"]),
            OpaqueProviderState.from_serialized(row["terminal_state_version"], row["terminal_state_json"])
            if row["terminal_state_version"] is not None else None,
        )

    def _raw_provider_record_from_row(self, row) -> RawProviderRecord:
        source_link = self.get_source_link(row["source_link_id"])
        return RawProviderRecord(
            row["raw_record_id"], row["user_id"], source_link.source_kind,
            source_link.provider, source_link.source_activity_id, row["fetched_at"],
            f"sqlite:raw-payload:{row['payload_sha256']}", row["payload_sha256"],
            row["source_schema_version"], row["record_kind"],
        )


def _deterministic_gzip(content: bytes) -> bytes:
    output = BytesIO()
    with gzip.GzipFile(fileobj=output, mode="wb", filename="", mtime=0) as archive:
        archive.write(content)
    return output.getvalue()


def raw_payload_sha256_from_bytes(content: bytes) -> str:
    import hashlib

    return hashlib.sha256(content).hexdigest()
