From 2e78b9ae509308333c1ed95fa28a7cbb44588ad4 Mon Sep 17 00:00:00 2001 From: Albrecht Degering Date: Sun, 2 Aug 2026 06:20:24 +0200 Subject: [PATCH] Add governed contact point snapshots --- README.md | 32 +- docs/ADDRESS_MODULE_ARCHITECTURE.md | 24 +- docs/IMPLEMENTATION_PLAN.md | 7 + .../backend/capabilities.py | 966 +++++++++++++++++- src/govoplan_addresses/backend/db/models.py | 31 + src/govoplan_addresses/backend/manifest.py | 33 +- .../a3b5c6d7e8f9_contact_point_snapshots.py | 63 ++ src/govoplan_addresses/backend/router.py | 170 ++- src/govoplan_addresses/backend/schemas.py | 132 +++ tests/test_addresses_service.py | 178 ++++ 10 files changed, 1620 insertions(+), 16 deletions(-) create mode 100644 src/govoplan_addresses/backend/migrations/versions/a3b5c6d7e8f9_contact_point_snapshots.py diff --git a/README.md b/README.md index ea3a99e..6e1eccd 100644 --- a/README.md +++ b/README.md @@ -74,13 +74,16 @@ It must not own: ## First Capabilities -The module exposes three core-mediated capabilities: +The module exposes four core-mediated capabilities: - `addresses.lookup`: read-only contact/recipient lookup for autocomplete. - `addresses.recipient_source`: immutable recipient snapshots for campaign, reporting, mail-build, forms, portal, and postbox workflows. - `addresses.contact_writer`: address-book-scoped write decisions and contact creation for local or otherwise writable sources. +- `addresses.contact_point_resolution`: purpose-aware, channel-neutral + resolution and immutable snapshots for email, postal, internal-mail, and + portal targets. `addresses.recipient_source` returns: @@ -98,10 +101,13 @@ address. Email-oriented consumers snapshot email targets and whole-contact entries with a usable email address; postal-only entries remain valid list members for later postal/document workflows. -Consumers must store their own immutable snapshot when they need historical -evidence. The addresses module remains the owner of the reusable source, not of -the consumer's historical records. Consumers must resolve these capabilities -through the platform registry and must not import address ORM/service internals. +Legacy `addresses.recipient_source` consumers must store their own immutable +snapshot when they need historical evidence. Channel-neutral consumers may use +the dedicated freeze operation described below. The addresses module remains +the owner of reusable sources; domain consumers remain responsible for linking +their own records to snapshot evidence. Consumers must resolve these +capabilities through the platform registry and must not import address +ORM/service internals. `addresses.contact_writer` returns an explicit decision before a consumer shows or executes write actions: allowed/blocked, reason, user-facing message, @@ -109,6 +115,22 @@ required scopes, source kind, read-only state, and provenance. The decision is address-book specific; broader policy modules may later contribute to the same decision path, but consumers should not import or duplicate policy logic. +For channel-neutral consumers, `addresses.contact_point_resolution` supersedes +the email-only shape without removing it. It accepts local contact IDs and +stable provider references such as `idm:identity:`, applies an effective +date, communication purpose, address purpose, fallback rule, locale, and +domestic/international postal formatting, and returns candidates plus excluded +targets with stable reason codes. Bounded previews are live. A frozen snapshot +stores the complete values, source and governance revisions, provenance, and a +deterministic hash in Addresses so later contact edits cannot rewrite evidence. + +The corresponding HTTP API is available below `/api/v1/addresses`: + +- `POST /contact-points/resolve` +- `POST /contact-point-sources/preview` +- `POST /contact-point-snapshots` +- `GET /contact-point-snapshots/{snapshot_id}` + ## Design Documents - [Address module architecture](docs/ADDRESS_MODULE_ARCHITECTURE.md) diff --git a/docs/ADDRESS_MODULE_ARCHITECTURE.md b/docs/ADDRESS_MODULE_ARCHITECTURE.md index 4ae104c..8195bbe 100644 --- a/docs/ADDRESS_MODULE_ARCHITECTURE.md +++ b/docs/ADDRESS_MODULE_ARCHITECTURE.md @@ -87,6 +87,8 @@ The first stable capabilities are: campaign, scheduling, postbox, portal, and case workflows. - `addresses.contact_writer`: provide address-book-scoped write target decisions and contact creation for local or otherwise writable sources. +- `addresses.contact_point_resolution` version 1.x: resolve channel-neutral + contact points and freeze immutable recipient evidence. Capabilities use DTOs and source IDs. Consumers must not receive ORM objects or write address tables directly. Consumers that need historical evidence must @@ -96,10 +98,24 @@ provenance; they must not treat live address records as historical evidence. `addresses.recipient_source` exposes both complete address books and classical address lists. Address-book sources use `addresses:address_book:`. Address-list sources use `addresses:address_list:` and include the -address-list entry ID in each recipient's provenance. The current snapshot DTO -is email-recipient oriented; postal-only list entries are valid address-list -members but are skipped by the email recipient-source path until postal -recipient DTOs are added. +address-list entry ID in each recipient's provenance. The legacy snapshot DTO +remains email-oriented for compatible campaign consumers. + +Channel-neutral consumers use `addresses.contact_point_resolution`, which +supports email, postal, internal-mail, and portal targets, including postal-only +address-list entries. Requests make effective date, communication purpose, +address purpose, fallback behavior, locale, and domestic/international postal +formatting explicit. Results retain stable subject/contact/contact-point IDs, +source, preference and consent revisions, provenance, and reasons for excluded +or unresolved candidates. + +Live previews are bounded to 500 rows per page and 20,000 source members per +request. Frozen snapshots persist resolved values and exclusions with a +deterministic hash; reading a snapshot never resolves the live contact again. +Mixed-audience expansion and final cross-provider Policy/channel decisions +remain owned by Distribution Lists and Policy. The contract is defined in Core, +and Addresses does not import IDM, Organizations, or Distribution Lists +implementations. The writer capability is intentionally address-book specific. It answers whether the current principal may perform an operation such as `create_contact`, diff --git a/docs/IMPLEMENTATION_PLAN.md b/docs/IMPLEMENTATION_PLAN.md index cfab057..7f3550a 100644 --- a/docs/IMPLEMENTATION_PLAN.md +++ b/docs/IMPLEMENTATION_PLAN.md @@ -73,6 +73,11 @@ Tasks: - [x] define immutable recipient snapshot DTOs - [x] expose source provenance in capability responses - [x] expose classical address lists as `addresses.recipient_source` sources +- [x] expose versioned channel-neutral contact-point resolution for local and + stable provider subject references +- [x] support purpose/address-purpose selection, deterministic fallback, + locale, and domestic/international postal rendering +- [x] add bounded source previews and immutable postal/email snapshots - [x] add module presence/capability tests - [x] document consumer rules for campaign, mail, scheduling, portal, postbox, and reporting @@ -82,6 +87,8 @@ Exit criteria: - [x] campaign can request a recipient source via core-mediated capability - [x] mail/scheduling can request autocomplete candidates via core-mediated lookup - [x] consumers do not import `govoplan_addresses` +- [x] postal-only contacts/list entries can be resolved without changing the + legacy email recipient-source contract ## Milestone 4: Campaign Integration diff --git a/src/govoplan_addresses/backend/capabilities.py b/src/govoplan_addresses/backend/capabilities.py index 6f2d4bb..f06bffd 100644 --- a/src/govoplan_addresses/backend/capabilities.py +++ b/src/govoplan_addresses/backend/capabilities.py @@ -1,12 +1,13 @@ from __future__ import annotations -from dataclasses import dataclass, field +from dataclasses import asdict, dataclass, field, replace from datetime import UTC, datetime import hashlib import json from typing import Any -from sqlalchemy import func +from sqlalchemy import func, or_ +from sqlalchemy.orm import selectinload from govoplan_core.auth import ApiPrincipal from govoplan_core.core.distribution_lists import ( @@ -16,6 +17,16 @@ from govoplan_core.core.distribution_lists import ( RecipientChannelFacts, RecipientChannelFactsRequest, ) +from govoplan_core.core.contact_points import ( + CONTACT_POINT_CONTRACT_VERSION, + ContactPointCandidate, + ContactPointResolution, + ContactPointResolutionProvider, + ContactPointResolutionRequest, + ContactPointSnapshotRef, + ContactPointSourcePreview, + ContactPointSourceRequest, +) from govoplan_core.db.base import utcnow from govoplan_core.core.people import ( PeopleSearchGroup, @@ -30,6 +41,7 @@ from govoplan_addresses.backend.db.models import ( ContactChannelRule, ContactEmail, ContactPhone, + ContactPointSnapshot, ContactPostalAddress, ) from govoplan_addresses.backend.schemas import ContactCreateRequest @@ -574,6 +586,198 @@ class AddressesChannelFactsCapability: ) +@dataclass(frozen=True, slots=True) +class _SourceContactContext: + contact: Contact + allowed_contact_point_ids: frozenset[str] | None = None + address_list_entry_ids: tuple[str, ...] = () + + +class AddressesContactPointResolutionCapability(ContactPointResolutionProvider): + """Resolve and freeze channel-neutral contact points without policy imports.""" + + def resolve_contact_points( + self, + session: Any, + principal: ApiPrincipal, + *, + request: ContactPointResolutionRequest, + ) -> ContactPointResolution: + self._assert_tenant(principal, request.tenant_id) + contacts = _contacts_for_subject(session, principal, request.subject) + if not contacts: + return _unresolved_contact_point_resolution( + request, + code="addresses.subject.unresolved", + message="No visible Addresses contact is linked to this subject reference.", + ) + if len(contacts) > 1: + return _unresolved_contact_point_resolution( + request, + code="addresses.subject.ambiguous", + message="More than one visible Addresses contact is linked to this subject reference.", + status="ambiguous", + provenance={"contact_ids": [item.id for item in contacts]}, + ) + return _resolve_contact_points_for_contact( + session, + principal, + request=request, + context=_SourceContactContext(contact=contacts[0]), + ) + + def preview_source( + self, + session: Any, + principal: ApiPrincipal, + *, + request: ContactPointSourceRequest, + offset: int = 0, + limit: int = 100, + ) -> ContactPointSourcePreview: + self._assert_tenant(principal, request.tenant_id) + bounded_offset = max(0, offset) + bounded_limit = max(1, min(limit, 500)) + source, source_revision, contexts = _contact_point_source_contexts( + session, + principal, + request.source_id, + ) + if len(contexts) > max(1, min(request.max_items, 20_000)): + raise ValueError( + f"Contact-point source exceeds the configured {request.max_items} item limit." + ) + selected = contexts[bounded_offset : bounded_offset + bounded_limit] + resolutions = tuple( + _resolve_contact_points_for_contact( + session, + principal, + request=_contact_point_resolution_request(request, item.contact), + context=item, + ) + for item in selected + ) + source_fingerprint = _source_fingerprint(request.source_id, source_revision) + return ContactPointSourcePreview( + contract_version=CONTACT_POINT_CONTRACT_VERSION, + source=source, + request=request, + resolutions=resolutions, + total_count=len(contexts), + usable_count=sum(len(item.candidates) for item in resolutions), + excluded_count=sum( + len(item.excluded) + int(not item.candidates and not item.excluded) + for item in resolutions + ), + offset=bounded_offset, + limit=bounded_limit, + has_more=bounded_offset + len(resolutions) < len(contexts), + source_revision=source_revision, + source_fingerprint=source_fingerprint, + generated_at=utcnow(), + provenance={ + "module": "addresses", + "contract_version": CONTACT_POINT_CONTRACT_VERSION, + "bounded": True, + }, + ) + + def freeze_source( + self, + session: Any, + principal: ApiPrincipal, + *, + request: ContactPointSourceRequest, + ) -> ContactPointSnapshotRef: + self._assert_tenant(principal, request.tenant_id) + source, source_revision, contexts = _contact_point_source_contexts( + session, + principal, + request.source_id, + ) + max_items = max(1, min(request.max_items, 20_000)) + if len(contexts) > max_items: + raise ValueError( + f"Contact-point source exceeds the configured {request.max_items} item limit." + ) + resolutions = tuple( + _resolve_contact_points_for_contact( + session, + principal, + request=_contact_point_resolution_request(request, item.contact), + context=item, + ) + for item in contexts + ) + generated_at = utcnow() + source_fingerprint = _source_fingerprint(request.source_id, source_revision) + request_payload = _json_value(asdict(request)) + resolution_payload = [_json_value(asdict(item)) for item in resolutions] + snapshot_hash = hashlib.sha256( + json.dumps( + { + "contract_version": CONTACT_POINT_CONTRACT_VERSION, + "source": _json_value(asdict(source)), + "request": request_payload, + "resolutions": resolution_payload, + "source_revision": source_revision, + "source_fingerprint": source_fingerprint, + }, + sort_keys=True, + separators=(",", ":"), + ).encode("utf-8") + ).hexdigest() + snapshot = ContactPointSnapshot( + tenant_id=request.tenant_id, + source_id=request.source_id, + contract_version=CONTACT_POINT_CONTRACT_VERSION, + source_revision=source_revision, + source_fingerprint=source_fingerprint, + purpose=request.purpose, + effective_at=request.effective_at, + generated_at=generated_at, + request_payload=request_payload, + resolution_payload=resolution_payload, + recipient_count=sum(len(item.candidates) for item in resolutions), + excluded_count=sum( + len(item.excluded) + int(not item.candidates and not item.excluded) + for item in resolutions + ), + snapshot_hash=snapshot_hash, + created_by_account_id=principal.account_id, + provenance={ + "module": "addresses", + "source": _json_value(asdict(source)), + "immutable": True, + }, + ) + session.add(snapshot) + session.flush() + return _contact_point_snapshot_ref(snapshot) + + def get_snapshot( + self, + session: Any, + principal: ApiPrincipal, + *, + snapshot_id: str, + ) -> ContactPointSnapshotRef | None: + snapshot = ( + session.query(ContactPointSnapshot) + .filter( + ContactPointSnapshot.id == snapshot_id, + ContactPointSnapshot.tenant_id == principal.tenant_id, + ) + .one_or_none() + ) + return _contact_point_snapshot_ref(snapshot) if snapshot is not None else None + + @staticmethod + def _assert_tenant(principal: ApiPrincipal, tenant_id: str) -> None: + if tenant_id != principal.tenant_id: + raise PermissionError("Contact points cannot be resolved across tenants.") + + def lookup_capability(_context: Any) -> AddressesLookupCapability: return AddressesLookupCapability() @@ -590,6 +794,12 @@ def channel_facts_capability(_context: Any) -> AddressesChannelFactsCapability: return AddressesChannelFactsCapability() +def contact_point_resolution_capability( + _context: Any, +) -> AddressesContactPointResolutionCapability: + return AddressesContactPointResolutionCapability() + + def contact_writer_capability(_context: Any) -> AddressesContactWriterCapability: return AddressesContactWriterCapability() @@ -912,6 +1122,12 @@ def _address_list_source_ref(address_list: AddressList, *, recipient_count: int, def _active_address_book_contacts(session: Any, address_book_id: str) -> list[Contact]: return ( session.query(Contact) + .options( + selectinload(Contact.emails), + selectinload(Contact.phones), + selectinload(Contact.postal_addresses), + selectinload(Contact.channel_rules), + ) .filter(Contact.address_book_id == address_book_id, Contact.deleted_at.is_(None)) .order_by(Contact.display_name.asc(), Contact.id.asc()) .all() @@ -989,8 +1205,12 @@ def _address_list_updated_at(session: Any, address_list_ids: list[str]) -> dict[ ) email_rows = ( session.query(AddressListEntry.address_list_id, func.max(ContactEmail.updated_at)) - .join(ContactEmail, AddressListEntry.contact_email_id == ContactEmail.id) - .filter(AddressListEntry.address_list_id.in_(address_list_ids)) + .join(Contact, AddressListEntry.contact_id == Contact.id) + .join(ContactEmail, ContactEmail.contact_id == Contact.id) + .filter( + AddressListEntry.address_list_id.in_(address_list_ids), + Contact.deleted_at.is_(None), + ) .group_by(AddressListEntry.address_list_id) .all() ) @@ -1001,7 +1221,27 @@ def _address_list_updated_at(session: Any, address_list_ids: list[str]) -> dict[ .group_by(AddressListEntry.address_list_id) .all() ) - for address_list_id, updated_at in [*entry_rows, *contact_rows, *email_rows, *rule_rows]: + postal_rows = ( + session.query( + AddressListEntry.address_list_id, + func.max(ContactPostalAddress.updated_at), + ) + .join(Contact, AddressListEntry.contact_id == Contact.id) + .join(ContactPostalAddress, ContactPostalAddress.contact_id == Contact.id) + .filter( + AddressListEntry.address_list_id.in_(address_list_ids), + Contact.deleted_at.is_(None), + ) + .group_by(AddressListEntry.address_list_id) + .all() + ) + for address_list_id, updated_at in [ + *entry_rows, + *contact_rows, + *email_rows, + *postal_rows, + *rule_rows, + ]: if updated_at is None: continue key = str(address_list_id) @@ -1168,3 +1408,719 @@ def _aware_datetime(value: datetime) -> datetime: if value.tzinfo is None: return value.replace(tzinfo=UTC) return value + + +def _contacts_for_subject( + session: Any, + principal: ApiPrincipal, + subject: DistributionSourceReference, +) -> list[Contact]: + direct_id = subject.metadata.get("contact_id") + if direct_id is None and subject.provider == "addresses" and subject.resource_type in { + "contact", + "address_contact", + "address_email", + }: + direct_id = subject.resource_id + if direct_id is not None: + try: + return [get_visible_contact(session, principal, str(direct_id))] + except AddressBookError: + return [] + + explicit_ref = subject.metadata.get("source_ref") + source_refs = { + str(explicit_ref).strip() if explicit_ref else "", + f"{subject.provider}:{subject.resource_type}:{subject.resource_id}", + } + source_refs.discard("") + rows = ( + session.query(Contact) + .options( + selectinload(Contact.emails), + selectinload(Contact.phones), + selectinload(Contact.postal_addresses), + selectinload(Contact.channel_rules), + ) + .filter( + Contact.source_ref.in_(source_refs), + or_(Contact.tenant_id == principal.tenant_id, Contact.tenant_id.is_(None)), + Contact.deleted_at.is_(None), + ) + .order_by(Contact.id.asc()) + .limit(3) + .all() + ) + visible: list[Contact] = [] + for row in rows: + try: + visible.append(get_visible_contact(session, principal, row.id)) + except AddressBookError: + continue + return visible + + +def _unresolved_contact_point_resolution( + request: ContactPointResolutionRequest, + *, + code: str, + message: str, + status: str = "unresolved", + provenance: dict[str, Any] | None = None, +) -> ContactPointResolution: + return ContactPointResolution( + contract_version=CONTACT_POINT_CONTRACT_VERSION, + subject=request.subject, + status=status, # type: ignore[arg-type] + explanations=( + DistributionExplanation( + code=code, + message=message, + severity="warning", + provider="addresses", + source=request.subject, + ), + ), + provenance={ + "module": "addresses", + "effective_at": request.effective_at.isoformat(), + "purpose": request.purpose, + **(provenance or {}), + }, + ) + + +def _resolve_contact_points_for_contact( + session: Any, + principal: ApiPrincipal, + *, + request: ContactPointResolutionRequest, + context: _SourceContactContext, +) -> ContactPointResolution: + contact = context.contact + facts = AddressesChannelFactsCapability().resolve_channel_facts( + session, + principal, + request=RecipientChannelFactsRequest( + tenant_id=request.tenant_id, + source=DistributionSourceReference( + provider="addresses", + resource_type="contact", + resource_id=contact.id, + revision=contact.source_revision, + label=contact.display_name, + metadata={ + "contact_id": contact.id, + "source_ref": contact.source_ref, + }, + ), + recipient_key=f"contact:{contact.id}", + effective_at=request.effective_at, + purpose=request.purpose, + requested_channels=request.requested_channels, + context=request.context, + ), + ) + preference_revision = _rule_revision( + contact, + decisions={"preferred"}, + ) + consent_revision = _rule_revision( + contact, + decisions={ + "allowed", + "opted_in", + "opted_out", + "suppressed", + "invalid", + "returned", + "temporarily_unavailable", + }, + ) + candidates: list[ContactPointCandidate] = [] + for fact in facts.candidates: + if ( + context.allowed_contact_point_ids is not None + and ( + fact.contact_point_id is None + or fact.contact_point_id not in context.allowed_contact_point_ids + ) + ): + continue + details = _contact_point_details( + contact, + channel=fact.channel, + contact_point_id=fact.contact_point_id, + fallback_target=fact.target, + fallback_target_key=fact.target_key, + postal_format=request.postal_format, + ) + if details is None: + continue + address_purpose, primary, target, target_key, value, order_index = details + preference_rank = fact.decision_provenance.get("preference_rank") + candidates.append( + ContactPointCandidate( + channel=fact.channel, + target=target, + target_key=target_key, + status=fact.status, + contact_point_id=fact.contact_point_id, + address_purpose=address_purpose, + locale=fact.locale or request.locale, + preferred=fact.preferred or primary, + preference_rank=( + int(preference_rank) + if isinstance(preference_rank, int) + else None + ), + reason_code=fact.reason_code, + explanation=fact.explanation, + source=fact.source, + source_revision=facts.source_revision, + preference_revision=preference_revision, + consent_revision=consent_revision, + value=value, + provenance={ + **dict(fact.decision_provenance), + "contact_id": contact.id, + "address_book_id": contact.address_book_id, + "source_kind": contact.source_kind, + "source_ref": contact.source_ref, + "source_revision": contact.source_revision, + "order_index": order_index, + "address_list_entry_ids": context.address_list_entry_ids, + }, + ) + ) + + selected, excluded = _select_contact_point_candidates(candidates, request) + if selected: + result_status = "usable" + elif excluded: + result_status = _resolution_status(excluded) + else: + result_status = "unresolved" + explanations = list(facts.explanations) + if not selected and not excluded: + explanations.append( + DistributionExplanation( + code="addresses.contact_point.none", + message="The contact has no contact point for the requested channels.", + severity="warning", + provider="addresses", + source=request.subject, + ) + ) + return ContactPointResolution( + contract_version=CONTACT_POINT_CONTRACT_VERSION, + subject=request.subject, + status=result_status, # type: ignore[arg-type] + contact_id=contact.id, + display_name=contact.display_name, + candidates=tuple(selected), + excluded=tuple(excluded), + explanations=tuple(explanations), + source_revision=facts.source_revision, + source_fingerprint=facts.source_fingerprint, + provenance={ + **dict(facts.provenance), + "provider_subject": _json_value(asdict(request.subject)), + "address_purpose": request.address_purpose, + "fallback_rule": request.fallback_rule, + "postal_format": request.postal_format, + "preference_revision": preference_revision, + "consent_revision": consent_revision, + "address_list_entry_ids": context.address_list_entry_ids, + }, + ) + + +def _contact_point_details( + contact: Contact, + *, + channel: str, + contact_point_id: str | None, + fallback_target: str, + fallback_target_key: str, + postal_format: str, +) -> tuple[str | None, bool, str, str, dict[str, Any], int] | None: + if channel == "email": + point = next((item for item in contact.emails if item.id == contact_point_id), None) + if point is None: + return None + target = point.email.strip() + return ( + point.label, + point.is_primary, + target, + f"email:{target.casefold()}", + {"email": target, "label": point.label}, + point.order_index, + ) + if channel == "postal": + point = next( + (item for item in contact.postal_addresses if item.id == contact_point_id), + None, + ) + if point is None: + return None + lines = [ + contact.display_name, + point.street, + " ".join(part for part in (point.postal_code, point.locality) if part), + point.region, + ] + if postal_format == "international": + lines.append(point.country) + normalized_lines = [str(item).strip() for item in lines if item and str(item).strip()] + canonical = "|".join( + str(item or "").strip().casefold() + for item in ( + point.street, + point.postal_code, + point.locality, + point.region, + point.country, + ) + ) + return ( + point.label, + point.is_primary, + "\n".join(normalized_lines), + f"postal:{canonical}", + { + "label": point.label, + "street": point.street, + "postal_code": point.postal_code, + "locality": point.locality, + "region": point.region, + "country": point.country, + "formatted_lines": normalized_lines, + "format": postal_format, + }, + point.order_index, + ) + if channel in {"internal_mail", "portal"}: + purpose = (contact.provenance or {}).get(f"{channel}_purpose") + return ( + str(purpose) if purpose else None, + True, + fallback_target, + fallback_target_key, + {"reference": fallback_target}, + 0, + ) + return None + + +def _select_contact_point_candidates( + candidates: list[ContactPointCandidate], + request: ContactPointResolutionRequest, +) -> tuple[list[ContactPointCandidate], list[ContactPointCandidate]]: + candidates = sorted(candidates, key=_contact_point_sort_key) + accepted: list[ContactPointCandidate] = [] + excluded: list[ContactPointCandidate] = [] + requested_purpose = (request.address_purpose or "").strip().casefold() + by_channel: dict[str, list[ContactPointCandidate]] = {} + for item in candidates: + by_channel.setdefault(item.channel, []).append(item) + for channel_items in by_channel.values(): + selected_ids: set[int] + if not requested_purpose: + selected_ids = {id(item) for item in channel_items} + else: + exact = [ + item + for item in channel_items + if (item.address_purpose or "").strip().casefold() == requested_purpose + ] + if exact: + selected_ids = {id(item) for item in exact} + elif request.fallback_rule == "primary": + primary = [item for item in channel_items if item.preferred] + selected_ids = {id(item) for item in primary[:1]} + elif request.fallback_rule == "any": + selected_ids = {id(channel_items[0])} if channel_items else set() + else: + selected_ids = set() + for item in channel_items: + if id(item) in selected_ids: + accepted.append(item) + else: + excluded.append( + item + if item.status not in {"usable", "stale"} + else replace( + item, + status="unresolved", + reason_code="addresses.address_purpose.not_selected", + explanation=( + "This contact point did not match the requested address purpose " + "or its configured fallback." + ), + ) + ) + deduplicated: list[ContactPointCandidate] = [] + seen: set[tuple[str, str]] = set() + for item in accepted: + key = (item.channel, item.target_key) + if key in seen and item.status in {"usable", "stale"}: + excluded.append( + replace( + item, + status="duplicate", + reason_code="addresses.contact_point.duplicate", + explanation="An earlier contact point resolves to the same channel target.", + ) + ) + continue + seen.add(key) + if item.status in {"usable", "stale"}: + deduplicated.append(item) + else: + excluded.append(item) + return deduplicated, sorted(excluded, key=_contact_point_sort_key) + + +def _contact_point_sort_key(item: ContactPointCandidate) -> tuple[Any, ...]: + order = item.provenance.get("order_index") + return ( + {"email": 0, "postal": 1, "internal_mail": 2, "portal": 3}.get( + item.channel, + 9, + ), + 0 if item.preferred else 1, + item.preference_rank if item.preference_rank is not None else 10001, + int(order) if isinstance(order, int) else 0, + item.contact_point_id or "", + item.target_key, + ) + + +def _resolution_status(excluded: list[ContactPointCandidate]) -> str: + for status in ( + "invalid", + "suppressed", + "ambiguous", + "duplicate", + "stale", + "unresolved", + ): + if any(item.status == status for item in excluded): + return status + return "unresolved" + + +def _rule_revision(contact: Contact, *, decisions: set[str]) -> str | None: + stamps = [ + _aware_datetime(item.updated_at) + for item in contact.channel_rules + if item.decision in decisions + ] + return max(stamps).isoformat() if stamps else None + + +def _contact_point_source_contexts( + session: Any, + principal: ApiPrincipal, + source_id: str, +) -> tuple[DistributionSourceReference, str, list[_SourceContactContext]]: + book_prefix = "addresses:address_book:" + list_prefix = "addresses:address_list:" + if source_id.startswith(book_prefix): + book = get_visible_address_book( + session, + principal, + source_id.removeprefix(book_prefix), + ) + contexts = [ + _SourceContactContext(contact=item) + for item in _active_address_book_contacts(session, book.id) + ] + revision = _recipient_source_revision( + book, + _address_book_contact_updated_at(session, [book.id]).get(book.id), + ) + return ( + DistributionSourceReference( + provider="addresses", + resource_type="address_book", + resource_id=book.id, + revision=revision, + label=book.name, + metadata={ + "source_id": source_id, + "scope_type": book.scope_type, + "scope_id": book.scope_id, + "source_kind": book.source_kind, + "source_ref": book.source_ref, + }, + ), + revision, + contexts, + ) + if source_id.startswith(list_prefix): + address_list = get_visible_address_list( + session, + principal, + source_id.removeprefix(list_prefix), + ) + entries = list_address_list_entries(session, principal, address_list.id) + grouped: dict[str, dict[str, Any]] = {} + for entry in entries: + item = grouped.setdefault( + entry.contact_id, + {"contact": entry.contact, "allowed": set(), "all": False, "entry_ids": []}, + ) + item["entry_ids"].append(entry.id) + point_id = entry.contact_email_id or entry.contact_postal_address_id + if point_id: + item["allowed"].add(point_id) + else: + item["all"] = True + contexts = [ + _SourceContactContext( + contact=item["contact"], + allowed_contact_point_ids=( + None if item["all"] else frozenset(item["allowed"]) + ), + address_list_entry_ids=tuple(item["entry_ids"]), + ) + for item in grouped.values() + ] + revision = _address_list_source_revision( + address_list, + _address_list_updated_at(session, [address_list.id]).get(address_list.id), + ) + return ( + DistributionSourceReference( + provider="addresses", + resource_type="address_list", + resource_id=address_list.id, + revision=revision, + label=address_list.name, + metadata={ + "source_id": source_id, + "address_book_id": address_list.address_book_id, + "scope_type": address_list.address_book.scope_type, + "scope_id": address_list.address_book.scope_id, + }, + ), + revision, + contexts, + ) + raise ValueError(f"Unsupported Addresses contact-point source id: {source_id}") + + +def _contact_point_resolution_request( + request: ContactPointSourceRequest, + contact: Contact, +) -> ContactPointResolutionRequest: + return ContactPointResolutionRequest( + tenant_id=request.tenant_id, + subject=DistributionSourceReference( + provider="addresses", + resource_type="contact", + resource_id=contact.id, + revision=contact.source_revision, + label=contact.display_name, + metadata={ + "contact_id": contact.id, + "source_kind": contact.source_kind, + "source_ref": contact.source_ref, + }, + ), + effective_at=request.effective_at, + purpose=request.purpose, + requested_channels=request.requested_channels, + address_purpose=request.address_purpose, + fallback_rule=request.fallback_rule, + locale=request.locale, + postal_format=request.postal_format, + context=request.context, + ) + + +def _source_fingerprint(source_id: str, revision: str) -> str: + return hashlib.sha256(f"{source_id}\0{revision}".encode("utf-8")).hexdigest() + + +def _json_value(value: Any) -> Any: + if isinstance(value, datetime): + return value.isoformat() + if isinstance(value, dict): + return {str(key): _json_value(item) for key, item in value.items()} + if isinstance(value, (list, tuple, set, frozenset)): + return [_json_value(item) for item in value] + return value + + +def _source_reference_from_payload(payload: dict[str, Any]) -> DistributionSourceReference: + return DistributionSourceReference( + provider=str(payload.get("provider") or "addresses"), + resource_type=str(payload.get("resource_type") or "unknown"), + resource_id=str(payload.get("resource_id") or ""), + revision=str(payload["revision"]) if payload.get("revision") else None, + fingerprint=str(payload["fingerprint"]) if payload.get("fingerprint") else None, + label=str(payload["label"]) if payload.get("label") else None, + metadata=dict(payload.get("metadata") or {}), + ) + + +def _candidate_from_payload(payload: dict[str, Any]) -> ContactPointCandidate: + source_payload = payload.get("source") + return ContactPointCandidate( + channel=str(payload.get("channel") or "email"), # type: ignore[arg-type] + target=str(payload.get("target") or ""), + target_key=str(payload.get("target_key") or ""), + status=str(payload.get("status") or "unresolved"), # type: ignore[arg-type] + contact_point_id=( + str(payload["contact_point_id"]) + if payload.get("contact_point_id") + else None + ), + address_purpose=( + str(payload["address_purpose"]) + if payload.get("address_purpose") + else None + ), + locale=str(payload["locale"]) if payload.get("locale") else None, + preferred=bool(payload.get("preferred")), + preference_rank=( + int(payload["preference_rank"]) + if isinstance(payload.get("preference_rank"), int) + else None + ), + reason_code=( + str(payload["reason_code"]) if payload.get("reason_code") else None + ), + explanation=( + str(payload["explanation"]) if payload.get("explanation") else None + ), + source=( + _source_reference_from_payload(dict(source_payload)) + if isinstance(source_payload, dict) + else None + ), + source_revision=( + str(payload["source_revision"]) + if payload.get("source_revision") + else None + ), + preference_revision=( + str(payload["preference_revision"]) + if payload.get("preference_revision") + else None + ), + consent_revision=( + str(payload["consent_revision"]) + if payload.get("consent_revision") + else None + ), + value=dict(payload.get("value") or {}), + provenance=dict(payload.get("provenance") or {}), + ) + + +def _explanation_from_payload(payload: dict[str, Any]) -> DistributionExplanation: + source_payload = payload.get("source") + return DistributionExplanation( + code=str(payload.get("code") or "addresses.snapshot"), + message=str(payload.get("message") or ""), + severity=str(payload.get("severity") or "warning"), # type: ignore[arg-type] + provider=str(payload["provider"]) if payload.get("provider") else None, + source=( + _source_reference_from_payload(dict(source_payload)) + if isinstance(source_payload, dict) + else None + ), + provenance=dict(payload.get("provenance") or {}), + ) + + +def _resolution_from_payload(payload: dict[str, Any]) -> ContactPointResolution: + return ContactPointResolution( + contract_version=str( + payload.get("contract_version") or CONTACT_POINT_CONTRACT_VERSION + ), + subject=_source_reference_from_payload(dict(payload.get("subject") or {})), + status=str(payload.get("status") or "unresolved"), # type: ignore[arg-type] + contact_id=str(payload["contact_id"]) if payload.get("contact_id") else None, + display_name=( + str(payload["display_name"]) if payload.get("display_name") else None + ), + candidates=tuple( + _candidate_from_payload(dict(item)) + for item in payload.get("candidates") or [] + if isinstance(item, dict) + ), + excluded=tuple( + _candidate_from_payload(dict(item)) + for item in payload.get("excluded") or [] + if isinstance(item, dict) + ), + explanations=tuple( + _explanation_from_payload(dict(item)) + for item in payload.get("explanations") or [] + if isinstance(item, dict) + ), + source_revision=( + str(payload["source_revision"]) + if payload.get("source_revision") + else None + ), + source_fingerprint=( + str(payload["source_fingerprint"]) + if payload.get("source_fingerprint") + else None + ), + provenance=dict(payload.get("provenance") or {}), + ) + + +def _source_request_from_payload(payload: dict[str, Any]) -> ContactPointSourceRequest: + effective_at = datetime.fromisoformat(str(payload["effective_at"])) + return ContactPointSourceRequest( + tenant_id=str(payload.get("tenant_id") or ""), + source_id=str(payload.get("source_id") or ""), + effective_at=effective_at, + purpose=str(payload["purpose"]) if payload.get("purpose") else None, + requested_channels=tuple(payload.get("requested_channels") or ()), # type: ignore[arg-type] + address_purpose=( + str(payload["address_purpose"]) + if payload.get("address_purpose") + else None + ), + fallback_rule=str(payload.get("fallback_rule") or "primary"), # type: ignore[arg-type] + locale=str(payload["locale"]) if payload.get("locale") else None, + postal_format=str(payload.get("postal_format") or "domestic"), # type: ignore[arg-type] + max_items=int(payload.get("max_items") or 5000), + context=dict(payload.get("context") or {}), + ) + + +def _contact_point_snapshot_ref(snapshot: ContactPointSnapshot) -> ContactPointSnapshotRef: + provenance = dict(snapshot.provenance or {}) + source_payload = provenance.get("source") + source = _source_reference_from_payload( + dict(source_payload) if isinstance(source_payload, dict) else {} + ) + return ContactPointSnapshotRef( + id=snapshot.id, + tenant_id=snapshot.tenant_id, + contract_version=snapshot.contract_version, + source=source, + request=_source_request_from_payload(dict(snapshot.request_payload or {})), + resolutions=tuple( + _resolution_from_payload(dict(item)) + for item in snapshot.resolution_payload or [] + ), + recipient_count=snapshot.recipient_count, + excluded_count=snapshot.excluded_count, + source_revision=snapshot.source_revision, + source_fingerprint=snapshot.source_fingerprint, + snapshot_hash=snapshot.snapshot_hash, + generated_at=_aware_datetime(snapshot.generated_at), + provenance=provenance, + ) diff --git a/src/govoplan_addresses/backend/db/models.py b/src/govoplan_addresses/backend/db/models.py index b351124..30b53cb 100644 --- a/src/govoplan_addresses/backend/db/models.py +++ b/src/govoplan_addresses/backend/db/models.py @@ -56,6 +56,11 @@ class Contact(Base, TimestampMixin): __table_args__ = ( Index("ix_addresses_contacts_book_name", "address_book_id", "display_name"), Index("ix_addresses_contacts_tenant_name", "tenant_id", "display_name"), + Index( + "ix_addresses_contacts_source_ref", + "source_ref", + postgresql_using="hash", + ), ) id: Mapped[str] = mapped_column(String(36), primary_key=True, default=new_uuid) @@ -191,6 +196,31 @@ class ContactChannelRule(Base, TimestampMixin): contact: Mapped[Contact] = relationship(back_populates="channel_rules") +class ContactPointSnapshot(Base, TimestampMixin): + __tablename__ = "addresses_contact_point_snapshots" + __table_args__ = ( + Index("ix_addresses_contact_point_snapshots_source", "tenant_id", "source_id", "created_at"), + Index("ix_addresses_contact_point_snapshots_hash", "tenant_id", "snapshot_hash"), + ) + + id: Mapped[str] = mapped_column(String(36), primary_key=True, default=new_uuid) + tenant_id: Mapped[str] = mapped_column(String(36), nullable=False, index=True) + source_id: Mapped[str] = mapped_column(String(255), nullable=False, index=True) + contract_version: Mapped[str] = mapped_column(String(20), nullable=False) + source_revision: Mapped[str] = mapped_column(String(255), nullable=False) + source_fingerprint: Mapped[str] = mapped_column(String(64), nullable=False) + purpose: Mapped[str | None] = mapped_column(String(120), nullable=True, index=True) + effective_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, index=True) + generated_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, index=True) + request_payload: Mapped[dict[str, Any]] = mapped_column(JSON, nullable=False) + resolution_payload: Mapped[list[dict[str, Any]]] = mapped_column(JSON, nullable=False) + recipient_count: Mapped[int] = mapped_column(Integer, nullable=False) + excluded_count: Mapped[int] = mapped_column(Integer, nullable=False) + snapshot_hash: Mapped[str] = mapped_column(String(64), nullable=False, index=True) + created_by_account_id: Mapped[str | None] = mapped_column(String(36), nullable=True, index=True) + provenance: Mapped[dict[str, Any]] = mapped_column(JSON, default=dict, nullable=False) + + class AddressList(Base, TimestampMixin): __tablename__ = "addresses_address_lists" __table_args__ = ( @@ -371,6 +401,7 @@ __all__ = [ "Contact", "ContactEmail", "ContactPhone", + "ContactPointSnapshot", "ContactPostalAddress", "new_uuid", ] diff --git a/src/govoplan_addresses/backend/manifest.py b/src/govoplan_addresses/backend/manifest.py index df1db97..e7b3c53 100644 --- a/src/govoplan_addresses/backend/manifest.py +++ b/src/govoplan_addresses/backend/manifest.py @@ -12,6 +12,7 @@ from govoplan_addresses.backend.capabilities import ( ) from govoplan_addresses.backend.db import models as addresses_models # noqa: F401 - populate address ORM metadata from govoplan_core.core.access import CAPABILITY_AUTH_PERMISSION_EVALUATOR, CAPABILITY_AUTH_PRINCIPAL_RESOLVER +from govoplan_core.core.contact_points import CAPABILITY_ADDRESSES_CONTACT_POINT_RESOLUTION from govoplan_core.core.module_guards import drop_table_retirement_provider, persistent_table_uninstall_guard from govoplan_core.core.people import CAPABILITY_ADDRESSES_PEOPLE_SEARCH from govoplan_core.core.distribution_lists import CAPABILITY_RECIPIENT_CHANNEL_FACTS @@ -42,6 +43,7 @@ from govoplan_addresses.backend.provider_state import ( _addresses_table_retirement_provider = drop_table_retirement_provider( + addresses_models.ContactPointSnapshot, addresses_models.AddressSyncDiagnostic, addresses_models.AddressSyncConflict, addresses_models.AddressSyncTombstone, @@ -147,12 +149,19 @@ ROLE_TEMPLATES = ( def _tenant_summary(session, tenant_id: str) -> dict[str, int]: - from govoplan_addresses.backend.db.models import AddressBook, AddressList, AddressSyncSource, Contact + from govoplan_addresses.backend.db.models import ( + AddressBook, + AddressList, + AddressSyncSource, + Contact, + ContactPointSnapshot, + ) return { "address_books": session.query(AddressBook).filter(AddressBook.tenant_id == tenant_id, AddressBook.deleted_at.is_(None)).count(), "address_lists": session.query(AddressList).filter(AddressList.tenant_id == tenant_id, AddressList.deleted_at.is_(None)).count(), "contacts": session.query(Contact).filter(Contact.tenant_id == tenant_id, Contact.deleted_at.is_(None)).count(), + "contact_point_snapshots": session.query(ContactPointSnapshot).filter(ContactPointSnapshot.tenant_id == tenant_id).count(), "sync_sources": session.query(AddressSyncSource).filter(AddressSyncSource.tenant_id == tenant_id, AddressSyncSource.enabled.is_(True)).count(), } @@ -226,6 +235,7 @@ manifest = ModuleManifest( ModuleInterfaceProvider(name=CAPABILITY_ADDRESSES_LOOKUP, version="0.1.8"), ModuleInterfaceProvider(name=CAPABILITY_ADDRESSES_PEOPLE_SEARCH, version="0.1.0"), ModuleInterfaceProvider(name=CAPABILITY_ADDRESSES_RECIPIENT_SOURCE, version="0.1.9"), + ModuleInterfaceProvider(name=CAPABILITY_ADDRESSES_CONTACT_POINT_RESOLUTION, version="1.0.0"), ModuleInterfaceProvider(name=CAPABILITY_ADDRESSES_CONTACT_WRITER, version="0.1.8"), ModuleInterfaceProvider(name=CAPABILITY_RECIPIENT_CHANNEL_FACTS, version="0.1.0"), ), @@ -263,6 +273,10 @@ manifest = ModuleManifest( "govoplan_addresses.backend.capabilities", fromlist=["channel_facts_capability"], ).channel_facts_capability(context), + CAPABILITY_ADDRESSES_CONTACT_POINT_RESOLUTION: lambda context: __import__( + "govoplan_addresses.backend.capabilities", + fromlist=["contact_point_resolution_capability"], + ).contact_point_resolution_capability(context), }, uninstall_guard_providers=( persistent_table_uninstall_guard( @@ -272,6 +286,7 @@ manifest = ModuleManifest( addresses_models.AddressSyncSource, addresses_models.AddressListEntry, addresses_models.AddressList, + addresses_models.ContactPointSnapshot, addresses_models.AddressBook, addresses_models.Contact, addresses_models.ContactEmail, @@ -297,6 +312,22 @@ manifest = ModuleManifest( related_modules=("campaigns", "mail", "forms", "reporting", "portal", "postbox"), order=30, ), + DocumentationTopic( + id="addresses.contact-point-resolution", + title="Contact-point resolution and snapshots", + summary="Resolve purpose-aware channel targets and freeze immutable recipient evidence.", + body=( + "Addresses exposes a versioned contact-point capability for email, postal, internal-mail, and portal targets. " + "Callers can request an effective date, communication purpose, address purpose, fallback rule, locale, and " + "postal format. Bounded previews remain live; frozen snapshots retain the resolved values, exclusions, " + "source and governance revisions, provenance, and a deterministic evidence hash even after contacts change." + ), + layer="configured", + documentation_types=("admin", "user"), + audience=("tenant_admin", "operator", "module_admin"), + related_modules=("dist_lists", "campaigns", "policy", "templates"), + order=31, + ), ), external_providers=(CARDDAV_PROVIDER,), external_provider_state_providers=( diff --git a/src/govoplan_addresses/backend/migrations/versions/a3b5c6d7e8f9_contact_point_snapshots.py b/src/govoplan_addresses/backend/migrations/versions/a3b5c6d7e8f9_contact_point_snapshots.py new file mode 100644 index 0000000..00a2523 --- /dev/null +++ b/src/govoplan_addresses/backend/migrations/versions/a3b5c6d7e8f9_contact_point_snapshots.py @@ -0,0 +1,63 @@ +"""Add immutable contact-point snapshots. + +Revision ID: a3b5c6d7e8f9 +Revises: f2a4b5c6d7e +""" + +from alembic import op +import sqlalchemy as sa + + +revision = "a3b5c6d7e8f9" +down_revision = "f2a4b5c6d7e" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + op.create_index( + "ix_addresses_contacts_source_ref", + "addresses_contacts", + ["source_ref"], + unique=False, + postgresql_using="hash", + ) + op.create_table( + "addresses_contact_point_snapshots", + sa.Column("id", sa.String(length=36), nullable=False), + sa.Column("tenant_id", sa.String(length=36), nullable=False), + sa.Column("source_id", sa.String(length=255), nullable=False), + sa.Column("contract_version", sa.String(length=20), nullable=False), + sa.Column("source_revision", sa.String(length=255), nullable=False), + sa.Column("source_fingerprint", sa.String(length=64), nullable=False), + sa.Column("purpose", sa.String(length=120), nullable=True), + sa.Column("effective_at", sa.DateTime(timezone=True), nullable=False), + sa.Column("generated_at", sa.DateTime(timezone=True), nullable=False), + sa.Column("request_payload", sa.JSON(), nullable=False), + sa.Column("resolution_payload", sa.JSON(), nullable=False), + sa.Column("recipient_count", sa.Integer(), nullable=False), + sa.Column("excluded_count", sa.Integer(), nullable=False), + sa.Column("snapshot_hash", sa.String(length=64), nullable=False), + sa.Column("created_by_account_id", sa.String(length=36), nullable=True), + sa.Column("provenance", sa.JSON(), nullable=False), + sa.Column("created_at", sa.DateTime(timezone=True), nullable=False), + sa.Column("updated_at", sa.DateTime(timezone=True), nullable=False), + sa.PrimaryKeyConstraint("id"), + ) + for name, columns in ( + ("ix_addresses_contact_point_snapshots_tenant_id", ["tenant_id"]), + ("ix_addresses_contact_point_snapshots_source_id", ["source_id"]), + ("ix_addresses_contact_point_snapshots_purpose", ["purpose"]), + ("ix_addresses_contact_point_snapshots_effective_at", ["effective_at"]), + ("ix_addresses_contact_point_snapshots_generated_at", ["generated_at"]), + ("ix_addresses_contact_point_snapshots_snapshot_hash", ["snapshot_hash"]), + ("ix_addresses_contact_point_snapshots_created_by_account_id", ["created_by_account_id"]), + ("ix_addresses_contact_point_snapshots_source", ["tenant_id", "source_id", "created_at"]), + ("ix_addresses_contact_point_snapshots_hash", ["tenant_id", "snapshot_hash"]), + ): + op.create_index(name, "addresses_contact_point_snapshots", columns, unique=False) + + +def downgrade() -> None: + op.drop_table("addresses_contact_point_snapshots") + op.drop_index("ix_addresses_contacts_source_ref", table_name="addresses_contacts") diff --git a/src/govoplan_addresses/backend/router.py b/src/govoplan_addresses/backend/router.py index 6dcc880..fe85996 100644 --- a/src/govoplan_addresses/backend/router.py +++ b/src/govoplan_addresses/backend/router.py @@ -8,6 +8,11 @@ from sqlalchemy.orm import Session from govoplan_core.audit.logging import audit_from_principal from govoplan_core.auth import ApiPrincipal, get_api_principal, has_scope +from govoplan_core.core.contact_points import ( + ContactPointResolutionRequest, + ContactPointSourceRequest, +) +from govoplan_core.core.distribution_lists import DistributionSourceReference from govoplan_core.db.session import get_session from govoplan_addresses.backend.carddav import AddressCardDAVError from govoplan_addresses.backend.db.models import ( @@ -22,7 +27,10 @@ from govoplan_addresses.backend.db.models import ( ContactChannelRule, ContactPostalAddress, ) -from govoplan_addresses.backend.capabilities import AddressesContactWriterCapability +from govoplan_addresses.backend.capabilities import ( + AddressesContactPointResolutionCapability, + AddressesContactWriterCapability, +) from govoplan_addresses.backend.schemas import ( AddressBookCreateRequest, AddressBookListResponse, @@ -66,6 +74,11 @@ from govoplan_addresses.backend.schemas import ( ContactChannelRuleCreateRequest, ContactChannelRuleListResponse, ContactChannelRuleResponse, + ContactPointResolveRequest, + ContactPointResolutionResponse, + ContactPointSnapshotResponse, + ContactPointSourcePreviewResponse, + ContactPointSourceRequestPayload, ContactListResponse, ContactResponse, ContactUpdateRequest, @@ -212,6 +225,43 @@ def _write_decision_response(decision) -> AddressBookWriteDecisionResponse: return AddressBookWriteDecisionResponse.model_validate(payload) +def _contact_point_resolution_request( + principal: ApiPrincipal, + payload: ContactPointResolveRequest, +) -> ContactPointResolutionRequest: + return ContactPointResolutionRequest( + tenant_id=principal.tenant_id, + subject=DistributionSourceReference(**payload.subject.model_dump()), + effective_at=payload.effective_at, + purpose=payload.purpose, + requested_channels=tuple(payload.requested_channels), + address_purpose=payload.address_purpose, + fallback_rule=payload.fallback_rule, + locale=payload.locale, + postal_format=payload.postal_format, + context=payload.context, + ) + + +def _contact_point_source_request( + principal: ApiPrincipal, + payload: ContactPointSourceRequestPayload, +) -> ContactPointSourceRequest: + return ContactPointSourceRequest( + tenant_id=principal.tenant_id, + source_id=payload.source_id, + effective_at=payload.effective_at, + purpose=payload.purpose, + requested_channels=tuple(payload.requested_channels), + address_purpose=payload.address_purpose, + fallback_rule=payload.fallback_rule, + locale=payload.locale, + postal_format=payload.postal_format, + max_items=payload.max_items, + context=payload.context, + ) + + def _sync_source_response(sync_source: AddressSyncSource) -> AddressSyncSourceResponse: return AddressSyncSourceResponse.model_validate( { @@ -493,6 +543,124 @@ def api_lookup_addresses( return AddressLookupResponse(contacts=[_contact_response(contact) for contact in contacts]) +@router.post("/contact-points/resolve", response_model=ContactPointResolutionResponse) +def api_resolve_contact_points( + payload: ContactPointResolveRequest, + principal: ApiPrincipal = Depends(get_api_principal), + session: Session = Depends(get_session), +): + _require_scope(principal, "addresses:contact:read") + _require_scope(principal, "addresses:governance:read") + try: + result = AddressesContactPointResolutionCapability().resolve_contact_points( + session, + principal, + request=_contact_point_resolution_request(principal, payload), + ) + return ContactPointResolutionResponse.model_validate(asdict(result)) + except (AddressBookError, ValueError) as exc: + raise HTTPException( + status_code=status.HTTP_422_UNPROCESSABLE_CONTENT, + detail=str(exc), + ) from exc + + +@router.post( + "/contact-point-sources/preview", + response_model=ContactPointSourcePreviewResponse, +) +def api_preview_contact_point_source( + payload: ContactPointSourceRequestPayload, + offset: int = Query(default=0, ge=0), + limit: int = Query(default=100, ge=1, le=500), + principal: ApiPrincipal = Depends(get_api_principal), + session: Session = Depends(get_session), +): + _require_scope(principal, "addresses:contact:read") + _require_scope(principal, "addresses:governance:read") + try: + result = AddressesContactPointResolutionCapability().preview_source( + session, + principal, + request=_contact_point_source_request(principal, payload), + offset=offset, + limit=limit, + ) + return ContactPointSourcePreviewResponse.model_validate(asdict(result)) + except (AddressBookError, ValueError) as exc: + raise HTTPException( + status_code=status.HTTP_422_UNPROCESSABLE_CONTENT, + detail=str(exc), + ) from exc + + +@router.post( + "/contact-point-snapshots", + response_model=ContactPointSnapshotResponse, + status_code=status.HTTP_201_CREATED, +) +def api_freeze_contact_point_source( + payload: ContactPointSourceRequestPayload, + principal: ApiPrincipal = Depends(get_api_principal), + session: Session = Depends(get_session), +): + _require_scope(principal, "addresses:contact:read") + _require_scope(principal, "addresses:governance:read") + try: + snapshot = AddressesContactPointResolutionCapability().freeze_source( + session, + principal, + request=_contact_point_source_request(principal, payload), + ) + audit_from_principal( + session, + principal, + action="addresses.contact_point_snapshot_created", + object_type="address_contact_point_snapshot", + object_id=snapshot.id, + details={ + "source_id": snapshot.request.source_id, + "source_revision": snapshot.source_revision, + "snapshot_hash": snapshot.snapshot_hash, + "recipient_count": snapshot.recipient_count, + "excluded_count": snapshot.excluded_count, + "purpose": snapshot.request.purpose, + }, + ) + session.commit() + return ContactPointSnapshotResponse.model_validate(asdict(snapshot)) + except (AddressBookError, ValueError) as exc: + session.rollback() + raise HTTPException( + status_code=status.HTTP_422_UNPROCESSABLE_CONTENT, + detail=str(exc), + ) from exc + + +@router.get( + "/contact-point-snapshots/{snapshot_id}", + response_model=ContactPointSnapshotResponse, +) +def api_get_contact_point_snapshot( + snapshot_id: str, + principal: ApiPrincipal = Depends(get_api_principal), + session: Session = Depends(get_session), +): + _require_scope(principal, "addresses:contact:read") + _require_scope(principal, "addresses:governance:read") + snapshot = AddressesContactPointResolutionCapability().get_snapshot( + session, + principal, + snapshot_id=snapshot_id, + ) + if snapshot is None: + raise HTTPException( + status_code=status.HTTP_404_NOT_FOUND, + detail="Contact-point snapshot not found.", + ) + return ContactPointSnapshotResponse.model_validate(asdict(snapshot)) + + @router.get("/address-lists", response_model=AddressListListResponse) def api_list_address_lists( address_book_id: str | None = Query(default=None), diff --git a/src/govoplan_addresses/backend/schemas.py b/src/govoplan_addresses/backend/schemas.py index b6ab2ff..a137c2a 100644 --- a/src/govoplan_addresses/backend/schemas.py +++ b/src/govoplan_addresses/backend/schemas.py @@ -25,6 +25,19 @@ AddressChannelDecision = Literal[ "returned", "temporarily_unavailable", ] +AddressDistributionOutcome = Literal[ + "usable", + "unresolved", + "invalid", + "suppressed", + "ambiguous", + "duplicate", + "policy_blocked", + "provider_unavailable", + "stale", +] +AddressContactPointFallbackRule = Literal["none", "primary", "any"] +AddressPostalFormat = Literal["domestic", "international"] class ContactEmailPayload(BaseModel): @@ -282,6 +295,125 @@ class ContactChannelRuleListResponse(BaseModel): rules: list[ContactChannelRuleResponse] = Field(default_factory=list) +class AddressSourceReferencePayload(BaseModel): + provider: str = Field(min_length=1, max_length=120) + resource_type: str = Field(min_length=1, max_length=120) + resource_id: str = Field(min_length=1, max_length=1000) + revision: str | None = Field(default=None, max_length=1000) + fingerprint: str | None = Field(default=None, max_length=255) + label: str | None = Field(default=None, max_length=500) + metadata: dict[str, Any] = Field(default_factory=dict) + + +class ContactPointResolveRequest(BaseModel): + model_config = ConfigDict(extra="forbid") + + subject: AddressSourceReferencePayload + effective_at: datetime + purpose: str | None = Field(default=None, max_length=120) + requested_channels: list[AddressDistributionChannel] = Field(default_factory=list) + address_purpose: str | None = Field(default=None, max_length=80) + fallback_rule: AddressContactPointFallbackRule = "primary" + locale: str | None = Field(default=None, max_length=20) + postal_format: AddressPostalFormat = "domestic" + context: dict[str, Any] = Field(default_factory=dict) + + +class ContactPointSourceRequestPayload(BaseModel): + model_config = ConfigDict(extra="forbid") + + source_id: str = Field(min_length=1, max_length=1000) + effective_at: datetime + purpose: str | None = Field(default=None, max_length=120) + requested_channels: list[AddressDistributionChannel] = Field(default_factory=list) + address_purpose: str | None = Field(default=None, max_length=80) + fallback_rule: AddressContactPointFallbackRule = "primary" + locale: str | None = Field(default=None, max_length=20) + postal_format: AddressPostalFormat = "domestic" + max_items: int = Field(default=5000, ge=1, le=20000) + context: dict[str, Any] = Field(default_factory=dict) + + +class ContactPointSourceRequestResponse(ContactPointSourceRequestPayload): + tenant_id: str + + +class ContactPointCandidateResponse(BaseModel): + channel: AddressDistributionChannel + target: str + target_key: str + status: AddressDistributionOutcome + contact_point_id: str | None = None + address_purpose: str | None = None + locale: str | None = None + preferred: bool = False + preference_rank: int | None = None + reason_code: str | None = None + explanation: str | None = None + source: AddressSourceReferencePayload | None = None + source_revision: str | None = None + preference_revision: str | None = None + consent_revision: str | None = None + value: dict[str, Any] = Field(default_factory=dict) + provenance: dict[str, Any] = Field(default_factory=dict) + + +class DistributionExplanationResponse(BaseModel): + code: str + message: str + severity: Literal["info", "warning", "error"] + provider: str | None = None + source: AddressSourceReferencePayload | None = None + provenance: dict[str, Any] = Field(default_factory=dict) + + +class ContactPointResolutionResponse(BaseModel): + contract_version: str + subject: AddressSourceReferencePayload + status: AddressDistributionOutcome + contact_id: str | None = None + display_name: str | None = None + candidates: list[ContactPointCandidateResponse] = Field(default_factory=list) + excluded: list[ContactPointCandidateResponse] = Field(default_factory=list) + explanations: list[DistributionExplanationResponse] = Field(default_factory=list) + source_revision: str | None = None + source_fingerprint: str | None = None + provenance: dict[str, Any] = Field(default_factory=dict) + + +class ContactPointSourcePreviewResponse(BaseModel): + contract_version: str + source: AddressSourceReferencePayload + request: ContactPointSourceRequestResponse + resolutions: list[ContactPointResolutionResponse] + total_count: int + usable_count: int + excluded_count: int + offset: int + limit: int + has_more: bool + source_revision: str + source_fingerprint: str + generated_at: datetime + provenance: dict[str, Any] = Field(default_factory=dict) + + +class ContactPointSnapshotResponse(BaseModel): + id: str + tenant_id: str + contract_version: str + source: AddressSourceReferencePayload + request: ContactPointSourceRequestResponse + resolutions: list[ContactPointResolutionResponse] + recipient_count: int + excluded_count: int + source_revision: str + source_fingerprint: str + snapshot_hash: str + generated_at: datetime + provenance: dict[str, Any] = Field(default_factory=dict) + + class AddressLookupResponse(BaseModel): contacts: list[ContactResponse] diff --git a/tests/test_addresses_service.py b/tests/test_addresses_service.py index 897542f..241bdf1 100644 --- a/tests/test_addresses_service.py +++ b/tests/test_addresses_service.py @@ -1,6 +1,7 @@ from __future__ import annotations import unittest +from dataclasses import asdict from datetime import timedelta from unittest.mock import patch @@ -13,6 +14,12 @@ from govoplan_core.core.distribution_lists import ( DistributionSourceReference, RecipientChannelFactsRequest, ) +from govoplan_core.core.contact_points import ( + CAPABILITY_ADDRESSES_CONTACT_POINT_RESOLUTION, + CONTACT_POINT_CONTRACT_VERSION, + ContactPointResolutionRequest, + ContactPointSourceRequest, +) from govoplan_core.core.people import CAPABILITY_ADDRESSES_PEOPLE_SEARCH, PeopleSearchProvider from govoplan_core.db.base import Base from govoplan_core.db.base import utcnow @@ -26,6 +33,7 @@ from govoplan_addresses.backend.capabilities import ( CAPABILITY_ADDRESSES_LOOKUP, CAPABILITY_ADDRESSES_RECIPIENT_SOURCE, AddressesChannelFactsCapability, + AddressesContactPointResolutionCapability, AddressesContactWriterCapability, AddressesLookupCapability, AddressesPeopleSearchProvider, @@ -44,6 +52,7 @@ from govoplan_addresses.backend.db.models import ( ContactChannelRule, ContactEmail, ContactPhone, + ContactPointSnapshot, ContactPostalAddress, ) from govoplan_addresses.backend.schemas import ( @@ -63,6 +72,7 @@ from govoplan_addresses.backend.schemas import ( ContactChannelRuleCreateRequest, ContactEmailPayload, ContactPostalAddressPayload, + ContactPointSnapshotResponse, ) from govoplan_addresses.backend.manifest import manifest from govoplan_addresses.backend.router import _sync_source_response @@ -196,6 +206,7 @@ class AddressServiceTest(unittest.TestCase): ContactPhone.__table__, ContactPostalAddress.__table__, ContactChannelRule.__table__, + ContactPointSnapshot.__table__, AddressListEntry.__table__, AddressSyncSource.__table__, AddressSyncTombstone.__table__, @@ -397,6 +408,11 @@ END:VCARD self.assertIn(CAPABILITY_ADDRESSES_RECIPIENT_SOURCE, provided) self.assertIn(CAPABILITY_ADDRESSES_CONTACT_WRITER, provided) self.assertIn(CAPABILITY_RECIPIENT_CHANNEL_FACTS, provided) + self.assertIn(CAPABILITY_ADDRESSES_CONTACT_POINT_RESOLUTION, provided) + self.assertIn( + CAPABILITY_ADDRESSES_CONTACT_POINT_RESOLUTION, + manifest.capability_factories, + ) book = create_address_book(self.session, self.principal, AddressBookCreateRequest(scope_type="user", name="Recipients")) self.session.commit() @@ -691,6 +707,168 @@ END:VCARD self.session.commit() self.assertEqual(list_address_list_entries(self.session, self.principal, address_list.id), []) + def test_contact_point_resolution_supports_external_refs_and_frozen_postal_snapshots(self) -> None: + book = create_address_book( + self.session, + self.principal, + AddressBookCreateRequest(scope_type="user", name="Official contacts"), + ) + self.session.commit() + self.session.refresh(book) + contact = create_contact( + self.session, + self.principal, + book.id, + ContactCreateRequest( + display_name="Ada Lovelace", + emails=[ + ContactEmailPayload( + label="private", + email="ada.private@example.local", + is_primary=True, + ) + ], + postal_addresses=[ + ContactPostalAddressPayload( + label="official", + street="Main Street 1", + postal_code="10115", + locality="Berlin", + country="Germany", + is_primary=True, + ), + ContactPostalAddressPayload( + label="private", + street="Side Street 2", + postal_code="10117", + locality="Berlin", + country="Germany", + ), + ], + ), + ) + contact.source_kind = "idm" + contact.source_ref = "idm:identity:identity-1" + self.session.flush() + create_contact_channel_rule( + self.session, + self.principal, + contact.id, + ContactChannelRuleCreateRequest( + channel="postal", + purpose="official_notice", + contact_point_id=contact.postal_addresses[0].id, + decision="preferred", + legal_basis="public_task", + evidence_ref="idm:function-assignment:17", + preference_rank=1, + locale="de-DE", + ), + ) + address_list = create_address_list( + self.session, + self.principal, + book.id, + AddressListCreateRequest(name="Postal recipients"), + ) + self.session.flush() + postal_entry = create_address_list_entry( + self.session, + self.principal, + address_list.id, + AddressListEntryCreateRequest( + contact_id=contact.id, + contact_postal_address_id=contact.postal_addresses[0].id, + ), + ) + self.session.commit() + + capability = AddressesContactPointResolutionCapability() + direct = capability.resolve_contact_points( + self.session, + self.principal, + request=ContactPointResolutionRequest( + tenant_id=self.principal.tenant_id, + subject=DistributionSourceReference( + provider="idm", + resource_type="identity", + resource_id="identity-1", + ), + effective_at=utcnow(), + purpose="official_notice", + requested_channels=("postal",), + address_purpose="official", + fallback_rule="none", + locale="de-DE", + postal_format="international", + ), + ) + self.assertEqual(CONTACT_POINT_CONTRACT_VERSION, direct.contract_version) + self.assertEqual("usable", direct.status) + self.assertEqual(contact.id, direct.contact_id) + self.assertEqual(1, len(direct.candidates)) + self.assertEqual(contact.postal_addresses[0].id, direct.candidates[0].contact_point_id) + self.assertEqual("official", direct.candidates[0].address_purpose) + self.assertIn("Ada Lovelace", direct.candidates[0].target) + self.assertIn("Germany", direct.candidates[0].target) + self.assertEqual("de-DE", direct.candidates[0].locale) + self.assertTrue(direct.candidates[0].preference_revision) + self.assertEqual(1, len(direct.excluded)) + self.assertEqual( + "addresses.address_purpose.not_selected", + direct.excluded[0].reason_code, + ) + + source_request = ContactPointSourceRequest( + tenant_id=self.principal.tenant_id, + source_id=f"addresses:address_list:{address_list.id}", + effective_at=utcnow(), + purpose="official_notice", + requested_channels=("email", "postal"), + address_purpose="official", + fallback_rule="none", + locale="de-DE", + postal_format="international", + ) + preview = capability.preview_source( + self.session, + self.principal, + request=source_request, + limit=1, + ) + self.assertEqual(1, preview.total_count) + self.assertEqual(1, preview.usable_count) + self.assertFalse(preview.has_more) + self.assertEqual(postal_entry.id, preview.resolutions[0].provenance["address_list_entry_ids"][0]) + self.assertEqual("postal", preview.resolutions[0].candidates[0].channel) + + snapshot = capability.freeze_source( + self.session, + self.principal, + request=source_request, + ) + self.session.commit() + original_target = snapshot.resolutions[0].candidates[0].target + contact.postal_addresses[0].street = "Changed Street 99" + contact.postal_addresses[0].country = "France" + self.session.commit() + + frozen = capability.get_snapshot( + self.session, + self.principal, + snapshot_id=snapshot.id, + ) + self.assertIsNotNone(frozen) + assert frozen is not None + self.assertEqual(original_target, frozen.resolutions[0].candidates[0].target) + self.assertIn("Main Street 1", frozen.resolutions[0].candidates[0].target) + self.assertNotIn("Changed Street 99", frozen.resolutions[0].candidates[0].target) + self.assertEqual(snapshot.snapshot_hash, frozen.snapshot_hash) + self.assertEqual(1, frozen.recipient_count) + response = ContactPointSnapshotResponse.model_validate(asdict(frozen)) + self.assertEqual(snapshot.id, response.id) + self.assertEqual("postal", response.resolutions[0].candidates[0].channel) + def test_sync_source_marks_read_only_books_and_can_be_made_writable(self) -> None: book = create_address_book(self.session, self.principal, AddressBookCreateRequest(scope_type="user", name="CardDAV")) self.session.commit()