Skip to content

Commit aa05aa9

Browse files
committed
feat: add net module
1. The Node to deal with gossip 2. add Node Delegate
1 parent 835c452 commit aa05aa9

3 files changed

Lines changed: 227 additions & 0 deletions

File tree

Lines changed: 63 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,63 @@
1+
package org.tron.core.net.node;
2+
3+
import org.tron.common.overlay.gossip.LocalNode;
4+
import org.tron.core.net.message.BlockMessage;
5+
import org.tron.core.net.message.Message;
6+
import org.tron.core.net.message.MessageTypes;
7+
import org.tron.core.net.message.TransationMessage;
8+
9+
import java.io.UnsupportedEncodingException;
10+
11+
public class Node {
12+
13+
private NodeDelegate nodeDel;
14+
private LocalNode localNode;
15+
16+
public void setNodeDelegate(NodeDelegate nodeDel) {
17+
this.nodeDel = nodeDel;
18+
}
19+
20+
public Node() {
21+
localNode = LocalNode.getInstance();
22+
}
23+
24+
public void start() {
25+
localNode.getGossipManager().registerSharedDataSubscriber((key, oldValue, newValue) -> {
26+
byte[] newValueBytes = null;
27+
try {
28+
newValueBytes = newValue.toString().getBytes("ISO-8859-1");
29+
} catch (UnsupportedEncodingException e) {
30+
e.printStackTrace();
31+
}
32+
33+
recieve(key, newValueBytes);
34+
});
35+
}
36+
37+
public void recieve(String key, byte[] msgStr) {
38+
switch (MessageTypes.valueOf(key)) {
39+
case BLOCK:
40+
handleBlock(new BlockMessage(msgStr));
41+
break;
42+
case TRX:
43+
handleTranscation(new TransationMessage(msgStr));
44+
break;
45+
default:
46+
throw new IllegalArgumentException("No such message");
47+
}
48+
}
49+
50+
public void broadcast(Message msg) {
51+
localNode.broadcast(msg);
52+
}
53+
54+
private void handleBlock(BlockMessage msg) {
55+
nodeDel.handleBlock(msg);
56+
}
57+
58+
private void handleTranscation(TransationMessage msg) {
59+
nodeDel.handleTransation(msg);
60+
}
61+
62+
63+
}
Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,37 @@
1+
package org.tron.core.net.node;
2+
3+
import org.tron.core.net.message.BlockMessage;
4+
import org.tron.core.net.message.Message;
5+
import org.tron.core.net.message.TransationMessage;
6+
7+
import java.util.ArrayList;
8+
9+
public interface NodeDelegate {
10+
11+
void reset();
12+
13+
void start();
14+
15+
void stop();
16+
17+
void handleBlock(BlockMessage blkMsg);
18+
19+
void handleTransation(TransationMessage trxMsg);
20+
21+
void handleMsg(Message msg);
22+
23+
boolean isIncludedBlock(int blkId);
24+
25+
ArrayList<Integer> getBlockIds(ArrayList<Integer> blockChainSynopsis);
26+
27+
ArrayList<Integer> getBlockChainSynopsis(int refPoint, int num);
28+
29+
void sync();
30+
31+
void getBlockNum(int blkId);
32+
33+
void getBlockTime(int blkId);
34+
35+
void getHeadBlockId();
36+
37+
}
Lines changed: 127 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,127 @@
1+
package org.tron.core.net.node;
2+
3+
import com.google.protobuf.InvalidProtocolBufferException;
4+
import org.slf4j.Logger;
5+
import org.slf4j.LoggerFactory;
6+
import org.tron.core.BlockUtils;
7+
import org.tron.core.TransactionUtils;
8+
import org.tron.core.db.BlockStore;
9+
import org.tron.core.net.message.BlockMessage;
10+
import org.tron.core.net.message.Message;
11+
import org.tron.core.net.message.TransationMessage;
12+
import org.tron.protos.core.TronBlock;
13+
import org.tron.protos.core.TronTransaction;
14+
15+
import java.util.ArrayList;
16+
17+
public class NodeImpl implements NodeDelegate {
18+
19+
private static final Logger logger = LoggerFactory.getLogger("NodeImpl");
20+
private Node p2pNode;
21+
private BlockStore blockdb;
22+
23+
// set seeds
24+
@java.lang.Override
25+
public void reset() {
26+
p2pNode = new Node();
27+
p2pNode.setNodeDelegate(this);
28+
}
29+
30+
@java.lang.Override
31+
public void start() {
32+
// init database
33+
logger.info("init database");
34+
blockdb = new BlockStore();
35+
blockdb.initBlockDbSource("database-test", "block");
36+
blockdb.initUnspendDbSource("database-test", "trx");
37+
38+
logger.info("reset p2p network");
39+
reset();
40+
41+
}
42+
43+
@java.lang.Override
44+
public void stop() {
45+
46+
}
47+
48+
@java.lang.Override
49+
public void handleBlock(BlockMessage blkMsg) {
50+
TronBlock.Block block = null;
51+
try {
52+
block = TronBlock.Block.parseFrom(blkMsg.getData());
53+
} catch (InvalidProtocolBufferException e) {
54+
e.printStackTrace();
55+
}
56+
System.out.println("handle block: ");
57+
System.out.println(BlockUtils.toPrintString(block));
58+
}
59+
60+
@java.lang.Override
61+
public void handleTransation(TransationMessage trxMsg) {
62+
TronTransaction.Transaction transaction = null;
63+
try {
64+
transaction = TronTransaction.Transaction.parseFrom(trxMsg.getData());
65+
} catch (InvalidProtocolBufferException e) {
66+
e.printStackTrace();
67+
}
68+
System.out.println("handle transaction: ");
69+
System.out.println(TransactionUtils.toPrintString(transaction));
70+
}
71+
72+
@java.lang.Override
73+
public void handleMsg(Message msg) {
74+
75+
}
76+
77+
@java.lang.Override
78+
public boolean isIncludedBlock(int blkId) {
79+
return false;
80+
}
81+
82+
@java.lang.Override
83+
public ArrayList<Integer> getBlockIds(ArrayList<Integer> blockChainSynopsis) {
84+
return null;
85+
}
86+
87+
@java.lang.Override
88+
public ArrayList<Integer> getBlockChainSynopsis(int refPoint, int num) {
89+
return null;
90+
}
91+
92+
@java.lang.Override
93+
public void sync() {
94+
95+
}
96+
97+
@java.lang.Override
98+
public void getBlockNum(int blkId) {
99+
100+
}
101+
102+
@java.lang.Override
103+
public void getBlockTime(int blkId) {
104+
105+
}
106+
107+
@java.lang.Override
108+
public void getHeadBlockId() {
109+
110+
}
111+
112+
public Node getP2pNode() {
113+
return p2pNode;
114+
}
115+
116+
public void setP2pNode(Node p2pNode) {
117+
this.p2pNode = p2pNode;
118+
}
119+
120+
public BlockStore getBlockdb() {
121+
return blockdb;
122+
}
123+
124+
public void setBlockdb(BlockStore blockdb) {
125+
this.blockdb = blockdb;
126+
}
127+
}

0 commit comments

Comments
 (0)