From 8a02dec1eb89aaf6516b378f92b85cdf96224765 Mon Sep 17 00:00:00 2001 From: "ren.ji" Date: Thu, 1 Oct 2026 18:53:02 +0800 Subject: [PATCH] Add layered embedded runtime diagnostics --- CHANGELOG.md | 5 + docs/diagnostics.md | 55 +++++++ platform/linux/Integration/main.cpp | 26 ++++ platform/linux/runtime/backend.cpp | 93 +++++++++++- platform/linux/runtime/backend.hpp | 2 + platform/macos/Integration/Sources/main.swift | 52 ++++++- .../EmbeddedRacketBackend.swift | 100 +++++++++++- .../macos/Sources/RivetRuntime/Client.swift | 112 +++++++++++++- .../Sources/RivetRuntime/Diagnostics.swift | 58 +++++++ .../RivetRuntimeTests/DiagnosticsTests.swift | 26 ++++ platform/windows/Integration/main.cpp | 28 ++++ platform/windows/runtime/backend.cpp | 93 +++++++++++- platform/windows/runtime/backend.hpp | 2 + rivet/backend.rkt | 143 ++++++++++++++---- runtime/include/rivet/diagnostics.hpp | 83 ++++++++++ runtime/tests/protocol_test.cpp | 10 ++ tests/backend-diagnostics.rkt | 36 ++++- 17 files changed, 883 insertions(+), 41 deletions(-) create mode 100644 platform/macos/Sources/RivetRuntime/Diagnostics.swift create mode 100644 platform/macos/Tests/RivetRuntimeTests/DiagnosticsTests.swift create mode 100644 runtime/include/rivet/diagnostics.hpp diff --git a/CHANGELOG.md b/CHANGELOG.md index 57f395d..6a23e1b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,11 @@ ## Unreleased +- Add provider-neutral, layer-attributed JSONL diagnostics across the Racket + backend, C++/Swift native clients, transports, and embedding bridges. Records + identify lifecycle/RPC boundaries, status, last RVT1 event, and request id; + stderr defaults and injectable sinks keep logging dependencies out of the + runtime and leave the RVT1 wire contract unchanged. - Launch packaged GUI artifacts as the final verification gate when a graphical session is available. The executable starts from an unrelated temporary directory, must remain alive for five seconds, and is then terminated; diff --git a/docs/diagnostics.md b/docs/diagnostics.md index dd27fc5..2763542 100644 --- a/docs/diagnostics.md +++ b/docs/diagnostics.md @@ -80,3 +80,58 @@ raco rivet build ``` For automated diagnostics, substitute `raco rivet doctor --json` and archive the resulting JSON with the failing build logs. + +## Embedded runtime diagnostics + +The embedded Windows, macOS, and Linux runtimes and the Racket backend emit one +JSON object per line for lifecycle and request-boundary events. This is separate +from RVT1: it does not add protocol frames, alter application payloads, or couple +the runtime to a logging provider. + +```json +{"schema":"rivet.diagnostic.v1","layer":"native-client","event":"rpc-dispatch","status":"success","last_protocol_event":"response","request_id":42} +``` + +Every record contains: + +- `schema`: always `rivet.diagnostic.v1`; +- `layer`: `native-runtime`, `abi-bridge`, `transport`, `protocol`, + `native-client`, or `racket-backend`; +- `event` and `status`: the lifecycle boundary and `begin`, `success`, or + `failure`; +- `last_protocol_event`: the most recently observed RVT1 message kind, or + `none` before the first frame; +- optional `request_id` and `message` fields. + +The event stream covers backend initialization, transport creation, Hello +handshake, per-RPC dispatch, cancellation, orderly shutdown, unexpected channel +closure, reader-loop failure, and backend exit. A failure record therefore says +whether the last known boundary was the Racket backend, the ABI bridge, RVT1 +validation, transport I/O, or the native client. + +By default, embedded applications write JSONL to standard error. Applications +can redirect records to their own logger or crash reporter without adding a +Rivet logging dependency: + +```cpp +rivet::windows::RacketRuntimeConfig config; +config.diagnostic_sink = [](rivet::DiagnosticRecord const& record) { + application_log(rivet::diagnostic_json_line(record)); +}; +``` + +The same `diagnostic_sink` field is available in the Linux runtime config. On +Apple platforms, pass `diagnosticSink:` to +`EmbeddedRacketConfiguration.resolvedDefault` or `RivetClient`. On the Racket +side, `serve-fds` uses `current-rivet-diagnostic-sink`; direct `serve` callers +can pass `#:diagnostic-sink` explicitly. + +Diagnostic sinks run on runtime and request threads. They should be fast, +thread-safe, non-blocking, and must not call back into the same runtime. C++ +sink exceptions are isolated from the application lifecycle. + +Rivet never records RPC argument or result payloads. RPC begin records may name +the called API, and failure messages may contain application exception text; +treat the JSONL stream as operational log data and apply the same redaction and +retention policy as other crash reports. Racket-side messages are bounded to +4096 characters and avoid invoking custom printers on arbitrary raised values. diff --git a/platform/linux/Integration/main.cpp b/platform/linux/Integration/main.cpp index 6caeea9..223b370 100644 --- a/platform/linux/Integration/main.cpp +++ b/platform/linux/Integration/main.cpp @@ -7,6 +7,7 @@ #include #include #include +#include #include #include #include @@ -202,6 +203,27 @@ int main(int argc, char** argv) { config.entry_symbol = "start"; config.max_pending_requests = 32; + std::mutex diagnostics_mutex; + std::vector diagnostics; + config.diagnostic_sink = [&](rivet::DiagnosticRecord const& record) { + std::lock_guard lock(diagnostics_mutex); + diagnostics.push_back(record); + }; + + auto require_diagnostic = [&](std::string const& layer, + std::string const& event, + std::string const& status) { + std::lock_guard lock(diagnostics_mutex); + for (auto const& record : diagnostics) { + if (record.layer == layer && record.event == event && + record.status == status) { + return; + } + } + throw std::runtime_error("missing diagnostic: " + layer + "/" + event + + "/" + status); + }; + progress("starting backend"); rivet::linux_runtime::Backend backend(std::move(config)); backend.start(); @@ -311,6 +333,10 @@ int main(int argc, char** argv) { } progress("backend stopped"); + require_diagnostic("abi-bridge", "backend-init", "success"); + require_diagnostic("protocol", "handshake", "success"); + require_diagnostic("native-client", "rpc-dispatch", "success"); + require_diagnostic("native-runtime", "backend-stop", "success"); std::cout << "Rivet embedded Linux round-trip passed\n"; return 0; } catch (std::exception const& error) { diff --git a/platform/linux/runtime/backend.cpp b/platform/linux/runtime/backend.cpp index ed45059..4751e7e 100644 --- a/platform/linux/runtime/backend.cpp +++ b/platform/linux/runtime/backend.cpp @@ -148,6 +148,17 @@ std::exception_ptr stopped_error() { return std::make_exception_ptr(std::runtime_error("Rivet backend stopped")); } +std::string exception_message(std::exception_ptr error) noexcept { + if (!error) return "unknown failure"; + try { + std::rethrow_exception(error); + } catch (std::exception const& exception) { + return exception.what(); + } catch (...) { + return "non-standard exception"; + } +} + ptr quoted_symbol(std::string const& name) { auto const quote = Sstring_to_symbol("quote"); auto const module = Sstring_to_symbol(name.c_str()); @@ -179,11 +190,13 @@ class Backend::Impl { throw std::logic_error("Rivet backend instances cannot be restarted"); } started_ = true; + emit_diagnostic("native-runtime", "backend-start", "begin"); auto endpoints = create_socket_endpoints(); auto server_read = std::move(endpoints.server_read); auto server_write = std::move(endpoints.server_write); transport_ = std::make_unique(std::move(endpoints.native)); + emit_diagnostic("transport", "channel-opened", "success"); ready_ = std::make_shared>(); ready_future_ = ready_->get_future(); @@ -200,7 +213,17 @@ class Backend::Impl { // The first valid frame from Racket is Hello. Waiting for it makes start // fail synchronously when boot files, core.zo, or the entry point are bad. - ready_future_.get(); + try { + ready_future_.get(); + emit_diagnostic("protocol", "handshake", "success"); + emit_diagnostic("native-runtime", "backend-start", "success"); + } catch (...) { + emit_diagnostic("protocol", "handshake", "failure", + exception_message(std::current_exception())); + emit_diagnostic("native-runtime", "backend-start", "failure", + exception_message(std::current_exception())); + throw; + } } void stop() { @@ -211,13 +234,16 @@ class Backend::Impl { if (!started_) { return; } + stopping_.store(true, std::memory_order_release); running_.store(false, std::memory_order_release); } + emit_diagnostic("native-runtime", "backend-stop", "begin"); try { std::lock_guard write_lock(write_mutex_); auto* transport = transport_.get(); if (transport != nullptr) { + note_protocol_event("shutdown"); write_frame(*transport, Frame{MessageType::Shutdown, 0, {}}); } } catch (...) { @@ -238,6 +264,7 @@ class Backend::Impl { } reject_all(stopped_error()); + emit_diagnostic("native-runtime", "backend-stop", "success"); } bool running() const noexcept { @@ -286,6 +313,9 @@ class Backend::Impl { auto* transport = transport_.get(); if (transport != nullptr) { + note_protocol_event("cancel"); + emit_diagnostic("native-client", "request-cancel", "begin", "", + request_id); write_frame(*transport, Frame{MessageType::Cancel, request_id, {}}); } } @@ -296,6 +326,31 @@ class Backend::Impl { } private: + void note_protocol_event(std::string event) noexcept { + std::lock_guard lock(diagnostic_mutex_); + last_protocol_event_ = std::move(event); + } + + void emit_diagnostic(std::string layer, + std::string event, + std::string status, + std::string message = {}, + std::optional request_id = std::nullopt) + noexcept { + DiagnosticRecord record; + { + std::lock_guard lock(diagnostic_mutex_); + record = DiagnosticRecord{std::move(layer), std::move(event), + std::move(status), last_protocol_event_, + request_id, std::move(message)}; + } + try { + if (config_.diagnostic_sink) config_.diagnostic_sink(record); + } catch (...) { + // Diagnostics must never become a new runtime failure path. + } + } + std::uint64_t submit_request(std::string rpc_name, Value::List arguments, PendingRequest pending) { @@ -323,6 +378,9 @@ class Backend::Impl { } } + note_protocol_event("request"); + emit_diagnostic("native-client", "rpc-dispatch", "begin", rpc_name, id); + Value::List request; request.reserve(arguments.size() + 1); request.emplace_back(std::move(rpc_name)); @@ -355,6 +413,8 @@ class Backend::Impl { write_frame(*transport, Frame{MessageType::Cancel, id, {}}); } } catch (...) { + emit_diagnostic("transport", "request-write", "failure", + exception_message(std::current_exception()), id); fail_request(id, std::current_exception()); } @@ -362,6 +422,8 @@ class Backend::Impl { } void racket_main(UniqueFd server_read, UniqueFd server_write) noexcept { + emit_diagnostic("abi-bridge", "backend-init", "begin"); + bool initialized = false; try { racket_boot_arguments_t boot{}; boot.boot1_path = config_.petite_boot.c_str(); @@ -386,13 +448,25 @@ class Backend::Impl { Scons(Sfixnum(server_read.get()), Scons(Sfixnum(server_write.get()), Snil)); + initialized = true; + emit_diagnostic("abi-bridge", "backend-init", "success"); + // The ports created by serve-fds own distinct descriptors referring to // the same full-duplex socket. They close them during server teardown. (void)server_read.release(); (void)server_write.release(); (void)racket_apply(procedure, args); Sscheme_deinit(); + emit_diagnostic( + "racket-backend", "backend-exit", + stopping_.load(std::memory_order_acquire) ? "success" : "failure", + stopping_.load(std::memory_order_acquire) + ? "" + : "backend returned before native shutdown"); } catch (...) { + emit_diagnostic(initialized ? "racket-backend" : "abi-bridge", + initialized ? "backend-exit" : "backend-init", "failure", + exception_message(std::current_exception())); // Owned descriptors are closed on native startup failures. Once they are // transferred, serve-fds owns their lifetime. } @@ -414,28 +488,40 @@ class Backend::Impl { auto frame = read_frame(*transport); if (!frame.has_value()) { + if (!stopping_.load(std::memory_order_acquire)) { + emit_diagnostic("transport", "channel-closed", "failure", + "Rivet transport closed unexpectedly"); + } break; } switch (frame->type) { case MessageType::Hello: + note_protocol_event("hello"); accept_hello(*frame); break; case MessageType::Response: + note_protocol_event("response"); + emit_diagnostic("native-client", "rpc-dispatch", "success", "", + frame->id); resolve_request(frame->id, decode_value(frame->payload)); break; case MessageType::Error: { + note_protocol_event("error"); auto error_value = decode_value(frame->payload); std::string message{"Rivet backend error"}; if (auto* text = std::get_if(&error_value.data)) { message = *text; } + emit_diagnostic("racket-backend", "rpc-dispatch", "failure", + message, frame->id); fail_request( frame->id, std::make_exception_ptr(std::runtime_error(std::move(message)))); break; } case MessageType::Event: + note_protocol_event("event"); deliver_event(decode_value(frame->payload)); break; default: @@ -447,6 +533,8 @@ class Backend::Impl { set_ready_exception(stopped_error()); } } catch (...) { + emit_diagnostic("protocol", "reader-loop", "failure", + exception_message(std::current_exception())); set_ready_exception(std::current_exception()); reject_all(std::current_exception()); } @@ -602,6 +690,7 @@ class Backend::Impl { std::mutex write_mutex_; std::mutex pending_mutex_; std::mutex event_mutex_; + std::mutex diagnostic_mutex_; std::unique_ptr transport_; std::thread racket_thread_; std::thread reader_thread_; @@ -610,6 +699,8 @@ class Backend::Impl { detail::RequestIdAllocator request_ids_; std::atomic running_{false}; std::atomic hello_seen_{false}; + std::atomic stopping_{false}; + std::string last_protocol_event_{"none"}; bool started_{false}; std::shared_ptr> ready_; std::future ready_future_; diff --git a/platform/linux/runtime/backend.hpp b/platform/linux/runtime/backend.hpp index 303d274..e58e0ad 100644 --- a/platform/linux/runtime/backend.hpp +++ b/platform/linux/runtime/backend.hpp @@ -11,6 +11,7 @@ #include #include +#include "rivet/diagnostics.hpp" #include "rivet/protocol.hpp" namespace rivet::linux_runtime { @@ -26,6 +27,7 @@ struct RacketRuntimeConfig { std::string module_name{"backend"}; std::string entry_symbol{"start"}; std::size_t max_pending_requests{1024}; + DiagnosticSink diagnostic_sink{default_diagnostic_sink()}; }; struct PendingCall { diff --git a/platform/macos/Integration/Sources/main.swift b/platform/macos/Integration/Sources/main.swift index 9d2d250..5fc44f7 100644 --- a/platform/macos/Integration/Sources/main.swift +++ b/platform/macos/Integration/Sources/main.swift @@ -3,6 +3,25 @@ import Foundation import RivetEmbedding import RivetRuntime +private final class DiagnosticCapture: @unchecked Sendable { + private let lock = NSLock() + private var records: [RivetDiagnosticRecord] = [] + + func append(_ record: RivetDiagnosticRecord) { + lock.lock() + records.append(record) + lock.unlock() + } + + func contains(layer: String, event: String, status: String) -> Bool { + lock.lock() + defer { lock.unlock() } + return records.contains { + $0.layer == layer && $0.event == event && $0.status == status + } + } +} + #if arch(arm64) private let benchmarkArchitecture = "arm64" #elseif arch(x86_64) @@ -120,9 +139,11 @@ private func runBenchmark( struct RivetIntegration { static func main() async throws { let benchmark = try benchmarkMode() + let diagnostics = DiagnosticCapture() let configuration = try EmbeddedRacketConfiguration.resolvedDefault( moduleName: "backend", - entryName: "start" + entryName: "start", + diagnosticSink: { diagnostics.append($0) } ) let backend = EmbeddedRacketBackend(configuration: configuration) @@ -145,6 +166,12 @@ struct RivetIntegration { startupMilliseconds: startupMilliseconds ) backend.stop() + try requireDiagnostic( + diagnostics, + layer: "native-runtime", + event: "backend-stop", + status: "success" + ) return } @@ -177,6 +204,15 @@ struct RivetIntegration { backend.stop() backend.stop() // stop is intentionally idempotent. + try requireDiagnostic( + diagnostics, layer: "abi-bridge", event: "backend-init", status: "success") + try requireDiagnostic( + diagnostics, layer: "protocol", event: "handshake", status: "success") + try requireDiagnostic( + diagnostics, layer: "native-client", event: "rpc-dispatch", status: "success") + try requireDiagnostic( + diagnostics, layer: "native-runtime", event: "backend-stop", status: "success") + do { try backend.start() throw IntegrationError.restartWasAllowed @@ -186,6 +222,17 @@ struct RivetIntegration { print("Rivet embedded macOS round-trip passed") } + + private static func requireDiagnostic( + _ diagnostics: DiagnosticCapture, + layer: String, + event: String, + status: String + ) throws { + guard diagnostics.contains(layer: layer, event: event, status: status) else { + throw IntegrationError.missingDiagnostic("\(layer)/\(event)/\(status)") + } + } } enum IntegrationError: Error, CustomStringConvertible { @@ -193,6 +240,7 @@ enum IntegrationError: Error, CustomStringConvertible { case restartWasAllowed case invalidArguments([String]) case benchmarkEncoding + case missingDiagnostic(String) case unexpectedWorkingDirectory(expected: String, actual: String) var description: String { @@ -205,6 +253,8 @@ enum IntegrationError: Error, CustomStringConvertible { return "usage: RivetIntegration [--benchmark]; received: \(arguments)" case .benchmarkEncoding: return "failed to encode benchmark report as UTF-8 JSON" + case .missingDiagnostic(let record): + return "missing runtime diagnostic: \(record)" case .unexpectedWorkingDirectory(let expected, let actual): return "embedded runtime working directory was \(actual), expected \(expected)" } diff --git a/platform/macos/Sources/RivetEmbedding/EmbeddedRacketBackend.swift b/platform/macos/Sources/RivetEmbedding/EmbeddedRacketBackend.swift index dd10ecf..869d719 100644 --- a/platform/macos/Sources/RivetEmbedding/EmbeddedRacketBackend.swift +++ b/platform/macos/Sources/RivetEmbedding/EmbeddedRacketBackend.swift @@ -3,6 +3,26 @@ import Foundation import CRivetRacket import RivetRuntime +private final class EmbeddedDiagnosticContext: @unchecked Sendable { + private let lock = NSLock() + private var lastProtocolEvent = "none" + + func forward(_ record: RivetDiagnosticRecord, to sink: RivetDiagnosticSink) { + if record.lastProtocolEvent != "none" { + lock.lock() + lastProtocolEvent = record.lastProtocolEvent + lock.unlock() + } + sink(record) + } + + func snapshot() -> String { + lock.lock() + defer { lock.unlock() } + return lastProtocolEvent + } +} + public struct EmbeddedRacketConfiguration: Sendable { public var executable: URL public var petiteBoot: URL @@ -15,6 +35,7 @@ public struct EmbeddedRacketConfiguration: Sendable { public var moduleName: String public var entryName: String public var maxPendingRequests: Int + public var diagnosticSink: RivetDiagnosticSink public init( executable: URL, @@ -25,7 +46,8 @@ public struct EmbeddedRacketConfiguration: Sendable { workingDirectory: URL? = nil, moduleName: String = "backend", entryName: String = "start", - maxPendingRequests: Int = 1024 + maxPendingRequests: Int = 1024, + diagnosticSink: @escaping RivetDiagnosticSink = RivetDiagnostics.standardError ) { precondition(maxPendingRequests > 0, "Rivet native pending request limit must be positive") self.executable = executable @@ -37,6 +59,7 @@ public struct EmbeddedRacketConfiguration: Sendable { self.moduleName = moduleName self.entryName = entryName self.maxPendingRequests = maxPendingRequests + self.diagnosticSink = diagnosticSink } /// Resolves the canonical Rivet runtime layout used by both packaged apps @@ -49,7 +72,8 @@ public struct EmbeddedRacketConfiguration: Sendable { public static func resolvedDefault( moduleName: String = "backend", entryName: String = "start", - maxPendingRequests: Int = 1024 + maxPendingRequests: Int = 1024, + diagnosticSink: @escaping RivetDiagnosticSink = RivetDiagnostics.standardError ) throws -> EmbeddedRacketConfiguration { let executable = Bundle.main.executableURL ?? URL(fileURLWithPath: CommandLine.arguments[0]).standardizedFileURL @@ -63,7 +87,8 @@ public struct EmbeddedRacketConfiguration: Sendable { candidateRoots: roots, moduleName: moduleName, entryName: entryName, - maxPendingRequests: maxPendingRequests + maxPendingRequests: maxPendingRequests, + diagnosticSink: diagnosticSink ) } @@ -73,6 +98,7 @@ public struct EmbeddedRacketConfiguration: Sendable { moduleName: String, entryName: String, maxPendingRequests: Int, + diagnosticSink: @escaping RivetDiagnosticSink, fileManager: FileManager = .default ) throws -> EmbeddedRacketConfiguration { var searchedRoots: [URL] = [] @@ -105,7 +131,8 @@ public struct EmbeddedRacketConfiguration: Sendable { workingDirectory: root, moduleName: moduleName, entryName: entryName, - maxPendingRequests: maxPendingRequests + maxPendingRequests: maxPendingRequests, + diagnosticSink: diagnosticSink ) } @@ -143,6 +170,7 @@ public final class EmbeddedRacketBackend: @unchecked Sendable { private let responsePipe = Pipe() private let lifecycle = NSCondition() private let serverExited = DispatchSemaphore(value: 0) + private let diagnosticContext = EmbeddedDiagnosticContext() private var state: Lifecycle = .created private var serverThread: Thread? @@ -150,7 +178,10 @@ public final class EmbeddedRacketBackend: @unchecked Sendable { public private(set) lazy var client = RivetClient( input: responsePipe.fileHandleForReading, output: requestPipe.fileHandleForWriting, - maxPendingRequests: configuration.maxPendingRequests + maxPendingRequests: configuration.maxPendingRequests, + diagnosticSink: { [diagnosticContext, sink = configuration.diagnosticSink] record in + diagnosticContext.forward(record, to: sink) + } ) public init(configuration: EmbeddedRacketConfiguration) { @@ -165,6 +196,7 @@ public final class EmbeddedRacketBackend: @unchecked Sendable { } state = .starting lifecycle.unlock() + diagnose(layer: "abi-bridge", event: "backend-init", status: "begin") do { if let workingDirectory = configuration.workingDirectory, @@ -191,9 +223,25 @@ public final class EmbeddedRacketBackend: @unchecked Sendable { let config = configuration let completion = serverExited + let diagnostics = diagnosticContext let server = Thread { defer { completion.signal() } - runEmbeddedRacket(config, inputFD: racketInput, outputFD: racketOutput) + let result = runEmbeddedRacket( + config, + inputFD: racketInput, + outputFD: racketOutput + ) + if result != 0 { + config.diagnosticSink( + RivetDiagnosticRecord( + layer: "abi-bridge", + event: "backend-exit", + status: "failure", + lastProtocolEvent: diagnostics.snapshot(), + message: "rivet_racket_run returned status \(result)" + ) + ) + } } server.name = "Rivet Racket CS" server.qualityOfService = .userInitiated @@ -205,6 +253,12 @@ public final class EmbeddedRacketBackend: @unchecked Sendable { // Hello is emitted only after the Racket module and server have loaded. try client.start(onEvent: onEvent) + diagnose( + layer: "abi-bridge", + event: "backend-init", + status: "success", + lastProtocolEvent: "hello" + ) lifecycle.lock() if state == .starting { @@ -228,6 +282,12 @@ public final class EmbeddedRacketBackend: @unchecked Sendable { closeNativePipeEndpoints() waitForServerExit() } + diagnose( + layer: "abi-bridge", + event: "backend-init", + status: "failure", + message: String(describing: error) + ) throw error } } @@ -251,6 +311,11 @@ public final class EmbeddedRacketBackend: @unchecked Sendable { case .starting, .running: state = .stopping lifecycle.unlock() + diagnose( + layer: "native-runtime", + event: "backend-stop", + status: "begin" + ) case .stopping: // The loop above consumes this state. lifecycle.unlock() @@ -265,6 +330,7 @@ public final class EmbeddedRacketBackend: @unchecked Sendable { state = .stopped lifecycle.broadcast() lifecycle.unlock() + diagnose(layer: "native-runtime", event: "backend-stop", status: "success") } deinit { @@ -288,13 +354,32 @@ public final class EmbeddedRacketBackend: @unchecked Sendable { serverThread = nil lifecycle.unlock() } + + private func diagnose( + layer: String, + event: String, + status: String, + lastProtocolEvent: String? = nil, + message: String? = nil + ) { + configuration.diagnosticSink( + RivetDiagnosticRecord( + layer: layer, + event: event, + status: status, + lastProtocolEvent: lastProtocolEvent ?? diagnosticContext.snapshot(), + message: message + ) + ) + } } +@discardableResult private func runEmbeddedRacket( _ config: EmbeddedRacketConfiguration, inputFD: Int32, outputFD: Int32 -) { +) -> Int32 { let result: Int32 = config.executable.path.withCString { executable in config.petiteBoot.path.withCString { petite in config.schemeBoot.path.withCString { scheme in @@ -326,6 +411,7 @@ private func runEmbeddedRacket( Darwin.close(inputFD) Darwin.close(outputFD) } + return result } public enum EmbeddedBackendError: Error, CustomStringConvertible { diff --git a/platform/macos/Sources/RivetRuntime/Client.swift b/platform/macos/Sources/RivetRuntime/Client.swift index 4a27851..647d6c3 100644 --- a/platform/macos/Sources/RivetRuntime/Client.swift +++ b/platform/macos/Sources/RivetRuntime/Client.swift @@ -121,24 +121,29 @@ public final class RivetClient: @unchecked Sendable { private let input: FileHandle private let output: FileHandle private let maxPendingRequests: Int + private let diagnosticSink: RivetDiagnosticSink private let stateLock = NSLock() private let writeLock = NSLock() + private let diagnosticLock = NSLock() private let readerQueue = DispatchQueue(label: "dev.rivet.protocol-reader") private var requestIDs = RequestIDAllocator() private var lifecycle = ClientLifecycleState() private var pending: [UInt64: CheckedContinuation] = [:] private var eventHandler: EventHandler? + private var lastProtocolEvent = "none" public init( input: FileHandle, output: FileHandle, - maxPendingRequests: Int = 1024 + maxPendingRequests: Int = 1024, + diagnosticSink: @escaping RivetDiagnosticSink = RivetDiagnostics.standardError ) { precondition(maxPendingRequests > 0, "Rivet native pending request limit must be positive") self.input = input self.output = output self.maxPendingRequests = maxPendingRequests + self.diagnosticSink = diagnosticSink } deinit { @@ -154,9 +159,11 @@ public final class RivetClient: @unchecked Sendable { stateLock.unlock() throw error } + diagnose(layer: "native-client", event: "client-start", status: "begin") do { let hello = try readFrame() + noteProtocolEvent("hello") try validateHello(hello) stateLock.lock() @@ -168,10 +175,24 @@ public final class RivetClient: @unchecked Sendable { stateLock.unlock() throw error } + diagnose(layer: "protocol", event: "handshake", status: "success") + diagnose(layer: "native-client", event: "client-start", status: "success") } catch { stateLock.lock() lifecycle.failStart() stateLock.unlock() + diagnose( + layer: "protocol", + event: "handshake", + status: "failure", + message: String(describing: error) + ) + diagnose( + layer: "native-client", + event: "client-start", + status: "failure", + message: String(describing: error) + ) throw error } @@ -205,6 +226,14 @@ public final class RivetClient: @unchecked Sendable { do { var values: [RivetValue] = [.string(name)] values.append(contentsOf: arguments) + noteProtocolEvent("request") + diagnose( + layer: "native-client", + event: "rpc-dispatch", + status: "begin", + requestID: id, + message: name + ) try write( RivetFrame( type: .request, @@ -216,6 +245,13 @@ public final class RivetClient: @unchecked Sendable { cancel(cancelledID) } } catch { + diagnose( + layer: "transport", + event: "request-write", + status: "failure", + requestID: id, + message: String(describing: error) + ) if let continuation = takePending(id) { continuation.resume(throwing: error) } @@ -250,6 +286,13 @@ public final class RivetClient: @unchecked Sendable { } do { + noteProtocolEvent("cancel") + diagnose( + layer: "native-client", + event: "request-cancel", + status: "begin", + requestID: requestID + ) try output.write(contentsOf: data) stateLock.unlock() writeLock.unlock() @@ -270,9 +313,14 @@ public final class RivetClient: @unchecked Sendable { stateLock.unlock() if wasRunning { + diagnose(layer: "native-client", event: "client-stop", status: "begin") + noteProtocolEvent("shutdown") try? write(RivetFrame(type: .shutdown, id: 0)) } failAll(ClientError.stopped) + if wasRunning { + diagnose(layer: "native-client", event: "client-stop", status: "success") + } } private func registerPending( @@ -351,9 +399,16 @@ public final class RivetClient: @unchecked Sendable { do { while isRunning { let frame = try readFrame() + noteProtocolEvent(protocolEventName(frame.type)) switch frame.type { case .response: let value = try decodeRivetValue(frame.payload) + diagnose( + layer: "native-client", + event: "rpc-dispatch", + status: "success", + requestID: frame.id + ) takePending(frame.id)?.resume(returning: value) case .error: let value = try decodeRivetValue(frame.payload) @@ -363,6 +418,13 @@ public final class RivetClient: @unchecked Sendable { } else { message = "Rivet backend error" } + diagnose( + layer: "racket-backend", + event: "rpc-dispatch", + status: "failure", + requestID: frame.id, + message: message + ) takePending(frame.id)?.resume(throwing: ClientError.backend(message)) case .event: deliverEvent(frame) @@ -417,10 +479,56 @@ public final class RivetClient: @unchecked Sendable { private func finishWithError(_ error: Error) { stateLock.lock() - lifecycle.stop() + let wasRunning = lifecycle.stop() stateLock.unlock() + if wasRunning { + diagnose( + layer: error is RivetProtocolError ? "protocol" : "transport", + event: "reader-loop", + status: "failure", + message: String(describing: error) + ) + } failAll(error) } + + private func protocolEventName(_ type: RivetMessageType) -> String { + switch type { + case .hello: return "hello" + case .request: return "request" + case .response: return "response" + case .error: return "error" + case .event: return "event" + case .cancel: return "cancel" + case .shutdown: return "shutdown" + } + } + + private func noteProtocolEvent(_ event: String) { + diagnosticLock.lock() + lastProtocolEvent = event + diagnosticLock.unlock() + } + + private func diagnose( + layer: String, + event: String, + status: String, + requestID: UInt64? = nil, + message: String? = nil + ) { + diagnosticLock.lock() + let record = RivetDiagnosticRecord( + layer: layer, + event: event, + status: status, + lastProtocolEvent: lastProtocolEvent, + requestID: requestID, + message: message + ) + diagnosticLock.unlock() + diagnosticSink(record) + } } public enum ClientError: Error, CustomStringConvertible { diff --git a/platform/macos/Sources/RivetRuntime/Diagnostics.swift b/platform/macos/Sources/RivetRuntime/Diagnostics.swift new file mode 100644 index 0000000..ea39160 --- /dev/null +++ b/platform/macos/Sources/RivetRuntime/Diagnostics.swift @@ -0,0 +1,58 @@ +import Foundation + +public struct RivetDiagnosticRecord: Sendable, Equatable { + public var layer: String + public var event: String + public var status: String + public var lastProtocolEvent: String + public var requestID: UInt64? + public var message: String? + + public init( + layer: String, + event: String, + status: String, + lastProtocolEvent: String = "none", + requestID: UInt64? = nil, + message: String? = nil + ) { + self.layer = layer + self.event = event + self.status = status + self.lastProtocolEvent = lastProtocolEvent + self.requestID = requestID + self.message = message + } + + public func jsonLine() -> String { + var object: [String: Any] = [ + "schema": "rivet.diagnostic.v1", + "layer": layer, + "event": event, + "status": status, + "last_protocol_event": lastProtocolEvent + ] + if let requestID { object["request_id"] = requestID } + if let message, !message.isEmpty { object["message"] = message } + guard let data = try? JSONSerialization.data( + withJSONObject: object, + options: [.sortedKeys] + ) else { + return "{\"schema\":\"rivet.diagnostic.v1\",\"layer\":\"native-runtime\",\"event\":\"diagnostic-encoding\",\"status\":\"failure\",\"last_protocol_event\":\"none\"}" + } + return String(decoding: data, as: UTF8.self) + } +} + +public typealias RivetDiagnosticSink = @Sendable (RivetDiagnosticRecord) -> Void + +public enum RivetDiagnostics { + private static let outputLock = NSLock() + + public static let standardError: RivetDiagnosticSink = { record in + let data = Data((record.jsonLine() + "\n").utf8) + outputLock.lock() + defer { outputLock.unlock() } + try? FileHandle.standardError.write(contentsOf: data) + } +} diff --git a/platform/macos/Tests/RivetRuntimeTests/DiagnosticsTests.swift b/platform/macos/Tests/RivetRuntimeTests/DiagnosticsTests.swift new file mode 100644 index 0000000..0f9993c --- /dev/null +++ b/platform/macos/Tests/RivetRuntimeTests/DiagnosticsTests.swift @@ -0,0 +1,26 @@ +import Foundation +import Testing +@testable import RivetRuntime + +@Test func diagnosticRecordUsesStableJSONLShape() throws { + let record = RivetDiagnosticRecord( + layer: "racket-backend", + event: "rpc-dispatch", + status: "failure", + lastProtocolEvent: "request\nread", + requestID: 42, + message: "quote: \" and control: \u{1}" + ) + + let data = try #require(record.jsonLine().data(using: .utf8)) + let object = try #require( + JSONSerialization.jsonObject(with: data) as? [String: Any] + ) + #expect(object["schema"] as? String == "rivet.diagnostic.v1") + #expect(object["layer"] as? String == "racket-backend") + #expect(object["event"] as? String == "rpc-dispatch") + #expect(object["status"] as? String == "failure") + #expect(object["last_protocol_event"] as? String == "request\nread") + #expect((object["request_id"] as? NSNumber)?.uint64Value == 42) + #expect(object["message"] as? String == "quote: \" and control: \u{1}") +} diff --git a/platform/windows/Integration/main.cpp b/platform/windows/Integration/main.cpp index 769d1db..647d7d1 100644 --- a/platform/windows/Integration/main.cpp +++ b/platform/windows/Integration/main.cpp @@ -10,8 +10,10 @@ #include #include #include +#include #include #include +#include #include "backend.hpp" @@ -186,6 +188,27 @@ int main(int argc, char** argv) { config.dll_dir = runtime.wstring(); config.max_pending_requests = 1; + std::mutex diagnostics_mutex; + std::vector diagnostics; + config.diagnostic_sink = [&](rivet::DiagnosticRecord const& record) { + std::lock_guard lock(diagnostics_mutex); + diagnostics.push_back(record); + }; + + auto require_diagnostic = [&](std::string const& layer, + std::string const& event, + std::string const& status) { + std::lock_guard lock(diagnostics_mutex); + for (auto const& record : diagnostics) { + if (record.layer == layer && record.event == event && + record.status == status) { + return; + } + } + throw std::runtime_error("missing diagnostic: " + layer + "/" + event + + "/" + status); + }; + progress("starting backend"); rivet::windows::Backend backend(std::move(config)); auto const startup_begin = BenchmarkClock::now(); @@ -196,6 +219,7 @@ int main(int argc, char** argv) { if (benchmark) { run_benchmark(backend, startup_ms); backend.stop(); + require_diagnostic("native-runtime", "backend-stop", "success"); return 0; } @@ -310,6 +334,10 @@ int main(int argc, char** argv) { } progress("backend stopped"); + require_diagnostic("abi-bridge", "backend-init", "success"); + require_diagnostic("protocol", "handshake", "success"); + require_diagnostic("native-client", "rpc-dispatch", "success"); + require_diagnostic("native-runtime", "backend-stop", "success"); std::cout << "Rivet embedded Windows round-trip passed\n"; return 0; } catch (std::exception const& error) { diff --git a/platform/windows/runtime/backend.cpp b/platform/windows/runtime/backend.cpp index 2313288..7da300c 100644 --- a/platform/windows/runtime/backend.cpp +++ b/platform/windows/runtime/backend.cpp @@ -85,6 +85,17 @@ std::exception_ptr stopped_error() { return std::make_exception_ptr(std::runtime_error("Rivet backend stopped")); } +std::string exception_message(std::exception_ptr error) noexcept { + if (!error) return "unknown failure"; + try { + std::rethrow_exception(error); + } catch (std::exception const& exception) { + return exception.what(); + } catch (...) { + return "non-standard exception"; + } +} + ptr quoted_symbol(std::string const& name) { auto const quote = Sstring_to_symbol("quote"); auto const module = Sstring_to_symbol(name.c_str()); @@ -116,6 +127,7 @@ class Backend::Impl { throw std::logic_error("Rivet backend instances cannot be restarted"); } started_ = true; + emit_diagnostic("native-runtime", "backend-start", "begin"); auto request_pipe = create_pipe(); // native -> Racket auto response_pipe = create_pipe(); // Racket -> native @@ -125,6 +137,7 @@ class Backend::Impl { transport_ = std::make_unique( response_pipe.read.release(), request_pipe.write.release()); + emit_diagnostic("transport", "channel-opened", "success"); ready_ = std::make_shared>(); ready_future_ = ready_->get_future(); @@ -141,7 +154,17 @@ class Backend::Impl { // The first valid frame from Racket is Hello. Waiting for it makes start // fail synchronously when boot files, core.zo, or the entry point are bad. - ready_future_.get(); + try { + ready_future_.get(); + emit_diagnostic("protocol", "handshake", "success"); + emit_diagnostic("native-runtime", "backend-start", "success"); + } catch (...) { + emit_diagnostic("protocol", "handshake", "failure", + exception_message(std::current_exception())); + emit_diagnostic("native-runtime", "backend-start", "failure", + exception_message(std::current_exception())); + throw; + } } void stop() { @@ -154,10 +177,12 @@ class Backend::Impl { if (!started_) { return; } + stopping_.store(true, std::memory_order_release); // Writers recheck this flag after acquiring write_mutex_. Once false, no // new Request/Cancel may be written after the Shutdown boundary below. running_.store(false, std::memory_order_release); } + emit_diagnostic("native-runtime", "backend-stop", "begin"); // Try Shutdown even when a worker already marked running_ false. A native // reader failure does not necessarily mean the Racket server stopped @@ -167,6 +192,7 @@ class Backend::Impl { std::lock_guard write_lock(write_mutex_); auto* transport = transport_.get(); if (transport != nullptr) { + note_protocol_event("shutdown"); write_frame(*transport, Frame{MessageType::Shutdown, 0, {}}); } } catch (...) { @@ -196,6 +222,7 @@ class Backend::Impl { } reject_all(stopped_error()); + emit_diagnostic("native-runtime", "backend-stop", "success"); } bool running() const noexcept { @@ -248,6 +275,9 @@ class Backend::Impl { auto* transport = transport_.get(); if (transport != nullptr) { + note_protocol_event("cancel"); + emit_diagnostic("native-client", "request-cancel", "begin", "", + request_id); write_frame(*transport, Frame{MessageType::Cancel, request_id, {}}); } } @@ -258,6 +288,31 @@ class Backend::Impl { } private: + void note_protocol_event(std::string event) noexcept { + std::lock_guard lock(diagnostic_mutex_); + last_protocol_event_ = std::move(event); + } + + void emit_diagnostic(std::string layer, + std::string event, + std::string status, + std::string message = {}, + std::optional request_id = std::nullopt) + noexcept { + DiagnosticRecord record; + { + std::lock_guard lock(diagnostic_mutex_); + record = DiagnosticRecord{std::move(layer), std::move(event), + std::move(status), last_protocol_event_, + request_id, std::move(message)}; + } + try { + if (config_.diagnostic_sink) config_.diagnostic_sink(record); + } catch (...) { + // Diagnostics must never become a new runtime failure path. + } + } + std::uint64_t submit_request(std::string rpc_name, Value::List arguments, PendingRequest pending) { @@ -287,6 +342,9 @@ class Backend::Impl { } } + note_protocol_event("request"); + emit_diagnostic("native-client", "rpc-dispatch", "begin", rpc_name, id); + Value::List request; request.reserve(arguments.size() + 1); request.emplace_back(std::move(rpc_name)); @@ -326,6 +384,8 @@ class Backend::Impl { write_frame(*transport, Frame{MessageType::Cancel, id, {}}); } } catch (...) { + emit_diagnostic("transport", "request-write", "failure", + exception_message(std::current_exception()), id); fail_request(id, std::current_exception()); } @@ -333,6 +393,8 @@ class Backend::Impl { } void racket_main(UniqueHandle server_read, UniqueHandle server_write) noexcept { + emit_diagnostic("abi-bridge", "backend-init", "begin"); + bool initialized = false; try { // `unsafe-file-descriptor->port` consumes a Rktio system descriptor. // On Windows that descriptor is the native HANDLE value, not a CRT fd. @@ -369,13 +431,25 @@ class Backend::Impl { auto const args = Scons(Sfixnum(in_handle), Scons(Sfixnum(out_handle), Snil)); + initialized = true; + emit_diagnostic("abi-bridge", "backend-init", "success"); + // From this point the Racket ports created by serve-fds own the native // handles and close them during server teardown. (void)server_read.release(); (void)server_write.release(); (void)racket_apply(procedure, args); Sscheme_deinit(); + emit_diagnostic( + "racket-backend", "backend-exit", + stopping_.load(std::memory_order_acquire) ? "success" : "failure", + stopping_.load(std::memory_order_acquire) + ? "" + : "backend returned before native shutdown"); } catch (...) { + emit_diagnostic(initialized ? "racket-backend" : "abi-bridge", + initialized ? "backend-exit" : "backend-init", "failure", + exception_message(std::current_exception())); // UniqueHandle closes endpoints for native startup failures. Once the // handles are transferred, serve-fds owns their lifetime on the Racket // side. Closing the server ends makes the native reader observe EOF. @@ -398,28 +472,40 @@ class Backend::Impl { auto frame = read_frame(*transport); if (!frame.has_value()) { + if (!stopping_.load(std::memory_order_acquire)) { + emit_diagnostic("transport", "channel-closed", "failure", + "Rivet transport closed unexpectedly"); + } break; } switch (frame->type) { case MessageType::Hello: + note_protocol_event("hello"); accept_hello(*frame); break; case MessageType::Response: + note_protocol_event("response"); + emit_diagnostic("native-client", "rpc-dispatch", "success", "", + frame->id); resolve_request(frame->id, decode_value(frame->payload)); break; case MessageType::Error: { + note_protocol_event("error"); auto error_value = decode_value(frame->payload); std::string message{"Rivet backend error"}; if (auto* text = std::get_if(&error_value.data)) { message = *text; } + emit_diagnostic("racket-backend", "rpc-dispatch", "failure", + message, frame->id); fail_request( frame->id, std::make_exception_ptr(std::runtime_error(std::move(message)))); break; } case MessageType::Event: + note_protocol_event("event"); deliver_event(decode_value(frame->payload)); break; default: @@ -431,6 +517,8 @@ class Backend::Impl { set_ready_exception(stopped_error()); } } catch (...) { + emit_diagnostic("protocol", "reader-loop", "failure", + exception_message(std::current_exception())); set_ready_exception(std::current_exception()); reject_all(std::current_exception()); } @@ -586,6 +674,7 @@ class Backend::Impl { std::mutex write_mutex_; std::mutex pending_mutex_; std::mutex event_mutex_; + std::mutex diagnostic_mutex_; std::unique_ptr transport_; std::thread racket_thread_; std::thread reader_thread_; @@ -594,6 +683,8 @@ class Backend::Impl { detail::RequestIdAllocator request_ids_; std::atomic running_{false}; std::atomic hello_seen_{false}; + std::atomic stopping_{false}; + std::string last_protocol_event_{"none"}; bool started_{false}; std::shared_ptr> ready_; std::future ready_future_; diff --git a/platform/windows/runtime/backend.hpp b/platform/windows/runtime/backend.hpp index 5367ba4..46a5980 100644 --- a/platform/windows/runtime/backend.hpp +++ b/platform/windows/runtime/backend.hpp @@ -11,6 +11,7 @@ #include #include +#include "rivet/diagnostics.hpp" #include "rivet/protocol.hpp" namespace rivet::windows { @@ -27,6 +28,7 @@ struct RacketRuntimeConfig { std::string config_dir; std::wstring dll_dir; std::size_t max_pending_requests{1024}; + DiagnosticSink diagnostic_sink{default_diagnostic_sink()}; }; struct PendingCall { diff --git a/rivet/backend.rkt b/rivet/backend.rkt index 279f70e..3599d73 100644 --- a/rivet/backend.rkt +++ b/rivet/backend.rkt @@ -1,6 +1,7 @@ #lang racket/base (require ffi/unsafe/port + json racket/async-channel racket/list racket/match @@ -19,6 +20,7 @@ enum-case serve serve-fds + current-rivet-diagnostic-sink registered-rpcs registered-events registered-states @@ -53,6 +55,24 @@ (define record-registry (make-hash)) (define enum-registry (make-hash)) (define current-event-emitter (make-parameter #f)) + +(define diagnostic-schema "rivet.diagnostic.v1") + +(define current-rivet-diagnostic-sink + (make-parameter + (lambda (record) + (write-json record (current-error-port)) + (newline (current-error-port)) + (flush-output (current-error-port))))) + +(define (safe-diagnostic-message raised) + (define message + (if (exn? raised) + (exn-message raised) + "non-exception value")) + (if (<= (string-length message) 4096) + message + (string-append (substring message 0 4080) "... [truncated]"))) ;; Request workers replace this identity wrapper with a cancellation barrier. ;; Keeping it private lets state-set! preserve the same behavior outside serve. (define current-state-commit-guard (make-parameter (lambda (thunk) (thunk)))) @@ -721,7 +741,8 @@ (define (serve in out #:max-pending-requests [max-pending-requests 1024] - #:max-outgoing-frames [max-outgoing-frames 64]) + #:max-outgoing-frames [max-outgoing-frames 64] + #:diagnostic-sink [diagnostic-sink void]) (unless (input-port? in) (raise-argument-error 'serve "input-port?" in)) (unless (output-port? out) @@ -732,6 +753,49 @@ (unless (and (exact-integer? max-outgoing-frames) (positive? max-outgoing-frames)) (raise-argument-error 'serve "positive exact integer" max-outgoing-frames)) + (unless (and (procedure? diagnostic-sink) + (procedure-arity-includes? diagnostic-sink 1)) + (raise-argument-error 'serve "procedure accepting one argument" diagnostic-sink)) + + (define diagnostic-lock (make-semaphore 1)) + (define last-protocol-event "none") + + (define (note-protocol-event! event) + (call-with-semaphore + diagnostic-lock + (lambda () (set! last-protocol-event event)))) + + (define (diagnose! layer event status + #:request-id [request-id #f] + #:message [message #f]) + (define record + (call-with-semaphore + diagnostic-lock + (lambda () + (define base + (hasheq 'schema diagnostic-schema + 'layer layer + 'event event + 'status status + 'last_protocol_event last-protocol-event)) + (define with-id + (if request-id (hash-set base 'request_id request-id) base)) + (if message (hash-set with-id 'message message) with-id)))) + (with-handlers ([(lambda (_) #t) void]) + (diagnostic-sink record))) + + (define (frame-event-name f) + (case (frame-type f) + [(1) "hello"] + [(2) "request"] + [(3) "response"] + [(4) "error"] + [(5) "event"] + [(6) "cancel"] + [(7) "shutdown"] + [else "unknown-frame"])) + + (diagnose! "racket-backend" "backend-init" "begin") ;; Keep application/request threads and the writer in separate custodians. ;; Graceful shutdown can stop every producer first, drain already accepted @@ -755,6 +819,8 @@ (lambda () (with-handlers ([exn? (lambda (e) + (diagnose! "transport" "response-writer" "failure" + #:message (safe-diagnostic-message e)) (set-box! writer-error e))]) (let loop () (define response (async-channel-get responses)) @@ -773,6 +839,8 @@ ;; A bounded channel applies output backpressure. Waiting producers also ;; observe writer death, so a broken transport cannot strand request/Event ;; threads forever behind a full queue. + (when (frame? value) + (note-protocol-event! (frame-event-name value))) (define outcome (sync (handle-evt (async-channel-put-evt responses value) @@ -980,6 +1048,9 @@ (define (run-request! id rpc-name args internal-state-request? info) (with-handlers ([(lambda (_) #t) (lambda (raised) + (diagnose! "racket-backend" "rpc-dispatch" "failure" + #:request-id id + #:message (safe-diagnostic-message raised)) (request-error! id raised))]) (define result (if internal-state-request? @@ -1011,6 +1082,8 @@ (rpc-info-result-type info) value) (typed->wire (rpc-info-result-type info) value)))) + (diagnose! "racket-backend" "rpc-dispatch" "success" + #:request-id id) (finish! id message:response result))) (define (start-request! f) @@ -1028,8 +1101,14 @@ ;; failures local. (with-handlers ((exn:fail? (lambda (e) + (diagnose! "racket-backend" "rpc-dispatch" "failure" + #:request-id id + #:message (safe-diagnostic-message e)) (reject-request! id (exn-message e))))) (define-values (rpc-name args) (request->call (frame-payload f))) + (diagnose! "racket-backend" "rpc-dispatch" "begin" + #:request-id id + #:message (symbol->string rpc-name)) (define internal-state-request? (memq rpc-name '($state/get $state/set))) (define info @@ -1073,6 +1152,7 @@ (frame message:error id (error-message->payload "request cancelled"))))) (define (dispatch! f) + (note-protocol-event! (frame-event-name f)) (case (frame-type f) [(2) (start-request! f)] [(6) (cancel! (frame-id f))] @@ -1105,8 +1185,10 @@ (raise value) (error 'rivet/backend "backend requested exit with non-exception value")))]) + (diagnose! "racket-backend" "backend-init" "success") (send! (frame message:hello 0 (encode-value (list "rivet" protocol-version)))) + (diagnose! "protocol" "handshake" "success") (let loop () (unless stopped? (define f (read-frame in)) @@ -1162,31 +1244,40 @@ [reader-problem (raise reader-problem)] [else (void)])) - (dynamic-wind - void - (lambda () - (define first-exit - (sync - (handle-evt reader-dead-evt (lambda (_) 'reader)) - (handle-evt writer-dead-evt (lambda (_) 'writer)))) - (cond - [(eq? first-exit 'writer) - (define failure (unbox writer-error)) - (abort-server!) - (if failure - (raise failure) - (error 'serve "response writer terminated unexpectedly"))] - [else - (finish-after-reader!)])) - (lambda () - (unless cleanup-done? - (abort-server!)))) + (with-handlers ([(lambda (_) #t) + (lambda (raised) + (diagnose! "racket-backend" "backend-exit" "failure" + #:message (safe-diagnostic-message raised)) + (raise raised))]) + (dynamic-wind + void + (lambda () + (define first-exit + (sync + (handle-evt reader-dead-evt (lambda (_) 'reader)) + (handle-evt writer-dead-evt (lambda (_) 'writer)))) + (cond + [(eq? first-exit 'writer) + (define failure (unbox writer-error)) + (abort-server!) + (if failure + (raise failure) + (error 'serve "response writer terminated unexpectedly"))] + [else + (finish-after-reader!)])) + (lambda () + (unless cleanup-done? + (abort-server!))))) + + (diagnose! "racket-backend" "backend-exit" "success") (void)) (define (serve-fds in-fd out-fd #:max-pending-requests [max-pending-requests 1024] - #:max-outgoing-frames [max-outgoing-frames 64]) + #:max-outgoing-frames [max-outgoing-frames 64] + #:diagnostic-sink + [diagnostic-sink (current-rivet-diagnostic-sink)]) (unless (exact-integer? in-fd) (raise-argument-error 'serve-fds "exact-integer?" in-fd)) (unless (exact-integer? out-fd) @@ -1197,16 +1288,12 @@ void (lambda () (with-handlers ([exn? - (lambda (e) - ((error-display-handler) - (format "Rivet backend terminated: ~a" - (exn-message e)) - e) - (void))]) + (lambda (_e) (void))]) (serve in out #:max-pending-requests max-pending-requests - #:max-outgoing-frames max-outgoing-frames))) + #:max-outgoing-frames max-outgoing-frames + #:diagnostic-sink diagnostic-sink))) (lambda () (unless (port-closed? in) (close-input-port in)) (unless (port-closed? out) (close-output-port out))))) diff --git a/runtime/include/rivet/diagnostics.hpp b/runtime/include/rivet/diagnostics.hpp new file mode 100644 index 0000000..bcac9ea --- /dev/null +++ b/runtime/include/rivet/diagnostics.hpp @@ -0,0 +1,83 @@ +#pragma once + +#include +#include +#include +#include +#include +#include +#include +#include + +namespace rivet { + +// Dependency-light, provider-neutral lifecycle diagnostics shared by the +// embedded desktop runtimes. Applications can forward these records to their +// own structured logger or crash reporter without coupling Rivet to it. +struct DiagnosticRecord { + std::string layer; + std::string event; + std::string status; + std::string last_protocol_event{"none"}; + std::optional request_id; + std::string message; +}; + +using DiagnosticSink = std::function; + +inline std::string diagnostic_json_escape(std::string_view value) { + std::ostringstream out; + static constexpr char hex[] = "0123456789abcdef"; + for (unsigned char byte : value) { + switch (byte) { + case '"': out << "\\\""; break; + case '\\': out << "\\\\"; break; + case '\b': out << "\\b"; break; + case '\f': out << "\\f"; break; + case '\n': out << "\\n"; break; + case '\r': out << "\\r"; break; + case '\t': out << "\\t"; break; + default: + if (byte < 0x20) { + out << "\\u00" << hex[(byte >> 4) & 0x0f] << hex[byte & 0x0f]; + } else { + out << static_cast(byte); + } + } + } + return out.str(); +} + +inline std::string diagnostic_json_line(DiagnosticRecord const& record) { + std::ostringstream out; + out << "{\"schema\":\"rivet.diagnostic.v1\",\"layer\":\"" + << diagnostic_json_escape(record.layer) << "\",\"event\":\"" + << diagnostic_json_escape(record.event) << "\",\"status\":\"" + << diagnostic_json_escape(record.status) + << "\",\"last_protocol_event\":\"" + << diagnostic_json_escape(record.last_protocol_event) << '"'; + if (record.request_id.has_value()) { + out << ",\"request_id\":" << *record.request_id; + } + if (!record.message.empty()) { + out << ",\"message\":\"" << diagnostic_json_escape(record.message) << '"'; + } + out << '}'; + return out.str(); +} + +inline void write_diagnostic_to_stderr(DiagnosticRecord const& record) { + // A process can own several runtime threads. Serialize whole JSONL records + // so concurrent failures cannot corrupt the diagnostic stream. + static std::mutex output_mutex; + std::lock_guard lock(output_mutex); + std::cerr << diagnostic_json_line(record) << '\n'; +} + +inline DiagnosticSink default_diagnostic_sink() { + return [](DiagnosticRecord const& record) { + write_diagnostic_to_stderr(record); + }; +} + +} // namespace rivet diff --git a/runtime/tests/protocol_test.cpp b/runtime/tests/protocol_test.cpp index 55c08f2..ce8206b 100644 --- a/runtime/tests/protocol_test.cpp +++ b/runtime/tests/protocol_test.cpp @@ -1,3 +1,4 @@ +#include "rivet/diagnostics.hpp" #include "rivet/protocol.hpp" #include @@ -254,5 +255,14 @@ int main() { } assert(rejected_invalid_utf8_encode); + auto diagnostic = rivet::diagnostic_json_line(rivet::DiagnosticRecord{ + "racket-backend", "rpc-dispatch", "failure", "request\nread", 42, + "quote: \" and control: \x01"}); + assert(diagnostic == + "{\"schema\":\"rivet.diagnostic.v1\",\"layer\":\"racket-backend\"," + "\"event\":\"rpc-dispatch\",\"status\":\"failure\"," + "\"last_protocol_event\":\"request\\nread\",\"request_id\":42," + "\"message\":\"quote: \\\" and control: \\u0001\"}"); + return 0; } diff --git a/tests/backend-diagnostics.rkt b/tests/backend-diagnostics.rkt index 847e012..1e0196f 100644 --- a/tests/backend-diagnostics.rkt +++ b/tests/backend-diagnostics.rkt @@ -6,6 +6,24 @@ (define result-writer-called? (box #f)) (define exit-writer-called? (box #f)) +(define diagnostic-records (box '())) +(define diagnostic-lock (make-semaphore 1)) + +(define (capture-diagnostic! record) + (call-with-semaphore + diagnostic-lock + (lambda () + (set-box! diagnostic-records + (cons record (unbox diagnostic-records)))))) + +(define (diagnostic-record-exists? layer event status [request-id #f]) + (for/or ([record (in-list (unbox diagnostic-records))]) + (and (equal? (hash-ref record 'schema) "rivet.diagnostic.v1") + (equal? (hash-ref record 'layer) layer) + (equal? (hash-ref record 'event) event) + (equal? (hash-ref record 'status) status) + (or (not request-id) + (= (hash-ref record 'request_id) request-id))))) (struct explosive-result () #:property prop:custom-write @@ -45,7 +63,7 @@ (channel-put server-result (with-handlers ([exn? values]) - (serve server-in server-out) + (serve server-in server-out #:diagnostic-sink capture-diagnostic!) 'completed))))) (define hello (read-frame/timeout client-in)) @@ -87,3 +105,19 @@ (write-frame (frame message:shutdown 0 #"") client-out) (check-equal? (sync/timeout 2 server-result) 'completed) (thread-wait server-thread) + +(check-true + (diagnostic-record-exists? "protocol" "handshake" "success")) +(check-true + (diagnostic-record-exists? "racket-backend" "rpc-dispatch" "failure" 3)) +(check-true + (diagnostic-record-exists? "racket-backend" "backend-exit" "success")) +(define rpc-failure + (for/first ([record (in-list (unbox diagnostic-records))] + #:when (and (equal? (hash-ref record 'event) "rpc-dispatch") + (equal? (hash-ref record 'status) "failure") + (= (hash-ref record 'request_id -1) 3))) + record)) +(check-equal? (hash-ref rpc-failure 'last_protocol_event) "request") +(check-false (regexp-match? #rx"EXPLOSIVE-EXIT" + (hash-ref rpc-failure 'message)))