Implement address quality and reversible contact merges

This commit is contained in:
2026-08-02 07:03:27 +02:00
parent 2e78b9ae50
commit 19e9096572
15 changed files with 4369 additions and 31 deletions
@@ -0,0 +1,281 @@
"""Add address quality, provenance, merge evidence, and redirects.
Revision ID: b4c6d7e8f9a0
Revises: a3b5c6d7e8f9
"""
from __future__ import annotations
import re
from alembic import op
import sqlalchemy as sa
revision = "b4c6d7e8f9a0"
down_revision = "a3b5c6d7e8f9"
branch_labels = None
depends_on = None
_JSON_OBJECT = sa.text("'{}'")
def upgrade() -> None:
with op.batch_alter_table("addresses_contact_emails") as batch:
batch.add_column(sa.Column("original_email", sa.String(length=320), nullable=False, server_default=""))
batch.add_column(sa.Column("normalized_email", sa.String(length=320), nullable=False, server_default=""))
batch.add_column(sa.Column("provenance", sa.JSON(), nullable=False, server_default=_JSON_OBJECT))
with op.batch_alter_table("addresses_contact_phones") as batch:
batch.add_column(sa.Column("original_phone", sa.String(length=100), nullable=False, server_default=""))
batch.add_column(sa.Column("normalized_phone", sa.String(length=100), nullable=False, server_default=""))
batch.add_column(sa.Column("provenance", sa.JSON(), nullable=False, server_default=_JSON_OBJECT))
with op.batch_alter_table("addresses_contact_postal_addresses") as batch:
batch.add_column(sa.Column("original_value", sa.JSON(), nullable=False, server_default=_JSON_OBJECT))
batch.add_column(sa.Column("normalized_value", sa.JSON(), nullable=False, server_default=_JSON_OBJECT))
batch.add_column(sa.Column("provenance", sa.JSON(), nullable=False, server_default=_JSON_OBJECT))
bind = op.get_bind()
bind.execute(
sa.text(
"UPDATE addresses_contact_emails "
"SET original_email = email, normalized_email = lower(trim(email))"
)
)
phone_rows = bind.execute(
sa.text("SELECT id, phone FROM addresses_contact_phones")
).mappings().all()
for row in phone_rows:
bind.execute(
sa.text(
"UPDATE addresses_contact_phones "
"SET original_phone = :original, normalized_phone = :normalized "
"WHERE id = :id"
),
{
"id": row["id"],
"original": row["phone"],
"normalized": _normalized_phone(str(row["phone"] or "")),
},
)
postal = sa.table(
"addresses_contact_postal_addresses",
sa.column("id", sa.String()),
sa.column("label", sa.String()),
sa.column("street", sa.String()),
sa.column("postal_code", sa.String()),
sa.column("locality", sa.String()),
sa.column("region", sa.String()),
sa.column("country", sa.String()),
sa.column("original_value", sa.JSON()),
sa.column("normalized_value", sa.JSON()),
)
postal_rows = bind.execute(
sa.select(
postal.c.id,
postal.c.label,
postal.c.street,
postal.c.postal_code,
postal.c.locality,
postal.c.region,
postal.c.country,
)
).mappings().all()
for row in postal_rows:
original = {
key: row[key]
for key in ("label", "street", "postal_code", "locality", "region", "country")
}
normalized = {
key: _normalized_text(row[key])
for key in ("label", "street", "postal_code", "locality", "region", "country")
}
bind.execute(
postal.update()
.where(postal.c.id == row["id"])
.values(original_value=original, normalized_value=normalized)
)
op.create_index(
"ix_addresses_contact_emails_normalized_email",
"addresses_contact_emails",
["normalized_email"],
)
op.create_index(
"ix_addresses_contact_phones_normalized_phone",
"addresses_contact_phones",
["normalized_phone"],
)
op.create_table(
"addresses_contact_point_quality_decisions",
sa.Column("id", sa.String(length=36), nullable=False),
sa.Column("tenant_id", sa.String(length=36), nullable=True),
sa.Column("contact_id", sa.String(length=36), nullable=False),
sa.Column("channel", sa.String(length=30), nullable=False),
sa.Column("contact_point_id", sa.String(length=36), nullable=True),
sa.Column("state", sa.String(length=30), nullable=False),
sa.Column("reason_code", sa.String(length=120), nullable=False),
sa.Column("reason", sa.Text(), nullable=True),
sa.Column("evidence_ref", sa.String(length=1000), nullable=True),
sa.Column("effective_from", sa.DateTime(timezone=True), nullable=False),
sa.Column("effective_until", sa.DateTime(timezone=True), nullable=True),
sa.Column("created_by_account_id", sa.String(length=36), nullable=True),
sa.Column("metadata", sa.JSON(), nullable=False, server_default=_JSON_OBJECT),
sa.Column("created_at", sa.DateTime(timezone=True), nullable=False),
sa.Column("updated_at", sa.DateTime(timezone=True), nullable=False),
sa.ForeignKeyConstraint(["contact_id"], ["addresses_contacts.id"], ondelete="CASCADE"),
sa.PrimaryKeyConstraint("id"),
)
for name, columns in (
("ix_addresses_contact_point_quality_decisions_tenant_id", ["tenant_id"]),
("ix_addresses_contact_point_quality_decisions_contact_id", ["contact_id"]),
("ix_addresses_contact_point_quality_decisions_channel", ["channel"]),
("ix_addresses_contact_point_quality_decisions_contact_point_id", ["contact_point_id"]),
("ix_addresses_contact_point_quality_decisions_state", ["state"]),
("ix_addresses_contact_point_quality_decisions_effective_from", ["effective_from"]),
("ix_addresses_contact_point_quality_decisions_effective_until", ["effective_until"]),
("ix_addresses_contact_point_quality_decisions_created_by_account_id", ["created_by_account_id"]),
("ix_addresses_quality_current", ["tenant_id", "contact_id", "channel", "contact_point_id", "effective_until"]),
("ix_addresses_quality_state", ["tenant_id", "state", "effective_until"]),
):
op.create_index(name, "addresses_contact_point_quality_decisions", columns)
op.create_table(
"addresses_contact_merge_records",
sa.Column("id", sa.String(length=36), nullable=False),
sa.Column("tenant_id", sa.String(length=36), nullable=True),
sa.Column("address_book_id", sa.String(length=36), nullable=False),
sa.Column("winner_contact_id", sa.String(length=36), nullable=False),
sa.Column("loser_contact_ids", sa.JSON(), nullable=False),
sa.Column("status", sa.String(length=30), nullable=False),
sa.Column("reason", sa.Text(), nullable=False),
sa.Column("survivorship", sa.JSON(), nullable=False, server_default=_JSON_OBJECT),
sa.Column("decisions", sa.JSON(), nullable=False, server_default="[]"),
sa.Column("before_payload", sa.JSON(), nullable=False),
sa.Column("after_payload", sa.JSON(), nullable=False),
sa.Column("before_hash", sa.String(length=64), nullable=False),
sa.Column("after_hash", sa.String(length=64), nullable=False),
sa.Column("created_by_account_id", sa.String(length=36), nullable=True),
sa.Column("recovered_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("recovered_by_account_id", sa.String(length=36), nullable=True),
sa.Column("recovery_action", sa.String(length=30), nullable=True),
sa.Column("recovery_reason", sa.Text(), nullable=True),
sa.Column("provenance", sa.JSON(), nullable=False, server_default=_JSON_OBJECT),
sa.Column("created_at", sa.DateTime(timezone=True), nullable=False),
sa.Column("updated_at", sa.DateTime(timezone=True), nullable=False),
sa.ForeignKeyConstraint(["address_book_id"], ["addresses_address_books.id"], ondelete="CASCADE"),
sa.ForeignKeyConstraint(["winner_contact_id"], ["addresses_contacts.id"], ondelete="RESTRICT"),
sa.PrimaryKeyConstraint("id"),
)
for name, columns in (
("ix_addresses_contact_merge_records_tenant_id", ["tenant_id"]),
("ix_addresses_contact_merge_records_address_book_id", ["address_book_id"]),
("ix_addresses_contact_merge_records_winner_contact_id", ["winner_contact_id"]),
("ix_addresses_contact_merge_records_status", ["status"]),
("ix_addresses_contact_merge_records_created_by_account_id", ["created_by_account_id"]),
("ix_addresses_merge_winner", ["tenant_id", "winner_contact_id", "created_at"]),
("ix_addresses_merge_status", ["tenant_id", "status", "created_at"]),
):
op.create_index(name, "addresses_contact_merge_records", columns)
op.create_table(
"addresses_contact_redirects",
sa.Column("id", sa.String(length=36), nullable=False),
sa.Column("tenant_id", sa.String(length=36), nullable=True),
sa.Column("source_contact_id", sa.String(length=36), nullable=False),
sa.Column("target_contact_id", sa.String(length=36), nullable=False),
sa.Column("merge_record_id", sa.String(length=36), nullable=False),
sa.Column("ended_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("created_at", sa.DateTime(timezone=True), nullable=False),
sa.Column("updated_at", sa.DateTime(timezone=True), nullable=False),
sa.ForeignKeyConstraint(["merge_record_id"], ["addresses_contact_merge_records.id"], ondelete="CASCADE"),
sa.ForeignKeyConstraint(["source_contact_id"], ["addresses_contacts.id"], ondelete="CASCADE"),
sa.ForeignKeyConstraint(["target_contact_id"], ["addresses_contacts.id"], ondelete="RESTRICT"),
sa.PrimaryKeyConstraint("id"),
)
for name, columns in (
("ix_addresses_contact_redirects_tenant_id", ["tenant_id"]),
("ix_addresses_contact_redirects_source_contact_id", ["source_contact_id"]),
("ix_addresses_contact_redirects_target_contact_id", ["target_contact_id"]),
("ix_addresses_contact_redirects_merge_record_id", ["merge_record_id"]),
("ix_addresses_contact_redirects_ended_at", ["ended_at"]),
("ix_addresses_contact_redirects_target", ["tenant_id", "target_contact_id", "ended_at"]),
):
op.create_index(name, "addresses_contact_redirects", columns)
op.create_index(
"uq_addresses_contact_redirects_active_source",
"addresses_contact_redirects",
["tenant_id", "source_contact_id"],
unique=True,
sqlite_where=sa.text("ended_at IS NULL"),
postgresql_where=sa.text("ended_at IS NULL"),
)
op.create_table(
"addresses_contact_field_provenance",
sa.Column("id", sa.String(length=36), nullable=False),
sa.Column("tenant_id", sa.String(length=36), nullable=True),
sa.Column("contact_id", sa.String(length=36), nullable=False),
sa.Column("field_path", sa.String(length=255), nullable=False),
sa.Column("value", sa.JSON(), nullable=True),
sa.Column("source_kind", sa.String(length=40), nullable=False),
sa.Column("source_ref", sa.String(length=1000), nullable=True),
sa.Column("source_revision", sa.String(length=255), nullable=True),
sa.Column("precedence", sa.Integer(), nullable=False),
sa.Column("selected", sa.Boolean(), nullable=False),
sa.Column("reason_code", sa.String(length=120), nullable=False),
sa.Column("explanation", sa.Text(), nullable=True),
sa.Column("visibility", sa.String(length=30), nullable=False),
sa.Column("merge_record_id", sa.String(length=36), nullable=True),
sa.Column("created_by_account_id", sa.String(length=36), nullable=True),
sa.Column("metadata", sa.JSON(), nullable=False, server_default=_JSON_OBJECT),
sa.Column("created_at", sa.DateTime(timezone=True), nullable=False),
sa.Column("updated_at", sa.DateTime(timezone=True), nullable=False),
sa.ForeignKeyConstraint(["contact_id"], ["addresses_contacts.id"], ondelete="CASCADE"),
sa.ForeignKeyConstraint(["merge_record_id"], ["addresses_contact_merge_records.id"], ondelete="SET NULL"),
sa.PrimaryKeyConstraint("id"),
)
for name, columns in (
("ix_addresses_contact_field_provenance_tenant_id", ["tenant_id"]),
("ix_addresses_contact_field_provenance_contact_id", ["contact_id"]),
("ix_addresses_contact_field_provenance_field_path", ["field_path"]),
("ix_addresses_contact_field_provenance_selected", ["selected"]),
("ix_addresses_contact_field_provenance_merge_record_id", ["merge_record_id"]),
("ix_addresses_contact_field_provenance_created_by_account_id", ["created_by_account_id"]),
("ix_addresses_field_provenance_contact", ["contact_id", "field_path", "created_at"]),
("ix_addresses_field_provenance_selected", ["tenant_id", "contact_id", "selected"]),
):
op.create_index(name, "addresses_contact_field_provenance", columns)
def downgrade() -> None:
op.drop_table("addresses_contact_field_provenance")
op.drop_table("addresses_contact_redirects")
op.drop_table("addresses_contact_merge_records")
op.drop_table("addresses_contact_point_quality_decisions")
with op.batch_alter_table("addresses_contact_postal_addresses") as batch:
batch.drop_column("provenance")
batch.drop_column("normalized_value")
batch.drop_column("original_value")
with op.batch_alter_table("addresses_contact_phones") as batch:
batch.drop_index("ix_addresses_contact_phones_normalized_phone")
batch.drop_column("provenance")
batch.drop_column("normalized_phone")
batch.drop_column("original_phone")
with op.batch_alter_table("addresses_contact_emails") as batch:
batch.drop_index("ix_addresses_contact_emails_normalized_email")
batch.drop_column("provenance")
batch.drop_column("normalized_email")
batch.drop_column("original_email")
def _normalized_text(value: object) -> str | None:
if value is None:
return None
normalized = " ".join(str(value).strip().casefold().split())
return normalized or None
def _normalized_phone(value: str) -> str:
prefix = "+" if value.strip().startswith("+") else ""
return prefix + re.sub(r"\D", "", value)