db-mode config changes apply live in-process
settings and upstream writes now follow a prepare, commit, publish, retire contract: candidates are built and validated before the database transaction, published as infallible pointer swaps, and old generations retire after their readers drain. per-query policy values snapshot once per query; upstream pool, cache, rate limiter, sessions, api limiter, log sink, blocklist scheduler and the query-log queue each gained one named live operation. restart_required shrinks from every scalar key to the bind keys and web.enabled; the admin ui drops its restart notices for everything else. file mode is unchanged.
This commit is contained in:
+274
-8
@@ -109,10 +109,13 @@ pub const CertStore = struct {
|
||||
pub const Kind = enum { doh, dot };
|
||||
|
||||
gpa: std.mem.Allocator,
|
||||
/// Borrowed from the config; must outlive the store.
|
||||
cert_path: []const u8,
|
||||
/// Borrowed from the config; must outlive the store.
|
||||
key_path: []const u8,
|
||||
/// Owned. A settings apply replaces both paths, so a borrowed config slice
|
||||
/// would dangle the moment the row it came from went away. Read and
|
||||
/// written only under `reload_mutex`, which every apply and every reload
|
||||
/// holds for its whole read-build-publish sequence.
|
||||
cert_path: []u8,
|
||||
/// Owned; see `cert_path`.
|
||||
key_path: []u8,
|
||||
/// Passed through to every `ServerContext.init`; the same lifetime rule
|
||||
/// applies — a comptime-constant NULL-terminated array, never memory that
|
||||
/// can go away before the store.
|
||||
@@ -136,7 +139,8 @@ pub const CertStore = struct {
|
||||
/// Test seam: runs inside `reload` between a successful `load` and the
|
||||
/// publish, i.e. inside `reload_mutex`. Lets a test occupy the window
|
||||
/// where an unserialized reload could be overtaken. Must not call
|
||||
/// `reload` synchronously (that would self-deadlock on `reload_mutex`).
|
||||
/// `reload`, `pollOnce` or `preparePathChange` synchronously — all four
|
||||
/// take `reload_mutex`, which is not reentrant.
|
||||
after_load_hook: ?ReloadHook,
|
||||
|
||||
/// Set by the composition root right after `init`, with `diagnostics`.
|
||||
@@ -168,10 +172,17 @@ pub const CertStore = struct {
|
||||
alpn: ?[*:null]const ?[*:0]const u8,
|
||||
) ReloadError!CertStore {
|
||||
const first = try load(gpa, io, cert_path, key_path, alpn);
|
||||
errdefer {
|
||||
first.entry.ctx.deinit(gpa);
|
||||
gpa.destroy(first.entry);
|
||||
}
|
||||
const owned_cert = try gpa.dupe(u8, cert_path);
|
||||
errdefer gpa.free(owned_cert);
|
||||
const owned_key = try gpa.dupe(u8, key_path);
|
||||
return .{
|
||||
.gpa = gpa,
|
||||
.cert_path = cert_path,
|
||||
.key_path = key_path,
|
||||
.cert_path = owned_cert,
|
||||
.key_path = owned_key,
|
||||
.alpn = alpn,
|
||||
.mutex = .init,
|
||||
.reload_mutex = .init,
|
||||
@@ -193,6 +204,8 @@ pub const CertStore = struct {
|
||||
std.debug.assert(current.refs == 0);
|
||||
self.mutex.unlock(io);
|
||||
self.destroyEntry(current);
|
||||
self.gpa.free(self.cert_path);
|
||||
self.gpa.free(self.key_path);
|
||||
self.* = undefined;
|
||||
}
|
||||
|
||||
@@ -226,7 +239,14 @@ pub const CertStore = struct {
|
||||
pub fn reload(self: *CertStore, io: std.Io) ReloadError!void {
|
||||
self.reload_mutex.lockUncancelable(io);
|
||||
defer self.reload_mutex.unlock(io);
|
||||
return self.reloadLocked(io);
|
||||
}
|
||||
|
||||
/// `reload`'s body, for callers that already hold `reload_mutex` —
|
||||
/// `pollOnce` does, because it must read `cert_path`/`key_path` under the
|
||||
/// same lock an apply replaces them under. `std.Io.Mutex` is not
|
||||
/// reentrant, so this exists rather than a recursive `reload` call.
|
||||
fn reloadLocked(self: *CertStore, io: std.Io) ReloadError!void {
|
||||
const next = load(self.gpa, io, self.cert_path, self.key_path, self.alpn) catch |err| {
|
||||
_ = self.reload_failures.fetchAdd(1, .monotonic);
|
||||
return err;
|
||||
@@ -247,6 +267,93 @@ pub const CertStore = struct {
|
||||
self.last_reload_unix.store(std.Io.Clock.real.now(io).toSeconds(), .monotonic);
|
||||
}
|
||||
|
||||
/// A candidate certificate loaded from new paths, not yet published.
|
||||
/// Holding one means holding `reload_mutex`: exactly one of
|
||||
/// `publishPathChange` or `abortPathChange` must follow, and it releases
|
||||
/// the lock.
|
||||
pub const PreparedPaths = struct {
|
||||
cert_path: []u8,
|
||||
key_path: []u8,
|
||||
loaded: Loaded,
|
||||
/// Read at prepare so publish reads no clock: publish must touch
|
||||
/// nothing outside memory it already owns.
|
||||
loaded_at_unix: i64,
|
||||
};
|
||||
|
||||
/// Prepare half of a `doh_server`/`dot_server` cert-path change: takes
|
||||
/// `reload_mutex` and loads the certificate and key from the NEW paths.
|
||||
/// Nothing is published, so a failure leaves the store exactly as it was —
|
||||
/// the old certificate keeps serving and the caller writes no DB row. The
|
||||
/// lock is released on failure and held on success, which is what makes
|
||||
/// the whole apply serialized against `reload` and `pollOnce`.
|
||||
pub fn preparePathChange(
|
||||
self: *CertStore,
|
||||
io: std.Io,
|
||||
cert_path: []const u8,
|
||||
key_path: []const u8,
|
||||
) ReloadError!PreparedPaths {
|
||||
self.reload_mutex.lockUncancelable(io);
|
||||
errdefer self.reload_mutex.unlock(io);
|
||||
|
||||
const owned_cert = try self.gpa.dupe(u8, cert_path);
|
||||
errdefer self.gpa.free(owned_cert);
|
||||
const owned_key = try self.gpa.dupe(u8, key_path);
|
||||
errdefer self.gpa.free(owned_key);
|
||||
|
||||
const next = load(self.gpa, io, cert_path, key_path, self.alpn) catch |err| {
|
||||
_ = self.reload_failures.fetchAdd(1, .monotonic);
|
||||
return err;
|
||||
};
|
||||
return .{
|
||||
.cert_path = owned_cert,
|
||||
.key_path = owned_key,
|
||||
.loaded = next,
|
||||
.loaded_at_unix = std.Io.Clock.real.now(io).toSeconds(),
|
||||
};
|
||||
}
|
||||
|
||||
/// Publish half: infallible and I/O-free. The paths and the generation
|
||||
/// they were loaded from are installed together — the generation `mutex`
|
||||
/// is taken only for that swap, inside `reload_mutex`, the same lock order
|
||||
/// `reload` uses. Releases `reload_mutex`.
|
||||
///
|
||||
/// Connections that pinned the old generation finish on the old
|
||||
/// certificate; the old entry is freed once its last reader releases.
|
||||
pub fn publishPathChange(self: *CertStore, io: std.Io, prepared: PreparedPaths) void {
|
||||
const old_cert_path = self.cert_path;
|
||||
const old_key_path = self.key_path;
|
||||
|
||||
self.mutex.lockUncancelable(io);
|
||||
const old = self.current;
|
||||
self.cert_path = prepared.cert_path;
|
||||
self.key_path = prepared.key_path;
|
||||
self.current = prepared.loaded.entry;
|
||||
self.loaded = prepared.loaded.sig;
|
||||
old.retired = true;
|
||||
const free_old = old.refs == 0;
|
||||
self.mutex.unlock(io);
|
||||
|
||||
// The branch runs before the counter bump, never after: a `bool` still
|
||||
// live across an atomic read-modify-write is the zig 0.16.0 Debug
|
||||
// miscompile AGENTS.md documents.
|
||||
if (free_old) self.destroyEntry(old);
|
||||
self.gpa.free(old_cert_path);
|
||||
self.gpa.free(old_key_path);
|
||||
|
||||
_ = self.reloads.fetchAdd(1, .monotonic);
|
||||
self.last_reload_unix.store(prepared.loaded_at_unix, .monotonic);
|
||||
self.reload_mutex.unlock(io);
|
||||
}
|
||||
|
||||
/// Discards a prepared candidate — the commit that would have published it
|
||||
/// failed, or a sibling owner's prepare did. Releases `reload_mutex`.
|
||||
pub fn abortPathChange(self: *CertStore, io: std.Io, prepared: PreparedPaths) void {
|
||||
self.gpa.free(prepared.cert_path);
|
||||
self.gpa.free(prepared.key_path);
|
||||
self.destroyEntry(prepared.loaded.entry);
|
||||
self.reload_mutex.unlock(io);
|
||||
}
|
||||
|
||||
/// Sleep first: `init` just loaded the files this poll would compare
|
||||
/// against. `.boot` so a suspended box still sees the interval elapse.
|
||||
pub fn watch(self: *CertStore, io: std.Io) std.Io.Cancelable!void {
|
||||
@@ -265,7 +372,15 @@ pub const CertStore = struct {
|
||||
/// changed, and the old one keeps serving either way. A failed reload
|
||||
/// warns and counts (`reload_failures`); the signature stays at the loaded
|
||||
/// pair, so every subsequent poll retries until the files parse.
|
||||
///
|
||||
/// The whole pass runs under `reload_mutex`: the paths it stats are the
|
||||
/// ones an apply replaces, and a poll that read a path outside the lock
|
||||
/// could stat a freed slice or reload a pair that was never published
|
||||
/// together.
|
||||
pub fn pollOnce(self: *CertStore, io: std.Io, now_s: i64) void {
|
||||
self.reload_mutex.lockUncancelable(io);
|
||||
defer self.reload_mutex.unlock(io);
|
||||
|
||||
const cert_sig = statSig(io, self.cert_path) catch |err| {
|
||||
log.warn("stat {s} failed; keeping the loaded certificate", .{self.cert_path});
|
||||
self.reportReload(io, now_s, "stat of the certificate failed", @errorName(err));
|
||||
@@ -290,7 +405,7 @@ pub const CertStore = struct {
|
||||
return;
|
||||
}
|
||||
|
||||
if (self.reload(io)) {
|
||||
if (self.reloadLocked(io)) {
|
||||
log.info("certificate reloaded from {s}", .{self.cert_path});
|
||||
if (self.diagnostics) |store| store.resolve(io, now_s, .certificate_reload, @tagName(self.kind));
|
||||
} else |err| {
|
||||
@@ -993,3 +1108,154 @@ test "an unchanged poll closes the episode a transient stat failure opened" {
|
||||
try fx.count("SELECT count(*) FROM operational_events WHERE resolved_at IS NULL"),
|
||||
);
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// path apply (milestone-34 S3.4)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
test "a bad candidate is refused at prepare and the store keeps serving" {
|
||||
var env: TestEnv = undefined;
|
||||
try env.init();
|
||||
defer env.deinit();
|
||||
const io = env.io();
|
||||
|
||||
var store = try CertStore.init(testing.allocator, io, env.cert_path, env.key_path, null);
|
||||
defer store.deinit(io);
|
||||
|
||||
const before = store.acquire(io);
|
||||
store.release(io, before);
|
||||
const before_paths_cert = store.cert_path;
|
||||
|
||||
try env.tmp.dir.writeFile(io, .{ .sub_path = "bad.pem", .data = "not a certificate" });
|
||||
var bad_buf: [128]u8 = undefined;
|
||||
const bad_path = try std.fmt.bufPrint(&bad_buf, ".zig-cache/tmp/{s}/bad.pem", .{env.tmp.sub_path});
|
||||
|
||||
try testing.expectError(
|
||||
error.CertParse,
|
||||
store.preparePathChange(io, bad_path, env.key_path),
|
||||
);
|
||||
|
||||
// Nothing published: same generation, same paths, and `reload_mutex` was
|
||||
// released — a second prepare would deadlock otherwise.
|
||||
const after = store.acquire(io);
|
||||
store.release(io, after);
|
||||
try testing.expectEqual(before, after);
|
||||
try testing.expectEqual(before_paths_cert.ptr, store.cert_path.ptr);
|
||||
try testing.expectEqualStrings(env.cert_path, store.cert_path);
|
||||
try testing.expectEqual(@as(u64, 0), store.snapshotStats().reloads);
|
||||
try testing.expectEqual(@as(u64, 1), store.snapshotStats().reload_failures);
|
||||
|
||||
// A missing candidate path is refused the same way.
|
||||
try testing.expectError(
|
||||
error.CertUnreadable,
|
||||
store.preparePathChange(io, "./nxdns-no-such-cert-4a11.pem", env.key_path),
|
||||
);
|
||||
}
|
||||
|
||||
test "a published path change installs the new pair and retires the old generation" {
|
||||
var env: TestEnv = undefined;
|
||||
try env.init();
|
||||
defer env.deinit();
|
||||
const io = env.io();
|
||||
|
||||
var store = try CertStore.init(testing.allocator, io, env.cert_path, env.key_path, null);
|
||||
defer store.deinit(io);
|
||||
|
||||
// A second, byte-different copy of the same valid pair under new names.
|
||||
const grown = try std.mem.concat(testing.allocator, u8, &.{ fixtures.cert_pem, "\n" });
|
||||
defer testing.allocator.free(grown);
|
||||
try env.tmp.dir.writeFile(io, .{ .sub_path = "next-cert.pem", .data = grown });
|
||||
try env.tmp.dir.writeFile(io, .{ .sub_path = "next-key.pem", .data = fixtures.key_pem });
|
||||
var cert_buf: [128]u8 = undefined;
|
||||
var key_buf: [128]u8 = undefined;
|
||||
const next_cert = try std.fmt.bufPrint(&cert_buf, ".zig-cache/tmp/{s}/next-cert.pem", .{env.tmp.sub_path});
|
||||
const next_key = try std.fmt.bufPrint(&key_buf, ".zig-cache/tmp/{s}/next-key.pem", .{env.tmp.sub_path});
|
||||
|
||||
// A connection pinned to the old generation finishes on it.
|
||||
const pinned = store.acquire(io);
|
||||
|
||||
const prepared = try store.preparePathChange(io, next_cert, next_key);
|
||||
store.publishPathChange(io, prepared);
|
||||
|
||||
try testing.expectEqualStrings(next_cert, store.cert_path);
|
||||
try testing.expectEqualStrings(next_key, store.key_path);
|
||||
try testing.expectEqual(@as(u64, 1), store.snapshotStats().reloads);
|
||||
try testing.expect(pinned.retired);
|
||||
|
||||
const serving = store.acquire(io);
|
||||
try testing.expect(serving != pinned);
|
||||
store.release(io, serving);
|
||||
store.release(io, pinned);
|
||||
|
||||
// The watcher now measures the new pair, so an untouched pair polls clean
|
||||
// and a rewritten one reloads.
|
||||
store.pollOnce(io, 1_000);
|
||||
try testing.expectEqual(@as(u64, 1), store.snapshotStats().reloads);
|
||||
try env.tmp.dir.writeFile(io, .{ .sub_path = "next-cert.pem", .data = fixtures.cert_pem });
|
||||
store.pollOnce(io, 1_100);
|
||||
try testing.expectEqual(@as(u64, 2), store.snapshotStats().reloads);
|
||||
}
|
||||
|
||||
test "an aborted path change frees the candidate and leaves the store untouched" {
|
||||
var env: TestEnv = undefined;
|
||||
try env.init();
|
||||
defer env.deinit();
|
||||
const io = env.io();
|
||||
|
||||
var store = try CertStore.init(testing.allocator, io, env.cert_path, env.key_path, null);
|
||||
defer store.deinit(io);
|
||||
|
||||
const before = store.acquire(io);
|
||||
store.release(io, before);
|
||||
|
||||
// The commit this candidate was built for failed; the testing allocator
|
||||
// proves the abort frees everything the prepare took.
|
||||
const prepared = try store.preparePathChange(io, env.cert_path, env.key_path);
|
||||
store.abortPathChange(io, prepared);
|
||||
|
||||
const after = store.acquire(io);
|
||||
store.release(io, after);
|
||||
try testing.expectEqual(before, after);
|
||||
try testing.expectEqual(@as(u64, 0), store.snapshotStats().reloads);
|
||||
|
||||
// `reload_mutex` came back, so the store still reloads.
|
||||
try store.reload(io);
|
||||
}
|
||||
|
||||
test "a reload racing a path apply is serialized behind it" {
|
||||
var env: TestEnv = undefined;
|
||||
try env.init();
|
||||
defer env.deinit();
|
||||
const io = env.io();
|
||||
|
||||
var store = try CertStore.init(testing.allocator, io, env.cert_path, env.key_path, null);
|
||||
defer store.deinit(io);
|
||||
|
||||
const grown = try std.mem.concat(testing.allocator, u8, &.{ fixtures.cert_pem, "\n" });
|
||||
defer testing.allocator.free(grown);
|
||||
try env.tmp.dir.writeFile(io, .{ .sub_path = "next-cert.pem", .data = grown });
|
||||
try env.tmp.dir.writeFile(io, .{ .sub_path = "next-key.pem", .data = fixtures.key_pem });
|
||||
var cert_buf: [128]u8 = undefined;
|
||||
var key_buf: [128]u8 = undefined;
|
||||
const next_cert = try std.fmt.bufPrint(&cert_buf, ".zig-cache/tmp/{s}/next-cert.pem", .{env.tmp.sub_path});
|
||||
const next_key = try std.fmt.bufPrint(&key_buf, ".zig-cache/tmp/{s}/next-key.pem", .{env.tmp.sub_path});
|
||||
|
||||
// Prepare holds `reload_mutex` across the whole apply.
|
||||
const prepared = try store.preparePathChange(io, next_cert, next_key);
|
||||
try testing.expect(!store.reload_mutex.tryLock());
|
||||
|
||||
// The concurrent reload cannot start, so it cannot publish the OLD paths
|
||||
// over the new generation.
|
||||
var racing = try io.concurrent(CertStore.reload, .{ &store, io });
|
||||
store.publishPathChange(io, prepared);
|
||||
try racing.await(io);
|
||||
|
||||
// Two publications, and the last word is the apply's pair: the racing
|
||||
// reload reread the paths the apply installed.
|
||||
try testing.expectEqual(@as(u64, 2), store.snapshotStats().reloads);
|
||||
try testing.expectEqualStrings(next_cert, store.cert_path);
|
||||
store.mutex.lockUncancelable(io);
|
||||
const final = store.loaded;
|
||||
store.mutex.unlock(io);
|
||||
try testing.expectEqual(@as(u64, grown.len), final.cert.size);
|
||||
}
|
||||
|
||||
+59
-18
@@ -27,6 +27,7 @@ const db = @import("../storage/db.zig");
|
||||
const disk_monitor = @import("../storage/disk_monitor.zig");
|
||||
const events = @import("../storage/events.zig");
|
||||
const logger = @import("../storage/logger.zig");
|
||||
const retention = @import("../storage/retention.zig");
|
||||
|
||||
const log = std.log.scoped(.clients);
|
||||
|
||||
@@ -59,7 +60,7 @@ pub const Tracker = struct {
|
||||
/// Guards `pending`, `count`, `passes` and `stats`. Every field below is
|
||||
/// written under it, so a reader takes it too; see `snapshotStats`.
|
||||
mutex: std.Io.Mutex,
|
||||
retention_days: u16,
|
||||
retention_days: *const retention.RetentionDays,
|
||||
pending: [max_pending]Pending,
|
||||
count: u32,
|
||||
passes: u64,
|
||||
@@ -70,8 +71,9 @@ pub const Tracker = struct {
|
||||
|
||||
/// `retention_days` is `logging.retention_days`, the same knob the query log
|
||||
/// prunes by (milestone-7 ruling 16). A client silent for that long is as
|
||||
/// uninteresting as a query that old.
|
||||
pub fn init(retention_days: u16) Tracker {
|
||||
/// uninteresting as a query that old — so both consumers share ONE cell
|
||||
/// and a settings apply moves them together.
|
||||
pub fn init(retention_days: *const retention.RetentionDays) Tracker {
|
||||
return .{
|
||||
.mutex = .init,
|
||||
.retention_days = retention_days,
|
||||
@@ -225,7 +227,7 @@ pub const Tracker = struct {
|
||||
self.mutex.unlock(io);
|
||||
|
||||
if (due) {
|
||||
const cutoff = now_s - @as(i64, self.retention_days) * 86_400;
|
||||
const cutoff = now_s - self.retention_days.seconds();
|
||||
if (clients_repo.pruneStale(database, cutoff)) |deleted| {
|
||||
self.mutex.lockUncancelable(io);
|
||||
self.stats.pruned += deleted;
|
||||
@@ -328,7 +330,8 @@ test "a client tracked twice before a flush yields one row at the later time" {
|
||||
var database = try openMigrated();
|
||||
defer database.close();
|
||||
|
||||
var tracker: Tracker = .init(30);
|
||||
var days: retention.RetentionDays = .init(30);
|
||||
var tracker: Tracker = .init(&days);
|
||||
// A table with room reports no drops.
|
||||
try testing.expectEqual(@as(u64, 0), tracker.trackAt(io, parsed("192.168.1.10"), 1700000000));
|
||||
try testing.expectEqual(@as(u64, 0), tracker.trackAt(io, parsed("192.168.1.10"), 1700000030));
|
||||
@@ -354,7 +357,8 @@ test "distinct clients each get a row and ipv6 text is canonical" {
|
||||
var database = try openMigrated();
|
||||
defer database.close();
|
||||
|
||||
var tracker: Tracker = .init(30);
|
||||
var days: retention.RetentionDays = .init(30);
|
||||
var tracker: Tracker = .init(&days);
|
||||
_ = tracker.trackAt(io, parsed("192.168.1.10"), 1700000000);
|
||||
_ = tracker.trackAt(io, parsed("192.168.1.11"), 1700000001);
|
||||
_ = tracker.trackAt(io, parsed("fd00:0:0:0:0:0:0:1"), 1700000002);
|
||||
@@ -378,7 +382,8 @@ test "a full table drops further clients and counts them" {
|
||||
var database = try openMigrated();
|
||||
defer database.close();
|
||||
|
||||
var tracker: Tracker = .init(30);
|
||||
var days: retention.RetentionDays = .init(30);
|
||||
var tracker: Tracker = .init(&days);
|
||||
for (0..Tracker.max_pending) |i| {
|
||||
var octets: [4]u8 = undefined;
|
||||
std.mem.writeInt(u32, &octets, @intCast(i), .big);
|
||||
@@ -421,7 +426,8 @@ test "a flush touches a hand-edited row without changing what the operator set"
|
||||
\\VALUES ('192.168.1.10', 'laptop', 2, 1, 1690000000, 1690000000);
|
||||
);
|
||||
|
||||
var tracker: Tracker = .init(30);
|
||||
var days: retention.RetentionDays = .init(30);
|
||||
var tracker: Tracker = .init(&days);
|
||||
_ = tracker.trackAt(io, parsed("192.168.1.10"), 1700000000);
|
||||
tracker.flushOnce(io, &database, true, null);
|
||||
|
||||
@@ -449,7 +455,8 @@ test "a gated pass writes nothing and keeps the pending clients" {
|
||||
monitor.state_raw.store(@intFromEnum(disk_monitor.State.critical), .monotonic);
|
||||
try testing.expect(!monitor.writesAllowed());
|
||||
|
||||
var tracker: Tracker = .init(30);
|
||||
var days: retention.RetentionDays = .init(30);
|
||||
var tracker: Tracker = .init(&days);
|
||||
_ = tracker.trackAt(io, parsed("192.168.1.10"), 1700000000);
|
||||
tracker.flushOnce(io, &database, monitor.writesAllowed(), null);
|
||||
|
||||
@@ -477,7 +484,8 @@ test "a failing upsert counts and leaves the client to be tracked again" {
|
||||
\\BEGIN SELECT RAISE(ABORT, 'refused'); END;
|
||||
);
|
||||
|
||||
var tracker: Tracker = .init(30);
|
||||
var days: retention.RetentionDays = .init(30);
|
||||
var tracker: Tracker = .init(&days);
|
||||
_ = tracker.trackAt(io, parsed("192.168.1.10"), 1700000000);
|
||||
_ = tracker.trackAt(io, parsed("192.168.1.11"), 1700000000);
|
||||
tracker.flushOnce(io, &database, true, null);
|
||||
@@ -507,7 +515,8 @@ test "the pass that comes due prunes the clients that went quiet" {
|
||||
try clients_repo.upsertSeen(&database, "10.0.0.1", now - 40 * day);
|
||||
try clients_repo.upsertSeen(&database, "10.0.0.2", now - 29 * day);
|
||||
|
||||
var tracker: Tracker = .init(30);
|
||||
var days: retention.RetentionDays = .init(30);
|
||||
var tracker: Tracker = .init(&days);
|
||||
// Every pass before the due one leaves both rows alone.
|
||||
for (0..Tracker.prune_every_passes - 1) |_| {
|
||||
tracker.flushOnce(io, &database, true, null);
|
||||
@@ -534,7 +543,8 @@ test "a shorter retention prunes what the default keeps" {
|
||||
const now = std.Io.Clock.real.now(io).toSeconds();
|
||||
try clients_repo.upsertSeen(&database, "10.0.0.1", now - 3 * 86_400);
|
||||
|
||||
var tracker: Tracker = .init(1);
|
||||
var days: retention.RetentionDays = .init(1);
|
||||
var tracker: Tracker = .init(&days);
|
||||
tracker.passes = Tracker.prune_every_passes - 1;
|
||||
tracker.flushOnce(io, &database, true, null);
|
||||
|
||||
@@ -542,6 +552,31 @@ test "a shorter retention prunes what the default keeps" {
|
||||
try testing.expectEqual(@as(u64, 1), tracker.snapshotStats(io).pruned);
|
||||
}
|
||||
|
||||
test "setRetentionDays changes the cutoff the next prune pass uses" {
|
||||
var threaded: std.Io.Threaded = .init(testing.allocator, .{});
|
||||
defer threaded.deinit();
|
||||
const io = threaded.io();
|
||||
|
||||
var database = try openMigrated();
|
||||
defer database.close();
|
||||
|
||||
const now = std.Io.Clock.real.now(io).toSeconds();
|
||||
try clients_repo.upsertSeen(&database, "10.0.0.1", now - 3 * 86_400);
|
||||
|
||||
var days: retention.RetentionDays = .init(7);
|
||||
var tracker: Tracker = .init(&days);
|
||||
tracker.passes = Tracker.prune_every_passes - 1;
|
||||
tracker.flushOnce(io, &database, true, null);
|
||||
try testing.expectEqual(@as(i64, 1), try clients_repo.countClients(&database));
|
||||
|
||||
// The shared cell, not a copy taken at construction.
|
||||
days.setRetentionDays(1);
|
||||
tracker.passes = Tracker.prune_every_passes - 1;
|
||||
tracker.flushOnce(io, &database, true, null);
|
||||
try testing.expectEqual(@as(i64, 0), try clients_repo.countClients(&database));
|
||||
try testing.expectEqual(@as(u64, 1), tracker.snapshotStats(io).pruned);
|
||||
}
|
||||
|
||||
test "the run loop flushes on its interval and returns on cancel" {
|
||||
var threaded: std.Io.Threaded = .init(testing.allocator, .{});
|
||||
defer threaded.deinit();
|
||||
@@ -550,7 +585,8 @@ test "the run loop flushes on its interval and returns on cancel" {
|
||||
var database = try openMigrated();
|
||||
defer database.close();
|
||||
|
||||
var tracker: Tracker = .init(30);
|
||||
var days: retention.RetentionDays = .init(30);
|
||||
var tracker: Tracker = .init(&days);
|
||||
_ = tracker.trackAt(io, parsed("192.168.1.10"), 1700000000);
|
||||
|
||||
var future = try io.concurrent(Tracker.run, .{
|
||||
@@ -622,7 +658,8 @@ test "the drain lands before any exchange, and attempts stop at the cap" {
|
||||
names.exchange_fn = CountingExchange.exchange;
|
||||
CountingExchange.reset(&database);
|
||||
|
||||
var tracker: Tracker = .init(30);
|
||||
var days: retention.RetentionDays = .init(30);
|
||||
var tracker: Tracker = .init(&days);
|
||||
const pending = clients_repo.max_per_pass + 4;
|
||||
for (0..pending) |i| {
|
||||
_ = tracker.trackAt(io, .{ .ip4 = .{ 192, 168, 2, @intCast(i) } }, 1700000000);
|
||||
@@ -656,7 +693,8 @@ test "a row the due pass prunes is never asked about" {
|
||||
const now = std.Io.Clock.real.now(io).toSeconds();
|
||||
try clients_repo.upsertSeen(&database, "192.168.1.10", now - 40 * 86_400);
|
||||
|
||||
var tracker: Tracker = .init(30);
|
||||
var days: retention.RetentionDays = .init(30);
|
||||
var tracker: Tracker = .init(&days);
|
||||
tracker.passes = Tracker.prune_every_passes - 1;
|
||||
tracker.flushOnce(io, &database, true, &names);
|
||||
|
||||
@@ -681,7 +719,8 @@ test "a gated pass attempts no naming either" {
|
||||
CountingExchange.reset(&database);
|
||||
|
||||
try clients_repo.upsertSeen(&database, "192.168.1.10", 1700000000);
|
||||
var tracker: Tracker = .init(30);
|
||||
var days: retention.RetentionDays = .init(30);
|
||||
var tracker: Tracker = .init(&days);
|
||||
tracker.flushOnce(io, &database, false, &names);
|
||||
|
||||
try testing.expectEqual(@as(usize, 0), CountingExchange.calls);
|
||||
@@ -700,7 +739,8 @@ test "a failing materialise opens one episode per pass and a clean pass closes i
|
||||
try fx.init(io, 1000);
|
||||
defer fx.deinit();
|
||||
|
||||
var tracker: Tracker = .init(30);
|
||||
var days: retention.RetentionDays = .init(30);
|
||||
var tracker: Tracker = .init(&days);
|
||||
tracker.diagnostics = &fx.store;
|
||||
|
||||
try database.exec(
|
||||
@@ -742,7 +782,8 @@ test "a failing prune opens its own episode the next due pass closes" {
|
||||
try fx.init(io, 1000);
|
||||
defer fx.deinit();
|
||||
|
||||
var tracker: Tracker = .init(30);
|
||||
var days: retention.RetentionDays = .init(30);
|
||||
var tracker: Tracker = .init(&days);
|
||||
tracker.diagnostics = &fx.store;
|
||||
// One pass short of due, so the pass below is the pruning one.
|
||||
tracker.passes = Tracker.prune_every_passes - 1;
|
||||
|
||||
@@ -36,6 +36,7 @@ const listener = @import("listener.zig");
|
||||
const model = @import("../config/model.zig");
|
||||
const tls_server = @import("../platform/tls_server.zig");
|
||||
const transport = @import("../upstream/transport.zig");
|
||||
const upstream_owner = @import("../upstream/owner.zig");
|
||||
|
||||
pub const dns_query_path = "/dns-query";
|
||||
|
||||
@@ -619,6 +620,7 @@ const Harness = struct {
|
||||
store: cert_store.CertStore,
|
||||
tables: local_tables_mod.LocalTables,
|
||||
upstream: FailingUpstream,
|
||||
upstream_owner: upstream_owner.Borrowed,
|
||||
h: handler.Handler,
|
||||
server: DohServer,
|
||||
group: std.Io.Group,
|
||||
@@ -646,10 +648,10 @@ const Harness = struct {
|
||||
errdefer hx.tables.deinit(testing.allocator);
|
||||
|
||||
hx.upstream = .{};
|
||||
hx.upstream_owner = .{};
|
||||
hx.h = .{
|
||||
.upstream = hx.upstream.client(),
|
||||
.blocking = test_blocking,
|
||||
.forward_read_timeout = test_forward_timeout,
|
||||
.upstream = hx.upstream_owner.client(hx.upstream.client()),
|
||||
.policy = .{ .blocking = test_blocking, .forward_read_timeout = test_forward_timeout },
|
||||
.local_tables = &hx.tables,
|
||||
};
|
||||
|
||||
|
||||
+103
-9
@@ -28,6 +28,7 @@ const handler = @import("handler.zig");
|
||||
const listener = @import("listener.zig");
|
||||
const tls_server = @import("../platform/tls_server.zig");
|
||||
const transport = @import("../upstream/transport.zig");
|
||||
const upstream_owner = @import("../upstream/owner.zig");
|
||||
|
||||
/// Plaintext staging for `ServerStream`: the framing bytes and the decrypted
|
||||
/// record tail pass through here, while whole messages go straight to
|
||||
@@ -312,11 +313,10 @@ const forward_timeout: std.Io.Clock.Duration = .{
|
||||
/// An upstream and nothing else optional: no filtering, no cache, no log. The
|
||||
/// listener is what these tests exercise, so the handler is the same bare one
|
||||
/// its own tests use.
|
||||
fn bareHandler(client: transport.Client) handler.Handler {
|
||||
fn bareHandler(up: *upstream_owner.Owner) handler.Handler {
|
||||
return .{
|
||||
.upstream = client,
|
||||
.blocking = blocking,
|
||||
.forward_read_timeout = forward_timeout,
|
||||
.upstream = up,
|
||||
.policy = .{ .blocking = blocking, .forward_read_timeout = forward_timeout },
|
||||
};
|
||||
}
|
||||
|
||||
@@ -567,7 +567,8 @@ test "dot: two framed queries share one TLS connection" {
|
||||
defer env.deinit(io);
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = response_bytes };
|
||||
var h = bareHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = bareHandler(h_owner.client(fake.client()));
|
||||
|
||||
const listen_address: std.Io.net.IpAddress = try .parse("127.0.0.1", 0);
|
||||
var server = try DotServer.listen(gpa, io, listen_address, &h, &env.store, .{ .max_connections = 2 });
|
||||
@@ -607,7 +608,8 @@ test "dot: a transport EOF without close_notify is a connection error, not a cra
|
||||
defer env.deinit(io);
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = response_bytes };
|
||||
var h = bareHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = bareHandler(h_owner.client(fake.client()));
|
||||
|
||||
const listen_address: std.Io.net.IpAddress = try .parse("127.0.0.1", 0);
|
||||
var server = try DotServer.listen(gpa, io, listen_address, &h, &env.store, .{ .max_connections = 2 });
|
||||
@@ -645,7 +647,8 @@ test "dot: plain TCP bytes fail the handshake and are counted" {
|
||||
defer env.deinit(io);
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = response_bytes };
|
||||
var h = bareHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = bareHandler(h_owner.client(fake.client()));
|
||||
|
||||
const listen_address: std.Io.net.IpAddress = try .parse("127.0.0.1", 0);
|
||||
var server = try DotServer.listen(gpa, io, listen_address, &h, &env.store, .{ .max_connections = 2 });
|
||||
@@ -731,7 +734,8 @@ test "dot: a reload serves new handshakes without breaking the old connection" {
|
||||
defer env.deinit(io);
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = response_bytes };
|
||||
var h = bareHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = bareHandler(h_owner.client(fake.client()));
|
||||
|
||||
const listen_address: std.Io.net.IpAddress = try .parse("127.0.0.1", 0);
|
||||
var server = try DotServer.listen(gpa, io, listen_address, &h, &env.store, .{ .max_connections = 2 });
|
||||
@@ -759,6 +763,95 @@ test "dot: a reload serves new handshakes without breaking the old connection" {
|
||||
};
|
||||
}
|
||||
|
||||
/// S3.4 across a live listener: a cert PATH change — new files, not rewritten
|
||||
/// ones — serves the new certificate on the next handshake while the
|
||||
/// connection pinned to the old generation finishes on it.
|
||||
fn dotPathChangeServesNewCert(io: std.Io, address_: std.Io.net.IpAddress, env: *CertEnv) anyerror!void {
|
||||
var first: TestTls = undefined;
|
||||
try first.connect(io, address_);
|
||||
defer first.close(io);
|
||||
|
||||
try first.sendQuery();
|
||||
try expectAnswersQuery(try first.readReply());
|
||||
|
||||
const old_entry = env.store.acquire(io);
|
||||
defer env.store.release(io, old_entry);
|
||||
|
||||
try env.tmp.dir.writeFile(io, .{ .sub_path = "cert2.pem", .data = fixtures.cert2_pem });
|
||||
try env.tmp.dir.writeFile(io, .{ .sub_path = "key2.pem", .data = fixtures.key2_pem });
|
||||
var cert_buf: [128]u8 = undefined;
|
||||
var key_buf: [128]u8 = undefined;
|
||||
const next_cert = try std.fmt.bufPrint(&cert_buf, ".zig-cache/tmp/{s}/cert2.pem", .{env.tmp.sub_path});
|
||||
const next_key = try std.fmt.bufPrint(&key_buf, ".zig-cache/tmp/{s}/key2.pem", .{env.tmp.sub_path});
|
||||
|
||||
const prepared = try env.store.preparePathChange(io, next_cert, next_key);
|
||||
env.store.publishPathChange(io, prepared);
|
||||
|
||||
const new_entry = env.store.acquire(io);
|
||||
defer env.store.release(io, new_entry);
|
||||
try testing.expect(old_entry != new_entry);
|
||||
|
||||
var second: TestTls = undefined;
|
||||
try second.connect(io, address_);
|
||||
defer second.close(io);
|
||||
try second.sendQuery();
|
||||
try expectAnswersQuery(try second.readReply());
|
||||
|
||||
try first.sendQuery();
|
||||
try expectAnswersQuery(try first.readReply());
|
||||
|
||||
try second.client.end();
|
||||
try second.net_writer.interface.flush();
|
||||
var second_tail: [1]u8 = undefined;
|
||||
try testing.expectError(error.EndOfStream, second.client.reader.readSliceAll(&second_tail));
|
||||
|
||||
try first.client.end();
|
||||
try first.net_writer.interface.flush();
|
||||
var first_tail: [1]u8 = undefined;
|
||||
try testing.expectError(error.EndOfStream, first.client.reader.readSliceAll(&first_tail));
|
||||
}
|
||||
|
||||
test "dot: a cert path change serves the new certificate on the next handshake" {
|
||||
const build_options = @import("build_options");
|
||||
if (!build_options.integration) return error.SkipZigTest;
|
||||
|
||||
const gpa = testing.allocator;
|
||||
var threaded: std.Io.Threaded = .init(gpa, .{});
|
||||
defer threaded.deinit();
|
||||
const io = threaded.io();
|
||||
|
||||
var env: CertEnv = undefined;
|
||||
try env.init(io);
|
||||
defer env.deinit(io);
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = response_bytes };
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = bareHandler(h_owner.client(fake.client()));
|
||||
|
||||
const listen_address: std.Io.net.IpAddress = try .parse("127.0.0.1", 0);
|
||||
var server = try DotServer.listen(gpa, io, listen_address, &h, &env.store, .{ .max_connections = 2 });
|
||||
const server_address = server.boundAddress();
|
||||
|
||||
var group: std.Io.Group = .init;
|
||||
try group.concurrent(io, DotServer.serve, .{ &server, io });
|
||||
|
||||
try bounded(io, dotPathChangeServesNewCert, .{ io, server_address, &env });
|
||||
|
||||
const stats = server.snapshotStats();
|
||||
try testing.expectEqual(@as(u64, 2), stats.connections);
|
||||
try testing.expectEqual(@as(u64, 0), stats.tls_handshake_failures);
|
||||
try testing.expectEqual(@as(u64, 0), stats.connection_errors);
|
||||
|
||||
const store_stats = env.store.snapshotStats();
|
||||
try testing.expectEqual(@as(u64, 1), store_stats.reloads);
|
||||
try testing.expectEqual(@as(u64, 0), store_stats.reload_failures);
|
||||
|
||||
server.deinit(io);
|
||||
group.await(io) catch |err| switch (err) {
|
||||
error.Canceled => unreachable,
|
||||
};
|
||||
}
|
||||
|
||||
test "dot: an idle connection is closed with close_notify and counted" {
|
||||
const build_options = @import("build_options");
|
||||
if (!build_options.integration) return error.SkipZigTest;
|
||||
@@ -773,7 +866,8 @@ test "dot: an idle connection is closed with close_notify and counted" {
|
||||
defer env.deinit(io);
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = response_bytes };
|
||||
var h = bareHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = bareHandler(h_owner.client(fake.client()));
|
||||
|
||||
const listen_address: std.Io.net.IpAddress = try .parse("127.0.0.1", 0);
|
||||
var server = try DotServer.listen(gpa, io, listen_address, &h, &env.store, .{
|
||||
|
||||
+532
-142
File diff suppressed because it is too large
Load Diff
@@ -32,6 +32,7 @@ const forward_zones = @import("../local/forward_zones.zig");
|
||||
const handler = @import("handler.zig");
|
||||
const header = @import("../dns/header.zig");
|
||||
const local_tables = @import("local_tables.zig");
|
||||
const logger_controller = @import("../storage/logger_controller.zig");
|
||||
const logger_mod = @import("../storage/logger.zig");
|
||||
const provenance = @import("../storage/provenance.zig");
|
||||
const manager = @import("../filter/manager.zig");
|
||||
@@ -47,8 +48,10 @@ const rate_limiter = @import("rate_limiter.zig");
|
||||
const record = @import("../dns/record.zig");
|
||||
const records = @import("../local/records.zig");
|
||||
const response = @import("../filter/response.zig");
|
||||
const retention_mod = @import("../storage/retention.zig");
|
||||
const shutdown = @import("shutdown.zig");
|
||||
const transport = @import("../upstream/transport.zig");
|
||||
const upstream_owner = @import("../upstream/owner.zig");
|
||||
const types = @import("../dns/types.zig");
|
||||
const udp_server = @import("udp_server.zig");
|
||||
|
||||
@@ -79,11 +82,10 @@ const zone_ttl: u32 = 120;
|
||||
|
||||
/// The handler every case starts from: an upstream, the blocking options and
|
||||
/// the empty local tables. Each case wires in the collaborators it exercises.
|
||||
fn baseHandler(client: transport.Client) handler.Handler {
|
||||
fn baseHandler(up: *upstream_owner.Owner) handler.Handler {
|
||||
return .{
|
||||
.upstream = client,
|
||||
.blocking = blocking,
|
||||
.forward_read_timeout = forward_timeout,
|
||||
.upstream = up,
|
||||
.policy = .{ .blocking = blocking, .forward_read_timeout = forward_timeout },
|
||||
};
|
||||
}
|
||||
|
||||
@@ -265,6 +267,11 @@ fn fixtureManager(m: *manager.Manager, snapshot: *matcher.Snapshot) void {
|
||||
.paths = undefined,
|
||||
.fetcher = undefined,
|
||||
.update = .{},
|
||||
.schedule_mutex = .init,
|
||||
.schedule_version = 0,
|
||||
.schedule_anchor_s = null,
|
||||
.schedule_event = .unset,
|
||||
.schedule_clock = .real,
|
||||
.total_budget = forward_timeout,
|
||||
.lock = .init,
|
||||
.writer_lock = .init,
|
||||
@@ -317,10 +324,12 @@ test "S7 case 1: a blocked domain is answered with the zero address and logged"
|
||||
|
||||
var queue_buf: [log_queue_len]logger_mod.Entry = undefined;
|
||||
var lg: logger_mod.Logger = .init(.{}, &queue_buf);
|
||||
var sink: query_sink.QuerySink = .init(&lg, null);
|
||||
var log_owner: logger_controller.Borrowed = .{};
|
||||
var sink: query_sink.QuerySink = .init(log_owner.over(&lg), null);
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = .a };
|
||||
var h = baseHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = baseHandler(h_owner.client(fake.client()));
|
||||
h.manager = &mgr;
|
||||
h.sink = &sink;
|
||||
|
||||
@@ -374,7 +383,8 @@ test "S7 case 2: an allow rule beats the blocklist and the upstream answers" {
|
||||
fixtureManager(&mgr, &snapshot);
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = .a };
|
||||
var h = baseHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = baseHandler(h_owner.client(fake.client()));
|
||||
h.manager = &mgr;
|
||||
|
||||
var loop = try Loop.bind(gpa, io, &h);
|
||||
@@ -411,7 +421,8 @@ test "S7 case 3: a local record answers authoritatively without an upstream" {
|
||||
defer table.deinit(gpa);
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = .a };
|
||||
var h = baseHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = baseHandler(h_owner.client(fake.client()));
|
||||
var tables: local_tables.LocalTables = .{ .records = table };
|
||||
h.local_tables = &tables;
|
||||
|
||||
@@ -479,12 +490,13 @@ test "S7 case 4: a forward zone reaches its resolver, bypasses the blocklist and
|
||||
defer cache.deinit();
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = .a };
|
||||
var h = baseHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = baseHandler(h_owner.client(fake.client()));
|
||||
h.manager = &mgr;
|
||||
var tables: local_tables.LocalTables = .{ .zones = zones };
|
||||
h.local_tables = &tables;
|
||||
h.cache = &cache;
|
||||
h.negative_ttl_max = 3600;
|
||||
h.policy.negative_ttl_max = 3600;
|
||||
|
||||
var loop = try Loop.bind(gpa, io, &h);
|
||||
defer loop.stop(gpa, io);
|
||||
@@ -533,12 +545,14 @@ test "S7 case 5: a cached answer comes back with a fresh id, an aged ttl and a l
|
||||
|
||||
var queue_buf: [log_queue_len]logger_mod.Entry = undefined;
|
||||
var lg: logger_mod.Logger = .init(.{}, &queue_buf);
|
||||
var sink: query_sink.QuerySink = .init(&lg, null);
|
||||
var log_owner: logger_controller.Borrowed = .{};
|
||||
var sink: query_sink.QuerySink = .init(log_owner.over(&lg), null);
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = .a };
|
||||
var h = baseHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = baseHandler(h_owner.client(fake.client()));
|
||||
h.cache = &cache;
|
||||
h.negative_ttl_max = 3600;
|
||||
h.policy.negative_ttl_max = 3600;
|
||||
h.sink = &sink;
|
||||
|
||||
var loop = try Loop.bind(gpa, io, &h);
|
||||
@@ -622,10 +636,12 @@ test "S7 case 6: a cname into a blocked target blocks the original question" {
|
||||
|
||||
var queue_buf: [log_queue_len]logger_mod.Entry = undefined;
|
||||
var lg: logger_mod.Logger = .init(.{}, &queue_buf);
|
||||
var sink: query_sink.QuerySink = .init(&lg, null);
|
||||
var log_owner: logger_controller.Borrowed = .{};
|
||||
var sink: query_sink.QuerySink = .init(log_owner.over(&lg), null);
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = .{ .cname = "tracker.example.org" } };
|
||||
var h = baseHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = baseHandler(h_owner.client(fake.client()));
|
||||
h.manager = &mgr;
|
||||
h.sink = &sink;
|
||||
|
||||
@@ -681,7 +697,8 @@ test "S7 case 7: safe search answers the original question with a cname to the t
|
||||
fixtureManager(&mgr, &snapshot);
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = .a };
|
||||
var h = baseHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = baseHandler(h_owner.client(fake.client()));
|
||||
h.manager = &mgr;
|
||||
|
||||
var loop = try Loop.bind(gpa, io, &h);
|
||||
@@ -733,10 +750,12 @@ test "S7 case 8: the third query inside the window is refused" {
|
||||
|
||||
var queue_buf: [log_queue_len]logger_mod.Entry = undefined;
|
||||
var lg: logger_mod.Logger = .init(.{}, &queue_buf);
|
||||
var sink: query_sink.QuerySink = .init(&lg, null);
|
||||
var log_owner: logger_controller.Borrowed = .{};
|
||||
var sink: query_sink.QuerySink = .init(log_owner.over(&lg), null);
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = .a };
|
||||
var h = baseHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = baseHandler(h_owner.client(fake.client()));
|
||||
h.limiter = &limiter;
|
||||
h.sink = &sink;
|
||||
|
||||
@@ -785,7 +804,8 @@ test "S7 case 9: pause lifts filtering and unpause restores it" {
|
||||
paused.pauseFor(std.Io.Clock.real.now(io).toSeconds(), null);
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = .a };
|
||||
var h = baseHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = baseHandler(h_owner.client(fake.client()));
|
||||
h.manager = &mgr;
|
||||
h.pause = &paused;
|
||||
|
||||
@@ -831,10 +851,12 @@ test "S7 case 10: the querying client is materialised as a row" {
|
||||
try db.applyPragmas(&database, .{});
|
||||
_ = try migrations.migrate(&database);
|
||||
|
||||
var tracker: clients.Tracker = .init(30);
|
||||
var retention_days: retention_mod.RetentionDays = .init(30);
|
||||
var tracker: clients.Tracker = .init(&retention_days);
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = .a };
|
||||
var h = baseHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = baseHandler(h_owner.client(fake.client()));
|
||||
h.tracker = &tracker;
|
||||
|
||||
var loop = try Loop.bind(gpa, io, &h);
|
||||
|
||||
@@ -13,24 +13,36 @@
|
||||
const std = @import("std");
|
||||
|
||||
const logger = @import("../storage/logger.zig");
|
||||
const logger_controller = @import("../storage/logger_controller.zig");
|
||||
const sse = @import("../web/sse.zig");
|
||||
|
||||
pub const QuerySink = struct {
|
||||
logger: *logger.Logger,
|
||||
/// The controller, not a `Logger`: `logging.query_log_buffer_max` can
|
||||
/// change while the server runs, and the generation a producer enqueues
|
||||
/// into has to be the one that is live at that moment.
|
||||
controller: *logger_controller.Controller,
|
||||
/// Null when `web.enabled` is false: nothing subscribes, so nothing needs
|
||||
/// a hub, and the DNS path pays one null check.
|
||||
hub: ?*sse.Hub,
|
||||
|
||||
pub fn init(query_logger: *logger.Logger, hub: ?*sse.Hub) QuerySink {
|
||||
return .{ .logger = query_logger, .hub = hub };
|
||||
pub fn init(controller: *logger_controller.Controller, hub: ?*sse.Hub) QuerySink {
|
||||
return .{ .controller = controller, .hub = hub };
|
||||
}
|
||||
|
||||
/// Transforms once, publishes, then enqueues. Never blocks the query path
|
||||
/// and never fails: both consumers drop rather than wait.
|
||||
///
|
||||
/// The borrow spans both halves. A resize that lands between them would
|
||||
/// otherwise leave this entry going into a queue retirement has already
|
||||
/// closed, and that is exactly the drop window the controller exists to
|
||||
/// make impossible.
|
||||
pub fn log(self: *QuerySink, io: std.Io, entry: logger.Entry) void {
|
||||
const transformed = self.logger.transformed(entry);
|
||||
const generation = self.controller.acquire(io);
|
||||
defer self.controller.release(io, generation);
|
||||
|
||||
const transformed = generation.logger.transformed(entry);
|
||||
if (self.hub) |hub| hub.publish(io, transformed);
|
||||
self.logger.logTransformed(io, transformed);
|
||||
generation.logger.logTransformed(io, transformed);
|
||||
}
|
||||
};
|
||||
|
||||
@@ -60,7 +72,8 @@ test "the sink publishes and logs the same entry" {
|
||||
|
||||
var queue_buf: [4]logger.Entry = undefined;
|
||||
var query_logger: logger.Logger = .init(.{}, &queue_buf);
|
||||
var sink: QuerySink = .init(&query_logger, hub);
|
||||
var owner: logger_controller.Borrowed = .{};
|
||||
var sink: QuerySink = .init(owner.over(&query_logger), hub);
|
||||
|
||||
const id = hub.subscribe(io).?;
|
||||
defer hub.unsubscribe(io, id);
|
||||
@@ -87,7 +100,8 @@ test "fanout does not depend on the entry reaching the queue" {
|
||||
|
||||
var queue_buf: [4]logger.Entry = undefined;
|
||||
var query_logger: logger.Logger = .init(.{}, &queue_buf);
|
||||
var sink: QuerySink = .init(&query_logger, hub);
|
||||
var owner: logger_controller.Borrowed = .{};
|
||||
var sink: QuerySink = .init(owner.over(&query_logger), hub);
|
||||
|
||||
const id = hub.subscribe(io).?;
|
||||
defer hub.unsubscribe(io, id);
|
||||
@@ -115,7 +129,8 @@ test "the privacy transforms run once, before both consumers" {
|
||||
.{ .hide_domains = true, .hide_client_ips = true },
|
||||
&queue_buf,
|
||||
);
|
||||
var sink: QuerySink = .init(&query_logger, hub);
|
||||
var owner: logger_controller.Borrowed = .{};
|
||||
var sink: QuerySink = .init(owner.over(&query_logger), hub);
|
||||
|
||||
const id = hub.subscribe(io).?;
|
||||
defer hub.unsubscribe(io, id);
|
||||
@@ -138,7 +153,8 @@ test "a sink without a hub still logs" {
|
||||
|
||||
var queue_buf: [4]logger.Entry = undefined;
|
||||
var query_logger: logger.Logger = .init(.{}, &queue_buf);
|
||||
var sink: QuerySink = .init(&query_logger, null);
|
||||
var owner: logger_controller.Borrowed = .{};
|
||||
var sink: QuerySink = .init(owner.over(&query_logger), null);
|
||||
|
||||
sink.log(io, sampleEntry(5, "nohub.example"));
|
||||
|
||||
|
||||
@@ -26,6 +26,7 @@ const types = @import("../dns/types.zig");
|
||||
const health = @import("../upstream/health.zig");
|
||||
const pool = @import("../upstream/pool.zig");
|
||||
const transport = @import("../upstream/transport.zig");
|
||||
const upstream_owner = @import("../upstream/owner.zig");
|
||||
|
||||
const testing = std.testing;
|
||||
|
||||
@@ -42,11 +43,10 @@ const forward_timeout: std.Io.Clock.Duration = .{
|
||||
/// An upstream and nothing else optional: no filtering, no cache, no log. The
|
||||
/// listeners and the pool are what this test exercises, so the handler is the
|
||||
/// same bare one its own tests use.
|
||||
fn bareHandler(client: transport.Client) handler.Handler {
|
||||
fn bareHandler(up: *upstream_owner.Owner) handler.Handler {
|
||||
return .{
|
||||
.upstream = client,
|
||||
.blocking = blocking,
|
||||
.forward_read_timeout = forward_timeout,
|
||||
.upstream = up,
|
||||
.policy = .{ .blocking = blocking, .forward_read_timeout = forward_timeout },
|
||||
};
|
||||
}
|
||||
|
||||
@@ -275,7 +275,9 @@ test "the whole resolver answers over udp and tcp and fails over to a healthy up
|
||||
};
|
||||
var upstreams: pool.Pool = .init(&entries, test_cfg, pool_timeouts, 1);
|
||||
|
||||
var h = bareHandler(upstreams.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
|
||||
var h = bareHandler(h_owner.client(upstreams.client()));
|
||||
|
||||
const listen_address: net.IpAddress = try .parse("127.0.0.1", 0);
|
||||
var udp = try udp_server.UdpServer.bind(gpa, io, listen_address, &h, .{ .max_in_flight = 4 });
|
||||
|
||||
@@ -21,6 +21,7 @@ const header = @import("../dns/header.zig");
|
||||
const packet = @import("../dns/packet.zig");
|
||||
const types = @import("../dns/types.zig");
|
||||
const transport = @import("../upstream/transport.zig");
|
||||
const upstream_owner = @import("../upstream/owner.zig");
|
||||
|
||||
const testing = std.testing;
|
||||
|
||||
@@ -37,11 +38,10 @@ const forward_timeout: std.Io.Clock.Duration = .{
|
||||
/// An upstream and nothing else optional: no filtering, no cache, no log. The
|
||||
/// listener is what these tests exercise, so the handler is the same bare one
|
||||
/// its own tests use.
|
||||
fn bareHandler(client: transport.Client) handler.Handler {
|
||||
fn bareHandler(up: *upstream_owner.Owner) handler.Handler {
|
||||
return .{
|
||||
.upstream = client,
|
||||
.blocking = blocking,
|
||||
.forward_read_timeout = forward_timeout,
|
||||
.upstream = up,
|
||||
.policy = .{ .blocking = blocking, .forward_read_timeout = forward_timeout },
|
||||
};
|
||||
}
|
||||
|
||||
@@ -175,7 +175,8 @@ test "two length-prefixed queries share one connection" {
|
||||
const io = threaded.io();
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = response_bytes };
|
||||
var h = bareHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = bareHandler(h_owner.client(fake.client()));
|
||||
|
||||
const listen_address: net.IpAddress = try .parse("127.0.0.1", 0);
|
||||
var server = try tcp_server.TcpServer.listen(gpa, io, listen_address, &h, .{ .max_connections = 2 });
|
||||
@@ -205,7 +206,8 @@ test "the claimed slot records the connecting client" {
|
||||
const io = threaded.io();
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = response_bytes };
|
||||
var h = bareHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = bareHandler(h_owner.client(fake.client()));
|
||||
|
||||
const listen_address: net.IpAddress = try .parse("127.0.0.1", 0);
|
||||
var server = try tcp_server.TcpServer.listen(gpa, io, listen_address, &h, .{ .max_connections = 2 });
|
||||
@@ -282,7 +284,8 @@ test "a canceled serve does not wait for a live connection" {
|
||||
const io = threaded.io();
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = response_bytes };
|
||||
var h = bareHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = bareHandler(h_owner.client(fake.client()));
|
||||
|
||||
const listen_address: net.IpAddress = try .parse("127.0.0.1", 0);
|
||||
// The idle budget is the whole time a drain would have to wait out, so it
|
||||
@@ -344,7 +347,8 @@ test "an idle connection is closed and counted" {
|
||||
const io = threaded.io();
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = response_bytes };
|
||||
var h = bareHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = bareHandler(h_owner.client(fake.client()));
|
||||
|
||||
const listen_address: net.IpAddress = try .parse("127.0.0.1", 0);
|
||||
var server = try tcp_server.TcpServer.listen(gpa, io, listen_address, &h, .{
|
||||
@@ -377,7 +381,8 @@ test "a zero-length message is a connection error" {
|
||||
const io = threaded.io();
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = response_bytes };
|
||||
var h = bareHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = bareHandler(h_owner.client(fake.client()));
|
||||
|
||||
const listen_address: net.IpAddress = try .parse("127.0.0.1", 0);
|
||||
var server = try tcp_server.TcpServer.listen(gpa, io, listen_address, &h, .{
|
||||
@@ -429,7 +434,8 @@ test "deinit ends a serve loop that is blocked on accept" {
|
||||
const io = threaded.io();
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = response_bytes };
|
||||
var h = bareHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = bareHandler(h_owner.client(fake.client()));
|
||||
|
||||
const listen_address: net.IpAddress = try .parse("127.0.0.1", 0);
|
||||
var server = try tcp_server.TcpServer.listen(gpa, io, listen_address, &h, .{ .max_connections = 2 });
|
||||
|
||||
@@ -20,6 +20,7 @@ const header = @import("../dns/header.zig");
|
||||
const packet = @import("../dns/packet.zig");
|
||||
const types = @import("../dns/types.zig");
|
||||
const transport = @import("../upstream/transport.zig");
|
||||
const upstream_owner = @import("../upstream/owner.zig");
|
||||
|
||||
const testing = std.testing;
|
||||
|
||||
@@ -36,11 +37,10 @@ const forward_timeout: std.Io.Clock.Duration = .{
|
||||
/// An upstream and nothing else optional: no filtering, no cache, no log. The
|
||||
/// listener is what these tests exercise, so the handler is the same bare one
|
||||
/// its own tests use.
|
||||
fn bareHandler(client: transport.Client) handler.Handler {
|
||||
fn bareHandler(up: *upstream_owner.Owner) handler.Handler {
|
||||
return .{
|
||||
.upstream = client,
|
||||
.blocking = blocking,
|
||||
.forward_read_timeout = forward_timeout,
|
||||
.upstream = up,
|
||||
.policy = .{ .blocking = blocking, .forward_read_timeout = forward_timeout },
|
||||
};
|
||||
}
|
||||
|
||||
@@ -112,7 +112,8 @@ test "a udp query is answered on the loopback" {
|
||||
const io = threaded.io();
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = response_bytes };
|
||||
var h = bareHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = bareHandler(h_owner.client(fake.client()));
|
||||
|
||||
const listen_address: net.IpAddress = try .parse("127.0.0.1", 0);
|
||||
var server = try udp_server.UdpServer.bind(gpa, io, listen_address, &h, .{ .max_in_flight = 4 });
|
||||
@@ -150,7 +151,8 @@ test "a runt datagram is dropped and no reply is sent" {
|
||||
const io = threaded.io();
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = response_bytes };
|
||||
var h = bareHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = bareHandler(h_owner.client(fake.client()));
|
||||
|
||||
const listen_address: net.IpAddress = try .parse("127.0.0.1", 0);
|
||||
var server = try udp_server.UdpServer.bind(gpa, io, listen_address, &h, .{ .max_in_flight = 4 });
|
||||
@@ -185,7 +187,8 @@ test "an oversize datagram arrives truncated and is dropped" {
|
||||
const io = threaded.io();
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = response_bytes };
|
||||
var h = bareHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = bareHandler(h_owner.client(fake.client()));
|
||||
|
||||
const listen_address: net.IpAddress = try .parse("127.0.0.1", 0);
|
||||
var server = try udp_server.UdpServer.bind(gpa, io, listen_address, &h, .{ .max_in_flight = 4 });
|
||||
@@ -226,7 +229,8 @@ test "deinit ends a serve loop that is blocked on receive" {
|
||||
const io = threaded.io();
|
||||
|
||||
var fake: FakeUpstream = .{ .reply = response_bytes };
|
||||
var h = bareHandler(fake.client());
|
||||
var h_owner: upstream_owner.Borrowed = .{};
|
||||
var h = bareHandler(h_owner.client(fake.client()));
|
||||
|
||||
const listen_address: net.IpAddress = try .parse("127.0.0.1", 0);
|
||||
var server = try udp_server.UdpServer.bind(gpa, io, listen_address, &h, .{ .max_in_flight = 4 });
|
||||
|
||||
Reference in New Issue
Block a user