personal memory agent
0

Configure Feed

Select the types of activity you want to include in your feed.

solstone / tests / test_journal_caught_up.py
12 kB 441 lines
1# SPDX-License-Identifier: AGPL-3.0-only 2# Copyright (c) 2026 sol pbc 3 4from __future__ import annotations 5 6import builtins 7import json 8import os 9 10import pytest 11 12from solstone.think import catchup_state, pipeline_health 13from solstone.think.pipeline_health import ( 14 BACKLOG_STATE_COMPLETE, 15 BACKLOG_STATE_PENDING, 16 BACKLOG_STATE_UNKNOWN, 17 BacklogDay, 18 BacklogError, 19 BacklogView, 20) 21 22 23@pytest.fixture 24def doctor(): 25 from solstone.think import doctor as doctor_module 26 27 return doctor_module 28 29 30def args(doctor): 31 return doctor.Args(verbose=False, json=False, jsonl=False, port=5015) 32 33 34def clean_view() -> BacklogView: 35 return BacklogView( 36 window=30, 37 days=(), 38 pending_days=0, 39 stuck_days=0, 40 oldest_pending_day=None, 41 errors=(), 42 ) 43 44 45def unknown_day(day: str = "20200228") -> tuple[BacklogDay, BacklogError]: 46 error = BacklogError(day=day, stage="terminal_states", message="boom") 47 backlog_day = BacklogDay( 48 day=day, 49 state=BACKLOG_STATE_UNKNOWN, 50 segments=0, 51 units=0, 52 not_sensed=0, 53 why=(), 54 reason=None, 55 reason_code=None, 56 provider=None, 57 model=None, 58 error=error, 59 ) 60 return backlog_day, error 61 62 63def backlog_day(day: str, state: str) -> BacklogDay: 64 return BacklogDay( 65 day=day, 66 state=state, 67 segments=0, 68 units=0, 69 not_sensed=0, 70 why=(), 71 reason=None, 72 reason_code=None, 73 provider=None, 74 model=None, 75 error=None, 76 ) 77 78 79def test_journal_caught_up_ok_only_when_fully_clean(doctor, monkeypatch): 80 monkeypatch.setattr(pipeline_health, "read_backlog_view", clean_view) 81 82 result = doctor.journal_caught_up_check(args(doctor)) 83 84 assert result.status == "ok" 85 assert result.detail == "caught up" 86 assert result.severity == "advisory" 87 88 89def test_journal_caught_up_unknown_day_warns_before_false_green( 90 doctor, 91 monkeypatch, 92): 93 day, error = unknown_day() 94 view = BacklogView( 95 window=30, 96 days=(day,), 97 pending_days=0, 98 stuck_days=0, 99 oldest_pending_day=None, 100 errors=(error,), 101 ) 102 monkeypatch.setattr(pipeline_health, "read_backlog_view", lambda: view) 103 104 result = doctor.journal_caught_up_check(args(doctor)) 105 106 assert result.status == "warn" 107 assert "couldn't fully determine" in result.detail 108 assert result.status != "ok" 109 110 111def test_journal_caught_up_pending_and_stuck_warn_with_distinct_counts( 112 doctor, 113 monkeypatch, 114): 115 view = BacklogView( 116 window=30, 117 days=( 118 backlog_day("20200229", "pending"), 119 backlog_day("20200301", "pending"), 120 backlog_day("20200302", "stuck"), 121 ), 122 pending_days=2, 123 stuck_days=1, 124 oldest_pending_day="20200229", 125 errors=(), 126 ) 127 monkeypatch.setattr(pipeline_health, "read_backlog_view", lambda: view) 128 129 result = doctor.journal_caught_up_check(args(doctor)) 130 131 assert result.status == "warn" 132 assert result.severity == "advisory" 133 assert "2 day(s) pending" in result.detail 134 assert "1 day(s) stuck" in result.detail 135 assert "oldest outstanding 20200229" in result.detail 136 assert str(2 + 1) not in result.detail 137 138 139def test_journal_caught_up_warns_for_degraded_segment_repair_day( 140 doctor, 141 monkeypatch, 142): 143 day = BacklogDay( 144 day="20200229", 145 state=BACKLOG_STATE_PENDING, 146 segments=0, 147 units=0, 148 not_sensed=0, 149 why=(), 150 reason="segment_repair_degraded", 151 reason_code="segment_repair_degraded", 152 provider=None, 153 model=None, 154 error=None, 155 segment_repair_status="degraded", 156 ) 157 view = BacklogView( 158 window=30, 159 days=(day,), 160 pending_days=1, 161 stuck_days=0, 162 oldest_pending_day="20200229", 163 errors=(), 164 ) 165 monkeypatch.setattr(pipeline_health, "read_backlog_view", lambda: view) 166 167 result = doctor.journal_caught_up_check(args(doctor)) 168 169 assert result.status == "warn" 170 assert result.status != "ok" 171 assert "1 day(s) pending" in result.detail 172 173 174def test_journal_caught_up_warns_for_unknown_segment_repair_day( 175 doctor, 176 monkeypatch, 177): 178 error = BacklogError( 179 day="20200229", 180 stage="segment_repair", 181 message="segment-repair state unreadable", 182 ) 183 day = BacklogDay( 184 day="20200229", 185 state=BACKLOG_STATE_UNKNOWN, 186 segments=0, 187 units=0, 188 not_sensed=0, 189 why=(), 190 reason="segment_repair_unknown", 191 reason_code="segment_repair_unknown", 192 provider=None, 193 model=None, 194 error=error, 195 segment_repair_status="unknown", 196 ) 197 view = BacklogView( 198 window=30, 199 days=(day,), 200 pending_days=0, 201 stuck_days=0, 202 oldest_pending_day=None, 203 errors=(), 204 degraded=True, 205 ) 206 monkeypatch.setattr(pipeline_health, "read_backlog_view", lambda: view) 207 208 result = doctor.journal_caught_up_check(args(doctor)) 209 210 assert result.status == "warn" 211 assert result.status != "ok" 212 assert "1 day(s) unknown" in result.detail 213 214 215def test_journal_caught_up_ok_for_all_complete_view(doctor, monkeypatch): 216 view = BacklogView( 217 window=30, 218 days=(backlog_day("20200229", BACKLOG_STATE_COMPLETE),), 219 pending_days=0, 220 stuck_days=0, 221 oldest_pending_day=None, 222 errors=(), 223 ) 224 monkeypatch.setattr(pipeline_health, "read_backlog_view", lambda: view) 225 226 result = doctor.journal_caught_up_check(args(doctor)) 227 228 assert result.status == "ok" 229 assert result.detail == "caught up" 230 231 232def test_journal_caught_up_ok_with_capped_complete_days_detail(doctor, monkeypatch): 233 day = BacklogDay( 234 day="20200229", 235 state=BACKLOG_STATE_COMPLETE, 236 segments=0, 237 units=0, 238 not_sensed=0, 239 why=(), 240 reason=None, 241 reason_code=None, 242 provider=None, 243 model=None, 244 error=None, 245 capped_daily_unit_count=1, 246 capped_daily_unit={ 247 "name": "entities:entity_observer", 248 "facet": "vconic", 249 "reason_code": "context_window_exceeded", 250 "count": 2, 251 }, 252 ) 253 view = BacklogView( 254 window=30, 255 days=(day,), 256 pending_days=0, 257 stuck_days=0, 258 oldest_pending_day=None, 259 errors=(), 260 ) 261 monkeypatch.setattr(pipeline_health, "read_backlog_view", lambda: view) 262 263 result = doctor.journal_caught_up_check(args(doctor)) 264 265 assert result.status == "ok" 266 assert result.detail == "caught up; 1 day(s) completed with capped daily unit(s)" 267 268 269def test_journal_caught_up_plain_caught_up_without_capped_complete_days( 270 doctor, monkeypatch 271): 272 view = BacklogView( 273 window=30, 274 days=(backlog_day("20200229", BACKLOG_STATE_COMPLETE),), 275 pending_days=0, 276 stuck_days=0, 277 oldest_pending_day=None, 278 errors=(), 279 ) 280 monkeypatch.setattr(pipeline_health, "read_backlog_view", lambda: view) 281 282 result = doctor.journal_caught_up_check(args(doctor)) 283 284 assert result.status == "ok" 285 assert result.detail == "caught up" 286 287 288def test_journal_caught_up_reports_backoff_stuck_day(doctor, tmp_path, monkeypatch): 289 journal = tmp_path / "journal" 290 day = "20990401" 291 health_dir = journal / "chronicle" / day / "health" 292 health_dir.mkdir(parents=True, exist_ok=True) 293 daily = health_dir / "daily.updated" 294 stream = health_dir / "stream.updated" 295 daily.touch() 296 stream.touch() 297 os.utime(daily, (100.0, 100.0)) 298 os.utime(stream, (200.0, 200.0)) 299 state_path = journal / "health" / "catchup-state.json" 300 state_path.parent.mkdir(parents=True, exist_ok=True) 301 state_path.write_text( 302 json.dumps( 303 { 304 "version": catchup_state.STATE_VERSION, 305 "entries": { 306 f"{day}:{catchup_state.KIND_DAILY_CATCHUP}": { 307 "day": day, 308 "command_kind": catchup_state.KIND_DAILY_CATCHUP, 309 "attempts": 3, 310 "consecutive_non_completion": 3, 311 "last_attempt_at": 1000.0, 312 "last_outcome": "interrupted", 313 "next_retry_at": 1600.0, 314 "entered_backoff_at": 1200.0, 315 "notified_at": 1200.0, 316 "fingerprint": "fingerprint", 317 "active": None, 318 } 319 }, 320 } 321 ), 322 encoding="utf-8", 323 ) 324 monkeypatch.setenv("SOLSTONE_JOURNAL", str(journal)) 325 326 result = doctor.journal_caught_up_check(args(doctor)) 327 328 assert result.status == "warn" 329 assert "0 day(s) pending" in result.detail 330 assert "1 day(s) stuck" in result.detail 331 332 333def test_journal_caught_up_invokes_backlog_reader_and_derives_result( 334 doctor, 335 monkeypatch, 336): 337 calls = [] 338 view = BacklogView( 339 window=30, 340 days=( 341 backlog_day("20200229", "pending"), 342 backlog_day("20200301", "stuck"), 343 ), 344 pending_days=1, 345 stuck_days=1, 346 oldest_pending_day="20200229", 347 errors=(), 348 ) 349 350 def fake_read_backlog_view(): 351 calls.append("read") 352 return view 353 354 monkeypatch.setattr(pipeline_health, "read_backlog_view", fake_read_backlog_view) 355 356 result = doctor.journal_caught_up_check(args(doctor)) 357 358 assert calls == ["read"] 359 assert result.status == "warn" 360 assert "1 day(s) pending" in result.detail 361 assert "1 day(s) stuck" in result.detail 362 assert "oldest outstanding 20200229" in result.detail 363 364 365def test_journal_caught_up_warns_when_backlog_read_raises(doctor, monkeypatch): 366 def fake_read_backlog_view(): 367 raise RuntimeError("boom") 368 369 monkeypatch.setattr(pipeline_health, "read_backlog_view", fake_read_backlog_view) 370 371 result = doctor.journal_caught_up_check(args(doctor)) 372 373 assert result.status == "warn" 374 assert result.status != "fail" 375 assert result.severity == "advisory" 376 assert "backlog read failed: RuntimeError: boom" in result.detail 377 378 379def test_journal_caught_up_warns_when_backlog_reader_import_raises( 380 doctor, 381 monkeypatch, 382): 383 real_import = builtins.__import__ 384 385 def fail_pipeline_health_import(name, *args, **kwargs): 386 if name == "solstone.think.pipeline_health": 387 raise ImportError("pipeline unavailable") 388 return real_import(name, *args, **kwargs) 389 390 monkeypatch.setattr(builtins, "__import__", fail_pipeline_health_import) 391 392 result = doctor.journal_caught_up_check(args(doctor)) 393 394 assert result.status == "warn" 395 assert result.status != "fail" 396 assert result.severity == "advisory" 397 assert "backlog read failed: ImportError: pipeline unavailable" in result.detail 398 399 400def test_journal_caught_up_warn_is_visible_but_non_blocking(doctor): 401 results = [ 402 doctor.make_result(doctor.JOURNAL_CAUGHT_UP_CHECK, "warn", "..."), 403 doctor.make_result(doctor.SERVICE_RUNNING_CHECK, "ok", "..."), 404 ] 405 406 assert doctor.jsonl_summary_status(results) == "warning" 407 assert not any( 408 result.severity == "blocker" and result.status == "fail" for result in results 409 ) 410 411 412def test_journal_caught_up_is_only_in_journal_battery(doctor): 413 assert "journal_caught_up" in {check.name for check, _ in doctor.JOURNAL_CHECKS} 414 assert "journal_caught_up" not in { 415 check.name for check, _ in doctor.UNIVERSAL_CHECKS 416 } 417 assert "journal_caught_up" not in { 418 check.name for check, _ in doctor.READINESS_CHECKS 419 } 420 421 422def test_journal_caught_up_rides_existing_emission_paths(doctor, capsys): 423 warn_result = doctor.make_result(doctor.JOURNAL_CAUGHT_UP_CHECK, "warn", "...") 424 ok_result = doctor.make_result(doctor.JOURNAL_CAUGHT_UP_CHECK, "ok", "caught up") 425 426 doctor.emit_json([warn_result]) 427 assert "journal_caught_up" in capsys.readouterr().out 428 429 doctor.emit_jsonl( 430 [warn_result], 431 started_at_iso="x", 432 duration_ms=0, 433 summary_status="warning", 434 ) 435 assert "journal_caught_up" in capsys.readouterr().out 436 437 doctor.emit_text([ok_result], verbose=False) 438 assert "journal_caught_up" not in capsys.readouterr().out 439 440 doctor.emit_text([ok_result], verbose=True) 441 assert "journal_caught_up" in capsys.readouterr().out