From e4e32bc71897c51397927aee0a2aa739504ad380 Mon Sep 17 00:00:00 2001 From: Makisuo Date: Wed, 26 Aug 2026 01:10:02 +0200 Subject: [PATCH 1/5] feat(ingest): mirror Tinybird writes into a second workspace MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The Tinybird workspace is on GCP Frankfurt while the ingest fleet runs in AWS us-east-1, so every exported byte crosses cloud and continent at the public-internet egress rate instead of the same-region one. Moving the workspace to us-east-1 in one step would strand every read behind an empty warehouse, so the gateway needs to write both for a while: the new workspace's materialized views only build rollups from rows inserted into it, and a month of live dual-emit is what makes it readable at cutover. Adds `ExportDestination::TinybirdMirror` as a third export lane rather than a fan-out inside the export call. The lane machinery already exists for exactly this reason — its own WAL file, cursor, bounded channel and worker, so a stalled destination can only back-pressure itself. Appending to `ExportDestination::ALL` keeps the existing lane ordinals, and WAL files are addressed by destination name, so a rolling deploy does not renumber lanes under frames already committed. Wire tag 3 is reserved permanently: a mirror-written WAL file has to decode under a rolled-back binary. The mirror is best-effort, and that is enforced structurally rather than by convention. `commit_frames` splits into a fail-closed primary commit and a mirror commit that cannot return an error — a full lane, an exhausted org byte budget or a WAL failure is metered and dropped. The primary commits first, so it always wins the race for the last reservable bytes. A ClickHouse-routed org is never mirrored: its rows never reached the workspace being migrated and must not start reaching its replacement. Mirror credentials carry a smaller retry budget and a per-attempt timeout the primary deliberately lacks, so a sick workspace releases its lane worker in seconds instead of holding a batch through the full budget. A sample-percent knob ramps by org hash, so an org is consistently in or out and no org's stream is ever half-mirrored — a backfill boundary can only be described if that holds. Export metrics gain a `destination` label. Comparing per-datasource row counts between the two workspaces is the acceptance test for the whole window, and `ingest_tinybird_mirror_dropped_total` is the only signal that the mirror has a gap, since a mirror shed is invisible to clients and to every other counter. Feature is off unless both halves of the mirror config are set; half a config fails at startup rather than 401-ing a lane forever. Also adds a us-east-1 leg to the Tinybird CD matrix — both production workspaces must receive every schema deploy for the duration of the window, since a month of drift between them is the quiet way this migration fails. --- .env.example | 10 + .github/workflows/tinybird-cd.yml | 7 + CONTRIBUTING.md | 1 + apps/ingest/alchemy.run.ts | 38 +- apps/ingest/benches/ingest_bench.rs | 1 + apps/ingest/src/main.rs | 52 +++ apps/ingest/src/metrics.rs | 70 +++- apps/ingest/src/telemetry.rs | 600 +++++++++++++++++++++++++++- 8 files changed, 752 insertions(+), 27 deletions(-) diff --git a/.env.example b/.env.example index ff496b7c5..1c140a8ea 100644 --- a/.env.example +++ b/.env.example @@ -8,6 +8,16 @@ TINYBIRD_TOKEN=your-tinybird-token # Optional Tinybird-side ceiling shared by raw SQL from the API and alerting, # with an independent rate-limit bucket for each org. Unset means no JWT RPS limit. # TINYBIRD_RAW_SQL_JWT_RPS_LIMIT=100 +# A second Tinybird workspace the ingest gateway mirrors writes into, for the +# duration of a workspace migration. Best-effort: the mirror has its own WAL +# lane and can never slow or fail ingest. Set BOTH or neither — the gateway +# refuses to start on half a config. Reads always stay on TINYBIRD_HOST. +# TINYBIRD_MIRROR_HOST=https://api.us-east.aws.tinybird.co +# TINYBIRD_MIRROR_TOKEN=your-mirror-workspace-token +# Ramp knob: percentage of orgs to mirror, by org hash (default 100). Start at 1. +# INGEST_TINYBIRD_MIRROR_SAMPLE_PERCENT=1 +# INGEST_TINYBIRD_MIRROR_MAX_ATTEMPTS=5 +# INGEST_TINYBIRD_MIRROR_TIMEOUT_MS=3000 # ClickHouse CLICKHOUSE_URL=http://localhost:9000 diff --git a/.github/workflows/tinybird-cd.yml b/.github/workflows/tinybird-cd.yml index 1994f2cec..b5326a2ae 100644 --- a/.github/workflows/tinybird-cd.yml +++ b/.github/workflows/tinybird-cd.yml @@ -36,6 +36,13 @@ jobs: environment: tinybird-cd - label: staging environment: tinybird-cd-stg + # The us-east-1 workspace being migrated to. Both production + # legs must receive every schema deploy for the whole + # dual-emit window — a month of drift between the two + # workspaces is the quiet way that migration fails. No-ops + # until the environment's secrets exist. + - label: production-us + environment: tinybird-cd-us environment: ${{ matrix.target.environment }} steps: - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6 diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 94af6f3df..508ba2d7f 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -238,6 +238,7 @@ port and keep `VITE_*` / `MAPLE_INGEST_PUBLIC_URL` consistent. | `MAPLE_ORG_ID_OVERRIDE` | with static | Must match `MAPLE_DEFAULT_ORG_ID` | | `MAPLE_PG_URL` | postgres store | `postgres://maple:maple@localhost:5499/maple` if not using static store | | `TINYBIRD_HOST` / `TINYBIRD_TOKEN` | tinybird mode | When `INGEST_WRITE_MODE=tinybird` or `dual` | +| `TINYBIRD_MIRROR_HOST` / `_TOKEN` | migration only | Mirrors writes into a second workspace; best-effort, set both or neither | | `INGEST_PORT` | optional | Default from port / env | | `INGEST_REQUIRE_TLS` | optional | `false` locally | diff --git a/apps/ingest/alchemy.run.ts b/apps/ingest/alchemy.run.ts index d2bf487c6..ad3ac3ee4 100644 --- a/apps/ingest/alchemy.run.ts +++ b/apps/ingest/alchemy.run.ts @@ -62,14 +62,21 @@ const INGEST_PORT = 3474 * default (`INGEST_QUEUE_MAX_BYTES` = 20 GiB, `apps/ingest/src/main.rs`) would * sit exactly on the line. 8 GiB still buys hours of buffering at current * volume; raising it means paying for ephemeral storage beyond the free tier. + * + * Note the per-lane budget is `INGEST_QUEUE_MAX_BYTES / (WAL_SHARDS * lanes)`, + * so enabling the Tinybird mirror (a third lane) cuts every lane's share by a + * third at the same moment Tinybird-bound traffic doubles. Raise this — within + * the ephemeral ceiling — for the duration of a mirrored window. */ const WAL_MAX_BYTES = 8 * 1024 * 1024 * 1024 /** * Pinned rather than derived. The gateway defaults to `num_cpus * 2`, which * makes on-disk layout and fd count a function of task size — so a cpu bump, or - * a move to another capacity provider, would silently reshape the WAL. Two lanes per shard - * (Tinybird + ClickHouse) means this is 8 open WAL files. + * a move to another capacity provider, would silently reshape the WAL. Three + * lanes per shard (Tinybird + ClickHouse + Tinybird mirror) means this is 12 + * open WAL files; the mirror lane is present but idle unless + * `TINYBIRD_MIRROR_HOST`/`TINYBIRD_MIRROR_TOKEN` are set. */ const WAL_SHARDS = 4 @@ -229,6 +236,22 @@ export const createMapleIngest = ({ stage, domains, region, replayBlobs }: Creat const autumnKey = process.env.AUTUMN_SECRET_KEY?.trim() const autumnSecret = autumnKey ? yield* secret("autumn-secret-key", autumnKey) : undefined + // Second Tinybird workspace to mirror writes into during a workspace + // migration. Both halves must be set together; the gateway rejects a + // half-configured mirror at startup rather than 401-ing its lane forever. + // + // Resolved through `optionalPlain`, not `process.env`: alchemy reads + // `--env-file`/`.env` through its own ConfigProvider and never copies + // those values into `process.env`, so a bare read would see the var in + // CI and miss it locally. The host is plain (not a secret); the token is + // resolved here only to mint the Secrets Manager entry below, exactly as + // TINYBIRD_TOKEN is, and never reaches `env`. + const tinybirdMirrorHostEntry = yield* optionalPlain("TINYBIRD_MIRROR_HOST") + const tinybirdMirrorTokenValue = (yield* optionalPlain("TINYBIRD_MIRROR_TOKEN")).TINYBIRD_MIRROR_TOKEN + const tinybirdMirrorToken = tinybirdMirrorTokenValue + ? yield* secret("tinybird-mirror-token", tinybirdMirrorTokenValue) + : undefined + // Both halves are stack-minted (`createReplayBlobStore`), so there is no // half-set config left to guard against. The access key id is not secret, // but it only exists once the token does and `env` takes plan-time strings @@ -431,6 +454,9 @@ export const createMapleIngest = ({ stage, domains, region, replayBlobs }: Creat MAPLE_INGEST_KEY_ENCRYPTION_KEY: keyEncryptionKey.secretArn, MAPLE_INGEST_KEY_LOOKUP_HMAC_KEY: keyLookupHmacKey.secretArn, ...(autumnSecret ? { AUTUMN_SECRET_KEY: autumnSecret.secretArn } : undefined), + ...(tinybirdMirrorToken + ? { TINYBIRD_MIRROR_TOKEN: tinybirdMirrorToken.secretArn } + : undefined), ...(replayR2Secret && replayR2AccessKeyId ? { INGEST_REPLAY_R2_SECRET_ACCESS_KEY: replayR2Secret.secretArn, @@ -496,6 +522,14 @@ export const createMapleIngest = ({ stage, domains, region, replayBlobs }: Creat ...(yield* optionalPlain("INGEST_MAX_REQUEST_BODY_BYTES")), ...(yield* optionalPlain("INGEST_EXPORT_MAX_ATTEMPTS")), ...(yield* optionalPlain("INGEST_TINYBIRD_CONCURRENCY_PER_SHARD")), + // Tinybird mirror. The host is plain (it is not a secret); the token + // goes through Secrets Manager above. Ramp with SAMPLE_PERCENT: start + // at 1, confirm rows land and `ingest_tinybird_mirror_dropped_total` + // stays flat, then climb to 100. + ...tinybirdMirrorHostEntry, + ...(yield* optionalPlain("INGEST_TINYBIRD_MIRROR_SAMPLE_PERCENT")), + ...(yield* optionalPlain("INGEST_TINYBIRD_MIRROR_MAX_ATTEMPTS")), + ...(yield* optionalPlain("INGEST_TINYBIRD_MIRROR_TIMEOUT_MS")), ...(yield* optionalPlain("INGEST_REPLAY_MAX_SESSION_BYTES")), // The org Maple's own telemetry is filed under. Required here and in // the gateway (`AppConfig::from_env`), with no fallback on either diff --git a/apps/ingest/benches/ingest_bench.rs b/apps/ingest/benches/ingest_bench.rs index 66e523ce4..c9164cce1 100644 --- a/apps/ingest/benches/ingest_bench.rs +++ b/apps/ingest/benches/ingest_bench.rs @@ -105,6 +105,7 @@ impl BenchFixture { TinybirdConfig { endpoint: format!("http://{addr}"), token: "bench-token".to_owned(), + mirror: None, queue_dir: queue_dir.clone(), // Effectively uncapped: this benchmark measures accept latency // (encode + WAL append + ack), not back-pressure. A single org diff --git a/apps/ingest/src/main.rs b/apps/ingest/src/main.rs index 7879f9b6c..67b9de8fa 100644 --- a/apps/ingest/src/main.rs +++ b/apps/ingest/src/main.rs @@ -56,6 +56,7 @@ use maple_ingest::telemetry::{ AttributeMappingRule, ClickHouseBreakerConfig, ClickHouseTarget, ClickHouseTargetProvider, DatasourceNames, ExportDestination, HttpClient, MappingOperation, MappingSourceContext, PipelineError, SamplingPolicy, TelemetryPipeline, TelemetrySignal, TinybirdConfig, + TinybirdMirrorConfig, }; use maple_ingest::usage_metrics::{billable_gb, usage_cardinality_view, UsageMetrics}; use moka::future::Cache; @@ -262,6 +263,51 @@ impl AppConfig { 10_000, )?; + // A second Tinybird workspace to mirror writes into during a workspace + // migration. Present only when both halves are set, so the feature is + // off everywhere that does not opt in; setting exactly one half is a + // deploy mistake and `validate_for_pipeline` rejects it rather than + // letting the lane 401 in silence. + let mirror_host = std::env::var("TINYBIRD_MIRROR_HOST") + .unwrap_or_default() + .trim() + .trim_end_matches('/') + .to_owned(); + let mirror_token = std::env::var("TINYBIRD_MIRROR_TOKEN") + .unwrap_or_default() + .trim() + .to_owned(); + let mirror = if mirror_host.is_empty() && mirror_token.is_empty() { + None + } else { + Some(TinybirdMirrorConfig { + endpoint: mirror_host, + token: mirror_token, + // Deliberately far below the primary's 20: the mirror holds a + // lane worker for the whole budget, and losing a mirror batch is + // cheaper than stalling the lane behind a sick workspace. + max_attempts: parse_u32( + "INGEST_TINYBIRD_MIRROR_MAX_ATTEMPTS", + std::env::var("INGEST_TINYBIRD_MIRROR_MAX_ATTEMPTS").ok(), + 5, + )?, + export_timeout: Duration::from_millis(parse_u64( + "INGEST_TINYBIRD_MIRROR_TIMEOUT_MS", + std::env::var("INGEST_TINYBIRD_MIRROR_TIMEOUT_MS").ok(), + 3_000, + )?), + // The ramp knob: start at 1, confirm rows land and nothing is + // shed, then climb. Sampling is by org hash, so an org is + // consistently in or out and its stream is never half-mirrored. + sample_percent: u8::try_from(parse_u32( + "INGEST_TINYBIRD_MIRROR_SAMPLE_PERCENT", + std::env::var("INGEST_TINYBIRD_MIRROR_SAMPLE_PERCENT").ok(), + 100, + )?) + .map_err(|_| "INGEST_TINYBIRD_MIRROR_SAMPLE_PERCENT must be 0..=100".to_owned())?, + }) + }; + let tinybird = TinybirdConfig { endpoint: std::env::var("TINYBIRD_HOST") .unwrap_or_default() @@ -272,6 +318,7 @@ impl AppConfig { .unwrap_or_default() .trim() .to_owned(), + mirror, queue_dir: PathBuf::from( std::env::var("INGEST_QUEUE_DIR") .unwrap_or_else(|_| "/var/lib/maple-ingest/wal".to_owned()), @@ -353,6 +400,10 @@ impl AppConfig { }; if write_mode.uses_tinybird() { tinybird.validate()?; + } else { + // The mirror is a Tinybird destination, so a half-configured one is + // still a deploy mistake in forward-only mode. + tinybird.validate_for_pipeline(false)?; } let max_request_body_bytes = parse_usize( @@ -7022,6 +7073,7 @@ mod tests { TinybirdConfig { endpoint: String::new(), token: String::new(), + mirror: None, queue_dir, queue_max_bytes: 1024 * 1024, org_queue_max_bytes: 1024 * 1024, diff --git a/apps/ingest/src/metrics.rs b/apps/ingest/src/metrics.rs index 26014025f..e4010089e 100644 --- a/apps/ingest/src/metrics.rs +++ b/apps/ingest/src/metrics.rs @@ -175,6 +175,13 @@ static TINYBIRD_EXPORT_RETRIES_TOTAL: LazyLock> = LazyLock::new(|| .build() }); +static TINYBIRD_MIRROR_DROPPED_TOTAL: LazyLock> = LazyLock::new(|| { + METER + .u64_counter("ingest_tinybird_mirror_dropped_total") + .with_description("Rows dropped before reaching the Tinybird mirror lane") + .build() +}); + static CLICKHOUSE_EXPORT_ROWS_TOTAL: LazyLock> = LazyLock::new(|| { METER .u64_counter("ingest_clickhouse_export_rows_total") @@ -553,14 +560,30 @@ pub fn org_queue_bytes(org_id: &str, bytes: u64) { } /// Latency and exported-byte size of a completed WAL export batch. -pub fn export_batch_completed(shard: usize, signal: &str, duration_secs: f64, exported_bytes: u64) { - EXPORT_BATCH_DURATION_SECONDS - .record(duration_secs, &[KeyValue::new("shard", shard.to_string())]); +/// +/// `destination` matters here: without it a mirror lane's drain is +/// indistinguishable from the primary's in the same shard, which is exactly the +/// comparison a workspace migration needs. +pub fn export_batch_completed( + shard: usize, + destination: &str, + signal: &str, + duration_secs: f64, + exported_bytes: u64, +) { + EXPORT_BATCH_DURATION_SECONDS.record( + duration_secs, + &[ + KeyValue::new("shard", shard.to_string()), + KeyValue::new("destination", destination.to_owned()), + ], + ); WAL_EXPORTED_BYTES.record( exported_bytes, &[ KeyValue::new("signal", signal.to_owned()), KeyValue::new("shard", shard.to_string()), + KeyValue::new("destination", destination.to_owned()), ], ); } @@ -605,22 +628,38 @@ pub fn native_sampled_dropped(signal: &str, count: u64) { } /// A successful Tinybird export: latency and exported row count. -pub fn tinybird_export_succeeded(datasource: &str, duration_secs: f64, rows: u64) { +/// +/// `destination` is `tinybird` or `tinybird_mirror`. Comparing the two row +/// counters per datasource is how a mirrored workspace is proved complete. +pub fn tinybird_export_succeeded( + destination: &str, + datasource: &str, + duration_secs: f64, + rows: u64, +) { TINYBIRD_EXPORT_DURATION_SECONDS.record( duration_secs, &[ + KeyValue::new("destination", destination.to_owned()), KeyValue::new("datasource", datasource.to_owned()), KeyValue::new("status", "2xx"), ], ); - TINYBIRD_EXPORT_ROWS_TOTAL.add(rows, &[KeyValue::new("datasource", datasource.to_owned())]); + TINYBIRD_EXPORT_ROWS_TOTAL.add( + rows, + &[ + KeyValue::new("destination", destination.to_owned()), + KeyValue::new("datasource", datasource.to_owned()), + ], + ); } /// Rows dropped while exporting to Tinybird (`status` is an HTTP code or `retries_exhausted`). -pub fn tinybird_export_dropped(datasource: &str, status: &str, rows: u64) { +pub fn tinybird_export_dropped(destination: &str, datasource: &str, status: &str, rows: u64) { TINYBIRD_EXPORT_DROPPED_TOTAL.add( rows, &[ + KeyValue::new("destination", destination.to_owned()), KeyValue::new("datasource", datasource.to_owned()), KeyValue::new("status", status.to_owned()), ], @@ -628,16 +667,33 @@ pub fn tinybird_export_dropped(datasource: &str, status: &str, rows: u64) { } /// A Tinybird export attempt was retried (`status` is an HTTP code or `transport`). -pub fn tinybird_export_retry(datasource: &str, status: &str) { +pub fn tinybird_export_retry(destination: &str, datasource: &str, status: &str) { TINYBIRD_EXPORT_RETRIES_TOTAL.add( 1, &[ + KeyValue::new("destination", destination.to_owned()), KeyValue::new("datasource", datasource.to_owned()), KeyValue::new("status", status.to_owned()), ], ); } +/// Rows shed on the commit path before they ever reached the mirror lane +/// (`reason` is `lane_full`, `org_quota`, `wal_error`, or `no_target`). +/// +/// The mirror is best-effort, so these drops are invisible to clients and to +/// every other counter. This is the only signal that a mirrored workspace has a +/// gap, and therefore the only thing that tells you a backfill window is needed. +pub fn tinybird_mirror_dropped(datasource: &str, reason: &str, rows: u64) { + TINYBIRD_MIRROR_DROPPED_TOTAL.add( + rows, + &[ + KeyValue::new("datasource", datasource.to_owned()), + KeyValue::new("reason", reason.to_owned()), + ], + ); +} + /// A successful ClickHouse export: latency and exported row count. pub fn clickhouse_export_succeeded(datasource: &str, status: &str, duration_secs: f64, rows: u64) { let attrs = [ diff --git a/apps/ingest/src/telemetry.rs b/apps/ingest/src/telemetry.rs index acf7c8d20..87fb27ecd 100644 --- a/apps/ingest/src/telemetry.rs +++ b/apps/ingest/src/telemetry.rs @@ -57,18 +57,26 @@ const WAL_V3_HEADER_LEN: usize = 23; pub enum ExportDestination { Tinybird, ClickHouse, + /// A second Tinybird workspace, written in addition to `Tinybird` while a + /// workspace migration is in flight (see `TinybirdMirrorConfig`). Strictly + /// best-effort: it has its own lane, so it can never back-pressure the + /// primary, and `commit_mirror_frames` swallows every failure rather than + /// failing a request that the primary already accepted. + TinybirdMirror, } impl ExportDestination { /// Every export destination, in lane order. A frame's lane within its WAL /// shard is `shard * LANES_PER_SHARD + destination.lane_ordinal()`, so the - /// position in this slice is load-bearing — do not reorder. - pub const ALL: [Self; 2] = [Self::Tinybird, Self::ClickHouse]; + /// position in this slice is load-bearing — do not reorder. Appending is + /// safe (existing ordinals keep their value); inserting is not. + pub const ALL: [Self; 3] = [Self::Tinybird, Self::ClickHouse, Self::TinybirdMirror]; pub fn as_str(self) -> &'static str { match self { Self::Tinybird => "tinybird", Self::ClickHouse => "clickhouse", + Self::TinybirdMirror => "tinybird_mirror", } } @@ -77,6 +85,7 @@ impl ExportDestination { match self { Self::Tinybird => 0, Self::ClickHouse => 1, + Self::TinybirdMirror => 2, } } } @@ -424,10 +433,34 @@ impl DatasourceNames { } } +/// A second Tinybird workspace to mirror writes into, for the duration of a +/// workspace migration. +/// +/// Deliberately not modelled as "two equal endpoints": the mirror is the one +/// that is allowed to lose rows. Its retry budget is smaller than the primary's +/// so a sick workspace releases its lane worker in seconds rather than holding a +/// batch through the full budget, and `sample_percent` exists so the mirror can +/// be canaried on a slice of orgs before it carries the whole fleet. +#[derive(Clone, Debug)] +pub struct TinybirdMirrorConfig { + pub endpoint: String, + pub token: String, + /// Retry budget for a mirror batch. Lower than `export_max_attempts`. + pub max_attempts: u32, + /// Per-attempt timeout. The primary has none; the mirror does, because a + /// hanging mirror host must not occupy its lane worker indefinitely. + pub export_timeout: Duration, + /// Percentage of orgs to mirror, by `hash64(org_id) % 100`. The same org is + /// consistently in or out, so a ramp never splits one org's stream. + pub sample_percent: u8, +} + #[derive(Clone, Debug)] pub struct TinybirdConfig { pub endpoint: String, pub token: String, + /// `None` disables mirroring entirely — no mirror lane traffic, no clones. + pub mirror: Option, pub queue_dir: PathBuf, pub queue_max_bytes: u64, pub org_queue_max_bytes: u64, @@ -457,7 +490,62 @@ impl TinybirdConfig { self.validate_for_pipeline(true) } + /// Endpoint, token, retry budget and per-attempt timeout for one Tinybird + /// destination. `TinybirdMirror` resolves to the mirror workspace; a missing + /// mirror config means the lane is idle, so `None` here is not an error. + pub(crate) fn tinybird_target( + &self, + destination: ExportDestination, + ) -> Option<(&str, &str, u32, Option)> { + match destination { + ExportDestination::Tinybird => Some(( + self.endpoint.as_str(), + self.token.as_str(), + self.export_max_attempts, + None, + )), + ExportDestination::TinybirdMirror => self.mirror.as_ref().map(|mirror| { + ( + mirror.endpoint.as_str(), + mirror.token.as_str(), + mirror.max_attempts, + Some(mirror.export_timeout), + ) + }), + ExportDestination::ClickHouse => None, + } + } + + /// Whether `org_id`'s Tinybird-bound frames should also be mirrored. + pub(crate) fn mirrors_org(&self, org_id: &str) -> bool { + self.mirror.as_ref().is_some_and(|mirror| { + mirror.sample_percent >= 100 + || (mirror.sample_percent > 0 + && hash64(org_id) % 100 < u64::from(mirror.sample_percent)) + }) + } + pub fn validate_for_pipeline(&self, require_tinybird_credentials: bool) -> Result<(), String> { + if let Some(mirror) = &self.mirror { + // A half-configured mirror is always a deploy mistake, and it fails + // silently otherwise: rows go to a lane whose POSTs 401 forever. + if mirror.endpoint.is_empty() { + return Err( + "TINYBIRD_MIRROR_HOST is required when TINYBIRD_MIRROR_TOKEN is set".to_owned(), + ); + } + if mirror.token.is_empty() { + return Err( + "TINYBIRD_MIRROR_TOKEN is required when TINYBIRD_MIRROR_HOST is set".to_owned(), + ); + } + if mirror.max_attempts == 0 { + return Err("INGEST_TINYBIRD_MIRROR_MAX_ATTEMPTS must be greater than 0".to_owned()); + } + if mirror.sample_percent > 100 { + return Err("INGEST_TINYBIRD_MIRROR_SAMPLE_PERCENT must be 0..=100".to_owned()); + } + } if self.endpoint.is_empty() && require_tinybird_credentials { return Err( "TINYBIRD_HOST is required when INGEST_WRITE_MODE uses tinybird".to_owned(), @@ -888,11 +976,142 @@ impl TelemetryPipeline { Ok(stats) } - #[hotpath::measure] + /// Commit a batch to its destination lane, then — if a Tinybird mirror is + /// configured — to the mirror lane. + /// + /// Order matters. The primary commits first and is fail-closed: its errors + /// are what a client sees. The mirror runs second, is fail-open, and shares + /// the per-org byte budget, so the primary always wins the race for the last + /// reservable bytes. async fn commit_frames( &self, frames: Vec, destination: ExportDestination, + ) -> Result<(), PipelineError> { + // Cloned before the primary consumes `frames`, which means a mirrored + // batch is held twice in memory for the length of the primary commit. + // Bounded by the request body, and only paid while mirroring is on. + let mirrored = self.clone_frames_for_mirror(&frames, destination); + self.commit_frames_to_lane(frames, destination).await?; + if !mirrored.is_empty() { + self.commit_mirror_frames(mirrored).await; + } + Ok(()) + } + + /// The mirror copy of a batch, or empty when mirroring is off, the batch is + /// not Tinybird-bound, or the org is outside the sampling ramp. + /// + /// ClickHouse-bound frames are never mirrored: a BYO-ClickHouse org's rows + /// never reached the workspace being migrated, and must not start reaching + /// its replacement. + fn clone_frames_for_mirror( + &self, + frames: &[EncodedFrame], + destination: ExportDestination, + ) -> Vec { + if destination != ExportDestination::Tinybird || self.inner.cfg.mirror.is_none() { + return Vec::new(); + } + frames + .iter() + .filter(|frame| self.inner.cfg.mirrors_org(&frame.org_id)) + .map(|frame| EncodedFrame { + routing_key: frame.routing_key, + org_id: frame.org_id.clone(), + signal: frame.signal, + destination: ExportDestination::TinybirdMirror, + datasource: frame.datasource.clone(), + row_count: frame.row_count, + payload: frame.payload.clone(), + }) + .collect() + } + + /// Best-effort commit to the mirror lane. + /// + /// Every failure mode — a full lane, an exhausted org byte budget, a WAL + /// write error — is metered and dropped. Nothing here can return an error, + /// by design: the request was already accepted on the strength of the + /// primary commit, and the mirror is a migration aid, not a durability + /// guarantee. Gaps it leaves are repaired by the same backfill that carries + /// pre-mirror history. + /// + /// `record_failing_frame` is deliberately not called: that annotates the + /// 429 diagnosis path, and a mirror shed never produces a 429. + async fn commit_mirror_frames(&self, frames: Vec) { + let span = wal_commit_internal_span(ExportDestination::TinybirdMirror.as_str()); + let frame_count = frames.len(); + let mut committed_bytes = 0u64; + let mut frames_committed = 0usize; + let source_span = tracing::Span::current() + .context() + .span() + .span_context() + .clone(); + let source_span = source_span.is_sampled().then_some(source_span); + + async { + for mut frame in frames { + frame.destination = ExportDestination::TinybirdMirror; + #[expect( + clippy::cast_possible_truncation, + reason = "routing_key is a hash; a truncated hash still spreads uniformly across shards" + )] + let shard = (frame.routing_key as usize) % self.inner.cfg.wal_shards; + let lane = lane_index(shard, ExportDestination::TinybirdMirror); + let sender = self.inner.lane_senders[lane].clone(); + let Ok(permit) = sender.try_reserve_owned() else { + metrics::tinybird_mirror_dropped(&frame.datasource, "lane_full", frame.row_count as u64); + continue; + }; + let queued_bytes = frame.payload.len() as u64; + if self + .reserve_org_queue_bytes(&frame.org_id, queued_bytes) + .is_err() + { + metrics::tinybird_mirror_dropped(&frame.datasource, "org_quota", frame.row_count as u64); + continue; + } + let (start, end) = match self.inner.wal.append(lane, &frame).await { + Ok(offsets) => offsets, + Err(error) => { + self.release_org_queue_bytes(&frame.org_id, queued_bytes); + metrics::tinybird_mirror_dropped(&frame.datasource, "wal_error", frame.row_count as u64); + warn!(error = %error, "Dropping mirror frame after WAL append failure"); + continue; + } + }; + committed_bytes += queued_bytes; + frames_committed += 1; + permit.send(QueuedFrame { + shard, + start, + end, + org_id: frame.org_id, + queued_bytes, + signal: frame.signal, + destination: frame.destination, + datasource: frame.datasource, + row_count: frame.row_count, + payload: frame.payload, + source_span: source_span.clone(), + }); + } + } + .instrument(span.clone()) + .await; + + span.record("maple.ingest.frame_count", frame_count); + span.record("maple.ingest.frames_committed", frames_committed); + span.record("maple.ingest.queued_bytes", committed_bytes); + } + + #[hotpath::measure] + async fn commit_frames_to_lane( + &self, + frames: Vec, + destination: ExportDestination, ) -> Result<(), PipelineError> { if frames.is_empty() { return Ok(()); @@ -1653,10 +1872,15 @@ fn signal_from_tag(tag: u8) -> Option { } } +/// Wire tags are permanent: a WAL file written by one binary is replayed by the +/// next one, so a tag may be added but never reassigned. Tag 3 stays reserved +/// for the Tinybird mirror even after the mirror is removed — a rollback that +/// meets a mirror frame must decode it, not fail the shard. fn destination_tag(destination: ExportDestination) -> u8 { match destination { ExportDestination::Tinybird => 1, ExportDestination::ClickHouse => 2, + ExportDestination::TinybirdMirror => 3, } } @@ -1664,6 +1888,7 @@ fn destination_from_tag(tag: u8) -> Option { match tag { 1 => Some(ExportDestination::Tinybird), 2 => Some(ExportDestination::ClickHouse), + 3 => Some(ExportDestination::TinybirdMirror), _ => None, } } @@ -1799,7 +2024,10 @@ impl ExportWorker { let mut by_clickhouse: BTreeMap<(String, String), Vec<&QueuedFrame>> = BTreeMap::new(); for frame in &frames { match frame.destination { - ExportDestination::Tinybird => { + // Both Tinybird destinations group by datasource alone; which + // workspace the batch goes to is the lane's property, not the + // frame's, and a lane only ever carries one destination. + ExportDestination::Tinybird | ExportDestination::TinybirdMirror => { by_tinybird .entry(frame.datasource.clone()) .or_default() @@ -1844,6 +2072,7 @@ impl ExportWorker { } metrics::export_batch_completed( frames[0].shard, + self.destination.as_str(), &format!("{first_signal:?}"), start.elapsed().as_secs_f64(), end.saturating_sub(first_offset), @@ -1854,6 +2083,7 @@ impl ExportWorker { #[hotpath::measure] #[expect( clippy::cognitive_complexity, + clippy::too_many_lines, reason = "one export attempt plus the retry and error classification around it" )] async fn post_tinybird( @@ -1862,9 +2092,27 @@ impl ExportWorker { body: Vec, rows: usize, ) -> Result<(), String> { + // Credentials come from the lane's destination, not unconditionally from + // the primary config, so the mirror lane posts to the mirror workspace. + let Some((endpoint, token, max_attempts, timeout)) = + self.cfg.tinybird_target(self.destination) + else { + // Only reachable if a mirror lane drained frames after the mirror + // config was removed (a WAL replay across a config change). Dropping + // is right: the workspace those rows were meant for is gone. + metrics::tinybird_mirror_dropped(datasource, "no_target", rows as u64); + warn!( + datasource, + rows, + destination = self.destination.as_str(), + "Dropping Tinybird batch with no configured target" + ); + return Ok(()); + }; + let destination = self.destination.as_str(); let url = format!( "{}/v0/events?name={}", - self.cfg.endpoint.trim_end_matches('/'), + endpoint.trim_end_matches('/'), datasource ); let compressed = bytes::Bytes::from(gzip(&body)?); @@ -1873,7 +2121,6 @@ impl ExportWorker { // an upstream outage the retry loop below runs for minutes, and holding // it that long, once per in-flight lane, is real memory. drop(body); - let max_attempts = self.cfg.export_max_attempts; let mut attempt = 0u32; // The real host, not a hardcoded one — self-hosted and local Tinybird // endpoints were previously all reported as `api.tinybird.co`, which @@ -1904,16 +2151,21 @@ impl ExportWorker { span.record("db.operation.name", "INSERT"); span.record("db.collection.name", datasource); let span_handle = span.clone(); - let response = self + let request = self .http .post(&url) - .bearer_auth(&self.cfg.token) + .bearer_auth(token) .header(reqwest::header::CONTENT_TYPE, "application/x-ndjson") .header(reqwest::header::CONTENT_ENCODING, "gzip") - .body(compressed.clone()) - .send() - .instrument(span) - .await; + .body(compressed.clone()); + // Only the mirror sets one: a hanging mirror host must release its + // lane worker, whereas the primary is fail-closed and would rather + // wait than lose rows. + let request = match timeout { + Some(timeout) => request.timeout(timeout), + None => request, + }; + let response = request.send().instrument(span).await; let last_status: String; match response { @@ -1921,6 +2173,7 @@ impl ExportWorker { span_handle.record("http.response.status_code", response.status().as_u16()); span_handle.record("maple.ingest.outcome", "delivered"); metrics::tinybird_export_succeeded( + destination, datasource, started.elapsed().as_secs_f64(), rows as u64, @@ -1935,7 +2188,12 @@ impl ExportWorker { // Terminal: these rows are gone. 4xx from Tinybird means we // sent bad data, so this one is genuinely ours. record_stage_error(&span_handle, "non_retryable", &body, true); - metrics::tinybird_export_dropped(datasource, &status.to_string(), rows as u64); + metrics::tinybird_export_dropped( + destination, + datasource, + &status.to_string(), + rows as u64, + ); warn!(datasource, status, body = %body, rows, "Dropping non-retryable Tinybird batch"); return Ok(()); } @@ -1945,16 +2203,19 @@ impl ExportWorker { span_handle.record("http.response.status_code", status); span_handle.record("maple.ingest.outcome", "retry"); span_handle.record("error.type", "upstream_5xx"); - metrics::tinybird_export_retry(datasource, &last_status); - warn!(datasource, status, attempt, "Retrying Tinybird batch"); + metrics::tinybird_export_retry(destination, datasource, &last_status); + warn!( + datasource, + destination, status, attempt, "Retrying Tinybird batch" + ); } Err(error) => { last_status = "transport".to_owned(); span_handle.record("maple.ingest.outcome", "retry"); span_handle.record("error.type", "transport"); span_handle.record("otel.status_description", error.to_string().as_str()); - metrics::tinybird_export_retry(datasource, &last_status); - warn!(datasource, attempt, error = %error, "Retrying Tinybird batch after transport error"); + metrics::tinybird_export_retry(destination, datasource, &last_status); + warn!(datasource, destination, attempt, error = %error, "Retrying Tinybird batch after transport error"); } } @@ -1967,9 +2228,15 @@ impl ExportWorker { &format!("{attempt} attempts, last status {last_status}"), true, ); - metrics::tinybird_export_dropped(datasource, "retries_exhausted", rows as u64); + metrics::tinybird_export_dropped( + destination, + datasource, + "retries_exhausted", + rows as u64, + ); error!( datasource, + destination, rows, attempts = attempt, last_status = %last_status, @@ -3263,6 +3530,7 @@ mod tests { TinybirdConfig { endpoint: "http://tinybird.test".to_owned(), token: "token".to_owned(), + mirror: None, queue_dir: std::env::temp_dir(), queue_max_bytes: 1024 * 1024, org_queue_max_bytes: 1024 * 1024, @@ -4316,6 +4584,302 @@ mod tests { } } + #[test] + fn destination_tag_round_trips_all_variants() { + // WAL on-disk format. A frame written by one binary is replayed by the + // next, so every destination — the mirror included — must survive the + // tag encode/decode or the whole lane fails to replay. + for destination in ExportDestination::ALL { + assert_eq!( + destination_from_tag(destination_tag(destination)), + Some(destination) + ); + } + } + + #[test] + fn mirror_lane_ordinals_stay_stable_when_a_destination_is_appended() { + // `lane_index` is `shard * LANES_PER_SHARD + lane_ordinal`, and WAL files + // are named by destination. Appending TinybirdMirror must leave the two + // existing ordinals alone, or a rolling deploy renumbers lanes out from + // under frames already committed. + assert_eq!(lane_index(0, ExportDestination::Tinybird), 0); + assert_eq!(lane_index(0, ExportDestination::ClickHouse), 1); + assert_eq!(lane_index(0, ExportDestination::TinybirdMirror), 2); + assert_eq!(lane_index(1, ExportDestination::Tinybird), LANES_PER_SHARD); + } + + fn test_mirror_cfg(endpoint: &str, sample_percent: u8) -> TinybirdMirrorConfig { + TinybirdMirrorConfig { + endpoint: endpoint.to_owned(), + token: "mirror-token".to_owned(), + max_attempts: 2, + export_timeout: Duration::from_millis(500), + sample_percent, + } + } + + #[test] + fn mirror_sampling_is_all_or_nothing_per_org() { + // The ramp knob picks orgs by hash, so an org is consistently in or out. + // A percentage that split one org's stream would leave that org's rows + // half-mirrored, which no backfill boundary can describe. + let mut cfg = test_cfg(); + cfg.mirror = Some(test_mirror_cfg("http://mirror.test", 0)); + assert!(!cfg.mirrors_org("org_a"), "0% must mirror nothing"); + + cfg.mirror = Some(test_mirror_cfg("http://mirror.test", 100)); + assert!(cfg.mirrors_org("org_a")); + assert!(cfg.mirrors_org("org_b"), "100% must mirror every org"); + + cfg.mirror = None; + assert!(!cfg.mirrors_org("org_a"), "no mirror config, no mirroring"); + + // Same org, same answer, every time — a ramp never flips mid-stream. + cfg.mirror = Some(test_mirror_cfg("http://mirror.test", 50)); + let first = cfg.mirrors_org("org_ramp"); + assert!((0..64).all(|_| cfg.mirrors_org("org_ramp") == first)); + } + + #[test] + fn mirror_config_rejects_half_configuration() { + // Host without token (or the reverse) is always a deploy mistake, and it + // fails silently otherwise: the lane POSTs and 401s forever. + let mut cfg = test_cfg(); + cfg.mirror = Some(TinybirdMirrorConfig { + token: String::new(), + ..test_mirror_cfg("http://mirror.test", 100) + }); + assert!(cfg + .validate() + .unwrap_err() + .contains("TINYBIRD_MIRROR_TOKEN")); + + cfg.mirror = Some(test_mirror_cfg("", 100)); + assert!(cfg.validate().unwrap_err().contains("TINYBIRD_MIRROR_HOST")); + + cfg.mirror = Some(TinybirdMirrorConfig { + sample_percent: 101, + ..test_mirror_cfg("http://mirror.test", 100) + }); + assert!(cfg.validate().unwrap_err().contains("SAMPLE_PERCENT")); + + // A well-formed mirror is accepted even in ClickHouse-only mode, where + // the primary Tinybird credentials are absent by design. + cfg.endpoint = String::new(); + cfg.token = String::new(); + cfg.mirror = Some(test_mirror_cfg("http://mirror.test", 100)); + assert!(cfg.validate_for_pipeline(false).is_ok()); + } + + #[test] + fn tinybird_target_resolves_per_destination() { + let mut cfg = test_cfg(); + cfg.mirror = Some(test_mirror_cfg("http://mirror.test", 100)); + + let (endpoint, token, attempts, timeout) = cfg + .tinybird_target(ExportDestination::Tinybird) + .expect("primary target"); + assert_eq!(endpoint, "http://tinybird.test"); + assert_eq!(token, "token"); + assert_eq!(attempts, cfg.export_max_attempts); + assert!(timeout.is_none(), "the primary waits rather than lose rows"); + + let (endpoint, token, attempts, timeout) = cfg + .tinybird_target(ExportDestination::TinybirdMirror) + .expect("mirror target"); + assert_eq!(endpoint, "http://mirror.test"); + assert_eq!(token, "mirror-token"); + assert_eq!(attempts, 2); + assert_eq!(timeout, Some(Duration::from_millis(500))); + + assert!(cfg.tinybird_target(ExportDestination::ClickHouse).is_none()); + // A mirror lane draining replayed frames after the mirror config was + // removed has nowhere to send them, and says so rather than panicking. + cfg.mirror = None; + assert!(cfg + .tinybird_target(ExportDestination::TinybirdMirror) + .is_none()); + } + + /// A fake Tinybird that reports every import it receives. + async fn spawn_fake_tinybird() -> (String, mpsc::UnboundedReceiver) { + let (tx, rx) = mpsc::unbounded_channel(); + let listener = tokio::net::TcpListener::bind(("127.0.0.1", 0)) + .await + .unwrap(); + let addr = listener.local_addr().unwrap(); + let app = Router::new() + .route("/v0/events", post(fake_tinybird_import)) + .with_state(tx); + tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + (format!("http://{addr}"), rx) + } + + #[tokio::test] + async fn mirror_lane_receives_the_same_rows_as_the_primary() { + // The whole premise of the dual-emit window: both workspaces see an + // identical body, each authenticated with its own token. + let (primary_url, mut primary_rx) = spawn_fake_tinybird().await; + let (mirror_url, mut mirror_rx) = spawn_fake_tinybird().await; + + let queue_dir = unique_test_dir("mirror-parity"); + let mut cfg = test_cfg(); + cfg.endpoint = primary_url; + cfg.mirror = Some(test_mirror_cfg(&mirror_url, 100)); + cfg.queue_dir = queue_dir.clone(); + cfg.wal_shards = 1; + cfg.batch_max_wait = Duration::from_millis(1); + + let pipeline = TelemetryPipeline::new( + cfg, + Client::builder() + .timeout(Duration::from_secs(5)) + .build() + .unwrap(), + ) + .await + .unwrap(); + + pipeline + .accept_logs("org_mirror", &populated_log_request()) + .await + .unwrap(); + + let primary = tokio::time::timeout(Duration::from_secs(2), primary_rx.recv()) + .await + .expect("primary should receive an import") + .unwrap(); + let mirror = tokio::time::timeout(Duration::from_secs(2), mirror_rx.recv()) + .await + .expect("mirror should receive the same import") + .unwrap(); + + assert_eq!(primary.datasource, mirror.datasource); + assert_eq!(primary.body, mirror.body); + assert_eq!(primary.content_encoding, "gzip"); + assert_eq!(mirror.content_encoding, "gzip"); + assert_eq!(primary.authorization, "Bearer token"); + assert_eq!( + mirror.authorization, "Bearer mirror-token", + "the mirror lane must authenticate against the mirror workspace" + ); + + drop(std::fs::remove_dir_all(queue_dir)); + } + + #[tokio::test] + async fn clickhouse_bound_rows_are_never_mirrored() { + // A BYO-ClickHouse org's rows never reached the workspace being + // migrated, so they must not start reaching its replacement. + let (mirror_url, mut mirror_rx) = spawn_fake_tinybird().await; + + let queue_dir = unique_test_dir("mirror-skips-clickhouse"); + let mut cfg = test_cfg(); + cfg.mirror = Some(test_mirror_cfg(&mirror_url, 100)); + cfg.queue_dir = queue_dir.clone(); + cfg.wal_shards = 1; + cfg.batch_max_wait = Duration::from_millis(1); + + let provider = Arc::new(StaticClickHouseTargetProvider { + target: ClickHouseTarget { + endpoint: "http://127.0.0.1:1".to_owned(), + user: "ingest".to_owned(), + password: String::new(), + database: "maple".to_owned(), + }, + }); + let pipeline = TelemetryPipeline::new_with_clickhouse( + cfg, + Client::builder() + .timeout(Duration::from_millis(200)) + .build() + .unwrap(), + Some(provider), + ) + .await + .unwrap(); + + pipeline + .accept_logs_to( + "org_byo", + &populated_log_request(), + ExportDestination::ClickHouse, + ) + .await + .unwrap(); + + sleep(Duration::from_millis(300)).await; + assert!( + mirror_rx.try_recv().is_err(), + "ClickHouse-routed rows must never reach the Tinybird mirror" + ); + // The mirror lane's WAL should not even exist as a non-empty file. + let mirror_wal = queue_dir.join("shard-000-tinybird_mirror.wal"); + let mirror_bytes = std::fs::metadata(&mirror_wal).map_or(0, |meta| meta.len()); + assert_eq!( + mirror_bytes, 0, + "nothing should have been committed to the mirror lane" + ); + + drop(std::fs::remove_dir_all(queue_dir)); + } + + #[tokio::test] + async fn a_stalled_mirror_never_fails_the_accept_path() { + // The load-bearing guarantee of the whole design. The mirror host here + // refuses every connection, so its lane worker sits in the retry budget + // and its bounded channel fills. Accept must keep returning Ok and the + // primary must keep receiving every row — the mirror is allowed to lose + // rows, never to cost a request. + let (primary_url, mut primary_rx) = spawn_fake_tinybird().await; + + let queue_dir = unique_test_dir("mirror-stalled"); + let mut cfg = test_cfg(); + cfg.endpoint = primary_url; + // Port 1 is reserved and unbound: every mirror POST fails to connect. + cfg.mirror = Some(test_mirror_cfg("http://127.0.0.1:1", 100)); + cfg.queue_dir = queue_dir.clone(); + cfg.wal_shards = 1; + cfg.batch_max_wait = Duration::from_millis(1); + // Small enough that the mirror channel fills while its worker is stuck + // in the retry budget, but not so small that the primary — which drains + // every batch in milliseconds — sheds too. The primary is fail-closed by + // design, so a shed there is a real 429, not a mirror concern. + cfg.queue_channel_capacity = 8; + + let pipeline = TelemetryPipeline::new( + cfg, + Client::builder() + .timeout(Duration::from_millis(100)) + .build() + .unwrap(), + ) + .await + .unwrap(); + + // Paced so the primary lane keeps up; the mirror worker is held ~500ms + // per batch by its backoff, so its channel fills partway through. + for _ in 0..20u32 { + pipeline + .accept_logs("org_stalled", &populated_log_request()) + .await + .expect("a stalled mirror must never fail the accept path"); + sleep(Duration::from_millis(25)).await; + } + + // And the primary still got the rows. + let primary = tokio::time::timeout(Duration::from_secs(2), primary_rx.recv()) + .await + .expect("primary should still receive imports while the mirror is down") + .unwrap(); + assert_eq!(primary.datasource, "logs"); + + drop(std::fs::remove_dir_all(queue_dir)); + } + #[test] fn breaker_default_failure_threshold_is_two() { // Lowered from 4 → 2 so a stalled BYO-ClickHouse target enters fast-shed From f1a7d90fb3d0627b6c1b8dee8f50dd4f14c11792 Mon Sep 17 00:00:00 2001 From: Makisuo Date: Wed, 26 Aug 2026 01:18:04 +0200 Subject: [PATCH 2/5] fix(ci): target a Tinybird workspace by config baseUrl, not a --host flag MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `tinybird deploy --host X` exits `unknown option '--host'` on the pinned @tinybirdco/sdk 0.0.78, so the CD job's HOST_FLAG has never worked: an environment with TINYBIRD_HOST set fails the deploy outright, and one without it silently deploys to the tinybird.json default. Nobody noticed because the only leg with a token also had no host secret. Found while validating the us-east-1 leg added in the previous commit, which depends entirely on being able to point a deploy at a second workspace. The CLI takes its host from the config file's `baseUrl` and resolves config files in priority order, `tinybird.config.json` ahead of `tinybird.json`. Writing one in the job overrides the repo default for that deploy only, without making TINYBIRD_HOST mandatory for `tb build` locally — which is what putting `${TINYBIRD_HOST}` in the committed tinybird.json would do, since interpolation throws on an unset variable. Verified by reproducing the step against a local Tinybird: the deploy reaches the config-file host rather than api.tinybird.co. --- .github/workflows/tinybird-cd.yml | 14 ++++++++++---- .gitignore | 3 +++ 2 files changed, 13 insertions(+), 4 deletions(-) diff --git a/.github/workflows/tinybird-cd.yml b/.github/workflows/tinybird-cd.yml index b5326a2ae..36c4f0f02 100644 --- a/.github/workflows/tinybird-cd.yml +++ b/.github/workflows/tinybird-cd.yml @@ -64,12 +64,18 @@ jobs: echo "::notice::TINYBIRD_TOKEN not set for the '${{ matrix.target.environment }}' environment; skipping Tinybird deploy." exit 0 fi - HOST_FLAG=() + # `tinybird deploy` has NO --host flag (verified against the pinned + # @tinybirdco/sdk 0.0.78: it exits `unknown option '--host'`). The + # host comes from the config file's `baseUrl`, and the CLI resolves + # config files in priority order, `tinybird.config.json` ahead of + # `tinybird.json` — so writing one here overrides the repo default + # without making TINYBIRD_HOST mandatory for `tb build` locally. if [ -n "${TINYBIRD_HOST:-}" ]; then - HOST_FLAG=(--host "$TINYBIRD_HOST") + jq --arg host "$TINYBIRD_HOST" '.baseUrl = $host' tinybird.json > tinybird.config.json + echo "::notice::Deploying to $TINYBIRD_HOST" fi if [ "$ALLOW_DESTRUCTIVE" = "true" ]; then - bunx tinybird deploy "${HOST_FLAG[@]}" --allow-destructive-operations + bunx tinybird deploy --allow-destructive-operations else - bunx tinybird deploy "${HOST_FLAG[@]}" + bunx tinybird deploy fi diff --git a/.gitignore b/.gitignore index dcd957ded..aa8af3c41 100644 --- a/.gitignore +++ b/.gitignore @@ -1,4 +1,7 @@ .tinyb +# Written by the Tinybird CD job to point `deploy` at a specific workspace host +# (the CLI has no --host flag); takes priority over the committed tinybird.json. +tinybird.config.json .tinybird-generated .tinybird-entities-*.mjs # D1→Postgres migration dumps — contain prod data, never commit From a9c99bc4b4b5de45d370df8f13a9abac128e6813 Mon Sep 17 00:00:00 2001 From: Makisuo Date: Wed, 26 Aug 2026 01:28:30 +0200 Subject: [PATCH 3/5] docs(ingest): a mirror drop is permanent loss, not a backfill signal The migration is now mirror-and-cut: losing data older than 30 days is acceptable, so there is no backfill phase. That inverts what this counter means operationally. The old comment told whoever is watching the ramp that a non-zero value marks a window to repair later; nothing repairs it, and the rows are gone from the new workspace for good. Says so instead, because this is the number the ramp is gated on. --- apps/ingest/src/metrics.rs | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/apps/ingest/src/metrics.rs b/apps/ingest/src/metrics.rs index e4010089e..5b3b6aa15 100644 --- a/apps/ingest/src/metrics.rs +++ b/apps/ingest/src/metrics.rs @@ -682,8 +682,9 @@ pub fn tinybird_export_retry(destination: &str, datasource: &str, status: &str) /// (`reason` is `lane_full`, `org_quota`, `wal_error`, or `no_target`). /// /// The mirror is best-effort, so these drops are invisible to clients and to -/// every other counter. This is the only signal that a mirrored workspace has a -/// gap, and therefore the only thing that tells you a backfill window is needed. +/// every other counter — and nothing backfills the mirrored workspace, so a +/// non-zero value here is permanent loss, not a gap to be repaired later. It is +/// the one number that has to stay at zero for the whole migration window. pub fn tinybird_mirror_dropped(datasource: &str, reason: &str, rows: u64) { TINYBIRD_MIRROR_DROPPED_TOTAL.add( rows, From 69431141e31380e67e9a6a626a7ce62bf76921d8 Mon Sep 17 00:00:00 2001 From: Makisuo Date: Wed, 26 Aug 2026 01:43:57 +0200 Subject: [PATCH 4/5] fix(ingest): give the mirror lane its own per-org byte counter MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The mirror reserved from the same `org_queue_bytes` map as the primary, and bytes are only released once a frame exports. So a stalled mirror lane held an org's bytes for as long as it stayed stuck, and once they reached `org_queue_max_bytes` that org's PRIMARY commits failed to reserve — 429s on an org whose primary lane was perfectly healthy. That is precisely the coupling the per-destination lanes exist to prevent. A best-effort destination must not be able to reject a request for a fail-closed one, and "the primary reserves first" only holds within a single batch, not across the minutes a lane can stay stuck. Separate counter, and the export worker releases into whichever map its lane reserved from. The accompanying test fails with `Throttled("Telemetry org queue byte limit exceeded")` when the two share a counter. Also removes the need to raise INGEST_ORG_QUEUE_MAX_BYTES for the migration, which was a workaround for this bug rather than a real capacity requirement. --- apps/ingest/src/telemetry.rs | 99 +++++++++++++++++++++++++++++++++--- 1 file changed, 93 insertions(+), 6 deletions(-) diff --git a/apps/ingest/src/telemetry.rs b/apps/ingest/src/telemetry.rs index 87fb27ecd..6a2e12e30 100644 --- a/apps/ingest/src/telemetry.rs +++ b/apps/ingest/src/telemetry.rs @@ -732,6 +732,16 @@ struct PipelineInner { /// co-sharded Tinybird lane. lane_senders: Vec>, org_queue_bytes: Arc>>, + /// Per-org queued bytes for the Tinybird MIRROR lane, deliberately a + /// SEPARATE counter from `org_queue_bytes`. + /// + /// Sharing one counter breaks the isolation the lanes exist to provide. + /// Bytes are only released once a frame exports, so a stalled mirror lane + /// holds an org's bytes for as long as it is stuck — and once that reaches + /// `org_queue_max_bytes`, the org's PRIMARY commits start failing to reserve + /// and the org gets 429s even though its primary lane is perfectly healthy. + /// A best-effort destination must not be able to do that. + mirror_org_queue_bytes: Arc>>, } #[derive(Clone, Debug)] @@ -798,6 +808,7 @@ impl TelemetryPipeline { // the new per-destination lanes before workers start draining them. wal.migrate_legacy_shards(&cfg).await; let org_queue_bytes = Arc::new(DashMap::new()); + let mirror_org_queue_bytes = Arc::new(DashMap::new()); let clickhouse_breakers = Arc::new(ClickHouseBreakerRegistry::new(cfg.clickhouse_breaker)); let mut lane_senders = Vec::with_capacity(cfg.wal_shards * LANES_PER_SHARD); @@ -815,7 +826,13 @@ impl TelemetryPipeline { destination, cfg: Arc::clone(&cfg), wal: Arc::clone(&wal), - org_queue_bytes: Arc::clone(&org_queue_bytes), + // The worker releases into whichever counter its lane + // reserved from, so the mirror lane gets the mirror map. + org_queue_bytes: if destination == ExportDestination::TinybirdMirror { + Arc::clone(&mirror_org_queue_bytes) + } else { + Arc::clone(&org_queue_bytes) + }, clickhouse_breakers: Arc::clone(&clickhouse_breakers), clickhouse_targets: clickhouse_targets.clone(), http: http.clone(), @@ -831,6 +848,7 @@ impl TelemetryPipeline { wal, lane_senders, org_queue_bytes, + mirror_org_queue_bytes, }), }; pipeline.replay_committed_frames().await; @@ -1067,7 +1085,11 @@ impl TelemetryPipeline { }; let queued_bytes = frame.payload.len() as u64; if self - .reserve_org_queue_bytes(&frame.org_id, queued_bytes) + .reserve_org_bytes_in( + &self.inner.mirror_org_queue_bytes, + &frame.org_id, + queued_bytes, + ) .is_err() { metrics::tinybird_mirror_dropped(&frame.datasource, "org_quota", frame.row_count as u64); @@ -1076,7 +1098,11 @@ impl TelemetryPipeline { let (start, end) = match self.inner.wal.append(lane, &frame).await { Ok(offsets) => offsets, Err(error) => { - self.release_org_queue_bytes(&frame.org_id, queued_bytes); + release_org_queue_bytes( + &self.inner.mirror_org_queue_bytes, + &frame.org_id, + queued_bytes, + ); metrics::tinybird_mirror_dropped(&frame.datasource, "wal_error", frame.row_count as u64); warn!(error = %error, "Dropping mirror frame after WAL append failure"); continue; @@ -1201,13 +1227,20 @@ impl TelemetryPipeline { } fn reserve_org_queue_bytes(&self, org_id: &str, bytes: u64) -> Result<(), PipelineError> { + self.reserve_org_bytes_in(&self.inner.org_queue_bytes, org_id, bytes) + } + + fn reserve_org_bytes_in( + &self, + counters: &Arc>>, + org_id: &str, + bytes: u64, + ) -> Result<(), PipelineError> { if org_id.is_empty() || bytes == 0 { return Ok(()); } - let counter = self - .inner - .org_queue_bytes + let counter = counters .entry(org_id.to_owned()) .or_insert_with(|| Arc::new(AtomicU64::new(0))) .clone(); @@ -4827,6 +4860,60 @@ mod tests { drop(std::fs::remove_dir_all(queue_dir)); } + #[tokio::test] + async fn a_stalled_mirror_cannot_exhaust_the_org_byte_budget() { + // Regression: the mirror used to share `org_queue_bytes` with the primary. + // Bytes are only released once a frame exports, so a stalled mirror lane + // held an org's bytes indefinitely, and once they reached + // `org_queue_max_bytes` the org's PRIMARY commits failed to reserve — + // 429s for an org whose primary lane was perfectly healthy. Separate + // counters are what make the lane isolation actually hold. + let (primary_url, mut primary_rx) = spawn_fake_tinybird().await; + + let queue_dir = unique_test_dir("mirror-org-budget"); + let mut cfg = test_cfg(); + cfg.endpoint = primary_url; + cfg.mirror = Some(test_mirror_cfg("http://127.0.0.1:1", 100)); + cfg.queue_dir = queue_dir.clone(); + cfg.wal_shards = 1; + cfg.batch_max_wait = Duration::from_millis(1); + // Sized against the frame this request encodes to: a handful of stuck + // mirror frames must exceed the budget, while the primary — which + // exports and releases within milliseconds — never holds more than one. + cfg.org_queue_max_bytes = 4_096; + // Room for the mirror lane to accumulate rather than shedding at the + // channel before it ever reserves any bytes. + cfg.queue_channel_capacity = 64; + + let pipeline = TelemetryPipeline::new( + cfg, + Client::builder() + .timeout(Duration::from_millis(100)) + .build() + .unwrap(), + ) + .await + .unwrap(); + + for _ in 0..40u32 { + pipeline + .accept_logs("org_budget", &populated_log_request()) + .await + .expect("mirror bytes must never consume the primary's org budget"); + sleep(Duration::from_millis(10)).await; + } + + // The primary kept exporting throughout, which is what proves its + // reservations were never starved by the stuck mirror lane. + let primary = tokio::time::timeout(Duration::from_secs(2), primary_rx.recv()) + .await + .expect("primary should keep exporting while the mirror lane is stuck") + .unwrap(); + assert_eq!(primary.datasource, "logs"); + + drop(std::fs::remove_dir_all(queue_dir)); + } + #[tokio::test] async fn a_stalled_mirror_never_fails_the_accept_path() { // The load-bearing guarantee of the whole design. The mirror host here From b7d5072151ae5d676a721b3ad8d5ed8b52d176cc Mon Sep 17 00:00:00 2001 From: Makisuo Date: Wed, 26 Aug 2026 01:50:04 +0200 Subject: [PATCH 5/5] fix(ingest): size the WAL for three lanes, not two MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The per-lane WAL budget is INGEST_QUEUE_MAX_BYTES / (WAL_SHARDS * lanes). The mirror adds a third lane per shard, so at 8 GiB every lane would have dropped from 1 GiB to 683 MiB — including the primary Tinybird lane, at the exact moment its traffic doubled and a second full copy of every payload started competing for the same disk. 12 GiB over 12 lanes restores the 1 GiB per lane that 8 GiB gave across 8, and leaves ~8 GB of Fargate's 20 GB ephemeral allowance for the image and OS. --- apps/ingest/alchemy.run.ts | 14 +++++++++----- 1 file changed, 9 insertions(+), 5 deletions(-) diff --git a/apps/ingest/alchemy.run.ts b/apps/ingest/alchemy.run.ts index ad3ac3ee4..dff17f7be 100644 --- a/apps/ingest/alchemy.run.ts +++ b/apps/ingest/alchemy.run.ts @@ -63,12 +63,16 @@ const INGEST_PORT = 3474 * sit exactly on the line. 8 GiB still buys hours of buffering at current * volume; raising it means paying for ephemeral storage beyond the free tier. * - * Note the per-lane budget is `INGEST_QUEUE_MAX_BYTES / (WAL_SHARDS * lanes)`, - * so enabling the Tinybird mirror (a third lane) cuts every lane's share by a - * third at the same moment Tinybird-bound traffic doubles. Raise this — within - * the ephemeral ceiling — for the duration of a mirrored window. + * The per-lane budget is `INGEST_QUEUE_MAX_BYTES / (WAL_SHARDS * lanes)`, so + * the Tinybird mirror's third lane would have cut every lane's share from 1 GiB + * to 683 MiB at exactly the moment Tinybird-bound traffic doubled. 12 GiB over + * 12 lanes restores the 1 GiB per lane that 8 GiB gave across 8, and still + * leaves ~8 GB of the ephemeral allowance for the image and the OS. + * + * Drop back to 8 GiB once the mirror is removed, or every lane silently gains + * headroom nobody sized for. */ -const WAL_MAX_BYTES = 8 * 1024 * 1024 * 1024 +const WAL_MAX_BYTES = 12 * 1024 * 1024 * 1024 /** * Pinned rather than derived. The gateway defaults to `num_cpus * 2`, which