mj_controller/worker_client/
relay.rs1use super::*;
2
3impl RelayClient {
4 pub async fn jev_decisions(
5 &mut self,
6 decision_id: Option<String>,
7 ) -> Result<mj_core::jev::DecisionPage> {
8 let request = RelayRequest::JevDecisions { decision_id };
9 anyhow::ensure!(
10 request.supported_at(self.protocol_version),
11 "This worker does not support Jev decision details; restart it with a newer mj build."
12 );
13 match self.call(request).await? {
14 RelayResponsePayload::JevDecisions(page) => Ok(page),
15 other => bail!("unexpected Jev diagnostics response: {other:?}"),
16 }
17 }
18
19 pub async fn history_requests(&mut self) -> Result<Vec<mj_core::history::HistoryRequest>> {
20 if !RelayRequest::HistoryRequests.supported_at(self.protocol_version) {
21 return Ok(Vec::new());
22 }
23 match self.call(RelayRequest::HistoryRequests).await? {
24 RelayResponsePayload::HistoryRequests { requests } => Ok(requests),
25 _ => bail!("relay returned an unexpected history queue response"),
26 }
27 }
28
29 pub async fn complete_history_request(
30 &mut self,
31 result: mj_core::history::HistoryResult,
32 ) -> Result<()> {
33 match self
34 .call(RelayRequest::CompleteHistoryRequest { result })
35 .await?
36 {
37 RelayResponsePayload::HistoryRequestCompleted => Ok(()),
38 _ => bail!("relay returned an unexpected history completion response"),
39 }
40 }
41
42 pub fn session_id(&self) -> &str {
43 &self.session_id
44 }
45
46 pub fn supports_project_memory_sync(&self) -> bool {
47 RelayRequest::ProjectMemorySnapshot.supported_at(self.protocol_version)
48 }
49
50 pub fn relay_version(&self) -> &str {
51 &self.relay_version
52 }
53
54 pub fn worker_build(&self) -> Option<&str> {
58 self.worker_build.as_deref()
59 }
60
61 pub fn protocol_version(&self) -> u32 {
62 self.protocol_version
63 }
64
65 pub fn latest_ordinal(&self) -> u64 {
66 self.latest_ordinal
67 }
68
69 pub fn latest_digest(&self) -> &str {
70 &self.latest_digest
71 }
72
73 pub async fn attach(
74 &mut self,
75 after_ordinal: u64,
76 after_digest: impl Into<String>,
77 ) -> Result<RelayAttachment> {
78 let after_digest = after_digest.into();
79 match self
80 .call_with_timeout(
81 RelayRequest::Attach {
82 after_ordinal,
83 after_digest: after_digest.clone(),
84 },
85 RELAY_HISTORY_TIMEOUT,
86 )
87 .await?
88 {
89 RelayResponsePayload::Attached {
90 mut state,
91 events,
92 through_ordinal,
93 through_digest,
94 } => {
95 let mut cursor = RelayCursor {
96 ordinal: after_ordinal,
97 digest: after_digest,
98 };
99 for event in &events {
100 validate_relay_event(cursor.ordinal, &cursor.digest, event)
101 .context("verify relay attachment event chain")?;
102 cursor.ordinal = event.ordinal;
103 cursor.digest.clone_from(&event.digest);
104 }
105 if cursor.ordinal != through_ordinal || cursor.digest != through_digest {
106 bail!("relay attachment frontier does not match its event chain");
107 }
108 state.relay_protocol_version = Some(self.protocol_version);
109 self.latest_ordinal = state.latest_ordinal;
110 self.latest_digest = state.latest_digest.clone();
111 Ok(RelayAttachment {
112 state,
113 events,
114 through_ordinal,
115 through_digest,
116 })
117 }
118 _ => bail!("relay returned an unexpected attach response"),
119 }
120 }
121
122 pub async fn begin_catch_up(
127 &mut self,
128 after_ordinal: u64,
129 after_digest: impl Into<String>,
130 ) -> Result<RelayCatchUp> {
131 let after_digest = after_digest.into();
132 let first = self.attach(after_ordinal, after_digest.clone()).await?;
133 let frontier = RelayCursor {
134 ordinal: first.state.latest_ordinal,
135 digest: first.state.latest_digest.clone(),
136 };
137 let previous = RelayCursor {
138 ordinal: after_ordinal,
139 digest: after_digest,
140 };
141 let state = first.state.clone();
142 let first_page = clip_catch_up_page(first, &previous, &frontier)?;
143 Ok(RelayCatchUp {
144 state,
145 frontier,
146 first_page,
147 })
148 }
149
150 pub async fn next_catch_up_page(
154 &mut self,
155 previous: &RelayCursor,
156 frontier: &RelayCursor,
157 ) -> Result<RelayEventPage> {
158 if previous.ordinal >= frontier.ordinal {
159 bail!("relay catch-up is already at its fixed frontier");
160 }
161 let attachment = self
162 .attach(previous.ordinal, previous.digest.clone())
163 .await?;
164 clip_catch_up_page(attachment, previous, frontier)
165 }
166
167 pub async fn acknowledge(
168 &mut self,
169 through_ordinal: u64,
170 through_digest: impl Into<String>,
171 ) -> Result<RelayCursor> {
172 match self
173 .call_with_timeout(
174 RelayRequest::Acknowledge {
175 through_ordinal,
176 through_digest: through_digest.into(),
177 },
178 RELAY_ACKNOWLEDGE_TIMEOUT,
179 )
180 .await?
181 {
182 RelayResponsePayload::Acknowledged {
183 through_ordinal,
184 through_digest,
185 } => Ok(RelayCursor {
186 ordinal: through_ordinal,
187 digest: through_digest,
188 }),
189 _ => bail!("relay returned an unexpected acknowledgement response"),
190 }
191 }
192
193 pub async fn status(&mut self) -> Result<RelayOperationalState> {
194 match self.call(RelayRequest::Status).await? {
195 RelayResponsePayload::Status(mut status) => {
196 status.relay_protocol_version = Some(self.protocol_version);
197 self.latest_ordinal = status.latest_ordinal;
198 self.latest_digest = status.latest_digest.clone();
199 Ok(status)
200 }
201 _ => bail!("relay returned an unexpected status response"),
202 }
203 }
204
205 pub async fn credential_state(&mut self) -> Result<CredentialSnapshot> {
208 credential_snapshot(self.call(RelayRequest::CredentialState).await?)
209 }
210
211 pub async fn read_credentials(&mut self) -> Result<Vec<u8>> {
214 match self.call(RelayRequest::ReadCredentials).await? {
215 RelayResponsePayload::Credentials { data } => BASE64
216 .decode(data.as_bytes())
217 .context("decode relay credential payload"),
218 _ => bail!("relay returned an unexpected credential response"),
219 }
220 }
221
222 pub async fn install_credentials(&mut self, bytes: &[u8]) -> Result<CredentialSnapshot> {
225 credential_snapshot(
226 self.call(RelayRequest::InstallCredentials {
227 data: BASE64.encode(bytes),
228 })
229 .await?,
230 )
231 }
232
233 pub async fn github_token_state(
234 &mut self,
235 ) -> Result<mj_core::credentials::GithubTokenSnapshot> {
236 github_token_snapshot(self.call(RelayRequest::GithubTokenState).await?)
237 }
238
239 pub async fn install_github_token(
240 &mut self,
241 token: &str,
242 ) -> Result<mj_core::credentials::GithubTokenSnapshot> {
243 github_token_snapshot(
244 self.call(RelayRequest::InstallGithubToken {
245 data: BASE64.encode(token.as_bytes()),
246 })
247 .await?,
248 )
249 }
250
251 pub async fn remove_github_token(
252 &mut self,
253 ) -> Result<mj_core::credentials::GithubTokenSnapshot> {
254 github_token_snapshot(self.call(RelayRequest::RemoveGithubToken).await?)
255 }
256
257 pub async fn skills_state(&mut self) -> Result<mj_core::skills::SkillsSyncState> {
260 skills_sync_state(self.call(RelayRequest::SkillsState).await?)
261 }
262
263 pub async fn install_prompt_context(&mut self, text: String) -> Result<()> {
266 let request = RelayRequest::InstallPromptContext { text };
267 match self.call(request).await? {
268 RelayResponsePayload::PromptContextInstalled => Ok(()),
269 _ => bail!("relay returned an unexpected prompt-context response"),
270 }
271 }
272
273 pub async fn project_memory_snapshot(
274 &mut self,
275 ) -> Result<(
276 mj_core::project_memory::ProjectMemorySnapshot,
277 mj_core::project_memory::ProjectMemorySnapshot,
278 )> {
279 let request = RelayRequest::ProjectMemorySnapshot;
280 match self.call(request).await? {
281 RelayResponsePayload::ProjectMemorySnapshot { baseline, replica } => {
282 Ok((baseline, replica))
283 }
284 _ => bail!("relay returned an unexpected project-memory response"),
285 }
286 }
287
288 pub async fn install_project_memory_snapshot(
289 &mut self,
290 snapshot: mj_core::project_memory::ProjectMemorySnapshot,
291 ) -> Result<()> {
292 let request = RelayRequest::InstallProjectMemorySnapshot { snapshot };
293 match self.call(request).await? {
294 RelayResponsePayload::ProjectMemorySnapshotInstalled => Ok(()),
295 _ => bail!("relay returned an unexpected project-memory install response"),
296 }
297 }
298
299 pub async fn install_skills(
303 &mut self,
304 archive_bytes: &[u8],
305 ) -> Result<mj_core::skills::SkillsSyncState> {
306 skills_sync_state(
307 self.call(RelayRequest::InstallSkills {
308 data: BASE64.encode(archive_bytes),
309 })
310 .await?,
311 )
312 }
313
314 pub async fn ensure_attachment(
316 &mut self,
317 reference: &mj_core::attachment::AttachmentRef,
318 ) -> Result<()> {
319 match self
320 .call(RelayRequest::AttachmentPresent {
321 reference: reference.clone(),
322 })
323 .await?
324 {
325 RelayResponsePayload::AttachmentPresent { present: true } => return Ok(()),
326 RelayResponsePayload::AttachmentPresent { present: false } => {}
327 _ => bail!("unexpected image presence response"),
328 }
329 let store = mj_core::attachment::AttachmentStore::controller(&self.session_id)?;
330 let reference_copy = reference.clone();
331 let bytes = tokio::task::spawn_blocking(move || store.read(&reference_copy))
332 .await
333 .context("image loading task failed")??;
334 match self
335 .call(RelayRequest::InstallAttachment {
336 reference: reference.clone(),
337 data: BASE64.encode(bytes),
338 })
339 .await?
340 {
341 RelayResponsePayload::AttachmentInstalled => Ok(()),
342 _ => bail!("unexpected image upload response"),
343 }
344 }
345
346 pub async fn cache_attachment(
348 &mut self,
349 reference: &mj_core::attachment::AttachmentRef,
350 ) -> Result<()> {
351 let store = mj_core::attachment::AttachmentStore::controller(&self.session_id)?;
352 let local = store.clone();
353 let reference_copy = reference.clone();
354 if tokio::task::spawn_blocking(move || local.contains(&reference_copy))
355 .await
356 .context("image lookup task failed")??
357 {
358 return Ok(());
359 }
360 let RelayResponsePayload::AttachmentData { data } = self
361 .call(RelayRequest::ReadAttachment {
362 reference: reference.clone(),
363 })
364 .await?
365 else {
366 bail!("unexpected image download response")
367 };
368 let reference = reference.clone();
369 tokio::task::spawn_blocking(move || {
370 anyhow::ensure!(
371 data.len() <= mj_core::attachment::MAX_IMAGE_BYTES.div_ceil(3) * 4,
372 "image download is too large"
373 );
374 store.install(&reference, &BASE64.decode(data)?)
375 })
376 .await
377 .context("image caching task failed")?
378 }
379
380 pub async fn submit(
381 &mut self,
382 command_id: impl Into<String>,
383 command: RelayCommand,
384 ) -> Result<u64> {
385 let command_id = command_id.into();
386 if let RelayCommand::Prompt { prompt } = &command {
387 for reference in mj_core::attachment::references(prompt)? {
388 self.ensure_attachment(&reference).await?;
389 }
390 }
391 match self
392 .call(RelayRequest::Submit {
393 command_id: command_id.clone(),
394 command,
395 })
396 .await?
397 {
398 RelayResponsePayload::Accepted {
399 command_id: accepted_id,
400 ordinal,
401 } if accepted_id == command_id => Ok(ordinal),
402 RelayResponsePayload::Accepted {
403 command_id: accepted_id,
404 ..
405 } => bail!("relay accepted command under ID {accepted_id}, expected {command_id}"),
406 _ => bail!("relay returned an unexpected command response"),
407 }
408 }
409}