package notify import ( "context" "net" "sync" ) type bulkDedicatedSidecar struct { laneID uint32 conn net.Conn closeOnce sync.Once connMu sync.Mutex closed bool senderMu sync.Mutex sender *bulkDedicatedLaneSender } type bulkDedicatedLane struct { id uint32 route clientSessionRoute activeBulks int sidecar *bulkDedicatedSidecar attachFlight *bulkDedicatedAttachFlight } type bulkDedicatedAttachFlight struct { done chan struct{} once sync.Once err error } func normalizeBulkDedicatedLaneID(laneID uint32) uint32 { if laneID == 0 { return 1 } return laneID } func newBulkDedicatedSidecar(conn net.Conn, laneID uint32) *bulkDedicatedSidecar { if conn == nil { return nil } return &bulkDedicatedSidecar{ laneID: normalizeBulkDedicatedLaneID(laneID), conn: conn, } } func newBulkDedicatedAttachFlight() *bulkDedicatedAttachFlight { return &bulkDedicatedAttachFlight{ done: make(chan struct{}), } } func (f *bulkDedicatedAttachFlight) finish(err error) { if f == nil { return } f.once.Do(func() { f.err = err close(f.done) }) } func (f *bulkDedicatedAttachFlight) wait(ctx context.Context) error { if f == nil { return nil } if ctx == nil { ctx = context.Background() } select { case <-ctx.Done(): return ctx.Err() case <-f.done: return f.err } } func (s *bulkDedicatedSidecar) close() { if s == nil { return } s.closeOnce.Do(func() { s.connMu.Lock() s.closed = true conn := s.conn s.connMu.Unlock() if conn != nil { _ = conn.Close() } if sender := s.laneSenderSnapshot(); sender != nil { sender.stop() } }) } // withConn serializes logical attachment with sidecar cleanup. Without this // boundary a cleanup can close the socket after lookup but before the bulk // handle installs it. func (s *bulkDedicatedSidecar) withConn(fn func(net.Conn) error) error { if s == nil || fn == nil { return errTransportDetached } s.connMu.Lock() defer s.connMu.Unlock() if s.closed || s.conn == nil { return errTransportDetached } return fn(s.conn) } func (s *bulkDedicatedSidecar) laneSenderSnapshot() *bulkDedicatedLaneSender { if s == nil { return nil } s.senderMu.Lock() defer s.senderMu.Unlock() return s.sender } func (s *bulkDedicatedSidecar) laneSenderWithFactory(factory func(net.Conn) *bulkDedicatedLaneSender) *bulkDedicatedLaneSender { if s == nil || factory == nil { return nil } s.senderMu.Lock() defer s.senderMu.Unlock() if s.sender != nil { return s.sender } s.connMu.Lock() if s.closed || s.conn == nil { s.connMu.Unlock() return nil } s.sender = factory(s.conn) s.connMu.Unlock() return s.sender } func (c *ClientCommon) clientDedicatedSidecarSnapshot() *bulkDedicatedSidecar { if c == nil { return nil } route := c.clientSessionRouteSnapshot() c.bulkDedicatedSidecarMu.Lock() defer c.bulkDedicatedSidecarMu.Unlock() var selected *bulkDedicatedSidecar var bestID uint32 for laneID, lane := range c.bulkDedicatedLanes { if lane == nil || lane.sidecar == nil || !sameClientDedicatedLaneRoute(lane, route) { continue } if selected == nil || laneID < bestID { selected = lane.sidecar bestID = laneID } } return selected } func (c *ClientCommon) reserveBulkDedicatedLane() uint32 { laneID, _ := c.reserveBulkDedicatedLaneAtRoute(c.clientSessionRouteSnapshot()) return laneID } func (c *ClientCommon) reserveBulkDedicatedLaneAtRoute(route clientSessionRoute) (uint32, error) { if c == nil { return 0, errBulkClientNil } if route.bound() { if err := c.ensureClientSessionRouteSendReady(route); err != nil { return 0, err } } c.bulkDedicatedSidecarMu.Lock() if route.bound() && !c.clientSessionRouteCurrent(route) { c.bulkDedicatedSidecarMu.Unlock() return 0, transportDetachedSessionEpochError() } if c.bulkDedicatedLanes == nil { c.bulkDedicatedLanes = make(map[uint32]*bulkDedicatedLane) } limit := c.bulkDedicatedLaneLimitSnapshot() var best *bulkDedicatedLane for _, lane := range c.bulkDedicatedLanes { if lane == nil || !sameClientDedicatedLaneRoute(lane, route) { continue } if best == nil || lane.activeBulks < best.activeBulks || (lane.activeBulks == best.activeBulks && lane.id < best.id) { best = lane } } if best == nil || ((limit <= 0 || len(c.bulkDedicatedLanes) < limit) && best.activeBulks > 0) { laneID := c.bulkDedicatedNextLaneID for { laneID++ laneID = normalizeBulkDedicatedLaneID(laneID) if _, exists := c.bulkDedicatedLanes[laneID]; !exists { break } } c.bulkDedicatedNextLaneID = laneID best = &bulkDedicatedLane{id: laneID, route: route} c.bulkDedicatedLanes[laneID] = best } best.activeBulks++ c.bulkDedicatedSidecarMu.Unlock() return best.id, nil } func (c *ClientCommon) retainBulkDedicatedLane(laneID uint32) uint32 { _ = c.retainBulkDedicatedLaneAtRoute(laneID, c.clientSessionRouteSnapshot()) return normalizeBulkDedicatedLaneID(laneID) } func (c *ClientCommon) retainBulkDedicatedLaneAtRoute(laneID uint32, route clientSessionRoute) error { laneID = normalizeBulkDedicatedLaneID(laneID) if c == nil { return errBulkClientNil } if route.bound() { if err := c.ensureClientSessionRouteSendReady(route); err != nil { return err } } c.bulkDedicatedSidecarMu.Lock() if route.bound() && !c.clientSessionRouteCurrent(route) { c.bulkDedicatedSidecarMu.Unlock() return transportDetachedSessionEpochError() } if c.bulkDedicatedLanes == nil { c.bulkDedicatedLanes = make(map[uint32]*bulkDedicatedLane) } lane := c.bulkDedicatedLanes[laneID] var retiredSidecar *bulkDedicatedSidecar var retiredFlight *bulkDedicatedAttachFlight if lane == nil || !sameClientDedicatedLaneRoute(lane, route) { if lane != nil { retiredSidecar = lane.sidecar retiredFlight = lane.attachFlight } lane = &bulkDedicatedLane{id: laneID, route: route} c.bulkDedicatedLanes[laneID] = lane } lane.activeBulks++ c.bulkDedicatedSidecarMu.Unlock() if retiredSidecar != nil { retiredSidecar.close() } if retiredFlight != nil { retiredFlight.finish(errTransportDetached) } return nil } func (c *ClientCommon) releaseBulkDedicatedLane(laneID uint32) { c.releaseBulkDedicatedLaneAtRoute(laneID, c.clientSessionRouteSnapshot()) } func (c *ClientCommon) releaseBulkDedicatedLaneAtRoute(laneID uint32, route clientSessionRoute) { if c == nil { return } laneID = normalizeBulkDedicatedLaneID(laneID) c.bulkDedicatedSidecarMu.Lock() defer c.bulkDedicatedSidecarMu.Unlock() lane := c.bulkDedicatedLanes[laneID] if lane == nil || !sameClientDedicatedLaneRoute(lane, route) { return } if lane.activeBulks > 0 { lane.activeBulks-- } if lane.activeBulks == 0 && lane.sidecar == nil && lane.attachFlight == nil { delete(c.bulkDedicatedLanes, laneID) } } func sameClientDedicatedLaneRoute(lane *bulkDedicatedLane, route clientSessionRoute) bool { if lane == nil { return false } if lane.route.bound() || route.bound() { return sameClientSessionRoute(lane.route, route) } return true } func (c *ClientCommon) clientDedicatedSidecarSnapshotForLane(laneID uint32) *bulkDedicatedSidecar { if c == nil { return nil } route := c.clientSessionRouteSnapshot() return c.clientDedicatedSidecarSnapshotForLaneAtRoute(laneID, route) } func (c *ClientCommon) clientDedicatedSidecarSnapshotForLaneAtRoute(laneID uint32, route clientSessionRoute) *bulkDedicatedSidecar { if c == nil { return nil } if route.bound() && !c.clientSessionRouteCurrent(route) { return nil } laneID = normalizeBulkDedicatedLaneID(laneID) c.bulkDedicatedSidecarMu.Lock() defer c.bulkDedicatedSidecarMu.Unlock() if route.bound() && !c.clientSessionRouteCurrent(route) { return nil } if lane := c.bulkDedicatedLanes[laneID]; lane != nil { if !sameClientDedicatedLaneRoute(lane, route) { return nil } return lane.sidecar } return nil } func (c *ClientCommon) beginClientDedicatedSidecarAttach(laneID uint32, route clientSessionRoute) (*bulkDedicatedSidecar, *bulkDedicatedAttachFlight, bool, error) { if c == nil { return nil, nil, false, errBulkClientNil } if err := c.ensureClientSessionRouteSendReady(route); err != nil { return nil, nil, false, err } laneID = normalizeBulkDedicatedLaneID(laneID) c.bulkDedicatedSidecarMu.Lock() if !c.clientSessionRouteCurrent(route) { c.bulkDedicatedSidecarMu.Unlock() return nil, nil, false, transportDetachedSessionEpochError() } if c.bulkDedicatedLanes == nil { c.bulkDedicatedLanes = make(map[uint32]*bulkDedicatedLane) } lane := c.bulkDedicatedLanes[laneID] var retiredSidecar *bulkDedicatedSidecar var retiredFlight *bulkDedicatedAttachFlight if lane == nil || !sameClientDedicatedLaneRoute(lane, route) { if lane != nil { retiredSidecar = lane.sidecar retiredFlight = lane.attachFlight } lane = &bulkDedicatedLane{id: laneID, route: route} c.bulkDedicatedLanes[laneID] = lane } if lane.sidecar != nil { activeSidecar := lane.sidecar c.bulkDedicatedSidecarMu.Unlock() if retiredSidecar != nil { retiredSidecar.close() } if retiredFlight != nil { retiredFlight.finish(errTransportDetached) } return activeSidecar, nil, false, nil } if lane.attachFlight != nil { pendingFlight := lane.attachFlight c.bulkDedicatedSidecarMu.Unlock() if retiredSidecar != nil { retiredSidecar.close() } if retiredFlight != nil { retiredFlight.finish(errTransportDetached) } return nil, pendingFlight, false, nil } flight := newBulkDedicatedAttachFlight() lane.attachFlight = flight c.bulkDedicatedSidecarMu.Unlock() if retiredSidecar != nil { retiredSidecar.close() } if retiredFlight != nil { retiredFlight.finish(errTransportDetached) } return nil, flight, true, nil } func (c *ClientCommon) finishClientDedicatedSidecarAttach(laneID uint32, flight *bulkDedicatedAttachFlight, err error) { if c == nil || flight == nil { return } laneID = normalizeBulkDedicatedLaneID(laneID) c.bulkDedicatedSidecarMu.Lock() if lane := c.bulkDedicatedLanes[laneID]; lane != nil && lane.attachFlight == flight { lane.attachFlight = nil if lane.activeBulks == 0 && lane.sidecar == nil { delete(c.bulkDedicatedLanes, laneID) } } c.bulkDedicatedSidecarMu.Unlock() flight.finish(err) } func (c *ClientCommon) installClientDedicatedSidecar(laneID uint32, sidecar *bulkDedicatedSidecar) (*bulkDedicatedSidecar, bool) { active, installed, _ := c.installClientDedicatedSidecarAtRoute(laneID, sidecar, c.clientSessionRouteSnapshot()) return active, installed } func (c *ClientCommon) installClientDedicatedSidecarAtRoute(laneID uint32, sidecar *bulkDedicatedSidecar, route clientSessionRoute) (*bulkDedicatedSidecar, bool, error) { if c == nil || sidecar == nil { return nil, false, errBulkClientNil } if route.bound() { if err := c.ensureClientSessionRouteSendReady(route); err != nil { return nil, false, err } } laneID = normalizeBulkDedicatedLaneID(laneID) c.bulkDedicatedSidecarMu.Lock() if route.bound() && !c.clientSessionRouteCurrent(route) { c.bulkDedicatedSidecarMu.Unlock() return nil, false, transportDetachedSessionEpochError() } if c.bulkDedicatedLanes == nil { c.bulkDedicatedLanes = make(map[uint32]*bulkDedicatedLane) } lane := c.bulkDedicatedLanes[laneID] var retiredSidecar *bulkDedicatedSidecar var retiredFlight *bulkDedicatedAttachFlight if lane == nil || !sameClientDedicatedLaneRoute(lane, route) { if lane != nil { retiredSidecar = lane.sidecar retiredFlight = lane.attachFlight } lane = &bulkDedicatedLane{id: laneID, route: route} c.bulkDedicatedLanes[laneID] = lane } if lane.sidecar != nil { activeSidecar := lane.sidecar c.bulkDedicatedSidecarMu.Unlock() if retiredSidecar != nil { retiredSidecar.close() } if retiredFlight != nil { retiredFlight.finish(errTransportDetached) } return activeSidecar, false, nil } lane.sidecar = sidecar c.bulkDedicatedSidecarMu.Unlock() if retiredSidecar != nil { retiredSidecar.close() } if retiredFlight != nil { retiredFlight.finish(errTransportDetached) } return sidecar, true, nil } func (c *ClientCommon) clearClientDedicatedSidecar(laneID uint32, sidecar *bulkDedicatedSidecar) bool { if c == nil || sidecar == nil { return false } laneID = normalizeBulkDedicatedLaneID(laneID) c.bulkDedicatedSidecarMu.Lock() defer c.bulkDedicatedSidecarMu.Unlock() lane := c.bulkDedicatedLanes[laneID] if lane == nil || lane.sidecar != sidecar { return false } lane.sidecar = nil if lane.activeBulks == 0 && lane.attachFlight == nil { delete(c.bulkDedicatedLanes, laneID) } return true } func (c *ClientCommon) closeClientDedicatedSidecar() { c.closeClientDedicatedSidecarWithError(errServiceShutdown) } func (c *ClientCommon) closeClientDedicatedSidecarWithError(closeErr error) { if c == nil { return } if closeErr == nil { closeErr = errServiceShutdown } c.bulkDedicatedSidecarMu.Lock() lanes := c.bulkDedicatedLanes c.bulkDedicatedLanes = make(map[uint32]*bulkDedicatedLane) c.bulkDedicatedSidecarMu.Unlock() for _, lane := range lanes { if lane == nil { continue } if lane.sidecar != nil { lane.sidecar.close() } if lane.attachFlight != nil { lane.attachFlight.finish(closeErr) } } } func (c *ClientCommon) handleClientDedicatedSidecarFailure(sidecar *bulkDedicatedSidecar, err error) { if c == nil || sidecar == nil { return } if !c.clearClientDedicatedSidecar(sidecar.laneID, sidecar) { return } runtime := c.getBulkRuntime() if runtime != nil && sidecar.conn != nil { runtime.handleDedicatedReadErrorByConn(clientFileScope(), sidecar.conn, err) } sidecar.close() } func (s *ServerCommon) serverDedicatedSidecarSnapshot(logical *LogicalConn) *bulkDedicatedSidecar { if s == nil || logical == nil { return nil } s.bulkDedicatedSidecarMu.Lock() defer s.bulkDedicatedSidecarMu.Unlock() return firstServerDedicatedSidecarLocked(s.bulkDedicatedSidecars[logical]) } func firstServerDedicatedSidecarLocked(lanes map[uint32]*bulkDedicatedSidecar) *bulkDedicatedSidecar { var ( selected *bulkDedicatedSidecar bestID uint32 ) for laneID, sidecar := range lanes { if sidecar == nil { continue } if selected == nil || laneID < bestID { selected = sidecar bestID = laneID } } return selected } func (s *ServerCommon) serverDedicatedSidecarSnapshotForLane(logical *LogicalConn, laneID uint32) *bulkDedicatedSidecar { if s == nil || logical == nil { return nil } laneID = normalizeBulkDedicatedLaneID(laneID) s.bulkDedicatedSidecarMu.Lock() defer s.bulkDedicatedSidecarMu.Unlock() if lanes := s.bulkDedicatedSidecars[logical]; lanes != nil { return lanes[laneID] } return nil } func (s *ServerCommon) serverDedicatedSidecarCurrent(logical *LogicalConn, sidecar *bulkDedicatedSidecar) bool { if s == nil || logical == nil || sidecar == nil { return false } return s.serverDedicatedSidecarSnapshotForLane(logical, sidecar.laneID) == sidecar } func bulkDedicatedSidecarConnCurrent(bulk *bulkHandle, sidecar *bulkDedicatedSidecar) bool { if bulk == nil || sidecar == nil || sidecar.conn == nil { return false } return bulk.dedicatedConnSnapshot() == sidecar.conn } func (s *ServerCommon) installServerDedicatedSidecar(logical *LogicalConn, laneID uint32, sidecar *bulkDedicatedSidecar) *bulkDedicatedSidecar { if s == nil || logical == nil || sidecar == nil { return nil } laneID = normalizeBulkDedicatedLaneID(laneID) s.bulkDedicatedSidecarMu.Lock() defer s.bulkDedicatedSidecarMu.Unlock() lanes := s.bulkDedicatedSidecars[logical] if lanes == nil { lanes = make(map[uint32]*bulkDedicatedSidecar) s.bulkDedicatedSidecars[logical] = lanes } prev := lanes[laneID] lanes[laneID] = sidecar return prev } func (s *ServerCommon) clearServerDedicatedSidecar(logical *LogicalConn, laneID uint32, sidecar *bulkDedicatedSidecar) bool { if s == nil || logical == nil || sidecar == nil { return false } laneID = normalizeBulkDedicatedLaneID(laneID) s.bulkDedicatedSidecarMu.Lock() defer s.bulkDedicatedSidecarMu.Unlock() lanes := s.bulkDedicatedSidecars[logical] if lanes == nil || lanes[laneID] != sidecar { return false } delete(lanes, laneID) if len(lanes) == 0 { delete(s.bulkDedicatedSidecars, logical) } return true } func (s *ServerCommon) closeServerDedicatedSidecar(logical *LogicalConn) { if s == nil || logical == nil { return } s.bulkDedicatedSidecarMu.Lock() lanes := s.bulkDedicatedSidecars[logical] delete(s.bulkDedicatedSidecars, logical) s.bulkDedicatedSidecarMu.Unlock() for _, sidecar := range lanes { if sidecar != nil { sidecar.close() } } } func (s *ServerCommon) closeAllServerDedicatedSidecars() { if s == nil { return } s.bulkDedicatedSidecarMu.Lock() lanesByLogical := s.bulkDedicatedSidecars s.bulkDedicatedSidecars = make(map[*LogicalConn]map[uint32]*bulkDedicatedSidecar) s.bulkDedicatedSidecarMu.Unlock() for _, lanes := range lanesByLogical { for _, sidecar := range lanes { if sidecar != nil { sidecar.close() } } } } func (s *ServerCommon) handleServerDedicatedSidecarFailure(logical *LogicalConn, sidecar *bulkDedicatedSidecar, err error) { if s == nil || logical == nil || sidecar == nil { return } if !s.clearServerDedicatedSidecar(logical, sidecar.laneID, sidecar) { return } runtime := s.getBulkRuntime() if runtime != nil && sidecar.conn != nil { runtime.handleDedicatedReadErrorByConn(serverFileScope(logical), sidecar.conn, err) } sidecar.close() } func (s *ServerCommon) attachServerDedicatedSidecarIfExists(logical *LogicalConn, bulk *bulkHandle) { if s == nil || logical == nil || bulk == nil || !bulk.Dedicated() { return } sidecar := s.serverDedicatedSidecarSnapshotForLane(logical, bulk.dedicatedLaneIDSnapshot()) if sidecar == nil || sidecar.conn == nil { return } _ = sidecar.withConn(func(conn net.Conn) error { return bulk.attachDedicatedConnShared(conn) }) }