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