1use 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
20pub const DEFAULT_FORK_PROVIDER: &str = "git-worktree";
22
23impl Runtime {
24 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 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 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 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
112pub const REMOTE_RUNNER_FORK_PROVIDER_ID: &str = "remote-runner";
115
116pub struct RemoteRunnerForkAdapter {
127 provider: Arc<dyn RemoteRunnerProvider>,
128 destination: RunnerDestination,
129 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 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 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 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}