Coverage for src/lilbee/data/ingest/skip_marker.py: 100%

121 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-09-28 17:20 +0000

1"""Sidecar records of files a sync holds out: ingestion failures and user removals. 

2 

3``skipped_sources.json`` maps a filename to the file hash it is held out at. 

4``_plan_file_changes`` treats a file whose current hash matches its marker as 

5unchanged, so a failed extract is paid once and a removal stays out. Editing the 

6file changes its hash and re-arms it. ``skip_reasons.json`` records why each 

7file is held out, and ``skip_kinds.json`` records its ``SkipKind`` with the 

8hash and reason it was written for. A stored kind counts only while both still 

9match; otherwise the record reads as a removal when its reason is 

10``REMOVED_SKIP_REASON`` and as a failure when it is not. Production writes go 

11through ``update_skip_records`` and ``clear_skip_markers``, both under one 

12cross-process lock. 

13""" 

14 

15from __future__ import annotations 

16 

17import contextlib 

18import json 

19import logging 

20import os 

21from collections.abc import Callable, Iterable, Mapping 

22from dataclasses import dataclass, field 

23from enum import StrEnum 

24from pathlib import Path 

25from typing import TypedDict 

26 

27from lilbee.core.security import file_lock_or_warn 

28from lilbee.data.types import SkippedSource 

29 

30log = logging.getLogger(__name__) 

31 

32SKIP_MARKER_FILENAME = "skipped_sources.json" 

33SKIP_REASON_FILENAME = "skip_reasons.json" 

34SKIP_KIND_FILENAME = "skip_kinds.json" 

35DEFAULT_SKIP_REASON = "held out by an earlier sync" 

36REMOVED_SKIP_REASON = "removed via remove (re-add the source or run retry-skipped to restore)" 

37# A sync, a /delete, and a reset from another process all change these records. 

38_RECORDS_LOCK_TIMEOUT_S = 10.0 

39 

40 

41class SkipKind(StrEnum): 

42 """Why a skip marker holds a file out of the sync.""" 

43 

44 FAILED = "failed" 

45 REMOVED = "removed" 

46 

47 

48class _StoredKind(TypedDict): 

49 """One ``skip_kinds.json`` entry: the kind and the marker it was written for.""" 

50 

51 kind: SkipKind 

52 hash: str 

53 reason: str | None 

54 

55 

56@dataclass 

57class SkipRecords: 

58 """The skip markers with the reason and the kind of each.""" 

59 

60 markers: dict[str, str] = field(default_factory=dict) 

61 reasons: dict[str, str] = field(default_factory=dict) 

62 kinds: dict[str, SkipKind] = field(default_factory=dict) 

63 

64 

65def _load_json_map(path: Path) -> dict[str, object]: 

66 """Load a JSON object, or empty dict on any read/parse error.""" 

67 if not path.exists(): 

68 return {} 

69 try: 

70 raw = json.loads(path.read_text(encoding="utf-8")) 

71 except (OSError, json.JSONDecodeError, UnicodeDecodeError) as exc: 

72 log.debug("Sidecar %s unreadable, treating as empty: %s", path.name, exc) 

73 return {} 

74 if not isinstance(raw, dict): # the file is untyped JSON 

75 return {} 

76 return {str(k): v for k, v in raw.items()} 

77 

78 

79def _load_str_map(path: Path) -> dict[str, str]: 

80 """Load a ``{str: str}`` JSON file, or empty dict on any read/parse error.""" 

81 return {k: v for k, v in _load_json_map(path).items() if isinstance(v, str)} 

82 

83 

84def _write_json_map(path: Path, data: Mapping[str, object]) -> None: 

85 """Replace *path* atomically with a JSON object. Best-effort.""" 

86 tmp = path.with_suffix(path.suffix + ".tmp") 

87 try: 

88 path.parent.mkdir(parents=True, exist_ok=True) 

89 tmp.write_text(json.dumps(data, sort_keys=True), encoding="utf-8") 

90 os.replace(tmp, path) 

91 except OSError as exc: 

92 log.warning("Failed to persist %s: %s", path, exc) 

93 with contextlib.suppress(OSError): 

94 tmp.unlink() 

95 

96 

97def _unlink(path: Path) -> None: 

98 try: 

99 path.unlink(missing_ok=True) 

100 except OSError as exc: 

101 log.debug("Could not remove %s: %s", path, exc) 

102 

103 

104def load_skip_markers(data_root: Path) -> dict[str, str]: 

105 """Load the filename → failed-hash map, or empty dict on any read error.""" 

106 return _load_str_map(data_root / SKIP_MARKER_FILENAME) 

107 

108 

109def write_skip_markers(data_root: Path, markers: dict[str, str]) -> None: 

110 """Replace the marker file atomically. Best-effort: errors are logged, not raised.""" 

111 _write_json_map(data_root / SKIP_MARKER_FILENAME, markers) 

112 

113 

114def load_skip_reasons(data_root: Path) -> dict[str, str]: 

115 """Load the filename → skip-reason map (informational), empty on any read error.""" 

116 return _load_str_map(data_root / SKIP_REASON_FILENAME) 

117 

118 

119def write_skip_reasons(data_root: Path, reasons: dict[str, str]) -> None: 

120 """Replace the reasons sidecar atomically. Best-effort: errors are logged, not raised.""" 

121 _write_json_map(data_root / SKIP_REASON_FILENAME, reasons) 

122 

123 

124def _kind_of(stored: object, marker: str, reason: str | None) -> SkipKind: 

125 """The stored kind while it names this marker and reason; else REMOVED for the removal text.""" 

126 # the kinds file is untyped JSON 

127 if isinstance(stored, dict) and (stored.get("hash"), stored.get("reason")) == (marker, reason): 

128 with contextlib.suppress(ValueError): 

129 return SkipKind(stored.get("kind", "")) 

130 return SkipKind.REMOVED if reason == REMOVED_SKIP_REASON else SkipKind.FAILED 

131 

132 

133def _load_records(data_root: Path) -> SkipRecords: 

134 """Read the three sidecars; a marker without a matching stored kind reads by its reason.""" 

135 markers = load_skip_markers(data_root) 

136 reasons = load_skip_reasons(data_root) 

137 stored = _load_json_map(data_root / SKIP_KIND_FILENAME) 

138 kinds = { 

139 name: _kind_of(stored.get(name), marker, reasons.get(name)) 

140 for name, marker in markers.items() 

141 } 

142 return SkipRecords(markers, reasons, kinds) 

143 

144 

145def write_skip_records(data_root: Path, records: SkipRecords) -> None: 

146 """Replace the three sidecars; each kind names the hash and reason it is written for.""" 

147 write_skip_markers(data_root, records.markers) 

148 write_skip_reasons(data_root, records.reasons) 

149 stored = { 

150 name: _StoredKind(kind=kind, hash=marker, reason=records.reasons.get(name)) 

151 for name, kind in records.kinds.items() 

152 if (marker := records.markers.get(name)) is not None 

153 } 

154 _write_json_map(data_root / SKIP_KIND_FILENAME, stored) 

155 

156 

157def load_skip_kinds(data_root: Path) -> dict[str, SkipKind]: 

158 """The kind of every marked file.""" 

159 return _load_records(data_root).kinds 

160 

161 

162def update_skip_records(data_root: Path, change: Callable[[SkipRecords], None]) -> None: 

163 """Apply *change* to the records as they are on disk, under a cross-process lock. 

164 

165 Reasons and kinds whose marker is gone are dropped in the same write. 

166 """ 

167 with file_lock_or_warn(data_root / SKIP_MARKER_FILENAME, _RECORDS_LOCK_TIMEOUT_S): 

168 records = _load_records(data_root) 

169 before = SkipRecords(dict(records.markers), dict(records.reasons), dict(records.kinds)) 

170 change(records) 

171 records.reasons = {k: v for k, v in records.reasons.items() if k in records.markers} 

172 records.kinds = {k: v for k, v in records.kinds.items() if k in records.markers} 

173 if records != before: 

174 write_skip_records(data_root, records) 

175 

176 

177def held_out_names(data_root: Path) -> list[str]: 

178 """Every file held out by an ingestion failure, sorted; removed sources are not listed.""" 

179 kinds = load_skip_kinds(data_root) 

180 return sorted(name for name, kind in kinds.items() if kind is SkipKind.FAILED) 

181 

182 

183def mark_removed( 

184 data_root: Path, hashes: Mapping[str, str], unreachable: Iterable[str] = () 

185) -> None: 

186 """Hold files out of every sync as user removals. 

187 

188 Each file in *hashes* is held at the given hash; each marked file in 

189 *unreachable* keeps the hash its marker already records. 

190 """ 

191 

192 def _hold(records: SkipRecords) -> None: 

193 kept = {name: records.markers[name] for name in unreachable if name in records.markers} 

194 held = {**kept, **hashes} 

195 records.markers.update(held) 

196 records.reasons.update(dict.fromkeys(held, REMOVED_SKIP_REASON)) 

197 records.kinds.update(dict.fromkeys(held, SkipKind.REMOVED)) 

198 

199 update_skip_records(data_root, _hold) 

200 

201 

202def clear_failed_markers(data_root: Path) -> list[str]: 

203 """Drop the record of every file an ingestion failure holds out; removals stay. Returns them.""" 

204 dropped: list[str] = [] 

205 

206 def _drop(records: SkipRecords) -> None: 

207 dropped.extend(name for name, kind in records.kinds.items() if kind is SkipKind.FAILED) 

208 for name in dropped: 

209 records.markers.pop(name) 

210 

211 update_skip_records(data_root, _drop) 

212 return dropped 

213 

214 

215def describe_skips(data_root: Path, names: Iterable[str]) -> list[SkippedSource]: 

216 """Pair each name with its recorded reason, in order; ``DEFAULT_SKIP_REASON`` when none.""" 

217 reasons = load_skip_reasons(data_root) 

218 return [ 

219 SkippedSource(filename=name, reason=reasons.get(name, DEFAULT_SKIP_REASON)) 

220 for name in names 

221 ] 

222 

223 

224def clear_skip_markers(data_root: Path) -> None: 

225 """Delete the marker file and both sidecars. No-op if absent.""" 

226 with file_lock_or_warn(data_root / SKIP_MARKER_FILENAME, _RECORDS_LOCK_TIMEOUT_S): 

227 _unlink(data_root / SKIP_MARKER_FILENAME) 

228 _unlink(data_root / SKIP_REASON_FILENAME) 

229 _unlink(data_root / SKIP_KIND_FILENAME)