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