From f8af8df6de3a00d01fae8235fef6af59098f6ff4 Mon Sep 17 00:00:00 2001 From: Jer Miller Date: Tue, 14 Jul 2026 03:20:20 -0600 Subject: [PATCH] Route mutation lanes through entity resolution boundary Move the audited journal-writing name-resolution lanes onto record_entity_resolution so tier 5-8 matches persist ambiguity before any domain write and then reuse each lane's existing unresolved representation. This prevents low-confidence names from creating duplicates, attaching speaker labels or voiceprints, or firing irreversible identity merges. In particular merge.py now stages ambiguous source entities, speakers.merge_names refuses before merge_entity(commit=True), and discovery exits before entity creation, voiceprint saves, labels, resolved-cluster writes, and candidate retroactive confirmation. The import bridge checks resolved choices before the identical-content shortcut, so a staged low-confidence entity can be resolved and then applied on re-ingest without leaving received, id_map, or staging state incoherent. --- solstone/apps/activities/routes.py | 37 +++++-- solstone/apps/entities/routes.py | 24 +++- .../apps/entities/talent/entity_observer.py | 59 ++++++++-- solstone/apps/entities/tests/test_routes.py | 50 +++++++++ solstone/apps/import/ingest.py | 81 +++++++++++--- solstone/apps/speakers/attribution.py | 57 ++++++++-- solstone/apps/speakers/bootstrap.py | 104 ++++++++++++++++-- solstone/apps/speakers/discovery.py | 33 +++++- .../apps/speakers/tests/test_attribution.py | 31 +++++- .../apps/speakers/tests/test_discovery.py | 59 ++++++++++ .../apps/speakers/tests/test_merge_names.py | 21 ++++ .../apps/speakers/tests/test_seed_imports.py | 95 +++++++++++++++- solstone/talent/participation.py | 57 +++++++++- solstone/talent/schedule.py | 31 +++++- solstone/talent/speaker_attribution.py | 28 ++++- solstone/talent/story.py | 79 +++++++++++-- solstone/think/entities/seeding.py | 30 ++++- solstone/think/merge.py | 78 ++++++++++++- tests/test_activities_cli_create.py | 33 ++++++ tests/test_entity_ingest.py | 52 ++++++++- tests/test_entity_observer_context.py | 93 ++++++++++++++++ tests/test_importer_granola.py | 64 +++++++++++ tests/test_journal_merge.py | 62 +++++++++++ tests/test_participation_corroboration.py | 42 +++++++ tests/test_schedule_hook.py | 82 ++++++++++++++ tests/test_speaker_attribution_hook.py | 71 ++++++++---- tests/test_story_hook.py | 54 +++++++++ 27 files changed, 1388 insertions(+), 119 deletions(-) diff --git a/solstone/apps/activities/routes.py b/solstone/apps/activities/routes.py index 2269370ec..17ef3ad2c 100644 --- a/solstone/apps/activities/routes.py +++ b/solstone/apps/activities/routes.py @@ -42,7 +42,12 @@ from solstone.think.activities import ( update_activity_record, ) from solstone.think.entities.loading import load_entities -from solstone.think.entities.matching import find_matching_entity +from solstone.think.entities.matching import ( + EntityResolutionOutcome, + ResolutionOrigin, + ResolutionScope, + record_entity_resolution, +) from solstone.think.facets import get_facets, log_call_action from solstone.think.journal_io import LockTimeout from solstone.think.utils import now_ms, segment_parse @@ -116,15 +121,33 @@ def _record_payload(record: dict[str, Any]) -> dict[str, Any]: def _resolve_participation_entity_ids( - entries: list[dict[str, Any]], *, facet: str, day: str + entries: list[dict[str, Any]], *, facet: str, day: str, record_id: str ) -> list[dict[str, Any]]: entities_list = load_entities(facet=facet, day=day) + scope = ResolutionScope.facet_scope(facet) + origin = ResolutionOrigin( + lane="apps.activities.create", + facet=facet, + day=day, + record_id=record_id, + field="participation.name", + ) resolved_entries = [] for entry in entries: resolved = dict(entry) - match = find_matching_entity(resolved["name"], entities_list) - resolved["entity_id"] = match.get("id") if match else None + resolution = record_entity_resolution( + str(resolved.get("name") or ""), + entities_list, + scope=scope, + origin=origin, + ) + resolved["entity_id"] = ( + resolution.entity.get("id") + if resolution.outcome == EntityResolutionOutcome.RESOLVED + and resolution.entity + else None + ) resolved_entries.append(resolved) return resolved_entries @@ -187,15 +210,15 @@ def activities_create_record(day: str) -> Any: details = str(body.get("details") or "") participation_provided = "participation" in body participation: list[dict[str, Any]] = [] + actor = "cogitate:activities" if source == "cogitate" else "cli:create" + span_id = make_activity_id(activity_type, anchor) if participation_provided: raw_participation = body.get("participation") participation = raw_participation if isinstance(raw_participation, list) else [] participation = _resolve_participation_entity_ids( - participation, facet=facet, day=day + participation, facet=facet, day=day, record_id=span_id ) - actor = "cogitate:activities" if source == "cogitate" else "cli:create" - span_id = make_activity_id(activity_type, anchor) record: dict[str, Any] = { "id": span_id, "activity": activity_type, diff --git a/solstone/apps/entities/routes.py b/solstone/apps/entities/routes.py index 3d91cab0e..7a4fed3f3 100644 --- a/solstone/apps/entities/routes.py +++ b/solstone/apps/entities/routes.py @@ -57,6 +57,9 @@ from solstone.think.entities import ( EntityDict, EntityExistsError, EntityNotFoundError, + EntityResolutionOutcome, + ResolutionOrigin, + ResolutionScope, add_entity_aka, add_observation, attach_or_reactivate_entity, @@ -77,6 +80,7 @@ from solstone.think.entities import ( load_facet_relationship, load_observations, merge_entity, + record_entity_resolution, resolve_entity, resolve_journal_entity, save_detected_entity, @@ -494,15 +498,29 @@ def detect_entity_route(facet_name: str) -> Any: detail=f"Invalid entity type '{type_}'", ) - resolved, _ = resolve_entity(facet_name, entity) - if resolved is None: + resolution = record_entity_resolution( + entity, + load_entities(facet_name), + scope=ResolutionScope.facet_scope(facet_name), + origin=ResolutionOrigin( + lane="apps.entities.detect", + facet=facet_name, + day=day, + field="entity", + ), + ) + if resolution.outcome != EntityResolutionOutcome.RESOLVED: blocked_match, _ = resolve_entity(facet_name, entity, include_blocked=True) if blocked_match and blocked_match.get("blocked"): return error_response( ENTITY_BLOCKED, detail=str(blocked_match.get("name") or entity), ) - name = str(resolved.get("name", entity)) if resolved else entity + name = ( + str(resolution.entity.get("name", entity)) + if resolution.outcome == EntityResolutionOutcome.RESOLVED and resolution.entity + else entity + ) try: save_detected_entity(facet_name, day, type_, name, description) diff --git a/solstone/apps/entities/talent/entity_observer.py b/solstone/apps/entities/talent/entity_observer.py index 6b1b338dd..e2073a736 100644 --- a/solstone/apps/entities/talent/entity_observer.py +++ b/solstone/apps/entities/talent/entity_observer.py @@ -18,7 +18,12 @@ import logging from solstone.talent.story import ALLOWED_RELATION_KINDS from solstone.think.entities.context import assemble_observer_context from solstone.think.entities.loading import detected_entities_path, load_entities -from solstone.think.entities.matching import find_matching_entity +from solstone.think.entities.matching import ( + EntityResolutionOutcome, + ResolutionOrigin, + ResolutionScope, + record_entity_resolution, +) from solstone.think.entities.observations import record_observation_ops from solstone.think.journal_io import LockTimeout from solstone.think.utils import now_ms @@ -78,6 +83,9 @@ def _clean_relation( value: object, op: str, entities: list[dict], + *, + scope: ResolutionScope, + origin: ResolutionOrigin, ) -> tuple[dict | None, str | None]: if value is None or op in {"drop", "keep"}: return None, None @@ -102,8 +110,14 @@ def _clean_relation( logger.warning("entity_observer: relation kind 'other' requires note") return None, "skipped" - match = find_matching_entity(target_name, entities, fuzzy_threshold=90) - if not match: + resolution = record_entity_resolution( + target_name, + entities, + scope=scope, + origin=origin, + fuzzy_threshold=90, + ) + if resolution.outcome != EntityResolutionOutcome.RESOLVED or not resolution.entity: logger.warning( "entity_observer: unresolved relation target %r for %s op", target_name, @@ -118,14 +132,19 @@ def _clean_relation( return { "kind": kind, - "target_entity_id": match["id"], + "target_entity_id": resolution.entity["id"], "target_name": target_name, "note": note, }, None def _clean_operation( - item: object, seen_indexes: set[int], entities: list[dict] + item: object, + seen_indexes: set[int], + entities: list[dict], + *, + scope: ResolutionScope, + origin: ResolutionOrigin, ) -> tuple[dict | None, str | None]: if not isinstance(item, dict): return None, "skipped" @@ -135,7 +154,13 @@ def _clean_operation( content = item.get("content") if not isinstance(content, str) or not content.strip(): return None, "skipped" - relation, status = _clean_relation(item.get("relation"), op, entities) + relation, status = _clean_relation( + item.get("relation"), + op, + entities, + scope=scope, + origin=origin, + ) if status == "skipped": return None, status cleaned = {"op": "add", "content": content.strip()} @@ -170,7 +195,13 @@ def _clean_operation( if content is not None: cleaned["content"] = content - relation, status = _clean_relation(item.get("relation"), op, entities) + relation, status = _clean_relation( + item.get("relation"), + op, + entities, + scope=scope, + origin=origin, + ) if status == "skipped": return None, status if relation is not None: @@ -212,6 +243,7 @@ def post_process(result: str, context: dict) -> str | None: return None attached_entities = load_entities(facet) + resolution_scope = ResolutionScope.facet_scope(facet) valid_entity_ids = { entity.get("id") for entity in attached_entities if entity.get("id") } @@ -238,9 +270,20 @@ def post_process(result: str, context: dict) -> str | None: clean_ops: list[dict] = [] seen_indexes: set[int] = set() + relation_origin = ResolutionOrigin( + lane="apps.entities.entity_observer", + facet=facet, + day=day, + record_id=entity_id, + field="relation.target_name", + ) for item in operations: clean_op, status = _clean_operation( - item, seen_indexes, attached_entities + item, + seen_indexes, + attached_entities, + scope=resolution_scope, + origin=relation_origin, ) if status is not None: counts[status] += 1 diff --git a/solstone/apps/entities/tests/test_routes.py b/solstone/apps/entities/tests/test_routes.py index 589dde837..3529e75ec 100644 --- a/solstone/apps/entities/tests/test_routes.py +++ b/solstone/apps/entities/tests/test_routes.py @@ -8,11 +8,14 @@ from __future__ import annotations from pathlib import Path from solstone.think.entities import ( + ResolutionScope, attach_or_reactivate_entity, detach_facet_entity, + load_ambiguities, load_entities, load_facet_relationship, load_journal_entity, + record_ambiguity_choice, save_entities, ) from solstone.think.journal_io import LockTimeout @@ -151,6 +154,53 @@ def test_delete_detected_returns_days_modified(client): } +def test_detected_ambiguous_name_keeps_submitted_name_without_id(client): + save_entities( + "personal", + [ + {"type": "Person", "name": "Ambig Route Alpha"}, + {"type": "Person", "name": "Ambig Route Beta"}, + ], + ) + + response = client.post( + "/app/entities/api/personal/detected", + json={ + "day": "20240102", + "type": "Person", + "entity": "Ambig", + "description": "Detected ambiguous person", + }, + ) + + assert response.status_code == 200 + detected = load_entities("personal", "20240102") + assert [entity["name"] for entity in detected] == ["Ambig"] + rows = load_ambiguities() + assert len(rows) == 1 + assert rows[0]["normalized_query"] == "ambig" + + record_ambiguity_choice( + "Ambig", + "ambig_route_alpha", + load_entities("personal"), + scope=ResolutionScope.facet_scope("personal"), + ) + response = client.post( + "/app/entities/api/personal/detected", + json={ + "day": "20240103", + "type": "Person", + "entity": "Ambig", + "description": "Detected resolved person", + }, + ) + + assert response.status_code == 200 + detected = load_entities("personal", "20240103") + assert [entity["name"] for entity in detected] == ["Ambig Route Alpha"] + + def test_owner_lock_timeout_maps_to_entity_busy(client, monkeypatch): def raise_busy(*args, **kwargs): raise LockTimeout(Path("busy"), 0.01) diff --git a/solstone/apps/import/ingest.py b/solstone/apps/import/ingest.py index 65d8b8b2b..d4d30150b 100644 --- a/solstone/apps/import/ingest.py +++ b/solstone/apps/import/ingest.py @@ -29,13 +29,22 @@ from solstone.observe.utils import ( compute_file_sha256, find_available_segment, ) +from solstone.think.entities.ambiguities import ( + load_resolved_ambiguity_choice, + normalize_resolution_query, +) from solstone.think.entities.core import EntityDict, entity_slug from solstone.think.entities.journal import ( has_journal_principal, load_all_journal_entities, save_journal_entity, ) -from solstone.think.entities.matching import find_matching_entity +from solstone.think.entities.matching import ( + EntityResolutionOutcome, + ResolutionOrigin, + ResolutionScope, + record_entity_resolution, +) from solstone.think.journal_io import append_text, atomic_replace from solstone.think.utils import DEFAULT_STREAM, STREAM_RE, day_path @@ -382,6 +391,8 @@ def register_ingest_routes(bp) -> None: staged = 0 skipped = 0 errors: list[dict[str, str]] = [] + state_dirty = False + resolution_scope = ResolutionScope.journal() for entity_data in entities: try: @@ -401,8 +412,25 @@ def register_ingest_routes(bp) -> None: json.dumps(entity_data, sort_keys=True, ensure_ascii=False).encode() ).hexdigest() + staged_path = staged_dir / f"{source_id}.json" existing_hash = entity_state["received"].get(source_id) - if existing_hash == content_hash: + resolved_choice = load_resolved_ambiguity_choice( + resolution_scope, + normalize_resolution_query(entity_data["name"]), + ) + resolved_target_id = ( + str(resolved_choice.get("resolved_entity_id") or "") + if resolved_choice + else "" + ) + needs_resolved_choice_apply = bool( + resolved_target_id + and ( + entity_state["id_map"].get(source_id) != resolved_target_id + or staged_path.exists() + ) + ) + if existing_hash == content_hash and not needs_resolved_choice_apply: skipped += 1 entity_state["received"][source_id] = content_hash _append_decision( @@ -421,12 +449,23 @@ def register_ingest_routes(bp) -> None: ) continue - match = find_matching_entity( - entity_data["name"], list(target_entities.values()) + resolution = record_entity_resolution( + entity_data["name"], + list(target_entities.values()), + scope=resolution_scope, + origin=ResolutionOrigin( + lane="apps.import.ingest", + source_id=source_id, + path=key_prefix, + field="name", + ), ) - if match is not None and match.is_high_confidence: - target_id = str(match["id"]) + if ( + resolution.outcome == EntityResolutionOutcome.RESOLVED + and resolution.entity + ): + target_id = str(resolution.entity["id"]) target_entity: EntityDict = dict(target_entities[target_id]) pre_merge_snapshot = dict(target_entity) @@ -485,7 +524,10 @@ def register_ingest_routes(bp) -> None: if pre_merge_snapshot.get(key) != target_entity.get(key) ) entity_state["id_map"][source_id] = target_id + if staged_path.exists(): + staged_path.unlink() auto_merged += 1 + match_tier = int(resolution.tier) if resolution.tier else None _append_decision( log_path, { @@ -493,29 +535,30 @@ def register_ingest_routes(bp) -> None: "action": "auto_merged", "item_type": "entity", "item_id": source_id, - "match_tier": int(match.tier), - "reason": "high_confidence_match", + "match_tier": match_tier, + "reason": ( + "high_confidence_match" + if match_tier is not None + else "resolved_ambiguity_choice" + ), "source": entity_data, "target": target_entity, "fields_changed": fields_changed, }, ) - elif match is not None and not match.is_high_confidence: + elif resolution.outcome == EntityResolutionOutcome.AMBIGUOUS: staged_dir.mkdir(parents=True, exist_ok=True) staged_payload = { "source_entity": entity_data, "match_candidates": [ - { - "id": match["id"], - "name": match["name"], - "tier": int(match.tier), - } + candidate.to_dict() for candidate in resolution.candidates ], "reason": "low_confidence_match", + "ambiguity_id": resolution.ambiguity_id, "staged_at": datetime.now(timezone.utc).isoformat(), } atomic_replace( - staged_dir / f"{source_id}.json", + staged_path, json.dumps(staged_payload, indent=2, ensure_ascii=False) + "\n", ) staged += 1 @@ -526,7 +569,9 @@ def register_ingest_routes(bp) -> None: "action": "staged", "item_type": "entity", "item_id": source_id, - "match_tier": int(match.tier), + "match_tier": int(resolution.tier) + if resolution.tier + else None, "reason": "low_confidence_match", "source": entity_data, "target": None, @@ -620,13 +665,15 @@ def register_ingest_routes(bp) -> None: ) entity_state["received"][source_id] = content_hash + state_dirty = True except Exception as exc: entity_id = ( entity_data.get("id", "") if isinstance(entity_data, dict) else "" ) errors.append({"entity_id": entity_id, "error": str(exc)}) - _write_state_atomic(state_path, entity_state) + if state_dirty: + _write_state_atomic(state_path, entity_state) written = auto_merged + created if written > 0: diff --git a/solstone/apps/speakers/attribution.py b/solstone/apps/speakers/attribution.py index 2b84a7305..24db7448c 100644 --- a/solstone/apps/speakers/attribution.py +++ b/solstone/apps/speakers/attribution.py @@ -41,7 +41,12 @@ from solstone.apps.speakers.encoder_config import ( VP_OUTLIER_MIN_SIMILARITY, ) from solstone.apps.speakers.owner import load_owner_centroid -from solstone.think.entities import find_matching_entity +from solstone.think.entities import ( + EntityResolutionOutcome, + ResolutionOrigin, + ResolutionScope, + record_entity_resolution, +) from solstone.think.entities.journal import ( get_journal_principal, load_all_journal_entities, @@ -377,33 +382,65 @@ def attribute_segment( # Resolve candidates to entities candidate_entities: dict[str, dict] = {} + resolution_scope = ResolutionScope.journal() + resolution_origin = ResolutionOrigin( + lane="apps.speakers.attribution", + day=day, + segment_id=segment_key, + field="candidate_name", + ) for name in candidate_names: - entity = find_matching_entity(name, entities_list) - if entity: - candidate_entities[entity["id"]] = entity + resolution = record_entity_resolution( + name, + entities_list, + scope=resolution_scope, + origin=resolution_origin, + ) + if resolution.outcome == EntityResolutionOutcome.RESOLVED and resolution.entity: + candidate_entities[resolution.entity["id"]] = resolution.entity # 2a: single-listed-speaker — all non-owner sentences belong to them if len(speakers) == 1: - entity = find_matching_entity(speakers[0], entities_list) - if entity: + resolution = record_entity_resolution( + speakers[0], + entities_list, + scope=resolution_scope, + origin=ResolutionOrigin( + lane="apps.speakers.attribution", + day=day, + segment_id=segment_key, + field="structural_single_speaker", + ), + ) + if resolution.outcome == EntityResolutionOutcome.RESOLVED and resolution.entity: for sid in non_owner_sids: if labels[sid]["speaker"] is None: labels[sid] = { "sentence_id": sid, - "speaker": entity["id"], + "speaker": resolution.entity["id"], "confidence": "high", "method": "structural_single_speaker", } # 2b: single setting-field participant (import segments without speakers.json) elif not speakers and len(setting_names) == 1: - entity = find_matching_entity(setting_names[0], entities_list) - if entity: + resolution = record_entity_resolution( + setting_names[0], + entities_list, + scope=resolution_scope, + origin=ResolutionOrigin( + lane="apps.speakers.attribution", + day=day, + segment_id=segment_key, + field="structural_setting", + ), + ) + if resolution.outcome == EntityResolutionOutcome.RESOLVED and resolution.entity: for sid in non_owner_sids: if labels[sid]["speaker"] is None: labels[sid] = { "sentence_id": sid, - "speaker": entity["id"], + "speaker": resolution.entity["id"], "confidence": "high", "method": "structural_setting", } diff --git a/solstone/apps/speakers/bootstrap.py b/solstone/apps/speakers/bootstrap.py index a352048f2..7db908686 100644 --- a/solstone/apps/speakers/bootstrap.py +++ b/solstone/apps/speakers/bootstrap.py @@ -30,9 +30,14 @@ from typing import Any from solstone.apps.speakers.owner import load_owner_centroid from solstone.think.entities import ( + EntityResolution, + EntityResolutionOutcome, + ResolutionOrigin, + ResolutionScope, entity_slug, find_matching_entity, is_name_variant_match, + record_entity_resolution, ) from solstone.think.entities.journal import ( create_journal_entity, @@ -48,6 +53,14 @@ logger = logging.getLogger(__name__) NAME_MERGE_THRESHOLD = 0.90 +def _ambiguity_payload(resolution: EntityResolution) -> dict[str, Any]: + """Return a compact API payload for an ambiguous resolution.""" + return { + "ambiguity_id": resolution.ambiguity_id, + "candidates": [candidate.to_dict() for candidate in resolution.candidates], + } + + def bootstrap_voiceprints(dry_run: bool = False) -> dict[str, Any]: """Bootstrap voiceprints from 1-listed-speaker segments across the full journal. @@ -103,6 +116,7 @@ def bootstrap_voiceprints(dry_run: bool = False) -> dict[str, Any]: "embeddings_saved": 0, "embeddings_skipped_owner": 0, "embeddings_skipped_duplicate": 0, + "speakers_unmatched": [], "errors": [], } @@ -110,6 +124,8 @@ def bootstrap_voiceprints(dry_run: bool = False) -> dict[str, Any]: entity_embeddings: dict[str, list[tuple[np.ndarray, dict]]] = defaultdict(list) entity_existing: dict[str, set] = {} entity_names: dict[str, str] = {} + unmatched_speakers: set[str] = set() + resolution_scope = ResolutionScope.journal() days = sorted(day_dirs().keys()) @@ -130,10 +146,28 @@ def bootstrap_voiceprints(dry_run: bool = False) -> dict[str, Any]: speaker_name = speakers[0] # Match speaker name to an existing journal entity - entity = find_matching_entity(speaker_name, entities_list) - if entity: - entity_id = entity["id"] - entity_name = entity.get("name", speaker_name) + resolution = record_entity_resolution( + speaker_name, + entities_list, + scope=resolution_scope, + origin=ResolutionOrigin( + lane="apps.speakers.bootstrap_voiceprints", + day=day, + segment_id=seg_key, + field="speaker", + ), + ) + if resolution.outcome == EntityResolutionOutcome.AMBIGUOUS: + if speaker_name not in unmatched_speakers: + unmatched_speakers.add(speaker_name) + stats["speakers_unmatched"].append(speaker_name) + continue + if ( + resolution.outcome == EntityResolutionOutcome.RESOLVED + and resolution.entity + ): + entity_id = resolution.entity["id"] + entity_name = resolution.entity.get("name", speaker_name) else: # Create a new entity for this speaker entity_id = entity_slug(speaker_name) @@ -229,13 +263,46 @@ def merge_names(alias_name: str, canonical_name: str) -> dict[str, Any]: journal_entities = load_all_journal_entities() entities_list = list(journal_entities.values()) - - alias_entity = find_matching_entity(alias_name, entities_list) - if not alias_entity: + resolution_scope = ResolutionScope.journal() + + alias_resolution = record_entity_resolution( + alias_name, + entities_list, + scope=resolution_scope, + origin=ResolutionOrigin( + lane="apps.speakers.merge_names", + field="alias_name", + ), + ) + if alias_resolution.outcome == EntityResolutionOutcome.AMBIGUOUS: + return { + "error": f"Ambiguous entity for alias: {alias_name}", + "ambiguous": {"alias": _ambiguity_payload(alias_resolution)}, + } + if alias_resolution.outcome != EntityResolutionOutcome.RESOLVED: + return {"error": f"No entity found for alias: {alias_name}"} + alias_entity = alias_resolution.entity + if alias_entity is None: return {"error": f"No entity found for alias: {alias_name}"} - canonical_entity = find_matching_entity(canonical_name, entities_list) - if not canonical_entity: + canonical_resolution = record_entity_resolution( + canonical_name, + entities_list, + scope=resolution_scope, + origin=ResolutionOrigin( + lane="apps.speakers.merge_names", + field="canonical_name", + ), + ) + if canonical_resolution.outcome == EntityResolutionOutcome.AMBIGUOUS: + return { + "error": f"Ambiguous entity for canonical: {canonical_name}", + "ambiguous": {"canonical": _ambiguity_payload(canonical_resolution)}, + } + if canonical_resolution.outcome != EntityResolutionOutcome.RESOLVED: + return {"error": f"No entity found for canonical: {canonical_name}"} + canonical_entity = canonical_resolution.entity + if canonical_entity is None: return {"error": f"No entity found for canonical: {canonical_name}"} alias_id = alias_entity["id"] @@ -658,6 +725,7 @@ def seed_from_imports(dry_run: bool = False) -> dict[str, Any]: entity_existing: dict[str, set] = {} unmatched_set: set[str] = set() speaker_entity_cache: dict[str, Any] = {} + resolution_scope = ResolutionScope.journal() days = sorted(day_dirs().keys()) @@ -718,8 +786,22 @@ def seed_from_imports(dry_run: bool = False) -> dict[str, Any]: continue if speaker_name not in speaker_entity_cache: - entity = find_matching_entity(speaker_name, entities_list) - speaker_entity_cache[speaker_name] = entity + resolution = record_entity_resolution( + speaker_name, + entities_list, + scope=resolution_scope, + origin=ResolutionOrigin( + lane="apps.speakers.seed_from_imports", + day=day, + segment_id=seg_key, + field="speaker", + ), + ) + speaker_entity_cache[speaker_name] = ( + resolution.entity + if resolution.outcome == EntityResolutionOutcome.RESOLVED + else None + ) entity = speaker_entity_cache[speaker_name] if entity is None: diff --git a/solstone/apps/speakers/discovery.py b/solstone/apps/speakers/discovery.py index e409504f7..4d20b6d1c 100644 --- a/solstone/apps/speakers/discovery.py +++ b/solstone/apps/speakers/discovery.py @@ -364,7 +364,13 @@ def identify_cluster( return {"error": f"Cluster {cluster_id} not found in scan results."} from solstone.apps.speakers.routes import _load_speaker_corrections - from solstone.think.entities import entity_slug, find_matching_entity + from solstone.think.entities import ( + EntityResolutionOutcome, + ResolutionOrigin, + ResolutionScope, + entity_slug, + record_entity_resolution, + ) from solstone.think.entities.journal import ( create_journal_entity, load_all_journal_entities, @@ -385,10 +391,27 @@ def identify_cluster( entity for entity in journal_entities.values() if not entity.get("blocked") ] - entity = find_matching_entity(name, entities_list) - if entity: - entity_id = entity["id"] - entity_name = entity.get("name", name) + resolution = record_entity_resolution( + name, + entities_list, + scope=ResolutionScope.journal(), + origin=ResolutionOrigin( + lane="apps.speakers.discovery.identify_cluster", + record_id=str(cluster_id), + field="name", + ), + ) + if resolution.outcome == EntityResolutionOutcome.AMBIGUOUS: + return { + "status": "ambiguous", + "ambiguity_id": resolution.ambiguity_id, + "candidates": [ + candidate.to_dict() for candidate in resolution.candidates + ], + } + if resolution.outcome == EntityResolutionOutcome.RESOLVED and resolution.entity: + entity_id = resolution.entity["id"] + entity_name = resolution.entity.get("name", name) else: entity_id = entity_slug(name) existing = load_journal_entity(entity_id) diff --git a/solstone/apps/speakers/tests/test_attribution.py b/solstone/apps/speakers/tests/test_attribution.py index b7d73e5b3..c8bb616e1 100644 --- a/solstone/apps/speakers/tests/test_attribution.py +++ b/solstone/apps/speakers/tests/test_attribution.py @@ -291,6 +291,35 @@ def test_layer2_single_speaker(speakers_env): assert result["unmatched"] == [] +def test_layer2_single_speaker_ambiguous_name_stays_unmatched(speakers_env): + from solstone.apps.speakers.attribution import attribute_segment + from solstone.think.entities import load_ambiguities + + env = speakers_env() + _setup_owner(env) + env.create_entity("Sarah Connor") + env.create_entity("Sarah Lee") + + owner_emb = _normalized([0.95, 0.05]) + other_emb = _normalized([0.1, 0.99]) + embeddings = np.vstack([owner_emb, other_emb]) + + seg_dir = _write_controlled_segment(env, "20240101", "090000_300", embeddings) + agents_dir = seg_dir / "talents" + agents_dir.mkdir(parents=True, exist_ok=True) + (agents_dir / "speakers.json").write_text(json.dumps(["Sarah"])) + + result = attribute_segment("20240101", STREAM, "090000_300") + + assert result["labels"][1]["speaker"] is None + assert result["labels"][1]["confidence"] is None + assert result["labels"][1]["method"] is None + assert result["unmatched"] == [2] + rows = load_ambiguities() + assert len(rows) == 1 + assert rows[0]["normalized_query"] == "sarah" + + # --------------------------------------------------------------------------- # Layer 2: Setting field # --------------------------------------------------------------------------- @@ -316,7 +345,7 @@ def test_layer2_setting_field(speakers_env): header = { "raw": "imported_audio.flac", "model": "medium.en", - "setting": "Jer and Jack at coffee", + "setting": "Jer and Jack Andersohn at coffee", } lines = [json.dumps(header)] lines.append(json.dumps({"start": "09:00:00", "text": "Owner talking"})) diff --git a/solstone/apps/speakers/tests/test_discovery.py b/solstone/apps/speakers/tests/test_discovery.py index 37a8eaed4..3ed0af08a 100644 --- a/solstone/apps/speakers/tests/test_discovery.py +++ b/solstone/apps/speakers/tests/test_discovery.py @@ -271,6 +271,65 @@ def test_identify_matches_existing(speakers_env): assert (env.journal / "entities" / "bob_smith" / "voiceprints.npz").exists() +def test_identify_ambiguous_name_returns_before_writes(speakers_env): + from solstone.think.entities import ( + ResolutionScope, + load_all_journal_entities, + load_ambiguities, + record_ambiguity_choice, + ) + + env = speakers_env() + _setup_owner_centroid(env.journal, [0.0, 1.0]) + env.create_entity("Sarah Connor") + env.create_entity("Sarah Lee") + embeddings = _make_speaker_embeddings([1.0, 0.0], 5) + segments = _create_cluster_segments(env, embeddings) + + scan_result = discover_unknown_speakers() + cluster_id = scan_result["clusters"][0]["cluster_id"] + result = identify_cluster(cluster_id, "Sarah") + + assert result["status"] == "ambiguous" + assert result["ambiguity_id"] + assert {candidate["id"] for candidate in result["candidates"]} == { + "sarah_connor", + "sarah_lee", + } + assert not (env.journal / "entities" / "sarah" / "entity.json").exists() + assert not (env.journal / "entities" / "sarah_connor" / "voiceprints.npz").exists() + for day, segment_key, _sentence_count in segments: + labels_path = ( + env.journal / day / "test" / segment_key / "talents" / "speaker_labels.json" + ) + corrections_path = ( + env.journal + / day + / "test" + / segment_key + / "talents" + / "speaker_corrections.json" + ) + assert not labels_path.exists() + assert not corrections_path.exists() + assert load_ambiguities()[0]["normalized_query"] == "sarah" + + record_ambiguity_choice( + "Sarah", + "sarah_connor", + list(load_all_journal_entities().values()), + scope=ResolutionScope.journal(), + ) + + resolved = identify_cluster(cluster_id, "Sarah") + + assert resolved["entity_id"] == "sarah_connor" + assert resolved["voiceprints_saved"] == 20 + assert (env.journal / "entities" / "sarah_connor" / "voiceprints.npz").exists() + row = load_ambiguities()[0] + assert row["status"] == "resolved" + + def test_identify_idempotent(speakers_env): env = speakers_env() _setup_owner_centroid(env.journal, [0.0, 1.0]) diff --git a/solstone/apps/speakers/tests/test_merge_names.py b/solstone/apps/speakers/tests/test_merge_names.py index 6a5a2d788..de647202b 100644 --- a/solstone/apps/speakers/tests/test_merge_names.py +++ b/solstone/apps/speakers/tests/test_merge_names.py @@ -171,6 +171,27 @@ def test_deep_merge_full(speakers_env): assert c.get("corrected_speaker") != "alice_alias" +def test_merge_names_ambiguous_alias_returns_error_without_merge(speakers_env): + env = speakers_env() + env.create_entity("Sarah Connor") + env.create_entity("Sarah Lee") + env.create_entity("Alice Canonical") + + result = merge_names("Sarah", "Alice Canonical") + + assert "error" in result + assert "Ambiguous entity for alias" in result["error"] + assert result["ambiguous"]["alias"]["ambiguity_id"] + assert { + candidate["id"] for candidate in result["ambiguous"]["alias"]["candidates"] + } == { + "sarah_connor", + "sarah_lee", + } + assert load_journal_entity("sarah_connor") is not None + assert load_journal_entity("sarah_lee") is not None + + def test_alias_entity_deleted(speakers_env): """After merge, scan_journal_entities does not return alias_id.""" env = speakers_env() diff --git a/solstone/apps/speakers/tests/test_seed_imports.py b/solstone/apps/speakers/tests/test_seed_imports.py index d861393ad..43b9605e7 100644 --- a/solstone/apps/speakers/tests/test_seed_imports.py +++ b/solstone/apps/speakers/tests/test_seed_imports.py @@ -9,7 +9,11 @@ import json import numpy as np -from solstone.apps.speakers.bootstrap import link_import, seed_from_imports +from solstone.apps.speakers.bootstrap import ( + bootstrap_voiceprints, + link_import, + seed_from_imports, +) # --- link-import tests --- @@ -163,6 +167,95 @@ def test_seed_from_imports_skips_unmatched_speakers(speakers_env): assert result["embeddings_saved"] == 0 +def test_seed_from_imports_skips_ambiguous_speakers(speakers_env): + from solstone.think.entities import ( + ResolutionScope, + load_all_journal_entities, + load_ambiguities, + record_ambiguity_choice, + ) + + env = speakers_env() + _create_owner_centroid(env) + env.create_entity("Sarah Connor") + env.create_entity("Sarah Lee") + + embs = np.zeros((1, 256), dtype=np.float32) + embs[0, 1] = 1.0 + + env.create_import_segment( + "20240101", + "100000_300", + [("Sarah", "Hello")], + embeddings=embs, + ) + + result = seed_from_imports(dry_run=True) + assert result["embeddings_saved"] == 0 + assert result["speakers_unmatched"] == ["Sarah"] + rows = load_ambiguities() + assert len(rows) == 1 + assert rows[0]["normalized_query"] == "sarah" + + record_ambiguity_choice( + "Sarah", + "sarah_connor", + list(load_all_journal_entities().values()), + scope=ResolutionScope.journal(), + ) + resolved = seed_from_imports(dry_run=True) + assert resolved["embeddings_saved"] == 1 + assert resolved["speakers_unmatched"] == [] + assert "Sarah Connor" in resolved["speakers_found"] + + +def test_bootstrap_voiceprints_skips_ambiguous_single_speaker(speakers_env): + from solstone.think.entities import ( + ResolutionScope, + load_all_journal_entities, + load_ambiguities, + record_ambiguity_choice, + ) + + env = speakers_env() + _create_owner_centroid(env) + env.create_entity("Sarah Connor") + env.create_entity("Sarah Lee") + + embs = np.zeros((2, 256), dtype=np.float32) + embs[0, 1] = 1.0 + embs[1, 2] = 1.0 + + env.create_segment( + "20240101", + "100000_300", + ["mic_audio"], + embeddings=embs, + ) + env.create_speakers_json("20240101", "100000_300", ["Sarah"]) + + result = bootstrap_voiceprints(dry_run=True) + + assert result["embeddings_saved"] == 0 + assert result["entities_created"] == 0 + assert result["speakers_unmatched"] == ["Sarah"] + rows = load_ambiguities() + assert len(rows) == 1 + assert rows[0]["normalized_query"] == "sarah" + + record_ambiguity_choice( + "Sarah", + "sarah_connor", + list(load_all_journal_entities().values()), + scope=ResolutionScope.journal(), + ) + resolved = bootstrap_voiceprints(dry_run=True) + assert resolved["embeddings_saved"] == 2 + assert resolved["entities_created"] == 0 + assert resolved["speakers_unmatched"] == [] + assert "Sarah Connor" in resolved["speakers_found"] + + def test_seed_from_imports_owner_contamination(speakers_env): """seed-from-imports skips embeddings too similar to owner.""" env = speakers_env() diff --git a/solstone/talent/participation.py b/solstone/talent/participation.py index d243d0313..0e5f2b7df 100644 --- a/solstone/talent/participation.py +++ b/solstone/talent/participation.py @@ -9,7 +9,12 @@ import logging from solstone.think.activities import update_record_fields from solstone.think.cluster import _find_segment_dir from solstone.think.entities.loading import load_entities -from solstone.think.entities.matching import find_matching_entity +from solstone.think.entities.matching import ( + EntityResolutionOutcome, + ResolutionOrigin, + ResolutionScope, + record_entity_resolution, +) logger = logging.getLogger(__name__) @@ -106,6 +111,14 @@ def post_process(result: str, context: dict) -> str | None: return None entities_list = load_entities(facet=facet, day=day) + scope = ResolutionScope.facet_scope(facet) + participation_origin = ResolutionOrigin( + lane="talent.participation", + facet=facet, + day=day, + record_id=str(record_id), + field="participation.name", + ) resolved_entries = [] for entry in participation: @@ -114,8 +127,18 @@ def post_process(result: str, context: dict) -> str | None: continue resolved_entry = dict(entry) - match = find_matching_entity(resolved_entry.get("name", ""), entities_list) - resolved_entry["entity_id"] = match.get("id") if match else None + resolution = record_entity_resolution( + str(resolved_entry.get("name") or ""), + entities_list, + scope=scope, + origin=participation_origin, + ) + resolved_entry["entity_id"] = ( + resolution.entity.get("id") + if resolution.outcome == EntityResolutionOutcome.RESOLVED + and resolution.entity + else None + ) resolved_entries.append(resolved_entry) segments = activity.get("segments") or [] @@ -143,12 +166,34 @@ def post_process(result: str, context: dict) -> str | None: def _name_resolves_to(entity_id: str | None, entry_name: str) -> bool: if not named_speakers: return False + speaker_origin = ResolutionOrigin( + lane="talent.participation", + facet=facet, + day=day, + record_id=str(record_id), + field="speaker.name", + ) for name in named_speakers: + resolution = record_entity_resolution( + name, + entities_list, + scope=scope, + origin=speaker_origin, + ) + if resolution.outcome == EntityResolutionOutcome.AMBIGUOUS: + continue if entity_id: - match = find_matching_entity(name, entities_list) - if match and match.get("id") == entity_id: + if ( + resolution.outcome == EntityResolutionOutcome.RESOLVED + and resolution.entity + and resolution.entity.get("id") == entity_id + ): return True - if entry_name and name.casefold() == entry_name.casefold(): + if ( + resolution.outcome == EntityResolutionOutcome.NO_MATCH + and entry_name + and name.casefold() == entry_name.casefold() + ): return True return False diff --git a/solstone/talent/schedule.py b/solstone/talent/schedule.py index 88e9e8a8b..afb8e0daa 100644 --- a/solstone/talent/schedule.py +++ b/solstone/talent/schedule.py @@ -19,7 +19,12 @@ from solstone.think.activities import ( mute_activity_record, ) from solstone.think.entities.loading import load_entities -from solstone.think.entities.matching import find_matching_entity +from solstone.think.entities.matching import ( + EntityResolutionOutcome, + ResolutionOrigin, + ResolutionScope, + record_entity_resolution, +) from solstone.think.facets import get_facets logger = logging.getLogger(__name__) @@ -107,15 +112,32 @@ def post_process(result: str, context: dict) -> None: resolved_participation: list[dict[str, Any]] = [] active_entities: list[str] = [] seen_active_entities: set[str] = set() + new_id = make_anticipation_id(activity, start, target_date) + resolution_origin = ResolutionOrigin( + lane="talent.schedule", + facet=facet, + day=cache_key[1], + record_id=new_id, + field="participation.name", + ) + resolution_scope = ResolutionScope.facet_scope(facet) for entry in participation: if not isinstance(entry, dict): continue resolved_entry = dict(entry) - match = find_matching_entity( - resolved_entry.get("name", ""), entities_list + resolution = record_entity_resolution( + str(resolved_entry.get("name") or ""), + entities_list, + scope=resolution_scope, + origin=resolution_origin, + ) + entity_id = ( + resolution.entity.get("id") + if resolution.outcome == EntityResolutionOutcome.RESOLVED + and resolution.entity + else None ) - entity_id = match.get("id") if match else None resolved_entry["entity_id"] = entity_id resolved_participation.append(resolved_entry) @@ -126,7 +148,6 @@ def post_process(result: str, context: dict) -> None: seen_active_entities.add(entity_id) active_entities.append(entity_id) - new_id = make_anticipation_id(activity, start, target_date) record = { "id": new_id, "activity": activity, diff --git a/solstone/talent/speaker_attribution.py b/solstone/talent/speaker_attribution.py index a1b63f291..973b88e9d 100644 --- a/solstone/talent/speaker_attribution.py +++ b/solstone/talent/speaker_attribution.py @@ -114,7 +114,12 @@ def post_process(result: str, context: dict) -> str | None: accumulate_voiceprints, save_speaker_labels, ) - from solstone.think.entities import find_matching_entity + from solstone.think.entities import ( + EntityResolutionOutcome, + ResolutionOrigin, + ResolutionScope, + record_entity_resolution, + ) from solstone.think.entities.journal import load_all_journal_entities from solstone.think.utils import segment_path @@ -148,6 +153,13 @@ def post_process(result: str, context: dict) -> str | None: entities_list = [ e for e in journal_entities.values() if not e.get("blocked") ] + resolution_scope = ResolutionScope.journal() + resolution_origin = ResolutionOrigin( + lane="talent.speaker_attribution", + day=str(day), + segment_id=str(segment), + field="layer4.speaker", + ) for item in items: if not isinstance(item, dict): @@ -157,10 +169,18 @@ def post_process(result: str, context: dict) -> str | None: if sid is None or not speaker_name: continue - entity = find_matching_entity(speaker_name, entities_list) - if entity: + resolution = record_entity_resolution( + str(speaker_name), + entities_list, + scope=resolution_scope, + origin=resolution_origin, + ) + if ( + resolution.outcome == EntityResolutionOutcome.RESOLVED + and resolution.entity + ): layer4[int(sid)] = { - "speaker": entity["id"], + "speaker": resolution.entity["id"], "confidence": "medium", "method": "contextual", } diff --git a/solstone/talent/story.py b/solstone/talent/story.py index 6f9b3cac8..39cdeb894 100644 --- a/solstone/talent/story.py +++ b/solstone/talent/story.py @@ -12,7 +12,12 @@ from typing import Any from solstone.think.activities import merge_story_fields from solstone.think.entities.loading import load_entities -from solstone.think.entities.matching import find_matching_entity +from solstone.think.entities.matching import ( + EntityResolutionOutcome, + ResolutionOrigin, + ResolutionScope, + record_entity_resolution, +) logger = logging.getLogger(__name__) @@ -73,9 +78,23 @@ def _normalize_confidence(value: Any) -> float | None: return clamped -def _resolve_entity_id(name: str, entities: list[dict[str, Any]]) -> str | None: - match = find_matching_entity(name, entities, fuzzy_threshold=90) - return match.get("id") if match else None +def _resolve_entity_id( + name: str, + entities: list[dict[str, Any]], + *, + scope: ResolutionScope, + origin: ResolutionOrigin, +) -> str | None: + resolution = record_entity_resolution( + name, + entities, + scope=scope, + origin=origin, + fuzzy_threshold=90, + ) + if resolution.outcome == EntityResolutionOutcome.RESOLVED and resolution.entity: + return str(resolution.entity.get("id") or "") or None + return None def _validate_fields( @@ -149,6 +168,16 @@ def post_process(result: str, context: dict) -> str: return "" entities = load_entities(facet=facet, day=day) + resolution_scope = ResolutionScope.facet_scope(facet) + + def origin(field: str) -> ResolutionOrigin: + return ResolutionOrigin( + lane="talent.story", + facet=facet, + day=day, + record_id=record_id, + field=field, + ) resolved_commitments: list[dict[str, Any]] = [] for index, entry in enumerate(commitments): @@ -168,10 +197,16 @@ def post_process(result: str, context: dict) -> str: continue resolved_commitment = dict(normalized) resolved_commitment["owner_entity_id"] = _resolve_entity_id( - normalized["owner"], entities + normalized["owner"], + entities, + scope=resolution_scope, + origin=origin("commitments.owner"), ) resolved_commitment["counterparty_entity_id"] = _resolve_entity_id( - normalized["counterparty"], entities + normalized["counterparty"], + entities, + scope=resolution_scope, + origin=origin("commitments.counterparty"), ) resolved_commitments.append(resolved_commitment) @@ -198,10 +233,16 @@ def post_process(result: str, context: dict) -> str: continue resolved_closure = dict(normalized) resolved_closure["owner_entity_id"] = _resolve_entity_id( - normalized["owner"], entities + normalized["owner"], + entities, + scope=resolution_scope, + origin=origin("closures.owner"), ) resolved_closure["counterparty_entity_id"] = _resolve_entity_id( - normalized["counterparty"], entities + normalized["counterparty"], + entities, + scope=resolution_scope, + origin=origin("closures.counterparty"), ) resolved_closures.append(resolved_closure) @@ -227,10 +268,18 @@ def post_process(result: str, context: dict) -> str: resolved_decision = dict(normalized) resolved_decision["counterparty"] = counterparty resolved_decision["owner_entity_id"] = _resolve_entity_id( - normalized["owner"], entities + normalized["owner"], + entities, + scope=resolution_scope, + origin=origin("decisions.owner"), ) resolved_decision["counterparty_entity_id"] = ( - _resolve_entity_id(counterparty, entities) + _resolve_entity_id( + counterparty, + entities, + scope=resolution_scope, + origin=origin("decisions.counterparty"), + ) if isinstance(counterparty, str) and counterparty.strip() else None ) @@ -271,10 +320,16 @@ def post_process(result: str, context: dict) -> str: resolved_relation = dict(normalized) resolved_relation["quote"] = quote resolved_relation["from_entity_id"] = _resolve_entity_id( - normalized["from"], entities + normalized["from"], + entities, + scope=resolution_scope, + origin=origin("relations.from"), ) resolved_relation["to_entity_id"] = _resolve_entity_id( - normalized["to"], entities + normalized["to"], + entities, + scope=resolution_scope, + origin=origin("relations.to"), ) resolved_relations.append(resolved_relation) diff --git a/solstone/think/entities/seeding.py b/solstone/think/entities/seeding.py index 9650e33bc..bb5fc8704 100644 --- a/solstone/think/entities/seeding.py +++ b/solstone/think/entities/seeding.py @@ -49,8 +49,11 @@ def seed_entities( save_journal_entity, ) from solstone.think.entities.matching import ( + EntityResolutionOutcome, + ResolutionOrigin, + ResolutionScope, find_entity_by_email, - find_matching_entity, + record_entity_resolution, ) from solstone.think.entities.observations import add_observation, load_observations @@ -60,6 +63,8 @@ def seed_entities( resolved: list[EntityDict] = [] facet_ensured = False + ambiguous = 0 + scope = ResolutionScope.journal() for ent in entities: name = ent.get("name", "").strip() @@ -77,7 +82,25 @@ def seed_entities( # Fall back to name match if not matched: - matched = find_matching_entity(name, entity_list) + resolution = record_entity_resolution( + name, + entity_list, + scope=scope, + origin=ResolutionOrigin( + lane="think.entities.seed_entities", + facet=facet, + day=day, + field="name", + ), + ) + if resolution.outcome == EntityResolutionOutcome.AMBIGUOUS: + ambiguous += 1 + continue + if ( + resolution.outcome == EntityResolutionOutcome.RESOLVED + and resolution.entity + ): + matched = resolution.entity if matched: # Merge email into existing entity if new @@ -126,4 +149,7 @@ def seed_entities( break existing_contents.add(obs_content) + if ambiguous: + logger.warning("seed_entities: skipped %d ambiguous entities", ambiguous) + return resolved diff --git a/solstone/think/merge.py b/solstone/think/merge.py index 82fb59a71..f243dd44d 100644 --- a/solstone/think/merge.py +++ b/solstone/think/merge.py @@ -21,7 +21,12 @@ from solstone.think.entities.journal import ( load_all_journal_entities, save_journal_entity, ) -from solstone.think.entities.matching import find_matching_entity +from solstone.think.entities.matching import ( + EntityResolutionOutcome, + ResolutionOrigin, + ResolutionScope, + record_entity_resolution, +) from solstone.think.entities.merge import ( _dedupe_akas, _dedupe_emails, @@ -317,7 +322,76 @@ def _merge_entities( raise ValueError("missing entity id") source_entity["id"] = entity_id - match = find_matching_entity(source_name, list(target_entities.values())) + resolution = record_entity_resolution( + source_name, + list(target_entities.values()), + scope=ResolutionScope.journal(), + origin=ResolutionOrigin( + lane="think.merge.entities", + record_id=entity_id, + path=str(entity_path), + field="name", + ), + ) + if resolution.outcome == EntityResolutionOutcome.AMBIGUOUS: + 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": "ambiguous_name_match", + "source": source_entity, + "target": None, + "ambiguity_id": resolution.ambiguity_id, + "match_candidates": [ + candidate.to_dict() + for candidate in resolution.candidates + ], + "staging_path": str( + staging_path / entity_id / "entity.json" + ), + }, + ) + else: + summary.entities_skipped += 1 + _log_decision( + log_path, + { + "action": "entity_skipped", + "item_type": "entity", + "item_id": entity_id, + "reason": "ambiguous_name_no_staging", + "source": source_entity, + "target": None, + "ambiguity_id": resolution.ambiguity_id, + "match_candidates": [ + candidate.to_dict() + for candidate in resolution.candidates + ], + }, + ) + continue + + match = ( + resolution.entity + if resolution.outcome == EntityResolutionOutcome.RESOLVED + else None + ) if match is None: if entity_id in target_entities: if staging_path is not None: diff --git a/tests/test_activities_cli_create.py b/tests/test_activities_cli_create.py index 2ab664c96..fea01727c 100644 --- a/tests/test_activities_cli_create.py +++ b/tests/test_activities_cli_create.py @@ -144,6 +144,39 @@ def test_create_resolves_participation_entity_ids(tmp_path, monkeypatch): think_utils._journal_path_cache = None +def test_create_ambiguous_participation_keeps_entity_id_none(tmp_path, monkeypatch): + from solstone.think.entities import load_ambiguities + + _configure_cli_env(tmp_path, monkeypatch) + + _write_detected_entities( + tmp_path, + "work", + "20260418", + [ + {"id": "sarah_connor", "type": "Person", "name": "Sarah Connor"}, + {"id": "sarah_lee", "type": "Person", "name": "Sarah Lee"}, + ], + ) + + payload = _base_payload() + payload["participation"] = [_valid_participation_entry(name="Sarah")] + + result = _invoke_create(payload) + + assert result.exit_code == 0 + record = _read_written_record(tmp_path) + assert record["participation"][0]["name"] == "Sarah" + assert record["participation"][0]["entity_id"] is None + rows = load_ambiguities() + assert len(rows) == 1 + assert rows[0]["normalized_query"] == "sarah" + + import solstone.think.utils as think_utils + + think_utils._journal_path_cache = None + + def test_create_omits_participation_when_not_provided(tmp_path, monkeypatch): _configure_cli_env(tmp_path, monkeypatch) diff --git a/tests/test_entity_ingest.py b/tests/test_entity_ingest.py index 52e50a655..185200afa 100644 --- a/tests/test_entity_ingest.py +++ b/tests/test_entity_ingest.py @@ -11,6 +11,10 @@ import pytest from flask import Blueprint, Flask import solstone.convey.state as convey_state +from solstone.think.entities import ( + ResolutionScope, + record_ambiguity_choice, +) from solstone.think.entities.core import entity_slug from solstone.think.entities.journal import ( load_all_journal_entities, @@ -263,14 +267,56 @@ def test_stage_low_confidence(ingest_env): staged = _read_staged(env["key_prefix"], "alce_jonson") assert staged["reason"] == "low_confidence_match" - assert staged["match_candidates"] == [ - {"id": "alice_johnson", "name": "Alice Johnson", "tier": 8} - ] + assert staged["match_candidates"][0]["id"] == "alice_johnson" + assert staged["match_candidates"][0]["name"] == "Alice Johnson" + assert staged["match_candidates"][0]["tier"] == 8 + assert "score" in staged["match_candidates"][0] state = _read_state(env["key_prefix"]) assert "alce_jonson" not in state["id_map"] assert "alce_jonson" in state["received"] +def test_resolved_ambiguity_applies_on_identical_reingest(ingest_env): + env = ingest_env + save_journal_entity( + {"id": "alice_johnson", "name": "Alice Johnson", "type": "Person"} + ) + + source = {"name": "Alce Jonson", "type": "Person", "aka": ["AJ"]} + first = _post_entities(env["client"], env["key"], env["key_prefix"], [source]) + + assert first.status_code == 200 + assert first.get_json()["staged"] == 1 + staged_path = ( + get_state_directory(env["key_prefix"]) + / "entities" + / "staged" + / "alce_jonson.json" + ) + assert staged_path.exists() + + record_ambiguity_choice( + "Alce Jonson", + "alice_johnson", + list(load_all_journal_entities().values()), + scope=ResolutionScope.journal(), + ) + second = _post_entities(env["client"], env["key"], env["key_prefix"], [source]) + + assert second.status_code == 200 + assert second.get_json()["auto_merged"] == 1 + assert second.get_json()["skipped"] == 0 + assert not staged_path.exists() + state = _read_state(env["key_prefix"]) + assert state["id_map"]["alce_jonson"] == "alice_johnson" + assert state["received"]["alce_jonson"] == _entity_hash( + {**source, "id": "alce_jonson"} + ) + merged = load_journal_entity("alice_johnson") + assert merged is not None + assert merged["aka"] == ["AJ"] + + def test_stage_id_collision(ingest_env): env = ingest_env save_journal_entity({"id": "test", "name": "Test Entity", "type": "Tool"}) diff --git a/tests/test_entity_observer_context.py b/tests/test_entity_observer_context.py index a4ea087c1..bdd5a6444 100644 --- a/tests/test_entity_observer_context.py +++ b/tests/test_entity_observer_context.py @@ -989,6 +989,99 @@ def test_post_process_preserves_op_with_unresolvable_relation_target( } +def test_post_process_preserves_op_with_ambiguous_relation_target( + tmp_path, monkeypatch +): + from solstone.think.entities import ( + ResolutionScope, + load_ambiguities, + load_entities, + record_ambiguity_choice, + ) + + _set_journal(monkeypatch, str(tmp_path)) + facet = "work" + day = "20260304" + _attach_entity(tmp_path, facet, "alice_johnson", "Alice Johnson") + _attach_entity(tmp_path, facet, "sarah_connor", "Sarah Connor") + _attach_entity(tmp_path, facet, "sarah_lee", "Sarah Lee") + + post_process( + json.dumps( + { + "entities": [ + { + "entity_id": "alice_johnson", + "operations": [ + { + "op": "add", + "content": "Works with Sarah on a confidential project", + "reasoning": "Relational, but Sarah is ambiguous.", + "relation": { + "kind": "works-with", + "target_name": "Sarah", + "note": "", + }, + } + ], + } + ], + "summary": "ambiguous relation target", + } + ), + {"facet": facet, "day": day}, + ) + + observations = load_observations(facet, "alice_johnson") + assert len(observations) == 1 + assert observations[0]["relation"] == { + "kind": "works-with", + "target_entity_id": None, + "target_name": "Sarah", + "note": "", + } + outcome = _load_outcome(tmp_path, facet, day) + assert outcome["relation_unresolved"] == 1 + rows = load_ambiguities() + assert len(rows) == 1 + assert rows[0]["normalized_query"] == "sarah" + + record_ambiguity_choice( + "Sarah", + "sarah_connor", + load_entities(facet), + scope=ResolutionScope.facet_scope(facet), + ) + post_process( + json.dumps( + { + "entities": [ + { + "entity_id": "alice_johnson", + "operations": [ + { + "op": "add", + "content": "Coordinates with Sarah after review", + "reasoning": "Choice resolved.", + "relation": { + "kind": "works-with", + "target_name": "Sarah", + "note": "", + }, + } + ], + } + ], + "summary": "resolved relation target", + } + ), + {"facet": facet, "day": day}, + ) + + observations = load_observations(facet, "alice_johnson") + assert observations[-1]["relation"]["target_entity_id"] == "sarah_connor" + + def test_post_process_drops_op_with_other_relation_kind_and_no_note( tmp_path, monkeypatch ): diff --git a/tests/test_importer_granola.py b/tests/test_importer_granola.py index 160514d5c..5e9fc5fc5 100644 --- a/tests/test_importer_granola.py +++ b/tests/test_importer_granola.py @@ -732,6 +732,70 @@ def test_seed_entities_without_observations(tmp_path, monkeypatch): assert result[0]["name"] == "Test Person" +def test_seed_entities_skips_ambiguous_names_without_observations( + tmp_path, monkeypatch +): + """seed_entities() does not mint duplicates or observations for ambiguous names.""" + from solstone.think.entities import ( + ResolutionScope, + load_all_journal_entities, + load_ambiguities, + record_ambiguity_choice, + ) + from solstone.think.entities.journal import save_journal_entity + from solstone.think.entities.observations import load_observations + from solstone.think.entities.seeding import seed_entities + + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + save_journal_entity( + {"id": "sarah_connor", "name": "Sarah Connor", "type": "Person"} + ) + save_journal_entity({"id": "sarah_lee", "name": "Sarah Lee", "type": "Person"}) + + result = seed_entities( + "test.facet", + "20251028", + [ + { + "name": "Sarah", + "type": "Person", + "observations": ["Should not be written."], + } + ], + ) + + assert result == [] + assert not (tmp_path / "entities" / "sarah" / "entity.json").exists() + assert load_observations("test.facet", "Sarah") == [] + rows = load_ambiguities() + assert len(rows) == 1 + assert rows[0]["normalized_query"] == "sarah" + + record_ambiguity_choice( + "Sarah", + "sarah_connor", + list(load_all_journal_entities().values()), + scope=ResolutionScope.journal(), + ) + resolved = seed_entities( + "test.facet", + "20251028", + [ + { + "name": "Sarah", + "type": "Person", + "observations": ["Written after choice."], + } + ], + ) + + assert [entity["id"] for entity in resolved] == ["sarah_connor"] + assert [ + observation["content"] + for observation in load_observations("test.facet", "sarah_connor") + ] == ["Written after choice."] + + def test_seed_entities_observation_formatting(tmp_path, monkeypatch): """seed_entities() creates observations with correct formatting for all field combos.""" from solstone.think.entities.observations import load_observations diff --git a/tests/test_journal_merge.py b/tests/test_journal_merge.py index 1980d8fe4..b6f551d30 100644 --- a/tests/test_journal_merge.py +++ b/tests/test_journal_merge.py @@ -256,6 +256,68 @@ def test_entity_id_collision(merge_journals_fixture, monkeypatch): assert "staged" in result.output +def test_ambiguous_entity_name_stages_instead_of_merging( + merge_journals_fixture, monkeypatch +): + from solstone.think.entities import load_ambiguities + + paths = merge_journals_fixture + _mock_indexer(monkeypatch) + _write_json( + paths["target"] / "entities" / "sarah_connor" / "entity.json", + { + "id": "sarah_connor", + "name": "Sarah Connor", + "type": "person", + "created_at": 1000, + }, + ) + _write_json( + paths["target"] / "entities" / "sarah_lee" / "entity.json", + { + "id": "sarah_lee", + "name": "Sarah Lee", + "type": "person", + "created_at": 1000, + }, + ) + _write_json( + paths["source"] / "entities" / "source_sarah" / "entity.json", + { + "id": "source_sarah", + "name": "Sarah", + "type": "person", + "aka": ["S"], + "created_at": 3000, + }, + ) + staging_path = paths["target"].parent / "manual-staging" + + summary = merge_journals( + paths["source"], + paths["target"], + dry_run=False, + staging_path=staging_path, + ) + + assert summary.entities_staged == 1 + assert (staging_path / "source_sarah" / "entity.json").exists() + assert not (paths["target"] / "entities" / "source_sarah").exists() + assert ( + _read_json(paths["target"] / "entities" / "sarah_connor" / "entity.json")[ + "name" + ] + == "Sarah Connor" + ) + assert ( + _read_json(paths["target"] / "entities" / "sarah_lee" / "entity.json")["name"] + == "Sarah Lee" + ) + rows = load_ambiguities() + assert len(rows) == 1 + assert rows[0]["normalized_query"] == "sarah" + + def test_facet_copy_new(merge_journals_fixture, monkeypatch): paths = merge_journals_fixture _mock_indexer(monkeypatch) diff --git a/tests/test_participation_corroboration.py b/tests/test_participation_corroboration.py index 9f0376fc9..ff79f44b8 100644 --- a/tests/test_participation_corroboration.py +++ b/tests/test_participation_corroboration.py @@ -212,6 +212,48 @@ def test_ac5_preserves_unresolved_casefold_name_match(tmp_path, monkeypatch): assert record["participation"][0]["role"] == "attendee" +def test_ambiguous_speakers_json_name_does_not_corroborate_and_dedupes_origin( + tmp_path, + monkeypatch, +): + from solstone.think.entities import load_ambiguities + + facet = "work" + day = "20260418" + segment = "090000_300" + _write_detected_entities( + tmp_path, + facet, + day, + [ + {"id": "sarah_connor", "type": "Person", "name": "Sarah Connor"}, + {"id": "sarah_lee", "type": "Person", "name": "Sarah Lee"}, + ], + ) + _write_segment_talent_files( + tmp_path, + day, + segment, + sense=_sense_payload(), + speakers=["Sarah", "Sarah"], + ) + + record = _run_participation( + tmp_path, + monkeypatch, + segments=[segment], + result=_participation_result(name="Sarah", source="voice"), + seed_entities=False, + ) + + assert record["participation"][0]["entity_id"] is None + assert record["participation"][0]["role"] == "mentioned" + rows = load_ambiguities() + assert len(rows) == 1 + assert rows[0]["normalized_query"] == "sarah" + assert rows[0]["occurrence_count"] == 2 + + def test_ac6_ignores_transcript_attendee(tmp_path, monkeypatch): segment = "090000_300" _write_segment_talent_files(tmp_path, "20260418", segment, sense=_sense_payload()) diff --git a/tests/test_schedule_hook.py b/tests/test_schedule_hook.py index c6c59ff4d..f2108e102 100644 --- a/tests/test_schedule_hook.py +++ b/tests/test_schedule_hook.py @@ -106,6 +106,88 @@ def test_schedule_post_process_writes_record_and_resolves_entities( assert record["edits"][-1]["note"] == "created by schedule" +def test_schedule_ambiguous_participant_keeps_unresolved_form( + tmp_path, + monkeypatch, +): + from solstone.talent.schedule import post_process + from solstone.think.activities import load_activity_records + from solstone.think.entities import ( + ResolutionScope, + load_ambiguities, + load_entities, + record_ambiguity_choice, + ) + + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + _write_facet(tmp_path, "work") + _write_detected_entities( + tmp_path, + "work", + "20260420", + [ + {"id": "sarah_connor", "type": "Person", "name": "Sarah Connor"}, + {"id": "sarah_lee", "type": "Person", "name": "Sarah Lee"}, + ], + ) + + payload = [ + { + "activity": "meeting", + "target_date": "2026-04-20", + "start": "16:30:00", + "end": "17:30:00", + "title": "Sarah intro call", + "description": "Intro call with Sarah.", + "details": "Google Meet", + "participation": [ + { + "name": "Sarah", + "role": "attendee", + "source": "screen", + "confidence": 0.95, + "context": "calendar invite", + } + ], + "participation_confidence": 0.88, + "facet": "work", + "cancelled": False, + } + ] + + assert post_process(json.dumps(payload), {"day": "20260418"}) is None + + records = load_activity_records("work", "20260420", include_hidden=True) + assert len(records) == 1 + record = records[0] + assert record["active_entities"] == [] + assert record["participation"][0]["entity_id"] is None + rows = load_ambiguities() + assert len(rows) == 1 + assert rows[0]["normalized_query"] == "sarah" + + record_ambiguity_choice( + "Sarah", + "sarah_connor", + load_entities("work", "20260420"), + scope=ResolutionScope.facet_scope("work"), + ) + resolved_payload = json.loads(json.dumps(payload)) + resolved_payload[0]["start"] = "17:00:00" + resolved_payload[0]["end"] = "18:00:00" + resolved_payload[0]["title"] = "Sarah follow-up call" + assert post_process(json.dumps(resolved_payload), {"day": "20260418"}) is None + + records = load_activity_records("work", "20260420", include_hidden=True) + assert len(records) == 2 + record = records[1] + assert record["active_entities"] == ["sarah_connor"] + assert record["participation"][0]["entity_id"] == "sarah_connor" + rows = load_ambiguities() + assert len(rows) == 1 + assert rows[0]["status"] == "resolved" + + def test_schedule_post_process_accepts_wrapped_events(tmp_path, monkeypatch): from solstone.talent.schedule import post_process from solstone.think.activities import load_activity_records diff --git a/tests/test_speaker_attribution_hook.py b/tests/test_speaker_attribution_hook.py index 7c861f921..2849dc381 100644 --- a/tests/test_speaker_attribution_hook.py +++ b/tests/test_speaker_attribution_hook.py @@ -165,12 +165,6 @@ def _post_process_context(): } -def _match_entity(name, _entities): - if name == "Alice": - return {"id": "alice"} - return None - - class TestPostProcess: def test_bare_list_merges_layer4_attributions(self, tmp_path): result = json.dumps( @@ -185,13 +179,9 @@ class TestPostProcess: patch( "solstone.apps.speakers.attribution.accumulate_voiceprints" ) as accumulate_mock, - patch( - "solstone.think.entities.find_matching_entity", - side_effect=_match_entity, - ), patch( "solstone.think.entities.journal.load_all_journal_entities", - return_value={"alice": {"id": "alice"}}, + return_value={"alice": {"id": "alice", "name": "Alice"}}, ), patch("solstone.think.utils.segment_path", return_value=tmp_path), ): @@ -235,13 +225,9 @@ class TestPostProcess: patch( "solstone.apps.speakers.attribution.accumulate_voiceprints" ) as accumulate_mock, - patch( - "solstone.think.entities.find_matching_entity", - side_effect=_match_entity, - ), patch( "solstone.think.entities.journal.load_all_journal_entities", - return_value={"alice": {"id": "alice"}}, + return_value={"alice": {"id": "alice", "name": "Alice"}}, ), patch("solstone.think.utils.segment_path", return_value=tmp_path), ): @@ -264,6 +250,52 @@ class TestPostProcess: } accumulate_mock.assert_not_called() + def test_ambiguous_layer4_attribution_leaves_speaker_unmatched( + self, tmp_path, monkeypatch + ): + from solstone.think.entities import load_ambiguities + + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + import solstone.think.utils as think_utils + + think_utils._journal_path_cache = None + result = json.dumps( + [{"sentence_id": 1, "speaker": "Sarah", "reasoning": "said her name"}] + ) + context = _post_process_context() + + with ( + patch( + "solstone.apps.speakers.attribution.save_speaker_labels" + ) as save_mock, + patch( + "solstone.apps.speakers.attribution.accumulate_voiceprints" + ) as accumulate_mock, + patch( + "solstone.think.entities.journal.load_all_journal_entities", + return_value={ + "sarah_connor": {"id": "sarah_connor", "name": "Sarah Connor"}, + "sarah_lee": {"id": "sarah_lee", "name": "Sarah Lee"}, + }, + ), + patch("solstone.think.utils.segment_path", return_value=tmp_path), + ): + from solstone.talent.speaker_attribution import post_process + + post_process(result, context) + + saved_labels = save_mock.call_args[0][1] + assert saved_labels[0] == { + "sentence_id": 1, + "speaker": None, + "confidence": None, + "method": None, + } + accumulate_mock.assert_not_called() + rows = load_ambiguities() + assert len(rows) == 1 + assert rows[0]["normalized_query"] == "sarah" + def test_non_list_non_dict_yields_zero_merges_and_warns(self, tmp_path, caplog): context = _post_process_context() @@ -272,13 +304,9 @@ class TestPostProcess: "solstone.apps.speakers.attribution.save_speaker_labels" ) as save_mock, patch("solstone.apps.speakers.attribution.accumulate_voiceprints"), - patch( - "solstone.think.entities.find_matching_entity", - side_effect=_match_entity, - ) as match_mock, patch( "solstone.think.entities.journal.load_all_journal_entities", - return_value={"alice": {"id": "alice"}}, + return_value={"alice": {"id": "alice", "name": "Alice"}}, ) as load_mock, patch("solstone.think.utils.segment_path", return_value=tmp_path), caplog.at_level(logging.WARNING), @@ -302,4 +330,3 @@ class TestPostProcess: } assert "expected JSON array, got int" in caplog.text load_mock.assert_not_called() - match_mock.assert_not_called() diff --git a/tests/test_story_hook.py b/tests/test_story_hook.py index 505a359e2..4cb654aed 100644 --- a/tests/test_story_hook.py +++ b/tests/test_story_hook.py @@ -342,6 +342,60 @@ def test_story_hook_resolves_entities(tmp_path, monkeypatch): assert record["closures"][0]["counterparty"] == "Mina" +def test_story_hook_ambiguous_entities_keep_names_and_null_ids(tmp_path, monkeypatch): + from solstone.talent.story import post_process + from solstone.think.activities import append_activity_record + from solstone.think.entities import load_ambiguities + + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + _write_detected_entities( + tmp_path, + "work", + "20260418", + [ + {"id": "sarah_connor", "type": "Person", "name": "Sarah Connor"}, + {"id": "sarah_lee", "type": "Person", "name": "Sarah Lee"}, + ], + ) + append_activity_record("work", "20260418", _activity_record()) + + post_process( + _valid_result( + commitments=[ + { + "owner": "Sarah", + "action": "send the revised deck", + "counterparty": "Sarah", + "when": "Friday morning", + "context": "Ambiguous Sarah should not attach.", + } + ], + closures=[], + decisions=[], + relations=[ + { + "from": "Sarah", + "to": "Sarah", + "kind": "works-with", + "note": "Ambiguous relationship.", + } + ], + ), + _context(tmp_path), + ) + + record = _load_record("work", "20260418") + assert record["commitments"][0]["owner"] == "Sarah" + assert record["commitments"][0]["owner_entity_id"] is None + assert record["commitments"][0]["counterparty_entity_id"] is None + assert record["relations"][0]["from"] == "Sarah" + assert record["relations"][0]["from_entity_id"] is None + assert record["relations"][0]["to_entity_id"] is None + rows = load_ambiguities() + assert len(rows) == 1 + assert rows[0]["normalized_query"] == "sarah" + + def test_story_hook_skips_invalid_relations(tmp_path, monkeypatch, caplog): from solstone.talent.story import post_process from solstone.think.activities import append_activity_record -- 2.51.2