use std::io::{self, Write};
use std::path::{Path, PathBuf};
use std::process::Stdio;
use serde::{Deserialize, Serialize};
use serde_json::{json, Value};
use tokio::io::AsyncWriteExt;
use crate::mail_route::Caller;
use crate::mailbox::{mail_root, Envelope, MailAddress, MailKind, Mailbox, ReplyVia};
#[derive(Clone, Serialize, Deserialize)]
struct Binding {
session: MailAddress,
creator: MailAddress,
pane: String,
#[serde(default)]
socket_path: Option<String>,
#[serde(default)]
transcript_path: Option<PathBuf>,
}
#[derive(Clone, Serialize, Deserialize)]
struct Question {
binding: Binding,
native: Value,
envelope: Envelope,
}
#[derive(Serialize, Deserialize)]
struct Reply {
envelope: Envelope,
answers: Value,
}
fn directory() -> PathBuf {
mail_root().join("native-questions")
}
fn binding_path(address: &MailAddress) -> PathBuf {
directory().join(format!(
"binding-{}.json",
blake3::hash(address.to_string().as_bytes())
))
}
fn question_path(id: &str, suffix: &str) -> PathBuf {
directory().join(format!("{id}.{suffix}"))
}
fn read<T: serde::de::DeserializeOwned>(path: &Path) -> io::Result<T> {
serde_json::from_slice(&std::fs::read(path)?).map_err(io::Error::other)
}
fn create<T: Serialize>(path: &Path, value: &T) -> io::Result<()> {
std::fs::create_dir_all(directory())?;
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(directory(), std::fs::Permissions::from_mode(0o700))?;
}
let mut options = std::fs::OpenOptions::new();
options.write(true).create_new(true);
#[cfg(unix)]
{
use std::os::unix::fs::OpenOptionsExt;
options.mode(0o600);
}
let mut file = options.open(path)?;
file.write_all(&serde_json::to_vec(value).map_err(io::Error::other)?)?;
file.sync_all()
}
fn question_id(address: &MailAddress, native: &Value) -> String {
let identity = json!([
address,
native["threadId"],
native["turnId"],
native["itemId"]
]);
format!(
"q-{}",
&blake3::hash(identity.to_string().as_bytes()).to_hex()[..24]
)
}
pub fn creator_for(address: &MailAddress) -> Option<MailAddress> {
read::<Binding>(&binding_path(address))
.ok()
.map(|binding| binding.creator)
}
pub fn transcript_for(address: &MailAddress) -> Option<PathBuf> {
read::<Binding>(&binding_path(address))
.ok()?
.transcript_path
}
fn bind(seed: &Path, input: &Value) -> io::Result<Binding> {
let seed: Value = read(seed)?;
let harness = seed["kind"].as_str().unwrap_or_default();
if harness != "codex" && harness != "claude-code" {
return Err(io::Error::other("unsupported native question harness"));
}
let session = MailAddress::new(
crate::mailbox::local_machine_name(),
harness,
input["session_id"].as_str().unwrap_or_default(),
)
.map_err(io::Error::other)?;
let creator = match seed["by"].as_str() {
Some(address) => MailAddress::parse(address).map_err(io::Error::other)?,
None => crate::mail_transcript::user_address(&session.machine)?,
};
let binding = Binding {
session,
creator,
pane: seed["pane"].as_str().unwrap_or_default().to_string(),
socket_path: seed["socketPath"].as_str().map(str::to_string),
transcript_path: input["transcript_path"].as_str().map(PathBuf::from),
};
match create(&binding_path(&binding.session), &binding) {
Ok(()) => Ok(binding),
Err(error) if error.kind() == io::ErrorKind::AlreadyExists => {
let mut existing: Binding = read(&binding_path(&binding.session))?;
if existing.pane != binding.pane
|| existing.socket_path != binding.socket_path
|| binding.transcript_path.is_some()
&& existing.transcript_path != binding.transcript_path
{
existing.pane = binding.pane;
existing.socket_path = binding.socket_path;
if binding.transcript_path.is_some() {
existing.transcript_path = binding.transcript_path;
}
let staged = binding_path(&existing.session)
.with_extension(format!("{}.tmp", std::process::id()));
create(&staged, &existing)?;
std::fs::rename(staged, binding_path(&existing.session))?;
}
Ok(existing)
}
Err(error) => Err(error),
}
}
fn file_question(binding: &Binding, native: Value) -> io::Result<Question> {
let id = question_id(&binding.session, &native);
let mut body = String::new();
for question in native["questions"].as_array().into_iter().flatten() {
body.push_str(&format!(
"{}: {}\n",
question["id"].as_str().unwrap_or_default(),
question["question"].as_str().unwrap_or_default()
));
for option in question["options"].as_array().into_iter().flatten() {
body.push_str(&format!(
"- {}\n",
option["label"].as_str().unwrap_or_default()
));
}
}
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.");
let mut envelope = Envelope::new(
binding.session.clone(),
binding.session.to_string(),
MailKind::Question,
ReplyVia::Command,
body,
)?;
envelope.id = id.clone();
let question = Question {
binding: binding.clone(),
native,
envelope,
};
let question = match create(&question_path(&id, "json"), &question) {
Ok(()) => question,
Err(error) if error.kind() == io::ErrorKind::AlreadyExists => {
read(&question_path(&id, "json"))?
}
Err(error) => return Err(error),
};
Mailbox::open(&mail_root(), &binding.session)?.deliver_read(&question.envelope)?;
crate::mailbox::deliver_to(&binding.creator, &question.envelope)?;
if binding.creator.machine != binding.session.machine {
Mailbox::open(&mail_root(), &binding.creator)?.deliver_read(&question.envelope)?;
}
Ok(question)
}
async fn codex(request: Value) -> io::Result<Value> {
let entry = crate::teams_entry().map_err(io::Error::other)?;
let adapter = entry
.parent()
.and_then(Path::parent)
.ok_or_else(|| io::Error::other("invalid Teams entry"))?
.join("native-questions.mjs");
let mut child = tokio::process::Command::new(
std::env::var(crate::orchestrator_door::NODE_BIN_ENV).unwrap_or_else(|_| "node".into()),
)
.arg(adapter)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::null())
.kill_on_drop(true)
.spawn()?;
child
.stdin
.take()
.ok_or_else(|| io::Error::other("native question stdin missing"))?
.write_all(request.to_string().as_bytes())
.await?;
let output = tokio::time::timeout(std::time::Duration::from_secs(25), child.wait_with_output())
.await
.map_err(io::Error::other)??;
serde_json::from_slice(&output.stdout).map_err(io::Error::other)
}
pub async fn poll() {
let Ok(entries) = std::fs::read_dir(mail_root().join("native-launches")) else {
return;
};
for entry in entries.flatten() {
let Ok(seed) = read::<Value>(&entry.path()) else {
continue;
};
if seed["kind"] != "codex" {
continue;
}
let Some(socket) = seed["socketPath"].as_str() else {
continue;
};
#[cfg(unix)]
if seed["ownerPid"]
.as_u64()
.is_none_or(|pid| unsafe { libc::kill(pid as i32, 0) } != 0)
{
continue;
}
if !Path::new(socket).exists() {
continue;
}
let Ok(reply) = codex(json!({"socketPath":socket,"discover":true})).await else {
continue;
};
for thread in reply["threads"].as_array().into_iter().flatten() {
if let Err(error) = bind(
&entry.path(),
&json!({"session_id":thread["id"], "transcript_path":thread["path"]}),
) {
eprintln!("supercode: native session creator could not be recorded: {error}");
}
}
for native in reply["questions"].as_array().into_iter().flatten() {
match bind(&entry.path(), &json!({"session_id":native["threadId"]}))
.and_then(|binding| file_question(&binding, native.clone()))
{
Ok(_) => {}
Err(error) => {
eprintln!("supercode: native question mail could not be filed: {error}")
}
}
}
}
}
fn answers(question: &Question, body: &str) -> io::Result<Value> {
let questions = question.native["questions"]
.as_array()
.ok_or_else(|| io::Error::other("question has no native inputs"))?;
let values: Value = if questions.len() == 1 && !body.trim_start().starts_with('{') {
json!({questions[0]["id"].as_str().unwrap_or_default(): body})
} else {
serde_json::from_str(body).map_err(io::Error::other)?
};
let mut output = serde_json::Map::new();
for question in questions {
let id = question["id"]
.as_str()
.ok_or_else(|| io::Error::other("native question id missing"))?;
let value = values[id]
.as_str()
.filter(|value| !value.trim().is_empty())
.ok_or_else(|| io::Error::other(format!("answer {id} with nonempty text")))?;
output.insert(id.to_string(), json!({"answers":[value]}));
}
if values.as_object().map(|object| object.len()) != Some(output.len()) {
return Err(io::Error::other("answer only the question ids shown"));
}
Ok(Value::Object(output))
}
pub async fn reply(
caller: &Caller,
to: &MailAddress,
id: &str,
body: &str,
) -> io::Result<crate::mail_send::Outcome> {
use crate::mail_send::{Outcome, EXIT_REFUSED, EXIT_STORED};
let question: Question = read(&question_path(id, "json"))?;
if &question.binding.session != to || question.binding.creator != caller.address {
return Ok(Outcome::new(
EXIT_REFUSED,
"Not answered: this question belongs to a different session or recipient.",
));
}
if to.harness == "claude-code" {
let waiting: Value = match read(&question_path(id, "waiting")) {
Ok(waiting) => waiting,
Err(_) => {
return Ok(Outcome::new(
EXIT_REFUSED,
"Not answered: the native question hook is no longer waiting.",
))
}
};
#[cfg(unix)]
if waiting["pid"]
.as_u64()
.is_none_or(|pid| unsafe { libc::kill(pid as i32, 0) } != 0)
{
return Ok(Outcome::new(
EXIT_REFUSED,
"Not answered: the native question hook has ended.",
));
}
}
let answers = answers(&question, body)?;
let mut envelope = Envelope::new(
caller.address.clone(),
caller.name.clone(),
MailKind::User,
ReplyVia::None,
body,
)?;
envelope.id = format!(
"u-{}",
&blake3::hash(format!("{id}\n{}", caller.address).as_bytes()).to_hex()[..24]
);
envelope.in_reply_to = Some(id.to_string());
envelope.thread = Some(id.to_string());
let reply = Reply { envelope, answers };
match create(&question_path(id, "claim"), &reply) {
Ok(()) => {}
Err(error) if error.kind() == io::ErrorKind::AlreadyExists => {
return Ok(Outcome::new(EXIT_REFUSED, "An answer is already recorded for this question. Nothing was submitted again; inspect its thread for delivery."));
}
Err(error) => return Err(error),
}
let mailbox = Mailbox::open(&mail_root(), to)?;
mailbox.deliver(&reply.envelope)?;
std::fs::hard_link(question_path(id, "claim"), question_path(id, "reply"))?;
if to.harness == "claude-code" {
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))));
}
let current: Binding = read(&binding_path(to))?;
if current.socket_path.is_none() {
return Ok(Outcome::new(EXIT_STORED, "Answer recorded but NOT delivered: this launch has no native question endpoint. Nothing was typed."));
}
let result = codex(json!({"socketPath":current.socket_path,"threadId": to.session_id, "answer":{"question":question.native,"answers":reply.answers}})).await?;
if result["delivered"] == true {
confirm(id, &reply)?;
Ok(Outcome::new(
0,
format!(
"Answer {} delivered to native question {id}.",
crate::mailbox::short_id(&reply.envelope.id)
),
))
} else {
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"])))
}
}
fn confirm(id: &str, reply: &Reply) -> io::Result<()> {
let question: Question = read(&question_path(id, "json"))?;
let mailbox = Mailbox::open(&mail_root(), &question.binding.session)?;
if let Some(stored) = mailbox.find(&reply.envelope.id)? {
if stored.state == crate::mailbox::MailState::Unread {
mailbox.mark_read(&stored)?;
}
}
let at = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_millis() as u64;
match create(&question_path(id, "delivered"), &at) {
Ok(()) => Ok(()),
Err(error) if error.kind() == io::ErrorKind::AlreadyExists => Ok(()),
Err(error) => Err(error),
}
}
pub fn delivered_at(envelope: &Envelope) -> Option<u64> {
let id = envelope
.in_reply_to
.as_deref()
.filter(|id| id.starts_with("q-"))?;
let reply: Reply = read(&question_path(id, "reply")).ok()?;
if reply.envelope.id != envelope.id {
return None;
}
read(&question_path(id, "delivered")).ok()
}
pub async fn hook(seed: &Path, input: Value) -> io::Result<Value> {
let binding = bind(seed, &input)?;
if input["hook_event_name"] == "SessionStart" {
if let Some(path) = &binding.transcript_path {
let receipt = json!({"path": path, "offset": std::fs::metadata(path).map(|m| m.len()).unwrap_or(0)});
match create(&seed.with_extension("session.json"), &receipt) {
Ok(()) => {}
Err(error) if error.kind() == io::ErrorKind::AlreadyExists => {}
Err(error) => return Err(error),
}
}
}
if input["tool_name"] != "AskUserQuestion" {
return Ok(json!({}));
}
let tool = &input["tool_input"];
let questions: Vec<Value> = tool["questions"]
.as_array()
.into_iter()
.flatten()
.map(|question| {
let mut question = question.clone();
question["id"] = question["question"].clone();
question
})
.collect();
let native = json!({"threadId":binding.session.session_id,"turnId":"claude","itemId":input["tool_use_id"],"questions":questions});
if !native["itemId"].is_string() {
return Err(io::Error::other("native question tool_use_id missing"));
}
let id = question_id(&binding.session, &native);
if input["hook_event_name"] == "PostToolUse" {
if let Ok(reply) = read::<Reply>(&question_path(&id, "reply")) {
let expected: serde_json::Map<String, Value> = reply
.answers
.as_object()
.into_iter()
.flatten()
.map(|(key, value)| (key.clone(), value["answers"][0].clone()))
.collect();
if input["tool_response"]["answers"] == Value::Object(expected) {
confirm(&id, &reply)?;
}
}
return Ok(json!({}));
}
if input["hook_event_name"] != "PreToolUse" {
return Ok(json!({}));
}
let waiting_path = question_path(&id, "waiting");
create(&waiting_path, &json!({"pid": std::process::id()}))?;
struct Waiting(PathBuf);
impl Drop for Waiting {
fn drop(&mut self) {
let _ = std::fs::remove_file(&self.0);
}
}
let _waiting = Waiting(waiting_path);
file_question(&binding, native)?;
loop {
if let Ok(reply) = read::<Reply>(&question_path(&id, "reply")) {
let answers: serde_json::Map<String, Value> = reply
.answers
.as_object()
.into_iter()
.flatten()
.map(|(key, value)| (key.clone(), value["answers"][0].clone()))
.collect();
let mut updated = tool.clone();
updated["answers"] = Value::Object(answers);
return Ok(
json!({"hookSpecificOutput":{"hookEventName":"PreToolUse","permissionDecision":"allow","updatedInput":updated}}),
);
}
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
}
}