task-5-report.md 16 KB

Task 5 Report — Governed Polars cross-source execution and artifact handoff

Date: 2026-07-23 Branch: codex/data-rule-execution-m3a-m5 Implementation commits before the third remediation: 69d83ed, f4d798b, 8bffdb6, c152e74. The third-remediation commit is the commit containing this report revision.

8bffdb6 is the final hash of the second remediation commit. It replaces the intermediate a17bcd0 after restoring the historical M3A acceptance document byte-for-byte; no production or test behavior changed in that amend.

Lifecycle boundary: physical Polars plans are created as compiled; Task 5 does not add a production promotion API. The real integration inserts a published fixture only to exercise the existing Runner publication gate.

Final outcome

Task 5 now provides a fail-closed PostgreSQL/MySQL → Parquet/MinIO → isolated Polars → cataloged Parquet path through the production repository, resolver, executor, HTTP token, and durable task-ledger combination.

The final implementation includes:

  • closed JSON Polars plans and shared compile/runtime semantic validation;
  • typed cast targets carrying nullability, Decimal precision/scale, and timestamptz timezone;
  • real rule-timezone semantics for date(), timestamp(), timestamp cast, and timestamptz cast;
  • server-owned Parquet refs, digest validation, TTL enforcement, and bounded metadata;
  • streamed MinIO staging to temporary files and trusted PyArrow footer preflight before Polars scans;
  • all allocation-heavy Polars scans, expression evaluation, regex, join, group/aggregate, counts, and output write in an isolated spawn child;
  • parent RSS monitoring, child OS RLIMIT_AS, bounded execution time, and fail-closed worker errors;
  • durable pending → object upload → ready catalog handoff keyed by (correlation_id, binding_id, artifact_kind), with the exact attested binding_hash;
  • bounded bidirectional reconciliation for pending, ready, missing, invalid, and orphaned objects, plus the existing expiry cleanup;
  • minimal public artifact results without full schema fields;
  • standards-valid CycloneDX 1.5 runtime SBOM with deterministic offline schema validation.

The bound plan contains canonical RuleVersion, SchemaSnapshot, and DatasetBinding hashes plus dataops-polars-1.42.1 provenance. It contains no Python source, pickle, callable, module, caller-supplied path, arbitrary URL, secret, or serialized Polars internal plan.

Second-remediation RED/GREEN evidence

The second review was implemented as fail-first slices. Representative RED evidence captured during the work:

  • Immutable handoff tests failed before the catalog uniqueness changed from digest-based uniqueness to the correlation/binding/kind invariant and before retry cleanup/reuse existed.
  • Public-metadata tests failed while MinIO metadata still carried the base64 full schema contract and public results still exposed schema_fields.
  • Cleanup tests failed before a bounded production catalog cleanup entrypoint deleted both the MinIO object and the exact PostgreSQL row.
  • Compression-bomb and tempfile-staging tests failed before PyArrow footer row-group preflight and disk streaming existed.
  • Worker isolation tests failed before app/runner/polars_worker.py existed; the first adapter reconstruction run then exposed one schema-order failure (1 failed) because canonical fields are name-sorted while Parquet retains physical column order. Validation was corrected to compare the closed field set and exact per-field types.
  • Typed-cast RED: test_polars_cast_operations_carry_complete_typed_targets failed (1 failed) with missing target_field.
  • Timezone RED: test_polars_expression_date_and_timestamp_use_plan_timezone failed (1 failed) because timestamp() produced timezone-aware output for a timezone-free timestamp contract.
  • The final full-suite pre-close run found one stale test monkeypatch (1 failed, 553 passed) after in-memory BytesIO staging was removed; the regression check was changed to assert the stage implementation contains no BytesIO.
  • The first Ruff closeout found two style-only failures (unused io, nested with); both were fixed before final verification.

Second-remediation GREEN:

  • Focused compiler/adapter/worker/artifact/bootstrap/schema/SBOM command: 57 passed in 4.17s.
  • Final artifact/worker/runtime/SBOM recheck after the last portability and timestamp-cast changes: 30 passed in 3.92s.
  • Final real cross-source + HTTP/ledger integration: 1 passed in 2.03s.
  • Final full suite: 554 passed, 26 skipped, 59 subtests passed in 8.39s.
  • Ruff across every changed Python file: All checks passed!.
  • git diff --check: passed.
  • docker compose -f deploy/docker/docker-compose.yml config -q: passed.

The third remediation supersedes the second-remediation migration and handoff implementation described below wherever they conflict.

Third-remediation RED/GREEN evidence

The third review was also implemented in fail-first slices:

  • schema source tests failed until committed migration 140 was restored byte-for-byte from f4d798b and the state change moved to new migration 20260723_150;
  • reservation-order tests failed until the resolver inserted a durable pending row before the first MinIO PUT;
  • crash-injection tests failed until reserve, upload, and finalize had distinct failure semantics and commit-acknowledgement loss was rechecked;
  • stale-binding tests failed until reservation and finalization locked and attested the canonical binding in the same database transaction;
  • reconciliation tests failed until missing pending rows, valid pending objects, missing ready objects, invalid objects, old orphans, fresh objects, and unsafe keys had separate bounded policies;
  • worker-start injection initially leaked raw startup errors and pipe/process handles; it now returns a safe error and closes every created endpoint in a unified finally.

Third-remediation GREEN:

  • focused schema/artifact/handoff/bootstrap/worker/adapter plus both real integrations: 49 passed in 5.62s;
  • isolated temporary PostgreSQL migration test: real old 140 schema, old-format row insertion, upgrade to head, migrated ready row and binding_hash, same-digest retry, different-digest CAS conflict: passed;
  • real source PostgreSQL + source MySQL + platform PostgreSQL + MinIO + isolated Polars + Flask HTTP + durable ledger: passed;
  • final full suite: 564 passed, 26 skipped, 59 subtests passed in 9.01s;
  • Ruff across every changed Python file: passed;
  • git diff --check: passed;
  • docker compose -f deploy/docker/docker-compose.yml config -q: passed;
  • git diff f4d798b -- migrations/versions/20260723_140_rule_run_artifacts.py: empty.

Durable artifact handoff, attestation, and reconciliation

Migration 20260723_140 is historical and unchanged. It retains its original unique constraint:

UNIQUE (correlation_id, binding_id, artifact_digest)

New forward-only migration 20260723_150:

  • adds binding_hash, handoff_status, ready_at, failed_at, failure_code, and updated_at;
  • migrates old rows to ready with the current canonical binding hash;
  • refuses unattested or conflicting old rows;
  • discovers and drops only the exact old-140 digest uniqueness constraint;
  • adds named uniqueness on (correlation_id, binding_id, artifact_kind) and a reconciliation index.

PostgresArtifactResolver.publish_path implements the durable protocol:

  1. validate the worker path and reserve a server-generated exact object key;
  2. in one transaction, lock the canonical binding with FOR SHARE, attest its hash/object/access contract, and insert the exact pending catalog row;
  3. upload only to that reserved key and revalidate the stored object;
  4. in a new transaction, re-lock and re-attest the binding, then atomically finalize only the matching pending row to ready.

Resolution sees only ready rows whose persisted binding hash still equals the canonical binding hash. Same-triple/same-digest retries reuse the ready object without uploading; different digests fail before upload. Upload failure removes the exact object and pending row. A lost finalize acknowledgement is rechecked: confirmed ready is success; unconfirmed state returns commit_outcome="unknown" and deliberately does not delete the possibly cataloged object.

reconcile-rule-artifacts --limit --grace-seconds processes only aged, bounded candidates. It deletes missing pending rows, finalizes valid pending objects only while the binding still matches, marks ready rows failed when their object is missing or invalid, deletes invalid pending objects safely, and deletes only old safe-key orphans after the grace period. Fresh objects, referenced objects, and malformed/unowned keys are retained.

cleanup-rule-artifacts --limit uses a bounded FOR UPDATE SKIP LOCKED selection, removes each expired store object, and deletes its exact catalog row in the same database transaction scope. Unit tests cover the CLI and object+row cleanup behavior.

The local database had an earlier applied empty copy of migration 140 with a manually altered constraint. It was explicitly not accepted as migration evidence. After verifying zero rows, only that empty table was rebuilt by rewinding its Alembic marker to 130 and running the formal 130 → original-140 → 150 chain. Independently, the automated acceptance test creates and drops a fresh temporary database and proves old140 → head without manual ALTER.

Resource enforcement and worker isolation

Pre-allocation artifact checks

ArtifactStore.stage:

  1. validates the server-owned ref, object stat, content type, compressed size, metadata digest, row count, schema digest, and TTL;
  2. streams MinIO response chunks to a NamedTemporaryFile while hashing;
  3. validates the downloaded digest;
  4. reads the Parquet footer with pyarrow.parquet.ParquetFile;
  5. rejects excessive footer row count or summed row-group uncompressed size;
  6. checks only the lazy physical schema in the main process;
  7. yields the path and always removes the temporary file.

Full schema contracts are no longer base64-encoded into MinIO metadata. MinIO metadata contains only bounded digest/count/expiry/size fields, with a 2,048-character aggregate ceiling. PostgreSQL remains the canonical store for the complete schema contract.

Isolated execution

The adapter stages paths but does not collect input frames. A spawn child:

  • warms Polars before establishing its runtime baseline;
  • applies RLIMIT_AS relative to baseline virtual memory;
  • validates the plan again;
  • scans input and lookup Parquet;
  • executes every allowlisted operation;
  • validates output row/type/nullability contracts;
  • writes the output Parquet path.

The parent independently watches incremental RSS with psutil, kills the child on limit breach, enforces a hard timeout, and returns only bounded metrics. Worker PID is deliberately not exposed in the public result.

Pipe creation, process construction, process.start(), message handling, and resource enforcement now share one guarded lifecycle. Startup errors are sanitized, both pipe endpoints are closed when created, and a started process is killed/joined/closed from the unified finally.

Tests prove:

  • worker PID differs from the Runner PID;
  • a 128 MiB allocation fails under an 8 MiB allowance;
  • a highly compressed Parquet payload is rejected from footer bounds before scan exposure;
  • regex expansion, high-cardinality group/aggregate, and wide lookup join peak allocations fail inside the isolated worker.

Portability boundary: spawn and psutil are portable, but the OS hard limit requires Unix resource.RLIMIT_AS. Linux/macOS are supported targets. On a platform without RLIMIT_AS, the worker fails closed instead of silently running without an OS memory boundary.

Typed casts and timezone semantics

Each compiled cast operation now contains an exact target_field. The shared validator rejects a missing field, an unknown/extra shape, a target name/type mismatch, incomplete Decimal precision/scale, or missing timestamptz timezone. The worker no longer falls back to inferring cast details from output fields.

Executable tests verify:

  • string → Decimal(12, 2);
  • string → Datetime(time_zone="Asia/Shanghai");
  • offset timestamp → rule-local timezone-free timestamp;
  • date(offset_value) converts to the rule timezone before selecting the calendar date;
  • timestamp(offset_value) converts to the rule timezone and removes timezone metadata to match the platform's timezone-free timestamp type.

Real production combination

The integration discovers runtime values from deploy/docker/docker-compose.yml; credentials are not copied into source.

  • source PostgreSQL: live customer input;
  • source MySQL: live segment lookup;
  • MinIO: input, lookup, and output Parquet objects;
  • platform PostgreSQL: snapshots, bindings, published plan, artifact catalog, and durable task ledger;
  • Polars spawn worker: normalize → lookup join → assert → deduplicate;
  • Flask /v1/tasks/execute: signed short-lived task token;
  • PostgresTaskLedger: single-use claim and committed success record.

Evidence:

  • 4 input rows → 2 output rows;
  • 1 assertion reject;
  • 1 deduplicated row;
  • exact enriched values verified after reread;
  • direct repeat reuses the stable output ref;
  • conflicting output digest fails before a new object is uploaded;
  • HTTP result returns the same stable output ref;
  • replaying the same token returns HTTP 409;
  • durable ledger records success / committed;
  • public result contains no schema_fields;
  • exactly three catalog rows and three correlation-scoped MinIO objects remain before cleanup.

The test finally removes only its exact correlation prefix, source tables, ledger record, catalog entries, plan/binding/deployment/rule rows, and schema snapshots, and asserts the MinIO prefix is empty.

CycloneDX 1.5 SBOM

docs/security/data-rule-runtime-sbom.json now uses standard CycloneDX fields:

  • licenses;
  • externalReferences;
  • properties;
  • component bom-ref and dependency relationships.

Pinned runtime components:

  • polars==1.42.1 — MIT;
  • polars-runtime-32==1.42.1 — MIT;
  • pyarrow==21.0.0 — Apache-2.0;
  • psutil==5.9.8 — BSD-3-Clause;
  • minio==7.2.10 — Apache-2.0.

The test validates the document with jsonschema against vendored official CycloneDX 1.5, SPDX, and JSF schemas. The official files' SHA-256 hashes are asserted before validation, so the check is deterministic and offline.

Key files in the final remediation

Production:

  • app/core/data_rules/compilers/polars.py
  • app/runner/artifacts.py
  • app/runner/polars_worker.py
  • app/runner/rule_polars.py
  • app/runner/bootstrap.py
  • migrations/versions/20260723_140_rule_run_artifacts.py
  • migrations/versions/20260723_150_rule_artifact_handoff_state.py
  • requirements.txt
  • docs/security/data-rule-runtime-sbom.json
  • docs/security/cyclonedx-1.5-schema/

Tests:

  • tests/core/data_rules/test_polars_compiler.py
  • tests/runner/test_artifacts.py
  • tests/runner/test_artifact_handoff.py
  • tests/runner/test_polars_worker.py
  • tests/runner/test_rule_polars.py
  • tests/runner/test_bootstrap.py
  • tests/integration/test_data_rule_polars_execution.py
  • tests/integration/test_rule_artifact_migration_upgrade.py
  • tests/test_data_rule_runtime_sbom.py
  • tests/test_data_rule_schema.py

Remaining boundaries

  • Task 7 still owns test-evidence recording and plan publication. RulePlanExecutor continues to refuse compiled plans in the published runtime path.
  • quality_check remains fail-closed through the existing adapter; Task 5 does not claim an artifact-backed quality-check implementation.
  • ArtifactStore.read remains a compatibility/helper API that can return a bounded in-process frame. The production Polars adapter does not call it; it uses staged files plus the isolated worker.
  • Windows requires a separate Job Object implementation before Polars batch execution can be enabled there; the current implementation fails closed when RLIMIT_AS is unavailable.