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#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
37#[expect(
38 clippy::large_enum_variant,
39 reason = "commands are transient write-path values; boxing `Stream` would \
40 add a heap allocation to every replicated write"
41)]
42pub enum GroupWriteCommand {
43 Stream(StreamCommand),
44 Batch { commands: Vec<StreamCommand> },
45}
46
47impl From<StreamCommand> for GroupWriteCommand {
48 fn from(command: StreamCommand) -> Self {
49 Self::Stream(command)
50 }
51}
52
53impl From<crate::request::ImportGroupStateRequest> for StreamCommand {
54 fn from(request: crate::request::ImportGroupStateRequest) -> Self {
55 Self::ImportSnapshot {
56 snapshot: request.snapshot,
57 }
58 }
59}
60
61impl From<CreateStreamRequest> for StreamCommand {
62 fn from(request: CreateStreamRequest) -> Self {
63 Self::CreateStream {
64 stream_id: request.stream_id,
65 content_type: request.content_type,
66 initial_payload: request.initial_payload,
67 close_after: request.close_after,
68 stream_seq: request.stream_seq,
69 producer: request.producer,
70 stream_ttl_seconds: request.stream_ttl_seconds,
71 stream_expires_at_ms: request.stream_expires_at_ms,
72 attrs: request.attrs,
73 now_ms: request.now_ms,
74 }
75 }
76}
77
78impl From<CreateStreamExternalRequest> for StreamCommand {
79 fn from(request: CreateStreamExternalRequest) -> Self {
80 Self::CreateExternal {
81 stream_id: request.stream_id,
82 content_type: request.content_type,
83 initial_payload: request.initial_payload,
84 record_ends: request.record_ends,
85 close_after: request.close_after,
86 stream_seq: request.stream_seq,
87 producer: request.producer,
88 stream_ttl_seconds: request.stream_ttl_seconds,
89 stream_expires_at_ms: request.stream_expires_at_ms,
90 attrs: request.attrs,
91 now_ms: request.now_ms,
92 }
93 }
94}
95
96impl From<UpdateStreamAttrsRequest> for StreamCommand {
97 fn from(request: UpdateStreamAttrsRequest) -> Self {
98 Self::UpdateStreamAttrs {
99 stream_id: request.stream_id,
100 attrs: request.attrs,
101 now_ms: request.now_ms,
102 }
103 }
104}
105
106impl From<AppendRequest> for StreamCommand {
107 fn from(request: AppendRequest) -> Self {
108 Self::Append {
109 stream_id: request.stream_id,
110 content_type: Some(request.content_type),
111 payload: request.payload,
112 close_after: request.close_after,
113 stream_seq: request.stream_seq,
114 producer: request.producer,
115 now_ms: request.now_ms,
116 record_match: request.record_match,
117 }
118 }
119}
120
121impl From<AppendExternalRequest> for StreamCommand {
122 fn from(request: AppendExternalRequest) -> Self {
123 Self::AppendExternal {
124 stream_id: request.stream_id,
125 content_type: Some(request.content_type),
126 payload: request.payload,
127 record_ends: request.record_ends,
128 close_after: request.close_after,
129 stream_seq: request.stream_seq,
130 producer: request.producer,
131 now_ms: request.now_ms,
132 record_match: request.record_match,
133 }
134 }
135}
136
137impl From<AppendBatchRequest> for StreamCommand {
138 fn from(request: AppendBatchRequest) -> Self {
139 Self::AppendBatch {
140 stream_id: request.stream_id,
141 content_type: Some(request.content_type),
142 payloads: request.payloads,
143 producer: request.producer,
144 now_ms: request.now_ms,
145 }
146 }
147}
148
149impl From<PublishSnapshotRequest> for StreamCommand {
150 fn from(request: PublishSnapshotRequest) -> Self {
151 Self::PublishSnapshot {
152 stream_id: request.stream_id,
153 snapshot_offset: request.snapshot_offset,
154 content_type: request.content_type,
155 payload: request.payload,
156 expected_digest: request.expected_digest,
157 now_ms: request.now_ms,
158 }
159 }
160}
161
162impl From<AdvanceRetentionRequest> for StreamCommand {
163 fn from(request: AdvanceRetentionRequest) -> Self {
164 Self::AdvanceRetention {
165 stream_id: request.stream_id,
166 retained_offset: request.retained_offset,
167 now_ms: request.now_ms,
168 }
169 }
170}
171
172impl From<SetBucketQuotaRequest> for StreamCommand {
173 fn from(request: SetBucketQuotaRequest) -> Self {
174 Self::SetBucketQuota {
175 bucket_id: request.bucket_id,
176 max_streams: request.max_streams,
177 max_retained_bytes: request.max_retained_bytes,
178 }
179 }
180}
181
182impl From<CloseStreamRequest> for StreamCommand {
183 fn from(request: CloseStreamRequest) -> Self {
184 Self::Close {
185 stream_id: request.stream_id,
186 stream_seq: request.stream_seq,
187 producer: request.producer,
188 now_ms: request.now_ms,
189 }
190 }
191}
192
193impl From<DeleteStreamRequest> for StreamCommand {
194 fn from(request: DeleteStreamRequest) -> Self {
195 Self::DeleteStream {
196 stream_id: request.stream_id,
197 }
198 }
199}
200
201impl From<FlushColdRequest> for StreamCommand {
202 fn from(request: FlushColdRequest) -> Self {
203 Self::FlushCold {
204 stream_id: request.stream_id,
205 chunk: request.chunk,
206 }
207 }
208}
209
210impl From<CompactColdRequest> for StreamCommand {
211 fn from(request: CompactColdRequest) -> Self {
212 Self::CompactCold {
213 stream_id: request.stream_id,
214 old_chunks: request.old_chunks,
215 replacement: request.replacement,
216 gc_not_before_ms: request.gc_not_before_ms,
217 }
218 }
219}
220
221macro_rules! group_write_from_request {
222 ($($request:ty),+ $(,)?) => {
223 $(impl From<$request> for GroupWriteCommand {
224 fn from(request: $request) -> Self {
225 Self::Stream(StreamCommand::from(request))
226 }
227 })+
228 };
229}
230
231group_write_from_request!(
232 CreateStreamRequest,
233 CreateStreamExternalRequest,
234 UpdateStreamAttrsRequest,
235 AppendRequest,
236 AppendExternalRequest,
237 AppendBatchRequest,
238 PublishSnapshotRequest,
239 AdvanceRetentionRequest,
240 SetBucketQuotaRequest,
241 CloseStreamRequest,
242 DeleteStreamRequest,
243 FlushColdRequest,
244 CompactColdRequest,
245);
246
247impl fmt::Display for GroupWriteCommand {
248 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
249 match self {
250 Self::Stream(command) => command.fmt(f),
251 Self::Batch { commands } => write!(f, "batch:{} commands", commands.len()),
252 }
253 }
254}