Skip to content

Commit b78beea

Browse files
committed
feat(ocsf): emit full JSON records to supervisor stderr
Signed-off-by: John Myers <9696606+johntmyers@users.noreply.github.com>
1 parent a4954d3 commit b78beea

13 files changed

Lines changed: 634 additions & 82 deletions

File tree

‎crates/openshell-core/src/settings.rs‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -115,9 +115,9 @@ pub const PROPOSAL_APPROVAL_MODE_VALUES: &[&str] = &["manual", "auto"];
115115
pub const OCSF_SCHEMA_VERSION_VALUES: &[&str] = &["", "1.1", "1.3"];
116116

117117
pub const REGISTERED_SETTINGS: &[RegisteredSetting] = &[
118-
// When true the sandbox writes OCSF v1.8.0 JSONL records to
119-
// `/var/log/openshell-ocsf*.log` (daily rotation, 3 files) in addition
120-
// to the human-readable shorthand log. Defaults to false (no JSONL written).
118+
// When true the supervisor emits OCSF-JSON records to stderr and writes
119+
// JSONL to `/var/log/openshell-ocsf*.log` when file logging is available.
120+
// Defaults to false; shorthand output remains unchanged.
121121
RegisteredSetting {
122122
key: "ocsf_json_enabled",
123123
kind: SettingValueKind::Bool,

‎crates/openshell-ocsf/src/tracing_layers/jsonl_layer.rs‎

Lines changed: 21 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,14 @@
11
// SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
22
// SPDX-License-Identifier: Apache-2.0
33

4-
//! Tracing layer that writes OCSF JSONL to a writer.
4+
//! Tracing layer that writes OCSF JSONL or marked console records to a writer.
55
66
use std::io::Write;
77
use std::sync::Arc;
88
use std::sync::Mutex;
99
use std::sync::atomic::{AtomicBool, Ordering};
1010

11+
use chrono::Utc;
1112
use tracing::Subscriber;
1213
use tracing_subscriber::Layer;
1314
use tracing_subscriber::layer::Context;
@@ -33,6 +34,7 @@ pub struct OcsfJsonlLayer<W: Write + Send + 'static> {
3334
writer: Mutex<W>,
3435
enabled: Option<Arc<AtomicBool>>,
3536
target_version: Option<Arc<Mutex<String>>>,
37+
console_format: bool,
3638
}
3739

3840
impl<W: Write + Send + 'static> OcsfJsonlLayer<W> {
@@ -43,6 +45,7 @@ impl<W: Write + Send + 'static> OcsfJsonlLayer<W> {
4345
writer: Mutex::new(writer),
4446
enabled: None,
4547
target_version: None,
48+
console_format: false,
4649
}
4750
}
4851

@@ -65,6 +68,16 @@ impl<W: Write + Send + 'static> OcsfJsonlLayer<W> {
6568
self.target_version = Some(version);
6669
self
6770
}
71+
72+
/// Prefix each compact JSON record with a UTC timestamp and `OCSF-JSON`.
73+
///
74+
/// Console collectors can match the marker and decode the remainder as
75+
/// JSON. File sinks retain plain JSONL unless this option is selected.
76+
#[must_use]
77+
pub fn with_console_format(mut self) -> Self {
78+
self.console_format = true;
79+
self
80+
}
6881
}
6982

7083
impl<S, W> Layer<S> for OcsfJsonlLayer<W>
@@ -108,6 +121,13 @@ where
108121
Err(_) => return,
109122
}
110123
};
124+
let line = if self.console_format {
125+
let ts = Utc::now().format("%Y-%m-%dT%H:%M:%S%.3fZ");
126+
format!("{ts} OCSF-JSON {line}")
127+
} else {
128+
line
129+
};
130+
// Submit the entire record in one write, including its newline.
111131
let _ = w.write_all(line.as_bytes());
112132
}
113133
}

‎crates/openshell-ocsf/src/tracing_layers/shorthand_layer.rs‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -60,7 +60,8 @@ where
6060
if let Some(ocsf_event) = clone_current_event() {
6161
let line = ocsf_event.format_shorthand();
6262
if let Ok(mut w) = self.writer.lock() {
63-
let _ = writeln!(w, "{ts} OCSF {line}");
63+
let line = format!("{ts} OCSF {line}\n");
64+
let _ = w.write_all(line.as_bytes());
6465
}
6566
}
6667
} else if self.include_non_ocsf {
@@ -71,7 +72,8 @@ where
7172
let mut message = String::new();
7273
event.record(&mut MessageVisitor(&mut message));
7374
if let Ok(mut w) = self.writer.lock() {
74-
let _ = writeln!(w, "{ts} {level} {target}: {message}");
75+
let line = format!("{ts} {level} {target}: {message}\n");
76+
let _ = w.write_all(line.as_bytes());
7577
}
7678
}
7779
}
Lines changed: 131 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,131 @@
1+
// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
2+
// SPDX-License-Identifier: Apache-2.0
3+
4+
//! Collector framing, live settings, and whole-line console writes.
5+
6+
use std::io::Write;
7+
use std::sync::atomic::{AtomicBool, Ordering};
8+
use std::sync::{Arc, Mutex};
9+
10+
use openshell_ocsf::{
11+
EventContext, NetworkActivityBuilder, OcsfJsonlLayer, OcsfShorthandLayer, ocsf_emit,
12+
};
13+
use tracing_subscriber::layer::SubscriberExt;
14+
15+
/// Capture individual writes to verify queue entries are whole lines.
16+
#[derive(Clone, Default)]
17+
struct Records(Arc<Mutex<Vec<Vec<u8>>>>);
18+
19+
impl Write for Records {
20+
fn write(&mut self, bytes: &[u8]) -> std::io::Result<usize> {
21+
self.0.lock().unwrap().push(bytes.to_vec());
22+
Ok(bytes.len())
23+
}
24+
25+
fn flush(&mut self) -> std::io::Result<()> {
26+
Ok(())
27+
}
28+
}
29+
30+
fn payload(record: &[u8]) -> serde_json::Value {
31+
let line = std::str::from_utf8(record).unwrap();
32+
assert_eq!(line.lines().count(), 1);
33+
assert!(line.ends_with('\n'));
34+
assert!(!line.contains('\x1b'));
35+
let (timestamp, json) = line.trim_end().split_once(" OCSF-JSON ").unwrap();
36+
chrono::DateTime::parse_from_rfc3339(timestamp).unwrap();
37+
serde_json::from_str(json).unwrap()
38+
}
39+
40+
fn context() -> EventContext {
41+
EventContext {
42+
sandbox_id: "sb-console".to_string(),
43+
sandbox_name: "console-test".to_string(),
44+
container_image: "test-image".to_string(),
45+
hostname: "test-host".to_string(),
46+
product_version: "test".to_string(),
47+
proxy_ip: std::net::IpAddr::V4(std::net::Ipv4Addr::LOCALHOST),
48+
proxy_port: 3128,
49+
origin: openshell_ocsf::EventOrigin::Supervisor,
50+
}
51+
}
52+
53+
#[test]
54+
fn console_json_toggles_and_downgrades_without_changing_event_identity() {
55+
let records = Records::default();
56+
let enabled = Arc::new(AtomicBool::new(false));
57+
let version = Arc::new(Mutex::new(String::new()));
58+
let subscriber = tracing_subscriber::registry().with(
59+
OcsfJsonlLayer::new(records.clone())
60+
.with_console_format()
61+
.with_enabled_flag(enabled.clone())
62+
.with_target_version(version.clone()),
63+
);
64+
let event = NetworkActivityBuilder::new(&context())
65+
.dst_endpoint(openshell_ocsf::Endpoint::from_domain("example.com", 443))
66+
.message("first line\nsecond line")
67+
.build();
68+
let expected = event.to_json().unwrap();
69+
tracing::subscriber::with_default(subscriber, || {
70+
ocsf_emit!(event.clone());
71+
tracing::warn!("diagnostics are not JSON audit events");
72+
enabled.store(true, Ordering::Relaxed);
73+
ocsf_emit!(event.clone());
74+
*version.lock().unwrap() = "1.3".to_string();
75+
ocsf_emit!(event.clone());
76+
enabled.store(false, Ordering::Relaxed);
77+
ocsf_emit!(event);
78+
});
79+
let records = records.0.lock().unwrap();
80+
assert_eq!(records.len(), 2);
81+
assert_eq!(payload(&records[0]), expected);
82+
let downgraded = payload(&records[1]);
83+
assert_eq!(downgraded["metadata"]["version"], "1.3");
84+
assert_eq!(downgraded["metadata"]["uid"], expected["metadata"]["uid"]);
85+
assert_eq!(downgraded["message"], expected["message"]);
86+
}
87+
88+
#[test]
89+
fn concurrent_console_formats_submit_complete_records() {
90+
let records = Records::default();
91+
let subscriber = tracing_subscriber::registry()
92+
.with(OcsfShorthandLayer::new(records.clone()))
93+
.with(OcsfJsonlLayer::new(records.clone()).with_console_format());
94+
let dispatch = tracing::Dispatch::new(subscriber);
95+
std::thread::scope(|scope| {
96+
for worker in 0..8 {
97+
let dispatch = dispatch.clone();
98+
scope.spawn(move || {
99+
tracing::dispatcher::with_default(&dispatch, || {
100+
for index in 0..50 {
101+
tracing::info!("diagnostic {worker}:{index}");
102+
ocsf_emit!(
103+
NetworkActivityBuilder::new(&context())
104+
.dst_endpoint(openshell_ocsf::Endpoint::from_domain(
105+
"example.com",
106+
443
107+
))
108+
.message(format!("connection {worker}:{index}"))
109+
.build()
110+
);
111+
}
112+
});
113+
});
114+
}
115+
});
116+
let records = records.0.lock().unwrap();
117+
assert_eq!(records.len(), 8 * 50 * 3);
118+
let mut json_count = 0;
119+
for record in records.iter() {
120+
let line = std::str::from_utf8(record).unwrap();
121+
assert_eq!(line.lines().count(), 1);
122+
assert!(line.ends_with('\n'));
123+
if line.contains(" OCSF-JSON ") {
124+
payload(record);
125+
json_count += 1;
126+
} else {
127+
assert!(line.contains(" OCSF ") || line.contains(" INFO "));
128+
}
129+
}
130+
assert_eq!(json_count, 8 * 50);
131+
}

‎crates/openshell-supervisor/README.md‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,3 +15,9 @@ The supervisor admits policy and prepares credentials before constructing and at
1515
The client receives the supervisor's live provider state, bearer-token slot, and CA-path slot. Provider refresh, token rotation, and later CA publication must remain visible through those shared handles. Startup does not create independent copies of their current values.
1616

1717
The setup interface stays private to the supervisor. It adds no runtime backend registration, endpoint configuration, or public factory API. The public `run_sandbox` signature and standard backend selection remain unchanged.
18+
19+
## Console logging
20+
21+
Diagnostics and OCSF shorthand share a bounded, nonblocking stderr writer. With `ocsf_json_enabled=true`, the same writer also receives timestamped `OCSF-JSON` records. Each formatter submits a complete line in one write so concurrent producers cannot interleave records in the queue. The 1,024-line queue drops new lines when full rather than waiting for stderr.
22+
23+
Console JSON is installed independently of optional file appenders and uses an INFO filter independent of the diagnostic filter. Its runtime enabled flag and target schema version are shared with the existing JSONL file layer. The supervisor retains the writer guards until shutdown. The gateway log push layer continues to emit shorthand only.

0 commit comments

Comments
 (0)