diff --git a/.idea/.gitignore b/.idea/.gitignore
new file mode 100644
index 0000000..1c2fda5
--- /dev/null
+++ b/.idea/.gitignore
@@ -0,0 +1,8 @@
+# Default ignored files
+/shelf/
+/workspace.xml
+# Editor-based HTTP Client requests
+/httpRequests/
+# Datasource local storage ignored files
+/dataSources/
+/dataSources.local.xml
diff --git a/.idea/distributed-framework.iml b/.idea/distributed-framework.iml
new file mode 100644
index 0000000..338a266
--- /dev/null
+++ b/.idea/distributed-framework.iml
@@ -0,0 +1,9 @@
+
+
+
+
+
+
+
+
+
\ No newline at end of file
diff --git a/.idea/modules.xml b/.idea/modules.xml
new file mode 100644
index 0000000..1c4fd81
--- /dev/null
+++ b/.idea/modules.xml
@@ -0,0 +1,8 @@
+
+
+
+
+
+
+
+
\ No newline at end of file
diff --git a/.idea/vcs.xml b/.idea/vcs.xml
new file mode 100644
index 0000000..c8397c9
--- /dev/null
+++ b/.idea/vcs.xml
@@ -0,0 +1,6 @@
+
+
+
+
+
+
\ No newline at end of file
diff --git a/Makefile b/Makefile
index 19cecd0..579b77f 100644
--- a/Makefile
+++ b/Makefile
@@ -26,6 +26,7 @@ leader_undirected_ring_example:
leader_directed_clique_example:
go run example/leader_clique_async_loui_matsushita_west.go 10
+ go run example/leader_clique_async_loui_matsushita_west_2.go 10
leader_undirected_clique_example:
go run example/leader_clique_async_korach_moran_zaks.go 10
@@ -66,6 +67,7 @@ graphs_mst_example:
graphs_mis_example:
go run example/graphs_mis_sync_luby.go 20 0.25
go run example/graphs_mis_sync_metivier_c.go 20 0.25
+ go run example/graphs_mis_sync_metivier_a.go 20 0.25
graphs_ds_example:
go run example/graphs_ds_sync_lrg.go 10 0.70
diff --git a/example/graphs_mis_sync_metivier_a.go b/example/graphs_mis_sync_metivier_a.go
new file mode 100644
index 0000000..18867ca
--- /dev/null
+++ b/example/graphs_mis_sync_metivier_a.go
@@ -0,0 +1,15 @@
+package main
+
+import (
+ "fmt"
+ "github.com/krzysztof-turowski/distributed-framework/graphs/mis/sync_metivier_a"
+ "os"
+ "strconv"
+)
+
+func main() {
+ n, _ := strconv.Atoi(os.Args[len(os.Args)-2])
+ p, _ := strconv.ParseFloat(os.Args[len(os.Args)-1], 64)
+ fmt.Println(n, p)
+ sync_metivier_a.Run(n, p)
+}
diff --git a/example/leader_clique_async_loui_matsushita_west_2.go b/example/leader_clique_async_loui_matsushita_west_2.go
new file mode 100644
index 0000000..443f5b3
--- /dev/null
+++ b/example/leader_clique_async_loui_matsushita_west_2.go
@@ -0,0 +1,23 @@
+package main
+
+import (
+ "fmt"
+ "github.com/krzysztof-turowski/distributed-framework/leader/clique/async_loui_matsushita_west_2"
+ "github.com/krzysztof-turowski/distributed-framework/lib"
+ "os"
+ "strconv"
+)
+
+func main() {
+ if len(os.Args) < 2 {
+ fmt.Println("No size specified.")
+ return
+ }
+ n, err := strconv.Atoi(os.Args[1])
+ if err != nil || n < 2 {
+ fmt.Println("Invalid size: ", n)
+ return
+ }
+ nodes, runner := lib.BuildHamiltonianOrientedCompleteGraph(n)
+ async_loui_matsushita_west_2.Run(nodes, runner)
+}
diff --git a/graphs/mis/sync_metivier_a/sync_metivier_a.go b/graphs/mis/sync_metivier_a/sync_metivier_a.go
new file mode 100644
index 0000000..e62d35d
--- /dev/null
+++ b/graphs/mis/sync_metivier_a/sync_metivier_a.go
@@ -0,0 +1,241 @@
+package sync_metivier_a
+
+import (
+ "encoding/json"
+ "fmt"
+ "log"
+ "math"
+ "math/rand"
+
+ "github.com/krzysztof-turowski/distributed-framework/lib"
+)
+
+const (
+ None = iota
+)
+
+type State struct {
+ InGraph bool
+ InMIS bool
+ NeighborInMIS bool
+ SubPhaseNum int
+ RandomMin uint64
+ ReceivedMin uint64
+ AliveNeighbors map[int]bool
+}
+
+const (
+ phaseMin = iota
+ phaseClaim
+ phaseDecline
+ subPhases = 3
+)
+
+type Message struct {
+ MessageType int
+ Info uint64
+}
+
+const (
+ SendingMinimum = iota
+ JoiningMIS
+ NeverJoiningMIS
+)
+
+func setState(node lib.Node, state State) {
+ newState, _ := json.Marshal(state)
+ node.SetState(newState)
+}
+
+func getState(node lib.Node) State {
+ encodedState := node.GetState()
+ var state State
+ json.Unmarshal(encodedState, &state)
+ return state
+}
+
+func prepareMessage(message Message) []byte {
+ encoded, _ := json.Marshal(message)
+ return encoded
+}
+
+func receiveMessage(node lib.Node, idx int) Message {
+ var message Message
+ receivedMessage := node.ReceiveMessage(idx)
+ err := json.Unmarshal(receivedMessage, &message)
+ if err != nil {
+ return Message{}
+ }
+ return message
+}
+
+func sendNullMessage(node lib.Node, i int) {
+ node.SendMessage(i, nil)
+}
+
+func initialize(node lib.Node) {
+ state := State{
+ true,
+ false,
+ false,
+ None,
+ None,
+ None,
+ map[int]bool{},
+ }
+ for i := range node.GetInNeighbors() {
+ state.AliveNeighbors[i] = true
+ }
+ setState(node, state)
+}
+
+func phasedSend(
+ node lib.Node,
+ state *State,
+ subPhaseNum int,
+ idx int,
+) {
+ switch subPhaseNum {
+ case phaseMin:
+ node.SendMessage(
+ idx,
+ prepareMessage(
+ Message{
+ SendingMinimum,
+ state.RandomMin,
+ },
+ ))
+ case phaseClaim:
+ if state.InMIS {
+ node.SendMessage(
+ idx,
+ prepareMessage(
+ Message{
+ JoiningMIS,
+ state.RandomMin,
+ },
+ ))
+ } else {
+ sendNullMessage(node, idx)
+ }
+ case phaseDecline:
+ if state.NeighborInMIS {
+ node.SendMessage(
+ idx,
+ prepareMessage(
+ Message{
+ NeverJoiningMIS,
+ state.RandomMin,
+ },
+ ))
+ state.InGraph = false
+ } else {
+ sendNullMessage(node, idx)
+ }
+ }
+}
+
+func phasedReceive(
+ node lib.Node,
+ nodeState *State,
+ subPhaseNum int,
+ idx int,
+) {
+ message := receiveMessage(node, idx)
+ switch subPhaseNum {
+ case phaseMin:
+ nodeState.ReceivedMin = min(nodeState.ReceivedMin, message.Info)
+ case phaseClaim:
+ if message.MessageType == JoiningMIS {
+ nodeState.NeighborInMIS = true
+ }
+ case phaseDecline:
+ if message.MessageType == NeverJoiningMIS {
+ nodeState.AliveNeighbors[idx] = false
+ }
+ }
+}
+
+func process(node lib.Node, subPhaseNum int) bool {
+ state := getState(node)
+ for i := range node.GetInNeighbors() {
+ if state.AliveNeighbors[i] {
+ phasedSend(node, &state, subPhaseNum, i)
+ }
+
+ }
+ for i := range node.GetInNeighbors() {
+ if state.AliveNeighbors[i] {
+ phasedReceive(node, &state, subPhaseNum, i)
+ }
+
+ }
+ if subPhaseNum == phaseMin {
+ state.InMIS = state.RandomMin < state.ReceivedMin
+ }
+ if subPhaseNum == phaseDecline {
+ state.InGraph = !state.InMIS && !state.NeighborInMIS
+ }
+ setState(node, state)
+ return !state.InGraph
+}
+
+func uniform() uint64 {
+ return rand.Uint64()
+}
+
+func check(nodes []lib.Node) {
+ for i, node := range nodes {
+ state := getState(node)
+ if state.InMIS != true && state.InMIS != false {
+ panic(fmt.Sprint("Incorrect location at ", i))
+ }
+
+ outCnt := 0
+ inCnt := 0
+
+ for j, neighbor := range node.GetInNeighbors() {
+ if getState(neighbor).InMIS == true {
+ inCnt++
+ } else {
+ log.Println(j, " neighbor of ", i)
+ outCnt++
+ }
+ }
+
+ if state.InMIS == true && inCnt > 0 {
+ panic(fmt.Sprint(i, " and his neighbor ended up in MIS"))
+ }
+ if state.InMIS == false && inCnt == 0 {
+ panic(fmt.Sprint(i, " and his neighbor did not end up in MIS"))
+ }
+ }
+}
+
+func Run(n int, p float64) (int, int) {
+ nodes, synchronizer := lib.BuildSynchronizedRandomGraph(n, p)
+ for _, node := range nodes {
+ go run(node)
+ }
+ synchronizer.Synchronize(0)
+ check(nodes)
+ return synchronizer.GetStats()
+}
+
+func run(node lib.Node) {
+ node.StartProcessing()
+ initialize(node)
+ finish := false
+ node.FinishProcessing(finish)
+ for !finish {
+ node.StartProcessing()
+ state := getState(node)
+ state.RandomMin = uniform()
+ state.ReceivedMin = math.MaxUint64
+ setState(node, state)
+ for i := 0; i < subPhases; i++ {
+ finish = finish || process(node, i)
+ }
+ node.FinishProcessing(finish)
+ }
+}
diff --git a/leader/clique/async_loui_matsushita_west_2/async_loui_matsushita_west_2.go b/leader/clique/async_loui_matsushita_west_2/async_loui_matsushita_west_2.go
new file mode 100644
index 0000000..acdf1bd
--- /dev/null
+++ b/leader/clique/async_loui_matsushita_west_2/async_loui_matsushita_west_2.go
@@ -0,0 +1,177 @@
+package async_loui_matsushita_west_2
+
+import (
+ "encoding/json"
+ "fmt"
+ "github.com/krzysztof-turowski/distributed-framework/lib"
+)
+
+const (
+ None = iota
+)
+
+type State struct {
+ Status int
+ Leader int
+ N int
+ D int
+ E int
+}
+
+const (
+ Active = iota
+ Passive
+)
+
+type Message struct {
+ Sender int
+ MessageType int
+ Distance int
+}
+
+const (
+ DeclaringLeader = iota
+ Electing
+)
+
+func setState(node lib.Node, state *State) {
+ encodedState, _ := json.Marshal(*state)
+ node.SetState(encodedState)
+}
+
+func getState(node lib.Node) State {
+ var state State
+ err := json.Unmarshal(node.GetState(), &state)
+ if err != nil {
+ panic(fmt.Sprint("Error while decoding data on node ", node.GetIndex()))
+ return State{}
+ }
+ return state
+}
+
+func prepareMessage(message Message) []byte {
+ encoded, _ := json.Marshal(message)
+ return encoded
+}
+
+func receiveMessage(node lib.Node) Message {
+ var message Message
+ _, receivedMessage := node.ReceiveAnyMessage()
+ err := json.Unmarshal(receivedMessage, &message)
+ if err != nil {
+ panic(fmt.Sprint("Error while decoding data on node ", node.GetIndex()))
+ return Message{}
+ }
+ return message
+}
+
+func sendNullMessage(node lib.Node, i int) {
+ node.SendMessage(i, nil)
+}
+
+func initialize(node lib.Node) {
+ n := node.GetSize() + 1
+ setState(
+ node,
+ &State{
+ Active,
+ None,
+ n,
+ None,
+ None,
+ },
+ )
+
+ id := node.GetIndex()
+ fmt.Println("Node: ", id)
+ node.SendMessage(0, prepareMessage(Message{id, Electing, 1}))
+}
+
+func process(node lib.Node) bool {
+ id := node.GetIndex()
+ message := receiveMessage(node)
+ state := getState(node)
+ defer setState(node, &state)
+
+ switch message.MessageType {
+ case DeclaringLeader:
+ state.Leader = message.Sender
+ case Electing:
+ switch state.Status {
+ case Active:
+ if message.Sender == id {
+ state.Leader = id
+ encoded := prepareMessage(
+ Message{
+ id,
+ DeclaringLeader,
+ state.N - message.Distance,
+ },
+ )
+ for i := 0; i < node.GetOutChannelsCount(); i++ {
+ node.SendMessage(
+ i,
+ encoded,
+ )
+ }
+ } else if message.Sender > id {
+ state.Status = Passive
+ state.E = message.Distance
+ } else {
+ node.SendMessage(
+ state.N-message.Distance-1,
+ prepareMessage(
+ Message{
+ id,
+ Electing,
+ state.N - message.Distance,
+ },
+ ),
+ )
+ }
+ case Passive:
+ node.SendMessage(
+ state.N-state.E-1,
+ prepareMessage(
+ Message{
+ message.Sender,
+ Electing,
+ message.Distance - state.E,
+ },
+ ),
+ )
+ }
+ }
+
+ return state.Leader != None
+}
+
+func checkSingleLeaderElected(nodes []lib.Node) {
+ leader := getState(nodes[0]).Leader
+ if leader == None {
+ panic(fmt.Sprint("Error due to lack of leader"))
+ }
+ for _, node := range nodes {
+ if getState(node).Leader != leader {
+ panic(fmt.Sprint("Error due to leader ambiguity"))
+ }
+ }
+ fmt.Println("Leader: ", leader)
+}
+
+func Run(nodes []lib.Node, runner lib.Runner) (int, int) {
+ for _, node := range nodes {
+ go run(node)
+ }
+ runner.Run(true)
+ checkSingleLeaderElected(nodes)
+ return runner.GetStats()
+}
+
+func run(node lib.Node) {
+ node.StartProcessing()
+ initialize(node)
+ for process(node) {
+ }
+ node.FinishProcessing(true)
+}
diff --git a/test/graphs_mis_test.go b/test/graphs_mis_test.go
index ff64beb..104ba35 100644
--- a/test/graphs_mis_test.go
+++ b/test/graphs_mis_test.go
@@ -1,6 +1,7 @@
package test
import (
+ "github.com/krzysztof-turowski/distributed-framework/graphs/mis/sync_metivier_a"
"io"
"log"
"testing"
@@ -32,3 +33,15 @@ func BenchmarkMetivierC(b *testing.B) {
sync_metivier_c.Run(100, 0.75)
}
}
+
+func TestMetivierA(t *testing.T) {
+ checkLogOutput()
+ sync_metivier_a.Run(100, 0.75)
+}
+
+func BenchmarkMetivierA(b *testing.B) {
+ log.SetOutput(io.Discard)
+ for iteration := 0; iteration < b.N; iteration++ {
+ sync_metivier_a.Run(100, 0.75)
+ }
+}
diff --git a/test/leader_clique_test.go b/test/leader_clique_test.go
index 8fa963d..a8d21c2 100644
--- a/test/leader_clique_test.go
+++ b/test/leader_clique_test.go
@@ -1,6 +1,7 @@
package test
import (
+ "github.com/krzysztof-turowski/distributed-framework/leader/clique/async_loui_matsushita_west_2"
"io/ioutil"
"log"
"math/rand"
@@ -47,6 +48,14 @@ func TestAfekGafniB(t *testing.T) {
}
}
+func TestLouiMatsushitaWest2(t *testing.T) {
+ checkLogOutput()
+ for n := 2; n <= 100; n++ {
+ nodes, runner := lib.BuildCompleteGraph(n)
+ async_loui_matsushita_west_2.Run(nodes, runner)
+ }
+}
+
func BenchmarkHumblet(b *testing.B) {
log.SetOutput(ioutil.Discard)
@@ -86,3 +95,13 @@ func BenchmarkAfekGafniB(b *testing.B) {
async_afek_gafni_b.Run(nodes, runner)
}
}
+
+func BenchmarkLouiMatsushitaWest2(b *testing.B) {
+ log.SetOutput(ioutil.Discard)
+
+ for i := 0; i < b.N; i++ {
+ rand.Seed(0)
+ nodes, runner := lib.BuildCompleteGraph(100)
+ async_loui_matsushita_west_2.Run(nodes, runner)
+ }
+}