From ec96142049c37abc496ba74768477ebb6b847491 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Wed, 22 Jul 2026 01:08:35 -0500 Subject: [PATCH] match upstream logging controls --- README.md | 9 ++ docs/configuration-parity.md | 3 +- justfile | 6 ++ src/internal/logging.zig | 157 +++++++++++++++++++++++++++++++++++ src/main.zig | 52 +++++++++++- tests/logging_contract.py | 149 +++++++++++++++++++++++++++++++++ 6 files changed, 374 insertions(+), 2 deletions(-) create mode 100644 src/internal/logging.zig create mode 100644 tests/logging_contract.py diff --git a/README.md b/README.md index b0489d0..a60958d 100644 --- a/README.md +++ b/README.md @@ -121,6 +121,15 @@ As upstream does, a non-positive HTTP shutdown timeout selects 30 seconds, whereas a non-positive client-drain timeout is already expired: teardown is immediate and delivery of the best-effort close frame is not guaranteed. +Process logs follow Jetstream V2's canonical controls: `--log-level` accepts +debug/info/warn/error and their upstream aliases, while `--log-format` selects +`json` (the default) or `text`. Filtering remains runtime-configurable in +ReleaseSafe builds and applies to Stream plus its Zig dependencies. Root flags +remain persistent, so both `stream --log-level=debug serve` and +`stream serve --log-level=debug` work. `just logging-contract` proves the real +process formats, filtering, placement, and invalid-configuration failures +without external network access. + The subscriber's writer-owned hot log retains 256 MiB by default and never evicts rows above the durable archive watermark. Operators can tune the real byte budget with `--subscribe-read-log-retention-bytes`. The upstream slow diff --git a/docs/configuration-parity.md b/docs/configuration-parity.md index f3a87c8..ab938e9 100644 --- a/docs/configuration-parity.md +++ b/docs/configuration-parity.md @@ -27,7 +27,8 @@ the same runtime mechanism and an offline receipt observes that behavior. | timestamp import token and directory | closed | real HTTP/import, restart, and strict power-loss receipts | | public/debug bind addresses | closed | `serve --addr` binds the protocol-only public listener, an empty `--debug-addr` binds nothing, and a configured debug listener exclusively serves health/readiness/metrics; `just listener-contract` proves route isolation, public WebSocket upgrade, rejected debug upgrade, disabled-debug refusal, listener readiness, and the historical combined-port compatibility mode through real offline processes | | shutdown/client drain timeouts | closed | canonical signed Go-duration flags, 5s/10s defaults, net/http's non-positive 30s sentinel, the client drain's already-expired non-positive sentinel (immediate teardown with a best-effort close frame), listener-admission gating, concurrent public/debug/pipeline shutdown, close code 1001 for positive budgets, cooperative early exit, and one shared silent-client deadline are covered by unit and real ReleaseSafe process receipts. Per-client close hooks now fan out concurrently with borrowed-handler lifetime protection; a real-socket transport regression saturates the first peer's send buffer, proves a later peer receives 1001 before the deadline, and proves the blocked writer is interrupted at the shared deadline. | -| log level/format and OpenTelemetry | open | canonical metrics exist; upstream process configuration and tracing surface do not | +| log level/format | closed | `--log-level` retains debug calls in ReleaseSafe and applies the upstream case-insensitive aliases plus info default to every `std.log` producer; `--log-format` defaults to newline-delimited JSON and selects canonical text fields. `just logging-contract` proves default JSON, warn filtering, text/debug output, persistent root-flag placement, and fail-closed invalid values through real offline processes. | +| OpenTelemetry | open | canonical metrics exist, and the sibling `otel-zig` provides a real Zig 0.16 OTLP HTTP exporter/SDK, but Stream has not yet wired upstream's spans, standard `OTEL_*` configuration, service resource, or bounded shutdown flush | | `JETSTREAM_*` environment sources and unknown-variable rejection | open | Stream currently accepts CLI arguments only | | serve/version/inspect command surface | closed | `serve` is accepted explicitly while retaining Stream's historical flag-only invocation. `version` emits the upstream build-information shape. `inspect-segment` matches the pinned renderer byte-for-byte on an upstream-produced sealed fixture while also walking active files, surfacing partial tails, and labelling readable checksum corruption. `inspect-all` folds the real steady and bootstrap trees with upstream's missing-root, racing-tail, active-skip, aggregation, sorting, truncation, and text-rendering semantics; a copied upstream golden report is byte-exact. Listener and root logging controls remain tracked by their separate rows. | diff --git a/justfile b/justfile index 202f4fb..af34962 100644 --- a/justfile +++ b/justfile @@ -21,6 +21,12 @@ shutdown-contract: zig build -Doptimize=ReleaseSafe python3 tests/shutdown_contract.py +# Offline: canonical JSON/text handlers, runtime level filtering, persistent +# root-flag placement, and fail-closed invalid logging configuration. +logging-contract: + zig build -Doptimize=ReleaseSafe + python3 tests/logging_contract.py + # Offline: validate the vendored upstream Grafana dashboard and its parser. dashboard-test: python3 tools/adapt_dashboard.py --check diff --git a/src/internal/logging.zig b/src/internal/logging.zig new file mode 100644 index 0000000..a313a78 --- /dev/null +++ b/src/internal/logging.zig @@ -0,0 +1,157 @@ +const std = @import("std"); + +const Io = std.Io; + +pub const Format = enum(u8) { + text, + json, +}; + +var configured_level: std.atomic.Value(u8) = .init(@intFromEnum(std.log.Level.info)); +var configured_format: std.atomic.Value(u8) = .init(@intFromEnum(Format.json)); + +pub fn configure(level_raw: []const u8, format_raw: []const u8) !void { + const level = try parseLevel(level_raw); + const format = try parseFormat(format_raw); + configured_level.store(@intFromEnum(level), .release); + configured_format.store(@intFromEnum(format), .release); +} + +pub fn parseLevel(raw: []const u8) !std.log.Level { + const value = std.mem.trim(u8, raw, " \t\r\n"); + if (value.len == 0 or std.ascii.eqlIgnoreCase(value, "i") or std.ascii.eqlIgnoreCase(value, "info")) return .info; + if (std.ascii.eqlIgnoreCase(value, "d") or std.ascii.eqlIgnoreCase(value, "dbg") or std.ascii.eqlIgnoreCase(value, "debug")) return .debug; + if (std.ascii.eqlIgnoreCase(value, "w") or std.ascii.eqlIgnoreCase(value, "warn") or std.ascii.eqlIgnoreCase(value, "warning")) return .warn; + if (std.ascii.eqlIgnoreCase(value, "e") or std.ascii.eqlIgnoreCase(value, "err") or std.ascii.eqlIgnoreCase(value, "error")) return .err; + return error.InvalidLogLevel; +} + +pub fn parseFormat(raw: []const u8) !Format { + const value = std.mem.trim(u8, raw, " \t\r\n"); + if (value.len == 0 or std.ascii.eqlIgnoreCase(value, "json")) return .json; + if (std.ascii.eqlIgnoreCase(value, "text")) return .text; + return error.InvalidLogFormat; +} + +pub fn logFn( + comptime level: std.log.Level, + comptime scope: @EnumLiteral(), + comptime format: []const u8, + args: anytype, +) void { + if (@intFromEnum(level) > configured_level.load(.acquire)) return; + + var message: Io.Writer.Allocating = .init(std.heap.page_allocator); + defer message.deinit(); + message.writer.print(format, args) catch { + std.log.defaultLog(level, scope, format, args); + return; + }; + + const io = std.Options.debug_io; + const previous = io.swapCancelProtection(.blocked); + defer _ = io.swapCancelProtection(previous); + var buffer: [1024]u8 = undefined; + const terminal = std.debug.lockStderr(&buffer).terminal(); + defer std.debug.unlockStderr(); + + const selected: Format = @enumFromInt(configured_format.load(.acquire)); + writeRecord( + terminal.writer, + selected, + level, + @tagName(scope), + message.written(), + Io.Timestamp.now(io, .real).toMicroseconds(), + ) catch {}; +} + +fn writeRecord( + writer: *Io.Writer, + format: Format, + level: std.log.Level, + scope: []const u8, + message: []const u8, + now_us: i64, +) !void { + var time_buffer: [32]u8 = undefined; + const timestamp = formatRfc3339(&time_buffer, now_us); + const component = if (std.mem.eql(u8, scope, "default")) "main" else scope; + switch (format) { + .json => { + var json: std.json.Stringify = .{ .writer = writer }; + try json.beginObject(); + try json.objectField("time"); + try json.write(timestamp); + try json.objectField("level"); + try json.write(levelName(level)); + try json.objectField("msg"); + try json.write(message); + try json.objectField("component"); + try json.write(component); + try json.endObject(); + try writer.writeByte('\n'); + }, + .text => { + try writer.print("time={s} level={s} msg=", .{ timestamp, levelName(level) }); + var json: std.json.Stringify = .{ .writer = writer }; + try json.write(message); + try writer.print(" component={s}\n", .{component}); + }, + } + try writer.flush(); +} + +fn levelName(level: std.log.Level) []const u8 { + return switch (level) { + .debug => "DEBUG", + .info => "INFO", + .warn => "WARN", + .err => "ERROR", + }; +} + +fn formatRfc3339(buffer: []u8, micros: i64) []const u8 { + const seconds: u64 = @intCast(@divFloor(micros, std.time.us_per_s)); + const fraction: u64 = @intCast(@mod(micros, std.time.us_per_s)); + const epoch: std.time.epoch.EpochSeconds = .{ .secs = seconds }; + const year_day = epoch.getEpochDay().calculateYearDay(); + const month_day = year_day.calculateMonthDay(); + const day = epoch.getDaySeconds(); + return std.fmt.bufPrint(buffer, "{d:0>4}-{d:0>2}-{d:0>2}T{d:0>2}:{d:0>2}:{d:0>2}.{d:0>6}Z", .{ + year_day.year, + month_day.month.numeric(), + month_day.day_index + 1, + day.getHoursIntoDay(), + day.getMinutesIntoHour(), + day.getSecondsIntoMinute(), + fraction, + }) catch unreachable; +} + +test "canonical log level and format parsing" { + try std.testing.expectEqual(std.log.Level.info, try parseLevel("")); + try std.testing.expectEqual(std.log.Level.debug, try parseLevel("DBG")); + try std.testing.expectEqual(std.log.Level.warn, try parseLevel(" warning ")); + try std.testing.expectEqual(std.log.Level.err, try parseLevel("ERROR")); + try std.testing.expectError(error.InvalidLogLevel, parseLevel("loud")); + try std.testing.expectEqual(Format.json, try parseFormat("")); + try std.testing.expectEqual(Format.json, try parseFormat("JSON")); + try std.testing.expectEqual(Format.text, try parseFormat(" text ")); + try std.testing.expectError(error.InvalidLogFormat, parseFormat("xml")); +} + +test "JSON and text log records are structured and escaped" { + var output: Io.Writer.Allocating = .init(std.testing.allocator); + defer output.deinit(); + try writeRecord(&output.writer, .json, .info, "stream", "hello \"river\"", 1_720_000_000_123_456); + const parsed = try std.json.parseFromSlice(std.json.Value, std.testing.allocator, output.written(), .{}); + defer parsed.deinit(); + try std.testing.expectEqualStrings("INFO", parsed.value.object.get("level").?.string); + try std.testing.expectEqualStrings("hello \"river\"", parsed.value.object.get("msg").?.string); + try std.testing.expectEqualStrings("stream", parsed.value.object.get("component").?.string); + + output.clearRetainingCapacity(); + try writeRecord(&output.writer, .text, .warn, "stream", "slow peer", 1_720_000_000_123_456); + try std.testing.expect(std.mem.indexOf(u8, output.written(), "level=WARN msg=\"slow peer\" component=stream\n") != null); +} diff --git a/src/main.zig b/src/main.zig index 8270214..eb11aa2 100644 --- a/src/main.zig +++ b/src/main.zig @@ -29,10 +29,18 @@ const timestamp_jobs = @import("internal/timestamp/jobs.zig"); const timestamp_manager = @import("internal/timestamp/manager.zig"); const operational = @import("internal/operational.zig"); const inspect_all = @import("internal/inspect_all.zig"); +const logging = @import("internal/logging.zig"); const Io = std.Io; const log = std.log.scoped(.stream); +pub const std_options: std.Options = .{ + // Runtime filtering happens in logging.logFn so ReleaseSafe retains the + // canonical --log-level=debug surface instead of compiling it away. + .log_level = .debug, + .logFn = logging.logFn, +}; + const default_subscribe_read_log_retention_bytes: usize = 256 * 1024 * 1024; /// 8 MB: ReleaseSafe inlining makes TLS/CBOR/crypto call chains deep enough @@ -476,7 +484,23 @@ pub fn main(init: std.process.Init.Minimal) !void { var arg_it = init.args.iterate(); _ = arg_it.next(); // program name + var log_level: []const u8 = "info"; + var log_format: []const u8 = "json"; var first_arg = arg_it.next(); + // urfave/cli root flags are persistent: accept them before every command, + // with either of the spellings supported by the upstream CLI. + while (first_arg) |arg| { + if (std.mem.startsWith(u8, arg, "--log-level=")) { + log_level = arg["--log-level=".len..]; + } else if (std.mem.eql(u8, arg, "--log-level")) { + log_level = arg_it.next() orelse return error.BadArgs; + } else if (std.mem.startsWith(u8, arg, "--log-format=")) { + log_format = arg["--log-format=".len..]; + } else if (std.mem.eql(u8, arg, "--log-format")) { + log_format = arg_it.next() orelse return error.BadArgs; + } else break; + first_arg = arg_it.next(); + } var explicit_serve = false; if (first_arg) |command| { if (std.mem.eql(u8, command, "version") or @@ -485,7 +509,18 @@ pub fn main(init: std.process.Init.Minimal) !void { { var command_args: std.ArrayList([]const u8) = .empty; defer command_args.deinit(allocator); - while (arg_it.next()) |arg| try command_args.append(allocator, arg); + while (arg_it.next()) |arg| { + if (std.mem.startsWith(u8, arg, "--log-level=")) { + log_level = arg["--log-level=".len..]; + } else if (std.mem.eql(u8, arg, "--log-level")) { + log_level = arg_it.next() orelse return error.BadArgs; + } else if (std.mem.startsWith(u8, arg, "--log-format=")) { + log_format = arg["--log-format=".len..]; + } else if (std.mem.eql(u8, arg, "--log-format")) { + log_format = arg_it.next() orelse return error.BadArgs; + } else try command_args.append(allocator, arg); + } + try logging.configure(log_level, log_format); if (std.mem.eql(u8, command, "version")) { try writeVersion(io, command_args.items); } else if (std.mem.eql(u8, command, "inspect-segment")) { @@ -581,6 +616,14 @@ pub fn main(init: std.process.Init.Minimal) !void { shutdown_timeout = try parseGoDuration(arg["--shutdown-timeout=".len..]); } else if (std.mem.startsWith(u8, arg, "--client-drain-timeout=")) { client_drain_timeout = try parseGoDuration(arg["--client-drain-timeout=".len..]); + } else if (std.mem.startsWith(u8, arg, "--log-level=")) { + log_level = arg["--log-level=".len..]; + } else if (std.mem.eql(u8, arg, "--log-level")) { + log_level = arg_it.next() orelse return error.BadArgs; + } else if (std.mem.startsWith(u8, arg, "--log-format=")) { + log_format = arg["--log-format=".len..]; + } else if (std.mem.eql(u8, arg, "--log-format")) { + log_format = arg_it.next() orelse return error.BadArgs; } else if (std.mem.startsWith(u8, arg, "--plc=")) { plc_url = arg["--plc=".len..]; } else if (std.mem.startsWith(u8, arg, "--plc-url=")) { @@ -710,11 +753,18 @@ pub fn main(init: std.process.Init.Minimal) !void { repo_action_rate_limits = false; } else if (std.mem.eql(u8, arg, "--stdout")) { to_stdout = true; + } else if (std.mem.eql(u8, arg, "serve") and !explicit_serve) { + // Root flags are persistent upstream, so `--log-level=debug + // serve` is equivalent to `serve --log-level=debug`. + explicit_serve = true; + if (public_addr_raw == null) public_addr_raw = ":8080"; } else { log.err("unknown arg: {s}", .{arg}); return error.BadArgs; } } + try logging.configure(log_level, log_format); + log.debug("logging configured: level={s} format={s}", .{ log_level, log_format }); if (legacy_port_set and public_addr_raw != null) { log.err("--port cannot be combined with upstream --addr/--debug-addr or explicit serve mode", .{}); return error.BadArgs; diff --git a/tests/logging_contract.py b/tests/logging_contract.py new file mode 100644 index 0000000..bd6fd83 --- /dev/null +++ b/tests/logging_contract.py @@ -0,0 +1,149 @@ +#!/usr/bin/env python3 +"""Offline process receipt for Jetstream V2 log level and format semantics.""" + +from __future__ import annotations + +import json +import signal +import socket +import subprocess +import tempfile +import time +import urllib.request + + +BIN = "./zig-out/bin/stream" + + +def unused_port() -> int: + sock = socket.socket() + sock.bind(("127.0.0.1", 0)) + port = sock.getsockname()[1] + sock.close() + return port + + +def start(root: str, port: int, prefix: list[str]) -> subprocess.Popen[str]: + return subprocess.Popen( + [ + BIN, + *prefix, + f"--addr=127.0.0.1:{port}", + f"--data-dir={root}", + "--upstream=ws://127.0.0.1:17990", + "--relay-http=http://127.0.0.1:17990", + "--plc-url=http://127.0.0.1:17990", + "--compaction-interval=0", + "--retry-interval=0", + "--no-verify", + ], + stdout=subprocess.DEVNULL, + stderr=subprocess.PIPE, + text=True, + ) + + +def wait_ready(process: subprocess.Popen[str], port: int) -> None: + deadline = time.monotonic() + 10 + while time.monotonic() < deadline: + if process.poll() is not None: + stderr = process.stderr.read() if process.stderr else "" + raise AssertionError(f"stream exited with {process.returncode}:\n{stderr}") + try: + urllib.request.urlopen(f"http://127.0.0.1:{port}/", timeout=0.2).read() + return + except OSError: + time.sleep(0.025) + raise AssertionError(f"listener {port} did not start") + + +def stop(process: subprocess.Popen[str]) -> list[str]: + process.send_signal(signal.SIGTERM) + process.wait(timeout=5) + assert process.returncode == 0, process.stderr.read() if process.stderr else "" + return [line for line in process.stderr.read().splitlines() if line] + + +def default_json_case(root: str) -> None: + port = unused_port() + process = start(root, port, ["serve"]) + wait_ready(process, port) + records = [json.loads(line) for line in stop(process)] + assert records + assert all(set(("time", "level", "msg", "component")) <= record.keys() for record in records) + assert any(record["level"] == "INFO" for record in records) + assert not any(record["level"] == "DEBUG" for record in records) + + +def warn_json_case(root: str) -> None: + port = unused_port() + process = start(root, port, ["serve", "--log-level=warning"]) + wait_ready(process, port) + time.sleep(0.1) + records = [json.loads(line) for line in stop(process)] + assert records + assert any(record["level"] == "ERROR" for record in records) + assert all(record["level"] in ("WARN", "ERROR") for record in records) + + +def root_flags_text_debug_case(root: str) -> None: + port = unused_port() + process = start( + root, + port, + ["--log-level=DBG", "--log-format=text", "serve"], + ) + wait_ready(process, port) + lines = stop(process) + assert lines + assert all(line.startswith("time=") for line in lines) + assert any(" level=DEBUG " in line for line in lines) + assert any('msg="logging configured: level=DBG format=text"' in line for line in lines) + + +def invalid_config_case() -> None: + for flag, expected in ( + ("--log-level=loud", "InvalidLogLevel"), + ("--log-format=xml", "InvalidLogFormat"), + ): + result = subprocess.run( + [BIN, "serve", flag], + stdout=subprocess.DEVNULL, + stderr=subprocess.PIPE, + text=True, + timeout=3, + check=False, + ) + assert result.returncode != 0 + assert expected in result.stderr, result.stderr + + +def persistent_command_flags_case() -> None: + for args in ( + ["--log-level", "debug", "--log-format", "text", "version"], + ["version", "--log-level=debug", "--log-format=text"], + ): + result = subprocess.run( + [BIN, *args], + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + text=True, + timeout=3, + check=False, + ) + assert result.returncode == 0, result.stderr + assert result.stdout.startswith("jetstream version "), result.stdout + + +def main() -> None: + with tempfile.TemporaryDirectory() as root: + default_json_case(f"{root}/default") + warn_json_case(f"{root}/warn") + root_flags_text_debug_case(f"{root}/debug") + persistent_command_flags_case() + invalid_config_case() + print("logging contract: json-default warn-filter text-debug persistent-flags invalid-config passed") + + +if __name__ == "__main__": + main() -- 2.51.2