Skip to main content

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}