| 139 | } |
| 140 | |
| 141 | func (s *session) request(ctx context.Context, msg Message) (<-chan Message, error) { |
| 142 | if msg.ID == "" { |
| 143 | return nil, fmt.Errorf("wsrelay: message id is required") |
| 144 | } |
| 145 | if _, loaded := s.pending.LoadOrStore(msg.ID, &pendingRequest{ch: make(chan Message, 8)}); loaded { |
| 146 | return nil, fmt.Errorf("wsrelay: duplicate message id %s", msg.ID) |
| 147 | } |
| 148 | value, _ := s.pending.Load(msg.ID) |
| 149 | req := value.(*pendingRequest) |
| 150 | if err := s.send(ctx, msg); err != nil { |
| 151 | if actual, loaded := s.pending.LoadAndDelete(msg.ID); loaded { |
| 152 | req := actual.(*pendingRequest) |
| 153 | req.close() |
| 154 | } |
| 155 | return nil, err |
| 156 | } |
| 157 | go func() { |
| 158 | select { |
| 159 | case <-ctx.Done(): |
| 160 | if actual, loaded := s.pending.LoadAndDelete(msg.ID); loaded { |
| 161 | actual.(*pendingRequest).close() |
| 162 | } |
| 163 | case <-s.closed: |
| 164 | } |
| 165 | }() |
| 166 | return req.ch, nil |
| 167 | } |
| 168 | |
| 169 | func (s *session) cleanup(cause error) { |
| 170 | s.closeOnce.Do(func() { |