quit from interactive mode means quit
This commit is contained in:
parent
4709403eca
commit
2876358c5a
12 changed files with 458 additions and 94 deletions
|
|
@ -1187,7 +1187,7 @@ fn dividendsWorker(self: *PortfolioData, delay_ms: usize) void {
|
|||
const div_syms = arena_alloc.alloc([]const u8, summary_ref.allocations.len) catch break :warm;
|
||||
for (summary_ref.allocations, div_syms) |alloc, *dst| dst.* = alloc.symbol;
|
||||
self.io.checkCancel() catch return;
|
||||
self.svc.loadAllDividends(div_syms, self.fetch_options);
|
||||
self.svc.loadAllDividends(div_syms, self.fetch_options) catch return;
|
||||
}
|
||||
|
||||
var map = std.StringHashMap([]const Dividend).init(arena_alloc);
|
||||
|
|
|
|||
|
|
@ -212,7 +212,7 @@ pub fn run(ctx: *framework.RunCtx, parsed: ParsedArgs) !void {
|
|||
);
|
||||
defer div_syms.deinit(allocator);
|
||||
for (pf_data.summary.allocations) |a| div_syms.appendAssumeCapacity(a.symbol);
|
||||
svc.loadAllDividends(div_syms.items, cli.fetchOptionsFromPolicy(ctx.globals.refresh_policy));
|
||||
try svc.loadAllDividends(div_syms.items, cli.fetchOptionsFromPolicy(ctx.globals.refresh_policy));
|
||||
}
|
||||
for (pf_data.summary.allocations) |a| {
|
||||
if (svc.getCachedDividends(allocator, a.symbol)) |divs| {
|
||||
|
|
|
|||
|
|
@ -75,28 +75,27 @@ pub fn tryAcquire(self: *RateLimiter) bool {
|
|||
}
|
||||
|
||||
/// Acquire a token, blocking (sleeping) until one is available.
|
||||
pub fn acquire(self: *RateLimiter) void {
|
||||
///
|
||||
/// Returns `error.Canceled` when the calling task is canceled while
|
||||
/// waiting. Propagate it: std delivers a cancelation to exactly one
|
||||
/// cancelation point, so a caller that swallows it (this function used
|
||||
/// to loop back on an interrupted sleep) runs on uncancelable, and on a
|
||||
/// 4-per-minute bucket that turns a cancel into minutes of blocking.
|
||||
pub fn acquire(self: *RateLimiter) std.Io.Cancelable!void {
|
||||
while (!self.tryAcquire()) {
|
||||
// Sleep for the time needed to generate 1 token. An
|
||||
// interrupted sleep (cancelation propagating through the
|
||||
// Io) just loops back to tryAcquire - the next refill
|
||||
// covers whatever fraction of the wait elapsed.
|
||||
// Sleep for the time needed to generate 1 token.
|
||||
const wait_ns: u64 = @intFromFloat(1.0 / self.refill_rate_per_ns);
|
||||
std.Io.sleep(self.io, .{ .nanoseconds = @intCast(wait_ns) }, .awake) catch |err| {
|
||||
std.log.scoped(.rate_limiter).debug("acquire sleep interrupted: {t}", .{err});
|
||||
};
|
||||
try std.Io.sleep(self.io, .{ .nanoseconds = @intCast(wait_ns) }, .awake);
|
||||
}
|
||||
}
|
||||
|
||||
/// Sleep until a token is likely available, with a minimum 2-second floor.
|
||||
/// Use after receiving a server-side 429 to wait before retrying.
|
||||
pub fn backoff(self: *RateLimiter) void {
|
||||
/// Returns `error.Canceled` if the task is canceled during the wait; see
|
||||
/// `acquire`.
|
||||
pub fn backoff(self: *RateLimiter) std.Io.Cancelable!void {
|
||||
const wait_ns: u64 = @max(self.estimateWaitNs(), 2 * std.time.ns_per_s);
|
||||
// Interrupted backoff sleep degrades to a shorter wait; the
|
||||
// caller's retry may hit 429 again and re-backoff.
|
||||
std.Io.sleep(self.io, .{ .nanoseconds = @intCast(wait_ns) }, .awake) catch |err| {
|
||||
std.log.scoped(.rate_limiter).debug("backoff sleep interrupted: {t}", .{err});
|
||||
};
|
||||
try std.Io.sleep(self.io, .{ .nanoseconds = @intCast(wait_ns) }, .awake);
|
||||
}
|
||||
|
||||
/// Returns estimated wait time in nanoseconds until a token is available.
|
||||
|
|
@ -153,3 +152,26 @@ test "rate limiter exhaustion" {
|
|||
// Bucket should be empty now
|
||||
try std.testing.expect(!rl.tryAcquire());
|
||||
}
|
||||
|
||||
test "acquire: a cancel while waiting for a token returns Canceled" {
|
||||
// Used to loop back on the interrupted sleep, which made the cancel
|
||||
// a no-op: on a 1-per-minute bucket the waiter then blocked for the
|
||||
// full refill.
|
||||
const io = std.testing.io;
|
||||
var rl = RateLimiter.perMinute(io, 1);
|
||||
try rl.acquire(); // spend the only token
|
||||
|
||||
var waiter = try io.concurrent(RateLimiter.acquire, .{&rl});
|
||||
try io.sleep(.fromMilliseconds(20), .awake);
|
||||
try std.testing.expectError(error.Canceled, waiter.cancel(io));
|
||||
}
|
||||
|
||||
test "backoff: a cancel during the wait returns Canceled" {
|
||||
const io = std.testing.io;
|
||||
var rl = RateLimiter.perMinute(io, 1);
|
||||
try rl.acquire();
|
||||
|
||||
var waiter = try io.concurrent(RateLimiter.backoff, .{&rl});
|
||||
try io.sleep(.fromMilliseconds(20), .awake);
|
||||
try std.testing.expectError(error.Canceled, waiter.cancel(io));
|
||||
}
|
||||
|
|
|
|||
206
src/net/http.zig
206
src/net/http.zig
|
|
@ -186,19 +186,27 @@ pub const Client = struct {
|
|||
// "RequestFailed" because the per-attempt error was discarded
|
||||
// by `catch {}`. Now the caller's `@errorName(err)` reports
|
||||
// the real cause (e.g., `NoAddressReturned`).
|
||||
//
|
||||
// `error.Canceled` is never retried. It is the calling task
|
||||
// being stopped, not a transport failure, and std delivers a
|
||||
// cancelation to exactly one cancelation point: a retry would
|
||||
// run uncancelable and could block indefinitely on a server
|
||||
// that never answers. That is how quitting the TUI used to
|
||||
// wait out a whole dividend warm-up.
|
||||
var attempt: u8 = 0;
|
||||
var last_err: HttpError = HttpError.RequestFailed;
|
||||
while (true) : (attempt += 1) {
|
||||
const response = self.doRequest(method, url, body, extra_headers) catch |err| {
|
||||
if (err == error.Canceled) return err;
|
||||
last_err = err;
|
||||
if (attempt >= self.max_retries) return last_err;
|
||||
self.backoffSleep(attempt);
|
||||
try self.backoffSleep(attempt);
|
||||
continue;
|
||||
};
|
||||
return classifyResponse(response) catch |err| {
|
||||
if (err == HttpError.ServerError and attempt < self.max_retries) {
|
||||
last_err = err;
|
||||
self.backoffSleep(attempt);
|
||||
try self.backoffSleep(attempt);
|
||||
continue;
|
||||
}
|
||||
return err;
|
||||
|
|
@ -206,9 +214,9 @@ pub const Client = struct {
|
|||
}
|
||||
}
|
||||
|
||||
fn backoffSleep(self: *Client, attempt: u8) void {
|
||||
fn backoffSleep(self: *Client, attempt: u8) std.Io.Cancelable!void {
|
||||
const backoff = self.base_backoff_ms * std.math.shl(u64, 1, attempt);
|
||||
std.Io.sleep(self.io, std.Io.Duration.fromMilliseconds(@intCast(backoff)), .awake) catch |err| std.log.debug("backoff sleep interrupted: {t}", .{err});
|
||||
try std.Io.sleep(self.io, std.Io.Duration.fromMilliseconds(@intCast(backoff)), .awake);
|
||||
}
|
||||
|
||||
fn doRequest(self: *Client, method: std.http.Method, url: []const u8, body: ?[]const u8, extra_headers: []const std.http.Header) HttpError!Response {
|
||||
|
|
@ -240,14 +248,6 @@ pub const Client = struct {
|
|||
// doesn't produce nonsense elapsed values.
|
||||
const t_start = std.Io.Timestamp.now(self.io, .awake).nanoseconds;
|
||||
var t_stage = t_start;
|
||||
const stageElapsedMs = struct {
|
||||
fn f(prev: *i96, io: std.Io) i64 {
|
||||
const now = std.Io.Timestamp.now(io, .awake).nanoseconds;
|
||||
const delta_ns = now - prev.*;
|
||||
prev.* = now;
|
||||
return @intCast(@divTrunc(delta_ns, std.time.ns_per_ms));
|
||||
}
|
||||
}.f;
|
||||
|
||||
const uri = std.Uri.parse(url) catch |err| {
|
||||
log.warn("http {s}: stage=uri_parse err={s} url={f}", .{ @tagName(method), @errorName(err), redactUrl(url) });
|
||||
|
|
@ -284,8 +284,7 @@ pub const Client = struct {
|
|||
// TLS handshake. Logging at warn level (rather than debug)
|
||||
// because DNS / connectivity failures are exactly what
|
||||
// operators need to see immediately.
|
||||
log.warn("http {s}: stage=connect err={s} elapsed_ms={d} url={f}", .{ @tagName(method), @errorName(err), stageElapsedMs(&t_stage, self.io), redactUrl(url) });
|
||||
return err;
|
||||
return self.stageFailed(method, "connect", err, null, &t_stage, url);
|
||||
};
|
||||
defer req.deinit();
|
||||
const ms_connect = stageElapsedMs(&t_stage, self.io);
|
||||
|
|
@ -293,36 +292,18 @@ pub const Client = struct {
|
|||
if (body) |payload| {
|
||||
var send_buf: [4096]u8 = undefined;
|
||||
req.transfer_encoding = .{ .content_length = payload.len };
|
||||
var bw = req.sendBodyUnflushed(&send_buf) catch |err| {
|
||||
log.warn("http {s}: stage=send_body_open err={s} elapsed_ms={d} url={f}", .{ @tagName(method), @errorName(err), stageElapsedMs(&t_stage, self.io), redactUrl(url) });
|
||||
return err;
|
||||
};
|
||||
bw.writer.writeAll(payload) catch |err| {
|
||||
log.warn("http {s}: stage=send_body_write err={s} elapsed_ms={d} url={f}", .{ @tagName(method), @errorName(err), stageElapsedMs(&t_stage, self.io), redactUrl(url) });
|
||||
return err;
|
||||
};
|
||||
bw.end() catch |err| {
|
||||
log.warn("http {s}: stage=send_body_end err={s} elapsed_ms={d} url={f}", .{ @tagName(method), @errorName(err), stageElapsedMs(&t_stage, self.io), redactUrl(url) });
|
||||
return err;
|
||||
};
|
||||
req.connection.?.flush() catch |err| {
|
||||
log.warn("http {s}: stage=send_body_flush err={s} elapsed_ms={d} url={f}", .{ @tagName(method), @errorName(err), stageElapsedMs(&t_stage, self.io), redactUrl(url) });
|
||||
return err;
|
||||
};
|
||||
var bw = req.sendBodyUnflushed(&send_buf) catch |err| return self.stageFailed(method, "send_body_open", err, &req, &t_stage, url);
|
||||
bw.writer.writeAll(payload) catch |err| return self.stageFailed(method, "send_body_write", err, &req, &t_stage, url);
|
||||
bw.end() catch |err| return self.stageFailed(method, "send_body_end", err, &req, &t_stage, url);
|
||||
req.connection.?.flush() catch |err| return self.stageFailed(method, "send_body_flush", err, &req, &t_stage, url);
|
||||
} else {
|
||||
req.sendBodiless() catch |err| {
|
||||
log.warn("http {s}: stage=send_bodiless err={s} elapsed_ms={d} url={f}", .{ @tagName(method), @errorName(err), stageElapsedMs(&t_stage, self.io), redactUrl(url) });
|
||||
return err;
|
||||
};
|
||||
req.sendBodiless() catch |err| return self.stageFailed(method, "send_bodiless", err, &req, &t_stage, url);
|
||||
}
|
||||
const ms_send = stageElapsedMs(&t_stage, self.io);
|
||||
|
||||
// Matches the default redirect capacity in std.http.Client.fetch.
|
||||
var redirect_buffer: [8 * 1024]u8 = undefined;
|
||||
var response = req.receiveHead(&redirect_buffer) catch |err| {
|
||||
log.warn("http {s}: stage=receive_head err={s} elapsed_ms={d} url={f}", .{ @tagName(method), @errorName(err), stageElapsedMs(&t_stage, self.io), redactUrl(url) });
|
||||
return err;
|
||||
};
|
||||
var response = req.receiveHead(&redirect_buffer) catch |err| return self.stageFailed(method, "receive_head", err, &req, &t_stage, url);
|
||||
const ms_receive_head = stageElapsedMs(&t_stage, self.io);
|
||||
|
||||
// Capture the ETag (if any) from the response head BEFORE
|
||||
|
|
@ -358,10 +339,7 @@ pub const Client = struct {
|
|||
var decompress: std.http.Decompress = undefined;
|
||||
var decompress_buffer: [64 * 1024]u8 = undefined;
|
||||
const reader = response.readerDecompressing(&transfer_buffer, &decompress, &decompress_buffer);
|
||||
_ = reader.streamRemaining(&aw.writer) catch |err| {
|
||||
log.warn("http {s}: stage=stream_body err={s} elapsed_ms={d} url={f}", .{ @tagName(method), @errorName(err), stageElapsedMs(&t_stage, self.io), redactUrl(url) });
|
||||
return err;
|
||||
};
|
||||
_ = reader.streamRemaining(&aw.writer) catch |err| return self.stageFailed(method, "stream_body", err, &req, &t_stage, url);
|
||||
const ms_body = stageElapsedMs(&t_stage, self.io);
|
||||
|
||||
const resp_body = try aw.toOwnedSlice();
|
||||
|
|
@ -391,6 +369,40 @@ pub const Client = struct {
|
|||
};
|
||||
}
|
||||
|
||||
/// Report a failed request stage: log it and return the error the
|
||||
/// caller should see.
|
||||
///
|
||||
/// std's reader and writer interfaces turn every transport failure
|
||||
/// into a generic `ReadFailed` / `WriteFailed` and leave the real
|
||||
/// cause on the connection's socket stream. When that cause is
|
||||
/// `error.Canceled`, return it instead: callers must be able to tell
|
||||
/// "the task was stopped" from "the network failed", and the retry
|
||||
/// loop in `request` must not retry the former (see there).
|
||||
///
|
||||
/// A cancel is logged at debug level. It is the expected result of
|
||||
/// quitting mid-request, not something an operator needs to see.
|
||||
/// Other failures log at warn, except under `zig build test`, where
|
||||
/// the transport tests drive them on purpose (as `classifyResponse`
|
||||
/// does for statuses).
|
||||
fn stageFailed(
|
||||
self: *Client,
|
||||
method: std.http.Method,
|
||||
stage: []const u8,
|
||||
raw: HttpError,
|
||||
req: ?*const std.http.Client.Request,
|
||||
t_stage: *i96,
|
||||
url: []const u8,
|
||||
) HttpError {
|
||||
const err = if (req) |r| unmaskCancel(raw, r) else raw;
|
||||
const elapsed_ms = stageElapsedMs(t_stage, self.io);
|
||||
if (err == error.Canceled) {
|
||||
log.debug("http {s}: stage={s} canceled elapsed_ms={d} url={f}", .{ @tagName(method), stage, elapsed_ms, redactUrl(url) });
|
||||
} else if (!builtin.is_test) {
|
||||
log.warn("http {s}: stage={s} err={s} elapsed_ms={d} url={f}", .{ @tagName(method), stage, @errorName(err), elapsed_ms, redactUrl(url) });
|
||||
}
|
||||
return err;
|
||||
}
|
||||
|
||||
fn classifyResponse(response: Response) HttpError!Response {
|
||||
switch (response.status) {
|
||||
.ok => return response,
|
||||
|
|
@ -443,6 +455,37 @@ pub const Client = struct {
|
|||
}
|
||||
};
|
||||
|
||||
/// Milliseconds since `prev`, advancing `prev` to now. Per-stage
|
||||
/// request timing; `.awake` (monotonic) so a clock jump mid-request
|
||||
/// cannot produce a negative stage.
|
||||
fn stageElapsedMs(prev: *i96, io: std.Io) i64 {
|
||||
const now = std.Io.Timestamp.now(io, .awake).nanoseconds;
|
||||
const delta_ns = now - prev.*;
|
||||
prev.* = now;
|
||||
return @intCast(@divTrunc(delta_ns, std.time.ns_per_ms));
|
||||
}
|
||||
|
||||
/// `error.Canceled` when a generic `ReadFailed` / `WriteFailed` was a
|
||||
/// cancelation of the request's socket I/O, else `err` unchanged. See
|
||||
/// `Client.stageFailed`.
|
||||
///
|
||||
/// Reads the socket streams' recorded errors directly rather than
|
||||
/// through `Connection.getReadError`, which asserts one is set.
|
||||
fn unmaskCancel(err: HttpError, req: *const std.http.Client.Request) HttpError {
|
||||
switch (err) {
|
||||
error.ReadFailed, error.WriteFailed => {},
|
||||
else => return err,
|
||||
}
|
||||
const conn = req.connection orelse return err;
|
||||
if (conn.stream_reader.err) |e| {
|
||||
if (e == error.Canceled) return error.Canceled;
|
||||
}
|
||||
if (conn.stream_writer.err) |e| {
|
||||
if (e == error.Canceled) return error.Canceled;
|
||||
}
|
||||
return err;
|
||||
}
|
||||
|
||||
// ── URL redaction for logs ───────────────────────────────────
|
||||
//
|
||||
// Providers authenticate by query parameter, so a raw request URL is a
|
||||
|
|
@ -849,3 +892,82 @@ test "isCredentialParam: case-insensitive across provider spellings" {
|
|||
try std.testing.expect(!isCredentialParam("ticker"));
|
||||
try std.testing.expect(!isCredentialParam("startDate"));
|
||||
}
|
||||
|
||||
// ---- Cancelation ----
|
||||
|
||||
const SilentServer = @import("../testutil/silent_server.zig").SilentServer;
|
||||
|
||||
fn silentUrl(server: *const SilentServer, buf: []u8) ![]const u8 {
|
||||
var base_buf: [64]u8 = undefined;
|
||||
return std.fmt.bufPrint(buf, "{s}/SAMPLE/dividends", .{try server.baseUrl(&base_buf)});
|
||||
}
|
||||
|
||||
test "request: canceling a request blocked on an unresponsive server returns Canceled, with no retry" {
|
||||
// The bug this pins: std reports a canceled read as a generic
|
||||
// `ReadFailed`, the retry loop retried it, and the retry ran
|
||||
// uncancelable (std delivers a cancel only once) - so canceling a
|
||||
// request to a server that never answers blocked forever.
|
||||
const io = std.testing.io;
|
||||
var server: SilentServer = undefined;
|
||||
try server.start(io, .hold_open);
|
||||
defer server.stop();
|
||||
|
||||
var client = Client.init(io, std.testing.allocator);
|
||||
defer client.deinit();
|
||||
var url_buf: [96]u8 = undefined;
|
||||
const url = try silentUrl(&server, &url_buf);
|
||||
|
||||
var pending = try io.concurrent(Client.get, .{ &client, url });
|
||||
try server.waitForConnections(1, 5000);
|
||||
try std.testing.expectError(error.Canceled, pending.cancel(io));
|
||||
|
||||
// A retry would have connected again within its first backoff.
|
||||
try io.sleep(.fromMilliseconds(base_backoff_for_tests_ms * 4), .awake);
|
||||
try std.testing.expectEqual(@as(u32, 1), server.accepted.load(.acquire));
|
||||
}
|
||||
|
||||
test "request: a server that hangs up is still retried" {
|
||||
// The cancel check must not swallow real transport failures: a
|
||||
// closed connection is retried up to `max_retries` times.
|
||||
const io = std.testing.io;
|
||||
var server: SilentServer = undefined;
|
||||
try server.start(io, .hang_up);
|
||||
defer server.stop();
|
||||
|
||||
var client = Client.init(io, std.testing.allocator);
|
||||
defer client.deinit();
|
||||
client.base_backoff_ms = base_backoff_for_tests_ms;
|
||||
var url_buf: [96]u8 = undefined;
|
||||
const url = try silentUrl(&server, &url_buf);
|
||||
|
||||
if (client.get(url)) |resp| {
|
||||
var r = resp;
|
||||
r.deinit();
|
||||
return error.TestUnexpectedResult;
|
||||
} else |err| {
|
||||
try std.testing.expect(err != error.Canceled);
|
||||
}
|
||||
try std.testing.expectEqual(@as(u32, 1 + client.max_retries), server.accepted.load(.acquire));
|
||||
}
|
||||
|
||||
test "request: a cancel during the retry backoff returns Canceled" {
|
||||
const io = std.testing.io;
|
||||
var server: SilentServer = undefined;
|
||||
try server.start(io, .hang_up);
|
||||
defer server.stop();
|
||||
|
||||
var client = Client.init(io, std.testing.allocator);
|
||||
defer client.deinit();
|
||||
// Long enough that the cancel lands inside the first backoff sleep.
|
||||
client.base_backoff_ms = 60_000;
|
||||
var url_buf: [96]u8 = undefined;
|
||||
const url = try silentUrl(&server, &url_buf);
|
||||
|
||||
var pending = try io.concurrent(Client.get, .{ &client, url });
|
||||
try server.waitForConnections(1, 5000);
|
||||
try std.testing.expectError(error.Canceled, pending.cancel(io));
|
||||
try std.testing.expectEqual(@as(u32, 1), server.accepted.load(.acquire));
|
||||
}
|
||||
|
||||
/// Backoff for tests that exercise retries, so they take milliseconds.
|
||||
const base_backoff_for_tests_ms = 5;
|
||||
|
|
|
|||
|
|
@ -202,7 +202,7 @@ pub fn deinit(self: *Edgar) void {
|
|||
/// requires on every request and acquires a rate-limit token before
|
||||
/// issuing the call.
|
||||
fn httpGet(self: *Edgar, url: []const u8) !http.Response {
|
||||
self.rate_limiter.acquire();
|
||||
try self.rate_limiter.acquire();
|
||||
|
||||
var ua_buf: [256]u8 = undefined;
|
||||
const ua = std.fmt.bufPrint(&ua_buf, "zfin/0.1 ({s})", .{self.user_email}) catch return error.UserEmailTooLong;
|
||||
|
|
|
|||
|
|
@ -41,7 +41,7 @@ pub const Cboe = struct {
|
|||
allocator: std.mem.Allocator,
|
||||
symbol: []const u8,
|
||||
) ![]OptionsChain {
|
||||
self.rate_limiter.acquire();
|
||||
try self.rate_limiter.acquire();
|
||||
|
||||
// Build URL: {base_url}/{SYMBOL}.json
|
||||
const url = try buildCboeUrl(allocator, symbol);
|
||||
|
|
|
|||
|
|
@ -65,7 +65,7 @@ pub const Fmp = struct {
|
|||
allocator: std.mem.Allocator,
|
||||
symbol: []const u8,
|
||||
) ![]EarningsEvent {
|
||||
self.rate_limiter.acquire();
|
||||
try self.rate_limiter.acquire();
|
||||
|
||||
const url = try http.buildUrl(allocator, base_url ++ "/earnings", &.{
|
||||
.{ "symbol", symbol },
|
||||
|
|
|
|||
|
|
@ -56,7 +56,7 @@ pub const Polygon = struct {
|
|||
|
||||
// First request
|
||||
{
|
||||
self.rate_limiter.acquire();
|
||||
try self.rate_limiter.acquire();
|
||||
|
||||
var params: [5][2][]const u8 = undefined;
|
||||
var n: usize = 0;
|
||||
|
|
@ -94,7 +94,7 @@ pub const Polygon = struct {
|
|||
|
||||
// Paginate
|
||||
while (next_url) |cursor_url| {
|
||||
self.rate_limiter.acquire();
|
||||
try self.rate_limiter.acquire();
|
||||
|
||||
const authed = try appendApiKey(allocator, cursor_url, self.api_key);
|
||||
defer allocator.free(authed);
|
||||
|
|
@ -123,7 +123,7 @@ pub const Polygon = struct {
|
|||
allocator: std.mem.Allocator,
|
||||
symbol: []const u8,
|
||||
) ![]Split {
|
||||
self.rate_limiter.acquire();
|
||||
try self.rate_limiter.acquire();
|
||||
|
||||
const url = try http.buildUrl(allocator, base_url ++ "/v3/reference/splits", &.{
|
||||
.{ "ticker", symbol },
|
||||
|
|
|
|||
|
|
@ -153,7 +153,7 @@ pub const Tiingo = struct {
|
|||
|
||||
// Honor Tiingo's 50/hour free-tier cap. Blocks only once a
|
||||
// single run has spent its 50-token burst within the hour.
|
||||
self.rate_limiter.acquire();
|
||||
try self.rate_limiter.acquire();
|
||||
|
||||
var response = try self.client.get(url);
|
||||
defer response.deinit();
|
||||
|
|
@ -229,7 +229,7 @@ pub const Tiingo = struct {
|
|||
defer allocator.free(url);
|
||||
|
||||
// One batched call = one request against the hourly bucket.
|
||||
self.rate_limiter.acquire();
|
||||
try self.rate_limiter.acquire();
|
||||
|
||||
var response = try self.client.get(url);
|
||||
defer response.deinit();
|
||||
|
|
|
|||
|
|
@ -48,7 +48,7 @@ pub const TwelveData = struct {
|
|||
from: Date,
|
||||
to: Date,
|
||||
) ![]Candle {
|
||||
self.rate_limiter.acquire();
|
||||
try self.rate_limiter.acquire();
|
||||
|
||||
var from_buf: [10]u8 = undefined;
|
||||
var to_buf: [10]u8 = undefined;
|
||||
|
|
@ -82,7 +82,7 @@ pub const TwelveData = struct {
|
|||
allocator: std.mem.Allocator,
|
||||
symbol: []const u8,
|
||||
) !Quote {
|
||||
self.rate_limiter.acquire();
|
||||
try self.rate_limiter.acquire();
|
||||
|
||||
const url = try http.buildUrl(allocator, base_url ++ "/quote", &.{
|
||||
.{ "symbol", symbol },
|
||||
|
|
|
|||
178
src/service.zig
178
src/service.zig
|
|
@ -82,6 +82,13 @@ pub const DataError = error{
|
|||
/// user "this symbol isn't in the provider's catalog; mark it
|
||||
/// manually" instead of an opaque "fetch failed."
|
||||
NotFound,
|
||||
/// The calling task was canceled (e.g. the TUI is quitting and
|
||||
/// `PortfolioData.cancelLoad` stopped its workers). Not a data or
|
||||
/// provider fault: nothing is cached, negative-cached, or retried,
|
||||
/// and no fallback provider is tried. Callers return; looping on to
|
||||
/// the next symbol would run uncancelable, because std delivers a
|
||||
/// cancelation to exactly one cancelation point.
|
||||
Canceled,
|
||||
};
|
||||
|
||||
/// Per-call options controlling cache vs network behavior. Drives
|
||||
|
|
@ -717,7 +724,9 @@ pub const DataService = struct {
|
|||
}
|
||||
|
||||
// Try server sync before hitting providers (skipped on force_refresh).
|
||||
if (!opts.force_refresh and self.syncFromServer(symbol, data_type)) {
|
||||
// A canceled sync stops here: falling through to the provider would
|
||||
// be exactly the uncancelable extra work a cancel must prevent.
|
||||
if (!opts.force_refresh and (self.syncFromServer(symbol, data_type) catch return DataError.Canceled)) {
|
||||
if (s.read(self.allocator, T, symbol, postProcess, .fresh_only)) |cached| {
|
||||
// The hook applies here too. Without it a configured
|
||||
// ZFIN_SERVER defeats the whole mechanism: the sync
|
||||
|
|
@ -738,10 +747,14 @@ pub const DataService = struct {
|
|||
log.debug("{s}: fetching {s} from provider", .{ symbol, @tagName(data_type) });
|
||||
self.assertNetworkAllowed("fetchCached fetchFromProvider");
|
||||
const fetched = self.fetchFromProvider(T, symbol) catch |err| {
|
||||
// Before anything that caches or retries: a cancel is not a
|
||||
// statement about the symbol.
|
||||
if (err == error.Canceled) return DataError.Canceled;
|
||||
if (err == error.RateLimited) {
|
||||
// Wait and retry once
|
||||
self.rateLimitBackoff();
|
||||
self.rateLimitBackoff() catch return DataError.Canceled;
|
||||
const retried = self.fetchFromProvider(T, symbol) catch |retry_err| {
|
||||
if (retry_err == error.Canceled) return DataError.Canceled;
|
||||
log.warn("{s}: {s} fetch failed after rate-limit retry: {t}", .{ symbol, @tagName(data_type), retry_err });
|
||||
return DataError.FetchFailed;
|
||||
};
|
||||
|
|
@ -950,6 +963,8 @@ pub const DataService = struct {
|
|||
self.allocator.free(triple.candles);
|
||||
log.warn("{s}: Tiingo full history returned no bars, trying Yahoo", .{symbol});
|
||||
} else |err| {
|
||||
// Nor may a cancel: the Yahoo request would run uncancelable.
|
||||
if (err == error.Canceled) return DataError.Canceled;
|
||||
// Transient failures must not silently degrade to a
|
||||
// second provider - the caller needs to know to retry.
|
||||
if (err == error.RateLimited or isTransientError(err)) return DataError.TransientError;
|
||||
|
|
@ -988,6 +1003,7 @@ pub const DataService = struct {
|
|||
self.allocator.free(candles);
|
||||
log.warn("{s}: Yahoo full history returned no bars", .{symbol});
|
||||
} else |err| {
|
||||
if (err == error.Canceled) return DataError.Canceled;
|
||||
log.warn("{s}: Yahoo full history failed: {s}", .{ symbol, @errorName(err) });
|
||||
// Both providers affirmatively disclaim the symbol.
|
||||
if (tiingo_not_found and isPermanentProviderFailure(err)) return error.NotFound;
|
||||
|
|
@ -1156,6 +1172,7 @@ pub const DataService = struct {
|
|||
log.debug("{s}: candles from Yahoo (Tiingo backoff active)", .{symbol});
|
||||
return .{ .candles = candles, .provider = .yahoo, .tiingo_coverage = .unknown };
|
||||
} else |err| {
|
||||
if (err == error.Canceled) return DataError.Canceled;
|
||||
log.warn("{s}: Yahoo (Tiingo backoff active) failed: {s}", .{ symbol, @errorName(err) });
|
||||
}
|
||||
} else |_| {}
|
||||
|
|
@ -1170,6 +1187,9 @@ pub const DataService = struct {
|
|||
log.debug("{s}: candles from Tiingo", .{symbol});
|
||||
return .{ .candles = candles, .provider = .tiingo, .tiingo_coverage = .covered };
|
||||
} else |err| {
|
||||
// A cancel is not a Tiingo verdict, and every branch below
|
||||
// would make another (now uncancelable) request.
|
||||
if (err == error.Canceled) return DataError.Canceled;
|
||||
log.warn("{s}: Tiingo failed: {s}", .{ symbol, @errorName(err) });
|
||||
|
||||
if (err == error.Unauthorized) {
|
||||
|
|
@ -1180,19 +1200,22 @@ pub const DataService = struct {
|
|||
if (err == error.RateLimited) {
|
||||
// Rate limited: back off and retry - this is expected, not a failure
|
||||
log.info("{s}: Tiingo rate limited, backing off", .{symbol});
|
||||
self.rateLimitBackoff();
|
||||
self.rateLimitBackoff() catch return DataError.Canceled;
|
||||
if (tg.fetchCandles(self.allocator, symbol, from, to)) |candles| {
|
||||
log.debug("{s}: candles from Tiingo (after rate limit backoff)", .{symbol});
|
||||
return .{ .candles = candles, .provider = .tiingo, .tiingo_coverage = .covered };
|
||||
} else |retry_err| {
|
||||
if (retry_err == error.Canceled) return DataError.Canceled;
|
||||
log.warn("{s}: Tiingo retry after backoff failed: {s}", .{ symbol, @errorName(retry_err) });
|
||||
if (retry_err == error.RateLimited) {
|
||||
// Still rate limited after backoff - one more try
|
||||
self.rateLimitBackoff();
|
||||
self.rateLimitBackoff() catch return DataError.Canceled;
|
||||
if (tg.fetchCandles(self.allocator, symbol, from, to)) |candles| {
|
||||
log.debug("{s}: candles from Tiingo (after second backoff)", .{symbol});
|
||||
return .{ .candles = candles, .provider = .tiingo, .tiingo_coverage = .covered };
|
||||
} else |_| {}
|
||||
} else |final_err| {
|
||||
if (final_err == error.Canceled) return DataError.Canceled;
|
||||
}
|
||||
}
|
||||
// Exhausted rate limit retries - treat as transient
|
||||
return DataError.TransientError;
|
||||
|
|
@ -1230,6 +1253,7 @@ pub const DataService = struct {
|
|||
log.info("{s}: candles from Yahoo (Tiingo fallback)", .{symbol});
|
||||
return .{ .candles = candles, .provider = .yahoo, .tiingo_coverage = coverage };
|
||||
} else |err| {
|
||||
if (err == error.Canceled) return DataError.Canceled;
|
||||
log.warn("{s}: Yahoo fallback also failed: {s}", .{ symbol, @errorName(err) });
|
||||
}
|
||||
} else |_| {
|
||||
|
|
@ -1401,7 +1425,7 @@ pub const DataService = struct {
|
|||
// Stale - try server sync before incremental fetch.
|
||||
// (Force-refresh skips server sync too: the user explicitly
|
||||
// asked for fresh provider data.)
|
||||
if (!opts.force_refresh and self.syncCandlesFromServer(symbol)) {
|
||||
if (!opts.force_refresh and (self.syncCandlesFromServer(symbol) catch return DataError.Canceled)) {
|
||||
// Re-read meta: the sync wrote the server's bytes
|
||||
// verbatim, so its view of both freshness AND
|
||||
// adjustment basis is now ours. `serverBarRegression`
|
||||
|
|
@ -1448,6 +1472,9 @@ pub const DataService = struct {
|
|||
if (self.refetchFullHistory(symbol, today, now_s, now_s < m.tiingo_retry_after_s)) |candles| {
|
||||
return .{ .data = candles, .source = .fetched, .timestamp = std.Io.Timestamp.now(self.io, .real).toSeconds(), .allocator = self.allocator };
|
||||
} else |err| {
|
||||
// A cancel ends the call; falling through would
|
||||
// start the incremental fetch uncancelable.
|
||||
if (err == DataError.Canceled) return DataError.Canceled;
|
||||
// Restatement is best-effort. The existing series
|
||||
// is untouched and still usable, just understated
|
||||
// by the missed adjustment - fall through to the
|
||||
|
|
@ -1475,6 +1502,8 @@ pub const DataService = struct {
|
|||
// Incremental fetch from day after last cached candle
|
||||
self.assertNetworkAllowed("getCandles incremental fetchCandlesFromProviders");
|
||||
const result = self.fetchCandlesFromProviders(symbol, fetch_from, today, now_s < m.tiingo_retry_after_s) catch |err| {
|
||||
// Not a failure of the symbol: no fail_count, no stale fallback.
|
||||
if (err == DataError.Canceled) return DataError.Canceled;
|
||||
if (err == DataError.TransientError) {
|
||||
// Increment fail_count for this symbol
|
||||
const new_fail_count = m.fail_count +| 1; // saturating add
|
||||
|
|
@ -1540,7 +1569,7 @@ pub const DataService = struct {
|
|||
}
|
||||
|
||||
// No usable cache - try server sync first (skipped on force_refresh).
|
||||
if (!opts.force_refresh and self.syncCandlesFromServer(symbol)) {
|
||||
if (!opts.force_refresh and (self.syncCandlesFromServer(symbol) catch return DataError.Canceled)) {
|
||||
if (s.isCandleMetaFresh(symbol)) {
|
||||
log.debug("{s}: candles synced from server and fresh (no prior cache)", .{symbol});
|
||||
if (s.read(self.allocator, Candle, symbol, null, .any)) |r|
|
||||
|
|
@ -1569,6 +1598,8 @@ pub const DataService = struct {
|
|||
|
||||
const prior_backoff: i64 = if (meta_result) |mr| mr.meta.tiingo_retry_after_s else 0;
|
||||
const candles = self.refetchFullHistory(symbol, today, now_s, now_s < prior_backoff) catch |err| {
|
||||
// Not a failure of the symbol: no fail_count, no negative cache.
|
||||
if (err == DataError.Canceled) return DataError.Canceled;
|
||||
if (err == DataError.TransientError) {
|
||||
// Transient: increment fail_count on existing meta so
|
||||
// we know to back off if this keeps happening.
|
||||
|
|
@ -2008,7 +2039,7 @@ pub const DataService = struct {
|
|||
}
|
||||
|
||||
// Try server sync before hitting Wikidata.
|
||||
if (!opts.force_refresh and self.syncFromServer(symbol, .classification)) {
|
||||
if (!opts.force_refresh and (self.syncFromServer(symbol, .classification) catch return DataError.Canceled)) {
|
||||
if (s.read(self.allocator, Wikidata.ClassificationRecord, symbol, null, .fresh_only)) |cached| {
|
||||
log.debug("{s}: classification synced from server", .{symbol});
|
||||
return .{ .data = cached.data, .source = .cached, .timestamp = cached.timestamp, .allocator = self.allocator };
|
||||
|
|
@ -2021,8 +2052,9 @@ pub const DataService = struct {
|
|||
|
||||
const symbols = [_][]const u8{symbol};
|
||||
const fetched = wd.fetch(self.allocator, &symbols) catch |err| {
|
||||
if (err == error.Canceled) return DataError.Canceled;
|
||||
if (err == error.RateLimited) {
|
||||
self.rateLimitBackoff();
|
||||
self.rateLimitBackoff() catch return DataError.Canceled;
|
||||
if (wd.fetch(self.allocator, &symbols)) |retried| {
|
||||
return self.finalizeClassification(symbol, retried, opts);
|
||||
} else |_| {}
|
||||
|
|
@ -2327,7 +2359,7 @@ pub const DataService = struct {
|
|||
return DataError.FetchFailed;
|
||||
}
|
||||
|
||||
if (!opts.force_refresh and self.syncFromServer(cik, .entity_facts)) {
|
||||
if (!opts.force_refresh and (self.syncFromServer(cik, .entity_facts) catch return DataError.Canceled)) {
|
||||
if (s.read(self.allocator, Edgar.EntityFactRecord, cik, null, .fresh_only)) |cached| {
|
||||
log.debug("CIK {s}: entity_facts synced from server", .{cik});
|
||||
return .{ .data = cached.data, .source = .cached, .timestamp = cached.timestamp, .allocator = self.allocator };
|
||||
|
|
@ -2411,7 +2443,7 @@ pub const DataService = struct {
|
|||
return DataError.FetchFailed;
|
||||
}
|
||||
|
||||
if (!opts.force_refresh and self.syncFromServer(symbol, .etf_metrics)) {
|
||||
if (!opts.force_refresh and (self.syncFromServer(symbol, .etf_metrics) catch return DataError.Canceled)) {
|
||||
if (s.read(self.allocator, Edgar.EtfMetricRecord, symbol, null, .fresh_only)) |cached| {
|
||||
log.debug("{s}: etf_metrics synced from server", .{symbol});
|
||||
return .{
|
||||
|
|
@ -3194,7 +3226,7 @@ pub const DataService = struct {
|
|||
fn run(io: std.Io, svc: *DataService, slot: *ServerSyncResult, done: *AtomicCounter) std.Io.Cancelable!void {
|
||||
defer _ = done.increment();
|
||||
try io.checkCancel();
|
||||
slot.success = svc.syncCandlesFromServer(slot.symbol);
|
||||
slot.success = try svc.syncCandlesFromServer(slot.symbol);
|
||||
}
|
||||
};
|
||||
|
||||
|
|
@ -3651,11 +3683,13 @@ pub const DataService = struct {
|
|||
|
||||
/// Sleep before retrying after a rate limit error.
|
||||
/// Uses the provider's rate limiter if available, otherwise a fixed 10s backoff.
|
||||
fn rateLimitBackoff(self: *DataService) void {
|
||||
/// Wait out a provider's 429. Returns `error.Canceled` if the task
|
||||
/// is canceled while waiting; callers stop rather than retry.
|
||||
fn rateLimitBackoff(self: *DataService) std.Io.Cancelable!void {
|
||||
if (self.td) |*td| {
|
||||
td.rate_limiter.backoff();
|
||||
try td.rate_limiter.backoff();
|
||||
} else {
|
||||
std.Io.sleep(self.io, std.Io.Duration.fromSeconds(10), .awake) catch |err| log.debug("rate-limit backoff sleep interrupted: {t}", .{err});
|
||||
try std.Io.sleep(self.io, std.Io.Duration.fromSeconds(10), .awake);
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -3665,6 +3699,10 @@ pub const DataService = struct {
|
|||
/// Returns true if the file was successfully synced, false on any error.
|
||||
/// Silently returns false if no server is configured.
|
||||
///
|
||||
/// `error.Canceled` means the calling task was canceled mid-sync.
|
||||
/// Callers stop rather than fall back to a provider; see
|
||||
/// `DataError.Canceled`.
|
||||
///
|
||||
/// Applies a single retry with a short delay when the first attempt
|
||||
/// fails at the HTTP layer OR produces a torn body (integrity
|
||||
/// mismatch / `looksCompleteSrf` rejection). Motivation: refreshes
|
||||
|
|
@ -3676,7 +3714,7 @@ pub const DataService = struct {
|
|||
/// from the same refresh are the most valuable diagnostic signal
|
||||
/// we can produce (same body shape? same byte offset? same time
|
||||
/// delta? all answers we can't get from a single failure).
|
||||
fn syncFromServer(self: *DataService, symbol: []const u8, data_type: cache.DataType) bool {
|
||||
fn syncFromServer(self: *DataService, symbol: []const u8, data_type: cache.DataType) std.Io.Cancelable!bool {
|
||||
const server_url = self.config.server_url orelse return false;
|
||||
const endpoint = switch (data_type) {
|
||||
.candles_daily => "/candles",
|
||||
|
|
@ -3709,18 +3747,21 @@ pub const DataService = struct {
|
|||
"{s}: retrying {s} server sync (attempt {d}/{d}) after {d}ms delay",
|
||||
.{ symbol, @tagName(data_type), attempt + 1, max_attempts, retry_delay_ms },
|
||||
);
|
||||
std.Io.sleep(self.io, std.Io.Duration.fromMilliseconds(retry_delay_ms), .awake) catch |err| log.debug("syncFromServer retry-delay sleep interrupted: {t}", .{err});
|
||||
try std.Io.sleep(self.io, std.Io.Duration.fromMilliseconds(retry_delay_ms), .awake);
|
||||
}
|
||||
switch (self.tryOneSync(symbol, data_type, full_url)) {
|
||||
.ok => return true,
|
||||
// Torn or network error - retry if attempts remain.
|
||||
.torn, .net_err => {},
|
||||
.canceled => return error.Canceled,
|
||||
}
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
const SyncAttempt = enum { ok, torn, net_err };
|
||||
/// `canceled`: the request was stopped by a cancelation of the
|
||||
/// calling task, which must not be retried (see `http.Client.request`).
|
||||
const SyncAttempt = enum { ok, torn, net_err, canceled };
|
||||
|
||||
/// One attempt at syncing a file from the server. Archives a torn
|
||||
/// body when detected but does NOT retry - the caller decides that.
|
||||
|
|
@ -3746,6 +3787,11 @@ pub const DataService = struct {
|
|||
const extra_headers = serverAuthHeaders(self.config.server_api_key, &hdr_buf);
|
||||
var response = client.request(.GET, full_url, null, extra_headers) catch |err| {
|
||||
const elapsed_ms = @divTrunc(std.Io.Timestamp.now(self.io, .awake).nanoseconds - t_start, std.time.ns_per_ms);
|
||||
// Not a sync failure, so not a warning: the caller is stopping.
|
||||
if (err == error.Canceled) {
|
||||
log.debug("{s}: tryOneSync finished ({s}) result=canceled elapsed_ms={d}", .{ symbol, @tagName(data_type), elapsed_ms });
|
||||
return .canceled;
|
||||
}
|
||||
// Operator-visible: surfaces meaningful failures
|
||||
// (`NoAddressReturned`, `ConnectionRefused`,
|
||||
// `TlsInitializationFailed`, etc.) instead of swallowing
|
||||
|
|
@ -3983,9 +4029,9 @@ pub const DataService = struct {
|
|||
return market.staleCandleExpiry(now_s, kind, newest);
|
||||
}
|
||||
|
||||
fn syncCandlesFromServer(self: *DataService, symbol: []const u8) bool {
|
||||
const daily = self.syncFromServer(symbol, .candles_daily);
|
||||
const meta = self.syncFromServer(symbol, .candles_meta);
|
||||
fn syncCandlesFromServer(self: *DataService, symbol: []const u8) std.Io.Cancelable!bool {
|
||||
const daily = try self.syncFromServer(symbol, .candles_daily);
|
||||
const meta = try self.syncFromServer(symbol, .candles_meta);
|
||||
return daily and meta;
|
||||
}
|
||||
|
||||
|
|
@ -4113,11 +4159,14 @@ pub const DataService = struct {
|
|||
/// a whole portfolio load over. `fetchCached` logs the provider's own
|
||||
/// error (rate limit vs auth vs no-such-data) before collapsing it to
|
||||
/// `FetchFailed`, so a swallowed failure is still diagnosable.
|
||||
///
|
||||
/// The exception is cancelation, which stops the batch and returns
|
||||
/// `error.Canceled` whether it lands between symbols or mid-fetch.
|
||||
pub fn loadAllDividends(
|
||||
self: *DataService,
|
||||
syms: []const []const u8,
|
||||
opts: FetchOptions,
|
||||
) void {
|
||||
) std.Io.Cancelable!void {
|
||||
for (syms) |sym| {
|
||||
// Per symbol, not once before the loop. Every cache miss here is a
|
||||
// provider round trip under a 4-request/minute budget, so a batch
|
||||
|
|
@ -4126,8 +4175,14 @@ pub const DataService = struct {
|
|||
// (`PortfolioData.cancelLoad`), and cancelling waits for the
|
||||
// worker - so without a check inside the loop, quitting blocked
|
||||
// until the whole batch drained.
|
||||
self.io.checkCancel() catch return;
|
||||
const fr = self.getDividends(sym, opts) catch continue;
|
||||
try self.io.checkCancel();
|
||||
// The check above only catches a cancel that arrived between
|
||||
// symbols; one that lands mid-fetch comes back as this error,
|
||||
// and continuing would run every remaining symbol uncancelable.
|
||||
const fr = self.getDividends(sym, opts) catch |err| switch (err) {
|
||||
DataError.Canceled => return error.Canceled,
|
||||
else => continue,
|
||||
};
|
||||
fr.deinit();
|
||||
}
|
||||
}
|
||||
|
|
@ -4135,6 +4190,8 @@ pub const DataService = struct {
|
|||
|
||||
// ── Tests ─────────────────────────────────────────────────────────
|
||||
|
||||
const SilentServer = @import("testutil/silent_server.zig").SilentServer;
|
||||
|
||||
test "serverAuthHeaders: key present yields one X-API-Key header" {
|
||||
var buf: [1]std.http.Header = .{.{ .name = "", .value = "" }};
|
||||
const h = DataService.serverAuthHeaders("s3cret", &buf);
|
||||
|
|
@ -5677,7 +5734,7 @@ test "loadAllDividends: honors skip_network for every symbol, and one miss does
|
|||
|
||||
// The whole loop must stay offline, not just the first symbol.
|
||||
svc.panic_on_network_attempt = true;
|
||||
svc.loadAllDividends(&.{ "TSTA", "TSTB" }, .{ .skip_network = true });
|
||||
try svc.loadAllDividends(&.{ "TSTA", "TSTB" }, .{ .skip_network = true });
|
||||
|
||||
// The cached symbol survives the pass.
|
||||
const b = svc.getCachedDividends(allocator, "TSTB") orelse return error.TestUnexpectedResult;
|
||||
|
|
@ -5721,6 +5778,75 @@ test "getCachedSplits: reads cache only and reports absence" {
|
|||
try std.testing.expect(svc.getCachedSplits(allocator, "TSTA") == null);
|
||||
}
|
||||
|
||||
test "loadAllDividends: a cancel mid-sync stops the batch - no retry, no next symbol, no provider" {
|
||||
// The TUI exit hang. Quitting cancels the dividend warm; the cancel
|
||||
// used to surface as a generic read failure that the HTTP and sync
|
||||
// layers retried, and every retry and every later symbol then ran
|
||||
// uncancelable, so exit waited out the whole batch against a server
|
||||
// stuck on rate-limited provider fetches. Here the server never
|
||||
// answers at all: the batch must still return Canceled after one
|
||||
// connection.
|
||||
const allocator = std.testing.allocator;
|
||||
const io = std.testing.io;
|
||||
var tmp = std.testing.tmpDir(.{});
|
||||
defer tmp.cleanup();
|
||||
const dir_path = try tmp.dir.realPathFileAlloc(io, ".", allocator);
|
||||
defer allocator.free(dir_path);
|
||||
|
||||
var server: SilentServer = undefined;
|
||||
try server.start(io, .hold_open);
|
||||
defer server.stop();
|
||||
var url_buf: [64]u8 = undefined;
|
||||
|
||||
var svc = DataService.init(io, allocator, .{ .cache_dir = dir_path, .server_url = try server.baseUrl(&url_buf) });
|
||||
defer svc.deinit();
|
||||
// Falling back to a provider after a cancel would panic here.
|
||||
svc.panic_on_network_attempt = true;
|
||||
|
||||
const symbols = [_][]const u8{ "SMPLA", "SMPLB", "SMPLC" };
|
||||
var warm = try io.concurrent(DataService.loadAllDividends, .{ &svc, &symbols, FetchOptions{} });
|
||||
try server.waitForConnections(1, 5000);
|
||||
try std.testing.expectError(error.Canceled, warm.cancel(io));
|
||||
|
||||
// A retry (250ms sync delay) or the next symbol would have connected.
|
||||
try io.sleep(.fromMilliseconds(400), .awake);
|
||||
try std.testing.expectEqual(@as(u32, 1), server.accepted.load(.acquire));
|
||||
// Nothing was cached for the canceled symbol, not even a negative entry.
|
||||
try std.testing.expect(svc.getCachedDividends(allocator, "SMPLA") == null);
|
||||
}
|
||||
|
||||
test "getCandles / getClassification: a cancel during server sync returns Canceled without a provider fallback" {
|
||||
// Same contract as the dividend warm, through the other entry points
|
||||
// that sync from the server first. Cold caches, so each goes straight
|
||||
// to the sync.
|
||||
const allocator = std.testing.allocator;
|
||||
const io = std.testing.io;
|
||||
var tmp = std.testing.tmpDir(.{});
|
||||
defer tmp.cleanup();
|
||||
const dir_path = try tmp.dir.realPathFileAlloc(io, ".", allocator);
|
||||
defer allocator.free(dir_path);
|
||||
|
||||
var server: SilentServer = undefined;
|
||||
try server.start(io, .hold_open);
|
||||
defer server.stop();
|
||||
var url_buf: [64]u8 = undefined;
|
||||
|
||||
var svc = DataService.init(io, allocator, .{ .cache_dir = dir_path, .server_url = try server.baseUrl(&url_buf) });
|
||||
defer svc.deinit();
|
||||
svc.panic_on_network_attempt = true;
|
||||
|
||||
var pending_candles = try io.concurrent(DataService.getCandles, .{ &svc, "SMPLA", FetchOptions{} });
|
||||
try server.waitForConnections(1, 5000);
|
||||
try std.testing.expectError(DataError.Canceled, pending_candles.cancel(io));
|
||||
|
||||
var pending_classification = try io.concurrent(DataService.getClassification, .{ &svc, "SMPLB", FetchOptions{} });
|
||||
try server.waitForConnections(2, 5000);
|
||||
try std.testing.expectError(DataError.Canceled, pending_classification.cancel(io));
|
||||
|
||||
try io.sleep(.fromMilliseconds(400), .awake);
|
||||
try std.testing.expectEqual(@as(u32, 2), server.accepted.load(.acquire));
|
||||
}
|
||||
|
||||
test "loadAllDividends: empty symbol list is a no-op" {
|
||||
const allocator = std.testing.allocator;
|
||||
const io = std.testing.io;
|
||||
|
|
@ -5736,7 +5862,7 @@ test "loadAllDividends: empty symbol list is a no-op" {
|
|||
// No symbols means no fetches, so this must hold even with network
|
||||
// otherwise allowed.
|
||||
svc.panic_on_network_attempt = true;
|
||||
svc.loadAllDividends(&.{}, .{});
|
||||
try svc.loadAllDividends(&.{}, .{});
|
||||
}
|
||||
|
||||
test "getQuote offline mode returns FetchFailed (quotes never cached)" {
|
||||
|
|
|
|||
94
src/testutil/silent_server.zig
Normal file
94
src/testutil/silent_server.zig
Normal file
|
|
@ -0,0 +1,94 @@
|
|||
//! Test-only: a loopback TCP server that accepts connections and never
|
||||
//! answers.
|
||||
//!
|
||||
//! In `.hold_open` mode a request to it connects, sends, and then blocks
|
||||
//! in `receiveHead` for as long as the test likes - the shape of a
|
||||
//! zfin-server holding a request open while it fetches from a
|
||||
//! rate-limited provider. Tests use it to prove that canceling the task
|
||||
//! making the request stops it promptly. In `.hang_up` mode every
|
||||
//! connection is closed as soon as it is accepted, a real transport
|
||||
//! failure that the HTTP client should still retry.
|
||||
//!
|
||||
//! Either way `accepted` counts connections, so a retry, or a move on to
|
||||
//! the next symbol, shows up as an extra one.
|
||||
//!
|
||||
//! Loopback only, ephemeral port, nothing leaves the machine.
|
||||
|
||||
const std = @import("std");
|
||||
const net = std.Io.net;
|
||||
|
||||
pub const SilentServer = struct {
|
||||
pub const Behavior = enum {
|
||||
/// Accept and hold the connection open without responding.
|
||||
hold_open,
|
||||
/// Accept and close immediately without responding.
|
||||
hang_up,
|
||||
};
|
||||
|
||||
io: std.Io,
|
||||
server: net.Server,
|
||||
behavior: Behavior,
|
||||
/// Connections accepted so far.
|
||||
accepted: std.atomic.Value(u32) = .init(0),
|
||||
/// Streams held open in `.hold_open` mode, so a client waits on a
|
||||
/// live connection instead of seeing it closed. Only the first
|
||||
/// `accepted` entries (up to `max_streams`) are set.
|
||||
// SAFETY: entries are written by `acceptLoop` before `accepted` is
|
||||
// incremented past them, and only that many are read in `stop`.
|
||||
streams: [max_streams]net.Stream = undefined,
|
||||
acceptor: ?std.Io.Future(void) = null,
|
||||
|
||||
const max_streams = 8;
|
||||
|
||||
/// Listen on 127.0.0.1 with an ephemeral port and start accepting.
|
||||
/// The server must not move after `start`; the accept task holds a
|
||||
/// pointer to it.
|
||||
pub fn start(self: *SilentServer, io: std.Io, behavior: Behavior) !void {
|
||||
const address: net.IpAddress = try .parseIp4("127.0.0.1", 0);
|
||||
self.* = .{ .io = io, .server = try address.listen(io, .{}), .behavior = behavior };
|
||||
errdefer self.server.deinit(io);
|
||||
// `concurrent`, not `async`: accepting blocks, and the test
|
||||
// thread has to keep running alongside it.
|
||||
self.acceptor = try io.concurrent(acceptLoop, .{self});
|
||||
}
|
||||
|
||||
/// Stop accepting and close everything.
|
||||
pub fn stop(self: *SilentServer) void {
|
||||
if (self.acceptor) |*f| f.cancel(self.io);
|
||||
if (self.behavior == .hold_open) {
|
||||
const held = @min(self.accepted.load(.acquire), max_streams);
|
||||
for (self.streams[0..held]) |*s| s.close(self.io);
|
||||
}
|
||||
self.server.deinit(self.io);
|
||||
}
|
||||
|
||||
/// `http://127.0.0.1:<port>`, written into `buf`.
|
||||
pub fn baseUrl(self: *const SilentServer, buf: []u8) ![]const u8 {
|
||||
return std.fmt.bufPrint(buf, "http://127.0.0.1:{d}", .{self.server.socket.address.getPort()});
|
||||
}
|
||||
|
||||
/// Wait until at least `n` connections have been accepted, or fail
|
||||
/// after `timeout_ms`.
|
||||
pub fn waitForConnections(self: *SilentServer, n: u32, timeout_ms: u32) !void {
|
||||
var waited: u32 = 0;
|
||||
while (self.accepted.load(.acquire) < n) : (waited += 5) {
|
||||
if (waited >= timeout_ms) return error.ConnectionNeverArrived;
|
||||
try self.io.sleep(.fromMilliseconds(5), .awake);
|
||||
}
|
||||
}
|
||||
|
||||
fn acceptLoop(self: *SilentServer) void {
|
||||
while (true) {
|
||||
// Canceled by `stop`; any other accept failure also ends the
|
||||
// loop, and the test sees it as a missing connection.
|
||||
const stream = self.server.accept(self.io) catch return;
|
||||
const index = self.accepted.load(.acquire);
|
||||
if (self.behavior == .hold_open and index < max_streams) {
|
||||
self.streams[index] = stream;
|
||||
} else {
|
||||
stream.close(self.io);
|
||||
}
|
||||
_ = self.accepted.fetchAdd(1, .release);
|
||||
}
|
||||
}
|
||||
};
|
||||
Loading…
Add table
Reference in a new issue