MCPcopy Create free account
hub / github.com/dank/rlapi / readMessages

Method readMessages

psynetrpc.go:194–236  ·  view source on GitHub ↗
()

Source from the content-addressed store, hash-verified

192}
193
194func (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
238func (p *PsyNetRPC) sendRequestAsync(ctx context.Context, service string, data interface{}) (<-chan *PsyResponse, error) {
239 if !p.IsConnected() {

Calls 5

CloseMethod · 0.95
parseMessageMethod · 0.95
sendEventMethod · 0.95
ErrorMethod · 0.80
StringMethod · 0.45