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`]: either one per-stream command, or an atomic batch of
34/// them applied as a single raft entry. This enum (serde-encoded) is the raft
35/// log payload; there is no separate wire mirror.
36#[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}