(ctx sdk.Context, queuePrefix byte, blockTime time.Time, limit uint32, process func([]byte, time.Time) error)
| 17 | } |
| 18 | |
| 19 | func (k *keeper) processDueQueue(ctx sdk.Context, queuePrefix byte, blockTime time.Time, limit uint32, process func([]byte, time.Time) error) error { |
| 20 | if limit == 0 { |
| 21 | return nil |
| 22 | } |
| 23 | |
| 24 | store := storeprefix.NewStore(ctx.KVStore(k.skey), singletonKey(queuePrefix)) |
| 25 | iter := store.Iterator(nil, nil) |
| 26 | defer func() { |
| 27 | _ = iter.Close() |
| 28 | }() |
| 29 | |
| 30 | entries := make([]dueQueueEntry, 0, limit) |
| 31 | for ; iter.Valid() && uint32(len(entries)) < limit; iter.Next() { |
| 32 | dueTime, err := decodeQueueTime(iter.Key()) |
| 33 | if err != nil { |
| 34 | return err |
| 35 | } |
| 36 | if dueTime.After(blockTime) { |
| 37 | break |
| 38 | } |
| 39 | |
| 40 | entries = append(entries, dueQueueEntry{ |
| 41 | key: append([]byte(nil), iter.Key()...), |
| 42 | dueTime: dueTime, |
| 43 | }) |
| 44 | } |
| 45 | |
| 46 | for _, entry := range entries { |
| 47 | if err := process(entry.key, entry.dueTime); err != nil { |
| 48 | return err |
| 49 | } |
| 50 | store.Delete(entry.key) |
| 51 | } |
| 52 | |
| 53 | return nil |
| 54 | } |
| 55 | |
| 56 | func decodeQueueTime(key []byte) (time.Time, error) { |
| 57 | if len(key) < 8 { |
no test coverage detected