Skip to main content

roder_core/
forks.rs

1//! Runtime fork manager (roadmap phase 81, Task 2): resolves configured
2//! `ForkProvider`s from the extension registry and applies path/policy
3//! checks before provider calls. Thread attachment lives in
4//! `conversation_forks`; these are the provider-facing primitives shared by
5//! conversation forks and the `forks/*` app-server surface.
6
7use std::collections::HashMap;
8use std::path::PathBuf;
9use std::sync::Arc;
10
11use roder_api::forks::{
12    ForkCapabilities, ForkId, ForkProvenance, ForkProvider, ForkProviderDescriptor, ForkRequest,
13    ForkStatus, RemoveForkPolicy, RemoveForkResult, WorkspaceFork,
14};
15use roder_api::remote_runner::{RemoteRunnerProvider, RemoteRunnerSession, RunnerDestination};
16use tokio::sync::Mutex;
17
18use crate::Runtime;
19
20/// Default fork provider when a request does not name one.
21pub const DEFAULT_FORK_PROVIDER: &str = "git-worktree";
22
23impl Runtime {
24    /// Lists registered fork providers.
25    pub fn fork_providers(&self) -> Vec<ForkProviderDescriptor> {
26        self.registry
27            .fork_providers
28            .iter()
29            .map(|provider| provider.descriptor())
30            .collect()
31    }
32
33    /// Resolves a fork provider with an actionable error naming the
34    /// available ids.
35    pub fn fork_provider(&self, provider_id: &str) -> anyhow::Result<Arc<dyn ForkProvider>> {
36        self.registry.fork_provider(provider_id).ok_or_else(|| {
37            let available = self
38                .registry
39                .fork_providers
40                .iter()
41                .map(|provider| provider.descriptor().id)
42                .collect::<Vec<_>>()
43                .join(", ");
44            anyhow::anyhow!(
45                "fork provider {provider_id:?} is not installed (available: {})",
46                if available.is_empty() {
47                    "none"
48                } else {
49                    &available
50                }
51            )
52        })
53    }
54
55    /// Creates a workspace fork after host-side validation.
56    pub async fn create_workspace_fork(
57        &self,
58        provider_id: &str,
59        request: ForkRequest,
60    ) -> anyhow::Result<WorkspaceFork> {
61        anyhow::ensure!(
62            request.source_workspace.is_absolute(),
63            "fork source workspace must be an absolute path: {}",
64            request.source_workspace.display()
65        );
66        if let Some(name) = &request.name {
67            anyhow::ensure!(
68                !name.trim().is_empty() && name.len() <= 100,
69                "fork names must be 1-100 characters"
70            );
71        }
72        let provider = self.fork_provider(provider_id)?;
73        anyhow::ensure!(
74            provider.descriptor().capabilities.create,
75            "fork provider {provider_id:?} does not support creation"
76        );
77        provider.create_fork(request).await
78    }
79
80    pub async fn list_workspace_forks(
81        &self,
82        provider_id: &str,
83        source_workspace: &std::path::Path,
84    ) -> anyhow::Result<Vec<WorkspaceFork>> {
85        self.fork_provider(provider_id)?
86            .list_forks(source_workspace)
87            .await
88    }
89
90    pub async fn resume_workspace_fork(
91        &self,
92        provider_id: &str,
93        id: &ForkId,
94    ) -> anyhow::Result<WorkspaceFork> {
95        self.fork_provider(provider_id)?.resume_fork(id).await
96    }
97
98    /// Removes a fork; destructive and always path-confirmed by the
99    /// provider contract.
100    pub async fn remove_workspace_fork(
101        &self,
102        provider_id: &str,
103        id: &ForkId,
104        policy: RemoveForkPolicy,
105    ) -> anyhow::Result<RemoveForkResult> {
106        self.fork_provider(provider_id)?
107            .remove_fork(id, policy)
108            .await
109    }
110}
111
112/// Provider id used by the remote-runner fork adapter (roadmap phase 81,
113/// Task 5).
114pub const REMOTE_RUNNER_FORK_PROVIDER_ID: &str = "remote-runner";
115
116/**
117 * Represents a fresh remote-runner session as a `WorkspaceFork` with
118 * `remote_compute = true`. The fork layer owns only lifecycle, provenance,
119 * and attachment — file/process operations stay delegated to the
120 * `RemoteRunnerSession` contract. `remove_fork` closes the runner session
121 * (providers without snapshot deletion simply terminate the session, which
122 * is documented as the deterministic cleanup behavior); `resume_fork`
123 * re-opens it through the runner provider's `resume_session` path, with
124 * snapshot-backed restore when the provider recorded one.
125 */
126pub struct RemoteRunnerForkAdapter {
127    provider: Arc<dyn RemoteRunnerProvider>,
128    destination: RunnerDestination,
129    /// Absolute workspace path on the runner.
130    runner_workspace: PathBuf,
131    sessions: Mutex<HashMap<ForkId, Arc<dyn RemoteRunnerSession>>>,
132}
133
134impl RemoteRunnerForkAdapter {
135    pub fn new(
136        provider: Arc<dyn RemoteRunnerProvider>,
137        destination: RunnerDestination,
138        runner_workspace: PathBuf,
139    ) -> Self {
140        Self {
141            provider,
142            destination,
143            runner_workspace,
144            sessions: Mutex::new(HashMap::new()),
145        }
146    }
147
148    /// The live runner session backing a fork, for tool execution wiring.
149    pub async fn session(&self, id: &ForkId) -> Option<Arc<dyn RemoteRunnerSession>> {
150        self.sessions.lock().await.get(id).cloned()
151    }
152
153    fn fork_for(&self, session: &Arc<dyn RemoteRunnerSession>) -> WorkspaceFork {
154        let state = session.state();
155        WorkspaceFork {
156            id: state.session_id.clone(),
157            provider_id: REMOTE_RUNNER_FORK_PROVIDER_ID.to_string(),
158            source_workspace: self.runner_workspace.clone(),
159            workspace: self.runner_workspace.clone(),
160            status: ForkStatus::Active,
161            provenance: ForkProvenance {
162                branch: None,
163                source_branch: None,
164                source_commit: None,
165                snapshot_id: state
166                    .snapshot
167                    .as_ref()
168                    .map(|snapshot| snapshot.snapshot_id.clone()),
169                session_id: Some(state.session_id),
170                created_at: time::OffsetDateTime::now_utc(),
171            },
172            cleanup: Default::default(),
173            metadata: serde_json::json!({
174                "runnerProviderId": state.provider_id,
175                "destinationId": state.destination_id,
176            }),
177        }
178    }
179}
180
181#[async_trait::async_trait]
182impl ForkProvider for RemoteRunnerForkAdapter {
183    fn descriptor(&self) -> ForkProviderDescriptor {
184        ForkProviderDescriptor {
185            id: REMOTE_RUNNER_FORK_PROVIDER_ID.to_string(),
186            display_name: format!("Remote runner ({})", self.destination.provider_id),
187            capabilities: ForkCapabilities {
188                create: true,
189                list: false,
190                remove: true,
191                resume: true,
192                diff_summary: false,
193                merge_back: false,
194                copy_on_write: false,
195                remote_compute: true,
196            },
197        }
198    }
199
200    async fn create_fork(&self, request: ForkRequest) -> anyhow::Result<WorkspaceFork> {
201        anyhow::ensure!(
202            !request.policy.allow_dirty_source,
203            "remote-runner forks always start from the destination's own state; \
204             allow_dirty_source has no meaning here and must stay false"
205        );
206        let session = self
207            .provider
208            .create_session(self.destination.clone())
209            .await?;
210        let fork = self.fork_for(&session);
211        self.sessions.lock().await.insert(fork.id.clone(), session);
212        Ok(fork)
213    }
214
215    async fn list_forks(&self, _source: &std::path::Path) -> anyhow::Result<Vec<WorkspaceFork>> {
216        // Runner providers own session listing; the adapter only tracks the
217        // sessions it created in this process.
218        let sessions = self.sessions.lock().await;
219        Ok(sessions
220            .values()
221            .map(|session| self.fork_for(session))
222            .collect())
223    }
224
225    async fn resume_fork(&self, id: &ForkId) -> anyhow::Result<WorkspaceFork> {
226        if let Some(session) = self.session(id).await {
227            return Ok(self.fork_for(&session));
228        }
229        // Snapshot-backed resume through the runner provider.
230        let state = roder_api::remote_runner::RunnerSessionState {
231            provider_id: self.destination.provider_id.clone(),
232            session_id: id.clone(),
233            destination_id: self.destination.id.clone(),
234            snapshot: None,
235            metadata: self.destination.config.clone(),
236        };
237        let session = self.provider.resume_session(state).await?;
238        let fork = self.fork_for(&session);
239        self.sessions.lock().await.insert(fork.id.clone(), session);
240        Ok(fork)
241    }
242
243    async fn remove_fork(
244        &self,
245        id: &ForkId,
246        policy: RemoveForkPolicy,
247    ) -> anyhow::Result<RemoveForkResult> {
248        anyhow::ensure!(
249            policy.confirm_workspace == self.runner_workspace,
250            "removal is path-confirmed: confirm the runner workspace {}",
251            self.runner_workspace.display()
252        );
253        let session = self
254            .sessions
255            .lock()
256            .await
257            .remove(id)
258            .ok_or_else(|| anyhow::anyhow!("remote fork {id} is not active in this process"))?;
259        session.close().await?;
260        Ok(RemoveForkResult {
261            id: id.clone(),
262            removed: true,
263            workspace: self.runner_workspace.clone(),
264        })
265    }
266}