From a1d0b0cd9d12f1663126b600b1946f95464ffae2 Mon Sep 17 00:00:00 2001 From: katerieul Date: Sun, 28 Jan 2024 09:46:06 +0100 Subject: [PATCH] added MIS Metivier A and Loui Matsushita West Leader Election, no pdfs --- .idea/.gitignore | 8 + .idea/distributed-framework.iml | 9 + .idea/modules.xml | 8 + .idea/vcs.xml | 6 + Makefile | 2 + example/graphs_mis_sync_metivier_a.go | 15 ++ ...der_clique_async_loui_matsushita_west_2.go | 23 ++ graphs/mis/sync_metivier_a/sync_metivier_a.go | 241 ++++++++++++++++++ .../async_loui_matsushita_west_2.go | 177 +++++++++++++ test/graphs_mis_test.go | 13 + test/leader_clique_test.go | 19 ++ 11 files changed, 521 insertions(+) create mode 100644 .idea/.gitignore create mode 100644 .idea/distributed-framework.iml create mode 100644 .idea/modules.xml create mode 100644 .idea/vcs.xml create mode 100644 example/graphs_mis_sync_metivier_a.go create mode 100644 example/leader_clique_async_loui_matsushita_west_2.go create mode 100644 graphs/mis/sync_metivier_a/sync_metivier_a.go create mode 100644 leader/clique/async_loui_matsushita_west_2/async_loui_matsushita_west_2.go 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) + } +}