From 23ec9b4979cc656af456a1cd8a78a638d7ab31ac Mon Sep 17 00:00:00 2001 From: fatedier Date: Fri, 28 Aug 2026 10:56:05 +0800 Subject: [PATCH] vnet: serialize VirtualNet route lifecycle (#5512) * vnet: serialize virtual net route lifecycle * test: remove fixed virtual net helper parameters --- pkg/plugin/visitor/virtual_net.go | 98 ++++++++---- pkg/plugin/visitor/virtual_net_test.go | 198 +++++++++++++++++++++++++ pkg/vnet/controller_test.go | 7 + 3 files changed, 273 insertions(+), 30 deletions(-) diff --git a/pkg/plugin/visitor/virtual_net.go b/pkg/plugin/visitor/virtual_net.go index 87a67744..f95856c2 100644 --- a/pkg/plugin/visitor/virtual_net.go +++ b/pkg/plugin/visitor/virtual_net.go @@ -20,6 +20,7 @@ import ( "context" "errors" "fmt" + "io" "net" "sync" "time" @@ -33,10 +34,16 @@ func init() { Register(v1.VisitorPluginVirtualNet, NewVirtualNetPlugin) } +type clientRouteController interface { + RegisterClientRoute(context.Context, string, []net.IPNet, io.ReadWriteCloser) + UnregisterClientRoute(string, io.Writer) bool +} + type VirtualNetPlugin struct { pluginCtx PluginContext - routes []net.IPNet + routeController clientRouteController + routes []net.IPNet mu sync.Mutex controllerConn net.Conn @@ -60,6 +67,9 @@ func NewVirtualNetPlugin(pluginCtx PluginContext, options v1.VisitorPluginOption pluginCtx: pluginCtx, routes: make([]net.IPNet, 0), } + if pluginCtx.VnetController != nil { + p.routeController = pluginCtx.VnetController + } p.ctx, p.cancel = context.WithCancel(pluginCtx.Ctx) @@ -90,7 +100,7 @@ func (p *VirtualNetPlugin) Name() string { func (p *VirtualNetPlugin) Start() { xl := xlog.FromContextSafe(p.pluginCtx.Ctx) - if p.pluginCtx.VnetController == nil { + if p.routeController == nil { return } @@ -116,16 +126,17 @@ func (p *VirtualNetPlugin) run() { select { case <-p.ctx.Done(): xl.Infof("VirtualNetPlugin run loop for visitor [%s] stopping (context cancelled before pipe creation).", p.pluginCtx.Name) - p.cleanupControllerConn(xl) + p.cleanupCurrentControllerConn(xl) return default: } controllerConn, pluginConn := net.Pipe() - - p.mu.Lock() - p.controllerConn = controllerConn - p.mu.Unlock() + xl.Infof("attempting to register client route for visitor [%s]", p.pluginCtx.Name) + if !p.registerControllerConn(controllerConn, pluginConn) { + xl.Infof("VirtualNetPlugin run loop for visitor [%s] stopping (context cancelled before route registration).", p.pluginCtx.Name) + return + } // Wrap with CloseNotifyConn which supports both close notification and error recording var closeErr error @@ -134,8 +145,6 @@ func (p *VirtualNetPlugin) run() { close(currentCloseSignal) // Signal the run loop on close. }) - xl.Infof("attempting to register client route for visitor [%s]", p.pluginCtx.Name) - p.pluginCtx.VnetController.RegisterClientRoute(p.ctx, p.pluginCtx.Name, p.routes, controllerConn) xl.Infof("successfully registered client route for visitor [%s]. Starting connection handler with CloseNotifyConn.", p.pluginCtx.Name) // Pass the CloseNotifyConn to the visitor for handling. @@ -146,7 +155,7 @@ func (p *VirtualNetPlugin) run() { select { case <-p.ctx.Done(): xl.Infof("VirtualNetPlugin run loop stopping for visitor [%s] (context cancelled while waiting).", p.pluginCtx.Name) - p.cleanupControllerConn(xl) + p.cleanupControllerConn(xl, controllerConn) return case <-currentCloseSignal: // Determine reconnect delay based on error with exponential backoff @@ -171,7 +180,7 @@ func (p *VirtualNetPlugin) run() { } // The visitor closed the plugin side. Close the controller side. - p.cleanupControllerConn(xl) + p.cleanupControllerConn(xl, controllerConn) xl.Infof("waiting %v before attempting reconnection for visitor [%s]...", reconnectDelay, p.pluginCtx.Name) select { @@ -186,6 +195,23 @@ func (p *VirtualNetPlugin) run() { } } +// registerControllerConn publishes and registers controllerConn atomically with +// respect to Close. A canceled plugin cannot register a new route. +func (p *VirtualNetPlugin) registerControllerConn(controllerConn, pluginConn net.Conn) bool { + p.mu.Lock() + if p.ctx.Err() != nil || p.routeController == nil { + p.mu.Unlock() + _ = controllerConn.Close() + _ = pluginConn.Close() + return false + } + + p.controllerConn = controllerConn + p.routeController.RegisterClientRoute(p.ctx, p.pluginCtx.Name, p.routes, controllerConn) + p.mu.Unlock() + return true +} + // virtualNetReconnectDelay returns a bounded reconnect delay without allowing // the exponential shift to overflow for large consecutive error counts. func virtualNetReconnectDelay(consecutiveErrors int) time.Duration { @@ -198,16 +224,37 @@ func virtualNetReconnectDelay(consecutiveErrors int) time.Duration { return virtualNetReconnectBaseDelay * time.Duration(1<