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
11 changes: 8 additions & 3 deletions broker/client/fragment_store_health.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,9 +10,14 @@ import (
// FragmentStoreHealth queries for the latest health status on the specified fragment store.
// It returns the health check response or an error if the RPC fails with a status other
// than FRAGMENT_STORE_UNHEALTHY.
func FragmentStoreHealth(ctx context.Context, client pb.JournalClient, store pb.FragmentStore) (*pb.FragmentStoreHealthResponse, error) {
var req = pb.FragmentStoreHealthRequest{
FragmentStore: store,
//
// A non-nil checkDeletePrefix also probes delete permission under that prefix;
// a pointer to the empty string probes the store root.
func FragmentStoreHealth(ctx context.Context, client pb.JournalClient, store pb.FragmentStore, checkDeletePrefix *string) (*pb.FragmentStoreHealthResponse, error) {
var req = pb.FragmentStoreHealthRequest{FragmentStore: store}
if checkDeletePrefix != nil {
req.CheckDelete = true
req.CheckDeletePrefix = *checkDeletePrefix
}

if resp, err := client.FragmentStoreHealth(pb.WithDispatchDefault(ctx), &req); err != nil {
Expand Down
36 changes: 32 additions & 4 deletions broker/client/fragment_store_health_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ func TestCheckFragmentStoreHealth(t *testing.T) {
}, nil
}

resp, err := FragmentStoreHealth(ctx, broker.Client(), "s3://my-bucket/")
resp, err := FragmentStoreHealth(ctx, broker.Client(), "s3://my-bucket/", nil)
require.NoError(t, err)
require.Equal(t, pb.Status_OK, resp.Status)
require.Empty(t, resp.StoreHealthError)
Expand All @@ -38,7 +38,7 @@ func TestCheckFragmentStoreHealth(t *testing.T) {
}, nil
}

resp, err = FragmentStoreHealth(ctx, broker.Client(), "s3://my-bucket/")
resp, err = FragmentStoreHealth(ctx, broker.Client(), "s3://my-bucket/", nil)
require.NoError(t, err)
require.Equal(t, pb.Status_FRAGMENT_STORE_UNHEALTHY, resp.Status)
require.Equal(t, "store is unhealthy: connection timeout", resp.StoreHealthError)
Expand All @@ -48,7 +48,7 @@ func TestCheckFragmentStoreHealth(t *testing.T) {
return nil, errors.New("network error")
}

resp, err = FragmentStoreHealth(ctx, broker.Client(), "s3://my-bucket/")
resp, err = FragmentStoreHealth(ctx, broker.Client(), "s3://my-bucket/", nil)
require.Error(t, err)
require.Contains(t, err.Error(), "network error")
require.Nil(t, resp)
Expand All @@ -61,8 +61,36 @@ func TestCheckFragmentStoreHealth(t *testing.T) {
}, nil
}

resp, err = FragmentStoreHealth(ctx, broker.Client(), "s3://my-bucket/")
resp, err = FragmentStoreHealth(ctx, broker.Client(), "s3://my-bucket/", nil)
require.Error(t, err)
require.Equal(t, "JOURNAL_NOT_FOUND", err.Error())
require.Nil(t, resp)

// Case 5: nil leaves check_delete unset; a non-nil prefix sets both fields.
var gotReq *pb.FragmentStoreHealthRequest
broker.FragmentStoreHealthFunc = func(ctx context.Context, req *pb.FragmentStoreHealthRequest) (*pb.FragmentStoreHealthResponse, error) {
gotReq = req
return &pb.FragmentStoreHealthResponse{
Status: pb.Status_OK,
Header: *buildHeaderFixture(broker),
}, nil
}

_, err = FragmentStoreHealth(ctx, broker.Client(), "s3://my-bucket/", nil)
require.NoError(t, err)
require.False(t, gotReq.CheckDelete)
require.Empty(t, gotReq.CheckDeletePrefix)

var prefix = "recovery/"
_, err = FragmentStoreHealth(ctx, broker.Client(), "s3://my-bucket/", &prefix)
require.NoError(t, err)
require.True(t, gotReq.CheckDelete)
require.Equal(t, "recovery/", gotReq.CheckDeletePrefix)

// Case 6: a pointer to "" probes the store root: check_delete set, prefix empty.
var root = ""
_, err = FragmentStoreHealth(ctx, broker.Client(), "s3://my-bucket/", &root)
require.NoError(t, err)
require.True(t, gotReq.CheckDelete)
require.Empty(t, gotReq.CheckDeletePrefix)
}
9 changes: 9 additions & 0 deletions broker/fragment_store_health_api.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package broker
import (
"context"
"net"
"time"

log "github.com/sirupsen/logrus"
pb "go.gazette.dev/core/broker/protocol"
Expand Down Expand Up @@ -55,6 +56,14 @@ func (svc *Service) FragmentStoreHealth(ctx context.Context, claims pb.Claims, r
if healthErr != nil {
resp.Status = pb.Status_FRAGMENT_STORE_UNHEALTHY
resp.StoreHealthError = healthErr.Error()
} else if req.CheckDelete {
var probeCtx, cancel = context.WithTimeout(ctx, time.Minute)
defer cancel()

if deleteErr := activeStore.CheckDelete(probeCtx, req.CheckDeletePrefix); deleteErr != nil {
resp.Status = pb.Status_FRAGMENT_STORE_UNHEALTHY
resp.StoreHealthError = deleteErr.Error()
}
}
return resp, nil
}
57 changes: 57 additions & 0 deletions broker/fragment_store_health_api_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -55,3 +55,60 @@ func TestFragmentStoreHealthCases(t *testing.T) {
require.Error(t, err)
require.Contains(t, err.Error(), "context canceled")
}

func TestFragmentStoreHealthDeleteProbe(t *testing.T) {
var ctx, etcd = pb.WithDispatchDefault(context.Background()), etcdtest.TestClient()
defer etcdtest.Cleanup()

// The "no-delete" store denies delete but delegates all else to memory, so
// the periodic health check passes and the RPC reaches the delete probe.
stores.RegisterProviders(map[string]stores.Constructor{
"s3": func(ep *url.URL) (stores.Store, error) {
var mem = stores.NewMemoryStore(ep)
if ep.Host == "no-delete" {
return &stores.CallbackStore{
Fallback: mem,
RemoveFunc: func(_ stores.Store, ctx context.Context, _ string) error {
if _, ok := ctx.Deadline(); !ok {
return fmt.Errorf("delete probe ran without a deadline")
}
return fmt.Errorf("simulated missing delete permission")
},
}, nil
}
return mem, nil
},
})

var broker = newTestBroker(t, etcd, pb.ProcessSpec_ID{Zone: "local", Suffix: "broker"})
defer broker.cleanup()

// Healthy store grants delete: passes under a sub-prefix.
var resp, err = broker.client().FragmentStoreHealth(ctx, &pb.FragmentStoreHealthRequest{
FragmentStore: "s3://bucket/",
CheckDelete: true,
CheckDeletePrefix: "recovery/",
})
require.NoError(t, err)
require.Equal(t, pb.Status_OK, resp.Status)
require.Empty(t, resp.StoreHealthError)

// Store lacks delete permission: probe fails.
resp, err = broker.client().FragmentStoreHealth(ctx, &pb.FragmentStoreHealthRequest{
FragmentStore: "s3://no-delete/",
CheckDelete: true,
CheckDeletePrefix: "recovery/",
})
require.NoError(t, err)
require.Equal(t, pb.Status_FRAGMENT_STORE_UNHEALTHY, resp.Status)
require.Contains(t, resp.StoreHealthError, "delete-probe DELETE failed")
require.Contains(t, resp.StoreHealthError, "simulated missing delete permission")

// check_delete false: reports healthy (background-check behavior).
resp, err = broker.client().FragmentStoreHealth(ctx, &pb.FragmentStoreHealthRequest{
FragmentStore: "s3://no-delete/",
})
require.NoError(t, err)
require.Equal(t, pb.Status_OK, resp.Status)
require.Empty(t, resp.StoreHealthError)
}
Loading
Loading