| 782 | } |
| 783 | |
| 784 | func (sub *ClientSubscription) forward() (err error, unsubscribeServer bool) { |
| 785 | cases := []reflect.SelectCase{ |
| 786 | {Dir: reflect.SelectRecv, Chan: reflect.ValueOf(sub.quit)}, |
| 787 | {Dir: reflect.SelectRecv, Chan: reflect.ValueOf(sub.in)}, |
| 788 | {Dir: reflect.SelectSend, Chan: sub.channel}, |
| 789 | } |
| 790 | buffer := list.New() |
| 791 | defer buffer.Init() |
| 792 | for { |
| 793 | var chosen int |
| 794 | var recv reflect.Value |
| 795 | if buffer.Len() == 0 { |
| 796 | // Idle, omit send case. |
| 797 | chosen, recv, _ = reflect.Select(cases[:2]) |
| 798 | } else { |
| 799 | // Non-empty buffer, send the first queued item. |
| 800 | cases[2].Send = reflect.ValueOf(buffer.Front().Value) |
| 801 | chosen, recv, _ = reflect.Select(cases) |
| 802 | } |
| 803 | |
| 804 | switch chosen { |
| 805 | case 0: // <-sub.quit |
| 806 | return nil, false |
| 807 | case 1: // <-sub.in |
| 808 | val, err := sub.unmarshal(recv.Interface().(json.RawMessage)) |
| 809 | if err != nil { |
| 810 | return err, true |
| 811 | } |
| 812 | if buffer.Len() == maxClientSubscriptionBuffer { |
| 813 | return ErrSubscriptionQueueOverflow, true |
| 814 | } |
| 815 | buffer.PushBack(val) |
| 816 | case 2: // sub.channel<- |
| 817 | cases[2].Send = reflect.Value{} // Don't hold onto the value. |
| 818 | buffer.Remove(buffer.Front()) |
| 819 | } |
| 820 | } |
| 821 | } |
| 822 | |
| 823 | func (sub *ClientSubscription) unmarshal(result json.RawMessage) (interface{}, error) { |
| 824 | val := reflect.New(sub.etype) |