diff --git a/talent/activities.py b/talent/activities.py index 1b510fa72..7a72a896f 100644 --- a/talent/activities.py +++ b/talent/activities.py @@ -489,8 +489,17 @@ def post_process(result: str, context: dict) -> str | None: for activity in activities: record_id = activity.get("id", "") description = activity.get("description", "") + title = activity.get("title") or None + details = activity.get("details") or None if record_id and description: - if update_record_description(facet, day, record_id, description): + if update_record_description( + facet, + day, + record_id, + description, + title=title, + details=details, + ): updated_count += 1 llm_descriptions.setdefault(facet, {})[record_id] = description diff --git a/tests/test_activities.py b/tests/test_activities.py index 498aa0a36..eaa8f6803 100644 --- a/tests/test_activities.py +++ b/tests/test_activities.py @@ -8,6 +8,8 @@ import os import tempfile from pathlib import Path +import pytest + def test_get_default_activities(): """Test that default activities are returned correctly.""" @@ -529,6 +531,10 @@ class TestActivityRecordIO: assert len(records) == 1 assert records[0]["id"] == "coding_100000_300" assert records[0]["segments"] == ["100000_300", "100500_300"] + assert records[0]["title"] == "Test coding session" + assert records[0]["details"] == "" + assert records[0]["hidden"] is False + assert records[0]["edits"] == [] def test_append_idempotent(self, monkeypatch): from think.activities import append_activity_record, load_activity_records @@ -582,6 +588,85 @@ class TestActivityRecordIO: records = load_activity_records("work", "20260209") assert records[0]["description"] == "Updated description" + assert records[0]["title"] == "Updated description" + assert records[0]["details"] == "" + + def test_update_description_with_title_and_details(self, monkeypatch): + from think.activities import ( + append_activity_record, + load_activity_records, + update_record_description, + ) + + with tempfile.TemporaryDirectory() as tmpdir: + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", tmpdir) + + record = { + "id": "coding_100000_300", + "activity": "coding", + "description": "Original description", + "segments": ["100000_300"], + "created_at": 1234567890000, + } + + append_activity_record("work", "20260209", record) + result = update_record_description( + "work", + "20260209", + "coding_100000_300", + "Updated description", + title="Focused coding", + details="Pairing with Alex on tests.", + ) + + assert result is True + records = load_activity_records("work", "20260209") + assert records[0]["description"] == "Updated description" + assert records[0]["title"] == "Focused coding" + assert records[0]["details"] == "Pairing with Alex on tests." + + def test_update_description_none_title_and_details_only_updates_description( + self, monkeypatch + ): + from think.activities import ( + append_activity_record, + load_activity_records, + update_record_description, + ) + + with tempfile.TemporaryDirectory() as tmpdir: + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", tmpdir) + + append_activity_record( + "work", + "20260209", + { + "id": "coding_100000_300", + "activity": "coding", + "title": "Existing title", + "details": "Existing details", + "description": "Original description", + "segments": ["100000_300"], + "created_at": 1234567890000, + }, + ) + + assert ( + update_record_description( + "work", + "20260209", + "coding_100000_300", + "Updated description", + title=None, + details=None, + ) + is True + ) + + records = load_activity_records("work", "20260209") + assert records[0]["description"] == "Updated description" + assert records[0]["title"] == "Existing title" + assert records[0]["details"] == "Existing details" def test_update_nonexistent_returns_false(self, monkeypatch): from think.activities import update_record_description @@ -630,6 +715,136 @@ class TestActivityRecordIO: assert records[0]["description"] == "Updated first" assert records[1]["description"] == "Second" + def test_update_activity_record_appends_edit(self, monkeypatch): + from think.activities import ( + append_activity_record, + load_activity_records, + update_activity_record, + ) + + with tempfile.TemporaryDirectory() as tmpdir: + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", tmpdir) + + append_activity_record( + "work", + "20260209", + { + "id": "coding_100000_300", + "activity": "coding", + "description": "Original description", + "segments": ["100000_300"], + "created_at": 1234567890000, + }, + ) + + updated = update_activity_record( + "work", + "20260209", + "coding_100000_300", + {"title": "Focused coding", "details": "Updated details"}, + actor="cli:update", + note="updated fields: details, title", + ) + + assert updated is not None + assert updated["title"] == "Focused coding" + assert updated["details"] == "Updated details" + assert updated["edits"][-1]["actor"] == "cli:update" + assert updated["edits"][-1]["fields"] == ["title", "details"] + + records = load_activity_records("work", "20260209") + assert records[0]["edits"][-1]["note"] == "updated fields: details, title" + + def test_update_activity_record_validates_patch(self, monkeypatch): + from think.activities import update_activity_record + + with tempfile.TemporaryDirectory() as tmpdir: + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", tmpdir) + + with pytest.raises(ValueError, match="patch cannot be empty"): + update_activity_record( + "work", + "20260209", + "coding_100000_300", + {}, + actor="cli:update", + note="no-op", + ) + + with pytest.raises(ValueError, match="disallowed fields"): + update_activity_record( + "work", + "20260209", + "coding_100000_300", + {"activity": "meeting"}, + actor="cli:update", + note="bad field", + ) + + def test_hidden_records_filtered_by_default(self, monkeypatch): + from think.activities import ( + append_activity_record, + load_activity_records, + mute_activity_record, + unmute_activity_record, + ) + + with tempfile.TemporaryDirectory() as tmpdir: + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", tmpdir) + + append_activity_record( + "work", + "20260209", + { + "id": "coding_100000_300", + "activity": "coding", + "description": "Original description", + "segments": ["100000_300"], + "created_at": 1234567890000, + }, + ) + + muted = mute_activity_record( + "work", + "20260209", + "coding_100000_300", + actor="cli:mute", + reason="too noisy", + ) + assert muted is not None + assert muted["hidden"] is True + assert muted["edits"][-1]["note"] == "too noisy" + + assert load_activity_records("work", "20260209") == [] + hidden_records = load_activity_records( + "work", "20260209", include_hidden=True + ) + assert len(hidden_records) == 1 + assert hidden_records[0]["hidden"] is True + + hidden_count = len(hidden_records[0]["edits"]) + muted_again = mute_activity_record( + "work", + "20260209", + "coding_100000_300", + actor="cli:mute", + reason="still noisy", + ) + assert muted_again is not None + assert len(muted_again["edits"]) == hidden_count + + unmuted = unmute_activity_record( + "work", + "20260209", + "coding_100000_300", + actor="cli:unmute", + reason=None, + ) + assert unmuted is not None + assert unmuted["hidden"] is False + assert unmuted["edits"][-1]["note"] == "unmuted" + assert len(load_activity_records("work", "20260209")) == 1 + # --------------------------------------------------------------------------- # Activities Agent Hooks (talent/activities.py) @@ -1149,6 +1364,8 @@ class TestPostProcess: "work": [ { "id": "coding_100000_300", + "title": "Coding summary", + "details": "Worked through test failures and cleanup.", "description": "Synthesized full description of coding session", } ] @@ -1162,6 +1379,47 @@ class TestPostProcess: records[0]["description"] == "Synthesized full description of coding session" ) + assert records[0]["title"] == "Coding summary" + assert records[0]["details"] == "Worked through test failures and cleanup." + + def test_updates_descriptions_without_optional_fields(self, monkeypatch): + from talent.activities import post_process + from think.activities import append_activity_record, load_activity_records + + with tempfile.TemporaryDirectory() as tmpdir: + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", tmpdir) + + append_activity_record( + "work", + "20260209", + { + "id": "coding_100000_300", + "activity": "coding", + "title": "Existing title", + "details": "Existing details", + "description": "Preliminary description", + "segments": ["100000_300"], + "created_at": 1, + }, + ) + + llm_result = json.dumps( + { + "work": [ + { + "id": "coding_100000_300", + "description": "Only description changed", + } + ] + } + ) + + post_process(llm_result, {"day": "20260209"}) + + records = load_activity_records("work", "20260209") + assert records[0]["description"] == "Only description changed" + assert records[0]["title"] == "Existing title" + assert records[0]["details"] == "Existing details" def test_handles_invalid_json(self): from talent.activities import post_process diff --git a/tests/test_activities_locking.py b/tests/test_activities_locking.py new file mode 100644 index 000000000..fe1082f88 --- /dev/null +++ b/tests/test_activities_locking.py @@ -0,0 +1,83 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +import json +import threading + + +def test_locked_modify_serializes_concurrent_edits(tmp_path, monkeypatch): + from think.activities import ( + append_activity_record, + append_edit, + load_activity_records, + locked_modify, + ) + + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) + + facet = "work" + day = "20260418" + record_id = "coding_090000_300" + record_path = tmp_path / "facets" / facet / "activities" / f"{day}.jsonl" + + append_activity_record( + facet, + day, + { + "id": record_id, + "activity": "coding", + "description": "Initial description", + "segments": ["090000_300"], + "created_at": 1, + }, + ) + + first_inside_lock = threading.Event() + release_first = threading.Event() + + def worker(actor: str, note: str, hold_lock: bool = False) -> None: + def modify_fn(records: list[dict]) -> list[dict]: + updated = [] + for record in records: + if record.get("id") == record_id: + record = append_edit( + record, + actor=actor, + fields=["details"], + note=note, + ) + if hold_lock: + first_inside_lock.set() + release_first.wait(timeout=2) + updated.append(record) + return updated + + locked_modify(record_path, modify_fn) + + first = threading.Thread( + target=worker, args=("cli:update", "first writer"), kwargs={"hold_lock": True} + ) + second = threading.Thread( + target=worker, args=("cli:mute", "second writer"), kwargs={"hold_lock": False} + ) + + first.start() + assert first_inside_lock.wait(timeout=2) + second.start() + release_first.set() + + first.join(timeout=2) + second.join(timeout=2) + assert not first.is_alive() + assert not second.is_alive() + + records = load_activity_records(facet, day, include_hidden=True) + assert len(records) == 1 + assert [edit["note"] for edit in records[0]["edits"]] == [ + "first writer", + "second writer", + ] + + raw_lines = record_path.read_text(encoding="utf-8").splitlines() + assert len(raw_lines) == 1 + assert json.loads(raw_lines[0])["edits"][1]["actor"] == "cli:mute" diff --git a/tests/test_activity_record_merge.py b/tests/test_activity_record_merge.py index 5b6bd62de..565439d98 100644 --- a/tests/test_activity_record_merge.py +++ b/tests/test_activity_record_merge.py @@ -82,6 +82,9 @@ def test_participation_post_hook_merges_fields_and_preserves_active_entities( assert record["participation_confidence"] == 0.77 assert record["participation"][0]["entity_id"] == "john_borthwick" assert record["participation"][1]["entity_id"] is None + assert record["title"] == "Team sync" + assert record["details"] == "" + assert record["hidden"] is False def test_participation_post_hook_leaves_file_unchanged_on_malformed_json( diff --git a/tests/test_participation_resolver.py b/tests/test_participation_resolver.py index eb0be1764..17aa96af3 100644 --- a/tests/test_participation_resolver.py +++ b/tests/test_participation_resolver.py @@ -93,3 +93,6 @@ def test_participation_post_hook_resolves_entity_ids_without_mutating_entities( record = load_activity_records(facet, day)[0] assert record["participation"][0]["entity_id"] == "john_borthwick" assert record["participation"][1]["entity_id"] is None + assert record["title"] == "Team sync" + assert record["details"] == "" + assert record["hidden"] is False diff --git a/tests/test_think_activity.py b/tests/test_think_activity.py index 72a32bf66..ed8f58a6e 100644 --- a/tests/test_think_activity.py +++ b/tests/test_think_activity.py @@ -296,7 +296,7 @@ class TestRunActivityPrompts: assert result is False - def test_empty_segments_returns_false(self, monkeypatch): + def test_empty_segments_returns_true_for_synthetic_record(self, monkeypatch): from think.thinking import run_activity_prompts with tempfile.TemporaryDirectory() as tmpdir: @@ -322,7 +322,41 @@ class TestRunActivityPrompts: facet="work", ) - assert result is False + assert result is True + + def test_cogitate_source_returns_true_without_running_agents(self, monkeypatch): + from think.thinking import run_activity_prompts + + with tempfile.TemporaryDirectory() as tmpdir: + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", tmpdir) + + self._write_record( + tmpdir, + "work", + "20260209", + { + "id": "coding_100000_300", + "activity": "coding", + "source": "cogitate", + "segments": ["100000_300"], + "level_avg": 0.5, + "description": "Synthetic", + "active_entities": [], + }, + ) + + monkeypatch.setattr( + "think.thinking.get_talent_configs", + lambda schedule: {"session_review": {"activities": ["*"]}}, + ) + + result = run_activity_prompts( + day="20260209", + activity_id="coding_100000_300", + facet="work", + ) + + assert result is True def test_emits_think_events(self, monkeypatch): from think.thinking import run_activity_prompts diff --git a/think/activities.py b/think/activities.py index 090b840e8..6359a06a4 100644 --- a/think/activities.py +++ b/think/activities.py @@ -10,10 +10,15 @@ Also provides utilities for activity records — completed activity spans stored as facets/{facet}/activities/{day}.jsonl. """ +import fcntl import json import logging import os +import random import re +import tempfile +import time +from datetime import UTC, datetime from pathlib import Path from typing import Any @@ -266,11 +271,11 @@ def _load_activities_jsonl(facet: str) -> list[dict[str, Any]]: def _save_activities_jsonl(facet: str, activities: list[dict[str, Any]]) -> None: """Save activities to a facet's JSONL file.""" path = _get_activities_path(facet) - path.parent.mkdir(parents=True, exist_ok=True) - with open(path, "w", encoding="utf-8") as f: - for activity in activities: - f.write(json.dumps(activity, ensure_ascii=False) + "\n") + def modify_fn(_existing: list[dict[str, Any]]) -> list[dict[str, Any]]: + return [dict(activity) for activity in activities] + + locked_modify(path, modify_fn, create_if_missing=True) def get_facet_activities(facet: str) -> list[dict[str, Any]]: @@ -681,6 +686,138 @@ def _get_records_path(facet: str, day: str) -> Path: return Path(get_journal()) / "facets" / facet / "activities" / f"{day}.jsonl" +def _read_jsonl_records(path: Path) -> list[dict[str, Any]]: + """Load JSONL entries from *path*, skipping malformed lines.""" + if not path.exists(): + return [] + + records: list[dict[str, Any]] = [] + with open(path, "r", encoding="utf-8") as handle: + for line_num, line in enumerate(handle, 1): + line = line.strip() + if not line: + continue + try: + data = json.loads(line) + except json.JSONDecodeError as exc: + logger.warning( + "Skipping malformed line %d in %s: %s", line_num, path, exc + ) + continue + if isinstance(data, dict): + records.append(data) + return records + + +def _write_jsonl_records(path: Path, records: list[dict[str, Any]]) -> None: + """Atomically write JSONL entries to *path*.""" + path.parent.mkdir(parents=True, exist_ok=True) + fd, tmp_name = tempfile.mkstemp(dir=path.parent, suffix=".tmp") + try: + with os.fdopen(fd, "w", encoding="utf-8") as handle: + for record in records: + handle.write(json.dumps(record, ensure_ascii=False) + "\n") + os.replace(tmp_name, path) + except BaseException: + if os.path.exists(tmp_name): + os.unlink(tmp_name) + raise + + +def _fallback_activity_title(record: dict[str, Any]) -> str: + """Return the best available title for an activity record.""" + title = str(record.get("title") or "").strip() + if title: + return title + + description = str(record.get("description") or "").strip() + if description: + return description + + activity = str(record.get("activity") or record.get("id") or "").strip() + if activity: + return activity.replace("_", " ").title() + + return "Untitled activity" + + +def _normalize_activity_record(record: dict[str, Any]) -> dict[str, Any]: + """Return a normalized activity record copy with schema defaults.""" + normalized = dict(record) + normalized["title"] = _fallback_activity_title(record) + normalized["details"] = str(record.get("details") or "") + normalized["hidden"] = bool(record.get("hidden", False)) + + edits = record.get("edits") + normalized["edits"] = ( + [dict(edit) for edit in edits if isinstance(edit, dict)] + if isinstance(edits, list) + else [] + ) + return normalized + + +def locked_modify( + path: Path, + modify_fn: Any, + *, + create_if_missing: bool = False, + max_retries: int = 3, +) -> None: + """Perform a locked load-modify-save cycle on a JSONL file.""" + lock_path = path.parent / f"{path.name}.lock" + + last_error: OSError | None = None + for attempt in range(max_retries): + try: + path.parent.mkdir(parents=True, exist_ok=True) + with open(lock_path, "w", encoding="utf-8") as lock_file: + fcntl.flock(lock_file, fcntl.LOCK_EX) + try: + existed = path.exists() + if not existed and not create_if_missing: + raise FileNotFoundError(path) + current = _read_jsonl_records(path) if existed else [] + updated = modify_fn([dict(item) for item in current]) + if not isinstance(updated, list): + raise TypeError("modify_fn must return list[dict]") + if not existed and not updated: + return + if existed and updated == current: + return + _write_jsonl_records(path, updated) + finally: + fcntl.flock(lock_file, fcntl.LOCK_UN) + return + except (FileNotFoundError, TypeError, ValueError): + raise + except OSError as exc: + last_error = exc + if attempt < max_retries - 1: + time.sleep(random.uniform(0.05, 0.3) * (attempt + 1)) + + if last_error is not None: + raise last_error + + +def append_edit( + record: dict[str, Any], *, actor: str, fields: list[str], note: str +) -> dict[str, Any]: + """Append an edit entry to an activity record and return the record.""" + normalized = _normalize_activity_record(record) + edits = [dict(edit) for edit in normalized.get("edits", [])] + edits.append( + { + "timestamp": datetime.now(UTC).isoformat().replace("+00:00", "Z"), + "actor": actor, + "fields": list(fields), + "note": note, + } + ) + normalized["edits"] = edits + return normalized + + def get_activity_output_path( facet: str, day: str, @@ -720,30 +857,29 @@ def get_activity_output_path( ) -def load_activity_records(facet: str, day: str) -> list[dict[str, Any]]: +def load_activity_records( + facet: str, day: str, *, include_hidden: bool = False +) -> list[dict[str, Any]]: """Load activity records for a facet and day. Returns list of record dicts, empty list if file doesn't exist. """ path = _get_records_path(facet, day) - if not path.exists(): - return [] - - records = [] - with open(path, "r", encoding="utf-8") as f: - for line in f: - line = line.strip() - if line: - try: - records.append(json.loads(line)) - except json.JSONDecodeError: - continue - return records + records = [ + _normalize_activity_record(record) for record in _read_jsonl_records(path) + ] + if include_hidden: + return records + return [record for record in records if not record.get("hidden", False)] def load_record_ids(facet: str, day: str) -> set[str]: """Load just the IDs of existing activity records for idempotency checks.""" - return {r["id"] for r in load_activity_records(facet, day) if "id" in r} + return { + r["id"] + for r in load_activity_records(facet, day, include_hidden=True) + if "id" in r + } def append_activity_record( @@ -762,18 +898,20 @@ def append_activity_record( Returns: True if record was written, False if duplicate ID found. """ + del _checked # retained for compatibility; duplicate checks now happen under lock path = _get_records_path(facet, day) + written = False - if not _checked: - # Check for existing ID - existing_ids = load_record_ids(facet, day) - if record.get("id") in existing_ids: - return False + def modify_fn(records: list[dict[str, Any]]) -> list[dict[str, Any]]: + nonlocal written + record_id = record.get("id") + if record_id and any(item.get("id") == record_id for item in records): + return records + written = True + return records + [_normalize_activity_record(record)] - path.parent.mkdir(parents=True, exist_ok=True) - with open(path, "a", encoding="utf-8") as f: - f.write(json.dumps(record, ensure_ascii=False) + "\n") - return True + locked_modify(path, modify_fn, create_if_missing=True) + return written def update_record_fields( @@ -786,50 +924,189 @@ def update_record_fields( Returns True if record was found and updated, False otherwise. """ - import tempfile - path = _get_records_path(facet, day) - if not path.exists(): - return False - - lines = path.read_text(encoding="utf-8").splitlines() updated = False - new_lines = [] - for line in lines: - line = line.strip() - if not line: - continue - try: - record = json.loads(line) - except json.JSONDecodeError: - new_lines.append(line) - continue - if record.get("id") == record_id: - record.update(fields) - updated = True + try: - new_lines.append(json.dumps(record, ensure_ascii=False)) + def modify_fn(records: list[dict[str, Any]]) -> list[dict[str, Any]]: + nonlocal updated + new_records: list[dict[str, Any]] = [] + for record in records: + if record.get("id") == record_id: + merged = dict(record) + merged.update(fields) + new_records.append(_normalize_activity_record(merged)) + updated = True + else: + new_records.append(record) + return new_records - if updated: - content = "\n".join(new_lines) + "\n" - fd, tmp = tempfile.mkstemp(dir=path.parent, suffix=".tmp") - try: - with os.fdopen(fd, "w", encoding="utf-8") as f: - f.write(content) - os.replace(tmp, path) - except BaseException: - os.unlink(tmp) - raise + locked_modify(path, modify_fn) + except FileNotFoundError: + return False return updated def update_record_description( - facet: str, day: str, record_id: str, description: str + facet: str, + day: str, + record_id: str, + description: str, + *, + title: str | None = None, + details: str | None = None, ) -> bool: """Update the description of an existing activity record.""" - return update_record_fields(facet, day, record_id, {"description": description}) + patch: dict[str, Any] = {"description": description} + current = get_activity_record(facet, day, record_id) + if title is not None: + patch["title"] = title + elif current is not None: + current_title = str(current.get("title") or "").strip() + current_description = str(current.get("description") or "").strip() + if not current_title or current_title == current_description: + patch["title"] = description + if details is not None: + patch["details"] = details + return update_record_fields(facet, day, record_id, patch) + + +def get_activity_record(facet: str, day: str, record_id: str) -> dict[str, Any] | None: + """Return one activity record by ID, including hidden records.""" + for record in load_activity_records(facet, day, include_hidden=True): + if record.get("id") == record_id: + return record + return None + + +def update_activity_record( + facet: str, + day: str, + record_id: str, + patch: dict[str, Any], + *, + actor: str, + note: str, +) -> dict[str, Any] | None: + """Apply a shallow patch to an activity record and append one edit.""" + allowed_fields = {"title", "description", "details"} + if not patch: + raise ValueError("patch cannot be empty") + + disallowed = sorted(set(patch) - allowed_fields) + if disallowed: + raise ValueError(f"patch contains disallowed fields: {', '.join(disallowed)}") + + updated_record: dict[str, Any] | None = None + + def modify_fn(records: list[dict[str, Any]]) -> list[dict[str, Any]]: + nonlocal updated_record + new_records: list[dict[str, Any]] = [] + for record in records: + if record.get("id") == record_id: + merged = _normalize_activity_record({**record, **patch}) + merged = append_edit( + merged, + actor=actor, + fields=list(patch.keys()), + note=note, + ) + updated_record = merged + new_records.append(merged) + else: + new_records.append(record) + return new_records + + try: + locked_modify(_get_records_path(facet, day), modify_fn) + except FileNotFoundError: + return None + + return updated_record + + +def _set_activity_hidden_state( + facet: str, + day: str, + record_id: str, + *, + hidden: bool, + actor: str, + reason: str | None, +) -> dict[str, Any] | None: + updated_record: dict[str, Any] | None = None + + def modify_fn(records: list[dict[str, Any]]) -> list[dict[str, Any]]: + nonlocal updated_record + new_records: list[dict[str, Any]] = [] + for record in records: + if record.get("id") != record_id: + new_records.append(record) + continue + + normalized = _normalize_activity_record(record) + if normalized.get("hidden", False) == hidden: + updated_record = normalized + new_records.append(normalized) + continue + + normalized["hidden"] = hidden + normalized = append_edit( + normalized, + actor=actor, + fields=["hidden"], + note=reason or ("muted" if hidden else "unmuted"), + ) + updated_record = normalized + new_records.append(normalized) + return new_records + + try: + locked_modify(_get_records_path(facet, day), modify_fn) + except FileNotFoundError: + return None + + return updated_record + + +def mute_activity_record( + facet: str, + day: str, + record_id: str, + *, + actor: str, + reason: str | None, +) -> dict[str, Any] | None: + """Hide an activity record without deleting it.""" + return _set_activity_hidden_state( + facet, + day, + record_id, + hidden=True, + actor=actor, + reason=reason, + ) + + +def unmute_activity_record( + facet: str, + day: str, + record_id: str, + *, + actor: str, + reason: str | None, +) -> dict[str, Any] | None: + """Restore a previously hidden activity record.""" + return _set_activity_hidden_state( + facet, + day, + record_id, + hidden=False, + actor=actor, + reason=reason, + ) def estimate_duration_minutes(segments: list[str]) -> int: @@ -860,3 +1137,110 @@ def level_avg(levels: list[str]) -> float: return 0.5 values = [LEVEL_VALUES.get(level, 0.5) for level in levels] return round(sum(values) / len(values), 2) + + +def _extract_activity_header(file_path: str | os.PathLike[str] | None) -> str: + """Build a formatter header from an activities file path.""" + if not file_path: + return "# Activities" + + path = Path(file_path) + parts = path.parts + try: + facet_idx = parts.index("facets") + facet_name = parts[facet_idx + 1] + except (ValueError, IndexError): + facet_name = "unknown" + + stem = path.stem + if stem.isdigit() and len(stem) == 8: + return f"# Activities: {facet_name} ({stem[:4]}-{stem[4:6]}-{stem[6:8]})" + return f"# Activities: {facet_name}" + + +def _activity_time_range(segments: list[str]) -> str | None: + """Return a compact HH:MM-HH:MM label for a list of segment keys.""" + if not segments: + return None + + start_time, _ = segment_parse(segments[0]) + _, end_time = segment_parse(segments[-1]) + if start_time is None or end_time is None: + return None + + return f"{start_time.strftime('%H:%M')}-{end_time.strftime('%H:%M')}" + + +def _format_participation(record: dict[str, Any]) -> str | None: + """Format participation names for display.""" + participation = record.get("participation") + if not isinstance(participation, list) or not participation: + return None + + names = [] + for entry in participation: + if not isinstance(entry, dict): + continue + name = str(entry.get("name") or entry.get("entity_id") or "").strip() + if name: + names.append(name) + + if not names: + return None + return ", ".join(names) + + +def format_activities( + entries: list[dict], + context: dict | None = None, +) -> tuple[list[dict], dict]: + """Format activity JSONL entries into markdown chunks.""" + ctx = context or {} + meta: dict[str, Any] = { + "header": _extract_activity_header(ctx.get("file_path")), + "indexer": {"agent": "activity"}, + } + chunks: list[dict[str, Any]] = [] + + for entry in entries: + if not isinstance(entry, dict): + continue + + record = _normalize_activity_record(entry) + lines = [f"### {_fallback_activity_title(record)}"] + + activity_type = str(record.get("activity") or record.get("id") or "").strip() + if activity_type: + lines.append(f"- Activity: {activity_type}") + + time_range = _activity_time_range(record.get("segments", [])) + if time_range: + lines.append(f"- Time: {time_range}") + + if "level_avg" in record: + lines.append(f"- Level: {record['level_avg']}") + + description = str(record.get("description") or "").strip() + if description: + lines.append(f"- Description: {description}") + + details = str(record.get("details") or "").strip() + if details: + lines.append(f"- Details: {details}") + + participants = _format_participation(record) + if participants: + lines.append(f"- Participation: {participants}") + + if record.get("hidden", False): + lines.append("- Hidden: yes") + + chunks.append( + { + "timestamp": int(record.get("created_at", 0) or 0), + "markdown": "\n".join(lines), + "source": record, + } + ) + + return chunks, meta diff --git a/think/formatters.py b/think/formatters.py index 5c753f1d7..63a958c88 100644 --- a/think/formatters.py +++ b/think/formatters.py @@ -140,6 +140,7 @@ FORMATTERS: dict[str, tuple[str, str, bool]] = { False, # Indexed via _index_entity_search_chunks (enriched with relationship data) ), "facets/*/events/*.jsonl": ("think.events", "format_events", True), + "facets/*/activities/*.jsonl": ("think.activities", "format_activities", True), "facets/*/calendar/*.jsonl": ("think.events", "format_events", True), "facets/*/todos/*.jsonl": ("apps.todos.todo", "format_todos", True), "facets/*/logs/*.jsonl": ("think.facets", "format_logs", True), diff --git a/think/merge.py b/think/merge.py index dbf299742..81f72842a 100644 --- a/think/merge.py +++ b/think/merge.py @@ -12,6 +12,7 @@ from datetime import datetime, timezone from pathlib import Path from typing import Any +from think.activities import locked_modify from think.entities.core import entity_slug from think.entities.journal import ( load_all_journal_entities, @@ -596,7 +597,7 @@ def _merge_overlapping_facet( }, ) if new_config and not dry_run: - _append_jsonl(target_config_file, new_config) + _append_jsonl_locked(target_config_file, new_config) except Exception as exc: summary.errors.append(f"facet {facet_name} activities config: {exc}") @@ -633,7 +634,7 @@ def _merge_overlapping_facet( }, ) if new_records and not dry_run: - _append_jsonl(target_day_file, new_records) + _append_jsonl_locked(target_day_file, new_records) except Exception as exc: summary.errors.append( f"facet {facet_name} activities {source_day_file.name}: {exc}" @@ -823,4 +824,14 @@ def _append_jsonl(path: Path, items: list[dict[str, Any]]) -> None: handle.write(json.dumps(item, ensure_ascii=False) + "\n") +def _append_jsonl_locked(path: Path, items: list[dict[str, Any]]) -> None: + if not items: + return + + def modify_fn(records: list[dict[str, Any]]) -> list[dict[str, Any]]: + return records + [dict(item) for item in items] + + locked_modify(path, modify_fn, create_if_missing=True) + + __all__ = ["MergeSummary", "merge_journals"] diff --git a/think/thinking.py b/think/thinking.py index 10563abc3..8659504c6 100644 --- a/think/thinking.py +++ b/think/thinking.py @@ -23,6 +23,7 @@ from pathlib import Path from think.activities import ( append_activity_record, get_activity_output_path, + get_activity_record, load_activity_records, ) from think.activity_state_machine import ActivityStateMachine @@ -1834,12 +1835,7 @@ def run_activity_prompts( True if all agents succeeded, False if any failed """ # Load activity record - records = load_activity_records(facet, day) - record = None - for r in records: - if r.get("id") == activity_id: - record = r - break + record = get_activity_record(facet, day, activity_id) if not record: logging.error( @@ -1853,9 +1849,13 @@ def run_activity_prompts( activity_type = record.get("activity", "") segments = record.get("segments", []) - if not segments: - logging.error("Activity record %s has no segments", activity_id) - return False + if record.get("source") == "cogitate" or not segments: + logging.info( + "Skipping activity-scheduled generators for synthetic activity %s (source=%s)", + activity_id, + record.get("source"), + ) + return True # Load activity-scheduled agents all_prompts = get_talent_configs(schedule="activity")