2026-04-15 15:24:36 +08:00
|
|
|
package notify
|
|
|
|
|
|
|
|
|
|
import (
|
|
|
|
|
"context"
|
|
|
|
|
"errors"
|
2026-09-23 15:33:17 +08:00
|
|
|
"strings"
|
2026-04-15 15:24:36 +08:00
|
|
|
"time"
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
type StreamOpenRequest struct {
|
2026-04-18 16:05:57 +08:00
|
|
|
StreamID string
|
|
|
|
|
DataID uint64
|
|
|
|
|
FastPathVersion uint8
|
|
|
|
|
Channel StreamChannel
|
|
|
|
|
Metadata StreamMetadata
|
|
|
|
|
ReadTimeout time.Duration
|
|
|
|
|
WriteTimeout time.Duration
|
2026-04-15 15:24:36 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
type StreamOpenResponse struct {
|
|
|
|
|
StreamID string
|
|
|
|
|
DataID uint64
|
2026-04-18 16:05:57 +08:00
|
|
|
FastPathVersion uint8
|
2026-04-15 15:24:36 +08:00
|
|
|
Accepted bool
|
|
|
|
|
TransportGeneration uint64
|
2026-04-15 19:52:45 +08:00
|
|
|
Metadata StreamMetadata
|
2026-04-15 15:24:36 +08:00
|
|
|
Error string
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
type StreamCloseRequest struct {
|
|
|
|
|
StreamID string
|
2026-09-23 15:33:17 +08:00
|
|
|
DataID uint64
|
2026-04-15 15:24:36 +08:00
|
|
|
Full bool
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
type StreamCloseResponse struct {
|
|
|
|
|
StreamID string
|
|
|
|
|
Accepted bool
|
|
|
|
|
Error string
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
type StreamResetRequest struct {
|
2026-09-23 15:33:17 +08:00
|
|
|
StreamID string
|
|
|
|
|
DataID uint64
|
|
|
|
|
Error string
|
|
|
|
|
RecordFailure *RecordFailure
|
2026-04-15 15:24:36 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
type StreamResetResponse struct {
|
|
|
|
|
StreamID string
|
|
|
|
|
Accepted bool
|
|
|
|
|
Error string
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func bindClientStreamControl(c *ClientCommon) {
|
|
|
|
|
if c == nil {
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
c.SetLink(StreamOpenSignalKey, func(msg *Message) {
|
|
|
|
|
c.handleInboundStreamOpen(msg)
|
|
|
|
|
})
|
|
|
|
|
c.SetLink(StreamCloseSignalKey, func(msg *Message) {
|
|
|
|
|
c.handleInboundStreamClose(msg)
|
|
|
|
|
})
|
|
|
|
|
c.SetLink(StreamResetSignalKey, func(msg *Message) {
|
|
|
|
|
c.handleInboundStreamReset(msg)
|
|
|
|
|
})
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func bindServerStreamControl(s *ServerCommon) {
|
|
|
|
|
if s == nil {
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
s.SetLink(StreamOpenSignalKey, func(msg *Message) {
|
|
|
|
|
s.handleInboundStreamOpen(msg)
|
|
|
|
|
})
|
|
|
|
|
s.SetLink(StreamCloseSignalKey, func(msg *Message) {
|
|
|
|
|
s.handleInboundStreamClose(msg)
|
|
|
|
|
})
|
|
|
|
|
s.SetLink(StreamResetSignalKey, func(msg *Message) {
|
|
|
|
|
s.handleInboundStreamReset(msg)
|
|
|
|
|
})
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (c *ClientCommon) handleInboundStreamOpen(msg *Message) {
|
|
|
|
|
req, err := decodeStreamOpenRequest(msg)
|
|
|
|
|
resp := StreamOpenResponse{StreamID: req.StreamID, DataID: req.DataID}
|
|
|
|
|
if err != nil {
|
|
|
|
|
resp.Error = err.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-09-23 15:33:17 +08:00
|
|
|
route := msg.clientRoute
|
|
|
|
|
if !route.bound() {
|
|
|
|
|
route = c.clientSessionRouteSnapshot()
|
|
|
|
|
}
|
|
|
|
|
if err := c.ensureClientSessionRouteSendReady(route); err != nil {
|
|
|
|
|
resp.Error = err.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-04-15 15:24:36 +08:00
|
|
|
runtime := c.getStreamRuntime()
|
|
|
|
|
if runtime == nil {
|
|
|
|
|
resp.Error = errStreamRuntimeNil.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
scope := clientFileScope()
|
2026-09-23 15:33:17 +08:00
|
|
|
if existing, ok := runtime.lookup(scope, req.StreamID); ok && !existing.acceptsClientSessionRoute(route) {
|
|
|
|
|
existing.markReset(transportDetachedSessionEpochError())
|
|
|
|
|
}
|
2026-04-18 16:05:57 +08:00
|
|
|
req.FastPathVersion = negotiateStreamFastPathVersion(req.FastPathVersion)
|
|
|
|
|
resp.FastPathVersion = req.FastPathVersion
|
2026-04-15 19:52:45 +08:00
|
|
|
req.Metadata, resp.Metadata = negotiateRecordStreamOpenMetadata(req.Channel, req.Metadata)
|
2026-09-23 15:33:17 +08:00
|
|
|
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())
|
2026-04-15 15:24:36 +08:00
|
|
|
stream.setClientSnapshotOwner(c)
|
2026-09-23 15:33:17 +08:00
|
|
|
stream.setClientSessionRoute(route)
|
|
|
|
|
stream.setAddrSnapshot(c.clientStreamAddrSnapshotAtRoute(route))
|
|
|
|
|
if err := runtime.adoptInbound(scope, stream); err != nil {
|
|
|
|
|
resp.Error = err.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
if err := c.ensureClientSessionRouteSendReady(route); err != nil {
|
|
|
|
|
runtime.remove(scope, stream)
|
|
|
|
|
stream.markReset(err)
|
|
|
|
|
resp.Error = err.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
if !stream.acceptDispatchAllowed() || !stream.claimAcceptDispatch() {
|
|
|
|
|
err := stream.resetErrSnapshot()
|
|
|
|
|
if err == nil {
|
|
|
|
|
err = transportDetachedSessionEpochError()
|
|
|
|
|
stream.markReset(err)
|
|
|
|
|
}
|
2026-04-15 15:24:36 +08:00
|
|
|
resp.Error = err.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
if claimed, err := c.claimInboundRecordStream(stream); claimed {
|
|
|
|
|
if err != nil {
|
|
|
|
|
stream.markReset(err)
|
|
|
|
|
resp.Error = err.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
resp.Accepted = true
|
2026-09-23 15:33:17 +08:00
|
|
|
resp.DataID = stream.dataIDSnapshot()
|
2026-04-15 15:24:36 +08:00
|
|
|
resp.TransportGeneration = stream.TransportGeneration()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
if claimed, err := c.claimInboundTransferStream(stream); claimed {
|
|
|
|
|
if err != nil {
|
|
|
|
|
stream.markReset(err)
|
|
|
|
|
resp.Error = err.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
resp.Accepted = true
|
2026-09-23 15:33:17 +08:00
|
|
|
resp.DataID = stream.dataIDSnapshot()
|
2026-04-15 15:24:36 +08:00
|
|
|
resp.TransportGeneration = stream.TransportGeneration()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
handler := runtime.handlerSnapshot()
|
|
|
|
|
if handler == nil {
|
|
|
|
|
stream.markReset(errStreamHandlerNotConfigured)
|
|
|
|
|
resp.Error = errStreamHandlerNotConfigured.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
info := StreamAcceptInfo{
|
|
|
|
|
ID: stream.ID(),
|
|
|
|
|
DataID: stream.dataIDSnapshot(),
|
|
|
|
|
Channel: stream.Channel(),
|
|
|
|
|
Metadata: stream.Metadata(),
|
|
|
|
|
TransportGeneration: stream.TransportGeneration(),
|
|
|
|
|
Stream: stream,
|
|
|
|
|
}
|
|
|
|
|
if err := handler(info); err != nil {
|
|
|
|
|
stream.markReset(err)
|
|
|
|
|
resp.Error = err.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
resp.Accepted = true
|
|
|
|
|
resp.DataID = stream.dataIDSnapshot()
|
|
|
|
|
resp.TransportGeneration = stream.TransportGeneration()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (s *ServerCommon) handleInboundStreamOpen(msg *Message) {
|
|
|
|
|
req, err := decodeStreamOpenRequest(msg)
|
|
|
|
|
resp := StreamOpenResponse{StreamID: req.StreamID, DataID: req.DataID}
|
|
|
|
|
if err != nil {
|
|
|
|
|
resp.Error = err.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
runtime := s.getStreamRuntime()
|
|
|
|
|
if runtime == nil {
|
|
|
|
|
resp.Error = errStreamRuntimeNil.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
logical := messageLogicalConnSnapshot(msg)
|
|
|
|
|
if logical == nil {
|
|
|
|
|
resp.Error = errStreamLogicalConnNil.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
transport := messageTransportConnSnapshot(msg)
|
2026-09-23 15:33:17 +08:00
|
|
|
if transport != nil && !transport.IsCurrent() {
|
|
|
|
|
resp.Error = transportDetachedErrorForTransport(transport).Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-04-15 15:24:36 +08:00
|
|
|
scope := serverFileScope(logical)
|
2026-09-23 15:33:17 +08:00
|
|
|
if existing, ok := runtime.lookup(scope, req.StreamID); ok && !existing.acceptsTransportGeneration(transport) {
|
|
|
|
|
existing.markReset(transportDetachedGenerationMismatchError(existing.TransportGeneration(), transport))
|
|
|
|
|
}
|
2026-04-18 16:05:57 +08:00
|
|
|
req.FastPathVersion = negotiateStreamFastPathVersion(req.FastPathVersion)
|
|
|
|
|
resp.FastPathVersion = req.FastPathVersion
|
2026-04-15 19:52:45 +08:00
|
|
|
req.Metadata, resp.Metadata = negotiateRecordStreamOpenMetadata(req.Channel, req.Metadata)
|
2026-04-15 15:24:36 +08:00
|
|
|
stream := newStreamHandle(logical.stopContextSnapshot(), runtime, scope, req, 0, logical, transport, streamTransportGeneration(logical, transport), serverStreamCloseSender(s, logical, transport), serverStreamResetSender(s, logical, transport), serverStreamDataSender(s, transport), runtime.configSnapshot())
|
2026-09-23 15:33:17 +08:00
|
|
|
if err := runtime.adoptInbound(scope, stream); err != nil {
|
|
|
|
|
resp.Error = err.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
if transport != nil && !transport.IsCurrent() {
|
|
|
|
|
runtime.remove(scope, stream)
|
|
|
|
|
stream.markReset(transportDetachedErrorForTransport(transport))
|
|
|
|
|
resp.Error = transportDetachedErrorForTransport(transport).Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
if !stream.acceptDispatchAllowed() || !stream.claimAcceptDispatch() {
|
|
|
|
|
err := stream.resetErrSnapshot()
|
|
|
|
|
if err == nil {
|
|
|
|
|
err = transportDetachedErrorForTransport(transport)
|
|
|
|
|
stream.markReset(err)
|
|
|
|
|
}
|
2026-04-15 15:24:36 +08:00
|
|
|
resp.Error = err.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
if claimed, err := s.claimInboundRecordStream(logical, transport, stream); claimed {
|
|
|
|
|
if err != nil {
|
|
|
|
|
stream.markReset(err)
|
|
|
|
|
resp.Error = err.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
resp.Accepted = true
|
2026-09-23 15:33:17 +08:00
|
|
|
resp.DataID = stream.dataIDSnapshot()
|
2026-04-15 15:24:36 +08:00
|
|
|
resp.TransportGeneration = stream.TransportGeneration()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
if claimed, err := s.claimInboundTransferStream(logical, transport, stream); claimed {
|
|
|
|
|
if err != nil {
|
|
|
|
|
stream.markReset(err)
|
|
|
|
|
resp.Error = err.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
resp.Accepted = true
|
2026-09-23 15:33:17 +08:00
|
|
|
resp.DataID = stream.dataIDSnapshot()
|
2026-04-15 15:24:36 +08:00
|
|
|
resp.TransportGeneration = stream.TransportGeneration()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
handler := runtime.handlerSnapshot()
|
|
|
|
|
if handler == nil {
|
|
|
|
|
stream.markReset(errStreamHandlerNotConfigured)
|
|
|
|
|
resp.Error = errStreamHandlerNotConfigured.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
info := StreamAcceptInfo{
|
|
|
|
|
ID: stream.ID(),
|
|
|
|
|
DataID: stream.dataIDSnapshot(),
|
|
|
|
|
Channel: stream.Channel(),
|
|
|
|
|
Metadata: stream.Metadata(),
|
|
|
|
|
LogicalConn: logical,
|
|
|
|
|
TransportConn: transport,
|
|
|
|
|
TransportGeneration: stream.TransportGeneration(),
|
|
|
|
|
Stream: stream,
|
|
|
|
|
}
|
|
|
|
|
if err := handler(info); err != nil {
|
|
|
|
|
stream.markReset(err)
|
|
|
|
|
resp.Error = err.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
resp.Accepted = true
|
|
|
|
|
resp.DataID = stream.dataIDSnapshot()
|
|
|
|
|
resp.TransportGeneration = stream.TransportGeneration()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (c *ClientCommon) handleInboundStreamClose(msg *Message) {
|
|
|
|
|
req, err := decodeStreamCloseRequest(msg)
|
|
|
|
|
resp := StreamCloseResponse{StreamID: req.StreamID}
|
|
|
|
|
if err != nil {
|
|
|
|
|
resp.Error = err.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
runtime := c.getStreamRuntime()
|
|
|
|
|
if runtime == nil {
|
|
|
|
|
resp.Error = errStreamRuntimeNil.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-09-23 15:33:17 +08:00
|
|
|
route := msg.clientRoute
|
|
|
|
|
if !route.bound() {
|
|
|
|
|
route = c.clientSessionRouteSnapshot()
|
|
|
|
|
}
|
|
|
|
|
stream, ok := runtime.lookupControl(clientFileScope(), req.StreamID, req.DataID)
|
2026-04-15 15:24:36 +08:00
|
|
|
if !ok {
|
|
|
|
|
resp.Error = errStreamNotFound.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-09-23 15:33:17 +08:00
|
|
|
if !stream.acceptsClientSessionRoute(route) {
|
|
|
|
|
resp.Error = transportDetachedSessionEpochError().Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-04-15 15:24:36 +08:00
|
|
|
if req.Full {
|
|
|
|
|
stream.markPeerClosed()
|
|
|
|
|
} else {
|
|
|
|
|
stream.markRemoteClosed()
|
|
|
|
|
}
|
|
|
|
|
resp.Accepted = true
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (s *ServerCommon) handleInboundStreamClose(msg *Message) {
|
|
|
|
|
req, err := decodeStreamCloseRequest(msg)
|
|
|
|
|
resp := StreamCloseResponse{StreamID: req.StreamID}
|
|
|
|
|
if err != nil {
|
|
|
|
|
resp.Error = err.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
runtime := s.getStreamRuntime()
|
|
|
|
|
if runtime == nil {
|
|
|
|
|
resp.Error = errStreamRuntimeNil.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
logical := messageLogicalConnSnapshot(msg)
|
|
|
|
|
scope := serverFileScope(logical)
|
2026-09-23 15:33:17 +08:00
|
|
|
stream, ok := runtime.lookupControl(scope, req.StreamID, req.DataID)
|
2026-04-15 15:24:36 +08:00
|
|
|
if !ok {
|
|
|
|
|
resp.Error = errStreamNotFound.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-09-23 15:33:17 +08:00
|
|
|
if !stream.acceptsTransportGeneration(messageTransportConnSnapshot(msg)) {
|
|
|
|
|
resp.Error = transportDetachedGenerationMismatchError(stream.TransportGeneration(), messageTransportConnSnapshot(msg)).Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-04-15 15:24:36 +08:00
|
|
|
if req.Full {
|
|
|
|
|
stream.markPeerClosed()
|
|
|
|
|
} else {
|
|
|
|
|
stream.markRemoteClosed()
|
|
|
|
|
}
|
|
|
|
|
resp.Accepted = true
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (c *ClientCommon) handleInboundStreamReset(msg *Message) {
|
|
|
|
|
req, err := decodeStreamResetRequest(msg)
|
|
|
|
|
resp := StreamResetResponse{StreamID: req.StreamID}
|
|
|
|
|
if err != nil {
|
|
|
|
|
resp.Error = err.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
runtime := c.getStreamRuntime()
|
|
|
|
|
if runtime == nil {
|
|
|
|
|
resp.Error = errStreamRuntimeNil.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-09-23 15:33:17 +08:00
|
|
|
route := msg.clientRoute
|
|
|
|
|
if !route.bound() {
|
|
|
|
|
route = c.clientSessionRouteSnapshot()
|
2026-04-15 15:24:36 +08:00
|
|
|
}
|
2026-09-23 15:33:17 +08:00
|
|
|
stream, ok := runtime.lookupControl(clientFileScope(), req.StreamID, req.DataID)
|
2026-04-15 15:24:36 +08:00
|
|
|
if !ok {
|
|
|
|
|
resp.Error = errStreamNotFound.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-09-23 15:33:17 +08:00
|
|
|
if !stream.acceptsClientSessionRoute(route) {
|
|
|
|
|
resp.Error = transportDetachedSessionEpochError().Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-04-15 15:24:36 +08:00
|
|
|
if resp.StreamID == "" {
|
|
|
|
|
resp.StreamID = stream.ID()
|
|
|
|
|
}
|
2026-09-23 15:33:17 +08:00
|
|
|
stream.markReset(req.resetError(stream.Channel()))
|
2026-04-15 15:24:36 +08:00
|
|
|
resp.Accepted = true
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (s *ServerCommon) handleInboundStreamReset(msg *Message) {
|
|
|
|
|
req, err := decodeStreamResetRequest(msg)
|
|
|
|
|
resp := StreamResetResponse{StreamID: req.StreamID}
|
|
|
|
|
if err != nil {
|
|
|
|
|
resp.Error = err.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
runtime := s.getStreamRuntime()
|
|
|
|
|
if runtime == nil {
|
|
|
|
|
resp.Error = errStreamRuntimeNil.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
logical := messageLogicalConnSnapshot(msg)
|
|
|
|
|
scope := serverFileScope(logical)
|
2026-09-23 15:33:17 +08:00
|
|
|
stream, ok := runtime.lookupControl(scope, req.StreamID, req.DataID)
|
2026-04-15 15:24:36 +08:00
|
|
|
if !ok {
|
|
|
|
|
resp.Error = errStreamNotFound.Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-09-23 15:33:17 +08:00
|
|
|
if !stream.acceptsTransportGeneration(messageTransportConnSnapshot(msg)) {
|
|
|
|
|
resp.Error = transportDetachedGenerationMismatchError(stream.TransportGeneration(), messageTransportConnSnapshot(msg)).Error()
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-04-15 15:24:36 +08:00
|
|
|
if resp.StreamID == "" {
|
|
|
|
|
resp.StreamID = stream.ID()
|
|
|
|
|
}
|
2026-09-23 15:33:17 +08:00
|
|
|
stream.markReset(req.resetError(stream.Channel()))
|
2026-04-15 15:24:36 +08:00
|
|
|
resp.Accepted = true
|
|
|
|
|
replyStreamControlIfNeeded(msg, resp)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func replyStreamControlIfNeeded(msg *Message, value interface{}) {
|
|
|
|
|
if msg == nil || !requiresSignalReplyWait(msg.TransferMsg) {
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
_ = msg.ReplyObj(value)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func sendStreamOpenClient(ctx context.Context, c Client, req StreamOpenRequest) (StreamOpenResponse, error) {
|
|
|
|
|
if c == nil {
|
|
|
|
|
return StreamOpenResponse{}, errStreamClientNil
|
|
|
|
|
}
|
|
|
|
|
msg, err := c.SendObjCtx(ctx, StreamOpenSignalKey, req)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return StreamOpenResponse{}, err
|
|
|
|
|
}
|
|
|
|
|
return decodeStreamOpenResponse(msg)
|
|
|
|
|
}
|
|
|
|
|
|
2026-09-23 15:33:17 +08:00
|
|
|
func sendStreamOpenClientAtRoute(ctx context.Context, c *ClientCommon, route clientSessionRoute, req StreamOpenRequest) (StreamOpenResponse, error) {
|
|
|
|
|
if c == nil {
|
|
|
|
|
return StreamOpenResponse{}, errStreamClientNil
|
|
|
|
|
}
|
|
|
|
|
msg, err := c.sendObjCtxAtRoute(ctx, route, StreamOpenSignalKey, req)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return StreamOpenResponse{}, err
|
|
|
|
|
}
|
|
|
|
|
return decodeStreamOpenResponse(msg)
|
|
|
|
|
}
|
|
|
|
|
|
2026-04-15 15:24:36 +08:00
|
|
|
func sendStreamOpenServerLogical(ctx context.Context, s Server, logical *LogicalConn, req StreamOpenRequest) (StreamOpenResponse, error) {
|
|
|
|
|
if s == nil {
|
|
|
|
|
return StreamOpenResponse{}, errStreamServerNil
|
|
|
|
|
}
|
|
|
|
|
if logical == nil {
|
|
|
|
|
return StreamOpenResponse{}, errStreamLogicalConnNil
|
|
|
|
|
}
|
|
|
|
|
msg, err := s.SendObjCtxLogical(ctx, logical, StreamOpenSignalKey, req)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return StreamOpenResponse{}, err
|
|
|
|
|
}
|
|
|
|
|
return decodeStreamOpenResponse(msg)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func sendStreamOpenServerTransport(ctx context.Context, s Server, transport *TransportConn, req StreamOpenRequest) (StreamOpenResponse, error) {
|
|
|
|
|
if s == nil {
|
|
|
|
|
return StreamOpenResponse{}, errStreamServerNil
|
|
|
|
|
}
|
|
|
|
|
if transport == nil {
|
|
|
|
|
return StreamOpenResponse{}, errStreamTransportNil
|
|
|
|
|
}
|
|
|
|
|
msg, err := s.SendObjCtxTransport(ctx, transport, StreamOpenSignalKey, req)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return StreamOpenResponse{}, err
|
|
|
|
|
}
|
|
|
|
|
return decodeStreamOpenResponse(msg)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func sendStreamCloseClient(ctx context.Context, c Client, req StreamCloseRequest) (StreamCloseResponse, error) {
|
|
|
|
|
if c == nil {
|
|
|
|
|
return StreamCloseResponse{}, errStreamClientNil
|
|
|
|
|
}
|
|
|
|
|
msg, err := c.SendObjCtx(ctx, StreamCloseSignalKey, req)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return StreamCloseResponse{}, err
|
|
|
|
|
}
|
|
|
|
|
return decodeStreamCloseResponse(msg)
|
|
|
|
|
}
|
|
|
|
|
|
2026-09-23 15:33:17 +08:00
|
|
|
func sendStreamCloseClientAtRoute(ctx context.Context, c *ClientCommon, route clientSessionRoute, req StreamCloseRequest) (StreamCloseResponse, error) {
|
|
|
|
|
if c == nil {
|
|
|
|
|
return StreamCloseResponse{}, errStreamClientNil
|
|
|
|
|
}
|
|
|
|
|
msg, err := c.sendObjCtxAtRoute(ctx, route, StreamCloseSignalKey, req)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return StreamCloseResponse{}, err
|
|
|
|
|
}
|
|
|
|
|
return decodeStreamCloseResponse(msg)
|
|
|
|
|
}
|
|
|
|
|
|
2026-04-15 15:24:36 +08:00
|
|
|
func sendStreamCloseServerLogical(ctx context.Context, s Server, logical *LogicalConn, req StreamCloseRequest) (StreamCloseResponse, error) {
|
|
|
|
|
if s == nil {
|
|
|
|
|
return StreamCloseResponse{}, errStreamServerNil
|
|
|
|
|
}
|
|
|
|
|
if logical == nil {
|
|
|
|
|
return StreamCloseResponse{}, errStreamLogicalConnNil
|
|
|
|
|
}
|
|
|
|
|
msg, err := s.SendObjCtxLogical(ctx, logical, StreamCloseSignalKey, req)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return StreamCloseResponse{}, err
|
|
|
|
|
}
|
|
|
|
|
return decodeStreamCloseResponse(msg)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func sendStreamCloseServerTransport(ctx context.Context, s Server, transport *TransportConn, req StreamCloseRequest) (StreamCloseResponse, error) {
|
|
|
|
|
if s == nil {
|
|
|
|
|
return StreamCloseResponse{}, errStreamServerNil
|
|
|
|
|
}
|
|
|
|
|
if transport == nil {
|
|
|
|
|
return StreamCloseResponse{}, errStreamTransportNil
|
|
|
|
|
}
|
|
|
|
|
msg, err := s.SendObjCtxTransport(ctx, transport, StreamCloseSignalKey, req)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return StreamCloseResponse{}, err
|
|
|
|
|
}
|
|
|
|
|
return decodeStreamCloseResponse(msg)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func sendStreamResetClient(ctx context.Context, c Client, req StreamResetRequest) (StreamResetResponse, error) {
|
|
|
|
|
if c == nil {
|
|
|
|
|
return StreamResetResponse{}, errStreamClientNil
|
|
|
|
|
}
|
|
|
|
|
msg, err := c.SendObjCtx(ctx, StreamResetSignalKey, req)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return StreamResetResponse{}, err
|
|
|
|
|
}
|
|
|
|
|
return decodeStreamResetResponse(msg)
|
|
|
|
|
}
|
|
|
|
|
|
2026-09-23 15:33:17 +08:00
|
|
|
func sendStreamResetClientAtRoute(ctx context.Context, c *ClientCommon, route clientSessionRoute, req StreamResetRequest) (StreamResetResponse, error) {
|
|
|
|
|
if c == nil {
|
|
|
|
|
return StreamResetResponse{}, errStreamClientNil
|
|
|
|
|
}
|
|
|
|
|
msg, err := c.sendObjCtxAtRoute(ctx, route, StreamResetSignalKey, req)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return StreamResetResponse{}, err
|
|
|
|
|
}
|
|
|
|
|
return decodeStreamResetResponse(msg)
|
|
|
|
|
}
|
|
|
|
|
|
2026-04-15 15:24:36 +08:00
|
|
|
func sendStreamResetServerLogical(ctx context.Context, s Server, logical *LogicalConn, req StreamResetRequest) (StreamResetResponse, error) {
|
|
|
|
|
if s == nil {
|
|
|
|
|
return StreamResetResponse{}, errStreamServerNil
|
|
|
|
|
}
|
|
|
|
|
if logical == nil {
|
|
|
|
|
return StreamResetResponse{}, errStreamLogicalConnNil
|
|
|
|
|
}
|
|
|
|
|
msg, err := s.SendObjCtxLogical(ctx, logical, StreamResetSignalKey, req)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return StreamResetResponse{}, err
|
|
|
|
|
}
|
|
|
|
|
return decodeStreamResetResponse(msg)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func sendStreamResetServerTransport(ctx context.Context, s Server, transport *TransportConn, req StreamResetRequest) (StreamResetResponse, error) {
|
|
|
|
|
if s == nil {
|
|
|
|
|
return StreamResetResponse{}, errStreamServerNil
|
|
|
|
|
}
|
|
|
|
|
if transport == nil {
|
|
|
|
|
return StreamResetResponse{}, errStreamTransportNil
|
|
|
|
|
}
|
|
|
|
|
msg, err := s.SendObjCtxTransport(ctx, transport, StreamResetSignalKey, req)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return StreamResetResponse{}, err
|
|
|
|
|
}
|
|
|
|
|
return decodeStreamResetResponse(msg)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func decodeStreamOpenRequest(msg *Message) (StreamOpenRequest, error) {
|
|
|
|
|
var req StreamOpenRequest
|
|
|
|
|
if msg == nil {
|
|
|
|
|
return StreamOpenRequest{}, errStreamIDEmpty
|
|
|
|
|
}
|
|
|
|
|
if err := msg.Value.Orm(&req); err != nil {
|
|
|
|
|
return StreamOpenRequest{}, err
|
|
|
|
|
}
|
|
|
|
|
req = normalizeStreamOpenRequest(req)
|
|
|
|
|
if req.StreamID == "" {
|
|
|
|
|
return StreamOpenRequest{}, errStreamIDEmpty
|
|
|
|
|
}
|
|
|
|
|
return req, nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func decodeStreamCloseRequest(msg *Message) (StreamCloseRequest, error) {
|
|
|
|
|
var req StreamCloseRequest
|
|
|
|
|
if msg == nil {
|
|
|
|
|
return StreamCloseRequest{}, errStreamIDEmpty
|
|
|
|
|
}
|
|
|
|
|
if err := msg.Value.Orm(&req); err != nil {
|
|
|
|
|
return StreamCloseRequest{}, err
|
|
|
|
|
}
|
|
|
|
|
if req.StreamID == "" {
|
|
|
|
|
return StreamCloseRequest{}, errStreamIDEmpty
|
|
|
|
|
}
|
|
|
|
|
return req, nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func decodeStreamResetRequest(msg *Message) (StreamResetRequest, error) {
|
|
|
|
|
var req StreamResetRequest
|
|
|
|
|
if msg == nil {
|
|
|
|
|
return StreamResetRequest{}, errStreamIDEmpty
|
|
|
|
|
}
|
|
|
|
|
if err := msg.Value.Orm(&req); err != nil {
|
|
|
|
|
return StreamResetRequest{}, err
|
|
|
|
|
}
|
2026-09-23 15:33:17 +08:00
|
|
|
if req.StreamID == "" && req.DataID == 0 {
|
2026-04-15 15:24:36 +08:00
|
|
|
return StreamResetRequest{}, errStreamIDEmpty
|
|
|
|
|
}
|
|
|
|
|
return req, nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func decodeStreamOpenResponse(msg Message) (StreamOpenResponse, error) {
|
|
|
|
|
var resp StreamOpenResponse
|
|
|
|
|
if err := msg.Value.Orm(&resp); err != nil {
|
|
|
|
|
return StreamOpenResponse{}, err
|
|
|
|
|
}
|
|
|
|
|
return resp, streamControlResultError("open", resp.Accepted, resp.Error, nil)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func decodeStreamCloseResponse(msg Message) (StreamCloseResponse, error) {
|
|
|
|
|
var resp StreamCloseResponse
|
|
|
|
|
if err := msg.Value.Orm(&resp); err != nil {
|
|
|
|
|
return StreamCloseResponse{}, err
|
|
|
|
|
}
|
|
|
|
|
return resp, streamControlResultError("close", resp.Accepted, resp.Error, nil)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func decodeStreamResetResponse(msg Message) (StreamResetResponse, error) {
|
|
|
|
|
var resp StreamResetResponse
|
|
|
|
|
if err := msg.Value.Orm(&resp); err != nil {
|
|
|
|
|
return StreamResetResponse{}, err
|
|
|
|
|
}
|
|
|
|
|
return resp, streamControlResultError("reset", resp.Accepted, resp.Error, nil)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func streamControlResultError(op string, accepted bool, message string, callErr error) error {
|
|
|
|
|
if callErr != nil {
|
|
|
|
|
return callErr
|
|
|
|
|
}
|
|
|
|
|
if message != "" {
|
|
|
|
|
return streamControlMessageError(message)
|
|
|
|
|
}
|
|
|
|
|
if accepted {
|
|
|
|
|
return nil
|
|
|
|
|
}
|
|
|
|
|
if op == "open" {
|
|
|
|
|
return errStreamRejected
|
|
|
|
|
}
|
|
|
|
|
return errors.New("stream " + op + " rejected")
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func streamControlMessageError(message string) error {
|
2026-09-23 15:33:17 +08:00
|
|
|
if message == errTransportDetached.Error() || strings.HasPrefix(message, errTransportDetached.Error()+":") {
|
|
|
|
|
return errTransportDetached
|
|
|
|
|
}
|
2026-04-15 15:24:36 +08:00
|
|
|
switch message {
|
|
|
|
|
case errStreamNotFound.Error():
|
|
|
|
|
return errStreamNotFound
|
|
|
|
|
case errStreamAlreadyExists.Error():
|
|
|
|
|
return errStreamAlreadyExists
|
|
|
|
|
case errStreamHandlerNotConfigured.Error():
|
|
|
|
|
return errStreamHandlerNotConfigured
|
|
|
|
|
case errStreamLogicalConnNil.Error():
|
|
|
|
|
return errStreamLogicalConnNil
|
|
|
|
|
case errStreamTransportNil.Error():
|
|
|
|
|
return errStreamTransportNil
|
|
|
|
|
case errStreamRuntimeNil.Error():
|
|
|
|
|
return errStreamRuntimeNil
|
|
|
|
|
case errStreamIDEmpty.Error():
|
|
|
|
|
return errStreamIDEmpty
|
|
|
|
|
default:
|
|
|
|
|
return errors.New(message)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func streamTransportGeneration(logical *LogicalConn, transport *TransportConn) uint64 {
|
|
|
|
|
if transport != nil {
|
|
|
|
|
return transport.TransportGeneration()
|
|
|
|
|
}
|
|
|
|
|
if logical != nil {
|
|
|
|
|
return logical.transportGenerationSnapshot()
|
|
|
|
|
}
|
|
|
|
|
return 0
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func streamRemoteResetError(message string) error {
|
|
|
|
|
if message == "" {
|
|
|
|
|
return errStreamReset
|
|
|
|
|
}
|
|
|
|
|
return errors.New(message)
|
|
|
|
|
}
|