polyc-state-connect 2026.10.2

State plane transport adapter: capability-specific Connect clients and server-trait glue mapping the generated wire types onto the polyc-state kernel — typed outcomes, per-call admission, and the conformance surface the authenticated shell proves itself against (docs/proposals/separated-planes.md).
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
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
use connectrpc::client::{ClientConfig, ClientTransport};
use polyc_proto::proto::polychrome::state::v1 as pb;
use polyc_state::{
    consistency::Watermark,
    digest::ContentDigest,
    error::StateError,
    id::{CommandId, NamespaceId},
    page::{Cursor, Page},
    query_audit::{
        AuditSource, AuditSourceHead, HistoryReceipt, QueryAuditHistoryEntry, QueryCompletion,
        QueryId,
    },
    revision::{JournalPosition, PartitionIncarnation},
};

use crate::{
    MAX_QUERY_AUDIT_WIRE_MESSAGE_BYTES,
    error::{TransportFallback, from_connect_error},
    trace::bounded_traced_options,
    wire::{DeclaredCall, Kernel},
};

/// Capability-specific State client for the query-audit history.
pub struct QueryAuditHistoryClient<T> {
    inner: pb::StateQueryAuditHistoryServiceClient<T>,
}

/// Refuses a reply this reader cannot stand on.
fn malformed(field: &str, reason: &str) -> StateError {
    StateError::Malformed {
        field: field.to_owned(),
        reason: reason.to_owned(),
    }
}

/// Decodes the authority a reply names.
fn source_of(source: Option<pb::AuditSource>) -> Result<AuditSource, StateError> {
    let source =
        source.ok_or_else(|| malformed("source", "a reply names the authority it read"))?;
    if source.partition != polyc_state::query_audit::SOURCE_PARTITION {
        return Err(malformed(
            "partition",
            "the query-audit authority has one partition",
        ));
    }
    let incarnation: [u8; PartitionIncarnation::LEN] = source
        .incarnation
        .try_into()
        .map_err(|_| malformed("incarnation", "a lineage is exactly 32 bytes"))?;
    Ok(AuditSource::new(PartitionIncarnation::from_bytes(
        incarnation,
    )))
}

impl<T> QueryAuditHistoryClient<T>
where
    T: ClientTransport,
    <T::ResponseBody as connectrpc::http_body::Body>::Error: std::fmt::Display,
{
    /// Builds a client for one State listener.
    #[must_use]
    pub fn new(transport: T, config: ClientConfig) -> Self {
        Self {
            inner: pb::StateQueryAuditHistoryServiceClient::new(
                transport,
                config.with_default_max_message_size(MAX_QUERY_AUDIT_WIRE_MESSAGE_BYTES),
            ),
        }
    }

    fn fallback(declared: &DeclaredCall, attempted: usize) -> TransportFallback {
        TransportFallback::new(
            polyc_state::query_audit::family(),
            MAX_QUERY_AUDIT_WIRE_MESSAGE_BYTES as u64,
            attempted as u64,
            declared.budget,
        )
    }

    /// Reads which lineage the authority is on and how far it has issued.
    ///
    /// # Errors
    ///
    /// Returns a typed State or transport refusal. A reply naming another
    /// partition, or one whose lineage is not exactly 32 bytes, is malformed.
    pub async fn source_head(
        &self,
        declared: &DeclaredCall,
    ) -> Result<AuditSourceHead, StateError> {
        let message = pb::GetQueryAuditSourceRequest {
            context: buffa::MessageField::some(Kernel(declared).into()),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        };
        let attempted = buffa::Message::encoded_len(&message) as usize;
        let reply = self
            .inner
            .get_source_with_options(message, bounded_traced_options(declared))
            .await
            .map_err(|error| from_connect_error(&error, &Self::fallback(declared, attempted)))?
            .into_owned();
        Ok(AuditSourceHead::new(
            source_of(reply.source.into_option())?,
            JournalPosition::new(reply.head),
        ))
    }

    /// Reads one bounded page of the history after `after`, in `expected_lineage`.
    ///
    /// A resume passes the lineage its cursor counts ordinals in. A wiped and
    /// refilled authority then refuses the cursor rather than serving its own
    /// entries as a continuation of the log the cursor came from.
    ///
    /// # Errors
    ///
    /// Returns a typed State or transport refusal. The page is re-validated
    /// here rather than trusted: a server that returned a descending page, a
    /// duplicate ordinal, a completion at or before its intent, or a watermark
    /// below its own last ordinal would make a folding reader skip a phase or
    /// stop early, and neither is visible later.
    pub async fn history(
        &self,
        declared: &DeclaredCall,
        after: Option<JournalPosition>,
        expected_lineage: Option<PartitionIncarnation>,
        limit: u32,
    ) -> Result<(Page<QueryAuditHistoryEntry>, AuditSource), StateError> {
        let message = pb::ListQueryAuditHistoryRequest {
            context: buffa::MessageField::some(Kernel(declared).into()),
            after: after.map(JournalPosition::get),
            limit,
            expected_lineage: expected_lineage.map(|lineage| lineage.as_bytes().to_vec()),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        };
        let attempted = buffa::Message::encoded_len(&message) as usize;
        let reply = self
            .inner
            .list_history_with_options(message, bounded_traced_options(declared))
            .await
            .map_err(|error| from_connect_error(&error, &Self::fallback(declared, attempted)))?
            .into_owned();
        let source = source_of(reply.source.into_option())?;
        let mut records = Vec::with_capacity(reply.entries.len());
        let mut last = after.map_or(0, JournalPosition::get);
        for entry in reply.entries {
            let decoded = decode_entry(entry)?;
            let ordinal = decoded_ordinal(&decoded);
            // Strictly ascending, and past the cursor the caller resumed from.
            // A repeated or descending ordinal would let a fold count one
            // phase twice or miss one entirely.
            if ordinal.get() <= last {
                return Err(malformed(
                    "entries",
                    "a history page is strictly ascending past its cursor",
                ));
            }
            last = ordinal.get();
            records.push(decoded);
        }
        // The server's own cursor is checked, not adopted. A fold resumes from
        // it, so a cursor above the page's last ordinal would make the next
        // read skip every entry in between — the exact loss this validation
        // exists to prevent.
        if let Some(cursor) = reply.next.as_option()
            && cursor.position != last
        {
            return Err(malformed(
                "next",
                "a resume cursor names the page's own last ordinal",
            ));
        }
        // The bound the caller asked for is the bound the page must respect.
        if u32::try_from(records.len()).unwrap_or(u32::MAX) > limit {
            return Err(malformed(
                "entries",
                "a page holds no more entries than the bound asked for",
            ));
        }
        if reply.watermark < last {
            return Err(malformed(
                "watermark",
                "a watermark is at least the page's own last ordinal",
            ));
        }
        let completeness = crate::wire::completeness("completeness", reply.completeness)?;
        let consistency = crate::wire::consistency("consistency", reply.consistency)?;
        let next = reply
            .next
            .into_option()
            .map(|cursor| Kernel::<Cursor>::from(cursor).into_inner());
        Ok((
            Page::new(records, next, completeness, consistency)
                .with_watermark(Watermark::new(reply.watermark)),
            source,
        ))
    }
}

/// Returns the ordinal one decoded entry occupies.
const fn decoded_ordinal(entry: &QueryAuditHistoryEntry) -> JournalPosition {
    match entry {
        QueryAuditHistoryEntry::Intent { ordinal, .. }
        | QueryAuditHistoryEntry::Completion { ordinal, .. } => *ordinal,
    }
}

/// Decodes one entry, refusing a phase this authority never records.
fn decode_entry(entry: pb::QueryAuditHistoryEntry) -> Result<QueryAuditHistoryEntry, StateError> {
    use polyc_proto::proto::polychrome::state::v1::query_audit_history_entry::Phase;
    let pb::QueryAuditHistoryEntry {
        ordinal,
        command_id,
        receipt_digest,
        phase,
        __buffa_unknown_fields: _,
    } = entry;
    let digest: [u8; ContentDigest::LEN] = receipt_digest
        .try_into()
        .map_err(|_| malformed("receipt_digest", "a digest is exactly 32 bytes"))?;
    let receipt = HistoryReceipt::new(
        CommandId::new(command_id),
        ContentDigest::from_bytes(digest),
    );
    let ordinal = JournalPosition::new(ordinal);
    match phase.ok_or_else(|| malformed("phase", "an entry records which phase it is"))? {
        Phase::Intent(intent) => Ok(QueryAuditHistoryEntry::Intent {
            ordinal,
            intent: Kernel::<polyc_state::query_audit::AuditIntent>::try_from(*intent)?
                .into_inner(),
            receipt,
        }),
        Phase::Completion(completion) => {
            let pb::QueryAuditHistoryCompletion {
                query,
                namespace,
                intent_ordinal,
                completion,
                __buffa_unknown_fields: _,
            } = *completion;
            let intent_ordinal = JournalPosition::new(intent_ordinal);
            // A completion is recorded after the intent it completes, so it
            // took a later ordinal from the same counter.
            if intent_ordinal >= ordinal {
                return Err(malformed(
                    "intent_ordinal",
                    "a completion is ordered after the intent it completes",
                ));
            }
            let completion =
                Kernel::<QueryCompletion>::try_from(completion.into_option().ok_or_else(
                    || malformed("completion", "a completion entry carries its completion"),
                )?)?
                .into_inner();
            Ok(QueryAuditHistoryEntry::Completion {
                ordinal,
                query: QueryId::new(query),
                namespace: NamespaceId::new(namespace),
                intent_ordinal,
                completion,
                receipt,
            })
        }
    }
}

#[cfg(test)]
mod tests {
    use std::{
        pin::Pin,
        sync::Arc,
        task::{Context, Poll},
        time::Duration,
    };

    use connectrpc::{
        client::{ClientBody, ClientConfig, ClientTransport},
        http_body::{Body, Frame},
    };
    use futures::future::BoxFuture;
    use polyc_proto::proto::polychrome::state::v1 as pb;
    use polyc_state::error::StateError;

    use super::QueryAuditHistoryClient;
    use crate::{state_audience, wire::DeclaredCall};

    type Bytes = bytes::Bytes;

    struct CannedBody(std::vec::IntoIter<Bytes>);

    impl Body for CannedBody {
        type Data = Bytes;
        type Error = std::io::Error;

        fn poll_frame(
            mut self: Pin<&mut Self>,
            _context: &mut Context<'_>,
        ) -> Poll<Option<Result<Frame<Bytes>, Self::Error>>> {
            Poll::Ready(self.0.next().map(|bytes| Ok(Frame::data(bytes))))
        }
    }

    #[derive(Clone)]
    struct CannedTransport {
        frames: Arc<Vec<Bytes>>,
    }

    impl ClientTransport for CannedTransport {
        type ResponseBody = CannedBody;
        type Error = std::io::Error;

        fn send(
            &self,
            _request: http::Request<ClientBody>,
        ) -> BoxFuture<'static, Result<http::Response<Self::ResponseBody>, Self::Error>> {
            let frames = self.frames.as_ref().clone();
            Box::pin(async move {
                Ok(http::Response::builder()
                    .status(http::StatusCode::OK)
                    .header(http::header::CONTENT_TYPE, "application/proto")
                    .body(CannedBody(frames.into_iter()))
                    .unwrap())
            })
        }
    }

    fn client(reply: &pb::ListQueryAuditHistoryReply) -> QueryAuditHistoryClient<CannedTransport> {
        QueryAuditHistoryClient::new(
            CannedTransport {
                frames: Arc::new(vec![Bytes::from(buffa::Message::encode_to_vec(reply))]),
            },
            ClientConfig::new("http://audit.invalid".parse().unwrap()),
        )
    }

    fn audit_source() -> pb::AuditSource {
        pb::AuditSource {
            partition: polyc_state::query_audit::SOURCE_PARTITION.to_owned(),
            incarnation: vec![7; 32],
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        }
    }

    fn intent_entry(ordinal: u64) -> pb::QueryAuditHistoryEntry {
        use polyc_proto::proto::polychrome::state::v1::query_audit_history_entry::Phase;
        pb::QueryAuditHistoryEntry {
            ordinal,
            command_id: format!("c-{ordinal}"),
            receipt_digest: vec![3; 32],
            phase: Some(Phase::Intent(Box::new(pb::QueryAuditIntent {
                query: format!("q-{ordinal}"),
                requester: "r".to_owned(),
                shape: vec![5; 32],
                recorded_at_nanos: 1,
                position: ordinal,
                source: buffa::MessageField::some(pb::QuerySourceSnapshot::default()),
                namespace: "n".to_owned(),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            }))),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        }
    }

    fn page(
        entries: Vec<pb::QueryAuditHistoryEntry>,
        watermark: u64,
    ) -> pb::ListQueryAuditHistoryReply {
        pb::ListQueryAuditHistoryReply {
            entries,
            next: buffa::MessageField::default(),
            completeness: pb::PageCompleteness::PAGE_COMPLETENESS_COMPLETE.into(),
            consistency: pb::Consistency::CONSISTENCY_LINEARIZABLE_CURRENT.into(),
            watermark,
            source: buffa::MessageField::some(audit_source()),
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        }
    }

    async fn refusal(reply: &pb::ListQueryAuditHistoryReply) -> StateError {
        client(reply)
            .history(
                &DeclaredCall::live(state_audience(), Duration::MAX),
                None,
                None,
                8,
            )
            .await
            .expect_err("a hostile page was accepted")
    }

    /// An honest page is accepted.
    ///
    /// Without this, every refusal below would also pass for a client that
    /// refused everything.
    #[tokio::test]
    async fn an_ascending_page_is_accepted() {
        let (page, source) = client(&page(vec![intent_entry(1), intent_entry(2)], 2))
            .history(
                &DeclaredCall::live(state_audience(), Duration::MAX),
                None,
                None,
                8,
            )
            .await
            .expect("an honest page was refused");
        assert_eq!(page.records().len(), 2);
        assert_eq!(
            source.partition().as_str(),
            polyc_state::query_audit::SOURCE_PARTITION
        );
    }

    /// A descending page is refused rather than sorted.
    ///
    /// A fold resumes from the last ordinal it saw. Sorting here would hide
    /// that the server broke the order the cursor depends on.
    #[tokio::test]
    async fn a_descending_page_is_refused() {
        let error = refusal(&page(vec![intent_entry(2), intent_entry(1)], 2)).await;
        assert!(
            matches!(&error, StateError::Malformed { field, .. } if field == "entries"),
            "expected an entries refusal, got {error:?}"
        );
    }

    /// A repeated ordinal is refused.
    ///
    /// Two entries at one ordinal would make a fold count one phase twice.
    #[tokio::test]
    async fn a_repeated_ordinal_is_refused() {
        let error = refusal(&page(vec![intent_entry(1), intent_entry(1)], 2)).await;
        assert!(
            matches!(&error, StateError::Malformed { field, .. } if field == "entries"),
            "expected an entries refusal, got {error:?}"
        );
    }

    /// A page holding more entries than the bound is refused.
    ///
    /// The bound is the caller's, and a server that ignored it could return a
    /// page no caller sized a budget for.
    #[tokio::test]
    async fn a_page_past_the_bound_is_refused() {
        let over = page(vec![intent_entry(1), intent_entry(2), intent_entry(3)], 3);
        let error = client(&over)
            .history(
                &DeclaredCall::live(state_audience(), Duration::MAX),
                None,
                None,
                2,
            )
            .await
            .expect_err("a page past the bound was accepted");
        assert!(
            matches!(&error, StateError::Malformed { field, .. } if field == "entries"),
            "expected an entries refusal, got {error:?}"
        );
    }

    /// A page whose cursor names its own last ordinal is accepted.
    ///
    /// Without this the cursor rule would pass for a client that refused every
    /// present cursor, and the honest-page control above sends none.
    #[tokio::test]
    async fn a_cursor_naming_the_last_ordinal_is_accepted() {
        let mut reply = page(vec![intent_entry(1), intent_entry(2)], 2);
        reply.next = buffa::MessageField::some(pb::Cursor {
            position: 2,
            snapshot: None,
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        });
        let (accepted, _) = client(&reply)
            .history(
                &DeclaredCall::live(state_audience(), Duration::MAX),
                None,
                None,
                8,
            )
            .await
            .expect("an honest cursor was refused");
        assert_eq!(accepted.records().len(), 2);
    }

    /// A resume cursor above the page's own last ordinal is refused.
    ///
    /// A fold resumes from the cursor the page returns. One above the last
    /// entry makes the next read skip everything in between, and nothing
    /// downstream can see that a phase was never folded.
    #[tokio::test]
    async fn a_cursor_past_the_page_is_refused() {
        let mut reply = page(vec![intent_entry(1), intent_entry(2)], 9);
        reply.next = buffa::MessageField::some(pb::Cursor {
            position: 8,
            snapshot: None,
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        });
        let error = refusal(&reply).await;
        assert!(
            matches!(&error, StateError::Malformed { field, .. } if field == "next"),
            "expected a next-cursor refusal, got {error:?}"
        );
    }

    /// A watermark below the page's own last ordinal is refused.
    ///
    /// A consumer compares its last ordinal against the watermark to decide it
    /// is caught up. One below the page it arrived with describes no authority
    /// state a reader could act on.
    #[tokio::test]
    async fn a_watermark_below_the_page_is_refused() {
        let error = refusal(&page(vec![intent_entry(1), intent_entry(2)], 1)).await;
        assert!(
            matches!(&error, StateError::Malformed { field, .. } if field == "watermark"),
            "expected a watermark refusal, got {error:?}"
        );
    }

    /// A completion ordered at or before its intent is refused.
    ///
    /// A completion is recorded after the intent it completes, so it took a
    /// later ordinal from the same counter. The reverse names a query that
    /// finished before it began.
    #[tokio::test]
    async fn a_completion_before_its_intent_is_refused() {
        use polyc_proto::proto::polychrome::state::v1::query_audit_history_entry::Phase;
        let mut entry = intent_entry(2);
        entry.phase = Some(Phase::Completion(Box::new(
            pb::QueryAuditHistoryCompletion {
                query: "q".to_owned(),
                namespace: "n".to_owned(),
                intent_ordinal: 2,
                completion: buffa::MessageField::some(pb::QueryCompletion::default()),
                __buffa_unknown_fields: buffa::UnknownFields::default(),
            },
        )));
        let error = refusal(&page(vec![entry], 2)).await;
        assert!(
            matches!(&error, StateError::Malformed { field, .. } if field == "intent_ordinal"),
            "expected an intent_ordinal refusal, got {error:?}"
        );
    }

    /// A reply naming another partition is refused.
    #[tokio::test]
    async fn a_reply_about_another_authority_is_refused() {
        let mut reply = page(vec![intent_entry(1)], 1);
        reply.source = buffa::MessageField::some(pb::AuditSource {
            partition: "state.other".to_owned(),
            incarnation: vec![7; 32],
            __buffa_unknown_fields: buffa::UnknownFields::default(),
        });
        let error = refusal(&reply).await;
        assert!(
            matches!(&error, StateError::Malformed { field, .. } if field == "partition"),
            "expected a partition refusal, got {error:?}"
        );
    }
}