diff --git a/solstone/observe/sense.py b/solstone/observe/sense.py index 910e5391d..592d5f46d 100644 --- a/solstone/observe/sense.py +++ b/solstone/observe/sense.py @@ -568,8 +568,9 @@ class FileSensor: f"Segment fully observed{note_str}: {day}/{segment} ({duration}s)" ) - # Touch stream.updated marker for downstream consumers - if day: + # Touch stream.updated marker for downstream consumers -- live ingest + # only; batch (re-process / importer) segments must not advance it. + if day and not batch: try: health_dir = day_path(day) / "health" health_dir.mkdir(parents=True, exist_ok=True) diff --git a/solstone/think/importers/cli.py b/solstone/think/importers/cli.py index 8a982336d..d93443102 100644 --- a/solstone/think/importers/cli.py +++ b/solstone/think/importers/cli.py @@ -1136,6 +1136,7 @@ def _import_one_from_args(args: argparse.Namespace) -> dict[str, Any] | None: f"transcribed successfully ({total_elapsed}s)" ) + _touch_health_marker(day) _callosum.emit("supervisor", "drain", day=day) # Complete processing metadata diff --git a/solstone/think/thinking.py b/solstone/think/thinking.py index cb2e406ac..fa7517d32 100644 --- a/solstone/think/thinking.py +++ b/solstone/think/thinking.py @@ -3437,13 +3437,6 @@ def main() -> None: skip_talents=skip_talents, live=False, ) - # Touch stream.updated marker after each segment - try: - health_dir = day_path(day) / "health" - health_dir.mkdir(parents=True, exist_ok=True) - (health_dir / "stream.updated").touch() - except Exception: - pass batch_success += success batch_failed += failed _update_status(segments_completed=i, segments_total=total) @@ -3614,15 +3607,6 @@ def main() -> None: _run_result["success"] = success_count _run_result["failed"] = fail_count - # Touch stream.updated marker after segment processing - if args.segment: - try: - health_dir = day_path(day) / "health" - health_dir.mkdir(parents=True, exist_ok=True) - (health_dir / "stream.updated").touch() - except Exception: - pass - # POST-PHASE: Final indexing and stats (daily only) if not args.segment: logging.info("Running post-phase: indexer rescan") diff --git a/tests/test_importer.py b/tests/test_importer.py index 2c9983d97..d2f7c9f33 100644 --- a/tests/test_importer.py +++ b/tests/test_importer.py @@ -1462,6 +1462,15 @@ def test_import_one_skips_wait_when_disabled(tmp_path, monkeypatch): assert result.get("segments") assert "failed_segments" not in result + marker = tmp_path / "chronicle" / "20260303" / "health" / "stream.updated" + assert marker.exists() + drain_days = [ + call.kwargs["day"] + for call in callosum.emit.call_args_list + if call.args[:2] == ("supervisor", "drain") + ] + assert drain_days == ["20260303"] + def test_import_one_audio_reimport_is_deduped(tmp_path, monkeypatch): mod = importlib.import_module("solstone.think.importers.cli") diff --git a/tests/test_sense.py b/tests/test_sense.py index ba4917b04..d8f6d8e39 100644 --- a/tests/test_sense.py +++ b/tests/test_sense.py @@ -1086,6 +1086,24 @@ def test_file_sensor_observing_without_batch_emits_live_observed( assert "batch" not in observed +def test_file_sensor_live_observed_touches_stream_updated( + tmp_path, monkeypatch, mock_callosum +): + _observed_event_from_observing(tmp_path, monkeypatch, mock_callosum, batch=False) + + marker = tmp_path / "chronicle" / "20250101" / "health" / "stream.updated" + assert marker.exists() + + +def test_file_sensor_batch_observed_does_not_touch_stream_updated( + tmp_path, monkeypatch, mock_callosum +): + _observed_event_from_observing(tmp_path, monkeypatch, mock_callosum, batch=True) + + marker = tmp_path / "chronicle" / "20250101" / "health" / "stream.updated" + assert not marker.exists() + + def test_file_sensor_segment_observed_no_handlers(tmp_path, monkeypatch, mock_callosum): """Test that observe.observed is emitted immediately for segments with no matching handlers. diff --git a/tests/test_updated_days.py b/tests/test_updated_days.py index e8c956875..d6a09a26c 100644 --- a/tests/test_updated_days.py +++ b/tests/test_updated_days.py @@ -33,6 +33,17 @@ def test_updated_days_clean(tmp_path, monkeypatch): assert updated_days() == [] +def test_updated_days_stream_newer_than_daily_is_updated(tmp_path, monkeypatch): + """Day with stream.updated newer than daily.updated is updated.""" + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + day_dir = tmp_path / "chronicle" / "20260101" / "health" + day_dir.mkdir(parents=True) + (day_dir / "daily.updated").touch() + time.sleep(0.05) + (day_dir / "stream.updated").touch() + assert updated_days() == ["20260101"] + + def test_updated_days_no_stream(tmp_path, monkeypatch): """Day without stream.updated is not updated (no stream data).""" monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path))