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}