| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931193219331934193519361937193819391940194119421943194419451946194719481949195019511952195319541955195619571958195919601961196219631964196519661967196819691970197119721973197419751976197719781979198019811982198319841985198619871988198919901991199219931994199519961997199819992000200120022003200420052006200720082009201020112012201320142015201620172018201920202021202220232024202520262027202820292030203120322033203420352036203720382039204020412042204320442045204620472048204920502051205220532054205520562057205820592060206120622063206420652066206720682069207020712072207320742075207620772078207920802081208220832084208520862087208820892090209120922093209420952096209720982099210021012102210321042105210621072108210921102111211221132114211521162117211821192120212121222123212421252126212721282129213021312132213321342135213621372138213921402141214221432144214521462147214821492150215121522153215421552156215721582159216021612162216321642165216621672168216921702171217221732174217521762177217821792180218121822183218421852186218721882189219021912192219321942195219621972198219922002201220222032204220522062207220822092210221122122213221422152216221722182219222022212222222322242225222622272228222922302231223222332234 |
- from __future__ import annotations
- import hashlib
- import json
- import re
- from datetime import UTC, datetime, timedelta
- from pathlib import Path
- import polars as pl
- import pytest
- from minio import Minio
- from sqlalchemy import create_engine, text
- from app.core.common.identifiers import new_governance_uid
- from app.core.data_rules.contracts import rule_spec_hash, validate_rule_spec
- from app.core.data_rules.execution_contracts import canonical_schema_hash
- COMPOSE = (
- Path(__file__).resolve().parents[2]
- / "deploy"
- / "docker"
- / "docker-compose.yml"
- )
- def _compose_value(pattern):
- source = COMPOSE.read_text(encoding="utf-8")
- match = re.search(pattern, source, flags=re.DOTALL)
- assert match is not None
- return match.group(1)
- def _schema(schema_ref, fields):
- normalized = [
- {"name": name, "type": field_type, "nullable": nullable}
- for name, field_type, nullable in fields
- ]
- return {
- "id": new_governance_uid(),
- "schema_ref": schema_ref,
- "schema_hash": canonical_schema_hash(normalized),
- "fields": normalized,
- "source_revision": "task5:real-cross-source",
- }
- def _binding(schema, *, source_uid, access_mode, object_ref):
- return {
- "id": new_governance_uid(),
- "data_source_uid": source_uid,
- "object_kind": "parquet_artifact",
- "object_ref": object_ref,
- "schema_snapshot_id": schema["id"],
- "access_mode": access_mode,
- "dialect": "parquet",
- "write_mode": "append",
- }
- def test_real_postgres_mysql_minio_polars_cross_source_execution(tmp_path):
- from app.core.data_rules.compilers.polars import PolarsRuleCompiler
- from app.runner.artifacts import ArtifactStore, PostgresArtifactResolver
- from app.runner.rule_polars import PolarsRulePlanAdapter
- from app.runner.rules import PostgresRulePlanRepository, RulePlanExecutor
- source_user = _compose_value(
- r"source-postgres:.*?POSTGRES_USER:\s*([^\s]+)"
- )
- source_password = _compose_value(
- r"source-postgres:.*?POSTGRES_PASSWORD:\s*([^\s]+)"
- )
- platform_user = _compose_value(
- r"\n postgres:.*?POSTGRES_USER:\s*([^\s]+)"
- )
- platform_password = _compose_value(
- r"\n postgres:.*?POSTGRES_PASSWORD:\s*([^\s]+)"
- )
- postgres_port = _compose_value(r'"(25432):5432"')
- mysql_port = _compose_value(r'"(23306):3306"')
- minio_user = _compose_value(r"MINIO_ROOT_USER:\s*([^\s]+)")
- minio_password = _compose_value(r"MINIO_ROOT_PASSWORD:\s*([^\s]+)")
- minio_port = _compose_value(r'"(19000):9000"')
- platform_port = _compose_value(r'"(15432):5432"')
- bucket = _compose_value(r"mc mb --ignore-existing local/([^\s]+)")
- postgres = create_engine(
- f"postgresql+psycopg2://{source_user}:{source_password}"
- f"@127.0.0.1:{postgres_port}/acceptance",
- pool_pre_ping=True,
- )
- mysql = create_engine(
- f"mysql+pymysql://{source_user}:{source_password}"
- f"@127.0.0.1:{mysql_port}/acceptance",
- pool_pre_ping=True,
- )
- platform = create_engine(
- f"postgresql+psycopg2://{platform_user}:{platform_password}"
- f"@127.0.0.1:{platform_port}/dataops",
- pool_pre_ping=True,
- )
- minio = Minio(
- f"127.0.0.1:{minio_port}",
- access_key=minio_user,
- secret_key=minio_password,
- secure=False,
- )
- store = ArtifactStore(
- minio,
- bucket=bucket,
- max_artifact_bytes=4 * 1024 * 1024,
- max_rows=1_000,
- memory_limit_bytes=256 * 1024 * 1024,
- max_ttl_seconds=3600,
- )
- correlation_id = new_governance_uid()
- failure_correlation_id = new_governance_uid()
- unknown_correlation_id = new_governance_uid()
- lease_correlation_id = new_governance_uid()
- sample_crash_correlation_id = new_governance_uid()
- receipt_correlation_id = new_governance_uid()
- failed_receipt_correlation_id = new_governance_uid()
- no_run_crash_correlation_id = new_governance_uid()
- evidence_crash_correlation_id = new_governance_uid()
- active_retry_correlation_id = new_governance_uid()
- sql_binding_id = new_governance_uid()
- prefix = f"rules/{correlation_id}/"
- customer_table = "task5_polars_customers"
- segment_table = "task5_polars_segments"
- rule_uid = new_governance_uid()
- rule_id = new_governance_uid()
- downstream_rule_uid = new_governance_uid()
- downstream_rule_id = new_governance_uid()
- dataflow_uid = new_governance_uid()
- dataflow_version_id = new_governance_uid()
- deployment_id = new_governance_uid()
- component_binding_id = new_governance_uid()
- plan_id = new_governance_uid()
- downstream_component_binding_id = new_governance_uid()
- downstream_plan_id = new_governance_uid()
- ledger_jti = None
- retry_ledger_jti = None
- no_run_crash_jti = None
- evidence_crash_jti = None
- active_retry_jti = None
- gateway_candidate_id = None
- try:
- with postgres.begin() as connection:
- connection.execute(text(f"DROP TABLE IF EXISTS {customer_table}"))
- connection.execute(
- text(
- f"CREATE TABLE {customer_table} ("
- "customer_id BIGINT NOT NULL, "
- "name VARCHAR(100), mobile VARCHAR(30), "
- "segment_code VARCHAR(20), version_no BIGINT NOT NULL)"
- )
- )
- connection.execute(
- text(
- f"INSERT INTO {customer_table} "
- "(customer_id, name, mobile, segment_code, version_no) "
- "VALUES "
- "(1, ' Alice ', '13800138000', 'A', 1), "
- "(1, ' Alice Updated ', '13800138000', 'A', 2), "
- "(2, ' Bad ', 'invalid', 'B', 1), "
- "(3, ' Carol ', '13900139000', 'C', 1)"
- )
- )
- with mysql.begin() as connection:
- connection.execute(text(f"DROP TABLE IF EXISTS {segment_table}"))
- connection.execute(
- text(
- f"CREATE TABLE {segment_table} ("
- "code VARCHAR(20) PRIMARY KEY, "
- "segment_name VARCHAR(100) NOT NULL)"
- )
- )
- connection.execute(
- text(
- f"INSERT INTO {segment_table} (code, segment_name) "
- "VALUES ('A', 'Gold'), ('B', 'Basic'), ('C', 'Silver')"
- )
- )
- with postgres.connect() as connection:
- customer_rows = [
- dict(row)
- for row in connection.execute(
- text(
- f"SELECT customer_id, name, mobile, "
- f"segment_code, version_no FROM {customer_table}"
- )
- ).mappings()
- ]
- with mysql.connect() as connection:
- segment_rows = [
- dict(row)
- for row in connection.execute(
- text(
- f"SELECT code, segment_name FROM {segment_table}"
- )
- ).mappings()
- ]
- input_schema = _schema(
- "bd:task5:customer:raw",
- [
- ("customer_id", "integer", False),
- ("name", "string", True),
- ("mobile", "string", True),
- ("segment_code", "string", True),
- ("version_no", "integer", False),
- ],
- )
- lookup_schema = _schema(
- "bd:task5:segment:lookup",
- [
- ("code", "string", False),
- ("segment_name", "string", False),
- ],
- )
- output_schema = _schema(
- "bd:task5:customer:enriched",
- [
- ("customer_id", "integer", False),
- ("name", "string", True),
- ("mobile", "string", True),
- ("segment_code", "string", True),
- ("version_no", "integer", False),
- ("segment_name", "string", True),
- ],
- )
- input_binding = _binding(
- input_schema,
- source_uid=new_governance_uid(),
- access_mode="read",
- object_ref="postgres-customer-artifact",
- )
- lookup_binding = _binding(
- lookup_schema,
- source_uid=new_governance_uid(),
- access_mode="read",
- object_ref="mysql-segment-artifact",
- )
- output_binding = _binding(
- output_schema,
- source_uid=new_governance_uid(),
- access_mode="read_write",
- object_ref="polars-output-artifact",
- )
- downstream_output_binding = _binding(
- output_schema,
- source_uid=new_governance_uid(),
- access_mode="write",
- object_ref="polars-downstream-output-artifact",
- )
- spec = validate_rule_spec(
- {
- "schema_version": "2.0",
- "rule_uid": new_governance_uid(),
- "name": "task5_real_cross_source",
- "input_schema_ref": input_schema["schema_ref"],
- "output_schema_ref": output_schema["schema_ref"],
- "steps": [
- {
- "id": "normalize_name",
- "op": "normalize_text",
- "column": "name",
- "trim": True,
- },
- {
- "id": "join_segment",
- "op": "lookup_join",
- "lookup": {
- "binding_id": lookup_binding["id"],
- "left_on": ["segment_code"],
- "right_on": ["code"],
- "select": {
- "segment_name": "segment_name"
- },
- "how": "left",
- },
- },
- {
- "id": "valid_mobile",
- "op": "assert",
- "expression": "matches(mobile, '^[0-9]{11}$')",
- "on_failure": "reject",
- "severity": "error",
- },
- {
- "id": "latest_customer",
- "op": "deduplicate",
- "keys": ["customer_id"],
- "order_by": ["version_no"],
- "keep": "last",
- },
- ],
- "null_policy": "explicit",
- "timezone": "Asia/Shanghai",
- }
- )
- rule = {
- "id": rule_id,
- "status": "published",
- "rule_spec": spec,
- "spec_hash": rule_spec_hash(spec),
- }
- compiled = PolarsRuleCompiler().compile(
- rule_version=rule,
- input_schema=input_schema,
- output_schema=output_schema,
- input_binding=input_binding,
- output_binding=output_binding,
- backend={
- "max_rows": 1_000,
- "max_artifact_bytes": 4 * 1024 * 1024,
- "memory_limit_bytes": 256 * 1024 * 1024,
- "masking_policies": {},
- "lookup_bindings": {
- lookup_binding["id"]: {
- "binding": lookup_binding,
- "schema": lookup_schema,
- }
- },
- },
- )
- downstream_spec = validate_rule_spec(
- {
- "schema_version": "2.0",
- "rule_uid": new_governance_uid(),
- "name": "task6_real_artifact_handoff",
- "input_schema_ref": output_schema["schema_ref"],
- "output_schema_ref": output_schema["schema_ref"],
- "steps": [
- {
- "id": "normalize_downstream_name",
- "op": "normalize_text",
- "column": "name",
- "trim": True,
- }
- ],
- "null_policy": "explicit",
- "timezone": "Asia/Shanghai",
- }
- )
- downstream_rule = {
- "id": downstream_rule_id,
- "status": "published",
- "rule_spec": downstream_spec,
- "spec_hash": rule_spec_hash(downstream_spec),
- }
- downstream_compiled = PolarsRuleCompiler().compile(
- rule_version=downstream_rule,
- input_schema=output_schema,
- output_schema=output_schema,
- input_binding=output_binding,
- output_binding=downstream_output_binding,
- backend={
- "max_rows": 1_000,
- "max_artifact_bytes": 4 * 1024 * 1024,
- "memory_limit_bytes": 256 * 1024 * 1024,
- "masking_policies": {},
- "lookup_bindings": {},
- },
- )
- lookup_operation = compiled["plan"]["operations"][1]
- schema_hashes = {
- "rule_spec_hash": compiled["plan"]["rule_spec_hash"],
- "input_schema_snapshot_id": input_schema["id"],
- "input_schema_hash": input_schema["schema_hash"],
- "output_schema_snapshot_id": output_schema["id"],
- "output_schema_hash": output_schema["schema_hash"],
- }
- with platform.begin() as connection:
- for schema in (input_schema, lookup_schema, output_schema):
- connection.execute(
- text(
- """
- INSERT INTO public.data_schema_snapshots
- (id, schema_ref, schema_hash, fields, source_revision)
- VALUES (CAST(:id AS uuid), :schema_ref, :schema_hash,
- CAST(:fields AS jsonb), :source_revision)
- """
- ),
- {**schema, "fields": json.dumps(schema["fields"])},
- )
- connection.execute(
- text(
- """
- INSERT INTO public.data_rules
- (id, rule_uid, name, category, status)
- VALUES (CAST(:id AS uuid), CAST(:rule_uid AS uuid),
- :name, 'general', 'active')
- """
- ),
- {
- "id": new_governance_uid(),
- "rule_uid": rule_uid,
- "name": spec["name"],
- },
- )
- connection.execute(
- text(
- """
- INSERT INTO public.data_rules
- (id, rule_uid, name, category, status)
- VALUES (CAST(:id AS uuid), CAST(:rule_uid AS uuid),
- :name, 'general', 'active')
- """
- ),
- {
- "id": new_governance_uid(),
- "rule_uid": downstream_rule_uid,
- "name": downstream_spec["name"],
- },
- )
- connection.execute(
- text(
- """
- INSERT INTO public.data_rule_versions
- (id, rule_uid, version_no, source_text, source_language,
- rule_spec, spec_hash, generated_kind, status, published_at)
- VALUES (CAST(:id AS uuid), CAST(:rule_uid AS uuid), 1,
- :source_text, 'en', CAST(:rule_spec AS jsonb),
- :spec_hash, 'polars', 'published',
- CURRENT_TIMESTAMP)
- """
- ),
- {
- "id": rule_id,
- "rule_uid": rule_uid,
- "source_text": "Task 5 real cross-source integration",
- "rule_spec": json.dumps(spec),
- "spec_hash": rule["spec_hash"],
- },
- )
- connection.execute(
- text(
- """
- INSERT INTO public.data_rule_versions
- (id, rule_uid, version_no, source_text,
- source_language, rule_spec, spec_hash,
- generated_kind, status, published_at)
- VALUES (
- CAST(:id AS uuid), CAST(:rule_uid AS uuid), 1,
- :source_text, 'en', CAST(:rule_spec AS jsonb),
- :spec_hash, 'polars', 'published',
- CURRENT_TIMESTAMP
- )
- """
- ),
- {
- "id": downstream_rule_id,
- "rule_uid": downstream_rule_uid,
- "source_text": (
- "Task 6 real two-node artifact handoff"
- ),
- "rule_spec": json.dumps(downstream_spec),
- "spec_hash": downstream_rule["spec_hash"],
- },
- )
- connection.execute(
- text(
- """
- INSERT INTO public.dataflow_versions
- (id, dataflow_uid, version_no, name, dataflow_spec,
- input_schema_hashes, output_schema_hash, status,
- released_at)
- VALUES (CAST(:id AS uuid), CAST(:dataflow_uid AS uuid), 1,
- :name, '{}'::jsonb,
- CAST(:input_schema_hashes AS jsonb),
- :output_schema_hash, 'released', CURRENT_TIMESTAMP)
- """
- ),
- {
- "id": dataflow_version_id,
- "dataflow_uid": dataflow_uid,
- "name": "Task 5 real cross-source integration",
- "input_schema_hashes": json.dumps(
- [
- input_schema["schema_hash"],
- lookup_schema["schema_hash"],
- ]
- ),
- "output_schema_hash": output_schema["schema_hash"],
- },
- )
- connection.execute(
- text(
- """
- INSERT INTO public.dataflow_deployments
- (id, dataflow_version_id, environment, deployment_config,
- status)
- VALUES (CAST(:id AS uuid),
- CAST(:dataflow_version_id AS uuid), 'test',
- '{}'::jsonb, 'disabled')
- """
- ),
- {
- "id": deployment_id,
- "dataflow_version_id": dataflow_version_id,
- },
- )
- for logical_ref, binding, binding_hash in (
- (
- "customers",
- input_binding,
- compiled["plan"]["input_binding_hash"],
- ),
- (
- "segments",
- lookup_binding,
- lookup_operation["lookup_binding_hash"],
- ),
- (
- "enriched",
- output_binding,
- compiled["plan"]["output_binding_hash"],
- ),
- (
- "downstream",
- downstream_output_binding,
- downstream_compiled["plan"][
- "output_binding_hash"
- ],
- ),
- ):
- connection.execute(
- text(
- """
- INSERT INTO public.dataflow_dataset_bindings
- (id, dataflow_deployment_id, logical_ref,
- data_source_uid, object_kind, object_ref,
- schema_snapshot_id, dialect, access_mode, write_mode,
- binding_hash)
- VALUES (CAST(:id AS uuid), CAST(:deployment_id AS uuid),
- :logical_ref, CAST(:source_uid AS uuid),
- 'parquet_artifact', :object_ref,
- CAST(:schema_snapshot_id AS uuid), 'parquet',
- :access_mode, 'append', :binding_hash)
- """
- ),
- {
- "id": binding["id"],
- "deployment_id": deployment_id,
- "logical_ref": logical_ref,
- "source_uid": binding["data_source_uid"],
- "object_ref": binding["object_ref"],
- "schema_snapshot_id": binding["schema_snapshot_id"],
- "access_mode": binding["access_mode"],
- "binding_hash": binding_hash,
- },
- )
- connection.execute(
- text(
- """
- INSERT INTO public.dataflow_component_bindings
- (id, dataflow_version_id, component_id, component_kind,
- rule_version_id, stage, order_no, idempotency, provenance)
- VALUES (CAST(:id AS uuid),
- CAST(:dataflow_version_id AS uuid),
- 'task5_real_polars', 'rule.apply',
- CAST(:rule_version_id AS uuid), 'transform', 0,
- CAST(:idempotency AS jsonb), '{}'::jsonb)
- """
- ),
- {
- "id": component_binding_id,
- "dataflow_version_id": dataflow_version_id,
- "rule_version_id": rule_id,
- "idempotency": json.dumps(
- {
- "strategy": "deduplication_key",
- "key": "customer_id",
- }
- ),
- },
- )
- connection.execute(
- text(
- """
- INSERT INTO public.dataflow_component_bindings
- (id, dataflow_version_id, component_id,
- component_kind, rule_version_id, stage, order_no,
- idempotency, provenance)
- VALUES (
- CAST(:id AS uuid),
- CAST(:dataflow_version_id AS uuid),
- 'task6_real_handoff', 'rule.apply',
- CAST(:rule_version_id AS uuid), 'transform', 1,
- CAST(:idempotency AS jsonb), '{}'::jsonb
- )
- """
- ),
- {
- "id": downstream_component_binding_id,
- "dataflow_version_id": dataflow_version_id,
- "rule_version_id": downstream_rule_id,
- "idempotency": json.dumps(
- {
- "strategy": "deduplication_key",
- "key": "customer_id",
- }
- ),
- },
- )
- connection.execute(
- text(
- """
- INSERT INTO public.rule_execution_plans
- (id, component_binding_id, backend, compiler_version, plan,
- plan_hash, schema_hashes, status)
- VALUES (CAST(:id AS uuid),
- CAST(:component_binding_id AS uuid),
- 'polars_batch', :compiler_version,
- CAST(:plan AS jsonb), :plan_hash,
- CAST(:schema_hashes AS jsonb), 'published')
- """
- ),
- {
- "id": plan_id,
- "component_binding_id": component_binding_id,
- "compiler_version": compiled["compiler_version"],
- "plan": json.dumps(compiled["plan"]),
- "plan_hash": compiled["plan_hash"],
- "schema_hashes": json.dumps(schema_hashes),
- },
- )
- connection.execute(
- text(
- """
- INSERT INTO public.rule_execution_plans
- (id, component_binding_id, backend,
- compiler_version, plan, plan_hash,
- schema_hashes, status)
- VALUES (
- CAST(:id AS uuid),
- CAST(:component_binding_id AS uuid),
- 'polars_batch', :compiler_version,
- CAST(:plan AS jsonb), :plan_hash,
- CAST(:schema_hashes AS jsonb), 'published'
- )
- """
- ),
- {
- "id": downstream_plan_id,
- "component_binding_id": (
- downstream_component_binding_id
- ),
- "compiler_version": downstream_compiled[
- "compiler_version"
- ],
- "plan": json.dumps(downstream_compiled["plan"]),
- "plan_hash": downstream_compiled["plan_hash"],
- "schema_hashes": json.dumps(
- {
- "rule_spec_hash": downstream_compiled[
- "plan"
- ]["rule_spec_hash"],
- "input_schema_snapshot_id": output_schema[
- "id"
- ],
- "input_schema_hash": output_schema[
- "schema_hash"
- ],
- "output_schema_snapshot_id": output_schema[
- "id"
- ],
- "output_schema_hash": output_schema[
- "schema_hash"
- ],
- }
- ),
- },
- )
- customer_path = tmp_path / "customers.parquet"
- segment_path = tmp_path / "segments.parquet"
- pl.DataFrame(customer_rows).write_parquet(customer_path)
- pl.DataFrame(segment_rows).write_parquet(segment_path)
- resolver = PostgresArtifactResolver(platform, store)
- customer_artifact = resolver.publish_path(
- str(customer_path),
- binding_id=input_binding["id"],
- binding_hash=compiled["plan"]["input_binding_hash"],
- correlation_id=correlation_id,
- kind="input",
- ttl_seconds=900,
- schema_fields=input_schema["fields"],
- limits=compiled["plan"]["resource_limits"],
- )
- segment_artifact = resolver.publish_path(
- str(segment_path),
- binding_id=lookup_binding["id"],
- binding_hash=lookup_operation["lookup_binding_hash"],
- correlation_id=correlation_id,
- kind="lookup",
- ttl_seconds=900,
- schema_fields=lookup_schema["fields"],
- limits=compiled["plan"]["resource_limits"],
- )
- node = {
- "id": "task5_real_polars",
- "type": "rule.apply",
- "purpose": "write",
- "idempotency": {
- "strategy": "deduplication_key",
- "key": "customer_id",
- },
- "config": {
- "component_binding_id": component_binding_id,
- "rule_version_id": rule["id"],
- "execution_plan_hash": compiled["plan_hash"],
- },
- }
- from app.core.mcp.gateway import SchedulingGateway
- from app.core.mcp.identity import AgentIdentity
- from app.core.mcp.persistence import PostgresSchedulingPlanStore
- from app.runner.api import create_runner_app
- from app.runner.auth import TaskTokenIssuer, TaskTokenVerifier
- from app.runner.ledger import PostgresTaskLedger
- from app.runner.nodes import NodeRegistry
- from app.runner.rule_evidence import PostgresRuleEvidenceWriter
- class Audit:
- def __init__(self):
- self.events = []
- def record(self, event):
- self.events.append(event)
- class CanaryEngine:
- def __init__(self):
- self.calls = []
- def deploy_disabled(self, _definition):
- return {"revision": "task6-real-gateway"}
- def activate(self, namespace, flow_id):
- self.calls.append(("activate", namespace, flow_id))
- def execute(self, namespace, flow_id, inputs=None, **_options):
- self.calls.append(
- ("execute", namespace, flow_id, dict(inputs or {}))
- )
- return {"id": f"task6-canary-{correlation_id}"}
- def deactivate(self, namespace, flow_id):
- self.calls.append(("deactivate", namespace, flow_id))
- task_secret = "task5-real-http-secret-value-32-bytes"
- verifier = TaskTokenVerifier(task_secret)
- gateway_engine = CanaryEngine()
- gateway_store = PostgresSchedulingPlanStore(platform)
- gateway = SchedulingGateway(
- plans=gateway_store,
- audit=Audit(),
- engine=gateway_engine,
- token_issuer=TaskTokenIssuer(task_secret),
- )
- gateway_identity = AgentIdentity(
- subject="task6-real-scheduler",
- roles=frozenset({"scheduler"}),
- business_domains=frozenset({"sales"}),
- environments=frozenset({"test"}),
- correlation_id=correlation_id,
- )
- candidate = gateway.create_candidate_plan(
- gateway_identity,
- business_domain="sales",
- environment="test",
- workflow_spec={
- "schema_version": "1.0",
- "dataflow_uid": dataflow_uid,
- "name": "Task 6 governed rule canary",
- "nodes": [node],
- "edges": [],
- "parameters": {},
- },
- schedule_plan={
- "schema_version": "1.0",
- "timezone": "Asia/Shanghai",
- "triggers": [{"type": "manual"}],
- "max_concurrency": 1,
- "conflict_policy": "skip",
- "timeout_seconds": 600,
- "retry": {
- "max_attempts": 1,
- "delay_seconds": 1,
- },
- "backfill": {"max_days": 1, "max_runs": 1},
- },
- )
- gateway_candidate_id = candidate["candidate_id"]
- with platform.begin() as connection:
- connection.execute(
- text(
- """
- UPDATE public.dataflow_workflow_versions
- SET write_authorized = TRUE
- WHERE id = CAST(:id AS uuid)
- """
- ),
- {"id": gateway_candidate_id},
- )
- deployed = gateway.deploy_disabled_version(
- gateway_identity,
- candidate_id=gateway_candidate_id,
- business_domain="sales",
- environment="test",
- )
- assert deployed["deployment_id"] == deployment_id
- assert deployed["candidate_id"] != deployed["deployment_id"]
- gateway.run_canary(
- gateway_identity,
- candidate_id=gateway_candidate_id,
- business_domain="sales",
- environment="test",
- inputs={},
- )
- execution_call = next(
- call for call in gateway_engine.calls if call[0] == "execute"
- )
- task_token = execution_call[3]["dataops_task_tokens"][node["id"]]
- issued_claims = verifier.verify(task_token, node=node)
- assert issued_claims.deployment_id == deployment_id
- assert issued_claims.workflow_version == 1
- assert issued_claims.environment == "test"
- executor = RulePlanExecutor(
- PostgresRulePlanRepository(platform),
- adapters={
- "polars_batch": PolarsRulePlanAdapter(
- artifact_store=store,
- artifact_resolver=resolver,
- artifact_ttl_seconds=900,
- )
- },
- )
- result = executor.execute(
- node,
- {},
- write_authorized=True,
- correlation_id=correlation_id,
- )
- repeated = executor.execute(
- node,
- {},
- write_authorized=True,
- correlation_id=correlation_id,
- )
- assert result["rows_in"] == 4
- assert result["rows_out"] == 2
- assert result["rows_rejected"] == 1
- assert result["rows_deduplicated"] == 1
- assert result["rows_filtered"] == 0
- assert result["rows_join_dropped"] == 0
- assert result["rows_aggregated"] == 0
- assert result["violation_count"] == 1
- assert result["violations"] == [
- {"step_id": "valid_mobile", "count": 1}
- ]
- output = store.read(
- result["artifact_ref"],
- result["digest"],
- expected_schema_fields=output_schema["fields"],
- limits=compiled["plan"]["resource_limits"],
- ).collect()
- assert output.sort("customer_id").to_dicts() == [
- {
- "customer_id": 1,
- "mobile": "13800138000",
- "name": "Alice Updated",
- "segment_code": "A",
- "segment_name": "Gold",
- "version_no": 2,
- },
- {
- "customer_id": 3,
- "mobile": "13900139000",
- "name": "Carol",
- "segment_code": "C",
- "segment_name": "Silver",
- "version_no": 1,
- },
- ]
- assert all(
- item.object_name.startswith(prefix)
- for item in minio.list_objects(
- bucket, prefix=prefix, recursive=True
- )
- )
- assert repeated["rows_out"] == 2
- assert repeated["artifact_ref"] == result["artifact_ref"]
- assert "schema_fields" not in result
- ledger_jti = verifier.verify(task_token, node=node).jti
- retry_token = TaskTokenIssuer(task_secret).issue(
- task_uid=new_governance_uid(),
- dataflow_uid=dataflow_uid,
- deployment_id=deployment_id,
- environment="test",
- workflow_version=1,
- correlation_id=correlation_id,
- node=node,
- write_authorized=True,
- )
- retry_ledger_jti = verifier.verify(
- retry_token, node=node
- ).jti
- real_evidence_writer = PostgresRuleEvidenceWriter(
- platform,
- store,
- sample_ttl_seconds=900,
- )
- evidenced_executor = RulePlanExecutor(
- PostgresRulePlanRepository(platform),
- adapters={
- "polars_batch": PolarsRulePlanAdapter(
- artifact_store=store,
- artifact_resolver=resolver,
- artifact_ttl_seconds=900,
- )
- },
- evidence_writer=real_evidence_writer,
- )
- real_ledger = PostgresTaskLedger(platform, lease_seconds=30)
- runner_app = create_runner_app(
- verifier=verifier,
- ledger=real_ledger,
- registry=NodeRegistry({"rule.apply": evidenced_executor}),
- )
- no_run_crash_token = TaskTokenIssuer(task_secret).issue(
- task_uid=new_governance_uid(),
- dataflow_uid=dataflow_uid,
- deployment_id=deployment_id,
- environment="test",
- workflow_version=1,
- correlation_id=no_run_crash_correlation_id,
- node=node,
- write_authorized=True,
- )
- no_run_claims = verifier.verify(no_run_crash_token, node=node)
- no_run_crash_jti = no_run_claims.jti
- evidence_crash_token = TaskTokenIssuer(task_secret).issue(
- task_uid=new_governance_uid(),
- dataflow_uid=dataflow_uid,
- deployment_id=deployment_id,
- environment="test",
- workflow_version=1,
- correlation_id=evidence_crash_correlation_id,
- node=node,
- write_authorized=True,
- )
- evidence_crash_claims = verifier.verify(
- evidence_crash_token,
- node=node,
- )
- evidence_crash_jti = evidence_crash_claims.jti
- active_retry_token = TaskTokenIssuer(task_secret).issue(
- task_uid=new_governance_uid(),
- dataflow_uid=dataflow_uid,
- deployment_id=deployment_id,
- environment="test",
- workflow_version=1,
- correlation_id=active_retry_correlation_id,
- node=node,
- write_authorized=True,
- )
- active_retry_claims = verifier.verify(
- active_retry_token,
- node=node,
- )
- active_retry_jti = active_retry_claims.jti
- def task_binding(claims):
- return {
- "task_uid": claims.task_uid,
- "dataflow_uid": claims.dataflow_uid,
- "deployment_id": claims.deployment_id,
- "environment": claims.environment,
- "workflow_version": claims.workflow_version,
- "correlation_id": claims.correlation_id,
- "node_id": claims.node_id,
- "node_type": claims.node_type,
- "data_source_uid": None,
- "idempotency_key": "customer_id",
- }
- assert real_ledger.claim(
- no_run_crash_jti,
- task_binding(no_run_claims),
- expires_at=no_run_claims.expires_at,
- )
- assert real_ledger.claim(
- evidence_crash_jti,
- task_binding(evidence_crash_claims),
- expires_at=evidence_crash_claims.expires_at,
- )
- assert real_ledger.claim(
- active_retry_jti,
- task_binding(active_retry_claims),
- expires_at=active_retry_claims.expires_at,
- )
- real_evidence_writer.start(
- component_binding_id=component_binding_id,
- rule_version_id=rule_id,
- plan_hash=compiled["plan_hash"],
- correlation_id=evidence_crash_correlation_id,
- dataflow_uid=dataflow_uid,
- deployment_id=deployment_id,
- environment="test",
- workflow_version=1,
- node_id=node["id"],
- lease_owner=evidence_crash_jti,
- )
- real_evidence_writer.start(
- component_binding_id=component_binding_id,
- rule_version_id=rule_id,
- plan_hash=compiled["plan_hash"],
- correlation_id=active_retry_correlation_id,
- dataflow_uid=dataflow_uid,
- deployment_id=deployment_id,
- environment="test",
- workflow_version=1,
- node_id=node["id"],
- lease_owner=active_retry_jti,
- )
- with platform.begin() as connection:
- connection.execute(
- text(
- """
- UPDATE public.runner_task_executions
- SET lease_expires_at =
- CURRENT_TIMESTAMP - INTERVAL '1 second'
- WHERE token_jti IN (
- CAST(:no_run_jti AS uuid),
- CAST(:evidence_jti AS uuid)
- )
- """
- ),
- {
- "no_run_jti": no_run_crash_jti,
- "evidence_jti": evidence_crash_jti,
- },
- )
- connection.execute(
- text(
- """
- UPDATE public.rule_runs
- SET lease_expires_at =
- CURRENT_TIMESTAMP - INTERVAL '1 second'
- WHERE lease_owner = CAST(:jti AS uuid)
- """
- ),
- {"jti": evidence_crash_jti},
- )
- with runner_app.test_client() as client:
- http_result = client.post(
- "/v1/tasks/execute",
- json={
- "task_token": task_token,
- "node": node,
- "parameters": {},
- },
- )
- replay = client.post(
- "/v1/tasks/execute",
- json={
- "task_token": task_token,
- "node": node,
- "parameters": {},
- },
- )
- retried = client.post(
- "/v1/tasks/execute",
- json={
- "task_token": retry_token,
- "node": node,
- "parameters": {},
- },
- )
- no_run_crash = client.post(
- "/v1/tasks/execute",
- json={
- "task_token": no_run_crash_token,
- "node": node,
- "parameters": {},
- },
- )
- evidence_crash = client.post(
- "/v1/tasks/execute",
- json={
- "task_token": evidence_crash_token,
- "node": node,
- "parameters": {},
- },
- )
- active_retry = client.post(
- "/v1/tasks/execute",
- json={
- "task_token": active_retry_token,
- "node": node,
- "parameters": {},
- },
- )
- assert http_result.status_code == 200, http_result.get_json()
- assert http_result.get_json()["output_artifact"] == result[
- "artifact_ref"
- ]
- assert http_result.get_json()["result"]["artifact_ref"] == result[
- "artifact_ref"
- ]
- assert replay.status_code == 200
- assert replay.headers["X-Idempotent-Replay"] == "true"
- assert replay.get_json() == http_result.get_json()
- assert retried.status_code == 200
- assert retried.get_json()["output_artifact"] == result[
- "artifact_ref"
- ]
- assert no_run_crash.status_code == 409
- assert no_run_crash.get_json() == {
- "error": "task execution outcome is unknown"
- }
- assert evidence_crash.status_code == 409
- assert evidence_crash.get_json() == {
- "error": "task execution outcome is unknown"
- }
- assert active_retry.status_code == 202
- assert active_retry.headers["Retry-After"] == "2"
- assert real_ledger.get(active_retry_jti).status == "running"
- assert real_ledger.get(active_retry_jti).replay_body is None
- assert real_ledger.get(no_run_crash_jti).status == "unknown"
- assert real_ledger.get(evidence_crash_jti).status == "unknown"
- with platform.connect() as connection:
- assert connection.execute(
- text(
- """
- SELECT COUNT(*)
- FROM public.rule_runs
- WHERE correlation_id =
- CAST(:correlation_id AS uuid)
- """
- ),
- {"correlation_id": no_run_crash_correlation_id},
- ).scalar_one() == 0
- crash_evidence = connection.execute(
- text(
- """
- SELECT status, commit_outcome
- FROM public.rule_runs
- WHERE correlation_id =
- CAST(:correlation_id AS uuid)
- """
- ),
- {"correlation_id": evidence_crash_correlation_id},
- ).mappings().one()
- assert dict(crash_evidence) == {
- "status": "unknown",
- "commit_outcome": "unknown",
- }
- terminal_replay = real_evidence_writer.reconcile_expired_lease(
- lease_owner=ledger_jti,
- deployment_id=deployment_id,
- correlation_id=correlation_id,
- component_binding_id=component_binding_id,
- rule_version_id=rule_id,
- plan_hash=compiled["plan_hash"],
- )
- assert terminal_replay["state"] == "terminal"
- assert terminal_replay["status"] == "success"
- assert terminal_replay["result_digest"] == hashlib.sha256(
- json.dumps(
- terminal_replay["result"],
- sort_keys=True,
- separators=(",", ":"),
- ensure_ascii=False,
- ).encode("utf-8")
- ).hexdigest()
- assert re.fullmatch(
- r"[0-9a-f]{64}",
- terminal_replay["evidence_digest"],
- )
- ledger_record = PostgresTaskLedger(platform).get(ledger_jti)
- assert ledger_record is not None
- assert ledger_record.status == "success"
- assert ledger_record.commit_outcome == "committed"
- with platform.connect() as connection:
- run_evidence = connection.execute(
- text(
- """
- SELECT status, commit_outcome, rows_in, rows_out,
- rows_rejected, rows_quarantined, public_result
- FROM public.rule_runs
- WHERE correlation_id = CAST(:correlation_id AS uuid)
- AND component_binding_id =
- CAST(:component_binding_id AS uuid)
- """
- ),
- {
- "correlation_id": correlation_id,
- "component_binding_id": component_binding_id,
- },
- ).mappings().one()
- sample_evidence = connection.execute(
- text(
- """
- SELECT s.artifact_ref, s.artifact_digest,
- s.sample_count, s.redaction_policy,
- s.expires_at, s.handoff_status
- FROM public.rule_violation_samples s
- JOIN public.rule_runs r ON r.id = s.rule_run_id
- WHERE r.correlation_id =
- CAST(:correlation_id AS uuid)
- """
- ),
- {"correlation_id": correlation_id},
- ).mappings().one()
- assert run_evidence["status"] == "success"
- assert run_evidence["commit_outcome"] == "committed"
- assert run_evidence["rows_in"] == 4
- assert run_evidence["rows_out"] == 2
- assert run_evidence["rows_rejected"] == 1
- assert run_evidence["rows_quarantined"] == 0
- assert run_evidence["public_result"]["output_artifact"] == result[
- "artifact_ref"
- ]
- assert sample_evidence["artifact_digest"]
- assert sample_evidence["sample_count"] == 1
- assert (
- sample_evidence["redaction_policy"]
- == "rule-violation-default-v1"
- )
- assert sample_evidence["handoff_status"] == "ready"
- assert store.read(
- sample_evidence["artifact_ref"],
- sample_evidence["artifact_digest"],
- expected_schema_fields=[
- {
- "name": name,
- "type": "string",
- "nullable": True,
- }
- for name in (
- "customer_id",
- "mobile",
- "name",
- "segment_code",
- "segment_name",
- "version_no",
- )
- ],
- ).collect().to_dicts() == [
- {
- "customer_id": "[REDACTED]",
- "mobile": "[REDACTED]",
- "name": "[REDACTED]",
- "segment_code": "[REDACTED]",
- "segment_name": "[REDACTED]",
- "version_no": "[REDACTED]",
- }
- ]
- orphan = store.write(
- pl.DataFrame({"value": ["orphan"]}),
- correlation_id,
- 900,
- schema_fields=[
- {
- "name": "value",
- "type": "string",
- "nullable": True,
- }
- ],
- )
- original_clock = store.clock
- store.clock = lambda: datetime.now(UTC) + timedelta(seconds=60)
- try:
- reconciliation = resolver.reconcile(
- limit=20,
- grace_seconds=30,
- )
- finally:
- store.clock = original_clock
- assert reconciliation["orphans_deleted"] >= 1
- assert store.describe_optional(orphan["artifact_ref"]) is None
- assert store.describe(sample_evidence["artifact_ref"])[
- "digest"
- ] == sample_evidence["artifact_digest"]
- downstream_node = {
- "id": "task6_real_handoff",
- "type": "rule.apply",
- "purpose": "write",
- "idempotency": {
- "strategy": "deduplication_key",
- "key": "customer_id",
- },
- "config": {
- "component_binding_id": (
- downstream_component_binding_id
- ),
- "rule_version_id": downstream_rule_id,
- "execution_plan_hash": downstream_compiled[
- "plan_hash"
- ],
- },
- }
- downstream_result = evidenced_executor.execute(
- downstream_node,
- {"input_artifact": result["artifact_ref"]},
- write_authorized=True,
- correlation_id=correlation_id,
- dataflow_uid=dataflow_uid,
- deployment_id=deployment_id,
- environment="test",
- workflow_version=1,
- node_id="task6_real_handoff",
- task_jti=new_governance_uid(),
- )
- assert downstream_result["rows_in"] == 2
- assert downstream_result["rows_out"] == 2
- assert downstream_result["output_artifact"] != result[
- "artifact_ref"
- ]
- assert store.read(
- downstream_result["artifact_ref"],
- downstream_result["digest"],
- expected_schema_fields=output_schema["fields"],
- limits=downstream_compiled["plan"]["resource_limits"],
- ).collect().sort("customer_id").to_dicts() == (
- output.sort("customer_id").to_dicts()
- )
- replayed_downstream = evidenced_executor.execute(
- downstream_node,
- {"input_artifact": result["artifact_ref"]},
- write_authorized=True,
- correlation_id=correlation_id,
- dataflow_uid=dataflow_uid,
- deployment_id=deployment_id,
- environment="test",
- workflow_version=1,
- node_id="task6_real_handoff",
- task_jti=new_governance_uid(),
- )
- assert replayed_downstream["artifact_ref"] == downstream_result[
- "artifact_ref"
- ]
- with platform.connect() as connection:
- assert connection.execute(
- text(
- """
- SELECT COUNT(*)
- FROM public.rule_runs
- WHERE correlation_id =
- CAST(:correlation_id AS uuid)
- """
- ),
- {"correlation_id": correlation_id},
- ).scalar_one() == 2
- from app.runner.nodes import NodeExecutionError
- class FailingAdapter:
- def execute(self, **_kwargs):
- raise NodeExecutionError(
- "safe downstream failure",
- commit_outcome="not_committed",
- )
- failed_executor = RulePlanExecutor(
- PostgresRulePlanRepository(platform),
- adapters={"polars_batch": FailingAdapter()},
- evidence_writer=PostgresRuleEvidenceWriter(
- platform,
- store,
- sample_ttl_seconds=900,
- ),
- )
- with pytest.raises(NodeExecutionError, match="safe downstream"):
- failed_executor.execute(
- node,
- {},
- write_authorized=True,
- correlation_id=failure_correlation_id,
- dataflow_uid=dataflow_uid,
- deployment_id=deployment_id,
- environment="test",
- workflow_version=1,
- node_id="task5_real_polars",
- task_jti=new_governance_uid(),
- )
- with platform.connect() as connection:
- failed_evidence = connection.execute(
- text(
- """
- SELECT status, commit_outcome
- FROM public.rule_runs
- WHERE correlation_id =
- CAST(:correlation_id AS uuid)
- """
- ),
- {"correlation_id": failure_correlation_id},
- ).mappings().one()
- assert dict(failed_evidence) == {
- "status": "failed",
- "commit_outcome": "not_committed",
- }
- class UnknownAdapter:
- def execute(self, **_kwargs):
- raise NodeExecutionError(
- "safe uncertain commit",
- commit_outcome="unknown",
- )
- unknown_executor = RulePlanExecutor(
- PostgresRulePlanRepository(platform),
- adapters={"polars_batch": UnknownAdapter()},
- evidence_writer=PostgresRuleEvidenceWriter(
- platform,
- store,
- sample_ttl_seconds=900,
- ),
- )
- with pytest.raises(NodeExecutionError, match="uncertain commit"):
- unknown_executor.execute(
- node,
- {},
- write_authorized=True,
- correlation_id=unknown_correlation_id,
- dataflow_uid=dataflow_uid,
- deployment_id=deployment_id,
- environment="test",
- workflow_version=1,
- node_id="task5_real_polars",
- task_jti=new_governance_uid(),
- )
- with platform.connect() as connection:
- unknown_evidence = connection.execute(
- text(
- """
- SELECT status, commit_outcome
- FROM public.rule_runs
- WHERE correlation_id =
- CAST(:correlation_id AS uuid)
- """
- ),
- {"correlation_id": unknown_correlation_id},
- ).mappings().one()
- assert dict(unknown_evidence) == {
- "status": "unknown",
- "commit_outcome": "unknown",
- }
- evidence_writer = PostgresRuleEvidenceWriter(
- platform,
- store,
- sample_ttl_seconds=900,
- lease_seconds=30,
- )
- lease_owner = new_governance_uid()
- lease_run_id = evidence_writer.start(
- component_binding_id=component_binding_id,
- rule_version_id=rule_id,
- plan_hash=compiled["plan_hash"],
- correlation_id=lease_correlation_id,
- dataflow_uid=dataflow_uid,
- deployment_id=deployment_id,
- environment="test",
- workflow_version=1,
- node_id="task5_real_polars",
- lease_owner=lease_owner,
- )
- with pytest.raises(ValueError, match="not owned"):
- evidence_writer.heartbeat(
- lease_run_id,
- new_governance_uid(),
- )
- evidence_writer.heartbeat(lease_run_id, lease_owner)
- with platform.begin() as connection:
- connection.execute(
- text(
- """
- UPDATE public.rule_runs
- SET lease_expires_at =
- CURRENT_TIMESTAMP - INTERVAL '1 second'
- WHERE id = CAST(:id AS uuid)
- """
- ),
- {"id": lease_run_id},
- )
- assert evidence_writer.start(
- component_binding_id=component_binding_id,
- rule_version_id=rule_id,
- plan_hash=compiled["plan_hash"],
- correlation_id=lease_correlation_id,
- dataflow_uid=dataflow_uid,
- deployment_id=deployment_id,
- environment="test",
- workflow_version=1,
- node_id="task5_real_polars",
- lease_owner=new_governance_uid(),
- ) == lease_run_id
- assert evidence_writer.replay(lease_run_id)["status"] == "unknown"
- class FailOnceConnection:
- def __init__(self, connection, state):
- self.connection = connection
- self.state = state
- def execute(self, statement, parameters=None):
- sql = str(statement)
- if (
- self.state["armed"]
- and "UPDATE public.rule_runs" in sql
- and "rows_in = :rows_in" in sql
- ):
- self.state["armed"] = False
- raise RuntimeError("simulated response loss")
- return self.connection.execute(statement, parameters)
- class FailOnceBegin:
- def __init__(self, context, state):
- self.context = context
- self.state = state
- def __enter__(self):
- return FailOnceConnection(
- self.context.__enter__(),
- self.state,
- )
- def __exit__(self, *args):
- return self.context.__exit__(*args)
- class FailOnceEngine:
- def __init__(self, engine):
- self.engine = engine
- self.state = {"armed": True}
- def connect(self):
- return self.engine.connect()
- def begin(self):
- return FailOnceBegin(
- self.engine.begin(),
- self.state,
- )
- sample_crash_writer = PostgresRuleEvidenceWriter(
- FailOnceEngine(platform),
- store,
- sample_ttl_seconds=900,
- )
- sample_crash_run_id = sample_crash_writer.start(
- component_binding_id=component_binding_id,
- rule_version_id=rule_id,
- plan_hash=compiled["plan_hash"],
- correlation_id=sample_crash_correlation_id,
- dataflow_uid=dataflow_uid,
- deployment_id=deployment_id,
- environment="test",
- workflow_version=1,
- node_id="task5_real_polars",
- lease_owner=new_governance_uid(),
- )
- sample_crash_evidence = {
- "status": "success",
- "rows_in": 1,
- "rows_out": 0,
- "rows_rejected": 1,
- "rows_quarantined": 0,
- "commit_outcome": "committed",
- "timings": {"duration_ms": 1},
- "violation_sample": [{"mobile": "[REDACTED]"}],
- "sample_count": 1,
- "redaction_policy": "rule-violation-default-v1",
- }
- with pytest.raises(RuntimeError, match="outcome is unknown"):
- sample_crash_writer.finish(
- sample_crash_run_id,
- sample_crash_evidence,
- )
- with platform.connect() as connection:
- assert connection.execute(
- text(
- """
- SELECT handoff_status
- FROM public.rule_violation_samples
- WHERE rule_run_id = CAST(:id AS uuid)
- """
- ),
- {"id": sample_crash_run_id},
- ).scalar_one() == "ready"
- with platform.begin() as connection:
- connection.execute(
- text(
- """
- UPDATE public.rule_violation_samples
- SET handoff_status = 'unknown',
- cleanup_claim = CAST(:cleanup_claim AS uuid),
- cleanup_claim_expires_at =
- CURRENT_TIMESTAMP + INTERVAL '5 minutes'
- WHERE rule_run_id = CAST(:id AS uuid)
- """
- ),
- {
- "id": sample_crash_run_id,
- "cleanup_claim": new_governance_uid(),
- },
- )
- before_claim_expiry = evidence_writer.reconcile_samples(
- limit=10
- )
- assert before_claim_expiry["claimed"] == 0
- with platform.begin() as connection:
- connection.execute(
- text(
- """
- UPDATE public.rule_violation_samples
- SET cleanup_claim_expires_at =
- CURRENT_TIMESTAMP - INTERVAL '1 second'
- WHERE rule_run_id = CAST(:id AS uuid)
- """
- ),
- {"id": sample_crash_run_id},
- )
- sample_reconciliation = evidence_writer.reconcile_samples(
- limit=10
- )
- assert sample_reconciliation["claimed"] == 1
- assert sample_reconciliation["ready"] >= 1
- evidence_writer.finish(
- sample_crash_run_id,
- sample_crash_evidence,
- )
- assert evidence_writer.replay(sample_crash_run_id)[
- "status"
- ] == "success"
- with platform.connect() as connection:
- missing_sample_ref = connection.execute(
- text(
- """
- SELECT artifact_ref
- FROM public.rule_violation_samples
- WHERE rule_run_id = CAST(:id AS uuid)
- """
- ),
- {"id": sample_crash_run_id},
- ).scalar_one()
- store.delete(missing_sample_ref)
- with platform.begin() as connection:
- connection.execute(
- text(
- """
- UPDATE public.rule_violation_samples
- SET handoff_status = 'unknown',
- cleanup_claim = NULL,
- cleanup_claim_expires_at = NULL
- WHERE rule_run_id = CAST(:id AS uuid)
- """
- ),
- {"id": sample_crash_run_id},
- )
- missing_reconciliation = evidence_writer.reconcile_samples(
- limit=10
- )
- assert missing_reconciliation["failed"] == 1
- with platform.connect() as connection:
- assert connection.execute(
- text(
- """
- SELECT handoff_status
- FROM public.rule_violation_samples
- WHERE rule_run_id = CAST(:id AS uuid)
- """
- ),
- {"id": sample_crash_run_id},
- ).scalar_one() == "failed"
- class TransientArtifactStore:
- def __init__(self, wrapped):
- self.wrapped = wrapped
- def __getattr__(self, name):
- return getattr(self.wrapped, name)
- def describe_optional(self, _ref):
- raise TimeoutError("temporary MinIO timeout")
- with platform.begin() as connection:
- transient_sample_id = connection.execute(
- text(
- """
- UPDATE public.rule_violation_samples s
- SET handoff_status = 'unknown',
- cleanup_claim = NULL,
- cleanup_claim_expires_at = NULL
- FROM public.rule_runs r
- WHERE s.rule_run_id = r.id
- AND r.correlation_id =
- CAST(:correlation_id AS uuid)
- RETURNING s.id::text
- """
- ),
- {"correlation_id": correlation_id},
- ).scalar_one()
- transient_writer = PostgresRuleEvidenceWriter(
- platform,
- TransientArtifactStore(store),
- sample_ttl_seconds=900,
- )
- transient_result = transient_writer.reconcile_samples(
- limit=10
- )
- assert transient_result["claimed"] == 1
- assert transient_result["ready"] == 0
- assert transient_result["failed"] == 0
- with platform.connect() as connection:
- transient_state = connection.execute(
- text(
- """
- SELECT handoff_status,
- cleanup_claim,
- cleanup_claim_expires_at
- FROM public.rule_violation_samples
- WHERE id = CAST(:id AS uuid)
- """
- ),
- {"id": transient_sample_id},
- ).mappings().one()
- assert transient_state["handoff_status"] == "unknown"
- assert transient_state["cleanup_claim"] is None
- assert transient_state["cleanup_claim_expires_at"] is None
- assert evidence_writer.reconcile_samples(limit=10)["ready"] == 1
- with platform.begin() as connection:
- connection.execute(
- text(
- """
- INSERT INTO public.dataflow_dataset_bindings
- (id, dataflow_deployment_id, logical_ref,
- data_source_uid, object_kind, object_ref,
- schema_snapshot_id, dialect, access_mode, write_mode,
- binding_hash)
- VALUES (
- CAST(:id AS uuid), CAST(:deployment_id AS uuid),
- 'sql_receipt', CAST(:source_uid AS uuid), 'table',
- 'public.task6_sql_receipt',
- CAST(:schema_snapshot_id AS uuid), 'postgresql',
- 'read_write', 'upsert', :binding_hash
- )
- """
- ),
- {
- "id": sql_binding_id,
- "deployment_id": deployment_id,
- "source_uid": output_binding["data_source_uid"],
- "schema_snapshot_id": output_schema["id"],
- "binding_hash": "b" * 64,
- },
- )
- receipt_run_id = evidence_writer.start(
- component_binding_id=component_binding_id,
- rule_version_id=rule_id,
- plan_hash=compiled["plan_hash"],
- correlation_id=receipt_correlation_id,
- dataflow_uid=dataflow_uid,
- deployment_id=deployment_id,
- environment="test",
- workflow_version=1,
- node_id="task5_real_polars",
- lease_owner=new_governance_uid(),
- )
- receipt = evidence_writer.stage_sql_output(
- receipt_run_id,
- output_binding_id=sql_binding_id,
- )
- assert re.fullmatch(
- r"dataops-staging://[0-9a-f-]{36}",
- receipt,
- )
- with pytest.raises(ValueError, match="not executable"):
- evidence_writer.resolve_sql_staging(
- receipt,
- deployment_id=deployment_id,
- correlation_id=receipt_correlation_id,
- input_binding_id=sql_binding_id,
- )
- receipt_evidence = {
- "status": "success",
- "rows_in": 2,
- "rows_out": 2,
- "rows_rejected": 0,
- "rows_quarantined": 0,
- "commit_outcome": "committed",
- "timings": {"duration_ms": 1},
- "public_result": {
- "rows_in": 2,
- "rows_out": 2,
- "rows_rejected": 0,
- "rows_quarantined": 0,
- "commit_outcome": "committed",
- },
- }
- evidence_writer.finish(receipt_run_id, receipt_evidence)
- evidence_writer.finish(receipt_run_id, receipt_evidence)
- with pytest.raises(ValueError, match="immutable"):
- evidence_writer.finish(
- receipt_run_id,
- {
- **receipt_evidence,
- "rows_out": 1,
- },
- )
- resolved_receipt = evidence_writer.resolve_sql_staging(
- receipt,
- deployment_id=deployment_id,
- correlation_id=receipt_correlation_id,
- input_binding_id=sql_binding_id,
- )
- assert resolved_receipt["relation_ref"] == (
- "public.task6_sql_receipt"
- )
- assert len(resolved_receipt["relation_digest"]) == 64
- with pytest.raises(ValueError, match="not executable"):
- evidence_writer.resolve_sql_staging(
- receipt,
- deployment_id=deployment_id,
- correlation_id=new_governance_uid(),
- input_binding_id=sql_binding_id,
- )
- with pytest.raises(ValueError, match="invalid"):
- evidence_writer.resolve_sql_staging(
- "dataops-staging://public.task6_sql_receipt",
- deployment_id=deployment_id,
- correlation_id=receipt_correlation_id,
- input_binding_id=sql_binding_id,
- )
- failed_receipt_run_id = evidence_writer.start(
- component_binding_id=component_binding_id,
- rule_version_id=rule_id,
- plan_hash=compiled["plan_hash"],
- correlation_id=failed_receipt_correlation_id,
- dataflow_uid=dataflow_uid,
- deployment_id=deployment_id,
- environment="test",
- workflow_version=1,
- node_id="task5_real_polars",
- lease_owner=new_governance_uid(),
- )
- failed_receipt = evidence_writer.stage_sql_output(
- failed_receipt_run_id,
- output_binding_id=sql_binding_id,
- )
- evidence_writer.finish(
- failed_receipt_run_id,
- {
- "status": "failed",
- "commit_outcome": "not_committed",
- "timings": {"duration_ms": 1},
- },
- )
- with pytest.raises(ValueError, match="not executable"):
- evidence_writer.resolve_sql_staging(
- failed_receipt,
- deployment_id=deployment_id,
- correlation_id=failed_receipt_correlation_id,
- input_binding_id=sql_binding_id,
- )
- with platform.begin() as connection:
- connection.execute(
- text(
- """
- UPDATE public.rule_sql_staging_receipts
- SET expires_at =
- CURRENT_TIMESTAMP - INTERVAL '1 second',
- cleanup_claim = CASE
- WHEN correlation_id =
- CAST(:ready AS uuid)
- THEN CAST(:cleanup_claim AS uuid)
- ELSE NULL
- END,
- cleanup_claim_expires_at = CASE
- WHEN correlation_id =
- CAST(:ready AS uuid)
- THEN CURRENT_TIMESTAMP
- + INTERVAL '5 minutes'
- ELSE NULL
- END
- WHERE correlation_id IN (
- CAST(:ready AS uuid),
- CAST(:failed AS uuid)
- )
- """
- ),
- {
- "ready": receipt_correlation_id,
- "failed": failed_receipt_correlation_id,
- "cleanup_claim": new_governance_uid(),
- },
- )
- assert evidence_writer.cleanup_sql_staging(limit=10) == 1
- with platform.connect() as connection:
- assert connection.execute(
- text(
- """
- SELECT status
- FROM public.rule_sql_staging_receipts
- WHERE correlation_id = CAST(:ready AS uuid)
- """
- ),
- {"ready": receipt_correlation_id},
- ).scalar_one() == "ready"
- with platform.begin() as connection:
- connection.execute(
- text(
- """
- UPDATE public.rule_sql_staging_receipts
- SET cleanup_claim_expires_at =
- CURRENT_TIMESTAMP - INTERVAL '1 second'
- WHERE correlation_id = CAST(:ready AS uuid)
- """
- ),
- {"ready": receipt_correlation_id},
- )
- assert evidence_writer.cleanup_sql_staging(limit=10) == 1
- with platform.connect() as connection:
- assert connection.execute(
- text(
- """
- SELECT COUNT(*)
- FROM public.rule_sql_staging_receipts
- WHERE correlation_id IN (
- CAST(:ready AS uuid),
- CAST(:failed AS uuid)
- )
- AND status = 'expired'
- """
- ),
- {
- "ready": receipt_correlation_id,
- "failed": failed_receipt_correlation_id,
- },
- ).scalar_one() == 2
- conflict_path = tmp_path / "conflict.parquet"
- pl.DataFrame(
- {
- "customer_id": [99],
- "mobile": ["13800138000"],
- "name": ["Conflict"],
- "segment_code": ["A"],
- "segment_name": ["Gold"],
- "version_no": [1],
- }
- ).write_parquet(conflict_path)
- object_count_before_conflict = len(
- list(minio.list_objects(bucket, prefix=prefix, recursive=True))
- )
- with pytest.raises(ValueError, match="immutable|digest"):
- resolver.publish_path(
- str(conflict_path),
- binding_id=output_binding["id"],
- binding_hash=compiled["plan"]["output_binding_hash"],
- correlation_id=correlation_id,
- kind="output",
- ttl_seconds=900,
- schema_fields=output_schema["fields"],
- limits=compiled["plan"]["resource_limits"],
- )
- assert len(
- list(minio.list_objects(bucket, prefix=prefix, recursive=True))
- ) == object_count_before_conflict
- assert customer_artifact["digest"]
- assert segment_artifact["digest"]
- with platform.connect() as connection:
- catalog_rows = connection.execute(
- text(
- """
- SELECT artifact_kind, artifact_ref, handoff_status,
- binding_hash
- FROM public.rule_run_artifacts
- WHERE correlation_id = CAST(:correlation_id AS uuid)
- ORDER BY artifact_kind
- """
- ),
- {"correlation_id": correlation_id},
- ).mappings().all()
- # The repeated deterministic output has the same digest and is
- # idempotently retained as one stable catalog handoff.
- assert len(catalog_rows) == 4
- assert all(row["handoff_status"] == "ready" for row in catalog_rows)
- assert all(len(row["binding_hash"]) == 64 for row in catalog_rows)
- assert {
- row["artifact_ref"]
- for row in catalog_rows
- if row["artifact_kind"] == "output"
- } == {
- result["artifact_ref"],
- downstream_result["artifact_ref"],
- }
- assert len(
- list(minio.list_objects(bucket, prefix=prefix, recursive=True))
- ) == 5
- finally:
- for item in list(
- minio.list_objects(bucket, prefix=prefix, recursive=True)
- ):
- minio.remove_object(bucket, item.object_name)
- assert list(
- minio.list_objects(bucket, prefix=prefix, recursive=True)
- ) == []
- with postgres.begin() as connection:
- connection.execute(text(f"DROP TABLE IF EXISTS {customer_table}"))
- with mysql.begin() as connection:
- connection.execute(text(f"DROP TABLE IF EXISTS {segment_table}"))
- with platform.begin() as connection:
- if ledger_jti is not None:
- connection.execute(
- text(
- "DELETE FROM public.runner_task_executions "
- "WHERE token_jti = CAST(:jti AS uuid)"
- ),
- {"jti": ledger_jti},
- )
- if retry_ledger_jti is not None:
- connection.execute(
- text(
- "DELETE FROM public.runner_task_executions "
- "WHERE token_jti = CAST(:jti AS uuid)"
- ),
- {"jti": retry_ledger_jti},
- )
- for crash_jti in (
- no_run_crash_jti,
- evidence_crash_jti,
- active_retry_jti,
- ):
- if crash_jti is not None:
- connection.execute(
- text(
- "DELETE FROM public.runner_task_executions "
- "WHERE token_jti = CAST(:jti AS uuid)"
- ),
- {"jti": crash_jti},
- )
- connection.execute(
- text(
- """
- DELETE FROM public.rule_sql_staging_receipts
- WHERE correlation_id IN (
- CAST(:receipt_correlation_id AS uuid),
- CAST(:failed_receipt_correlation_id AS uuid)
- )
- """
- ),
- {
- "receipt_correlation_id": receipt_correlation_id,
- "failed_receipt_correlation_id": (
- failed_receipt_correlation_id
- ),
- },
- )
- connection.execute(
- text(
- """
- DELETE FROM public.rule_violation_samples s
- USING public.rule_runs r
- WHERE s.rule_run_id = r.id
- AND r.correlation_id IN (
- CAST(:correlation_id AS uuid),
- CAST(:sample_crash_correlation_id AS uuid)
- )
- """
- ),
- {
- "correlation_id": correlation_id,
- "sample_crash_correlation_id": (
- sample_crash_correlation_id
- ),
- },
- )
- connection.execute(
- text(
- "DELETE FROM public.rule_runs "
- "WHERE correlation_id IN ("
- "CAST(:correlation_id AS uuid), "
- "CAST(:failure_correlation_id AS uuid), "
- "CAST(:unknown_correlation_id AS uuid), "
- "CAST(:lease_correlation_id AS uuid), "
- "CAST(:sample_crash_correlation_id AS uuid), "
- "CAST(:receipt_correlation_id AS uuid), "
- "CAST(:failed_receipt_correlation_id AS uuid), "
- "CAST(:no_run_crash_correlation_id AS uuid), "
- "CAST(:evidence_crash_correlation_id AS uuid), "
- "CAST(:active_retry_correlation_id AS uuid))"
- ),
- {
- "correlation_id": correlation_id,
- "failure_correlation_id": failure_correlation_id,
- "unknown_correlation_id": unknown_correlation_id,
- "lease_correlation_id": lease_correlation_id,
- "sample_crash_correlation_id": (
- sample_crash_correlation_id
- ),
- "receipt_correlation_id": receipt_correlation_id,
- "failed_receipt_correlation_id": (
- failed_receipt_correlation_id
- ),
- "no_run_crash_correlation_id": (
- no_run_crash_correlation_id
- ),
- "evidence_crash_correlation_id": (
- evidence_crash_correlation_id
- ),
- "active_retry_correlation_id": (
- active_retry_correlation_id
- ),
- },
- )
- connection.execute(
- text(
- "DELETE FROM public.dataflow_dataset_bindings "
- "WHERE id = CAST(:id AS uuid)"
- ),
- {"id": sql_binding_id},
- )
- connection.execute(
- text(
- "DELETE FROM public.rule_run_artifacts "
- "WHERE correlation_id = CAST(:correlation_id AS uuid)"
- ),
- {"correlation_id": correlation_id},
- )
- connection.execute(
- text(
- "DELETE FROM public.rule_execution_plans "
- "WHERE id IN (CAST(:id AS uuid), "
- "CAST(:downstream_id AS uuid))"
- ),
- {
- "id": plan_id,
- "downstream_id": downstream_plan_id,
- },
- )
- connection.execute(
- text(
- "DELETE FROM public.dataflow_dataset_bindings "
- "WHERE dataflow_deployment_id = CAST(:id AS uuid)"
- ),
- {"id": deployment_id},
- )
- connection.execute(
- text(
- "DELETE FROM public.dataflow_component_bindings "
- "WHERE id IN (CAST(:id AS uuid), "
- "CAST(:downstream_id AS uuid))"
- ),
- {
- "id": component_binding_id,
- "downstream_id": downstream_component_binding_id,
- },
- )
- connection.execute(
- text(
- "DELETE FROM public.dataflow_deployments "
- "WHERE id = CAST(:id AS uuid)"
- ),
- {"id": deployment_id},
- )
- if gateway_candidate_id is not None:
- connection.execute(
- text(
- "DELETE FROM public.workflow_canary_evidence "
- "WHERE workflow_version_id = CAST(:id AS uuid)"
- ),
- {"id": gateway_candidate_id},
- )
- connection.execute(
- text(
- "DELETE FROM public.workflow_schedules "
- "WHERE workflow_version_id = CAST(:id AS uuid)"
- ),
- {"id": gateway_candidate_id},
- )
- connection.execute(
- text(
- "DELETE FROM public.workflow_engine_bindings "
- "WHERE workflow_version_id = CAST(:id AS uuid)"
- ),
- {"id": gateway_candidate_id},
- )
- connection.execute(
- text(
- "DELETE FROM public.dataflow_workflow_versions "
- "WHERE id = CAST(:id AS uuid)"
- ),
- {"id": gateway_candidate_id},
- )
- connection.execute(
- text(
- "DELETE FROM public.dataflow_versions "
- "WHERE id = CAST(:id AS uuid)"
- ),
- {"id": dataflow_version_id},
- )
- connection.execute(
- text(
- "DELETE FROM public.data_rule_versions "
- "WHERE id IN (CAST(:id AS uuid), "
- "CAST(:downstream_id AS uuid))"
- ),
- {
- "id": rule_id,
- "downstream_id": downstream_rule_id,
- },
- )
- connection.execute(
- text(
- "DELETE FROM public.data_rules "
- "WHERE rule_uid IN (CAST(:rule_uid AS uuid), "
- "CAST(:downstream_rule_uid AS uuid))"
- ),
- {
- "rule_uid": rule_uid,
- "downstream_rule_uid": downstream_rule_uid,
- },
- )
- for schema in (input_schema, lookup_schema, output_schema):
- connection.execute(
- text(
- "DELETE FROM public.data_schema_snapshots "
- "WHERE id = CAST(:id AS uuid)"
- ),
- {"id": schema["id"]},
- )
- postgres.dispose()
- mysql.dispose()
- platform.dispose()
|