(msg *Message)
| 251 | } |
| 252 | |
| 253 | func (q *queue) sendAsyn(msg *Message) error { |
| 254 | if q.isClosed() { |
| 255 | return types.ErrChannelClosed |
| 256 | } |
| 257 | sub := q.chanSub(msg.Topic) |
| 258 | if sub.isClose == 1 { |
| 259 | return types.ErrChannelClosed |
| 260 | } |
| 261 | select { |
| 262 | case sub.low <- msg: |
| 263 | return nil |
| 264 | default: |
| 265 | qlog.Error("send asyn err", "msg", msg, "err", ErrQueueChannelFull) |
| 266 | return ErrQueueChannelFull |
| 267 | } |
| 268 | } |
| 269 | |
| 270 | func (q *queue) sendLowTimeout(msg *Message, timeout time.Duration) error { |
| 271 | if q.isClosed() { |
no test coverage detected