86 lines
2.8 KiB
Python
86 lines
2.8 KiB
Python
"""Storage contract for local Wiki pages, vectors, revisions and sync leases."""
|
|
|
|
from __future__ import annotations
|
|
|
|
from abc import ABC, abstractmethod
|
|
from collections.abc import AsyncIterator
|
|
from typing import Any
|
|
|
|
|
|
class LlmWikiIndexStore(ABC):
|
|
@abstractmethod
|
|
async def get_page_by_slug_hash(self, mapping_id: str, slug_hash: str) -> dict[str, Any] | None: ...
|
|
|
|
@abstractmethod
|
|
async def get_page(self, mapping_id: str, slug: str) -> dict[str, Any] | None: ...
|
|
|
|
@abstractmethod
|
|
async def touch_page(self, page_id: str, sync_id: str) -> None: ...
|
|
|
|
@abstractmethod
|
|
async def replace_page_vectors(
|
|
self,
|
|
*,
|
|
page: dict[str, Any],
|
|
vectors: list[dict[str, Any]],
|
|
fingerprint: str,
|
|
sync_id: str,
|
|
) -> dict[str, Any]: ...
|
|
|
|
@abstractmethod
|
|
async def record_page_failure(self, *, page: dict[str, Any], sync_id: str, error: str) -> None: ...
|
|
|
|
@abstractmethod
|
|
async def mark_missing_pages_deleted(self, mapping_id: str, sync_id: str) -> int: ...
|
|
|
|
@abstractmethod
|
|
async def load_vector_snapshot(self, mapping_id: str, fingerprint: str) -> dict[str, Any]: ...
|
|
|
|
async def iter_vector_snapshot_rows(
|
|
self,
|
|
mapping_id: str,
|
|
fingerprint: str,
|
|
*,
|
|
batch_size: int = 1000,
|
|
) -> AsyncIterator[dict[str, Any]]:
|
|
"""Iterate vector rows for streaming export.
|
|
|
|
Implementations with a database backend should override this method so
|
|
very large vector packages do not need to be materialized in memory.
|
|
"""
|
|
snapshot = await self.load_vector_snapshot(mapping_id, fingerprint)
|
|
for row in snapshot.get("rows") or []:
|
|
if isinstance(row, dict):
|
|
yield row
|
|
|
|
@abstractmethod
|
|
async def get_index_revision(self, mapping_id: str) -> tuple[int, str | None, str]: ...
|
|
|
|
@abstractmethod
|
|
async def try_acquire_sync_lease(self, mapping_id: str, *, owner: str, sync_id: str, lease_seconds: int, fingerprint: str) -> bool: ...
|
|
|
|
@abstractmethod
|
|
async def renew_sync_lease(self, mapping_id: str, *, owner: str, lease_seconds: int) -> bool: ...
|
|
|
|
@abstractmethod
|
|
async def update_sync_progress(
|
|
self,
|
|
mapping_id: str,
|
|
*,
|
|
owner: str,
|
|
sync_id: str,
|
|
remote_page_count: int,
|
|
) -> bool: ...
|
|
|
|
@abstractmethod
|
|
async def finish_sync(self, mapping_id: str, *, owner: str, sync_id: str, success: bool, remote_page_count: int, error: str | None = None) -> None: ...
|
|
|
|
@abstractmethod
|
|
async def list_index_status(self, mapping_ids: list[str] | None = None) -> list[dict[str, Any]]: ...
|
|
|
|
@abstractmethod
|
|
async def list_page_index_status(self, mapping_id: str) -> list[dict[str, Any]]: ...
|
|
|
|
@abstractmethod
|
|
async def mark_rebuild_required(self, mapping_ids: list[str], fingerprint: str) -> int: ...
|