(ctx context.Context)
| 158 | ) |
| 159 | |
| 160 | func (ps *PubSub) subscribeTask(ctx context.Context) error { |
| 161 | logger := log.FromContext(ctx) |
| 162 | ch := ps.sub.Channel() |
| 163 | store := &PubSubStore{ |
| 164 | PubSub: ps, |
| 165 | // NOTE: only for reading; no additional settings needed. |
| 166 | } |
| 167 | for { |
| 168 | select { |
| 169 | case <-ctx.Done(): |
| 170 | return ctx.Err() |
| 171 | case msg, ok := <-ch: |
| 172 | if !ok { |
| 173 | return errChannelClosed.New() |
| 174 | } |
| 175 | var evtPB *ttnpb.Event |
| 176 | switch { |
| 177 | case strings.HasPrefix(msg.Payload, protoEncodingPrefix): |
| 178 | evtPB = &ttnpb.Event{} |
| 179 | err := decodeEventData(msg.Payload, evtPB) |
| 180 | if err != nil { |
| 181 | logger.WithError(err).Warn("Failed to decode event payload") |
| 182 | continue |
| 183 | } |
| 184 | case strings.HasPrefix(msg.Payload, metaEncodingPrefix): |
| 185 | m := strings.Split(strings.TrimPrefix(msg.Payload, metaEncodingPrefix), " ") |
| 186 | var err error |
| 187 | evtPB, err = store.LoadEvent(ctx, m[0]) |
| 188 | if err != nil { |
| 189 | logger.WithError(err).Warn("Failed to load event payload") |
| 190 | continue |
| 191 | } |
| 192 | default: |
| 193 | logger.Warn("Skip decoding event with unexpected encoding") |
| 194 | continue |
| 195 | } |
| 196 | evt, err := events.FromProto(evtPB) |
| 197 | if err != nil { |
| 198 | logger.WithError(err).Warn("Failed to convert event from protobuf") |
| 199 | continue |
| 200 | } |
| 201 | ps.PubSub.Publish(&patternEvent{Event: evt, pattern: msg.Pattern}) |
| 202 | } |
| 203 | } |
| 204 | } |
| 205 | |
| 206 | type patternEvent struct { |
| 207 | events.Event |
nothing calls this directly
no test coverage detected