krafka/admin/share_group_offsets.rs
1//! AdminClient operation group: share-group offsets (KIP-932 / KIP-1226).
2//!
3//! krafka has shipped a [`ShareConsumer`](crate::share_consumer::ShareConsumer)
4//! since share groups landed, but until these operations existed a share group
5//! could be *run* and not *operated*: there was no way to read its
6//! share-partition start offsets, reset them, or clean up state for a retired
7//! topic. Those are the three things an on-call engineer needs at 3 a.m., and
8//! all three live behind API keys 90–92.
9
10use tracing::{info, warn};
11
12use crate::error::{KrafkaError, ProtocolErrorKind, Result};
13use crate::protocol::{
14 AlterShareGroupOffsetsRequest, AlterShareGroupOffsetsRequestPartition,
15 AlterShareGroupOffsetsRequestTopic, AlterShareGroupOffsetsResponse,
16 DeleteShareGroupOffsetsRequest, DeleteShareGroupOffsetsResponse,
17 DescribeShareGroupOffsetsRequest, DescribeShareGroupOffsetsRequestGroup,
18 DescribeShareGroupOffsetsRequestTopic, DescribeShareGroupOffsetsResponse, VersionedDecode,
19 VersionedEncode, validate_topic_name, versions,
20};
21
22#[allow(clippy::wildcard_imports)]
23use super::*;
24
25/// One share partition's offset state, from
26/// [`AdminClient::describe_share_group_offsets`].
27#[non_exhaustive]
28#[derive(Debug, Clone)]
29pub struct ShareGroupPartitionOffset {
30 /// Topic name.
31 pub topic: String,
32 /// Partition index.
33 pub partition: i32,
34 /// Share-partition start offset — the earliest offset the group may still
35 /// deliver.
36 pub start_offset: i64,
37 /// Leader epoch of the partition.
38 pub leader_epoch: i32,
39 /// Share-partition lag, or `None` when the broker did not report it
40 /// (KIP-1226 added `Lag` in `DescribeShareGroupOffsets` v1; Kafka 4.2 and
41 /// earlier answer at v0).
42 pub lag: Option<i64>,
43 /// Partition-level error, or `None` on success.
44 pub error: Option<String>,
45}
46
47/// Result of [`AdminClient::describe_share_group_offsets`].
48#[non_exhaustive]
49#[derive(Debug, Clone)]
50pub struct DescribeShareGroupOffsetsResult {
51 /// Share group identifier.
52 pub group_id: String,
53 /// Group-level error, or `None` on success. When set, `partitions` is
54 /// typically empty.
55 pub error: Option<String>,
56 /// Per-partition offset state.
57 pub partitions: Vec<ShareGroupPartitionOffset>,
58}
59
60/// One partition's outcome from
61/// [`AdminClient::alter_share_group_offsets`].
62#[non_exhaustive]
63#[derive(Debug, Clone)]
64pub struct ShareGroupOffsetAlteration {
65 /// Topic name.
66 pub topic: String,
67 /// Partition index.
68 pub partition: i32,
69 /// Partition-level error, or `None` on success.
70 pub error: Option<String>,
71}
72
73/// One topic's outcome from
74/// [`AdminClient::delete_share_group_offsets`].
75#[non_exhaustive]
76#[derive(Debug, Clone)]
77pub struct ShareGroupOffsetDeletion {
78 /// Topic name.
79 pub topic: String,
80 /// Topic-level error, or `None` on success.
81 pub error: Option<String>,
82}
83
84impl AdminClient {
85 /// Read a share group's share-partition start offsets (KIP-932).
86 ///
87 /// Pass `topics = None` to describe **every** topic-partition the group
88 /// holds state for. Passing `Some(&[])` describes nothing — the wire
89 /// protocol distinguishes a null topics array from an empty one, and so
90 /// does this method.
91 ///
92 /// [`lag`](ShareGroupPartitionOffset::lag) is `Some(_)` only when the
93 /// coordinator supports `DescribeShareGroupOffsets` v1 (KIP-1226, Kafka
94 /// 4.3+); against an older broker it is `None` rather than a misleading
95 /// zero.
96 ///
97 /// The request goes to the group coordinator, which is where share-group
98 /// state lives.
99 ///
100 /// # Example
101 ///
102 /// ```ignore
103 /// // Everything the group knows about.
104 /// let all = admin.describe_share_group_offsets("orders-share", None).await?;
105 /// for p in &all.partitions {
106 /// println!("{}-{} start={} lag={:?}", p.topic, p.partition, p.start_offset, p.lag);
107 /// }
108 ///
109 /// // Just two partitions of one topic.
110 /// let some = admin
111 /// .describe_share_group_offsets("orders-share", Some(&[("orders", &[0, 1][..])]))
112 /// .await?;
113 /// ```
114 ///
115 /// # Errors
116 ///
117 /// Returns an error if the client is closed, a topic name is invalid, the
118 /// coordinator cannot be found, or the broker supports no compatible
119 /// `DescribeShareGroupOffsets` version.
120 pub async fn describe_share_group_offsets(
121 &self,
122 group_id: &str,
123 topics: Option<&[(&str, &[i32])]>,
124 ) -> Result<DescribeShareGroupOffsetsResult> {
125 self.check_not_closed()?;
126 if let Some(topics) = topics {
127 for (name, _) in topics {
128 validate_topic_name(name)?;
129 }
130 }
131
132 let coordinator = self.find_group_coordinator(group_id).await?;
133
134 let request = DescribeShareGroupOffsetsRequest {
135 groups: vec![DescribeShareGroupOffsetsRequestGroup {
136 group_id: group_id.to_string(),
137 topics: topics.map(|topics| {
138 topics
139 .iter()
140 .map(|(name, partitions)| DescribeShareGroupOffsetsRequestTopic {
141 topic_name: (*name).to_string(),
142 partitions: partitions.to_vec(),
143 })
144 .collect()
145 }),
146 }],
147 };
148
149 let version = coordinator
150 .negotiate_api_version(
151 ApiKey::DescribeShareGroupOffsets,
152 versions::DESCRIBE_SHARE_GROUP_OFFSETS_MAX,
153 versions::DESCRIBE_SHARE_GROUP_OFFSETS_MIN,
154 )
155 .ok_or_else(|| {
156 KrafkaError::protocol_kind(
157 ProtocolErrorKind::UnknownApiVersion,
158 "no mutually supported DescribeShareGroupOffsets API version; \
159 share groups require Kafka 4.2 or later",
160 )
161 })?;
162
163 let response_bytes = coordinator
164 .send_request(ApiKey::DescribeShareGroupOffsets, version, |buf| {
165 request.encode_versioned(version, buf)
166 })
167 .await?;
168
169 let mut buf = response_bytes;
170 let response = DescribeShareGroupOffsetsResponse::decode_versioned(version, &mut buf)?;
171
172 // One group was requested, so one group is expected back. A response
173 // with none is a broker-side contract violation, not an empty result.
174 let group = response.groups.into_iter().next().ok_or_else(|| {
175 KrafkaError::protocol_kind(
176 ProtocolErrorKind::Malformed,
177 format!(
178 "DescribeShareGroupOffsets returned no entry for group '{group_id}'; \
179 exactly one was requested"
180 ),
181 )
182 })?;
183
184 // `lag` is only meaningful from v1; below that the decoder reports the
185 // -1 sentinel, which must not be handed to callers as a real lag.
186 let lag_supported = version >= 1;
187
188 let partitions = group
189 .topics
190 .into_iter()
191 .flat_map(|topic| {
192 let topic_name = topic.topic_name;
193 topic
194 .partitions
195 .into_iter()
196 .map(move |p| ShareGroupPartitionOffset {
197 topic: topic_name.clone(),
198 partition: p.partition_index,
199 start_offset: p.start_offset,
200 leader_epoch: p.leader_epoch,
201 lag: if lag_supported && p.lag >= 0 {
202 Some(p.lag)
203 } else {
204 None
205 },
206 error: error_text(p.error_code, p.error_message),
207 })
208 })
209 .collect();
210
211 Ok(DescribeShareGroupOffsetsResult {
212 group_id: group.group_id,
213 error: error_text(group.error_code, group.error_message),
214 partitions,
215 })
216 }
217
218 /// Reset a share group's share-partition start offsets (KIP-932).
219 ///
220 /// **This is a destructive operation.** Moving the start offset backwards
221 /// re-delivers records the group already processed; moving it forwards
222 /// skips records permanently.
223 ///
224 /// The group must be **empty**. A group with a live member is answered
225 /// with `NON_EMPTY_GROUP`, for the same reason
226 /// [`alter_consumer_group_offsets`](Self::alter_consumer_group_offsets)
227 /// requires it: rewriting the start offset under an active member would
228 /// hand it records it has already acquired.
229 ///
230 /// # Example
231 ///
232 /// ```ignore
233 /// // Rewind two partitions to the beginning of the log.
234 /// let results = admin
235 /// .alter_share_group_offsets("orders-share", &[("orders", &[(0, 0), (1, 0)][..])])
236 /// .await?;
237 /// for r in &results {
238 /// if let Some(e) = &r.error {
239 /// eprintln!("{}-{} failed: {e}", r.topic, r.partition);
240 /// }
241 /// }
242 /// ```
243 ///
244 /// # Errors
245 ///
246 /// Returns an error if the client is closed, a topic name is invalid, the
247 /// coordinator cannot be found, the broker supports no compatible
248 /// `AlterShareGroupOffsets` version, or the request fails at the top level
249 /// (including `NON_EMPTY_GROUP`).
250 pub async fn alter_share_group_offsets(
251 &self,
252 group_id: &str,
253 topic_offsets: &[(&str, &[(i32, i64)])],
254 ) -> Result<Vec<ShareGroupOffsetAlteration>> {
255 self.check_not_closed()?;
256 for (name, _) in topic_offsets {
257 validate_topic_name(name)?;
258 }
259
260 let coordinator = self.find_group_coordinator(group_id).await?;
261
262 let request = AlterShareGroupOffsetsRequest {
263 group_id: group_id.to_string(),
264 topics: topic_offsets
265 .iter()
266 .map(|(name, partitions)| AlterShareGroupOffsetsRequestTopic {
267 topic_name: (*name).to_string(),
268 partitions: partitions
269 .iter()
270 .map(|&(partition_index, start_offset)| {
271 AlterShareGroupOffsetsRequestPartition {
272 partition_index,
273 start_offset,
274 }
275 })
276 .collect(),
277 })
278 .collect(),
279 };
280
281 let version = coordinator
282 .negotiate_api_version(
283 ApiKey::AlterShareGroupOffsets,
284 versions::ALTER_SHARE_GROUP_OFFSETS_MAX,
285 versions::ALTER_SHARE_GROUP_OFFSETS_MIN,
286 )
287 .ok_or_else(|| {
288 KrafkaError::protocol_kind(
289 ProtocolErrorKind::UnknownApiVersion,
290 "no mutually supported AlterShareGroupOffsets API version; \
291 share groups require Kafka 4.2 or later",
292 )
293 })?;
294
295 let response_bytes = coordinator
296 .send_request(ApiKey::AlterShareGroupOffsets, version, |buf| {
297 request.encode_versioned(version, buf)
298 })
299 .await?;
300
301 let mut buf = response_bytes;
302 let response = AlterShareGroupOffsetsResponse::decode_versioned(version, &mut buf)?;
303
304 // A top-level failure means nothing was applied. Surfacing it as an
305 // error — with the broker's own code, so `is_retriable()` governs the
306 // retry — beats returning an empty success list that reads like "no
307 // partitions were requested".
308 if !response.error_code.is_ok() {
309 let msg = response
310 .error_message
311 .unwrap_or_else(|| format!("{:?}", response.error_code));
312 return Err(KrafkaError::broker(response.error_code, msg));
313 }
314
315 let results: Vec<_> = response
316 .responses
317 .into_iter()
318 .flat_map(|topic| {
319 let topic_name = topic.topic_name;
320 topic
321 .partitions
322 .into_iter()
323 .map(move |p| ShareGroupOffsetAlteration {
324 topic: topic_name.clone(),
325 partition: p.partition_index,
326 error: error_text(p.error_code, p.error_message),
327 })
328 })
329 .collect();
330
331 let failed = results.iter().filter(|r| r.error.is_some()).count();
332 if failed > 0 {
333 warn!(
334 group = group_id,
335 failed,
336 total = results.len(),
337 "AlterShareGroupOffsets completed with partition-level failures"
338 );
339 } else {
340 info!(
341 group = group_id,
342 partitions = results.len(),
343 "AlterShareGroupOffsets completed"
344 );
345 }
346 Ok(results)
347 }
348
349 /// Delete a share group's offset state for whole topics (KIP-932).
350 ///
351 /// **This is a destructive operation** — the deleted state cannot be
352 /// recovered, and the group restarts those topics from its configured
353 /// reset policy. The group must be **empty**.
354 ///
355 /// Use it after retiring a topic, so the coordinator stops carrying state
356 /// for partitions that no longer exist.
357 ///
358 /// # Errors
359 ///
360 /// Returns an error if the client is closed, a topic name is invalid, the
361 /// coordinator cannot be found, the broker supports no compatible
362 /// `DeleteShareGroupOffsets` version, or the request fails at the top level
363 /// (including `NON_EMPTY_GROUP`).
364 pub async fn delete_share_group_offsets(
365 &self,
366 group_id: &str,
367 topics: &[&str],
368 ) -> Result<Vec<ShareGroupOffsetDeletion>> {
369 self.check_not_closed()?;
370 for name in topics {
371 validate_topic_name(name)?;
372 }
373
374 let coordinator = self.find_group_coordinator(group_id).await?;
375
376 let request = DeleteShareGroupOffsetsRequest {
377 group_id: group_id.to_string(),
378 topics: topics.iter().map(|t| (*t).to_string()).collect(),
379 };
380
381 let version = coordinator
382 .negotiate_api_version(
383 ApiKey::DeleteShareGroupOffsets,
384 versions::DELETE_SHARE_GROUP_OFFSETS_MAX,
385 versions::DELETE_SHARE_GROUP_OFFSETS_MIN,
386 )
387 .ok_or_else(|| {
388 KrafkaError::protocol_kind(
389 ProtocolErrorKind::UnknownApiVersion,
390 "no mutually supported DeleteShareGroupOffsets API version; \
391 share groups require Kafka 4.2 or later",
392 )
393 })?;
394
395 let response_bytes = coordinator
396 .send_request(ApiKey::DeleteShareGroupOffsets, version, |buf| {
397 request.encode_versioned(version, buf)
398 })
399 .await?;
400
401 let mut buf = response_bytes;
402 let response = DeleteShareGroupOffsetsResponse::decode_versioned(version, &mut buf)?;
403
404 if !response.error_code.is_ok() {
405 let msg = response
406 .error_message
407 .unwrap_or_else(|| format!("{:?}", response.error_code));
408 return Err(KrafkaError::broker(response.error_code, msg));
409 }
410
411 Ok(response
412 .responses
413 .into_iter()
414 .map(|t| ShareGroupOffsetDeletion {
415 topic: t.topic_name,
416 error: error_text(t.error_code, t.error_message),
417 })
418 .collect())
419 }
420}
421
422/// Render a `(code, message)` pair as `Some(text)` on failure, `None` on
423/// success — the shape every other admin result in this crate uses.
424///
425/// Falls back to the debug form of the code when the broker sent no message,
426/// so a failure is never reported as an empty string.
427fn error_text(code: crate::error::ErrorCode, message: Option<String>) -> Option<String> {
428 if code.is_ok() {
429 None
430 } else {
431 Some(message.unwrap_or_else(|| format!("{code:?}")))
432 }
433}
434
435#[cfg(test)]
436#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
437mod tests {
438 use super::*;
439 use crate::error::ErrorCode;
440
441 #[test]
442 fn error_text_is_none_on_success() {
443 assert_eq!(error_text(ErrorCode::None, None), None);
444 assert_eq!(error_text(ErrorCode::None, Some("ignored".into())), None);
445 }
446
447 /// A failure without a broker message must still say *something*; an empty
448 /// `Some("")` would render as a silent failure in operator tooling.
449 #[test]
450 fn error_text_falls_back_to_the_code() {
451 let text = error_text(ErrorCode::UnknownServerError, None).expect("failure is Some");
452 assert!(!text.is_empty());
453 assert!(text.contains("UnknownServerError"), "got: {text}");
454
455 assert_eq!(
456 error_text(ErrorCode::UnknownServerError, Some("boom".into())),
457 Some("boom".to_string())
458 );
459 }
460}