File size: 4,437 Bytes
afa0cbf | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 | use std::path::Path;
use std::process::Command;
use std::process::Stdio;
use std::sync::Arc;
use std::sync::Mutex;
use std::time::Duration;
use anyhow::Context;
use anyhow::Result;
use serde_json::Value;
use serde_json::json;
use tokio::sync::Notify;
#[derive(Clone, Default)]
pub(crate) struct JsonLogCapture {
lines: Arc<Mutex<Vec<String>>>,
updated: Arc<Notify>,
}
impl JsonLogCapture {
pub(crate) fn record(&self, line: String) {
self.lines
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push(line);
self.updated.notify_one();
}
pub(crate) async fn wait_for_event(&self, event_name: &str) -> Result<Value> {
let mut events = self.wait_for_events(event_name, /*count*/ 1).await?;
Ok(events.remove(0))
}
pub(crate) async fn wait_for_events(
&self,
event_name: &str,
count: usize,
) -> Result<Vec<Value>> {
let result = tokio::time::timeout(Duration::from_secs(10), async {
loop {
let updated = self.updated.notified();
let events = self
.events()?
.into_iter()
.filter(|event| event["fields"]["event.name"].as_str() == Some(event_name))
.collect::<Vec<_>>();
if events.len() >= count {
return Ok(events);
}
updated.await;
}
})
.await;
match result {
Ok(result) => result,
Err(_) => {
let lines = self
.lines
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.join("\n");
anyhow::bail!(
"timed out waiting for {count} JSON log event(s) named `{event_name}`; captured stderr:\n{lines}"
)
}
}
}
pub(crate) fn events(&self) -> Result<Vec<Value>> {
let lines = self
.lines
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
json_log_events(lines.iter().map(String::as_str))
}
}
pub fn app_server_json_shutdown_event(
binary: &str,
args: &[&str],
codex_home: &Path,
) -> Result<Value> {
std::fs::write(
codex_home.join("config.toml"),
"[features]\nplugins = false\n",
)?;
let output = Command::new(codex_utils_cargo_bin::cargo_bin(binary)?)
.stdin(Stdio::null())
.env("CODEX_HOME", codex_home)
.env(
"CODEX_APP_SERVER_MANAGED_CONFIG_PATH",
codex_home.join("managed_config.toml"),
)
.env("LOG_FORMAT", "json")
.env("RUST_LOG", "codex_app_server=info")
.args(args)
.output()?;
let stderr = String::from_utf8(output.stderr)?;
anyhow::ensure!(output.status.success(), "app-server failed: {stderr}");
let events = json_log_events(stderr.lines())
.with_context(|| format!("app-server stderr was not valid JSONL: {stderr}"))?;
let event = events
.iter()
.find(|event| event["fields"]["message"] == "processor task exited")
.context("missing INFO shutdown event in app-server JSON logs")?;
Ok(json!({
"level": event["level"],
"fields": event["fields"],
"target": event["target"],
}))
}
fn json_log_events<'a>(lines: impl IntoIterator<Item = &'a str>) -> Result<Vec<Value>> {
lines
.into_iter()
.filter(|line| !line.is_empty())
.map(|line| {
let event = serde_json::from_str::<Value>(line)
.with_context(|| format!("log line was not JSON: {line}"))?;
anyhow::ensure!(
event["level"].is_string()
&& event["fields"].is_object()
&& event["target"].is_string(),
"JSON log event did not include level, fields, and target: {line}"
);
let timestamp = event["timestamp"]
.as_str()
.with_context(|| format!("JSON log event did not include a timestamp: {line}"))?;
chrono::DateTime::parse_from_rfc3339(timestamp).with_context(|| {
format!("JSON log event timestamp was not RFC 3339: {timestamp}")
})?;
Ok(event)
})
.collect()
}
|