diff --git a/.github/workflows/build.yml b/.github/workflows/build.yml index 331c6f16d17..c5684924cee 100644 --- a/.github/workflows/build.yml +++ b/.github/workflows/build.yml @@ -7,6 +7,7 @@ on: types: - opened - synchronize + - ready_for_review workflow_run: workflows: ["Fork PR Gate"] types: [completed] @@ -303,6 +304,128 @@ jobs: PATH=".venv/bin:$PATH" script/bqetl format --check \ $(git ls-tree -d HEAD --name-only) + flag-sensitive-flows: + name: Flag sensitive data flows + runs-on: ubuntu-latest + permissions: + contents: read + pull-requests: write # to request @mozilla/dataplatform-wg review + environment: *build-env + needs: [build, decide-runs] + if: needs.decide-runs.outputs.validate-sql == 'true' + steps: + - *checkout-with-history + - *setup-python + - *restore-venv + # Resolve PR number / base SHA / draft across event types. On the + # workflow_run path (fork PRs, post gate-approval) github.event.pull_request + # is null, so resolve from the workflow_run payload / API instead — mirrors + # resolve-pr-base and Determine context elsewhere in this workflow. + - id: pr-context + env: + GH_TOKEN: ${{ secrets.GITHUB_TOKEN }} + REPO: ${{ github.repository }} + HEAD_SHA: ${{ github.event.workflow_run.head_sha || github.event.pull_request.head.sha || github.sha }} + PR_FROM_WR: ${{ github.event.workflow_run.pull_requests[0].number }} + EVENT_PR_NUMBER: ${{ github.event.pull_request.number }} + EVENT_BASE_SHA: ${{ github.event.pull_request.base.sha }} + EVENT_DRAFT: ${{ github.event.pull_request.draft }} + run: | + set -euo pipefail + pr_number=""; base_sha=""; draft="" + if [[ "$GITHUB_EVENT_NAME" == "pull_request" ]]; then + pr_number="$EVENT_PR_NUMBER" + base_sha="$EVENT_BASE_SHA" + draft="$EVENT_DRAFT" + elif [[ "$GITHUB_EVENT_NAME" == "workflow_run" ]]; then + pr_number="$PR_FROM_WR" + if [[ -z "$pr_number" || "$pr_number" == "null" ]]; then + pr_number=$(gh api "/repos/${REPO}/commits/${HEAD_SHA}/pulls" \ + --jq '.[0].number' 2>/dev/null || echo "") + fi + if [[ -n "$pr_number" && "$pr_number" != "null" ]]; then + base_sha=$(gh api "/repos/${REPO}/pulls/${pr_number}" \ + --jq '.base.sha' 2>/dev/null || echo "") + draft=$(gh api "/repos/${REPO}/pulls/${pr_number}" \ + --jq '.draft' 2>/dev/null || echo "") + fi + fi + echo "pr_number=${pr_number}" >> "$GITHUB_OUTPUT" + echo "base_sha=${base_sha}" >> "$GITHUB_OUTPUT" + echo "draft=${draft}" >> "$GITHUB_OUTPUT" + - id: changed-queries + uses: tj-actions/changed-files@9426d40962ed5378910ee2e21d5f8c6fcbf2dd96 # v47.0.6 + with: + files: sql/**/query.sql + base_sha: ${{ steps.pr-context.outputs.base_sha || github.event.pull_request.base.sha || github.event.before }} + # A changed query reads sensitive (restricted/workgroup-gated) data and + # writes to a more broadly readable destination that no team owns. This is + # advisory: rather than block or make the author edit CODEOWNERS, it just + # requests Data Platform review and posts a comment. Flows on paths already + # owned by a team are left to that team's normal review. + - name: Flag sensitive-data flows widening read access + if: steps.changed-queries.outputs.any_changed == 'true' + env: + CHANGED_QUERIES: ${{ steps.changed-queries.outputs.all_changed_files }} + GH_TOKEN: ${{ github.token }} + PR: ${{ steps.pr-context.outputs.pr_number }} + REPO: ${{ github.repository }} + DRAFT: ${{ steps.pr-context.outputs.draft }} + run: | + set +e + advisory=$(PATH=".venv/bin:$PATH" \ + script/bqetl data_governance sensitivity $CHANGED_QUERIES) + code=$? + set -e + echo "$advisory" + # Detection always runs (visible above), but hold off on requesting + # review / commenting while the PR is a draft — the ready_for_review + # trigger re-runs this and posts once it's marked ready. + if [ "$code" -eq 2 ] && [ -n "$PR" ] && [ "$DRAFT" != "true" ]; then + # The @mozilla/dataplatform-wg mention in the comment below is the + # reliable notification. Also try to add the team to the Reviewers + # box, but this is best-effort: the repo-scoped Actions GITHUB_TOKEN + # can't resolve an org team as a reviewer (needs read:org, which the + # permissions: block can't grant), so it 422s — a PAT/App token + # secret with read:org would make it work. + gh api --method POST \ + "repos/$REPO/pulls/$PR/requested_reviewers" \ + -f "team_reviewers[]=dataplatform-wg" \ + || echo "Note: could not add the team to the Reviewers box with the" \ + "Actions token; the @mention in the comment notifies them." + # Upsert a single advisory comment (marker-keyed) so re-runs on each + # push update it in place instead of piling up duplicates. + marker='' + body="$(mktemp)" + { + echo "$marker" + echo "⚠️ **Sensitive-data flow detected** — cc @mozilla/dataplatform-wg, please review." + echo + echo "A changed query reads restricted / workgroup-gated data and writes it to a more broadly readable destination that no team owns:" + echo + echo '```' + echo "$advisory" + echo '```' + echo + echo "This is advisory and does not block merge." + echo "A Data Platform reviewer should confirm the widened access is intended, or narrow the destination's \`workgroup_access\`." + } > "$body" + existing=$(gh api "repos/$REPO/issues/$PR/comments" \ + --jq "[.[] | select(.body | contains(\"$marker\"))][0].id" 2>/dev/null || echo "") + if [ -n "$existing" ] && [ "$existing" != "null" ]; then + gh api --method PATCH "repos/$REPO/issues/comments/$existing" \ + -F "body=@$body" >/dev/null || echo "Could not update existing comment" + else + gh pr comment "$PR" -R "$REPO" --body-file "$body" \ + || echo "Could not post comment (fork PR token is read-only?)" + fi + elif [ "$code" -eq 2 ] && [ "$DRAFT" = "true" ]; then + echo "::notice::Sensitive-data flow(s) detected; PR is a draft, deferring @mozilla/dataplatform-wg review until it's marked ready." + elif [ "$code" -ne 0 ] && [ "$code" -ne 2 ]; then + echo "::warning::sensitivity check errored (exit $code); skipping" + fi + exit $code + test-bqetl: name: Test bqetl runs-on: ubuntu-latest diff --git a/bigquery_etl/data_governance/cli.py b/bigquery_etl/data_governance/cli.py index f74c5700904..992925c34e3 100644 --- a/bigquery_etl/data_governance/cli.py +++ b/bigquery_etl/data_governance/cli.py @@ -1,11 +1,14 @@ """bigquery-etl CLI data_governance command.""" +import json +import sys from datetime import datetime, timezone from pathlib import Path import rich_click as click from ..cli.utils import sql_dir_option +from ..sensitivity import CODEOWNERS_FILE, check_paths, format_findings from ..util.common import block_coding_agents from .classification import runner, upstream from .classification.config import ( @@ -18,6 +21,10 @@ ) from .classification.runner import TargetKey +# exit code CI keys on to request review; distinct from a tool crash (exit 1) so +# the job knows to request DPE review rather than fail. +SENSITIVITY_UNGATED_EXIT_CODE = 2 + def _parse_target(value: str) -> TargetKey: """Turn `project.dataset[.table]` into the triple the library takes.""" @@ -299,3 +306,50 @@ def classify( targets, refresh=refresh, ) + + +@data_governance.command() +@click.argument("paths", nargs=-1, required=True, type=click.Path()) +@sql_dir_option +@click.option( + "--codeowners", + default=CODEOWNERS_FILE, + help="CODEOWNERS file used to decide whether a flagged flow is already " + "owned/reviewed.", +) +@click.option( + "--json", + "as_json", + is_flag=True, + default=False, + help="Emit findings as JSON instead of the human-readable advisory.", +) +def sensitivity(paths, sql_dir, codeowners, as_json): + """Flag queries that read sensitive data and write it somewhere broader. + + Scans the given query.sql files (or directories of them) for flows where a + restricted / narrowly workgroup-gated source is written to a more broadly + readable destination that the source's readers don't already cover. + + Advisory and read-only: exits 2 when there are ungated flows so CI can + request @mozilla/dataplatform-wg review; it never blocks a merge. + """ + findings = check_paths(list(paths), sql_dir, codeowners_file=codeowners) + ungated = [f for f in findings if f.get("gated") is False] + + if as_json: + click.echo(json.dumps(findings, indent=2)) + elif findings: + click.echo(format_findings(findings)) + + if ungated: + n = len({(f["source"], f["query"]) for f in ungated}) + click.echo( + f"::warning::{n} sensitive-data flow(s) widen read access beyond the " + "source's authorized readers and aren't owned by any team. Requesting " + "@mozilla/dataplatform-wg review (advisory — this check does not " + "block).", + err=True, + ) + sys.exit(SENSITIVITY_UNGATED_EXIT_CODE) + click.echo("no ungated sensitive-data flows", err=True) diff --git a/bigquery_etl/sensitivity.py b/bigquery_etl/sensitivity.py new file mode 100644 index 00000000000..713b961ccbd --- /dev/null +++ b/bigquery_etl/sensitivity.py @@ -0,0 +1,367 @@ +"""Flag risky data flows in SQL queries. + +Detects when a query reads *sensitive* (restricted / narrowly workgroup-gated) +data and writes it to a *more broadly readable* destination. + +Only handles SQL (`query.sql`); `query.py` builds table names dynamically and +can't be resolved statically. +""" + +import logging +import os +import re +from pathlib import Path +from typing import Dict, List, Optional, Set, Tuple + +import pathspec + +from bigquery_etl.config import ConfigLoader +from bigquery_etl.dependency import extract_table_references +from bigquery_etl.metadata.parse_metadata import ( + DATASET_METADATA_FILE, + METADATA_FILE, + DatasetMetadata, + Metadata, +) +from bigquery_etl.util.common import render + +logger = logging.getLogger(__name__) + +# a parsed CODEOWNERS entry: a compiled matcher and the owners for that pattern +CodeownersEntry = Tuple[pathspec.PathSpec, List[str]] + + +def _config(key: str, default): + """Read a `sensitivity.` value from bqetl_project.yaml, or `default`.""" + value = ConfigLoader.get("sensitivity", key, fallback=None) + return value if value is not None else default + + +# roles that grant the ability to read row-level data (dataEditor/dataOwner also +# confer read, so they count as readers even though they're rare here) +READ_ROLES = set( + _config( + "read_roles", + [ + "roles/bigquery.dataViewer", + "roles/bigquery.dataEditor", + "roles/bigquery.dataOwner", + ], + ) +) +# the org-wide "everyone with confidential access" reader; anything scoped only +# to this (or broader) is not considered sensitive +BROAD_READERS = set( + _config("broad_readers", ["workgroup:mozilla-confidential/data-viewers"]) +) +# projects whose datasets are public/non-sensitive and safe to read from. +# mozfun is the shared public UDF library; every other unresolved source is +# treated as sensitive (can't be verified) rather than skipped. +SAFE_PROJECTS = set(_config("safe_projects", ["mozfun"])) +# stable/live dataset suffixes: unless listed in `gated_datasets`, a stable/live +# dataset in an ingestion project is granted the broad default and treated as +# non-sensitive. +INGESTION_SUFFIXES = tuple(_config("ingestion_suffixes", ["_stable", "_live"])) +# CODEOWNERS file used to decide whether a flagged flow is already owned/reviewed. +CODEOWNERS_FILE = _config("codeowners_file", "CODEOWNERS") +# version suffix on stable/live table names (e.g. main_v5); gated_datasets keys +# tables by the version-stripped doctype. +_VERSION_RE = re.compile(r"_v[0-9]+$") + + +def load_gated_datasets() -> Dict[str, Dict]: + """Load the `sensitivity.gated_datasets` map (keyed by project, then dataset).""" + return ConfigLoader.get("sensitivity", "gated_datasets", fallback={}) or {} + + +def gated_readers( + gated: Dict[str, Dict], project: str, dataset: str, table: str +) -> Optional[Set[str]]: + """Resolve a gated dataset/table's dataViewer members, or None if unlisted. + + Table-level entries (keyed by version-stripped doctype) override the dataset + `default`; a listed dataset with no matching table falls back to `default` + (empty when the dataset grants no broad read access). + """ + entry = (gated.get(project) or {}).get(dataset) + if entry is None: + return None + tables = entry.get("tables") or {} + doctype = _VERSION_RE.sub("", table) + members = tables.get(doctype, tables.get(table, entry.get("default", []))) + return set(members) + + +class DatasetAccess: + """Effective read access for a dataset, from its dataset_metadata.yaml.""" + + def __init__(self, readers: Set[str], base_acl: str): + """Hold a dataset's read-role members and its base ACL archetype.""" + self.readers = readers + self.base_acl = base_acl + + @property + def sensitive(self) -> bool: + """Restricted base ACL, no authorized readers, or a non-broad workgroup.""" + if "restricted" in self.base_acl: + return True + if not self.readers: + # empty readers == nobody authorized, strictly narrower than broad + # (note: `set() <= BROAD_READERS` is True, so this must be explicit) + return True + return not self.readers <= BROAD_READERS + + +def dataset_access(sql_dir: str, project: str, dataset: str) -> Optional[DatasetAccess]: + """Resolve a dataset's effective read access, or None if unresolved.""" + metadata_path = Path(sql_dir) / project / dataset / DATASET_METADATA_FILE + if not metadata_path.exists(): + return None + try: + metadata = DatasetMetadata.from_file(metadata_path) + except Exception: + return None + readers: Set[str] = set() + for entry in metadata.workgroup_access or []: + if entry.get("role") in READ_ROLES: + readers.update(entry.get("members", [])) + return DatasetAccess(readers=readers, base_acl=metadata.dataset_base_acl or "") + + +def table_readers(sql_dir: str, project: str, dataset: str, table: str) -> Set[str]: + """Read-role members granted at the table level. + + A table's own metadata.yaml `workgroup_access` only widens the dataset's + grant, so these are unioned into the table's effective readers. + """ + metadata_path = Path(sql_dir) / project / dataset / table / METADATA_FILE + if not metadata_path.exists(): + return set() + try: + metadata = Metadata.from_file(metadata_path) + except Exception: + return set() + readers: Set[str] = set() + for entry in metadata.workgroup_access or []: + if entry.role in READ_ROLES: + readers.update(entry.members) + return readers + + +def effective_access( + sql_dir: str, project: str, dataset: str, table: str +) -> Optional[DatasetAccess]: + """Dataset access widened by the table's own workgroup_access.""" + da = dataset_access(sql_dir, project, dataset) + if da is None: + return None + return DatasetAccess( + readers=da.readers | table_readers(sql_dir, project, dataset, table), + base_acl=da.base_acl, + ) + + +def source_access( + sql_dir: str, + project: str, + dataset: str, + table: str, + gated: Dict[str, Dict], +) -> Optional[DatasetAccess]: + """Resolve a source table's effective read access. + + Resolution order: + 1. `gated_datasets` config: authoritative for ingestion (stable/live) + ACLs, which have no dataset_metadata.yaml in this repo. + 2. in-repo dataset_metadata.yaml (derived/view datasets). + 3. an unlisted stable/live dataset in an ingestion project (a project keyed + in `gated_datasets`): granted the broad default, so non-sensitive. + 4. None: genuinely unresolved (caller treats as sensitive). + """ + readers = gated_readers(gated, project, dataset, table) + if readers is not None: + return DatasetAccess(readers=readers, base_acl="gated") + + da = effective_access(sql_dir, project, dataset, table) + if da is not None: + return da + + if project in gated and dataset.endswith(INGESTION_SUFFIXES): + return DatasetAccess(readers=set(BROAD_READERS), base_acl="ingestion") + + return None + + +def _resolve_ref(ref: str, default_project: str) -> Optional[Tuple[str, str, str]]: + """Return (project, dataset, table) for a table ref, or None if unusable.""" + parts = ref.split(".") + if len(parts) == 3: + project, dataset, table = parts + elif len(parts) == 2: + project, dataset, table = default_project, parts[0], parts[1] + else: + return None + if dataset == "INFORMATION_SCHEMA": + return None # query/job metadata, not a data source + return project, dataset, table + + +def load_codeowners(codeowners_file: str) -> List[CodeownersEntry]: + """Parse CODEOWNERS into ordered (compiled matcher, owners) entries. + + The matcher is compiled once here rather than per path so that checking many + query files against the same patterns stays cheap. + """ + entries: List[CodeownersEntry] = [] + for line in Path(codeowners_file).read_text().splitlines(): + line = line.strip() + if not line or line.startswith("#"): + continue + pattern, *owners = line.split() + spec = pathspec.PathSpec.from_lines("gitwildmatch", [pattern]) + entries.append((spec, owners)) + return entries + + +def path_owners(rel_path: str, entries: List[CodeownersEntry]) -> List[str]: + """Owners for a repo-relative path, honoring CODEOWNERS last-match-wins.""" + owners: List[str] = [] + for spec, pattern_owners in entries: + if spec.match_file(rel_path): + owners = pattern_owners + return owners + + +def check_query( + query_file: str, + sql_dir: str, + codeowners: Optional[List[CodeownersEntry]] = None, + gated_datasets: Optional[Dict[str, Dict]] = None, +) -> List[Dict]: + """Flag sensitive-source -> broader-destination flows for one query.sql. + + Returns one finding per sensitive source whose readers don't already cover + the destination's readers (i.e. the write widens access). When `codeowners` + is given, each finding records whether the query path is owned (i.e. review + is required). + + A query whose destination access can't be resolved, or that can't be + rendered/parsed, is logged and skipped (returns []) — visible rather than + silently passing, since an un-analyzable query is exactly what a reviewer + should hear about for a guardrail. + """ + query_file_path = Path(query_file) + dest_table = query_file_path.parent.name + dest_dataset = query_file_path.parent.parent.name + dest_project = query_file_path.parent.parent.parent.name + dest = effective_access(sql_dir, dest_project, dest_dataset, dest_table) + if dest is None: + logger.warning( + "%s: destination %s.%s has no resolvable dataset_metadata.yaml; " + "skipping sensitivity analysis", + query_file_path, + dest_project, + dest_dataset, + ) + return [] + + # query.sql is a Jinja template; render before parsing + try: + sql = render(query_file_path.name, template_folder=query_file_path.parent) + refs = extract_table_references(sql) + except Exception as exc: + logger.warning( + "%s: could not render/parse for sensitivity analysis (%s); skipping", + query_file_path, + exc, + ) + return [] + + owners = ( + path_owners(os.path.relpath(query_file_path), codeowners) + if codeowners is not None + else None + ) + gated = load_gated_datasets() if gated_datasets is None else gated_datasets + + findings: List[Dict] = [] + for ref in refs: + resolved = _resolve_ref(ref, dest_project) + if resolved is None: + continue + src_project, src_dataset, src_table = resolved + if src_project in SAFE_PROJECTS: + continue # mozfun etc.: public, safe to read + if (src_project, src_dataset) == (dest_project, dest_dataset): + continue # self-reference + src = source_access(sql_dir, src_project, src_dataset, src_table, gated) + if src is None: + # unresolved and not a known-safe project (e.g. other-project refs): + # can't verify its access, so treat it as sensitive. + src_readers: Set[str] = set() + src_base_acl = "unresolved (treated as sensitive)" + elif src.sensitive: + src_readers = src.readers + src_base_acl = src.base_acl + else: + continue # resolved and not sensitive + # readers the destination grants that the source does not authorize + widened = dest.readers - src_readers + if widened: + findings.append( + { + "query_path": str(query_file_path), + "query": f"{dest_project}.{dest_dataset}", + "source": f"{src_project}.{src_dataset}", + "source_readers": sorted(src_readers), + "source_base_acl": src_base_acl, + "destination_readers": sorted(dest.readers), + "extra_readers": sorted(widened), + "owners": owners, + # gated == a reviewer is required for this path + "gated": bool(owners) if owners is not None else None, + } + ) + return findings + + +def check_paths( + paths: List[str], + sql_dir: str, + codeowners_file: Optional[str] = None, +) -> List[Dict]: + """Run check_query over query.sql files under the given paths.""" + codeowners = load_codeowners(codeowners_file) if codeowners_file else None + findings: List[Dict] = [] + for p in paths: + path = Path(p) + if path.name == "query.sql": + query_files = [path] + elif path.is_dir(): + query_files = sorted(path.rglob("query.sql")) + else: + continue # non-SQL artifact (e.g. query.py) + for query_file_path in query_files: + findings.extend(check_query(str(query_file_path), sql_dir, codeowners)) + return findings + + +def format_findings(findings: List[Dict]) -> str: + """Human-readable advisory, one block per source -> destination pair.""" + pairs: Dict[Tuple[str, str], Dict] = {} + for f in findings: + pairs.setdefault((f["source"], f["query"]), f) + lines = [] + for (source, dest), f in sorted(pairs.items()): + gate = ( + "needs Data Platform (@mozilla/dataplatform-wg) review" + if not f["gated"] + else f"already reviewed by {', '.join(f['owners'])}" + ) + lines.append( + f"- {source} ({f['source_base_acl'] or 'gated'}) -> {dest}\n" + f" grants read to {', '.join(f['extra_readers'])} " + f"not authorized on the source\n" + f" {gate}" + ) + return "\n".join(lines) diff --git a/bqetl_project.yaml b/bqetl_project.yaml index 28b8f1fbe18..d55146022b1 100644 --- a/bqetl_project.yaml +++ b/bqetl_project.yaml @@ -755,3 +755,97 @@ retention_exclusion_list: - sql/moz-fx-data-shared-prod/firefox_desktop_background_defaultagent/baseline_clients_city_seen_v1 - sql/moz-fx-data-shared-prod/firefox_desktop_background_tasks/baseline_clients_city_seen_v1 - sql/moz-fx-data-shared-prod/firefox_desktop_derived/metrics_clients_first_seen_v1 + +sensitivity: + # Configuration for the sensitive-data-flow check (bigquery_etl/sensitivity.py), + # which flags queries that read narrowly-gated data and write it somewhere more + # broadly readable. + + # Roles that grant the ability to read row-level data. dataEditor/dataOwner are + # rare in workgroup_access here but also confer read, so count them as readers + # (omitting them would understate a source's readers -> false positives). + read_roles: + - roles/bigquery.dataViewer + - roles/bigquery.dataEditor + - roles/bigquery.dataOwner + broad_readers: + - workgroup:mozilla-confidential/data-viewers + safe_projects: + - mozfun + ingestion_suffixes: + - _stable + - _live + codeowners_file: CODEOWNERS + + # Datasets/tables whose read access is narrower than the broad_readers above. + # Omitted/empty means restricted (no broad read access). + gated_datasets: + moz-fx-data-shared-prod: + ads_backend_live: &ads_backend + default: [workgroup:ads/data-viewers] + ads_backend_stable: *ads_backend + + contextual_services_live: &contextual_services + default: [workgroup:contextual-services/data-viewers] + tables: + quicksuggest_impression: + - workgroup:contextual-services/data-viewers + - workgroup:search-terms/sanitized-writer + - workgroup:search-terms/unsanitized + contextual_services_stable: *contextual_services + + firefox_enterprise_desktop_live: &firefox_enterprise_desktop + default: [workgroup:fx-enterprise/data-viewers] + firefox_enterprise_desktop_stable: *firefox_enterprise_desktop + + search_terms: + default: [] + search_terms_derived: + default: [] + + telemetry_live: &telemetry + default: [workgroup:mozilla-confidential/data-viewers] + tables: + xfocsp_error_report: [] + account_ecosystem: [] + telemetry_stable: *telemetry + + firefox_desktop_live: &firefox_desktop + default: [workgroup:mozilla-confidential/data-viewers] + tables: + microsurvey: [workgroup:microsurvey/data-viewers] + fx_accounts: [workgroup:client-association/data-viewers] + quick_suggest: [workgroup:contextual-services/data-viewers] + top_sites: [workgroup:contextual-services/data-viewers] + serp_categorization: [workgroup:revenue/cat2] + user_characteristics: [workgroup:user-characteristics/data-viewers] + firefox_desktop_stable: *firefox_desktop + + org_mozilla_firefox_live: &glean_client + default: [workgroup:mozilla-confidential/data-viewers] + tables: + fx_accounts: [workgroup:client-association/data-viewers] + fx_suggest: [workgroup:contextual-services/data-viewers] + topsites_impression: [workgroup:contextual-services/data-viewers] + client_deduplication: [] + user_characteristics: [workgroup:user-characteristics/data-viewers] + org_mozilla_firefox_stable: *glean_client + org_mozilla_firefox_beta_live: *glean_client + org_mozilla_firefox_beta_stable: *glean_client + org_mozilla_fenix_live: *glean_client + org_mozilla_fenix_stable: *glean_client + org_mozilla_fenix_nightly_live: *glean_client + org_mozilla_fenix_nightly_stable: *glean_client + org_mozilla_fennec_aurora_live: *glean_client + org_mozilla_fennec_aurora_stable: *glean_client + + org_mozilla_ios_firefox_live: &ios_client + default: [workgroup:mozilla-confidential/data-viewers] + tables: + fx_accounts: [workgroup:client-association/data-viewers] + topsites_impression: [workgroup:contextual-services/data-viewers] + org_mozilla_ios_firefox_stable: *ios_client + org_mozilla_ios_firefoxbeta_live: *ios_client + org_mozilla_ios_firefoxbeta_stable: *ios_client + org_mozilla_ios_fennec_live: *ios_client + org_mozilla_ios_fennec_stable: *ios_client diff --git a/tests/test_sensitivity.py b/tests/test_sensitivity.py new file mode 100644 index 00000000000..0f3bc876266 --- /dev/null +++ b/tests/test_sensitivity.py @@ -0,0 +1,410 @@ +"""Tests for the sensitive-data-flow check (bigquery_etl/sensitivity.py).""" + +from pathlib import Path + +from click.testing import CliRunner + +# Import the bqetl CLI before the data_governance group: the group takes an +# option from bigquery_etl.cli.utils and bigquery_etl.cli imports the group +# back, so importing the group first finds it partially initialized. +from bigquery_etl.cli import cli as _bqetl_cli # noqa: F401 +from bigquery_etl.data_governance.cli import sensitivity as sensitivity_cmd +from bigquery_etl.sensitivity import ( + check_query, + dataset_access, + effective_access, + gated_readers, + load_codeowners, + path_owners, + source_access, +) + +PROJECT = "moz-fx-data-shared-prod" +BROAD = "workgroup:mozilla-confidential/data-viewers" + + +def _dataset(sql_dir, name, base_acl="derived", members=(BROAD,), project=PROJECT): + ds = sql_dir / project / name + ds.mkdir(parents=True, exist_ok=True) + wa = "workgroup_access:\n- role: roles/bigquery.dataViewer\n members:\n" + "".join( + f" - {m}\n" for m in members + ) + (ds / "dataset_metadata.yaml").write_text( + f"friendly_name: {name}\ndescription: d\n" + f"dataset_base_acl: {base_acl}\nuser_facing: false\n{wa}" + ) + return ds + + +def _query(ds, table, sql, table_members=None): + t = ds / table + t.mkdir(parents=True, exist_ok=True) + (t / "query.sql").write_text(sql) + if table_members is not None: + wa = ( + "workgroup_access:\n- role: roles/bigquery.dataViewer\n members:\n" + + "".join(f" - {m}\n" for m in table_members) + ) + (t / "metadata.yaml").write_text( + f"friendly_name: {table}\ndescription: d\nowners:\n - a@b.org\n{wa}" + ) + return str(t / "query.sql") + + +def _select(dataset, table="t", project=PROJECT): + return f"SELECT * FROM `{project}.{dataset}.{table}`" + + +class TestAccessResolution: + def test_sensitivity_classification(self, tmp_path): + sql = tmp_path / "sql" + _dataset(sql, "restricted_ds", base_acl="restricted") + _dataset(sql, "broad_ds", base_acl="derived", members=[BROAD]) + _dataset(sql, "narrow_ds", base_acl="derived", members=["workgroup:acme/x"]) + + assert dataset_access(str(sql), PROJECT, "restricted_ds").sensitive is True + # broad, non-restricted -> not sensitive + assert dataset_access(str(sql), PROJECT, "broad_ds").sensitive is False + # scoped to a narrow workgroup -> sensitive + assert dataset_access(str(sql), PROJECT, "narrow_ds").sensitive is True + + def test_unresolved_dataset_returns_none(self, tmp_path): + sql = tmp_path / "sql" + assert dataset_access(str(sql), PROJECT, "does_not_exist") is None + + def test_effective_access_unions_table_readers(self, tmp_path): + sql = tmp_path / "sql" + ds = _dataset( + sql, "foo_external", base_acl="restricted", members=["workgroup:a/x"] + ) + _query(ds, "tbl", "SELECT 1", table_members=["workgroup:a/x", "workgroup:b/y"]) + acc = effective_access(str(sql), PROJECT, "foo_external", "tbl") + # table-level grant (b/y) is unioned on top of the dataset grant (a/x) + assert acc.readers == {"workgroup:a/x", "workgroup:b/y"} + + +class TestFlowCheck: + def test_flags_restricted_source_to_broad_dest(self, tmp_path): + sql = tmp_path / "sql" + _dataset( + sql, "src_external", base_acl="restricted", members=["workgroup:acme/x"] + ) + dest = _dataset(sql, "broad_derived", members=[BROAD]) + q = _query(dest, "out_v1", _select("src_external")) + + findings = check_query(q, str(sql)) + assert len(findings) == 1 + assert findings[0]["source"] == f"{PROJECT}.src_external" + assert findings[0]["extra_readers"] == [BROAD] + + def test_no_flag_when_dest_same_workgroup(self, tmp_path): + sql = tmp_path / "sql" + _dataset( + sql, "src_external", base_acl="restricted", members=["workgroup:acme/x"] + ) + dest = _dataset( + sql, "acme_derived", base_acl="restricted", members=["workgroup:acme/x"] + ) + q = _query(dest, "out_v1", _select("src_external")) + + assert check_query(q, str(sql)) == [] + + def test_no_flag_when_source_table_grants_dest_workgroup(self, tmp_path): + sql = tmp_path / "sql" + # dataset is narrow, but the source *table* widens to the dest's workgroup + src = _dataset( + sql, "src_external", base_acl="restricted", members=["workgroup:acme/x"] + ) + _query( + src, + "wide_tbl", + "SELECT 1", + table_members=["workgroup:acme/x", "workgroup:team/y"], + ) + dest = _dataset( + sql, "team_derived", base_acl="restricted", members=["workgroup:team/y"] + ) + q = _query(dest, "out_v1", _select("src_external", table="wide_tbl")) + + assert check_query(q, str(sql)) == [] + + def test_mozfun_source_is_safe(self, tmp_path): + sql = tmp_path / "sql" + dest = _dataset(sql, "broad_derived", members=[BROAD]) + q = _query(dest, "out_v1", _select("stats", table="foo", project="mozfun")) + + assert check_query(q, str(sql)) == [] + + def test_unresolved_source_treated_as_sensitive(self, tmp_path): + sql = tmp_path / "sql" + dest = _dataset(sql, "broad_derived", members=[BROAD]) + # an unknown other-project dataset can't be resolved -> unsafe + q = _query( + dest, + "out_v1", + _select("other_dataset", table="t", project="other-project"), + ) + + findings = check_query(q, str(sql), gated_datasets={}) + assert len(findings) == 1 + assert "unresolved" in findings[0]["source_base_acl"] + assert findings[0]["extra_readers"] == [BROAD] + + def test_information_schema_ignored(self, tmp_path): + sql = tmp_path / "sql" + dest = _dataset(sql, "broad_derived", members=[BROAD]) + q = _query( + dest, + "out_v1", + "SELECT * FROM `region-us`.INFORMATION_SCHEMA.JOBS", + ) + + assert check_query(q, str(sql)) == [] + + +class TestCodeowners: + def _codeowners(self, tmp_path): + (tmp_path / "CODEOWNERS").write_text( + "* @mozilla/dataplatform-wg\n" + "/sql/\n" # explicitly unowned + f"/sql/{PROJECT}/owned_derived/** @mozilla/team\n" + ) + return load_codeowners(str(tmp_path / "CODEOWNERS")) + + def test_path_owners_last_match_wins(self, tmp_path): + entries = self._codeowners(tmp_path) + assert path_owners(f"sql/{PROJECT}/plain_derived/x/query.sql", entries) == [] + assert path_owners(f"sql/{PROJECT}/owned_derived/x/query.sql", entries) == [ + "@mozilla/team" + ] + + def test_finding_gated_flag_reflects_ownership(self, tmp_path, monkeypatch): + # CODEOWNERS patterns are repo-relative, so run from the repo root and + # use a relative sql dir/query path + monkeypatch.chdir(tmp_path) + sql = Path("sql") + _dataset( + sql, "src_external", base_acl="restricted", members=["workgroup:acme/x"] + ) + dest = _dataset(sql, "owned_derived", members=[BROAD]) + q = _query(dest, "out_v1", _select("src_external")) + + entries = self._codeowners(tmp_path) + findings = check_query(q, "sql", codeowners=entries) + assert len(findings) == 1 + assert findings[0]["gated"] is True + assert findings[0]["owners"] == ["@mozilla/team"] + + +# A minimal gated_datasets config mirroring the bqetl_project.yaml shape: +# keyed by project, then dataset -- a narrow dataset, a broad dataset with +# per-doctype gated tables, and a second project. +GATED = { + PROJECT: { + "acme_stable": {"default": ["workgroup:acme/data-viewers"]}, + "telemetry_stable": { + "default": [BROAD], + "tables": { + "account_ecosystem": [], + "serp_categorization": ["workgroup:revenue/cat2"], + }, + }, + }, + "other-project": { + "glam_stable": {"default": ["workgroup:glam/data-viewers"]}, + }, +} + + +class TestGatedDatasets: + def test_table_override_beats_default_and_strips_version(self): + # version suffix (_v5) stripped before matching the doctype key + assert gated_readers( + GATED, PROJECT, "telemetry_stable", "serp_categorization_v5" + ) == {"workgroup:revenue/cat2"} + + def test_falls_back_to_default_when_table_unlisted(self): + assert gated_readers(GATED, PROJECT, "telemetry_stable", "main_v5") == {BROAD} + + def test_empty_members_is_restricted(self): + assert ( + gated_readers(GATED, PROJECT, "telemetry_stable", "account_ecosystem_v1") + == set() + ) + + def test_unlisted_dataset_returns_none(self): + assert gated_readers(GATED, PROJECT, "not_listed_stable", "t_v1") is None + + def test_resolves_per_project(self): + assert gated_readers(GATED, "other-project", "glam_stable", "t_v1") == { + "workgroup:glam/data-viewers" + } + + def test_dataset_not_matched_across_projects(self): + # acme_stable is listed under PROJECT, not other-project + assert gated_readers(GATED, "other-project", "acme_stable", "t_v1") is None + assert gated_readers(GATED, "unlisted-project", "acme_stable", "t_v1") is None + + def test_source_access_uses_gated_config(self, tmp_path): + acc = source_access( + str(tmp_path / "sql"), PROJECT, "acme_stable", "t_v1", GATED + ) + assert acc.readers == {"workgroup:acme/data-viewers"} + assert acc.sensitive is True + + def test_source_access_unlisted_ingestion_is_broad(self, tmp_path): + # a stable/live dataset absent from a listed project gets the broad default + acc = source_access( + str(tmp_path / "sql"), PROJECT, "unlisted_stable", "t_v1", GATED + ) + assert acc.readers == {BROAD} + assert acc.sensitive is False + + def test_source_access_unlisted_project_unresolved(self, tmp_path): + # a stable dataset in a project not keyed in gated_datasets -> unresolved + acc = source_access( + str(tmp_path / "sql"), "unlisted-project", "some_stable", "t_v1", GATED + ) + assert acc is None + + def test_source_access_second_ingestion_project_is_broad(self, tmp_path): + # other-project is keyed in gated_datasets, so its unlisted stable dataset + # is broad rather than unresolved + acc = source_access( + str(tmp_path / "sql"), "other-project", "some_stable", "t_v1", GATED + ) + assert acc.readers == {BROAD} + assert acc.sensitive is False + + def test_flags_gated_dataset_source(self, tmp_path): + sql = tmp_path / "sql" + dest = _dataset(sql, "broad_derived", members=[BROAD]) + q = _query(dest, "out_v1", _select("acme_stable", table="t_v1")) + + findings = check_query(q, str(sql), gated_datasets=GATED) + assert len(findings) == 1 + assert findings[0]["source"] == f"{PROJECT}.acme_stable" + assert findings[0]["extra_readers"] == [BROAD] + + def test_flags_gated_table_in_broad_dataset(self, tmp_path): + sql = tmp_path / "sql" + dest = _dataset(sql, "broad_derived", members=[BROAD]) + q = _query( + dest, "out_v1", _select("telemetry_stable", table="serp_categorization_v1") + ) + + findings = check_query(q, str(sql), gated_datasets=GATED) + assert len(findings) == 1 + assert findings[0]["extra_readers"] == [BROAD] + + def test_no_flag_for_broad_table_in_broad_dataset(self, tmp_path): + sql = tmp_path / "sql" + dest = _dataset(sql, "broad_derived", members=[BROAD]) + # telemetry_stable.main is broad (default), dest is broad -> no widening + q = _query(dest, "out_v1", _select("telemetry_stable", table="main_v5")) + + assert check_query(q, str(sql), gated_datasets=GATED) == [] + + def test_no_flag_for_unlisted_ingestion_source(self, tmp_path): + sql = tmp_path / "sql" + dest = _dataset(sql, "broad_derived", members=[BROAD]) + q = _query(dest, "out_v1", _select("unlisted_stable", table="t_v1")) + + assert check_query(q, str(sql), gated_datasets=GATED) == [] + + def test_empty_readers_source_is_sensitive(self, tmp_path): + # `[]` readers means nobody is authorized -> strictly narrower than broad + acc = source_access( + str(tmp_path / "sql"), + PROJECT, + "telemetry_stable", + "account_ecosystem_v1", + GATED, + ) + assert acc.readers == set() + assert acc.sensitive is True + + def test_flags_empty_reader_table_written_to_broad_dest(self, tmp_path): + # regression: a `default: []` / empty-member gated table must be flagged + sql = tmp_path / "sql" + dest = _dataset(sql, "broad_derived", members=[BROAD]) + q = _query( + dest, "out_v1", _select("telemetry_stable", table="account_ecosystem_v1") + ) + + findings = check_query(q, str(sql), gated_datasets=GATED) + assert len(findings) == 1 + assert findings[0]["extra_readers"] == [BROAD] + + +class TestSensitivityCommand: + """`./bqetl data_governance sensitivity` CLI wrapper.""" + + def _widening_query(self, tmp_path): + sql = tmp_path / "sql" + _dataset( + sql, "src_external", base_acl="restricted", members=["workgroup:acme/x"] + ) + dest = _dataset(sql, "broad_derived", members=[BROAD]) + q = _query(dest, "out_v1", _select("src_external")) + return sql, q + + def _codeowners(self, tmp_path, body): + p = tmp_path / "CODEOWNERS" + p.write_text(body) + return str(p) + + def test_ungated_flow_exits_2_and_requests_review(self, tmp_path): + sql, q = self._widening_query(tmp_path) + # CODEOWNERS that owns nothing here -> the flow is ungated + codeowners = self._codeowners(tmp_path, "/sql/owned/** @mozilla/team\n") + + result = CliRunner().invoke( + sensitivity_cmd, + [q, "--sql_dir", str(sql), "--codeowners", codeowners], + ) + assert result.exit_code == 2 + assert f"{PROJECT}.src_external" in result.output + assert "needs Data Platform" in result.output + assert "::warning::" in result.stderr + + def test_gated_flow_exits_0(self, tmp_path): + sql, q = self._widening_query(tmp_path) + # CODEOWNERS owns everything -> the flow is already reviewed, not ungated + codeowners = self._codeowners(tmp_path, "* @mozilla/team\n") + + result = CliRunner().invoke( + sensitivity_cmd, + [q, "--sql_dir", str(sql), "--codeowners", codeowners], + ) + assert result.exit_code == 0 + assert "already reviewed by @mozilla/team" in result.output + + def test_clean_query_exits_0(self, tmp_path): + sql = tmp_path / "sql" + dest = _dataset(sql, "broad_derived", members=[BROAD]) + q = _query(dest, "out_v1", "SELECT 1") + codeowners = self._codeowners(tmp_path, "* @mozilla/team\n") + + result = CliRunner().invoke( + sensitivity_cmd, + [q, "--sql_dir", str(sql), "--codeowners", codeowners], + ) + assert result.exit_code == 0 + assert "no ungated sensitive-data flows" in result.stderr + + def test_json_output(self, tmp_path): + import json + + sql, q = self._widening_query(tmp_path) + codeowners = self._codeowners(tmp_path, "* @mozilla/team\n") + + result = CliRunner().invoke( + sensitivity_cmd, + [q, "--sql_dir", str(sql), "--codeowners", codeowners, "--json"], + ) + assert result.exit_code == 0 + # the JSON array is on stdout; a stderr status line may follow it + payload = result.output[: result.output.rindex("]") + 1] + findings = json.loads(payload) + assert findings[0]["source"] == f"{PROJECT}.src_external"