Skip to main content

mj_controller/import/
native_import.rs

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
15/// Build, verify, and install a local archive for one already-located,
16/// already-parsed native session, for any harness, then update the in-memory
17/// state. The caller saves `state` only after this returns successfully.
18pub 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        // The preview may precede a user's confirmation by minutes. Never
58        // pair its old transcript with a newer native conversation.
59        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(&current.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            // An imported history is a new session: when it is resumed into a
150            // container it gets its own workspace, like any session created
151            // now.
152            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
215/// Target that hosts raw project sessions on this machine.
216pub(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
231/// A session that only wrote to its own repository can keep working in that
232/// directory, so import it as a raw project session instead of a bundle
233/// session. `session_edit_targets` always records the cwd root, so a single
234/// durable root is that root.
235pub 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            // A local source remains identified by its configured path even
260            // after its checkout gains a network origin. GitHub-origin
261            // detection is still used for configured network sources.
262            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        // Indexed parallel iteration keeps repository and manifest order
292        // identical to the configured bundle.
293        .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            // Isolated imports restore the native session's committed work on
320            // top of the source's available network baseline. Raw imports
321            // continue to use the live local checkout, including repositories
322            // without any remote.
323            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
442/// The upstream revision an imported repository deltas from. A repository
443/// without remote-tracking refs cannot tell us which ancestry a newly
444/// provisioned clone has, and Hel never bundles full history, so it fails here.
445pub(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}