1use super::*;
2
3pub struct NativeImportRequest<'a> {
4 pub harness: HarnessKind,
5 pub harness_home: &'a Path,
6 pub native_session_id: &'a str,
7 pub source_path: &'a Path,
8 pub transcript: &'a ClaudeTranscript,
9 pub bundle_id: &'a str,
10 pub profile_id: Option<&'a str>,
11 pub title: Option<&'a str>,
12 pub archive_directory: &'a Path,
13}
14
15pub fn import_native_session(
19 config: &Config,
20 state: &mut State,
21 request: NativeImportRequest<'_>,
22 control: Option<&ImportControl<'_>>,
23) -> Result<ImportedClaudeSession> {
24 let NativeImportRequest {
25 harness,
26 harness_home,
27 native_session_id,
28 source_path,
29 transcript,
30 bundle_id,
31 profile_id,
32 title,
33 archive_directory,
34 } = request;
35 let bundle = config
36 .bundles
37 .get(bundle_id)
38 .with_context(|| format!("unknown bundle {bundle_id:?}"))?;
39 let session_title_override = title.map(str::to_owned);
40 let title = match session_title_override.as_deref() {
41 Some(title) if !title.trim().is_empty() => title.to_owned(),
42 Some(_) => bail!("import title must not be empty"),
43 None => harness_session_title(&transcript.events).unwrap_or_else(|| {
44 format!(
45 "Imported {} session {native_session_id}",
46 harness.display_name()
47 )
48 }),
49 };
50 let targets = session_edit_targets(transcript, harness_home)?;
51 let raw_project = raw_project_import(config, &targets);
52 let repositories =
53 collect_local_repositories(bundle, &targets.git_roots, raw_project.is_none(), control)?;
54 let native_artifacts =
55 collect_import_native_artifacts(harness, harness_home, native_session_id, source_path)?;
56 if harness == HarnessKind::Muse {
57 if let Some(control) = control {
60 control.check_cancelled()?;
61 }
62 let current = read_native_transcript(harness, source_path)?;
63 ensure!(
64 current.cwd == transcript.cwd
65 && current.edited_paths == transcript.edited_paths
66 && serde_json::to_value(¤t.events)?
67 == serde_json::to_value(&transcript.events)?,
68 "native session changed after it was selected; select it again"
69 );
70 ensure!(
71 native_artifacts
72 == collect_import_native_artifacts(
73 harness,
74 harness_home,
75 native_session_id,
76 source_path
77 )?,
78 "native session changed while being imported; stop its harness and retry"
79 );
80 }
81 let session_id = new_session_id()?;
82 let canonical_session =
83 canonical_import_session(session_id.as_str(), &transcript.events, source_path)?;
84 let timestamp = timestamp();
85 let profile_id = import_profile_id(config, profile_id, harness, harness_home)?;
86 let target_id = default_import_target_id(config);
87 let archive_path = archive_directory.join(format!("{session_id}.hel.zip"));
88 if let Some(control) = control {
89 control.report(ImportArchiveProgress::WritingArchive)?;
90 }
91 let verified = write_archive_atomic(
92 &archive_path,
93 &ArchiveInput {
94 session: mj_checkpoint::archive::SessionManifest {
95 id: session_id.clone(),
96 title: title.clone(),
97 harness_kind: harness,
98 profile_id: profile_id.clone(),
99 native_session_id: native_session_id.to_owned(),
100 created_at: timestamp.clone(),
101 checkpointed_at: timestamp.clone(),
102 hel_version: env!("CARGO_PKG_VERSION").into(),
103 relay_version: env!("CARGO_PKG_VERSION").into(),
104 adapter_version: "acp-v1".into(),
105 },
106 target: TargetManifest {
107 template_id: target_id.clone(),
108 target_kind: "import".into(),
109 details: BTreeMap::from([("source".into(), format!("{}-import", harness.id()))]),
110 },
111 bundle: BundleManifest {
112 id: bundle_id.to_owned(),
113 primary_repository: bundle.primary_repo.clone(),
114 },
115 canonical_session,
116 native_artifacts,
117 repositories,
118 },
119 )?;
120 if let Some(control) = control
121 && let Err(error) = control.check_cancelled()
122 {
123 let _ = fs::remove_file(&archive_path);
124 return Err(error);
125 }
126 let checkpoint = CheckpointMetadata {
127 archive_path: archive_path.clone(),
128 sha256: verified.archive_sha256,
129 created_at: timestamp.clone(),
130 event_frontier: transcript.events.last().map_or(0, |event| event.seq),
131 };
132 state.sessions.insert(
133 session_id.clone(),
134 SessionRecord {
135 project: Some(crate::project_catalog::snapshot(
136 bundle,
137 &mj_core::targets::CancellableProcessExecutor::with_timeout(Duration::from_secs(
138 15,
139 )),
140 raw_project.is_none(),
141 )?),
142 target_runtime: None,
143 launch_base: None,
144 launch_branch: None,
145 checkout: None,
146 publication: None,
147 build_cache: None,
148 subagents: None,
149 container_workspace: Some(mj_core::targets::new_container_workspace(&session_id)?),
153 create_managed_worktree: None,
154 workspace_id: mj_core::workspace::DEFAULT_WORKSPACE_ID.to_owned(),
155 archived: false,
156 container_cpus: None,
157 container_memory: None,
158 id: session_id.clone(),
159 title,
160 harness_kind: harness,
161 last_profile: profile_id,
162 bundle_id: bundle_id.to_owned(),
163 project_directory: raw_project.as_ref().map(|(directory, _)| directory.clone()),
164 managed_worktree: None,
165 review: None,
166 target_template_id: raw_project.map_or(target_id, |(_, raw_target_id)| raw_target_id),
167 resource_allocation: None,
168 additional_mounts: Vec::new(),
169 state: SessionState::Stopped,
170 target: None,
171 native_session_id: Some(native_session_id.to_owned()),
172 acp_session_title: None,
173 session_title_override,
174 created_at: timestamp.clone(),
175 updated_at: timestamp,
176 viewed_through_event_ordinal: 0,
177 draft_input: String::new(),
178 last_error: None,
179 last_checkpoint_error: None,
180 checkpoint: Some(checkpoint),
181 },
182 );
183 Ok(ImportedClaudeSession {
184 session_id,
185 native_session_id: native_session_id.to_owned(),
186 source_jsonl: source_path.to_path_buf(),
187 source_cwd: transcript.cwd.clone(),
188 bundle_id: bundle_id.to_owned(),
189 archive_path,
190 })
191}
192
193pub(super) fn default_import_target_id(config: &Config) -> String {
194 config
195 .targets
196 .get_key_value("podman")
197 .map(|(id, _)| id)
198 .or_else(|| {
199 config.targets.iter().find_map(|(id, target)| {
200 matches!(
201 target,
202 TargetTemplate::LocalPodman { .. }
203 | TargetTemplate::LocalDocker { .. }
204 | TargetTemplate::SshPodman { .. }
205 | TargetTemplate::SshDocker { .. }
206 )
207 .then_some(id)
208 })
209 })
210 .or_else(|| config.targets.keys().next())
211 .cloned()
212 .unwrap_or_else(|| "import".into())
213}
214
215pub(super) fn raw_import_target_id(config: &Config) -> Option<String> {
217 let local_bare = |template: &TargetTemplate| matches!(template, TargetTemplate::LocalBare);
218 config
219 .targets
220 .get_key_value("localhost")
221 .filter(|(_, template)| local_bare(template))
222 .map(|(id, _)| id.clone())
223 .or_else(|| {
224 config
225 .targets
226 .iter()
227 .find_map(|(id, template)| local_bare(template).then(|| id.clone()))
228 })
229}
230
231pub fn raw_project_import(
236 config: &Config,
237 targets: &SessionEditTargets,
238) -> Option<(PathBuf, String)> {
239 let [cwd_root] = targets.git_roots.as_slice() else {
240 return None;
241 };
242 Some((cwd_root.clone(), raw_import_target_id(config)?))
243}
244
245pub(super) fn collect_local_repositories(
246 bundle: &ProjectBundle,
247 detected_roots: &[PathBuf],
248 isolated: bool,
249 control: Option<&ImportControl<'_>>,
250) -> Result<Vec<mj_checkpoint::archive::RepositorySnapshot>> {
251 let detected = detected_roots
252 .iter()
253 .map(|root| Ok((root_identity(root)?, root.clone())))
254 .collect::<Result<BTreeMap<_, _>>>()?;
255 let repository_paths = bundle
256 .repositories
257 .iter()
258 .map(|repository| {
259 let path = if let Some(configured_path) = repository.local.as_ref() {
263 let configured_path =
264 fs::canonicalize(configured_path).unwrap_or_else(|_| configured_path.clone());
265 detected_roots
266 .iter()
267 .find(|root| {
268 fs::canonicalize(root).unwrap_or_else(|_| (*root).clone())
269 == configured_path
270 })
271 .cloned()
272 } else {
273 let identity = configured_repository_identity(repository)?.with_context(|| {
274 format!("repository {:?} has no usable source", repository.id)
275 })?;
276 detected.get(&identity).cloned()
277 };
278 let path = path.with_context(|| {
279 format!(
280 "repository {:?} was not detected in the native session",
281 repository.id
282 )
283 })?;
284 Ok((repository.id.clone(), path))
285 })
286 .collect::<Result<BTreeMap<_, _>>>()?;
287 let git = SystemGit;
288 let repository_count = bundle.repositories.len();
289 bundle
290 .repositories
291 .par_iter()
294 .enumerate()
295 .map(|(index, repository)| {
296 if let Some(control) = control {
297 control.report(ImportArchiveProgress::Repository {
298 current: index + 1,
299 total: repository_count,
300 id: repository.id.clone(),
301 })?;
302 }
303 let path = repository_paths
304 .get(&repository.id)
305 .expect("repository paths cover the validated bundle")
306 .clone();
307 ensure!(
308 path.is_dir(),
309 "local repository {:?} is missing at {}",
310 repository.id,
311 path.display()
312 );
313 let source = isolated
314 .then(|| {
315 resolve_repository(repository, &ProcessExecutor)
316 .with_context(|| format!("resolve network source for {:?}", repository.id))
317 })
318 .transpose()?;
319 let history = if isolated {
324 GitHistoryMode::DeltaFrom(import_delta_base(
325 &path,
326 &source.as_ref().expect("isolated source resolved").fetch_url,
327 )?)
328 } else {
329 GitHistoryMode::NoBundle
330 };
331 let origin_override = if let Some(source) = &source {
332 Some(source.fetch_url.clone())
333 } else {
334 Some(path.to_string_lossy().into_owned())
335 };
336 let mut snapshot = collect_git_snapshot_with_progress(
337 &git,
338 &path,
339 &GitCollectionSpec {
340 id: repository.id.clone(),
341 relative_destination: repository.destination.clone(),
342 history,
343 origin_override,
344 },
345 control.is_none_or(|control| control.include_untracked),
346 &|progress| {
347 let Some(control) = control else {
348 return Ok(());
349 };
350 match progress {
351 GitSnapshotProgress::UntrackedFile {
352 current,
353 total,
354 path,
355 } => control.report(ImportArchiveProgress::UntrackedFile {
356 repository_id: repository.id.clone(),
357 current,
358 total,
359 path,
360 }),
361 }
362 },
363 )
364 .with_context(|| format!("collect local repository {:?}", repository.id))?;
365 if let Some(source) = source {
366 snapshot.metadata.push_urls = source
367 .push_urls
368 .iter()
369 .map(|url| mj_checkpoint::archive::redact_origin_credentials(url))
370 .collect::<Result<Vec<_>>>()?;
371 snapshot.metadata.remote_workspace = true;
372 } else {
373 snapshot.metadata.push_urls.clear();
374 }
375 Ok(snapshot)
376 })
377 .collect()
378}
379
380pub(super) fn canonical_import_session(
381 session_id: &str,
382 events: &[SequencedEvent],
383 source_path: &Path,
384) -> Result<mj_checkpoint::archive::CanonicalSessionSnapshot> {
385 let mut events = events.to_vec();
386 finalize_import_event_times(&mut events, source_path)?;
387 let mut materialized =
388 mj_transcript::projection::imported_materialized_session(session_id, &events);
389 materialized.session_title = harness_session_title(&events);
390 if let Some(last_activity_at_ms) = events.iter().filter_map(|event| event.recorded_at_ms).max()
391 {
392 materialized.last_activity_at_ms = Some(
393 materialized
394 .last_activity_at_ms
395 .map_or(last_activity_at_ms, |current| {
396 current.max(last_activity_at_ms)
397 }),
398 );
399 }
400 canonical_session_from_materialized(&materialized)
401}
402
403pub(super) fn default_profile(config: &Config, harness: HarnessKind, home: &Path) -> String {
404 let source = fs::canonicalize(home).unwrap_or_else(|_| home.to_path_buf());
405 config
406 .enabled_profiles()
407 .find(|(_, profile)| {
408 profile.kind == harness
409 && fs::canonicalize(&profile.home).unwrap_or_else(|_| profile.home.clone())
410 == source
411 })
412 .or_else(|| {
413 config
414 .enabled_profiles()
415 .find(|(_, profile)| profile.kind == harness)
416 })
417 .map(|(id, _)| id.to_owned())
418 .unwrap_or_else(|| format!("{}-import", harness.id()))
419}
420
421pub(super) fn import_profile_id(
422 config: &Config,
423 requested: Option<&str>,
424 harness: HarnessKind,
425 home: &Path,
426) -> Result<String> {
427 let Some(requested) = requested else {
428 return Ok(default_profile(config, harness, home));
429 };
430 let profile = config
431 .profiles
432 .get(requested)
433 .with_context(|| format!("unknown import profile {requested:?}"))?;
434 ensure!(profile.enabled, "import profile {requested:?} is disabled");
435 ensure!(
436 profile.kind == harness,
437 "import profile {requested:?} does not use {harness:?}"
438 );
439 Ok(requested.to_owned())
440}
441
442pub(super) fn import_delta_base(path: &Path, fetch_url: &str) -> Result<String> {
446 let remotes = git_optional_text(path, ["remote"])?.unwrap_or_default();
447 let upstream_ref = git_optional_text(
448 path,
449 [
450 "rev-parse",
451 "--symbolic-full-name",
452 "--verify",
453 "--quiet",
454 "@{upstream}",
455 ],
456 )?;
457 for remote in remotes.lines() {
458 let Some(url) = git_optional_text(path, ["remote", "get-url", remote])? else {
459 continue;
460 };
461 let same_source = url == fetch_url
462 || mj_core::state::ProjectSourceIdentity::git_remote(&url).is_some_and(|identity| {
463 Some(identity) == mj_core::state::ProjectSourceIdentity::git_remote(fetch_url)
464 });
465 if !same_source {
466 continue;
467 }
468 let prefix = format!("refs/remotes/{remote}/");
469 let revision = upstream_ref
470 .as_deref()
471 .filter(|reference| reference.starts_with(&prefix))
472 .map(str::to_owned)
473 .unwrap_or_else(|| format!("{prefix}HEAD"));
474 if let Some(base) =
475 git_optional_text(path, ["rev-parse", "--verify", "--quiet", &revision])?
476 {
477 return Ok(base);
478 }
479 }
480 bail!(
481 "repository {} has no remote-tracking refs to import against for its selected network source; fetch its remote first",
482 path.display()
483 )
484}
485
486pub(super) fn git_optional_text<const N: usize>(
487 cwd: &Path,
488 arguments: [&str; N],
489) -> Result<Option<String>> {
490 let output = Command::new("git")
491 .args(arguments)
492 .current_dir(cwd)
493 .output()
494 .with_context(|| format!("start git in {}", cwd.display()))?;
495 if !output.status.success() {
496 return Ok(None);
497 }
498 let text = String::from_utf8(output.stdout).context("decode Git output")?;
499 Ok((!text.trim().is_empty()).then(|| text.trim().to_owned()))
500}
501
502pub(super) fn timestamp() -> String {
503 Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true)
504}
505
506pub fn persist_imported_session_locally(session: &SessionRecord) -> Result<()> {
507 crate::database::save_session(session)?;
508 let checkpoint = session
509 .checkpoint
510 .as_ref()
511 .context("imported session has no checkpoint")?;
512 let canonical = mj_checkpoint::archive::verify_archive_streaming(&checkpoint.archive_path)?
513 .canonical_session;
514 let materialized = mj_transcript::projection::materialized_session_from_canonical(
515 session.id.clone(),
516 &canonical,
517 )?;
518 crate::database::save_materialized_session(&materialized)
519}