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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions .github/workflows/iac-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -137,6 +137,15 @@ jobs:
-q -p no:randomly --junitxml=state-drift.xml
.venv/bin/python scripts/ci/assert_lane_coverage.py state-drift.xml \
--require "state drift=tests/iac/test_iac_state_drift_moto.py"
# The default state key names the provider, and the first apply after the
# upgrade moves the old key's state with `tofu init -migrate-state`. Only a
# real apply, move and plan prove the plan after it changes nothing.
- name: Stage 1 — per-provider state key and its migration (tofu vs moto, creds-free)
run: |
.venv/bin/python -m pytest tests/iac/test_iac_state_key_migration_moto.py \
-q -p no:randomly --junitxml=state-key-migration.xml
.venv/bin/python scripts/ci/assert_lane_coverage.py state-key-migration.xml \
--require "state key migration=tests/iac/test_iac_state_key_migration_moto.py"

# -------------------------------------------------------------------
# Stage 2 — docker emulators. PR + push.
Expand Down
2 changes: 1 addition & 1 deletion .secrets.baseline

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

15 changes: 15 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,21 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### Upgrade notes

- **The first `fluid schedule-sync` after upgrading must retire the old DAG
of each env.** An env's Airflow DAGs now live in `<product-id>__<env>/` with
dag id `<product>__<env>__<build>`; 0.16.7 and earlier wrote them to
`<product-id>/` as `<product>__<build>`. Left in place, the old DAG runs
beside the new one, and both apply the same product against the same state.
With `--delete-scope product` (the default) to a local path or a `git+ssh`
repository, the sync retires them itself: it deletes from `<product-id>/`
only the DAGs rendered for the same product and env under the old id, and
keeps any other file. For every other transport it prints the step to take:
delete those DAG files at the destination once. `--delete-scope destination`
removes the old directory with the rest of what the sync does not ship.
The report's `superseded_scopes` records which case applied.

## [0.16.7] — 2026-09-28

A Lake Formation grant that hides columns now applies on AWS, not only plans.
Expand Down
15 changes: 10 additions & 5 deletions fluid_build/build_runners/_bigquery_load.py
Original file line number Diff line number Diff line change
Expand Up @@ -69,7 +69,11 @@ def bigquery_load_target(
"""The table a binding loads into, or None when it is not a BigQuery table.

Resolved with the IaC's own helpers, so the load names the dataset, table
and location ``_emit_bigquery`` created.
and location ``_emit_bigquery`` created. A binding that names no region
gives ``location: None``, never a guessed ``US``: :func:`load_file` then
runs the job where the table itself is (its ``location``), which is the
only place a load job can run. The guessed default could send a job for
an EU table to the US multi-region.
"""
from ..iac.providers.gcp import BIGQUERY_TABLE, _bq_table_name, resolve_gcp_target

Expand All @@ -83,7 +87,7 @@ def bigquery_load_target(
"project": project,
"dataset": loc.get("dataset") or "default",
"table": _bq_table_name(expose, loc),
"location": loc.get("region") or loc.get("location") or "US",
"location": loc.get("region") or loc.get("location") or None,
}


Expand Down Expand Up @@ -139,10 +143,11 @@ def load_file(
create_disposition="CREATE_NEVER",
schema=table.schema,
)
# The job runs where the table is: the binding's region when it names
# one, otherwise the table's own location, read above, never a default.
location = target.get("location") or getattr(table, "location", None)
with open(path, "rb") as fh:
job = client.load_table_from_file(
fh, table_id, job_config=job_config, location=target["location"]
)
job = client.load_table_from_file(fh, table_id, job_config=job_config, location=location)
job.result()
loaded = int(job.output_rows or 0)
if loaded != expected_rows:
Expand Down
69 changes: 69 additions & 0 deletions fluid_build/build_runners/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -723,6 +723,54 @@ def _run_env(args: argparse.Namespace, plan_data: Optional[Dict[str, Any]] = Non
return None


def _runs_dir(contract_dir: Path, product_id: str, build_id: str) -> Optional[Path]:
"""Where a build's run records are (``FileStateStore``), for ids the store accepts."""
from ._ids import IdentifierViolation, validate_identifier

try:
validate_identifier(product_id, kind="contract.id")
validate_identifier(build_id, kind="build.id")
except IdentifierViolation:
return None
return contract_dir / ".fluid" / "runs" / product_id / build_id / "runs"


def _run_ids(contract_dir: Path, product_id: str, build_id: str) -> Set[str]:
"""The run ids a build has recorded so far."""
runs = _runs_dir(contract_dir, product_id, build_id)
try:
return {p.stem for p in runs.glob("*.json")} if runs and runs.is_dir() else set()
except OSError:
return set()


def _report_build(
report: Any,
contract_dir: Path,
product_id: str,
build_id: str,
result: int,
runs_before: Optional[Set[str]],
) -> None:
"""Record one build on ``report``, with the newest run record it wrote, if any.

Run ids sort by time (``generate_run_id``), so the newest id this build
added is its run. Nothing is read when no ``fluid apply`` report is open.
"""
if report is None:
return
run: Optional[Dict[str, Any]] = None
runs = _runs_dir(contract_dir, product_id, build_id)
added = sorted(_run_ids(contract_dir, product_id, build_id) - (runs_before or set()))
if runs is not None and added:
try:
loaded = json.loads((runs / f"{added[-1]}.json").read_text(encoding="utf-8"))
run = loaded if isinstance(loaded, dict) else None
except (OSError, ValueError):
run = None
report.record_build(build_id=build_id, status="succeeded" if result == 0 else "failed", run=run)


def run_builds_from_args(
args: argparse.Namespace,
logger: logging.Logger,
Expand Down Expand Up @@ -750,6 +798,10 @@ def run_builds_from_args(
root): the contract itself, the contract a bundle's MANIFEST records, or
the contract a plan records (through its bundle when it was planned from
one). See :func:`fluid_build._contract_loader.source_contract_path`.

Each build is recorded on the running ``fluid apply``'s Command Center
report, when there is one (``observability.apply_run``): its status and
the run record it wrote, and why the build phase failed.
"""
# Deferred imports to avoid circular import at module-load time:
# base.py -> python.runner -> base.py (for _resolve_env_placeholders).
Expand Down Expand Up @@ -909,11 +961,19 @@ def run_builds_from_args(
if _b.get("id"):
validate_identifier(_b["id"], kind="build.id")

# The running ``fluid apply``'s run report, if any: each build is recorded on it.
from fluid_build.observability.apply_run import current_apply_run

report = current_apply_run()
product_id = str(contract.get("id") or "")

# Filter builds if specific ID requested
if args.build_id:
builds = [b for b in builds if b.get("id") == args.build_id]
if not builds:
LOG.error(f"Build not found: {args.build_id}")
if report is not None:
report.build_failed(f"build_not_found:{args.build_id}")
return 1

# The overlay env the contract above was loaded with, decided once and
Expand All @@ -936,6 +996,7 @@ def run_builds_from_args(

for build in builds:
build_id = build.get("id", "unknown")
runs_before = _run_ids(contract_path.parent, product_id, build_id) if report else None

if is_acquisition_build(build):
sample_rows = getattr(args, "sample_rows", None)
Expand All @@ -946,6 +1007,7 @@ def run_builds_from_args(
dry_run=args.dry_run,
sample_rows=sample_rows,
)
_report_build(report, contract_path.parent, product_id, build_id, result, runs_before)
if result == 0:
total_executed += 1
else:
Expand Down Expand Up @@ -987,6 +1049,8 @@ def run_builds_from_args(
expected = (contract_path.parent / repository / "dbt_project.yml").resolve()
cprint(f"\n⚠️ Build '{build_id}' - dbt project not found: {expected}")
total_skipped += 1
if report is not None:
report.record_build(build_id=build_id, status="skipped")
continue

result = execute_dbt_build(
Expand Down Expand Up @@ -1019,6 +1083,8 @@ def run_builds_from_args(
"For Python builds, create the script at the expected path above."
)
total_skipped += 1
if report is not None:
report.record_build(build_id=build_id, status="skipped")
continue

# Execute build
Expand All @@ -1033,6 +1099,7 @@ def run_builds_from_args(
force_run=force_run,
)

_report_build(report, contract_path.parent, product_id, build_id, result, runs_before)
if result == 0:
total_executed += 1
else:
Expand Down Expand Up @@ -1067,6 +1134,8 @@ def run_builds_from_args(
total_skipped,
)
return 0
if report is not None:
report.build_failed("builds_all_skipped")
console_error(
f"Every build was skipped ({total_skipped}/{len(builds)}) — nothing "
"was transformed and no rows were produced. Fix the missing build "
Expand Down
Loading
Loading