Skip to main content

roder_core/
conversation_forks.rs

1//! Conversation forks (roadmap phases 90 + 81).
2//!
3//! Forks an existing thread into a child thread backed by a workspace fork
4//! from any registered `ForkProvider` (default: `git-worktree`): the child
5//! starts from the parent transcript (no side-effectful tool replay — only
6//! conversation history records are copied) and all subsequent tool
7//! execution resolves against the fork workspace because the child's
8//! `ThreadMetadata.workspace` points at it. Cleanup is explicit and
9//! path-confirmed; the parent workspace is never modified.
10
11use 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    /// User-facing fork name; the provider sanitizes it into its naming
34    /// scheme (directories, branches, snapshot names).
35    pub name: String,
36    /// Fork at a specific parent turn; `None` forks at the latest turn.
37    pub from_turn_id: Option<TurnId>,
38    /// Fork provider id; `None` uses [`DEFAULT_FORK_PROVIDER`].
39    pub provider_id: Option<String>,
40    /// Provider-specific options (never secrets).
41    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    /// Seeds a long-lived collaboration agent with a safe subset of the parent
64    /// conversation. Unlike `fork_thread`, this keeps the same workspace: only
65    /// transcript/lifecycle records are copied, never executable tool or approval
66    /// events.
67    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    /// Forks `parent_thread_id` into a new child thread backed by a fresh
112    /// workspace fork of the parent workspace.
113    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        // Materialize the workspace fork first; thread creation only
160        // proceeds once an isolated workspace exists.
161        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            // The fork workspace lives outside registered workspace roots.
198            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            // Local workspace forks never inherit runner bindings.
207            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                // Best-effort cleanup so a failed fork does not leak a workspace.
231                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    /**
282     * Removes the workspace fork behind a forked thread. Destructive and
283     * explicit: `confirm_path` must match the fork workspace exactly. The
284     * thread itself is kept (status flips to `Removed`) so the conversation
285     * stays readable.
286     */
287    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}