Skip to main content

mj_controller/server/
actions.rs

1use super::*;
2
3/// The complete set of operations a phone may ask the controller to perform.
4/// Secret/config editing is intentionally not representable here, and the one
5/// destructive variant, `ForceClose`, is not representable on the wire: it is
6/// `#[serde(skip)]` so only in-process callers such as the HTTP API can build
7/// it.
8#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
9#[serde(tag = "action", rename_all = "kebab-case", deny_unknown_fields)]
10pub enum ControllerAction {
11    TurnControl {
12        session_id: String,
13        command: mj_core::relay::RelayCommand,
14    },
15    New {
16        #[serde(default)]
17        create_managed_worktree: Option<bool>,
18        /// Full commit object ID to start the bundle's primary repository at.
19        #[serde(default, skip_serializing_if = "Option::is_none")]
20        at: Option<String>,
21        /// With `at`, the new branch created there; otherwise an existing branch.
22        #[serde(default, skip_serializing_if = "Option::is_none")]
23        branch: Option<String>,
24        /// Diff base; defaults to `at`. Without `at`, also the starting
25        /// revision for a raw managed worktree.
26        #[serde(default, skip_serializing_if = "Option::is_none")]
27        base: Option<String>,
28        /// Omitted uses the selected profile's subagent setting.
29        #[serde(default)]
30        subagents: Option<mj_core::subagent::SubagentPolicy>,
31        /// Turn review for this session. Omitted follows `[review]`.
32        #[serde(default, skip_serializing_if = "Option::is_none")]
33        review: Option<mj_core::config::SessionReview>,
34        /// Which workspace the session belongs to. Optional on the wire so a
35        /// viewer cached from before workspaces reached the phone still parses,
36        /// but a controller holding more than one workspace refuses an empty
37        /// one rather than guessing.
38        #[serde(default)]
39        workspace_id: String,
40        profile_id: String,
41        bundle_id: String,
42        target_id: String,
43        /// Explicit container or EC2 sizing selected in the create wizard.
44        /// Omitted preserves the historical target defaults for older clients.
45        #[serde(default, skip_serializing_if = "Option::is_none")]
46        resource_allocation: Option<Box<mj_core::state::SessionResourceAllocation>>,
47        /// Absent means "derive it", which is what the terminal does.
48        #[serde(default)]
49        title: Option<String>,
50        #[serde(default)]
51        project_directory: Option<PathBuf>,
52        /// The repositories the person was shown as having uncommitted changes
53        /// and chose to launch over anyway.
54        ///
55        /// This names them rather than being a bare yes, so an acknowledgement
56        /// cannot be replayed against a set the person never saw: if a
57        /// different repository has gone dirty since the preflight, the launch
58        /// stops and asks again.
59        #[serde(default, skip_serializing_if = "Vec::is_empty")]
60        dirty_ack: Vec<String>,
61    },
62    /// Give a session a new title. The terminal calls this a rename.
63    Rename {
64        session_id: String,
65        title: String,
66    },
67    /// Stop the turn the agent is working on, leaving the session alive. This
68    /// is not `Cancel`, which stops a provision, resume or stop.
69    InterruptTurn {
70        session_id: String,
71    },
72    /// Change one setting the harness advertised, such as `model` or `effort`.
73    SetConfig {
74        session_id: String,
75        key: String,
76        value: String,
77    },
78    /// Turn plan mode on or off. The harness decides how, which is why this
79    /// carries an intent rather than a mode id.
80    SetPlanMode {
81        session_id: String,
82        active: bool,
83    },
84    /// Move a session, and the sub-agents under it, to another workspace.
85    ChangeWorkspace {
86        session_id: String,
87        workspace_id: String,
88    },
89    /// Record the session's container size overrides and attached
90    /// directories. An empty or absent size clears that override, and the
91    /// mounts list replaces the whole list. Nothing changes inside a running
92    /// container: the values are read the next time it is created.
93    SetContainerSettings {
94        session_id: String,
95        #[serde(default)]
96        cpus: Option<String>,
97        #[serde(default)]
98        memory: Option<String>,
99        #[serde(default)]
100        mounts: Vec<AdditionalMount>,
101    },
102    /// Suspend the session if it is live, then resume it with the profile,
103    /// target and mounts it last ran with.
104    Restart {
105        session_id: String,
106    },
107    /// Stop the turn this session is working on and the turns of every
108    /// sub-agent under it, leaving the sessions alive.
109    InterruptAll {
110        session_id: String,
111    },
112    RefreshQuota {
113        profile_id: String,
114    },
115    RefreshCapacity {
116        target_id: String,
117    },
118    Resume {
119        session_id: String,
120        workspace_id: String,
121        profile_id: String,
122        target_id: String,
123        queue: ResumeQueueDisposition,
124        /// A failed Move supplies the settings recorded before source
125        /// teardown. Ordinary Resume requests leave these absent and retain
126        /// the historical inheritance behavior.
127        #[serde(default)]
128        additional_mounts: Option<Vec<AdditionalMount>>,
129        #[serde(default)]
130        resource_allocation: Option<SessionResourceAllocation>,
131    },
132    /// Confirm a previously prepared move. Preparation is a separate
133    /// authenticated request so changing the destination cannot be smuggled
134    /// into a confirmation from an older browser form.
135    Move {
136        request: Box<MoveSessionRequest>,
137    },
138    Open {
139        session_id: String,
140    },
141    Prompt {
142        #[serde(default, skip_serializing_if = "Option::is_none")]
143        command_id: Option<String>,
144        session_id: String,
145        text: String,
146        /// Images to send with the prompt. The controller turns each one into
147        /// the ACP image content block its prompt path already speaks.
148        #[serde(default, skip_serializing_if = "Vec::is_empty")]
149        images: Vec<ViewerPromptImage>,
150    },
151    RunShell {
152        #[serde(default, skip_serializing_if = "Option::is_none")]
153        command_id: Option<String>,
154        session_id: String,
155        command: String,
156    },
157    CancelShell {
158        session_id: String,
159        shell_command_id: String,
160    },
161    Suspend {
162        session_id: String,
163        #[serde(default)]
164        acknowledge_unpublished_work: bool,
165    },
166    /// Destroy a session without checkpointing it: the live target is torn
167    /// down, the recovery archive is removed, and sub-agent children are
168    /// destroyed first. This is irreversible.
169    ///
170    /// Skipped by serde on purpose. The browser viewer posts this enum to
171    /// `/actions`, so a wire request must never be able to name this variant;
172    /// it is reachable only from the HTTP API, which builds it in process.
173    #[serde(skip)]
174    Destroy {
175        session_id: String,
176        /// Whether the managed worktree's branch goes with the session.
177        /// Destruction keeps it unless the request asks for the deletion.
178        delete_branch: bool,
179    },
180    Cancel {
181        session_id: String,
182    },
183    /// Review the turn this session just finished.
184    StartReview {
185        session_id: String,
186    },
187    /// Forward the findings, dismiss them, or cancel the open review.
188    ResolveReview {
189        session_id: String,
190        /// `forward`, `dismiss`, or `cancel`.
191        resolution: String,
192    },
193    RemoveQueuedPrompt {
194        session_id: String,
195        queue_id: String,
196    },
197    /// Answer one of the session's pending form questions.
198    RespondElicitation {
199        session_id: String,
200        elicitation_id: String,
201        response: ElicitationResponse,
202    },
203}
204
205/// One image a phone attached to a prompt. Legacy callers may send inline
206/// base64 data; the server normalizes it into an attachment before dispatch.
207#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
208#[serde(deny_unknown_fields)]
209pub struct ViewerPromptImage {
210    /// Legacy inline image bytes. New browser uploads and normalized inline
211    /// prompts carry an attachment reference and leave this empty.
212    #[serde(default)]
213    pub data_base64: String,
214    pub mime_type: String,
215    pub width: u32,
216    pub height: u32,
217    /// Session-scoped, immutable image bytes. The worker resolves this just
218    /// before dispatch, keeping browser actions and durable commands small.
219    #[serde(default, skip_serializing_if = "Option::is_none")]
220    pub attachment: Option<AttachmentRef>,
221}
222
223/// The controller's answer to one phone action.
224///
225/// The answer means "accepted", not "finished": provisioning, resume and close
226/// run for minutes, and a phone on a mobile network drops a request held open
227/// that long. How the action then goes travels in snapshots — session state,
228/// queued prompts, transcripts, and `has_error`.
229///
230/// Only the outcome crosses this boundary. The controller's own failure text
231/// names profile homes, project paths and SSH hosts, so it stays on the
232/// controller. A caller therefore gets one of two things: a [`Refusal`], whose
233/// sentence was written for it at the place the failure was produced, or a
234/// generic internal failure carrying a reference that also appears in the
235/// daemon log.
236#[derive(Debug, Clone, PartialEq, Eq)]
237pub enum ActionOutcome {
238    /// Admitted and now running; watch the snapshot for what happens next.
239    ///
240    /// A `new` action carries the published session id, which is the only way
241    /// its caller learns what it just created.
242    Accepted { session_id: Option<String> },
243    /// The controller already runs as many phone actions as it allows.
244    /// `running` is the admission control's own count of in-flight actions and
245    /// `limit` the most it admits, so the message never keeps a second copy.
246    Busy { running: usize, limit: usize },
247    /// This session already has an operation running.
248    SessionBusy,
249    /// A cancel found no operation to cancel.
250    NotCancellable,
251    /// The action was refused for a reason the caller can act on, and the
252    /// refusal says what it is.
253    Refused(Refusal),
254    /// The controller could not start the action, for a reason that stays
255    /// server-side. `reference` is logged with the failure, so the person who
256    /// owns the daemon can find the entry that explains it.
257    Failed { reference: String },
258}
259
260impl ActionOutcome {
261    /// Admitted, with no session id to report.
262    pub const fn accepted() -> Self {
263        Self::Accepted { session_id: None }
264    }
265
266    /// The published session id, when this outcome carries one.
267    pub fn session_id(&self) -> Option<&str> {
268        match self {
269            Self::Accepted { session_id } => session_id.as_deref(),
270            _ => None,
271        }
272    }
273
274    /// The reply an outcome owes the phone, or `None` when it was accepted.
275    pub(super) fn rejection(&self) -> Option<ApiError> {
276        match self {
277            Self::Accepted { .. } => None,
278            Self::Busy { running, limit } => Some(
279                ApiError::new(
280                    StatusCode::TOO_MANY_REQUESTS,
281                    format!(
282                        "the daemon is at its limit of {limit} concurrent actions ({running} running; a session that is still starting holds one until it is ready); retry shortly"
283                    ),
284                )
285                .with_busy(*running, *limit),
286            ),
287            Self::SessionBusy => Some(ApiError::new(
288                StatusCode::CONFLICT,
289                "another operation is already running for this session",
290            )),
291            Self::NotCancellable => Some(ApiError::new(
292                StatusCode::CONFLICT,
293                "the session has no cancellable operation",
294            )),
295            // A refusal is a precondition the caller can fix, so it answers
296            // 4xx with the sentence written for it: 409 for a state that has
297            // to change first, 422 for a request naming something unusable.
298            Self::Refused(refusal) => Some(
299                ApiError::new(
300                    match refusal.kind() {
301                        RefusalKind::Precondition => StatusCode::CONFLICT,
302                        RefusalKind::Unusable => StatusCode::UNPROCESSABLE_ENTITY,
303                    },
304                    refusal.message().to_owned(),
305                )
306                .with_code(refusal.code()),
307            ),
308            Self::Failed { reference } => Some(ApiError::new(
309                StatusCode::INTERNAL_SERVER_ERROR,
310                format!(
311                    "the controller could not start this action; \
312                     the daemon log records the reason under reference {reference}"
313                ),
314            )),
315        }
316    }
317}
318
319#[derive(Debug)]
320pub struct ControllerRequest {
321    pub action: ControllerAction,
322    pub reply: tokio::sync::oneshot::Sender<ActionOutcome>,
323}
324
325/// A phone request to create or reuse a quick project bundle. This has its
326/// own channel because bundle creation returns a durable id and must publish a
327/// config snapshot before the HTTP request can succeed; [`ControllerAction`]
328/// intentionally carries only action admission outcomes.
329#[derive(Debug)]
330pub struct BundleRequest {
331    /// Legacy primary-repository matching. Empty when `exact_sources` is set.
332    pub source: String,
333    /// The requested full repository set, with its primary repository first.
334    pub exact_sources: Option<Vec<String>>,
335    pub reply: tokio::sync::oneshot::Sender<Result<String, BundleFailure>>,
336}
337
338/// Safe failure classes for bundle creation. Detailed controller errors stay
339/// in daemon logs; a browser only needs to know whether to fix its source or
340/// report a server-side failure.
341#[derive(Debug, Clone, Copy, PartialEq, Eq)]
342pub enum BundleFailure {
343    InvalidSource,
344    Controller,
345}
346
347/// A phone acknowledging how far it has read a conversation.
348///
349/// This deliberately is not a `ControllerAction`: the viewer posts it after
350/// every conversation fetch, and a fetch follows every revision. Routing it
351/// through the action pipeline made each receipt reload the controller, bump
352/// the revision and broadcast a snapshot, which triggered the next fetch, so
353/// viewer and controller never went quiet; it also consumed the session's
354/// single action slot, intermittently rejecting real actions. A receipt
355/// therefore travels on its own channel and only persists one cursor field.
356/// A phone asking whether a session it is about to create would launch
357/// cleanly, and which network sources it will use first.
358///
359/// This is not a `ControllerAction`: it starts nothing, it takes no session
360/// slot, and it must answer before the person has decided anything. It also
361/// needs the controller, because resolving a local repository's configured
362/// remotes is a fact about the disk rather than about the projection.
363///
364/// Resume preflights share this channel, and so the concurrency cap on it,
365/// because they do the same kind of work on the same disk.
366#[derive(Debug)]
367pub enum PreflightRequest {
368    New(NewPreflightRequest),
369    Resume(ResumePreflightRequest),
370    CompletePath(PathCompletionRequest),
371    DiscoverProjects(ProjectDiscoveryPreflight),
372    ProjectCatalog {
373        refresh: bool,
374        retry: bool,
375        reply: tokio::sync::oneshot::Sender<
376            Result<mj_core::project_catalog::ProjectCatalogView, String>,
377        >,
378    },
379}
380
381/// A project picker lookup sharing the preflight supervision and concurrency cap.
382#[derive(Debug)]
383pub struct ProjectDiscoveryPreflight {
384    pub request: crate::project_picker::ProjectDiscoveryRequest,
385    pub reply:
386        tokio::sync::oneshot::Sender<Result<crate::project_picker::ProjectDiscovery, &'static str>>,
387}
388
389/// A browser asking what a half-typed path could be. It shares the preflight
390/// channel because it does the same kind of work: one short-lived, cancellable
391/// look at a local or remote filesystem, under the same concurrency cap.
392#[derive(Debug)]
393pub struct PathCompletionRequest {
394    pub host: CompletionHost,
395    pub prefix: String,
396    pub kind: CompletionKind,
397    pub reply: tokio::sync::oneshot::Sender<Result<PathCompletion, String>>,
398}
399
400#[derive(Debug)]
401pub struct NewPreflightRequest {
402    pub bundle_id: String,
403    pub target_id: String,
404    pub project_directory: Option<PathBuf>,
405    pub remote_repairs: Vec<mj_core::local_git::LocalRemoteRepair>,
406    pub reply: tokio::sync::oneshot::Sender<Result<PreflightNew, PreflightFailure>>,
407}
408
409/// A resume preflight for one stopped session and one destination target. It
410/// travels on the same channel and under the same concurrency cap as the
411/// new-session preflight because it does the same kind of work: reading a
412/// working tree and asking a remote about itself.
413#[derive(Debug)]
414pub struct ResumePreflightRequest {
415    pub session_id: String,
416    pub target_id: String,
417    pub reply: tokio::sync::oneshot::Sender<Result<PreflightResume, PreflightFailure>>,
418}
419
420/// What a resume preflight found.
421///
422/// `Ready` covers every resume that changes nothing about where repository
423/// content comes from. `ConvertingRawCheckout` means this resume moves a
424/// local checkout into an isolated workspace, and carries the preview the
425/// person has to confirm. `Unavailable` reports why the conversion cannot be
426/// planned, in the plan's own words, because that message says what to do
427/// about it (add a remote, commit a submodule) and the browser has no other
428/// way to learn it.
429#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
430#[serde(tag = "kind", rename_all = "kebab-case")]
431pub enum PreflightResume {
432    Ready,
433    ConvertingRawCheckout {
434        preview: Box<mj_core::state::RawConversionPreview>,
435    },
436    Unavailable {
437        detail: String,
438    },
439}
440
441/// A move preparation is intentionally separate from action admission. It
442/// performs read-only compatibility checks and returns the exact fingerprint
443/// the later confirmation must echo; it never interrupts the source session.
444#[derive(Debug)]
445pub struct MovePreparationRequest {
446    pub selection: MoveSelection,
447    pub reply: tokio::sync::oneshot::Sender<Result<MovePreparation, String>>,
448}
449
450/// A preflight can fail because the requested bare directory is unusable, an
451/// isolated repository lacks a usable network source, or the controller-side
452/// check itself could not complete. The HTTP surface keeps those outcomes
453/// distinct without carrying filesystem, Git, or SSH details to the phone.
454#[derive(Debug)]
455pub enum PreflightFailure {
456    Validation,
457    /// A configured isolated-session repository cannot be used as a network
458    /// source. The detail is safe for the phone and tells the person how to
459    /// choose the supported raw-local path instead.
460    InvalidRepository(String),
461    Controller(String),
462}
463
464/// One configured repository's network clone and publication destinations.
465/// URLs have already been passed through the shared display sanitizer before
466/// they reach a phone.
467#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
468#[serde(deny_unknown_fields)]
469pub struct PreflightRepository {
470    pub id: String,
471    pub fetch_url: String,
472    pub default_branch: String,
473    pub push_urls: Vec<String>,
474}
475
476/// What a preflight found. Isolated sessions expose their complete network
477/// source plan so the person can review it before creation. Raw-local targets
478/// leave the plan empty because they use the selected checkout directly;
479/// isolated targets set `local_changes_excluded` to make the copy boundary
480/// explicit.
481#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
482#[serde(deny_unknown_fields)]
483pub struct PreflightNew {
484    #[serde(default, skip_serializing_if = "Option::is_none")]
485    pub project_directory: Option<PathBuf>,
486    #[serde(default)]
487    pub managed_worktree: mj_core::state::ManagedWorktreeOptions,
488    #[serde(default)]
489    pub remote_repairs: Vec<mj_core::local_git::LocalRemoteRepair>,
490    #[serde(default)]
491    pub dirty_repositories: Vec<String>,
492    #[serde(default)]
493    pub remote_repositories: Vec<PreflightRepository>,
494    pub local_changes_excluded: bool,
495}
496
497/// What a phone asks about, or stores against, its own identity.
498///
499/// These travel on their own channel rather than as actions, for the reason a
500/// read receipt does: they are frequent, they start nothing, and routing them
501/// through the action pipeline would consume the session's single action slot
502/// and reload the controller on every keystroke.
503#[derive(Debug)]
504pub enum ClientStateRequest {
505    Read {
506        client_id: String,
507        session_id: String,
508        reply: tokio::sync::oneshot::Sender<Result<ViewerClientState, String>>,
509    },
510    SaveDraft {
511        client_id: String,
512        session_id: String,
513        draft: String,
514        reply: tokio::sync::oneshot::Sender<Result<(), String>>,
515    },
516    MarkWorkspaceRead {
517        client_id: String,
518        workspace_id: String,
519        reply: tokio::sync::oneshot::Sender<Result<(), String>>,
520    },
521    History {
522        session_id: String,
523        query: String,
524        scope: String,
525        reply: tokio::sync::oneshot::Sender<Result<ViewerPromptHistory, String>>,
526    },
527}
528
529#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
530#[serde(deny_unknown_fields)]
531pub struct ViewerClientState {
532    pub draft: String,
533    pub through_event_ordinal: u64,
534}
535
536#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
537#[serde(deny_unknown_fields)]
538pub struct ViewerPromptHistory {
539    pub entries: Vec<String>,
540    /// Whether the search stopped before it ran out of history, so a phone can
541    /// say the answer is partial rather than presenting it as complete.
542    pub truncated: bool,
543}
544
545#[derive(Debug)]
546pub struct ReadReceiptRequest {
547    pub client_id: String,
548    pub session_id: String,
549    pub through: u64,
550    pub reply: tokio::sync::oneshot::Sender<Result<(), String>>,
551}
552
553/// A phone request to stop one currently projected background task.
554///
555/// This is intentionally not a [`ControllerAction`]. The request is already
556/// validated against the current operational snapshot by the HTTP handler,
557/// then the controller resolves the live session handle and waits for the
558/// provider acknowledgement in a supervised task.
559#[derive(Debug)]
560pub struct BackgroundTaskStopRequest {
561    pub session_id: String,
562    pub background_task_id: String,
563    pub reply: tokio::sync::oneshot::Sender<Result<(), BackgroundTaskStopFailure>>,
564}
565
566#[derive(Debug, Clone, Copy, PartialEq, Eq)]
567pub enum BackgroundTaskStopFailure {
568    /// The session manager could not resolve the live session handle.
569    SessionUnavailable,
570    /// The provider or relay rejected the stop request.
571    Provider,
572    /// The stop task itself failed before reaching the provider.
573    Internal,
574}
575
576#[cfg(test)]
577mod tests {
578    use super::ControllerAction;
579    use mj_core::state::SessionResourceAllocation;
580
581    #[test]
582    fn new_action_resource_allocation_round_trips_and_remains_optional() {
583        let with_allocation = serde_json::json!({
584            "action": "new",
585            "profile_id": "codex",
586            "bundle_id": "project",
587            "target_id": "podman",
588            "resource_allocation": {
589                "kind": "container",
590                "cpus": 4,
591                "memory_bytes": 8 * 1024 * 1024 * 1024_u64,
592            },
593        });
594        let action: ControllerAction = serde_json::from_value(with_allocation).unwrap();
595        assert!(matches!(
596            &action,
597            ControllerAction::New {
598                resource_allocation: Some(allocation),
599                ..
600            } if matches!(allocation.as_ref(), SessionResourceAllocation::Container {
601                cpus: 4,
602                memory_bytes: 8_589_934_592,
603            })
604        ));
605        let encoded = serde_json::to_value(&action).unwrap();
606        assert_eq!(
607            serde_json::from_value::<ControllerAction>(encoded).unwrap(),
608            action
609        );
610
611        let older_client: ControllerAction = serde_json::from_value(serde_json::json!({
612            "action": "new",
613            "profile_id": "codex",
614            "bundle_id": "project",
615            "target_id": "podman",
616        }))
617        .unwrap();
618        assert!(matches!(
619            &older_client,
620            ControllerAction::New {
621                resource_allocation: None,
622                ..
623            }
624        ));
625        assert!(
626            serde_json::to_value(older_client)
627                .unwrap()
628                .get("resource_allocation")
629                .is_none()
630        );
631    }
632}