""" Work-record UUID ↔ external issue id mapping. Fleet bookkeeping for the connector role (INTENT / work-record canon). Storage is a dedicated SQLite file so mapping works regardless of whether the active CRUD backend is local or Gitea. """ from __future__ import annotations import sqlite3 import uuid from dataclasses import dataclass from datetime import datetime, timezone from pathlib import Path from typing import Any, Dict, List, Optional ACTIVE = "active" DETACHED = "detached" SUPERSEDED = "superseded" @dataclass class WorkRecordMapping: """One projection link between a work record and a backend issue.""" id: str work_record_uuid: str work_record_id: Optional[str] work_record_kind: Optional[str] backend: str external_id: str external_url: Optional[str] target_repo: Optional[str] direction: str # outward | linked status: str # active | detached | superseded created_at: str updated_at: str last_pushed_at: Optional[str] = None last_pulled_at: Optional[str] = None def to_dict(self) -> Dict[str, Any]: return { "id": self.id, "work_record_uuid": self.work_record_uuid, "work_record_id": self.work_record_id, "work_record_kind": self.work_record_kind, "backend": self.backend, "external_id": self.external_id, "external_url": self.external_url, "target_repo": self.target_repo, "direction": self.direction, "status": self.status, "created_at": self.created_at, "updated_at": self.updated_at, "last_pushed_at": self.last_pushed_at, "last_pulled_at": self.last_pulled_at, } # Work-record task status → tracker IssueState (outward v1). # See docs/boundary-sync-and-status-mapping.md TASK_STATUS_TO_ISSUE_STATE = { "wait": "open", "todo": "open", "progress": "in_progress", "done": "closed", "cancel": "closed", # Already tracker-ish values (pass-through) "open": "open", "closed": "closed", "in_progress": "in_progress", "blocked": "blocked", } def map_work_record_status_to_issue_state(status: str) -> str: """Map fleet work-record status (or tracker state) to IssueState value.""" key = (status or "").strip().lower().replace("-", "_") if key not in TASK_STATUS_TO_ISSUE_STATE: raise ValueError( f"Unknown status {status!r}; expected one of " f"{sorted(set(TASK_STATUS_TO_ISSUE_STATE))}" ) return TASK_STATUS_TO_ISSUE_STATE[key] class MappingService: """CRUD for work_record_issue_map rows.""" def __init__(self, db_path: Optional[str] = None): self.db_path = db_path or "mappings.db" self.connection: Optional[sqlite3.Connection] = None def connect(self, db_path: Optional[str] = None) -> None: if db_path: self.db_path = db_path Path(self.db_path).parent.mkdir(parents=True, exist_ok=True) self.connection = sqlite3.connect(self.db_path) self.connection.row_factory = sqlite3.Row self.connection.execute("PRAGMA foreign_keys = ON") self._initialize_schema() def disconnect(self) -> None: if self.connection: self.connection.close() self.connection = None def _initialize_schema(self) -> None: assert self.connection is not None self.connection.executescript( """ CREATE TABLE IF NOT EXISTS work_record_issue_map ( id TEXT PRIMARY KEY, work_record_uuid TEXT NOT NULL, work_record_id TEXT, work_record_kind TEXT, backend TEXT NOT NULL, external_id TEXT NOT NULL, external_url TEXT, target_repo TEXT, direction TEXT NOT NULL DEFAULT 'outward', status TEXT NOT NULL DEFAULT 'active' CHECK (status IN ('active', 'detached', 'superseded')), created_at TEXT NOT NULL, updated_at TEXT NOT NULL, last_pushed_at TEXT, last_pulled_at TEXT ); CREATE UNIQUE INDEX IF NOT EXISTS idx_map_active_uuid_backend ON work_record_issue_map (work_record_uuid, backend) WHERE status = 'active'; CREATE UNIQUE INDEX IF NOT EXISTS idx_map_active_backend_external ON work_record_issue_map (backend, external_id) WHERE status = 'active'; CREATE INDEX IF NOT EXISTS idx_map_work_record_id ON work_record_issue_map (work_record_id); """ ) self.connection.commit() def _require_conn(self) -> sqlite3.Connection: if not self.connection: raise RuntimeError("MappingService is not connected") return self.connection def _row_to_mapping(self, row: sqlite3.Row) -> WorkRecordMapping: return WorkRecordMapping( id=row["id"], work_record_uuid=row["work_record_uuid"], work_record_id=row["work_record_id"], work_record_kind=row["work_record_kind"], backend=row["backend"], external_id=row["external_id"], external_url=row["external_url"], target_repo=row["target_repo"], direction=row["direction"], status=row["status"], created_at=row["created_at"], updated_at=row["updated_at"], last_pushed_at=row["last_pushed_at"], last_pulled_at=row["last_pulled_at"], ) def get_active_by_uuid( self, work_record_uuid: str, backend: Optional[str] = None ) -> Optional[WorkRecordMapping]: conn = self._require_conn() if backend: cur = conn.execute( """ SELECT * FROM work_record_issue_map WHERE work_record_uuid = ? AND backend = ? AND status = 'active' """, (work_record_uuid, backend), ) else: cur = conn.execute( """ SELECT * FROM work_record_issue_map WHERE work_record_uuid = ? AND status = 'active' ORDER BY updated_at DESC """, (work_record_uuid,), ) row = cur.fetchone() return self._row_to_mapping(row) if row else None def get_active_by_canonical_id( self, work_record_id: str, backend: Optional[str] = None ) -> Optional[WorkRecordMapping]: conn = self._require_conn() if backend: cur = conn.execute( """ SELECT * FROM work_record_issue_map WHERE work_record_id = ? AND backend = ? AND status = 'active' """, (work_record_id, backend), ) else: cur = conn.execute( """ SELECT * FROM work_record_issue_map WHERE work_record_id = ? AND status = 'active' ORDER BY updated_at DESC """, (work_record_id,), ) row = cur.fetchone() return self._row_to_mapping(row) if row else None def get_active_by_external( self, backend: str, external_id: str ) -> Optional[WorkRecordMapping]: conn = self._require_conn() cur = conn.execute( """ SELECT * FROM work_record_issue_map WHERE backend = ? AND external_id = ? AND status = 'active' """, (backend, str(external_id)), ) row = cur.fetchone() return self._row_to_mapping(row) if row else None def resolve( self, *, work_record_uuid: Optional[str] = None, work_record_id: Optional[str] = None, backend: Optional[str] = None, external_id: Optional[str] = None, ) -> Optional[WorkRecordMapping]: if work_record_uuid: return self.get_active_by_uuid(work_record_uuid, backend=backend) if work_record_id: return self.get_active_by_canonical_id(work_record_id, backend=backend) if backend and external_id is not None: return self.get_active_by_external(backend, external_id) raise ValueError("Provide work_record_uuid, work_record_id, or backend+external_id") def upsert( self, *, work_record_uuid: str, backend: str, external_id: str, work_record_id: Optional[str] = None, work_record_kind: Optional[str] = None, external_url: Optional[str] = None, target_repo: Optional[str] = None, direction: str = "outward", ) -> WorkRecordMapping: """Idempotent active mapping for (uuid, backend).""" conn = self._require_conn() now = datetime.now(timezone.utc).isoformat() existing = self.get_active_by_uuid(work_record_uuid, backend=backend) if existing: if existing.external_id != str(external_id): raise ValueError( f"Active mapping for {work_record_uuid} on {backend} already " f"points at external_id={existing.external_id!r}, not {external_id!r}" ) conn.execute( """ UPDATE work_record_issue_map SET work_record_id = COALESCE(?, work_record_id), work_record_kind = COALESCE(?, work_record_kind), external_url = COALESCE(?, external_url), target_repo = COALESCE(?, target_repo), updated_at = ? WHERE id = ? """, ( work_record_id, work_record_kind, external_url, target_repo, now, existing.id, ), ) conn.commit() refreshed = self.get_active_by_uuid(work_record_uuid, backend=backend) assert refreshed is not None return refreshed # External id already mapped to a different UUID? conflict = self.get_active_by_external(backend, external_id) if conflict and conflict.work_record_uuid != work_record_uuid: raise ValueError( f"external_id={external_id!r} on {backend} already mapped to " f"{conflict.work_record_uuid}" ) mapping_id = str(uuid.uuid4()) conn.execute( """ INSERT INTO work_record_issue_map ( id, work_record_uuid, work_record_id, work_record_kind, backend, external_id, external_url, target_repo, direction, status, created_at, updated_at ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, 'active', ?, ?) """, ( mapping_id, work_record_uuid, work_record_id, work_record_kind, backend, str(external_id), external_url, target_repo, direction, now, now, ), ) conn.commit() created = self.get_active_by_uuid(work_record_uuid, backend=backend) assert created is not None return created def detach( self, *, work_record_uuid: Optional[str] = None, work_record_id: Optional[str] = None, backend: Optional[str] = None, ) -> Optional[WorkRecordMapping]: conn = self._require_conn() mapping = self.resolve( work_record_uuid=work_record_uuid, work_record_id=work_record_id, backend=backend, ) if not mapping: return None now = datetime.now(timezone.utc).isoformat() conn.execute( """ UPDATE work_record_issue_map SET status = 'detached', updated_at = ? WHERE id = ? """, (now, mapping.id), ) conn.commit() cur = conn.execute( "SELECT * FROM work_record_issue_map WHERE id = ?", (mapping.id,) ) row = cur.fetchone() return self._row_to_mapping(row) if row else None def mark_pushed(self, mapping_id: str) -> None: conn = self._require_conn() now = datetime.now(timezone.utc).isoformat() conn.execute( """ UPDATE work_record_issue_map SET last_pushed_at = ?, updated_at = ? WHERE id = ? """, (now, now, mapping_id), ) conn.commit() def list_active(self, limit: int = 100) -> List[WorkRecordMapping]: conn = self._require_conn() cur = conn.execute( """ SELECT * FROM work_record_issue_map WHERE status = 'active' ORDER BY updated_at DESC LIMIT ? """, (limit,), ) return [self._row_to_mapping(r) for r in cur.fetchall()]