| 192 | } |
| 193 | |
| 194 | func (p *PsyNetRPC) readMessages() { |
| 195 | defer func() { |
| 196 | _ = p.Close() |
| 197 | }() |
| 198 | |
| 199 | for { |
| 200 | _, message, err := p.wsConn.ReadMessage() |
| 201 | if err != nil { |
| 202 | p.logger.Error("failed to read websocket message", slog.Any("err", err)) |
| 203 | break |
| 204 | } |
| 205 | |
| 206 | if strings.HasPrefix(string(message), "PsyPong:") { |
| 207 | select { |
| 208 | case p.pongChan <- struct{}{}: |
| 209 | default: |
| 210 | } |
| 211 | continue |
| 212 | } |
| 213 | |
| 214 | p.logger.Debug("received websocket response", slog.String("message", string(message))) |
| 215 | |
| 216 | response, err := p.parseMessage(string(message)) |
| 217 | if err != nil { |
| 218 | p.logger.Error("failed to parse psynet message", slog.Any("err", err), slog.String("message", string(message))) |
| 219 | p.sendEvent(EventTypeMessage, string(message)) |
| 220 | continue |
| 221 | } |
| 222 | |
| 223 | if response.ResponseID != "" { |
| 224 | p.mu.Lock() |
| 225 | ch, exists := p.pendingReqs[response.ResponseID] |
| 226 | p.mu.Unlock() |
| 227 | |
| 228 | if exists { |
| 229 | ch <- response |
| 230 | continue |
| 231 | } |
| 232 | } |
| 233 | |
| 234 | p.sendEvent(EventTypeMessage, string(message)) |
| 235 | } |
| 236 | } |
| 237 | |
| 238 | func (p *PsyNetRPC) sendRequestAsync(ctx context.Context, service string, data interface{}) (<-chan *PsyResponse, error) { |
| 239 | if !p.IsConnected() { |