MCPcopy Create free account
hub / github.com/0xUnixIO/pulse / CallStream

Method CallStream

internal/nodehub/stream.go:119–181  ·  view source on GitHub ↗

CallStream 发起一次流式调用:发送 ServerMessage{Method, Body} 后,把后续 node 推回的同 reqID push 帧(event="log" / "traceroute_hop")通过 Stream.Frames() 暴露。node 端在流自然结束时回一条普通响应(Ok=true 或 Ok=false+Error), 触发 Stream.Done();caller 也可通过 Stream.Close() 主动取消(向 node 发送 cancel_id 帧)。 reqBody 用 json 编码(nil 表示空 body)。

(ctx context.Context, nodeID, method string, reqBody any)

Source from the content-addressed store, hash-verified

117//
118// reqBody 用 json 编码(nil 表示空 body)。
119func (h *Hub) CallStream(ctx context.Context, nodeID, method string, reqBody any) (*Stream, error) {
120 h.mu.RLock()
121 c, ok := h.conns[nodeID]
122 h.mu.RUnlock()
123 if !ok {
124 return nil, ErrNodeOffline
125 }
126
127 var bodyBytes []byte
128 if reqBody != nil {
129 b, err := json.Marshal(reqBody)
130 if err != nil {
131 return nil, err
132 }
133 bodyBytes = b
134 }
135
136 reqID, err := newReqID()
137 if err != nil {
138 return nil, err
139 }
140
141 s := &Stream{
142 conn: c,
143 reqID: reqID,
144 frames: make(chan StreamFrame, streamFramesBuffer),
145 done: make(chan struct{}),
146 }
147
148 c.streamSubs.Store(reqID, s)
149
150 if err := c.send(&nodev1.ServerMessage{
151 Id: reqID,
152 Method: method,
153 Body: bodyBytes,
154 }); err != nil {
155 c.streamSubs.Delete(reqID)
156 return nil, err
157 }
158
159 // ctx 取消 → 主动 Close(best-effort 通知 node)
160 go func() {
161 select {
162 case <-ctx.Done():
163 // 仅在尚未结束时 Close
164 s.mu.Lock()
165 closed := s.closed
166 s.mu.Unlock()
167 if !closed {
168 // 标记 ctx 错误,再 Close
169 s.mu.Lock()
170 if s.err == nil {
171 s.err = ctx.Err()
172 }
173 s.mu.Unlock()
174 s.Close()
175 }
176 case <-s.done:

Callers 4

callStreamAdapterFunction · 0.80
TestCallStream_LogFramesFunction · 0.80

Calls 6

newReqIDFunction · 0.85
DeleteMethod · 0.65
DoneMethod · 0.65
ErrMethod · 0.65
CloseMethod · 0.65
sendMethod · 0.45

Tested by 3

TestCallStream_LogFramesFunction · 0.64