Skip to content

popoto.backends.postgres.validity

popoto.backends.postgres.validity

Validity intervals and supersession on Postgres (#759 M3; plan §2 group D).

This module is the Postgres half of :class:~popoto.fields.validity_field. ValidityField and :class:~popoto.fields.supersession.SupersessionProtocol: the interval columns, supersede (SUPERSEDE_LUA phase for phase), chain (WITH RECURSIVE), the open-claim pointers, the exclusion rule every retrieval gate applies, and the field_call adapters the field layer reads through. :class:PostgresValidityOps is mixed into :class:~popoto.backends.postgres.PostgresBackend (through :class:~.memory.PostgresMemoryOps).

Storage (plan §3, "Field → column type mapping"; a recorded departure)

The plan proposed <f> tstzrange. It cannot hold what the Redis indexes hold, measured before this module was written (PostgreSQL 18.6):

  • Precision. timestamptz keeps microseconds, so the epoch 1700000000.1234567 comes back 1700000000.123457: a close instant one ulp after an as_of reads as one at it, and the gate's bound is no longer bit-exact (invalid_at <= as_of flips). The Redis score is the full double -- the M2a clock decision, for the same reason.
  • A close at the record's own start. SUPERSEDE_LUA lets a record close at exactly its valid_from (the check is close < start), and tstzrange(t, t) is the canonical empty range, which keeps neither bound: the start the record had, and the close the chain recorded, are both lost. A lower > upper pair raises outright.
  • Partial intervals are not the reason, though they need care: a member can sit in one index and not the other (an import_state shape), and the exclusion rule treats an absent end as "never excludes" and +inf as "open" -- they differ at as_of = +inf. A range can say that (upper_inf('[t,)') is true, upper_inf('[t,infinity)') false), so this is no argument against tstzrange; it is how the columns below spell it (NULL vs 'Infinity').

So the interval is three double precision columns beside the field's own column, which keeps the declared value as the Redis hash does:

  • <f>__valid_from -- the valid_from ZSET score
  • <f>__invalid_at -- the invalid_at ZSET score ('Infinity' = open)
  • <f>__ingested_at -- the ingested_at ZSET score
  • <f>__supersedes / <f>__superseded_by (text) -- the chain:rev / chain:fwd HASH entries for this record

NULL is "absent from that index". B-trees on the two gate columns serve filter(validity__as_of=…); a supersede rewrites them only on close.

The open-claim pointers are a companion table <table>__<f>__open (digest text PRIMARY KEY, member text REFERENCES <table>(_pk) ON DELETE CASCADE) with an index on member, not the plan's partial UNIQUE on an identity column: a record can be the open claim of several identities at once, and an invalidate leaves the pointer naming the record it closed (the next supersede on that identity reads it as "already closed"), neither of which a per-row identity column can say. The cascade is on_delete's pointer cleanup, matched on the exact key, so a never takes ab's pointer with it (#750's drop_validity prefix over-match). The table is created on first use, like popoto_recall_proposal.

Lock order (plan §2 supersede, §6 TD-2)

Every supersede on one (model, field) takes pg_advisory_xact_lock(hashtext('popoto:validity:<schema>.<table>.<f>')) first -- Redis's single thread, per model and field -- and then locks the rows it reads FOR UPDATE in _pk order (COLLATE "C"). Two supersedes therefore never interleave: two pointer writers, two explicit writers, a pointer and an explicit writer, and the crossing chains that deadlocked the POC (#750 B1: d1 -> X superseded by Y while d2 -> Y is superseded by X) all run one after the other.

That is the head of the backend's one lock order (plan §6, TD-2; M2b's :meth:~popoto.backends.postgres.PostgresBackend._record_locked): the (model, field) lock, then the record-key advisory locks in _pk byte order, then the row locks in _pk order. Every validity writer follows it: supersede (record keys of the successor and the incumbent before its FOR UPDATE), save_and_supersede / save_and_invalidate (all of the supersede's locks before the save, through the lock adapter -- the save would otherwise take the successor's key and row first and cross a concurrent supersede naming it), import_state (a pointer writer), and ObservationProtocol's batch (which locks a contradicted record's successor with the batch). A plain save does not take the field lock; it takes its record key and meets a supersede there, and its upsert re-reads the row it waited on, so it cannot reopen a record the supersede closed (plan Race 2). What remains is a caller's own transaction() that locks a record before a supersede on it: Postgres detects that cycle and one side is aborted -- an owned transaction retries; in the caller's transaction it is :class:~popoto.backends.BackendRetryableError at once.

Never imports redis.

NOT_HANDLED = object() module-attribute

What :meth:PostgresValidityOps._validity_field_call returns for an op it does not register.

PostgresValidityOps

supersede, chain and the validity field_call adapters for :class:~popoto.backends.postgres.PostgresBackend. Relies on the backend's _table, _run and _atomically.

Source code in src/popoto/backends/postgres/validity.py
 440
 441
 442
 443
 444
 445
 446
 447
 448
 449
 450
 451
 452
 453
 454
 455
 456
 457
 458
 459
 460
 461
 462
 463
 464
 465
 466
 467
 468
 469
 470
 471
 472
 473
 474
 475
 476
 477
 478
 479
 480
 481
 482
 483
 484
 485
 486
 487
 488
 489
 490
 491
 492
 493
 494
 495
 496
 497
 498
 499
 500
 501
 502
 503
 504
 505
 506
 507
 508
 509
 510
 511
 512
 513
 514
 515
 516
 517
 518
 519
 520
 521
 522
 523
 524
 525
 526
 527
 528
 529
 530
 531
 532
 533
 534
 535
 536
 537
 538
 539
 540
 541
 542
 543
 544
 545
 546
 547
 548
 549
 550
 551
 552
 553
 554
 555
 556
 557
 558
 559
 560
 561
 562
 563
 564
 565
 566
 567
 568
 569
 570
 571
 572
 573
 574
 575
 576
 577
 578
 579
 580
 581
 582
 583
 584
 585
 586
 587
 588
 589
 590
 591
 592
 593
 594
 595
 596
 597
 598
 599
 600
 601
 602
 603
 604
 605
 606
 607
 608
 609
 610
 611
 612
 613
 614
 615
 616
 617
 618
 619
 620
 621
 622
 623
 624
 625
 626
 627
 628
 629
 630
 631
 632
 633
 634
 635
 636
 637
 638
 639
 640
 641
 642
 643
 644
 645
 646
 647
 648
 649
 650
 651
 652
 653
 654
 655
 656
 657
 658
 659
 660
 661
 662
 663
 664
 665
 666
 667
 668
 669
 670
 671
 672
 673
 674
 675
 676
 677
 678
 679
 680
 681
 682
 683
 684
 685
 686
 687
 688
 689
 690
 691
 692
 693
 694
 695
 696
 697
 698
 699
 700
 701
 702
 703
 704
 705
 706
 707
 708
 709
 710
 711
 712
 713
 714
 715
 716
 717
 718
 719
 720
 721
 722
 723
 724
 725
 726
 727
 728
 729
 730
 731
 732
 733
 734
 735
 736
 737
 738
 739
 740
 741
 742
 743
 744
 745
 746
 747
 748
 749
 750
 751
 752
 753
 754
 755
 756
 757
 758
 759
 760
 761
 762
 763
 764
 765
 766
 767
 768
 769
 770
 771
 772
 773
 774
 775
 776
 777
 778
 779
 780
 781
 782
 783
 784
 785
 786
 787
 788
 789
 790
 791
 792
 793
 794
 795
 796
 797
 798
 799
 800
 801
 802
 803
 804
 805
 806
 807
 808
 809
 810
 811
 812
 813
 814
 815
 816
 817
 818
 819
 820
 821
 822
 823
 824
 825
 826
 827
 828
 829
 830
 831
 832
 833
 834
 835
 836
 837
 838
 839
 840
 841
 842
 843
 844
 845
 846
 847
 848
 849
 850
 851
 852
 853
 854
 855
 856
 857
 858
 859
 860
 861
 862
 863
 864
 865
 866
 867
 868
 869
 870
 871
 872
 873
 874
 875
 876
 877
 878
 879
 880
 881
 882
 883
 884
 885
 886
 887
 888
 889
 890
 891
 892
 893
 894
 895
 896
 897
 898
 899
 900
 901
 902
 903
 904
 905
 906
 907
 908
 909
 910
 911
 912
 913
 914
 915
 916
 917
 918
 919
 920
 921
 922
 923
 924
 925
 926
 927
 928
 929
 930
 931
 932
 933
 934
 935
 936
 937
 938
 939
 940
 941
 942
 943
 944
 945
 946
 947
 948
 949
 950
 951
 952
 953
 954
 955
 956
 957
 958
 959
 960
 961
 962
 963
 964
 965
 966
 967
 968
 969
 970
 971
 972
 973
 974
 975
 976
 977
 978
 979
 980
 981
 982
 983
 984
 985
 986
 987
 988
 989
 990
 991
 992
 993
 994
 995
 996
 997
 998
 999
1000
1001
1002
1003
1004
1005
1006
1007
1008
1009
1010
1011
1012
1013
1014
1015
1016
1017
1018
1019
1020
1021
1022
1023
1024
1025
1026
1027
1028
class PostgresValidityOps:
    """``supersede``, ``chain`` and the validity ``field_call`` adapters for
    :class:`~popoto.backends.postgres.PostgresBackend`. Relies on the
    backend's ``_table``, ``_run`` and ``_atomically``."""

    # Provided by PostgresBackend / PostgresMemoryOps.
    schema: str
    _table: Callable[..., TableSpec]
    _run: Callable[..., Any]
    _atomically: Callable[..., Any]
    _record_locked: Callable[..., tuple[str, list[Any]]]

    # -- plumbing ---------------------------------------------------------------

    def _validity_field(self, spec: ModelSpec, field: str) -> None:
        fs = spec.fields.get(field)
        if fs is None or fs.kind != VALIDITY_KIND:
            raise BackendCapabilityError(f"{spec.name}.{field} is not a ValidityField")

    def _validity_lock(
        self, ts: TableSpec, fields: Iterable[str], uow: Optional[UnitOfWork]
    ) -> None:
        """The ``(model, field)`` advisory locks, in name order: the first
        lock any supersede on the field takes."""
        for field in sorted(fields):
            self._run(
                "SELECT pg_advisory_xact_lock(hashtext(%s))",
                [f"popoto:validity:{ts.schema}.{ts.table}.{field}"],
                uow=uow,
                write=True,
            )

    # -- D. supersede -------------------------------------------------------------

    def supersede(
        self,
        spec: ModelSpec,
        field: str,
        *,
        successor: Optional[RecordId],
        incumbent: Optional[RecordId],
        identity: Optional[str],
        mode: str,
        valid_from: Optional[float],
        invalid_at: Optional[float],
        now: float,
        uow: Optional[UnitOfWork] = None,
        ingested_at: Optional[float] = None,
        assert_valid_from: bool = False,
        **options: Any,
    ) -> Optional[str]:
        """``SUPERSEDE_LUA``, phase for phase, in one transaction.

        ``mode`` is ``'open'`` (open ``successor``; close nothing),
        ``'supersede'`` (close the incumbent -- ``incumbent``, else the record
        ``identity``'s pointer names -- chain it to ``successor`` and open
        that) or ``'invalidate'`` (as ``'supersede'``; without a successor
        nothing is chained or opened). ``invalid_at`` is the close instant,
        and every instant the caller leaves out is ``now``.

        Phases, after the ``(model, field)`` lock and the ``_pk``-ordered row
        locks: **validation** (reads and refusals only) -- a successor that
        does not exist, an *asserted* incumbent that does not exist (a
        pointer-resolved one that does not is "no incumbent"), a close
        before the incumbent's own start, an asserted ``valid_from`` that
        disagrees with the stored one -- raises the typed error the script's
        reply maps to, with the same text, having written nothing; then
        **mutation** -- close (only an open incumbent: closing is
        idempotent), both chain links, the ``NX`` open of an open successor,
        and the pointer repoint. Returns the closed record's key, or ``None``.

        A record that does not exist holds no interval here (the interval is
        its row), so mode ``'open'`` on an absent ``successor`` writes
        nothing, where the script's ``ZADD NX`` would index a member with no
        record (only a direct ``execute_supersede`` call can ask for that).
        """
        from ...fields.validity_field import (
            CLOSE_BEFORE_START_ERROR,
            MEMBER_ABSENT_ERROR,
            VALID_FROM_CONFLICT_ERROR,
            VALID_MODES,
        )

        self._validity_field(spec, field)
        if mode not in VALID_MODES:
            raise ValueError(
                f"ValidityField mode must be one of {sorted(VALID_MODES)}, got {mode!r}"
            )
        clock = float(now)
        start = clock if valid_from is None else float(valid_from)
        ingest = clock if ingested_at is None else float(ingested_at)
        close_at = clock if invalid_at is None else float(invalid_at)
        for label, value in (
            ("now", clock),
            ("valid_from", start),
            ("ingested_at", ingest),
            ("close_at", close_at),
        ):
            if math.isnan(value):
                # Redis: the script's ZADD replies ``value is not a valid
                # float`` (a ResponseError; with only valid_from NaN, after
                # the incumbent's close -- #778). Here: the same text, a
                # ValueError, and nothing written.
                raise ValueError(f"{NAN_SCORE_ERROR} ({label} is NaN)")
        new = successor.canonical if successor is not None else ""
        named_old = incumbent.canonical if incumbent is not None else ""
        ts = self._table(spec, write=True)
        vf, ia, ig = (_col(field, s) for s in (VALID_FROM, INVALID_AT, INGESTED_AT))
        sup, supby = _col(field, SUPERSEDES), _col(field, SUPERSEDED_BY)

        def work(tx: UnitOfWork) -> Optional[str]:
            self._validity_lock(ts, [field], tx)
            old = named_old
            if mode != "open" and not old and identity:
                old = self._pointer_get(ts, field, identity, uow=tx) or ""
            keys = sorted({k for k in (new, old) if k}, key=_sort_key)
            rows: dict[str, tuple[Any, ...]] = {}
            if keys:
                # The backend's one lock order (plan §6, TD-2): the
                # (model, field) lock above, then the record-key locks in
                # ``_pk`` byte order, then the row locks in ``_pk`` order.
                sql, params = self._record_locked(
                    ts,
                    keys,
                    f'SELECT "_pk", {vf}, {ia}, {ig}, {sup}, {supby} '
                    f'FROM {ts.qualified} WHERE "_pk" = ANY(%s::text[])'
                    f"{_and_live(ts)} "
                    'ORDER BY "_pk" COLLATE "C" FOR UPDATE',
                    [keys],
                )
                found, _ = self._run(sql, params, uow=tx, write=True)
                rows = {row[0]: tuple(row[1:]) for row in found}

            # -- VALIDATION PHASE: reads and refusals only --------------------
            will_close = False
            if mode != "open":
                if new and new not in rows:
                    raise _refuse(f"{MEMBER_ABSENT_ERROR} successor {new}")
                if old and old not in rows:
                    if named_old:
                        raise _refuse(f"{MEMBER_ABSENT_ERROR} incumbent {old}")
                    old = ""  # a pointer naming a deleted record: a hint
                if old and old != new:
                    old_start, old_close = rows[old][0], rows[old][1]
                    if old_close is not None and old_close == INF:
                        if old_start is not None and close_at < old_start:
                            raise _refuse(CLOSE_BEFORE_START_ERROR)
                        will_close = True
            if new and assert_valid_from and new in rows:
                stored = rows[new][0]
                if stored is not None and stored != start:
                    raise _refuse(
                        f"{VALID_FROM_CONFLICT_ERROR} {_lua_number(stored)} "
                        f"{_lua_number(start)}"
                    )

            # -- MUTATION PHASE: every check above has passed -----------------
            updates: dict[str, dict[str, Any]] = {}
            if will_close:
                updates[old] = {"ia": close_at}
                if new:
                    updates[old]["supby"] = new
                    updates.setdefault(new, {})["sup"] = old
            new_open = False
            if new and new in rows:
                state = rows[new]
                if state[1] is None or state[1] == INF:
                    new_open = True
                    target = updates.setdefault(new, {})
                    target["vf"] = state[0] if state[0] is not None else start
                    target["ig"] = state[2] if state[2] is not None else ingest
                    target["ia"] = INF if state[1] is None else state[1]
            if updates:
                self._write_states(ts, field, rows, updates, tx)
            if new_open and identity:
                self._run(
                    f"INSERT INTO {pointer_table(ts, field)} (digest, member) "
                    "VALUES (%s, %s) ON CONFLICT (digest) DO UPDATE "
                    "SET member = EXCLUDED.member",
                    [identity, new],
                    uow=tx,
                    write=True,
                )
            return old if will_close else None

        return self._atomically(work, uow=uow)

    def _write_states(
        self,
        ts: TableSpec,
        field: str,
        rows: Mapping[str, tuple[Any, ...]],
        updates: Mapping[str, Mapping[str, Any]],
        uow: UnitOfWork,
    ) -> None:
        """One ``UPDATE … FROM (VALUES …)`` with each locked row's final
        state; columns a row does not change keep the value read under its
        lock."""
        order = ("vf", "ia", "ig", "sup", "supby")
        values = []
        params: list[Any] = []
        for pk in sorted(updates, key=_sort_key):
            current = dict(zip(order, rows[pk]))
            current.update(updates[pk])
            values.append(
                "(%s::text, %s::float8, %s::float8, %s::float8, %s::text, %s::text)"
            )
            params.extend([pk] + [current[k] for k in order])
        cols = {
            "vf": quote_ident(field + VALID_FROM),
            "ia": quote_ident(field + INVALID_AT),
            "ig": quote_ident(field + INGESTED_AT),
            "sup": quote_ident(field + SUPERSEDES),
            "supby": quote_ident(field + SUPERSEDED_BY),
        }
        sets = ", ".join(f"{cols[k]} = v.{k}" for k in order)
        self._run(
            f'UPDATE {ts.qualified} AS t SET {sets}, "_updated_at" = now() '
            f"FROM (VALUES {', '.join(values)}) AS v(pk, {', '.join(order)}) "
            f'WHERE t."_pk" = v.pk',
            params,
            uow=uow,
            write=True,
        )

    def _pointer_get(
        self,
        ts: TableSpec,
        field: str,
        digest: str,
        *,
        uow: Optional[UnitOfWork] = None,
    ) -> Optional[str]:
        rows, _ = self._run(
            f"SELECT member FROM {pointer_table(ts, field)} WHERE digest = %s",
            [digest],
            uow=uow,
        )
        return rows[0][0] if rows else None

    def _supersede_locks(
        self,
        spec: ModelSpec,
        field: str,
        successor: str,
        incumbent: str,
        identity: str,
        *,
        uow: Optional[UnitOfWork] = None,
    ) -> None:
        """Take, inside the caller's transaction ``uow``, every lock a
        ``supersede`` of ``successor`` over ``incumbent`` (else the record
        ``identity``'s pointer names) will take, in the backend's order: the
        ``(model, field)`` lock, then the record-key locks in ``_pk`` byte
        order. ``save_and_*`` takes them before its save, which would
        otherwise take the successor's key and row first and so cross a
        concurrent supersede that holds the field lock and names it (#777
        review). Advisory locks stack, so the supersede retaking them is a
        no-op."""
        from . import _pg_uow

        if _pg_uow(uow) is None:
            raise BackendCapabilityError(
                "ValidityField 'lock' runs inside a transaction() only"
            )
        ts = self._table(spec, write=True)
        self._validity_lock(ts, [field], uow)
        old = incumbent
        if not old and identity:
            old = self._pointer_get(ts, field, identity, uow=uow) or ""
        keys = [k for k in (successor, old) if k]
        if keys:
            sql, params = self._record_locked(ts, keys, "SELECT 1", [])
            self._run(sql, params, uow=uow, write=True)

    # -- D. chain -----------------------------------------------------------------

    def chain(self, spec: ModelSpec, field: str, id: RecordId) -> list[RecordId]:
        """The supersession chain through ``id``, oldest first, in one
        ``WITH RECURSIVE``: back along ``supersedes``, then forward along
        ``superseded_by``. ``[]`` when ``id`` has no ``valid_from`` (unsaved,
        or never opened). A walk stops at a missing link, at a record it has
        already visited (a cycle -- the forward walk counts the backward
        walk's records as visited, as the Redis walk's shared ``seen`` set
        does), and at a link naming a record with no ``valid_from`` (a hard
        delete leaves the neighbour's link naming it)."""
        self._validity_field(spec, field)
        ts = self._table(spec)
        vf = _col(field, VALID_FROM, "n")
        sup = _col(field, SUPERSEDES, "s")
        supby = _col(field, SUPERSEDED_BY, "s")
        anchor_vf = _col(field, VALID_FROM)
        t = ts.qualified
        # M5: an expired record is not on the chain -- the walk stops at it
        # as at a hard delete.
        live_n = _and_live(ts, "n")
        sql = (
            "WITH RECURSIVE "
            f'anchor AS (SELECT "_pk" AS pk FROM {t} WHERE "_pk" = %s '
            f"AND {anchor_vf} IS NOT NULL{_and_live(ts)}), "
            "older(pk, depth, path) AS ("
            f'SELECT n."_pk", 1, ARRAY[a.pk, n."_pk"] FROM anchor a '
            f'JOIN {t} s ON s."_pk" = a.pk JOIN {t} n ON n."_pk" = {sup} '
            f'WHERE {vf} IS NOT NULL AND n."_pk" <> a.pk{live_n} '
            f'UNION ALL SELECT n."_pk", o.depth + 1, o.path || n."_pk" '
            f'FROM older o JOIN {t} s ON s."_pk" = o.pk '
            f'JOIN {t} n ON n."_pk" = {sup} '
            f'WHERE {vf} IS NOT NULL AND n."_pk" <> ALL(o.path){live_n}), '
            "seen AS (SELECT a.pk AS pk FROM anchor a UNION ALL "
            "SELECT pk FROM older), "
            "newer(pk, depth, path) AS ("
            f'SELECT n."_pk", 1, '
            '(SELECT array_agg(pk) FROM seen) || n."_pk" FROM anchor a '
            f'JOIN {t} s ON s."_pk" = a.pk JOIN {t} n ON n."_pk" = {supby} '
            f'WHERE {vf} IS NOT NULL AND n."_pk" <> ALL(SELECT pk FROM seen){live_n} '
            f'UNION ALL SELECT n."_pk", w.depth + 1, w.path || n."_pk" '
            f'FROM newer w JOIN {t} s ON s."_pk" = w.pk '
            f'JOIN {t} n ON n."_pk" = {supby} '
            f'WHERE {vf} IS NOT NULL AND n."_pk" <> ALL(w.path){live_n}) '
            "SELECT pk, -depth FROM older UNION ALL SELECT pk, 0 FROM anchor "
            "UNION ALL SELECT pk, depth FROM newer ORDER BY 2"
        )
        rows, _ = self._run(sql, [id.canonical])
        return [RecordId(spec.name, (), pk, native=pk) for pk, _depth in rows]

    # -- field_call adapters ---------------------------------------------------

    def _validity_field_call(
        self,
        spec: ModelSpec,
        field: str,
        op: str,
        args: tuple[Any, ...],
        kwargs: dict[str, Any],
        uow: Optional[UnitOfWork],
    ) -> Any:
        """The ``ValidityField`` adapters; :data:`NOT_HANDLED` when ``op`` is
        not one of them."""
        handlers: dict[str, Callable[..., Any]] = {
            "interval": self._interval_of,
            "members": self._interval_members,
            "pointers": self._pointers_naming,
            "pointer": self._pointer_of,
            "export": self._export_state,
            "import": self._import_state,
            "dump": self._dump_state,
            "lock": self._supersede_locks,
        }
        handler = handlers.get(op)
        if handler is None:
            return NOT_HANDLED
        self._validity_field(spec, field)
        return handler(spec, field, *args, uow=uow, **kwargs)

    def _interval_of(
        self,
        spec: ModelSpec,
        field: str,
        id: RecordId,
        *,
        uow: Optional[UnitOfWork] = None,
    ) -> Optional[dict[str, Any]]:
        """The record's ``valid_from``/``invalid_at``/``ingested_at`` and
        chain links (``None`` for a record that does not exist): the five
        ZSET scores and HASH entries Redis keeps for the member."""
        ts = self._table(spec)
        cols = ", ".join(_col(field, s) for s, _t in VALIDITY_SUFFIXES)
        rows, _ = self._run(
            f'SELECT {cols} FROM {ts.qualified} WHERE "_pk" = %s{_and_live(ts)}',
            [id.canonical],
            uow=uow,
        )
        if not rows:
            return None
        valid_from, invalid_at, ingested_at, supersedes, superseded_by = rows[0]
        return {
            "valid_from": valid_from,
            "invalid_at": invalid_at,
            "ingested_at": ingested_at,
            "supersedes": supersedes,
            "superseded_by": superseded_by,
        }

    def _dump_state(
        self, spec: ModelSpec, field: str, *, uow: Optional[UnitOfWork] = None
    ) -> dict[str, list[tuple[str, Any]]]:
        """Every interval score, chain link and pointer of the field, as the
        six Redis keys would list them: ``{index: sorted [(member, value)]}``
        for ``valid_from``, ``invalid_at``, ``ingested_at``, ``chain_fwd``
        and ``chain_rev``, and ``pointers`` as ``[(digest, member)]``. For
        tests and the parity probe (an admin read: it scans the table)."""
        ts = self._table(spec)
        cols = ", ".join(_col(field, s) for s, _t in VALIDITY_SUFFIXES)
        rows, _ = self._run(f'SELECT "_pk", {cols} FROM {ts.qualified}', [], uow=uow)
        out: dict[str, list[tuple[str, Any]]] = {
            "valid_from": [],
            "invalid_at": [],
            "ingested_at": [],
            "chain_fwd": [],
            "chain_rev": [],
        }
        for pk, vf, ia, ig, sup, supby in rows:
            for name, value in (
                ("valid_from", vf),
                ("invalid_at", ia),
                ("ingested_at", ig),
                ("chain_fwd", supby),
                ("chain_rev", sup),
            ):
                if value is not None:
                    out[name].append((pk, value))
        pointers, _ = self._run(
            f"SELECT digest, member FROM {pointer_table(ts, field)}", [], uow=uow
        )
        out["pointers"] = [(digest, member) for digest, member in pointers]
        return {name: sorted(items) for name, items in out.items()}

    def _interval_members(
        self,
        spec: ModelSpec,
        field: str,
        as_of: float,
        *,
        select: str,
        uow: Optional[UnitOfWork] = None,
    ) -> set[str]:
        """``resolve_valid_keys`` (``select="valid"``: ``valid_from <= t AND
        invalid_at > t``, both present) or ``resolve_excluded_keys``
        (``select="excluded"``: the exclusion rule)."""
        ts = self._table(spec)
        t = range_bound(as_of)
        if select == "valid":
            clause = validity_cond_sql(field, (t, True))
        elif select == "excluded":
            clause = excluded_sql(field, t, alias="")
        else:
            raise ValueError(f"select must be 'valid' or 'excluded', got {select!r}")
        rows, _ = self._run(
            f'SELECT "_pk" FROM {ts.qualified} WHERE {clause}{_and_live(ts)}',
            [],
            uow=uow,
        )
        return {row[0] for row in rows}

    def _pointers_naming(
        self,
        spec: ModelSpec,
        field: str,
        id: RecordId,
        *,
        uow: Optional[UnitOfWork] = None,
    ) -> list[str]:
        """The identity digests whose open-claim pointer names ``id``."""
        ts = self._table(spec)
        expired = expired_pks_sql(ts)
        rows, _ = self._run(
            f"SELECT digest FROM {pointer_table(ts, field)} WHERE member = %s "
            + (f"AND member NOT IN {expired} " if expired else "")
            + 'ORDER BY digest COLLATE "C"',
            [id.canonical],
            uow=uow,
        )
        return [row[0] for row in rows]

    def _pointer_of(
        self,
        spec: ModelSpec,
        field: str,
        digest: str,
        *,
        uow: Optional[UnitOfWork] = None,
    ) -> Optional[str]:
        """The record ``digest``'s open-claim pointer names, or ``None`` --
        also when it names a record that has expired (M5)."""
        ts = self._table(spec)
        member = self._pointer_get(ts, field, digest, uow=uow)
        if member is not None and ts.ttl:
            live, _ = self._run(
                f'SELECT 1 FROM {ts.qualified} WHERE "_pk" = %s{_and_live(ts)}',
                [member],
                uow=uow,
            )
            if not live:
                return None
        return member

    def _export_state(
        self,
        spec: ModelSpec,
        field: str,
        id: RecordId,
        *,
        uow: Optional[UnitOfWork] = None,
    ) -> Optional[dict[str, Any]]:
        """``ValidityField.export_state``'s dict, from the row."""
        state = self._interval_of(spec, field, id, uow=uow)
        if state is None or (
            state["valid_from"] is None
            and state["invalid_at"] is None
            and state["ingested_at"] is None
        ):
            return None
        invalid_at = state["invalid_at"]
        return {
            "valid_from": state["valid_from"],
            "invalid_at": (
                "+inf" if invalid_at is None or invalid_at == INF else invalid_at
            ),
            "ingested_at": state["ingested_at"],
            "chain_fwd": state["superseded_by"],
            "chain_rev": state["supersedes"],
            "open_pointers": self._pointers_naming(spec, field, id, uow=uow),
        }

    def _import_state(
        self,
        spec: ModelSpec,
        field: str,
        id: RecordId,
        state: Mapping[str, Any],
        *,
        uow: Optional[UnitOfWork] = None,
    ) -> None:
        """``ValidityField.import_state``: overwrite (never ``NX``) the
        carried scores and links, and repoint each carried identity."""
        ts = self._table(spec, write=True)
        sets: list[str] = []
        params: list[Any] = []
        for key, suffix in (
            ("valid_from", VALID_FROM),
            ("invalid_at", INVALID_AT),
            ("ingested_at", INGESTED_AT),
        ):
            value = state.get(key)
            if value is None:
                continue
            if key == "invalid_at" and value == "+inf":
                value = INF
            sets.append(f"{_col(field, suffix)} = %s")
            params.append(float(value))
        for key, suffix in (("chain_fwd", SUPERSEDED_BY), ("chain_rev", SUPERSEDES)):
            if state.get(key):
                sets.append(f"{_col(field, suffix)} = %s")
                params.append(str(state[key]))
        digests = [str(d) for d in state.get("open_pointers") or []]
        if not sets and not digests:
            return

        def work(tx: UnitOfWork) -> None:
            # A pointer writer, so the supersede lock order: the
            # (model, field) lock, the record-key lock, then the row (the
            # UPDATE's, or the pointer's foreign-key check).
            self._validity_lock(ts, [field], tx)
            if sets:
                sql = (
                    f"UPDATE {ts.qualified} SET {', '.join(sets)}, "
                    '"_updated_at" = now() WHERE "_pk" = %s'
                )
                args: list[Any] = params + [id.canonical]
            else:
                sql, args = "SELECT 1", []
            sql, args = self._record_locked(ts, [id.canonical], sql, args)
            self._run(sql, args, uow=tx, write=True)
            if digests:
                self._run(
                    f"INSERT INTO {pointer_table(ts, field)} (digest, member) "
                    "SELECT d, %s FROM unnest(%s::text[]) AS d "
                    "ON CONFLICT (digest) DO UPDATE SET member = EXCLUDED.member",
                    [id.canonical, digests],
                    uow=tx,
                    write=True,
                )

        import psycopg.errors

        try:
            self._atomically(work, uow=uow)
        except psycopg.errors.ForeignKeyViolation as e:
            # Redis stores a pointer naming a record that is not there (a plain
            # ``SET``) and reads it later as "no incumbent"; the pointer table's
            # foreign key refuses it. Typed as the validity layer's "member
            # absent" error, not a raw driver error (M3 review).
            from ...fields.validity_field import ValidityMemberAbsentError

            raise ValidityMemberAbsentError(
                f"import_state: {id.canonical!r} is not stored, so a carried "
                "open-claim pointer cannot name it (save the record before "
                "importing its validity state)"
            ) from e

supersede(spec, field, *, successor, incumbent, identity, mode, valid_from, invalid_at, now, uow=None, ingested_at=None, assert_valid_from=False, **options)

SUPERSEDE_LUA, phase for phase, in one transaction.

mode is 'open' (open successor; close nothing), 'supersede' (close the incumbent -- incumbent, else the record identity's pointer names -- chain it to successor and open that) or 'invalidate' (as 'supersede'; without a successor nothing is chained or opened). invalid_at is the close instant, and every instant the caller leaves out is now.

Phases, after the (model, field) lock and the _pk-ordered row locks: validation (reads and refusals only) -- a successor that does not exist, an asserted incumbent that does not exist (a pointer-resolved one that does not is "no incumbent"), a close before the incumbent's own start, an asserted valid_from that disagrees with the stored one -- raises the typed error the script's reply maps to, with the same text, having written nothing; then mutation -- close (only an open incumbent: closing is idempotent), both chain links, the NX open of an open successor, and the pointer repoint. Returns the closed record's key, or None.

A record that does not exist holds no interval here (the interval is its row), so mode 'open' on an absent successor writes nothing, where the script's ZADD NX would index a member with no record (only a direct execute_supersede call can ask for that).

Source code in src/popoto/backends/postgres/validity.py
def supersede(
    self,
    spec: ModelSpec,
    field: str,
    *,
    successor: Optional[RecordId],
    incumbent: Optional[RecordId],
    identity: Optional[str],
    mode: str,
    valid_from: Optional[float],
    invalid_at: Optional[float],
    now: float,
    uow: Optional[UnitOfWork] = None,
    ingested_at: Optional[float] = None,
    assert_valid_from: bool = False,
    **options: Any,
) -> Optional[str]:
    """``SUPERSEDE_LUA``, phase for phase, in one transaction.

    ``mode`` is ``'open'`` (open ``successor``; close nothing),
    ``'supersede'`` (close the incumbent -- ``incumbent``, else the record
    ``identity``'s pointer names -- chain it to ``successor`` and open
    that) or ``'invalidate'`` (as ``'supersede'``; without a successor
    nothing is chained or opened). ``invalid_at`` is the close instant,
    and every instant the caller leaves out is ``now``.

    Phases, after the ``(model, field)`` lock and the ``_pk``-ordered row
    locks: **validation** (reads and refusals only) -- a successor that
    does not exist, an *asserted* incumbent that does not exist (a
    pointer-resolved one that does not is "no incumbent"), a close
    before the incumbent's own start, an asserted ``valid_from`` that
    disagrees with the stored one -- raises the typed error the script's
    reply maps to, with the same text, having written nothing; then
    **mutation** -- close (only an open incumbent: closing is
    idempotent), both chain links, the ``NX`` open of an open successor,
    and the pointer repoint. Returns the closed record's key, or ``None``.

    A record that does not exist holds no interval here (the interval is
    its row), so mode ``'open'`` on an absent ``successor`` writes
    nothing, where the script's ``ZADD NX`` would index a member with no
    record (only a direct ``execute_supersede`` call can ask for that).
    """
    from ...fields.validity_field import (
        CLOSE_BEFORE_START_ERROR,
        MEMBER_ABSENT_ERROR,
        VALID_FROM_CONFLICT_ERROR,
        VALID_MODES,
    )

    self._validity_field(spec, field)
    if mode not in VALID_MODES:
        raise ValueError(
            f"ValidityField mode must be one of {sorted(VALID_MODES)}, got {mode!r}"
        )
    clock = float(now)
    start = clock if valid_from is None else float(valid_from)
    ingest = clock if ingested_at is None else float(ingested_at)
    close_at = clock if invalid_at is None else float(invalid_at)
    for label, value in (
        ("now", clock),
        ("valid_from", start),
        ("ingested_at", ingest),
        ("close_at", close_at),
    ):
        if math.isnan(value):
            # Redis: the script's ZADD replies ``value is not a valid
            # float`` (a ResponseError; with only valid_from NaN, after
            # the incumbent's close -- #778). Here: the same text, a
            # ValueError, and nothing written.
            raise ValueError(f"{NAN_SCORE_ERROR} ({label} is NaN)")
    new = successor.canonical if successor is not None else ""
    named_old = incumbent.canonical if incumbent is not None else ""
    ts = self._table(spec, write=True)
    vf, ia, ig = (_col(field, s) for s in (VALID_FROM, INVALID_AT, INGESTED_AT))
    sup, supby = _col(field, SUPERSEDES), _col(field, SUPERSEDED_BY)

    def work(tx: UnitOfWork) -> Optional[str]:
        self._validity_lock(ts, [field], tx)
        old = named_old
        if mode != "open" and not old and identity:
            old = self._pointer_get(ts, field, identity, uow=tx) or ""
        keys = sorted({k for k in (new, old) if k}, key=_sort_key)
        rows: dict[str, tuple[Any, ...]] = {}
        if keys:
            # The backend's one lock order (plan §6, TD-2): the
            # (model, field) lock above, then the record-key locks in
            # ``_pk`` byte order, then the row locks in ``_pk`` order.
            sql, params = self._record_locked(
                ts,
                keys,
                f'SELECT "_pk", {vf}, {ia}, {ig}, {sup}, {supby} '
                f'FROM {ts.qualified} WHERE "_pk" = ANY(%s::text[])'
                f"{_and_live(ts)} "
                'ORDER BY "_pk" COLLATE "C" FOR UPDATE',
                [keys],
            )
            found, _ = self._run(sql, params, uow=tx, write=True)
            rows = {row[0]: tuple(row[1:]) for row in found}

        # -- VALIDATION PHASE: reads and refusals only --------------------
        will_close = False
        if mode != "open":
            if new and new not in rows:
                raise _refuse(f"{MEMBER_ABSENT_ERROR} successor {new}")
            if old and old not in rows:
                if named_old:
                    raise _refuse(f"{MEMBER_ABSENT_ERROR} incumbent {old}")
                old = ""  # a pointer naming a deleted record: a hint
            if old and old != new:
                old_start, old_close = rows[old][0], rows[old][1]
                if old_close is not None and old_close == INF:
                    if old_start is not None and close_at < old_start:
                        raise _refuse(CLOSE_BEFORE_START_ERROR)
                    will_close = True
        if new and assert_valid_from and new in rows:
            stored = rows[new][0]
            if stored is not None and stored != start:
                raise _refuse(
                    f"{VALID_FROM_CONFLICT_ERROR} {_lua_number(stored)} "
                    f"{_lua_number(start)}"
                )

        # -- MUTATION PHASE: every check above has passed -----------------
        updates: dict[str, dict[str, Any]] = {}
        if will_close:
            updates[old] = {"ia": close_at}
            if new:
                updates[old]["supby"] = new
                updates.setdefault(new, {})["sup"] = old
        new_open = False
        if new and new in rows:
            state = rows[new]
            if state[1] is None or state[1] == INF:
                new_open = True
                target = updates.setdefault(new, {})
                target["vf"] = state[0] if state[0] is not None else start
                target["ig"] = state[2] if state[2] is not None else ingest
                target["ia"] = INF if state[1] is None else state[1]
        if updates:
            self._write_states(ts, field, rows, updates, tx)
        if new_open and identity:
            self._run(
                f"INSERT INTO {pointer_table(ts, field)} (digest, member) "
                "VALUES (%s, %s) ON CONFLICT (digest) DO UPDATE "
                "SET member = EXCLUDED.member",
                [identity, new],
                uow=tx,
                write=True,
            )
        return old if will_close else None

    return self._atomically(work, uow=uow)

chain(spec, field, id)

The supersession chain through id, oldest first, in one WITH RECURSIVE: back along supersedes, then forward along superseded_by. [] when id has no valid_from (unsaved, or never opened). A walk stops at a missing link, at a record it has already visited (a cycle -- the forward walk counts the backward walk's records as visited, as the Redis walk's shared seen set does), and at a link naming a record with no valid_from (a hard delete leaves the neighbour's link naming it).

Source code in src/popoto/backends/postgres/validity.py
def chain(self, spec: ModelSpec, field: str, id: RecordId) -> list[RecordId]:
    """The supersession chain through ``id``, oldest first, in one
    ``WITH RECURSIVE``: back along ``supersedes``, then forward along
    ``superseded_by``. ``[]`` when ``id`` has no ``valid_from`` (unsaved,
    or never opened). A walk stops at a missing link, at a record it has
    already visited (a cycle -- the forward walk counts the backward
    walk's records as visited, as the Redis walk's shared ``seen`` set
    does), and at a link naming a record with no ``valid_from`` (a hard
    delete leaves the neighbour's link naming it)."""
    self._validity_field(spec, field)
    ts = self._table(spec)
    vf = _col(field, VALID_FROM, "n")
    sup = _col(field, SUPERSEDES, "s")
    supby = _col(field, SUPERSEDED_BY, "s")
    anchor_vf = _col(field, VALID_FROM)
    t = ts.qualified
    # M5: an expired record is not on the chain -- the walk stops at it
    # as at a hard delete.
    live_n = _and_live(ts, "n")
    sql = (
        "WITH RECURSIVE "
        f'anchor AS (SELECT "_pk" AS pk FROM {t} WHERE "_pk" = %s '
        f"AND {anchor_vf} IS NOT NULL{_and_live(ts)}), "
        "older(pk, depth, path) AS ("
        f'SELECT n."_pk", 1, ARRAY[a.pk, n."_pk"] FROM anchor a '
        f'JOIN {t} s ON s."_pk" = a.pk JOIN {t} n ON n."_pk" = {sup} '
        f'WHERE {vf} IS NOT NULL AND n."_pk" <> a.pk{live_n} '
        f'UNION ALL SELECT n."_pk", o.depth + 1, o.path || n."_pk" '
        f'FROM older o JOIN {t} s ON s."_pk" = o.pk '
        f'JOIN {t} n ON n."_pk" = {sup} '
        f'WHERE {vf} IS NOT NULL AND n."_pk" <> ALL(o.path){live_n}), '
        "seen AS (SELECT a.pk AS pk FROM anchor a UNION ALL "
        "SELECT pk FROM older), "
        "newer(pk, depth, path) AS ("
        f'SELECT n."_pk", 1, '
        '(SELECT array_agg(pk) FROM seen) || n."_pk" FROM anchor a '
        f'JOIN {t} s ON s."_pk" = a.pk JOIN {t} n ON n."_pk" = {supby} '
        f'WHERE {vf} IS NOT NULL AND n."_pk" <> ALL(SELECT pk FROM seen){live_n} '
        f'UNION ALL SELECT n."_pk", w.depth + 1, w.path || n."_pk" '
        f'FROM newer w JOIN {t} s ON s."_pk" = w.pk '
        f'JOIN {t} n ON n."_pk" = {supby} '
        f'WHERE {vf} IS NOT NULL AND n."_pk" <> ALL(w.path){live_n}) '
        "SELECT pk, -depth FROM older UNION ALL SELECT pk, 0 FROM anchor "
        "UNION ALL SELECT pk, depth FROM newer ORDER BY 2"
    )
    rows, _ = self._run(sql, [id.canonical])
    return [RecordId(spec.name, (), pk, native=pk) for pk, _depth in rows]

validity_field_names(spec)

The model's ValidityField names, in declaration-sorted order.

Source code in src/popoto/backends/postgres/validity.py
def validity_field_names(spec: ModelSpec) -> list[str]:
    """The model's ``ValidityField`` names, in declaration-sorted order."""
    return sorted(n for n, fs in spec.fields.items() if fs.kind == VALIDITY_KIND)

validity_columns(spec)

The interval and chain-link columns for each ValidityField.

Source code in src/popoto/backends/postgres/validity.py
def validity_columns(spec: ModelSpec) -> list[Column]:
    """The interval and chain-link columns for each ``ValidityField``."""
    out: list[Column] = []
    for name in validity_field_names(spec):
        for suffix, sql_type in VALIDITY_SUFFIXES:
            out.append(Column(name + suffix, sql_type, name, "state"))
    return out

validity_indexes(spec, index_name)

A B-tree on each gate column: filter(validity__as_of=…) and the exclusion reads (invalid_at <= t, valid_from > t) are range scans on them.

Source code in src/popoto/backends/postgres/validity.py
def validity_indexes(
    spec: ModelSpec, index_name: Callable[..., str]
) -> list[tuple[str, str, bool]]:
    """A B-tree on each gate column: ``filter(validity__as_of=…)`` and the
    exclusion reads (``invalid_at <= t``, ``valid_from > t``) are range
    scans on them."""
    out = []
    for name in validity_field_names(spec):
        for suffix in (VALID_FROM, INVALID_AT):
            col = name + suffix
            out.append((index_name(col), f"({quote_ident(col)})", False))
    return out

pointer_table(ts, field)

<schema>.<table>__<f>__open: the open-claim pointers (digest -> member) of one ValidityField.

Source code in src/popoto/backends/postgres/validity.py
def pointer_table(ts: TableSpec, field: str) -> str:
    """``<schema>.<table>__<f>__open``: the open-claim pointers (digest ->
    member) of one ``ValidityField``."""
    from .schema import _bounded

    name = _bounded(f"{ts.table}__{field}__open")
    return f"{quote_ident(ts.schema)}.{quote_ident(name)}"

ensure_validity_tables(conn, ts, spec)

Create each ValidityField's pointer table if it is missing, on conn and in its own transaction, under the DDL advisory lock: run at the model's first use, right after its table is created or checked, so no supersede ever has to create one inside a transaction that already holds the model table's locks. CREATE … IF NOT EXISTS is idempotent, so a process that finds the table does nothing.

Source code in src/popoto/backends/postgres/validity.py
def ensure_validity_tables(conn: Any, ts: TableSpec, spec: ModelSpec) -> None:
    """Create each ``ValidityField``'s pointer table if it is missing, on
    ``conn`` and in its own transaction, under the DDL advisory lock: run at
    the model's first use, right after its table is created or checked, so
    no supersede ever has to create one inside a transaction that already
    holds the model table's locks. ``CREATE … IF NOT EXISTS`` is idempotent,
    so a process that finds the table does nothing."""
    from . import _schema_auto
    from .schema import _bounded, table_lock_key

    for field in validity_field_names(spec):
        qualified = pointer_table(ts, field)
        name = _bounded(f"{ts.table}__{field}__open")
        with conn.transaction():
            cur = conn.cursor()
            exists = cur.execute(
                "SELECT to_regclass(%s) IS NOT NULL", (qualified,)
            ).fetchone()[0]
            if exists:
                continue
            if not _schema_auto():
                from ..types import SchemaDriftError

                raise SchemaDriftError(
                    f"{qualified} does not exist and POPOTO_SCHEMA_AUTO=0"
                )
            cur.execute(
                "SELECT pg_advisory_xact_lock(hashtext(%s))",
                (table_lock_key(ts.schema, ts.table),),
            )
            cur.execute(
                f"CREATE TABLE IF NOT EXISTS {qualified} ("
                "digest text PRIMARY KEY, member text NOT NULL "
                f'REFERENCES {ts.qualified} ("_pk") ON DELETE CASCADE)'
            )
            index = quote_ident(_bounded(f"{name}__member"))
            cur.execute(f"CREATE INDEX IF NOT EXISTS {index} ON {qualified} (member)")

excluded_sql(field, as_of, alias='t')

invalid_at <= as_of OR valid_from > as_of: the rule every gate applies (DECAY_SCORE_LUA's, the composite mask's, resolve_excluded_keys'). An absent end (NULL) never excludes; an open record (+inf) is closed only at as_of = +inf.

Source code in src/popoto/backends/postgres/validity.py
def excluded_sql(field: str, as_of: float, alias: str = "t") -> str:
    """``invalid_at <= as_of OR valid_from > as_of``: the rule every gate
    applies (``DECAY_SCORE_LUA``'s, the composite mask's,
    ``resolve_excluded_keys``'). An absent end (``NULL``) never excludes; an
    open record (``+inf``) is closed only at ``as_of = +inf``."""
    a = _lit(as_of)
    ia = _col(field, INVALID_AT, alias)
    vf = _col(field, VALID_FROM, alias)
    return f"(coalesce({ia} <= {a}, false) OR coalesce({vf} > {a}, false))"

included_sql(field, as_of, alias='t')

The gate as a WHERE term: TRUE when there is no gate (no as-of, or a NaN one -- which excludes nothing in the Lua, where every comparison with NaN is false), else NOT the exclusion rule.

Source code in src/popoto/backends/postgres/validity.py
def included_sql(field: str, as_of: Optional[float], alias: str = "t") -> str:
    """The gate as a ``WHERE`` term: ``TRUE`` when there is no gate (no
    as-of, or a NaN one -- which excludes nothing in the Lua, where every
    comparison with NaN is false), else ``NOT`` the exclusion rule."""
    if as_of is None or math.isnan(float(as_of)):
        return "TRUE"
    return f"NOT {excluded_sql(field, float(as_of), alias)}"

range_bound(as_of)

as_of as a range-read bound, refused as Redis refuses it: a NaN bound makes ZRANGEBYSCORE/ZRANGESTORE reply min or max is not a float, so the reads that are range reads on Redis (the filters, the resolvers, the composite mask) raise QueryException with that text here -- the same text, a different class (M1's divergence (v)). The decay ranking's gate is a Lua comparison instead, and a NaN there excludes nothing on both (:func:included_sql).

Source code in src/popoto/backends/postgres/validity.py
def range_bound(as_of: Any) -> float:
    """``as_of`` as a range-read bound, refused as Redis refuses it: a NaN
    bound makes ``ZRANGEBYSCORE``/``ZRANGESTORE`` reply ``min or max is not
    a float``, so the reads that are range reads on Redis (the filters, the
    resolvers, the composite mask) raise ``QueryException`` with that text
    here -- the same text, a different class (M1's divergence (v)). The decay
    ranking's gate is a Lua comparison instead, and a NaN there excludes
    nothing on both (:func:`included_sql`)."""
    t = float(as_of)
    if math.isnan(t):
        from ...models.query import QueryException

        raise QueryException("min or max is not a float")
    return t

validity_cond_sql(field, value)

filter(validity__as_of=t) / validity__current=…, compiled by :mod:popoto.backends.planning to Cond(field, VALID_AT, (t, valid)).

valid=True: valid_from <= t AND invalid_at > t -- both ends present, the two ZRANGEBYSCOREs ValidityField._members_valid_at intersects. valid=False (__current=False): the members of either index that are not valid at t, the Redis complement.

Source code in src/popoto/backends/postgres/validity.py
def validity_cond_sql(field: str, value: Any) -> str:
    """``filter(validity__as_of=t)`` / ``validity__current=…``, compiled by
    :mod:`popoto.backends.planning` to ``Cond(field, VALID_AT, (t, valid))``.

    ``valid=True``: ``valid_from <= t AND invalid_at > t`` -- both ends
    present, the two ``ZRANGEBYSCORE``s ``ValidityField._members_valid_at``
    intersects. ``valid=False`` (``__current=False``): the members of either
    index that are not valid at ``t``, the Redis complement."""
    t, valid = value
    a = _lit(range_bound(t))
    vf = _col(field, VALID_FROM)
    ia = _col(field, INVALID_AT)
    is_valid = f"({vf} <= {a} AND {ia} > {a})"
    if valid:
        return is_valid
    return (
        f"(({vf} IS NOT NULL OR {ia} IS NOT NULL) AND NOT coalesce({is_valid}, false))"
    )

save_parts(ts, spec, obj, names)

What ValidityField.on_save does on Redis, as parts of the save upsert: SUPERSEDE_LUA in mode 'open'.

Returns (insert columns, ON CONFLICT overrides, DO UPDATE guards). A new row opens at the declared valid_from (else the save clock), ingested now, invalid_at = +inf. An existing row keeps its interval (NX): an absent end is filled, and a closed record is never reopened -- the guard is the row the upsert re-read after waiting on any concurrent supersede (plan Race 2). A declared valid_from that disagrees with the stored start fails the guard: nothing is written and the caller raises ValidityValidFromConflictError, the script's ARGV[8] check.

Source code in src/popoto/backends/postgres/validity.py
def save_parts(
    ts: TableSpec, spec: ModelSpec, obj: Any, names: Sequence[str]
) -> tuple[dict[str, Any], dict[str, str], list[str]]:
    """What ``ValidityField.on_save`` does on Redis, as parts of the save
    upsert: ``SUPERSEDE_LUA`` in mode ``'open'``.

    Returns ``(insert columns, ON CONFLICT overrides, DO UPDATE guards)``.
    A new row opens at the declared ``valid_from`` (else the save clock),
    ingested now, ``invalid_at = +inf``. An existing row keeps its interval
    (``NX``): an absent end is filled, and a closed record is never reopened
    -- the guard is the row the upsert re-read after waiting on any
    concurrent supersede (plan Race 2). A *declared* ``valid_from`` that
    disagrees with the stored start fails the guard: nothing is written and
    the caller raises ``ValidityValidFromConflictError``, the script's
    ``ARGV[8]`` check.
    """
    cols: dict[str, Any] = {}
    overrides: dict[str, str] = {}
    guards: list[str] = []
    fields = [n for n in validity_field_names(spec) if n in names]
    if not fields:
        return cols, overrides, guards
    now = time.time()
    from ...fields.validity_field import ValidityField

    for name in fields:
        # Plan D9, as ValidityField.on_save warns on Redis: since M5 a
        # Meta.ttl model stores here too, and an expired record drops out of
        # the chain walk while its neighbours' links still name it.
        ValidityField.warn_if_ttl(obj, name)
        declared = getattr(obj, name, None)
        valid_from = now
        asserted = False
        if declared is not None:
            try:
                valid_from = float(declared)
                asserted = True
            except (TypeError, ValueError):
                valid_from = now
        if math.isnan(valid_from):
            # SUPERSEDE_LUA's ``ZADD`` refuses a NaN score, so the Redis save
            # fails with ``ResponseError: value is not a valid float`` (after
            # MULTI/EXEC has written the hash, which nothing rolls back). Here
            # the save is refused before anything is written: the same text,
            # popoto's save error (the class difference is documented, as
            # M1's "min or max is not a float" is). Stored, a NaN start would
            # hide the record from every gate for good -- Postgres orders NaN
            # above every float.
            from ...exceptions import ModelException

            raise ModelException(NAN_SCORE_ERROR)
        vf, ia, ig = name + VALID_FROM, name + INVALID_AT, name + INGESTED_AT
        cols[vf] = valid_from
        cols[ig] = now
        cols[ia] = INF
        t = quote_ident(ts.table)
        qvf, qia, qig = quote_ident(vf), quote_ident(ia), quote_ident(ig)
        is_open = f"({t}.{qia} IS NULL OR {t}.{qia} = 'Infinity'::float8)"
        overrides[vf] = (
            f"CASE WHEN {is_open} THEN coalesce({t}.{qvf}, EXCLUDED.{qvf}) "
            f"ELSE {t}.{qvf} END"
        )
        overrides[ig] = (
            f"CASE WHEN {is_open} THEN coalesce({t}.{qig}, EXCLUDED.{qig}) "
            f"ELSE {t}.{qig} END"
        )
        overrides[ia] = f"coalesce({t}.{qia}, 'Infinity'::float8)"
        if asserted:
            guards.append(f"({t}.{qvf} IS NULL OR {t}.{qvf} = EXCLUDED.{qvf})")
    return cols, overrides, guards

refuse_valid_from_conflict(backend, spec, obj, *, uow=None)

Raise the VALID_FROM_CONFLICT error for a save whose upsert guard refused it (:func:save_parts), with the numbers the script's reply carries: the stored start, then the declared one, each as Lua's tostring prints it.

Source code in src/popoto/backends/postgres/validity.py
def refuse_valid_from_conflict(
    backend: Any, spec: ModelSpec, obj: Any, *, uow: Optional[UnitOfWork] = None
) -> None:
    """Raise the ``VALID_FROM_CONFLICT`` error for a save whose upsert guard
    refused it (:func:`save_parts`), with the numbers the script's reply
    carries: the stored start, then the declared one, each as Lua's
    ``tostring`` prints it."""
    from ...fields.validity_field import VALID_FROM_CONFLICT_ERROR

    key = obj.db_key.redis_key
    for name in validity_field_names(spec):
        declared: Any = getattr(obj, name, None)
        try:
            requested = float(declared)
        except (TypeError, ValueError):
            continue
        rows, _ = backend._run(
            f"SELECT {_col(name, VALID_FROM)} FROM "
            f'{backend._table(spec).qualified} WHERE "_pk" = %s',
            [key],
            uow=uow,
        )
        stored = rows[0][0] if rows else None
        if stored is not None and stored != requested:
            raise _refuse(
                f"{VALID_FROM_CONFLICT_ERROR} {_lua_number(stored)} "
                f"{_lua_number(requested)}"
            )
    raise _refuse(f"{VALID_FROM_CONFLICT_ERROR} {key}")  # pragma: no cover