================== Usage ================== This page covers the query and tagging API exposed by :py:class:`~herostools.actor.archiver.PostgresArchiver` and its base class :py:class:`~herostools.actor.archiver.PostgresDatabase`. All methods are callable locally or via :py:class:`~heros.RemoteHERO`. Artifact storage ---------------- Class responsibilities: - ``PostgresDatabase`` provides SQL record and tag primitives. - ``PostgresRecordStore`` provides current record state with artifact handling. It has no HEROS subscription, queue, or PostgreSQL notification connection. - ``PostgresArchiver`` receives HEROS events and publishes change notifications. Use ``PostgresRecordStore`` from analysis software that queries current state directly: .. code:: python from herostools.actor.archiver import PostgresRecordStore store = PostgresRecordStore( db_url="postgresql://heros:heros@db/herosdb", artifact_storage_kwargs={ "endpoint_url": "http://rustfs:9000", "bucket": "artifacts", "access_key": "analysis-reader", "secret_key": "...", "create_buckets": False, }, ) try: records = store.get(["shot_0001"], resolve_artifacts=False) image_ref = records["shot_0001"]["image"] image = store.get_artifact(image_ref["artifact_id"], bucket=image_ref.get("bucket")) finally: store.close() This client uses PostgreSQL and S3 directly. It does not open a HEROS/Zenoh connection or subscribe to PostgreSQL record-change notifications. The S3 credential requires object-read access. ``create_buckets=False`` prevents bucket validation and creation during initialization. ``get()`` resolves artifact references to numpy arrays by default (``resolve_artifacts=True``). Pass ``resolve_artifacts=False`` to retrieve the raw reference dict without an S3 round-trip: .. code:: python # resolved (default) data = archiver.get(["shot_0001"])["shot_0001"] print(data["image"]) # numpy array # raw reference data_raw = archiver.get(["shot_0001"], resolve_artifacts=False)["shot_0001"] print(data_raw["image"]) # {"__type__": "artifact_ref", "artifact_id": "..."} Pass ``size_limit_bytes`` (bytes) to skip fetching large artifacts and leave them as raw reference dicts instead. Only artifacts whose stored ``size_bytes`` is at or below the limit are fetched: .. code:: python # resolve only artifacts up to 1 MB data = archiver.get(["shot_0001"], size_limit_bytes=1_000_000) print(data["shot_0001"]["thumbnail"]) # numpy array (<= 1 MB) print(data["shot_0001"]["raw_scan"]) # artifact_ref dict (> 1 MB) ``size_bytes`` reflects the array's element data size (``ndarray.nbytes``) and excludes the fixed ~128-byte npy format header. Records stored before this feature was introduced have no ``size_bytes`` and are always resolved. New S3 artifacts are written to UTC daily buckets named ``-YYYY-MM-DD``. The bucket stored in each reference is used for retrieval and purge; the configured static bucket is used only when a legacy reference has no bucket field. Pass ``use_single_bucket=True`` when constructing the archiver to write new artifacts to the configured static bucket instead. See :doc:`schema` for details on how artifact references are stored. Array key template ------------------ By default a bare numpy-array payload is stored in ``data_json`` under the source name (e.g. ``camera``). When the same source emits one array per frame this key is overwritten on every event. Pass ``array_key_template`` to derive a unique key per event from the payload metadata: .. code:: python archiver = PostgresArchiver( db_url="postgresql://heros:heros@localhost/herosdb", array_key_template="{{ source_name }}_frame_{{ '%04d' % frame }}", ) The template is a `Jinja2 `_ expression rendered against the per-event metadata (merged with ``default_metadata``). ``source_name`` is always injected into the template context, so it is available regardless of what the metadata contains. For source ``camera`` and ``frame=42`` the key above becomes ``camera_frame_0042``. Using the default template ``"{{ source_name }}"``, the key is just the source name. If a variable referenced in the template is absent from the context, the archiver drops the record and logs a warning. A missing ``identifier`` has the same result. BOSS JSON example: .. code:: json { "_id": "my-camera-archiver", "classname": "herostools.archiver.PostgresArchiver", "arguments": { "db_url": "postgresql://heros:heros@localhost/herosdb", "object_selector": "my-camera", "event_name": "acquisition_data", "array_key_template": "{{ source_name }}_frame_{{ '%04d' % frame }}" } } Querying records ---------------- **get(ids)** Fetch the merged record for a list of identifiers. The result is ``dict[str, dict[str, Any]]`` keyed by existing identifiers. Missing requested identifiers are omitted, so partial results are expected. Do not rely on result iteration order to match input order. Duplicate input identifiers produce one result entry, and an empty input list returns ``{}``. Each inner dict combines ``data_json`` fields with post-hoc tags from ``identifier_groups``. Identifiers with no tags are returned without tag fields. If a tag key collides with a ``data_json`` key, the tag wins. Post-hoc annotations take priority over ingest-time payload fields. .. code:: python result = archiver.get(["shot_0001", "shot_0002"]) data = result["shot_0001"] print(data["temperature"]) # scalar from data_json print(data["image"]) # numpy array resolved from artifact storage print(data.get("discard")) # post-hoc tag, or None if not set # skip artifact resolution for a metadata-only fetch result_raw = archiver.get(["shot_0001"], resolve_artifacts=False) print(result_raw["shot_0001"]["image"]) # {"__type__": "artifact_ref", ...} **record_changed event** ``PostgresArchiver`` publishes a ``record_changed`` event whenever one or more rows in ``records`` are inserted or updated, or when tags are written through ``tag_ids()``. The payload is a deduplicated list of affected identifiers. A 50 ms batching window combines a burst of writes into one event. Subscribers receive the list and call ``get()`` to fetch the current state, including any updated tags: .. code:: python def on_record_changed(source_name, identifiers): data = archiver.get(identifiers) for identifier, record in data.items(): ... # add or update table row The event is driven by two PostgreSQL triggers defined in ``provision.sql``: ``records_notify`` fires on every INSERT or UPDATE to ``records``; ``identifier_groups_notify`` fires on every INSERT or UPDATE to ``identifier_groups``. Both call the same ``notify_record_changed()`` function and require no application-level code. They work across concurrent archiver instances. **get_ids(tag, values)** Return identifiers where ``tag`` matches any of the given ``values``, searching both ``identifier_groups`` (post-hoc tags) and the top-level keys of ``records.data_json`` (payload fields). Values are cast to strings before comparison. Returns a sorted, deduplicated list. .. code:: python shots = archiver.get_ids("run_id", [10, 23]) Post-hoc tagging ---------------- The ``identifier_groups`` table lets you attach labels to identifiers after they have been archived -- for example to mark a subset as belonging to a particular measurement campaign or quality tier. Use :py:meth:`~herostools.actor.archiver.PostgresDatabase.tag_ids` to tag a batch of identifiers in one call. Calling it again with the same tag key overwrites the previous value: .. code:: python archiver.tag_ids(["shot_0001", "shot_0002"], "measurement_campaign", "demo_2026") # re-tag a single identifier archiver.tag_ids(["shot_0001"], "measurement_campaign", "demo_2027") The primary key ``(identifier, tag_key)`` enforces one value per tag key per identifier. Tags written via ``tag_ids()`` are immediately visible in ``get()`` and trigger a ``record_changed`` event so connected subscribers refresh automatically. On reconnect, calling ``get(ids)`` returns the full merged record including all tags set in previous sessions -- no separate tag fetch is required. Retrieve identifiers by tag: .. code:: python shots = archiver.get_ids("measurement_campaign", ["demo_2026"]) # ["shot_0001", "shot_0002"] Fetch records including their tags in one call: .. code:: python data = archiver.get(shots) print(data["shot_0001"]["measurement_campaign"]) # "demo_2026" print(data["shot_0001"]["discard"]) # tag value if set Direct database queries ----------------------- ``get_ids`` covers equality and membership checks. For anything more expressive -- range filters, nested key access, multi-condition joins -- query the database directly with psycopg. ``data_json`` is a JSONB column, so the full PostgreSQL JSONB operator set is available. **Range filter on a scalar field** The ``->>`` operator extracts a top-level key as text; cast it to the target type before comparing: .. code:: python import psycopg with psycopg.connect("postgresql://heros:heros@localhost/herosdb") as conn: with conn.cursor() as cur: cur.execute( "SELECT identifier FROM records" " WHERE (data_json->>'temperature')::float BETWEEN %s AND %s" " ORDER BY identifier", (low, high), ) ids = [row[0] for row in cur.fetchall()] .. note:: Rows where ``temperature`` is absent or not castable to ``float`` are silently excluded by the cast. Add ``AND data_json ? 'temperature'`` to make the presence check explicit if needed. **Combining a payload filter with a post-hoc tag** .. code:: python with psycopg.connect("postgresql://heros:heros@localhost/herosdb") as conn: with conn.cursor() as cur: cur.execute( "SELECT r.identifier FROM records r" " JOIN identifier_groups g ON g.identifier = r.identifier" " WHERE g.tag_key = 'measurement_campaign'" " AND g.tag_value = %s" " AND (r.data_json->>'temperature')::float > %s" " ORDER BY r.identifier", (campaign, min_temp), ) ids = [row[0] for row in cur.fetchall()] Purging records --------------- ``purge(ids)`` deletes rows from both ``records`` and ``identifier_groups``, and by default also removes all referenced artifact objects from the configured storage backend. It requires ``allow_purge=True`` at instantiation time and raises :py:exc:`PermissionError` otherwise. .. code:: python archiver.purge(["shot_0001", "shot_0002"]) # DB rows + artifacts archiver.purge(["shot_0001"], remove_artifacts=False) # DB rows only .. warning:: Artifact objects whose ``artifact_id`` was overwritten by a later write (i.e. intermediate arrays for identifiers that received multiple events) are not tracked and will not be removed by ``purge``. Only the artifact referenced by the current ``data_json`` entry is deleted. Retention and cleanup --------------------- ``cleanup()`` deletes all records whose ``updated_at`` is older than a given number of days, removing associated artifacts from storage by default. Records tagged with ``_save = '1'`` (or ``'true'`` / ``'True'``) are exempt and will never be deleted by ``cleanup()``. Like ``purge()``, it requires ``allow_purge=True`` at construction time. ``cleanup()`` defaults to ``dry_run=True``: it returns the list of expired identifiers without deleting anything. Pass ``dry_run=False`` to actually purge: .. code:: python # inspect what would be deleted (no changes made) expired = archiver.cleanup(retention_days=30) print(expired) # ["shot_0001", "shot_0002"] # delete DB rows and artifacts archiver.cleanup(retention_days=30, dry_run=False) # delete DB rows only, keep artifacts in storage archiver.cleanup(retention_days=30, dry_run=False, remove_artifacts=False) Set a default at construction time so call sites need no argument: .. code:: python db = PostgresDatabase( db_url="postgresql://heros:heros@localhost/herosdb", allow_purge=True, retention_days=30, ) db.cleanup(dry_run=False) # uses the 30-day default ``retention_days=0`` (the default) means keep forever -- ``cleanup()`` returns ``[]`` unless an explicit value is provided at call time. Mark a record as permanent: .. code:: python archiver.tag_ids(["shot_0001"], "_save", "1") To run cleanup automatically once a day via a Docker Compose stack see :ref:`auto-cleanup` in the setup guide.