From e513358c5782a24471dc37c1b4c773120f2c9ec9 Mon Sep 17 00:00:00 2001 From: Altagos Date: Sat, 8 Nov 2025 17:59:36 +0100 Subject: [PATCH] create rate-limited http client for spacetraders --- src/st/http.zig | 302 ++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 302 insertions(+) create mode 100644 src/st/http.zig diff --git a/src/st/http.zig b/src/st/http.zig new file mode 100644 index 0000000..d44bfec --- /dev/null +++ b/src/st/http.zig @@ -0,0 +1,302 @@ +const std = @import("std"); +const HTTPClient = std.http.Client; +const Io = std.Io; +const json = std.json; + +const models = @import("models.zig"); + +const TIME_SLEEP_FACTOR: f64 = 0.98; + +const log = std.log.scoped(.SpaceTraders); + +const Semaphore = struct { + mutex: Io.Mutex = .{ .state = .unlocked }, + cond: Io.Condition = .{}, + /// It is OK to initialise this field to any value. + permits: u64 = 0, + + pub fn post(sem: *Semaphore, io: Io) void { + sem.mutex.lockUncancelable(io); + defer sem.mutex.unlock(io); + + sem.permits += 1; + sem.cond.signal(io); + } + + pub fn set(sem: *Semaphore, io: Io, permits: u64) void { + sem.mutex.lockUncancelable(io); + defer sem.mutex.unlock(io); + + sem.permits = permits; + sem.cond.signal(io); + } + + pub fn wait(sem: *Semaphore, io: Io) !void { + sem.mutex.lockUncancelable(io); + defer sem.mutex.unlock(io); + + while (sem.permits == 0) + try sem.cond.wait(io, &sem.mutex); + + sem.permits -= 1; + if (sem.permits > 0) + sem.cond.signal(io); + } + + pub fn available(sem: *Semaphore, io: Io) u64 { + sem.mutex.lockUncancelable(io); + defer sem.mutex.unlock(io); + + return sem.permits; + } +}; + +pub const Limiter = struct { + points: u64, + duration: i64, + time: ?Io.Timestamp, + + mutex: Io.Mutex = .{ .state = .unlocked }, + semaphor: Semaphore, + + pub fn init(opts: struct { points: u54 = 2, duration: i64 = 1000 }) Limiter { + return .{ + .points = opts.points, + .duration = opts.duration, + .time = null, + .semaphor = Semaphore{ .permits = opts.points }, + }; + } + + pub fn checkReset(l: *Limiter, io: Io) bool { + l.mutex.lock(io) catch return false; + defer l.mutex.unlock(io); + + if (l.time) |t| { + const dur = t.durationTo(Io.Clock.now(.real, io) catch return false); + if (dur.toSeconds() > 0) { + l.semaphor.set(io, l.points); + l.time = null; + return true; + } + } + + return false; + } + + pub fn aquire(l: *Limiter, io: Io) !void { + try l.mutex.lock(io); + defer l.mutex.unlock(io); + + if (l.time == null) { + const now = try Io.Clock.now(.real, io); + l.time = now.addDuration(.fromMilliseconds(l.duration)); + } + + return l.semaphor.wait(io); + } + + pub fn timeToReset(l: *Limiter, io: Io) i64 { + if (l.time) |t| { + return t.durationTo(Io.Clock.now(.real, io) catch return 0).raw.toMilliseconds(); + } + return 0; + } + + pub fn available(l: *Limiter, io: Io) bool { + return l.semaphor.available(io) > 0; + } +}; + +pub const BurstyLimiter = struct { + static: Limiter, + burst: Limiter, + + pub fn wait(bl: *BurstyLimiter, io: Io) !bool { + _ = bl.static.checkReset(io); + _ = bl.burst.checkReset(io); + + if (!bl.static.available(io)) { + if (bl.burst.available(io)) { + log.debug("Using Burst", .{}); + try bl.burst.aquire(io); + return true; + } else { + log.warn("No request available, waiting", .{}); + } + } + + try bl.static.aquire(io); + return true; + } +}; + +pub const AuthType = enum { account, agent, none }; + +pub const Auth = struct { + account: []const u8 = "", + agent: []const u8 = "", +}; + +pub const RequestOptions = struct { + method: std.http.Method = .GET, + auth: AuthType = .none, + body: Body = .empty, + + pub const Body = union(enum) { + empty: void, + buffer: []u8, + }; + + pub fn authorization(opts: *const RequestOptions, client: *const Client) HTTPClient.Request.Headers.Value { + switch (opts.auth) { + .account => return .{ .override = client.auth.account }, + .agent => return .{ .override = client.auth.agent }, + .none => return .{ .omit = {} }, + } + } +}; + +pub const RequestError = error{ + OutOfMemory, + InvalidResponse, + RateLimiterError, +} || HTTPClient.RequestError || HTTPClient.Request.ReceiveHeadError || std.Uri.ParseError; + +pub fn RawResponse(comptime T: type) type { + return Io.Future(RequestError!json.Parsed(T)); +} + +pub fn Response(comptime T: type) type { + return RawResponse(models.Wrapper(T)); +} + +pub const Client = struct { + allocator: std.mem.Allocator, + io: Io, + limiter: BurstyLimiter, + + base_url: []const u8, + auth: Auth, + + http: HTTPClient, + + pub fn init( + allocator: std.mem.Allocator, + io: std.Io, + opts: struct { + base_url: []const u8 = "https://api.spacetraders.io/v2", + auth: Auth = .{}, + }, + ) Client { + return .{ + .allocator = allocator, + .io = io, + .limiter = .{ + .static = .init(.{}), + .burst = .init(.{ .points = 30, .duration = 60_000 }), + }, + .base_url = opts.base_url, + .auth = opts.auth, + .http = .{ .allocator = allocator, .io = io }, + }; + } + + pub fn deinit(client: *Client) void { + client.http.deinit(); + } + + pub fn request( + client: *Client, + comptime T: type, + comptime path: []const u8, + args: anytype, + opts: RequestOptions, + ) !RawResponse(T) { + const path_fmt = try std.fmt.allocPrint(client.allocator, path, args); + defer client.allocator.free(path_fmt); + + const url = try std.fmt.allocPrint(client.allocator, "{s}{s}", .{ client.base_url, path_fmt }); + + const Wrapper = struct { + fn call( + cl: *Client, + url_param: []const u8, + opts_param: *const RequestOptions, + ) RequestError!json.Parsed(T) { + defer cl.allocator.free(url_param); + if (cl.limiter.wait(cl.io) catch return error.RateLimiterError) + return Client.requestRaw(cl, T, url_param, opts_param); + return error.RateLimiterError; + } + }; + + return client.io.concurrent( + Wrapper.call, + .{ client, url, &opts }, + ); + } + + pub fn requestRaw( + client: *Client, + comptime T: type, + url: []const u8, + opts: *const RequestOptions, + ) RequestError!json.Parsed(T) { + const uri = std.Uri.parse(url) catch |err| { + log.err("Error parsing url: {} - url = {s}", .{ err, url }); + return err; + }; + + var req = try client.http.request(opts.method, uri, .{ + .headers = .{ + .authorization = opts.authorization(client), + .user_agent = .{ .override = "All your codebases are belong to us" }, + }, + }); + defer req.deinit(); + + log.debug("requesting: {s}", .{uri.path.percent_encoded}); + + switch (opts.body) { + .empty => try req.sendBodiless(), + .buffer => |body| try req.sendBodyComplete(body), + } + + var redirect_buffer: [1024]u8 = undefined; + + var response = try req.receiveHead(&redirect_buffer); + const colour = blk: { + if (std.mem.eql(u8, response.head.reason, "OK")) { + break :blk "\x1b[92m"; + } else { + break :blk "\x1b[1m\x1b[91m"; + } + }; + log.debug( + "\x1b[2m[path = {s}]\x1b[0m received {s}{d} {s}\x1b[0m", + .{ url[client.base_url.len..], colour, response.head.status, response.head.reason }, + ); + + // var header_iter = response.head.iterateHeaders(); + // while (header_iter.next()) |header| { + // log.debug("{s}: {s}", .{ header.name, header.value }); + // } + + var decompress_buffer: [std.compress.flate.max_window_len]u8 = undefined; + var transfer_buffer: [64]u8 = undefined; + var decompress: std.http.Decompress = undefined; + + const decompressed_body_reader = response.readerDecompressing(&transfer_buffer, &decompress, &decompress_buffer); + + var json_reader: json.Reader = .init(client.allocator, decompressed_body_reader); + defer json_reader.deinit(); + + return json.parseFromTokenSource(T, client.allocator, &json_reader, .{ + .ignore_unknown_fields = true, + }) catch |err| { + log.err("Error parsing response: {}", .{err}); + return RequestError.InvalidResponse; + }; + } +}; -- 2.51.2