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()
}