Skip to main content

mj_controller/worker_client/
relay.rs

1use 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    /// Content address of the executable serving this connection, or `None`
40    /// from a worker too old to report one. A controller reads `None` as
41    /// outdated: it predates the field, so it predates this controller.
42    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    /// Start a bounded catch-up by capturing the relay frontier before the
108    /// caller applies anything. Callers persist `first_page`, request further
109    /// pages with [`Self::next_catch_up_page`], and may acknowledge the fixed
110    /// frontier after all of those pages are durable.
111    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    /// Fetch the next bounded page without chasing events that arrived after
136    /// `frontier` was captured. A response may contain such newer events; the
137    /// returned page is clipped at the exact ordinal-and-digest frontier.
138    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    /// Return the fingerprint and freshness of this session's harness
191    /// credentials without exposing the credential bytes.
192    pub async fn credential_state(&mut self) -> Result<CredentialSnapshot> {
193        credential_snapshot(self.call(RelayRequest::CredentialState).await?)
194    }
195
196    /// Read this session's credential file. Callers must keep these bytes out
197    /// of durable relay observations, logs, and archives.
198    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    /// Install credentials into the harness home fixed by this session's
208    /// launch config.
209    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    /// Return the fingerprint of this session's synced skills trees without
243    /// transferring the tree itself.
244    pub async fn skills_state(&mut self) -> Result<mj_core::skills::SkillsSyncState> {
245        skills_sync_state(self.call(RelayRequest::SkillsState).await?)
246    }
247
248    /// Install background text that only the target harness sees, prepended
249    /// to the next real prompt without creating a synthetic transcript turn.
250    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    /// Replace this session's synced skills trees with an encoded
285    /// `skills::SkillsArchive`. The destination directories are fixed by
286    /// the session's launch config and the harness skills whitelist.
287    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    /// Copy a verified controller blob to this session before admitting its reference.
300    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    /// Recover the local copy needed for queue editing and resubmission.
332    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}