Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -89,6 +89,8 @@ async fn check_store_health(
) -> Option<String> {
let fut = client.fragment_store_health(broker::FragmentStoreHealthRequest {
fragment_store: fragment_store.to_string(),
check_delete: true,
check_delete_prefix: "recovery/".to_string(),
});

match tokio::time::timeout(HEALTH_CHECK_TIMEOUT, fut).await {
Expand Down
8 changes: 8 additions & 0 deletions crates/proto-gazette/src/protocol.rs
Original file line number Diff line number Diff line change
Expand Up @@ -875,6 +875,14 @@ pub struct FragmentStoreHealthRequest {
/// Fragment store to check the health of.
#[prost(string, tag = "1")]
pub fragment_store: ::prost::alloc::string::String,
/// If true, also verify delete permission with a throwaway probe object under
/// check_delete_prefix. Defaults to false (delete is not exercised).
#[prost(bool, tag = "2")]
pub check_delete: bool,
/// Prefix under which the delete probe is written (e.g. "recovery/"); empty
/// probes the store root. Only used when check_delete is true.
#[prost(string, tag = "3")]
pub check_delete_prefix: ::prost::alloc::string::String,
}
/// FragmentStoreHealthResponse is the unary response message of the broker FragmentStoreHealth RPC.
#[derive(Clone, PartialEq, ::prost::Message)]
Expand Down
36 changes: 36 additions & 0 deletions crates/proto-gazette/src/protocol.serde.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1357,10 +1357,22 @@ impl serde::Serialize for FragmentStoreHealthRequest {
if !self.fragment_store.is_empty() {
len += 1;
}
if self.check_delete {
len += 1;
}
if !self.check_delete_prefix.is_empty() {
len += 1;
}
let mut struct_ser = serializer.serialize_struct("protocol.FragmentStoreHealthRequest", len)?;
if !self.fragment_store.is_empty() {
struct_ser.serialize_field("fragmentStore", &self.fragment_store)?;
}
if self.check_delete {
struct_ser.serialize_field("checkDelete", &self.check_delete)?;
}
if !self.check_delete_prefix.is_empty() {
struct_ser.serialize_field("checkDeletePrefix", &self.check_delete_prefix)?;
}
struct_ser.end()
}
}
Expand All @@ -1373,11 +1385,17 @@ impl<'de> serde::Deserialize<'de> for FragmentStoreHealthRequest {
const FIELDS: &[&str] = &[
"fragment_store",
"fragmentStore",
"check_delete",
"checkDelete",
"check_delete_prefix",
"checkDeletePrefix",
];

#[allow(clippy::enum_variant_names)]
enum GeneratedField {
FragmentStore,
CheckDelete,
CheckDeletePrefix,
__SkipField__,
}
impl<'de> serde::Deserialize<'de> for GeneratedField {
Expand All @@ -1401,6 +1419,8 @@ impl<'de> serde::Deserialize<'de> for FragmentStoreHealthRequest {
{
match value {
"fragmentStore" | "fragment_store" => Ok(GeneratedField::FragmentStore),
"checkDelete" | "check_delete" => Ok(GeneratedField::CheckDelete),
"checkDeletePrefix" | "check_delete_prefix" => Ok(GeneratedField::CheckDeletePrefix),
_ => Ok(GeneratedField::__SkipField__),
}
}
Expand All @@ -1421,6 +1441,8 @@ impl<'de> serde::Deserialize<'de> for FragmentStoreHealthRequest {
V: serde::de::MapAccess<'de>,
{
let mut fragment_store__ = None;
let mut check_delete__ = None;
let mut check_delete_prefix__ = None;
while let Some(k) = map_.next_key()? {
match k {
GeneratedField::FragmentStore => {
Expand All @@ -1429,13 +1451,27 @@ impl<'de> serde::Deserialize<'de> for FragmentStoreHealthRequest {
}
fragment_store__ = Some(map_.next_value()?);
}
GeneratedField::CheckDelete => {
if check_delete__.is_some() {
return Err(serde::de::Error::duplicate_field("checkDelete"));
}
check_delete__ = Some(map_.next_value()?);
}
GeneratedField::CheckDeletePrefix => {
if check_delete_prefix__.is_some() {
return Err(serde::de::Error::duplicate_field("checkDeletePrefix"));
}
check_delete_prefix__ = Some(map_.next_value()?);
}
GeneratedField::__SkipField__ => {
let _ = map_.next_value::<serde::de::IgnoredAny>()?;
}
}
}
Ok(FragmentStoreHealthRequest {
fragment_store: fragment_store__.unwrap_or_default(),
check_delete: check_delete__.unwrap_or_default(),
check_delete_prefix: check_delete_prefix__.unwrap_or_default(),
})
}
}
Expand Down
2 changes: 1 addition & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ require (
github.com/stretchr/testify v1.11.1
go.etcd.io/etcd/api/v3 v3.6.5
go.etcd.io/etcd/client/v3 v3.6.5
go.gazette.dev/core v0.103.1-0.20260715211535-42c2cc160b2c
go.gazette.dev/core v0.103.1-0.20260722193110-e54beb5c6e64
golang.org/x/net v0.55.0
google.golang.org/api v0.251.0
google.golang.org/grpc v1.79.3
Expand Down
4 changes: 2 additions & 2 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -278,8 +278,8 @@ go.etcd.io/etcd/client/pkg/v3 v3.6.5 h1:Duz9fAzIZFhYWgRjp/FgNq2gO1jId9Yae/rLn3Rr
go.etcd.io/etcd/client/pkg/v3 v3.6.5/go.mod h1:8Wx3eGRPiy0qOFMZT/hfvdos+DjEaPxdIDiCDUv/FQk=
go.etcd.io/etcd/client/v3 v3.6.5 h1:yRwZNFBx/35VKHTcLDeO7XVLbCBFbPi+XV4OC3QJf2U=
go.etcd.io/etcd/client/v3 v3.6.5/go.mod h1:ZqwG/7TAFZ0BJ0jXRPoJjKQJtbFo/9NIY8uoFFKcCyo=
go.gazette.dev/core v0.103.1-0.20260715211535-42c2cc160b2c h1:6oRaNmlNhFD4IaAkXSfOgXUxayPskHBTndhuNG7VSyw=
go.gazette.dev/core v0.103.1-0.20260715211535-42c2cc160b2c/go.mod h1:fbzVmqkkTxqOnpcvdXuSWJ28XoMMm08ovNaATv1Zuas=
go.gazette.dev/core v0.103.1-0.20260722193110-e54beb5c6e64 h1:YX2MhHISLxvz1xoS+YQntf4FMZ24Q02+97nmTHImrYI=
go.gazette.dev/core v0.103.1-0.20260722193110-e54beb5c6e64/go.mod h1:fbzVmqkkTxqOnpcvdXuSWJ28XoMMm08ovNaATv1Zuas=
go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64=
go.opentelemetry.io/auto/sdk v1.2.1/go.mod h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y=
go.opentelemetry.io/contrib/detectors/gcp v1.39.0 h1:kWRNZMsfBHZ+uHjiH4y7Etn2FK26LAGkNFw7RHv1DhE=
Expand Down
Loading