1use aequora_types::{
4 ActorId, AuthorityEpoch, AuthorityId, Cursor, DeviceId, EntityRef, EntityVersion, EventId,
5 HybridTimestamp, LineageContext, OperationId, ProtocolVersion, RegionId, RequestId,
6 SchemaVersion, Sequence, SessionId, SnapshotId, SyncScopeId, TenantId,
7};
8use serde::{Deserialize, Serialize};
9use smallvec::SmallVec;
10
11pub mod wire_limits {
14 use aequora_types::OperationId;
15 use serde::{
16 Deserialize,
17 de::{self, Deserializer, SeqAccess, Visitor},
18 };
19 use smallvec::SmallVec;
20 use std::{fmt, marker::PhantomData};
21
22 pub const OPERATIONS: usize = 4_096;
24 pub const DEPENDENCIES: usize = 1_024;
26 pub const PARTITIONS: usize = 512;
28 pub const CAPABILITIES: usize = 64;
30 pub const PAYLOAD_BYTES: usize = 16 * 1_024 * 1_024;
32 pub const RESULTS: usize = 8_192;
34 pub const SNAPSHOT_ENTITIES: usize = 8_192;
36
37 struct BoundedSequence<T, const MAX: usize>(PhantomData<T>);
38
39 impl<'de, T, const MAX: usize> Visitor<'de> for BoundedSequence<T, MAX>
40 where
41 T: Deserialize<'de>,
42 {
43 type Value = Vec<T>;
44
45 fn expecting(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
46 write!(formatter, "a sequence containing at most {MAX} elements")
47 }
48
49 fn visit_seq<A>(self, mut sequence: A) -> Result<Self::Value, A::Error>
50 where
51 A: SeqAccess<'de>,
52 {
53 if sequence.size_hint().is_some_and(|length| length > MAX) {
54 return Err(de::Error::invalid_length(MAX.saturating_add(1), &self));
55 }
56 let mut items = Vec::with_capacity(sequence.size_hint().unwrap_or(0).min(MAX));
57 while let Some(item) = sequence.next_element()? {
58 if items.len() == MAX {
59 return Err(de::Error::invalid_length(MAX.saturating_add(1), &self));
60 }
61 items.push(item);
62 }
63 Ok(items)
64 }
65 }
66
67 fn bounded_vec<'de, D, T, const MAX: usize>(deserializer: D) -> Result<Vec<T>, D::Error>
68 where
69 D: Deserializer<'de>,
70 T: Deserialize<'de>,
71 {
72 deserializer.deserialize_seq(BoundedSequence::<T, MAX>(PhantomData))
73 }
74
75 pub(crate) fn operations<'de, D, T>(deserializer: D) -> Result<Vec<T>, D::Error>
76 where
77 D: Deserializer<'de>,
78 T: Deserialize<'de>,
79 {
80 bounded_vec::<D, T, OPERATIONS>(deserializer)
81 }
82
83 pub(crate) fn partitions<'de, D, T>(deserializer: D) -> Result<Vec<T>, D::Error>
84 where
85 D: Deserializer<'de>,
86 T: Deserialize<'de>,
87 {
88 bounded_vec::<D, T, PARTITIONS>(deserializer)
89 }
90
91 pub(crate) fn capabilities<'de, D, T>(deserializer: D) -> Result<Vec<T>, D::Error>
92 where
93 D: Deserializer<'de>,
94 T: Deserialize<'de>,
95 {
96 bounded_vec::<D, T, CAPABILITIES>(deserializer)
97 }
98
99 pub(crate) fn payload<'de, D>(deserializer: D) -> Result<Vec<u8>, D::Error>
100 where
101 D: Deserializer<'de>,
102 {
103 bounded_vec::<D, u8, PAYLOAD_BYTES>(deserializer)
104 }
105
106 pub(crate) fn results<'de, D, T>(deserializer: D) -> Result<Vec<T>, D::Error>
107 where
108 D: Deserializer<'de>,
109 T: Deserialize<'de>,
110 {
111 bounded_vec::<D, T, RESULTS>(deserializer)
112 }
113
114 pub(crate) fn snapshot_entities<'de, D, T>(deserializer: D) -> Result<Vec<T>, D::Error>
115 where
116 D: Deserializer<'de>,
117 T: Deserialize<'de>,
118 {
119 bounded_vec::<D, T, SNAPSHOT_ENTITIES>(deserializer)
120 }
121
122 struct BoundedDependencies;
123
124 impl<'de> Visitor<'de> for BoundedDependencies {
125 type Value = SmallVec<[OperationId; 4]>;
126
127 fn expecting(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
128 write!(
129 formatter,
130 "a dependency sequence containing at most {DEPENDENCIES} elements"
131 )
132 }
133
134 fn visit_seq<A>(self, mut sequence: A) -> Result<Self::Value, A::Error>
135 where
136 A: SeqAccess<'de>,
137 {
138 if sequence
139 .size_hint()
140 .is_some_and(|length| length > DEPENDENCIES)
141 {
142 return Err(de::Error::invalid_length(
143 DEPENDENCIES.saturating_add(1),
144 &self,
145 ));
146 }
147 let mut items = SmallVec::new();
148 while let Some(item) = sequence.next_element()? {
149 if items.len() == DEPENDENCIES {
150 return Err(de::Error::invalid_length(
151 DEPENDENCIES.saturating_add(1),
152 &self,
153 ));
154 }
155 items.push(item);
156 }
157 Ok(items)
158 }
159 }
160
161 pub(crate) fn dependencies<'de, D>(
162 deserializer: D,
163 ) -> Result<SmallVec<[OperationId; 4]>, D::Error>
164 where
165 D: Deserializer<'de>,
166 {
167 deserializer.deserialize_seq(BoundedDependencies)
168 }
169
170 #[cfg(test)]
171 mod tests {
172 use super::*;
173 use serde::{Deserialize, Serialize};
174
175 #[derive(Debug, Deserialize, Eq, PartialEq, Serialize)]
176 struct TinySequence(#[serde(deserialize_with = "tiny")] Vec<u8>);
177
178 fn tiny<'de, D>(deserializer: D) -> Result<Vec<u8>, D::Error>
179 where
180 D: Deserializer<'de>,
181 {
182 bounded_vec::<D, u8, 2>(deserializer)
183 }
184
185 #[test]
186 fn declared_collection_length_is_rejected_by_the_deserializer() {
187 let encoded = postcard::to_stdvec(&TinySequence(vec![1, 2, 3]))
188 .unwrap_or_else(|error| panic!("{error}"));
189 let decoded = postcard::from_bytes::<TinySequence>(&encoded);
190 assert!(decoded.is_err());
191 }
192 }
193}
194
195#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, PartialEq, Serialize)]
197#[serde(transparent)]
198pub struct OperationKind(pub u16);
199
200#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
202pub struct OperationMetadata {
203 pub trace_id: Option<String>,
205 #[serde(deserialize_with = "wire_limits::dependencies")]
207 pub dependencies: SmallVec<[OperationId; 4]>,
208 #[serde(default = "legacy_lineage")]
210 pub lineage: LineageContext,
211}
212
213fn legacy_lineage() -> LineageContext {
214 LineageContext::legacy_missing()
215}
216
217#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
219pub struct OperationEnvelope {
220 pub protocol_version: ProtocolVersion,
222 pub operation_id: OperationId,
224 pub tenant_id: TenantId,
226 pub actor_id: ActorId,
228 pub device_id: DeviceId,
230 pub entity: EntityRef,
232 pub base_version: Option<EntityVersion>,
234 pub created_at: HybridTimestamp,
236 pub schema_version: SchemaVersion,
238 pub operation_kind: OperationKind,
240 #[serde(deserialize_with = "wire_limits::payload")]
242 pub payload: Vec<u8>,
243 pub metadata: OperationMetadata,
245}
246
247#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
249pub struct SessionMetadata {
250 pub session_id: SessionId,
252 pub device_id: DeviceId,
254 pub actor_id: ActorId,
256 pub tenant_id: TenantId,
258 pub scope_id: SyncScopeId,
260 #[serde(deserialize_with = "wire_limits::partitions")]
262 pub partitions: Vec<Partition>,
263}
264
265#[derive(Clone, Debug, Deserialize, Eq, Hash, PartialEq, Serialize)]
267pub struct Partition {
268 pub kind: u16,
270 #[serde(deserialize_with = "wire_limits::payload")]
272 pub value: Vec<u8>,
273}
274
275#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, PartialEq, Serialize)]
277#[non_exhaustive]
278pub enum Capability {
279 PostcardV1,
281 Zstd,
283 SnapshotV1,
285 Tombstones,
287 StreamingSnapshots,
289 PushHints,
291 Quic,
293 MultiRegion,
295 LineageV1,
297 IntegrityV1,
299 ScopeV1,
301 LiveV1,
303 SignedSnapshotV1,
305 EncryptedSnapshotV1,
307 DeviceSignatureV1,
309 AuthorityEpochV1,
311 ResourceConstrainedV1,
313 CompatibilityNegotiationV1,
315}
316
317impl Capability {
318 #[must_use]
320 pub const fn stable_id(self) -> u32 {
321 match self {
322 Self::PostcardV1 => 1,
323 Self::Zstd => 2,
324 Self::SnapshotV1 => 3,
325 Self::Tombstones => 4,
326 Self::StreamingSnapshots => 5,
327 Self::PushHints => 6,
328 Self::Quic => 7,
329 Self::MultiRegion => 8,
330 Self::LineageV1 => 9,
331 Self::IntegrityV1 => 10,
332 Self::ScopeV1 => 11,
333 Self::LiveV1 => 12,
334 Self::SignedSnapshotV1 => 13,
335 Self::EncryptedSnapshotV1 => 14,
336 Self::DeviceSignatureV1 => 15,
337 Self::AuthorityEpochV1 => 16,
338 Self::ResourceConstrainedV1 => 17,
339 Self::CompatibilityNegotiationV1 => 18,
340 }
341 }
342}
343
344#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
346pub struct ClientLimits {
347 pub max_changes: u32,
349 pub max_response_bytes: u32,
351}
352
353impl Default for ClientLimits {
354 fn default() -> Self {
355 Self {
356 max_changes: 1_024,
357 max_response_bytes: 4 * 1_024 * 1_024,
358 }
359 }
360}
361
362#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
364pub struct SyncRequest {
365 pub protocol: ProtocolVersion,
367 pub request_id: RequestId,
369 pub session: SessionMetadata,
371 pub cursor: Option<Cursor>,
373 #[serde(deserialize_with = "wire_limits::operations")]
375 pub operations: Vec<OperationEnvelope>,
376 pub limits: ClientLimits,
378 #[serde(deserialize_with = "wire_limits::capabilities")]
380 pub capabilities: Vec<Capability>,
381}
382
383#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
385pub struct OperationAck {
386 pub operation_id: OperationId,
388 pub event_id: EventId,
390 pub lineage: LineageContext,
392 pub entity_version: EntityVersion,
394 pub sequence: Sequence,
396 pub duplicate: bool,
398}
399
400#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
402#[non_exhaustive]
403pub enum RejectionCode {
404 IdentityMismatch,
406 Unauthorized,
408 InvalidOperation,
410 BusinessRule,
412 Dependency,
414 SchemaIncompatible,
416}
417
418#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
420pub struct OperationRejection {
421 pub operation_id: OperationId,
423 pub code: RejectionCode,
425 pub message: String,
427}
428
429#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
431#[non_exhaustive]
432pub enum ConflictPolicy {
433 Reject,
435 ServerWins,
437 ClientWins,
439 CustomMerge,
441 ManualResolution,
443 FieldMerge,
445 CommutativeOperation,
447 Crdt,
449 LastWriterWins,
451}
452
453#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
455pub struct Conflict {
456 pub operation_id: OperationId,
458 pub entity: EntityRef,
460 pub client_base: Option<EntityVersion>,
462 pub server_version: Option<EntityVersion>,
464 pub policy: ConflictPolicy,
466 pub message: String,
468}
469
470#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
472pub enum ChangeKind {
473 Upsert,
475 Tombstone,
477}
478
479#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
481pub struct RemoteChange {
482 pub tenant_id: TenantId,
484 pub scope_id: SyncScopeId,
486 pub sequence: Sequence,
488 pub operation_id: OperationId,
490 pub event_id: EventId,
492 pub lineage: LineageContext,
494 pub entity: EntityRef,
496 pub version: EntityVersion,
498 pub change_kind: ChangeKind,
500 #[serde(deserialize_with = "wire_limits::payload")]
502 pub payload: Vec<u8>,
503 pub timestamp: HybridTimestamp,
505}
506
507#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
509#[non_exhaustive]
510pub enum ResyncReason {
511 CursorExpired,
513 ScopeChanged,
515 SchemaIncompatible,
517 DeviceInactive,
519 CorruptionDetected,
521}
522
523#[derive(Clone, Copy, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
525pub enum SyncDirective {
526 #[default]
528 Continue,
529 UpgradeRequired {
531 minimum: ProtocolVersion,
533 current: ProtocolVersion,
535 },
536 ResyncRequired {
538 reason: ResyncReason,
540 },
541 AuthorityChanged {
543 authority_id: AuthorityId,
545 previous_epoch: AuthorityEpoch,
547 current_epoch: AuthorityEpoch,
549 },
550}
551
552#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
554pub struct SyncResponse {
555 pub protocol: ProtocolVersion,
557 pub directive: SyncDirective,
559 #[serde(deserialize_with = "wire_limits::results")]
561 pub acknowledged: Vec<OperationAck>,
562 #[serde(deserialize_with = "wire_limits::results")]
564 pub rejected: Vec<OperationRejection>,
565 #[serde(deserialize_with = "wire_limits::results")]
567 pub conflicts: Vec<Conflict>,
568 #[serde(deserialize_with = "wire_limits::results")]
570 pub changes: Vec<RemoteChange>,
571 pub next_cursor: Cursor,
573 pub has_more: bool,
575 pub server_time: HybridTimestamp,
577}
578
579#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
581pub struct SnapshotLimits {
582 pub max_entities: u32,
584 pub max_payload_bytes: u32,
586}
587
588impl Default for SnapshotLimits {
589 fn default() -> Self {
590 Self {
591 max_entities: 512,
592 max_payload_bytes: 4 * 1_024 * 1_024,
593 }
594 }
595}
596
597#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
599pub struct BootstrapRequest {
600 pub protocol: ProtocolVersion,
602 pub request_id: RequestId,
604 pub session: SessionMetadata,
606 pub snapshot_id: Option<SnapshotId>,
608 pub offset: u64,
610 pub limits: SnapshotLimits,
612 #[serde(deserialize_with = "wire_limits::capabilities")]
614 pub capabilities: Vec<Capability>,
615}
616
617#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
619pub struct SnapshotEntity {
620 pub entity: EntityRef,
622 pub version: EntityVersion,
624 #[serde(deserialize_with = "wire_limits::payload")]
626 pub payload: Vec<u8>,
627 pub tombstone: bool,
629}
630
631#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
633pub struct BootstrapResponse {
634 pub protocol: ProtocolVersion,
636 pub snapshot_id: SnapshotId,
638 pub cursor: Cursor,
640 pub offset: u64,
642 #[serde(deserialize_with = "wire_limits::snapshot_entities")]
644 pub entities: Vec<SnapshotEntity>,
645 pub next_offset: u64,
647 pub has_more: bool,
649 pub server_time: HybridTimestamp,
651}
652
653#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
655#[non_exhaustive]
656pub enum PushHintReason {
657 JournalAdvanced,
659 SnapshotInvalidated,
661 RegionChanged,
663}
664
665#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
668pub struct PushHint {
669 pub protocol: ProtocolVersion,
671 pub tenant_id: TenantId,
673 pub scope_id: SyncScopeId,
675 pub sequence: Sequence,
677 pub reason: PushHintReason,
679 pub region_id: Option<RegionId>,
681}
682
683#[cfg(test)]
684mod compatibility_tests {
685 use super::*;
686
687 #[test]
688 fn conflict_policy_wire_discriminants_remain_append_only() {
689 let policies = [
690 ConflictPolicy::Reject,
691 ConflictPolicy::ServerWins,
692 ConflictPolicy::ClientWins,
693 ConflictPolicy::CustomMerge,
694 ConflictPolicy::ManualResolution,
695 ConflictPolicy::FieldMerge,
696 ConflictPolicy::CommutativeOperation,
697 ConflictPolicy::Crdt,
698 ConflictPolicy::LastWriterWins,
699 ];
700 for (discriminant, policy) in policies.into_iter().enumerate() {
701 assert_eq!(
702 postcard::to_stdvec(&policy).unwrap_or_else(|error| panic!("{error}")),
703 vec![u8::try_from(discriminant).unwrap_or(u8::MAX)]
704 );
705 }
706 }
707
708 #[test]
709 fn capability_wire_discriminants_and_registry_ids_remain_append_only() {
710 let capabilities = [
711 Capability::PostcardV1,
712 Capability::Zstd,
713 Capability::SnapshotV1,
714 Capability::Tombstones,
715 Capability::StreamingSnapshots,
716 Capability::PushHints,
717 Capability::Quic,
718 Capability::MultiRegion,
719 Capability::LineageV1,
720 Capability::IntegrityV1,
721 Capability::ScopeV1,
722 Capability::LiveV1,
723 Capability::SignedSnapshotV1,
724 Capability::EncryptedSnapshotV1,
725 Capability::DeviceSignatureV1,
726 Capability::AuthorityEpochV1,
727 Capability::ResourceConstrainedV1,
728 Capability::CompatibilityNegotiationV1,
729 ];
730 for (discriminant, capability) in capabilities.into_iter().enumerate() {
731 assert_eq!(
732 postcard::to_stdvec(&capability).unwrap_or_else(|error| panic!("{error}")),
733 vec![u8::try_from(discriminant).unwrap_or(u8::MAX)]
734 );
735 assert_eq!(
736 capability.stable_id(),
737 u32::try_from(discriminant).unwrap_or(u32::MAX) + 1
738 );
739 }
740 }
741}