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)]
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}