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