Skip to content

Process new month #1108

Process new month

Process new month #1108

Workflow file for this run

name: Process new month
on:
workflow_dispatch:
inputs:
args:
description: "Args to pass to `ctbk` (default: auto-detect month and run `update`)"
required: false
quiet:
description: "Skip Slack pings (validation/test dispatches — watch via the GHA UI instead)"
type: boolean
default: false
workflow_run:
workflows: ["Sync tripdata bucket"]
types: [completed]
branches: [main]
# Mid-month rehearsal: re-run the LAST processed month through the
# full pipeline (all stages idempotent / no-op) so pipeline rot on
# HEAD is caught ~2 weeks before the next tripdata drop, instead of
# failing the real ingest (as happened repeatedly pre-2026-08). Also
# doubles as the monthly `station-luc.json` refresh.
schedule:
- cron: '0 12 20 * *'
permissions:
contents: write
env:
AWS_ACCESS_KEY_ID: ${{ secrets.AWS_ACCESS_KEY_ID }}
AWS_SECRET_ACCESS_KEY: ${{ secrets.AWS_SECRET_ACCESS_KEY }}
AWS_DEFAULT_REGION: us-east-1
# R2 creds for the rides-pyramid builds in `ctbk update` (`rides-v1-build`
# writes shards to R2; `r2_client()` prefers these over the local `cf` profile)
CLOUDFLARE_ACCOUNT_ID: ${{ secrets.CLOUDFLARE_ACCOUNT_ID }}
R2_ACCESS_KEY_ID: ${{ secrets.R2_ACCESS_KEY_ID }}
R2_SECRET_ACCESS_KEY: ${{ secrets.R2_SECRET_ACCESS_KEY }}
jobs:
process:
runs-on: ubuntu-latest
if: >-
github.event_name == 'workflow_dispatch' ||
github.event_name == 'schedule' ||
(github.event_name == 'workflow_run' && github.event.workflow_run.conclusion == 'success')
steps:
- uses: actions/checkout@v5
with:
fetch-depth: 0
ssh-key: ${{ secrets.WWW_DEPLOY_KEY }}
- name: Detect args
id: args
run: |
if [ "${{ github.event_name }}" = "schedule" ]; then
# Rehearsal: latest already-processed month, full pipeline.
ym=$(ls s3/ctbk/normalized/*.parquet.dvc | grep -oE '[0-9]{6}' | sort | tail -1)
echo "ym=$ym" >> "$GITHUB_OUTPUT"
echo "rehearsal=true" >> "$GITHUB_OUTPUT"
echo "args=update -S -ccda $ym" >> "$GITHUB_OUTPUT"
elif [ -n "${{ inputs.args }}" ]; then
echo "args=${{ inputs.args }}" >> "$GITHUB_OUTPUT"
# Surface the month in Slack notifies too (else "month ?")
ym=$(echo "${{ inputs.args }}" | grep -oE '20[0-9]{4}' | head -1)
[ -n "$ym" ] && echo "ym=$ym" >> "$GITHUB_OUTPUT"
else
# Auto-detect month from most recent tripdata import commit
ym=$(git log -1 --format=%s | sed -n 's/.*\(20[0-9]\{4\}\)-citibike.*/\1/p' | head -1)
if [ -z "$ym" ]; then
echo "No new month detected in recent commits; nothing to do"
exit 0
fi
if [ -f "s3/ctbk/normalized/${ym}.dvc" ]; then
echo "Month $ym already processed; nothing to do"
exit 0
fi
echo "ym=$ym" >> "$GITHUB_OUTPUT"
# NYC + JC zips publish independently (e.g. 202607: JC landed
# 2026-08-06, NYC later); `ctbk norm create` needs both. An
# incomplete month is an expected transient — annotate + notify
# and defer (the daily sync re-triggers when the rest arrives)
# instead of failing on a FileNotFoundError deep in `norm`.
missing=""
ls "s3/tripdata/${ym}-"*.zip.dvc >/dev/null 2>&1 || missing="NYC"
ls "s3/tripdata/JC-${ym}"*.zip.dvc >/dev/null 2>&1 || missing="${missing:+$missing,}JC"
if [ -n "$missing" ]; then
echo "missing=$missing" >> "$GITHUB_OUTPUT"
echo "::warning title=Tripdata month $ym incomplete::$missing zip(s) not yet published for $ym. Deferring processing; the daily sync will re-trigger this workflow when the missing region(s) arrive."
exit 0
fi
echo "args=update -S -ccda $ym" >> "$GITHUB_OUTPUT"
fi
- uses: astral-sh/setup-uv@v7
if: steps.args.outputs.args || steps.args.outputs.missing
with:
python-version: '3.12'
enable-cache: true
- name: Notify Slack — month incomplete
if: steps.args.outputs.missing
env:
SLACK_BOT_TOKEN: ${{ secrets.SLACK_BOT_TOKEN }}
YM: ${{ steps.args.outputs.ym }}
MISSING: ${{ steps.args.outputs.missing }}
RUN_URL: ${{ github.server_url }}/${{ github.repository }}/actions/runs/${{ github.run_id }}
run: |
uv run --no-project --with 'thrds==0.4.0' --python 3.12 python - <<'PY'
import os
from thrds import SlackClient
ym, missing, run_url = os.environ['YM'], os.environ['MISSING'], os.environ['RUN_URL']
SlackClient(
token=os.environ['SLACK_BOT_TOKEN'],
channel='C0B5MKF28NP',
username='ctbk-bot',
icon_emoji=':bike:',
).post(
f":hourglass_flowing_sand: *tripdata {ym} incomplete* — {missing} zip(s) not yet published; "
f"processing deferred until the month completes (daily sync will re-trigger). <{run_url}|run>"
)
PY
- name: Notify Slack — processing started
id: slack_start
if: steps.args.outputs.args && inputs.quiet != true
env:
SLACK_BOT_TOKEN: ${{ secrets.SLACK_BOT_TOKEN }}
ARGS: ${{ steps.args.outputs.args }}
YM: ${{ steps.args.outputs.ym }}
REHEARSAL: ${{ steps.args.outputs.rehearsal }}
RUN_URL: ${{ github.server_url }}/${{ github.repository }}/actions/runs/${{ github.run_id }}
run: |
uv run --no-project --with 'thrds==0.4.0' --python 3.12 python - <<'PY'
import os
from thrds import SlackClient
ym = os.environ.get('YM') or '?'
args, run_url = os.environ['ARGS'], os.environ['RUN_URL']
verb = 'Rehearsing' if os.environ.get('REHEARSAL') else 'Processing new'
msg = SlackClient(
token=os.environ['SLACK_BOT_TOKEN'],
channel='C0B5MKF28NP',
username='ctbk-bot',
icon_emoji=':bike:',
).post(f":inbox_tray: *{verb} tripdata month {ym}* (`ctbk {args}`)… <{run_url}|run>")
with open(os.environ['GITHUB_OUTPUT'], 'a') as f:
f.write(f'ts={msg.id}\n')
PY
- name: Sync deps + activate venv
if: steps.args.outputs.args
run: |
uv sync --frozen
echo "$GITHUB_WORKSPACE/.venv/bin" >> "$GITHUB_PATH"
- name: Configure Git
if: steps.args.outputs.args
run: |
git config --global user.name 'GitHub Actions'
git config --global user.email 'ryan-williams@users.noreply.github.com'
- name: Pull pipeline inputs
if: steps.args.outputs.args
run: |
# Pull any tripdata zips needed by the command
ym=$(echo "${{ steps.args.outputs.args }}" | grep -oE '20[0-9]{4}' | head -1)
if [ -n "$ym" ]; then
dvc pull s3/tripdata/*${ym}*.dvc || true
# `cons create` depends on EVERY prior normalized dir (any source
# month can hold rides ending in M; `consolidated.dep_artifacts`),
# reading each dir's `.dir` manifest from the local cache. Hydrate
# just those manifests (~40 KB total) — not the ~190 GiB of data —
# via `dvx pull --meta-only` (dvx `93e790a6c`). Without this the
# guard raises "N normalized directory manifests missing".
# Skip the current month ($ym): on a resume its `.dvc` is already
# committed (a prior run git-pushed before its R2 push) but its
# manifest isn't on R2 until THIS run's push, so pulling it would
# 404. `cons` exempts the current month from the guard anyway, and
# `norm create` rebuilds it locally (deterministic → same hash).
git ls-files 's3/ctbk/normalized/*.dvc' \
| grep -v "/${ym}\.dvc$" \
| xargs dvx pull -m
# Whole-history / prev-month inputs for `ctbk update` steps that
# would otherwise silently rebuild full-history artifacts from a
# partial local set (2026-08-14 truncation incident, #183):
# - prev month's consolidated parquet: rides base-tier refold
# (`rides-v1-build -O` re-buckets prev month + its spillover)
# - ymrgtb{s,e} aggregates (~20MB total): `station-trips-json`
# whole-artifact rebuild
# - station-observations: lat/lng fallback in rides base build
prev=$(date -u -d "${ym}01 -1 month" +%Y%m)
dvc pull s3/ctbk/normalized/${prev}.parquet.dvc
# Exclude the current month here too: on a resume its committed
# aggregate `.dvc`s reference blobs not yet on R2 (same torn-run
# cause as above); `agg` rebuilds them this run.
git ls-files 's3/ctbk/aggregated/ymrgtbs_cd_*.parquet.dvc' 's3/ctbk/aggregated/ymrgtbe_cd_*.parquet.dvc' \
| grep -v "_${ym}\.parquet\.dvc$" \
| xargs dvc pull
dvc pull s3/ctbk/stations/station-observations.parquet.dvc
# (station-history.parquet — the luc refresh's historical
# canonicals — is git-tracked; no pull needed)
fi
# ─── Station registry refresh (spec cadence item 4) ─────────────
# Regen `station-luc.json` from the GBFS snapshot union BEFORE any
# pyramid build, so the month's rides-v3 rebuild (in `ctbk update`)
# and rides-v5 Batch fills both attribute rides at NEW stations to
# their identity keys instead of coordinate fallback. Writes the
# www asset (committed by `update`'s git step) + uploads to R2 (the
# Batch factory reads it live). New stations are additive-safe;
# MOVED LUCs leave historical rows keyed under the old cell until
# their range is invalidated — nonzero churn pings Slack for manual
# review rather than blocking the month. `vocab check` then guards
# against a new station landing outside the frozen vocab's
# coverage (its rides would be invisible to bbox queries; fixing
# that is a deliberate vocab extension + re-key).
- name: Refresh station-luc + vocab coverage check
id: luc
if: steps.args.outputs.args
run: |
ctbk station-luc-build
ctbk gbfs vocab check
- name: Notify Slack — LUC churn needs review
if: steps.args.outputs.args && inputs.quiet != true && steps.luc.outputs.luc_moved != '0' && steps.luc.outputs.luc_moved != ''
env:
SLACK_BOT_TOKEN: ${{ secrets.SLACK_BOT_TOKEN }}
MOVED: ${{ steps.luc.outputs.luc_moved }}
MOVED_LIST: ${{ steps.luc.outputs.luc_moved_list }}
RUN_URL: ${{ github.server_url }}/${{ github.repository }}/actions/runs/${{ github.run_id }}
run: |
uv run --no-project --with 'thrds==0.4.0' --python 3.12 python - <<'PY'
import os
from thrds import SlackClient
moved, lst = os.environ['MOVED'], os.environ['MOVED_LIST']
SlackClient(
token=os.environ['SLACK_BOT_TOKEN'],
channel='C0B5MKF28NP',
username='ctbk-bot',
icon_emoji=':bike:',
).post(
f":warning: *station LUC churn* — {moved} station(s) moved ({lst}); historical pyramid rows "
f"stay keyed under the old cells until their ranges are invalidated (manual review). "
f"<{os.environ['RUN_URL']}|run>"
)
PY
- uses: actions/setup-node@v5
if: steps.args.outputs.args
with:
node-version: 20
- name: Install pnpm
if: steps.args.outputs.args
# Pin to 9 to match www.yml's `pnpm/action-setup` version. Bare
# `npm i -g pnpm` pulls latest (v11), which chases breaking migrations
# the v9-authored lockfile/config predate (unsafe `..` importer paths,
# `pnpm.onlyBuiltDependencies` moved out of package.json → IGNORED_BUILDS).
run: npm install -g pnpm@9
- name: Install Node dependencies
if: steps.args.outputs.args
run: cd www && pnpm install --frozen-lockfile
- name: Run ctbk
if: steps.args.outputs.args
run: ctbk ${{ steps.args.outputs.args }}
# ─── R2 dual-push (specs/s3-to-r2-migration.md) ────────────────
# The FE reads the DVX cache from R2 (data.ctbk.dev), so each
# month's new content-addressed blobs must land on R2 too. `ctbk
# update` above pushed them to the `s3` remote; mirror to `r2`.
# R2 creds differ from the AWS creds the `s3` remote uses, so
# inject them per-remote into config.local at runtime.
- name: Push DVX cache to R2
if: steps.args.outputs.args
run: |
dvx remote modify --local r2 access_key_id "$R2_ACCESS_KEY_ID"
dvx remote modify --local r2 secret_access_key "$R2_SECRET_ACCESS_KEY"
dvx push -r r2
# ─── rides-v5 monthly cadence (specs/rides-v5.md; #185) ─────────
# Prod's Home chart reads rides-v5, so the month isn't "done" until
# the v5 pyramid has it. `rides-v5-extend` mirrors the consolidated
# parquet to its public plain key (the `dvc push` above uploaded the
# CA blob), journals the prev-month spillback invalidation, submits
# + watches the Batch fills for both anchors (scoped `batch:SubmitJob`
# on the CI IAM user), and backfills the RG manifest. Idempotent on
# re-runs. Runs BEFORE the www push so a fill failure blocks the
# deploy and pings `:x:` top-level. Station-map/vocab regen for
# brand-new stations is NOT covered (spec cadence item 4 — manual).
- name: Install gbfs/api deps (wrangler, for relic sweep)
if: steps.args.outputs.args && steps.args.outputs.ym
working-directory: gbfs/api
run: pnpm install --frozen-lockfile
- name: Extend rides-v5
if: steps.args.outputs.args && steps.args.outputs.ym
env:
CTBK_REGISTRY_SECRET: ${{ secrets.CTBK_REGISTRY_SECRET }}
CLOUDFLARE_API_TOKEN: ${{ secrets.CLOUDFLARE_API_TOKEN }}
run: ctbk gbfs rides-v5-extend ${{ steps.args.outputs.ym }}
- name: Deploy to www branch
if: steps.args.outputs.args
run: git push origin HEAD:www
# (No rides-v3 step: the rollback pyramid is frozen at its last
# build — a v5 rollback missing only the newest month is an
# acceptable degraded mode, and the v3 'all'-shard merge-patches
# were the runner's OOM-flakiest phase. `ctbk rides-v3-extend`
# remains as a manual tool until rides<5 is GC'd wholesale.)
- name: Notify Slack — outcome
if: always() && steps.slack_start.outputs.ts
env:
SLACK_BOT_TOKEN: ${{ secrets.SLACK_BOT_TOKEN }}
YM: ${{ steps.args.outputs.ym }}
STATUS: ${{ job.status }}
REHEARSAL: ${{ steps.args.outputs.rehearsal }}
THREAD_TS: ${{ steps.slack_start.outputs.ts }}
RUN_URL: ${{ github.server_url }}/${{ github.repository }}/actions/runs/${{ github.run_id }}
run: |
uv run --no-project --with 'thrds==0.4.0' --python 3.12 python - <<'PY'
import os
from thrds import SlackClient
ym = os.environ.get('YM') or '?'
ok = os.environ['STATUS'] == 'success'
run_url = os.environ['RUN_URL']
rehearsal = bool(os.environ.get('REHEARSAL'))
if ok:
text = (
f":white_check_mark: *tripdata {ym} rehearsal passed* — pipeline healthy for the next drop"
if rehearsal else
f":white_check_mark: *tripdata {ym} processed* — pushed to `www` (site deploy follows)"
)
else:
text = (
f":x: *tripdata {ym} rehearsal FAILED* — pipeline rot on `main`; fix before the next drop — <{run_url}|logs>"
if rehearsal else
f":x: *tripdata {ym} processing failed* — <{run_url}|logs>"
)
# Failures post top-level: a threaded :x: is invisible at channel
# level (the 202607 failures sat unnoticed in threads for a day).
SlackClient(
token=os.environ['SLACK_BOT_TOKEN'],
channel='C0B5MKF28NP',
username='ctbk-bot',
icon_emoji=':bike:',
).post(text, thread_id=os.environ['THREAD_TS'] if ok else None)
PY