mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-11 16:57:45 +02:00
fix(filer): end metadata subscriptions before gRPC GracefulStop on shutdown (#11663)
* fix(filer): end metadata subscriptions before gRPC GracefulStop on shutdown GracefulStop waits for every open stream. The filer's own MetaAggregator subscription and, under -s3, the S3 gateway's IAM subscription live in the same process and only end when the filer shuts down, which happens after GracefulStop - so every shutdown sat out the full 15s timeout. StopSubscriptions ends them (and any started later) right after leaving the lock ring; the handlers derive their context from the stream and this signal. Measured on weed server -s3: stop 25.6s -> 10.7s, and 0.6s with -volume.preStopSeconds=0. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FTq6bDpgdQfqQagQsUUvdw * fix(filer): end stopped subscriptions with Unavailable, not cleanly A subscription that StopSubscriptions cut short returned nil. The client reads that as io.EOF, "caught up, done", and util.RetryUntil - which the S3 gateway and mount follow with - stops on nil. A separate S3 gateway therefore never resubscribed after its filer restarted (measured: 0 subscriptions in 25s after the restart; upstream master and this fix: it is back within a second). endOfSubscription now turns that clean end into codes.Unavailable, which is what a dropped connection looks like, so followers reconnect. Errors pass through, and so does the end of a stream the client closed. A subscriber that has stopped reading is still left to GracefulStop's timeout: its handler sits in stream.Send, which only the end of the stream releases, and grpc-go does not allow Send after the handler returns. Documented at StopSubscriptions. Shutdown under weed server -s3 is unchanged at 10.7s. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FTq6bDpgdQfqQagQsUUvdw * fix(filer): interrupt in-flight subscription reads and end them Unavailable StopSubscriptions only reached the wait points: a handler replaying the ring or the persisted log kept streaming until the pass ended, and a pass cut short at a ctx.Err() boundary (the replay semaphore, chunk ref collection, ref sends) returned a wrapped context.Canceled, which reads as codes.Canceled on the wire - a code transient-error classifiers treat as non-retryable. eachLogEntryFn and sendRefsBatched now check the subscription context per entry/batch, and endOfSubscription reports a context.Canceled cut short by the stop signal as codes.Unavailable, same as a clean end it interrupted. * filer: honor the subscription context in the pipelined sender A subscriber that stops reading wedges sendLoop in stream.Send; Send and Close then blocked past StopSubscriptions, so GracefulStop still waited out its timeout. Both now return when the subscription context ends, letting the handler return and gRPC tear the stream down. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * filer: stop sendLoop from starting sends after the subscription ends With Send and Close released by the subscription context, sendLoop could still issue a stream.Send after the handler returned. Gate every send on the context so at most the in-flight Send overlaps teardown - and the stream's own teardown is what unblocks it. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * filer: assert sendLoop exits in the cancel test Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --------- Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com> Co-authored-by: Chris Lu <chris.lu@gmail.com> Co-authored-by: Chris Lu <chrislusf@users.noreply.github.com> Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
9 files changed
+409
-38
No files matched your search
@@ -574,7 +574,7 @@ func (fo *FilerOptions) startFiler() {
|
||||
}
|
||||
httpS := newHttpServer(defaultHandler, tlsConfig)
|
||||
httpServers = append(httpServers, httpS)
|
||||
shutdown := newFilerShutdown(fs.LeaveLockRing, stopGrpcServer, fs.Shutdown, httpServers...)
|
||||
shutdown := newFilerShutdown(fs.LeaveLockRing, fs.StopSubscriptions, stopGrpcServer, fs.Shutdown, httpServers...)
|
||||
|
||||
grace.OnInterrupt(shutdown)
|
||||
|
||||
@@ -603,7 +603,7 @@ func (fo *FilerOptions) startFiler() {
|
||||
}
|
||||
httpS := newHttpServer(defaultHandler, nil)
|
||||
httpServers = append(httpServers, httpS)
|
||||
shutdown := newFilerShutdown(fs.LeaveLockRing, stopGrpcServer, fs.Shutdown, httpServers...)
|
||||
shutdown := newFilerShutdown(fs.LeaveLockRing, fs.StopSubscriptions, stopGrpcServer, fs.Shutdown, httpServers...)
|
||||
|
||||
grace.OnInterrupt(shutdown)
|
||||
|
||||
@@ -625,10 +625,13 @@ func (fo *FilerOptions) startFiler() {
|
||||
// newFilerShutdown joins shutdown callers while gRPC and HTTP drain concurrently.
|
||||
// The filer leaves the lock ring first: peers and S3 gateways route keys to it
|
||||
// until the ring changes, and would hit refused connections once gRPC stops.
|
||||
func newFilerShutdown(leaveLockRing, stopGrpc, shutdownFiler func(), httpServers ...*http.Server) func() {
|
||||
// Then it ends its metadata subscriptions, so gRPC GracefulStop does not wait
|
||||
// out its timeout on streams whose subscribers sit in this same process.
|
||||
func newFilerShutdown(leaveLockRing, stopSubscriptions, stopGrpc, shutdownFiler func(), httpServers ...*http.Server) func() {
|
||||
drain := newGracefulShutdown(stopGrpc, shutdownFiler, httpServers...)
|
||||
return sync.OnceFunc(func() {
|
||||
leaveLockRing()
|
||||
stopSubscriptions()
|
||||
drain()
|
||||
})
|
||||
}
|
||||
@@ -56,7 +56,7 @@ func TestFilerShutdownJoinsServeDuringParallelDrain(t *testing.T) {
|
||||
grpcStopped := make(chan struct{})
|
||||
filerClosed := make(chan struct{})
|
||||
var closes atomic.Int32
|
||||
shutdown := newFilerShutdown(func() {}, func() {
|
||||
shutdown := newFilerShutdown(func() {}, func() {}, func() {
|
||||
close(grpcStarted)
|
||||
<-grpcRelease
|
||||
close(grpcStopped)
|
||||
@@ -138,7 +138,7 @@ func TestFilerShutdownWaitsForHTTPAfterGrpcStops(t *testing.T) {
|
||||
server.Config.RegisterOnShutdown(func() { close(httpClosing) })
|
||||
grpcStopped := make(chan struct{})
|
||||
filerClosed := make(chan struct{})
|
||||
shutdown := newFilerShutdown(func() {}, func() { close(grpcStopped) }, func() { close(filerClosed) }, server.Config)
|
||||
shutdown := newFilerShutdown(func() {}, func() {}, func() { close(grpcStopped) }, func() { close(filerClosed) }, server.Config)
|
||||
joined := make(chan struct{})
|
||||
go func() { shutdown(); close(joined) }()
|
||||
<-httpClosing
|
||||
@@ -172,7 +172,7 @@ func TestFilerShutdownLeavesLockRingBeforeDraining(t *testing.T) {
|
||||
shutdown := newFilerShutdown(func() {
|
||||
close(leaving)
|
||||
<-releaseLeave
|
||||
}, func() { close(grpcStopping) }, func() {}, server.Config)
|
||||
}, func() {}, func() { close(grpcStopping) }, func() {}, server.Config)
|
||||
joined := make(chan struct{})
|
||||
go func() { shutdown(); close(joined) }()
|
||||
|
||||
@@ -196,3 +196,31 @@ func TestFilerShutdownLeavesLockRingBeforeDraining(t *testing.T) {
|
||||
t.Error("gRPC was not stopped after leaving the lock ring")
|
||||
}
|
||||
}
|
||||
|
||||
// The filer's own MetaAggregator and an in-process S3 gateway keep metadata
|
||||
// subscription streams open until the filer is shut down, which happens after
|
||||
// gRPC GracefulStop - so the subscriptions have to end before it, or every
|
||||
// shutdown sits out the full graceful-stop timeout.
|
||||
func TestFilerShutdownStopsSubscriptionsBeforeGrpc(t *testing.T) {
|
||||
var order []string
|
||||
var mu sync.Mutex
|
||||
record := func(step string) func() {
|
||||
return func() {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
order = append(order, step)
|
||||
}
|
||||
}
|
||||
shutdown := newFilerShutdown(record("leave"), record("subscriptions"), record("grpc"), record("filer"))
|
||||
shutdown()
|
||||
|
||||
want := []string{"leave", "subscriptions", "grpc", "filer"}
|
||||
if len(order) != len(want) {
|
||||
t.Fatalf("shutdown steps %v, want %v", order, want)
|
||||
}
|
||||
for i := range want {
|
||||
if order[i] != want[i] {
|
||||
t.Fatalf("shutdown steps %v, want %v", order, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -11,6 +11,8 @@ import (
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/stats"
|
||||
|
||||
"google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/status"
|
||||
"google.golang.org/protobuf/proto"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/filer"
|
||||
@@ -109,14 +111,16 @@ const (
|
||||
// current time (backlog catch-up), multiple events are packed into a single
|
||||
// stream.Send() using the Events field. Otherwise events are sent one-by-one.
|
||||
type pipelinedSender struct {
|
||||
ctx context.Context
|
||||
sendCh chan *filer_pb.SubscribeMetadataResponse
|
||||
errCh chan error
|
||||
done chan struct{}
|
||||
canBatch bool // true only if client set ClientSupportsBatching
|
||||
}
|
||||
|
||||
func newPipelinedSender(stream metadataStreamSender, bufSize int, clientSupportsBatching bool) *pipelinedSender {
|
||||
func newPipelinedSender(ctx context.Context, stream metadataStreamSender, bufSize int, clientSupportsBatching bool) *pipelinedSender {
|
||||
s := &pipelinedSender{
|
||||
ctx: ctx,
|
||||
sendCh: make(chan *filer_pb.SubscribeMetadataResponse, bufSize),
|
||||
errCh: make(chan error, 1),
|
||||
done: make(chan struct{}),
|
||||
@@ -128,6 +132,21 @@ func newPipelinedSender(stream metadataStreamSender, bufSize int, clientSupports
|
||||
|
||||
func (s *pipelinedSender) sendLoop(stream metadataStreamSender) {
|
||||
defer close(s.done)
|
||||
// No Send may start after the subscription context ends: the handler is
|
||||
// returning, and gRPC tears the stream down behind it. A Send already in
|
||||
// flight is released by that teardown (the stream context unblocks it).
|
||||
send := func(msg *filer_pb.SubscribeMetadataResponse) bool {
|
||||
if s.ctx.Err() != nil {
|
||||
return false
|
||||
}
|
||||
if err := stream.Send(msg); err != nil {
|
||||
if s.ctx.Err() == nil {
|
||||
s.reportErr(err)
|
||||
}
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
for msg := range s.sendCh {
|
||||
// LogFileRefs messages are unbatchable: the client recognizes them by
|
||||
// the top-level field and skips the rest of the response, so a refs
|
||||
@@ -142,8 +161,7 @@ func (s *pipelinedSender) sendLoop(stream metadataStreamSender) {
|
||||
|
||||
if !shouldBatch {
|
||||
// Real-time: send immediately for low latency
|
||||
if err := stream.Send(msg); err != nil {
|
||||
s.reportErr(err)
|
||||
if !send(msg) {
|
||||
return
|
||||
}
|
||||
continue
|
||||
@@ -181,18 +199,14 @@ func (s *pipelinedSender) sendLoop(stream metadataStreamSender) {
|
||||
toSend = batch[0]
|
||||
toSend.Events = batch[1:]
|
||||
}
|
||||
if err := stream.Send(toSend); err != nil {
|
||||
s.reportErr(err)
|
||||
if !send(toSend) {
|
||||
return
|
||||
}
|
||||
if toSend.Events != nil {
|
||||
toSend.Events = nil
|
||||
}
|
||||
if trailingSolo != nil {
|
||||
if err := stream.Send(trailingSolo); err != nil {
|
||||
s.reportErr(err)
|
||||
return
|
||||
}
|
||||
if trailingSolo != nil && !send(trailingSolo) {
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -207,11 +221,16 @@ func (s *pipelinedSender) reportErr(err error) {
|
||||
}
|
||||
|
||||
func (s *pipelinedSender) Send(msg *filer_pb.SubscribeMetadataResponse) error {
|
||||
if err := s.ctx.Err(); err != nil {
|
||||
return err
|
||||
}
|
||||
select {
|
||||
case s.sendCh <- msg:
|
||||
return nil
|
||||
case err := <-s.errCh:
|
||||
return err
|
||||
case <-s.ctx.Done():
|
||||
return s.ctx.Err()
|
||||
case <-s.done:
|
||||
// Sender goroutine exited (stream error or shutdown).
|
||||
select {
|
||||
@@ -225,7 +244,12 @@ func (s *pipelinedSender) Send(msg *filer_pb.SubscribeMetadataResponse) error {
|
||||
|
||||
func (s *pipelinedSender) Close() error {
|
||||
close(s.sendCh)
|
||||
<-s.done
|
||||
// A sendLoop stuck in stream.Send only unblocks once the handler's return
|
||||
// ends the stream, so stop waiting when the subscription ends.
|
||||
select {
|
||||
case <-s.done:
|
||||
case <-s.ctx.Done():
|
||||
}
|
||||
select {
|
||||
case err := <-s.errCh:
|
||||
return err
|
||||
@@ -602,7 +626,7 @@ func (p *gapPass) park(ctx context.Context, cursor *log_buffer.MessagePosition,
|
||||
return gapContinue
|
||||
}
|
||||
|
||||
func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest, stream filer_pb.SeaweedFiler_SubscribeMetadataServer) error {
|
||||
func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest, stream filer_pb.SeaweedFiler_SubscribeMetadataServer) (err error) {
|
||||
// A filer that has not learned remote peers yet serves the local log and
|
||||
// upgrades when the first one appears. RemotePeerArrivedChan takes the
|
||||
// arrival channel under the same lock as the peer check, so a peer
|
||||
@@ -616,7 +640,9 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
|
||||
return fs.subscribeLocalMetadata(req, stream, nil)
|
||||
}
|
||||
|
||||
ctx := stream.Context()
|
||||
ctx, cancelSubscription := fs.subscriptionContext(stream.Context())
|
||||
defer cancelSubscription()
|
||||
defer func() { err = fs.endOfSubscription(stream.Context(), err) }()
|
||||
peerAddress := findClientAddress(ctx, 0)
|
||||
|
||||
isReplacing, alreadyKnown, clientName := fs.addClient("", req.ClientName, peerAddress, req.PathPrefix, req.ClientId, req.ClientEpoch)
|
||||
@@ -643,7 +669,7 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
|
||||
// had already arrived and been delivered).
|
||||
diskAnchorTsNs := req.SinceNs
|
||||
|
||||
sender := newPipelinedSender(stream, 1024, req.ClientSupportsBatching)
|
||||
sender := newPipelinedSender(ctx, stream, 1024, req.ClientSupportsBatching)
|
||||
defer sender.Close()
|
||||
|
||||
// Register for instant notification when new data arrives in the aggregated log buffer.
|
||||
@@ -669,7 +695,7 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
|
||||
// written from this single goroutine, so no synchronization is needed.
|
||||
var lastSeenTsNs int64
|
||||
var lastHeartbeatNs int64
|
||||
baseEachLogEntryFn := eachLogEntryFn(req, sender, eachEventNotificationFn, &unsyncedEvents)
|
||||
baseEachLogEntryFn := eachLogEntryFn(ctx, req, sender, eachEventNotificationFn, &unsyncedEvents)
|
||||
// heldAtTsNs remembers the entry a read was held at (for the log line);
|
||||
// diskHeldAtTsNs is the same marker for the disk pass alone: a pending
|
||||
// disk hold keeps the pass re-reading until the entry is served.
|
||||
@@ -1004,6 +1030,59 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
|
||||
|
||||
}
|
||||
|
||||
// StopSubscriptions ends every metadata subscription and any started later.
|
||||
// Shutdown calls it before gRPC GracefulStop, which waits on open streams;
|
||||
// subscribers living in this process - the filer's own MetaAggregator and an
|
||||
// in-process S3 gateway - are otherwise torn down only by the later
|
||||
// fs.Shutdown, so every stop waited out the graceful-stop timeout. Ended
|
||||
// subscribers get codes.Unavailable (see endOfSubscription) and reconnect
|
||||
// with their usual retry.
|
||||
//
|
||||
// A subscriber blocked in stream.Send is released once its handler returns:
|
||||
// the sender's Send and Close honor the subscription context, so the handler
|
||||
// exits and gRPC tears the stream down.
|
||||
func (fs *FilerServer) StopSubscriptions() {
|
||||
if fs.stopSubscriptions != nil {
|
||||
fs.stopSubscriptions()
|
||||
}
|
||||
}
|
||||
|
||||
// subscriptionContext ends with the stream or with StopSubscriptions,
|
||||
// whichever comes first. It keeps the stream's values (peer address).
|
||||
func (fs *FilerServer) subscriptionContext(stream context.Context) (context.Context, context.CancelFunc) {
|
||||
ctx, cancel := context.WithCancel(stream)
|
||||
if fs.subscriptionsStopped == nil {
|
||||
return ctx, cancel
|
||||
}
|
||||
stop := context.AfterFunc(fs.subscriptionsStopped, cancel)
|
||||
return ctx, func() {
|
||||
stop()
|
||||
cancel()
|
||||
}
|
||||
}
|
||||
|
||||
// endOfSubscription turns the end of a subscription that StopSubscriptions
|
||||
// cut short into codes.Unavailable. Followers read a clean end as "caught up,
|
||||
// done": the client returns nil on io.EOF, and util.RetryUntil - which the S3
|
||||
// gateway and mount follow with - stops on nil. A context.Canceled from the
|
||||
// subscription context is the stop signal surfacing, not a real error, and
|
||||
// reads as non-retryable to transient-error classifiers. Unavailable is what
|
||||
// a dropped connection looks like, so followers reconnect with their usual
|
||||
// retry. Handler errors pass through, and so does the end of a stream the
|
||||
// client closed.
|
||||
func (fs *FilerServer) endOfSubscription(stream context.Context, err error) error {
|
||||
if stream.Err() != nil {
|
||||
return err
|
||||
}
|
||||
if fs.subscriptionsStopped == nil || fs.subscriptionsStopped.Err() == nil {
|
||||
return err
|
||||
}
|
||||
if err == nil || errors.Is(err, context.Canceled) {
|
||||
return status.Error(codes.Unavailable, "filer is shutting down")
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
func (fs *FilerServer) SubscribeLocalMetadata(req *filer_pb.SubscribeMetadataRequest, stream filer_pb.SeaweedFiler_SubscribeLocalMetadataServer) error {
|
||||
return fs.subscribeLocalMetadata(req, stream, nil)
|
||||
}
|
||||
@@ -1012,9 +1091,11 @@ func (fs *FilerServer) SubscribeLocalMetadata(req *filer_pb.SubscribeMetadataReq
|
||||
// aggregation streams pass upgradeOnRemotePeer == nil; the SubscribeMetadata
|
||||
// delegation passes the aggregator's arrival channel so the stream ends
|
||||
// when a remote peer appears and the client reconnects to the aggregated path.
|
||||
func (fs *FilerServer) subscribeLocalMetadata(req *filer_pb.SubscribeMetadataRequest, stream metadataLocalStream, upgradeOnRemotePeer <-chan struct{}) error {
|
||||
func (fs *FilerServer) subscribeLocalMetadata(req *filer_pb.SubscribeMetadataRequest, stream metadataLocalStream, upgradeOnRemotePeer <-chan struct{}) (err error) {
|
||||
|
||||
ctx := stream.Context()
|
||||
ctx, cancelSubscription := fs.subscriptionContext(stream.Context())
|
||||
defer cancelSubscription()
|
||||
defer func() { err = fs.endOfSubscription(stream.Context(), err) }()
|
||||
peerAddress := findClientAddress(ctx, 0)
|
||||
|
||||
// use negative client id to differentiate from addClient()/deleteClient() used in SubscribeMetadata()
|
||||
@@ -1033,7 +1114,7 @@ func (fs *FilerServer) subscribeLocalMetadata(req *filer_pb.SubscribeMetadataReq
|
||||
lastReadTime := log_buffer.NewMessagePosition(req.SinceNs, gapResumeCursorOffset)
|
||||
glog.V(0).Infof(" + %v local subscribe %s from %+v clientId:%d", clientName, req.PathPrefix, lastReadTime, req.ClientId)
|
||||
|
||||
sender := newPipelinedSender(stream, 1024, req.ClientSupportsBatching)
|
||||
sender := newPipelinedSender(ctx, stream, 1024, req.ClientSupportsBatching)
|
||||
defer sender.Close()
|
||||
|
||||
// Bounded gap waits use the buffer's subscriber notification plus a retry
|
||||
@@ -1059,7 +1140,7 @@ func (fs *FilerServer) subscribeLocalMetadata(req *filer_pb.SubscribeMetadataReq
|
||||
var lastSeenTsNs int64
|
||||
var lastHeartbeatNs int64
|
||||
var lastFlushReportNs int64
|
||||
baseEachLogEntryFn := eachLogEntryFn(req, sender, eachEventNotificationFn, &unsyncedEvents)
|
||||
baseEachLogEntryFn := eachLogEntryFn(ctx, req, sender, eachEventNotificationFn, &unsyncedEvents)
|
||||
eachLogEntryFn := func(logEntry *filer_pb.LogEntry) (bool, error) {
|
||||
if upgradeOnRemotePeer != nil {
|
||||
select {
|
||||
@@ -1224,12 +1305,17 @@ func (fs *FilerServer) subscribeLocalMetadata(req *filer_pb.SubscribeMetadataReq
|
||||
|
||||
}
|
||||
|
||||
func eachLogEntryFn(req *filer_pb.SubscribeMetadataRequest, sender metadataStreamSender, eachEventNotificationFn func(dirPath string, eventNotification *filer_pb.EventNotification, tsNs int64) error, filtered *int64) log_buffer.EachLogEntryFuncType {
|
||||
func eachLogEntryFn(ctx context.Context, req *filer_pb.SubscribeMetadataRequest, sender metadataStreamSender, eachEventNotificationFn func(dirPath string, eventNotification *filer_pb.EventNotification, tsNs int64) error, filtered *int64) log_buffer.EachLogEntryFuncType {
|
||||
// A shallow scan of the path fields skips unmarshaling chunk-heavy events
|
||||
// this subscriber would filter out anyway; scan surprises fall back to the
|
||||
// full decode. Only a delivery resets the shared unsynced-events counter.
|
||||
prefilter := req.PathPrefix != "" || len(req.PathPrefixes) > 0 || len(req.Directories) > 0
|
||||
return func(logEntry *filer_pb.LogEntry) (bool, error) {
|
||||
// A cancelled context ends the pass here: neither loop checks it
|
||||
// per entry, and the send path alone cannot cover filtered entries.
|
||||
if ctx.Err() != nil {
|
||||
return true, nil
|
||||
}
|
||||
if prefilter {
|
||||
if skeleton, ok := filer_pb.ScanMetadataEventSkeleton(logEntry.Data); ok &&
|
||||
!filer_pb.MetadataEventMatchesSubscription(skeleton, req.PathPrefix, req.PathPrefixes, req.Directories) {
|
||||
@@ -1370,7 +1456,7 @@ func (fs *FilerServer) chunkDiskPass(ctx context.Context, sender metadataStreamS
|
||||
if len(refs) == 0 {
|
||||
return startPos.Time.UnixNano(), false, nil
|
||||
}
|
||||
if err := fs.sendRefsBatched(sender, refs, upgradeOnRemotePeer); err != nil {
|
||||
if err := fs.sendRefsBatched(ctx, sender, refs, upgradeOnRemotePeer); err != nil {
|
||||
return 0, false, err
|
||||
}
|
||||
if upgradeOnRemotePeer != nil {
|
||||
@@ -1430,9 +1516,12 @@ func (fs *FilerServer) chunkDiskPass(ctx context.Context, sender metadataStreamS
|
||||
// sendRefsBatched sends refs through the pipelined sender, which keeps them
|
||||
// out of Events batches; gRPC allows one sending goroutine per stream and the
|
||||
// sender's goroutine is it.
|
||||
func (fs *FilerServer) sendRefsBatched(sender metadataStreamSender, refs []*filer_pb.LogFileChunkRef, upgradeOnRemotePeer <-chan struct{}) error {
|
||||
func (fs *FilerServer) sendRefsBatched(ctx context.Context, sender metadataStreamSender, refs []*filer_pb.LogFileChunkRef, upgradeOnRemotePeer <-chan struct{}) error {
|
||||
const maxRefsPerMessage = 64
|
||||
for i := 0; i < len(refs); i += maxRefsPerMessage {
|
||||
if err := ctx.Err(); err != nil {
|
||||
return err
|
||||
}
|
||||
if upgradeOnRemotePeer != nil {
|
||||
select {
|
||||
case <-upgradeOnRemotePeer:
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package weed_server
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
"google.golang.org/protobuf/proto"
|
||||
@@ -37,7 +38,7 @@ func TestEachLogEntryFnPrefilterSkipsDecode(t *testing.T) {
|
||||
sender := &recordingSender{}
|
||||
var decoded int
|
||||
var unsyncedEvents int64
|
||||
fn := eachLogEntryFn(req, sender, func(dirPath string, eventNotification *filer_pb.EventNotification, tsNs int64) error {
|
||||
fn := eachLogEntryFn(context.Background(), req, sender, func(dirPath string, eventNotification *filer_pb.EventNotification, tsNs int64) error {
|
||||
decoded++
|
||||
unsyncedEvents = 0 // emulate a delivery, like the notification fn after a send
|
||||
return nil
|
||||
@@ -83,7 +84,7 @@ func TestEachLogEntryFnNoFilterDecodesEverything(t *testing.T) {
|
||||
req := &filer_pb.SubscribeMetadataRequest{}
|
||||
var decoded int
|
||||
var unsyncedEvents int64
|
||||
fn := eachLogEntryFn(req, &recordingSender{}, func(dirPath string, eventNotification *filer_pb.EventNotification, tsNs int64) error {
|
||||
fn := eachLogEntryFn(context.Background(), req, &recordingSender{}, func(dirPath string, eventNotification *filer_pb.EventNotification, tsNs int64) error {
|
||||
decoded++
|
||||
return nil
|
||||
}, &unsyncedEvents)
|
||||
|
||||
@@ -1,6 +1,8 @@
|
||||
package weed_server
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"fmt"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
@@ -131,7 +133,7 @@ func TestPipelinedSenderThroughput(t *testing.T) {
|
||||
var batchedRate float64
|
||||
t.Run("pipelined_batched_send", func(t *testing.T) {
|
||||
stream := &slowStream{sendDelay: sendDelay}
|
||||
sender := newPipelinedSender(stream, 1024, true)
|
||||
sender := newPipelinedSender(context.Background(), stream, 1024, true)
|
||||
|
||||
start := time.Now()
|
||||
for _, file := range files {
|
||||
@@ -217,7 +219,7 @@ func TestBatchingAdaptive(t *testing.T) {
|
||||
|
||||
t.Run("old_events_are_batched", func(t *testing.T) {
|
||||
stream := &slowStream{sendDelay: 10 * time.Microsecond}
|
||||
sender := newPipelinedSender(stream, 1024, true)
|
||||
sender := newPipelinedSender(context.Background(), stream, 1024, true)
|
||||
|
||||
// Push all events at once (no read delay) so the sender can batch aggressively
|
||||
for _, ev := range makeOldEvents(numEvents) {
|
||||
@@ -237,7 +239,7 @@ func TestBatchingAdaptive(t *testing.T) {
|
||||
|
||||
t.Run("recent_events_sent_individually", func(t *testing.T) {
|
||||
stream := &slowStream{sendDelay: 10 * time.Microsecond}
|
||||
sender := newPipelinedSender(stream, 1024, true)
|
||||
sender := newPipelinedSender(context.Background(), stream, 1024, true)
|
||||
|
||||
for _, ev := range makeRecentEvents(numEvents) {
|
||||
sender.Send(ev)
|
||||
@@ -281,7 +283,7 @@ func TestPipelinedSenderErrorPropagation(t *testing.T) {
|
||||
t.Run("send_returns_error", func(t *testing.T) {
|
||||
// Stream fails after 5 successful sends
|
||||
stream := &errorStreamImpl{failAfter: 5, err: sendErr}
|
||||
sender := newPipelinedSender(stream, 4, true)
|
||||
sender := newPipelinedSender(context.Background(), stream, 4, true)
|
||||
|
||||
var lastErr error
|
||||
for i := 0; i < 100; i++ {
|
||||
@@ -303,7 +305,7 @@ func TestPipelinedSenderErrorPropagation(t *testing.T) {
|
||||
// since Send may have already returned before the sender goroutine
|
||||
// processes the message.
|
||||
stream := &errorStreamImpl{failAfter: 0, err: sendErr}
|
||||
sender := newPipelinedSender(stream, 1024, true)
|
||||
sender := newPipelinedSender(context.Background(), stream, 1024, true)
|
||||
|
||||
ev := makeOldEvents(1)[0]
|
||||
sender.Send(ev)
|
||||
@@ -352,7 +354,7 @@ func TestPipelinedSingleVsParallelStreams(t *testing.T) {
|
||||
// simulatePipeline: read files with I/O delay, push events, send via pipelinedSender
|
||||
simulatePipeline := func(files []logFile) (eventsSent, sends int64, elapsed time.Duration, err error) {
|
||||
stream := &slowStream{sendDelay: sendDelay}
|
||||
sender := newPipelinedSender(stream, 1024, true)
|
||||
sender := newPipelinedSender(context.Background(), stream, 1024, true)
|
||||
|
||||
start := time.Now()
|
||||
outer:
|
||||
|
||||
@@ -1,6 +1,8 @@
|
||||
package weed_server
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"errors"
|
||||
"sync"
|
||||
"testing"
|
||||
@@ -55,7 +57,7 @@ func (s *gatedRecordingStream) Send(msg *filer_pb.SubscribeMetadataResponse) err
|
||||
// reads past it.
|
||||
func TestPipelinedSenderControlMessagesNeverNested(t *testing.T) {
|
||||
stream := &gatedRecordingStream{gate: make(chan struct{})}
|
||||
sender := newPipelinedSender(stream, 16, true)
|
||||
sender := newPipelinedSender(context.Background(), stream, 16, true)
|
||||
|
||||
oldTs := time.Now().Add(-time.Hour).UnixNano()
|
||||
if err := sender.Send(makeEvent("/d", "e1", oldTs)); err != nil {
|
||||
|
||||
@@ -1,6 +1,9 @@
|
||||
package weed_server
|
||||
|
||||
import (
|
||||
"context"
|
||||
"sync/atomic"
|
||||
|
||||
"fmt"
|
||||
"sync"
|
||||
"testing"
|
||||
@@ -46,7 +49,7 @@ func (s *recordingStream) snapshot() []*filer_pb.SubscribeMetadataResponse {
|
||||
// event. Everything must arrive, in order, whatever the batcher does.
|
||||
func TestPipelinedSenderRefsNeverBatched(t *testing.T) {
|
||||
stream := &recordingStream{slow: 2 * time.Millisecond} // let the queue back up so batching engages
|
||||
sender := newPipelinedSender(stream, 64, true)
|
||||
sender := newPipelinedSender(context.Background(), stream, 64, true)
|
||||
|
||||
oldTs := time.Now().Add(-time.Hour).UnixNano() // far behind: the batch heuristic fires
|
||||
var wantOrder []string
|
||||
@@ -102,3 +105,89 @@ func TestPipelinedSenderRefsNeverBatched(t *testing.T) {
|
||||
t.Fatalf("delivery order/count changed:\n got %v\nwant %v", gotOrder, wantOrder)
|
||||
}
|
||||
}
|
||||
|
||||
type blockedStream struct {
|
||||
release chan struct{}
|
||||
entered chan struct{}
|
||||
sends atomic.Int32
|
||||
}
|
||||
|
||||
func (s *blockedStream) Send(*filer_pb.SubscribeMetadataResponse) error {
|
||||
s.sends.Add(1)
|
||||
select {
|
||||
case s.entered <- struct{}{}:
|
||||
default:
|
||||
}
|
||||
<-s.release
|
||||
return nil
|
||||
}
|
||||
|
||||
// A subscriber that stops reading wedges stream.Send. The subscription
|
||||
// context must release Send and Close, so the handler can return and end the
|
||||
// stream during filer shutdown.
|
||||
func TestPipelinedSenderUnblocksOnCancel(t *testing.T) {
|
||||
stream := &blockedStream{release: make(chan struct{})}
|
||||
defer close(stream.release)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
sender := newPipelinedSender(ctx, stream, 1, false)
|
||||
|
||||
msg := &filer_pb.SubscribeMetadataResponse{EventNotification: &filer_pb.EventNotification{}}
|
||||
if err := sender.Send(msg); err != nil { // consumed by the blocked sendLoop
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := sender.Send(msg); err != nil { // fills the one-deep queue
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
cancel()
|
||||
|
||||
done := make(chan error, 1)
|
||||
go func() { done <- sender.Send(msg) }()
|
||||
select {
|
||||
case err := <-done:
|
||||
if err != context.Canceled {
|
||||
t.Fatalf("Send returned %v, want context.Canceled", err)
|
||||
}
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("Send still blocked after cancel")
|
||||
}
|
||||
if err := sender.Close(); err != nil {
|
||||
t.Fatalf("Close returned %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// Once the subscription context ends, sendLoop must not start another
|
||||
// stream.Send: the handler returns and gRPC ends the stream behind it, so a
|
||||
// send started then would race the stream's teardown. The one Send already
|
||||
// in flight is released by that teardown.
|
||||
func TestPipelinedSenderStopsSendingOnCancel(t *testing.T) {
|
||||
stream := &blockedStream{release: make(chan struct{}), entered: make(chan struct{}, 1)}
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
sender := newPipelinedSender(ctx, stream, 4, false)
|
||||
|
||||
msg := &filer_pb.SubscribeMetadataResponse{EventNotification: &filer_pb.EventNotification{}}
|
||||
if err := sender.Send(msg); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
<-stream.entered // sendLoop is inside stream.Send
|
||||
for i := 0; i < 2; i++ { // queued behind the blocked send
|
||||
if err := sender.Send(msg); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
cancel()
|
||||
close(stream.release)
|
||||
|
||||
if err := sender.Close(); err != nil {
|
||||
t.Fatalf("Close returned %v", err)
|
||||
}
|
||||
select {
|
||||
case <-sender.done:
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("sendLoop did not exit after cancel")
|
||||
}
|
||||
if n := stream.sends.Load(); n != 1 {
|
||||
t.Fatalf("stream.Send ran %d times, want 1; a new send started after cancel", n)
|
||||
}
|
||||
}
|
||||
@@ -170,6 +170,12 @@ type FilerServer struct {
|
||||
posixLocks *posixlock.Manager
|
||||
// posixLockSweeperStop stops the lease-reaping sweeper goroutine on Shutdown.
|
||||
posixLockSweeperStop chan struct{}
|
||||
|
||||
// subscriptionsStopped is cancelled by StopSubscriptions when the filer
|
||||
// starts shutting down; every metadata subscription derives its context
|
||||
// from it as well as from its stream.
|
||||
subscriptionsStopped context.Context
|
||||
stopSubscriptions context.CancelFunc
|
||||
// posixLockReadyAt is the unix-nanos when this filer began serving POSIX
|
||||
// locks. For posixLockWarmup after it, the owner defers would-be grants while
|
||||
// mounts re-assert, so a (re)started owner does not double-grant from empty
|
||||
@@ -220,6 +226,7 @@ func NewFilerServer(defaultMux, readonlyMux *http.ServeMux, option *FilerOption)
|
||||
entryLockTable: util.NewLockTable[util.FullPath](),
|
||||
posixLocks: posixlock.NewManager(),
|
||||
}
|
||||
fs.subscriptionsStopped, fs.stopSubscriptions = context.WithCancel(context.Background())
|
||||
fs.startPosixLockSweeper()
|
||||
fs.mountPeerRegistry = filer.NewMountPeerRegistry()
|
||||
go fs.runMountPeerRegistrySweeper()
|
||||
|
||||
@@ -0,0 +1,150 @@
|
||||
package weed_server
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/status"
|
||||
)
|
||||
|
||||
func newStoppableFilerServer() *FilerServer {
|
||||
fs := &FilerServer{}
|
||||
fs.subscriptionsStopped, fs.stopSubscriptions = context.WithCancel(context.Background())
|
||||
return fs
|
||||
}
|
||||
|
||||
func TestStopSubscriptionsEndsOpenSubscriptions(t *testing.T) {
|
||||
fs := newStoppableFilerServer()
|
||||
ctx, cancel := fs.subscriptionContext(context.Background())
|
||||
defer cancel()
|
||||
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
t.Fatal("subscription ended before StopSubscriptions")
|
||||
default:
|
||||
}
|
||||
fs.StopSubscriptions()
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("subscription still open after StopSubscriptions")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSubscriptionStartedAfterStopEndsAtOnce(t *testing.T) {
|
||||
fs := newStoppableFilerServer()
|
||||
fs.StopSubscriptions()
|
||||
ctx, cancel := fs.subscriptionContext(context.Background())
|
||||
defer cancel()
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("a subscription started during shutdown stayed open")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSubscriptionStillEndsWithItsStream(t *testing.T) {
|
||||
fs := newStoppableFilerServer()
|
||||
stream, closeStream := context.WithCancel(context.WithValue(context.Background(), struct{}{}, "peer"))
|
||||
ctx, cancel := fs.subscriptionContext(stream)
|
||||
defer cancel()
|
||||
if ctx.Value(struct{}{}) != "peer" {
|
||||
t.Error("subscription context lost the stream's values")
|
||||
}
|
||||
closeStream()
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("subscription outlived its stream")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSubscriptionContextWithoutStopSignal(t *testing.T) {
|
||||
fs := &FilerServer{}
|
||||
fs.StopSubscriptions()
|
||||
ctx, cancel := fs.subscriptionContext(context.Background())
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
t.Fatal("subscription ended although nothing stopped it")
|
||||
default:
|
||||
}
|
||||
cancel()
|
||||
<-ctx.Done()
|
||||
}
|
||||
|
||||
// A follower reads a clean end as "done" and stops for good (util.RetryUntil
|
||||
// returns on nil), so a subscription that StopSubscriptions ended must not
|
||||
// end cleanly.
|
||||
func TestStoppedSubscriptionEndsUnavailable(t *testing.T) {
|
||||
fs := newStoppableFilerServer()
|
||||
fs.StopSubscriptions()
|
||||
err := fs.endOfSubscription(context.Background(), nil)
|
||||
if status.Code(err) != codes.Unavailable {
|
||||
t.Fatalf("stopped subscription ended with %v, want Unavailable", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSubscriptionEndsCleanlyWithoutStop(t *testing.T) {
|
||||
fs := newStoppableFilerServer()
|
||||
if err := fs.endOfSubscription(context.Background(), nil); err != nil {
|
||||
t.Fatalf("subscription without shutdown ended with %v, want nil", err)
|
||||
}
|
||||
if err := (&FilerServer{}).endOfSubscription(context.Background(), nil); err != nil {
|
||||
t.Fatalf("server without stop signal ended with %v, want nil", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSubscriptionErrorPassesThroughStop(t *testing.T) {
|
||||
fs := newStoppableFilerServer()
|
||||
fs.StopSubscriptions()
|
||||
want := errors.New("reading from persisted logs")
|
||||
if err := fs.endOfSubscription(context.Background(), want); !errors.Is(err, want) {
|
||||
t.Fatalf("got %v, want the handler's own error", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestClosedStreamEndsCleanlyDuringStop(t *testing.T) {
|
||||
fs := newStoppableFilerServer()
|
||||
fs.StopSubscriptions()
|
||||
stream, closeStream := context.WithCancel(context.Background())
|
||||
closeStream()
|
||||
if err := fs.endOfSubscription(stream, nil); err != nil {
|
||||
t.Fatalf("stream the client closed ended with %v, want nil", err)
|
||||
}
|
||||
}
|
||||
|
||||
// A read pass surfaces StopSubscriptions as a wrapped context.Canceled; on the
|
||||
// wire that must still read as Unavailable so followers retry.
|
||||
func TestCanceledReadEndsUnavailableDuringStop(t *testing.T) {
|
||||
fs := newStoppableFilerServer()
|
||||
fs.StopSubscriptions()
|
||||
canceled := fmt.Errorf("reading from persisted logs: %w", context.Canceled)
|
||||
if err := fs.endOfSubscription(context.Background(), canceled); status.Code(err) != codes.Unavailable {
|
||||
t.Fatalf("stopped read ended with %v, want Unavailable", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCanceledReadPassesThroughWithoutStop(t *testing.T) {
|
||||
fs := newStoppableFilerServer()
|
||||
canceled := fmt.Errorf("reading from persisted logs: %w", context.Canceled)
|
||||
if err := fs.endOfSubscription(context.Background(), canceled); !errors.Is(err, context.Canceled) {
|
||||
t.Fatalf("got %v, want the handler's own error", err)
|
||||
}
|
||||
}
|
||||
|
||||
// Neither read loop checks ctx per entry; the entry callback does.
|
||||
func TestEachLogEntryFnEndsOnCancel(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
cancel()
|
||||
var filtered int64
|
||||
fn := eachLogEntryFn(ctx, &filer_pb.SubscribeMetadataRequest{}, nil, nil, &filtered)
|
||||
done, err := fn(&filer_pb.LogEntry{})
|
||||
if !done || err != nil {
|
||||
t.Fatalf("cancelled subscription returned done=%v err=%v, want a clean end", done, err)
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user