Skip to main content

supercode_harness/
mail_question.rs

1//! Native questions as mail. Native request identities, never screen text or
2//! terminal keys, decide which question a reply answers.
3
4use std::io::{self, Write};
5use std::path::{Path, PathBuf};
6use std::process::Stdio;
7
8use serde::{Deserialize, Serialize};
9use serde_json::{json, Value};
10use tokio::io::AsyncWriteExt;
11
12use crate::mail_route::Caller;
13use crate::mailbox::{mail_root, Envelope, MailAddress, MailKind, Mailbox, ReplyVia};
14
15#[derive(Clone, Serialize, Deserialize)]
16struct Binding {
17    session: MailAddress,
18    creator: MailAddress,
19    pane: String,
20    #[serde(default)]
21    socket_path: Option<String>,
22    #[serde(default)]
23    transcript_path: Option<PathBuf>,
24}
25
26#[derive(Clone, Serialize, Deserialize)]
27struct Question {
28    binding: Binding,
29    native: Value,
30    envelope: Envelope,
31}
32
33#[derive(Serialize, Deserialize)]
34struct Reply {
35    envelope: Envelope,
36    answers: Value,
37}
38
39fn directory() -> PathBuf {
40    mail_root().join("native-questions")
41}
42
43fn binding_path(address: &MailAddress) -> PathBuf {
44    directory().join(format!(
45        "binding-{}.json",
46        blake3::hash(address.to_string().as_bytes())
47    ))
48}
49
50fn question_path(id: &str, suffix: &str) -> PathBuf {
51    directory().join(format!("{id}.{suffix}"))
52}
53
54fn read<T: serde::de::DeserializeOwned>(path: &Path) -> io::Result<T> {
55    serde_json::from_slice(&std::fs::read(path)?).map_err(io::Error::other)
56}
57
58fn create<T: Serialize>(path: &Path, value: &T) -> io::Result<()> {
59    std::fs::create_dir_all(directory())?;
60    #[cfg(unix)]
61    {
62        use std::os::unix::fs::PermissionsExt;
63        std::fs::set_permissions(directory(), std::fs::Permissions::from_mode(0o700))?;
64    }
65    let mut options = std::fs::OpenOptions::new();
66    options.write(true).create_new(true);
67    #[cfg(unix)]
68    {
69        use std::os::unix::fs::OpenOptionsExt;
70        options.mode(0o600);
71    }
72    let mut file = options.open(path)?;
73    file.write_all(&serde_json::to_vec(value).map_err(io::Error::other)?)?;
74    file.sync_all()
75}
76
77fn question_id(address: &MailAddress, native: &Value) -> String {
78    let identity = json!([
79        address,
80        native["threadId"],
81        native["turnId"],
82        native["itemId"]
83    ]);
84    format!(
85        "q-{}",
86        &blake3::hash(identity.to_string().as_bytes()).to_hex()[..24]
87    )
88}
89
90/// Creator recorded at this native session's first launch, when known.
91pub fn creator_for(address: &MailAddress) -> Option<MailAddress> {
92    read::<Binding>(&binding_path(address))
93        .ok()
94        .map(|binding| binding.creator)
95}
96
97/// Native transcript named by the session's own endpoint or hook.
98pub fn transcript_for(address: &MailAddress) -> Option<PathBuf> {
99    read::<Binding>(&binding_path(address))
100        .ok()?
101        .transcript_path
102}
103
104fn bind(seed: &Path, input: &Value) -> io::Result<Binding> {
105    let seed: Value = read(seed)?;
106    let harness = seed["kind"].as_str().unwrap_or_default();
107    if harness != "codex" && harness != "claude-code" {
108        return Err(io::Error::other("unsupported native question harness"));
109    }
110    let session = MailAddress::new(
111        crate::mailbox::local_machine_name(),
112        harness,
113        input["session_id"].as_str().unwrap_or_default(),
114    )
115    .map_err(io::Error::other)?;
116    let creator = match seed["by"].as_str() {
117        Some(address) => MailAddress::parse(address).map_err(io::Error::other)?,
118        None => crate::mail_transcript::user_address(&session.machine)?,
119    };
120    let binding = Binding {
121        session,
122        creator,
123        pane: seed["pane"].as_str().unwrap_or_default().to_string(),
124        socket_path: seed["socketPath"].as_str().map(str::to_string),
125        transcript_path: input["transcript_path"].as_str().map(PathBuf::from),
126    };
127    match create(&binding_path(&binding.session), &binding) {
128        Ok(()) => Ok(binding),
129        Err(error) if error.kind() == io::ErrorKind::AlreadyExists => {
130            let mut existing: Binding = read(&binding_path(&binding.session))?;
131            if existing.pane != binding.pane
132                || existing.socket_path != binding.socket_path
133                || binding.transcript_path.is_some()
134                    && existing.transcript_path != binding.transcript_path
135            {
136                existing.pane = binding.pane;
137                existing.socket_path = binding.socket_path;
138                if binding.transcript_path.is_some() {
139                    existing.transcript_path = binding.transcript_path;
140                }
141                let staged = binding_path(&existing.session)
142                    .with_extension(format!("{}.tmp", std::process::id()));
143                create(&staged, &existing)?;
144                std::fs::rename(staged, binding_path(&existing.session))?;
145            }
146            Ok(existing)
147        }
148        Err(error) => Err(error),
149    }
150}
151
152fn file_question(binding: &Binding, native: Value) -> io::Result<Question> {
153    let id = question_id(&binding.session, &native);
154    let mut body = String::new();
155    for question in native["questions"].as_array().into_iter().flatten() {
156        body.push_str(&format!(
157            "{}: {}\n",
158            question["id"].as_str().unwrap_or_default(),
159            question["question"].as_str().unwrap_or_default()
160        ));
161        for option in question["options"].as_array().into_iter().flatten() {
162            body.push_str(&format!(
163                "- {}\n",
164                option["label"].as_str().unwrap_or_default()
165            ));
166        }
167    }
168    body.push_str("\nReply with the answer text for one question, or a JSON object mapping each question id to its answer for several questions.");
169    let mut envelope = Envelope::new(
170        binding.session.clone(),
171        binding.session.to_string(),
172        MailKind::Question,
173        ReplyVia::Command,
174        body,
175    )?;
176    envelope.id = id.clone();
177    let question = Question {
178        binding: binding.clone(),
179        native,
180        envelope,
181    };
182    let question = match create(&question_path(&id, "json"), &question) {
183        Ok(()) => question,
184        Err(error) if error.kind() == io::ErrorKind::AlreadyExists => {
185            read(&question_path(&id, "json"))?
186        }
187        Err(error) => return Err(error),
188    };
189    Mailbox::open(&mail_root(), &binding.session)?.deliver_read(&question.envelope)?;
190    crate::mailbox::deliver_to(&binding.creator, &question.envelope)?;
191    // Keep an outbound mirror for remote recipients, so the session's mail
192    // timeline has the same question even when its creator is on another box.
193    if binding.creator.machine != binding.session.machine {
194        Mailbox::open(&mail_root(), &binding.creator)?.deliver_read(&question.envelope)?;
195    }
196    Ok(question)
197}
198
199async fn codex(request: Value) -> io::Result<Value> {
200    let entry = crate::teams_entry().map_err(io::Error::other)?;
201    let adapter = entry
202        .parent()
203        .and_then(Path::parent)
204        .ok_or_else(|| io::Error::other("invalid Teams entry"))?
205        .join("native-questions.mjs");
206    let mut child = tokio::process::Command::new(
207        std::env::var(crate::orchestrator_door::NODE_BIN_ENV).unwrap_or_else(|_| "node".into()),
208    )
209    .arg(adapter)
210    .stdin(Stdio::piped())
211    .stdout(Stdio::piped())
212    .stderr(Stdio::null())
213    .kill_on_drop(true)
214    .spawn()?;
215    child
216        .stdin
217        .take()
218        .ok_or_else(|| io::Error::other("native question stdin missing"))?
219        .write_all(request.to_string().as_bytes())
220        .await?;
221    let output = tokio::time::timeout(std::time::Duration::from_secs(25), child.wait_with_output())
222        .await
223        .map_err(io::Error::other)??;
224    serde_json::from_slice(&output.stdout).map_err(io::Error::other)
225}
226
227/// File native Codex questions for registered sessions. Each failed native
228/// connection is left unconfirmed; it is never replaced with terminal input.
229pub async fn poll() {
230    let Ok(entries) = std::fs::read_dir(mail_root().join("native-launches")) else {
231        return;
232    };
233    for entry in entries.flatten() {
234        let Ok(seed) = read::<Value>(&entry.path()) else {
235            continue;
236        };
237        if seed["kind"] != "codex" {
238            continue;
239        }
240        let Some(socket) = seed["socketPath"].as_str() else {
241            continue;
242        };
243        // Terminal teardown may kill the wrapper before its finally block.
244        // Retained socket files are not a live endpoint or a reason to poll it.
245        #[cfg(unix)]
246        if seed["ownerPid"]
247            .as_u64()
248            .is_none_or(|pid| unsafe { libc::kill(pid as i32, 0) } != 0)
249        {
250            continue;
251        }
252        // Only a socket belonging to this launched native pane is inspected.
253        // Never guess which thread a global shared daemon or a cwd belongs to.
254        if !Path::new(socket).exists() {
255            continue;
256        }
257        let Ok(reply) = codex(json!({"socketPath":socket,"discover":true})).await else {
258            continue;
259        };
260        for thread in reply["threads"].as_array().into_iter().flatten() {
261            if let Err(error) = bind(
262                &entry.path(),
263                &json!({"session_id":thread["id"], "transcript_path":thread["path"]}),
264            ) {
265                eprintln!("supercode: native session creator could not be recorded: {error}");
266            }
267        }
268        for native in reply["questions"].as_array().into_iter().flatten() {
269            match bind(&entry.path(), &json!({"session_id":native["threadId"]}))
270                .and_then(|binding| file_question(&binding, native.clone()))
271            {
272                Ok(_) => {}
273                Err(error) => {
274                    eprintln!("supercode: native question mail could not be filed: {error}")
275                }
276            }
277        }
278    }
279}
280
281fn answers(question: &Question, body: &str) -> io::Result<Value> {
282    let questions = question.native["questions"]
283        .as_array()
284        .ok_or_else(|| io::Error::other("question has no native inputs"))?;
285    let values: Value = if questions.len() == 1 && !body.trim_start().starts_with('{') {
286        json!({questions[0]["id"].as_str().unwrap_or_default(): body})
287    } else {
288        serde_json::from_str(body).map_err(io::Error::other)?
289    };
290    let mut output = serde_json::Map::new();
291    for question in questions {
292        let id = question["id"]
293            .as_str()
294            .ok_or_else(|| io::Error::other("native question id missing"))?;
295        let value = values[id]
296            .as_str()
297            .filter(|value| !value.trim().is_empty())
298            .ok_or_else(|| io::Error::other(format!("answer {id} with nonempty text")))?;
299        output.insert(id.to_string(), json!({"answers":[value]}));
300    }
301    if values.as_object().map(|object| object.len()) != Some(output.len()) {
302        return Err(io::Error::other("answer only the question ids shown"));
303    }
304    Ok(Value::Object(output))
305}
306
307/// Reply to a question through its native door. Only its recorded recipient
308/// can answer, and a durable claim prevents concurrent or automatic repeats.
309pub async fn reply(
310    caller: &Caller,
311    to: &MailAddress,
312    id: &str,
313    body: &str,
314) -> io::Result<crate::mail_send::Outcome> {
315    use crate::mail_send::{Outcome, EXIT_REFUSED, EXIT_STORED};
316    let question: Question = read(&question_path(id, "json"))?;
317    if &question.binding.session != to || question.binding.creator != caller.address {
318        return Ok(Outcome::new(
319            EXIT_REFUSED,
320            "Not answered: this question belongs to a different session or recipient.",
321        ));
322    }
323    if to.harness == "claude-code" {
324        let waiting: Value = match read(&question_path(id, "waiting")) {
325            Ok(waiting) => waiting,
326            Err(_) => {
327                return Ok(Outcome::new(
328                    EXIT_REFUSED,
329                    "Not answered: the native question hook is no longer waiting.",
330                ))
331            }
332        };
333        #[cfg(unix)]
334        if waiting["pid"]
335            .as_u64()
336            .is_none_or(|pid| unsafe { libc::kill(pid as i32, 0) } != 0)
337        {
338            return Ok(Outcome::new(
339                EXIT_REFUSED,
340                "Not answered: the native question hook has ended.",
341            ));
342        }
343    }
344    let answers = answers(&question, body)?;
345    let mut envelope = Envelope::new(
346        caller.address.clone(),
347        caller.name.clone(),
348        MailKind::User,
349        ReplyVia::None,
350        body,
351    )?;
352    envelope.id = format!(
353        "u-{}",
354        &blake3::hash(format!("{id}\n{}", caller.address).as_bytes()).to_hex()[..24]
355    );
356    envelope.in_reply_to = Some(id.to_string());
357    envelope.thread = Some(id.to_string());
358    let reply = Reply { envelope, answers };
359    match create(&question_path(id, "claim"), &reply) {
360        Ok(()) => {}
361        Err(error) if error.kind() == io::ErrorKind::AlreadyExists => {
362            return Ok(Outcome::new(EXIT_REFUSED, "An answer is already recorded for this question. Nothing was submitted again; inspect its thread for delivery."));
363        }
364        Err(error) => return Err(error),
365    }
366    let mailbox = Mailbox::open(&mail_root(), to)?;
367    mailbox.deliver(&reply.envelope)?;
368    // Publish to the hook only after its complete reply is durably in mail.
369    std::fs::hard_link(question_path(id, "claim"), question_path(id, "reply"))?;
370    if to.harness == "claude-code" {
371        return Ok(Outcome::new(0, format!("Answer {} recorded in reply to {id}; awaiting confirmation from the native question hook. It is not composer input.", crate::mailbox::short_id(&reply.envelope.id))));
372    }
373    let current: Binding = read(&binding_path(to))?;
374    if current.socket_path.is_none() {
375        return Ok(Outcome::new(EXIT_STORED, "Answer recorded but NOT delivered: this launch has no native question endpoint. Nothing was typed."));
376    }
377    let result = codex(json!({"socketPath":current.socket_path,"threadId": to.session_id, "answer":{"question":question.native,"answers":reply.answers}})).await?;
378    if result["delivered"] == true {
379        confirm(id, &reply)?;
380        Ok(Outcome::new(
381            0,
382            format!(
383                "Answer {} delivered to native question {id}.",
384                crate::mailbox::short_id(&reply.envelope.id)
385            ),
386        ))
387    } else {
388        Ok(Outcome::new(EXIT_STORED, format!("Answer recorded but NOT confirmed by the native question: {}. It will not be typed or automatically resent.", result["reason"])))
389    }
390}
391
392fn confirm(id: &str, reply: &Reply) -> io::Result<()> {
393    let question: Question = read(&question_path(id, "json"))?;
394    let mailbox = Mailbox::open(&mail_root(), &question.binding.session)?;
395    if let Some(stored) = mailbox.find(&reply.envelope.id)? {
396        if stored.state == crate::mailbox::MailState::Unread {
397            mailbox.mark_read(&stored)?;
398        }
399    }
400    let at = std::time::SystemTime::now()
401        .duration_since(std::time::UNIX_EPOCH)
402        .unwrap_or_default()
403        .as_millis() as u64;
404    match create(&question_path(id, "delivered"), &at) {
405        Ok(()) => Ok(()),
406        Err(error) if error.kind() == io::ErrorKind::AlreadyExists => Ok(()),
407        Err(error) => Err(error),
408    }
409}
410
411/// Receipt timestamp of a user's structured question reply, if confirmed.
412pub fn delivered_at(envelope: &Envelope) -> Option<u64> {
413    let id = envelope
414        .in_reply_to
415        .as_deref()
416        .filter(|id| id.starts_with("q-"))?;
417    let reply: Reply = read(&question_path(id, "reply")).ok()?;
418    if reply.envelope.id != envelope.id {
419        return None;
420    }
421    read(&question_path(id, "delivered")).ok()
422}
423
424/// Session-start binding and Claude's native AskUserQuestion hooks. Only the
425/// latter waits; its output is the documented updatedInput response contract.
426pub async fn hook(seed: &Path, input: Value) -> io::Result<Value> {
427    let binding = bind(seed, &input)?;
428    if input["hook_event_name"] == "SessionStart" {
429        // This launch's baseline precedes its native initial prompt, including
430        // on resume. Older identical prompts must never confirm a new launch.
431        if let Some(path) = &binding.transcript_path {
432            let receipt = json!({"path": path, "offset": std::fs::metadata(path).map(|m| m.len()).unwrap_or(0)});
433            match create(&seed.with_extension("session.json"), &receipt) {
434                Ok(()) => {}
435                Err(error) if error.kind() == io::ErrorKind::AlreadyExists => {}
436                Err(error) => return Err(error),
437            }
438        }
439    }
440    if input["tool_name"] != "AskUserQuestion" {
441        return Ok(json!({}));
442    }
443    let tool = &input["tool_input"];
444    let questions: Vec<Value> = tool["questions"]
445        .as_array()
446        .into_iter()
447        .flatten()
448        .map(|question| {
449            let mut question = question.clone();
450            question["id"] = question["question"].clone();
451            question
452        })
453        .collect();
454    let native = json!({"threadId":binding.session.session_id,"turnId":"claude","itemId":input["tool_use_id"],"questions":questions});
455    if !native["itemId"].is_string() {
456        return Err(io::Error::other("native question tool_use_id missing"));
457    }
458    let id = question_id(&binding.session, &native);
459    if input["hook_event_name"] == "PostToolUse" {
460        if let Ok(reply) = read::<Reply>(&question_path(&id, "reply")) {
461            let expected: serde_json::Map<String, Value> = reply
462                .answers
463                .as_object()
464                .into_iter()
465                .flatten()
466                .map(|(key, value)| (key.clone(), value["answers"][0].clone()))
467                .collect();
468            if input["tool_response"]["answers"] == Value::Object(expected) {
469                confirm(&id, &reply)?;
470            }
471        }
472        return Ok(json!({}));
473    }
474    if input["hook_event_name"] != "PreToolUse" {
475        return Ok(json!({}));
476    }
477    let waiting_path = question_path(&id, "waiting");
478    // One native hook owns this wait; its process identity expires it even if
479    // the native harness kills the hook before Rust can remove the marker.
480    create(&waiting_path, &json!({"pid": std::process::id()}))?;
481    struct Waiting(PathBuf);
482    impl Drop for Waiting {
483        fn drop(&mut self) {
484            let _ = std::fs::remove_file(&self.0);
485        }
486    }
487    let _waiting = Waiting(waiting_path);
488    file_question(&binding, native)?;
489    // The hook's native process owns this wait. A killed hook never submits
490    // a fallback answer; the native harness decides how to handle its exit.
491    loop {
492        if let Ok(reply) = read::<Reply>(&question_path(&id, "reply")) {
493            let answers: serde_json::Map<String, Value> = reply
494                .answers
495                .as_object()
496                .into_iter()
497                .flatten()
498                .map(|(key, value)| (key.clone(), value["answers"][0].clone()))
499                .collect();
500            let mut updated = tool.clone();
501            updated["answers"] = Value::Object(answers);
502            return Ok(
503                json!({"hookSpecificOutput":{"hookEventName":"PreToolUse","permissionDecision":"allow","updatedInput":updated}}),
504            );
505        }
506        tokio::time::sleep(std::time::Duration::from_millis(100)).await;
507    }
508}