Skip to main content

ursula_runtime/
command.rs

1use std::fmt;
2
3use serde::Deserialize;
4use serde::Serialize;
5use ursula_shard::ShardPlacement;
6use ursula_stream::StreamCommand;
7use ursula_stream::StreamSnapshot;
8
9use crate::request::AdvanceRetentionRequest;
10use crate::request::AppendBatchRequest;
11use crate::request::AppendExternalRequest;
12use crate::request::AppendRequest;
13use crate::request::CloseStreamRequest;
14use crate::request::CompactColdRequest;
15use crate::request::CreateStreamExternalRequest;
16use crate::request::CreateStreamRequest;
17use crate::request::DeleteStreamRequest;
18use crate::request::FlushColdRequest;
19use crate::request::PublishSnapshotRequest;
20use crate::request::SetBucketQuotaRequest;
21use crate::request::StreamAppendCount;
22use crate::request::UpdateStreamAttrsRequest;
23
24#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
25pub struct GroupSnapshot {
26    pub placement: ShardPlacement,
27    pub group_commit_index: u64,
28    pub stream_snapshot: StreamSnapshot,
29    pub stream_append_counts: Vec<StreamAppendCount>,
30}
31
32/// Replicated group-level write envelope around the canonical
33/// [`StreamCommand`]: one command, a throughput batch with independent item
34/// results, or an all-or-none append transaction. Batch and transaction
35/// envelopes each occupy one Raft entry. This enum (serde-encoded) is the Raft
36/// log payload; there is no separate wire mirror.
37#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
38#[expect(
39    clippy::large_enum_variant,
40    reason = "commands are transient write-path values; boxing `Stream` would \
41              add a heap allocation to every replicated write"
42)]
43pub enum GroupWriteCommand {
44    Stream(StreamCommand),
45    Batch { commands: Vec<StreamCommand> },
46    Transaction { commands: Vec<StreamCommand> },
47}
48
49impl From<StreamCommand> for GroupWriteCommand {
50    fn from(command: StreamCommand) -> Self {
51        Self::Stream(command)
52    }
53}
54
55impl From<crate::request::ImportGroupStateRequest> for StreamCommand {
56    fn from(request: crate::request::ImportGroupStateRequest) -> Self {
57        Self::ImportSnapshot {
58            snapshot: request.snapshot,
59        }
60    }
61}
62
63impl From<CreateStreamRequest> for StreamCommand {
64    fn from(request: CreateStreamRequest) -> Self {
65        Self::CreateStream {
66            stream_id: request.stream_id,
67            content_type: request.content_type,
68            initial_payload: request.initial_payload,
69            close_after: request.close_after,
70            stream_seq: request.stream_seq,
71            producer: request.producer,
72            stream_ttl_seconds: request.stream_ttl_seconds,
73            stream_expires_at_ms: request.stream_expires_at_ms,
74            attrs: request.attrs,
75            now_ms: request.now_ms,
76        }
77    }
78}
79
80impl From<CreateStreamExternalRequest> for StreamCommand {
81    fn from(request: CreateStreamExternalRequest) -> Self {
82        Self::CreateExternal {
83            stream_id: request.stream_id,
84            content_type: request.content_type,
85            initial_payload: request.initial_payload,
86            record_ends: request.record_ends,
87            close_after: request.close_after,
88            stream_seq: request.stream_seq,
89            producer: request.producer,
90            stream_ttl_seconds: request.stream_ttl_seconds,
91            stream_expires_at_ms: request.stream_expires_at_ms,
92            attrs: request.attrs,
93            now_ms: request.now_ms,
94        }
95    }
96}
97
98impl From<UpdateStreamAttrsRequest> for StreamCommand {
99    fn from(request: UpdateStreamAttrsRequest) -> Self {
100        Self::UpdateStreamAttrs {
101            stream_id: request.stream_id,
102            attrs: request.attrs,
103            now_ms: request.now_ms,
104        }
105    }
106}
107
108impl From<AppendRequest> for StreamCommand {
109    fn from(request: AppendRequest) -> Self {
110        Self::Append {
111            stream_id: request.stream_id,
112            content_type: Some(request.content_type),
113            payload: request.payload,
114            close_after: request.close_after,
115            stream_seq: request.stream_seq,
116            producer: request.producer,
117            now_ms: request.now_ms,
118            record_match: request.record_match,
119        }
120    }
121}
122
123impl From<AppendExternalRequest> for StreamCommand {
124    fn from(request: AppendExternalRequest) -> Self {
125        Self::AppendExternal {
126            stream_id: request.stream_id,
127            content_type: Some(request.content_type),
128            payload: request.payload,
129            record_ends: request.record_ends,
130            close_after: request.close_after,
131            stream_seq: request.stream_seq,
132            producer: request.producer,
133            now_ms: request.now_ms,
134            record_match: request.record_match,
135        }
136    }
137}
138
139impl From<AppendBatchRequest> for StreamCommand {
140    fn from(request: AppendBatchRequest) -> Self {
141        Self::AppendBatch {
142            stream_id: request.stream_id,
143            content_type: Some(request.content_type),
144            payloads: request.payloads,
145            producer: request.producer,
146            now_ms: request.now_ms,
147        }
148    }
149}
150
151impl From<PublishSnapshotRequest> for StreamCommand {
152    fn from(request: PublishSnapshotRequest) -> Self {
153        Self::PublishSnapshot {
154            stream_id: request.stream_id,
155            snapshot_offset: request.snapshot_offset,
156            content_type: request.content_type,
157            payload: request.payload,
158            expected_digest: request.expected_digest,
159            now_ms: request.now_ms,
160        }
161    }
162}
163
164impl From<AdvanceRetentionRequest> for StreamCommand {
165    fn from(request: AdvanceRetentionRequest) -> Self {
166        Self::AdvanceRetention {
167            stream_id: request.stream_id,
168            retained_offset: request.retained_offset,
169            now_ms: request.now_ms,
170        }
171    }
172}
173
174impl From<SetBucketQuotaRequest> for StreamCommand {
175    fn from(request: SetBucketQuotaRequest) -> Self {
176        Self::SetBucketQuota {
177            bucket_id: request.bucket_id,
178            max_streams: request.max_streams,
179            max_retained_bytes: request.max_retained_bytes,
180        }
181    }
182}
183
184impl From<CloseStreamRequest> for StreamCommand {
185    fn from(request: CloseStreamRequest) -> Self {
186        Self::Close {
187            stream_id: request.stream_id,
188            stream_seq: request.stream_seq,
189            producer: request.producer,
190            now_ms: request.now_ms,
191        }
192    }
193}
194
195impl From<DeleteStreamRequest> for StreamCommand {
196    fn from(request: DeleteStreamRequest) -> Self {
197        Self::DeleteStream {
198            stream_id: request.stream_id,
199        }
200    }
201}
202
203impl From<FlushColdRequest> for StreamCommand {
204    fn from(request: FlushColdRequest) -> Self {
205        Self::FlushCold {
206            stream_id: request.stream_id,
207            chunk: request.chunk,
208        }
209    }
210}
211
212impl From<CompactColdRequest> for StreamCommand {
213    fn from(request: CompactColdRequest) -> Self {
214        Self::CompactCold {
215            stream_id: request.stream_id,
216            old_chunks: request.old_chunks,
217            replacement: request.replacement,
218            gc_not_before_ms: request.gc_not_before_ms,
219        }
220    }
221}
222
223macro_rules! group_write_from_request {
224    ($($request:ty),+ $(,)?) => {
225        $(impl From<$request> for GroupWriteCommand {
226            fn from(request: $request) -> Self {
227                Self::Stream(StreamCommand::from(request))
228            }
229        })+
230    };
231}
232
233group_write_from_request!(
234    CreateStreamRequest,
235    CreateStreamExternalRequest,
236    UpdateStreamAttrsRequest,
237    AppendRequest,
238    AppendExternalRequest,
239    AppendBatchRequest,
240    PublishSnapshotRequest,
241    AdvanceRetentionRequest,
242    SetBucketQuotaRequest,
243    CloseStreamRequest,
244    DeleteStreamRequest,
245    FlushColdRequest,
246    CompactColdRequest,
247);
248
249impl fmt::Display for GroupWriteCommand {
250    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
251        match self {
252            Self::Stream(command) => command.fmt(f),
253            Self::Batch { commands } => write!(f, "batch:{} commands", commands.len()),
254            Self::Transaction { commands } => {
255                write!(f, "transaction:{} commands", commands.len())
256            }
257        }
258    }
259}