package notify import ( "context" "errors" "time" ) const bulkOpenRecoveryTimeout = 2 * time.Second func (c *ClientCommon) SetBulkHandler(fn func(BulkAcceptInfo) error) { runtime := c.getBulkRuntime() if runtime == nil { return } runtime.setHandler(fn) } func (c *ClientCommon) OpenSharedBulk(ctx context.Context, opt BulkOpenOptions) (Bulk, error) { opt.Mode = BulkOpenModeShared opt.Dedicated = false return c.OpenBulk(ctx, opt) } func (c *ClientCommon) OpenDedicatedBulk(ctx context.Context, opt BulkOpenOptions) (Bulk, error) { opt.Mode = BulkOpenModeDedicated opt.Dedicated = true return c.OpenBulk(ctx, opt) } func (c *ClientCommon) OpenBulk(ctx context.Context, opt BulkOpenOptions) (Bulk, error) { if normalizeBulkOpenMode(opt.Mode) == BulkOpenModeDefault && !opt.Dedicated { opt.Mode = c.bulkDefaultOpenModeSnapshot() } opt = normalizeBulkOpenOptions(opt) switch opt.Mode { case BulkOpenModeDedicated: opt.Dedicated = true return c.openBulkWithDedicatedMode(ctx, opt, false) case BulkOpenModeAuto: // Auto mode prefers dedicated path and falls back to shared if dedicated fails. if err := clientDedicatedBulkSupportError(c); err == nil { dedicatedOpt := opt dedicatedOpt.Mode = BulkOpenModeDedicated dedicatedOpt.Dedicated = true bulk, dedicatedErr := c.openBulkWithDedicatedMode(ctx, dedicatedOpt, true) if dedicatedErr == nil { return bulk, nil } sharedOpt := opt sharedOpt.Mode = BulkOpenModeShared sharedOpt.Dedicated = false sharedBulk, sharedErr := c.openBulkWithDedicatedMode(ctx, sharedOpt, false) if sharedErr == nil { c.bulkAttachFallbackCount.Add(1) return sharedBulk, nil } return nil, errors.Join(dedicatedErr, sharedErr) } opt.Mode = BulkOpenModeShared opt.Dedicated = false c.bulkAttachFallbackCount.Add(1) return c.openBulkWithDedicatedMode(ctx, opt, false) case BulkOpenModeShared, BulkOpenModeDefault: opt.Mode = BulkOpenModeShared opt.Dedicated = false return c.openBulkWithDedicatedMode(ctx, opt, false) default: opt.Mode = BulkOpenModeShared opt.Dedicated = false return c.openBulkWithDedicatedMode(ctx, opt, false) } } func (c *ClientCommon) openBulkWithDedicatedMode(ctx context.Context, opt BulkOpenOptions, waitForReset bool) (Bulk, error) { if c == nil { return nil, errBulkClientNil } route := c.clientSessionRouteSnapshot() if err := c.ensureClientSessionRouteSendReady(route); err != nil { return nil, err } opt = applyBulkOpenTuningDefaults(opt, c.bulkOpenTuningSnapshot()) runtime := c.getBulkRuntime() if runtime == nil { return nil, errBulkRuntimeNil } req := clientBulkRequest(runtime, opt) if req.BulkID == "" { return nil, errBulkIDEmpty } if req.Dedicated { if err := clientDedicatedBulkSupportError(c); err != nil { return nil, err } } if !validBulkRange(req.Range) { return nil, errBulkRangeInvalid } if existing, exists := runtime.lookup(clientFileScope(), req.BulkID); exists { if existing.acceptsClientSessionRoute(route) { return nil, errBulkAlreadyExists } existing.markReset(errTransportDetached) } if req.DataID == 0 { var reserveErr error req.DataID, reserveErr = runtime.reserveDataID(clientFileScope(), 0) if reserveErr != nil { return nil, reserveErr } } if req.Dedicated { var laneErr error req.DedicatedLaneID, laneErr = c.reserveBulkDedicatedLaneAtRoute(route) if laneErr != nil { runtime.releaseDataID(clientFileScope(), req.DataID) return nil, laneErr } if req.AttachToken == "" { req.AttachToken = newBulkAttachToken() } bulk := newBulkHandle(clientSessionRouteContext(route), runtime, clientFileScope(), req, route.epoch, nil, nil, 0, clientBulkCloseSender(c), clientBulkResetSender(c), clientBulkDataSender(c, route), clientBulkWriteSender(c, route), clientBulkReleaseSender(c)) bulk.setClientSnapshotOwner(c) bulk.setClientSessionRoute(route) bulk.markAcceptHandled() bulk.markDedicatedLaneReserved() if err := runtime.adoptReserved(clientFileScope(), bulk); err != nil { runtime.releaseDataID(clientFileScope(), req.DataID) return nil, err } resp, err := sendBulkOpenClientAtRoute(ctx, c, route, req) if err != nil { runtime.releaseDataID(clientFileScope(), req.DataID) cleanupErr := c.cleanupBulkResetAtRoute(ctx, route, BulkResetRequest{BulkID: req.BulkID, DataID: req.DataID, Error: err.Error()}, waitForReset) bulk.markReset(err) if cleanupErr != nil { return nil, errors.Join(err, cleanupErr) } return nil, err } if resp.DataID != 0 && resp.DataID != req.DataID { err = errBulkAlreadyExists cleanupErr := c.cleanupBulkResetAtRoute(ctx, route, BulkResetRequest{ BulkID: req.BulkID, Error: "bulk dedicated data id mismatch", }, waitForReset) bulk.markReset(err) if cleanupErr != nil { return nil, errors.Join(err, cleanupErr) } return nil, err } if resp.TransportGeneration != 0 { bulk.setTransportGeneration(resp.TransportGeneration) } if resp.FastPathVersion != 0 { bulk.setFastPathVersion(resp.FastPathVersion) } if resp.AttachToken != "" { req.AttachToken = resp.AttachToken bulk.setDedicatedAttachToken(resp.AttachToken) } if err := c.attachDedicatedBulkSidecar(ctx, bulk); err != nil { cleanupErr := c.cleanupBulkResetAtRoute(ctx, route, BulkResetRequest{ BulkID: req.BulkID, DataID: req.DataID, Error: err.Error(), }, waitForReset) bulk.markReset(err) if cleanupErr != nil { return nil, errors.Join(err, cleanupErr) } return nil, err } if err := bulk.waitAcceptReady(ctx); err != nil { var cleanupErr error if bulk.resetErrSnapshot() == nil { cleanupErr = c.cleanupBulkResetAtRoute(ctx, route, BulkResetRequest{ BulkID: req.BulkID, DataID: req.DataID, Error: err.Error(), }, waitForReset) } else { // A ready error already reset the remote handle. Keep the old // asynchronous cleanup for compatibility with that path and avoid // racing a concurrent dedicated attach teardown. c.bestEffortBulkResetAtRoute(route, BulkResetRequest{ BulkID: req.BulkID, DataID: req.DataID, Error: err.Error(), }) } bulk.markReset(err) if cleanupErr != nil { return nil, errors.Join(err, cleanupErr) } return nil, err } return bulk, nil } bulk := newBulkHandle(clientSessionRouteContext(route), runtime, clientFileScope(), req, route.epoch, nil, nil, 0, clientBulkCloseSender(c), clientBulkResetSender(c), clientBulkDataSender(c, route), clientBulkWriteSender(c, route), clientBulkReleaseSender(c)) bulk.setClientSnapshotOwner(c) bulk.setClientSessionRoute(route) bulk.markAcceptHandled() if err := runtime.adoptReserved(clientFileScope(), bulk); err != nil { runtime.releaseDataID(clientFileScope(), req.DataID) return nil, err } resp, err := sendBulkOpenClientAtRoute(ctx, c, route, req) if err != nil { c.bestEffortBulkResetAtRoute(route, BulkResetRequest{BulkID: req.BulkID, DataID: req.DataID, Error: err.Error()}) bulk.markReset(err) return nil, err } if resp.DataID != 0 && resp.DataID != req.DataID { err = errBulkAlreadyExists c.bestEffortBulkResetAtRoute(route, BulkResetRequest{BulkID: req.BulkID, Error: "bulk data id mismatch"}) bulk.markReset(err) return nil, err } if resp.FastPathVersion != 0 { bulk.setFastPathVersion(resp.FastPathVersion) } if resp.Dedicated { err = errBulkRejected c.bestEffortBulkResetAtRoute(route, BulkResetRequest{BulkID: req.BulkID, DataID: req.DataID, Error: "shared bulk upgraded to dedicated"}) bulk.markReset(err) return nil, err } if resp.AttachToken != "" { bulk.setDedicatedAttachToken(resp.AttachToken) } bulk.setTransportGeneration(resp.TransportGeneration) return bulk, nil } func (c *ClientCommon) bestEffortBulkReset(req BulkResetRequest) { if c == nil { return } c.bestEffortBulkResetAtRoute(c.clientSessionRouteSnapshot(), req) } func (c *ClientCommon) bestEffortBulkResetAtEpoch(epoch uint64, req BulkResetRequest) { route := c.clientSessionRouteSnapshot() route.epoch = epoch c.bestEffortBulkResetAtRoute(route, req) } func (c *ClientCommon) bestEffortBulkResetAtRoute(route clientSessionRoute, req BulkResetRequest) { if c == nil { return } task := newClientBulkResetRecoveryTaskAtRoute(c, route, req) q := c.bulkRecoveryQueue() if !q.enqueue(task) { c.handleBulkRecoveryOverflowAtRoute(route, req) } } func (c *ClientCommon) bulkRecoveryQueue() *bulkRecoveryQueue { if c == nil { return nil } c.bulkRecoveryMu.Lock() defer c.bulkRecoveryMu.Unlock() if c.bulkRecovery == nil { c.bulkRecovery = newBulkRecoveryQueue(c.reportBulkRecoveryError) } return c.bulkRecovery } func newClientBulkResetRecoveryTask(c *ClientCommon, epoch uint64, req BulkResetRequest) bulkRecoveryTask { route := c.clientSessionRouteSnapshot() route.epoch = epoch return newClientBulkResetRecoveryTaskAtRoute(c, route, req) } func newClientBulkResetRecoveryTaskAtRoute(c *ClientCommon, route clientSessionRoute, req BulkResetRequest) bulkRecoveryTask { return func(ctx context.Context) error { if !c.clientSessionRouteCurrent(route) { return transportDetachedSessionEpochError() } _, err := sendBulkResetClientAtRoute(ctx, c, route, req) if errors.Is(err, errBulkNotFound) { return nil } return err } } func clientBulkRequest(runtime *bulkRuntime, opt BulkOpenOptions) BulkOpenRequest { opt = normalizeBulkOpenOptions(opt) id := opt.ID if id == "" && runtime != nil { id = runtime.nextID() } return normalizeBulkOpenRequest(BulkOpenRequest{ BulkID: id, FastPathVersion: bulkFastPathVersionCurrent, Range: opt.Range, Metadata: cloneBulkMetadata(opt.Metadata), ReadTimeout: opt.ReadTimeout, WriteTimeout: opt.WriteTimeout, Dedicated: opt.Dedicated, ChunkSize: opt.ChunkSize, WindowBytes: opt.WindowBytes, MaxInFlight: opt.MaxInFlight, }) } func clientBulkCloseSender(c *ClientCommon) bulkCloseSender { return func(ctx context.Context, bulk *bulkHandle, full bool) error { if bulk != nil && bulk.Dedicated() { if err := bulk.waitDedicatedReady(ctx); err != nil { return err } return c.sendDedicatedBulkClose(ctx, bulk, full) } _, err := sendBulkCloseClientAtRoute(ctx, c, bulk.clientSessionRouteSnapshot(), BulkCloseRequest{ BulkID: bulk.ID(), DataID: bulk.dataIDSnapshot(), Full: full, }) return err } } func clientBulkResetSender(c *ClientCommon) bulkResetSender { return func(ctx context.Context, bulk *bulkHandle, message string) error { if bulk != nil && bulk.Dedicated() { if err := bulk.waitDedicatedReady(ctx); err != nil { return err } return c.sendDedicatedBulkReset(ctx, bulk, message) } _, err := sendBulkResetClientAtRoute(ctx, c, bulk.clientSessionRouteSnapshot(), BulkResetRequest{ BulkID: bulk.ID(), DataID: bulk.dataIDSnapshot(), Error: message, }) return err } } func clientBulkDataSender(c *ClientCommon, route clientSessionRoute) bulkDataSender { return func(ctx context.Context, bulk *bulkHandle, chunk []byte) error { if c == nil { return errBulkClientNil } if ctx != nil { select { case <-ctx.Done(): return ctx.Err() default: } } if bulk != nil && bulk.Dedicated() { if err := bulk.waitDedicatedReady(ctx); err != nil { return err } return c.sendDedicatedBulkData(ctx, bulk, chunk) } if !c.clientSessionRouteCurrent(route) { return errTransportDetached } dataID := bulk.dataIDSnapshot() if dataID == 0 { return errBulkDataPathNotReady } return c.sendFastBulkDataAtRoute(ctx, route, dataID, bulk.nextOutboundDataSeq(), chunk, bulk.fastPathVersionSnapshot()) } } func clientBulkWriteSender(c *ClientCommon, route clientSessionRoute) bulkWriteSender { return func(ctx context.Context, bulk *bulkHandle, startSeq uint64, payload []byte, payloadOwned bool) (int, error) { if c == nil { return 0, errBulkClientNil } if ctx != nil { select { case <-ctx.Done(): return 0, ctx.Err() default: } } if bulk != nil && bulk.Dedicated() { if err := bulk.waitDedicatedReady(ctx); err != nil { return 0, err } return c.sendDedicatedBulkWrite(ctx, bulk, startSeq, payload, payloadOwned) } if !c.clientSessionRouteCurrent(route) { return 0, errTransportDetached } if bulk == nil { return 0, errBulkRuntimeNil } dataID := bulk.dataIDSnapshot() if dataID == 0 { return 0, errBulkDataPathNotReady } return c.sendFastBulkWriteAtRoute(ctx, route, dataID, startSeq, bulk.chunkSize, bulk.fastPathVersionSnapshot(), payload, payloadOwned) } } func clientBulkReleaseSender(c *ClientCommon) bulkReleaseSender { fallbackRoute := c.clientSessionRouteSnapshot() return func(bulk *bulkHandle, bytes int64, chunks int) error { if c == nil || bulk == nil { return errBulkClientNil } if bytes <= 0 && chunks <= 0 { return nil } ctx, cancel, err := bulk.newWriteContext(bulk.Context(), bulk.writeTimeout) if err != nil { return err } defer cancel() if bulk.Dedicated() { return c.sendDedicatedBulkRelease(ctx, bulk, bytes, chunks) } route := bulk.clientSessionRouteSnapshot() if !route.bound() { route = fallbackRoute } if bulk.fastPathVersionSnapshot() >= bulkFastPathVersionV2 { payload, err := encodeBulkDedicatedReleasePayload(bytes, chunks) if err != nil { return err } return c.sendFastBulkControlAtRoute(ctx, route, bulkFastPayloadTypeRelease, 0, bulk.dataIDSnapshot(), 0, bulk.fastPathVersionSnapshot(), payload) } return sendBulkReleaseClientAtRoute(ctx, c, route, BulkReleaseRequest{ BulkID: bulk.ID(), DataID: bulk.dataIDSnapshot(), Bytes: bytes, Chunks: chunks, }) } }