Skip to main content

mj_controller/server/api/
subagent_backend.rs

1use super::*;
2
3/// Everything the API needs from the daemon: live session actors, the durable
4/// projection, and the target-side git operations.
5///
6/// The daemon's implementation lives in `server_runtime::api`; route tests
7/// supply a fake.
8pub trait SubagentBackend: Send + Sync {
9    /// Cheap token for this session's durable wait inputs. Backends without
10    /// a committed publication leave filtering to the detailed observation.
11    fn wait_revision(&self, _session_id: &str) -> AnyResult<Option<u64>> {
12        Ok(None)
13    }
14
15    /// Confirm before session admission that the configured bundle selects at
16    /// most one GitHub App installation.
17    fn validate_github_bundle(
18        &self,
19        _bundle_id: String,
20    ) -> BoxFuture<'_, Result<(), crate::controller::GithubBundleSelectionError>> {
21        Box::pin(async { Ok(()) })
22    }
23
24    /// Return one currently valid token for an owner installation or a set of
25    /// repositories that resolve to one installation.
26    fn github_token(
27        &self,
28        _owner: Option<String>,
29        _repositories: Vec<(String, String)>,
30    ) -> BoxFuture<'_, AnyResult<String>> {
31        Box::pin(async {
32            anyhow::bail!(
33                "GitHub App credentials are not configured; set [github.app] in config.toml"
34            )
35        })
36    }
37
38    fn transcript_history(
39        &self,
40        session_id: String,
41        before: Option<mj_core::storage::TranscriptCursor>,
42    ) -> BoxFuture<'_, AnyResult<mj_core::storage::TranscriptHistoryPage>> {
43        Box::pin(async move {
44            tokio::task::spawn_blocking(move || {
45                crate::database::load_transcript_history(&session_id, before.as_ref(), 128)?
46                    .context("session history is unavailable")
47            })
48            .await?
49        })
50    }
51
52    fn native_agent_history(
53        &self,
54        owner: String,
55        child: String,
56        before: Option<(u64, String)>,
57    ) -> BoxFuture<'_, AnyResult<mj_core::native_agent::NativeAgentHistoryPage>> {
58        Box::pin(async move {
59            tokio::task::spawn_blocking(move || {
60                crate::database::native_agent_history(&owner, &child, before)
61            })
62            .await?
63        })
64    }
65    fn events(
66        &self,
67        filter: crate::database::ApiEventFilter,
68        after_seq: Option<u64>,
69    ) -> BoxFuture<'_, AnyResult<crate::database::ApiEventPage>> {
70        events::load_events(filter, after_seq)
71    }
72
73    fn profile_config(
74        &self,
75        profile: String,
76        model: Option<String>,
77        refresh: bool,
78    ) -> BoxFuture<'_, AnyResult<mj_core::worker_launch::ProfileConfig>> {
79        Box::pin(crate::controller::profile_config::discover(
80            profile, model, refresh,
81        ))
82    }
83    /// The profiles a parent running on `parent_profile` may start a sub-agent
84    /// on, each with its discovered choices and remaining quota. It waits for
85    /// a discovery the daemon has not finished; a profile whose discovery
86    /// fails is reported as unavailable rather than failing the whole answer.
87    fn subagent_candidates(
88        &self,
89        _parent_profile: String,
90    ) -> BoxFuture<'_, AnyResult<SubagentCandidates>> {
91        Box::pin(async { anyhow::bail!("sub-agent profiles are unavailable") })
92    }
93    /// Every enabled profile usable for a user-created session, with its
94    /// discovered choices and remaining quota. This is independent of the
95    /// profile set allowed for sub-agent use.
96    fn session_profile_candidates(&self) -> BoxFuture<'_, AnyResult<SubagentCandidates>> {
97        Box::pin(async { anyhow::bail!("session profile selection is unavailable") })
98    }
99    fn start_subagent(
100        &self,
101        _request: crate::controller::RegisterSubagentRequest,
102    ) -> BoxFuture<'_, AnyResult<mj_core::subagent::SubagentRecord>> {
103        Box::pin(async { anyhow::bail!("sub-agent creation is unavailable") })
104    }
105    fn list_subagents(
106        &self,
107        parent_session_id: String,
108    ) -> BoxFuture<'_, AnyResult<Vec<mj_core::subagent::SubagentRecord>>> {
109        Box::pin(async move {
110            tokio::task::spawn_blocking(move || crate::database::list_subagents(&parent_session_id))
111                .await?
112        })
113    }
114    /// Whether a sub-agent child has handed back its report for its parent's
115    /// newest task; see [`crate::controller::subagent_has_handed_back`].
116    fn subagent_handed_back(&self, child_session_id: String) -> BoxFuture<'_, AnyResult<bool>> {
117        Box::pin(async move {
118            tokio::task::spawn_blocking(move || {
119                crate::controller::subagent_has_handed_back(&child_session_id)
120            })
121            .await?
122        })
123    }
124    fn read_context_file(
125        &self,
126        session_id: String,
127        path: PathBuf,
128    ) -> BoxFuture<'_, std::result::Result<Vec<u8>, ExportError>> {
129        self.read_file(session_id, path)
130    }
131    /// The workspaces the store holds, in the order the terminal's tabs and the
132    /// viewer's list show them.
133    fn list_workspaces(
134        &self,
135    ) -> BoxFuture<'_, AnyResult<Vec<mj_core::workspace::WorkspaceRecord>>> {
136        Box::pin(async { tokio::task::spawn_blocking(crate::database::list_workspaces).await? })
137    }
138    /// The workspace with this name, creating it when the store holds none.
139    ///
140    /// This is the daemon's own `CreateWorkspace` operation: create-or-get, so
141    /// two callers that both saw an empty list attach to the same normalized
142    /// name instead of one of them meeting a SQLite conflict. The daemon
143    /// overrides it to republish the list afterwards.
144    fn create_workspace(
145        &self,
146        name: String,
147    ) -> BoxFuture<'_, AnyResult<mj_core::workspace::WorkspaceRecord>> {
148        Box::pin(async move {
149            tokio::task::spawn_blocking(move || crate::database::create_or_get_workspace(&name))
150                .await?
151        })
152    }
153    fn set_config(
154        &self,
155        session_id: String,
156        key: String,
157        value: String,
158    ) -> BoxFuture<'_, AnyResult<()>> {
159        Box::pin(async move {
160            self.session_handle(session_id)
161                .await?
162                .ok_or_else(|| anyhow::anyhow!("session has no live actor"))?
163                .set_config(key, value)
164                .await
165        })
166    }
167    fn cancel_start(&self, _session_id: String) -> BoxFuture<'_, AnyResult<()>> {
168        Box::pin(async { Ok(()) })
169    }
170    /// The live actor for a session, or `None` when none holds it.
171    fn session_handle(&self, session_id: String)
172    -> BoxFuture<'_, AnyResult<Option<SessionHandle>>>;
173
174    /// Submit a prompt, returning its relay acceptance ordinal.
175    fn prompt(&self, session_id: String, text: String) -> BoxFuture<'_, AnyResult<u64>>;
176
177    /// Deliver a session message using the daemon's shared authorization and
178    /// mailbox/turn routing operation.
179    fn deliver_message(
180        &self,
181        _sender_session_id: Option<String>,
182        _target_session_id: String,
183        _text: String,
184        _request_id: String,
185        _created_at_ms: i64,
186    ) -> BoxFuture<'_, AnyResult<SessionMessageResponse>> {
187        Box::pin(async { anyhow::bail!("session messaging is unavailable") })
188    }
189
190    /// Submit under a producer-owned identity, preserving the receipt on retries.
191    fn prompt_with_id(
192        &self,
193        session_id: String,
194        text: String,
195        command_id: Option<String>,
196    ) -> BoxFuture<'_, AnyResult<u64>> {
197        let Some(command_id) = command_id else {
198            return self.prompt(session_id, text);
199        };
200        Box::pin(async move {
201            let handle = self
202                .session_handle(session_id)
203                .await?
204                .ok_or_else(|| anyhow::anyhow!("session has no live worker"))?;
205            handle
206                .submit(
207                    command_id,
208                    mj_core::relay::RelayCommand::Prompt {
209                        prompt: vec![agent_client_protocol::schema::v1::ContentBlock::Text(
210                            agent_client_protocol::schema::v1::TextContent::new(text),
211                        )],
212                    },
213                )
214                .await
215        })
216    }
217
218    /// Durable turn state for a session with no live actor.
219    fn turn_state(&self, session_id: String) -> BoxFuture<'_, AnyResult<Option<TurnState>>>;
220
221    /// A sub-agent child's recorded report, and whether it has the handback
222    /// tool. `None` for a session that is not a Mjolnir sub-agent.
223    fn subagent_report(
224        &self,
225        _session_id: String,
226    ) -> BoxFuture<'_, AnyResult<Option<(bool, mj_core::subagent::SubagentReport)>>> {
227        Box::pin(async { Ok(None) })
228    }
229
230    /// Summarize the turn that covers these transcript positions.
231    fn turn_summary(
232        &self,
233        session_id: String,
234        turn: TurnSpan,
235    ) -> BoxFuture<'_, AnyResult<TurnSummary>>;
236
237    /// Apply model, effort, and the first prompt once a new session is ready.
238    fn start_followup(
239        &self,
240        session_id: String,
241        followup: StartFollowup,
242    ) -> BoxFuture<'_, AnyResult<()>>;
243
244    /// How far a created session's follow-up has got.
245    fn start_status(&self, session_id: String) -> BoxFuture<'_, AnyResult<Option<StartStatus>>>;
246
247    /// A page of transcript items after `after_seq`.
248    fn transcript(
249        &self,
250        session_id: String,
251        after_seq: u64,
252        limit: usize,
253        roles: Vec<mj_core::transcript::TranscriptRole>,
254        finished_only: bool,
255    ) -> BoxFuture<'_, AnyResult<Option<TranscriptPage>>>;
256
257    fn usage(
258        &self,
259        session_id: String,
260        after_seq: u64,
261        limit: usize,
262    ) -> BoxFuture<'_, AnyResult<Option<crate::database::UsagePage>>> {
263        Box::pin(async move {
264            tokio::task::spawn_blocking(move || {
265                crate::database::load_session_usage(&session_id, after_seq, limit)
266            })
267            .await?
268        })
269    }
270
271    fn usage_tree(
272        &self,
273        parent: String,
274    ) -> BoxFuture<'_, AnyResult<Option<crate::database::UsageTree>>> {
275        Box::pin(async move {
276            tokio::task::spawn_blocking(move || crate::database::load_usage_tree(&parent)).await?
277        })
278    }
279
280    /// A unified diff of the session's work.
281    fn diff(
282        &self,
283        session_id: String,
284        options: DiffOptions,
285    ) -> BoxFuture<'_, Result<String, ExportError>>;
286
287    /// One file from the session's workspace.
288    fn read_file(
289        &self,
290        session_id: String,
291        path: PathBuf,
292    ) -> BoxFuture<'_, Result<Vec<u8>, ExportError>>;
293
294    fn write_file(
295        &self,
296        _session_id: String,
297        _path: PathBuf,
298        _bytes: Vec<u8>,
299        _overwrite: bool,
300    ) -> BoxFuture<'_, Result<(), ExportError>> {
301        Box::pin(async { Err(ExportError::Refused("file injection is unavailable".into())) })
302    }
303
304    /// Push the session's branch to its repository's default remote.
305    fn push_branch(
306        &self,
307        session_id: String,
308        branch: String,
309    ) -> BoxFuture<'_, Result<PushedBranch, ExportError>>;
310
311    /// A git bundle of the session's committed work.
312    fn bundle(&self, session_id: String) -> BoxFuture<'_, Result<BundleExport, ExportError>>;
313
314    /// Whether the index was synced recently enough that a query need not ask
315    /// for one.
316    fn wiki_sync_is_stale(&self) -> bool {
317        false
318    }
319
320    /// Ask for a background sync. It is never waited for: a query answers from
321    /// what the index holds now.
322    fn wiki_request_sync(&self) {}
323
324    fn wiki_search(
325        &self,
326        _query: String,
327        _limit: usize,
328    ) -> BoxFuture<'_, AnyResult<mj_client::daemon::WikiSearchPage>> {
329        Box::pin(async { anyhow::bail!("SessionWiki search is unavailable") })
330    }
331
332    /// `None` when the index holds no session with that id.
333    fn wiki_brief(
334        &self,
335        _wiki_id: String,
336        _max_chars: usize,
337    ) -> BoxFuture<'_, AnyResult<Option<String>>> {
338        Box::pin(async { anyhow::bail!("SessionWiki briefings are unavailable") })
339    }
340
341    /// The passages of one indexed session that match a query. `None` when the
342    /// index holds no session with that id.
343    fn wiki_hits(
344        &self,
345        _wiki_id: String,
346        _query: String,
347        _context_messages: usize,
348        _per_message_chars: usize,
349    ) -> BoxFuture<'_, AnyResult<Option<mj_client::daemon::WikiHitTranscript>>> {
350        Box::pin(async { anyhow::bail!("SessionWiki transcript hits are unavailable") })
351    }
352
353    /// What one indexed session is, and what continuing it would mean.
354    /// `None` when the index holds no session with that id.
355    fn wiki_session(
356        &self,
357        _wiki_id: String,
358    ) -> BoxFuture<'_, AnyResult<Option<mj_client::daemon::WikiSessionInfo>>> {
359        Box::pin(async { anyhow::bail!("SessionWiki lookups are unavailable") })
360    }
361
362    /// Start a session from an archived transcript, answering with its id, or
363    /// `None` when the index holds no session with that id.
364    fn wiki_restore(
365        &self,
366        _request: mj_client::daemon::WikiRestoreRequest,
367    ) -> BoxFuture<'_, AnyResult<Option<String>>> {
368        Box::pin(async { anyhow::bail!("SessionWiki restore is unavailable") })
369    }
370}
371
372pub(super) fn backend(state: &ServerState) -> Result<&Arc<dyn SubagentBackend>, ApiFailure> {
373    state
374        .subagent
375        .as_ref()
376        .ok_or_else(|| ApiFailure::unavailable("this server has no subagent backend installed"))
377}
378
379// ---------------------------------------------------------------------------
380// Wait resolution
381// ---------------------------------------------------------------------------