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)
| 117 | // |
| 118 | // reqBody 用 json 编码(nil 表示空 body)。 |
| 119 | func (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: |