From b317914689cdf28283abe366f9c5a53bd150785f Mon Sep 17 00:00:00 2001 From: Jer Miller Date: Mon, 20 Jul 2026 21:31:31 -0600 Subject: [PATCH] feat(speakers): expose identify operations and suppression Add HTTP and sol call surfaces for identify replay/undo, operation list/show, cluster dismissals, and keep-separate summaries. Recoverable, repair_required, conflict, and not_found results now map to non-2xx standard error envelopes with operation state in extra, while genuine identify and undo success keep 2xx responses; the old status:"partial" response shape is not retained. Wire the shared keep-separate predicate through name-variant record-time storage, curation read-time loading, and suggest re-detection so suppressed pairs stop nagging until detection_count exceeds the watermark. Add the dismissal read-side filter, generated OpenAPI and command references, and route/CLI/contract coverage. Co-Authored-By: Claude Opus 4.8 (1M context) --- docs/openapi/convey-clients.json | 4 + .../observer-client-contract/manifest.json | 8 +- .../projection.openapi.json | 4 + solstone/apps/speakers/call.py | 145 ++++++++- solstone/apps/speakers/maintenance.py | 8 +- solstone/apps/speakers/routes.py | 292 +++++++++++++++++- solstone/apps/speakers/suggest.py | 9 + solstone/apps/speakers/tests/test_routes.py | 232 ++++++++++++-- solstone/apps/speakers/tests/test_suggest.py | 23 ++ solstone/convey/contract/observer_bundle.py | 2 +- solstone/convey/reasons.py | 20 ++ solstone/talent/sol/references/commands.md | 4 +- solstone/think/curation.py | 8 + solstone/think/speaker_review_candidates.py | 58 +++- tests/test_curation.py | 22 ++ tests/test_observer_client_bundle.py | 30 +- tests/test_speaker_review_candidates.py | 67 +++- tests/test_speakers_call_parity.py | 159 ++++++++++ 18 files changed, 1016 insertions(+), 79 deletions(-) diff --git a/docs/openapi/convey-clients.json b/docs/openapi/convey-clients.json index 8f88be01c..f89984e61 100644 --- a/docs/openapi/convey-clients.json +++ b/docs/openapi/convey-clients.json @@ -111,6 +111,10 @@ "settings_operation_failed", "speaker_attribution_state_invalid", "speaker_command_failed", + "speaker_identify_conflict", + "speaker_identify_operation_not_found", + "speaker_identify_recoverable", + "speaker_identify_repair_required", "speaker_labels_busy", "speaker_not_found", "speaker_owner_centroid_required", diff --git a/docs/openapi/observer-client-contract/manifest.json b/docs/openapi/observer-client-contract/manifest.json index 10bd33736..aadcd4098 100644 --- a/docs/openapi/observer-client-contract/manifest.json +++ b/docs/openapi/observer-client-contract/manifest.json @@ -14,7 +14,7 @@ } ], "bundle_schema_identity": "solstone.observer-client-contract-bundle.schema.v1", - "bundle_semver": "1.0.2", + "bundle_semver": "2.0.0", "component_closure": [ "CallosumEvent", "Error", @@ -42,7 +42,7 @@ }, { "path": "projection.openapi.json", - "sha256": "8a2b7037552edf710597f2ffa6fdc5aa715311df4ea8cf168e70abe4231c64ca" + "sha256": "28c055279ab7d80c809a43c5f710ccebc4da643fec3f8e6beae48823c17c46c5" }, { "path": "vectors.json", @@ -55,7 +55,7 @@ "id": "bundle.projection_builder", "path": "solstone/convey/contract/observer_bundle.py", "role": "projection_builder", - "sha256": "00918bb64972340f6f8ec55fe0ccba5630056a2a6f5c975a6717a95ca6aa474b" + "sha256": "fc055febc441b931ef0605e61443e3f84a0be25f17996372ae948c4a0eb3c8b0" }, { "id": "bundle.recording", @@ -157,7 +157,7 @@ "id": "reason_codes", "path": "solstone/convey/reasons.py", "role": "vocabulary_source", - "sha256": "ea857c2c4eb2308c05ad02e33a0ead9f1c68c5e7e08c4c0f7aebf52f8a7e8750" + "sha256": "69564b0ce7a9a202ef14d7e7c2aecc63bb48a901176313fcc6fd519a8f604946" } ], "observer_protocol_version": 2, diff --git a/docs/openapi/observer-client-contract/projection.openapi.json b/docs/openapi/observer-client-contract/projection.openapi.json index b7dbb44c6..2877fb037 100644 --- a/docs/openapi/observer-client-contract/projection.openapi.json +++ b/docs/openapi/observer-client-contract/projection.openapi.json @@ -111,6 +111,10 @@ "settings_operation_failed", "speaker_attribution_state_invalid", "speaker_command_failed", + "speaker_identify_conflict", + "speaker_identify_operation_not_found", + "speaker_identify_recoverable", + "speaker_identify_repair_required", "speaker_labels_busy", "speaker_not_found", "speaker_owner_centroid_required", diff --git a/solstone/apps/speakers/call.py b/solstone/apps/speakers/call.py index 1729c160d..852454b0c 100644 --- a/solstone/apps/speakers/call.py +++ b/solstone/apps/speakers/call.py @@ -25,7 +25,13 @@ Commands: sol call speakers wipe [--commit] [--json] sol call speakers discover [--json] sol call speakers presence [--json] - sol call speakers identify [--entity-id ID] + sol call speakers identify [--entity-id ID] [--request-id ID] + sol call speakers identify-undo + sol call speakers identify-operations + sol call speakers identify-operation + sol call speakers dismiss-cluster --disposition + sol call speakers dismissals + sol call speakers keep-separate-list sol call speakers merge-names sol call speakers link-import --entity-id sol call speakers seed-from-imports [--commit] [--json] @@ -52,6 +58,10 @@ from solstone.convey.reasons import ( INVALID_DAY, INVALID_SEGMENT_OR_STREAM, SPEAKER_COMMAND_FAILED, + SPEAKER_IDENTIFY_CONFLICT, + SPEAKER_IDENTIFY_OPERATION_NOT_FOUND, + SPEAKER_IDENTIFY_RECOVERABLE, + SPEAKER_IDENTIFY_REPAIR_REQUIRED, SPEAKER_OWNER_CENTROID_REQUIRED, SPEAKER_OWNER_IDENTITY_REQUIRED, SPEAKER_REVIEW_UNAVAILABLE, @@ -86,6 +96,40 @@ def _exit_speaker_command_failed(err: ConveyClientError) -> None: raise typer.Exit(1) from err +_IDENTIFY_OPERATION_FAILURE_CODES = { + SPEAKER_IDENTIFY_RECOVERABLE.code, + SPEAKER_IDENTIFY_REPAIR_REQUIRED.code, + SPEAKER_IDENTIFY_CONFLICT.code, + SPEAKER_IDENTIFY_OPERATION_NOT_FOUND.code, +} + + +def _exit_identify_operation_failed( + err: ConveyClientError, + *, + request_id: str | None = None, +) -> None: + payload = err.payload if isinstance(err.payload, dict) else {} + operation_id = payload.get("operation_id") + retry_request_id = request_id or payload.get("request_id") + typer.echo(err.detail or err.error, err=True) + if retry_request_id: + typer.echo( + f"Retry with the same --request-id {retry_request_id}.", + err=True, + ) + typer.echo( + "Inspect operations with: sol call speakers identify-operations", + err=True, + ) + if operation_id: + typer.echo( + f"Inspect this operation with: sol call speakers identify-operation {operation_id}", + err=True, + ) + raise typer.Exit(1) from err + + @app.command("status") @convey_cli def status( @@ -728,6 +772,14 @@ def identify( resolve_only: bool = typer.Option( False, "--resolve-only", help="Resolve without writing (dry run)." ), + request_id: str | None = typer.Option( + None, "--request-id", help="Stable request id for retry/resume." + ), + reviewed_near_match_entity_id: list[str] | None = typer.Option( + None, + "--reviewed-near-match-entity-id", + help="Reviewed near-match entity id to keep separate when creating.", + ), ) -> None: """Identify a discovered unknown speaker cluster.""" if not name and not entity_id: @@ -743,8 +795,83 @@ def identify( "create_new": create, "entity_type": entity_type, "resolve_only": resolve_only, + "request_id": request_id, + "reviewed_near_match_entity_ids": reviewed_near_match_entity_id, }, ) + except ConveyClientError as err: + if err.reason_code in _IDENTIFY_OPERATION_FAILURE_CODES: + _exit_identify_operation_failed(err, request_id=request_id) + if err.reason_code == SPEAKER_COMMAND_FAILED.code: + _exit_speaker_command_failed(err) + raise + typer.echo(json.dumps(result, indent=2, default=str)) + + +@app.command("identify-undo") +@convey_cli +def identify_undo( + operation_id: str = typer.Argument(..., help="Identify operation id to undo."), +) -> None: + """Undo a committed speaker identify operation.""" + try: + result = _request( + "POST", + "/app/speakers/api/discovery/identify/undo", + json_body={"operation_id": operation_id}, + ) + except ConveyClientError as err: + if err.reason_code in _IDENTIFY_OPERATION_FAILURE_CODES: + _exit_identify_operation_failed(err) + if err.reason_code == SPEAKER_COMMAND_FAILED.code: + _exit_speaker_command_failed(err) + raise + typer.echo(json.dumps(result, indent=2, default=str)) + + +@app.command("identify-operations") +@convey_cli +def identify_operations() -> None: + """List redacted speaker identify operations.""" + result = _request("GET", "/app/speakers/api/discovery/identify/operations") + typer.echo(json.dumps(result, indent=2, default=str)) + + +@app.command("identify-operation") +@convey_cli +def identify_operation( + operation_id: str = typer.Argument(..., help="Identify operation id to inspect."), +) -> None: + """Show one redacted speaker identify operation.""" + try: + result = _request( + "GET", + f"/app/speakers/api/discovery/identify/operations/{operation_id}", + ) + except ConveyClientError as err: + if err.reason_code in _IDENTIFY_OPERATION_FAILURE_CODES: + _exit_identify_operation_failed(err) + raise + typer.echo(json.dumps(result, indent=2, default=str)) + + +@app.command("dismiss-cluster") +@convey_cli +def dismiss_cluster( + cluster_id: int = typer.Argument(..., help="Cluster ID from discovery output."), + disposition: str = typer.Option( + ..., + "--disposition", + help="Dismissal disposition: not_a_person or quiet.", + ), +) -> None: + """Dismiss a current discovery cluster from listing surfaces.""" + try: + result = _request( + "POST", + "/app/speakers/api/discovery/dismiss", + json_body={"cluster_id": cluster_id, "disposition": disposition}, + ) except ConveyClientError as err: if err.reason_code == SPEAKER_COMMAND_FAILED.code: _exit_speaker_command_failed(err) @@ -752,6 +879,22 @@ def identify( typer.echo(json.dumps(result, indent=2, default=str)) +@app.command("dismissals") +@convey_cli +def dismissals() -> None: + """List folded speaker cluster dismissals.""" + result = _request("GET", "/app/speakers/api/discovery/dismissals") + typer.echo(json.dumps(result, indent=2, default=str)) + + +@app.command("keep-separate-list") +@convey_cli +def keep_separate_list() -> None: + """List folded speaker keep-separate assertions.""" + result = _request("GET", "/app/speakers/api/name-variants/keep-separate") + typer.echo(json.dumps(result, indent=2, default=str)) + + @app.command("merge-names") @convey_cli def merge_names_cmd( diff --git a/solstone/apps/speakers/maintenance.py b/solstone/apps/speakers/maintenance.py index c2747c335..b8a1f725c 100644 --- a/solstone/apps/speakers/maintenance.py +++ b/solstone/apps/speakers/maintenance.py @@ -213,8 +213,9 @@ def run_name_variants(args: list[str]) -> int: detection = detect_name_variant_candidates() created = 0 updated = 0 + suppressed = 0 for candidate in detection.get("candidates", []): - _, was_created = record_name_variant_candidate( + _, was_created, was_suppressed = record_name_variant_candidate( source_id=candidate["source_id"], source_label=candidate["source_label"], target_id=candidate["target_id"], @@ -222,17 +223,20 @@ def run_name_variants(args: list[str]) -> int: similarity=candidate["similarity"], readiness=candidate["readiness"], ) + if was_suppressed: + suppressed += 1 if was_created: created += 1 else: updated += 1 logger.info( - "speaker name variant candidates refreshed: journal=%s found=%d created=%d updated=%d path=%s", + "speaker name variant candidates refreshed: journal=%s found=%d created=%d updated=%d suppressed=%d path=%s", journal, len(detection.get("candidates", [])), created, updated, + suppressed, review_candidates_path(), ) return 0 diff --git a/solstone/apps/speakers/routes.py b/solstone/apps/speakers/routes.py index d1fa19236..2ad30b92f 100644 --- a/solstone/apps/speakers/routes.py +++ b/solstone/apps/speakers/routes.py @@ -53,6 +53,8 @@ from solstone.apps.speakers.discovery import ( discover_unknown_speakers, get_cluster_presence, identify_cluster, + load_discovery_cache, + undo_identify_operation, ) from solstone.apps.speakers.encoder_config import ( OWNER_BOOTSTRAP_MIN_STMTS, @@ -96,6 +98,10 @@ from solstone.convey.reasons import ( MISSING_REQUIRED_FIELD, SPEAKER_ATTRIBUTION_STATE_INVALID, SPEAKER_COMMAND_FAILED, + SPEAKER_IDENTIFY_CONFLICT, + SPEAKER_IDENTIFY_OPERATION_NOT_FOUND, + SPEAKER_IDENTIFY_RECOVERABLE, + SPEAKER_IDENTIFY_REPAIR_REQUIRED, SPEAKER_LABELS_BUSY, SPEAKER_NOT_FOUND, SPEAKER_OWNER_CENTROID_REQUIRED, @@ -123,6 +129,15 @@ from solstone.think.entities.journal import ( from solstone.think.journal_io.errors import LockTimeout from solstone.think.journal_io.npz import load_npz, update_npz from solstone.think.media import MIME_TYPES +from solstone.think.speaker_cluster_dismissals import ( + list_dismissals, + record_cluster_dismissal, +) +from solstone.think.speaker_identify_operations import ( + fold_all_operations, + fold_operation, +) +from solstone.think.speaker_keep_separate import list_assertions from solstone.think.utils import ( STREAM_RE, day_dirs, @@ -2226,25 +2241,154 @@ def api_cluster_presence(cluster_id: int) -> Any: return jsonify(presence) +def _operation_result_summary(result: dict[str, Any] | None) -> dict[str, Any] | None: + if not isinstance(result, dict): + return None + return { + key: result[key] + for key in ( + "status", + "operation_id", + "entity_id", + "entity_created", + "voiceprints_saved", + "retroactive_voiceprints_saved", + "segments_updated", + "sentences_attributed", + "keep_separate_assertions_recorded", + ) + if key in result + } + + +def _undo_report_summary(report: dict[str, Any] | None) -> dict[str, Any] | None: + if not isinstance(report, dict): + return None + summary = { + "status": report.get("status"), + "operation_id": report.get("operation_id"), + } + categories = report.get("undo_report") + if isinstance(categories, dict): + summary["undo_report"] = { + category: data + for category, data in categories.items() + if isinstance(data, dict) + } + return summary + + +def _repair_summary(repair: dict[str, Any] | None) -> dict[str, Any] | None: + if not isinstance(repair, dict): + return None + return { + key: repair[key] + for key in ( + "phase", + "repair_code", + "repair_categories", + "partial_report", + "undo_report", + ) + if key in repair + } + + +def _identify_operation_summary(state: Any) -> dict[str, Any]: + return { + "operation_id": state.operation_id, + "request_id": state.request_id, + "status": state.terminal_status, + "target_entity_id": state.target_entity_id, + "will_create": state.will_create, + "entity_type": state.entity_type, + "reviewed_near_match_entity_ids": list(state.reviewed_near_match_entity_ids), + "cluster_member_count": len(state.cluster_member_set), + "completed_phases": list(state.completed_phases), + "pending_phases": list(state.pending_phases), + "checkpoints": { + "forward": list(state.phase_checkpoints.keys()), + "undo": list(state.undo_phase_checkpoints.keys()), + }, + "result": _operation_result_summary(state.result), + "undo_report": _undo_report_summary(state.undo_report), + "repair": _repair_summary(state.repair_required) + or _repair_summary(state.undo_repair_required), + } + + +def _operation_failure_extra(result: dict[str, Any]) -> dict[str, Any]: + extra: dict[str, Any] = { + "status": result.get("status"), + "operation_id": result.get("operation_id"), + "operation_state": result.get("status"), + } + for key in ( + "request_id", + "completed_phases", + "pending_phases", + "repair_categories", + "repair_code", + "phase", + "conflict_code", + "conflicting_operation_id", + "list_command", + "undo_report", + ): + if key in result: + extra[key] = result[key] + return extra + + def _identify_result_response(result: dict) -> Any: status = result.get("status") - if status in {"identified", "resolved", "ambiguous", "no_match"}: + if status in { + "identified", + "resolved", + "ambiguous", + "no_match", + "undone", + "already_undone", + }: return jsonify(result) - if status == "partial": + if status == "recoverable": return error_response( - SPEAKER_COMMAND_FAILED, - status=409, + SPEAKER_IDENTIFY_RECOVERABLE, detail=result.get("detail"), - extra={ - "status": "partial", - "completed": result.get("completed", []), - "failed": result.get("failed", []), - }, + extra=_operation_failure_extra(result), + ) + if status in {"repair_required", "undo_repair_required"}: + return error_response( + SPEAKER_IDENTIFY_REPAIR_REQUIRED, + detail=result.get("repair_code") or result.get("detail"), + extra=_operation_failure_extra(result), + ) + if status in {"conflict", "operation_already_undone"}: + return error_response( + SPEAKER_IDENTIFY_CONFLICT, + detail=result.get("conflict_code") or status, + extra=_operation_failure_extra(result), + ) + if status == "not_found": + return error_response( + SPEAKER_IDENTIFY_OPERATION_NOT_FOUND, + detail=f"Operation {result.get('operation_id')} was not found.", + extra=_operation_failure_extra(result), ) if result.get("not_found"): return error_response(SPEAKER_NOT_FOUND, detail=result["error"]) if result.get("invalid_entity_type"): return error_response(INVALID_ENTITY_TYPE, detail=result["error"]) + if status == "invalid_request": + return error_response( + INVALID_REQUEST_VALUE, + detail=result.get("error"), + extra={ + key: result[key] + for key in ("invalid_reviewed_near_match_entity_ids",) + if key in result + }, + ) if "error" in result: return error_response( INVALID_REQUEST_VALUE, @@ -2254,6 +2398,18 @@ def _identify_result_response(result: dict) -> Any: return jsonify(result) +def _optional_request_id(data: dict[str, Any]) -> tuple[str | None, Any | None]: + request_id_value = data.get("request_id") + if request_id_value is None: + return None, None + if not isinstance(request_id_value, str) or not request_id_value.strip(): + return None, error_response( + INVALID_REQUEST_VALUE, + detail="request_id must be a non-empty string", + ) + return request_id_value.strip(), None + + @speakers_bp.route("/api/discovery/identify", methods=["POST"]) def api_discovery_identify() -> Any: """Identify a discovered unknown speaker cluster by naming it.""" @@ -2266,6 +2422,8 @@ def api_discovery_identify() -> Any: resolve_only = bool(data.get("resolve_only", False)) create_new = bool(data.get("create_new", False)) entity_type = data.get("entity_type") or "Person" + request_id, request_id_error = _optional_request_id(data) + reviewed_near_match_entity_ids = data.get("reviewed_near_match_entity_ids") if cluster_id is None: return error_response( @@ -2285,6 +2443,8 @@ def api_discovery_identify() -> Any: INVALID_REQUEST_VALUE, detail="cluster_id must be an integer", ) + if request_id_error is not None: + return request_id_error try: result = identify_cluster( @@ -2294,6 +2454,8 @@ def api_discovery_identify() -> Any: resolve_only=resolve_only, create_new=create_new, entity_type=entity_type, + request_id=request_id, + reviewed_near_match_entity_ids=reviewed_near_match_entity_ids, ) except LockTimeout as exc: if exc.path.name in ("speaker_labels.json", "speaker_corrections.json"): @@ -2480,6 +2642,8 @@ def api_cli_discovery_identify() -> Any: resolve_only = bool(data.get("resolve_only", False)) create_new = bool(data.get("create_new", False)) entity_type = data.get("entity_type") or "Person" + request_id, request_id_error = _optional_request_id(data) + reviewed_near_match_entity_ids = data.get("reviewed_near_match_entity_ids") if cluster_id is None: return error_response( @@ -2499,6 +2663,8 @@ def api_cli_discovery_identify() -> Any: INVALID_REQUEST_VALUE, detail="cluster_id must be an integer", ) + if request_id_error is not None: + return request_id_error try: result = identify_cluster( @@ -2508,6 +2674,8 @@ def api_cli_discovery_identify() -> Any: resolve_only=resolve_only, create_new=create_new, entity_type=entity_type, + request_id=request_id, + reviewed_near_match_entity_ids=reviewed_near_match_entity_ids, ) except LockTimeout as exc: if exc.path.name in ("speaker_labels.json", "speaker_corrections.json"): @@ -2517,6 +2685,112 @@ def api_cli_discovery_identify() -> Any: return _identify_result_response(result) +@speakers_bp.route("/api/discovery/identify/undo", methods=["POST"]) +def api_discovery_identify_undo() -> Any: + """Undo a committed discovery-cluster identify operation.""" + data = request.get_json(silent=True) or {} + operation_id_value = data.get("operation_id") + if not isinstance(operation_id_value, str) or not operation_id_value.strip(): + return error_response( + MISSING_REQUIRED_FIELD, + detail="operation_id is required", + ) + + try: + result = undo_identify_operation(operation_id_value.strip()) + except LockTimeout as exc: + if exc.path.name in ("speaker_labels.json", "speaker_corrections.json"): + return _labels_busy_response(exc) + return _voiceprint_busy_response(exc) + return _identify_result_response(result) + + +@speakers_bp.route("/api/discovery/identify/operations", methods=["GET"]) +def api_discovery_identify_operations() -> Any: + """Return redacted identify operation summaries.""" + operations = [_identify_operation_summary(state) for state in fold_all_operations()] + return jsonify({"operations": operations, "total": len(operations)}) + + +@speakers_bp.route( + "/api/discovery/identify/operations/", + methods=["GET"], +) +def api_discovery_identify_operation(operation_id: str) -> Any: + """Return one redacted identify operation summary.""" + state = fold_operation(operation_id) + if state is None: + return error_response( + SPEAKER_IDENTIFY_OPERATION_NOT_FOUND, + detail=f"Operation {operation_id} was not found.", + extra={ + "status": "not_found", + "operation_id": operation_id, + "list_command": "sol call speakers identify-operations", + }, + ) + return jsonify({"operation": _identify_operation_summary(state)}) + + +@speakers_bp.route("/api/discovery/dismiss", methods=["POST"]) +def api_discovery_dismiss() -> Any: + """Record a read-side suppression dismissal for the current cluster members.""" + data = request.get_json(silent=True) or {} + cluster_id = data.get("cluster_id") + disposition = data.get("disposition") + if cluster_id is None or not disposition: + return error_response( + MISSING_REQUIRED_FIELD, + detail="cluster_id and disposition are required", + ) + try: + cluster_id = int(cluster_id) + except (TypeError, ValueError): + return error_response( + INVALID_REQUEST_VALUE, + detail="cluster_id must be an integer", + ) + cache = load_discovery_cache() + members = cache.get("clusters", {}).get(str(cluster_id)) if cache else None + if not members: + return error_response( + SPEAKER_REVIEW_UNAVAILABLE, + detail=f"Cluster {cluster_id} was not found. Run a discovery scan first.", + ) + try: + event = record_cluster_dismissal(members, str(disposition)) + except ValueError as exc: + return error_response(INVALID_REQUEST_VALUE, detail=str(exc)) + except LockTimeout as exc: + return error_response( + SPEAKER_COMMAND_FAILED, + detail=str(exc), + status=503, + ) + return jsonify( + { + "status": "dismissed", + "dismiss_event_id": event["dismiss_event_id"], + "disposition": event["disposition"], + "member_count": event["member_count"], + } + ) + + +@speakers_bp.route("/api/discovery/dismissals", methods=["GET"]) +def api_discovery_dismissals() -> Any: + """Return folded cluster dismissal summaries.""" + dismissals = list_dismissals() + return jsonify({"dismissals": dismissals, "total": len(dismissals)}) + + +@speakers_bp.route("/api/name-variants/keep-separate", methods=["GET"]) +def api_name_variants_keep_separate() -> Any: + """Return folded keep-separate assertion summaries.""" + assertions = list_assertions() + return jsonify({"assertions": assertions, "total": len(assertions)}) + + @speakers_bp.route("/api/merge-names", methods=["POST"]) def api_cli_merge_names() -> Any: """Merge two speaker names for the CLI.""" diff --git a/solstone/apps/speakers/suggest.py b/solstone/apps/speakers/suggest.py index f304c60c5..4a06e82cf 100644 --- a/solstone/apps/speakers/suggest.py +++ b/solstone/apps/speakers/suggest.py @@ -12,6 +12,8 @@ from datetime import time from pathlib import Path from typing import Any +from solstone.think.speaker_keep_separate import name_variant_pair_suppressed +from solstone.think.speaker_review_candidates import detection_count_for_pair from solstone.think.utils import day_dirs, get_journal, iter_segments, segment_parse logger = logging.getLogger(__name__) @@ -258,6 +260,13 @@ def _name_variant() -> list[dict[str, Any]]: suggestions: list[dict[str, Any]] = [] for pair in detect_name_variant_candidates().get("candidates", []): + detection_count = detection_count_for_pair(pair["source_id"], pair["target_id"]) + if name_variant_pair_suppressed( + pair["source_id"], + pair["target_id"], + detection_count, + ): + continue suggestions.append( { "type": "name_variant", diff --git a/solstone/apps/speakers/tests/test_routes.py b/solstone/apps/speakers/tests/test_routes.py index 3941c47fa..062e64d13 100644 --- a/solstone/apps/speakers/tests/test_routes.py +++ b/solstone/apps/speakers/tests/test_routes.py @@ -3,6 +3,7 @@ """Tests for speakers app - sentence-based embeddings.""" +import hashlib import json from datetime import datetime from pathlib import Path @@ -94,6 +95,22 @@ def _read_action_entries(journal_root): ] +def _journal_file_hashes(journal_root: Path) -> dict[str, str]: + return { + path.relative_to(journal_root).as_posix(): hashlib.sha256( + path.read_bytes() + ).hexdigest() + for path in sorted(journal_root.rglob("*")) + if path.is_file() and not path.name.endswith(".lock") + } + + +def _changed_paths(before: dict[str, str], after: dict[str, str]) -> set[str]: + return { + path for path in set(before) | set(after) if before.get(path) != after.get(path) + } + + def _speakers_client(): from solstone.apps.speakers.routes import speakers_bp @@ -1273,6 +1290,8 @@ def test_discovery_identify_route_entity_id_only_and_result_mapping( "resolve_only": False, "create_new": False, "entity_type": "Person", + "request_id": None, + "reviewed_near_match_entity_ids": None, } @@ -1342,76 +1361,221 @@ def test_discovery_identify_route_invalid_entity_type(speakers_env): def test_discovery_identify_route_partial_and_retry(speakers_env, monkeypatch): from solstone.apps.speakers import discovery + from solstone.think.entities.voiceprints import load_entity_voiceprints_file env = speakers_env() env.create_entity("Bob Smith") _write_discovery_cluster(env, 36, "110500_300") _write_discovery_cluster(env, 37, "111000_300") client = _convey_client(env.journal) + request_id = "route-retry-request" - def fail_after_voiceprints(stage: str) -> None: - if stage == "after_voiceprints": - raise RuntimeError("forced after voiceprints") + def fail_after_prepared(stage: str) -> None: + if stage == "after_prepared": + raise RuntimeError("forced after prepared") monkeypatch.setattr( discovery, "_maybe_inject_identify_fault", - fail_after_voiceprints, + fail_after_prepared, ) - partial_voice = client.post( + recoverable = client.post( "/app/speakers/api/discovery/identify", - json={"cluster_id": 36, "name": "Bob Smith"}, + json={"cluster_id": 36, "name": "Bob Smith", "request_id": request_id}, ) - assert partial_voice.status_code == 409 - body = partial_voice.get_json() - assert body["reason_code"] == "speaker_command_failed" - assert body["status"] == "partial" - assert body["completed"] == ["voiceprints"] - assert body["failed"] == ["segments", "sentinel"] + assert recoverable.status_code == 409 + body = recoverable.get_json() + assert body["reason_code"] == "speaker_identify_recoverable" + assert body["status"] == "recoverable" + operation_id = body["operation_id"] + assert body["request_id"] == request_id + assert body["completed_phases"] == [] + assert "entity" in body["pending_phases"] monkeypatch.setattr( discovery, "_maybe_inject_identify_fault", lambda stage: None, ) - retry_voice = client.post( + retry = client.post( "/app/speakers/api/discovery/identify", - json={"cluster_id": 36, "name": "Bob Smith"}, + json={"cluster_id": 36, "name": "Bob Smith", "request_id": request_id}, ) - assert retry_voice.status_code == 200 - assert retry_voice.get_json()["status"] == "identified" + assert retry.status_code == 200 + retry_body = retry.get_json() + assert retry_body["status"] == "identified" + assert retry_body["operation_id"] == operation_id + assert len(load_entity_voiceprints_file("bob_smith")[1]) == 1 - def fail_after_segments(stage: str) -> None: - if stage == "after_segments": - raise RuntimeError("forced after segments") + replay = client.post( + "/app/speakers/api/discovery/identify", + json={"cluster_id": 36, "name": "Bob Smith", "request_id": request_id}, + ) + assert replay.status_code == 200 + assert replay.get_json() == retry_body + assert len(load_entity_voiceprints_file("bob_smith")[1]) == 1 - monkeypatch.setattr( - discovery, - "_maybe_inject_identify_fault", - fail_after_segments, + conflict = client.post( + "/app/speakers/api/discovery/identify", + json={"cluster_id": 37, "name": "Bob Smith", "request_id": request_id}, ) - partial_segments = client.post( + assert conflict.status_code == 409 + conflict_body = conflict.get_json() + assert conflict_body["reason_code"] == "speaker_identify_conflict" + assert conflict_body["conflict_code"] == "request_fingerprint_mismatch" + + legacy_success = client.post( "/app/speakers/api/discovery/identify", json={"cluster_id": 37, "name": "Bob Smith"}, ) + assert legacy_success.status_code == 200 + assert legacy_success.get_json()["status"] == "identified" - assert partial_segments.status_code == 409 - body = partial_segments.get_json() - assert body["completed"] == ["voiceprints", "segments"] - assert body["failed"] == ["sentinel"] - + _write_discovery_cluster(env, 38, "111500_300") monkeypatch.setattr( discovery, "_maybe_inject_identify_fault", - lambda stage: None, + fail_after_prepared, ) - retry_segments = client.post( + legacy_recoverable = client.post( "/app/speakers/api/discovery/identify", - json={"cluster_id": 37, "name": "Bob Smith"}, + json={"cluster_id": 38, "name": "Bob Smith"}, + ) + assert legacy_recoverable.status_code == 409 + assert ( + legacy_recoverable.get_json()["reason_code"] == "speaker_identify_recoverable" + ) + + +def test_discovery_identify_operations_and_undo_routes(speakers_env): + env = speakers_env() + env.create_entity("Bob Smith") + _write_discovery_cluster(env, 39, "112000_300") + client = _convey_client(env.journal) + + identify = client.post( + "/app/speakers/api/discovery/identify-cli", + json={ + "cluster_id": 39, + "name": "Bob Smith", + "request_id": "route-operation-roundtrip", + }, + ) + + assert identify.status_code == 200 + operation_id = identify.get_json()["operation_id"] + + listing = client.get("/app/speakers/api/discovery/identify/operations") + assert listing.status_code == 200 + listed = listing.get_json() + assert listed["total"] == 1 + assert listed["operations"][0]["operation_id"] == operation_id + assert "prepared_plan" not in json.dumps(listed) + assert "Bob Smith" not in json.dumps(listed) + + shown = client.get( + f"/app/speakers/api/discovery/identify/operations/{operation_id}" + ) + assert shown.status_code == 200 + shown_body = shown.get_json() + assert shown_body["operation"]["operation_id"] == operation_id + assert "prepared_plan" not in json.dumps(shown_body) + assert "Bob Smith" not in json.dumps(shown_body) + + undo = client.post( + "/app/speakers/api/discovery/identify/undo", + json={"operation_id": operation_id}, ) - assert retry_segments.status_code == 200 - assert retry_segments.get_json()["status"] == "identified" + assert undo.status_code == 200 + assert undo.get_json()["status"] == "undone" + + second_undo = client.post( + "/app/speakers/api/discovery/identify/undo", + json={"operation_id": operation_id}, + ) + assert second_undo.status_code == 200 + assert second_undo.get_json()["status"] == "already_undone" + + missing = client.get("/app/speakers/api/discovery/identify/operations/idop_missing") + assert missing.status_code == 404 + assert missing.get_json()["reason_code"] == "speaker_identify_operation_not_found" + + +def test_discovery_dismissals_and_keep_separate_list_routes_are_store_only( + speakers_env, +): + from solstone.think.speaker_keep_separate import record_keep_separate_assertion + + env = speakers_env() + _write_discovery_cluster(env, 40, "112500_300") + client = _convey_client(env.journal) + before = _journal_file_hashes(env.journal) + + dismissed = client.post( + "/app/speakers/api/discovery/dismiss", + json={"cluster_id": 40, "disposition": "quiet"}, + ) + record_keep_separate_assertion( + "alice", + "alice_johnson", + source_kind="explicit_create_near_match", + operation_id="idop_route_test", + detection_count=1, + ) + dismissals = client.get("/app/speakers/api/discovery/dismissals") + keep_separate = client.get("/app/speakers/api/name-variants/keep-separate") + after = _journal_file_hashes(env.journal) + + assert dismissed.status_code == 200 + assert dismissed.get_json()["member_count"] == 1 + assert dismissals.status_code == 200 + assert dismissals.get_json()["total"] == 1 + assert keep_separate.status_code == 200 + assert keep_separate.get_json()["total"] == 1 + assert _changed_paths(before, after) <= { + "speakers/cluster-dismissals.jsonl", + "speakers/keep-separate.jsonl", + } + + +def test_api_suggest_filters_keep_separate_name_variant(speakers_env): + from solstone.think.speaker_keep_separate import record_keep_separate_assertion + + env = speakers_env() + env.create_entity("Alice") + env.create_entity("Alice Test") + base = env.create_embedding([1.0, 0.0, 0.0]) + similar = env.create_embedding([1.0, 0.01, 0.0]) + _write_voiceprints( + env, + "alice", + [ + (base, "20240101", "100000_300", "mic_audio", 1), + (similar, "20240101", "100005_300", "mic_audio", 2), + ], + ) + _write_voiceprints( + env, + "alice_test", + [ + (similar, "20240101", "101000_300", "mic_audio", 1), + (base, "20240101", "101005_300", "mic_audio", 2), + ], + ) + record_keep_separate_assertion( + "alice", + "alice_test", + source_kind="explicit_create_near_match", + operation_id="idop_suggest_route", + detection_count=1, + ) + client = _convey_client(env.journal) + + response = client.get("/app/speakers/api/suggest") + + assert response.status_code == 200 + assert all(item["type"] != "name_variant" for item in response.get_json()["items"]) def test_workspace_discovery_identify_requests_create_new(): diff --git a/solstone/apps/speakers/tests/test_suggest.py b/solstone/apps/speakers/tests/test_suggest.py index 7bceb7949..adf977e6a 100644 --- a/solstone/apps/speakers/tests/test_suggest.py +++ b/solstone/apps/speakers/tests/test_suggest.py @@ -16,6 +16,7 @@ from solstone.apps.speakers.suggest import ( from solstone.think.speaker_candidate_pair_review_candidates import ( record_candidate_pair, ) +from solstone.think.speaker_keep_separate import record_keep_separate_assertion def create_meetings_md(env, day: str, content: str) -> Path: @@ -136,6 +137,28 @@ def test_suggest_name_variant(speakers_env): assert suggestion["readiness"] == "ready" +def test_suggest_name_variant_respects_keep_separate(speakers_env): + env = speakers_env() + alice_dir = env.create_entity("Alice") + alice_test_dir = env.create_entity("Alice Test") + + base = env.create_embedding([1.0, 0.0, 0.0]) + similar = env.create_embedding([1.0, 0.01, 0.0]) + _write_voiceprints(alice_dir, [base, similar]) + _write_voiceprints(alice_test_dir, [similar, base]) + record_keep_separate_assertion( + "alice", + "alice_test", + source_kind="explicit_create_near_match", + operation_id="idop_test", + detection_count=1, + ) + + results = suggest_opportunities() + + assert all(item["type"] != "name_variant" for item in results) + + def test_suggest_speaker_candidate_pair_and_formats_it(speakers_env): speakers_env() anchor_a = '["20260101","090000_300","test","mic_audio",1]' diff --git a/solstone/convey/contract/observer_bundle.py b/solstone/convey/contract/observer_bundle.py index 0f7750d58..9fbf2ae85 100644 --- a/solstone/convey/contract/observer_bundle.py +++ b/solstone/convey/contract/observer_bundle.py @@ -20,7 +20,7 @@ from solstone.convey.contract.assemble import CALLOSUM_REGISTRY, build_document from solstone.observe import protocol INITIAL_BUNDLE_SEMVER = "1.0.0" -BUNDLE_SEMVER = "1.0.2" +BUNDLE_SEMVER = "2.0.0" GENERATOR_IDENTITY = "solstone.convey.contract.observer_bundle.v1" BUNDLE_SCHEMA_IDENTITY = "solstone.observer-client-contract-bundle.schema.v1" SCHEMA_DIALECT_URI = "https://json-schema.org/draft/2020-12/schema" diff --git a/solstone/convey/reasons.py b/solstone/convey/reasons.py index 4e5e061c0..6bbbfdb1e 100644 --- a/solstone/convey/reasons.py +++ b/solstone/convey/reasons.py @@ -412,6 +412,26 @@ SPEAKER_COMMAND_FAILED = Reason( "I couldn't finish that speaker command.", 400, ) +SPEAKER_IDENTIFY_RECOVERABLE = Reason( + "speaker_identify_recoverable", + "I couldn't finish that speaker identify operation, but it can be retried.", + 409, +) +SPEAKER_IDENTIFY_REPAIR_REQUIRED = Reason( + "speaker_identify_repair_required", + "I couldn't safely finish that speaker identify operation without repair.", + 409, +) +SPEAKER_IDENTIFY_CONFLICT = Reason( + "speaker_identify_conflict", + "I couldn't run that speaker identify operation because it conflicts with existing state.", + 409, +) +SPEAKER_IDENTIFY_OPERATION_NOT_FOUND = Reason( + "speaker_identify_operation_not_found", + "I couldn't find that speaker identify operation.", + 404, +) AWARENESS_BUSY = Reason( "awareness_busy", "I couldn't update what I know right now because it was busy. Try again in a moment.", diff --git a/solstone/talent/sol/references/commands.md b/solstone/talent/sol/references/commands.md index 1a885f3a3..ed9537a4d 100644 --- a/solstone/talent/sol/references/commands.md +++ b/solstone/talent/sol/references/commands.md @@ -40,11 +40,11 @@ Guidance: `solstone/apps/health/talent/health/SKILL.md` Triggers: `speaker`, `voice`, `who was talking`, `identify speaker`, `voiceprint`, `tag-owner`, `build-from-tags`, `rebuild-owner`, `owner-ready` -Read: `discover`, `resolve-names`, `status` +Read: `discover`, `keep-separate-list`, `resolve-names`, `status` Write: `backfill`, `backfill-last-seen`, `bootstrap`, `link-import`, `merge-names`, `rebuild-owner`, `seed-from-imports` -Other: `attribute-segment`, `build-from-tags`, `confirm-owner`, `correct`, `day-segments`, `detect`, `identify`, `owner-ready`, `presence`, `propagate-correction`, `reject-owner`, `sentences`, `suggest`, `tag-owner`, `wipe` +Other: `attribute-segment`, `build-from-tags`, `confirm-owner`, `correct`, `day-segments`, `detect`, `dismiss-cluster`, `dismissals`, `identify`, `identify-operation`, `identify-operations`, `identify-undo`, `owner-ready`, `presence`, `propagate-correction`, `reject-owner`, `sentences`, `suggest`, `tag-owner`, `wipe` Guidance: `solstone/apps/speakers/talent/speakers/SKILL.md` diff --git a/solstone/think/curation.py b/solstone/think/curation.py index 5065d146b..95148e73f 100644 --- a/solstone/think/curation.py +++ b/solstone/think/curation.py @@ -13,6 +13,7 @@ import solstone.think.facet_review_candidates as facet_store from solstone.think import ( speaker_candidate_pair_review_candidates as speaker_pair_store, ) +from solstone.think import speaker_keep_separate as keep_separate_store from solstone.think import speaker_review_candidates as speaker_store from solstone.think.entities import review_candidates as entity_store from solstone.think.entities.ambiguities import load_ambiguities @@ -363,6 +364,13 @@ def load_open_items() -> list[CurationItem]: evidence = {} source_id = str(row.get("source_id") or "") target_id = str(row.get("target_id") or "") + detection_count = _int_value(evidence.get("detection_count")) or 1 + if keep_separate_store.name_variant_pair_suppressed( + source_id, + target_id, + detection_count, + ): + continue similarity = float(row["similarity"]) speaker_evidence = dict(evidence) speaker_evidence["similarity"] = similarity diff --git a/solstone/think/speaker_review_candidates.py b/solstone/think/speaker_review_candidates.py index 715526276..ff09128c6 100644 --- a/solstone/think/speaker_review_candidates.py +++ b/solstone/think/speaker_review_candidates.py @@ -15,6 +15,7 @@ from pathlib import Path from typing import Any, Callable from solstone.think.journal_io import atomic_replace, hold_lock +from solstone.think.speaker_keep_separate import name_variant_pair_suppressed from solstone.think.utils import get_journal logger = logging.getLogger(__name__) @@ -102,6 +103,23 @@ def find_candidate( return None +def _detection_count_from_row(row: dict[str, Any] | None) -> int: + if not isinstance(row, dict): + return 1 + evidence = row.get("evidence") + if not isinstance(evidence, dict): + return 1 + try: + return max(1, int(evidence.get("detection_count", 1))) + except (TypeError, ValueError): + return 1 + + +def detection_count_for_pair(id_a: str, id_b: str) -> int: + """Return the current detection count for a candidate pair, defaulting to 1.""" + return _detection_count_from_row(find_candidate(load_candidates(), id_a, id_b)) + + def locked_modify_candidates( fn: Callable[[list[dict[str, Any]]], list[dict[str, Any]]], ) -> list[dict[str, Any]]: @@ -151,23 +169,30 @@ def record_name_variant_candidate( target_label: str, similarity: float, readiness: str = "ready", -) -> tuple[dict[str, Any], bool]: +) -> tuple[dict[str, Any], bool, bool]: """Create or update one speaker name-variant candidate.""" row: dict[str, Any] | None = None created = False + suppressed = False def mutate(rows: list[dict[str, Any]]) -> list[dict[str, Any]]: - nonlocal row, created + nonlocal row, created, suppressed existing = find_candidate(rows, source_id, target_id) now = utc_now_iso() if existing is None: + detection_count = 1 + suppressed = name_variant_pair_suppressed( + source_id, + target_id, + detection_count, + ) # Deterministic pairs reaching the recorder already met threshold and gates; "waiting" is reserved. row = { "source_id": source_id, "source_label": source_label, "target_id": target_id, "target_label": target_label, - "status": "open", + "status": "suppressed" if suppressed else "open", "similarity": similarity, "readiness": readiness, "evidence": _evidence( @@ -175,13 +200,16 @@ def record_name_variant_candidate( target_label, similarity, readiness, - 1, + detection_count, ), "first_surfaced": now, "last_surfaced": now, "created_at": now, "updated_at": now, } + if suppressed: + row["suppressed_by_keep_separate"] = True + row["suppressed_detection_count"] = detection_count created = True return list(rows) + [row] @@ -201,18 +229,36 @@ def record_name_variant_candidate( existing["similarity"] = similarity existing["readiness"] = readiness next_evidence = dict(evidence) + next_detection_count = detection_count + 1 next_evidence.update( _evidence( source_label, target_label, similarity, readiness, - detection_count + 1, + next_detection_count, ) ) existing["evidence"] = next_evidence existing["last_surfaced"] = now existing["updated_at"] = now + status = str(existing.get("status") or "") + suppressed_now = name_variant_pair_suppressed( + source_id, + target_id, + next_detection_count, + ) + if status not in {"accepted", "dismissed"}: + if suppressed_now: + existing["status"] = "suppressed" + existing["suppressed_by_keep_separate"] = True + existing["suppressed_detection_count"] = next_detection_count + suppressed = True + elif existing.get("suppressed_by_keep_separate") is True: + existing["status"] = "open" + existing.pop("suppressed_by_keep_separate", None) + existing.pop("suppressed_detection_count", None) + suppressed = False row = existing created = False return rows @@ -222,7 +268,7 @@ def record_name_variant_candidate( if row is None: # pragma: no cover - defensive assertion raise RuntimeError("record-name-variant-candidate produced no row") - return row, created + return row, created, suppressed def touch_updated(row: dict[str, Any]) -> None: diff --git a/tests/test_curation.py b/tests/test_curation.py index 2af5a216c..259a9c07c 100644 --- a/tests/test_curation.py +++ b/tests/test_curation.py @@ -59,6 +59,7 @@ from solstone.think.speaker_candidate_pair_review_candidates import ( from solstone.think.speaker_candidate_pair_review_candidates import ( record_candidate_pair, ) +from solstone.think.speaker_keep_separate import record_keep_separate_assertion from solstone.think.speaker_review_candidates import record_name_variant_candidate @@ -343,6 +344,27 @@ def test_load_open_items_includes_speaker_name_variant(curation_journal): assert item.strength == 93 +def test_load_open_items_filters_keep_separate_speaker_name_variant(curation_journal): + record_name_variant_candidate( + source_id="alice", + source_label="Alice", + target_id="alice_johnson", + target_label="Alice Johnson", + similarity=0.934, + ) + record_keep_separate_assertion( + "alice", + "alice_johnson", + source_kind="explicit_create_near_match", + operation_id="idop_test", + detection_count=1, + ) + + items = load_open_items() + + assert [item for item in items if item.kind == KIND_SPEAKER_NAME_VARIANT] == [] + + def test_load_open_items_includes_speaker_candidate_pair(curation_journal): left = _candidate_profile(1, _unit([1.0, 0.0])) right = _candidate_profile(2, _unit([0.62, np.sqrt(1.0 - 0.62**2)])) diff --git a/tests/test_observer_client_bundle.py b/tests/test_observer_client_bundle.py index 5da219f2d..c98fc3620 100644 --- a/tests/test_observer_client_bundle.py +++ b/tests/test_observer_client_bundle.py @@ -180,7 +180,7 @@ EXPECTED_BUNDLE_PAYLOAD_SHA256 = { "9749a50daba9b4a270da045d350bc5edb7a42c9723fa0bf420c8fb8a4a0415f8" ), observer_bundle.PROJECTION_REL: ( - "8a2b7037552edf710597f2ffa6fdc5aa715311df4ea8cf168e70abe4231c64ca" + "28c055279ab7d80c809a43c5f710ccebc4da643fec3f8e6beae48823c17c46c5" ), observer_bundle.VECTORS_REL: ( "7a5132c57b61e2a615a22719abc77e40b708d4a6636c45690cc522dc26c36dec" @@ -707,7 +707,7 @@ def test_observer_client_bundle_preserves_distinct_version_concepts() -> None: }, ) - assert manifest["bundle_semver"] == "1.0.2" + assert manifest["bundle_semver"] == observer_bundle.BUNDLE_SEMVER assert projection["info"]["version"] == "9.8.7" assert manifest["openapi_document_version"] == "9.8.7" assert manifest["openapi_spec_version"] == "3.1.9" @@ -715,7 +715,7 @@ def test_observer_client_bundle_preserves_distinct_version_concepts() -> None: manifest["bundle_semver"], manifest["openapi_document_version"], manifest["openapi_spec_version"], - } == {"1.0.2", "9.8.7", "3.1.9"} + } == {observer_bundle.BUNDLE_SEMVER, "9.8.7", "3.1.9"} def test_observer_client_bundle_manifest_file_inventory( @@ -729,7 +729,7 @@ def test_observer_client_bundle_manifest_file_inventory( "vectors.json", ] - assert manifest["bundle_semver"] == "1.0.2" + assert manifest["bundle_semver"] == observer_bundle.BUNDLE_SEMVER assert manifest["openapi_document_version"] == "1.0.0" assert manifest["openapi_spec_version"] == "3.1.0" assert manifest["observer_protocol_version"] == 2 @@ -2236,14 +2236,14 @@ def test_observer_client_bundle_history_equal_and_downgrade_failures( _commit_bundle(equal_repo, bundle_files, "bundle") _assert_history_failure( equal_repo, - _description_patch_candidate(bundle_files, "1.0.2"), + _description_patch_candidate(bundle_files, observer_bundle.BUNDLE_SEMVER), "without a version bump", ) downgrade_repo = tmp_path / "downgrade" _init_git_repo(downgrade_repo) previous = _clone_bundle_files(bundle_files) - _set_bundle_version(previous, "1.1.0") + _set_bundle_version(previous, "2.1.0") _commit_bundle(downgrade_repo, previous, "minor bundle") _assert_history_failure(downgrade_repo, bundle_files, "downgraded") @@ -2252,7 +2252,7 @@ def test_observer_client_bundle_history_equal_and_downgrade_failures( ("version", "expected_failure"), [ ("1.0.1", "without a version bump"), - ("1.0.2", None), + ("2.0.0", None), ], ) def test_observer_client_bundle_history_fixture_input_removal_requires_patch_bump( @@ -2296,7 +2296,7 @@ def test_observer_client_bundle_history_major_and_mixed_insufficient_bumps_fail( _commit_bundle(closed_repo, bundle_files, "bundle") _assert_history_failure( closed_repo, - _closed_status_add_candidate(bundle_files, "1.1.0"), + _closed_status_add_candidate(bundle_files, "2.1.0"), "major change", enforce_current_contract=False, ) @@ -2304,7 +2304,7 @@ def test_observer_client_bundle_history_major_and_mixed_insufficient_bumps_fail( mixed_repo = tmp_path / "mixed" _init_git_repo(mixed_repo) _commit_bundle(mixed_repo, bundle_files, "bundle") - mixed = _operation_removed_candidate(bundle_files, "1.1.0") + mixed = _operation_removed_candidate(bundle_files, "2.1.0") projection = _payload(mixed, observer_bundle.PROJECTION_REL) projection["paths"]["/app/observer/ingest"]["post"]["responses"]["200"]["content"][ "application/json" @@ -2329,7 +2329,7 @@ def test_observer_client_bundle_history_extensible_addition_accepts_minor( assert ( observer_bundle_compatibility.check_bundle_compatibility( repo, - _extensible_chat_add_candidate(bundle_files, "1.1.0"), + _extensible_chat_add_candidate(bundle_files, "2.1.0"), enforce_current_contract=False, ) == [] @@ -2351,7 +2351,7 @@ def test_observer_client_bundle_history_vector_removed_requires_major( if vector["id"] != "observer.auth.handle" ] _set_payload(removed, observer_bundle.VECTORS_REL, vectors) - _set_bundle_version(removed, "1.1.0") + _set_bundle_version(removed, "2.1.0") _assert_history_failure(repo, removed, "major change") @@ -2378,7 +2378,7 @@ def test_observer_client_bundle_history_does_not_apply_current_vector_policy_to_ _set_payload(baseline, observer_bundle.VECTORS_REL, vectors) _commit_bundle(repo, baseline, "legacy vector policy") candidate = _clone_bundle_files(bundle_files) - _set_bundle_version(candidate, "2.0.0") + _set_bundle_version(candidate, "3.0.0") assert ( observer_bundle_compatibility.check_bundle_compatibility( @@ -2455,7 +2455,7 @@ def test_observer_client_bundle_history_new_independent_vector_accepts_minor( ) vectors["vectors"] = sorted(vectors["vectors"], key=lambda item: item["id"]) _set_payload(added, observer_bundle.VECTORS_REL, vectors) - _set_bundle_version(added, "1.1.0") + _set_bundle_version(added, "2.1.0") assert ( observer_bundle_compatibility.check_bundle_compatibility( @@ -2477,14 +2477,14 @@ def test_observer_client_bundle_history_semantic_vector_change_is_never_patch( _assert_history_failure( repo, - _semantic_status_fixture_candidate(bundle_files, "1.1.0"), + _semantic_status_fixture_candidate(bundle_files, "2.1.0"), "major change", enforce_current_contract=False, ) assert ( observer_bundle_compatibility.check_bundle_compatibility( repo, - _semantic_status_fixture_candidate(bundle_files, "2.0.0"), + _semantic_status_fixture_candidate(bundle_files, "3.0.0"), enforce_current_contract=False, ) == [] diff --git a/tests/test_speaker_review_candidates.py b/tests/test_speaker_review_candidates.py index 1a3913d31..0268eb0fa 100644 --- a/tests/test_speaker_review_candidates.py +++ b/tests/test_speaker_review_candidates.py @@ -10,9 +10,11 @@ from pathlib import Path import pytest import solstone.think.speaker_review_candidates as mod +from solstone.think.speaker_keep_separate import record_keep_separate_assertion from solstone.think.speaker_review_candidates import ( accept_candidate, candidate_key, + detection_count_for_pair, dismiss_candidate, find_candidate, load_candidates, @@ -134,7 +136,7 @@ def test_record_name_variant_candidate_creates_exact_open_row( ): monkeypatch.setattr(mod, "utc_now_iso", lambda: "2026-06-03T17:30:00Z") - row, created = record_name_variant_candidate( + row, created, suppressed = record_name_variant_candidate( source_id="alice", source_label="Alice", target_id="alice_johnson", @@ -143,6 +145,7 @@ def test_record_name_variant_candidate_creates_exact_open_row( ) assert created is True + assert suppressed is False assert row == { "source_id": "alice", "source_label": "Alice", @@ -182,7 +185,7 @@ def test_record_name_variant_candidate_upserts_opposite_order_and_label_change( target_label="Alice Johnson", similarity=0.91, ) - row, created = record_name_variant_candidate( + row, created, suppressed = record_name_variant_candidate( source_id="alice_johnson", source_label="Alice J.", target_id="alice", @@ -191,6 +194,7 @@ def test_record_name_variant_candidate_upserts_opposite_order_and_label_change( ) assert created is False + assert suppressed is False rows = load_candidates() assert rows == [row] assert row["source_id"] == "alice_johnson" @@ -229,7 +233,7 @@ def test_record_name_variant_candidate_preserves_status_and_unknown_keys( similarity=0.91, ) accept_candidate("alice", "alice_johnson") - accepted, _ = record_name_variant_candidate( + accepted, _, _ = record_name_variant_candidate( source_id="alice", source_label="Alice A.", target_id="alice_johnson", @@ -239,7 +243,7 @@ def test_record_name_variant_candidate_preserves_status_and_unknown_keys( accepted["custom"] = "keep-me" accepted["evidence"]["custom_evidence"] = "keep-me-too" save_candidates([accepted]) - dismissed, _ = record_name_variant_candidate( + dismissed, _, _ = record_name_variant_candidate( source_id="bob", source_label="Bob", target_id="bob_smith", @@ -247,7 +251,7 @@ def test_record_name_variant_candidate_preserves_status_and_unknown_keys( similarity=0.94, ) dismiss_candidate("bob", "bob_smith") - dismissed, _ = record_name_variant_candidate( + dismissed, _, _ = record_name_variant_candidate( source_id="bob", source_label="Bobby", target_id="bob_smith", @@ -268,6 +272,59 @@ def test_record_name_variant_candidate_preserves_status_and_unknown_keys( assert dismissed_after["evidence"]["summary"].startswith("Bobby") +def test_record_name_variant_candidate_suppresses_then_resurfaces( + candidate_journal, monkeypatch +): + times = iter( + [ + "2026-06-03T17:30:00Z", + "2026-06-04T17:30:00Z", + "2026-06-05T17:30:00Z", + "2026-06-06T17:30:00Z", + ] + ) + monkeypatch.setattr(mod, "utc_now_iso", lambda: next(times)) + record_keep_separate_assertion( + "alice_johnson", + "alice", + source_kind="explicit_create_near_match", + operation_id="idop_test", + detection_count=2, + ) + + first, created, suppressed = record_name_variant_candidate( + source_id="alice", + source_label="Alice", + target_id="alice_johnson", + target_label="Alice Johnson", + similarity=0.91, + ) + second, _, still_suppressed = record_name_variant_candidate( + source_id="alice", + source_label="Alice", + target_id="alice_johnson", + target_label="Alice Johnson", + similarity=0.92, + ) + third, _, resurfaced_suppressed = record_name_variant_candidate( + source_id="alice", + source_label="Alice", + target_id="alice_johnson", + target_label="Alice Johnson", + similarity=0.93, + ) + + assert created is True + assert suppressed is True + assert first["status"] == "suppressed" + assert second["status"] == "suppressed" + assert still_suppressed is True + assert third["status"] == "open" + assert resurfaced_suppressed is False + assert "suppressed_by_keep_separate" not in third + assert detection_count_for_pair("alice_johnson", "alice") == 3 + + def test_accept_candidate_sets_status_and_updated_at(candidate_journal, monkeypatch): monkeypatch.setattr(mod, "utc_now_iso", lambda: "2026-06-03T17:30:00Z") save_candidates( diff --git a/tests/test_speakers_call_parity.py b/tests/test_speakers_call_parity.py index 2201c879b..75e0f9c47 100644 --- a/tests/test_speakers_call_parity.py +++ b/tests/test_speakers_call_parity.py @@ -5,6 +5,7 @@ from __future__ import annotations import json from pathlib import Path +from types import SimpleNamespace from typing import Any import pytest @@ -769,6 +770,21 @@ def test_identify_forwards_entity_id_create_resolve_only_and_no_match( ["identify", "5", "Alice", "--create", "--entity-type", "Organization"], ) resolve = runner.invoke(app, ["identify", "6", "Alice", "--resolve-only"]) + with_request = runner.invoke( + app, + [ + "identify", + "8", + "Alice", + "--create", + "--request-id", + "req-identify", + "--reviewed-near-match-entity-id", + "bob", + "--reviewed-near-match-entity-id", + "carol", + ], + ) missing_target = runner.invoke(app, ["identify", "7"]) _assert_json_stdout( @@ -794,6 +810,10 @@ def test_identify_forwards_entity_id_create_resolve_only_and_no_match( "has_voice": False, }, ) + _assert_json_stdout( + with_request, + {"status": "identified", "entity_id": "alice", "entity_created": True}, + ) assert missing_target.exit_code != 0 assert seen == [ { @@ -803,6 +823,8 @@ def test_identify_forwards_entity_id_create_resolve_only_and_no_match( "resolve_only": False, "create_new": False, "entity_type": "Person", + "request_id": None, + "reviewed_near_match_entity_ids": None, }, { "cluster_id": 4, @@ -811,6 +833,8 @@ def test_identify_forwards_entity_id_create_resolve_only_and_no_match( "resolve_only": False, "create_new": False, "entity_type": "Person", + "request_id": None, + "reviewed_near_match_entity_ids": None, }, { "cluster_id": 5, @@ -819,6 +843,8 @@ def test_identify_forwards_entity_id_create_resolve_only_and_no_match( "resolve_only": False, "create_new": True, "entity_type": "Organization", + "request_id": None, + "reviewed_near_match_entity_ids": None, }, { "cluster_id": 6, @@ -827,10 +853,143 @@ def test_identify_forwards_entity_id_create_resolve_only_and_no_match( "resolve_only": True, "create_new": False, "entity_type": "Person", + "request_id": None, + "reviewed_near_match_entity_ids": None, + }, + { + "cluster_id": 8, + "name": "Alice", + "entity_id": None, + "resolve_only": False, + "create_new": True, + "entity_type": "Person", + "request_id": "req-identify", + "reviewed_near_match_entity_ids": ["bob", "carol"], }, ] +def test_identify_recoverable_cli_exits_with_retry_guidance( + runner: CliRunner, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setattr( + speakers_routes, + "identify_cluster", + lambda *_args, **_kwargs: { + "status": "recoverable", + "operation_id": "idop_retry", + "request_id": "req-retry", + "completed_phases": [], + "pending_phases": ["entity"], + "detail": "forced", + }, + ) + + result = runner.invoke( + app, + ["identify", "9", "Alice", "--request-id", "req-retry"], + ) + + assert result.exit_code == 1 + assert result.stdout == "" + assert "forced" in result.stderr + assert "Retry with the same --request-id req-retry." in result.stderr + assert "sol call speakers identify-operations" in result.stderr + assert "sol call speakers identify-operation idop_retry" in result.stderr + + +def test_identify_operation_dismissal_and_keep_separate_cli_verbs( + runner: CliRunner, monkeypatch: pytest.MonkeyPatch +) -> None: + state = SimpleNamespace( + operation_id="idop_123", + request_id="req-123", + terminal_status="committed", + target_entity_id="alice", + will_create=False, + entity_type="Person", + reviewed_near_match_entity_ids=(), + cluster_member_set=frozenset({("20240101", "test", "120000_300", "audio", 1)}), + completed_phases=("entity", "sentinel"), + pending_phases=(), + phase_checkpoints={"entity": {}, "sentinel": {}}, + undo_phase_checkpoints={}, + result={"status": "identified", "operation_id": "idop_123"}, + undo_report=None, + repair_required=None, + undo_repair_required=None, + ) + member = { + "day": "20240101", + "stream": "test", + "segment_key": "120000_300", + "source": "audio", + "sentence_id": 1, + } + monkeypatch.setattr(speakers_routes, "fold_all_operations", lambda: [state]) + monkeypatch.setattr( + speakers_routes, + "fold_operation", + lambda operation_id: state if operation_id == "idop_123" else None, + ) + monkeypatch.setattr( + speakers_routes, + "undo_identify_operation", + lambda operation_id: {"status": "undone", "operation_id": operation_id}, + ) + monkeypatch.setattr( + speakers_routes, + "load_discovery_cache", + lambda: {"clusters": {"7": [member]}}, + ) + monkeypatch.setattr( + speakers_routes, + "record_cluster_dismissal", + lambda members, disposition: { + "dismiss_event_id": "cdev_123", + "disposition": disposition, + "member_count": len(members), + }, + ) + monkeypatch.setattr( + speakers_routes, + "list_dismissals", + lambda: [{"dismissal_id": "cdev_123", "member_count": 1}], + ) + monkeypatch.setattr( + speakers_routes, + "list_assertions", + lambda: [{"assertion_id": "ksep_123", "source_count": 1}], + ) + + undo = runner.invoke(app, ["identify-undo", "idop_123"]) + operations = runner.invoke(app, ["identify-operations"]) + operation = runner.invoke(app, ["identify-operation", "idop_123"]) + dismiss = runner.invoke( + app, + ["dismiss-cluster", "7", "--disposition", "quiet"], + ) + dismissals = runner.invoke(app, ["dismissals"]) + keep_separate = runner.invoke(app, ["keep-separate-list"]) + + _assert_json_stdout(undo, {"status": "undone", "operation_id": "idop_123"}) + assert json.loads(operations.stdout)["operations"][0]["operation_id"] == "idop_123" + assert json.loads(operation.stdout)["operation"]["operation_id"] == "idop_123" + _assert_json_stdout( + dismiss, + { + "status": "dismissed", + "dismiss_event_id": "cdev_123", + "disposition": "quiet", + "member_count": 1, + }, + ) + assert json.loads(dismissals.stdout)["dismissals"][0]["dismissal_id"] == "cdev_123" + assert ( + json.loads(keep_separate.stdout)["assertions"][0]["assertion_id"] == "ksep_123" + ) + + def test_merge_names_success_simple_error_and_multi_key_error( runner: CliRunner, monkeypatch: pytest.MonkeyPatch ) -> None: -- 2.51.2