Coverage for src/lilbee/cli/sync.py: 100%
103 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-09-28 17:20 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-09-28 17:20 +0000
1"""Background sync, executor management, and sync status for chat mode."""
3from __future__ import annotations
5import asyncio
6import threading
7from collections.abc import Callable
8from concurrent.futures import Future, ThreadPoolExecutor
9from typing import TYPE_CHECKING
11from rich.text import Text
13from lilbee.cli import theme
14from lilbee.cli.helpers import print_prefixed
15from lilbee.data.ingest import sync
16from lilbee.runtime.asyncio_loop import is_executor_shutdown
17from lilbee.runtime.console import PlainConsole
18from lilbee.runtime.progress import (
19 EventType,
20 ExtractEvent,
21 FileStartEvent,
22 OcrStartEvent,
23 ProgressEvent,
24 SyncDoneEvent,
25)
27if TYPE_CHECKING:
28 from lilbee.runtime.progress import DetailedProgressCallback
31def _format_sync_summary(
32 added: int, updated: int, removed: int, failed: int, skipped: int = 0, relocated: int = 0
33) -> str | None:
34 """Format sync counts into a human-readable summary, or None if nothing changed."""
35 counts = {
36 "added": added,
37 "updated": updated,
38 "removed": removed,
39 "relocated": relocated,
40 "skipped": skipped,
41 "failed": failed,
42 }
43 parts = [f"{n} {label}" for label, n in counts.items() if n]
44 return ", ".join(parts) if parts else None
47def _print_file_start(con: PlainConsole, data: ProgressEvent) -> None:
48 if not isinstance(data, FileStartEvent):
49 raise TypeError(f"Expected FileStartEvent, got {type(data).__name__}")
50 con.print(
51 Text(
52 f"Syncing [{data.current_file}/{data.total_files}]: {data.file}",
53 style=theme.MUTED,
54 ),
55 soft_wrap=True,
56 )
59def _print_done(con: PlainConsole, data: ProgressEvent) -> None:
60 if not isinstance(data, SyncDoneEvent):
61 raise TypeError(f"Expected SyncDoneEvent, got {type(data).__name__}")
62 summary = _format_sync_summary(
63 data.added, data.updated, data.removed, data.failed, data.skipped, data.relocated
64 )
65 if summary:
66 con.print(f"Synced: {summary}", style=theme.MUTED)
69def _sync_progress_printer(con: PlainConsole) -> DetailedProgressCallback:
70 """Return a callback that prints one-line status for FILE_START and SYNC_DONE events."""
71 handlers: dict[EventType, Callable[[PlainConsole, ProgressEvent], None]] = {
72 EventType.FILE_START: _print_file_start,
73 EventType.SYNC_DONE: _print_done,
74 }
76 def _callback(event_type: EventType, data: ProgressEvent) -> None:
77 handler = handlers.get(event_type)
78 if handler is not None:
79 handler(con, data)
81 return _callback
84_bg_executor: ThreadPoolExecutor | None = None
87def _get_executor() -> ThreadPoolExecutor:
88 """Lazy-init a single-worker executor."""
89 global _bg_executor
90 if _bg_executor is None:
91 _bg_executor = ThreadPoolExecutor(max_workers=1)
92 return _bg_executor
95def shutdown_executor() -> None:
96 """Shut down the background executor without blocking.
97 Uses wait=False + cancel_futures to avoid blocking the main thread.
98 """
99 global _bg_executor
100 if _bg_executor is None:
101 return
103 _bg_executor.shutdown(wait=False, cancel_futures=True)
104 _bg_executor = None
107def _on_sync_done(con: PlainConsole, future: Future[object], *, chat_mode: bool = False) -> None:
108 """Callback attached to background sync futures: logs errors."""
109 exc = future.exception()
110 if exc is None:
111 return
112 if isinstance(exc, asyncio.CancelledError):
113 return
114 if is_executor_shutdown(exc):
115 return
116 if chat_mode:
117 print(f"Background sync error: {exc}")
118 else:
119 print_prefixed(con, "Background sync error: ", exc, style=theme.ERROR)
122class SyncStatus:
123 """Thread-safe holder for background sync status text.
124 The background sync callback writes here; prompt_toolkit's
125 ``bottom_toolbar`` reads it on every render cycle: no cursor
126 manipulation, no flickering.
127 """
129 def __init__(self) -> None:
130 self.text: str = ""
131 self.pending: int = 0
132 self._pending_lock = threading.Lock()
134 def clear(self) -> None:
135 self.text = ""
137 def adjust_pending(self, delta: int) -> None:
138 """Atomically change the queued-sync counter (mutated from two threads)."""
139 with self._pending_lock:
140 self.pending += delta
143def _tesseract_status(data: ProgressEvent) -> str:
144 """The status line for an OCR_START event: Tesseract running on the file's pages."""
145 if not isinstance(data, OcrStartEvent):
146 raise TypeError(f"Expected OcrStartEvent, got {type(data).__name__}")
147 return f"⟳ {data.status_text}"
150def _chat_sync_callback(status: SyncStatus) -> DetailedProgressCallback:
151 """Return a progress callback for chat-mode background sync.
152 FILE_START updates *status.text* (rendered by prompt_toolkit's bottom
153 toolbar). On DONE the status is cleared and the summary is printed via
154 ``print()`` (goes through StdoutProxy → appears above the prompt).
155 """
156 status.clear()
158 def _callback(event_type: EventType, data: ProgressEvent) -> None:
159 queue_suffix = f" (+{status.pending} queued)" if status.pending > 0 else ""
160 if event_type == EventType.FILE_START:
161 if not isinstance(data, FileStartEvent):
162 raise TypeError(f"Expected FileStartEvent, got {type(data).__name__}")
163 status.text = (
164 f"⟳ Syncing [{data.current_file}/{data.total_files}]: {data.file}{queue_suffix}"
165 )
166 elif event_type == EventType.EXTRACT:
167 if not isinstance(data, ExtractEvent):
168 raise TypeError(f"Expected ExtractEvent, got {type(data).__name__}")
169 status.text = (
170 f"⟳ {data.step} [{data.page}/{data.total_pages}]: {data.file}{queue_suffix}"
171 )
172 elif event_type == EventType.OCR_START:
173 status.text = _tesseract_status(data) + queue_suffix
174 elif event_type == EventType.SYNC_DONE:
175 status.clear()
176 if not isinstance(data, SyncDoneEvent):
177 raise TypeError(f"Expected SyncDoneEvent, got {type(data).__name__}")
178 summary = _format_sync_summary(
179 data.added, data.updated, data.removed, data.failed, data.skipped, data.relocated
180 )
181 if summary:
182 print(f"✓ Synced: {summary}")
184 return _callback
187def run_sync_background(
188 con: PlainConsole,
189 *,
190 chat_mode: bool = False,
191 sync_status: SyncStatus | None = None,
192) -> Future[object]:
193 """Submit sync to a background thread. Returns the Future."""
194 status = sync_status or SyncStatus()
196 callback = _chat_sync_callback(status) if chat_mode else _sync_progress_printer(con)
198 def _run() -> object:
199 if chat_mode:
200 status.adjust_pending(-1)
201 return asyncio.run(sync(quiet=True, on_progress=callback))
203 if chat_mode:
204 status.adjust_pending(1)
206 future = _get_executor().submit(_run)
207 future.add_done_callback(lambda f: _on_sync_done(con, f, chat_mode=chat_mode))
208 return future