Files
nxdns/specs/milestone-7.md

46 KiB

Milestone 7: serving pipeline composition (PLAN Phase 7)

Goal: nxdns run serves DNS end to end — rate limit → parse → group → local records → forward zones → filtering → safe-search → cache → upstream → CNAME uncloaking → cache store → async query log — with client auto-materialization, pause/resume, and a composition root that starts every background loop and shuts down cleanly on SIGINT/SIGTERM.

Ground truth: PLAN.md §4 (pipeline diagram, lines 151-173), §6 (DNS behavior), §7.2 (auto-materialization), §8 (cache key), §10 (rate limit), §11.4/§11.6 (logging, disk), Phase 7 text at PLAN.md:606-608. Verified Zig facts: specs/research/zig-0.16-api-notes.md plus the scratchpad note m7-explore-zig.md (signals, Io.Group, cancelable net ops) — the load-bearing facts are restated inline below where a session needs them.

Rulings (engineering calls resolving PLAN ambiguities — binding)

  1. Scope: Phase 7 owns nxdns run end to end: src/app.zig (composition root), src/server/shutdown.zig (signal → event), the cli.runRun body, and both listener edits.
  2. Signals: no signalfd/epoll/timerfd. std.posix.sigaction handler for INT and TERM that calls std.Io.Event.set on a file-scope event (set is async-signal-safe on the Threaded Linux backend: raw futex wake, no locks — Threaded.zig:1105,17235).
  3. Locking: DnsCache and RateLimiter each get one std.Io.Mutex owned by the handler. Household scale makes contention irrelevant; sharding is over-engineering.
  4. Null snapshot (Manager.acquire returns null before the first successful reload): fail open — skip group/filter/safe-search, use the default group semantics, count unfiltered_queries. DNS availability beats filtering for a household. The app attempts one synchronous reload before serving; a failure warns and serving starts anyway.
  5. Cache scope: one global cache, no group in the key (PLAN §8). Safe because blocked responses are never cached. Safe-search responses are not cached either (13 domains, rare; caching them under the rewritten name would force answer re-encoding on hit). Paused-mode upstream answers cache normally (they are ordinary upstream responses). Forward-zone responses cache normally (PLAN §6.5) under the same global key.
  6. Local vs forward zone: local records win (checked first; PLAN §6.4 puts them before everything but rate limit/parse).
  7. Forward zones bypass filtering, safe-search, and uncloaking (PLAN §6.5). They are logged with upstream = the resolver's canonical text.
  8. Rate limit position: after the 12-byte header read (the REFUSED reply needs the ID), before full parse. Refused queries are counted (Stats.refused), never query-logged. Localhost is not exempt (PLAN §10's localhost relaxation is the Phase 8 API limiter only).
  9. Non-IN qclass: skip filtering/cache/safe-search/uncloak, go straight upstream, log normally. Rare and harmless; special-casing CHAOS is scope creep.
  10. Safe-search mechanics (per safesearch.zig:56 doc): the outgoing query's QNAME becomes the target; the reply to the client keeps the original question, starts with a synthesized CNAME original→target, then copies only the A/AAAA answer records from the upstream response verbatim (their RDATA is raw bytes, so no compression-pointer re-encoding; other rtypes are skipped). Applies only on the upstream path — never to local records or forward zones.
  11. CNAME uncloaking (PLAN §6.3): walk the answer-section CNAME chain of the upstream response — no re-resolution. Depth 8 = at most 8 CNAME links followed. Each target is normalized and evaluated for the same group; a blocked target synthesizes a blocked response for the ORIGINAL qname/qtype. block_reason = "cname:" ++ @tagName(reason) (fits the 32-byte Entry field). Uncloaked-blocked responses are not cached. Skipped when paused, for forward zones, and for non-IN.
  12. Local CNAME records are returned as-is (writeAnswers emits the CNAME and stops, milestone-5 As-built). The client re-queries the target, which then flows through the full pipeline — that is how §6.4's "targets resolve through the normal pipeline" is satisfied.
  13. ECS (PLAN §6.1): strip (default) removes a client-sent ECS option from the outgoing query by rebuilding it (header + question + OPT sans option 8); other options are preserved. forward sends the query unchanged and puts the raw ECS option payload (≤ dns_cache.max_ecs_len = 40 bytes) into the cache key. If a rewritten query cannot fit the 512-byte build buffer (freak OPT payloads), the original is forwarded unchanged, Stats.ecs_strip_failed counts it, and the query is excluded from the cache in both directions (the subnet reached the resolver; the strip-mode key says nothing about it).
  14. Client auto-materialization (PLAN §7.2): never a hot-path DB write. The handler records the client in a mutex-guarded fixed-capacity pending table; a background loop flushes every 60 s with its OWN config.db connection: INSERT ... hand_edited=0 ON CONFLICT(ip) DO UPDATE last_seen. At most one write per client per flush. Table full → drop silently (counted).
  15. Snapshot staleness: a materialized client reaches the matcher snapshot only on the next manager reload. Correct behavior does not depend on it: groupForClient already falls back to prefix/default, and materialized rows exist for UI visibility (Phase 8 reads the DB directly). No reload is triggered per client.
  16. Client pruning (PLAN §7.2 + milestone-6 ruling 1): the same background loop prunes daily: DELETE FROM clients WHERE hand_edited=0 AND last_seen < now - retention_days (logging.retention_days, same knob as the query log — documented). retention.zig stays frozen and querylog-only. last_seen is the only signal (§3.6 forbids cross-DB joins).
  17. Disk gate: the client flush/prune loop skips its pass when writesAllowed() is false; the logger holds its batches. As built: milestone-6.md:410 assigned the manager's consumption to Phase 7 and no seam existed, so S6 added Manager.monitor: ?*disk_monitor.Monitor (set by the composition root after init, before runScheduler) plus a refreshes_gated counter read through refreshesGated(). Only scheduled passes are gated: runScheduler's periodic pass and the source-refresh half of the startup pass. reload and refreshAll stay ungated (operator actions). The startup pass's reloadLocked runs ahead of the gate: publishing an already-compiled snapshot is a read, and a full disk must not cost the household its filtering as well as its downloads. A gated pass logs one warn line and counts.
  18. Pause/resume (PLAN §13.1 mechanism): suspends FILTERING only (evaluate + safe-search + uncloak are skipped); everything else — local records, forward zones, cache, upstream, logging — runs unchanged. Global, not per-group. In-memory only; a restart resumes filtering (safe default). State: one std.atomic.Value(i64) — 0 = active, -1 = paused indefinitely, else = paused until that unix second (auto-resume by compare, no timer task). Phase 8's endpoints will call pauseFor/unpause on the shared pointer.
  19. Paused queries log blocked=false, block_reason=""; Stats.paused_queries counts them. Evaluation is skipped entirely (not evaluated-then-ignored).
  20. Log field semantics: response_time_us = wall time from handle entry to reply ready. cache_hit: true/false on the upstream path, null for local/forward/refused-never-logged; upstream: "" on cache hit, "local" for local records, resolver text for forward zones, the pool endpoint's canonical text for upstream answers (pool already tracks which entry answered — if unavailable without new plumbing, "" and a spec note; do not rebuild the pool for it: upstream may be "pool" as the compromise, S2 decides and records). block_reason = @tagName(matcher.Reason) truncated to 32, "cname:"-prefixed for uncloaking. Malformed/dropped queries are never query-logged (handler stats only).
  21. Connections: config.db — one opened by the app (migrate/bootstrap/readConfig, then handed to the manager for reloads), plus one dedicated to the client loop. querylog.db — one for the log writer, one for retention (retention.zig:79 contract). Four total.
  22. Shutdown order: shutdown event fires → logger.shutdown(io) (closes the queue; the writer drains and its task returns; a writer stuck in the disk gate holds its final batch — known, documented) → group.cancel(io) (requests cancel and joins every task) → deinit servers/pool/manager/state → close DBs → return. Canceled UDP recv / TCP accept return error.Canceled on the Threaded backend (Threaded.zig:12439,12877).
  23. No runtime config re-read in Phase 7. The blocklist scheduler reloads snapshots; everything else is restart-required until Phase 8 mutates the DB and calls reload.
  24. No new config fields, no schema change. Pause is in-memory; pruning reuses retention_days; ECS uses the existing edns.ecs_mode.
  25. Performance (PLAN §18): no formal benchmark in this milestone; the integration test asserts correctness only. The targets stay in PLAN for a later measurement pass.

Sessions

Dependency graph: S1 + S4 parallel (wave 1) → S2 (needs S1's edns API) → S3 (needs S2's signature in code) → S6 (needs S2, S3, S4) → S7 (needs S6). The orchestrator wires src/tests.zig after each wave.


Session S1: ECS helpers in dns/edns.zig

Owns: src/dns/edns.zig (extend in place; existing API is frozen — additive only).

S1.1 API

/// The raw payload of the ECS option (code 8) in `opt`, or null when absent.
/// `pkt` is the packet `opt` was parsed from. error.BadOption on malformed option lengths.
pub fn ecsPayload(pkt: packet.Packet, opt: OptRecord) error{BadOption}!?[]const u8

/// Rebuilds `query` without its ECS option into `out`: header + question copied verbatim,
/// then the OPT re-encoded with every option except code 8. Returns `.unchanged` when the
/// query has no OPT or no ECS option (nothing written), else `.rewritten` with the slice
/// into `out`. `pkt`/`opt` are the caller's already-parsed views of `query`.
/// error.Overflow when `out` is too small.
pub const StripResult = union(enum) { unchanged, rewritten: []u8 };
pub fn stripEcs(query: []const u8, pkt: packet.Packet, opt: OptRecord, out: []u8)
    error{ BadOption, Overflow }!StripResult

Constraints: pure, allocation-free, no Io. The rebuilt query keeps the original ID, flags and question bytes verbatim (copy query[0..opt_record_offset_of_question_end] — concretely: copy everything up to the OPT record's start; the OPT is the sole additional record in a query per the existing handler validation, and packet.parse yields its offset). ARCOUNT stays 1 when the OPT is re-emitted (it always is — only the option list shrinks). Use the existing encodeOpt for the re-emission. The RDATA walk follows the same bounds checks findOption uses.

S1.2 Acceptance

  • zig fmt --check + zig ast-check clean.
  • Tests in-file: ECS present → payload returned, stripped query parses cleanly and re-parse finds no option 8, other options survive byte-exact, DO bit and payload size survive; no OPT → .unchanged; OPT without ECS → .unchanged; malformed option length → error.BadOption; out one byte short → error.Overflow.
  • Temp-root run of the file's tests passes.

S1 As built

Two spec claims were wrong and the implementation corrects them. (1) The handler does NOT guarantee the OPT is the sole additional record (handler.zig:105 rejects only misplaced OPTs; packet.parse rejects only a second OPT) — stripEcs therefore copies the bytes after the OPT verbatim and all four section counts survive. Caveat, documented on the function: a compression pointer in a trailing record that targets a byte inside the OPT's option list shifts after the rewrite — such a pointer is already malformed (RFC 1035 §4.1.4) and the consequence is an upstream FORMERR. (2) packet.parse does not yield the OPT's start offset; a private optRecordStart re-walks the sections and matches the record on its full RDATA span, so a hand-built OptRecord cannot aim it at a different record (deriving the start backwards is forgeable: a 0x00xx transaction ID makes a pointer's low byte decode as a root name). encodeOpt takes one contiguous options slice, so stripEcs emits it with an empty list, streams the survivors, and patches RDLENGTH at opt_rdlength_offset (9); a guard test pins that offset against encodeOpt's real output. All options with code 8 are removed (a malformed query can carry several). Walk failures map to error.BadOption (frozen error set; reachable only for hand-assembled Packets, documented). After review: stripEcs asserts both that query aliases pkt.bytes and that out is pointer-range disjoint from query in both directions; the caller owns keeping the rebuild buffer separate from the received query buffer (S2's scratch.build satisfies this). The import is aliased packet_mod (seven existing declarations named packet). ecsPayload was built here, then removed in the review round when the handler's subnetForKey single-walk replaced its only caller (see S2 As-built); the forward-mode clamp policy lives there — a payload longer than dns_cache.max_ecs_len (40) makes the query uncacheable (truncating the key would merge distinct subnets). 7 S1 tests remain, 25 in file.


Session S2: pipeline in server/handler.zig (+ pause.zig)

Owns: src/server/handler.zig (rewrite around the existing skeleton), src/server/pause.zig (new). The biggest session. Everything below is frozen interface — S3 and S6 build against it.

S2.1 pause.zig

pub const Pause = struct {
    until: std.atomic.Value(i64) = .init(0),
    pub fn isPaused(self: *const Pause, now_s: i64) bool  // 0→false, -1→true, else now_s < until
    pub fn pauseFor(self: *Pause, now_s: i64, duration_s: ?u32) void  // null → -1 (indefinite)
    pub fn unpause(self: *Pause) void  // store 0
};

Auto-resume is the compare in isPaused — no timer. .monotonic ordering is sufficient (a lone i64 flag, no dependent data). Tests: the three states, expiry boundary (now == until is resumed), pauseFor overwrites a previous pause.

S2.2 Handler struct (frozen for S3/S6)

pub const Handler = struct {
    upstream: transport.Client,
    blocking: response.Options,
    ecs_mode: model.EcsMode = .strip,
    forward_read_timeout: std.Io.Clock.Duration,
    manager: ?*manager.Manager = null,
    records: *const records.Records,
    zones: *const forward_zones.Zones,
    cache: ?*dns_cache.DnsCache = null,
    cache_mutex: std.Io.Mutex = .init,
    limiter: ?*rate_limiter.RateLimiter = null,
    limiter_mutex: std.Io.Mutex = .init,
    logger: ?*logger_mod.Logger = null,
    pause: ?*pause.Pause = null,
    tracker: ?*clients.Tracker = null,
    stats: Stats = .{},
};

Every ?* defaults to null so existing tests and the check command construct a bare handler exactly as today (upstream + blocking + records/zones + timeout are the only required fields; records/zones may point at the shared Records.empty/Zones.empty constants — if Zones.empty does not exist yet, S2 adds it mirroring Records.empty).

pub const Scratch = struct {
    normalize: [types.max_name_len]u8,
    key: [dns_cache.max_key_len]u8,
    build: [512]u8,                          // rewritten outgoing query (ECS strip / safe-search)
    frame: [forward_client.min_frame_buf]u8, // forward-zone client frame buffer
};

pub fn handle(self: *Handler, io: std.Io, which: Transport, from: address.NetAddress,
    query: []const u8, response_buf: []u8, scratch: *Scratch) Outcome

Outcome unchanged. Still no error union — every failure is a synthesized rcode or a counted drop. New Stats counters (all std.atomic.Value(u64)): refused, blocked, uncloak_blocked, local_answers, forward_zone_answers, cache_hits, paused_queries, unfiltered_queries, safesearch_rewrites, ecs_strip_failed, tracker_full.

S2.3 Pipeline order inside handle (binding; rulings 4-13, 18-20 apply)

  1. Read the 12-byte header (existing code). QR/opcode/qdcount checks as today.
  2. Rate limit: limiter_mutex.lock(io), check(now, from.key()), unlock. Refused → synthesize REFUSED (reuse synthesize), count, return. now from std.Io.Clock.awake.now(io) (limiter windows are awake-anchored, milestone-6 As-built).
  3. Full parse + OPT validation (existing code). Record start_us once from std.Io.Clock.real.now(io) for both response_time_us and cache now_s.
  4. tracker.track(from) (nonblocking; full table counts tracker_full).
  5. Snapshot: if (self.manager) |m| m.acquire(io) else null; released once on every exit path (single defer). Null snapshot → unfiltered_queries, group features off (ruling 4).
  6. Non-IN qclass → step 12 (upstream direct), no cache (ruling 9).
  7. Normalize qname into scratch.normalize.
  8. Local records: hasNamelookup; answer via ResponseBuilder + writeAnswers + addOptEcho, authoritative. A name that exists with no records of the qtype → NODATA, authoritative. Log (upstream="local", cache_hit=null) and return.
  9. Forward zones: zones.match → construct ForwardClient on the stack (init(zone.resolver, &scratch.frame, self.forward_read_timeout)), cache get → miss → exchange into response_buf → cache put (classify), log (upstream=resolver text), return. Filtering/safe-search/uncloak bypassed (ruling 7). Exchange failure → SERVFAIL.
  10. Pause check (isPaused(now_s)) → skip 11 and the uncloak in 14; count paused_queries.
  11. Filtering: snapshot.evaluate(group, domain); blocked → response.writeBlocked into response_buf, log (blocked=true, reason), return. Then safe-search: snapshot.safeSearch(group) and safesearch.rewrite(domain) → remember the target; outgoing QNAME = target (query rebuilt in scratch.build; count safesearch_rewrites).
  12. ECS (ruling 13): .stripedns.stripEcs into scratch.build (compose with the safe-search rebuild: build once with both the final QNAME and the filtered OPT — one rebuild, not two); .forwardedns.ecsPayload for the key. Overflow → ecs_strip_failed, forward original.
  13. Cache get (skip for safe-search, ruling 5): key = buildKey(outgoing qname, qtype, qclass, do_bit, ecs_payload_or_null) under cache_mutex; hit → copy is already in response_buf-sized out (pass response_buf), packet.setId, UDP truncation check as today, log (cache_hit=true, upstream=""), return.
  14. Upstream: self.upstream.exchange(io, outgoing, response_buf) (existing error mapping to SERVFAIL/drop stays). Then uncloak (ruling 11) unless paused/non-IN: iterate packet.answers, follow CNAMEs from the outgoing qname via record.rdataCname, depth ≤ 8, normalize each target, evaluate; blocked → writeBlocked for the ORIGINAL question, count uncloak_blocked, log (blocked=true, "cname:" reason), return.
  15. Safe-search synthesis (ruling 10): new ResponseBuilder with the original question, addAnswer(original, CNAME, target), then each A/AAAA answer record of the upstream response verbatim (owner = target name), addOptEcho. The upstream bytes are read from response_buf while the builder writes into... the builder needs its own buffer: build into a second half — concretely handle requires response_buf.len >= 1024 already via assert; S2 adds a synth: [4096]u8 field to Scratch for the safe-search build, then copies the finished reply into response_buf. A synthesized reply that cannot fit 4096 → SERVFAIL (safe-search answers are tiny in practice).
  16. Cache put: upstream path only, not blocked/safe-search (ruling 5): classifycache_mutexput. Allocation failure inside put → skip caching (already Allocator.Error; count nothing, put's stats cover it).
  17. UDP truncation check (existing). Log the entry: Entry.init(.{ .timestamp = now_s, .domain = normalized domain, .client_ip = formatted from, .qtype, .blocked=false, .response_time_us, .cache_hit=false, .upstream = per ruling 20 }), logger.log(io, entry). Return reply.

The query-log write happens on every replied path (blocked, local, forward, cache, upstream, paused) — one call site per return is acceptable; a tiny private logEntry helper keeps it readable. Refused/malformed/dropped are never logged (rulings 8, 20).

S2.4 Acceptance

  • zig fmt --check + zig ast-check clean; temp-root tests pass.
  • Existing handler tests still pass (bare handler, null optionals).
  • New in-file tests (fake transport.Client fixtures, as the existing tests do): REFUSED on limit exhaustion; local A answer + NODATA; forward-zone match exchanges via the zone resolver fake and bypasses a blocklisted name; blocked domain → zero-IP answer with blocking TTL; allow-over-block passes; cache miss→put→hit with aged TTL + fresh ID; blocked responses not cached; uncloak: CNAME chain to a blocked target at depth 1 and at depth 8 blocks, depth 9 does not, cname: reason logged; safe-search: rewritten QNAME goes upstream, reply carries original question + CNAME + A records only; ECS strip: option 8 gone upstream, others preserved; ECS forward: payload lands in the cache key (two subnets → two entries); paused: blocklisted name answers normally and counts; non-IN bypasses filtering; null manager → unfiltered count; every path's log entry fields per ruling 20 (capture via a Logger with a small queue, drain and assert).

S2 As built

Handler gained one additive field: negative_ttl_max: u32 = 0 (set from cfg.cache.negative_ttl_maxDnsCache keeps no copy of its config and classify needs the value; S6 must wire it or negative caching is silently off). Scratch gained synth: [4096]u8 (spec-blessed) and uncloak: [types.max_name_len]u8 (target normalization must not clobber normalize, which still holds the logged name). Zones.empty already existed (forward_zones.zig:34) — the S2.2 parenthetical was stale. Both mutexes and the tracker use lockUncancelablestd.Io.Mutex.lock returns Cancelable!void (Io.zig:1602) and handle has no error union. Safe-search + ECS strip compose via two buffers: rebuild into scratch.build, strip into scratch.synth (stripEcs requires non-overlapping output); synth is reused for the reply after the upstream answers. Uncloaking is also skipped after a safe-search rewrite (the answers belong to the provider's target, which the client never asked about). safesearch_rewrites counts at the rebuild, not the table lookup. SERVFAIL paths are counted but never query-logged (the Entry schema has no status field; a logged row would read as success — accepted ruling). Forward zones are exempt from the ECS strip (accepted ruling: the RFC 7871 concern is third-party upstreams, not the LAN resolver). Ruling 20's open choice: upstream = "pool" for pool answers; forward-zone rows log cache_hit false/true (ruling 20's "null" contradicted ruling 5 — null is now local records and blocked answers only). The safe-search CNAME TTL is the minimum TTL of the copied A/AAAA records, zero when none. handle keeps the 512-byte floor assert (the spec's "1024" was wrong); the copy out of synth is length-checked with SERVFAIL fallback. tracker_full mirrors Tracker.snapshotStats().dropped_full after each track (track returns nothing; an always-zero counter would lie). max_cname_depth = 8 public. queries still means pool-answered queries. 23 new handler tests + 25 adapted + 6 pause tests. After review, two ECS cache holes fixed at the root: (1) a failed ECS strip forwards the client's subnet upstream, so the answer may be subnet-specific while the .strip-mode key names no subnet — outgoingQuery returns Outgoing { bytes, cacheable } and a failed strip sets cacheable = false, skipping both the cache get and the put; (2) .forward mode forwards the whole OPT but edns.findOption reports only the first ECS option, so a query carrying several would key on one subnet while the resolver answered for another — a local subnetForKey walks the option list once and returns .uncacheable for a repeated option 8, for a payload over dns_cache.max_ecs_len, and for a list that will not walk; cacheKey then declines to cache. RFC 7871 §6 permits one ECS option, but a query is client-controlled, so a duplicate is treated as hostile. Both fixes are mutation-checked. edns.ecsPayload lost its only caller (the key needs the option count as well as the payload; one walk yields both) and is removed from edns.zig as dead surface. NOTE: the plain zig build test does not typecheck the listener call sites (serve is never analyzed there) — only -Dintegration validates S3's edit.


Session S3: listeners pass the client address and scratch

Owns: src/server/udp_server.zig, src/server/tcp_server.zig, and the three integration test files whose Handler{…} literals the S2 struct change broke: src/server/udp_server_integration_test.zig, src/server/tcp_server_integration_test.zig, src/server/resolver_integration_test.zig (S2 measured nine stale literals missing blocking/forward_read_timeout/records/zones at udp:90/128/163/204, tcp:153/183/216/268, resolver:229, plus the handle-arity change once those are fixed).

  • Slot (udp) gains scratch: handler.Scratch; the call becomes self.handler.handle(io, .udp, address.NetAddress.fromIp(slot.from), slot.query[0..slot.len], &slot.reply, &slot.scratch).
  • Conn (tcp) gains peer: std.Io.net.IpAddress (captured at claim from stream.socket.address) and scratch: handler.Scratch; call updated the same way.
  • No other behavior change. Slot/Conn are pool-allocated already, so the added ~5 KiB per slot changes no lifetime.
  • Acceptance: fmt/ast clean; both files' existing tests updated and passing; integration tests (-Dintegration udp/tcp/resolver) still pass unchanged in behavior.

S3 As built

Conn.peer is captured in claim inside the same mutex-held section that stores the stream — written once per connection, read only by that connection's task; verified from source that accept fills Socket.address with the peer (net.zig:1442 → Threaded.zig:12439 addressFromPosix), and a new TCP integration test asserts the captured peer is the real loopback client (an uninitialized field would typecheck and silently corrupt rate limiting/grouping/logging). Importing the address module forced renaming the bind/listen address parameters to bind_address/listen_address (Zig shadowing rule; call sites are positional, nothing else changed). The nine stale integration literals became a per-file bareHandler fixture deriving blocking from model.Blocking{} and forward_read_timeout from model.readTimeout(.{}) (3000 ms, .awake), so the fixtures cannot drift from config defaults. Memory-cost comments updated for the scratch: UDP ≈74 KiB/slot (≈4.6 MiB at 64), TCP ≈137 KiB/slot (≈8.8 MiB at 64); the UDP <8 MiB budget assertion still holds. After review: TcpServer.serve distinguishes its two exit paths, because the composition root cancels its task group before any listener deinit runs (the deinit defers are declared first, so group.cancel executes first). acceptLoop returns a Stop enum: .closing when deinit stopped it (error.SocketNotListening, a .shutting_down claim, or the state flag) and .canceled when the task is being canceled (error.Canceled from accept or the retry sleep). On .closing the connection group drains under .blocked protection as before — beginShutdown already unblocked every live stream, so the wait is bounded by the shutdown. On .canceled the group is canceled instead: nothing has shut the streams down, deinit cannot run until serve returns, and RFC 7766 §6.2.1.1 lets a client hold a connection open indefinitely by asking again inside the idle budget — draining would let one client stall process shutdown. The cost is the single reply mid-write. The old await under .blocked protection was what made cancel invisible: Syscall.start (Threaded.zig:1349) never arms cancellation while protection is blocked. Per-connection cleanup is unaffected (serveConn's defer finish runs on the canceled path; finish already used lockUncancelable + .blocked around the close). Pinned by "a canceled serve does not wait for a live connection" (blocked read, 600 s idle budget, cancel must return; reinstating drain-always hangs the suite). The same await shape remains in UdpServer.serve deliberately: a UDP task is one bounded upstream exchange, not a client-paced session.


Session S4: client tracker + repo additions

Owns: src/server/clients.zig (new), src/storage/repositories/clients_repo.zig (extend).

S4.1 clients_repo additions (SQL; existing functions frozen)

/// INSERT OR CONFLICT: materialize an unseen client (hand_edited=0, first_seen=last_seen=now)
/// or touch last_seen on an existing row (hand-edited rows get last_seen touched too — the
/// operator sees liveness; group/name/hand_edited never change here).
pub fn upsertSeen(database: *db.Db, ip: []const u8, now_s: i64) db.Error!void
/// DELETE hand_edited=0 AND last_seen < cutoff. Returns rows deleted.
pub fn pruneStale(database: *db.Db, cutoff_s: i64) db.Error!u32

ON CONFLICT(ip) DO UPDATE SET last_seen=excluded.last_seen. Check the actual clients schema (unique key on ip) in storage/config_schema.zig and use the real column names. The existing insertClient (hand_edited=1, import path) stays untouched.

S4.2 Tracker

pub const Tracker = struct {
    pub const max_pending = 512;
    pub const flush_interval_s = 60;
    pub const prune_every_passes = 1440;   // daily at 60 s per pass
    mutex: std.Io.Mutex, // guards the fixed table below
    // fixed-capacity map, same house pattern as rate_limiter: slot array + index
    pub fn init(retention_days: u32) Tracker
    pub fn track(self: *Tracker, io: std.Io, addr: address.NetAddress) void  // full → drop, count stats.dropped_full
    pub fn run(self: *Tracker, io: std.Io, database: *db.Db, monitor: ?*disk_monitor.Monitor)
        std.Io.Cancelable!void
    pub const Stats = struct { tracked: u64, flushed: u64, dropped_full: u64, pruned: u64, flush_failures: u64 };
};

run is the house loop shape (manager.zig:946): startup no-op, then every flush_interval_s on the .boot clock: skip the pass when monitor says writes are gated (ruling 17); else drain the pending table under the mutex into a local copy (release before I/O), one upsertSeen per client with real wall-clock seconds; every prune_every_passes-th pass also pruneStale(now - retention_days * 86_400). A failed upsert warns once per pass and counts flush_failures; entries stay dropped (next query re-tracks). The database handle is the loop's own dedicated config.db connection (ruling 21) — document on run like retention does. IP text via address.NetAddress.format into a stack buffer.

Tests: track dedupes per flush (same client twice → one row, last_seen = latest), table full drops and counts, flush inserts hand_edited=0 rows and touches existing hand-edited rows' last_seen only, prune removes only stale hand_edited=0 rows, gated pass writes nothing (in-memory db + a Monitor fixture — Monitor's classify/state accessors allow constructing a gated state; if not constructible without I/O, test the gate branch through a bool seam writes_allowed parameter on a private flushOnce and have run pass the monitor read).

Acceptance: fmt/ast clean; temp-root (with -lc -lsqlite3) green; migration untouched.

S4 As built

init(retention_days: u16) (matches model.Logging.retention_days). The pending table is a fixed 512-entry array with a linear scan, not a hash map — the frozen init takes no allocator, and rate_limiter's actual pattern is a pre-reserved hash map, so "same house pattern" did not apply; a linear scan over ≤512 entries of 17-byte keys is fine at household scale. Added: trackAt(io, addr, now_s) (timestamp seam; track = trackAt with real wall-clock), flushOnce(io, database, writes_allowed: bool) public (the gate seam run feeds from monitor.writesAllowed(); S7 case 10 forces passes with it), snapshotStats(io) and pendingClients(io)Stats stays plain u64 and every field plus passes is written under mutex, so readers go through the accessors. last_seen is stamped at track time (it means "last query", and it is what makes latest-wins dedupe observable); the flush writes each entry's own stamp. A gated pass does not increment passes (a skipped pass must not advance the prune schedule). track uses lockUncancelablehandle has no error union, so track is not a cancelation point. upsertSeen resolves group_id via (SELECT id FROM groups WHERE name = 'default') — the column is NOT NULL REFERENCES groups (config_schema.zig:26), which ruling 14's sketch omitted; a database with no default group fails the whole statement with error.Constraint and writes nothing (validate.zig:461 already rejects such configs, so this is a warn-and-continue path in the flush loop, never live). The upsert's DO UPDATE touches last_seen only — name, group_id, hand_edited, first_seen survive on every existing row, hand-edited or not, one test per column. pruneStale uses strict < (matches queries_repo.pruneOlderThan) and clamps sqlite3_changes to u32. 16 tests (9 tracker, 7 repo).


Session S6: composition root — app.zig, shutdown.zig, cli wiring

Owns: src/app.zig (new), src/server/shutdown.zig (new), src/cli.zig (runRun body, DataDir), src/main.zig (only if a signature forces it — expected unchanged).

S6.1 shutdown.zig

var event: std.Io.Event = .unset;
var event_io: ?std.Io = null;  // set once by install before handlers exist
pub fn install(io: std.Io) void  // sigaction INT+TERM → handler calls event.set(event_io.?)
pub fn wait(io: std.Io) std.Io.Cancelable!void  // event.wait
pub fn trigger(io: std.Io) void  // tests and future admin use: event.set

Handler fn: callconv(.c), calls set only (async-signal-safe per ruling 2). flags = 0, empty mask (matches Threaded's own handlers, Threaded.zig:1653). Do not restore handlers on exit — the process is leaving anyway; install is once-per-process (assert).

S6.2 app.zig — pub fn run(runner: cli.Runner, paths: cli.Paths) u8

Startup order (each failure prints to runner.err and returns cli.exit_runtime):

  1. DataDir.open (S6 changes its openDir flags to .{ .iterate = true } — the monitor needs it; nothing else cares). openConfigDb. migrations.migrate. bootstrap.bootstrap(...) (first run seeds; diags printed like runCheck does). export.readConfig into an arena.
  2. logging.install(io, cfg.logging, ...) — the sink goes live (main.zig's std_options already routes).
  3. Build state, in dependency order, all heap-allocated where pinning is required (Logger, Manager, Pool entries and their DoH/DoT client structs must not move — same pattern as cli.probeUpstreams but with owned arrays that outlive the pool): records + zones (local_repo listers → build), fetcher + manager (Manager.init), pool (upstreams_repo → clients array → Pool.init), cache (DnsCache.init), limiter (RateLimiter.init from cfg.dns), pause, tracker, logger (Logger.init with a heap queue_buf of cfg.logging.query_log_buffer_max entries), monitor (Monitor.init with the data dir), retention (Retention.init), querylog connections: add DataDir.openQuerylogDb(self, io) !db.Db mirroring openConfigDb (same pragmas + 0600); open TWO (writer, retention) plus ONE extra config.db connection for the tracker (ruling 21).
  4. One synchronous manager.reload(io); failure warns and continues (ruling 4).
  5. Handler construction (S2 struct) + UdpServer.bind + TcpServer.listen on cfg.dns.bind_ipv4/bind_ipv6/port — v4 and v6 both, matching PLAN §7.2 parity; a bind failure is fatal.
  6. shutdown.install(io). Then one std.Io.Group: udp.serve, tcp.serve, logger.runWriter(io, &querylog_writer_db, &monitor), retention.run(io, &querylog_retention_db), monitor.run(io), manager.runScheduler(io), tracker.run(io, &tracker_db, &monitor), and a maintenance loop (private fn in app.zig, house shape, .boot, 60 s): cache sweep + limiter sweep under their mutexes... the mutexes live in the Handler; the maintenance loop takes *Handler and locks the same mutexes.
  7. shutdown.wait(io) (a Canceled result also proceeds to teardown). Then ruling 22 order: logger.shutdown(io)group.cancel(io) → server deinits → pool/manager/state deinits → DB closes → DataDir.close → return cli.exit_ok.

Log a one-line startup summary (scope .nxdns, info): bind addresses, snapshot generation or "unfiltered", upstream count.

S6.3 cli.zig

runRun body: return app.run(r, paths); — keep the doc comment honest. DataDir.open flag change + openQuerylogDb as above (both used by app; runCheck untouched).

S6.4 Acceptance

  • fmt/ast clean on all owned files; zig build (exe) succeeds — app.zig is reached from main, so full semantic analysis covers it.
  • zig build test green (existing cli tests untouched; app.zig gets a smoke test only if expressible without sockets — otherwise S7 owns runtime coverage; state that choice).
  • Manual smoke (orchestrator runs it): nxdns run against a scratch data dir starts, answers dig @127.0.0.1 -p <port>, and exits 0 on SIGTERM promptly.

S6 As built

Startup order as specified, with these corrections. Binding: IPv6 binds first; an AddressInUse from the IPv4 bind that follows a wildcard IPv6 bind is dual-stack coverage (info line), not a failure — on Linux net.ipv6.bindv6only=0 makes the v6 wildcard own the v4 port, and BindOptions.ip6_only cannot prevent it because Threaded.zig:12274 sets IPV6_V6ONLY to 0 when the option is true (inverted vs its doc); a foreign IPv4 holder would have failed the v6 wildcard bind too, so the inference is sound. NetAddress.fromIp unmaps ::ffff: clients, so v4 clients keep their real identity. A kernel without IPv6 (AddressFamilyUnsupported and friends) warns and serves IPv4; every other bind failure is fatal. Exit codes: configuration faults (NoUsableUpstreams, BadBindAddress, BadRateLimit — zero rate fields are rejected before RateLimiter.init asserts) exit cli.exit_check with a "run 'nxdns check'" hint; everything else exits cli.exit_runtime. Querylog open: openQuerylogDb wraps querylog_schema.open (creation/repair lives there; chmod of the file plus both WAL sidecars happens after, missing sidecars ignored); reopenQuerylogDb opens the extra connections without repeating the O(size) quick_check. Two std.http.Clients: blocklist downloads must not queue DoH queries behind a multi-megabyte stream. Teardown is defer-driven; the invariants held are: group.cancel joins first, config.db closes after manager.deinit, DataDir closes last. negative_ttl_max and forward_read_timeout are wired (the S2 trap). shutdown.zig has reset() for tests and 2 in-file tests (orchestrator wired it into tests.zig — app.zig analysis via cli does not collect tests). app.zig itself has no unit test; S7 case 11 is the runtime proof. Smoke test passed: boot, migration 0->2, dual-stack answers on 15353/15354, SERVFAIL from the dead upstream, SIGTERM exit 0 in 0.007 s, DBs 0600, zero err lines. The manager disk gate was wired in a follow-up S6 edit — see ruling 17's as-built text; the gate test drives Monitor.state_raw directly (no I/O) and pins no-monitor/ok/warn pass, critical gates-and-counts, recovery resumes. After review: the composition root takes one synchronous monitor.sample(io) after Monitor.init, before the Io.Group spawns anything — Monitor initializes to .ok and run samples inside its own task, so without it the writer/tracker/scheduler would consult a gate still saying "writes allowed" during the boot window on a critically full disk. A failed first sample warns and leaves .ok (an unreadable filesystem is not evidence of a full disk — the monitor's own documented policy). run's own first sample repeats it (one extra statvfs + scan at boot; run's sample-then-sleep contract stays untouched). The ordering is not unit-testable (serve binds sockets and opens four databases); S7's boot case plus the statement order is the evidence.


Session S7: phase 7 integration tests

Owns: src/server/phase7_integration_test.zig (new; gated on build_options.integration like phase6). Fixtures may reuse the S8/M6 helper patterns (in-memory dbs, fake upstream sockets).

Cases (each end-to-end through a real Handler with real UDP sockets on 127.0.0.1 ephemeral ports; a fake upstream = a UDP socket task answering canned responses):

  1. Blocked domain answers 0.0.0.0 with blocking TTL; log entry captured with reason.
  2. Allow-rule-over-blocklist resolves via upstream.
  3. Local record answers authoritatively without touching the fake upstream.
  4. Forward zone: matching suffix goes to the zone resolver socket (not the pool fake), a blocklisted name inside the zone still resolves (bypass), response cached (second query does not hit the zone socket).
  5. Cache: miss→hit, TTL aged, transaction ID fresh, cache_hit logged on the hit.
  6. Uncloak: fake upstream returns CNAME chain to a blocked target → zero-IP answer with cname: reason.
  7. Safe-search group: query for a safesearch.zig table domain returns CNAME + the fake upstream's A for the target; question is the original name.
  8. Rate limit: limit=2 window=60 → third query REFUSED.
  9. Pause: pauseFor(indefinite) lets a blocked domain resolve; unpause blocks it again.
  10. Tracker: after queries + one forced flush (flushOnce seam or short interval), the clients table holds the source IP with hand_edited=0.
  11. Full app boot: app.run in a task against a temp data dir (bootstrap from a minimal config.zon on an ephemeral port), one real resolve through it, shutdown.trigger, task returns 0. This is the lifecycle proof.

Acceptance: zig build test -Dintegration green, zero std.log.err emitted (binding logging policy), cases deterministic (no external network — the pool's upstream is the fake).

S7 As built

Case 11 runs against a std.testing.tmpDir rather than a fixed scratch path (a committed test cannot hold a machine-specific absolute path; the phase6 fixture already resolves cwd-relative paths through both std.Io.Dir and SQLite). Port 15455, IPv4+IPv6 binds, the dead https://127.0.0.1:9/dns-query upstream, one seeded local record; app.serve installs and deinstalls its own log sink, so installForTest is unnecessary and the case asserts the error writer stayed empty; shutdown.reset() brackets the case. Case 5 proves TTL ageing with an entry planted ten seconds in the past under the handler's own key — the miss-hit-fresh-ID half is a real round trip; no test can advance the clock. Case 4's zone resolver is a UDP socket task answering with ResponseBuilder against the received query; the second query proves the cache by leaving the resolver's datagram count at 1. The fake pool upstream builds each answer from the question it receives (safe search and ECS rewrite the outgoing question, so canned bytes would not bind). A shared Loop fixture binds the listener and client before the serve task starts, so no task holds a pointer into a value that later moves. 11 tests, all gated on build_options.integration; the integration suite emits zero err-level lines.


Module layout (new/changed)

file change
src/dns/edns.zig +stripEcs (S1; ecsPayload added then removed in review)
src/server/handler.zig pipeline rewrite (S2)
src/server/pause.zig new (S2)
src/server/udp_server.zig, tcp_server.zig address + scratch plumbing (S3)
src/server/clients.zig new tracker (S4)
src/storage/repositories/clients_repo.zig +upsertSeen, +pruneStale (S4)
src/app.zig new composition root (S6)
src/server/shutdown.zig new (S6)
src/cli.zig runRun body, DataDir.iterate, openQuerylogDb (S6)
src/server/phase7_integration_test.zig new (S7)

File ownership

S1: dns/edns.zig. S2: server/handler.zig, server/pause.zig. S3: server/udp_server.zig, server/tcp_server.zig. S4: server/clients.zig, storage/repositories/clients_repo.zig. S6: app.zig, server/shutdown.zig, cli.zig, main.zig(if forced). S7: its test file. Orchestrator: src/tests.zig, specs. No parallel sessions share a file (S1∥S4; the rest are sequential).

Acceptance (milestone complete)

  • zig build test and zig build test -Dintegration green; zig fmt --check clean; zig build cross produces both static exes.
  • All S7 cases pass; the app boots, serves, and shuts down cleanly (case 11).
  • No std.log.err in any new code path exercised by tests (binding logging policy: err is reserved for swallowed failures).
  • Spec As-built sections synced per session.

Anti-requirements

  • No HTTP/SSE/metrics/auth (Phase 8). No DoH/DoT server (Phase 9). No packaging (Phase 10).
  • No regex rules, DHCP, DNSSEC validation, DoQ, clustering (PLAN §2.2).
  • No per-group cache, no cache persistence, no query re-resolution in uncloaking.
  • No new config fields, no schema migration, no pause persistence.
  • No benchmark harness this milestone.
  • dns/ stays pure (S1's additions are pure); filter/local/cache modules are not edited.