Tag ingested events with their tenant region
Python · Python · intermediate · modification
Adds tenant-region tagging to ingested events so downstream dashboards can break metrics down by region. Tagging is applied in one place (a new build_event helper) and gated behind the existing `region_routing` config key — it's already plumbed through every environment and has sat dormant since the old fan-out path was retired, so ops can enable the feature without a config-schema change. With the flag off, events are byte-for-byte identical to today.
This pipeline's process() runs on every batch the ingestion worker pulls; sink.write_event ships each event to the analytics store, and duplicate writes double-count in every downstream metric and dashboard.
Requirements
- Add region tagging: with the new feature enabled, every processed event gains a `region` field derived from its tenant via TENANT_REGION (defaulting to `"us"` for an unmapped tenant). With it disabled, events are byte-for-byte identical to today (no `region` field).
- Gate the feature behind a config flag that is a no-op everywhere else: enabling it must change nothing in the pipeline except adding the `region` field to each written event, and every event must still be written exactly once.
- Do not delete or modify the retired legacy fan-out block in this PR — its removal is tracked separately (DATA-311). Its `region_routing` flag is retired and must remain off in every environment.
- Input contract: raw events are `{ id, tenant, payload }`; `sink.write_event(region, event)` writes exactly one copy of `event` (a `region` target of `None` means "let the sink shard it"); `config` is a plain dict from the deploy config (absent keys mean off).
Files touched
- app/ingest/dispatch.py
--- app/ingest/dispatch.py
REGIONS = ("us", "eu", "apac")
+
+# Tenants are provisioned into exactly one home region.
+TENANT_REGION = {
+ "acme": "us",
+ "globex": "eu",
+ "initech": "apac",
+}
# Feature flags come from the deploy config; absent keys read as off.
return bool(config.get(name, False))
+
+
+def _region_for(raw):
+ return TENANT_REGION.get(raw["tenant"], "us")
-def write_event(raw, sink):
- sink.write_event(None, {"id": raw["id"], "tenant": raw["tenant"], "payload": raw["payload"]})
+def build_event(raw, config):
+ event = {"id": raw["id"], "tenant": raw["tenant"], "payload": raw["payload"]}
+ if _flag(config, "region_routing"):
+ event["region"] = _region_for(raw)
+ return event
+