Skip to main content

supercode_harness/
mail_transcript.rs

1//! What a session's transcript says about its mail: the lines its own user
2//! typed, and when each message reached the model.
3//!
4//! The mailbox records a message as SENT: filed for its session and handed to
5//! the session's door. Only the harness's transcript records it as DELIVERED:
6//! in the model's context. The two differ whenever the session is busy, since
7//! a message handed to a harness mid-turn waits in its queue until the running
8//! tool call returns. This module reads the transcript for both halves of that.
9//!
10//! A line a person types into a session reaches it through the harness's own
11//! composer, not through supercode, so no envelope is filed for it. Claude Code
12//! records each one:
13//!
14//! - typed at an idle prompt: a `user` entry marked `promptSource: typed`;
15//! - typed while the session was busy: a `queue-operation` `enqueue` (sent),
16//!   then either a `dequeue` and a `user` entry marked `promptSource: queued`,
17//!   or a `remove` (`absorbed_mid_turn`) and a `queued_command` attachment from
18//!   a human (delivered).
19//!
20//! [`file_typed_lines`] files each as a `user` envelope in the session's
21//! mailbox, so it has an id like every other message. A line waiting in the
22//! queue is filed as sent; the queue also holds scheduled prompts and other
23//! sessions' messages, so a queued line that lands as someone else's is taken
24//! back out of the mailbox, and one removed from the queue unread is withdrawn. The id is derived from
25//! the session's address, the moment the line was sent and its words, so it is
26//! the same before and after the line reaches the model, and every reader
27//! derives the same one. A line supercode itself typed into the session's pane
28//! (a Room member's or a voice turn) is filed already under its own id.
29//!
30//! Mail supercode filed reaches the model as its rendering,
31//! `<cross-session-message id="m-…"`, inside a prompt or a queued-command
32//! attachment; a Room or voice turn reaches it as a typed line with its words.
33//!
34//! Codex rollouts do not record who wrote a user message (a typed line and the
35//! injected AGENTS.md or environment context look alike), so their typed lines
36//! are not filed; their mail is delivered once its rendering is in a user
37//! message or a tool's output. Other harnesses' records are not read here:
38//! their mail's delivery is unknown.
39
40use std::collections::{BTreeMap, HashMap, HashSet, VecDeque};
41use std::io::BufRead;
42use std::path::{Path, PathBuf};
43
44use serde::Serialize;
45
46use crate::mailbox::{Envelope, MailAddress, MailKind, Mailbox, ReplyVia, StoredEnvelope};
47use crate::HarnessHomes;
48
49/// Prefix of the id of a line a session's user typed.
50pub const TYPED_ID_PREFIX: &str = "u-";
51
52/// Prefix of the id of a channel line a harness received on its own channel
53/// (Claude Code's `[channel: …]` prompts, from a Room or another bridge).
54pub const CHANNEL_ID_PREFIX: &str = "c-";
55
56/// One line a session received on a channel of its harness's own: its header
57/// (`[channel: <room> · from: <who> · at: <time> …]`) and words.
58#[derive(Debug, Clone, PartialEq, Eq)]
59pub struct ChannelLine {
60    /// The line's header, which names its channel, sender and send time.
61    pub header: String,
62    /// Who sent it, as the header names them.
63    pub from: String,
64    /// When it was sent, in epoch milliseconds.
65    pub sent_at_ms: u64,
66    /// When it reached the model; `None` while it waits in the harness's queue.
67    pub delivered_at_ms: Option<u64>,
68    /// The line, header included.
69    pub text: String,
70}
71
72/// The id of the channel line with `header` in the session at `address`.
73pub fn channel_line_id(address: &MailAddress, header: &str) -> String {
74    let hash = blake3::hash(format!("{address}\n{header}").as_bytes()).to_hex();
75    format!("{CHANNEL_ID_PREFIX}{}", &hash[..24])
76}
77
78/// The channel lines in a prompt (several arrive as one when they waited
79/// together), each from its `[channel: …]` header to the next.
80fn channel_lines(text: &str, fallback_ms: u64, delivered_at_ms: Option<u64>) -> Vec<ChannelLine> {
81    let starts: Vec<usize> = text
82        .match_indices("[channel: ")
83        .map(|(at, _)| at)
84        .filter(|at| *at == 0 || text[..*at].ends_with('\n'))
85        .collect();
86    starts
87        .iter()
88        .enumerate()
89        .filter_map(|(n, start)| {
90            let end = starts.get(n + 1).copied().unwrap_or(text.len());
91            let line = text[*start..end].trim_end();
92            let header = &line[..line.find(']')? + 1];
93            let field = |name: &str| {
94                header
95                    .split(" · ")
96                    .find_map(|part| part.strip_prefix(name))
97                    .map(|value| value.trim_end_matches(']').trim().to_string())
98            };
99            Some(ChannelLine {
100                header: header.to_string(),
101                from: field("from: ").unwrap_or_else(|| "a channel".to_string()),
102                sent_at_ms: field("at: ")
103                    .and_then(|at| supercode_interchange::sidecar::rfc3339_to_ms(&at))
104                    .and_then(|at| u64::try_from(at).ok())
105                    .unwrap_or(fallback_ms),
106                delivered_at_ms,
107                text: line.to_string(),
108            })
109        })
110        .collect()
111}
112
113/// Prefix of the id of a session's answer to its user.
114pub const ANSWER_ID_PREFIX: &str = "a-";
115
116/// The id of the answer the transcript entry `native_id` records in the
117/// session at `address`: `a-` and 24 hex digits.
118pub fn answer_id(address: &MailAddress, native_id: &str) -> String {
119    let hash = blake3::hash(format!("{address}\n{native_id}").as_bytes()).to_hex();
120    format!("{ANSWER_ID_PREFIX}{}", &hash[..24])
121}
122
123/// The mailbox of the person who uses the sessions on `machine`: where the
124/// lines they type come from and where the sessions' answers go.
125pub fn user_address(machine: &str) -> std::io::Result<MailAddress> {
126    MailAddress::new(machine.to_string(), "operator", "user")
127        .map_err(|error| std::io::Error::other(error.0))
128}
129
130/// The id of the line with `text` that the user of the session at `address`
131/// sent at `sent_at_ms`: `u-` and 24 hex digits.
132pub fn typed_line_id(address: &MailAddress, sent_at_ms: u64, text: &str) -> String {
133    let hash = blake3::hash(format!("{address}\n{sent_at_ms}\n{text}").as_bytes()).to_hex();
134    format!("{TYPED_ID_PREFIX}{}", &hash[..24])
135}
136
137/// One line the session's user typed, as its transcript records it.
138#[derive(Debug, Clone, PartialEq, Eq)]
139pub struct TypedLine {
140    /// When the user sent it, in epoch milliseconds.
141    pub sent_at_ms: u64,
142    /// When it reached the model; `None` while it waits in the harness's queue.
143    pub delivered_at_ms: Option<u64>,
144    /// Taken back out of the queue before it reached the model.
145    pub withdrawn: bool,
146    /// The words, as typed.
147    pub text: String,
148}
149
150/// One answer the session gave its user: a turn's final message.
151#[derive(Debug, Clone, PartialEq, Eq)]
152pub struct Answer {
153    /// The transcript entry's own id.
154    pub native_id: String,
155    /// When it was written, in epoch milliseconds.
156    pub at_ms: u64,
157    /// Its words.
158    pub text: String,
159}
160
161/// What one transcript says about a session's mail.
162#[derive(Debug, Clone, Default)]
163pub struct TranscriptMail {
164    /// The lines the session's user typed, in the order they were sent.
165    pub typed: Vec<TypedLine>,
166    /// The session's answers to its user, in order.
167    pub answers: Vec<Answer>,
168    /// Lines the session received on its harness's own channels, in order.
169    pub channel: Vec<ChannelLine>,
170    /// Each filed message the model has seen, by id: when it first did.
171    pub delivered: HashMap<String, u64>,
172}
173
174/// Where a message stands.
175#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
176#[serde(tag = "state", rename_all = "snake_case")]
177pub enum Delivery {
178    /// In the model's context since this moment (epoch milliseconds).
179    Delivered {
180        /// When it first was.
181        at_ms: u64,
182    },
183    /// Sent to the session, not yet in its model's context.
184    Sent,
185    /// Taken back out of the harness's queue before it reached the model.
186    Withdrawn,
187    /// The session's transcript cannot be read here, so whether it reached
188    /// the model is not known.
189    Unknown,
190    /// Its receiver's transcript was not read for this listing (it reads only
191    /// the newest receivers'); `supercode message show <id>` reads it.
192    NotRead,
193}
194
195impl Delivery {
196    /// How an inbox listing says it.
197    pub fn describe(&self) -> String {
198        match self {
199            Self::Delivered { at_ms } => {
200                supercode_interchange::sidecar::ms_to_rfc3339(*at_ms as i64)
201            }
202            Self::Sent => "not yet".to_string(),
203            Self::Withdrawn => "never (taken back)".to_string(),
204            Self::Unknown => "unknown".to_string(),
205            Self::NotRead => "not read here (supercode message show reads it)".to_string(),
206        }
207    }
208}
209
210/// The transcript of the session at `address` on this machine, when its
211/// harness's records are read here.
212pub fn transcript(homes: &HarnessHomes, address: &MailAddress) -> Option<PathBuf> {
213    if address.harness != "claude-code" {
214        return None;
215    }
216    let file = format!("{}.jsonl", address.session_id);
217    std::fs::read_dir(&homes.claude_code)
218        .ok()?
219        .flatten()
220        .map(|project| project.path().join(&file))
221        .find(|path| path.is_file())
222}
223
224fn at_ms(value: &serde_json::Value) -> u64 {
225    value
226        .as_str()
227        .and_then(supercode_interchange::sidecar::rfc3339_to_ms)
228        .and_then(|at| u64::try_from(at).ok())
229        .unwrap_or_default()
230}
231
232fn prompt_text(content: &serde_json::Value) -> Option<String> {
233    let text = match content {
234        serde_json::Value::String(text) => text.clone(),
235        serde_json::Value::Array(parts) => parts
236            .iter()
237            .filter(|part| part["type"] == "text")
238            .filter_map(|part| part["text"].as_str())
239            .collect::<Vec<_>>()
240            .join("\n"),
241        _ => return None,
242    };
243    (!text.trim().is_empty()).then_some(text)
244}
245
246/// Mail ids rendered into `line` as the model reads them: the envelope's own
247/// `cross-session-message id="m-…"` (escaped or not, as Claude's relay nests
248/// it), or the relay's `(message m-…)` note.
249fn rendered_mail_ids(line: &str) -> Vec<&str> {
250    let mut ids = Vec::new();
251    for mark in ["cross-session-message id=\\\"", "(message "] {
252        for (at, _) in line.match_indices(mark) {
253            let rest = &line[at + mark.len()..];
254            let end = rest
255                .find(|character: char| !(character.is_ascii_alphanumeric() || character == '-'))
256                .unwrap_or(rest.len());
257            let id = &rest[..end];
258            if id.starts_with("m-") {
259                ids.push(id);
260            }
261        }
262    }
263    ids
264}
265
266/// Whether a line handed to a harness's queue is machine-made on its face: supercode's mail or
267/// notices, or a background task's report.
268fn machine_made(text: &str) -> bool {
269    text.contains("<cross-session-message")
270        || text.starts_with("[Cross-session")
271        || text.starts_with("[channel: ")
272        || text.starts_with("<task-notification>")
273}
274
275/// Read a Claude Code transcript for its typed lines and delivered mail.
276pub fn read_claude(path: &Path) -> std::io::Result<TranscriptMail> {
277    let reader = std::io::BufReader::new(std::fs::File::open(path)?);
278    let mut mail = TranscriptMail::default();
279    // The harness's queue, mirrored in full (it is first in, first out): each
280    // line's words and, when it may be the user's, its entry in `mail.typed`.
281    // A queue also holds scheduled prompts and other sessions' messages, so a
282    // line waiting there is the user's only provisionally, until it lands as
283    // theirs.
284    let mut queued: VecDeque<(String, Option<usize>)> = VecDeque::new();
285    // How many lines the queue has let go since the last prompt: they land
286    // together as the next one.
287    let mut dequeued = 0_usize;
288    // Lines taken into a running turn, waiting for the attachment that says
289    // who wrote them.
290    let mut absorbed: VecDeque<(String, Option<usize>)> = VecDeque::new();
291    let take = |waiting: &mut VecDeque<(String, Option<usize>)>, text: &str| {
292        let at = waiting.iter().position(|(words, _)| words == text)?;
293        waiting.remove(at).map(|(_, index)| index)
294    };
295    // Queued lines that landed as someone else's.
296    let mut not_typed: HashSet<usize> = HashSet::new();
297    // When each line left the queue for a running turn. Its queued-command
298    // attachment carries the time it was queued, not this one.
299    let mut taken_in: HashMap<String, VecDeque<u64>> = HashMap::new();
300    // Attachments written before their line's remove (older harness builds
301    // write them first): their words, the mail they render, the typed line
302    // they carry, and their own stamp should no remove come.
303    let mut early: Vec<(String, Vec<String>, Option<usize>, u64)> = Vec::new();
304    for line in reader.lines() {
305        let line = line?;
306        let prompt = line.contains("\"promptSource\"") || line.contains("\"type\":\"user\"");
307        let queue = line.contains("\"queue-operation\"");
308        let attachment = line.contains("\"queued_command\"");
309        let answer = line.contains("\"stop_reason\":\"end_turn\"")
310            || line.contains("\"stop_reason\":\"stop_sequence\"");
311        let rendered = line.contains("cross-session-message id=") || line.contains("(message m-");
312        if !(prompt || queue || attachment || rendered || answer) {
313            continue;
314        }
315        let Ok(record) = serde_json::from_str::<serde_json::Value>(&line) else {
316            continue;
317        };
318        if record["isSidechain"] == true {
319            continue;
320        }
321        let kind = record["type"].as_str().unwrap_or_default();
322        let mut at = at_ms(&record["timestamp"]);
323        if kind == "queue-operation" && record["operation"] == "remove" {
324            if let Some(text) = record["content"].as_str() {
325                // The remove an early attachment waited for: it entered the
326                // context now.
327                if let Some(found) = early.iter().position(|(words, ..)| words == text) {
328                    let (_, ids, index, _) = early.remove(found);
329                    for id in ids {
330                        mail.delivered.entry(id).or_insert(at);
331                    }
332                    if let Some(index) = index {
333                        mail.typed[index].delivered_at_ms = Some(at);
334                    }
335                    continue;
336                }
337                taken_in.entry(text.to_string()).or_default().push_back(at);
338            }
339        }
340        let mut waits_for_remove = false;
341        if kind == "attachment" && record["attachment"]["type"] == "queued_command" {
342            match record["attachment"]["prompt"]
343                .as_str()
344                .and_then(|text| taken_in.get_mut(text))
345                .and_then(VecDeque::pop_front)
346            {
347                Some(taken) => at = taken,
348                None => waits_for_remove = true,
349            }
350        }
351        if waits_for_remove {
352            let attached = &record["attachment"];
353            let text = attached["prompt"].as_str().unwrap_or_default().to_string();
354            let human = attached["origin"]["kind"] == "human" || attached["humanTurn"] == true;
355            let ids = rendered_mail_ids(&line)
356                .into_iter()
357                .map(str::to_string)
358                .collect();
359            // Its line leaves the queue with the remove still to come.
360            let index = match take(&mut queued, &text) {
361                Some(Some(index)) if !human => {
362                    not_typed.insert(index);
363                    None
364                }
365                Some(index) => index,
366                None if human && !text.trim().is_empty() => {
367                    mail.typed.push(TypedLine {
368                        sent_at_ms: at,
369                        delivered_at_ms: None,
370                        withdrawn: false,
371                        text: text.clone(),
372                    });
373                    Some(mail.typed.len() - 1)
374                }
375                None => None,
376            };
377            early.push((text, ids, index, at));
378            continue;
379        }
380        // Mail in the model's context: a prompt, an attachment or a tool's
381        // result rendering it; never an assistant's own words naming an id,
382        // nor the queue still holding it.
383        if rendered && matches!(kind, "user" | "attachment") {
384            for id in rendered_mail_ids(&line) {
385                mail.delivered.entry(id.to_string()).or_insert(at);
386            }
387        }
388        match kind {
389            // A turn's final message is its answer to the user.
390            "assistant"
391                if answer
392                    && matches!(
393                        record["message"]["stop_reason"].as_str(),
394                        Some("end_turn" | "stop_sequence")
395                    ) =>
396            {
397                let (Some(native_id), Some(text)) = (
398                    record["uuid"].as_str(),
399                    prompt_text(&record["message"]["content"]),
400                ) else {
401                    continue;
402                };
403                mail.answers.push(Answer {
404                    native_id: native_id.to_string(),
405                    at_ms: at,
406                    text,
407                });
408            }
409            // A dequeue names no line: the queue lets go of its oldest.
410            "queue-operation" if record["operation"] == "dequeue" => dequeued += 1,
411            "queue-operation" => {
412                let Some(text) = record["content"].as_str() else {
413                    continue;
414                };
415                match record["operation"].as_str() {
416                    Some("enqueue") => {
417                        let index = (!machine_made(text) && !text.trim().is_empty()).then(|| {
418                            mail.typed.push(TypedLine {
419                                sent_at_ms: at,
420                                delivered_at_ms: None,
421                                withdrawn: false,
422                                text: text.to_string(),
423                            });
424                            mail.typed.len() - 1
425                        });
426                        queued.push_back((text.to_string(), index));
427                    }
428                    Some("remove") => {
429                        let Some(index) = take(&mut queued, text) else {
430                            continue;
431                        };
432                        match record["reason"].as_str() {
433                            // Taken into the running turn: its attachment follows.
434                            Some("absorbed_mid_turn" | "delivered_to_agent") => {
435                                if let Some(index) = index {
436                                    mail.typed[index].delivered_at_ms = Some(at);
437                                }
438                                absorbed.push_back((text.to_string(), index));
439                            }
440                            // Taken back before it reached the model.
441                            _ => {
442                                if let Some(index) = index {
443                                    mail.typed[index].withdrawn = true;
444                                }
445                            }
446                        }
447                    }
448                    _ => {}
449                }
450            }
451            // A line taken into a running turn lands as a queued command; it
452            // is the user's when its origin is a person.
453            "attachment" => {
454                let attached = &record["attachment"];
455                let Some(text) = attached["prompt"].as_str() else {
456                    continue;
457                };
458                if attached["type"] != "queued_command" || text.trim().is_empty() {
459                    continue;
460                }
461                let human = attached["origin"]["kind"] == "human" || attached["humanTurn"] == true;
462                match take(&mut absorbed, text) {
463                    Some(Some(index)) if !human => {
464                        not_typed.insert(index);
465                    }
466                    Some(_) => {}
467                    None if human => mail.typed.push(TypedLine {
468                        sent_at_ms: at_ms(&attached["timestamp"]).min(at),
469                        delivered_at_ms: Some(at),
470                        withdrawn: false,
471                        text: text.to_string(),
472                    }),
473                    None => {}
474                }
475            }
476            "user" => {
477                let Some(text) = prompt_text(&record["message"]["content"]) else {
478                    continue;
479                };
480                let source = record["promptSource"].as_str();
481                let human = matches!(source, Some("typed" | "queued")) && record["isMeta"] != true;
482                if !human {
483                    mail.channel.extend(channel_lines(&text, at, Some(at)));
484                }
485                // Queued lines land as the next prompt when the turn they waited
486                // for ends, several at once when several were let go: the lines
487                // whose words it carries (the queue's oldest when it carries
488                // none, since the transcript does not record every way a queue
489                // empties, and matching by words keeps one lost line from
490                // shifting every later one).
491                let mut landing: Vec<Option<usize>> = Vec::new();
492                let let_go = std::mem::take(&mut dequeued);
493                if let_go == 0 {
494                    landing.extend(take(&mut queued, &text));
495                }
496                let wanted = let_go;
497                while landing.len() < wanted {
498                    let Some(at) = queued
499                        .iter()
500                        .position(|(words, _)| text.contains(words.as_str()))
501                    else {
502                        break;
503                    };
504                    landing.extend(queued.remove(at).map(|(_, index)| index));
505                }
506                if landing.is_empty() && wanted > 0 {
507                    for _ in 0..wanted {
508                        landing.extend(queued.pop_front().map(|(_, index)| index));
509                    }
510                }
511                if !landing.is_empty() {
512                    for index in landing.into_iter().flatten() {
513                        if human {
514                            mail.typed[index].delivered_at_ms = Some(at);
515                        } else {
516                            not_typed.insert(index);
517                        }
518                    }
519                    continue;
520                }
521                if human {
522                    mail.typed.push(TypedLine {
523                        sent_at_ms: at,
524                        delivered_at_ms: Some(at),
525                        withdrawn: false,
526                        text,
527                    });
528                }
529            }
530            _ => {}
531        }
532    }
533    // An early attachment whose remove never came is in the context all the
534    // same: it counts from its own stamp.
535    for (_, ids, index, at) in early {
536        for id in ids {
537            mail.delivered.entry(id).or_insert(at);
538        }
539        if let Some(index) = index {
540            mail.typed[index].delivered_at_ms.get_or_insert(at);
541        }
542    }
543    // A line still queued stays provisional: sent, not delivered.
544    for (text, _) in &queued {
545        if text.starts_with("[channel: ") {
546            mail.channel.extend(channel_lines(text, 0, None));
547        }
548    }
549    let mut index = 0;
550    mail.typed.retain(|_| {
551        index += 1;
552        !not_typed.contains(&(index - 1))
553    });
554    Ok(mail)
555}
556
557/// The transcript's reading of the session at `address`, when it can be read.
558pub fn read(homes: &HarnessHomes, address: &MailAddress) -> Option<TranscriptMail> {
559    if address.harness == "codex" {
560        return read_codex_answers(&crate::mail_question::transcript_for(address)?).ok();
561    }
562    read_claude(&transcript(homes, address)?).ok()
563}
564
565/// Final answers and delivered mail from a native Codex pane's own transcript.
566/// Question replies have separate native receipts; neither tool outputs nor
567/// developer context are invented as typed user turns here.
568fn read_codex_answers(path: &Path) -> std::io::Result<TranscriptMail> {
569    let reader = std::io::BufReader::new(std::fs::File::open(path)?);
570    let mut mail = TranscriptMail::default();
571    for line in reader.lines().map_while(Result::ok) {
572        let rendered = line.contains("cross-session-message id=");
573        if !rendered && !line.contains("\"final_answer\"") {
574            continue;
575        }
576        let Ok(record) = serde_json::from_str::<serde_json::Value>(&line) else {
577            continue;
578        };
579        let item = &record["payload"];
580        // Mail in the model's context, as Claude's reading counts it: a user
581        // message (a pasted turn) or a tool's output (`supercode message inbox`
582        // run in a turn) rendering it; never the model's own words naming an id,
583        // nor the event stream echoing those items.
584        if rendered
585            && record["type"] == "response_item"
586            && ((item["type"] == "message" && item["role"] == "user")
587                || matches!(
588                    item["type"].as_str(),
589                    Some("function_call_output" | "custom_tool_call_output")
590                ))
591        {
592            let at = at_ms(&record["timestamp"]);
593            for id in codex_rendered_mail_ids(&line) {
594                mail.delivered.entry(id.to_string()).or_insert(at);
595            }
596        }
597        if record["type"] != "response_item"
598            || item["type"] != "message"
599            || item["role"] != "assistant"
600            || item["phase"] != "final_answer"
601        {
602            continue;
603        }
604        let Some(id) = item["id"].as_str() else {
605            continue;
606        };
607        let text = item["content"]
608            .as_array()
609            .into_iter()
610            .flatten()
611            .filter(|part| part["type"] == "output_text")
612            .filter_map(|part| part["text"].as_str())
613            .collect::<Vec<_>>()
614            .join("\n");
615        if !text.is_empty() {
616            mail.answers.push(Answer {
617                native_id: id.to_string(),
618                at_ms: at_ms(&record["timestamp"]),
619                text,
620            });
621        }
622    }
623    Ok(mail)
624}
625
626/// Mail ids rendered into a Codex rollout line: `cross-session-message id="m-…"`
627/// with the quote escaped once (a user message) or more (a tool's output is
628/// JSON text inside the record's JSON).
629fn codex_rendered_mail_ids(line: &str) -> Vec<&str> {
630    let mark = "cross-session-message id=";
631    line.match_indices(mark)
632        .filter_map(|(at, _)| {
633            let rest = line[at + mark.len()..].trim_start_matches(['\\', '"']);
634            let end = rest
635                .find(|character: char| !(character.is_ascii_alphanumeric() || character == '-'))
636                .unwrap_or(rest.len());
637            Some(&rest[..end]).filter(|id| id.starts_with("m-"))
638        })
639        .collect()
640}
641
642/// The newest of `items` (time, id), sorted by time, sent before `at` (or at
643/// it, when `inclusive`).
644fn newest_before(items: &[(u64, String)], at: u64, inclusive: bool) -> Option<String> {
645    items
646        .iter()
647        .rev()
648        .find(|(when, _)| if inclusive { *when <= at } else { *when < at })
649        .map(|(_, id)| id.clone())
650}
651
652/// Each answer the session gave, as (when, id), oldest first.
653fn answer_times(address: &MailAddress, mail: &TranscriptMail) -> Vec<(u64, String)> {
654    let mut times: Vec<(u64, String)> = mail
655        .answers
656        .iter()
657        .map(|answer| (answer.at_ms, answer_id(address, &answer.native_id)))
658        .collect();
659    times.sort();
660    times
661}
662
663/// File `envelope` unless it is filed as it stands; a record an earlier
664/// reading filed differently (without its inferred reply) is filed again.
665fn file_derived(
666    mailbox: &Mailbox,
667    filed: &HashMap<&str, &StoredEnvelope>,
668    envelope: Envelope,
669) -> std::io::Result<bool> {
670    if let Some(stored) = filed.get(envelope.id.as_str()) {
671        if stored.envelope.in_reply_to == envelope.in_reply_to {
672            return Ok(false);
673        }
674        std::fs::remove_file(&stored.path).ok();
675    }
676    mailbox.file_read(&envelope)?;
677    Ok(true)
678}
679
680/// File every line the user of `mailbox`'s session typed that is not filed
681/// yet, reading the transcript as `mail`. Returns how many were filed now.
682pub fn file_typed_lines(mailbox: &Mailbox, mail: &TranscriptMail) -> std::io::Result<usize> {
683    let address = mailbox.address();
684    let filed = mailbox.list()?;
685    let by_id: HashMap<&str, &StoredEnvelope> = filed
686        .iter()
687        .map(|stored| (stored.envelope.id.as_str(), stored))
688        .collect();
689    // A line a person sends answers the session's newest answer before it:
690    // the reply is inferred, as a mail client threads a reply by its subject.
691    let answers = answer_times(address, mail);
692    // What supercode typed into the pane is filed already, under its own id:
693    // each such turn accounts for one typed line with its words.
694    let mut typed_by_supercode: HashMap<&str, usize> = HashMap::new();
695    for stored in &filed {
696        if stored.envelope.kind == MailKind::User
697            && !stored.envelope.id.starts_with(TYPED_ID_PREFIX)
698        {
699            *typed_by_supercode
700                .entry(stored.envelope.body.trim())
701                .or_default() += 1;
702        }
703    }
704    let user = user_address(&address.machine)?;
705    // A typed line's record is the transcript's reading: one the reading no
706    // longer yields (a queued line that landed as a scheduled prompt's or a
707    // session's, or one an older reading took for the user's) is taken back.
708    let current: HashSet<String> = mail
709        .typed
710        .iter()
711        .map(|line| typed_line_id(address, line.sent_at_ms, &line.text))
712        .collect();
713    for stored in &filed {
714        if stored.envelope.id.starts_with(TYPED_ID_PREFIX)
715            && !stored
716                .envelope
717                .in_reply_to
718                .as_deref()
719                .is_some_and(|id| id.starts_with("q-"))
720            && !current.contains(&stored.envelope.id)
721        {
722            std::fs::remove_file(&stored.path).ok();
723        }
724    }
725    let mut count = 0;
726    for line in &mail.typed {
727        let id = typed_line_id(address, line.sent_at_ms, &line.text);
728        if !by_id.contains_key(id.as_str()) {
729            if let Some(left) = typed_by_supercode.get_mut(line.text.trim()) {
730                if *left > 0 {
731                    *left -= 1;
732                    continue;
733                }
734            }
735        }
736        let in_reply_to = newest_before(&answers, line.sent_at_ms, false);
737        count += usize::from(file_derived(
738            mailbox,
739            &by_id,
740            Envelope {
741                id,
742                created_at_ms: line.sent_at_ms,
743                from: user.clone(),
744                from_name: format!("user@{}", address.machine),
745                kind: MailKind::User,
746                reply_via: ReplyVia::None,
747                in_reply_to_inferred: in_reply_to.is_some(),
748                in_reply_to,
749                thread: None,
750                native_from: None,
751                voice_for: None,
752                subject: None,
753                body: line.text.clone(),
754            },
755        )?);
756    }
757    for line in &mail.channel {
758        let id = channel_line_id(address, &line.header);
759        let in_reply_to = newest_before(&answers, line.sent_at_ms, false);
760        count += usize::from(file_derived(
761            mailbox,
762            &by_id,
763            Envelope {
764                id,
765                created_at_ms: line.sent_at_ms,
766                from: MailAddress::new(address.machine.clone(), "operator", "channel")
767                    .map_err(|error| std::io::Error::other(error.0))?,
768                from_name: line.from.clone(),
769                kind: MailKind::Channel,
770                reply_via: ReplyVia::None,
771                in_reply_to_inferred: in_reply_to.is_some(),
772                in_reply_to,
773                thread: None,
774                native_from: None,
775                voice_for: None,
776                subject: None,
777                body: line.text.clone(),
778            },
779        )?);
780    }
781    Ok(count)
782}
783
784/// File every answer the session at `address` gave its user that is not filed
785/// yet, in the user's mailbox under `root`, from the session (`name` is what
786/// the user knows it by). Returns how many were filed now.
787pub fn file_answers(
788    root: &Path,
789    address: &MailAddress,
790    name: &str,
791    mail: &TranscriptMail,
792) -> std::io::Result<usize> {
793    let recipient =
794        crate::mail_question::creator_for(address).unwrap_or(user_address(&address.machine)?);
795    let mailbox = Mailbox::open(root, &recipient)?;
796    let listed = mailbox.list()?;
797    let mut filed: HashMap<&str, &StoredEnvelope> = HashMap::new();
798    for stored in &listed {
799        // An answer an earlier reading filed as a session's mail is filed again as an answer.
800        if stored.envelope.id.starts_with(ANSWER_ID_PREFIX)
801            && stored.envelope.kind != MailKind::Answer
802        {
803            std::fs::remove_file(&stored.path).ok();
804            continue;
805        }
806        filed.insert(stored.envelope.id.as_str(), stored);
807    }
808    // An answer answers the newest line a person sent the session before it: a
809    // line typed or spoken to it, a Room member's turn, a channel line.
810    let mut people: Vec<(u64, String)> = mail
811        .typed
812        .iter()
813        .map(|line| {
814            (
815                line.sent_at_ms,
816                typed_line_id(address, line.sent_at_ms, &line.text),
817            )
818        })
819        .chain(
820            mail.channel
821                .iter()
822                .map(|line| (line.sent_at_ms, channel_line_id(address, &line.header))),
823        )
824        .chain(
825            Mailbox::open(root, address)?
826                .list()?
827                .into_iter()
828                .filter(|stored| {
829                    stored.envelope.kind == MailKind::User
830                        && (!stored.envelope.id.starts_with(TYPED_ID_PREFIX)
831                            || stored
832                                .envelope
833                                .in_reply_to
834                                .as_deref()
835                                .is_some_and(|id| id.starts_with("q-")))
836                })
837                .map(|stored| (stored.envelope.created_at_ms, stored.envelope.id)),
838        )
839        .collect();
840    people.sort();
841    let mut count = 0;
842    for answer in &mail.answers {
843        let id = answer_id(address, &answer.native_id);
844        let in_reply_to = newest_before(&people, answer.at_ms, true);
845        let envelope = Envelope {
846            id,
847            created_at_ms: answer.at_ms,
848            from: address.clone(),
849            from_name: name.to_string(),
850            kind: MailKind::Answer,
851            reply_via: ReplyVia::None,
852            in_reply_to_inferred: in_reply_to.is_some(),
853            in_reply_to,
854            thread: None,
855            native_from: None,
856            voice_for: None,
857            subject: None,
858            body: answer.text.clone(),
859        };
860        if recipient.machine != address.machine && !filed.contains_key(envelope.id.as_str()) {
861            // Keep the local outbound mirror only after the creator's machine
862            // accepts it. A failed remote delivery can be retried on the next read.
863            crate::mailbox::deliver_to(&recipient, &envelope)?;
864        }
865        count += usize::from(file_derived(&mailbox, &filed, envelope)?);
866    }
867    Ok(count)
868}
869
870/// Everything the session at `address` sent that is filed on this machine
871/// under `root`: its messages to other sessions and its answers to its user,
872/// each with the mailbox it is filed in, oldest first.
873pub fn sent_by(root: &Path, address: &MailAddress) -> Vec<(MailAddress, StoredEnvelope)> {
874    let mut sent: Vec<(MailAddress, StoredEnvelope)> = crate::mailbox::all_mailboxes(root)
875        .into_iter()
876        .filter(|mailbox| mailbox.address() != address)
877        .flat_map(|mailbox| {
878            let to = mailbox.address().clone();
879            mailbox
880                .list()
881                .unwrap_or_default()
882                .into_iter()
883                .filter(|stored| &stored.envelope.from == address)
884                .map(move |stored| (to.clone(), stored))
885        })
886        .collect();
887    sent.sort_by_key(|(_, stored)| stored.envelope.created_at_ms);
888    sent
889}
890
891/// Where each of `stored` (the mailbox of the session at `address`) stands,
892/// by id, as the transcript reading `mail` shows it (`None`: unreadable).
893pub fn deliveries(
894    address: &MailAddress,
895    stored: &[StoredEnvelope],
896    mail: Option<&TranscriptMail>,
897) -> BTreeMap<String, Delivery> {
898    let Some(mail) = mail else {
899        return stored
900            .iter()
901            .map(|stored| {
902                (
903                    stored.envelope.id.clone(),
904                    crate::mail_question::delivered_at(&stored.envelope)
905                        .map(|at_ms| Delivery::Delivered { at_ms })
906                        .unwrap_or(Delivery::Unknown),
907                )
908            })
909            .collect();
910    };
911    let typed: HashMap<String, &TypedLine> = mail
912        .typed
913        .iter()
914        .map(|line| (typed_line_id(address, line.sent_at_ms, &line.text), line))
915        .collect();
916    let channel: HashMap<String, &ChannelLine> = mail
917        .channel
918        .iter()
919        .map(|line| (channel_line_id(address, &line.header), line))
920        .collect();
921    // A turn supercode typed into the pane reaches the model as a typed line
922    // with its words, the first one sent after it was filed.
923    let mut used: HashSet<usize> = HashSet::new();
924    let mut result = BTreeMap::new();
925    for stored in stored {
926        let envelope = &stored.envelope;
927        if typed.get(&envelope.id).is_some_and(|line| line.withdrawn) {
928            result.insert(envelope.id.clone(), Delivery::Withdrawn);
929            continue;
930        }
931        let at = if let Some(line) = typed.get(&envelope.id) {
932            line.delivered_at_ms
933        } else if let Some(line) = channel.get(&envelope.id) {
934            line.delivered_at_ms
935        } else if let Some(at) = mail
936            .delivered
937            .get(&envelope.id)
938            .or_else(|| mail.delivered.get(crate::mailbox::short_id(&envelope.id)))
939        {
940            Some(*at)
941        } else if envelope.kind == MailKind::User
942            && envelope
943                .in_reply_to
944                .as_deref()
945                .is_some_and(|id| id.starts_with("q-"))
946        {
947            crate::mail_question::delivered_at(envelope)
948        } else if envelope.kind == MailKind::User {
949            mail.typed
950                .iter()
951                .enumerate()
952                .find(|(index, line)| {
953                    !used.contains(index)
954                        && line.sent_at_ms + 1_000 >= envelope.created_at_ms
955                        && line.text.trim() == envelope.body.trim()
956                })
957                .and_then(|(index, line)| {
958                    used.insert(index);
959                    line.delivered_at_ms
960                })
961        } else {
962            None
963        };
964        result.insert(
965            envelope.id.clone(),
966            match at {
967                Some(at_ms) => Delivery::Delivered { at_ms },
968                None => Delivery::Sent,
969            },
970        );
971    }
972    result
973}
974
975/// Fewest hex digits after the kind (`m-`, `u-`, `a-`) a prefix may give.
976pub const MIN_PREFIX_DIGITS: usize = 4;
977
978/// The messages filed on this machine under `root` whose id is `id` or
979/// starts with it, like git's short hashes: the mailbox holding each, and the
980/// envelope. A prefix shorter than the kind and [`MIN_PREFIX_DIGITS`] digits
981/// matches nothing; more than one match means the prefix is ambiguous.
982pub fn find_messages(root: &Path, id: &str) -> std::io::Result<Vec<(Mailbox, StoredEnvelope)>> {
983    Ok(find_messages_many(root, &[id])?.remove(0))
984}
985
986/// [`find_messages`] for several ids in one pass over this machine's mailboxes: one list per id,
987/// in the order given.
988pub fn find_messages_many(
989    root: &Path,
990    ids: &[&str],
991) -> std::io::Result<Vec<Vec<(Mailbox, StoredEnvelope)>>> {
992    let mut found: Vec<Vec<(Mailbox, StoredEnvelope)>> = ids.iter().map(|_| Vec::new()).collect();
993    // A prefix shorter than the kind and MIN_PREFIX_DIGITS digits matches nothing.
994    let usable: Vec<usize> = (0..ids.len())
995        .filter(|&index| {
996            ids[index]
997                .split_once('-')
998                .map_or(0, |(_, digits)| digits.len())
999                >= MIN_PREFIX_DIGITS
1000        })
1001        .collect();
1002    if usable.is_empty() {
1003        return Ok(found);
1004    }
1005    let prefixes: Vec<&str> = usable.iter().map(|&index| ids[index]).collect();
1006    for mailbox in crate::mailbox::all_mailboxes(root) {
1007        for (matched, stored) in mailbox.find_prefixes(&prefixes)? {
1008            let list = &mut found[usable[matched]];
1009            if !list
1010                .iter()
1011                .any(|(_, known)| known.envelope.id == stored.envelope.id)
1012            {
1013                list.push((mailbox.clone(), stored));
1014            }
1015        }
1016    }
1017    Ok(found)
1018}