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