Something went wrong. Try again.
A Zig library for Bluesky activity and feed analytics, built on Zat.
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124const std = @import("std");const zat = @import("zat");
pub const Headers = zat.HttpTransport.RateLimitHeaders;pub const Shared = struct { mutex: std.Io.Mutex = .init, governor: Governor = .{},
pub fn wait(self: *Shared, io: std.Io) !void { while (true) { try self.mutex.lock(io); const now = std.Io.Clock.awake.now(io).toMilliseconds(); const delay = @max(0, self.governor.next_ms - now); if (delay == 0) { self.governor.next_ms = now + @as(i64, @intCast(self.governor.interval_ms)); self.mutex.unlock(io); return; } const excessive = delay > self.governor.max_wait_ms; self.mutex.unlock(io); if (excessive) return error.RateLimitWaitTooLong; try io.sleep(std.Io.Duration.fromMilliseconds(delay), .awake); } }
pub fn observe(self: *Shared, io: std.Io, headers: Headers, throttled: bool, retry: u8) void { self.mutex.lockUncancelable(io); defer self.mutex.unlock(io); self.governor.observe(headers, throttled, retry, std.Io.Clock.awake.now(io).toMilliseconds(), std.Io.Clock.real.now(io).toSeconds()); }};pub const Governor = struct { interval_ms: u64 = 0, next_ms: i64 = 0, max_wait_ms: u64 = 30_000,
pub fn setRate(self: *Governor, requests_per_second: f64) !void { if (!std.math.isFinite(requests_per_second) or requests_per_second <= 0 or requests_per_second > 1000) return error.InvalidRequestRate; self.interval_ms = @intFromFloat(@ceil(1000 / requests_per_second)); }
pub fn wait(self: *Governor, io: std.Io) !void { const now = std.Io.Clock.awake.now(io).toMilliseconds(); const delay: u64 = @intCast(@max(0, self.next_ms - now)); if (delay > self.max_wait_ms) return error.RateLimitWaitTooLong; if (delay > 0) try io.sleep(std.Io.Duration.fromMilliseconds(@intCast(delay)), .awake); self.next_ms = std.Io.Clock.awake.now(io).toMilliseconds() + @as(i64, @intCast(self.interval_ms)); }
pub fn observe(self: *Governor, headers: Headers, throttled: bool, retry: u8, now_ms: i64, wall_s: i64) void { const delay = self.delayMillis(headers, throttled, retry, wall_s); const bounded: i64 = @intCast(@min(delay, std.math.maxInt(i64))); self.next_ms = @max(self.next_ms, now_ms +| bounded); }
pub fn delayMillis(self: Governor, headers: Headers, throttled: bool, retry: u8, wall_s: i64) u64 { const wait_s = headers.delaySeconds(wall_s); if (throttled or headers.remaining == 0 or headers.retry_after != null or headers.retry_after_at != null) { const fallback = @as(u64, 1000) << @as(u6, @intCast(@min(retry, 5))); if (wait_s) |seconds| return @max(@max(self.interval_ms, 250), seconds *| 1000); return @max(self.interval_ms, fallback); } if (headers.remaining) |remaining| { if (remaining > 0) { if (wait_s) |seconds| { const window_ms = seconds *| 1000; const spaced = window_ms / remaining; return @max(self.interval_ms, spaced +| (spaced / 9)); } } } return self.interval_ms; }};
test "server cooldown is never shortened to the local wait budget" { const governor: Governor = .{}; try std.testing.expectEqual(@as(u64, 120_000), governor.delayMillis(.{ .retry_after = 120 }, true, 0, 0)); try std.testing.expectEqual(@as(u64, 120_000), governor.delayMillis(.{ .reset = 1_800_000_120, .remaining = 0 }, false, 0, 1_800_000_000));}
test "remaining budget paces successful requests and unknown budgets remain bounded" { const governor: Governor = .{}; try std.testing.expectEqual(@as(u64, 1111), governor.delayMillis(.{ .remaining = 10, .reset = 10 }, false, 0, 0)); try std.testing.expectEqual(@as(u64, 0), governor.delayMillis(.{}, false, 0, 0)); try std.testing.expectEqual(@as(u64, 4000), governor.delayMillis(.{}, true, 2, 0));}
test "HTTP-date Retry-After and delta reset use transport parsing" { const head = try std.http.Client.Response.Head.parse("HTTP/1.1 429 Too Many Requests\r\nRetry-After: Sun, 06 Nov 1994 08:49:37 GMT\r\n\r\n"); const headers = Headers.fromResponseHead(head); try std.testing.expectEqual(@as(u64, 10_000), (Governor{}).delayMillis(headers, true, 0, 784_111_767));}
test "independent clients do not spend one another's budgets" { var one: Governor = .{}; const two: Governor = .{}; one.observe(.{ .retry_after = 60 }, true, 0, 100, 0); try std.testing.expectEqual(@as(i64, 60_100), one.next_ms); try std.testing.expectEqual(@as(i64, 0), two.next_ms);}
test "concurrent callers share request spacing" { const io = std.testing.io; var shared: Shared = .{ .governor = .{ .interval_ms = 20 } }; var times: [3]i64 = undefined; var group: std.Io.Group = .init; defer group.cancel(io); const Probe = struct { fn run(governor: *Shared, clock: std.Io, result: *i64) std.Io.Cancelable!void { governor.wait(clock) catch |err| switch (err) { error.Canceled => return error.Canceled, else => unreachable, }; result.* = std.Io.Clock.awake.now(clock).toMilliseconds(); } }; for (×) |*result| try group.concurrent(io, Probe.run, .{ &shared, io, result }); try group.await(io); std.mem.sort(i64, ×, {}, std.sort.asc(i64)); try std.testing.expect(times[1] - times[0] >= 19); try std.testing.expect(times[2] - times[1] >= 19);}