package notify import ( "context" "net" ) func (c *ClientCommon) SetStreamHandler(fn func(StreamAcceptInfo) error) { runtime := c.getStreamRuntime() if runtime == nil { return } runtime.setHandler(fn) } func (c *ClientCommon) OpenStream(ctx context.Context, opt StreamOpenOptions) (Stream, error) { if c == nil { return nil, errStreamClientNil } runtime := c.getStreamRuntime() if runtime == nil { return nil, errStreamRuntimeNil } route := c.clientSessionRouteSnapshot() if err := c.ensureClientSessionRouteSendReady(route); err != nil { return nil, err } scope := clientFileScope() req := clientStreamRequest(runtime, opt) if req.StreamID == "" { return nil, errStreamIDEmpty } if existing, exists := runtime.lookup(scope, req.StreamID); exists { if existing.acceptsClientSessionRoute(route) { return nil, errStreamAlreadyExists } existing.markReset(transportDetachedSessionEpochError()) } dataID, err := runtime.reserveDataID(scope) if err != nil { return nil, err } req.DataID = dataID parent := clientSessionRouteContext(route) if parent == nil { parent = c.clientStopContextSnapshot() } stream := newStreamHandle(parent, runtime, scope, req, route.epoch, nil, nil, 0, clientStreamCloseSender(c), clientStreamResetSender(c), clientStreamDataSender(c, route), runtime.configSnapshot()) stream.setClientSnapshotOwner(c) stream.setClientSessionRoute(route) stream.setAddrSnapshot(c.clientStreamAddrSnapshotAtRoute(route)) if err := runtime.adoptReserved(scope, stream); err != nil { runtime.releaseDataID(scope, req.DataID) return nil, err } resp, err := sendStreamOpenClientAtRoute(ctx, c, route, req) if err != nil { c.bestEffortStreamResetAtRoute(route, StreamResetRequest{StreamID: req.StreamID, DataID: req.DataID, Error: err.Error()}) stream.markReset(err) return nil, err } if resp.DataID != 0 && resp.DataID != req.DataID { err = errStreamAlreadyExists c.bestEffortStreamResetAtRoute(route, StreamResetRequest{StreamID: req.StreamID, Error: "stream data id mismatch"}) stream.markReset(err) return nil, err } if resp.FastPathVersion != 0 { stream.setFastPathVersion(resp.FastPathVersion) } else { stream.setFastPathVersion(streamFastPathVersionV1) } stream.metadata = mergeStreamMetadata(req.Metadata, resp.Metadata) stream.setTransportGeneration(resp.TransportGeneration) return stream, nil } func (c *ClientCommon) clientStreamAddrSnapshot() (net.Addr, net.Addr) { return c.clientStreamAddrSnapshotAtRoute(c.clientSessionRouteSnapshot()) } func (c *ClientCommon) clientStreamAddrSnapshotAtRoute(route clientSessionRoute) (net.Addr, net.Addr) { if c == nil { return nil, nil } var conn net.Conn if route.binding != nil { conn = route.binding.connSnapshot() } if conn == nil { return nil, nil } return conn.LocalAddr(), conn.RemoteAddr() } func clientStreamRequest(runtime *streamRuntime, opt StreamOpenOptions) StreamOpenRequest { id := opt.ID if id == "" && runtime != nil { id = runtime.nextID() } return normalizeStreamOpenRequest(StreamOpenRequest{ StreamID: id, FastPathVersion: streamFastPathVersionCurrent, Channel: opt.Channel, Metadata: cloneStreamMetadata(opt.Metadata), ReadTimeout: opt.ReadTimeout, WriteTimeout: opt.WriteTimeout, }) } func clientStreamCloseSender(c *ClientCommon) streamCloseSender { return func(ctx context.Context, stream *streamHandle, full bool) error { _, err := sendStreamCloseClientAtRoute(ctx, c, stream.clientSessionRouteSnapshot(), StreamCloseRequest{ StreamID: stream.ID(), DataID: stream.dataIDSnapshot(), Full: full, }) return err } } func clientStreamResetSender(c *ClientCommon) streamResetSender { return func(ctx context.Context, stream *streamHandle, message string) error { _, err := sendStreamResetClientAtRoute(ctx, c, stream.clientSessionRouteSnapshot(), StreamResetRequest{ StreamID: stream.ID(), DataID: stream.dataIDSnapshot(), Error: message, RecordFailure: stream.recordResetFailure(), }) return err } } func clientStreamDataSender(c *ClientCommon, route clientSessionRoute) streamDataSender { return func(ctx context.Context, stream *streamHandle, chunk []byte) error { if c == nil { return errStreamClientNil } if err := c.ensureClientSessionRouteSendReady(route); err != nil { return err } if ctx != nil { select { case <-ctx.Done(): return ctx.Err() default: } } if dataID := stream.dataIDSnapshot(); dataID != 0 { return c.sendFastStreamDataAtRoute(ctx, route, stream, chunk) } return c.sendEnvelopeAtRoute(route, newStreamDataEnvelope(stream.ID(), chunk)) } } func (c *ClientCommon) bestEffortStreamResetAtRoute(route clientSessionRoute, req StreamResetRequest) { if c == nil { return } ctx, cancel := context.WithTimeout(context.Background(), streamDispatchRejectTimeout) defer cancel() _, _ = sendStreamResetClientAtRoute(ctx, c, route, req) }