package notify import ( "fmt" "net" "strings" "sync" "sync/atomic" ) type bulkRuntime struct { rolePrefix string seq atomic.Uint64 dataSeq uint64 dataStart uint64 dataStep uint64 mu sync.RWMutex handler func(BulkAcceptInfo) error bulks map[string]*bulkHandle inbound map[string]map[uint64]*bulkHandle outbound map[string]map[uint64]*bulkHandle reserved map[string]map[uint64]struct{} } type bulkDataIndexDirection uint8 const ( bulkDataIndexInbound bulkDataIndexDirection = 1 << iota bulkDataIndexOutbound bulkDataIndexBoth = bulkDataIndexInbound | bulkDataIndexOutbound ) func newBulkRuntime(rolePrefix string) *bulkRuntime { dataStart, dataStep := uint64(1), uint64(1) // Client- and server-originated IDs occupy disjoint wire namespaces. A // bulk is duplex, so every incoming frame must be routable by DataID alone; // partitioning the allocator prevents simultaneous opens from colliding. if rolePrefix == "cblk" { dataStep = 2 } else if rolePrefix == "sblk" { dataStart = 2 dataStep = 2 } return &bulkRuntime{ rolePrefix: rolePrefix, dataStart: dataStart, dataStep: dataStep, bulks: make(map[string]*bulkHandle), inbound: make(map[string]map[uint64]*bulkHandle), outbound: make(map[string]map[uint64]*bulkHandle), reserved: make(map[string]map[uint64]struct{}), } } func (r *bulkRuntime) nextID() string { if r == nil { return "" } return fmt.Sprintf("%s-%d", r.rolePrefix, r.seq.Add(1)) } func (r *bulkRuntime) nextDataID() uint64 { if r == nil { return 0 } r.mu.Lock() defer r.mu.Unlock() id, _ := r.nextDataIDLocked(defaultFileScope, nil) return id } // reserveDataID allocates a DataID before an open request is sent. The // reservation prevents another concurrent open from selecting the same ID // while the caller is still constructing and registering its bulk handle. func (r *bulkRuntime) reserveDataID(scope string, requested uint64) (uint64, error) { if r == nil { return 0, errBulkRuntimeNil } scope = normalizeFileScope(scope) r.mu.Lock() defer r.mu.Unlock() return r.nextDataIDLocked(scope, &requested) } func (r *bulkRuntime) releaseDataID(scope string, dataID uint64) { if r == nil || dataID == 0 { return } scope = normalizeFileScope(scope) r.mu.Lock() defer r.mu.Unlock() if reserved := r.reserved[scope]; reserved != nil { delete(reserved, dataID) if len(reserved) == 0 { delete(r.reserved, scope) } } } // nextDataIDLocked chooses an ID under r.mu. requested is nil for an // internal auto allocation and points to zero/non-zero for a caller that // wants a reservation for an outbound open. func (r *bulkRuntime) nextDataIDLocked(scope string, requested *uint64) (uint64, error) { if r == nil { return 0, errBulkRuntimeNil } var wanted uint64 if requested != nil { wanted = *requested } dataScope := r.outbound[scope] inboundScope := r.inbound[scope] reserved := r.reserved[scope] if wanted != 0 { if !r.localDataID(wanted) { return 0, errBulkDataIDEmpty } if dataScope != nil { if _, exists := dataScope[wanted]; exists { return 0, errBulkAlreadyExists } } if inboundScope != nil { if _, exists := inboundScope[wanted]; exists { return 0, errBulkAlreadyExists } } if _, exists := reserved[wanted]; exists { return 0, errBulkAlreadyExists } if wanted > r.dataSeq { r.dataSeq = wanted } if reserved == nil { reserved = make(map[uint64]struct{}) r.reserved[scope] = reserved } reserved[wanted] = struct{}{} return wanted, nil } for { candidate, ok := r.nextDataCandidateLocked() if !ok { return 0, errBulkDataIDExhausted } if dataScope != nil { if _, exists := dataScope[candidate]; exists { continue } } if inboundScope != nil { if _, exists := inboundScope[candidate]; exists { continue } } if _, exists := reserved[candidate]; exists { continue } if requested != nil { if reserved == nil { reserved = make(map[uint64]struct{}) r.reserved[scope] = reserved } reserved[candidate] = struct{}{} } return candidate, nil } } func (r *bulkRuntime) nextDataCandidateLocked() (uint64, bool) { if r == nil { return 0, false } step := r.dataStep if step == 0 { step = 1 } if r.dataSeq == 0 { candidate := r.dataStart if candidate == 0 { return 0, false } r.dataSeq = candidate return candidate, true } if r.dataSeq > ^uint64(0)-step { return 0, false } candidate := r.dataSeq + step if step == 2 && candidate%2 != r.dataStart%2 { if candidate == ^uint64(0) { return 0, false } candidate++ } r.dataSeq = candidate return candidate, true } func (r *bulkRuntime) localDataID(dataID uint64) bool { if r == nil || dataID == 0 { return false } if r.dataStep != 2 { return true } return dataID%2 == r.dataStart%2 } func (r *bulkRuntime) setHandler(fn func(BulkAcceptInfo) error) { if r == nil { return } r.mu.Lock() defer r.mu.Unlock() r.handler = fn } func (r *bulkRuntime) handlerSnapshot() func(BulkAcceptInfo) error { if r == nil { return nil } r.mu.RLock() defer r.mu.RUnlock() return r.handler } func (r *bulkRuntime) register(scope string, bulk *bulkHandle) error { // Keep the old helper usable by package-local callers and tests. Production // paths use registerInbound/registerOutbound so a DataID can exist once in // each direction without making inbound frame dispatch ambiguous. return r.registerWithDirections(scope, bulk, bulkDataIndexBoth, false) } func (r *bulkRuntime) registerInbound(scope string, bulk *bulkHandle) error { return r.registerWithDirections(scope, bulk, bulkDataIndexInbound, false) } // adoptInbound transfers ownership of a newly-created handle to the runtime. // A failed registration is terminal because the handle has already started // its background workers and must not be reused by the caller. func (r *bulkRuntime) adoptInbound(scope string, bulk *bulkHandle) error { err := r.registerInbound(scope, bulk) if err != nil && bulk != nil { bulk.markReset(err) } return err } func (r *bulkRuntime) registerOutbound(scope string, bulk *bulkHandle) error { return r.registerWithDirections(scope, bulk, bulkDataIndexOutbound, false) } func (r *bulkRuntime) registerReserved(scope string, bulk *bulkHandle) error { return r.registerWithDirections(scope, bulk, bulkDataIndexOutbound, true) } // adoptReserved is the outbound counterpart of adoptInbound. The caller must // still release an unconsumed DataID reservation when this method fails. func (r *bulkRuntime) adoptReserved(scope string, bulk *bulkHandle) error { err := r.registerReserved(scope, bulk) if err != nil && bulk != nil { bulk.markReset(err) } return err } func (r *bulkRuntime) registerWithDirections(scope string, bulk *bulkHandle, direction bulkDataIndexDirection, consumeReservation bool) error { if r == nil { return errBulkRuntimeNil } if bulk == nil || bulk.id == "" { return errBulkIDEmpty } scope = normalizeFileScope(scope) key := bulkRuntimeKey(scope, bulk.id) r.mu.Lock() defer r.mu.Unlock() if _, ok := r.bulks[key]; ok { return errBulkAlreadyExists } if direction == 0 { return errBulkDataIDEmpty } if bulk.dataID == 0 { dataID, err := r.nextDataIDLocked(scope, nil) if err != nil { return err } bulk.dataID = dataID } else if direction&bulkDataIndexOutbound != 0 && r.localDataID(bulk.dataID) && bulk.dataID > r.dataSeq { r.dataSeq = bulk.dataID } inbound := r.inbound[scope] outbound := r.outbound[scope] if direction&bulkDataIndexInbound != 0 && inbound != nil { if _, ok := inbound[bulk.dataID]; ok { return errBulkAlreadyExists } } if direction&bulkDataIndexOutbound != 0 && outbound != nil { if _, ok := outbound[bulk.dataID]; ok { return errBulkAlreadyExists } } // New peers use disjoint odd/even namespaces. Reject a legacy peer's // colliding explicit ID when both directions would otherwise share one // wire DataID; the inbound-frame router cannot disambiguate that case. if r.dataStep == 2 { if direction&bulkDataIndexInbound != 0 && outbound != nil { if _, ok := outbound[bulk.dataID]; ok { return errBulkAlreadyExists } } if direction&bulkDataIndexOutbound != 0 && inbound != nil { if _, ok := inbound[bulk.dataID]; ok { return errBulkAlreadyExists } } } if direction&bulkDataIndexOutbound != 0 { if reserved := r.reserved[scope]; reserved != nil { if _, exists := reserved[bulk.dataID]; exists { if !consumeReservation { return errBulkAlreadyExists } delete(reserved, bulk.dataID) if len(reserved) == 0 { delete(r.reserved, scope) } } } } r.bulks[key] = bulk if direction&bulkDataIndexInbound != 0 { if inbound == nil { inbound = make(map[uint64]*bulkHandle) r.inbound[scope] = inbound } inbound[bulk.dataID] = bulk } if direction&bulkDataIndexOutbound != 0 { if outbound == nil { outbound = make(map[uint64]*bulkHandle) r.outbound[scope] = outbound } outbound[bulk.dataID] = bulk } return nil } func (r *bulkRuntime) lookup(scope string, bulkID string) (*bulkHandle, bool) { if r == nil || bulkID == "" { return nil, false } key := bulkRuntimeKey(scope, bulkID) r.mu.RLock() defer r.mu.RUnlock() bulk, ok := r.bulks[key] return bulk, ok } func (r *bulkRuntime) lookupByDataID(scope string, dataID uint64) (*bulkHandle, bool) { return r.lookupByDataIDDirection(scope, dataID, bulkDataIndexBoth) } func (r *bulkRuntime) lookupInboundByDataID(scope string, dataID uint64) (*bulkHandle, bool) { return r.lookupByDataIDDirection(scope, dataID, bulkDataIndexInbound) } // lookupInboundFrame chooses the local-open or peer-open index from the // allocator partition. Both bulk kinds are duplex; the partition identifies // which handle owns a wire DataID before consulting the corresponding map. func (r *bulkRuntime) lookupInboundFrame(scope string, dataID uint64) (*bulkHandle, bool) { if r == nil { return nil, false } if r.dataStep != 2 { return r.lookupByDataID(scope, dataID) } if r.localDataID(dataID) { if bulk, ok := r.lookupOutboundByDataID(scope, dataID); ok { return bulk, true } // Legacy peers may send DataID=0 in the open request. The receiver // allocates an ID locally in that case, so accept the inbound index // when the preferred outbound slot is absent. return r.lookupInboundByDataID(scope, dataID) } if bulk, ok := r.lookupInboundByDataID(scope, dataID); ok { return bulk, true } // Keep old peers that selected the local namespace routable when no // inbound handle occupies the ID. return r.lookupOutboundByDataID(scope, dataID) } func (r *bulkRuntime) lookupOutboundByDataID(scope string, dataID uint64) (*bulkHandle, bool) { return r.lookupByDataIDDirection(scope, dataID, bulkDataIndexOutbound) } func (r *bulkRuntime) lookupByDataIDDirection(scope string, dataID uint64, direction bulkDataIndexDirection) (*bulkHandle, bool) { if r == nil || dataID == 0 { return nil, false } scope = normalizeFileScope(scope) r.mu.RLock() defer r.mu.RUnlock() var inbound, outbound *bulkHandle if direction&bulkDataIndexInbound != 0 { if dataScope := r.inbound[scope]; dataScope != nil { inbound = dataScope[dataID] } } if direction&bulkDataIndexOutbound != 0 { if dataScope := r.outbound[scope]; dataScope != nil { outbound = dataScope[dataID] } } if direction == bulkDataIndexBoth && inbound != nil && outbound != nil && inbound != outbound { return nil, false } if inbound != nil { return inbound, true } if outbound != nil { return outbound, true } return nil, false } // lookupControl resolves the identity carried by a control message. A // supplied BulkID is authoritative and, when present, DataID must agree. A // DataID-only message is accepted only when it maps to one direction; using a // colliding ID from the other direction would otherwise reset the wrong bulk. func (r *bulkRuntime) lookupControl(scope string, bulkID string, dataID uint64) (*bulkHandle, bool) { if r == nil { return nil, false } scope = normalizeFileScope(scope) r.mu.RLock() defer r.mu.RUnlock() if bulkID != "" { bulk, ok := r.bulks[bulkRuntimeKey(scope, bulkID)] if !ok || bulk == nil { return nil, false } if dataID != 0 && bulk.dataID != dataID { return nil, false } return bulk, true } if dataID == 0 { return nil, false } var inbound, outbound *bulkHandle if dataScope := r.inbound[scope]; dataScope != nil { inbound = dataScope[dataID] } if dataScope := r.outbound[scope]; dataScope != nil { outbound = dataScope[dataID] } if inbound != nil && outbound != nil && inbound != outbound { return nil, false } if inbound != nil { return inbound, true } return outbound, outbound != nil } func (r *bulkRuntime) remove(scope string, expected *bulkHandle) { if r == nil || expected == nil || expected.id == "" { return } scope = normalizeFileScope(scope) key := bulkRuntimeKey(scope, expected.id) r.mu.Lock() defer r.mu.Unlock() bulk := r.bulks[key] if bulk != expected { return } if bulk.dataID != 0 { if dataScope := r.inbound[scope]; dataScope != nil { if dataScope[bulk.dataID] == bulk { delete(dataScope, bulk.dataID) } if len(dataScope) == 0 { delete(r.inbound, scope) } } if dataScope := r.outbound[scope]; dataScope != nil { if dataScope[bulk.dataID] == bulk { delete(dataScope, bulk.dataID) } if len(dataScope) == 0 { delete(r.outbound, scope) } } } delete(r.bulks, key) } func (r *bulkRuntime) closeAll(err error) { r.closeMatching(func(string) bool { return true }, err) } func (r *bulkRuntime) closeScope(scope string, err error) { scope = normalizeFileScope(scope) r.closeMatching(func(key string) bool { return strings.HasPrefix(key, scope+"\x00") }, err) } func (r *bulkRuntime) closeClientRoute(route clientSessionRoute, err error) { if r == nil { return } if !r.mu.TryRLock() { go r.closeClientRouteBlocking(route, err) return } bulks := r.collectClientRouteLocked(route) r.mu.RUnlock() r.resetClientRouteHandles(bulks, err) } func (r *bulkRuntime) closeClientRouteBlocking(route clientSessionRoute, err error) { r.mu.RLock() bulks := r.collectClientRouteLocked(route) r.mu.RUnlock() r.resetClientRouteHandles(bulks, err) } func (r *bulkRuntime) collectClientRouteLocked(route clientSessionRoute) []*bulkHandle { bulks := make([]*bulkHandle, 0) for _, bulk := range r.bulks { if bulk == nil || !sameClientSessionRoute(bulk.clientRoute, route) { continue } bulks = append(bulks, bulk) } return bulks } func (r *bulkRuntime) resetClientRouteHandles(bulks []*bulkHandle, err error) { resetErr := bulkRuntimeCloseError(err) for _, bulk := range bulks { bulk.markReset(resetErr) } } func (r *bulkRuntime) closeMatching(match func(string) bool, err error) { if r == nil || match == nil { return } resetErr := bulkRuntimeCloseError(err) r.mu.RLock() bulks := make([]*bulkHandle, 0, len(r.bulks)) for key, bulk := range r.bulks { if bulk == nil || !match(key) { continue } bulks = append(bulks, bulk) } r.mu.RUnlock() for _, bulk := range bulks { bulk.markReset(resetErr) } } func (r *bulkRuntime) resetDedicatedByConn(scope string, conn net.Conn, err error) { if r == nil || conn == nil { return } scope = normalizeFileScope(scope) resetErr := bulkRuntimeCloseError(err) prefix := scope + "\x00" r.mu.RLock() bulks := make([]*bulkHandle, 0, len(r.bulks)) for key, bulk := range r.bulks { if bulk == nil { continue } if !strings.HasPrefix(key, prefix) { continue } if !bulk.Dedicated() { continue } if bulk.dedicatedConnSnapshot() != conn { continue } bulks = append(bulks, bulk) } r.mu.RUnlock() for _, bulk := range bulks { bulk.markReset(resetErr) } } func (r *bulkRuntime) handleDedicatedReadErrorByConn(scope string, conn net.Conn, err error) { if r == nil || conn == nil { return } scope = normalizeFileScope(scope) prefix := scope + "\x00" r.mu.RLock() bulks := make([]*bulkHandle, 0, len(r.bulks)) for key, bulk := range r.bulks { if bulk == nil { continue } if !strings.HasPrefix(key, prefix) { continue } if !bulk.Dedicated() { continue } if bulk.dedicatedConnSnapshot() != conn { continue } bulks = append(bulks, bulk) } r.mu.RUnlock() for _, bulk := range bulks { handleDedicatedBulkReadError(bulk, err) } } func (r *bulkRuntime) attachSharedDedicatedConn(scope string, laneID uint32, conn net.Conn) { if r == nil || conn == nil { return } scope = normalizeFileScope(scope) laneID = normalizeBulkDedicatedLaneID(laneID) prefix := scope + "\x00" r.mu.RLock() bulks := make([]*bulkHandle, 0, len(r.bulks)) for key, bulk := range r.bulks { if bulk == nil { continue } if !strings.HasPrefix(key, prefix) { continue } if !bulk.Dedicated() { continue } if bulk.dedicatedLaneIDSnapshot() != laneID { continue } bulks = append(bulks, bulk) } r.mu.RUnlock() for _, bulk := range bulks { current := bulk.dedicatedConnSnapshot() switch { case current == conn: _ = bulk.attachDedicatedConnShared(conn) case current == nil: _ = bulk.attachDedicatedConnShared(conn) default: oldConn, oldSender, err := bulk.replaceDedicatedConnShared(conn) if err != nil { bulk.markReset(err) continue } if oldSender != nil { oldSender.stop() } if oldConn != nil { _ = oldConn.Close() } } } } func (r *bulkRuntime) dedicatedBulksForConn(scope string, conn net.Conn) []*bulkHandle { if r == nil || conn == nil { return nil } scope = normalizeFileScope(scope) prefix := scope + "\x00" r.mu.RLock() bulks := make([]*bulkHandle, 0, len(r.bulks)) for key, bulk := range r.bulks { if bulk == nil { continue } if !strings.HasPrefix(key, prefix) { continue } if !bulk.Dedicated() { continue } if bulk.dedicatedConnSnapshot() != conn { continue } bulks = append(bulks, bulk) } r.mu.RUnlock() return bulks } func (r *bulkRuntime) snapshots() []BulkSnapshot { if r == nil { return nil } r.mu.RLock() snapshots := make([]BulkSnapshot, 0, len(r.bulks)) for _, bulk := range r.bulks { if bulk == nil { continue } snapshots = append(snapshots, bulk.snapshot()) } r.mu.RUnlock() sortBulkSnapshots(snapshots) return snapshots } func bulkRuntimeKey(scope string, bulkID string) string { return normalizeFileScope(scope) + "\x00" + bulkID } func bulkRuntimeCloseError(err error) error { if err != nil { return err } return errServiceShutdown } func (c *ClientCommon) getBulkRuntime() *bulkRuntime { if c == nil { return nil } return c.bulkRuntime } func (s *ServerCommon) getBulkRuntime() *bulkRuntime { if s == nil { return nil } return s.bulkRuntime }