Author SHA1 Message Date
ufarooqstatus 212933a4af GossipSub shadow simulation scripts 2023-09-28 01:10:36 +05:00
7 changed files with 5075 additions and 376 deletions
+7 -49
View File
@@ -2,21 +2,20 @@
* DST gossipsub test node
* incl shadow simulation setup
* incl awk scripts for detailed analysis
Three simulation cases are supportes
## Scenario1 - Basic Simulation Scenario
### Shadow example
## Shadow example
```sh
nimble install -dy
cd shadow
# the default shadow.yml will start 5k nodes, you might want to change that by removing
# lines and setting PEERS to the number of instances
./run.sh
# the output is a "latencies" file, or you can find each host output in the
# data.shadow folder
./run.sh <runs> <nodes>
# the first parameter <runs> tells the number of simulation runs, and second parameter <nodes> tells the
# number of nodes in simulation, for example ./run.sh 2 3000
# the output for each run creates latencies(X) and shadowlogX files. where X is the simulation number.
# the run script (run.sh) uses awk to summarize latencies(X) and shadowlogX files
# you can use the plotter tool to extract useful metrics & generate a graph
cd ../tools
@@ -26,44 +25,3 @@ nim -d:release c plotter
```
The dependencies will be installed in the `nimbledeps` folder, which enables easy tweaking
## Scenario2 - (P2P-Research Branch) -
Homogeneous nodes/links (100 Mbps bandwidth and 100ms latency)
supports variable network/message size, variable number of publishers and message fragments
automatically updates shadow.yaml file to accomodate different network sizes
awk scripts for detailed analysis
```sh
cd shadow
#the run.sh script is automated to meet different experiment needs, use ./run.sh <num_runs num_peers msg_size num_fragments>
#The below example runs the simulation twice for a 1000 node network. each publisher publishes a 15000 bytes messages, and every message is partitioned into 4 fragments. Total 10 publishers use
./run.sh 2 1000 15000 4 10
# The number of nodes is maintained in the shadow.yaml file, and automatically updated by run.sh.
# The output files latencies(x), stats(x) and shadowlog(x) carries the outputs for each simulation run.
# The summary_dontwant.awk, summary_latency.awk, summary_latency_large.awk, and summary_shadowlog.awk parse the output files.
# The run.sh script automatically calls these files to display the output
# a temperary data.shadow folder is created for each simulation and removed by the run.sh after the simulation is over
```
## Scenario3 - (Realistic-Scenarios Branch) -
Heterogeneous nodes/links (bandwidth/latency/packet_loss_ratio controllable through run.sh)
run.sh uses topogen.py to create a realistic gml file to emulate a real-world network scenario
supports variable network/message size, variable number of publishers and message fragments
automatically generates network_topology.gml and shadow.yaml files to accomodate different network sizes
awk scripts called from run.sh for detailed analysis
```sh
cd shadow
#The following sample command runs simulation 1 time, for a 1000 node network. Each published message size
#is 15KB (no-fragmentation), peer bandwidth varies between 50-130 Mbps, Latency between 60-160ms, and
#bandwidth,latency is roughly distributed in five different groups.
./run.sh 1 1000 15000 1 10 50 130 60 160 5 0.0
# The number of nodes is maintained in the shadow.yaml file, and automatically generated by run.sh.
# The output files latencies(x), stats(x) and shadowlog(x) carries the outputs for each simulation run.
# The summary_dontwant.awk, summary_latency.awk, summary_latency_large.awk, and summary_shadowlog.awk parse the output files.
# The run.sh script automatically calls these files to display the output
# a temperary data.shadow folder is created for each simulation and removed by the run.sh after the simulation is over
```
+18 -69
View File
@@ -1,22 +1,12 @@
import stew/endians2, stew/byteutils, tables, strutils, os
import libp2p, libp2p/protocols/pubsub/rpc/messages
import libp2p/muxers/mplex/lpchannel, libp2p/protocols/ping
#import libp2p/protocols/pubsub/pubsubpeer
import chronos
import sequtils, hashes, math, metrics
from times import getTime, toUnix, fromUnix, `-`, initTime, `$`, inMilliseconds
from nativesockets import getHostname
#These parameters are passed from yaml file, and each defined peer may receive different parameters (e.g. message size)
var
publisherCount = parseInt(getEnv("PUBLISHERS"))
msg_size = parseInt(getEnv("MSG_SIZE"))
chunks = parseInt(getEnv("FRAGMENTS"))
#we experiment with upto 10 fragments. 1 means, the messages are not fragmented
if chunks < 1 or chunks > 10:
chunks = 1
const chunks = 1
proc msgIdProvider(m: Message): Result[MessageId, ValidationResult] =
return ok(($m.data.hash).toBytes())
@@ -25,7 +15,9 @@ proc main {.async.} =
let
hostname = getHostname()
myId = parseInt(hostname[4..^1])
isPublisher = myId <= publisherCount #need to adjust is publishers ldont start from peer1
#publisherCount = client.param(int, "publisher_count")
publisherCount = 10
isPublisher = myId <= publisherCount
#isAttacker = (not isPublisher) and myId - publisherCount <= client.param(int, "attacker_count")
isAttacker = false
rng = libp2p.newRng()
@@ -53,7 +45,7 @@ proc main {.async.} =
anonymize = true,
)
pingProtocol = Ping.new(rng=rng)
gossipSub.parameters.floodPublish = false
gossipSub.parameters.floodPublish = false
#gossipSub.parameters.lazyPushThreshold = 1_000_000_000
#gossipSub.parameters.lazyPushThreshold = 0
gossipSub.parameters.opportunisticGraftThreshold = -10000
@@ -87,9 +79,6 @@ proc main {.async.} =
sentNanosecs = nanoseconds(sentMoment - seconds(sentMoment.seconds))
sentDate = initTime(sentMoment.seconds, sentNanosecs)
diff = getTime() - sentDate
# pubId = byte(data[11])
echo sentUint, " milliseconds: ", diff.inMilliseconds()
@@ -132,7 +121,7 @@ proc main {.async.} =
let connectTo = parseInt(getEnv("CONNECTTO"))
var connected = 0
for peerInfo in peersInfo:
if connected >= connectTo+2: break
if connected >= connectTo: break
let tAddress = "peer" & $peerInfo & ":5000"
echo tAddress
let addrs = resolveTAddress(tAddress).mapIt(MultiAddress.init(it).tryGet())
@@ -148,64 +137,24 @@ proc main {.async.} =
# warmupMessages = client.param(int, "warmup_messages")
#startOfTest = Moment.now() + milliseconds(warmupMessages * maxMessageDelay div 2)
await sleepAsync(12.seconds)
echo "Mesh size: ", gossipSub.mesh.getOrDefault("test").len,
", Total Peers Known : ", gossipSub.gossipsub.getOrDefault("test").len,
", Direct Peers : ", gossipSub.subscribedDirectPeers.getOrDefault("test").len,
", Fanout", gossipSub.fanout.getOrDefault("test").len,
", Heartbeat : ", gossipSub.parameters.heartbeatInterval.milliseconds
await sleepAsync(10.seconds)
echo "Mesh size: ", gossipSub.mesh.getOrDefault("test").len
await sleepAsync(5.seconds)
# Actual message publishing, one message published every 3 seconds
# First 1-2 messages take longer than expected time due to low cwnd.
# warmup_messages can set cwnd to a desired level. or alternatively, warmup messages can be set to 0
let
warmup_messages = 2
#shadow.yaml defines peers with changing latency/bandwith. In the current arrangement all the publishers
#will get different latency/bandwidth
pubStart = 4
pubEnd = pubStart + publisherCount + warmup_messages
#we send warmup_messages for adjusting TCP cwnd
for i in pubStart..<(pubStart + warmup_messages):
await sleepAsync(2.seconds)
if i == myId:
#two warmup messages for cwnd raising
let
now = getTime()
nowInt = seconds(now.toUnix()) + nanoseconds(times.nanosecond(now))
var nowBytes = @(toBytesLE(uint64(nowInt.nanoseconds))) & newSeq[byte](msg_size)
doAssert((await gossipSub.publish("test", nowBytes)) > 0)
#done sending warmup_messages , wait for short time
await sleepAsync(5.seconds)
#We now send publisher_count messages
for msg in (pubStart + warmup_messages) .. pubEnd:#client.param(int, "message_count"):
await sleepAsync(3.seconds)
if msg mod (pubEnd+1) == myId:
for msg in 0 ..< 10:#client.param(int, "message_count"):
await sleepAsync(12.seconds)
if msg mod publisherCount == myId - 1:
#if myId == 1:
let
now = getTime()
nowInt = seconds(now.toUnix()) + nanoseconds(times.nanosecond(now))
#[
if chunks == 1:
var nowBytes = @(toBytesLE(uint64(nowInt.nanoseconds))) & newSeq[byte](50000)
else:
var nowBytes = @(toBytesLE(uint64(nowInt.nanoseconds))) & newSeq[byte](500_000 div chunks)
]#
var nowBytes = @(toBytesLE(uint64(nowInt.nanoseconds))) & newSeq[byte](msg_size div chunks)
#var nowBytes = @(toBytesLE(uint64(nowInt.nanoseconds))) & newSeq[byte](500_000 div chunks)
var nowBytes = @(toBytesLE(uint64(nowInt.nanoseconds))) & newSeq[byte](50)
#echo "sending ", uint64(nowInt.nanoseconds)
for chunk in 0..<chunks:
nowBytes[10] = byte(chunk)
doAssert((await gossipSub.publish("test", nowBytes)) > 0)
echo "Done Publishing ", nowInt.nanoseconds
#we need to export these counters from gossipsub.nim
echo "statcounters: dup_during_validation ", libp2p_gossipsub_duplicate_during_validation.value(),
"\tidontwant_saves ", libp2p_gossipsub_idontwant_saved_messages.value(),
"\tdup_received ", libp2p_gossipsub_duplicate.value(),
"\tUnique_msg_received ", libp2p_gossipsub_received.value(),
"\tStaggered_Saves ", libp2p_gossipsub_staggerDontWantSave.value(),
"\tDontWant_IN_Stagger ", libp2p_gossipsub_staggerDontWantSave2.value()
#echo "BW: ", libp2p_protocols_bytes.value(labelValues=["/meshsub/1.1.0", "in"]) + libp2p_protocols_bytes.value(labelValues=["/meshsub/1.1.0", "out"])
#echo "DUPS: ", libp2p_gossipsub_duplicate.value(), " / ", libp2p_gossipsub_received.value()
waitFor(main())
+17 -49
View File
@@ -1,65 +1,33 @@
#!/bin/sh
set -e
if [ $# -ne 11 ]; then
echo "Usage: $0 <runs> <nodes> <Message_size> <num_fragment> <num_publishers>
<min_bandwidth> <max_bandwidth> <min_latency> <max_latency> <anchor_stages <packet_loss>>"
echo "The following sample command runs simulation 1 time, for a 1000 node network. Each published message size \
is 15KB (no-fragmentation), peer bandwidth varies between 50-130 Mbps, Latency between 60-160ms, and \
bandwidth,latency is roughly distributed in five different groups. \
see the generated network_topology.gml and shadow.yaml for peers/edges details"
echo "$0 1 1000 15000 1 10 50 130 60 160 5 0.0"
if [ $# -ne 2 ]; then
echo "Usage: $0 <runs> <nodes>"
exit 1
fi
runs="$1" #number of simulation runs
nodes="$2" #number of nodes to simulate
msg_size="$3" #message size to use (in bytes)
num_frag="$4" #number of fragments per message (1 for no fragmentation)
num_publishers="$5" #number of publishers
min_bandwidth="$6"
max_bandwidth="$7"
min_latency="$8"
max_latency="$9"
steps="${10}"
pkt_loss="${11}"
runs="$1" #number of simulation runs
nodes="$2" #number of nodes to simulate
shadow_file="shadow.yaml"
sed -i '/environment:/q' "$shadow_file"
sed -E -i "s/\"PEERS\": \"[0-9]+\"/\"PEERS\": \"$nodes\"/" "$shadow_file"
connect_to=5 #number of peers we connect with to form full message mesh
counter=2
while [ $counter -le $nodes ]; do
echo " peer$counter: *client_host" >> "$shadow_file"
counter=$((counter + 1))
done
#topogen.py uses networkx module from python to generate gml and yaml files
PYTHON=$(which python3 || which python)
if [ -z "$PYTHON" ]; then
echo "Error: Python, Networkx is required for topology files generation."
exit 1
fi
"$PYTHON" topogen.py $nodes $min_bandwidth $max_bandwidth $min_latency $max_latency $steps $pkt_loss $msg_size $num_frag $num_publishers
rm -f shadowlog* latencies* stats* main && rm -rf shadow.data/
rm -f shadowlog* latencies* main && rm -rf shadow.data/
nim c -d:chronicles_colors=None --threads:on -d:metrics -d:libp2p_network_protocols_metrics -d:release main
for i in $(seq $runs); do
echo "Running for turn "$i
shadow shadow.yaml > shadowlog$i &&
grep -rne 'milliseconds\|BW' shadow.data/ > latencies$i &&
grep -rne 'statcounters:' shadow.data/ > stats$i
#uncomment to to receive every nodes log in shadow data (only if runs == 1, or change data directory in yaml file)
#rm -rf shadow.data/
shadow shadow.yaml > shadowlog$i && grep -rne 'milliseconds\|BW' shadow.data/ > latencies$i
rm -rf shadow.data/
done
for i in $(seq $runs); do
echo "Summary for turn "$i
if [ "$msg_size" -lt 1000 ]; then
awk -f summary_latency.awk latencies$i #precise per hop coverage for short messages only
else
awk -f summary_latency_large.awk latencies$i #estimated coverage for large messages (TxTime adds to latency)
fi
awk -f summary_latency.awk latencies$i
awk -f summary_shadowlog.awk shadowlog$i
awk -f summary_dontwant.awk stats$i
done
done
+5033
View File
File diff suppressed because it is too large Load Diff
-30
View File
@@ -1,30 +0,0 @@
BEGIN {
FS = " "; #default column separator
idontwant_saves = min_idontwant = max_idontwant = 0;
dup_received = min_dup = max_dup = 0;
unique_msg_received = 0;
stagger_saves = 0;
stagger_DontWantSaves = 0;
}
{
#print $5, $7, $9
idontwant_saves += $5
if ($5 < min_idontwant || min_idontwant == 0) min_idontwant = $5
if ($5 > max_idontwant) max_idontwant = $5
dup_received += $7
if ($7 < min_dup || min_dup == 0) min_dup = $7
if ($7 > max_dup) max_dup = $7
unique_msg_received += $9
stagger_saves += $11
stagger_DontWantSaves += $13
}
END {
print "idontwant_saves min, max, avg, total : ", min_idontwant, "\t", max_idontwant, "\t", idontwant_saves/NR, "\t", idontwant_saves
print "dup_received min, max, avg, total : ", min_dup, "\t", max_dup, "\t", dup_received/NR, "\t", dup_received
print "Unique_msg_received: ", unique_msg_received, "\tStagger Saves : ", stagger_saves, "\tStaggerDontWantSaves", stagger_DontWantSaves
}
-73
View File
@@ -1,73 +0,0 @@
# we parse the latencies(x) file produced by run.sh to receive results summary (Max/Avg Latency --> per packet, overall)
# runs $awk -f result_summary.awk latencies(x)
BEGIN {
FS = " "; #default column separator
network_size = 0
max_nw_lat = sum_nw_lat = sum_max_delays = 0
hop_lat = 100 #should be consistent with shadow.yaml
}
{
clean_int = $3
gsub(/[^0-9]/, "", clean_int);
if ($3 == clean_int){ #get rid of unwanted rows
sum_nw_lat += $NF
if (max_nw_lat < $NF) {max_nw_lat = $NF}
if (split($1, arr, "peer|/main|:.*:")) {
#$3 = rx_latency, arr[4] = publish_time, arr[2] = peerID
#We compute network-wide dissemination latency for each message
if (max_msg_latency[arr[4]] < $NF) {max_msg_latency[arr[4]] = $NF}
#we round to values to nearest hop_lat to estimate hop coverage
rounded_RxTime = (int($3/hop_lat + 0.5)) * hop_lat
lat_arr[arr[4], rounded_RxTime]++;
msg_arr[arr[4]] = 1; #we maintain set of messages identified by their publish time
if (network_size < arr[2]) {network_size = arr[2]}
}
}
}
END {
print "Total Nodes : ", network_size, "Total Messages Published : ", length(msg_arr),
"Network Latency\t MAX : ", max_nw_lat, "\tAverage : ", sum_nw_lat/NR
print " Message ID \t Avg Latency \t Messages Received"
for (value in msg_arr) {
sum_rx_msgs = 0;
latency = 0;
spread[1] = spread[2] = spread[3] = spread[4] = spread[5] = spread[6] = spread[7] = spread[8] = spread[9] = 0
spread[10] = spread[11] = spread[12] = spread[13] = spread[14] = spread[15] = spread[16] = spread[17] = spread[18] = 0
for (key in lat_arr) {
split(key, parts, SUBSEP);
if (parts[1] == value) {
#parts[2] recv time
#10% 20% 30%....90% under parts[2]
sum_rx_msgs = sum_rx_msgs + lat_arr[key]; #total receives / message
latency = latency + (lat_arr[key] * parts[2])
spread[ int((parts[2]) / hop_lat) ] = lat_arr[key] #hop-by-hop spread count of messages
}
}
print value, "\t", latency/sum_rx_msgs, "\t ", sum_rx_msgs, "spread is",
spread[1], spread[2], spread[3], spread[4], spread[5], spread[6], spread[7], spread[8], spread[9],
spread[10], spread[11], spread[12], spread[13], spread[14], spread[15], spread[16], spread[17], spread[18],
spread[19], spread[20], spread[21], spread[22], spread[23], spread[24], spread[25], spread[26], spread[27],
spread[28], spread[29], spread[30], spread[31], spread[32], spread[33], spread[34], spread[35], spread[36],
spread[37], spread[38], spread[39], spread[40], spread[41], spread[42], spread[43], spread[44], spread[45],
spread[46], spread[47], spread[48], spread[49], spread[50], spread[51], spread[52], spread[53], spread[54]
delete spread
}
for (delay_val in max_msg_latency) {
print "MAX delay for ", delay_val, "is \t", max_msg_latency[delay_val]
sum_max_delays = sum_max_delays + max_msg_latency[delay_val]
}
print "Total Messages Published : ", length(max_msg_latency), "Average Max Message Dissemination Latency : ", sum_max_delays/length(max_msg_latency)
}
-106
View File
@@ -1,106 +0,0 @@
import sys, math, networkx as nx
args = sys.argv
if len(args) != 11:
print("Usage: python topogen.py <network_size> <min_bandwidth> <max_bandwidth> <min_latency> <max_latency> <anchor_stages> \
<packet_loss> <message_size> <num_frags> <num_publishers>")
print("Please note that bandwith and latency are integer values in Mbps and ms respectively")
print("Anchor stages represent the number of bandwidth and latency variations")
print("packet_loss [0-1], Message size [KB], num_frags [number of fragments/message 1-10], num_publishers [number of publishers]")
exit(-1)
print (args[1:])
try:
networkSize = int(args[1])
minBandwidth = int(args[2])
maxBandwidth = int(args[3])
minLatency = int(args[4])
maxLatency = int(args[5])
steps = int(args[6])
packetLoss = float(args[7])
messageSize = int(args[8])
numFrags = int(args[9])
numPublishers = int(args[10])
except ValueError:
print("Usage: python topogen.py <network_size> <min_bandwidth> <max_bandwidth> <min_latency> <max_latency> <anchor_stages> \
<packet_loss> <message_size> <num_frags> <num_publishers>")
print("Please note that bandwith and latency are integer values in Mbps and ms respectively")
print("Anchor stages represent the number of bandwidth and latency variations")
print("packet_loss [0-1], Message size [KB], num_frags [number of fragments/message 1-10], num_publishers [number of publishers]")
exit(-1)
gml_file = "network_topology.gml" #network topology layout in gml format, to be used by the yaml file
yaml_file = "shadow.yaml" #shadow simulator settings
connections = 5 #Initial connections to form full-message mesh
bandwidthJump = (maxBandwidth-minBandwidth)/(steps-1)
latencyJump = int((maxLatency-minLatency)/steps)
"""
We create network work graph, with 'steps' number of independent nodes. And all the nodes must be connected.
Shadow uses accumulative edge latencies to route traffic through the shortest paths (accumulative link latencies)
Multiple hosts can connect with a single node. The node must define 'host_bandwidth_up' and 'host_bandwidth_down'
bandwidths, and each connected host gets this bandwidth allocated (bandwidth is not shared between hosts)
We MUST have an edge connecting a node to itself. All the Intra-node communications (among the hosts connected to
the same node) happen by using that edge.
latency and packet loss are edge characteristics
"""
G=nx.complete_graph(steps)
for i in range(0, steps):
nodeBw = str(math.ceil(i * bandwidthJump + minBandwidth)) + " Mbit"
G.nodes[i]["hostbandwidthup"] = nodeBw
G.nodes[i]["hostbandwidthdown"] = nodeBw
G.add_edge(i,i)
G.edges[i,i]["latency"] = str( max((steps-i)*latencyJump, minLatency) ) + " ms"
G.edges[i,i]["packetloss"] = packetLoss
for j in range(i+1, steps):
edgeLatency = min(math.ceil((steps-j)*latencyJump + minLatency), maxLatency)
G.edges[i,j]["latency"] = str(edgeLatency) + " ms"
G.edges[i,j]["packetloss"] = packetLoss
nx.write_gml(G, gml_file)
#networkx package can not write underscores. so we created gml without underscores. Now we embed them underscores
with open(gml_file, 'r') as file:
gml_content = file.read()
modified_content = gml_content.replace("hostbandwidth", "host_bandwidth_")
modified_content = modified_content.replace("packetloss", "packet_loss")
with open(gml_file, "w") as file:
file.write(modified_content)
#we created the gml. now we create the yaml file required by shadow
m1 = "\n network_node_id: "
m2 = "\n processes:"
m3 = "\n - path: ./main"
m4 = "\n start_time: 5s"
with open(yaml_file, "w") as file:
file.write("general:\n bootstrap_end_time: 10s\n heartbeat_interval: 12s\n stop_time: 15m\n")
file.write(" progress: true\n\nexperimental:\n use_memory_manager: false\n\n")
file.write("network:\n graph:\n type: gml\n file:\n path: " + gml_file)
file.write("\n\nhosts:\n")
#we create 'steps' number of template peers, to be used by the remaining peers
for i in range(0,steps):
file.write(" peer" + str(i+1) + ": &client_host" + str(i))
file.write(m1 + str(i) + m2 + m3 + m4)
file.write("\n environment: {\"PEERS\": \"" + str(networkSize) +
"\", \"CONNECTTO\": \"" + str(connections) +
"\", \"MSG_SIZE\": \"" + str(messageSize) +
"\", \"FRAGMENTS\": \"" + str(numFrags) +
"\", \"PUBLISHERS\": \"" + str(numPublishers) + "\"}\n")
#we populate remaining peers on populated samples
for i in range(steps, networkSize):
file.write(" peer" + str(i+1) + ": *client_host" + str(i%steps) + "\n")