Skip to content

popoto.transfer.import_

popoto.transfer.import_

Import records from a JSON Lines export produced by :mod:popoto.transfer.

Keys are preserved by default. Records are saved one at a time -- deliberately not through bulk_create -- because an external pipeline makes save() return the pipeline for every record and destroys per-record observability, which the reconciliation ledger depends on.

preserve_keys=False is the opt-out: every record gets a freshly minted key and every field that declares a reference to another record is rewritten to point at the new one. That mode reads the record stream twice (see _spool_and_mint), because a reference on record 3 can point at record 4000 and the complete old-to-new map must exist before the first write. The default path is untouched by it: one forward pass, no buffering, no map.

BATCH_SIZE = 500 module-attribute

Records per conflict-check / reconciliation batch.

REMAPPED_REFERENCES_KEY = '__remapped_references__' module-attribute

In-memory slot on a record dict carrying its rewritten reference strings.

Set by :func:_remap_record and consumed by :func:_process_batch only when that batch was produced by the regenerating path. Export never writes this name, and the preserving path pops it and throws it away, so a file that happens to carry one cannot use it to set relationship values.

import_records(model_class, stream, on_conflict='error', on_write_gate='reject', on_embedding_mismatch='error', preserve_keys=True, key_map=None)

Import records from a JSON Lines export into model_class.

By default keys are preserved, so re-running an import converges rather than duplicating, and every Relationship value and application-level key pointer keeps pointing at the same record.

preserve_keys=False regenerates every key instead, and rewrites the references it can. Three properties of that mode matter before choosing it:

  • It is not idempotent. Every other mode converges on a re-run; this one mints fresh keys each time, so running it twice against the same destination leaves two copies of every record.
  • The remap is partial, by design. Only fields that declare a reference are rewritten -- Relationship is the only one Popoto ships. An application-level pointer stored in a plain Field(type=str) is indistinguishable from ordinary text, is never rewritten, and will dangle. Popoto does not guess: a heuristic that scanned strings for key-shaped values would turn a documented limitation into occasional silent corruption. The report names both counts so a dangling reference is visible in the run that created it.
  • It requires a mintable key. The destination model's key must be exactly one auto=True field; anything else is refused before a byte is written.

Parameters:

Name Type Description Default
model_class 'type[Model]'

The destination Model class.

required
stream TextIO

A text file-like object positioned at the manifest line.

required
on_conflict str

What to do when the destination already holds a key. "error" (default) refuses -- the only mode that cannot clobber; "skip" leaves the existing record untouched; "overwrite" replaces it, which makes a re-run idempotent.

'error'
on_write_gate str

"reject" (default) honors the destination model's WriteFilterMixin gate and reports every refusal; "bypass" writes around it. Bypass only disables Popoto's own gate -- an application's save() override returning falsy still surfaces as a rejection.

'reject'
on_embedding_mismatch str

"error" (default) refuses when the export's provider fingerprint differs from the destination's; "carry" imports the vectors anyway; "regenerate" drops them so on_save re-embeds.

'error'
preserve_keys bool

True (default) carries each record's key verbatim. False mints a new key per record and remaps declared references onto it -- see the three caveats above. Under False, on_conflict still applies but covers only the vanishing-probability case of a minted key colliding with one already on the destination, not the merge semantics it has when keys are preserved.

True
key_map 'dict[str, str] | None'

Only valid with preserve_keys=False: a seed of old-to-new redis_key mappings from a previous import of another model, so a multi-model migration can carry cross-model references across runs. Feed run N's ImportReport.key_map in as run N+1's seed. A reference whose target is in no seed and not in this file keeps its old key and is counted as dangling.

None

Returns:

Name Type Description
An ImportReport

class:ImportReport accounting for every record line as landed,

ImportReport

skipped, rejected, errored, or partial, with a reason per non-landed

ImportReport

record. Under preserve_keys=False its key_map holds every

ImportReport

mapping this run actually wrote, merged over the seed -- mints whose

ImportReport

record did not land are pruned; see :class:ImportReport.

Raises:

Type Description
ValueError

If a policy argument is not one of its allowed values, or if key_map is passed with preserve_keys=True (there is nothing to remap when no key changes).

ModelException

If the manifest is missing, the format version is unsupported, the model name does not match, an embedding provenance mismatch is refused, preserve_keys=False is asked of a model whose key is not exactly one auto=True field, or on_conflict="error" hits a collision. Every one of those but the last raises before anything is written.

Note

Import is not atomic across records and assumes the destination is not under concurrent write for this model. The recovery path for an interrupted run is to re-run with on_conflict="overwrite".

Example

with open("memories.jsonl") as fh: report = Memory.import_records(fh, on_conflict="overwrite") print(report.summary())

Source code in src/popoto/transfer/import_.py
def import_records(
    model_class: "type[Model]",
    stream: TextIO,
    on_conflict: str = "error",
    on_write_gate: str = "reject",
    on_embedding_mismatch: str = "error",
    preserve_keys: bool = True,
    key_map: "dict[str, str] | None" = None,
) -> ImportReport:
    """Import records from a JSON Lines export into ``model_class``.

    By default keys are preserved, so re-running an import converges rather
    than duplicating, and every ``Relationship`` value and application-level
    key pointer keeps pointing at the same record.

    ``preserve_keys=False`` regenerates every key instead, and rewrites the
    references it can. Three properties of that mode matter before choosing
    it:

    - **It is not idempotent.** Every other mode converges on a re-run; this
      one mints fresh keys each time, so running it twice against the same
      destination leaves two copies of every record.
    - **The remap is partial, by design.** Only fields that *declare* a
      reference are rewritten -- ``Relationship`` is the only one Popoto
      ships. An application-level pointer stored in a plain
      ``Field(type=str)`` is indistinguishable from ordinary text, is never
      rewritten, and will dangle. Popoto does not guess: a heuristic that
      scanned strings for key-shaped values would turn a documented
      limitation into occasional silent corruption. The report names both
      counts so a dangling reference is visible in the run that created it.
    - **It requires a mintable key.** The destination model's key must be
      exactly one ``auto=True`` field; anything else is refused before a byte
      is written.

    Args:
        model_class: The destination Model class.
        stream: A text file-like object positioned at the manifest line.
        on_conflict: What to do when the destination already holds a key.
            ``"error"`` (default) refuses -- the only mode that cannot
            clobber; ``"skip"`` leaves the existing record untouched;
            ``"overwrite"`` replaces it, which makes a re-run idempotent.
        on_write_gate: ``"reject"`` (default) honors the destination model's
            ``WriteFilterMixin`` gate and reports every refusal;
            ``"bypass"`` writes around it. Bypass only disables Popoto's own
            gate -- an application's ``save()`` override returning falsy still
            surfaces as a rejection.
        on_embedding_mismatch: ``"error"`` (default) refuses when the export's
            provider fingerprint differs from the destination's; ``"carry"``
            imports the vectors anyway; ``"regenerate"`` drops them so
            ``on_save`` re-embeds.
        preserve_keys: ``True`` (default) carries each record's key verbatim.
            ``False`` mints a new key per record and remaps declared
            references onto it -- see the three caveats above. Under
            ``False``, ``on_conflict`` still applies but covers only the
            vanishing-probability case of a minted key colliding with one
            already on the destination, not the merge semantics it has when
            keys are preserved.
        key_map: Only valid with ``preserve_keys=False``: a seed of old-to-new
            ``redis_key`` mappings from a previous import of *another* model,
            so a multi-model migration can carry cross-model references
            across runs. Feed run N's ``ImportReport.key_map`` in as run
            N+1's seed. A reference whose target is in no seed and not in
            this file keeps its old key and is counted as dangling.

    Returns:
        An :class:`ImportReport` accounting for every record line as landed,
        skipped, rejected, errored, or partial, with a reason per non-landed
        record. Under ``preserve_keys=False`` its ``key_map`` holds every
        mapping this run actually wrote, merged over the seed -- mints whose
        record did not land are pruned; see :class:`ImportReport`.

    Raises:
        ValueError: If a policy argument is not one of its allowed values, or
            if ``key_map`` is passed with ``preserve_keys=True`` (there is
            nothing to remap when no key changes).
        ModelException: If the manifest is missing, the format version is
            unsupported, the model name does not match, an embedding
            provenance mismatch is refused, ``preserve_keys=False`` is asked
            of a model whose key is not exactly one ``auto=True`` field, or
            ``on_conflict="error"`` hits a collision. Every one of those but
            the last raises before anything is written.

    Note:
        Import is not atomic across records and assumes the destination is not
        under concurrent write for this model. The recovery path for an
        interrupted run is to re-run with ``on_conflict="overwrite"``.

    Example:
        with open("memories.jsonl") as fh:
            report = Memory.import_records(fh, on_conflict="overwrite")
        print(report.summary())
    """
    _check_choice("on_conflict", on_conflict, _ON_CONFLICT)
    _check_choice("on_write_gate", on_write_gate, _ON_WRITE_GATE)
    _check_choice(
        "on_embedding_mismatch", on_embedding_mismatch, _ON_EMBEDDING_MISMATCH
    )
    if preserve_keys and key_map is not None:
        raise ValueError(
            "key_map is only meaningful with preserve_keys=False; when keys "
            "are preserved no reference needs remapping"
        )

    report = ImportReport(model=model_class.__name__)

    lines = iter_lines(stream)
    manifest = None
    for _line_number, raw in lines:
        try:
            manifest = parse_line(raw)
        except ValueError as exc:
            raise ModelException(
                f"first line of the export is not valid JSON: {exc}"
            ) from exc
        break

    manifest = _validate_manifest(model_class, manifest)
    report.source_matched_count = manifest.get("matched_count")
    report.fidelity = {
        **(manifest.get("fields") or {}),
        **(manifest.get("mixins") or {}),
    }
    drop_state = _resolve_embedding_provenance(
        model_class, manifest, on_embedding_mismatch, report
    )

    if not preserve_keys:
        working_map: "dict[str, str]" = dict(key_map or {})
        report.key_map = working_map
        _import_regenerating(
            model_class,
            lines,
            report,
            on_conflict,
            on_write_gate,
            drop_state,
            working_map,
        )
        return report

    batch: "list[dict[str, Any]]" = []
    for line_number, raw in lines:
        try:
            record = parse_line(raw)
        except ValueError as exc:
            report.add(f"line {line_number}", ERRORED, f"malformed JSON line: {exc}")
            continue
        key = record.get("key")
        if not isinstance(key, str) or not key:
            report.add(f"line {line_number}", ERRORED, _NO_KEY_REASON)
            continue
        batch.append(record)
        if len(batch) >= BATCH_SIZE:
            _process_batch(
                model_class, batch, report, on_conflict, on_write_gate, drop_state
            )
            batch = []

    if batch:
        _process_batch(
            model_class, batch, report, on_conflict, on_write_gate, drop_state
        )

    return report