diff --git a/README.md b/README.md index 0c4211f..4ea7896 100644 --- a/README.md +++ b/README.md @@ -27,6 +27,22 @@ filtered archive replay, and a gapless replay→live cutover. typed pre-upgrade errors (`CursorTooOld` / `UnknownZstdDictionary` / `InvalidRequest`), reconnect backoff that resumes at the highest delivered seq, TCP keepalive. +- **multi-host failover** (extension beyond the upstream Go client, which + binds a single `service`): `fallback_hosts` rotates to another instance + after `failover_after_errors` consecutive live failures with no events + between (or a sweep whose bounded retries exhaust). v2 seqs are + per-instance, so the switch re-anchors by witnessed time: the new host's + `listSegments` bounds place the sweep (`anchorForFloor`), and a delivery + time floor at `last_time_us − failover_rewind_us` trims the + segment-granularity over-fetch, bounding duplicates to the rewind margin + (default 10s, the v1 client's number). failback retries the primary after + any rotation that followed delivered events, mirroring zat's v1 + round-robin (`nextHostIndex`). cross-host semantics are at-least-once for + the live flow; a caller-supplied `after_seq`/`live_cursor` is bound to + the primary's seq space, so a primary that dies before delivering + anything fails loudly rather than silently restarting elsewhere. note: + re-anchoring reads the new host's `listSegments`, so token-gated + archives need `api_key` for failover, not just for sweeps. - **`ArchiveBackfill`** — the archive half standalone: `planSnapshot` (collections/dids, seq bounds, `api_key` for token-gated instances) → `getBlock`/`getSegment` → jss v1 columnar decode. two handler contracts: diff --git a/src/archive_backfill.zig b/src/archive_backfill.zig index fcc4df5..307f237 100644 --- a/src/archive_backfill.zig +++ b/src/archive_backfill.zig @@ -467,17 +467,42 @@ pub fn seqBoundsForWindow(spans: []const SegmentSpan, start_us: i64, end_us: ?i6 return .{ .after_seq = after, .before_seq = before, .covered = covered }; } -/// fetch every sealed segment's metadata from `host` and derive the seq -/// bounds covering the witnessed window [start_us, end_us]. end_us == null -/// means "through the sealed tip". -pub fn fetchSeqBounds(io: Io, allocator: Allocator, host: []const u8, start_us: i64, end_us: ?i64, api_key: ?[]const u8) !SeqBounds { +/// re-anchor point for entering a DIFFERENT instance's seq space from a +/// witnessed-time floor (cross-instance failover: seqs don't transfer, +/// witnessed time approximately does). `after_seq` is where the sweep +/// starts on the new host; rows older than the floor are the caller's to +/// trim (segment granularity is coarse — hours — so the seq bound decides +/// what to FETCH and the time floor decides what to DELIVER). +pub const Anchor = struct { + after_seq: ?u64, + /// false when no sealed segment reaches the floor: the window lives in + /// the unsealed tail (after_seq = sealed tip, live cutover covers the + /// rest) or the host has no sealed archive at all (after_seq = null, + /// start live from the tip). + covered: bool, +}; + +/// reduce sealed spans to a failover anchor for `floor_us`: the sweep must +/// start at or before the first row witnessed >= floor. covered → one seq +/// below the earliest intersecting segment; not covered but spans exist → +/// the sealed tip (the floor lies in the unsealed tail; the live cutover's +/// server-side cold replay covers it); no spans → null (live from tip). +pub fn anchorForFloor(spans: []const SegmentSpan, floor_us: i64) Anchor { + const bounds = seqBoundsForWindow(spans, floor_us, null); + if (bounds.covered) return .{ .after_seq = bounds.after_seq, .covered = true }; + var tip: ?u64 = null; + for (spans) |span| { + if (tip == null or span.max_seq > tip.?) tip = span.max_seq; + } + return .{ .after_seq = tip, .covered = false }; +} + +/// fetch every sealed segment's metadata from `host` into `arena` +pub fn fetchSegmentSpans(io: Io, allocator: Allocator, arena: Allocator, host: []const u8, api_key: ?[]const u8) ![]SegmentSpan { var transport = HttpTransport.init(io, allocator); defer transport.deinit(); const authorization = try bearerHeader(allocator, api_key); defer if (authorization) |a| allocator.free(a); - var arena_state = std.heap.ArenaAllocator.init(allocator); - defer arena_state.deinit(); - const arena = arena_state.allocator(); var spans: std.ArrayList(SegmentSpan) = .empty; var cursor: ?[]const u8 = null; @@ -515,7 +540,26 @@ pub fn fetchSeqBounds(io: Io, allocator: Allocator, host: []const u8, start_us: else => break, }; } - return seqBoundsForWindow(spans.items, start_us, end_us); + return spans.items; +} + +/// fetch every sealed segment's metadata from `host` and derive the seq +/// bounds covering the witnessed window [start_us, end_us]. end_us == null +/// means "through the sealed tip". +pub fn fetchSeqBounds(io: Io, allocator: Allocator, host: []const u8, start_us: i64, end_us: ?i64, api_key: ?[]const u8) !SeqBounds { + var arena_state = std.heap.ArenaAllocator.init(allocator); + defer arena_state.deinit(); + const spans = try fetchSegmentSpans(io, allocator, arena_state.allocator(), host, api_key); + return seqBoundsForWindow(spans, start_us, end_us); +} + +/// fetch segment metadata from `host` and derive the failover anchor for +/// `floor_us` (see anchorForFloor) +pub fn fetchAnchor(io: Io, allocator: Allocator, host: []const u8, floor_us: i64, api_key: ?[]const u8) !Anchor { + var arena_state = std.heap.ArenaAllocator.init(allocator); + defer arena_state.deinit(); + const spans = try fetchSegmentSpans(io, allocator, arena_state.allocator(), host, api_key); + return anchorForFloor(spans, floor_us); } fn parsePlan(arena: Allocator, body: []const u8) !Plan { @@ -1120,6 +1164,31 @@ test "plan request carries seq bounds only when set" { try std.testing.expect(mem.indexOf(u8, bounded, "\"beforeSeq\":200") != null); } +test "anchorForFloor: covered floor, unsealed-tail floor, empty archive" { + const spans = [_]SegmentSpan{ + .{ .min_seq = 1, .max_seq = 100, .min_witnessed_us = 1000, .max_witnessed_us = 2000 }, + .{ .min_seq = 101, .max_seq = 200, .min_witnessed_us = 1900, .max_witnessed_us = 3000 }, + }; + + // floor inside the sealed range: sweep from one below the earliest + // intersecting segment (the time floor trims the over-fetch) + const covered = anchorForFloor(&spans, 2500); + try std.testing.expect(covered.covered); + try std.testing.expectEqual(@as(?u64, 100), covered.after_seq); + + // floor past every sealed segment: the window lives in the unsealed + // tail — anchor at the sealed tip and let the live cutover's cold + // replay cover the rest + const tail = anchorForFloor(&spans, 9000); + try std.testing.expect(!tail.covered); + try std.testing.expectEqual(@as(?u64, 200), tail.after_seq); + + // no sealed archive at all: nothing to anchor; live from the tip + const empty = anchorForFloor(&.{}, 2500); + try std.testing.expect(!empty.covered); + try std.testing.expectEqual(@as(?u64, null), empty.after_seq); +} + test "seqBoundsForWindow admits intersecting segments, reports uncovered windows" { // three sealed segments; witnessed spans overlap at the edges the way // real segments do (per-connection skew) diff --git a/src/client.zig b/src/client.zig index 9c7d1d7..09ba342 100644 --- a/src/client.zig +++ b/src/client.zig @@ -21,6 +21,33 @@ //! regress it). re-backfill cycles that fail to advance the cursor are //! bounded (upstream maxRebackfillStalls = 5) — a non-advancing cycle is a //! pathological loop, not a real fall-behind. +//! +//! === multi-host failover (extension beyond the upstream Go client) === +//! +//! v2 seqs are per-instance, so a cursor cannot follow a host switch the +//! way zat's v1 client round-robins (v1 cursors are wall-clock time_us; +//! see zat src/internal/streaming/jetstream.zig). the bridge is witnessed +//! time: every delivered event carries time_us, so on failover we re-enter +//! the NEXT host's seq space through its own archive metadata — +//! +//! anchor: listSegments time bounds → the seq one below the earliest +//! sealed segment reaching (last_time_us − rewind margin); the normal +//! sweep→cutover loop runs from there (archive.anchorForFloor); +//! trim: segment granularity is coarse (hours), so the seq bound decides +//! what to FETCH and the time floor decides what to DELIVER — rows +//! witnessed before the floor are skipped, bounding duplicates to the +//! rewind margin; +//! failback: like the v1 client, a rotation after a connection that +//! delivered events retries the primary first (nextHostIndex). +//! +//! semantics across a switch are at-least-once for the live flow. two +//! honest caveats, both inherent to time re-anchoring: instances witness +//! independently (the margin covers skew, not resyncs the new host +//! witnessed at very different times), and a caller-supplied after_seq / +//! live_cursor is bound to the PRIMARY's seq space — if the primary dies +//! before any event is delivered there is no time anchor yet, and the +//! client fails with the primary's error rather than silently breaking +//! the caller's resume contract. const std = @import("std"); const zat = @import("zat"); @@ -37,8 +64,20 @@ const log = std.log.scoped(.jetstream); pub const max_rebackfill_stalls = 5; pub const Options = struct { - /// base URL of the instance, e.g. "https://stream.waow.tech" + /// base URL of the primary instance, e.g. "https://stream.waow.tech" host: []const u8, + /// additional instances to rotate through when the current host is + /// declared unhealthy. cross-instance resume re-anchors by witnessed + /// time (see module doc); empty = single-host behavior, terminal + /// errors propagate exactly as before. + fallback_hosts: []const []const u8 = &.{}, + /// consecutive live transient failures WITH NO EVENTS BETWEEN before + /// the current host is declared unhealthy and rotation kicks in. + /// only meaningful with fallback_hosts; 0 disables rotation. + failover_after_errors: u8 = 5, + /// witnessed-time rewind margin when re-anchoring on a new host, + /// covering inter-instance witness skew (v1 client rewinds 10s) + failover_rewind_us: i64 = 10 * std.time.us_per_s, /// replay start, exclusive: sweep the sealed archive from here and cut /// over to live. null = no replay — a pure live tail from the tip /// (set live_cursor to resume a live tail at a known seq instead). @@ -63,8 +102,36 @@ pub const Options = struct { /// live reconnect backoff overrides (tests); zero = defaults backoff_min_ns: u64 = 0, backoff_max_ns: u64 = 0, + + fn hostCount(self: Options) usize { + return 1 + self.fallback_hosts.len; + } + + fn hostAt(self: Options, index: usize) []const u8 { + return if (index == 0) self.host else self.fallback_hosts[index - 1]; + } }; +/// cross-host delivery state: the time anchor rotation re-enters from, +/// and the failback signal +const Progress = struct { + /// witnessed time of the newest delivered event, 0 = nothing yet + last_time_us: i64 = 0, + /// events delivered since the last rotation (failback: retry the + /// primary first, matching the v1 client) + delivered_since_rotate: bool = false, +}; + +/// v1 client's rotation rule (zat jetstream.zig nextHostIndex): a healthy +/// connection's drop retries the primary first; a failed attempt keeps +/// rotating +pub fn nextHostIndex(current: usize, delivered_events: bool, host_count: usize) usize { + if (delivered_events) return 0; + return (current + 1) % host_count; +} + +const HostOutcome = enum { stopped, rotate }; + /// run the unified stream through `handler`: /// fn onEvent(*H, livedecode.Event) bool — every event, archive and live, /// in seq order across the seam; false stops cleanly @@ -73,43 +140,156 @@ pub const Options = struct { /// blocks until the handler stops it, the snapshot completes /// (snapshot_only), or a terminal error. pub fn subscribe(io: Io, allocator: Allocator, options: Options, handler: anytype) !void { + const host_count = options.hostCount(); + const can_rotate = host_count > 1 and options.failover_after_errors > 0; + + var progress = Progress{}; + var host_idx: usize = 0; + var after = options.after_seq; + var live_cur = options.live_cursor; + var time_floor: i64 = 0; + + const rot_min: u64 = if (options.backoff_min_ns > 0) options.backoff_min_ns else 250 * std.time.ns_per_ms; + const rot_max: u64 = if (options.backoff_max_ns > 0) options.backoff_max_ns else 30 * std.time.ns_per_s; + var rot_backoff = rot_min; + + while (true) { + const host = options.hostAt(host_idx); + const outcome = runOnHost(io, allocator, options, host, after, live_cur, time_floor, can_rotate, &progress, handler) catch |err| { + // terminal for the whole stream (single-host keeps the exact + // pre-failover error surface; Canceled/OOM/InvalidRequest are + // always terminal) + return err; + }; + switch (outcome) { + .stopped => return, + .rotate => { + // no time anchor yet + a caller-supplied seq anchor: that + // anchor is bound to this host's seq space and cannot be + // mapped — fail loudly instead of silently restarting + // somewhere else (module doc) + if (progress.last_time_us == 0 and (after != null or live_cur != null)) + return error.HostUnhealthy; + + const delivered = progress.delivered_since_rotate; + const prev = host_idx; + host_idx = nextHostIndex(host_idx, delivered, host_count); + progress.delivered_since_rotate = false; + rot_backoff = if (delivered) rot_min else @min(rot_backoff * 2, rot_max); + log.warn("host {s} unhealthy; failing over to {s}", .{ options.hostAt(prev), options.hostAt(host_idx) }); + + if (progress.last_time_us > 0) { + const floor = progress.last_time_us - options.failover_rewind_us; + const anchor = archive.fetchAnchor(io, allocator, options.hostAt(host_idx), floor, options.api_key) catch |err| switch (err) { + error.Canceled, error.OutOfMemory => |e| return e, + else => |e| { + log.warn("anchor fetch on {s} failed: {s}; rotating on", .{ options.hostAt(host_idx), @errorName(e) }); + io.sleep(Io.Duration.fromNanoseconds(@intCast(rot_backoff)), .awake) catch return error.Canceled; + continue; + }, + }; + after = anchor.after_seq; + live_cur = null; + time_floor = floor; + } else { + // nothing promised yet and no caller anchor: fresh + // live tail from the new host's tip + after = null; + live_cur = null; + time_floor = 0; + } + io.sleep(Io.Duration.fromNanoseconds(@intCast(rot_backoff)), .awake) catch return error.Canceled; + }, + } + } +} + +/// the pre-failover single-host flow, verbatim, parameterized by host and +/// anchor. can_rotate=false propagates every error exactly as the +/// original did; can_rotate=true converts host-health errors (bounded +/// archive retries exhausted, live transient threshold) into .rotate. +fn runOnHost( + io: Io, + allocator: Allocator, + options: Options, + host: []const u8, + after_seq: ?u64, + live_cursor: ?u64, + time_floor_us: i64, + can_rotate: bool, + progress: *Progress, + handler: anytype, +) !HostOutcome { var dict_buf: ?[]u8 = null; defer if (dict_buf) |owned| allocator.free(owned); - if (options.zstd_compression) dict_buf = fetchZstdDict(io, allocator, options.host); + if (options.zstd_compression) dict_buf = fetchZstdDict(io, allocator, host); - if (options.after_seq == null) { + const H = @TypeOf(handler.*); + + if (after_seq == null) { // pure live tail (upstream runLiveOnly) - var client = live_mod.LiveClient.init(io, allocator, liveOptions(options, options.live_cursor, 0, dict_buf)); + var client = live_mod.LiveClient.init(io, allocator, liveOptions(options, host, live_cursor, 0, dict_buf)); defer client.deinit(); - var filtering = FilteringHandler(@TypeOf(handler.*)){ .inner = handler, .kinds = options.kinds }; - return client.run(&filtering) catch |err| switch (err) { - error.CursorTooOld => error.CursorTooOld, // live-only has no backfill to re-enter - else => |e| e, + var filtering = FilteringHandler(H){ + .inner = handler, + .kinds = options.kinds, + .time_floor_us = time_floor_us, + .progress = progress, + .rotate_threshold = if (can_rotate) options.failover_after_errors else 0, + }; + client.run(&filtering) catch |err| switch (err) { + error.CursorTooOld => return error.CursorTooOld, // live-only has no backfill to re-enter + else => |e| return e, }; + return if (filtering.rotate_pending) .rotate else .stopped; } - var cursor = options.after_seq.?; + var cursor = after_seq.?; var stalls: u8 = 0; while (true) { - var sweeper = SweepHandler(@TypeOf(handler.*)){ .inner = handler, .kinds = options.kinds }; - const result = try archive.run(io, allocator, .{ - .host = options.host, + var sweeper = SweepHandler(H){ + .inner = handler, + .kinds = options.kinds, + .time_floor_us = time_floor_us, + .progress = progress, + }; + const result = archive.run(io, allocator, .{ + .host = host, .collections = options.collections, .dids = options.dids, .api_key = options.api_key, .after_seq = cursor, .concurrency = options.concurrency, - }, &sweeper); - if (result.stopped) return; - if (options.snapshot_only) return; + }, &sweeper) catch |err| switch (err) { + error.Canceled, error.OutOfMemory => |e| return e, + else => |e| { + // the archive path retries transients internally with a + // bound (fetchWithRetry); exhaustion is a host-health + // verdict, not a stream verdict — when there is somewhere + // to go + if (can_rotate) { + log.warn("archive sweep on {s} failed: {s}", .{ host, @errorName(e) }); + return .rotate; + } + return e; + }, + }; + if (result.stopped) return .stopped; + if (options.snapshot_only) return .stopped; // cut over at the HIGHER of the freshly-learned sealed tip and the // cursor already processed through (monotonic non-decreasing; see // module doc) const cutover = @max(result.planned_through_seq, cursor); - var client = live_mod.LiveClient.init(io, allocator, liveOptions(options, cutover, cutover, dict_buf)); + var client = live_mod.LiveClient.init(io, allocator, liveOptions(options, host, cutover, cutover, dict_buf)); defer client.deinit(); - var filtering = FilteringHandler(@TypeOf(handler.*)){ .inner = handler, .kinds = options.kinds }; + var filtering = FilteringHandler(H){ + .inner = handler, + .kinds = options.kinds, + .time_floor_us = time_floor_us, + .progress = progress, + .rotate_threshold = if (can_rotate) options.failover_after_errors else 0, + }; client.run(&filtering) catch |err| switch (err) { error.CursorTooOld => { // §14: re-backfill from the last durably-processed seq, @@ -129,13 +309,13 @@ pub fn subscribe(io: Io, allocator: Allocator, options: Options, handler: anytyp }, else => |e| return e, }; - return; // clean stop: handler asked, or cancellation unwound + return if (filtering.rotate_pending) .rotate else .stopped; } } -fn liveOptions(options: Options, cursor: ?u64, dedup_floor: u64, dict: ?[]const u8) live_mod.Options { +fn liveOptions(options: Options, host: []const u8, cursor: ?u64, dedup_floor: u64, dict: ?[]const u8) live_mod.Options { var out: live_mod.Options = .{ - .host = options.host, + .host = host, .cursor = cursor, .dedup_floor = dedup_floor, .kinds = options.kinds, @@ -150,32 +330,55 @@ fn liveOptions(options: Options, cursor: ?u64, dedup_floor: u64, dict: ?[]const /// archive-half adapter: receives every row (with seq) from the sweep and /// applies the kinds filter — the plan prunes server-side where it can, -/// but the client-side matcher remains the correctness backstop +/// but the client-side matcher remains the correctness backstop. after a +/// failover the time floor trims the segment-granularity over-fetch. fn SweepHandler(comptime H: type) type { return struct { inner: *H, kinds: []const livedecode.Kind, + time_floor_us: i64 = 0, + progress: *Progress, pub fn onRow(self: *@This(), event: livedecode.Event) bool { + if (self.time_floor_us > 0 and event.time_us < self.time_floor_us) return true; if (!wantsKind(self.kinds, event)) return true; + noteDelivery(self.progress, event.time_us); return self.inner.onEvent(event); } }; } /// live-half adapter: same kinds backstop over the tail (the server also -/// filters; upstream wantsLive) +/// filters; upstream wantsLive). counts consecutive transports failures +/// with no events between toward the failover threshold. fn FilteringHandler(comptime H: type) type { return struct { inner: *H, kinds: []const livedecode.Kind, + time_floor_us: i64 = 0, + progress: *Progress, + /// 0 = never rotate (single-host / rotation disabled) + rotate_threshold: u8 = 0, + consecutive_errors: u8 = 0, + rotate_pending: bool = false, pub fn onEvent(self: *@This(), event: livedecode.Event) bool { + self.consecutive_errors = 0; + if (self.time_floor_us > 0 and event.time_us < self.time_floor_us) return true; if (!wantsKind(self.kinds, event)) return true; + noteDelivery(self.progress, event.time_us); return self.inner.onEvent(event); } pub fn onError(self: *@This(), err: anyerror) bool { + if (self.rotate_threshold > 0) { + self.consecutive_errors += 1; + if (self.consecutive_errors >= self.rotate_threshold) { + self.rotate_pending = true; + log.warn("live tail: {d} consecutive failures without progress (last: {s}); requesting failover", .{ self.consecutive_errors, @errorName(err) }); + return false; + } + } if (comptime @hasDecl(H, "onError")) return self.inner.onError(err); log.warn("live tail: {s} (reconnecting)", .{@errorName(err)}); return true; @@ -192,6 +395,11 @@ fn FilteringHandler(comptime H: type) type { }; } +fn noteDelivery(progress: *Progress, time_us: i64) void { + if (time_us > progress.last_time_us) progress.last_time_us = time_us; + progress.delivered_since_rotate = true; +} + fn wantsKind(kinds: []const livedecode.Kind, event: livedecode.Event) bool { if (kinds.len == 0) return true; const kind: livedecode.Kind = event.payload; @@ -246,3 +454,111 @@ test "kinds matcher: empty admits all, otherwise exact" { try testing.expect(!wantsKind(&.{.commit}, ev_sync)); try testing.expect(wantsKind(&.{ .commit, .sync }, ev_sync)); } + +test "nextHostIndex mirrors the v1 client: failback to primary after progress" { + // a connection that delivered events, then dropped: primary first + try testing.expectEqual(@as(usize, 0), nextHostIndex(2, true, 3)); + // failed attempts keep rotating, wrapping + try testing.expectEqual(@as(usize, 1), nextHostIndex(0, false, 3)); + try testing.expectEqual(@as(usize, 2), nextHostIndex(1, false, 3)); + try testing.expectEqual(@as(usize, 0), nextHostIndex(2, false, 3)); +} + +test "hostAt / hostCount: primary then fallbacks in order" { + const opts = Options{ + .host = "https://a.example", + .fallback_hosts = &.{ "https://b.example", "https://c.example" }, + }; + try testing.expectEqual(@as(usize, 3), opts.hostCount()); + try testing.expectEqualStrings("https://a.example", opts.hostAt(0)); + try testing.expectEqualStrings("https://b.example", opts.hostAt(1)); + try testing.expectEqualStrings("https://c.example", opts.hostAt(2)); +} + +const CountingHandler = struct { + events: usize = 0, + last_seq: u64 = 0, + + pub fn onEvent(self: *@This(), event: livedecode.Event) bool { + self.events += 1; + self.last_seq = event.seq; + return true; + } +}; + +fn commitEvent(seq: u64, time_us: i64) livedecode.Event { + return .{ .seq = seq, .time_us = time_us, .did = "did:plc:a", .payload = .{ .commit = .{ + .operation = .create, + .collection = "c", + .rkey = "k", + .rev = "r", + } } }; +} + +test "time floor trims below-floor rows in both halves and tracks progress" { + var inner = CountingHandler{}; + var progress = Progress{}; + + var sweep = SweepHandler(CountingHandler){ + .inner = &inner, + .kinds = &.{}, + .time_floor_us = 1000, + .progress = &progress, + }; + try testing.expect(sweep.onRow(commitEvent(1, 500))); // below floor: skipped, still continue + try testing.expectEqual(@as(usize, 0), inner.events); + try testing.expect(!progress.delivered_since_rotate); + try testing.expect(sweep.onRow(commitEvent(2, 1500))); + try testing.expectEqual(@as(usize, 1), inner.events); + try testing.expectEqual(@as(i64, 1500), progress.last_time_us); + try testing.expect(progress.delivered_since_rotate); + + var live = FilteringHandler(CountingHandler){ + .inner = &inner, + .kinds = &.{}, + .time_floor_us = 1000, + .progress = &progress, + }; + try testing.expect(live.onEvent(commitEvent(3, 900))); + try testing.expectEqual(@as(usize, 1), inner.events); + try testing.expect(live.onEvent(commitEvent(4, 2000))); + try testing.expectEqual(@as(usize, 2), inner.events); + try testing.expectEqual(@as(i64, 2000), progress.last_time_us); +} + +test "failover threshold: consecutive errors rotate, any event resets" { + var inner = CountingHandler{}; + var progress = Progress{}; + var live = FilteringHandler(CountingHandler){ + .inner = &inner, + .kinds = &.{}, + .progress = &progress, + .rotate_threshold = 3, + }; + + try testing.expect(live.onError(error.ReadFailed)); + try testing.expect(live.onError(error.ReadFailed)); + // an event between failures resets the count: the host is limping, + // not dead + try testing.expect(live.onEvent(commitEvent(1, 100))); + try testing.expect(live.onError(error.ReadFailed)); + try testing.expect(live.onError(error.ReadFailed)); + try testing.expect(!live.rotate_pending); + // third consecutive: rotation requested, run() stops + try testing.expect(!live.onError(error.ReadFailed)); + try testing.expect(live.rotate_pending); +} + +test "failover threshold disabled: errors fall through to default handling" { + var inner = CountingHandler{}; + var progress = Progress{}; + var live = FilteringHandler(CountingHandler){ + .inner = &inner, + .kinds = &.{}, + .progress = &progress, + .rotate_threshold = 0, + }; + // never rotates, never stops on the default path + for (0..20) |_| try testing.expect(live.onError(error.ReadFailed)); + try testing.expect(!live.rotate_pending); +}