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