1+ package p2p
2+
3+ import (
4+ "context"
5+ "fmt"
6+ "net/http"
7+ "time"
8+
9+ "github.com/libp2p/go-libp2p/core/peer"
10+ )
11+
12+ // DAHeightResponse represents the response from a peer's DA height query
13+ type DAHeightResponse struct {
14+ Height uint64 `json:"height"`
15+ Timestamp time.Time `json:"timestamp"`
16+ }
17+
18+ // QueryPeersDAHeight queries connected peers for their DA included height
19+ // Returns the maximum height found across all responsive peers
20+ func (c * Client ) QueryPeersDAHeight (ctx context.Context , timeout time.Duration ) (uint64 , error ) {
21+ peers , err := c .GetPeers ()
22+ if err != nil {
23+ return 0 , fmt .Errorf ("failed to get peers: %w" , err )
24+ }
25+
26+ c .logger .Debug ().Int ("peer_count" , len (peers )).Msg ("querying peers for DA height" )
27+
28+ if len (peers ) == 0 {
29+ return 0 , fmt .Errorf ("no connected peers available to query" )
30+ }
31+
32+ // Channel to collect results from peer queries
33+ type peerResult struct {
34+ peerID peer.ID
35+ height uint64
36+ err error
37+ }
38+
39+ resultCh := make (chan peerResult , len (peers ))
40+
41+ // Query each peer concurrently
42+ for _ , peerInfo := range peers {
43+ go func (pInfo peer.AddrInfo ) {
44+ height , err := c .queryPeerDAHeight (ctx , pInfo , timeout )
45+ resultCh <- peerResult {
46+ peerID : pInfo .ID ,
47+ height : height ,
48+ err : err ,
49+ }
50+ }(peerInfo )
51+ }
52+
53+ // Collect results with timeout
54+ var maxHeight uint64
55+ var successCount int
56+
57+ queryCtx , cancel := context .WithTimeout (ctx , timeout )
58+ defer cancel ()
59+
60+ for i := 0 ; i < len (peers ); i ++ {
61+ select {
62+ case result := <- resultCh :
63+ if result .err != nil {
64+ c .logger .Debug ().
65+ Str ("peer_id" , result .peerID .String ()).
66+ Err (result .err ).
67+ Msg ("failed to query peer for DA height" )
68+ continue
69+ }
70+
71+ successCount ++
72+ if result .height > maxHeight {
73+ maxHeight = result .height
74+ }
75+
76+ c .logger .Debug ().
77+ Str ("peer_id" , result .peerID .String ()).
78+ Uint64 ("height" , result .height ).
79+ Msg ("received DA height from peer" )
80+
81+ case <- queryCtx .Done ():
82+ c .logger .Warn ().
83+ Int ("responses" , successCount ).
84+ Int ("total_peers" , len (peers )).
85+ Msg ("timeout while querying peers for DA height" )
86+ break
87+ }
88+ }
89+
90+ if successCount == 0 {
91+ return 0 , fmt .Errorf ("no peers responded successfully to DA height query" )
92+ }
93+
94+ c .logger .Info ().
95+ Uint64 ("max_height" , maxHeight ).
96+ Int ("successful_queries" , successCount ).
97+ Int ("total_peers" , len (peers )).
98+ Msg ("completed peer DA height queries" )
99+
100+ return maxHeight , nil
101+ }
102+
103+ // queryPeerDAHeight queries a single peer for their DA included height
104+ func (c * Client ) queryPeerDAHeight (ctx context.Context , peerInfo peer.AddrInfo , timeout time.Duration ) (uint64 , error ) {
105+ // For now, we'll try to infer the HTTP RPC endpoint from the peer's multiaddr
106+ // This is a simplified approach - in a production system, peers might advertise their RPC endpoints
107+
108+ // Try common RPC ports - this is a heuristic approach
109+ // In practice, peers could advertise their RPC endpoints through peer discovery protocols
110+ rpcPorts := []string {"8080" , "8081" , "26657" }
111+
112+ for _ , addr := range peerInfo .Addrs {
113+ // Extract IP from multiaddr
114+ ip := extractIPFromMultiaddr (addr .String ())
115+ if ip == "" {
116+ continue
117+ }
118+
119+ // Try each potential RPC port
120+ for _ , port := range rpcPorts {
121+ endpoint := fmt .Sprintf ("http://%s:%s" , ip , port )
122+ height , err := c .queryEndpointDAHeight (ctx , endpoint , timeout )
123+ if err == nil {
124+ return height , nil
125+ }
126+ }
127+ }
128+
129+ return 0 , fmt .Errorf ("could not query DA height from peer %s" , peerInfo .ID .String ())
130+ }
131+
132+ // queryEndpointDAHeight queries a specific HTTP endpoint for DA height
133+ func (c * Client ) queryEndpointDAHeight (ctx context.Context , endpoint string , timeout time.Duration ) (uint64 , error ) {
134+ // Create HTTP client with timeout
135+ _ = & http.Client {
136+ Timeout : timeout ,
137+ }
138+
139+ // Try the store service endpoint to get DA included height
140+ // This uses the existing RPC infrastructure
141+ _ = fmt .Sprintf ("%s/evnode.v1.StoreService/GetMetadata" , endpoint )
142+
143+ // TODO: Implement proper RPC call to get DA included height from store
144+ // For now, return 0 to indicate this is a placeholder
145+ // In a complete implementation, this would make a proper gRPC/Connect call
146+ // to the store service to retrieve the DA included height from metadata
147+
148+ return 0 , fmt .Errorf ("DA height querying not fully implemented yet" )
149+ }
150+
151+ // extractIPFromMultiaddr extracts IP address from a multiaddr string
152+ func extractIPFromMultiaddr (multiaddr string ) string {
153+ // This is a simplified implementation
154+ // In practice, you'd use the multiaddr library to properly parse addresses
155+ // For now, return empty string to indicate we can't extract IP
156+ return ""
157+ }
0 commit comments