-
Notifications
You must be signed in to change notification settings - Fork 1
Expand file tree
/
Copy pathutils.go
More file actions
90 lines (79 loc) · 2.72 KB
/
Copy pathutils.go
File metadata and controls
90 lines (79 loc) · 2.72 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
// Copyright 2018-2025 Celer Network
package cnode
import (
"fmt"
"time"
"github.com/celer-network/agent-pay/common"
"github.com/celer-network/agent-pay/config"
"github.com/celer-network/agent-pay/ctype"
"github.com/celer-network/agent-pay/ledgerview"
"github.com/celer-network/agent-pay/rpc"
"github.com/celer-network/agent-pay/utils"
"github.com/celer-network/agent-pay/utils/lease"
"github.com/celer-network/goutils/log"
)
func (c *CNode) GetChannelIdForPeer(peer, tokenAddr ctype.Addr) (ctype.CidType, error) {
tokenInfo := utils.GetTokenInfoFromAddress(tokenAddr)
cid, found, err := c.dal.GetCidByPeerToken(peer, tokenInfo)
if err != nil {
return ctype.ZeroCid, err
}
if !found {
return ctype.ZeroCid, fmt.Errorf("No cid found for the peer and token pair")
}
return cid, nil
}
func (c *CNode) GetBalance(cid ctype.CidType) (*common.ChannelBalance, error) {
nowTs := uint64(time.Now().Unix())
return ledgerview.GetBalance(c.dal, cid, c.nodeConfig.GetOnChainAddr(), nowTs)
}
// GetJoinStatusForNode gets the join status of an endpoint
// CAVEAT: Note that this will break if we decide to set a default route
// so that LookupNextChannelOnToken always returns a nextHop for any query.
// TODO: May not rely on LookupNextChannelOnToken in the future(yilun)
func (c *CNode) GetJoinStatusForNode(dst, tokenAddr ctype.Addr) rpc.JoinCelerStatus {
// look up next hop channel
_, nxtHopAddr, err := c.routeForwarder.LookupNextChannelOnToken(dst, tokenAddr)
if err != nil {
return rpc.JoinCelerStatus_NOT_JOIN
}
if dst == nxtHopAddr {
return rpc.JoinCelerStatus_LOCAL
}
return rpc.JoinCelerStatus_REMOTE
}
func (c *CNode) registerEventListener() error {
log.Infoln("register event listener", c.nodeConfig.GetSvrName())
deadline := time.Now().Add(config.EventListenerLeaseTimeout)
var err error
for time.Now().Before(deadline) {
err = lease.Acquire(c.dal, config.EventListenerLeaseName, c.nodeConfig.GetSvrName(), config.EventListenerLeaseTimeout)
if err == nil {
return nil
}
log.Warnf("register event listener failed (%s), retry every 10 seconds until %s", err, deadline.UTC())
time.Sleep(10 * time.Second)
}
log.Error(err)
return fmt.Errorf("register event listener error: %w", err)
}
func (c *CNode) keepAliveEventListener() {
ticker := time.NewTicker(config.EventListenerLeaseRenewInterval)
defer ticker.Stop()
for {
select {
case <-c.quit:
return
case <-ticker.C:
err := lease.Renew(c.dal, config.EventListenerLeaseName, c.nodeConfig.GetSvrName())
if err != nil {
log.Fatalln("failed to renew event listener lease:", err)
}
}
}
}
// SignState signs the data using cnode crypto and return result
func (c *CNode) SignState(in []byte) []byte {
ret, _ := c.signer.SignEthMessage(in)
return ret
}