brokk-mj-controller 2.22.0

Daemon-side controller, session manager, and web server for Mjolnir
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
use super::*;

impl RelayClient {
    pub async fn history_requests(&mut self) -> Result<Vec<mj_core::history::HistoryRequest>> {
        if !RelayRequest::HistoryRequests.supported_at(self.protocol_version) {
            return Ok(Vec::new());
        }
        match self.call(RelayRequest::HistoryRequests).await? {
            RelayResponsePayload::HistoryRequests { requests } => Ok(requests),
            _ => bail!("relay returned an unexpected history queue response"),
        }
    }

    pub async fn complete_history_request(
        &mut self,
        result: mj_core::history::HistoryResult,
    ) -> Result<()> {
        match self
            .call(RelayRequest::CompleteHistoryRequest { result })
            .await?
        {
            RelayResponsePayload::HistoryRequestCompleted => Ok(()),
            _ => bail!("relay returned an unexpected history completion response"),
        }
    }

    pub fn session_id(&self) -> &str {
        &self.session_id
    }

    pub fn supports_project_memory_sync(&self) -> bool {
        RelayRequest::ProjectMemorySnapshot.supported_at(self.protocol_version)
    }

    /// The skills archive format this session's worker reads. A worker from
    /// before `RELAY_GZIP_SKILLS_PROTOCOL` reads only the uncompressed format.
    pub fn skills_archive_format(&self) -> mj_core::skills::SkillsArchiveFormat {
        if self.protocol_version >= mj_core::relay::RELAY_GZIP_SKILLS_PROTOCOL {
            mj_core::skills::SkillsArchiveFormat::Gzip
        } else {
            mj_core::skills::SkillsArchiveFormat::Plain
        }
    }

    pub fn relay_version(&self) -> &str {
        &self.relay_version
    }

    /// Content address of the executable serving this connection, or `None`
    /// from a worker too old to report one. A controller reads `None` as
    /// outdated: it predates the field, so it predates this controller.
    pub fn worker_build(&self) -> Option<&str> {
        self.worker_build.as_deref()
    }

    pub fn protocol_version(&self) -> u32 {
        self.protocol_version
    }

    pub fn latest_ordinal(&self) -> u64 {
        self.latest_ordinal
    }

    pub fn latest_digest(&self) -> &str {
        &self.latest_digest
    }

    pub async fn attach(
        &mut self,
        after_ordinal: u64,
        after_digest: impl Into<String>,
    ) -> Result<RelayAttachment> {
        let after_digest = after_digest.into();
        match self
            .call_with_timeout(
                RelayRequest::Attach {
                    after_ordinal,
                    after_digest: after_digest.clone(),
                },
                RELAY_HISTORY_TIMEOUT,
            )
            .await?
        {
            RelayResponsePayload::Attached {
                mut state,
                events,
                through_ordinal,
                through_digest,
            } => {
                let mut cursor = RelayCursor {
                    ordinal: after_ordinal,
                    digest: after_digest,
                };
                for event in &events {
                    validate_relay_event(cursor.ordinal, &cursor.digest, event)
                        .context("verify relay attachment event chain")?;
                    cursor.ordinal = event.ordinal;
                    cursor.digest.clone_from(&event.digest);
                }
                if cursor.ordinal != through_ordinal || cursor.digest != through_digest {
                    bail!("relay attachment frontier does not match its event chain");
                }
                state.relay_protocol_version = Some(self.protocol_version);
                self.latest_ordinal = state.latest_ordinal;
                self.latest_digest = state.latest_digest.clone();
                Ok(RelayAttachment {
                    state,
                    events,
                    through_ordinal,
                    through_digest,
                })
            }
            _ => bail!("relay returned an unexpected attach response"),
        }
    }

    /// Start a bounded catch-up by capturing the relay frontier before the
    /// caller applies anything. Callers persist `first_page`, request further
    /// pages with [`Self::next_catch_up_page`], and may acknowledge the fixed
    /// frontier after all of those pages are durable.
    pub async fn begin_catch_up(
        &mut self,
        after_ordinal: u64,
        after_digest: impl Into<String>,
    ) -> Result<RelayCatchUp> {
        let after_digest = after_digest.into();
        let first = self.attach(after_ordinal, after_digest.clone()).await?;
        let frontier = RelayCursor {
            ordinal: first.state.latest_ordinal,
            digest: first.state.latest_digest.clone(),
        };
        let previous = RelayCursor {
            ordinal: after_ordinal,
            digest: after_digest,
        };
        let state = first.state.clone();
        let first_page = clip_catch_up_page(first, &previous, &frontier)?;
        Ok(RelayCatchUp {
            state,
            frontier,
            first_page,
        })
    }

    /// Fetch the next bounded page without chasing events that arrived after
    /// `frontier` was captured. A response may contain such newer events; the
    /// returned page is clipped at the exact ordinal-and-digest frontier.
    pub async fn next_catch_up_page(
        &mut self,
        previous: &RelayCursor,
        frontier: &RelayCursor,
    ) -> Result<RelayEventPage> {
        if previous.ordinal >= frontier.ordinal {
            bail!("relay catch-up is already at its fixed frontier");
        }
        let attachment = self
            .attach(previous.ordinal, previous.digest.clone())
            .await?;
        clip_catch_up_page(attachment, previous, frontier)
    }

    pub async fn acknowledge(
        &mut self,
        through_ordinal: u64,
        through_digest: impl Into<String>,
    ) -> Result<RelayCursor> {
        match self
            .call_with_timeout(
                RelayRequest::Acknowledge {
                    through_ordinal,
                    through_digest: through_digest.into(),
                },
                RELAY_ACKNOWLEDGE_TIMEOUT,
            )
            .await?
        {
            RelayResponsePayload::Acknowledged {
                through_ordinal,
                through_digest,
            } => Ok(RelayCursor {
                ordinal: through_ordinal,
                digest: through_digest,
            }),
            _ => bail!("relay returned an unexpected acknowledgement response"),
        }
    }

    pub async fn status(&mut self) -> Result<RelayOperationalState> {
        match self.call(RelayRequest::Status).await? {
            RelayResponsePayload::Status(mut status) => {
                status.relay_protocol_version = Some(self.protocol_version);
                self.latest_ordinal = status.latest_ordinal;
                self.latest_digest = status.latest_digest.clone();
                Ok(status)
            }
            _ => bail!("relay returned an unexpected status response"),
        }
    }

    /// Return the fingerprint and freshness of this session's harness
    /// credentials without exposing the credential bytes.
    pub async fn credential_state(&mut self) -> Result<CredentialSnapshot> {
        credential_snapshot(self.call(RelayRequest::CredentialState).await?)
    }

    /// Read this session's credential file. Callers must keep these bytes out
    /// of durable relay observations, logs, and archives.
    pub async fn read_credentials(&mut self) -> Result<Vec<u8>> {
        match self.call(RelayRequest::ReadCredentials).await? {
            RelayResponsePayload::Credentials { data } => BASE64
                .decode(data.as_bytes())
                .context("decode relay credential payload"),
            _ => bail!("relay returned an unexpected credential response"),
        }
    }

    /// Install credentials into the harness home fixed by this session's
    /// launch config.
    pub async fn install_credentials(&mut self, bytes: &[u8]) -> Result<CredentialSnapshot> {
        credential_snapshot(
            self.call(RelayRequest::InstallCredentials {
                data: BASE64.encode(bytes),
            })
            .await?,
        )
    }

    pub async fn github_token_state(
        &mut self,
    ) -> Result<mj_core::credentials::GithubTokenSnapshot> {
        github_token_snapshot(self.call(RelayRequest::GithubTokenState).await?)
    }

    pub async fn install_github_token(
        &mut self,
        token: &str,
    ) -> Result<mj_core::credentials::GithubTokenSnapshot> {
        github_token_snapshot(
            self.call(RelayRequest::InstallGithubToken {
                data: BASE64.encode(token.as_bytes()),
            })
            .await?,
        )
    }

    pub async fn remove_github_token(
        &mut self,
    ) -> Result<mj_core::credentials::GithubTokenSnapshot> {
        github_token_snapshot(self.call(RelayRequest::RemoveGithubToken).await?)
    }

    /// Return the fingerprint of this session's synced skills trees without
    /// transferring the tree itself.
    pub async fn skills_state(&mut self) -> Result<mj_core::skills::SkillsSyncState> {
        skills_sync_state(self.call(RelayRequest::SkillsState).await?)
    }

    /// Install background text that only the target harness sees, prepended
    /// to the next real prompt without creating a synthetic transcript turn.
    pub async fn install_prompt_context(&mut self, text: String) -> Result<()> {
        let request = RelayRequest::InstallPromptContext { text };
        match self.call(request).await? {
            RelayResponsePayload::PromptContextInstalled => Ok(()),
            _ => bail!("relay returned an unexpected prompt-context response"),
        }
    }

    pub async fn project_memory_snapshot(
        &mut self,
    ) -> Result<(
        mj_core::project_memory::ProjectMemorySnapshot,
        mj_core::project_memory::ProjectMemorySnapshot,
    )> {
        let request = RelayRequest::ProjectMemorySnapshot;
        match self.call(request).await? {
            RelayResponsePayload::ProjectMemorySnapshot { baseline, replica } => {
                Ok((baseline, replica))
            }
            _ => bail!("relay returned an unexpected project-memory response"),
        }
    }

    pub async fn install_project_memory_snapshot(
        &mut self,
        snapshot: mj_core::project_memory::ProjectMemorySnapshot,
    ) -> Result<()> {
        let request = RelayRequest::InstallProjectMemorySnapshot { snapshot };
        match self.call(request).await? {
            RelayResponsePayload::ProjectMemorySnapshotInstalled => Ok(()),
            _ => bail!("relay returned an unexpected project-memory install response"),
        }
    }

    /// Replace this session's synced skills trees with an encoded
    /// `skills::SkillsArchive`. The destination directories are fixed by
    /// the session's launch config and the harness skills whitelist.
    pub async fn install_skills(
        &mut self,
        archive_bytes: &[u8],
    ) -> Result<mj_core::skills::SkillsSyncState> {
        skills_sync_state(
            self.call(RelayRequest::InstallSkills {
                data: BASE64.encode(archive_bytes),
            })
            .await?,
        )
    }

    /// Copy a verified controller blob to this session before admitting its reference.
    pub async fn ensure_attachment(
        &mut self,
        reference: &mj_core::attachment::AttachmentRef,
    ) -> Result<()> {
        match self
            .call(RelayRequest::AttachmentPresent {
                reference: reference.clone(),
            })
            .await?
        {
            RelayResponsePayload::AttachmentPresent { present: true } => return Ok(()),
            RelayResponsePayload::AttachmentPresent { present: false } => {}
            _ => bail!("unexpected image presence response"),
        }
        let store = mj_core::attachment::AttachmentStore::controller(&self.session_id)?;
        let reference_copy = reference.clone();
        let bytes = tokio::task::spawn_blocking(move || store.read(&reference_copy))
            .await
            .context("image loading task failed")??;
        match self
            .call(RelayRequest::InstallAttachment {
                reference: reference.clone(),
                data: BASE64.encode(bytes),
            })
            .await?
        {
            RelayResponsePayload::AttachmentInstalled => Ok(()),
            _ => bail!("unexpected image upload response"),
        }
    }

    /// Recover the local copy needed for queue editing and resubmission.
    pub async fn cache_attachment(
        &mut self,
        reference: &mj_core::attachment::AttachmentRef,
    ) -> Result<()> {
        let store = mj_core::attachment::AttachmentStore::controller(&self.session_id)?;
        let local = store.clone();
        let reference_copy = reference.clone();
        if tokio::task::spawn_blocking(move || local.contains(&reference_copy))
            .await
            .context("image lookup task failed")??
        {
            return Ok(());
        }
        let RelayResponsePayload::AttachmentData { data } = self
            .call(RelayRequest::ReadAttachment {
                reference: reference.clone(),
            })
            .await?
        else {
            bail!("unexpected image download response")
        };
        let reference = reference.clone();
        tokio::task::spawn_blocking(move || {
            anyhow::ensure!(
                data.len() <= mj_core::attachment::MAX_IMAGE_BYTES.div_ceil(3) * 4,
                "image download is too large"
            );
            store.install(&reference, &BASE64.decode(data)?)
        })
        .await
        .context("image caching task failed")?
    }

    pub async fn submit(
        &mut self,
        command_id: impl Into<String>,
        command: RelayCommand,
    ) -> Result<u64> {
        let command_id = command_id.into();
        if let RelayCommand::Prompt { prompt } = &command {
            for reference in mj_core::attachment::references(prompt)? {
                self.ensure_attachment(&reference).await?;
            }
        }
        match self
            .call(RelayRequest::Submit {
                command_id: command_id.clone(),
                command,
            })
            .await?
        {
            RelayResponsePayload::Accepted {
                command_id: accepted_id,
                ordinal,
            } if accepted_id == command_id => Ok(ordinal),
            RelayResponsePayload::Accepted {
                command_id: accepted_id,
                ..
            } => bail!("relay accepted command under ID {accepted_id}, expected {command_id}"),
            _ => bail!("relay returned an unexpected command response"),
        }
    }

    pub async fn reserve_idle(&mut self, command_id: String) -> Result<bool> {
        match self.call(RelayRequest::ReserveIdle { command_id }).await? {
            RelayResponsePayload::IdleReservation { ordinal } => Ok(ordinal.is_some()),
            _ => bail!("relay returned an unexpected idle reservation response"),
        }
    }
}