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

1"""Background sync, executor management, and sync status for chat mode.""" 

2 

3from __future__ import annotations 

4 

5import asyncio 

6import threading 

7from collections.abc import Callable 

8from concurrent.futures import Future, ThreadPoolExecutor 

9from typing import TYPE_CHECKING 

10 

11from rich.text import Text 

12 

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) 

26 

27if TYPE_CHECKING: 

28 from lilbee.runtime.progress import DetailedProgressCallback 

29 

30 

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 

45 

46 

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 ) 

57 

58 

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) 

67 

68 

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 } 

75 

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) 

80 

81 return _callback 

82 

83 

84_bg_executor: ThreadPoolExecutor | None = None 

85 

86 

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 

93 

94 

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 

102 

103 _bg_executor.shutdown(wait=False, cancel_futures=True) 

104 _bg_executor = None 

105 

106 

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) 

120 

121 

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 """ 

128 

129 def __init__(self) -> None: 

130 self.text: str = "" 

131 self.pending: int = 0 

132 self._pending_lock = threading.Lock() 

133 

134 def clear(self) -> None: 

135 self.text = "" 

136 

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 

141 

142 

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}" 

148 

149 

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() 

157 

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}") 

183 

184 return _callback 

185 

186 

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() 

195 

196 callback = _chat_sync_callback(status) if chat_mode else _sync_progress_printer(con) 

197 

198 def _run() -> object: 

199 if chat_mode: 

200 status.adjust_pending(-1) 

201 return asyncio.run(sync(quiet=True, on_progress=callback)) 

202 

203 if chat_mode: 

204 status.adjust_pending(1) 

205 

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