"""Persistence contract for the independent enterprise-research workbench.""" from __future__ import annotations from abc import ABC, abstractmethod from typing import Any class EnterpriseResearchTaskStore(ABC): """Owner-scoped task storage; evidence snapshots are DeerFlow-owned.""" @abstractmethod async def list_tasks(self, user_id: str, *, limit: int = 20) -> list[dict[str, Any]]: raise NotImplementedError @abstractmethod async def get_task(self, task_id: str, user_id: str) -> dict[str, Any] | None: raise NotImplementedError @abstractmethod async def create_task(self, data: dict[str, Any]) -> dict[str, Any]: raise NotImplementedError @abstractmethod async def update_task(self, task_id: str, user_id: str, data: dict[str, Any]) -> dict[str, Any] | None: raise NotImplementedError @abstractmethod async def delete_task(self, task_id: str, user_id: str) -> bool: raise NotImplementedError class EnterpriseResearchReportJobStore(ABC): """Owner-scoped, frozen-input jobs for enterprise report writing.""" @abstractmethod async def create_job(self, data: dict[str, Any]) -> dict[str, Any]: raise NotImplementedError @abstractmethod async def get_job(self, job_id: str, user_id: str) -> dict[str, Any] | None: raise NotImplementedError @abstractmethod async def get_active_for_task(self, task_id: str, user_id: str) -> dict[str, Any] | None: raise NotImplementedError @abstractmethod async def claim_job(self, job_id: str, user_id: str) -> bool: """Atomically move one queued job to running.""" raise NotImplementedError @abstractmethod async def requeue_job(self, job_id: str, user_id: str, *, expected_updated_at: str | None) -> bool: """Atomically requeue the exact running snapshot seen during recovery.""" raise NotImplementedError @abstractmethod async def transition_job( self, job_id: str, user_id: str, *, expected_statuses: set[str], data: dict[str, Any] ) -> dict[str, Any] | None: """Apply a terminal/state transition only from one of the expected states.""" raise NotImplementedError @abstractmethod async def update_job(self, job_id: str, user_id: str, data: dict[str, Any]) -> dict[str, Any] | None: raise NotImplementedError @abstractmethod async def list_jobs(self, task_id: str, user_id: str, *, limit: int = 10) -> list[dict[str, Any]]: raise NotImplementedError @abstractmethod async def list_recoverable(self, *, limit: int = 100) -> list[dict[str, Any]]: raise NotImplementedError