From 2876358c5a7d171356802ab7033a21a172c8e415 Mon Sep 17 00:00:00 2001 From: Emil Lerch Date: Sun, 4 Oct 2026 14:02:23 -0700 Subject: [PATCH] quit from interactive mode means quit --- src/PortfolioData.zig | 2 +- src/commands/review.zig | 2 +- src/net/RateLimiter.zig | 50 +++++--- src/net/http.zig | 206 ++++++++++++++++++++++++++------- src/providers/Edgar.zig | 2 +- src/providers/cboe.zig | 2 +- src/providers/fmp.zig | 2 +- src/providers/polygon.zig | 6 +- src/providers/tiingo.zig | 4 +- src/providers/twelvedata.zig | 4 +- src/service.zig | 178 +++++++++++++++++++++++----- src/testutil/silent_server.zig | 94 +++++++++++++++ 12 files changed, 458 insertions(+), 94 deletions(-) create mode 100644 src/testutil/silent_server.zig diff --git a/src/PortfolioData.zig b/src/PortfolioData.zig index d95d0ce..f8f434b 100644 --- a/src/PortfolioData.zig +++ b/src/PortfolioData.zig @@ -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); diff --git a/src/commands/review.zig b/src/commands/review.zig index dd585f2..527c8a3 100644 --- a/src/commands/review.zig +++ b/src/commands/review.zig @@ -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| { diff --git a/src/net/RateLimiter.zig b/src/net/RateLimiter.zig index ec0c6b2..47d38ca 100644 --- a/src/net/RateLimiter.zig +++ b/src/net/RateLimiter.zig @@ -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)); +} diff --git a/src/net/http.zig b/src/net/http.zig index 02f49d0..dd66ca3 100644 --- a/src/net/http.zig +++ b/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; diff --git a/src/providers/Edgar.zig b/src/providers/Edgar.zig index 4a949b6..65a0165 100644 --- a/src/providers/Edgar.zig +++ b/src/providers/Edgar.zig @@ -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; diff --git a/src/providers/cboe.zig b/src/providers/cboe.zig index 8132a6b..b96d5de 100644 --- a/src/providers/cboe.zig +++ b/src/providers/cboe.zig @@ -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); diff --git a/src/providers/fmp.zig b/src/providers/fmp.zig index 1ea89ad..bdb5ff4 100644 --- a/src/providers/fmp.zig +++ b/src/providers/fmp.zig @@ -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 }, diff --git a/src/providers/polygon.zig b/src/providers/polygon.zig index 22493dc..21f9848 100644 --- a/src/providers/polygon.zig +++ b/src/providers/polygon.zig @@ -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 }, diff --git a/src/providers/tiingo.zig b/src/providers/tiingo.zig index 2cac4ee..70e2eb3 100644 --- a/src/providers/tiingo.zig +++ b/src/providers/tiingo.zig @@ -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(); diff --git a/src/providers/twelvedata.zig b/src/providers/twelvedata.zig index e7e319a..f135957 100644 --- a/src/providers/twelvedata.zig +++ b/src/providers/twelvedata.zig @@ -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 }, diff --git a/src/service.zig b/src/service.zig index 78734f5..bb40916 100644 --- a/src/service.zig +++ b/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)" { diff --git a/src/testutil/silent_server.zig b/src/testutil/silent_server.zig new file mode 100644 index 0000000..d9a8339 --- /dev/null +++ b/src/testutil/silent_server.zig @@ -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:`, 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); + } + } +};