(ctx context.Context, sequence uint64, timeReceived channels.FeedTimestamp)
| 592 | } |
| 593 | |
| 594 | func (c *changeCache) releaseUnusedSequence(ctx context.Context, sequence uint64, timeReceived channels.FeedTimestamp) { |
| 595 | change := &LogEntry{ |
| 596 | Sequence: sequence, |
| 597 | TimeReceived: timeReceived, |
| 598 | UnusedSequence: true, |
| 599 | } |
| 600 | base.InfofCtx(ctx, base.KeyCache, "Received #%d (unused sequence)", sequence) |
| 601 | |
| 602 | // Since processEntry may unblock pending sequences, if there were any changed channels we need |
| 603 | // to notify any change listeners that are working changes feeds for these channels |
| 604 | var channelSet channels.Set |
| 605 | changedChannels := c.processEntry(ctx, change) |
| 606 | if changedChannels == nil { |
| 607 | channelSet = channels.SetOfNoValidate(unusedSeqChannelID) |
| 608 | } else { |
| 609 | channelSet = channels.SetFromArrayNoValidate(changedChannels) |
| 610 | channelSet.Add(unusedSeqChannelID) |
| 611 | } |
| 612 | if c.notifyChangeFunc != nil && len(channelSet) > 0 { |
| 613 | c.notifyChangeFunc(ctx, channelSet) |
| 614 | } |
| 615 | } |
| 616 | |
| 617 | // releaseUnusedSequenceRange will handle unused sequence range arriving over DCP. It will batch remove from skipped or |
| 618 | // push a range to pending sequences, or both. |
no test coverage detected