diff --git a/solstone/think/activity_state_machine.py b/solstone/think/activity_state_machine.py index 47a5798ab..8da230be7 100644 --- a/solstone/think/activity_state_machine.py +++ b/solstone/think/activity_state_machine.py @@ -121,6 +121,7 @@ class ActivityStateMachine: def _make_completed_record(self, entry: dict) -> dict: return { "id": entry["id"], + "facet": entry["facet"], "activity": entry["activity"], "segments": entry.get("segments", [entry["since"]]), "level_avg": LEVEL_VALUES.get(entry.get("level", "medium"), 0.5), diff --git a/solstone/think/thinking.py b/solstone/think/thinking.py index bf5352767..ee4f7fa71 100644 --- a/solstone/think/thinking.py +++ b/solstone/think/thinking.py @@ -222,12 +222,16 @@ def _persist_and_maybe_run_activity_prompts( segment: str, target_schedule: str, ended_triples: list[tuple[object, object, object]], - completed_lookup: dict[object, dict], + completed: list[dict], refresh: bool, verbose: bool, max_concurrency: int, skip_activity_prompts: bool, ) -> None: + completed_by_key: dict[tuple[str, str], dict] = {} + for rec in completed: + completed_by_key.setdefault((str(rec["facet"]), str(rec["id"])), rec) + written_by: dict[tuple[str, str], bool] = {} record_by: dict[tuple[str, str], dict] = {} @@ -245,7 +249,7 @@ def _persist_and_maybe_run_activity_prompts( state="ended", change=change, ) - rec = completed_lookup.get(activity_id) + rec = completed_by_key.get(key) if rec: record_by[key] = rec written_by[key] = append_activity_record(facet_str, routing_day, rec) @@ -348,9 +352,6 @@ def _run_activity_state_tail( for c in changes if c.get("state") == "ended" ] - completed_lookup = {} - for rec in state_machine.get_completed_activities(): - completed_lookup.setdefault(rec["id"], rec) if state_machine.journal_root is not None: try: snapshot = { @@ -373,7 +374,7 @@ def _run_activity_state_tail( segment=segment, target_schedule=target_schedule, ended_triples=ended_triples, - completed_lookup=completed_lookup, + completed=state_machine.get_completed_activities(), refresh=refresh, verbose=verbose, max_concurrency=max_concurrency, @@ -412,16 +413,13 @@ def _flush_batch_state_machines( ] if not ended_triples: continue - completed_lookup: dict = {} - for rec in sm.get_completed_activities(): - completed_lookup.setdefault(rec["id"], rec) _persist_and_maybe_run_activity_prompts( routing_day=sm.last_segment_day or day, log_day=day, segment=sm.last_segment_key, target_schedule="segment", ended_triples=ended_triples, - completed_lookup=completed_lookup, + completed=sm.get_completed_activities(), refresh=refresh, verbose=verbose, max_concurrency=max_concurrency, @@ -1067,16 +1065,13 @@ def run_segment_sense( for c in idle_changes if c.get("state") == "ended" ] - completed_lookup = {} - for rec in state_machine.get_completed_activities(): - completed_lookup.setdefault(rec["id"], rec) _persist_and_maybe_run_activity_prompts( routing_day=routing_day, log_day=day, segment=segment, target_schedule=target_schedule, ended_triples=ended_triples, - completed_lookup=completed_lookup, + completed=state_machine.get_completed_activities(), refresh=refresh, verbose=verbose, max_concurrency=max_concurrency, diff --git a/tests/test_activity_state_machine.py b/tests/test_activity_state_machine.py index 6ac60c0c6..74300976a 100644 --- a/tests/test_activity_state_machine.py +++ b/tests/test_activity_state_machine.py @@ -685,6 +685,7 @@ class TestCompletedRecordFields: rec = sm.get_completed_activities()[0] required = { "id", + "facet", "activity", "segments", "level_avg", diff --git a/tests/test_think_activity.py b/tests/test_think_activity.py index 156e2e3297..f69049c48 100644 --- a/tests/test_think_activity.py +++ b/tests/test_think_activity.py @@ -809,6 +809,54 @@ class TestActivityPersistenceRoundTrip: ) return sm + def _distinct_same_id_facets(self): + return [ + {"facet": "work", "activity": "deep parser work", "level": "high"}, + {"facet": "personal", "activity": "reading docs", "level": "medium"}, + ] + + def _shared_entities(self): + return [ + {"type": "Person", "name": "Alice", "context": "pairing"}, + {"type": "Tool", "name": "VS Code", "context": "editor"}, + ] + + def _same_id_multi_facet_machine(self, segments, day): + from solstone.think.activity_state_machine import ActivityStateMachine + + sm = ActivityStateMachine() + for segment in segments: + sm.update( + self._sense( + facets=self._distinct_same_id_facets(), + entities=self._shared_entities(), + ), + segment, + day, + ) + return sm + + def _assert_distinct_same_id_records( + self, load_activity_records, day, expected_segments + ): + work_records = load_activity_records("work", day) + personal_records = load_activity_records("personal", day) + assert len(work_records) == 1 + assert len(personal_records) == 1 + + work = work_records[0] + personal = personal_records[0] + assert work["description"] == "deep parser work" + assert work["level_avg"] == 1.0 + assert work["segments"] == expected_segments + assert work["active_entities"] == ["Alice", "VS Code"] + assert personal["description"] == "reading docs" + assert personal["level_avg"] == 0.5 + assert personal["segments"] == expected_segments + assert personal["active_entities"] == ["Alice", "VS Code"] + assert work["description"] != personal["description"] + assert work["level_avg"] != personal["level_avg"] + def test_multi_segment_round_trip(self, monkeypatch): """Multi-segment activity persists and loads with all segments intact.""" from solstone.think.activities import ( @@ -994,53 +1042,44 @@ class TestActivityPersistenceRoundTrip: assert loaded["created_at"] == rec["created_at"] def test_multi_facet_ending_persists_both(self, monkeypatch): - """Multiple facets ending simultaneously all persist correctly. - - This tests the ended_pairs fix: the old facet_by_id dict would overwrite - duplicate IDs, dropping all but one facet. The list-based approach preserves - all (id, facet) pairs. - """ - from solstone.think.activities import ( - append_activity_record, - load_activity_records, - ) - from solstone.think.activity_state_machine import ActivityStateMachine + """Tail persistence keeps same-id sibling records facet-specific.""" + from solstone.think import thinking as think + from solstone.think.activities import load_activity_records with tempfile.TemporaryDirectory() as tmpdir: monkeypatch.setenv("SOLSTONE_JOURNAL", tmpdir) - two = [ - {"facet": "work", "activity": "coding", "level": "high"}, - {"facet": "personal", "activity": "browsing", "level": "medium"}, - ] - sm = ActivityStateMachine() - sm.update(self._sense(facets=two), "090000_300", "20260304") - sm.update(self._sense(facets=two), "090500_300", "20260304") - # Both end via idle - changes = sm.update(self._sense(density="idle"), "091000_300", "20260304") + day = "20260304" + sm = self._same_id_multi_facet_machine(["090000_300", "090500_300"], day) + missing_facets = self._sense(facets=[], entities=self._shared_entities()) + think._run_activity_state_tail( + sm, + missing_facets, + "091000_300", + day, + "segment", + refresh=False, + verbose=False, + max_concurrency=1, + skip_activity_prompts=True, + ) + think._run_activity_state_tail( + sm, + missing_facets, + "091500_300", + day, + "segment", + refresh=False, + verbose=False, + max_concurrency=1, + skip_activity_prompts=True, + ) - # Use the fixed ended_pairs approach (matches thinking.py) - ended_pairs = [ - (c["id"], c["facet"]) for c in changes if c.get("state") == "ended" - ] - completed_lookup = {} - for rec in sm.get_completed_activities(): - completed_lookup.setdefault(rec["id"], rec) - for activity_id, facet in ended_pairs: - rec = completed_lookup.get(activity_id) - if rec: - append_activity_record(facet, "20260304", rec) - - work_records = load_activity_records("work", "20260304") - personal_records = load_activity_records("personal", "20260304") - assert len(work_records) == 1 - assert len(personal_records) == 1 - # Both facets use top-level content_type as activity - assert work_records[0]["activity"] == "coding" - assert personal_records[0]["activity"] == "coding" - # Both have 2 segments - assert work_records[0]["segments"] == ["090000_300", "090500_300"] - assert personal_records[0]["segments"] == ["090000_300", "090500_300"] + self._assert_distinct_same_id_records( + load_activity_records, + day, + ["090000_300", "090500_300", "091000_300"], + ) def test_close_active_flushes_dangling_activity(self, monkeypatch): """Dangling active state closes to one persisted clustered record.""" @@ -1099,6 +1138,33 @@ class TestActivityPersistenceRoundTrip: assert len(records) == 1 assert records[0]["segments"] == segments + def test_flush_same_id_siblings_persist_independent_records(self, monkeypatch): + """Batch flush writes each same-id sibling's own completed record.""" + from solstone.think.activities import load_activity_records + from solstone.think.thinking import _flush_batch_state_machines + + with tempfile.TemporaryDirectory() as tmpdir: + monkeypatch.setenv("SOLSTONE_JOURNAL", tmpdir) + + day = "20260412" + segments = ["162416_300", "162916_300"] + for _ in range(2): + sm = self._same_id_multi_facet_machine(segments, day) + _flush_batch_state_machines( + {"import.audio": sm}, + day, + refresh=False, + verbose=False, + max_concurrency=1, + skip_activity_prompts=True, + ) + + self._assert_distinct_same_id_records( + load_activity_records, + day, + segments, + ) + def test_flush_is_idempotent(self, monkeypatch): """Equivalent batch flushes for the same import day dedupe by id.""" from solstone.think.activities import load_activity_records diff --git a/tests/test_think_no_activity_prompts.py b/tests/test_think_no_activity_prompts.py index 235af25d6..bf4c7a4ff 100644 --- a/tests/test_think_no_activity_prompts.py +++ b/tests/test_think_no_activity_prompts.py @@ -68,6 +68,7 @@ class EndedActivityStateMachine: return [ { "id": ACTIVITY_ID, + "facet": FACET, "activity": "coding", "segments": [SEGMENT], "level_avg": 0.5, diff --git a/tests/test_think_segment.py b/tests/test_think_segment.py index 9ff59b1e6..5ee98e006 100644 --- a/tests/test_think_segment.py +++ b/tests/test_think_segment.py @@ -359,6 +359,94 @@ class TestRunSegmentSense: "active": {}, } + def test_idle_short_circuit_persists_same_id_facet_records( + self, segment_dir, monkeypatch + ): + from solstone.think import thinking as think + from solstone.think.activities import load_activity_records + from solstone.think.activity_state_machine import ActivityStateMachine + + spawned = [] + facets = [ + {"facet": "work", "activity": "deep parser work", "level": "high"}, + {"facet": "personal", "activity": "reading docs", "level": "medium"}, + ] + entities = [ + {"type": "Person", "name": "Alice", "context": "pairing"}, + {"type": "Tool", "name": "VS Code", "context": "editor"}, + ] + + def sense(density="active", facet_payload=None): + return { + "density": density, + "content_type": "coding", + "activity_summary": "Working on parser state.", + "entities": entities, + "facets": facets if facet_payload is None else facet_payload, + "meeting_detected": False, + "speakers": [], + "recommend": {"screen_record": True}, + } + + sm = ActivityStateMachine(journal_root=segment_dir.parents[3]) + sm.update(sense(), "115000_300", "20240115") + sm.update(sense(), "115500_300", "20240115") + _write_sense_output(segment_dir, sense(density="idle", facet_payload=[])) + + monkeypatch.setattr( + think, + "get_talent_configs", + lambda schedule=None, **kwargs: _segment_configs( + "sense", "entities", "screen" + ), + ) + monkeypatch.setattr( + think, + "cortex_request", + lambda prompt, name, config=None: spawned.append(name) or f"agent-{name}", + ) + monkeypatch.setattr( + think, + "wait_for_uses", + lambda agent_ids, timeout=600: ({aid: "finish" for aid in agent_ids}, []), + ) + monkeypatch.setattr(think, "_callosum", None) + + success, failed, failed_names = think.run_segment_sense( + "20240115", + "120000_300", + refresh=False, + verbose=False, + stream="default", + state_machine=sm, + skip_activity_prompts=True, + ) + + assert spawned == ["sense"] + assert success == 1 + assert failed == 0 + assert failed_names == [] + + work_records = load_activity_records("work", "20240115") + personal_records = load_activity_records("personal", "20240115") + assert len(work_records) == 1 + assert len(personal_records) == 1 + + work = work_records[0] + personal = personal_records[0] + expected_segments = ["115000_300", "115500_300"] + expected_entities = ["Alice", "VS Code"] + assert work["description"] == "deep parser work" + assert work["level_avg"] == 1.0 + assert work["segments"] == expected_segments + assert work["active_entities"] == expected_entities + assert personal["description"] == "reading docs" + assert personal["level_avg"] == 0.5 + assert personal["segments"] == expected_segments + assert personal["active_entities"] == expected_entities + assert work["description"] != personal["description"] + assert work["level_avg"] != personal["level_avg"] + def test_conditional_screen_dispatch(self, segment_dir, monkeypatch): from solstone.think import thinking as think @@ -988,6 +1076,7 @@ class TestRunSegmentSense: return [ { "id": "coding_120000_300", + "facet": "work", "activity": "coding", "segments": ["120000_300"], "level_avg": 0.5, diff --git a/tests/test_think_skip_talents.py b/tests/test_think_skip_talents.py index 38581b6d3..ec742ab5c 100644 --- a/tests/test_think_skip_talents.py +++ b/tests/test_think_skip_talents.py @@ -177,6 +177,7 @@ class EndedActivityStateMachine: return [ { "id": ACTIVITY_ID, + "facet": FACET, "activity": "coding", "segments": [SEGMENT], "level_avg": 0.5,