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/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/airflow/dags/common/public_suffix.py b/airflow/dags/common/public_suffix.py new file mode 100644 index 00000000..6643d8f2 --- /dev/null +++ b/airflow/dags/common/public_suffix.py @@ -0,0 +1,119 @@ +"""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_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)"}, + ) + with urllib.request.urlopen(req) as resp: + text = resp.read().decode("utf-8") + + current_section = None + rules: List[Dict[str, Any]] = [] + seen = set() + + for raw_line in text.splitlines(): + line = raw_line.strip() + 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 + + is_exception = line.startswith("!") + is_wildcard = line.startswith("*.") + + 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, is_exception, current_section) + if key not in seen: + seen.add(key) + 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_list( + project_id: str, + dataset_id: str = "urls", + table_name: str = "public_suffix_list", +) -> int: + """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}" + + job_config = bigquery.LoadJobConfig( + schema=[ + bigquery.SchemaField( + "suffix", + "STRING", + mode="REQUIRED", + description="Domain public suffix in ASCII punycode", + ), + bigquery.SchemaField( + "is_wildcard", + "BOOLEAN", + 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, + ) + + 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 d23e41cb..b2ac6e45 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_list 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 (ICANN + private) and sync into BigQuery + @task(task_id="sync_public_suffix_list") + def update_public_suffix_table() -> int: + """Download latest Public Suffix List and load into BigQuery.""" + return sync_public_suffix_list(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..54b26091 100644 --- a/definitions/declarations/httparchive.js +++ b/definitions/declarations/httparchive.js @@ -14,6 +14,12 @@ }) ) +// Public Suffix List synced via Airflow DAG (crawl_complete) +declare({ + schema: 'urls', + name: 'public_suffix_list' +}) + operate('httparchive_project_options').queries(` ALTER PROJECT httparchive SET OPTIONS ( \`region-US.default_sql_dialect_option\` = 'only_google_sql', 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