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    fn transcript_history(
10        &self,
11        session_id: String,
12        before: Option<mj_core::storage::TranscriptCursor>,
13    ) -> BoxFuture<'_, AnyResult<mj_core::storage::TranscriptHistoryPage>> {
14        Box::pin(async move {
15            tokio::task::spawn_blocking(move || {
16                crate::database::load_transcript_history(&session_id, before.as_ref(), 128)?
17                    .context("session history is unavailable")
18            })
19            .await?
20        })
21    }
22
23    fn native_agent_history(
24        &self,
25        owner: String,
26        child: String,
27        before: Option<(u64, String)>,
28    ) -> BoxFuture<'_, AnyResult<mj_core::native_agent::NativeAgentHistoryPage>> {
29        Box::pin(async move {
30            tokio::task::spawn_blocking(move || {
31                crate::database::native_agent_history(&owner, &child, before)
32            })
33            .await?
34        })
35    }
36    fn events(
37        &self,
38        filter: crate::database::ApiEventFilter,
39        after_seq: Option<u64>,
40    ) -> BoxFuture<'_, AnyResult<crate::database::ApiEventPage>> {
41        events::load_events(filter, after_seq)
42    }
43
44    fn profile_config(
45        &self,
46        profile: String,
47        model: Option<String>,
48        refresh: bool,
49    ) -> BoxFuture<'_, AnyResult<mj_core::worker_launch::ProfileConfig>> {
50        Box::pin(crate::controller::profile_config::discover(
51            profile, model, refresh,
52        ))
53    }
54    /// The profiles a parent running on `parent_profile` may start a sub-agent
55    /// on, each with its discovered choices and remaining quota. It waits for
56    /// a discovery the daemon has not finished; a profile whose discovery
57    /// fails is reported as unavailable rather than failing the whole answer.
58    fn subagent_candidates(
59        &self,
60        _parent_profile: String,
61    ) -> BoxFuture<'_, AnyResult<SubagentCandidates>> {
62        Box::pin(async { anyhow::bail!("sub-agent profiles are unavailable") })
63    }
64    fn start_subagent(
65        &self,
66        _request: crate::controller::RegisterSubagentRequest,
67    ) -> BoxFuture<'_, AnyResult<mj_core::subagent::SubagentRecord>> {
68        Box::pin(async { anyhow::bail!("sub-agent creation is unavailable") })
69    }
70    fn list_subagents(
71        &self,
72        parent_session_id: String,
73    ) -> BoxFuture<'_, AnyResult<Vec<mj_core::subagent::SubagentRecord>>> {
74        Box::pin(async move {
75            tokio::task::spawn_blocking(move || crate::database::list_subagents(&parent_session_id))
76                .await?
77        })
78    }
79    /// Whether a sub-agent child has handed back its report for its parent's
80    /// newest task; see [`crate::controller::subagent_has_handed_back`].
81    fn subagent_handed_back(&self, child_session_id: String) -> BoxFuture<'_, AnyResult<bool>> {
82        Box::pin(async move {
83            tokio::task::spawn_blocking(move || {
84                crate::controller::subagent_has_handed_back(&child_session_id)
85            })
86            .await?
87        })
88    }
89    fn read_context_file(
90        &self,
91        session_id: String,
92        path: PathBuf,
93    ) -> BoxFuture<'_, std::result::Result<Vec<u8>, ExportError>> {
94        self.read_file(session_id, path)
95    }
96    /// The workspaces the store holds, in the order the terminal's tabs and the
97    /// viewer's list show them.
98    fn list_workspaces(
99        &self,
100    ) -> BoxFuture<'_, AnyResult<Vec<mj_core::workspace::WorkspaceRecord>>> {
101        Box::pin(async { tokio::task::spawn_blocking(crate::database::list_workspaces).await? })
102    }
103    /// The workspace with this name, creating it when the store holds none.
104    ///
105    /// This is the daemon's own `CreateWorkspace` operation: create-or-get, so
106    /// two callers that both saw an empty list attach to the same normalized
107    /// name instead of one of them meeting a SQLite conflict. The daemon
108    /// overrides it to republish the list afterwards.
109    fn create_workspace(
110        &self,
111        name: String,
112    ) -> BoxFuture<'_, AnyResult<mj_core::workspace::WorkspaceRecord>> {
113        Box::pin(async move {
114            tokio::task::spawn_blocking(move || crate::database::create_or_get_workspace(&name))
115                .await?
116        })
117    }
118    fn set_config(
119        &self,
120        session_id: String,
121        key: String,
122        value: String,
123    ) -> BoxFuture<'_, AnyResult<()>> {
124        Box::pin(async move {
125            self.session_handle(session_id)
126                .await?
127                .ok_or_else(|| anyhow::anyhow!("session has no live actor"))?
128                .set_config(key, value)
129                .await
130        })
131    }
132    fn cancel_start(&self, _session_id: String) -> BoxFuture<'_, AnyResult<()>> {
133        Box::pin(async { Ok(()) })
134    }
135    /// The live actor for a session, or `None` when none holds it.
136    fn session_handle(&self, session_id: String)
137    -> BoxFuture<'_, AnyResult<Option<SessionHandle>>>;
138
139    /// Submit a prompt, returning its relay acceptance ordinal.
140    fn prompt(&self, session_id: String, text: String) -> BoxFuture<'_, AnyResult<u64>>;
141
142    /// Durable turn state for a session with no live actor.
143    fn turn_state(&self, session_id: String) -> BoxFuture<'_, AnyResult<Option<TurnState>>>;
144
145    /// A sub-agent child's recorded report, and whether it has the handback
146    /// tool. `None` for a session that is not a Mjolnir sub-agent.
147    fn subagent_report(
148        &self,
149        _session_id: String,
150    ) -> BoxFuture<'_, AnyResult<Option<(bool, mj_core::subagent::SubagentReport)>>> {
151        Box::pin(async { Ok(None) })
152    }
153
154    /// Summarize the turn that covers these transcript positions.
155    fn turn_summary(
156        &self,
157        session_id: String,
158        turn: TurnSpan,
159    ) -> BoxFuture<'_, AnyResult<TurnSummary>>;
160
161    /// Apply model, effort, and the first prompt once a new session is ready.
162    fn start_followup(
163        &self,
164        session_id: String,
165        followup: StartFollowup,
166    ) -> BoxFuture<'_, AnyResult<()>>;
167
168    /// How far a created session's follow-up has got.
169    fn start_status(&self, session_id: String) -> BoxFuture<'_, AnyResult<Option<StartStatus>>>;
170
171    /// A page of transcript items after `after_seq`.
172    fn transcript(
173        &self,
174        session_id: String,
175        after_seq: u64,
176        limit: usize,
177        role: Option<mj_core::transcript::TranscriptRole>,
178    ) -> BoxFuture<'_, AnyResult<Option<TranscriptPage>>>;
179
180    fn usage(
181        &self,
182        session_id: String,
183        after_seq: u64,
184        limit: usize,
185    ) -> BoxFuture<'_, AnyResult<Option<crate::database::UsagePage>>> {
186        Box::pin(async move {
187            tokio::task::spawn_blocking(move || {
188                crate::database::load_session_usage(&session_id, after_seq, limit)
189            })
190            .await?
191        })
192    }
193
194    fn usage_tree(
195        &self,
196        parent: String,
197    ) -> BoxFuture<'_, AnyResult<Option<crate::database::UsageTree>>> {
198        Box::pin(async move {
199            tokio::task::spawn_blocking(move || crate::database::load_usage_tree(&parent)).await?
200        })
201    }
202
203    /// A unified diff of the session's work.
204    fn diff(
205        &self,
206        session_id: String,
207        options: DiffOptions,
208    ) -> BoxFuture<'_, Result<String, ExportError>>;
209
210    /// One file from the session's workspace.
211    fn read_file(
212        &self,
213        session_id: String,
214        path: PathBuf,
215    ) -> BoxFuture<'_, Result<Vec<u8>, ExportError>>;
216
217    fn write_file(
218        &self,
219        _session_id: String,
220        _path: PathBuf,
221        _bytes: Vec<u8>,
222        _overwrite: bool,
223    ) -> BoxFuture<'_, Result<(), ExportError>> {
224        Box::pin(async { Err(ExportError::Refused("file injection is unavailable".into())) })
225    }
226
227    /// Push the session's branch to its repository's default remote.
228    fn push_branch(
229        &self,
230        session_id: String,
231        branch: String,
232    ) -> BoxFuture<'_, Result<PushedBranch, ExportError>>;
233
234    /// A git bundle of the session's committed work.
235    fn bundle(&self, session_id: String) -> BoxFuture<'_, Result<BundleExport, ExportError>>;
236
237    /// Whether the index was synced recently enough that a query need not ask
238    /// for one.
239    fn wiki_sync_is_stale(&self) -> bool {
240        false
241    }
242
243    /// Ask for a background sync. It is never waited for: a query answers from
244    /// what the index holds now.
245    fn wiki_request_sync(&self) {}
246
247    fn wiki_search(
248        &self,
249        _query: String,
250        _limit: usize,
251    ) -> BoxFuture<'_, AnyResult<mj_client::daemon::WikiSearchPage>> {
252        Box::pin(async { anyhow::bail!("SessionWiki search is unavailable") })
253    }
254
255    /// `None` when the index holds no session with that id.
256    fn wiki_brief(
257        &self,
258        _wiki_id: String,
259        _max_chars: usize,
260    ) -> BoxFuture<'_, AnyResult<Option<String>>> {
261        Box::pin(async { anyhow::bail!("SessionWiki briefings are unavailable") })
262    }
263
264    /// The passages of one indexed session that match a query. `None` when the
265    /// index holds no session with that id.
266    fn wiki_hits(
267        &self,
268        _wiki_id: String,
269        _query: String,
270        _context_messages: usize,
271        _per_message_chars: usize,
272    ) -> BoxFuture<'_, AnyResult<Option<mj_client::daemon::WikiHitTranscript>>> {
273        Box::pin(async { anyhow::bail!("SessionWiki transcript hits are unavailable") })
274    }
275
276    /// What one indexed session is, and what continuing it would mean.
277    /// `None` when the index holds no session with that id.
278    fn wiki_session(
279        &self,
280        _wiki_id: String,
281    ) -> BoxFuture<'_, AnyResult<Option<mj_client::daemon::WikiSessionInfo>>> {
282        Box::pin(async { anyhow::bail!("SessionWiki lookups are unavailable") })
283    }
284
285    /// Start a session from an archived transcript, answering with its id, or
286    /// `None` when the index holds no session with that id.
287    fn wiki_restore(
288        &self,
289        _request: mj_client::daemon::WikiRestoreRequest,
290    ) -> BoxFuture<'_, AnyResult<Option<String>>> {
291        Box::pin(async { anyhow::bail!("SessionWiki restore is unavailable") })
292    }
293}
294
295pub(super) fn backend(state: &ServerState) -> Result<&Arc<dyn SubagentBackend>, ApiFailure> {
296    state
297        .subagent
298        .as_ref()
299        .ok_or_else(|| ApiFailure::unavailable("this server has no subagent backend installed"))
300}
301
302// ---------------------------------------------------------------------------
303// Wait resolution
304// ---------------------------------------------------------------------------