Skip to content
Open
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
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,17 @@ sequenceDiagram
Finds GTFS datasets that are missing a validation report **or** have a report from an
older validator version, then triggers a GCP Workflow for each one.

By default "older version" means "not the exact version returned by the validator
endpoint". Set `filter_validator_version_prefix` (e.g. `"8."`) to instead target
datasets that lack **any** report matching that major-version prefix — useful when you
want to re-validate everything missing an `8.*` report regardless of the exact patch
version the endpoint reports.

By default the task only considers **each feed's latest dataset**. Set
`include_all_datasets: true` to consider **every** dataset of the feed — needed when
backfilling a historical window (e.g. so a feed's six-month reliability window is
complete), not just its most recent dataset.

The task is **resumable**: if it times out mid-loop, calling it again skips datasets
that were already triggered (tracked in `task_execution_log`).

Expand All @@ -74,8 +85,11 @@ that were already triggered (tracked in `task_execution_log`).
"validator_endpoint": "https://stg-gtfs-validator-web-mbzoxaljzq-ue.a.run.app",
"bypass_db_update": false,
"filter_after_in_days": 30,
"filter_downloaded_after": "2026-03-01",
"filter_downloaded_before": "2026-06-01",
"filter_statuses": ["active"],
"filter_op_statuses": ["published"],
"filter_validator_version_prefix": "8.",
"force_update": false,
"limit": 10,
"reports_bucket_name": "stg-gtfs-validator-results"
Expand All @@ -88,8 +102,12 @@ that were already triggered (tracked in `task_execution_log`).
| `validator_endpoint` | string | env-derived | Validator service URL to use and fetch version from |
| `bypass_db_update` | bool | `false` | When `true`, results are NOT written to DB/API (use for pre-release runs) |
| `filter_after_in_days` | int | `null` | Restrict to datasets downloaded within the last N days. Omit to include all datasets |
| `filter_downloaded_after` | string (ISO date) | `null` | Restrict to datasets downloaded **at or after** this date (inclusive), e.g. `"2026-03-01"`. Omit for no lower bound |
| `filter_downloaded_before` | string (ISO date) | `null` | Restrict to datasets downloaded **strictly before** this date (exclusive), e.g. `"2026-06-01"`. Omit for no upper bound |
| `filter_statuses` | list[str] | `null` | Filter feeds by status (e.g. `["active", "inactive"]`). Omit for all statuses |
| `filter_op_statuses` | list[str] | `["published"]` | Filter feeds by operational status. Accepted values: `"published"`, `"unpublished"`, `"wip"` |
| `filter_validator_version_prefix` | string | `null` | When set (e.g. `"8."`), only include datasets that do **not** already have a report whose `validator_version` starts with this prefix. Overrides the default exact-version comparison |
| `include_all_datasets` | bool | `false` | When `true`, consider **all** datasets of each feed, not just the latest one. Use for backfilling a historical window |
| `force_update` | bool | `false` | Re-trigger even when a current report already exists |
| `limit` | int | `null` | Cap the number of workflows triggered per call — useful for end-to-end testing |
| `reports_bucket_name` | string | env-derived | Override the GCS bucket where validator results are stored. Use when running in prod but pointing to the staging validator (e.g. `"stg-gtfs-validator-results"`) |
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -59,10 +59,20 @@ def rebuild_missing_validation_reports_handler(payload) -> dict:
"dry_run": bool, # [optional] If True, count only — do not trigger workflows. Default: True
"filter_after_in_days": int, # [optional] Restrict to datasets downloaded within the last N days.
# If omitted, all datasets are considered regardless of age.
"filter_downloaded_after": str, # [optional] ISO date (e.g. "2026-03-01"). Restrict to datasets
# downloaded at or after this date (inclusive).
"filter_downloaded_before": str, # [optional] ISO date (e.g. "2026-06-01"). Restrict to datasets
# downloaded strictly before this date (exclusive).
"filter_statuses": list[str],# [optional] Filter feeds by status
"filter_op_statuses": list[str],# [optional] Filter feeds by operational status.
# Default: ["published"]
# Accepted values: "published", "unpublished", "wip"
"filter_validator_version_prefix": str, # [optional] Only include datasets that do NOT already
# have a validation report whose validator_version starts with this
# prefix (e.g. "8." → datasets missing any 8.* report). When omitted,
# datasets are selected by comparing against the exact target version.
"include_all_datasets": bool,# [optional] If True, consider ALL datasets of each feed, not just the
# latest one. Use for backfilling historical windows. Default: False
"validator_endpoint": str, # [optional] Override validator URL (e.g. staging). Default: env-derived URL.
"bypass_db_update": bool, # [optional] If True, results are NOT written to the DB/API (pre-release runs).
Default: False
Expand All @@ -76,8 +86,12 @@ def rebuild_missing_validation_reports_handler(payload) -> dict:
(
dry_run,
filter_after_in_days,
filter_downloaded_after,
filter_downloaded_before,
filter_statuses,
filter_op_statuses,
filter_validator_version_prefix,
include_all_datasets,
prod_env,
validator_endpoint,
bypass_db_update,
Expand All @@ -91,8 +105,12 @@ def rebuild_missing_validation_reports_handler(payload) -> dict:
bypass_db_update=bypass_db_update,
dry_run=dry_run,
filter_after_in_days=filter_after_in_days,
filter_downloaded_after=filter_downloaded_after,
filter_downloaded_before=filter_downloaded_before,
filter_statuses=filter_statuses,
filter_op_statuses=filter_op_statuses,
filter_validator_version_prefix=filter_validator_version_prefix,
include_all_datasets=include_all_datasets,
prod_env=prod_env,
force_update=force_update,
limit=limit,
Expand All @@ -106,8 +124,12 @@ def rebuild_missing_validation_reports(
bypass_db_update: bool = False,
dry_run: bool = True,
filter_after_in_days: Optional[int] = None,
filter_downloaded_after: Optional[str] = None,
filter_downloaded_before: Optional[str] = None,
filter_statuses: List[str] | None = None,
filter_op_statuses: List[str] | None = None,
filter_validator_version_prefix: Optional[str] = None,
include_all_datasets: bool = False,
prod_env: bool = False,
force_update: bool = False,
limit: Optional[int] = None,
Expand All @@ -124,9 +146,19 @@ def rebuild_missing_validation_reports(
dry_run: If True, count only — do not trigger workflows. Default: True
filter_after_in_days: Restrict to datasets downloaded within the last N days.
If None (default), all datasets are considered regardless of age.
filter_downloaded_after: ISO date (e.g. "2026-03-01"). Restrict to datasets
downloaded at or after this date (inclusive). If None (default), no lower bound.
filter_downloaded_before: ISO date (e.g. "2026-06-01"). Restrict to datasets
downloaded strictly before this date (exclusive). If None (default), no upper bound.
filter_statuses: Filter feeds by status. Default: None (all)
filter_op_statuses: Filter feeds by operational status.
Default: ["published"]. Accepted: "published", "unpublished", "wip".
filter_validator_version_prefix: If set (e.g. "8."), only include datasets that do
NOT already have a validation report whose validator_version starts with this
prefix. When None (default), datasets are selected by comparing against the exact
target validator version.
include_all_datasets: If True, consider ALL datasets of each feed, not just the
latest one. Use for backfilling a historical window. Default: False.
prod_env: True if targeting the production environment. Default: False
force_update: Re-trigger even if a report already exists. Default: False
limit: Max datasets to trigger per call (for end-to-end testing). Default: unlimited
Expand All @@ -146,10 +178,14 @@ def rebuild_missing_validation_reports(
validator_version=validator_version,
force_update=force_update,
filter_after_in_days=filter_after_in_days,
filter_downloaded_after=filter_downloaded_after,
filter_downloaded_before=filter_downloaded_before,
filter_statuses=filter_statuses,
filter_op_statuses=(
filter_op_statuses if filter_op_statuses is not None else ["published"]
),
filter_validator_version_prefix=filter_validator_version_prefix,
include_all_datasets=include_all_datasets,
)
total_candidates = len(datasets)
logging.info("Found %s candidate datasets", total_candidates)
Expand Down Expand Up @@ -227,7 +263,11 @@ def rebuild_missing_validation_reports(
"validator_endpoint": validator_endpoint,
"bypass_db_update": bypass_db_update,
"filter_after_in_days": filter_after_in_days,
"filter_downloaded_after": filter_downloaded_after,
"filter_downloaded_before": filter_downloaded_before,
"filter_statuses": filter_statuses,
"filter_validator_version_prefix": filter_validator_version_prefix,
"include_all_datasets": include_all_datasets,
"prod_env": prod_env,
"force_update": force_update,
"limit": limit,
Expand All @@ -253,8 +293,12 @@ def _get_datasets_for_validation(
validator_version: str,
force_update: bool,
filter_after_in_days: Optional[int],
filter_downloaded_after: Optional[str],
filter_downloaded_before: Optional[str],
filter_statuses: Optional[List[str]],
filter_op_statuses: Optional[List[str]],
filter_validator_version_prefix: Optional[str],
include_all_datasets: bool = False,
) -> List[tuple]:
"""
Query datasets that need a (re)validation.
Expand All @@ -264,28 +308,66 @@ def _get_datasets_for_validation(
- Have a report from a different (older) validator version, OR
- force_update is True

By default only each feed's latest dataset is considered. When
include_all_datasets is True, every dataset of the feed is considered — use
this for backfilling a historical window (e.g. re-validating all datasets
published in a date range so the six-month reliability window is complete).

When filter_validator_version_prefix is set (e.g. "8."), the exact-version
comparison is replaced by "the dataset does not already have any validation
report whose validator_version starts with the prefix". This is what lets us
re-validate datasets missing an 8.* report regardless of the exact patch
version returned by the (prod) validator endpoint.

filter_after_in_days restricts to datasets downloaded within the last N days.
When None, all datasets are included regardless of age.
filter_downloaded_after / filter_downloaded_before restrict to a fixed date
window (inclusive lower bound, exclusive upper bound) on downloaded_at.
When all age filters are None, all datasets are included regardless of age.
filter_op_statuses filters by Feed.operational_status (e.g. ["published"]).
"""
query = (
db_session.query(Gtfsfeed.stable_id, Gtfsdataset.stable_id)
.select_from(Gtfsfeed)
.join(Gtfsdataset, Gtfsfeed.latest_dataset_id == Gtfsdataset.id)
.outerjoin(Validationreport, Gtfsdataset.validation_reports)
.filter(
query = db_session.query(Gtfsfeed.stable_id, Gtfsdataset.stable_id).select_from(
Gtfsfeed
)
if include_all_datasets:
query = query.join(Gtfsdataset, Gtfsdataset.feed_id == Gtfsfeed.id)
else:
query = query.join(Gtfsdataset, Gtfsfeed.latest_dataset_id == Gtfsdataset.id)

if filter_validator_version_prefix:
# Include datasets that do NOT already have a report matching the version
# prefix (using a correlated EXISTS over the dataset's reports).
has_matching_report = Gtfsdataset.validation_reports.any(
Validationreport.validator_version.like(
f"{filter_validator_version_prefix}%"
)
)
query = query.filter(or_(~has_matching_report, force_update))
else:
query = query.outerjoin(
Validationreport, Gtfsdataset.validation_reports
).filter(
or_(
Validationreport.id.is_(None),
Validationreport.validator_version != validator_version,
force_update,
)
)
.distinct(Gtfsfeed.stable_id, Gtfsdataset.stable_id)
.order_by(Gtfsdataset.stable_id, Gtfsfeed.stable_id)

query = query.distinct(Gtfsfeed.stable_id, Gtfsdataset.stable_id).order_by(
Gtfsdataset.stable_id, Gtfsfeed.stable_id
)

if filter_after_in_days is not None:
filter_after = datetime.today() - timedelta(days=filter_after_in_days)
query = query.filter(Gtfsdataset.downloaded_at >= filter_after)
if filter_downloaded_after is not None:
query = query.filter(
Gtfsdataset.downloaded_at >= datetime.fromisoformat(filter_downloaded_after)
)
if filter_downloaded_before is not None:
query = query.filter(
Gtfsdataset.downloaded_at < datetime.fromisoformat(filter_downloaded_before)
)
if filter_statuses:
query = query.filter(Gtfsfeed.status.in_(filter_statuses))
if filter_op_statuses:
Expand Down Expand Up @@ -333,8 +415,10 @@ def get_parameters(payload):
Args:
payload (dict): Task payload dict.
Returns:
Tuple of (dry_run, filter_after_in_days, filter_statuses, filter_op_statuses,
prod_env, validator_endpoint, bypass_db_update, force_update, limit,
Tuple of (dry_run, filter_after_in_days, filter_downloaded_after,
filter_downloaded_before, filter_statuses, filter_op_statuses,
filter_validator_version_prefix, include_all_datasets, prod_env,
validator_endpoint, bypass_db_update, force_update, limit,
reports_bucket_name)
"""
prod_env = os.getenv("ENVIRONMENT", "").lower() == "prod"
Expand All @@ -351,10 +435,24 @@ def get_parameters(payload):
else int(filter_after_in_days)
)

filter_downloaded_after = payload.get("filter_downloaded_after", None)
filter_downloaded_before = payload.get("filter_downloaded_before", None)

filter_statuses = payload.get("filter_statuses", None)

filter_op_statuses = payload.get("filter_op_statuses", None)

filter_validator_version_prefix = payload.get(
"filter_validator_version_prefix", None
)

include_all_datasets = payload.get("include_all_datasets", False)
include_all_datasets = (
include_all_datasets
if isinstance(include_all_datasets, bool)
else str(include_all_datasets).lower() == "true"
)

validator_endpoint = payload.get("validator_endpoint", default_endpoint)

bypass_db_update = payload.get("bypass_db_update", False)
Expand All @@ -380,8 +478,12 @@ def get_parameters(payload):
return (
dry_run,
filter_after_in_days,
filter_downloaded_after,
filter_downloaded_before,
filter_statuses,
filter_op_statuses,
filter_validator_version_prefix,
include_all_datasets,
prod_env,
validator_endpoint,
bypass_db_update,
Expand Down
Loading
Loading