diff --git a/CHANGELOG.d/voice-cutoff-reassertion.md b/CHANGELOG.d/voice-cutoff-reassertion.md new file mode 100644 index 000000000..63ab384c8 --- /dev/null +++ b/CHANGELOG.d/voice-cutoff-reassertion.md @@ -0,0 +1,5 @@ +# Additional perspective history + +Revising a connected perspective preserves the earlier evidence and evidence +status in historical views. Repeating an unchanged connection keeps its +original history. Previously overwritten evidence remains unavailable. diff --git a/CHANGELOG.md b/CHANGELOG.md index 7a2724eb8..024294cd9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -268,6 +268,10 @@ All notable changes to this project are documented here. Format follows ### Fixed +- Security dependency floors now require PyJWT 2.15.1 and urllib3 2.8.0, + regenerate the exact `uv.lock`, and test both source declarations and lock + selections against the CVEs reported on LineageWeave#1138. + - Full-corpus Event Lineage rebuilds now count candidate pairs before provider work and omit the optional LLM channel above the 5,000-pair ADR budget, preventing millions of synchronous orchestrator calls while retaining one diff --git a/backend/app/main.py b/backend/app/main.py index 122165990..09380ec4c 100644 --- a/backend/app/main.py +++ b/backend/app/main.py @@ -744,6 +744,7 @@ async def _load_post_voice_types( or ($2::timestamptz is not null and voice.effective_from <= $2 and (voice.effective_to is null or $2 < voice.effective_to))) + and ($2::timestamptz is null or voice.is_primary or voice.recorded_at <= $2) order by voice.is_primary desc, lookup.display_order, voice.voice_type_code """, post_id, diff --git a/backend/app/ontology_neighborhood_ingestion.py b/backend/app/ontology_neighborhood_ingestion.py index 51f56ff2c..5f981d248 100644 --- a/backend/app/ontology_neighborhood_ingestion.py +++ b/backend/app/ontology_neighborhood_ingestion.py @@ -914,6 +914,7 @@ async def _load_voice_assignments( or coalesce($2::timestamptz, $3::timestamptz) < voice.effective_to ) and voice.recorded_at <= $3::timestamptz + and (voice.is_primary or voice.recorded_at <= coalesce($2::timestamptz, $3::timestamptz)) order by voice.post_id, voice.is_primary desc, lookup.display_order, voice.voice_type_code """, diff --git a/backend/app/source_post_voice_ingestion.py b/backend/app/source_post_voice_ingestion.py index 609a71f4f..d975fe7e0 100644 --- a/backend/app/source_post_voice_ingestion.py +++ b/backend/app/source_post_voice_ingestion.py @@ -88,9 +88,17 @@ async def persist_additional_voice_assignment( truth_status_code: str, evidence_post_id: str, ) -> None: - """Atomically bind one additional Voice to an authorized evidence post.""" + """Bind an additional Voice without rewriting earlier cutoff evidence.""" assignment_iri = str(LW[f"voice-assignment/{post_id}/{voice_type_code}"]) async with conn.transaction(): + primary_code = await conn.fetchval( + "select voc_type_code from source_post where post_id = $1::uuid for update", + post_id, + ) + if primary_code == voice_type_code: + raise PrimaryVoiceAssignmentError( + "the imported primary Voice cannot be changed through the additional-voice path" + ) evidence_resource_id = await _post_resource_id(conn, evidence_post_id) assignment_resource_id = await conn.fetchval( """ @@ -139,28 +147,50 @@ async def persist_additional_voice_assignment( ) if assertion_id is None: raise RuntimeError("Voice evidence derivation was not persisted") - stored = await conn.fetchrow( + current = await conn.fetchrow( + """ + select voice_assignment_id, is_primary, truth_status_code, + provenance_assertion_id + from source_post_voice + where post_id = $1::uuid and voice_type_code = $2 + and effective_to is null + """, + post_id, + voice_type_code, + ) + if current is not None: + if current["is_primary"]: + raise PrimaryVoiceAssignmentError( + "the imported primary Voice cannot be changed through the additional-voice path" + ) + if ( + current["truth_status_code"] == truth_status_code + and current["provenance_assertion_id"] == assertion_id + ): + return + change_at = await conn.fetchval("select clock_timestamp()") + if current is not None: + await conn.execute( + """ + update source_post_voice set effective_to = $2 + where voice_assignment_id = $1::uuid and effective_to is null + """, + current["voice_assignment_id"], + change_at, + ) + await conn.execute( """ insert into source_post_voice (post_id, voice_type_code, is_primary, truth_status_code, provenance_assertion_id, effective_from, recorded_at) - values ($1::uuid, $2, false, $3, $4::uuid, now(), now()) - on conflict (post_id, voice_type_code) where effective_to is null do update - set truth_status_code = excluded.truth_status_code, - provenance_assertion_id = excluded.provenance_assertion_id, - recorded_at = now() - where not source_post_voice.is_primary - returning voice_type_code + values ($1::uuid, $2, false, $3, $4::uuid, $5, $5) """, post_id, voice_type_code, truth_status_code, assertion_id, + change_at, ) - if stored is None: - raise PrimaryVoiceAssignmentError( - "the imported primary Voice cannot be changed through the additional-voice path" - ) __all__ = ["PrimaryVoiceAssignmentError", "persist_additional_voice_assignment"] diff --git a/docs/adr/0256-evidence-bearing-voice-combinations.md b/docs/adr/0256-evidence-bearing-voice-combinations.md index 00cb2421b..f5f75d2ea 100644 --- a/docs/adr/0256-evidence-bearing-voice-combinations.md +++ b/docs/adr/0256-evidence-bearing-voice-combinations.md @@ -25,6 +25,23 @@ attributes rather than from one exhaustive industry-role list. Represent composition as rows in normalized `source_post_voice`, not as compound lookup codes. +Additional-assignment reassertion (2026-10-01): changing an additional Voice's +truth state or derivation evidence closes its current half-open interval and +inserts a new assignment interval. The closed row retains its original truth +state, assertion, start, and recording time; historical reads must not acquire +later evidence or lose an earlier assertion. Repeating the same truth state +and derivation is idempotent and retains the existing interval. The write locks +the carrying Post before reading its current Voice, serializing with other +assignments and imported-primary changes. The replacement boundary comes from +the database clock after that lock, not the transaction's possibly earlier +start time. A failed replacement rolls back the interval close and provenance +writes together. Previously overwritten evidence cannot be reconstructed: +an additional row recorded after the requested cutoff is omitted, even if its +old start predates that cutoff. The imported-primary source-time contract +remains governed by ADR 0252. +In-place upsert is rejected because it destroys cutoff evidence; an inferred +repair of old intervals is rejected because the overwritten evidence is absent. + - The existing `source_post.voc_type_code` remains the authoritative imported primary voice. A trigger mirrors it into exactly one primary association so existing import, filtering, and lineage behavior remains stable. @@ -63,8 +80,8 @@ compound lookup codes. - A `post_admin` may add an additional assignment by naming an ABAC-visible evidence Post, an atomic Voice code, and a governed truth state. The API does not accept a caller-supplied assertion identifier: one transaction binds the - evidence Post as a PROV Entity, records `prov:wasDerivedFrom`, and upserts the - assignment. It cannot replace or demote the imported primary Voice. + evidence Post as a PROV Entity, records `prov:wasDerivedFrom`, and records the + effective assignment interval. It cannot replace or demote the imported primary Voice. - In the live Post popup, a `post_admin` may choose one unassigned atomic Voice and one explicit truth state. The open Post is submitted as its own evidence, which covers a single record that contains several perspectives without diff --git a/docs/doctoring/lineageweave-dependency-security-20261001.md b/docs/doctoring/lineageweave-dependency-security-20261001.md new file mode 100644 index 000000000..cad5cb13c --- /dev/null +++ b/docs/doctoring/lineageweave-dependency-security-20261001.md @@ -0,0 +1,33 @@ +# LineageWeave dependency Security RCA — 2026-10-01 + +Status: Proposed on `ContextualWisdomLab/LineageWeave#1137`; protected +integration, exact-current-head Checks, and independent approval remain +mandatory. + +## Exact failure evidence + +`ContextualWisdomLab/LineageWeave#1138@e98a68ed99566f6c145b84c3ec816216dd720ebb` +failed Security Scan run `36811272599`, Trivy job `110206798976`. The exact +SARIF gate reported PyJWT CVE-2026-102265 through CVE-2026-102274 plus +CVE-2026-101917 and CVE-2026-101918 against `uv.lock`'s PyJWT 2.13.0. It also +reported urllib3 CVE-2026-97687, CVE-2026-97688, and CVE-2026-97689 against +urllib3 2.7.0. + +The Voice-derivation code changed by #1138 does not own dependency policy. +`ContextualWisdomLab/LineageWeave#1137` is the existing canonical dependency +owner and already selects PyJWT 2.15.1. The remaining root cause was that its +source floor still admitted earlier releases and urllib3 remained an unbounded +transitive dependency. + +## RED → GREEN repair + +The owner contract first failed three assertions: both PyJWT extras still +declared `>=2.14.0`, no direct urllib3 floor existed, and the lock selected +urllib3 2.7.0. The repair requires PyJWT 2.15.1 on both install surfaces, +declares urllib3 2.8.0 once in core dependencies, and regenerates `uv.lock` +with the repository's `uv` resolver. Only urllib3 moves in the resolved package +set; PyJWT was already resolved to 2.15.1. + +Consumers must ordinary-merge the accepted owner lineage. A terminal Security +gate on the owner and every consumer is required; skipped, queued, pending, or +predecessor results are not acceptance. diff --git a/docs/product-requirements.md b/docs/product-requirements.md index 27c254650..35bb03b3a 100644 --- a/docs/product-requirements.md +++ b/docs/product-requirements.md @@ -58,6 +58,8 @@ edge exposes the same authorized endpoints and evidence through API and UI. authorized Post as evidence and hide the write action on cutoff views. - Validate DB-to-RDF projections with SHACL, including complete reified ProjectMention subject/predicate/object chains. +- Preserve an additional perspective's earlier truth state and evidence when + it is revised; an unchanged retry retains its original availability time. - Keep SKOS broader/narrower distinct from OWL subclass semantics. Acceptance: Turtle, JSON-LD, N-Triples, SHACL, API payloads, persisted IRIs, diff --git a/pyproject.toml b/pyproject.toml index 6c24ecd05..5084138d6 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -17,9 +17,10 @@ dependencies = [ # Explicit CA bundle for http_client HTTPS posts -- some interpreter # distributions don't reliably inherit the OS trust store. "certifi>=2024.0.0", - "cryptography>=42.0", - # Explicit floor for the transport dependency after CVE-2026-97687/97689. + # Explicit security floor for transitive HTTP clients. Keep this aligned + # with the repository advisory contract and generated uv lock. "urllib3>=2.8.0", + "cryptography>=42.0", # The standard Python RDF/OWL library -- parses and validates # docs/ontology/lineageweave-kg.ttl and the standards-complete PROV-O # support profile (ADR 0011). Pure Python, no Rust/C toolchain. diff --git a/tests/test_pyjwt_advisory_floor.py b/tests/test_pyjwt_advisory_floor.py new file mode 100644 index 000000000..2e4269902 --- /dev/null +++ b/tests/test_pyjwt_advisory_floor.py @@ -0,0 +1,73 @@ +"""Keep dependency locks outside the PyJWT and urllib3 advisory ranges.""" + +from __future__ import annotations + +import re +import tomllib +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[1] +PYJWT_PATCHED_VERSION = (2, 15, 1) +URLLIB3_PATCHED_VERSION = (2, 8, 0) + + +def _parse_version_tuple(version: str) -> tuple[int, int, int]: + match = re.fullmatch(r"(\d+)\.(\d+)\.(\d+)", version) + assert match is not None, f"unexpected dependency version syntax: {version!r}" + return tuple(int(part) for part in match.groups()) + + +def test_dev_and_backend_require_the_patched_pyjwt_release() -> None: + project = tomllib.loads((ROOT / "pyproject.toml").read_text()) + extras = project["project"]["optional-dependencies"] + + for extra_name in ("dev", "backend"): + requirements = [ + requirement + for requirement in extras[extra_name] + if requirement.lower().startswith("pyjwt[crypto]") + ] + assert requirements == ["pyjwt[crypto]>=2.15.1"] + + +def test_lockfile_contains_only_patched_pyjwt_releases() -> None: + lock = tomllib.loads((ROOT / "uv.lock").read_text()) + versions = [ + package["version"] + for package in lock["package"] + if package["name"].lower() == "pyjwt" + ] + + assert versions, "uv.lock must contain PyJWT" + assert all( + _parse_version_tuple(version) >= PYJWT_PATCHED_VERSION + for version in versions + ) + + +def test_project_requires_the_patched_urllib3_release() -> None: + """Make the urllib3 advisory floor explicit instead of transitive.""" + project = tomllib.loads((ROOT / "pyproject.toml").read_text()) + requirements = [ + requirement + for requirement in project["project"]["dependencies"] + if requirement.lower().startswith("urllib3") + ] + + assert requirements == ["urllib3>=2.8.0"] + + +def test_lockfile_contains_only_patched_urllib3_releases() -> None: + """Reject urllib3 versions affected by the exact-head Trivy findings.""" + lock = tomllib.loads((ROOT / "uv.lock").read_text()) + versions = [ + package["version"] + for package in lock["package"] + if package["name"].lower() == "urllib3" + ] + + assert versions, "uv.lock must contain urllib3" + assert all( + _parse_version_tuple(version) >= URLLIB3_PATCHED_VERSION + for version in versions + ) diff --git a/tests/test_source_post_voice_history_live.py b/tests/test_source_post_voice_history_live.py index 4ae774d2e..d640e28ad 100644 --- a/tests/test_source_post_voice_history_live.py +++ b/tests/test_source_post_voice_history_live.py @@ -7,11 +7,13 @@ from __future__ import annotations +import asyncio import os import subprocess import threading import uuid from datetime import datetime, timedelta +from itertools import pairwise from pathlib import Path from urllib.parse import urlsplit, urlunsplit @@ -19,6 +21,8 @@ import psycopg2.errors import pytest +from backend.app.source_post_voice_ingestion import persist_additional_voice_assignment + _ADMIN_DSN = os.environ.get( "LINEAGEWEAVE_TEST_POSTGRES_ADMIN_DSN", "postgresql://localhost/postgres" ) @@ -100,6 +104,14 @@ def voice_history_dsn(): database_dsn = _database_dsn(database_name) try: _apply_migrations(database_dsn) + with _connect(database_dsn) as connection, connection.cursor() as cursor: + cursor.execute( + """ + insert into common_lookup_value (lookup_category, lookup_code, lookup_label) + values ('knowledge_graph_node_type', 'node_post', 'Post') + on conflict (lookup_code) do nothing + """ + ) yield database_dsn finally: admin_conn = psycopg2.connect(_ADMIN_DSN) @@ -180,6 +192,235 @@ def _primary_rows(cursor, post_id: str) -> list[tuple]: return cursor.fetchall() +def test_additional_voice_reassertion_preserves_cutoff_evidence( + voice_history_dsn, +) -> None: + """Later truth/evidence writes cannot rewrite an earlier additional Voice.""" + import asyncpg + + from backend.app.main import _load_post_voice_types + from backend.app.ontology_neighborhood_ingestion import _load_voice_assignments + + with _connect(voice_history_dsn) as connection, connection.cursor() as cursor: + post_id = _insert_synthetic_post(cursor) + first_evidence = _insert_synthetic_post(cursor) + later_evidence = _insert_synthetic_post(cursor) + + async def exercise(): + conn = await asyncpg.connect(voice_history_dsn) + try: + + async def assign(truth, evidence): + await persist_additional_voice_assignment( + conn, + post_id=post_id, + voice_type_code="vops", + truth_status_code=truth, + evidence_post_id=evidence, + ) + + await assign("truth_proposed", first_evidence) + first = await conn.fetchrow( + "select * from source_post_voice where post_id = $1::uuid and not is_primary", + post_id, + ) + cutoff = await conn.fetchval("select clock_timestamp()") + await assign("truth_proposed", first_evidence) + assert ( + await conn.fetchrow( + "select * from source_post_voice where voice_assignment_id = $1", + first["voice_assignment_id"], + ) + == first + ) + await assign("truth_observed", first_evidence) + await assign("truth_observed", later_evidence) + rows = await conn.fetch( + """ + select voice.*, binding.node_id as evidence_post_id + from source_post_voice voice + join provenance_assertion assertion + on assertion.assertion_id = voice.provenance_assertion_id + join provenance_resource_binding binding + on binding.resource_id = assertion.object_resource_id + where voice.post_id = $1::uuid and not voice.is_primary + order by voice.effective_from + """, + post_id, + ) + assert len(rows) == 3 + old, intermediate, new = rows + assert old["voice_assignment_id"] != new["voice_assignment_id"] + assert old["truth_status_code"] == "truth_proposed" + assert str(old["evidence_post_id"]) == first_evidence + assert old["recorded_at"] == first["recorded_at"] + assert old["effective_from"] == first["effective_from"] + assert old["effective_to"] == intermediate["effective_from"] + assert intermediate["truth_status_code"] == "truth_observed" + assert str(intermediate["evidence_post_id"]) == first_evidence + assert ( + intermediate["effective_to"] + == new["effective_from"] + == new["recorded_at"] + ) + assert new["truth_status_code"] == "truth_observed" + assert str(new["evidence_post_id"]) == later_evidence + assert new["effective_to"] is None + before_failure = await conn.fetchrow( + "select * from source_post_voice where voice_assignment_id = $1", + new["voice_assignment_id"], + ) + with pytest.raises(asyncpg.CheckViolationError): + await assign("truth_unsupported", first_evidence) + assert ( + await conn.fetchrow( + "select * from source_post_voice where voice_assignment_id = $1", + new["voice_assignment_id"], + ) + == before_failure + ) + assert ( + await conn.fetchval( + "select count(*) from source_post_voice where post_id=$1::uuid and not is_primary", + post_id, + ) + == 3 + ) + historical = await conn.fetchrow( + """ + select voice_assignment_id, truth_status_code, provenance_assertion_id + from source_post_voice + where post_id = $1::uuid and not is_primary + and effective_from <= $2 and recorded_at <= $2 + and (effective_to is null or $2 < effective_to) + """, + post_id, + cutoff, + ) + assert historical["voice_assignment_id"] == first["voice_assignment_id"] + assert historical["truth_status_code"] == first["truth_status_code"] + assert ( + historical["provenance_assertion_id"] + == first["provenance_assertion_id"] + ) + assert ( + await conn.fetchval( + "select voice_type_code from source_post_voice where post_id=$1::uuid and is_primary and effective_to is null", + post_id, + ) + == "voc" + ) + detail = await _load_post_voice_types(conn, post_id, cutoff) + assert [ + (voice["code"], voice["truth_status_code"]) for voice in detail + ] == [ + ("voc", "truth_observed"), + ("vops", "truth_proposed"), + ] + snapshot_at = await conn.fetchval("select clock_timestamp()") + assignments = await _load_voice_assignments( + conn, + [post_id, first_evidence, later_evidence], + knowledge_cutoff=cutoff, + snapshot_at=snapshot_at, + ) + earlier = next(voice for voice in assignments if not voice.is_primary) + assert earlier.evidence_post_id == first_evidence + assert earlier.truth_status_code == "truth_proposed" + # Emulate a legacy overwritten row whose original truth is lost. + await conn.execute( + "update source_post_voice set recorded_at=clock_timestamp() where voice_assignment_id=$1", + first["voice_assignment_id"], + ) + assert [ + voice["code"] + for voice in await _load_post_voice_types(conn, post_id, cutoff) + ] == ["voc"] + assignments = await _load_voice_assignments( + conn, + [post_id, first_evidence, later_evidence], + knowledge_cutoff=cutoff, + snapshot_at=await conn.fetchval("select clock_timestamp()"), + ) + assert all(voice.is_primary for voice in assignments) + finally: + await conn.close() + + asyncio.run(exercise()) + + +def test_waiting_additional_voice_writer_uses_post_lock_clock( + voice_history_dsn, +) -> None: + """An earlier-started transaction cannot backdate a replacement after waiting.""" + import asyncpg + + with _connect(voice_history_dsn) as connection, connection.cursor() as cursor: + post_id = _insert_synthetic_post(cursor) + evidence = _insert_synthetic_post(cursor) + + async def exercise(): + first = await asyncpg.connect(voice_history_dsn) + second = await asyncpg.connect(voice_history_dsn) + observer = await asyncpg.connect(voice_history_dsn) + pending = None + try: + + async def assign(conn, truth): + await persist_additional_voice_assignment( + conn, + post_id=post_id, + voice_type_code="vops", + truth_status_code=truth, + evidence_post_id=evidence, + ) + + await assign(first, "truth_proposed") + second_pid = await second.fetchval("select pg_backend_pid()") + async with first.transaction(): + await first.fetchval( + "select post_id from source_post where post_id=$1::uuid for update", + post_id, + ) + pending = asyncio.create_task(assign(second, "truth_authoritative")) + async with asyncio.timeout(5): + while not await observer.fetchval( + "select wait_event_type = 'Lock' from pg_stat_activity where pid=$1", + second_pid, + ): + await asyncio.sleep(0.001) + await assign(first, "truth_observed") + await asyncio.wait_for(pending, 5) + rows = await first.fetch( + """ + select truth_status_code, effective_from, effective_to, recorded_at + from source_post_voice where post_id=$1::uuid and not is_primary + order by effective_from + """, + post_id, + ) + assert [row["truth_status_code"] for row in rows] == [ + "truth_proposed", + "truth_observed", + "truth_authoritative", + ] + for old, new in pairwise(rows): + assert ( + old["effective_from"] < old["effective_to"] == new["effective_from"] + ) + assert new["recorded_at"] == new["effective_from"] + assert rows[-1]["effective_to"] is None + finally: + if pending is not None and not pending.done(): + pending.cancel() + await asyncio.gather(pending, return_exceptions=True) + await first.close() + await second.close() + await observer.close() + + asyncio.run(exercise()) + + def _api_primary(cursor, post_id: str, cutoff: datetime | None) -> list[str]: cursor.execute( _API_CUTOFF_SQL, @@ -251,7 +492,9 @@ def _insert_additional_voice(cursor, post_id: str, voice_type_code: str) -> None ) -def test_aba_primary_history_matches_api_and_ontology_cutoffs(voice_history_dsn: str) -> None: +def test_aba_primary_history_matches_api_and_ontology_cutoffs( + voice_history_dsn: str, +) -> None: """A → B → A is recoverable at before / between / after cutoffs.""" connection = _connect(voice_history_dsn) try: @@ -294,10 +537,16 @@ def test_aba_primary_history_matches_api_and_ontology_cutoffs(voice_history_dsn: snapshot_during_b = between snapshot_after = after_last - assert _ontology_primary(cursor, post_id, None, snapshot_during_b) == ["vops"] + assert _ontology_primary(cursor, post_id, None, snapshot_during_b) == [ + "vops" + ] assert _ontology_primary(cursor, post_id, None, snapshot_after) == ["voc"] - assert _ontology_primary(cursor, post_id, first_from, snapshot_after) == ["voc"] - assert _ontology_primary(cursor, post_id, between, snapshot_after) == ["vops"] + assert _ontology_primary(cursor, post_id, first_from, snapshot_after) == [ + "voc" + ] + assert _ontology_primary(cursor, post_id, between, snapshot_after) == [ + "vops" + ] cursor.execute( """ @@ -442,7 +691,7 @@ def _update(next_code: str) -> None: (next_code, post_id), ) connection.commit() - except Exception as exc: + except (psycopg2.Error, threading.BrokenBarrierError) as exc: errors.append(exc) connection.rollback() finally: diff --git a/tests/test_source_post_voice_ingestion.py b/tests/test_source_post_voice_ingestion.py index 52fdf0b1f..cdeb7d38b 100644 --- a/tests/test_source_post_voice_ingestion.py +++ b/tests/test_source_post_voice_ingestion.py @@ -4,6 +4,7 @@ import asyncio from contextlib import asynccontextmanager +from datetime import UTC, datetime from typing import Any import pytest @@ -18,19 +19,30 @@ class _Connection: """Record the ordered SQL contract without requiring a live database.""" def __init__( - self, *, primary_conflict: bool = False, existing_evidence: bool = False + self, + *, + existing_evidence: bool = False, + current: dict[str, object] | None = None, ) -> None: - self.primary_conflict = primary_conflict + self.current = current self.calls: list[tuple[str, tuple[object, ...]]] = [] self.fetchvals = iter( - ["evidence-resource", "assignment-resource", "assertion"] + [ + "voc", + "evidence-resource", + "assignment-resource", + "assertion", + datetime(2026, 10, 1, tzinfo=UTC), + ] if existing_evidence else [ + "voc", None, "evidence-resource", "evidence-resource", "assignment-resource", "assertion", + datetime(2026, 10, 1, tzinfo=UTC), ] ) @@ -48,10 +60,10 @@ async def fetchval(self, query: str, *args: object) -> Any: self.calls.append((query, args)) return next(self.fetchvals) - async def fetchrow(self, query: str, *args: object) -> dict[str, str] | None: - """Return no row only when the imported primary blocks the write.""" + async def fetchrow(self, query: str, *args: object) -> dict[str, object] | None: + """Return the existing additional interval, if supplied.""" self.calls.append((query, args)) - return None if self.primary_conflict else {"voice_type_code": str(args[1])} + return self.current def test_additional_voice_creates_prov_derivation_and_assignment_atomically() -> None: @@ -70,9 +82,9 @@ def test_additional_voice_creates_prov_derivation_and_assignment_atomically() -> sql = "\n".join(query for query, _args in conn.calls) assert "prov_was_derived_from" in sql - assert "where effective_to is null" in sql - assert "where not source_post_voice.is_primary" in sql - assert "where effective_to is null" in sql + assert "and effective_to is null" in sql + assert "for update" in sql + assert "select clock_timestamp()" in sql assert "voice-assignment/aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaa1/vops" in str( conn.calls ) @@ -80,7 +92,7 @@ def test_additional_voice_creates_prov_derivation_and_assignment_atomically() -> def test_additional_voice_cannot_demote_imported_primary() -> None: """The current primary remains owned by source_post.voc_type_code.""" - conn = _Connection(primary_conflict=True) + conn = _Connection() with pytest.raises(PrimaryVoiceAssignmentError): asyncio.run( @@ -112,3 +124,64 @@ def test_existing_evidence_binding_is_typed_as_a_prov_entity() -> None: "provenance_resource_type" in query and args == ("evidence-resource",) for query, args in conn.calls ) + + +def test_repeating_the_same_truth_and_evidence_retains_the_interval() -> None: + """A retry cannot move the original availability clock or create history.""" + conn = _Connection( + current={ + "voice_assignment_id": "current-assignment", + "is_primary": False, + "truth_status_code": "truth_observed", + "provenance_assertion_id": "assertion", + } + ) + asyncio.run( + persist_additional_voice_assignment( + conn, + post_id="aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaa1", + voice_type_code="vops", + truth_status_code="truth_observed", + evidence_post_id="aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaa2", + ) + ) + assert not any( + "update source_post_voice" in query or "insert into source_post_voice" in query + for query, _args in conn.calls + ) + + +def test_replacement_closes_only_the_previous_interval() -> None: + """A new truth state keeps the old assertion and its recording time intact.""" + conn = _Connection( + current={ + "voice_assignment_id": "old-assignment", + "is_primary": False, + "truth_status_code": "truth_proposed", + "provenance_assertion_id": "old-assertion", + } + ) + asyncio.run( + persist_additional_voice_assignment( + conn, + post_id="aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaa1", + voice_type_code="vops", + truth_status_code="truth_observed", + evidence_post_id="aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaa2", + ) + ) + close = next( + (query, args) + for query, args in conn.calls + if "update source_post_voice" in query + ) + insert = next( + (query, args) + for query, args in conn.calls + if "insert into source_post_voice" in query + ) + assert close[1][0] == "old-assignment" + assert close[1][1] == insert[1][4] + assert "recorded_at =" not in close[0] + assert "truth_status_code =" not in close[0] + assert "provenance_assertion_id =" not in close[0]