Skip to main content

mj_controller/worker_client/
relay.rs

1use 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    /// Content address of the executable serving this connection, or `None`
55    /// from a worker too old to report one. A controller reads `None` as
56    /// outdated: it predates the field, so it predates this controller.
57    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    /// Start a bounded catch-up by capturing the relay frontier before the
123    /// caller applies anything. Callers persist `first_page`, request further
124    /// pages with [`Self::next_catch_up_page`], and may acknowledge the fixed
125    /// frontier after all of those pages are durable.
126    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    /// Fetch the next bounded page without chasing events that arrived after
151    /// `frontier` was captured. A response may contain such newer events; the
152    /// returned page is clipped at the exact ordinal-and-digest frontier.
153    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    /// Return the fingerprint and freshness of this session's harness
206    /// credentials without exposing the credential bytes.
207    pub async fn credential_state(&mut self) -> Result<CredentialSnapshot> {
208        credential_snapshot(self.call(RelayRequest::CredentialState).await?)
209    }
210
211    /// Read this session's credential file. Callers must keep these bytes out
212    /// of durable relay observations, logs, and archives.
213    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    /// Install credentials into the harness home fixed by this session's
223    /// launch config.
224    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    /// Return the fingerprint of this session's synced skills trees without
258    /// transferring the tree itself.
259    pub async fn skills_state(&mut self) -> Result<mj_core::skills::SkillsSyncState> {
260        skills_sync_state(self.call(RelayRequest::SkillsState).await?)
261    }
262
263    /// Install background text that only the target harness sees, prepended
264    /// to the next real prompt without creating a synthetic transcript turn.
265    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    /// Replace this session's synced skills trees with an encoded
300    /// `skills::SkillsArchive`. The destination directories are fixed by
301    /// the session's launch config and the harness skills whitelist.
302    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    /// Copy a verified controller blob to this session before admitting its reference.
315    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    /// Recover the local copy needed for queue editing and resubmission.
347    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}