use std::path::{Component, Path, PathBuf};
use std::sync::Arc;
use std::time::Duration;
use serde_json::{json, Value};
use tokio::io::AsyncReadExt;
use crate::daemon::cline_prompt_backfill::BACKFILL_KEY;
use crate::envelope::{new_event_id, AgentType, EventEnvelope, HookEventType};
pub const BACKFILL_SOURCE: &str = "claude-code-plan";
const POLL_EVERY: Duration = Duration::from_millis(250);
const MAX_POLLS: usize = 20;
const MAX_READ_BYTES: u64 = 1024 * 1024;
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct PlanPromptJob {
pub session_id: String,
pub transcript_path: PathBuf,
pub cwd: Option<String>,
}
pub fn job_for(envelope: &EventEnvelope) -> Option<PlanPromptJob> {
if envelope.source != AgentType::ClaudeCode || envelope.type_ != HookEventType::SessionStart {
return None;
}
let session_id = envelope.subject.as_deref().filter(|s| !s.is_empty())?;
let data = envelope.data.as_ref()?;
if data.get("source").and_then(Value::as_str) != Some("clear") {
return None;
}
let transcript_path = Path::new(data.get("transcript_path").and_then(Value::as_str)?);
if !is_session_transcript(transcript_path, session_id) {
return None;
}
Some(PlanPromptJob {
session_id: session_id.to_string(),
transcript_path: transcript_path.to_path_buf(),
cwd: data.get("cwd").and_then(Value::as_str).map(str::to_string),
})
}
fn is_session_transcript(path: &Path, session_id: &str) -> bool {
let plain_name = !session_id.is_empty()
&& session_id != "."
&& session_id != ".."
&& !session_id.contains(['/', '\\'])
&& session_id.chars().all(|c| !c.is_control());
if !plain_name || !path.is_absolute() {
return false;
}
let only_normal = path.components().all(|c| {
matches!(
c,
Component::RootDir | Component::Prefix(_) | Component::Normal(_)
)
});
let named_for_session = path
.file_name()
.is_some_and(|name| name == format!("{session_id}.jsonl").as_str());
let under_projects = path
.parent()
.and_then(Path::parent)
.and_then(Path::file_name)
.is_some_and(|name| name == "projects");
only_normal && named_for_session && under_projects
}
#[derive(Debug, PartialEq, Eq)]
enum FirstUserMessage {
Pending,
Plan(i64, String),
NotAPlan,
}
fn first_user_message(transcript: &str) -> FirstUserMessage {
for line in transcript.lines() {
let Ok(entry) = serde_json::from_str::<Value>(line) else {
continue;
};
if entry.get("type").and_then(Value::as_str) != Some("user") {
continue;
}
if entry.get("planContent").is_none() {
return FirstUserMessage::NotAPlan;
}
let ts_ms = entry
.get("timestamp")
.and_then(Value::as_str)
.and_then(|t| chrono::DateTime::parse_from_rfc3339(t).ok())
.map(|t| t.timestamp_millis());
let prompt = entry
.pointer("/message/content")
.map(text_of)
.filter(|p| !p.is_empty());
return match (ts_ms, prompt) {
(Some(ts_ms), Some(prompt)) => FirstUserMessage::Plan(ts_ms, prompt),
_ => FirstUserMessage::NotAPlan,
};
}
FirstUserMessage::Pending
}
pub fn plan_prompt(transcript: &str) -> Option<(i64, String)> {
match first_user_message(transcript) {
FirstUserMessage::Plan(ts_ms, prompt) => Some((ts_ms, prompt)),
_ => None,
}
}
fn text_of(content: &Value) -> String {
match content {
Value::String(text) => text.clone(),
Value::Array(parts) => parts
.iter()
.filter(|p| p.get("type").and_then(Value::as_str) == Some("text"))
.filter_map(|p| p.get("text").and_then(Value::as_str))
.collect::<Vec<_>>()
.join("\n"),
_ => String::new(),
}
}
pub fn prompt_envelope(job: &PlanPromptJob, ts_ms: i64, prompt: &str) -> Option<EventEnvelope> {
let time = chrono::DateTime::from_timestamp_millis(ts_ms)?;
Some(EventEnvelope {
specversion: "1.0".into(),
id: new_event_id(),
source: AgentType::ClaudeCode,
type_: HookEventType::UserPromptSubmit,
time,
datacontenttype: Some("application/json".into()),
subject: Some(job.session_id.clone()),
data: Some(json!({
"session_id": job.session_id,
"transcript_path": job.transcript_path.to_string_lossy(),
"cwd": job.cwd,
"hook_event_name": "UserPromptSubmit",
"prompt": prompt,
BACKFILL_KEY: BACKFILL_SOURCE,
})),
os: Some(std::env::consts::OS.into()),
arch: Some(std::env::consts::ARCH.into()),
localipv4: None,
localipv6: None,
publicipv4: None,
publicipv6: None,
clientversion: None,
agentversion: None,
agentid: None,
osuser: None,
gitemail: None,
provideracct: None,
wireformat: None,
hostid: None,
hostname: None,
})
}
async fn read_head(path: &Path) -> Option<String> {
let file = tokio::fs::File::open(path).await.ok()?;
let mut raw = Vec::new();
file.take(MAX_READ_BYTES).read_to_end(&mut raw).await.ok()?;
Some(String::from_utf8_lossy(&raw).into_owned())
}
pub fn spawn(state: Arc<crate::daemon::AppState>, job: PlanPromptJob) {
tokio::spawn(async move {
for _ in 0..MAX_POLLS {
let answer = match read_head(&job.transcript_path).await {
Some(text) => first_user_message(&text),
None => FirstUserMessage::Pending,
};
match answer {
FirstUserMessage::Pending => tokio::time::sleep(POLL_EVERY).await,
FirstUserMessage::NotAPlan => return,
FirstUserMessage::Plan(ts_ms, prompt) => {
if let Some(envelope) = prompt_envelope(&job, ts_ms, &prompt) {
crate::daemon::handlers::process_envelope_for_replay(
state,
envelope,
std::time::Instant::now(),
)
.await;
}
return;
}
}
}
tracing::debug!(
path = %job.transcript_path.display(),
"claude plan prompt: no first user message in time"
);
});
}
#[cfg(test)]
mod tests {
use super::*;
const SESSION: &str = "29d570e8-0000-4000-8000-000000000000";
fn transcript_path() -> String {
let home = if cfg!(windows) {
r"C:\Users\u"
} else {
"/home/u"
};
Path::new(home)
.join(".claude")
.join("projects")
.join("-work-repo")
.join(format!("{SESSION}.jsonl"))
.display()
.to_string()
}
fn envelope(source: AgentType, type_: HookEventType, data: Value) -> EventEnvelope {
EventEnvelope {
specversion: "1.0".into(),
id: new_event_id(),
source,
type_,
time: chrono::Utc::now(),
datacontenttype: None,
subject: Some(SESSION.into()),
data: Some(data),
os: None,
arch: None,
localipv4: None,
localipv6: None,
publicipv4: None,
publicipv6: None,
clientversion: None,
agentversion: None,
agentid: None,
osuser: None,
gitemail: None,
provideracct: None,
wireformat: None,
hostid: None,
hostname: None,
}
}
fn session_start(source: &str, path: &str) -> EventEnvelope {
envelope(
AgentType::ClaudeCode,
HookEventType::SessionStart,
json!({"session_id": SESSION, "source": source, "transcript_path": path,
"cwd": "/work/repo", "hook_event_name": "SessionStart"}),
)
}
fn plan_line(content: Value) -> String {
json!({"type": "user", "timestamp": "2026-09-17T11:55:59.119Z",
"planContent": "# Plan", "message": {"role": "user", "content": content}})
.to_string()
}
fn typed_line(text: &str) -> String {
json!({"type": "user", "timestamp": "2026-09-17T11:55:58.000Z",
"message": {"role": "user", "content": text}})
.to_string()
}
#[test]
fn a_cleared_claude_session_with_its_transcript_is_a_job() {
let job = job_for(&session_start("clear", &transcript_path())).expect("job");
assert_eq!(job.session_id, SESSION);
assert_eq!(job.transcript_path, PathBuf::from(transcript_path()));
assert_eq!(job.cwd.as_deref(), Some("/work/repo"));
}
#[test]
fn other_session_starts_are_not_jobs() {
for source in ["startup", "resume", "compact"] {
assert_eq!(
job_for(&session_start(source, &transcript_path())),
None,
"{source}"
);
}
let mut cline = session_start("clear", &transcript_path());
cline.source = AgentType::Cline;
assert_eq!(job_for(&cline), None);
let mut prompt = session_start("clear", &transcript_path());
prompt.type_ = HookEventType::UserPromptSubmit;
assert_eq!(job_for(&prompt), None);
let no_path = envelope(
AgentType::ClaudeCode,
HookEventType::SessionStart,
json!({"source": "clear"}),
);
assert_eq!(job_for(&no_path), None);
let mut no_subject = session_start("clear", &transcript_path());
no_subject.subject = None;
assert_eq!(job_for(&no_subject), None);
}
#[test]
fn hostile_transcript_paths_are_not_jobs() {
for path in [
format!("/home/u/.claude/projects/../../etc/{SESSION}.jsonl"),
format!("/home/u/.claude/projects/p/../{SESSION}.jsonl"),
"/home/u/.claude/projects/p/other.jsonl".to_string(),
format!("/home/u/.claude/elsewhere/p/{SESSION}.jsonl"),
format!("projects/p/{SESSION}.jsonl"),
] {
assert_eq!(job_for(&session_start("clear", &path)), None, "{path}");
}
let mut odd_subject = session_start("clear", "/x/projects/p/a/b.jsonl");
odd_subject.subject = Some("a/b".into());
assert_eq!(job_for(&odd_subject), None);
}
#[test]
fn the_plan_is_the_first_user_message() {
let transcript = [
json!({"type": "permission-mode"}).to_string(),
json!({"type": "attachment"}).to_string(),
plan_line(json!("Implement the following plan:\n\n# Plan")),
typed_line("later"),
]
.join("\n");
assert_eq!(
plan_prompt(&transcript),
Some((
1_789_646_159_119,
"Implement the following plan:\n\n# Plan".to_string()
))
);
}
#[test]
fn a_typed_first_prompt_is_not_a_plan() {
assert_eq!(plan_prompt(&typed_line("hello")), None);
let typed_then_plan = [typed_line("hello"), plan_line(json!("# Plan"))].join("\n");
assert_eq!(plan_prompt(&typed_then_plan), None);
assert_eq!(
first_user_message(&typed_line("hello")),
FirstUserMessage::NotAPlan
);
}
#[test]
fn text_parts_are_joined() {
let content = json!([{"type": "text", "text": "one"}, {"type": "image"},
{"type": "text", "text": "two"}]);
assert_eq!(
plan_prompt(&plan_line(content)).map(|(_, p)| p),
Some("one\ntwo".to_string())
);
}
#[test]
fn garbage_and_partial_lines_are_skipped() {
let transcript = ["not json", "{\"type\":\"us", &plan_line(json!("# Plan"))].join("\n");
assert_eq!(
plan_prompt(&transcript).map(|(_, p)| p),
Some("# Plan".into())
);
assert_eq!(
first_user_message("{\"type\":\"us"),
FirstUserMessage::Pending
);
assert_eq!(first_user_message(""), FirstUserMessage::Pending);
}
#[test]
fn the_envelope_has_claudes_prompt_shape() {
let job = job_for(&session_start("clear", &transcript_path())).expect("job");
let envelope = prompt_envelope(&job, 1_789_646_159_119, "# Plan").expect("envelope");
assert_eq!(envelope.source, AgentType::ClaudeCode);
assert_eq!(envelope.type_, HookEventType::UserPromptSubmit);
assert_eq!(envelope.subject.as_deref(), Some(SESSION));
assert_eq!(envelope.time.timestamp_millis(), 1_789_646_159_119);
let data = envelope.data.expect("data");
assert_eq!(data["prompt"], "# Plan");
assert_eq!(data["hook_event_name"], "UserPromptSubmit");
assert_eq!(data["session_id"], SESSION);
assert_eq!(data["cwd"], "/work/repo");
assert_eq!(data["transcript_path"], transcript_path());
assert_eq!(data[BACKFILL_KEY], BACKFILL_SOURCE);
}
}