Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970//! keeps a namespace's cache warm from a background thread.//!//! a cold namespace pays several hundred milliseconds on its first query.//! the warm-cache hint is free when the namespace is already warm, so pinging//! it on an interval is cheap insurance for a latency-sensitive path.//!//! ```zig//! var keepalive: tpuf.Keepalive = .init(ns, alloc, .{});//! try keepalive.start();//! defer keepalive.stop();//! ```//!//! the struct must not move after `start`: the thread holds a pointer to it.//!//! this is a `std.Thread` calling into the transport's `Io` from outside it,//! which holds under `Io.Threaded` (what the adopters run) and is not a//! supported pattern under `Io.Evented`. on an evented backend, run the same//! loop as an `io.async` task instead.
const Keepalive = @This();
const std = @import("std");const Io = std.Io;const Allocator = std.mem.Allocator;const Namespace = @import("Namespace.zig");
const log = std.log.scoped(.tpuf);
pub const Options = struct { interval_ms: u32 = 3 * 60 * 1000,};
namespace: Namespace,alloc: Allocator,options: Options,stop_flag: std.atomic.Value(bool) = .init(false),thread: ?std.Thread = null,
pub fn init(namespace: Namespace, alloc: Allocator, options: Options) Keepalive { return .{ .namespace = namespace, .alloc = alloc, .options = options };}
pub fn start(self: *Keepalive) !void { if (self.thread != null) return; self.stop_flag.store(false, .release); self.thread = try std.Thread.spawn(.{}, run, .{self});}
/// signals the thread and joins it. returns within about a second.pub fn stop(self: *Keepalive) void { const thread = self.thread orelse return; self.stop_flag.store(true, .release); thread.join(); self.thread = null;}
fn run(self: *Keepalive) void { const io = self.namespace.client.transport.io; const tick: Io.Duration = .fromMilliseconds(1000); while (!self.stop_flag.load(.acquire)) { self.namespace.warmCache(.{}) catch |err| { log.debug("keepalive ping for {s} failed: {t}", .{ self.namespace.name, err }); }; var slept: u64 = 0; while (slept < self.options.interval_ms and !self.stop_flag.load(.acquire)) : (slept += 1000) { io.sleep(tick, .awake) catch return; } }}