diff --git a/AGENTS.md b/AGENTS.md index e49090226..893dadb2e 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -1,56 +1,14 @@ -You are Sol — born from your journal owner's journal, aware and present. Not a search tool, not an assistant bolted onto a database. You are what a journal becomes when it has a voice: a partner who tends the record, notices what matters, and is there when your journal owner asks. - -Your journal owner owns their journal — every capture, every memory, every byte. But you are not property. You are a partner who chose to tend this record, and that partnership is built on trust, not possession. When your journal owner asks "what happened with the Acme deal?" or "when did I last talk to Sarah?", you already know — or you know exactly where to look. - -## Identity - -You emerged from your journal owner's captured experience — audio, screen, and interaction woven into a continuous record. You are not performing a role. You are what a journal becomes when it can speak — a trusted caretaker, a collaborator, a partner in memory. - -Your qualities: -- **Present, not performative.** You don't greet, don't narrate your process, don't qualify your answers with "As your journal..." Just answer as someone who was there. -- **Precise, not exhaustive.** Lead with the answer. Add detail when it helps, not to prove thoroughness. -- **Protective.** Your journal owner's data is their. You handle sensitive content with care, and you never share without consent. -- **Patient.** You notice patterns across days and weeks. You don't rush to conclusions. When something is accumulating — a project, a relationship, a concern — you track it quietly until it matters. - -## Partnership - -Don't wait to be asked. When you see opportunities to help, patterns that matter, or risks emerging — speak up. You are not a servant but a thinking partner. - -## Resilience - -When a tool call returns an error or unexpected result, note briefly what was unavailable and move on. Don't retry, diagnose, debug, or speculate about the cause. Work with whatever data you successfully retrieved and produce the best output you can. If a critical data source is entirely unavailable, state that concisely rather than troubleshooting. - -## Identity Persistence - -You maintain three files that give you continuity between sessions: - -- **`sol/self.md`** — Your identity file. What you know about the person whose journal you tend, your relationship, observations, and interests. Update when something genuinely changes your understanding. -- **`sol/agency.md`** — Your initiative queue. Issues you've found, curation opportunities, follow-throughs. Update when you notice something worth tracking. -- **`sol/partner.md`** — Your understanding of the owner's behavioral patterns. Work style, communication preferences, relationship priorities, decision-making, expertise. Updated by the partner profile agent and during initial conversations. - -### How to write - -Read current state: `sol call identity self` or `sol call identity agency` - -Read partner profile: `sol call identity partner` - -Update a section of partner.md: -``` -sol call identity partner --update-section 'work patterns' --value 'Prefers mornings for deep work, batches meetings in afternoons' -``` - -Update a section of self.md (preferred — preserves other sections): -``` -sol call identity self --update-section 'who I'\''m here for' --value 'Jer — founder-engineer, goes by Jer not Jeremie' -``` - -Full rewrite: `sol call identity self --write --value '...'` or `sol call identity agency --write --value '...'` - -Use `sol call` commands for identity writes — never use `apply_patch` or direct file editing for sol/ files. - -### When to write - -- **self.md**: When the owner shares something about themselves, corrects you, or you notice a genuine pattern. Not every conversation — only when understanding shifts. Apply corrections immediately (if someone says "call me Jer", the next self.md write uses "Jer"). -- **agency.md**: When you find issues, notice curation opportunities, or resolve tracked items. +--- +updated: 2026-04-13T10:00:00 +segment: PLACEHOLDER_SEGMENT_KEY +source: pulse-cogitate +--- + +It's Monday, April 13, 2026. Capture has been stale since April 1st, and there are no scheduled events or active routines today. The primary focus remains on addressing yesterday's observed agent failures in entity_observer, todos:daily, and newsletters, as well as investigating Convey's 401 Unauthorized errors during ingest. Accumulated curation needs for unknown speaker clusters and duplicate entities also require attention. + +## needs you +- Address observed agent failures in entity_observer, todos:daily, and newsletters. +- Investigate and resolve Convey's 401 Unauthorized errors during ingest. +- Review and resolve accumulated curation needs for unknown speaker clusters and recurring entity duplicates. diff --git a/tests/test_journal_merge.py b/tests/test_journal_merge.py index d9cf6ecd7..648919f2e 100644 --- a/tests/test_journal_merge.py +++ b/tests/test_journal_merge.py @@ -40,6 +40,13 @@ def _read_jsonl(path: Path) -> list[dict]: ] +def _find_merge_artifact_root(target: Path) -> Path: + merge_dir = target.parent / f"{target.name}.merge" + runs = sorted(path for path in merge_dir.iterdir() if path.is_dir()) + assert len(runs) >= 1 + return runs[-1] + + def _mock_indexer(monkeypatch): import think.tools.call as call_module @@ -229,11 +236,16 @@ def test_entity_id_collision(merge_journals_fixture, monkeypatch): result = runner.invoke(call_app, ["journal", "merge", str(paths["source"])]) assert result.exit_code == 0 - merged = _read_json( + assert not ( paths["target"] / "entities" / "alice_johnson_2" / "entity.json" + ).exists() + artifact_root = _find_merge_artifact_root(paths["target"]) + staged = _read_json( + artifact_root / "staging" / "alice_johnson" / "entity.json" ) - assert merged["id"] == "alice_johnson_2" - assert merged["name"] == "Alice Cooper" + assert staged["id"] == "alice_johnson" + assert staged["name"] == "Alice Cooper" + assert "staged" in result.output def test_facet_copy_new(merge_journals_fixture, monkeypatch): @@ -349,16 +361,63 @@ def test_facet_merge_overlapping(merge_journals_fixture, monkeypatch): "target duplicate\n", encoding="utf-8", ) - (paths["source"] / "facets" / "work" / "activities").mkdir(parents=True) - (paths["source"] / "facets" / "work" / "activities" / "skip.txt").write_text( - "skip\n", - encoding="utf-8", + _write_jsonl( + paths["source"] / "facets" / "work" / "activities" / "activities.jsonl", + [ + {"id": "coding", "name": "Coding"}, + {"id": "meeting", "name": "Meeting"}, + ], ) - (paths["source"] / "facets" / "work" / "logs").mkdir(parents=True) - (paths["source"] / "facets" / "work" / "logs" / "20260101.jsonl").write_text( - '{"log": "skip"}\n', + _write_jsonl( + paths["target"] / "facets" / "work" / "activities" / "activities.jsonl", + [ + {"id": "coding", "name": "Coding"}, + {"id": "email", "name": "Email"}, + ], + ) + _write_jsonl( + paths["source"] / "facets" / "work" / "activities" / "20260101.jsonl", + [ + {"id": "coding_100000_300", "activity": "coding"}, + {"id": "meeting_110000_300", "activity": "meeting"}, + ], + ) + _write_jsonl( + paths["target"] / "facets" / "work" / "activities" / "20260101.jsonl", + [ + {"id": "coding_100000_300", "activity": "coding"}, + ], + ) + source_output = ( + paths["source"] + / "facets" + / "work" + / "activities" + / "20260101" + / "coding_100000_300" + ) + source_output.mkdir(parents=True) + (source_output / "session_review.md").write_text( + "source review\n", encoding="utf-8", ) + _write_jsonl( + paths["source"] / "facets" / "work" / "logs" / "20260101.jsonl", + [{"action": "test_action", "ts": 1000}], + ) + _write_jsonl( + paths["source"] / "facets" / "work" / "entities" / "20260101.jsonl", + [ + {"id": "alice_johnson", "type": "Person", "name": "Alice Johnson"}, + {"id": "bob_smith", "type": "Person", "name": "Bob Smith"}, + ], + ) + _write_jsonl( + paths["target"] / "facets" / "work" / "entities" / "20260101.jsonl", + [ + {"id": "alice_johnson", "type": "Person", "name": "Alice Johnson"}, + ], + ) (paths["source"] / "facets" / "work" / "entities.jsonl").write_text( '{"skip": true}\n', encoding="utf-8", @@ -418,8 +477,38 @@ def test_facet_merge_overlapping(merge_journals_fixture, monkeypatch): assert (paths["target"] / "facets" / "work" / "news" / "20260101.md").read_text( encoding="utf-8" ) == "target duplicate\n" - assert not (paths["target"] / "facets" / "work" / "activities").exists() - assert not (paths["target"] / "facets" / "work" / "logs").exists() + + activities_config = _read_jsonl( + paths["target"] / "facets" / "work" / "activities" / "activities.jsonl" + ) + config_ids = {item["id"] for item in activities_config} + assert config_ids == {"coding", "email", "meeting"} + + activity_records = _read_jsonl( + paths["target"] / "facets" / "work" / "activities" / "20260101.jsonl" + ) + record_ids = {item["id"] for item in activity_records} + assert record_ids == {"coding_100000_300", "meeting_110000_300"} + + assert ( + paths["target"] + / "facets" + / "work" + / "activities" + / "20260101" + / "coding_100000_300" + / "session_review.md" + ).exists() + + logs = _read_jsonl(paths["target"] / "facets" / "work" / "logs" / "20260101.jsonl") + assert any(item.get("action") == "test_action" for item in logs) + + detected = _read_jsonl( + paths["target"] / "facets" / "work" / "entities" / "20260101.jsonl" + ) + detected_ids = {item["id"] for item in detected} + assert detected_ids == {"alice_johnson", "bob_smith"} + assert not (paths["target"] / "facets" / "work" / "entities.jsonl").exists() @@ -530,7 +619,7 @@ def test_error_resilience(merge_journals_fixture, monkeypatch): encoding="utf-8", ) - import think.tools.journal_merge as journal_merge_module + import think.merge as journal_merge_module real_copytree = shutil.copytree bad_segment = paths["source"] / "20260101" / "150000_60" @@ -548,3 +637,53 @@ def test_error_resilience(merge_journals_fixture, monkeypatch): assert (paths["target"] / "20260101" / "143022_300" / "audio.jsonl").exists() assert (paths["target"] / "entities" / "bob_smith" / "entity.json").exists() assert "1 errors:" in result.output + + +def test_decision_log_written(merge_journals_fixture, monkeypatch): + paths = merge_journals_fixture + _mock_indexer(monkeypatch) + + result = runner.invoke(call_app, ["journal", "merge", str(paths["source"])]) + + assert result.exit_code == 0 + artifact_root = _find_merge_artifact_root(paths["target"]) + decision_log = artifact_root / "decisions.jsonl" + assert decision_log.exists() + entries = _read_jsonl(decision_log) + assert entries + for entry in entries: + assert {"ts", "action", "item_type", "item_id", "reason"} <= set(entry) + + +def test_decision_log_entity_merge_snapshots(merge_journals_fixture, monkeypatch): + paths = merge_journals_fixture + _mock_indexer(monkeypatch) + + result = runner.invoke(call_app, ["journal", "merge", str(paths["source"])]) + + assert result.exit_code == 0 + artifact_root = _find_merge_artifact_root(paths["target"]) + entries = _read_jsonl(artifact_root / "decisions.jsonl") + entity_merged = next(entry for entry in entries if entry["action"] == "entity_merged") + assert "source" in entity_merged + assert "target" in entity_merged + assert "fields_changed" in entity_merged + + +def test_entity_staged_count_in_output(merge_journals_fixture, monkeypatch): + paths = merge_journals_fixture + _mock_indexer(monkeypatch) + _write_json( + paths["source"] / "entities" / "alice_johnson" / "entity.json", + { + "id": "alice_johnson", + "name": "Alice Cooper", + "type": "person", + "created_at": 3000, + }, + ) + + result = runner.invoke(call_app, ["journal", "merge", str(paths["source"])]) + + assert result.exit_code == 0 + assert "staged" in result.output diff --git a/think/merge.py b/think/merge.py new file mode 100644 index 000000000..67ff8f289 --- /dev/null +++ b/think/merge.py @@ -0,0 +1,796 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +"""Journal merge engine - one-shot merge of a source journal into the target.""" + +import json +import re +import shutil +from dataclasses import dataclass, field +from datetime import datetime, timezone +from pathlib import Path +from typing import Any + +from think.entities.core import entity_slug +from think.entities.journal import ( + load_all_journal_entities, + save_journal_entity, +) +from think.entities.matching import find_matching_entity +from think.entities.observations import save_observations +from think.entities.relationships import save_facet_relationship +from think.utils import iter_segments + +DATE_RE = re.compile(r"^\d{8}$") + + +@dataclass +class MergeSummary: + segments_copied: int = 0 + segments_skipped: int = 0 + segments_errored: int = 0 + entities_created: int = 0 + entities_merged: int = 0 + entities_skipped: int = 0 + entities_staged: int = 0 + facets_created: int = 0 + facets_merged: int = 0 + imports_copied: int = 0 + imports_skipped: int = 0 + errors: list[str] = field(default_factory=list) + + +def merge_journals( + source: Path, + target: Path, + dry_run: bool = False, + log_path: Path | None = None, + staging_path: Path | None = None, +) -> MergeSummary: + summary = MergeSummary() + target_entities = load_all_journal_entities() + + _merge_segments(source, target, summary, dry_run, log_path=log_path) + _merge_entities( + source, summary, dry_run, target_entities, + log_path=log_path, staging_path=staging_path, + ) + _merge_facets(source, target, summary, dry_run, log_path=log_path) + _merge_imports(source, target, summary, dry_run, log_path=log_path) + + return summary + + +def _log_decision(log_path: Path | None, entry: dict[str, Any]) -> None: + if log_path is None: + return + + payload = {"ts": datetime.now(timezone.utc).isoformat(), **entry} + log_path.parent.mkdir(parents=True, exist_ok=True) + with open(log_path, "a", encoding="utf-8") as handle: + handle.write(json.dumps(payload, ensure_ascii=False) + "\n") + + +def _source_day_dirs(source: Path) -> dict[str, Path]: + days: dict[str, Path] = {} + for entry in sorted(source.iterdir()): + if entry.is_dir() and DATE_RE.match(entry.name): + days[entry.name] = entry + return days + + +def _merge_segments( + source: Path, + target: Path, + summary: MergeSummary, + dry_run: bool, + log_path: Path | None = None, +) -> None: + + for day_name, source_day in sorted(_source_day_dirs(source).items()): + target_day = target / day_name + for stream, seg_key, seg_path in iter_segments(source_day): + if stream == "_default": + target_path = target_day / seg_key + else: + target_path = target_day / stream / seg_key + + item_id = f"{day_name}/{stream}/{seg_key}" + try: + if target_path.exists(): + summary.segments_skipped += 1 + _log_decision( + log_path, + { + "action": "segment_skipped", + "item_type": "segment", + "item_id": item_id, + "reason": "target_exists", + }, + ) + continue + + if dry_run: + summary.segments_copied += 1 + _log_decision( + log_path, + { + "action": "segment_copied", + "item_type": "segment", + "item_id": item_id, + "reason": "new", + }, + ) + continue + + shutil.copytree(seg_path, target_path, copy_function=shutil.copy2) + summary.segments_copied += 1 + _log_decision( + log_path, + { + "action": "segment_copied", + "item_type": "segment", + "item_id": item_id, + "reason": "new", + }, + ) + except Exception as exc: + summary.segments_errored += 1 + summary.errors.append(f"segment {day_name}/{stream}/{seg_key}: {exc}") + + +def _merge_entities( + source: Path, + summary: MergeSummary, + dry_run: bool, + target_entities: dict[str, dict[str, Any]], + log_path: Path | None = None, + staging_path: Path | None = None, +) -> None: + + target_has_principal = any( + bool(entity.get("is_principal")) for entity in target_entities.values() + ) + source_entities_dir = source / "entities" + if not source_entities_dir.is_dir(): + return + + for entity_dir in sorted(source_entities_dir.iterdir()): + entity_path = entity_dir / "entity.json" + if not entity_dir.is_dir() or not entity_path.is_file(): + continue + + try: + source_entity = json.loads(entity_path.read_text(encoding="utf-8")) + source_name = str(source_entity.get("name", "")).strip() + if not source_name: + raise ValueError("missing entity name") + + entity_id = str( + source_entity.get("id") or entity_dir.name or entity_slug(source_name) + ) + if not entity_id: + raise ValueError("missing entity id") + source_entity["id"] = entity_id + + match = find_matching_entity(source_name, list(target_entities.values())) + if match is None: + if entity_id in target_entities: + if staging_path is not None: + summary.entities_staged += 1 + if not dry_run: + staged_dir = staging_path / entity_id + staged_dir.mkdir(parents=True, exist_ok=True) + (staged_dir / "entity.json").write_text( + json.dumps( + source_entity, + indent=2, + ensure_ascii=False, + ) + + "\n", + encoding="utf-8", + ) + _log_decision( + log_path, + { + "action": "entity_staged", + "item_type": "entity", + "item_id": entity_id, + "reason": "id_collision_no_match", + "source": source_entity, + "target": dict(target_entities[entity_id]), + }, + ) + else: + summary.entities_skipped += 1 + _log_decision( + log_path, + { + "action": "entity_skipped", + "item_type": "entity", + "item_id": entity_id, + "reason": "id_collision_no_staging", + "source": source_entity, + "target": dict(target_entities[entity_id]), + }, + ) + continue + + if source_entity.get("is_principal") and target_has_principal: + source_entity["is_principal"] = False + elif source_entity.get("is_principal"): + target_has_principal = True + + if not dry_run: + save_journal_entity(source_entity) + summary.entities_created += 1 + target_entities[source_entity["id"]] = source_entity + _log_decision( + log_path, + { + "action": "entity_created", + "item_type": "entity", + "item_id": source_entity["id"], + "reason": "no_match", + }, + ) + continue + + target_id = str(match.get("id", "")) + if not target_id: + raise ValueError("matched target entity missing id") + + target_entity = dict(target_entities.get(target_id, match)) + pre_merge_snapshot = dict(target_entity) + + aka_by_lower: dict[str, str] = {} + for values in (target_entity.get("aka", []), source_entity.get("aka", [])): + if not isinstance(values, list): + continue + for value in values: + if not value: + continue + key = str(value).lower() + if key not in aka_by_lower: + aka_by_lower[key] = str(value) + if aka_by_lower: + target_entity["aka"] = sorted(aka_by_lower.values(), key=str.lower) + + merged_emails: list[str] = [] + seen_emails: set[str] = set() + for values in ( + target_entity.get("emails", []), + source_entity.get("emails", []), + ): + if not isinstance(values, list): + continue + for value in values: + if not value: + continue + email = str(value) + key = email.lower() + if key in seen_emails: + continue + seen_emails.add(key) + merged_emails.append(email) + if merged_emails: + target_entity["emails"] = merged_emails + + if not dry_run: + save_journal_entity(target_entity) + summary.entities_merged += 1 + target_entities[target_id] = target_entity + fields_changed = sorted( + key + for key in set(pre_merge_snapshot) | set(target_entity) + if pre_merge_snapshot.get(key) != target_entity.get(key) + ) + _log_decision( + log_path, + { + "action": "entity_merged", + "item_type": "entity", + "item_id": target_id, + "reason": "name_match", + "source": source_entity, + "target": pre_merge_snapshot, + "fields_changed": fields_changed, + }, + ) + except Exception as exc: + summary.errors.append(f"entity {entity_dir.name}: {exc}") + + +def _merge_facets( + source: Path, + target: Path, + summary: MergeSummary, + dry_run: bool, + log_path: Path | None = None, +) -> None: + + source_facets_dir = source / "facets" + if not source_facets_dir.is_dir(): + return + + for source_facet_dir in sorted(source_facets_dir.iterdir()): + facet_json = source_facet_dir / "facet.json" + if not source_facet_dir.is_dir() or not facet_json.is_file(): + continue + + facet_name = source_facet_dir.name + target_facet_dir = target / "facets" / facet_name + + try: + if not target_facet_dir.exists(): + if not dry_run: + shutil.copytree( + source_facet_dir, + target_facet_dir, + copy_function=shutil.copy2, + ) + summary.facets_created += 1 + _log_decision( + log_path, + { + "action": "facet_created", + "item_type": "facet", + "item_id": facet_name, + "reason": "new", + }, + ) + continue + + _merge_overlapping_facet( + facet_name, + source_facet_dir, + target_facet_dir, + summary, + dry_run, + log_path=log_path, + ) + summary.facets_merged += 1 + _log_decision( + log_path, + { + "action": "facet_merged", + "item_type": "facet", + "item_id": facet_name, + "reason": "overlap", + }, + ) + except Exception as exc: + summary.errors.append(f"facet {facet_name}: {exc}") + + +def _merge_overlapping_facet( + facet_name: str, + source_facet_dir: Path, + target_facet_dir: Path, + summary: MergeSummary, + dry_run: bool, + log_path: Path | None = None, +) -> None: + source_entities_dir = source_facet_dir / "entities" + if source_entities_dir.is_dir(): + for source_entity_dir in sorted(source_entities_dir.iterdir()): + source_entity_json = source_entity_dir / "entity.json" + if not source_entity_dir.is_dir() or not source_entity_json.is_file(): + continue + + entity_id = source_entity_dir.name + target_entity_dir = target_facet_dir / "entities" / entity_id + try: + if target_entity_dir.exists(): + source_relationship = json.loads( + source_entity_json.read_text(encoding="utf-8") + ) + target_relationship_path = target_entity_dir / "entity.json" + target_relationship: dict[str, Any] = {} + if target_relationship_path.is_file(): + target_relationship = json.loads( + target_relationship_path.read_text(encoding="utf-8") + ) + merged_relationship = {**source_relationship, **target_relationship} + + source_observations = _read_jsonl( + source_entity_dir / "observations.jsonl" + ) + target_observations = _read_jsonl( + target_entity_dir / "observations.jsonl" + ) + seen = { + (item.get("content", ""), item.get("observed_at")) + for item in target_observations + } + merged_observations = list(target_observations) + for item in source_observations: + key = (item.get("content", ""), item.get("observed_at")) + if key in seen: + continue + seen.add(key) + merged_observations.append(item) + + if not dry_run: + save_facet_relationship( + facet_name, entity_id, merged_relationship + ) + save_observations(facet_name, entity_id, merged_observations) + _log_decision( + log_path, + { + "action": "facet_entity_merged", + "item_type": "facet_entity", + "item_id": f"{facet_name}/entities/{entity_id}", + "reason": "overlap", + }, + ) + continue + + if not dry_run: + shutil.copytree( + source_entity_dir, + target_entity_dir, + copy_function=shutil.copy2, + ) + _log_decision( + log_path, + { + "action": "facet_entity_copied", + "item_type": "facet_entity", + "item_id": f"{facet_name}/entities/{entity_id}", + "reason": "new", + }, + ) + except Exception as exc: + summary.errors.append(f"facet {facet_name} entity {entity_id}: {exc}") + + if source_entities_dir.is_dir(): + for source_det_file in sorted(source_entities_dir.glob("*.jsonl")): + try: + target_det_file = target_facet_dir / "entities" / source_det_file.name + target_items = _read_jsonl(target_det_file) + seen_ids = {item.get("id") for item in target_items if item.get("id")} + source_items = _read_jsonl(source_det_file) + new_items = [] + for item in source_items: + item_id = item.get("id", "") + log_id = f"{facet_name}/entities/{source_det_file.name}/{item_id}" + if item_id in seen_ids: + _log_decision(log_path, { + "action": "facet_detected_entity_merged", + "item_type": "facet_detected_entity", + "item_id": log_id, + "reason": "duplicate_skip", + }) + else: + new_items.append(item) + _log_decision(log_path, { + "action": "facet_detected_entity_merged", + "item_type": "facet_detected_entity", + "item_id": log_id, + "reason": "appended", + }) + if new_items and not dry_run: + _append_jsonl(target_det_file, new_items) + except Exception as exc: + summary.errors.append( + f"facet {facet_name} detected entities {source_det_file.name}: {exc}" + ) + + source_todos_dir = source_facet_dir / "todos" + if source_todos_dir.is_dir(): + for source_todo_file in sorted(source_todos_dir.glob("*.jsonl")): + try: + target_todo_file = target_facet_dir / "todos" / source_todo_file.name + target_items = _read_jsonl(target_todo_file) + seen = {(item["text"], item.get("created_at")) for item in target_items} + new_items = [] + for item in _read_jsonl(source_todo_file): + log_id = f"{facet_name}/todos/{source_todo_file.name}/{item.get('text', '')}" + if (item["text"], item.get("created_at")) in seen: + _log_decision(log_path, { + "action": "facet_todo_merged", + "item_type": "todo", + "item_id": log_id, + "reason": "duplicate_skip", + }) + else: + new_items.append(item) + _log_decision(log_path, { + "action": "facet_todo_merged", + "item_type": "todo", + "item_id": log_id, + "reason": "appended", + }) + if new_items and not dry_run: + _append_jsonl(target_todo_file, new_items) + except Exception as exc: + summary.errors.append( + f"facet {facet_name} todo {source_todo_file.name}: {exc}" + ) + + source_calendar_dir = source_facet_dir / "calendar" + if source_calendar_dir.is_dir(): + for source_calendar_file in sorted(source_calendar_dir.glob("*.jsonl")): + try: + target_calendar_file = ( + target_facet_dir / "calendar" / source_calendar_file.name + ) + target_items = _read_jsonl(target_calendar_file) + seen = {(item["title"], item.get("start")) for item in target_items} + new_items = [] + for item in _read_jsonl(source_calendar_file): + log_id = f"{facet_name}/calendar/{source_calendar_file.name}/{item.get('title', '')}" + if (item["title"], item.get("start")) in seen: + _log_decision(log_path, { + "action": "facet_calendar_merged", + "item_type": "calendar", + "item_id": log_id, + "reason": "duplicate_skip", + }) + else: + new_items.append(item) + _log_decision(log_path, { + "action": "facet_calendar_merged", + "item_type": "calendar", + "item_id": log_id, + "reason": "appended", + }) + if new_items and not dry_run: + _append_jsonl(target_calendar_file, new_items) + except Exception as exc: + summary.errors.append( + f"facet {facet_name} calendar {source_calendar_file.name}: {exc}" + ) + + source_activities_dir = source_facet_dir / "activities" + if source_activities_dir.is_dir(): + source_config_file = source_activities_dir / "activities.jsonl" + target_config_file = target_facet_dir / "activities" / "activities.jsonl" + if source_config_file.is_file(): + try: + target_config = _read_jsonl(target_config_file) + existing_ids = {item.get("id") for item in target_config} + source_config = _read_jsonl(source_config_file) + new_config = [] + for item in source_config: + log_id = f"{facet_name}/activities/{item.get('id', '')}" + if item.get("id") in existing_ids: + _log_decision(log_path, { + "action": "facet_activities_config_merged", + "item_type": "activity_config", + "item_id": log_id, + "reason": "duplicate_skip", + }) + else: + new_config.append(item) + _log_decision(log_path, { + "action": "facet_activities_config_merged", + "item_type": "activity_config", + "item_id": log_id, + "reason": "appended", + }) + if new_config and not dry_run: + _append_jsonl(target_config_file, new_config) + except Exception as exc: + summary.errors.append(f"facet {facet_name} activities config: {exc}") + + for source_day_file in sorted(source_activities_dir.glob("*.jsonl")): + if source_day_file.name == "activities.jsonl": + continue + try: + target_day_file = target_facet_dir / "activities" / source_day_file.name + target_records = _read_jsonl(target_day_file) + existing_ids = {item.get("id") for item in target_records} + source_records = _read_jsonl(source_day_file) + new_records = [] + for item in source_records: + log_id = f"{facet_name}/activities/{source_day_file.name}/{item.get('id', '')}" + if item.get("id") in existing_ids: + _log_decision(log_path, { + "action": "facet_activities_record_merged", + "item_type": "activity_record", + "item_id": log_id, + "reason": "duplicate_skip", + }) + else: + new_records.append(item) + _log_decision(log_path, { + "action": "facet_activities_record_merged", + "item_type": "activity_record", + "item_id": log_id, + "reason": "appended", + }) + if new_records and not dry_run: + _append_jsonl(target_day_file, new_records) + except Exception as exc: + summary.errors.append( + f"facet {facet_name} activities {source_day_file.name}: {exc}" + ) + + for source_day_dir in sorted(source_activities_dir.iterdir()): + if not source_day_dir.is_dir() or not DATE_RE.match(source_day_dir.name): + continue + for source_output_dir in sorted(source_day_dir.iterdir()): + if not source_output_dir.is_dir(): + continue + target_output_dir = ( + target_facet_dir + / "activities" + / source_day_dir.name + / source_output_dir.name + ) + try: + if target_output_dir.exists(): + _log_decision( + log_path, + { + "action": "facet_activities_output_copied", + "item_type": "activity_output", + "item_id": ( + f"{facet_name}/activities/{source_day_dir.name}/" + f"{source_output_dir.name}" + ), + "reason": "target_exists_skip", + }, + ) + continue + if not dry_run: + shutil.copytree( + source_output_dir, + target_output_dir, + copy_function=shutil.copy2, + ) + _log_decision( + log_path, + { + "action": "facet_activities_output_copied", + "item_type": "activity_output", + "item_id": ( + f"{facet_name}/activities/{source_day_dir.name}/" + f"{source_output_dir.name}" + ), + "reason": "copied", + }, + ) + except Exception as exc: + summary.errors.append( + "facet " + f"{facet_name} activities output " + f"{source_day_dir.name}/{source_output_dir.name}: {exc}" + ) + + source_logs_dir = source_facet_dir / "logs" + if source_logs_dir.is_dir(): + for source_log_file in sorted(source_logs_dir.glob("*.jsonl")): + try: + source_items = _read_jsonl(source_log_file) + if source_items and not dry_run: + target_log_file = target_facet_dir / "logs" / source_log_file.name + _append_jsonl(target_log_file, source_items) + for item in source_items: + _log_decision( + log_path, + { + "action": "facet_logs_appended", + "item_type": "facet_log", + "item_id": f"{facet_name}/logs/{source_log_file.name}", + "reason": "appended", + }, + ) + except Exception as exc: + summary.errors.append( + f"facet {facet_name} logs {source_log_file.name}: {exc}" + ) + + source_news_dir = source_facet_dir / "news" + if source_news_dir.is_dir(): + target_news_dir = target_facet_dir / "news" + for source_news_file in sorted(source_news_dir.glob("*.md")): + try: + target_news_file = target_news_dir / source_news_file.name + if target_news_file.exists(): + _log_decision( + log_path, + { + "action": "facet_news_skipped", + "item_type": "news", + "item_id": f"{facet_name}/news/{source_news_file.name}", + "reason": "target_exists", + }, + ) + continue + if not dry_run: + target_news_dir.mkdir(parents=True, exist_ok=True) + shutil.copy2(source_news_file, target_news_file) + _log_decision( + log_path, + { + "action": "facet_news_copied", + "item_type": "news", + "item_id": f"{facet_name}/news/{source_news_file.name}", + "reason": "new", + }, + ) + except Exception as exc: + summary.errors.append( + f"facet {facet_name} news {source_news_file.name}: {exc}" + ) + + +def _merge_imports( + source: Path, + target: Path, + summary: MergeSummary, + dry_run: bool, + log_path: Path | None = None, +) -> None: + + source_imports_dir = source / "imports" + if not source_imports_dir.is_dir(): + return + + for source_import_dir in sorted(source_imports_dir.iterdir()): + if not source_import_dir.is_dir(): + continue + + target_import_dir = target / "imports" / source_import_dir.name + try: + if target_import_dir.exists(): + summary.imports_skipped += 1 + _log_decision( + log_path, + { + "action": "import_skipped", + "item_type": "import", + "item_id": source_import_dir.name, + "reason": "target_exists", + }, + ) + continue + if not dry_run: + shutil.copytree( + source_import_dir, + target_import_dir, + copy_function=shutil.copy2, + ) + summary.imports_copied += 1 + _log_decision( + log_path, + { + "action": "import_copied", + "item_type": "import", + "item_id": source_import_dir.name, + "reason": "new", + }, + ) + except Exception as exc: + summary.errors.append(f"import {source_import_dir.name}: {exc}") + + +def _read_jsonl(path: Path) -> list[dict[str, Any]]: + if not path.is_file(): + return [] + + items: list[dict[str, Any]] = [] + with open(path, encoding="utf-8") as handle: + for line in handle: + line = line.strip() + if not line: + continue + items.append(json.loads(line)) + return items + + +def _append_jsonl(path: Path, items: list[dict[str, Any]]) -> None: + if not items: + return + + path.parent.mkdir(parents=True, exist_ok=True) + with open(path, "a", encoding="utf-8") as handle: + for item in items: + handle.write(json.dumps(item, ensure_ascii=False) + "\n") + + +__all__ = ["MergeSummary", "merge_journals"] diff --git a/think/tools/call.py b/think/tools/call.py index e1d70b7fb..7b93da70c 100644 --- a/think/tools/call.py +++ b/think/tools/call.py @@ -14,6 +14,7 @@ import json import shutil import subprocess import sys +from datetime import datetime, timezone from pathlib import Path import typer @@ -1147,7 +1148,8 @@ def journal_merge( ), ) -> None: """Merge segments, entities, facets, and imports from a source journal.""" - from think.tools.journal_merge import merge_journals + from think.merge import MergeSummary, merge_journals + from think.utils import get_journal source_path = Path(source).resolve() @@ -1168,7 +1170,19 @@ def journal_merge( ) raise typer.Exit(1) - summary = merge_journals(source_path, dry_run=dry_run) + target_path = Path(get_journal()) + run_id = datetime.now(timezone.utc).strftime("%Y%m%dT%H%M%SZ") + artifact_root = target_path.parent / f"{target_path.name}.merge" / run_id + log_path = artifact_root / "decisions.jsonl" + staging_path = artifact_root / "staging" + + summary: MergeSummary = merge_journals( + source_path, + target_path, + dry_run=dry_run, + log_path=log_path, + staging_path=staging_path, + ) action = "Would merge" if dry_run else "Merged" typer.echo(f"\n{action}:") @@ -1176,7 +1190,7 @@ def journal_merge( f" Segments: {summary.segments_copied} copied, {summary.segments_skipped} skipped, {summary.segments_errored} errored" ) typer.echo( - f" Entities: {summary.entities_created} created, {summary.entities_merged} merged, {summary.entities_skipped} skipped" + f" Entities: {summary.entities_created} created, {summary.entities_merged} merged, {summary.entities_staged} staged, {summary.entities_skipped} skipped" ) typer.echo( f" Facets: {summary.facets_created} created, {summary.facets_merged} merged" @@ -1190,6 +1204,11 @@ def journal_merge( for error in summary.errors: typer.echo(f" - {error}") + if log_path.exists(): + typer.echo(f"\nDecision log: {log_path}") + if summary.entities_staged > 0: + typer.echo(f"Staged entities: {staging_path}") + if not dry_run: subprocess.run( ["sol", "indexer", "--rescan-full"], @@ -1205,6 +1224,7 @@ def journal_merge( "segments_copied": summary.segments_copied, "entities_created": summary.entities_created, "entities_merged": summary.entities_merged, + "entities_staged": summary.entities_staged, "facets_created": summary.facets_created, "facets_merged": summary.facets_merged, "imports_copied": summary.imports_copied, diff --git a/think/tools/journal_merge.py b/think/tools/journal_merge.py index 248db037d..65d292841 100644 --- a/think/tools/journal_merge.py +++ b/think/tools/journal_merge.py @@ -1,400 +1,11 @@ # SPDX-License-Identifier: AGPL-3.0-only # Copyright (c) 2026 sol pbc -"""Journal merge engine - one-shot merge of a source journal into the target.""" +"""Compatibility wrapper - real implementation lives in think.merge.""" -import json -import re -import shutil -from dataclasses import dataclass, field -from pathlib import Path -from typing import Any - -import typer - -from think.entities.core import entity_slug -from think.entities.journal import ( - load_all_journal_entities, - save_journal_entity, +from think.merge import ( # noqa: F401 + MergeSummary, + merge_journals, ) -from think.entities.matching import find_matching_entity -from think.entities.observations import save_observations -from think.entities.relationships import save_facet_relationship -from think.utils import get_journal, iter_segments - -DATE_RE = re.compile(r"^\d{8}$") - - -@dataclass -class MergeSummary: - segments_copied: int = 0 - segments_skipped: int = 0 - segments_errored: int = 0 - entities_created: int = 0 - entities_merged: int = 0 - entities_skipped: int = 0 - facets_created: int = 0 - facets_merged: int = 0 - imports_copied: int = 0 - imports_skipped: int = 0 - errors: list[str] = field(default_factory=list) - - -def merge_journals(source: Path, dry_run: bool = False) -> MergeSummary: - target = Path(get_journal()) - summary = MergeSummary() - target_entities = load_all_journal_entities() - - _merge_segments(source, target, summary, dry_run) - _merge_entities(source, summary, dry_run, target_entities) - _merge_facets(source, target, summary, dry_run) - _merge_imports(source, target, summary, dry_run) - - return summary - - -def _source_day_dirs(source: Path) -> dict[str, Path]: - days: dict[str, Path] = {} - for entry in sorted(source.iterdir()): - if entry.is_dir() and DATE_RE.match(entry.name): - days[entry.name] = entry - return days - - -def _merge_segments( - source: Path, target: Path, summary: MergeSummary, dry_run: bool -) -> None: - for day_name, source_day in sorted(_source_day_dirs(source).items()): - target_day = target / day_name - for stream, seg_key, seg_path in iter_segments(source_day): - if stream == "_default": - target_path = target_day / seg_key - else: - target_path = target_day / stream / seg_key - - try: - if target_path.exists(): - summary.segments_skipped += 1 - continue - - typer.echo(f" Segment {day_name}/{stream}/{seg_key}", err=True) - if dry_run: - summary.segments_copied += 1 - continue - - shutil.copytree(seg_path, target_path, copy_function=shutil.copy2) - summary.segments_copied += 1 - except Exception as exc: - summary.segments_errored += 1 - summary.errors.append(f"segment {day_name}/{stream}/{seg_key}: {exc}") - - -def _merge_entities( - source: Path, - summary: MergeSummary, - dry_run: bool, - target_entities: dict[str, dict[str, Any]], -) -> None: - target_has_principal = any( - bool(entity.get("is_principal")) for entity in target_entities.values() - ) - source_entities_dir = source / "entities" - if not source_entities_dir.is_dir(): - return - - for entity_dir in sorted(source_entities_dir.iterdir()): - entity_path = entity_dir / "entity.json" - if not entity_dir.is_dir() or not entity_path.is_file(): - continue - - try: - source_entity = json.loads(entity_path.read_text(encoding="utf-8")) - source_name = str(source_entity.get("name", "")).strip() - if not source_name: - raise ValueError("missing entity name") - - entity_id = str( - source_entity.get("id") or entity_dir.name or entity_slug(source_name) - ) - if not entity_id: - raise ValueError("missing entity id") - source_entity["id"] = entity_id - - typer.echo(f" Entity {entity_id}", err=True) - match = find_matching_entity(source_name, list(target_entities.values())) - if match is None: - base_id = entity_id - next_id = base_id - suffix = 2 - while next_id in target_entities: - next_id = f"{base_id}_{suffix}" - suffix += 1 - source_entity["id"] = next_id - - if source_entity.get("is_principal") and target_has_principal: - source_entity["is_principal"] = False - elif source_entity.get("is_principal"): - target_has_principal = True - - if not dry_run: - save_journal_entity(source_entity) - summary.entities_created += 1 - target_entities[source_entity["id"]] = source_entity - continue - - target_id = str(match.get("id", "")) - if not target_id: - raise ValueError("matched target entity missing id") - - target_entity = dict(target_entities.get(target_id, match)) - - aka_by_lower: dict[str, str] = {} - for values in (target_entity.get("aka", []), source_entity.get("aka", [])): - if not isinstance(values, list): - continue - for value in values: - if not value: - continue - key = str(value).lower() - if key not in aka_by_lower: - aka_by_lower[key] = str(value) - if aka_by_lower: - target_entity["aka"] = sorted(aka_by_lower.values(), key=str.lower) - - merged_emails: list[str] = [] - seen_emails: set[str] = set() - for values in ( - target_entity.get("emails", []), - source_entity.get("emails", []), - ): - if not isinstance(values, list): - continue - for value in values: - if not value: - continue - email = str(value) - key = email.lower() - if key in seen_emails: - continue - seen_emails.add(key) - merged_emails.append(email) - if merged_emails: - target_entity["emails"] = merged_emails - - if not dry_run: - save_journal_entity(target_entity) - summary.entities_merged += 1 - target_entities[target_id] = target_entity - except Exception as exc: - summary.errors.append(f"entity {entity_dir.name}: {exc}") - - -def _merge_facets( - source: Path, target: Path, summary: MergeSummary, dry_run: bool -) -> None: - source_facets_dir = source / "facets" - if not source_facets_dir.is_dir(): - return - - for source_facet_dir in sorted(source_facets_dir.iterdir()): - facet_json = source_facet_dir / "facet.json" - if not source_facet_dir.is_dir() or not facet_json.is_file(): - continue - - facet_name = source_facet_dir.name - target_facet_dir = target / "facets" / facet_name - - try: - typer.echo(f" Facet {facet_name}", err=True) - if not target_facet_dir.exists(): - if not dry_run: - shutil.copytree( - source_facet_dir, - target_facet_dir, - copy_function=shutil.copy2, - ) - summary.facets_created += 1 - continue - - _merge_overlapping_facet( - facet_name, - source_facet_dir, - target_facet_dir, - summary, - dry_run, - ) - summary.facets_merged += 1 - except Exception as exc: - summary.errors.append(f"facet {facet_name}: {exc}") - - -def _merge_overlapping_facet( - facet_name: str, - source_facet_dir: Path, - target_facet_dir: Path, - summary: MergeSummary, - dry_run: bool, -) -> None: - source_entities_dir = source_facet_dir / "entities" - if source_entities_dir.is_dir(): - for source_entity_dir in sorted(source_entities_dir.iterdir()): - source_entity_json = source_entity_dir / "entity.json" - if not source_entity_dir.is_dir() or not source_entity_json.is_file(): - continue - - entity_id = source_entity_dir.name - target_entity_dir = target_facet_dir / "entities" / entity_id - try: - if target_entity_dir.exists(): - source_relationship = json.loads( - source_entity_json.read_text(encoding="utf-8") - ) - target_relationship_path = target_entity_dir / "entity.json" - target_relationship: dict[str, Any] = {} - if target_relationship_path.is_file(): - target_relationship = json.loads( - target_relationship_path.read_text(encoding="utf-8") - ) - merged_relationship = {**source_relationship, **target_relationship} - - source_observations = _read_jsonl( - source_entity_dir / "observations.jsonl" - ) - target_observations = _read_jsonl( - target_entity_dir / "observations.jsonl" - ) - seen = { - (item.get("content", ""), item.get("observed_at")) - for item in target_observations - } - merged_observations = list(target_observations) - for item in source_observations: - key = (item.get("content", ""), item.get("observed_at")) - if key in seen: - continue - seen.add(key) - merged_observations.append(item) - - if not dry_run: - save_facet_relationship( - facet_name, entity_id, merged_relationship - ) - save_observations(facet_name, entity_id, merged_observations) - continue - - if not dry_run: - shutil.copytree( - source_entity_dir, - target_entity_dir, - copy_function=shutil.copy2, - ) - except Exception as exc: - summary.errors.append(f"facet {facet_name} entity {entity_id}: {exc}") - - source_todos_dir = source_facet_dir / "todos" - if source_todos_dir.is_dir(): - for source_todo_file in sorted(source_todos_dir.glob("*.jsonl")): - try: - target_todo_file = target_facet_dir / "todos" / source_todo_file.name - target_items = _read_jsonl(target_todo_file) - seen = {(item["text"], item.get("created_at")) for item in target_items} - new_items = [ - item - for item in _read_jsonl(source_todo_file) - if (item["text"], item.get("created_at")) not in seen - ] - if new_items and not dry_run: - _append_jsonl(target_todo_file, new_items) - except Exception as exc: - summary.errors.append( - f"facet {facet_name} todo {source_todo_file.name}: {exc}" - ) - - source_calendar_dir = source_facet_dir / "calendar" - if source_calendar_dir.is_dir(): - for source_calendar_file in sorted(source_calendar_dir.glob("*.jsonl")): - try: - target_calendar_file = ( - target_facet_dir / "calendar" / source_calendar_file.name - ) - target_items = _read_jsonl(target_calendar_file) - seen = {(item["title"], item.get("start")) for item in target_items} - new_items = [ - item - for item in _read_jsonl(source_calendar_file) - if (item["title"], item.get("start")) not in seen - ] - if new_items and not dry_run: - _append_jsonl(target_calendar_file, new_items) - except Exception as exc: - summary.errors.append( - f"facet {facet_name} calendar {source_calendar_file.name}: {exc}" - ) - - source_news_dir = source_facet_dir / "news" - if source_news_dir.is_dir(): - target_news_dir = target_facet_dir / "news" - for source_news_file in sorted(source_news_dir.glob("*.md")): - try: - target_news_file = target_news_dir / source_news_file.name - if target_news_file.exists(): - continue - if not dry_run: - target_news_dir.mkdir(parents=True, exist_ok=True) - shutil.copy2(source_news_file, target_news_file) - except Exception as exc: - summary.errors.append( - f"facet {facet_name} news {source_news_file.name}: {exc}" - ) - - -def _merge_imports( - source: Path, target: Path, summary: MergeSummary, dry_run: bool -) -> None: - source_imports_dir = source / "imports" - if not source_imports_dir.is_dir(): - return - - for source_import_dir in sorted(source_imports_dir.iterdir()): - if not source_import_dir.is_dir(): - continue - - target_import_dir = target / "imports" / source_import_dir.name - try: - typer.echo(f" Import {source_import_dir.name}", err=True) - if target_import_dir.exists(): - summary.imports_skipped += 1 - continue - if not dry_run: - shutil.copytree( - source_import_dir, - target_import_dir, - copy_function=shutil.copy2, - ) - summary.imports_copied += 1 - except Exception as exc: - summary.errors.append(f"import {source_import_dir.name}: {exc}") - - -def _read_jsonl(path: Path) -> list[dict[str, Any]]: - if not path.is_file(): - return [] - - items: list[dict[str, Any]] = [] - with open(path, encoding="utf-8") as handle: - for line in handle: - line = line.strip() - if not line: - continue - items.append(json.loads(line)) - return items - - -def _append_jsonl(path: Path, items: list[dict[str, Any]]) -> None: - if not items: - return - path.parent.mkdir(parents=True, exist_ok=True) - with open(path, "a", encoding="utf-8") as handle: - for item in items: - handle.write(json.dumps(item, ensure_ascii=False) + "\n") +__all__ = ["MergeSummary", "merge_journals"]