| 36 | } |
| 37 | |
| 38 | func (h *Hub) doCall(ctx context.Context, nodeID, method string, reqBody any) (json.RawMessage, error) { |
| 39 | h.mu.RLock() |
| 40 | c, ok := h.conns[nodeID] |
| 41 | h.mu.RUnlock() |
| 42 | if !ok { |
| 43 | return nil, ErrNodeOffline |
| 44 | } |
| 45 | |
| 46 | var bodyBytes []byte |
| 47 | if reqBody != nil { |
| 48 | b, err := json.Marshal(reqBody) |
| 49 | if err != nil { |
| 50 | return nil, err |
| 51 | } |
| 52 | bodyBytes = b |
| 53 | } |
| 54 | |
| 55 | reqID, err := newReqID() |
| 56 | if err != nil { |
| 57 | return nil, err |
| 58 | } |
| 59 | |
| 60 | ch := make(chan *nodev1.NodeMessage, 1) |
| 61 | c.pending.Store(reqID, ch) |
| 62 | defer c.pending.Delete(reqID) |
| 63 | |
| 64 | if err := c.send(&nodev1.ServerMessage{ |
| 65 | Id: reqID, |
| 66 | Method: method, |
| 67 | Body: bodyBytes, |
| 68 | }); err != nil { |
| 69 | return nil, err |
| 70 | } |
| 71 | |
| 72 | select { |
| 73 | case msg := <-ch: |
| 74 | if !msg.GetOk() { |
| 75 | errStr := msg.GetError() |
| 76 | if errStr == "" { |
| 77 | errStr = "node returned not-ok with empty error" |
| 78 | } |
| 79 | return nil, errors.New(errStr) |
| 80 | } |
| 81 | return json.RawMessage(msg.GetBody()), nil |
| 82 | |
| 83 | case <-ctx.Done(): |
| 84 | // best-effort 取消下发 |
| 85 | _ = c.send(&nodev1.ServerMessage{CancelId: reqID}) |
| 86 | return nil, ctx.Err() |
| 87 | |
| 88 | case <-c.closed: |
| 89 | return nil, ErrNodeOffline |
| 90 | } |
| 91 | } |
| 92 | |
| 93 | func newReqID() (string, error) { |
| 94 | var buf [16]byte |