Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 19 additions & 0 deletions src/datajoint/adapters/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -1340,6 +1340,25 @@ def job_metadata_columns(self) -> list[str]:
"""
...

@abstractmethod
def provenance_columns(self) -> list[str]:
"""
Return the hidden extrinsic-provenance column for Entry tables.

Returns
-------
list[str]
List of column definition strings (fully formatted with quotes).

Examples
--------
MySQL:
["`_prov` json DEFAULT NULL"]
PostgreSQL:
['"_prov" jsonb DEFAULT NULL']
"""
...

# =========================================================================
# Error Translation
# =========================================================================
Expand Down
11 changes: 11 additions & 0 deletions src/datajoint/adapters/mysql.py
Original file line number Diff line number Diff line change
Expand Up @@ -1013,6 +1013,17 @@ def job_metadata_columns(self) -> list[str]:
"`_job_version` varchar(64) DEFAULT ''",
]

def provenance_columns(self) -> list[str]:
"""
Return the MySQL extrinsic-provenance column definition.

Examples
--------
>>> adapter.provenance_columns()
["`_prov` json DEFAULT NULL"]
"""
return ["`_prov` json DEFAULT NULL"]

# =========================================================================
# Error Translation
# =========================================================================
Expand Down
11 changes: 11 additions & 0 deletions src/datajoint/adapters/postgres.py
Original file line number Diff line number Diff line change
Expand Up @@ -1366,6 +1366,17 @@ def job_metadata_columns(self) -> list[str]:
"\"_job_version\" varchar(64) DEFAULT ''",
]

def provenance_columns(self) -> list[str]:
"""
Return the PostgreSQL extrinsic-provenance column definition.

Examples
--------
>>> adapter.provenance_columns()
['"_prov" jsonb DEFAULT NULL']
"""
return ['"_prov" jsonb DEFAULT NULL']

# =========================================================================
# Error Translation
# =========================================================================
Expand Down
9 changes: 9 additions & 0 deletions src/datajoint/autopopulate.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
import traceback
from typing import TYPE_CHECKING, Any, Generator

from . import provenance
from .errors import DataJointError, LostConnectionError
from .expression import AndList, QueryExpression

Expand Down Expand Up @@ -684,6 +685,13 @@ def _populate1(
self._upstream_key = dict(key)
self._upstream = None

# Rows this make() writes into Entry tables that carry no foreign key back
# here -- the fan-out ingestion pattern -- record the ingesting table and
# key, which is what makes such a row traceable without one.
from .jobs import _get_job_version

prov_token = provenance.set_ingesting(self.full_table_name, key, _get_job_version(self.connection._config))

try:
if not is_generator:
make(dict(key), **(make_kwargs or {}))
Expand Down Expand Up @@ -740,6 +748,7 @@ def _populate1(
jobs.complete(key, duration=duration)
return True
finally:
provenance.reset_ingesting(prov_token)
self.__class__._allow_insert = False
# Clear the per-make() upstream state: `_upstream = None` invalidates
# the memoized Diagram; `_upstream_key = None` restores the "outside
Expand Down
14 changes: 14 additions & 0 deletions src/datajoint/declare.py
Original file line number Diff line number Diff line change
Expand Up @@ -521,6 +521,20 @@ def declare(
job_metadata_sql = adapter.job_metadata_columns()
attribute_sql.extend(job_metadata_sql)

# Add the hidden extrinsic-provenance slot to Entry tables, where rows enter
# from outside the pipeline. Computed and Imported tables have no use for
# it -- their provenance is entailed by the foreign-key graph -- and a part
# inherits its master's.
# Matched against the Manual tier itself, not by excluding the other tiers'
# prefixes: enumerating exclusions makes every tier added later an Entry
# table by default, which is how job tables (`~`) first acquired the slot.
# Imported here rather than at module scope: user_tables imports table,
# which imports this module.
from .user_tables import Manual

if config.provenance.capture and re.fullmatch(Manual.tier_regexp, table_name):
attribute_sql.extend(adapter.provenance_columns())

if not primary_key:
# Singleton table: add hidden sentinel attribute
primary_key = ["_singleton"]
Expand Down
114 changes: 113 additions & 1 deletion src/datajoint/deploy.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,8 @@
``add_job_metadata_columns``, ``rebuild_lineage``.
- :mod:`datajoint.deploy` — configure an environment for a consumer's
requirements (CDC tools, replication, role grants, performance tuning).
Cadence: re-runnable, idempotent. Examples: :func:`set_replica_identity`.
Cadence: re-runnable, idempotent. Examples: :func:`set_replica_identity`,
:func:`add_prov_column`.

Functions in this module should be safe to call repeatedly from a deploy hook
without accumulating side effects.
Expand Down Expand Up @@ -183,3 +184,114 @@ def set_replica_identity(
connection.query(ddl)
result["tables_modified"] += 1
return result


def add_prov_column(target: "TargetType", dry_run: bool = True) -> dict:
"""
Add the hidden ``_prov`` attribute to Entry (``dj.Manual``) tables that lack it.

Capture defaults on, so tables declared from 2.3.4 onward already carry the
slot. Two populations do not: tables declared before 2.3.4, and tables
declared while ``config.provenance.capture`` was off. Inserts into those
record nothing, silently, and this brings them in line.

It belongs here rather than in :mod:`datajoint.migrate` because it is not a
one-shot correction of legacy state. It is idempotent — a table that already
has the column is reported and left alone — and it stays useful for as long
as capture can be turned off, which outlives the migration module.

Parameters
----------
target : Schema, Table class, or Table instance
Given a Schema, every Entry table in it is processed.
dry_run : bool, optional
If True, report what would change without altering anything. Default True.

Returns
-------
dict
``tables_analyzed``, ``tables_modified``, ``columns_added``, ``ddl``,
and ``details`` — a per-table list of dicts.

Examples
--------
>>> from datajoint.deploy import add_prov_column
>>> add_prov_column(schema, dry_run=True)["ddl"]
>>> add_prov_column(schema, dry_run=False)["tables_modified"]

Notes
-----
- Only Entry tables are touched. Computed and Imported tables have no use
for the slot — their provenance is entailed by the foreign-key graph — and
a part table inherits its master's.
- Rows already present keep ``NULL``. Provenance is recorded at insert and
is never reconstructed after the fact.
"""
import re

from . import provenance
from .schemas import _Schema
from .table import Table
from .user_tables import Manual

if isinstance(target, _Schema):
connection = target.connection
if connection is None or not target.database:
raise DataJointError("Schema is not activated. Call schema.activate(...) before add_prov_column().")
database = target.database
table_names = list(target.list_tables())
elif isinstance(target, type) and issubclass(target, Table):
instance = target()
connection = instance.connection
if connection is None:
raise DataJointError(f"Table {target.__name__} has no active connection.")
database, table_names = instance.database, [instance.table_name]
elif isinstance(target, Table):
connection = target.connection
if connection is None:
raise DataJointError(f"Table {type(target).__name__} has no active connection.")
database, table_names = target.database, [target.table_name]
else:
raise DataJointError(f"target must be a Schema or Table class/instance; got {type(target).__name__}")

if not database:
raise DataJointError("Cannot add the provenance column: the target has no database.")

adapter = connection.adapter
column_sql = adapter.provenance_columns()[0]

result: dict[str, Any] = {
"tables_analyzed": 0,
"tables_modified": 0,
"columns_added": 0,
"ddl": [],
"details": [],
}

for table_name in table_names:
if not re.fullmatch(Manual.tier_regexp, table_name):
continue
result["tables_analyzed"] += 1

existing = {
row[0]
for row in connection.query(
"SELECT COLUMN_NAME FROM information_schema.COLUMNS WHERE TABLE_SCHEMA = %s AND TABLE_NAME = %s",
args=(database, table_name),
).fetchall()
}
if provenance.PROV_ATTRIBUTE in existing:
result["details"].append({"table": f"{database}.{table_name}", "status": "already_present"})
continue

ddl = (
f"ALTER TABLE {adapter.quote_identifier(database)}.{adapter.quote_identifier(table_name)} ADD COLUMN {column_sql}"
)
result["ddl"].append(ddl)
result["details"].append({"table": f"{database}.{table_name}", "status": "pending" if dry_run else "added"})
if not dry_run:
connection.query(ddl)
result["tables_modified"] += 1
result["columns_added"] += 1

return result
6 changes: 5 additions & 1 deletion src/datajoint/heading.py
Original file line number Diff line number Diff line change
Expand Up @@ -397,7 +397,11 @@ def quote(name):
return adapter.quote_identifier(name) if adapter else f'"{name}"'

def render_field(name):
attr = self.attributes[name]
# `attributes` hides underscore-prefixed names, so a caller that asks
# for one by name -- copying `_prov` through an INSERT ... SELECT --
# falls back to the full set. Default field lists are unaffected:
# they are built from `attributes` and never contain hidden names.
attr = self.attributes.get(name) or self._attributes[name]
if attr.attribute_expression is None:
return quote(name)
else:
Expand Down
140 changes: 140 additions & 0 deletions src/datajoint/provenance.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,140 @@
"""Extrinsic provenance for rows that enter the pipeline from outside.

Inside the pipeline, provenance is structural: a Computed table's row cannot
exist unless its declared upstream exists, so the foreign-key graph *is* the
lineage and nothing has to be recorded for it to hold.

At the boundary the structure runs out. Rows arrive in Entry tables from a
person, an instrument, or a feed, and the framework has no way to say where
they came from. This module supplies the slot and fills it.

The attribute is **framework-owned: no author ever writes it.** ``insert``
takes no provenance argument, and its content comes from three places, none of
them the call site:

* **configuration** -- ``config.provenance.source``, set per deployment, naming
the external system this process draws from;
* **ambient connection state** -- the connecting user, host and database, the
insert time, and the code version;
* **ambient execution state** -- the ingesting table and key, when the insert
runs inside a ``make()``.

That ownership is the point. A field an operator can set is weaker evidence
than one the system sets, and nothing is left for a pipeline to neglect.
Anything an author wants to record deliberately belongs in the data model as a
visible attribute, where queries can reach it.
"""

import contextlib
import contextvars
import datetime
import json
from typing import Any

#: Name of the hidden attribute. Hidden attributes are excluded from
#: ``heading.attributes``, so this never appears in a query heading.
PROV_ATTRIBUTE = "_prov"

# Set by autopopulate around a make() call so that rows written to Entry tables
# from inside an ingesting make() record what wrote them, which is what makes a
# fanned-out row traceable without a foreign key.
_ingesting: contextvars.ContextVar = contextvars.ContextVar("dj_ingesting", default=None)


def set_ingesting(table_name, key, version=None):
"""Record the ``make()`` now executing; returns a token for ``reset_ingesting``.

Parameters
----------
table_name : str
Full table name of the ingesting table.
key : dict
The key ``make()`` was called with.
version : str, optional
Code version, as resolved for the job.
"""
value = {
"table": table_name,
"key": {k: _jsonable(v) for k, v in (key or {}).items()},
}
if version:
value["version"] = version
return _ingesting.set(value)


def reset_ingesting(token):
"""Restore the ingesting context saved by :func:`set_ingesting`."""
if token is not None:
_ingesting.reset(token)


@contextlib.contextmanager
def ingesting(table_name, key, version=None):
"""Scope :func:`set_ingesting` to a block."""
token = set_ingesting(table_name, key, version)
try:
yield
finally:
reset_ingesting(token)


def _jsonable(value):
"""Render a key value in a form ``json.dumps`` accepts."""
if isinstance(value, (str, int, float, bool)) or value is None:
return value
if isinstance(value, (datetime.datetime, datetime.date, datetime.time)):
return value.isoformat()
if isinstance(value, bytes):
return value.hex()
return str(value)


def build_payload(connection, config=None):
"""Assemble the provenance record for rows inserted on this connection.

Returns ``None`` when there is nothing worth recording, so that a row is
left with ``NULL`` rather than an empty object.
"""
if config is None:
from .settings import config as _config

config = _config

payload: dict[str, Any] = {"time": datetime.datetime.now(datetime.timezone.utc).isoformat()}

conn_info = getattr(connection, "conn_info", None) or {}
agent = {key: conn_info[key] for key in ("user", "host", "database_name") if conn_info.get(key) is not None}
if agent:
payload["agent"] = agent

try:
from .jobs import _get_job_version

version = _get_job_version(getattr(connection, "_config", None) or config)
except Exception: # version capture must never break an insert
version = ""
if version:
payload["version"] = version

source = config.provenance.source
if source:
payload["source"] = source

context = _ingesting.get()
if context:
payload["context"] = context

# Time alone says nothing about origin; without any of the other three this
# is noise rather than a record.
return payload if len(payload) > 1 else None


def serialize(payload):
"""Render a payload for the ``json`` column.

``default=str`` because ``config.provenance.source`` is deployment-supplied
and typed ``dict[str, Any]``: a ``date`` or a ``Path`` in it would otherwise
raise from inside every insert into every Entry table, with an error naming
neither provenance nor the setting that caused it.
"""
return json.dumps(payload, default=str)
Loading
Loading