MCPcopy Create free account
hub / github.com/akash-network/node / processDueQueue

Method processDueQueue

x/verification/keeper/queue.go:19–54  ·  view source on GitHub ↗
(ctx sdk.Context, queuePrefix byte, blockTime time.Time, limit uint32, process func([]byte, time.Time) error)

Source from the content-addressed store, hash-verified

17}
18
19func (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
56func decodeQueueTime(key []byte) (time.Time, error) {
57 if len(key) < 8 {

Calls 5

singletonKeyFunction · 0.85
decodeQueueTimeFunction · 0.85
AfterMethod · 0.80
CloseMethod · 0.65
DeleteMethod · 0.65

Tested by

no test coverage detected