1use std::path::PathBuf;
12use std::sync::Arc;
13
14use roder_api::events::{
15 EventEnvelope, RoderEvent, ThreadCreated, ThreadForkFailed, ThreadForkRemoved,
16 ThreadForkRequested, ThreadForked, ThreadId, TurnId,
17};
18use roder_api::forks::{
19 ForkPolicy, ForkReason, ForkRequest, ForkStatus, RemoveForkPolicy, WorkspaceFork,
20};
21use roder_api::thread::{ThreadMetadata, ThreadStore};
22use time::OffsetDateTime;
23
24use crate::Runtime;
25use crate::forks::DEFAULT_FORK_PROVIDER;
26
27mod history;
28use history::seed_events_for_child;
29
30#[derive(Debug, Clone)]
31pub struct ForkThreadRequest {
32 pub parent_thread_id: ThreadId,
33 pub name: String,
36 pub from_turn_id: Option<TurnId>,
38 pub provider_id: Option<String>,
40 pub provider_config: serde_json::Value,
42}
43
44impl ForkThreadRequest {
45 pub fn new(parent_thread_id: ThreadId, name: impl Into<String>) -> Self {
46 Self {
47 parent_thread_id,
48 name: name.into(),
49 from_turn_id: None,
50 provider_id: None,
51 provider_config: serde_json::json!({}),
52 }
53 }
54}
55
56#[derive(Debug, Clone)]
57pub struct ForkThreadOutcome {
58 pub child: ThreadMetadata,
59 pub warnings: Vec<String>,
60}
61
62impl Runtime {
63 pub(crate) async fn seed_agent_thread_history(
68 &self,
69 parent_thread_id: &ThreadId,
70 child_thread_id: &ThreadId,
71 fork_turns: &str,
72 ) -> anyhow::Result<()> {
73 if fork_turns == "none" {
74 return Ok(());
75 }
76 let Some(store) = self.thread_store.clone() else {
77 return Ok(());
78 };
79 let Some(parent) = store.load_thread(parent_thread_id).await? else {
80 return Ok(());
81 };
82 let mut events = seed_events_for_child(&parent.events, None)?;
83 if fork_turns != "all" {
84 let turn_count = fork_turns.parse::<usize>().map_err(|_| {
85 anyhow::anyhow!("fork_turns must be one of none, all, or a positive integer")
86 })?;
87 anyhow::ensure!(turn_count > 0, "fork_turns integer must be positive");
88 let mut ordered_turns = Vec::<TurnId>::new();
89 for event in &events {
90 if let Some(turn_id) = event.turn_id.as_ref()
91 && ordered_turns.last() != Some(turn_id)
92 {
93 ordered_turns.push(turn_id.clone());
94 }
95 }
96 let keep_from = ordered_turns.len().saturating_sub(turn_count);
97 let kept = &ordered_turns[keep_from..];
98 events.retain(|event| {
99 event
100 .turn_id
101 .as_ref()
102 .is_some_and(|turn_id| kept.contains(turn_id))
103 });
104 }
105 for event in &events {
106 store.append_event(child_thread_id, event).await?;
107 }
108 Ok(())
109 }
110
111 pub async fn fork_thread(
114 &self,
115 request: ForkThreadRequest,
116 ) -> anyhow::Result<ForkThreadOutcome> {
117 self.emit(RoderEvent::ThreadForkRequested(ThreadForkRequested {
118 parent_thread_id: request.parent_thread_id.clone(),
119 name: request.name.clone(),
120 timestamp: OffsetDateTime::now_utc(),
121 }))
122 .await;
123 match self.fork_thread_inner(&request).await {
124 Ok(outcome) => Ok(outcome),
125 Err(error) => {
126 self.emit(RoderEvent::ThreadForkFailed(ThreadForkFailed {
127 parent_thread_id: request.parent_thread_id.clone(),
128 name: request.name.clone(),
129 message: error.to_string(),
130 timestamp: OffsetDateTime::now_utc(),
131 }))
132 .await;
133 Err(error)
134 }
135 }
136 }
137
138 async fn fork_thread_inner(
139 &self,
140 request: &ForkThreadRequest,
141 ) -> anyhow::Result<ForkThreadOutcome> {
142 let store = self
143 .thread_store
144 .clone()
145 .ok_or_else(|| anyhow::anyhow!("conversation forks require a thread store"))?;
146 let parent = store
147 .load_thread(&request.parent_thread_id)
148 .await?
149 .ok_or_else(|| {
150 anyhow::anyhow!("parent thread {} was not found", request.parent_thread_id)
151 })?;
152 let parent_metadata = parent.metadata.clone().ok_or_else(|| {
153 anyhow::anyhow!(
154 "parent thread {} has no metadata to fork from",
155 request.parent_thread_id
156 )
157 })?;
158
159 let provider_id = request
162 .provider_id
163 .clone()
164 .unwrap_or_else(|| DEFAULT_FORK_PROVIDER.to_string());
165 let fork = self
166 .create_workspace_fork(
167 &provider_id,
168 ForkRequest {
169 source_workspace: PathBuf::from(&parent_metadata.workspace),
170 name: Some(request.name.clone()),
171 reason: ForkReason::ConversationFork,
172 policy: ForkPolicy::default(),
173 provider_config: request.provider_config.clone(),
174 },
175 )
176 .await?;
177
178 let now = OffsetDateTime::now_utc();
179 let seed_events = seed_events_for_child(&parent.events, request.from_turn_id.as_deref())?;
180 let mut warnings = Vec::new();
181 if request.from_turn_id.is_none() && seed_events.is_empty() && !parent.events.is_empty() {
182 warnings.push(
183 "parent thread has events but none were conversation records; the fork starts \
184 with an empty transcript"
185 .to_string(),
186 );
187 }
188
189 let child_id = uuid::Uuid::new_v4().to_string();
190 let child_metadata = ThreadMetadata {
191 thread_id: child_id.clone(),
192 title: Some(match &parent_metadata.title {
193 Some(title) => format!("{title} (fork: {})", request.name),
194 None => format!("fork: {}", request.name),
195 }),
196 workspace: fork.workspace.display().to_string(),
197 workspace_id: None,
199 root_id: None,
200 provider: parent_metadata.provider.clone(),
201 model: parent_metadata.model.clone(),
202 selection_mode: parent_metadata.selection_mode.clone(),
203 tool_allowlist: parent_metadata.tool_allowlist.clone(),
204 developer_instructions: parent_metadata.developer_instructions.clone(),
205 external_tools: parent_metadata.external_tools.clone(),
206 runner_destination: None,
208 runner_state: None,
209 runner_binding: None,
210 created_at: now,
211 updated_at: now,
212 message_count: 0,
213 usage: None,
214 parent_thread_id: Some(request.parent_thread_id.clone()),
215 forked_from_turn_id: request.from_turn_id.clone(),
216 workspace_fork: Some(fork.clone()),
217 };
218
219 let seed = async {
220 self.seed_child_thread(&store, child_metadata.clone(), &child_id, seed_events)
221 .await?;
222 self.goals
223 .inherit_thread_goal_snapshot(&request.parent_thread_id, &child_id)
224 .await
225 }
226 .await;
227 let inherited_goal = match seed {
228 Ok(goal) => goal,
229 Err(error) => {
230 let _ = self
232 .remove_workspace_fork(
233 &provider_id,
234 &fork.id,
235 RemoveForkPolicy {
236 confirm_workspace: fork.workspace.clone(),
237 },
238 )
239 .await;
240 return Err(error);
241 }
242 };
243
244 self.emit(RoderEvent::ThreadCreated(ThreadCreated {
245 thread_id: child_id.clone(),
246 timestamp: OffsetDateTime::now_utc(),
247 }))
248 .await;
249 if let Some(goal) = inherited_goal {
250 self.goals.emit_goal_updated(goal).await;
251 }
252 self.emit(RoderEvent::ThreadForked(ThreadForked {
253 parent_thread_id: request.parent_thread_id.clone(),
254 child_thread_id: child_id.clone(),
255 fork,
256 timestamp: OffsetDateTime::now_utc(),
257 }))
258 .await;
259
260 let child = store
261 .load_thread_metadata(&child_id)
262 .await?
263 .unwrap_or(child_metadata);
264 Ok(ForkThreadOutcome { child, warnings })
265 }
266
267 async fn seed_child_thread(
268 &self,
269 store: &Arc<dyn ThreadStore>,
270 child_metadata: ThreadMetadata,
271 child_id: &ThreadId,
272 seed_events: Vec<EventEnvelope>,
273 ) -> anyhow::Result<()> {
274 store.create_thread(child_metadata).await?;
275 for envelope in &seed_events {
276 store.append_event(child_id, envelope).await?;
277 }
278 Ok(())
279 }
280
281 pub async fn remove_thread_workspace_fork(
288 &self,
289 thread_id: &ThreadId,
290 confirm_path: &str,
291 ) -> anyhow::Result<WorkspaceFork> {
292 let store = self
293 .thread_store
294 .clone()
295 .ok_or_else(|| anyhow::anyhow!("conversation forks require a thread store"))?;
296 let mut metadata = store
297 .load_thread_metadata(thread_id)
298 .await?
299 .ok_or_else(|| anyhow::anyhow!("thread {thread_id} was not found"))?;
300 let mut fork = metadata
301 .workspace_fork
302 .clone()
303 .ok_or_else(|| anyhow::anyhow!("thread {thread_id} is not a workspace fork"))?;
304 anyhow::ensure!(
305 fork.status == ForkStatus::Active,
306 "fork {} was already removed",
307 fork.id
308 );
309 anyhow::ensure!(
310 std::path::Path::new(confirm_path) == fork.workspace,
311 "confirmation path does not match the fork workspace {}; removal is \
312 path-confirmed to prevent accidental deletion",
313 fork.workspace.display()
314 );
315
316 self.remove_workspace_fork(
317 &fork.provider_id.clone(),
318 &fork.id.clone(),
319 RemoveForkPolicy {
320 confirm_workspace: fork.workspace.clone(),
321 },
322 )
323 .await?;
324
325 fork.status = ForkStatus::Removed;
326 metadata.workspace_fork = Some(fork.clone());
327 metadata.updated_at = OffsetDateTime::now_utc();
328 store.update_thread_metadata(metadata).await?;
329
330 self.emit(RoderEvent::ThreadForkRemoved(ThreadForkRemoved {
331 thread_id: thread_id.clone(),
332 fork_id: fork.id.clone(),
333 worktree_path: fork.workspace.display().to_string(),
334 timestamp: OffsetDateTime::now_utc(),
335 }))
336 .await;
337 Ok(fork)
338 }
339}