use std::path::{Path, PathBuf};
use std::sync::Arc;
use dashmap::DashMap;
use serde_json::{json, Value};
use crate::envelope::{new_event_id, AgentType, EventEnvelope, HookEventType};
const RUN_START_SLACK_MS: i64 = 10_000;
const TYPED_PROMPT_MARKER: &str = "<user_input";
pub const BACKFILL_KEY: &str = "ai.openlatch.prompt.source";
pub const BACKFILL_SOURCE: &str = "cline-session-store";
const PRUNE_ABOVE: usize = 256;
const STALE_AFTER_MS: i64 = 60 * 60 * 1000;
#[derive(Clone, Copy, Debug, Default)]
struct SessionMarks {
run_started_ms: Option<i64>,
prompt_seen_ms: Option<i64>,
last_run_ended_ms: Option<i64>,
}
impl SessionMarks {
fn latest_ms(&self) -> i64 {
self.run_started_ms
.max(self.prompt_seen_ms)
.max(self.last_run_ended_ms)
.unwrap_or(i64::MIN)
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct BackfillJob {
pub subject: String,
pub root_session_id: String,
pub not_before_ms: i64,
}
#[derive(Default)]
pub struct PromptBackfill {
sessions: DashMap<String, SessionMarks>,
}
impl PromptBackfill {
pub fn observe(&self, envelope: &EventEnvelope, now_ms: i64) -> Option<BackfillJob> {
let subject = envelope.subject.clone().filter(|s| !s.is_empty())?;
if self.sessions.len() > PRUNE_ABOVE {
self.sessions
.retain(|_, marks| marks.latest_ms() >= now_ms - STALE_AFTER_MS);
}
if envelope.type_ == HookEventType::UserPromptSubmit {
if let Some(mut marks) = self.sessions.get_mut(&subject) {
marks.prompt_seen_ms = Some(now_ms);
} else {
self.sessions.insert(
subject,
SessionMarks {
prompt_seen_ms: Some(now_ms),
..SessionMarks::default()
},
);
}
return None;
}
let root_session_id = envelope
.data
.as_ref()?
.get("sessionContext")
.and_then(|c| c.get("rootSessionId"))
.and_then(Value::as_str)
.filter(|id| is_plain_folder_name(id))?
.to_string();
match envelope.type_ {
HookEventType::TaskStart | HookEventType::TaskResume => {
self.sessions.entry(subject).or_default().run_started_ms = Some(now_ms);
None
}
HookEventType::TaskComplete | HookEventType::TaskError | HookEventType::TaskCancel => {
let mut marks = self.sessions.get_mut(&subject)?;
let started = marks.run_started_ms.take()?;
let previous_end = marks.last_run_ended_ms.replace(now_ms);
let not_before_ms = match previous_end {
Some(ended) => (started - RUN_START_SLACK_MS).max(ended + 1),
None => started - RUN_START_SLACK_MS,
};
let covered = marks
.prompt_seen_ms
.is_some_and(|seen| seen >= not_before_ms);
drop(marks);
(!covered).then_some(BackfillJob {
subject,
root_session_id,
not_before_ms,
})
}
_ => None,
}
}
}
fn is_plain_folder_name(id: &str) -> bool {
!id.is_empty()
&& id != "."
&& id != ".."
&& !id.contains(['/', '\\'])
&& id.chars().all(|c| !c.is_control())
}
pub fn messages_path(data_root: &Path, root_session_id: &str) -> PathBuf {
data_root
.join("sessions")
.join(root_session_id)
.join(format!("{root_session_id}.messages.json"))
}
pub fn typed_prompts(messages_file: &Value, not_before_ms: i64) -> Vec<(i64, String)> {
let Some(messages) = messages_file.get("messages").and_then(Value::as_array) else {
return Vec::new();
};
let mut prompts: Vec<(i64, String)> = messages
.iter()
.filter(|m| m.get("role").and_then(Value::as_str) == Some("user"))
.filter(|m| m.pointer("/metadata/displayRole").and_then(Value::as_str) != Some("system"))
.filter_map(|m| {
let ts = m.get("ts").and_then(Value::as_i64)?;
let text = text_of(m.get("content")?);
(ts >= not_before_ms && text.contains(TYPED_PROMPT_MARKER)).then_some((ts, text))
})
.collect();
prompts.sort_by_key(|(ts, _)| *ts);
prompts
}
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(),
_ => String::new(),
}
}
pub fn prompt_envelope(job: &BackfillJob, 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::Cline,
type_: HookEventType::UserPromptSubmit,
time,
datacontenttype: Some("application/json".into()),
subject: Some(job.subject.clone()),
data: Some(json!({
"hookName": "prompt_submit",
"taskId": job.subject,
"sessionContext": { "rootSessionId": job.root_session_id },
"timestamp": time.to_rfc3339_opts(chrono::SecondsFormat::Millis, true),
"userPromptSubmit": { "prompt": prompt, "attachments": [] },
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,
})
}
pub fn spawn(state: Arc<crate::daemon::AppState>, job: BackfillJob) {
tokio::spawn(async move {
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
let Some((data_root, _)) = crate::hooks::cline::data_root() else {
return;
};
let path = messages_path(&data_root, &job.root_session_id);
let Ok(raw) = tokio::fs::read(&path).await else {
tracing::debug!(path = %path.display(), "cline prompt backfill: no messages file");
return;
};
let Ok(messages) = serde_json::from_slice::<Value>(&raw) else {
tracing::debug!(path = %path.display(), "cline prompt backfill: unreadable messages file");
return;
};
for (ts_ms, prompt) in typed_prompts(&messages, job.not_before_ms) {
if let Some(envelope) = prompt_envelope(&job, ts_ms, &prompt) {
crate::daemon::handlers::process_envelope_for_replay(
state.clone(),
envelope,
std::time::Instant::now(),
)
.await;
}
}
});
}
#[cfg(test)]
mod tests {
use super::*;
fn envelope(type_: HookEventType, subject: &str, data: Value) -> EventEnvelope {
EventEnvelope {
specversion: "1.0".into(),
id: new_event_id(),
source: AgentType::Cline,
type_,
time: chrono::Utc::now(),
datacontenttype: None,
subject: Some(subject.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 hook(type_: HookEventType) -> EventEnvelope {
envelope(
type_,
"conv_1",
json!({"taskId": "conv_1", "sessionContext": {"rootSessionId": "session_1"}}),
)
}
#[test]
fn a_run_without_a_prompt_event_is_backfilled() {
let backfill = PromptBackfill::default();
assert_eq!(
backfill.observe(&hook(HookEventType::TaskStart), 100_000),
None
);
assert_eq!(
backfill.observe(&hook(HookEventType::TaskComplete), 104_000),
Some(BackfillJob {
subject: "conv_1".into(),
root_session_id: "session_1".into(),
not_before_ms: 100_000 - RUN_START_SLACK_MS,
})
);
}
#[test]
fn a_run_with_a_prompt_event_is_left_alone() {
for prompt_first in [true, false] {
let backfill = PromptBackfill::default();
let prompt = envelope(
HookEventType::UserPromptSubmit,
"conv_1",
json!({"hookName": "prompt_submit", "taskId": "conv_1",
"userPromptSubmit": {"prompt": "hi"}}),
);
if prompt_first {
backfill.observe(&prompt, 99_990);
backfill.observe(&hook(HookEventType::TaskStart), 100_000);
} else {
backfill.observe(&hook(HookEventType::TaskStart), 100_000);
backfill.observe(&prompt, 100_010);
}
assert_eq!(
backfill.observe(&hook(HookEventType::TaskComplete), 104_000),
None,
"prompt_first={prompt_first}"
);
}
}
#[test]
fn a_follow_up_run_starts_after_the_previous_run_ended() {
let backfill = PromptBackfill::default();
backfill.observe(&hook(HookEventType::TaskStart), 100_000);
let first = backfill
.observe(&hook(HookEventType::TaskComplete), 101_000)
.expect("the first run is backfilled");
assert_eq!(first.not_before_ms, 100_000 - RUN_START_SLACK_MS);
backfill.observe(&hook(HookEventType::TaskStart), 106_000);
let second = backfill
.observe(&hook(HookEventType::TaskComplete), 108_000)
.expect("the follow-up run is backfilled");
assert_eq!(second.not_before_ms, 101_001);
}
#[test]
fn an_earlier_runs_prompt_does_not_cover_a_later_run() {
let backfill = PromptBackfill::default();
backfill.observe(&hook(HookEventType::TaskStart), 100_000);
backfill.observe(&hook(HookEventType::UserPromptSubmit), 100_010);
backfill.observe(&hook(HookEventType::TaskComplete), 104_000);
backfill.observe(&hook(HookEventType::TaskStart), 200_000);
assert!(backfill
.observe(&hook(HookEventType::TaskComplete), 204_000)
.is_some());
}
#[test]
fn stale_marks_are_pruned_and_live_ones_kept() {
let backfill = PromptBackfill::default();
for i in 0..=PRUNE_ABOVE {
let other = envelope(
HookEventType::UserPromptSubmit,
&format!("claude_{i}"),
json!({"prompt": "x"}),
);
backfill.observe(&other, 0);
}
let now = STALE_AFTER_MS + 10_000;
backfill.observe(&hook(HookEventType::TaskStart), now);
assert!(backfill.sessions.len() <= 2, "{}", backfill.sessions.len());
assert!(backfill
.observe(&hook(HookEventType::TaskComplete), now + 4_000)
.is_some());
}
#[test]
fn a_run_whose_start_was_not_seen_is_left_alone() {
let backfill = PromptBackfill::default();
assert_eq!(
backfill.observe(&hook(HookEventType::TaskComplete), 104_000),
None
);
}
#[test]
fn an_event_without_a_usable_root_session_is_ignored() {
for data in [
json!({"session_id": "conv_1"}),
json!({"sessionContext": {"rootSessionId": "../../etc"}}),
json!({"sessionContext": {"rootSessionId": "a/b"}}),
json!({"sessionContext": {"rootSessionId": ".."}}),
json!({"sessionContext": {"rootSessionId": ""}}),
] {
let backfill = PromptBackfill::default();
backfill.observe(
&envelope(HookEventType::TaskStart, "conv_1", data.clone()),
1,
);
assert_eq!(
backfill.observe(
&envelope(HookEventType::TaskComplete, "conv_1", data.clone()),
2
),
None,
"{data}"
);
}
}
#[test]
fn only_typed_prompts_since_the_run_start_are_taken() {
let file = json!({
"version": 1,
"sessionId": "session_1",
"messages": [
{"role": "user", "ts": 1_000, "content": [
{"type": "text", "text": "<user_input mode=\"act\">first</user_input>"}]},
{"role": "assistant", "ts": 1_500, "content": [{"type": "text", "text": "ok"}]},
{"role": "user", "ts": 2_000, "content": [
{"type": "text", "text": "<user_input mode=\"act\">sec"},
{"type": "image", "data": "..."},
{"type": "text", "text": "ond</user_input>"}]},
{"role": "user", "ts": 2_100, "content": [
{"type": "tool_result", "content": "a.txt"}]},
{"role": "user", "ts": 2_200, "metadata": {"displayRole": "system"}, "content": [
{"type": "text", "text": "<user_input>reminder</user_input>"}]},
{"role": "user", "ts": 2_300, "content": "[SYSTEM] finish with a tool"},
]
});
assert_eq!(
typed_prompts(&file, 1_800),
vec![(
2_000,
"<user_input mode=\"act\">second</user_input>".to_string()
)]
);
assert_eq!(typed_prompts(&file, 0).len(), 2);
assert!(typed_prompts(&json!({"unexpected": true}), 0).is_empty());
}
#[test]
fn the_prompt_envelope_is_clines_own_shape() {
let job = BackfillJob {
subject: "conv_1".into(),
root_session_id: "session_1".into(),
not_before_ms: 0,
};
let env = prompt_envelope(&job, 1_790_174_933_718, "<user_input>hi</user_input>")
.expect("a valid timestamp");
assert_eq!(env.type_, HookEventType::UserPromptSubmit);
assert_eq!(env.source, AgentType::Cline);
assert_eq!(env.subject.as_deref(), Some("conv_1"));
assert_eq!(env.time.timestamp_millis(), 1_790_174_933_718);
let data = env.data.expect("data");
assert_eq!(data["taskId"], "conv_1");
assert_eq!(data[BACKFILL_KEY], BACKFILL_SOURCE);
assert_eq!(
crate::core::envelope::normalize::prompt_of(&data),
Some("<user_input>hi</user_input>")
);
}
#[test]
fn the_messages_file_is_named_after_its_folder() {
assert_eq!(
messages_path(Path::new("/d"), "session_1"),
Path::new("/d/sessions/session_1/session_1.messages.json")
);
}
}