Skip to main content

sloop/sources/
exec.rs

1use std::io::Write;
2use std::path::PathBuf;
3use std::process::{Command, Stdio};
4
5use serde::Deserialize;
6use serde_json::json;
7
8use crate::frontmatter::Frontmatter;
9use crate::outcome::Outcome;
10
11use super::{AuthoredTicket, SourceError, TicketSource};
12
13pub struct ExecTicketSource {
14    root: PathBuf,
15    argv: Vec<String>,
16}
17
18impl ExecTicketSource {
19    pub fn new(root: impl Into<PathBuf>, argv: Vec<String>) -> Self {
20        Self {
21            root: root.into(),
22            argv,
23        }
24    }
25
26    fn invoke(&self, request: &serde_json::Value) -> Result<std::process::Output, SourceError> {
27        let (program, arguments) = self
28            .argv
29            .split_first()
30            .ok_or_else(|| SourceError::new("sources.tickets.exec must name a command"))?;
31        let mut child = Command::new(program)
32            .args(arguments)
33            .current_dir(&self.root)
34            .stdin(Stdio::piped())
35            .stdout(Stdio::piped())
36            .stderr(Stdio::piped())
37            .spawn()
38            .map_err(|error| {
39                SourceError::new(format!("cannot start ticket source `{program}`: {error}"))
40            })?;
41        let input = serde_json::to_vec(request)
42            .map_err(|error| SourceError::new(format!("cannot encode source request: {error}")))?;
43        let write_result = child
44            .stdin
45            .take()
46            .expect("piped source stdin is available")
47            .write_all(&input);
48        if let Err(error) = write_result {
49            let _ = child.kill();
50            let _ = child.wait();
51            return Err(SourceError::new(format!(
52                "cannot write to ticket source `{program}`: {error}"
53            )));
54        }
55        child.wait_with_output().map_err(|error| {
56            SourceError::new(format!(
57                "cannot wait for ticket source `{program}`: {error}"
58            ))
59        })
60    }
61}
62
63impl TicketSource for ExecTicketSource {
64    fn pull(&self) -> Result<Vec<AuthoredTicket>, SourceError> {
65        let output = self.invoke(&json!({ "verb": "pull" }))?;
66        if !output.status.success() {
67            return Err(command_failed(&self.argv, &output));
68        }
69        parse_tickets(&output.stdout)
70    }
71
72    fn report(&self, ticket_id: &str, outcome: &Outcome) -> Result<(), SourceError> {
73        let output = self.invoke(&json!({
74            "verb": "report",
75            "ticket": ticket_id,
76            "outcome": outcome,
77        }))?;
78        if output.status.success() {
79            Ok(())
80        } else {
81            Err(command_failed(&self.argv, &output))
82        }
83    }
84}
85
86#[derive(Debug, Deserialize)]
87#[serde(deny_unknown_fields)]
88struct ExecTicket {
89    id: Option<String>,
90    name: String,
91    project: Option<String>,
92    #[serde(default)]
93    blocked_by: Vec<String>,
94    target: Option<String>,
95    model: Option<String>,
96    effort: Option<String>,
97    flow: Option<String>,
98    body: String,
99}
100
101fn parse_tickets(output: &[u8]) -> Result<Vec<AuthoredTicket>, SourceError> {
102    let tickets: Vec<ExecTicket> = serde_json::from_slice(output)
103        .map_err(|error| SourceError::new(format!("invalid ticket source output: {error}")))?;
104    Ok(tickets
105        .into_iter()
106        .enumerate()
107        .map(|(index, ticket)| {
108            let source_ref = ticket.id.clone().unwrap_or_else(|| format!("row:{index}"));
109            let mut frontmatter = Frontmatter::sourced(ticket.name, ticket.blocked_by);
110            frontmatter.id = ticket.id;
111            frontmatter.project = ticket.project;
112            frontmatter.target = ticket.target;
113            frontmatter.model = ticket.model;
114            frontmatter.effort = ticket.effort;
115            frontmatter.flow = ticket.flow;
116            AuthoredTicket {
117                frontmatter,
118                body: ticket.body,
119                source: "exec".into(),
120                source_ref,
121                file_path: None,
122                original_content: None,
123                validation_error: None,
124            }
125        })
126        .collect())
127}
128
129fn command_failed(argv: &[String], output: &std::process::Output) -> SourceError {
130    let stderr = String::from_utf8_lossy(&output.stderr);
131    let detail = stderr.trim();
132    let detail = if detail.is_empty() {
133        String::new()
134    } else {
135        format!(": {detail}")
136    };
137    SourceError::new(format!(
138        "ticket source `{}` exited with {}{detail}",
139        argv.join(" "),
140        output.status
141    ))
142}
143
144#[cfg(test)]
145mod tests {
146    use super::parse_tickets;
147
148    #[test]
149    fn valid_rows_map_fields_and_apply_defaults() {
150        let tickets = parse_tickets(
151            br#"[
152                {"id":"EXT-1","name":"One","project":"docs","blocked_by":["EXT-0"],"target":"codex","model":"o3","effort":"high","flow":"release","body":"First"},
153                {"name":"Two","body":"Second"}
154            ]"#,
155        )
156        .unwrap();
157
158        assert_eq!(tickets[0].source_ref, "EXT-1");
159        assert_eq!(tickets[0].frontmatter.blocked_by, ["EXT-0"]);
160        assert_eq!(tickets[0].frontmatter.flow.as_deref(), Some("release"));
161        assert_eq!(tickets[1].source_ref, "row:1");
162        assert!(tickets[1].frontmatter.blocked_by.is_empty());
163        assert!(tickets[1].frontmatter.has_blocked_by());
164        assert_eq!(tickets[1].body, "Second");
165    }
166
167    #[test]
168    fn malformed_json_is_rejected() {
169        let error = parse_tickets(br#"[{"name":"One","body":"work"}"#).unwrap_err();
170        assert!(error.to_string().contains("invalid ticket source output"));
171    }
172
173    #[test]
174    fn unknown_fields_are_rejected() {
175        let error =
176            parse_tickets(br#"[{"name":"One","body":"work","status":"ready"}]"#).unwrap_err();
177        assert!(
178            error.to_string().contains("unknown field `status`"),
179            "{error}"
180        );
181    }
182}