From b22d387abd3e7600cc071a0fbc9a45d9a470d7ba Mon Sep 17 00:00:00 2001 From: Max Ostapenko <1611259+max-ostapenko@users.noreply.github.com> Date: Sun, 4 Oct 2026 01:20:03 +0200 Subject: [PATCH 1/3] feat: add sync_public_suffix_private task to crawl_complete DAG Signed-off-by: Max Ostapenko <1611259+max-ostapenko@users.noreply.github.com> --- .gitignore | 2 + airflow/dags/common/public_suffix.py | 74 +++++++++++++++++++++++++ airflow/dags/crawl_complete.py | 16 +++++- definitions/declarations/httparchive.js | 6 ++ 4 files changed, 95 insertions(+), 3 deletions(-) create mode 100644 airflow/dags/common/public_suffix.py diff --git a/.gitignore b/.gitignore index f287f375..dffa26a5 100644 --- a/.gitignore +++ b/.gitignore @@ -1,3 +1,5 @@ .DS_Store .vscode/ **/node_modules/ +__pycache__/ +*.pyc diff --git a/airflow/dags/common/public_suffix.py b/airflow/dags/common/public_suffix.py new file mode 100644 index 00000000..69d4d8f7 --- /dev/null +++ b/airflow/dags/common/public_suffix.py @@ -0,0 +1,74 @@ +"""Public Suffix List synchronization utilities for BigQuery.""" + +from typing import Any, Dict, List +import urllib.request +import idna +from google.cloud import bigquery + +PSL_URL = "https://publicsuffix.org/list/public_suffix_list.dat" + + +def fetch_psl_private_domains(url: str = PSL_URL) -> List[Dict[str, Any]]: + """Fetch the Public Suffix List and parse the PRIVATE DOMAINS section.""" + req = urllib.request.Request( + url, + headers={"User-Agent": "HTTPArchive-Airflow/1.0 (+https://httparchive.org)"}, + ) + with urllib.request.urlopen(req) as resp: + text = resp.read().decode("utf-8") + + begin = text.find("// ===BEGIN PRIVATE DOMAINS===") + end = text.find("// ===END PRIVATE DOMAINS===") + if begin == -1 or end == -1: + raise ValueError("PRIVATE DOMAINS section not found in Public Suffix List") + + rules: List[Dict[str, Any]] = [] + seen = set() + for raw_line in text[begin:end].splitlines(): + line = raw_line.strip() + if not line or line.startswith("//"): + continue + if line.startswith("!"): + raise ValueError(f"Exception rules are not supported: {line}") + is_wildcard = line.startswith("*.") + rule = line[2:] if is_wildcard else line + puny = idna.encode(rule).decode("ascii") + key = (puny, is_wildcard) + if key not in seen: + seen.add(key) + rules.append({"suffix": puny, "is_wildcard": is_wildcard}) + + return rules + + +def sync_public_suffix_private( + project_id: str, + dataset_id: str = "urls", + table_name: str = "public_suffix_private", +) -> int: + """Download and load private public suffix rules into BigQuery.""" + rules = fetch_psl_private_domains() + client = bigquery.Client(project=project_id) + table_id = f"{project_id}.{dataset_id}.{table_name}" + + job_config = bigquery.LoadJobConfig( + schema=[ + bigquery.SchemaField( + "suffix", + "STRING", + mode="REQUIRED", + description="Private domain public suffix (ASCII punycode)", + ), + bigquery.SchemaField( + "is_wildcard", + "BOOLEAN", + mode="REQUIRED", + description="Whether suffix is a wildcard rule (*.)", + ), + ], + write_disposition=bigquery.WriteDisposition.WRITE_TRUNCATE, + ) + + job = client.load_table_from_json(rules, table_id, job_config=job_config) + job.result() + return len(rules) diff --git a/airflow/dags/crawl_complete.py b/airflow/dags/crawl_complete.py index d23e41cb..b90d44a6 100644 --- a/airflow/dags/crawl_complete.py +++ b/airflow/dags/crawl_complete.py @@ -4,6 +4,7 @@ from airflow import DAG import sys from pathlib import Path +from airflow.decorators import task from airflow.providers.google.cloud.operators.dataform import ( DataformCreateCompilationResultOperator, DataformCreateWorkflowInvocationOperator, @@ -21,6 +22,7 @@ DEFAULT_DAG_ARGS, PROJECT_ID, ) +from common.public_suffix import sync_public_suffix_private with DAG( dag_id="crawl_complete", @@ -46,7 +48,15 @@ mode="reschedule", ) - # 2. Create compilation result from the Dataform production release configuration + # 2. Download latest Public Suffix List private domains and sync into BigQuery + @task(task_id="sync_public_suffix_private") + def update_public_suffix_table() -> int: + """Download latest Public Suffix List private domains and load into BigQuery.""" + return sync_public_suffix_private(project_id=PROJECT_ID) + + sync_psl = update_public_suffix_table() + + # 3. Create compilation result from the Dataform production release configuration create_compilation_result = DataformCreateCompilationResultOperator( task_id="create_compilation_result", project_id=PROJECT_ID, @@ -60,7 +70,7 @@ }, ) - # 3. Invoke Dataform workflow with crawl_complete and crawl_complete_reports tags + # 4. Invoke Dataform workflow with crawl_complete and crawl_complete_reports tags run_dataform_workflow = DataformCreateWorkflowInvocationOperator( task_id="run_dataform_workflow", project_id=PROJECT_ID, @@ -85,4 +95,4 @@ }, ) - wait_for_crawl_complete >> create_compilation_result >> run_dataform_workflow + wait_for_crawl_complete >> sync_psl >> create_compilation_result >> run_dataform_workflow diff --git a/definitions/declarations/httparchive.js b/definitions/declarations/httparchive.js index e374e702..1c200b67 100644 --- a/definitions/declarations/httparchive.js +++ b/definitions/declarations/httparchive.js @@ -14,6 +14,12 @@ }) ) +// Public Suffix List private domains synced via Airflow DAG (crawl_complete) +declare({ + schema: 'urls', + name: 'public_suffix_private' +}) + operate('httparchive_project_options').queries(` ALTER PROJECT httparchive SET OPTIONS ( \`region-US.default_sql_dialect_option\` = 'only_google_sql', From 4a16993897439c3ac7ad68c9ed5e4cd1ef4b6b8c Mon Sep 17 00:00:00 2001 From: Max Ostapenko <1611259+max-ostapenko@users.noreply.github.com> Date: Sun, 4 Oct 2026 01:34:23 +0200 Subject: [PATCH 2/3] refactor: expand public suffix list synchronization to include ICANN domains and exception rules Signed-off-by: Max Ostapenko <1611259+max-ostapenko@users.noreply.github.com> --- airflow/dags/common/public_suffix.py | 83 +++++++++++++++++++------ airflow/dags/crawl_complete.py | 10 +-- definitions/declarations/httparchive.js | 4 +- 3 files changed, 71 insertions(+), 26 deletions(-) diff --git a/airflow/dags/common/public_suffix.py b/airflow/dags/common/public_suffix.py index 69d4d8f7..6643d8f2 100644 --- a/airflow/dags/common/public_suffix.py +++ b/airflow/dags/common/public_suffix.py @@ -8,8 +8,8 @@ PSL_URL = "https://publicsuffix.org/list/public_suffix_list.dat" -def fetch_psl_private_domains(url: str = PSL_URL) -> List[Dict[str, Any]]: - """Fetch the Public Suffix List and parse the PRIVATE DOMAINS section.""" +def fetch_public_suffix_list(url: str = PSL_URL) -> List[Dict[str, Any]]: + """Fetch the full Public Suffix List and parse ICANN and PRIVATE domain sections.""" req = urllib.request.Request( url, headers={"User-Agent": "HTTPArchive-Airflow/1.0 (+https://httparchive.org)"}, @@ -17,37 +17,60 @@ def fetch_psl_private_domains(url: str = PSL_URL) -> List[Dict[str, Any]]: with urllib.request.urlopen(req) as resp: text = resp.read().decode("utf-8") - begin = text.find("// ===BEGIN PRIVATE DOMAINS===") - end = text.find("// ===END PRIVATE DOMAINS===") - if begin == -1 or end == -1: - raise ValueError("PRIVATE DOMAINS section not found in Public Suffix List") - + current_section = None rules: List[Dict[str, Any]] = [] seen = set() - for raw_line in text[begin:end].splitlines(): + + for raw_line in text.splitlines(): line = raw_line.strip() - if not line or line.startswith("//"): + if line.startswith("// ===BEGIN ICANN DOMAINS==="): + current_section = "ICANN" + continue + elif line.startswith("// ===END ICANN DOMAINS==="): + current_section = None + continue + elif line.startswith("// ===BEGIN PRIVATE DOMAINS==="): + current_section = "PRIVATE" + continue + elif line.startswith("// ===END PRIVATE DOMAINS==="): + current_section = None + continue + + if not current_section or not line or line.startswith("//"): continue - if line.startswith("!"): - raise ValueError(f"Exception rules are not supported: {line}") + + is_exception = line.startswith("!") is_wildcard = line.startswith("*.") - rule = line[2:] if is_wildcard else line + + if is_exception: + rule = line[1:] + elif is_wildcard: + rule = line[2:] + else: + rule = line + puny = idna.encode(rule).decode("ascii") - key = (puny, is_wildcard) + key = (puny, is_wildcard, is_exception, current_section) if key not in seen: seen.add(key) - rules.append({"suffix": puny, "is_wildcard": is_wildcard}) + rules.append({ + "suffix": puny, + "is_wildcard": is_wildcard, + "is_exception": is_exception, + "is_private": current_section == "PRIVATE", + "section": current_section, + }) return rules -def sync_public_suffix_private( +def sync_public_suffix_list( project_id: str, dataset_id: str = "urls", - table_name: str = "public_suffix_private", + table_name: str = "public_suffix_list", ) -> int: - """Download and load private public suffix rules into BigQuery.""" - rules = fetch_psl_private_domains() + """Download and load the full Public Suffix List (ICANN + private) into BigQuery.""" + rules = fetch_public_suffix_list() client = bigquery.Client(project=project_id) table_id = f"{project_id}.{dataset_id}.{table_name}" @@ -57,7 +80,7 @@ def sync_public_suffix_private( "suffix", "STRING", mode="REQUIRED", - description="Private domain public suffix (ASCII punycode)", + description="Domain public suffix in ASCII punycode", ), bigquery.SchemaField( "is_wildcard", @@ -65,6 +88,24 @@ def sync_public_suffix_private( mode="REQUIRED", description="Whether suffix is a wildcard rule (*.)", ), + bigquery.SchemaField( + "is_exception", + "BOOLEAN", + mode="REQUIRED", + description="Whether suffix is an exception rule (!)", + ), + bigquery.SchemaField( + "is_private", + "BOOLEAN", + mode="REQUIRED", + description="Whether suffix is in the PRIVATE section", + ), + bigquery.SchemaField( + "section", + "STRING", + mode="REQUIRED", + description="Public Suffix List section (ICANN or PRIVATE)", + ), ], write_disposition=bigquery.WriteDisposition.WRITE_TRUNCATE, ) @@ -72,3 +113,7 @@ def sync_public_suffix_private( job = client.load_table_from_json(rules, table_id, job_config=job_config) job.result() return len(rules) + + +# Backward-compatibility alias +sync_public_suffix_private = sync_public_suffix_list diff --git a/airflow/dags/crawl_complete.py b/airflow/dags/crawl_complete.py index b90d44a6..b2ac6e45 100644 --- a/airflow/dags/crawl_complete.py +++ b/airflow/dags/crawl_complete.py @@ -22,7 +22,7 @@ DEFAULT_DAG_ARGS, PROJECT_ID, ) -from common.public_suffix import sync_public_suffix_private +from common.public_suffix import sync_public_suffix_list with DAG( dag_id="crawl_complete", @@ -48,11 +48,11 @@ mode="reschedule", ) - # 2. Download latest Public Suffix List private domains and sync into BigQuery - @task(task_id="sync_public_suffix_private") + # 2. Download latest Public Suffix List (ICANN + private) and sync into BigQuery + @task(task_id="sync_public_suffix_list") def update_public_suffix_table() -> int: - """Download latest Public Suffix List private domains and load into BigQuery.""" - return sync_public_suffix_private(project_id=PROJECT_ID) + """Download latest Public Suffix List and load into BigQuery.""" + return sync_public_suffix_list(project_id=PROJECT_ID) sync_psl = update_public_suffix_table() diff --git a/definitions/declarations/httparchive.js b/definitions/declarations/httparchive.js index 1c200b67..54b26091 100644 --- a/definitions/declarations/httparchive.js +++ b/definitions/declarations/httparchive.js @@ -14,10 +14,10 @@ }) ) -// Public Suffix List private domains synced via Airflow DAG (crawl_complete) +// Public Suffix List synced via Airflow DAG (crawl_complete) declare({ schema: 'urls', - name: 'public_suffix_private' + name: 'public_suffix_list' }) operate('httparchive_project_options').queries(` From 069420676c10bef22f7b1faa6aba0b6771e2fdc2 Mon Sep 17 00:00:00 2001 From: Max Ostapenko <1611259+max-ostapenko@users.noreply.github.com> Date: Sun, 4 Oct 2026 02:13:49 +0200 Subject: [PATCH 3/3] docs: update documentation to reflect Airflow orchestration for pipelines Signed-off-by: Max Ostapenko <1611259+max-ostapenko@users.noreply.github.com> --- README.md | 34 ++++++++-------------------------- docs/dataform.md | 34 +++++++++++++++++++++++++++++----- 2 files changed, 37 insertions(+), 31 deletions(-) diff --git a/README.md b/README.md index 7f7b8b68..b167e77b 100644 --- a/README.md +++ b/README.md @@ -2,35 +2,17 @@ This repository handles the HTTP Archive data pipeline, which takes the results of the monthly HTTP Archive run and saves this to the `httparchive` dataset in BigQuery. -## Pipelines +## Pipelines & Orchestration -The pipelines are run in Dataform service in Google Cloud Platform (GCP) and are kicked off automatically on crawl completion and other events. The code in the `main` branch is used on each triggered pipeline run. +Pipelines are transformed via Dataform and orchestrated by Google Cloud Composer (Apache Airflow) in GCP. Airflow DAGs reside in [`airflow/dags/`](./airflow/dags/) and are synced to Cloud Composer on merge to `main`. -### HTTP Archive Crawl +| Pipeline | DAG | Dataform Tags | Primary Outputs | +| ---------------------- | -------------------------------------------------------------------- | ------------------------------------------ | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| **HTTP Archive Crawl** | [`airflow/dags/crawl_complete.py`](./airflow/dags/crawl_complete.py) | `crawl_complete`, `crawl_complete_reports` | `httparchive.crawl.*` ([Analytics Hub](https://console.cloud.google.com/bigquery/analytics-hub/discovery/projects/httparchive/locations/us/dataExchanges/httparchive/listings/crawl)), `httparchive.blink_features.usage` ([chromestatus.com](https://chromestatus.com/metrics/feature/timeline/popularity/2089)) | +| **Technology Report** | [`airflow/dags/crux_ready.py`](./airflow/dags/crux_ready.py) | `crux_ready`, `crux_ready_reports` | `httparchive.reports.cwv_tech_*`, `httparchive.reports.tech_*` ([Tech Report](https://httparchive.org/reports/techreport/landing)) | -Tag: `crawl_complete` - -- Crawl dataset `httparchive.crawl.*` - - Consumers: - - - public dataset and [BQ Sharing Listing](https://console.cloud.google.com/bigquery/analytics-hub/discovery/projects/httparchive/locations/us/dataExchanges/httparchive/listings/crawl) - -- Blink Features Report `httparchive.blink_features.usage` - - Consumers: - - - [chromestatus.com](https://chromestatus.com/metrics/feature/timeline/popularity/2089) - -### HTTP Archive Technology Report - -Tag: `crux_ready` - -- `httparchive.reports.cwv_tech_*` and `httparchive.reports.tech_*` - - Consumers: - - - [HTTP Archive Tech Report](https://httparchive.org/reports/techreport/landing) +For complete orchestration details (triggers, sensors, pre-tasks) and development workspace guidelines, see [Dataform Documentation](docs/dataform.md). +For overall GCP infrastructure and data flows, see [Infrastructure Overview](../tech-report-apis/docs/infra.md). ## Development Setup diff --git a/docs/dataform.md b/docs/dataform.md index b04988b6..897162e9 100644 --- a/docs/dataform.md +++ b/docs/dataform.md @@ -4,9 +4,31 @@ Runs the batch processing workflows. There are two Dataform repositories for [de The test repository is used [for development and testing purposes](https://cloud.google.com/dataform/docs/workspaces) and not connected to the rest of the pipeline infra. -Pipeline can be [run manually](https://cloud.google.com/dataform/docs/code-lifecycle) from the Dataform UI. +Pipelines can be [run manually](https://cloud.google.com/dataform/docs/code-lifecycle) from the Dataform UI or orchestrated automatically via Apache Airflow. -The infrastructure configurations (formerly located in `infra/`) have been migrated to the [tech-report-apis](../../tech-report-apis) monorepo. Refer to [terraform/dataform.tf](../../tech-report-apis/terraform/dataform.tf) in that repository for Dataform project IaC. +The infrastructure configurations are in the [tech-report-apis](../../tech-report-apis) monorepo. Refer to [terraform/dataform.tf](../../tech-report-apis/terraform/dataform.tf) and [terraform/airflow.tf](../../tech-report-apis/terraform/airflow.tf) for Dataform and Composer IaC. + +## Pipeline Orchestration (Apache Airflow) + +Production pipeline invocations are managed by Google Cloud Composer ([httparchive-pipelines](https://console.cloud.google.com/composer/environments/detail/us-central1/httparchive-pipelines/overview?authuser=2&project=httparchive)). The DAG definitions are maintained in [`airflow/dags/`](../airflow/dags/): + +1. **`crawl_complete` DAG** ([`airflow/dags/crawl_complete.py`](../airflow/dags/crawl_complete.py)): + - Triggered by the `crawl-complete` Pub/Sub topic via `PubSubPullSensor`. + - Executes `sync_public_suffix_list` task to fetch the full Public Suffix List (ICANN + private) into `httparchive.urls.public_suffix_list`. + - Creates a compilation result against the production release config and invokes Dataform with tags: + - `crawl_complete` + - `crawl_complete_reports` + +2. **`crux_ready` DAG** ([`airflow/dags/crux_ready.py`](../airflow/dags/crux_ready.py)): + - Scheduled at 08:00, 12:00, and 16:00 UTC during the CrUX release window (8th–14th of each month). + - Sensor queries BigQuery to verify that the previous month's `chrome-ux-report` table has been published. + - Creates a compilation result and invokes Dataform with tags: + - `crux_ready` + - `crux_ready_reports` + +### DAG Deployment + +Airflow DAGs are automatically synchronized from `airflow/dags/` to the Cloud Composer Cloud Storage bucket (`gs://us-central1-httparchive-pip-77e1b883-bucket/dags`) via the [Deploy Airflow DAGs GitHub Actions workflow](../.github/workflows/deploy_dags.yaml) on merge to `main`. ## Dataform Development Workspace @@ -27,9 +49,11 @@ _Some useful hints:_ ## Repository Structure +- `airflow/` - Cloud Composer / Apache Airflow DAG definitions and helper modules + - `dags/` - Production Airflow DAGs (`crawl_complete.py`, `crux_ready.py`) and shared utilities (`common/`) - `definitions/` - Contains the core Dataform SQL definitions and declarations - `output/` - Contains the main pipeline transformation logic - - `declarations/` - Contains referenced tables/views declarations and other resources definitions + - `declarations/` - Contains referenced tables/views declarations and external resources - `includes/` - Contains shared JavaScript utilities and constants - `docs/` - Additional documentation @@ -40,10 +64,10 @@ GitHub PAT saved to a [Secret Manager secret](https://console.cloud.google.com/s - repository: HTTPArchive/dataform - permissions: - Commit statuses: read - - Contents: read, write + - Contents: read ## Monitoring +- [Airflow Web UI](https://225068d09fdc4614b12e9aa283e4e4ab-dot-us-central1.composer.googleusercontent.com) - [Production Dataform workflow execution logs](https://console.cloud.google.com/bigquery/dataform/locations/us-central1/repositories/crawl-data/details/workflows?authuser=7&project=httparchive) - - [Dataform Workflow Invocation Failed](https://console.cloud.google.com/monitoring/alerting/policies/16526940745374967367?authuser=7&project=httparchive) policy