use polyc_proto::proto::polychrome::state::v1 as pb;
use polyc_state::{
command::{
CommandEnvelope, CommandMetadata, CommandScope, FencingToken, Precondition, ResourceBounds,
},
digest::ContentDigest,
error::StateError,
feed::{
self, AcknowledgeProjectorCursor, CommitEnvelope, CompactFeedPrefix, ConsumerPolicy,
CreateSnapshot, FeedAnchor, FeedChunk, FeedCompaction, FeedRecord, FeedRetention,
FeedSnapshot, ProjectorRegistration, ProjectorStatus, RegisterProjector, SubscribeCommits,
},
id::{
AggregateId, Audience, CommandId, ConsumerId, NamespaceId, PartitionId, Purpose, SnapshotId,
},
journal::JournalRecord,
page::{Cursor, Positioned as _, ReadStart},
receipt::Receipt,
revision::JournalPosition,
stream::StreamChunk,
};
use crate::wire::{Kernel, fixed_bytes, known, malformed, required};
impl From<Kernel<ConsumerPolicy>> for pb::ConsumerPolicy {
fn from(value: Kernel<ConsumerPolicy>) -> Self {
match value.0 {
ConsumerPolicy::Required => Self::CONSUMER_POLICY_REQUIRED,
ConsumerPolicy::Optional => Self::CONSUMER_POLICY_OPTIONAL,
}
}
}
fn consumer_policy(
field: &str,
value: buffa::EnumValue<pb::ConsumerPolicy>,
) -> Result<ConsumerPolicy, StateError> {
match known(field, value)? {
pb::ConsumerPolicy::CONSUMER_POLICY_REQUIRED => Ok(ConsumerPolicy::Required),
pb::ConsumerPolicy::CONSUMER_POLICY_OPTIONAL => Ok(ConsumerPolicy::Optional),
pb::ConsumerPolicy::CONSUMER_POLICY_UNSPECIFIED => Err(malformed(
field,
"a projector declares whether its lag holds the feed prefix",
)),
}
}
const fn feed_bounds() -> ResourceBounds {
ResourceBounds::new(crate::MAX_WIRE_MESSAGE_BYTES as u64, 1)
}
impl From<Kernel<&CommandMetadata>> for pb::FeedCommand {
fn from(value: Kernel<&CommandMetadata>) -> Self {
let metadata = value.0;
Self {
command_id: metadata.command_id().as_str().to_owned(),
partition: metadata.scope().partition().as_str().to_owned(),
aggregate: metadata.scope().aggregate().as_str().to_owned(),
namespace: metadata.scope().namespace().as_str().to_owned(),
purpose: metadata.envelope().purpose().as_str().to_owned(),
audience: metadata.envelope().audience().as_str().to_owned(),
digest: metadata.digest().as_bytes().to_vec(),
precondition: buffa::MessageField::some(pb::Precondition::from(Kernel(
metadata.precondition(),
))),
fence: metadata.fence().map(FencingToken::get),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::FeedCommand> for Kernel<CommandMetadata> {
type Error = StateError;
fn try_from(value: pb::FeedCommand) -> Result<Self, Self::Error> {
let digest = ContentDigest::from_bytes(fixed_bytes::<{ ContentDigest::LEN }>(
"digest",
&value.digest,
)?);
let precondition = Kernel::<Precondition>::try_from(required(
"precondition",
"a precondition names what durable state the command requires",
value.precondition,
)?)?
.into_inner();
let mut metadata = CommandMetadata::new(
CommandId::new(value.command_id),
feed::family(),
digest,
CommandScope::new(
AggregateId::new(value.aggregate),
PartitionId::new(value.partition),
NamespaceId::new(value.namespace),
),
CommandEnvelope::new(
Purpose::new(value.purpose),
Audience::new(value.audience),
feed_bounds(),
),
)
.with_precondition(precondition);
if let Some(fence) = value.fence {
metadata = metadata.with_fence(FencingToken::new(fence));
}
Ok(Self(metadata))
}
}
impl From<Kernel<&CreateSnapshot>> for pb::FeedCommand {
fn from(value: Kernel<&CreateSnapshot>) -> Self {
Self::from(Kernel(value.0.metadata()))
}
}
impl From<Kernel<&RegisterProjector>> for pb::FeedCommand {
fn from(value: Kernel<&RegisterProjector>) -> Self {
Self::from(Kernel(value.0.metadata()))
}
}
impl From<Kernel<&AcknowledgeProjectorCursor>> for pb::FeedCommand {
fn from(value: Kernel<&AcknowledgeProjectorCursor>) -> Self {
Self::from(Kernel(value.0.metadata()))
}
}
impl From<Kernel<&CompactFeedPrefix>> for pb::FeedCommand {
fn from(value: Kernel<&CompactFeedPrefix>) -> Self {
Self::from(Kernel(value.0.metadata()))
}
}
impl From<Kernel<&CommitEnvelope>> for pb::CommitEnvelope {
fn from(value: Kernel<&CommitEnvelope>) -> Self {
let envelope = value.0;
Self {
partition: envelope.partition().as_str().to_owned(),
command_id: envelope.command_id().as_str().to_owned(),
digest: envelope.digest().as_bytes().to_vec(),
fence: envelope.fence().map(FencingToken::get),
head_before: envelope.head_before().get(),
head_after: envelope.head_after().get(),
record_count: envelope.record_count(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::CommitEnvelope> for Kernel<CommitEnvelope> {
type Error = StateError;
fn try_from(value: pb::CommitEnvelope) -> Result<Self, Self::Error> {
let digest = ContentDigest::from_bytes(fixed_bytes::<{ ContentDigest::LEN }>(
"digest",
&value.digest,
)?);
let envelope = CommitEnvelope::new(
PartitionId::new(value.partition),
CommandId::new(value.command_id),
digest,
JournalPosition::new(value.head_before),
JournalPosition::new(value.head_after),
value.record_count,
);
Ok(Self(match value.fence {
Some(fence) => envelope.with_fence(FencingToken::new(fence)),
None => envelope,
}))
}
}
impl From<Kernel<&FeedRecord>> for pb::FeedRecord {
fn from(value: Kernel<&FeedRecord>) -> Self {
let entry = value.0;
Self {
position: entry.position().get(),
envelope: buffa::MessageField::some(pb::CommitEnvelope::from(Kernel(entry.envelope()))),
records: entry
.records()
.iter()
.map(|record| pb::JournalRecord::from(Kernel(record)))
.collect(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::FeedRecord> for Kernel<FeedRecord> {
type Error = StateError;
fn try_from(value: pb::FeedRecord) -> Result<Self, Self::Error> {
let envelope = Kernel::<CommitEnvelope>::try_from(required(
"envelope",
"a feed entry carries the commit it describes",
value.envelope,
)?)?
.into_inner();
let records = value
.records
.into_iter()
.map(|record| Kernel::<JournalRecord>::try_from(record).map(Kernel::into_inner))
.collect::<Result<Vec<_>, _>>()?;
Ok(Self(FeedRecord::new(
JournalPosition::new(value.position),
envelope,
records,
)))
}
}
impl From<Kernel<&FeedChunk>> for pb::FeedChunk {
fn from(value: Kernel<&FeedChunk>) -> Self {
let chunk = value.0;
Self {
records: chunk
.records()
.iter()
.map(|record| pb::FeedRecord::from(Kernel(record)))
.collect(),
next: chunk
.next_cursor()
.map_or_else(buffa::MessageField::default, |cursor| {
buffa::MessageField::some(pb::Cursor::from(Kernel(cursor)))
}),
end: pb::StreamEnd::from(Kernel(chunk.end())).into(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::FeedChunk> for Kernel<FeedChunk> {
type Error = StateError;
fn try_from(value: pb::FeedChunk) -> Result<Self, Self::Error> {
let end = crate::wire::stream_end("end", value.end)?;
let records = value
.records
.into_iter()
.map(|record| Kernel::<FeedRecord>::try_from(record).map(Kernel::into_inner))
.collect::<Result<Vec<_>, _>>()?;
let next = value
.next
.into_option()
.map(|cursor| Kernel::<Cursor>::from(cursor).into_inner());
Ok(Self(StreamChunk::new(records, next, end)))
}
}
impl From<Kernel<&FeedSnapshot>> for pb::FeedSnapshot {
fn from(value: Kernel<&FeedSnapshot>) -> Self {
let snapshot = value.0;
Self {
snapshot: snapshot.id().as_str().to_owned(),
partition: snapshot.anchor().partition().as_str().to_owned(),
position: snapshot.position().get(),
head: snapshot.head().get(),
receipt: buffa::MessageField::some(pb::Receipt::from(Kernel(snapshot.receipt()))),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::FeedSnapshot> for Kernel<FeedSnapshot> {
type Error = StateError;
fn try_from(value: pb::FeedSnapshot) -> Result<Self, Self::Error> {
let anchor = FeedAnchor::parse(&SnapshotId::new(value.snapshot))?;
if anchor.partition().as_str() != value.partition
|| anchor.position() != JournalPosition::new(value.position)
{
return Err(malformed(
"snapshot",
"a snapshot identity names the partition and position beside it",
));
}
let receipt = Kernel::<Receipt>::try_from(required(
"receipt",
"a recorded snapshot carries the receipt that recorded it",
value.receipt,
)?)?
.into_inner();
Ok(Self(FeedSnapshot::new(
anchor,
JournalPosition::new(value.head),
receipt,
)))
}
}
#[must_use]
pub fn subscription(request: &SubscribeCommits) -> (String, pb::ReadStart, u32, Option<String>) {
(
request.partition().as_str().to_owned(),
pb::ReadStart::from(Kernel(request.start())),
request.max_chunk_commits(),
request
.consumer()
.map(|consumer| consumer.as_str().to_owned()),
)
}
pub fn subscribe_request(
partition: String,
start: impl Into<Option<pb::ReadStart>>,
max_chunk_commits: u32,
consumer: Option<String>,
) -> Result<SubscribeCommits, StateError> {
let partition = PartitionId::new(partition);
if partition.is_empty() {
return Err(malformed(
"partition",
"a subscription names exactly one partition",
));
}
let start = start
.into()
.ok_or_else(|| malformed("start", "a subscription names where it begins"))?;
let start = Kernel::<ReadStart>::try_from(start)?.into_inner();
let request = SubscribeCommits::new(partition, start, max_chunk_commits);
Ok(match consumer {
Some(consumer) => request.on_behalf_of(ConsumerId::new(consumer)),
None => request,
})
}
impl From<Kernel<&ProjectorRegistration>> for pb::ProjectorRegistration {
fn from(value: Kernel<&ProjectorRegistration>) -> Self {
let registration = value.0;
Self {
consumer: registration.consumer().as_str().to_owned(),
partition: registration.partition().as_str().to_owned(),
policy: pb::ConsumerPolicy::from(Kernel(registration.policy())).into(),
declared_lag_commits: registration.declared_lag_commits(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::ProjectorRegistration> for Kernel<ProjectorRegistration> {
type Error = StateError;
fn try_from(value: pb::ProjectorRegistration) -> Result<Self, Self::Error> {
let policy = consumer_policy("policy", value.policy)?;
Ok(Self(ProjectorRegistration::new(
ConsumerId::new(value.consumer),
PartitionId::new(value.partition),
policy,
value.declared_lag_commits,
)))
}
}
impl From<Kernel<&ProjectorStatus>> for pb::ProjectorStatus {
fn from(value: Kernel<&ProjectorStatus>) -> Self {
let status = value.0;
Self {
registration: buffa::MessageField::some(pb::ProjectorRegistration::from(Kernel(
status.registration(),
))),
acknowledged: status.acknowledged().get(),
evicted: status.is_evicted(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::ProjectorStatus> for Kernel<ProjectorStatus> {
type Error = StateError;
fn try_from(value: pb::ProjectorStatus) -> Result<Self, Self::Error> {
let registration = Kernel::<ProjectorRegistration>::try_from(required(
"registration",
"a projector status carries what the projector declared",
value.registration,
)?)?
.into_inner();
let status = ProjectorStatus::new(registration, JournalPosition::new(value.acknowledged));
Ok(Self(if value.evicted {
status.evicted()
} else {
status
}))
}
}
impl From<Kernel<&FeedRetention>> for pb::FeedRetention {
fn from(value: Kernel<&FeedRetention>) -> Self {
let retention = value.0;
Self {
partition: retention.partition().as_str().to_owned(),
earliest: retention.earliest().get(),
head: retention.head().get(),
held_by: retention
.held_by_consumer()
.map(|consumer| consumer.as_str().to_owned()),
pinned_by: retention
.pinned_by_snapshot()
.map(|snapshot| snapshot.as_str().to_owned()),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl From<pb::FeedRetention> for Kernel<FeedRetention> {
fn from(value: pb::FeedRetention) -> Self {
let retention = FeedRetention::new(
PartitionId::new(value.partition),
JournalPosition::new(value.earliest),
JournalPosition::new(value.head),
);
let retention = match value.held_by {
Some(consumer) => retention.held_by(ConsumerId::new(consumer)),
None => retention,
};
Self(match value.pinned_by {
Some(snapshot) => retention.pinned_by(SnapshotId::new(snapshot)),
None => retention,
})
}
}
impl From<Kernel<&FeedCompaction>> for pb::FeedCompaction {
fn from(value: Kernel<&FeedCompaction>) -> Self {
let compaction = value.0;
Self {
retention: buffa::MessageField::some(pb::FeedRetention::from(Kernel(
compaction.retention(),
))),
evicted: compaction
.evicted()
.iter()
.map(|consumer| consumer.as_str().to_owned())
.collect(),
receipt: compaction
.receipt()
.map_or_else(buffa::MessageField::default, |receipt| {
buffa::MessageField::some(pb::Receipt::from(Kernel(receipt)))
}),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
}
impl TryFrom<pb::FeedCompaction> for Kernel<FeedCompaction> {
type Error = StateError;
fn try_from(value: pb::FeedCompaction) -> Result<Self, Self::Error> {
let retention = Kernel::<FeedRetention>::from(required(
"retention",
"a compaction reports what the partition now retains",
value.retention,
)?)
.into_inner();
let evicted = value.evicted.into_iter().map(ConsumerId::new).collect();
let receipt = value
.receipt
.into_option()
.map(|receipt| Kernel::<Receipt>::try_from(receipt).map(Kernel::into_inner))
.transpose()?;
Ok(Self(FeedCompaction::new(retention, evicted, receipt)))
}
}
#[cfg(test)]
mod tests {
#![allow(clippy::pedantic, clippy::nursery, missing_docs, clippy::unwrap_used)]
use super::*;
use polyc_state::{
consistency::Consistency,
journal::{RecordKind, RecordTrust},
receipt::CommitEvidence,
};
fn metadata() -> CommandMetadata {
CommandMetadata::new(
CommandId::new("cmd-1"),
feed::family(),
ContentDigest::from_bytes([3; ContentDigest::LEN]),
CommandScope::new(
AggregateId::new("conv-1"),
PartitionId::new("conv-1"),
NamespaceId::new("default"),
),
CommandEnvelope::new(
Purpose::new("projection"),
crate::state_audience(),
feed_bounds(),
),
)
.with_fence(FencingToken::new(4))
}
fn receipt() -> Receipt {
Receipt::committed(
&metadata(),
CommitEvidence::new()
.with_snapshot(SnapshotId::new("feed:conv-1@3"))
.with_position(JournalPosition::new(3)),
Consistency::OrderedPerAggregate,
)
}
fn entry() -> FeedRecord {
FeedRecord::new(
JournalPosition::new(3),
CommitEnvelope::new(
PartitionId::new("conv-1"),
CommandId::new("cmd-9"),
ContentDigest::from_bytes([8; ContentDigest::LEN]),
JournalPosition::new(5),
JournalPosition::new(7),
2,
)
.with_fence(FencingToken::new(2)),
vec![
JournalRecord::new(
JournalPosition::new(6),
RecordKind::new("turn_input"),
RecordTrust::TrustedUser,
b"one".to_vec(),
),
JournalRecord::new(
JournalPosition::new(7),
RecordKind::new("turn_output"),
RecordTrust::QuarantinedContent,
b"two".to_vec(),
),
],
)
}
#[test]
fn a_feed_command_round_trips_every_field_it_carries() {
let encoded = pb::FeedCommand::from(Kernel(&metadata()));
let back = Kernel::<CommandMetadata>::try_from(encoded).unwrap().0;
assert_eq!(back, metadata());
}
#[test]
fn every_command_shape_encodes_through_the_same_metadata() {
let expected = pb::FeedCommand::from(Kernel(&metadata()));
assert_eq!(
pb::FeedCommand::from(Kernel(&CreateSnapshot::new(metadata()))),
expected
);
assert_eq!(
pb::FeedCommand::from(Kernel(&CompactFeedPrefix::new(
metadata(),
JournalPosition::new(4)
))),
expected
);
assert_eq!(
pb::FeedCommand::from(Kernel(&AcknowledgeProjectorCursor::new(
metadata(),
ConsumerId::new("search"),
Cursor::at(JournalPosition::new(2))
))),
expected
);
assert_eq!(
pb::FeedCommand::from(Kernel(&RegisterProjector::new(
metadata(),
ProjectorRegistration::new(
ConsumerId::new("search"),
PartitionId::new("conv-1"),
ConsumerPolicy::Required,
4,
)
))),
expected
);
}
#[test]
fn a_feed_entry_round_trips_its_envelope_and_every_record() {
let encoded = pb::FeedRecord::from(Kernel(&entry()));
let back = Kernel::<FeedRecord>::try_from(encoded).unwrap().0;
assert_eq!(back, entry());
assert!(back.is_self_consistent());
}
#[test]
fn a_chunk_round_trips_its_records_cursor_and_reason_for_stopping() {
for end in [
polyc_state::stream::StreamEnd::More,
polyc_state::stream::StreamEnd::Exhausted,
polyc_state::stream::StreamEnd::Drained,
] {
let cursor =
Cursor::in_snapshot(SnapshotId::new("feed:conv-1@2"), JournalPosition::new(3));
let chunk = FeedChunk::new(vec![entry()], Some(cursor), end);
let back = Kernel::<FeedChunk>::try_from(pb::FeedChunk::from(Kernel(&chunk)))
.unwrap()
.0;
assert_eq!(back, chunk);
}
let cursorless =
FeedChunk::new(Vec::new(), None, polyc_state::stream::StreamEnd::Exhausted);
let back = Kernel::<FeedChunk>::try_from(pb::FeedChunk::from(Kernel(&cursorless)))
.unwrap()
.0;
assert_eq!(back, cursorless);
}
#[test]
fn a_chunk_that_does_not_say_why_it_stopped_is_malformed() {
let mut encoded = pb::FeedChunk::from(Kernel(&FeedChunk::new(
Vec::new(),
None,
polyc_state::stream::StreamEnd::Exhausted,
)));
encoded.end = pb::StreamEnd::STREAM_END_UNSPECIFIED.into();
assert!(matches!(
Kernel::<FeedChunk>::try_from(encoded).unwrap_err(),
StateError::Malformed { ref field, .. } if field == "end"
));
}
#[test]
fn a_snapshot_round_trips_and_refuses_an_identity_that_contradicts_itself() {
let snapshot = FeedSnapshot::new(
FeedAnchor::new(PartitionId::new("conv-1"), JournalPosition::new(3)),
JournalPosition::new(9),
receipt(),
);
let encoded = pb::FeedSnapshot::from(Kernel(&snapshot));
assert_eq!(encoded.snapshot, "feed:conv-1@3");
let back = Kernel::<FeedSnapshot>::try_from(encoded.clone()).unwrap().0;
assert_eq!(back.id(), snapshot.id());
assert_eq!(back.position(), snapshot.position());
assert_eq!(back.head(), snapshot.head());
let mut lying = encoded.clone();
lying.position = 4;
assert!(matches!(
Kernel::<FeedSnapshot>::try_from(lying).unwrap_err(),
StateError::Malformed { ref field, .. } if field == "snapshot"
));
let mut renamed = encoded;
renamed.partition = "conv-2".to_owned();
assert!(matches!(
Kernel::<FeedSnapshot>::try_from(renamed).unwrap_err(),
StateError::Malformed { ref field, .. } if field == "snapshot"
));
}
#[test]
fn a_subscription_round_trips_its_partition_start_bound_and_consumer() {
let request = SubscribeCommits::new(
PartitionId::new("conv-1"),
ReadStart::Resume(Cursor::at(JournalPosition::new(4))),
8,
)
.on_behalf_of(ConsumerId::new("search"));
let (partition, start, limit, consumer) = subscription(&request);
let back = subscribe_request(
partition,
buffa::MessageField::<_, buffa::Inline<_>>::some(start),
limit,
consumer,
)
.unwrap();
assert_eq!(back, request);
let anonymous = SubscribeCommits::new(
PartitionId::new("conv-1"),
ReadStart::Snapshot(SnapshotId::new("feed:conv-1@2")),
8,
);
let (partition, start, limit, consumer) = subscription(&anonymous);
assert!(consumer.is_none());
let back = subscribe_request(
partition,
buffa::MessageField::<_, buffa::Inline<_>>::some(start),
limit,
consumer,
)
.unwrap();
assert_eq!(back, anonymous);
}
#[test]
fn a_subscription_without_a_partition_or_a_start_is_malformed() {
assert!(matches!(
subscribe_request(
String::new(),
buffa::MessageField::<_, buffa::Inline<_>>::some(pb::ReadStart::from(Kernel(
&ReadStart::Resume(Cursor::at(JournalPosition::ORIGIN))
))),
4,
None,
)
.unwrap_err(),
StateError::Malformed { ref field, .. } if field == "partition"
));
assert!(matches!(
subscribe_request(
"conv-1".to_owned(),
buffa::MessageField::<_, buffa::Inline<_>>::default(),
4,
None
)
.unwrap_err(),
StateError::Malformed { ref field, .. } if field == "start"
));
}
#[test]
fn a_projector_status_round_trips_including_its_eviction() {
for policy in [ConsumerPolicy::Required, ConsumerPolicy::Optional] {
let registration = ProjectorRegistration::new(
ConsumerId::new("search"),
PartitionId::new("conv-1"),
policy,
6,
);
let live = ProjectorStatus::new(registration, JournalPosition::new(2));
for status in [live.clone(), live.evicted()] {
let back =
Kernel::<ProjectorStatus>::try_from(pb::ProjectorStatus::from(Kernel(&status)))
.unwrap()
.0;
assert_eq!(back, status);
}
}
}
#[test]
fn a_registration_that_declares_no_policy_is_malformed() {
let mut encoded = pb::ProjectorRegistration::from(Kernel(&ProjectorRegistration::new(
ConsumerId::new("search"),
PartitionId::new("conv-1"),
ConsumerPolicy::Required,
6,
)));
encoded.policy = pb::ConsumerPolicy::CONSUMER_POLICY_UNSPECIFIED.into();
assert!(matches!(
Kernel::<ProjectorRegistration>::try_from(encoded).unwrap_err(),
StateError::Malformed { ref field, .. } if field == "policy"
));
}
#[test]
fn retention_round_trips_with_and_without_the_consumer_holding_it() {
let bare = FeedRetention::new(
PartitionId::new("conv-1"),
JournalPosition::new(2),
JournalPosition::new(9),
);
for retention in [
bare.clone(),
bare.clone().held_by(ConsumerId::new("search")),
bare.clone().pinned_by(SnapshotId::new("feed:conv-1@2")),
bare.held_by(ConsumerId::new("search"))
.pinned_by(SnapshotId::new("feed:conv-1@2")),
] {
let back = Kernel::<FeedRetention>::from(pb::FeedRetention::from(Kernel(&retention))).0;
assert_eq!(back, retention);
}
}
#[test]
fn a_compaction_round_trips_with_and_without_a_receipt() {
let retention = FeedRetention::new(
PartitionId::new("conv-1"),
JournalPosition::new(4),
JournalPosition::new(9),
);
let held_back = FeedCompaction::new(retention.clone(), Vec::new(), None);
let back = Kernel::<FeedCompaction>::try_from(pb::FeedCompaction::from(Kernel(&held_back)))
.unwrap()
.0;
assert_eq!(back, held_back);
assert!(!back.is_recorded());
let recorded = FeedCompaction::new(
retention,
vec![ConsumerId::new("search"), ConsumerId::new("audit")],
Some(receipt()),
);
let back = Kernel::<FeedCompaction>::try_from(pb::FeedCompaction::from(Kernel(&recorded)))
.unwrap()
.0;
assert!(back.is_recorded());
assert_eq!(back.evicted().len(), 2);
assert_eq!(back.retention(), recorded.retention());
}
}