From 358e77793da2fd553e5bb57774301127dfd89635 Mon Sep 17 00:00:00 2001 From: Jer Miller Date: Tue, 14 Jul 2026 11:49:45 -0600 Subject: [PATCH] Promote stream and index maintenance helpers Move stream.updated touching into streams.py and segment index row deletion into the indexer owner, then update segment move, reprocess, and importer call sites to use the shared helpers. Add metadata-preserving stream-state tail repair so prune can repair tails without clobbering stream identity or lowering seq. --- solstone/think/importers/cli.py | 14 +++++--- solstone/think/indexer/journal.py | 34 ++++++++++++++++++++ solstone/think/reprocess.py | 4 +-- solstone/think/segment.py | 53 +++---------------------------- solstone/think/streams.py | 36 +++++++++++++++++++++ tests/test_streams.py | 26 +++++++++++++++ 6 files changed, 112 insertions(+), 55 deletions(-) diff --git a/solstone/think/importers/cli.py b/solstone/think/importers/cli.py index fdc0be782..044856739 100644 --- a/solstone/think/importers/cli.py +++ b/solstone/think/importers/cli.py @@ -27,8 +27,12 @@ from solstone.think.importers.utils import save_import_segments from solstone.think.indexer.journal import index_file from solstone.think.journal_io import atomic_replace from solstone.think.media import PDF_EXTENSIONS -from solstone.think.segment import _touch_health_marker -from solstone.think.streams import stream_name, update_stream, write_segment_stream +from solstone.think.streams import ( + stream_name, + touch_stream_health_marker, + update_stream, + write_segment_stream, +) from solstone.think.utils import ( day_path, get_journal, @@ -1064,7 +1068,7 @@ def _import_one_from_args(args: argparse.Namespace) -> dict[str, Any] | None: ) for _day in sorted({seg_day for seg_day, _seg_key in result.segments}): - _touch_health_marker(_day) + touch_stream_health_marker(_day) _callosum.emit("supervisor", "drain", day=_day) logger.info( @@ -1201,7 +1205,7 @@ def _import_one_from_args(args: argparse.Namespace) -> dict[str, Any] | None: ) logger.info(f"Emitted observe.observed for segment: {day}/{seg}") - _touch_health_marker(day) + touch_stream_health_marker(day) _callosum.emit("supervisor", "drain", day=day) else: @@ -1298,7 +1302,7 @@ def _import_one_from_args(args: argparse.Namespace) -> dict[str, Any] | None: f"transcribed successfully ({total_elapsed}s)" ) - _touch_health_marker(day) + touch_stream_health_marker(day) _callosum.emit("supervisor", "drain", day=day) # Complete processing metadata diff --git a/solstone/think/indexer/journal.py b/solstone/think/indexer/journal.py index 270f1884b..9460179a0 100644 --- a/solstone/think/indexer/journal.py +++ b/solstone/think/indexer/journal.py @@ -182,6 +182,40 @@ def prune_chunks_by_stream(stream: str, journal: str | None = None) -> dict: return {"chunks": count, "files": len(paths)} +def delete_segment_index_rows( + journal: str | None, rel_path: str +) -> dict[str, int | str | None]: + """Delete index rows that reference one chronicle segment path.""" + journal_root = Path(journal or get_journal()) + db_path = journal_root / INDEX_DIR / DB_NAME + if not db_path.exists(): + return {"chunks": 0, "files": 0, "error": None} + + try: + conn = sqlite3.connect(db_path) + try: + cur = conn.execute( + "DELETE FROM chunks WHERE path = ? OR path LIKE ?", + (rel_path, f"{rel_path}/%"), + ) + chunks_deleted = cur.rowcount + + cur = conn.execute( + "DELETE FROM files WHERE path LIKE ?", + (f"{rel_path}/%",), + ) + files_deleted = cur.rowcount + + conn.commit() + finally: + conn.close() + except sqlite3.Error as exc: + logger.warning("Segment index row delete failed for %s: %s", rel_path, exc) + return {"chunks": 0, "files": 0, "error": str(exc)} + + return {"chunks": chunks_deleted, "files": files_deleted, "error": None} + + def index_file(journal: str, file_path: str, verbose: bool = False) -> bool: """Index a single file into the journal index. diff --git a/solstone/think/reprocess.py b/solstone/think/reprocess.py index 4b118e7b6..e3e8d57ed 100644 --- a/solstone/think/reprocess.py +++ b/solstone/think/reprocess.py @@ -12,7 +12,7 @@ from datetime import date, datetime from enum import Enum from solstone.think.callosum import callosum_send -from solstone.think.segment import _touch_health_marker +from solstone.think.streams import touch_stream_health_marker from solstone.think.utils import ( DATE_RE, day_is_complete, @@ -95,7 +95,7 @@ def reprocess_day(day: str, flavor: str) -> ReprocessOutcome: if flavor == FLAVOR_MARK_UPDATED: # Durable effect first: advance stream.updated so updated_days() re-queues # the day even if it already completed. Then nudge a drain. - _touch_health_marker(day) + touch_stream_health_marker(day) ok = callosum_send("supervisor", "drain", day=day) return ReprocessOutcome( ReprocessCode.MARK_UPDATED_SUBMITTED if ok else ReprocessCode.UNREACHABLE diff --git a/solstone/think/segment.py b/solstone/think/segment.py index eba97a15a..260075437 100644 --- a/solstone/think/segment.py +++ b/solstone/think/segment.py @@ -16,10 +16,12 @@ import sys from datetime import datetime, timedelta from pathlib import Path +from solstone.think.indexer.journal import delete_segment_index_rows from solstone.think.streams import ( get_stream_state, read_segment_stream, rebuild_stream_state, + touch_stream_health_marker, write_segment_stream, ) from solstone.think.utils import ( @@ -403,51 +405,6 @@ def _rewrite_events_jsonl(seg_dir: Path, new_day: str, new_segment: str) -> int: return count -def _touch_health_marker(day: str) -> None: - """Touch health/stream.updated for a day.""" - health_dir = day_path(day) / "health" - health_dir.mkdir(parents=True, exist_ok=True) - (health_dir / "stream.updated").touch() - - -def _delete_index_rows(journal: str, rel_path: str) -> dict[str, int | str | None]: - """Delete all index rows referencing a segment path. - - Returns counts of deleted rows per table. - """ - db_path = Path(journal) / "indexer" / "journal.sqlite" - if not db_path.exists(): - return {"chunks": 0, "files": 0, "error": None} - - try: - conn = sqlite3.connect(db_path) - try: - cur = conn.execute( - "DELETE FROM chunks WHERE path = ? OR path LIKE ?", - (rel_path, f"{rel_path}/%"), - ) - chunks_deleted = cur.rowcount - - cur = conn.execute( - "DELETE FROM files WHERE path LIKE ?", - (f"{rel_path}/%",), - ) - files_deleted = cur.rowcount - - conn.commit() - finally: - conn.close() - except sqlite3.Error as exc: - logger.warning("Segment index row delete failed for %s: %s", rel_path, exc) - return {"chunks": 0, "files": 0, "error": str(exc)} - - return { - "chunks": chunks_deleted, - "files": files_deleted, - "error": None, - } - - def _reindex_segment(journal: str, seg_dir: Path) -> int: """Re-index all formattable files in a segment directory. @@ -613,7 +570,7 @@ def cmd_move(args: argparse.Namespace) -> None: f" pre-move index read errored: {index_info['error']} " "(attempting delete+reindex anyway)" ) - deleted = _delete_index_rows(journal, old_rel) + deleted = delete_segment_index_rows(journal, old_rel) if deleted.get("error"): print( f" index row delete failed: {deleted['error']} " @@ -633,8 +590,8 @@ def cmd_move(args: argparse.Namespace) -> None: elif verbose: print(" index not available, skipping reindex") - _touch_health_marker(src_day) - _touch_health_marker(to_day) + touch_stream_health_marker(src_day) + touch_stream_health_marker(to_day) print(f" touched health markers: {src_day}, {to_day}") if verbose: print(" think will re-run daily talents on both days") diff --git a/solstone/think/streams.py b/solstone/think/streams.py index e7327e419..90009c9c1 100644 --- a/solstone/think/streams.py +++ b/solstone/think/streams.py @@ -162,6 +162,15 @@ def delete_stream_state(name: str) -> bool: return True +def touch_stream_health_marker(day: str) -> None: + """Touch health/stream.updated for a chronicle day.""" + from solstone.think.utils import day_path + + health_dir = day_path(day) / "health" + health_dir.mkdir(parents=True, exist_ok=True) + (health_dir / "stream.updated").touch() + + def update_stream( name: str, day: str, @@ -244,6 +253,33 @@ def update_stream( return {"prev_day": prev_day, "prev_segment": prev_segment, "seq": seq} +def repair_stream_state_tail( + name: str, last_day: str | None, last_segment: str | None, *, max_seq: int +) -> dict: + """Update a stream's tail while preserving identity metadata and monotonic seq.""" + streams_dir = _streams_dir() + state_path = streams_dir / f"{name}.json" + state = get_stream_state(name) or { + "name": name, + "type": "unknown", + "host": None, + "platform": None, + "created_at": int(time.time()), + "last_day": None, + "last_segment": None, + "seq": 0, + } + previous_seq = state.get("seq", 0) + if isinstance(previous_seq, bool) or not isinstance(previous_seq, int): + previous_seq = 0 + state["name"] = state.get("name") or name + state["last_day"] = last_day + state["last_segment"] = last_segment + state["seq"] = max(previous_seq, max_seq) + write_json(state_path, state) + return state + + def write_segment_stream( segment_dir: str | Path, stream: str, diff --git a/tests/test_streams.py b/tests/test_streams.py index 844f75287..bff2ffa02 100644 --- a/tests/test_streams.py +++ b/tests/test_streams.py @@ -13,6 +13,7 @@ from solstone.think.streams import ( list_streams, read_segment_stream, rebuild_stream_state, + repair_stream_state_tail, stream_name, update_stream, write_segment_stream, @@ -148,6 +149,31 @@ def test_delete_stream_state(tmp_path, monkeypatch): assert delete_stream_state("import.share") is False +def test_repair_stream_state_tail_preserves_metadata_and_monotonic_seq( + tmp_path, monkeypatch +): + """Tail repair preserves stream metadata and never lowers seq.""" + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + update_stream( + "archon", + "20250119", + "142500_300", + type="observer", + host="desktop", + platform="linux", + ) + before = get_stream_state("archon") + + repaired = repair_stream_state_tail("archon", "20250119", "141500_300", max_seq=1) + + assert repaired["type"] == "observer" + assert repaired["host"] == "desktop" + assert repaired["platform"] == "linux" + assert repaired["created_at"] == before["created_at"] + assert repaired["last_segment"] == "141500_300" + assert repaired["seq"] == before["seq"] + + # --- write/read segment stream tests --- -- 2.51.2