Skip to content

Commit 6a3adc9

Browse files
Hidden _prov attribute for extrinsic provenance at Entry tables (#1555)
* feat(provenance): hidden `_prov` attribute on Entry tables Inside the pipeline provenance is structural: a Compute row cannot exist unless its declared upstream does, so the foreign-key graph is the lineage. At the boundary that runs out. Rows arrive in Entry tables from a person, an instrument, or a feed, and until now each pipeline invented its own record of where they came from -- a source_file column here, a notes varchar there, or nothing at all. Adds a hidden `_prov` JSON attribute to Entry (dj.Manual) tables, declared by default and filled on insert. The attribute is framework-owned: no author ever writes it. `insert` gains no argument, and the content comes from three places, none of them the call site: - configuration -- `config.provenance.source`, set per deployment - ambient connection state -- user, host, database, insert time, code version - ambient execution state -- the ingesting table and key, inside a `make()` That third source is what makes a fan-out write traceable: a row written into an Entry table from inside an ingesting `make()` records what wrote it without any foreign key. 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 model as a visible attribute, where queries can reach it. Entry tables only. Compute and Ingest have no use for the slot -- their provenance is entailed by the graph, and job metadata already records agent, time and version for them -- while a part inherits its master's. Capture defaults on. Off by default would leave no consumer able to assume the column exists and would turn "which Entry rows have no recorded origin" into a question conditional on each table's declaration-time config. `deploy.add_prov_column` adds the slot to tables declared earlier. It sits in deploy rather than migrate because it is idempotent and stays useful for as long as capture can be turned off, which outlives the migration module's scheduled removal. * fix(provenance): exclude job tables, and never let `source` break an insert Both from review of #1555. Job tables were getting `_prov`. `is_entry_table` excluded the other tiers' prefixes one at a time -- `_`, `#`, and `__` for parts -- and `~` was not among them, so `~~analysis`, `~jobs` and `~lineage` all passed. The column really landed and every job row carried a payload; on a large queue that is a lot of JSON nobody asked for. It only surfaced after `jobs.refresh()` materialises the table, which is why a plain populate() in the first round of tests missed it. Now matched against `Manual.tier_regexp`. Enumerating what to exclude makes every tier added later an Entry table until someone remembers this function; matching the tier definition inverts that, so a name is an Entry table only if the library says it is. A `source` that json could not render broke every insert. `config.provenance.source` is deployment-supplied and typed `dict[str, Any]`, so a `date` in it raised `TypeError: Object of type date is not JSON serializable` from inside every insert into every Entry table, naming neither provenance nor the setting responsible. Three layers close it: - `serialize` passes `default=str`, so ordinary values a deployment would actually set -- dates, paths -- record correctly rather than failing; - a validator on the field rejects what remains (non-string keys, cycles) at assignment, where the error belongs; - `_attach_provenance` catches and logs, so recording where a row came from can never stop the row being written. That guarantee previously covered only the version call. Also fixes the `datajoint.migrate.add_prov_column` reference in the settings description -- it lives in `deploy`, and that string shows up in config help. Regression tests fail against the previous code: 4 of them, including the integration one that drives a real job queue. * refactor(provenance): inline the tier test, carry _prov through INSERT ... SELECT Three follow-ups from review of #1555. The Entry-table test is now inlined at its two call sites rather than wrapped in `provenance.is_entry_table`. The import stays deferred inside the function: user_tables imports table, which imports declare, so a module-scope import is a cycle -- which is what the wrapper had been hiding. `insert(QueryExpression)` builds INSERT ... SELECT and returned before the provenance was attached, so copied rows landed with NULL. It now carries the source's `_prov` across. A copied row did not originate in the destination, so the source's record is the true one; re-stamping it here would claim an origin that is not where the data came from. `heading.as_sql` resolves an explicitly named hidden attribute against the full attribute set. Default field lists are built from `attributes` and still never contain hidden names, so nothing else changes. 681 passed, 14 skipped. * test: compare against str(Path), not a POSIX literal The provenance serialization test asserted the rendered Path equalled '/mnt/raw'. On Windows str(Path('/mnt/raw')) is '\\mnt\\raw', so the assertion failed on both Windows jobs while passing everywhere else. The behavior under test is unchanged and was always correct: default=str stringifies a Path rather than raising. Only the expectation was platform- specific.
1 parent f824faa commit 6a3adc9

12 files changed

Lines changed: 873 additions & 2 deletions

File tree

‎src/datajoint/adapters/base.py‎

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1340,6 +1340,25 @@ def job_metadata_columns(self) -> list[str]:
13401340
"""
13411341
...
13421342

1343+
@abstractmethod
1344+
def provenance_columns(self) -> list[str]:
1345+
"""
1346+
Return the hidden extrinsic-provenance column for Entry tables.
1347+
1348+
Returns
1349+
-------
1350+
list[str]
1351+
List of column definition strings (fully formatted with quotes).
1352+
1353+
Examples
1354+
--------
1355+
MySQL:
1356+
["`_prov` json DEFAULT NULL"]
1357+
PostgreSQL:
1358+
['"_prov" jsonb DEFAULT NULL']
1359+
"""
1360+
...
1361+
13431362
# =========================================================================
13441363
# Error Translation
13451364
# =========================================================================

‎src/datajoint/adapters/mysql.py‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1013,6 +1013,17 @@ def job_metadata_columns(self) -> list[str]:
10131013
"`_job_version` varchar(64) DEFAULT ''",
10141014
]
10151015

1016+
def provenance_columns(self) -> list[str]:
1017+
"""
1018+
Return the MySQL extrinsic-provenance column definition.
1019+
1020+
Examples
1021+
--------
1022+
>>> adapter.provenance_columns()
1023+
["`_prov` json DEFAULT NULL"]
1024+
"""
1025+
return ["`_prov` json DEFAULT NULL"]
1026+
10161027
# =========================================================================
10171028
# Error Translation
10181029
# =========================================================================

‎src/datajoint/adapters/postgres.py‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1366,6 +1366,17 @@ def job_metadata_columns(self) -> list[str]:
13661366
"\"_job_version\" varchar(64) DEFAULT ''",
13671367
]
13681368

1369+
def provenance_columns(self) -> list[str]:
1370+
"""
1371+
Return the PostgreSQL extrinsic-provenance column definition.
1372+
1373+
Examples
1374+
--------
1375+
>>> adapter.provenance_columns()
1376+
['"_prov" jsonb DEFAULT NULL']
1377+
"""
1378+
return ['"_prov" jsonb DEFAULT NULL']
1379+
13691380
# =========================================================================
13701381
# Error Translation
13711382
# =========================================================================

‎src/datajoint/autopopulate.py‎

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@
1111
import traceback
1212
from typing import TYPE_CHECKING, Any, Generator
1313

14+
from . import provenance
1415
from .errors import DataJointError, LostConnectionError
1516
from .expression import AndList, QueryExpression
1617

@@ -684,6 +685,13 @@ def _populate1(
684685
self._upstream_key = dict(key)
685686
self._upstream = None
686687

688+
# Rows this make() writes into Entry tables that carry no foreign key back
689+
# here -- the fan-out ingestion pattern -- record the ingesting table and
690+
# key, which is what makes such a row traceable without one.
691+
from .jobs import _get_job_version
692+
693+
prov_token = provenance.set_ingesting(self.full_table_name, key, _get_job_version(self.connection._config))
694+
687695
try:
688696
if not is_generator:
689697
make(dict(key), **(make_kwargs or {}))
@@ -740,6 +748,7 @@ def _populate1(
740748
jobs.complete(key, duration=duration)
741749
return True
742750
finally:
751+
provenance.reset_ingesting(prov_token)
743752
self.__class__._allow_insert = False
744753
# Clear the per-make() upstream state: `_upstream = None` invalidates
745754
# the memoized Diagram; `_upstream_key = None` restores the "outside

‎src/datajoint/declare.py‎

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -521,6 +521,20 @@ def declare(
521521
job_metadata_sql = adapter.job_metadata_columns()
522522
attribute_sql.extend(job_metadata_sql)
523523

524+
# Add the hidden extrinsic-provenance slot to Entry tables, where rows enter
525+
# from outside the pipeline. Computed and Imported tables have no use for
526+
# it -- their provenance is entailed by the foreign-key graph -- and a part
527+
# inherits its master's.
528+
# Matched against the Manual tier itself, not by excluding the other tiers'
529+
# prefixes: enumerating exclusions makes every tier added later an Entry
530+
# table by default, which is how job tables (`~`) first acquired the slot.
531+
# Imported here rather than at module scope: user_tables imports table,
532+
# which imports this module.
533+
from .user_tables import Manual
534+
535+
if config.provenance.capture and re.fullmatch(Manual.tier_regexp, table_name):
536+
attribute_sql.extend(adapter.provenance_columns())
537+
524538
if not primary_key:
525539
# Singleton table: add hidden sentinel attribute
526540
primary_key = ["_singleton"]

‎src/datajoint/deploy.py‎

Lines changed: 113 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,8 @@
1212
``add_job_metadata_columns``, ``rebuild_lineage``.
1313
- :mod:`datajoint.deploy` — configure an environment for a consumer's
1414
requirements (CDC tools, replication, role grants, performance tuning).
15-
Cadence: re-runnable, idempotent. Examples: :func:`set_replica_identity`.
15+
Cadence: re-runnable, idempotent. Examples: :func:`set_replica_identity`,
16+
:func:`add_prov_column`.
1617
1718
Functions in this module should be safe to call repeatedly from a deploy hook
1819
without accumulating side effects.
@@ -183,3 +184,114 @@ def set_replica_identity(
183184
connection.query(ddl)
184185
result["tables_modified"] += 1
185186
return result
187+
188+
189+
def add_prov_column(target: "TargetType", dry_run: bool = True) -> dict:
190+
"""
191+
Add the hidden ``_prov`` attribute to Entry (``dj.Manual``) tables that lack it.
192+
193+
Capture defaults on, so tables declared from 2.3.4 onward already carry the
194+
slot. Two populations do not: tables declared before 2.3.4, and tables
195+
declared while ``config.provenance.capture`` was off. Inserts into those
196+
record nothing, silently, and this brings them in line.
197+
198+
It belongs here rather than in :mod:`datajoint.migrate` because it is not a
199+
one-shot correction of legacy state. It is idempotent — a table that already
200+
has the column is reported and left alone — and it stays useful for as long
201+
as capture can be turned off, which outlives the migration module.
202+
203+
Parameters
204+
----------
205+
target : Schema, Table class, or Table instance
206+
Given a Schema, every Entry table in it is processed.
207+
dry_run : bool, optional
208+
If True, report what would change without altering anything. Default True.
209+
210+
Returns
211+
-------
212+
dict
213+
``tables_analyzed``, ``tables_modified``, ``columns_added``, ``ddl``,
214+
and ``details`` — a per-table list of dicts.
215+
216+
Examples
217+
--------
218+
>>> from datajoint.deploy import add_prov_column
219+
>>> add_prov_column(schema, dry_run=True)["ddl"]
220+
>>> add_prov_column(schema, dry_run=False)["tables_modified"]
221+
222+
Notes
223+
-----
224+
- Only Entry tables are touched. Computed and Imported tables have no use
225+
for the slot — their provenance is entailed by the foreign-key graph — and
226+
a part table inherits its master's.
227+
- Rows already present keep ``NULL``. Provenance is recorded at insert and
228+
is never reconstructed after the fact.
229+
"""
230+
import re
231+
232+
from . import provenance
233+
from .schemas import _Schema
234+
from .table import Table
235+
from .user_tables import Manual
236+
237+
if isinstance(target, _Schema):
238+
connection = target.connection
239+
if connection is None or not target.database:
240+
raise DataJointError("Schema is not activated. Call schema.activate(...) before add_prov_column().")
241+
database = target.database
242+
table_names = list(target.list_tables())
243+
elif isinstance(target, type) and issubclass(target, Table):
244+
instance = target()
245+
connection = instance.connection
246+
if connection is None:
247+
raise DataJointError(f"Table {target.__name__} has no active connection.")
248+
database, table_names = instance.database, [instance.table_name]
249+
elif isinstance(target, Table):
250+
connection = target.connection
251+
if connection is None:
252+
raise DataJointError(f"Table {type(target).__name__} has no active connection.")
253+
database, table_names = target.database, [target.table_name]
254+
else:
255+
raise DataJointError(f"target must be a Schema or Table class/instance; got {type(target).__name__}")
256+
257+
if not database:
258+
raise DataJointError("Cannot add the provenance column: the target has no database.")
259+
260+
adapter = connection.adapter
261+
column_sql = adapter.provenance_columns()[0]
262+
263+
result: dict[str, Any] = {
264+
"tables_analyzed": 0,
265+
"tables_modified": 0,
266+
"columns_added": 0,
267+
"ddl": [],
268+
"details": [],
269+
}
270+
271+
for table_name in table_names:
272+
if not re.fullmatch(Manual.tier_regexp, table_name):
273+
continue
274+
result["tables_analyzed"] += 1
275+
276+
existing = {
277+
row[0]
278+
for row in connection.query(
279+
"SELECT COLUMN_NAME FROM information_schema.COLUMNS WHERE TABLE_SCHEMA = %s AND TABLE_NAME = %s",
280+
args=(database, table_name),
281+
).fetchall()
282+
}
283+
if provenance.PROV_ATTRIBUTE in existing:
284+
result["details"].append({"table": f"{database}.{table_name}", "status": "already_present"})
285+
continue
286+
287+
ddl = (
288+
f"ALTER TABLE {adapter.quote_identifier(database)}.{adapter.quote_identifier(table_name)} ADD COLUMN {column_sql}"
289+
)
290+
result["ddl"].append(ddl)
291+
result["details"].append({"table": f"{database}.{table_name}", "status": "pending" if dry_run else "added"})
292+
if not dry_run:
293+
connection.query(ddl)
294+
result["tables_modified"] += 1
295+
result["columns_added"] += 1
296+
297+
return result

‎src/datajoint/heading.py‎

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -397,7 +397,11 @@ def quote(name):
397397
return adapter.quote_identifier(name) if adapter else f'"{name}"'
398398

399399
def render_field(name):
400-
attr = self.attributes[name]
400+
# `attributes` hides underscore-prefixed names, so a caller that asks
401+
# for one by name -- copying `_prov` through an INSERT ... SELECT --
402+
# falls back to the full set. Default field lists are unaffected:
403+
# they are built from `attributes` and never contain hidden names.
404+
attr = self.attributes.get(name) or self._attributes[name]
401405
if attr.attribute_expression is None:
402406
return quote(name)
403407
else:

‎src/datajoint/provenance.py‎

Lines changed: 140 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,140 @@
1+
"""Extrinsic provenance for rows that enter the pipeline from outside.
2+
3+
Inside the pipeline, provenance is structural: a Computed table's row cannot
4+
exist unless its declared upstream exists, so the foreign-key graph *is* the
5+
lineage and nothing has to be recorded for it to hold.
6+
7+
At the boundary the structure runs out. Rows arrive in Entry tables from a
8+
person, an instrument, or a feed, and the framework has no way to say where
9+
they came from. This module supplies the slot and fills it.
10+
11+
The attribute is **framework-owned: no author ever writes it.** ``insert``
12+
takes no provenance argument, and its content comes from three places, none of
13+
them the call site:
14+
15+
* **configuration** -- ``config.provenance.source``, set per deployment, naming
16+
the external system this process draws from;
17+
* **ambient connection state** -- the connecting user, host and database, the
18+
insert time, and the code version;
19+
* **ambient execution state** -- the ingesting table and key, when the insert
20+
runs inside a ``make()``.
21+
22+
That ownership is the point. A field an operator can set is weaker evidence
23+
than one the system sets, and nothing is left for a pipeline to neglect.
24+
Anything an author wants to record deliberately belongs in the data model as a
25+
visible attribute, where queries can reach it.
26+
"""
27+
28+
import contextlib
29+
import contextvars
30+
import datetime
31+
import json
32+
from typing import Any
33+
34+
#: Name of the hidden attribute. Hidden attributes are excluded from
35+
#: ``heading.attributes``, so this never appears in a query heading.
36+
PROV_ATTRIBUTE = "_prov"
37+
38+
# Set by autopopulate around a make() call so that rows written to Entry tables
39+
# from inside an ingesting make() record what wrote them, which is what makes a
40+
# fanned-out row traceable without a foreign key.
41+
_ingesting: contextvars.ContextVar = contextvars.ContextVar("dj_ingesting", default=None)
42+
43+
44+
def set_ingesting(table_name, key, version=None):
45+
"""Record the ``make()`` now executing; returns a token for ``reset_ingesting``.
46+
47+
Parameters
48+
----------
49+
table_name : str
50+
Full table name of the ingesting table.
51+
key : dict
52+
The key ``make()`` was called with.
53+
version : str, optional
54+
Code version, as resolved for the job.
55+
"""
56+
value = {
57+
"table": table_name,
58+
"key": {k: _jsonable(v) for k, v in (key or {}).items()},
59+
}
60+
if version:
61+
value["version"] = version
62+
return _ingesting.set(value)
63+
64+
65+
def reset_ingesting(token):
66+
"""Restore the ingesting context saved by :func:`set_ingesting`."""
67+
if token is not None:
68+
_ingesting.reset(token)
69+
70+
71+
@contextlib.contextmanager
72+
def ingesting(table_name, key, version=None):
73+
"""Scope :func:`set_ingesting` to a block."""
74+
token = set_ingesting(table_name, key, version)
75+
try:
76+
yield
77+
finally:
78+
reset_ingesting(token)
79+
80+
81+
def _jsonable(value):
82+
"""Render a key value in a form ``json.dumps`` accepts."""
83+
if isinstance(value, (str, int, float, bool)) or value is None:
84+
return value
85+
if isinstance(value, (datetime.datetime, datetime.date, datetime.time)):
86+
return value.isoformat()
87+
if isinstance(value, bytes):
88+
return value.hex()
89+
return str(value)
90+
91+
92+
def build_payload(connection, config=None):
93+
"""Assemble the provenance record for rows inserted on this connection.
94+
95+
Returns ``None`` when there is nothing worth recording, so that a row is
96+
left with ``NULL`` rather than an empty object.
97+
"""
98+
if config is None:
99+
from .settings import config as _config
100+
101+
config = _config
102+
103+
payload: dict[str, Any] = {"time": datetime.datetime.now(datetime.timezone.utc).isoformat()}
104+
105+
conn_info = getattr(connection, "conn_info", None) or {}
106+
agent = {key: conn_info[key] for key in ("user", "host", "database_name") if conn_info.get(key) is not None}
107+
if agent:
108+
payload["agent"] = agent
109+
110+
try:
111+
from .jobs import _get_job_version
112+
113+
version = _get_job_version(getattr(connection, "_config", None) or config)
114+
except Exception: # version capture must never break an insert
115+
version = ""
116+
if version:
117+
payload["version"] = version
118+
119+
source = config.provenance.source
120+
if source:
121+
payload["source"] = source
122+
123+
context = _ingesting.get()
124+
if context:
125+
payload["context"] = context
126+
127+
# Time alone says nothing about origin; without any of the other three this
128+
# is noise rather than a record.
129+
return payload if len(payload) > 1 else None
130+
131+
132+
def serialize(payload):
133+
"""Render a payload for the ``json`` column.
134+
135+
``default=str`` because ``config.provenance.source`` is deployment-supplied
136+
and typed ``dict[str, Any]``: a ``date`` or a ``Path`` in it would otherwise
137+
raise from inside every insert into every Entry table, with an error naming
138+
neither provenance nor the setting that caused it.
139+
"""
140+
return json.dumps(payload, default=str)

0 commit comments

Comments
 (0)