atproto pds in zig
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256const std = @import("std");const clock = @import("../core/clock.zig");const httpz = @import("httpz");const router = @import("../http/router.zig");
const Io = std.Io;const Route = router.Route;
const SAMPLE_COUNT = 512;const SLOW_COUNT = 32;const LABEL_LEN = 192;const ROUTE_COUNT = @typeInfo(Route).@"enum".fields.len;const METHOD_COUNT = @typeInfo(Method).@"enum".fields.len;
pub const Method = enum { get, head, post, options, other,
pub fn fromHttpz(method: httpz.Method) Method { return switch (method) { .GET => .get, .HEAD => .head, .POST => .post, .OPTIONS => .options, else => .other, }; }
pub fn label(method: Method) []const u8 { return switch (method) { .get => "GET", .head => "HEAD", .post => "POST", .options => "OPTIONS", .other => "OTHER", }; }};
pub const Class = enum { local, proxy, unknown,
pub fn label(class: Class) []const u8 { return switch (class) { .local => "local", .proxy => "proxy", .unknown => "unknown", }; }};
pub const EndpointStats = struct { route: Route, method: Method, count: u64, failure_count: u64, avg_ms: f64, p50_ms: f64, p95_ms: f64, max_ms: f64,};
pub const SlowRequest = struct { route: Route = .not_found, method: Method = .other, class: Class = .unknown, at_s: i64 = 0, elapsed_us: u32 = 0, status: u16 = 0, handler_failed: bool = false, label: [LABEL_LEN]u8 = [_]u8{0} ** LABEL_LEN, label_len: usize = 0,
pub fn labelText(self: *const SlowRequest) []const u8 { return self.label[0..self.label_len]; }
pub fn failed(self: *const SlowRequest) bool { return self.handler_failed or self.status >= 500; }};
pub const Snapshot = struct { started_at_s: i64, uptime_s: u64, total_requests: u64, total_failures: u64, routes: [ROUTE_COUNT][METHOD_COUNT]EndpointStats, slow_requests: [SLOW_COUNT]SlowRequest, slow_len: usize,};
const LatencyBuffer = struct { samples: [SAMPLE_COUNT]u32 = .{0} ** SAMPLE_COUNT, count: usize = 0, head: usize = 0, total_count: u64 = 0, failure_count: u64 = 0,
fn record(self: *LatencyBuffer, elapsed_us: u32, failed: bool) void { self.samples[self.head] = elapsed_us; self.head = (self.head + 1) % SAMPLE_COUNT; if (self.count < SAMPLE_COUNT) self.count += 1; self.total_count += 1; if (failed) self.failure_count += 1; }};
var global_io: ?Io = null;var mutex: Io.Mutex = .init;var started_at_s: i64 = 0;var buffers = [_][METHOD_COUNT]LatencyBuffer{[_]LatencyBuffer{.{}} ** METHOD_COUNT} ** ROUTE_COUNT;var slow_requests = [_]SlowRequest{.{}} ** SLOW_COUNT;var slow_next: usize = 0;var slow_len: usize = 0;var total_requests: u64 = 0;var total_failures: u64 = 0;
pub fn init(io: Io) void { global_io = io; started_at_s = clock.now();}
pub fn start() i64 { const io = global_io orelse return 0; return Io.Timestamp.now(io, .awake).toMicroseconds();}
pub fn record(route: Route, method: httpz.Method, class: Class, label: []const u8, start_us: i64, status: u16, handler_failed: bool) void { const io = global_io orelse return; const now_us = Io.Timestamp.now(io, .awake).toMicroseconds(); const elapsed_us: u32 = @intCast(@min(@max(0, now_us - start_us), std.math.maxInt(u32))); const method_key = Method.fromHttpz(method); const failed = handler_failed or status >= 500;
mutex.lockUncancelable(io); defer mutex.unlock(io);
buffers[@intFromEnum(route)][@intFromEnum(method_key)].record(elapsed_us, failed); total_requests += 1; if (failed) total_failures += 1;
if (elapsed_us >= 250 * std.time.us_per_ms or failed) { var slow: SlowRequest = .{ .route = route, .method = method_key, .class = class, .at_s = clock.now(), .elapsed_us = elapsed_us, .status = status, .handler_failed = handler_failed, }; slow.label_len = @min(label.len, LABEL_LEN); @memcpy(slow.label[0..slow.label_len], label[0..slow.label_len]); slow_requests[slow_next] = slow; slow_next = (slow_next + 1) % SLOW_COUNT; if (slow_len < SLOW_COUNT) slow_len += 1; }}
pub fn snapshot() Snapshot { const io = global_io orelse return emptySnapshot(); mutex.lockUncancelable(io); defer mutex.unlock(io);
var routes: [ROUTE_COUNT][METHOD_COUNT]EndpointStats = undefined; for (&routes, 0..) |*method_stats, route_idx| { for (method_stats, 0..) |*stats, method_idx| { stats.* = summarize(@enumFromInt(route_idx), @enumFromInt(method_idx), buffers[route_idx][method_idx]); } }
var slow: [SLOW_COUNT]SlowRequest = [_]SlowRequest{.{}} ** SLOW_COUNT; for (0..slow_len) |i| { const idx = (slow_next + SLOW_COUNT - slow_len + i) % SLOW_COUNT; slow[i] = slow_requests[idx]; }
const now_s = clock.now(); return .{ .started_at_s = started_at_s, .uptime_s = @intCast(@max(0, now_s - started_at_s)), .total_requests = total_requests, .total_failures = total_failures, .routes = routes, .slow_requests = slow, .slow_len = slow_len, };}
fn emptySnapshot() Snapshot { var routes: [ROUTE_COUNT][METHOD_COUNT]EndpointStats = undefined; for (&routes, 0..) |*method_stats, route_idx| { for (method_stats, 0..) |*stats, method_idx| { stats.* = .{ .route = @enumFromInt(route_idx), .method = @enumFromInt(method_idx), .count = 0, .failure_count = 0, .avg_ms = 0, .p50_ms = 0, .p95_ms = 0, .max_ms = 0, }; } } return .{ .started_at_s = 0, .uptime_s = 0, .total_requests = 0, .total_failures = 0, .routes = routes, .slow_requests = [_]SlowRequest{.{}} ** SLOW_COUNT, .slow_len = 0, };}
fn summarize(route: Route, method: Method, buffer: LatencyBuffer) EndpointStats { if (buffer.count == 0) { return .{ .route = route, .method = method, .count = 0, .failure_count = buffer.failure_count, .avg_ms = 0, .p50_ms = 0, .p95_ms = 0, .max_ms = 0, }; }
var sorted: [SAMPLE_COUNT]u32 = undefined; @memcpy(sorted[0..buffer.count], buffer.samples[0..buffer.count]); std.mem.sort(u32, sorted[0..buffer.count], {}, std.sort.asc(u32));
var sum: u64 = 0; for (sorted[0..buffer.count]) |value| sum += value; const p95_idx = @min(buffer.count - 1, (buffer.count * 95) / 100);
return .{ .route = route, .method = method, .count = buffer.total_count, .failure_count = buffer.failure_count, .avg_ms = @as(f64, @floatFromInt(sum)) / @as(f64, @floatFromInt(buffer.count)) / 1000.0, .p50_ms = @as(f64, @floatFromInt(sorted[buffer.count / 2])) / 1000.0, .p95_ms = @as(f64, @floatFromInt(sorted[p95_idx])) / 1000.0, .max_ms = @as(f64, @floatFromInt(sorted[buffer.count - 1])) / 1000.0, };}