From a46afb46a0e967d1d207365e187aac135ae5628b Mon Sep 17 00:00:00 2001 From: giongto35 Date: Sat, 20 Apr 2019 01:03:19 +0800 Subject: [PATCH 01/10] refactor with callback --- main.go | 362 ++++++++++++++++++++++++++++++------------- overlord.go | 28 ++++ overlord/overlord.go | 1 + static/js/ws.js | 12 +- ws.go | 93 +++++++++++ 5 files changed, 391 insertions(+), 105 deletions(-) create mode 100644 overlord.go create mode 100644 overlord/overlord.go create mode 100644 ws.go diff --git a/main.go b/main.go index 07ebb72f..57160d18 100644 --- a/main.go +++ b/main.go @@ -36,15 +36,9 @@ var indexFN = gameboyIndex var readWait = 30 * time.Second var writeWait = 30 * time.Second +var IsOverlord = false var upgrader = websocket.Upgrader{} -type WSPacket struct { - ID string `json:"id"` - Data string `json:"data"` - RoomID string `json:"room_id"` - PlayerIndex int `json:"player_index"` -} - // Room is a game session. multi webRTC sessions can connect to a same game. // A room stores all the channel for interaction between all webRTCs session and emulator type Room struct { @@ -68,6 +62,11 @@ func main() { indexFN = debugIndex fmt.Println("Use debug version") } + if len(os.Args) == 3 { + if os.Args[3] == "overlord" { + IsOverlord = true + } + } rand.Seed(time.Now().UTC().UnixNano()) fmt.Println("http://localhost:8000") @@ -164,116 +163,261 @@ func ws(w http.ResponseWriter, r *http.Request) { } defer c.Close() - log.Println("New ws connection") - webRTC := webrtc.NewWebRTC() - - // streaming game - - // start new games and webrtc stuff? - //isDone := false - var gameName string var roomID string var playerIndex int - for { - c.SetReadDeadline(time.Now().Add(readWait)) - mt, message, err := c.ReadMessage() + client := NewClient(c, webrtc.NewWebRTC()) + //&Client{ + //conn: c, + ////wsOverlord: createOverlordClient(), + //peerconnection: webrtc.NewWebRTC(), + //} + + client.syncReceive("initwebrtc", func(req WSPacket) WSPacket { + log.Println("Received user SDP") + localSession, err := client.peerconnection.StartClient(req.Data, width, height) if err != nil { - log.Println("[!] read:", err) - break + log.Fatalln(err) } - req := WSPacket{} - res := WSPacket{} + return WSPacket{ + ID: "sdp", + Data: localSession, + } + }) - err = json.Unmarshal(message, &req) - if err != nil { - log.Println("[!] json unmarshal:", err) - break + client.syncReceive("save", func(req WSPacket) (res WSPacket) { + log.Println("Saving game state") + res.ID = "save" + res.Data = "ok" + if roomID != "" { + err = rooms[roomID].director.SaveGame() + if err != nil { + log.Println("[!] Cannot save game state: ", err) + res.Data = "error" + } + } else { + res.Data = "error" } - // SDP connection initializations follows WebRTC convention - // https://developer.mozilla.org/en-US/docs/Web/API/WebRTC_API/Protocols - switch req.ID { - //case "ping": - //gameName = req.Data - //roomID = req.RoomID - //playerIndex = req.PlayerIndex + return res + }) + + client.syncReceive("load", func(req WSPacket) (res WSPacket) { + log.Println("Loading game state") + res.ID = "load" + res.Data = "ok" + if roomID != "" { + err = rooms[roomID].director.LoadGame() + if err != nil { + log.Println("[!] Cannot load game state: ", err) + res.Data = "error" + } + } else { + res.Data = "error" + } + + return res + }) + + client.syncReceive("start", func(req WSPacket) (res WSPacket) { + gameName = req.Data + roomID = req.RoomID + playerIndex = req.PlayerIndex //log.Println("Ping from server with game:", gameName) //res.ID = "pong" + log.Println("Starting game") + roomID = startSession(client.peerconnection, gameName, roomID, playerIndex) + res.ID = "start" + res.RoomID = roomID - case "initwebrtc": - log.Println("Received user SDP") - localSession, err := webRTC.StartClient(req.Data, width, height) - if err != nil { - log.Fatalln(err) - } + return res + }) - res.ID = "sdp" - res.Data = localSession - - case "candidate": - // Unuse code - hi := pionRTC.ICECandidateInit{} - err = json.Unmarshal([]byte(req.Data), &hi) - if err != nil { - log.Println("[!] Cannot parse candidate: ", err) - } else { - // webRTC.AddCandidate(hi) - } - res.ID = "candidate" - - case "start": - gameName = req.Data - roomID = req.RoomID - playerIndex = req.PlayerIndex - //log.Println("Ping from server with game:", gameName) - //res.ID = "pong" - log.Println("Starting game") - roomID = startSession(webRTC, gameName, roomID, playerIndex) - res.ID = "start" - res.RoomID = roomID - - case "save": - log.Println("Saving game state") - res.ID = "save" - res.Data = "ok" - if roomID != "" { - err = rooms[roomID].director.SaveGame() - if err != nil { - log.Println("[!] Cannot save game state: ", err) - res.Data = "error" - } - } else { - res.Data = "error" - } - - case "load": - log.Println("Loading game state") - res.ID = "load" - res.Data = "ok" - if roomID != "" { - err = rooms[roomID].director.LoadGame() - if err != nil { - log.Println("[!] Cannot load game state: ", err) - res.Data = "error" - } - } else { - res.Data = "error" - } - } - - stRes, err := json.Marshal(res) + client.syncReceive("candidate", func(req WSPacket) (res WSPacket) { + // Unuse code + hi := pionRTC.ICECandidateInit{} + err = json.Unmarshal([]byte(req.Data), &hi) if err != nil { - log.Println("json marshal:", err) + log.Println("[!] Cannot parse candidate: ", err) + } else { + // webRTC.AddCandidate(hi) } + res.ID = "candidate" - c.SetWriteDeadline(time.Now().Add(writeWait)) - err = c.WriteMessage(mt, []byte(stRes)) - } + return res + }) + + client.listen() } +//func wsold(w http.ResponseWriter, r *http.Request) { +//c, err := upgrader.Upgrade(w, r, nil) +//if err != nil { +//log.Print("[!] WS upgrade:", err) +//return +//} +//defer c.Close() + +//client := &Client{ +//ws: c, +//wsOverlord: createOverlordClient(), +//peerconnection: webrtc.NewWebRTC(), +//} + +//if IsOverlord { + +//} + +//log.Println("New ws connection") + +//// streaming game + +//// start new games and webrtc stuff? +////isDone := false + +//var gameName string +//var roomID string +//var playerIndex int + +//for { +//c.SetReadDeadline(time.Now().Add(readWait)) +//mt, message, err := c.ReadMessage() +//if err != nil { +//log.Println("[!] read:", err) +//break +//} + +//req := WSPacket{} +//res := WSPacket{} + +//err = json.Unmarshal(message, &req) +//if err != nil { +//log.Println("[!] json unmarshal:", err) +//break +//} + +//// SDP connection initializations follows WebRTC convention +//// https://developer.mozilla.org/en-US/docs/Web/API/WebRTC_API/Protocols +//switch req.ID { +////case "ping": +////gameName = req.Data +////roomID = req.RoomID +////playerIndex = req.PlayerIndex +////log.Println("Ping from server with game:", gameName) +////res.ID = "pong" + +//case "initwebrtc": +//log.Println("Received user SDP") +//localSession, err := client.peerconnection.StartClient(req.Data, width, height) +//if err != nil { +//log.Fatalln(err) +//} + +//res.ID = "sdp" +//res.Data = localSession + +//// Master received request Offer +////case "requestOffer": +//// Master ask the host connecting to master to request offer from browser + +////case "receivedRemoteSDP": +////log.Println("Received remote rtc") +//// Relay message to overlord, so overlord will send to the correct host +//// TODO: Include room to data +////res.id = "overlordrelaysdp" +////res.Data = req.Data + +////case "receivedRelaySDPFromOverlord": +////// initwebrtc +////log.Println("Received user SDP") +////localSession, err := client.peerconnection.StartClient(req.Data, width, height) +////if err != nil { +////log.Fatalln(err) +////} + +////res.ID = "sdp" +////res.Data = localSession + +//case "candidate": +//// Unuse code +//hi := pionRTC.ICECandidateInit{} +//err = json.Unmarshal([]byte(req.Data), &hi) +//if err != nil { +//log.Println("[!] Cannot parse candidate: ", err) +//} else { +//// webRTC.AddCandidate(hi) +//} +//res.ID = "candidate" + +//case "remoteStart": +//gameName = req.Data +//roomID = req.RoomID +//log.Println("Received Remote start", gameName, roomID) + +//roomID = startSession(client.peerconnection, gameName, roomID, playerIndex) +//res.ID = "start" +//res.RoomID = roomID + +//case "start": +//gameName = req.Data +//roomID = req.RoomID +//playerIndex = req.PlayerIndex +////log.Println("Ping from server with game:", gameName) +////res.ID = "pong" +//log.Println("Starting game") +//if overlord.isRemoteRoom(roomID) { +//res.ID = "overlordRequestOffer" +//} else { +//roomID = startSession(client.peerconnnection, gameName, roomID, playerIndex) +//} +//res.ID = "start" +//res.RoomID = roomID + +//case "save": +//log.Println("Saving game state") +//res.ID = "save" +//res.Data = "ok" +//if roomID != "" { +//err = rooms[roomID].director.SaveGame() +//if err != nil { +//log.Println("[!] Cannot save game state: ", err) +//res.Data = "error" +//} +//} else { +//res.Data = "error" +//} + +//case "load": +//log.Println("Loading game state") +//res.ID = "load" +//res.Data = "ok" +//if roomID != "" { +//err = rooms[roomID].director.LoadGame() +//if err != nil { +//log.Println("[!] Cannot load game state: ", err) +//res.Data = "error" +//} +//} else { +//res.Data = "error" +//} +//} + +//stRes, err := json.Marshal(res) +//if err != nil { +//log.Println("json marshal:", err) +//} + +//c.SetWriteDeadline(time.Now().Add(writeWait)) +//if strings.HasPrefix(res.ID, "overlord") { +//overlord.ws.WriteMessage(res.Data) +//} +//c.WriteMessage(res.Data) +//err = c.WriteMessage(mt, []byte(stRes)) +//} +//} + // generateRoomID generate a unique room ID containing 16 digits func generateRoomID() string { roomID := strconv.FormatInt(rand.Int63(), 16) @@ -324,7 +468,7 @@ func startWebRTCSession(room *Room, webRTC *webrtc.WebRTC, playerIndex int) { select { case <-webRTC.Done: fmt.Println("One session closed") - removeSession(room, webRTC) + removeSession(webRTC, room) default: } // Client stopped @@ -344,19 +488,19 @@ func startWebRTCSession(room *Room, webRTC *webrtc.WebRTC, playerIndex int) { } } -func cleanSession(webrtc *webrtc.WebRTC) { - room, ok := rooms[webrtc.RoomID] +func cleanSession(w *webrtc.WebRTC) { + room, ok := rooms[w.RoomID] if !ok { return } - removeSession(room, webrtc) + removeSession(w, room) } -func removeSession(room *Room, webrtc *webrtc.WebRTC) { +func removeSession(w *webrtc.WebRTC, room *Room) { room.sessionsLock.Lock() defer room.sessionsLock.Unlock() for i, s := range room.rtcSessions { - if s == webrtc { + if s == w { room.rtcSessions = append(room.rtcSessions[:i], room.rtcSessions[i+1:]...) break } @@ -366,3 +510,13 @@ func removeSession(room *Room, webrtc *webrtc.WebRTC) { room.Done <- struct{}{} } } + +//func (o *Overlord) isRemoteRoom(roomID string) bool { +//err := c.WriteMessage(websocket.TextMessage, []byte(stRes)) +//o.ws.WriteMessage() +//_, message, err := o.ws.ReadMessage() +//if message == "isRemoteRoom" { +//return true +//} +//return false +//} diff --git a/overlord.go b/overlord.go new file mode 100644 index 00000000..72b7f2a3 --- /dev/null +++ b/overlord.go @@ -0,0 +1,28 @@ +package main + +import ( + "github.com/gorilla/websocket" +) + +const overlordHost = "http://localhost:9000" + +type Overlord struct { + ws websocket.Conn +} + +//func createOverlordClient() websocket.Conn { +//signal.Notify(interrupt, os.Interrupt) + +////u := url.URL{Scheme: "ws", Host: *addr, Path: "/echo"} +////log.Printf("connecting to %s", u.String()) + +//c, _, err := websocket.DefaultDialer.Dial(overlordHost, nil) +//if err != nil { +//log.Fatal("dial:", err) +//} +//overlord := &Overlord{ +//ws: c, +//} + +//return overlord +//} diff --git a/overlord/overlord.go b/overlord/overlord.go new file mode 100644 index 00000000..8b137891 --- /dev/null +++ b/overlord/overlord.go @@ -0,0 +1 @@ + diff --git a/static/js/ws.js b/static/js/ws.js index afcf04c3..b50cd19e 100644 --- a/static/js/ws.js +++ b/static/js/ws.js @@ -27,6 +27,16 @@ conn.onmessage = e => { log("Got remote sdp"); pc.setRemoteDescription(new RTCSessionDescription(JSON.parse(atob(d["data"])))); break; + //case "requestOffer": + //pc.createOffer({offerToReceiveVideo: true, offerToReceiveAudio: false}).then(d => { + //pc.setLocalDescription(d).catch(log); + //}) + + //case "sdpremote": + //log("Got remote sdp"); + //pc.setRemoteDescription(new RTCSessionDescription(JSON.parse(atob(d["data"])))); + //conn.send(JSON.stringify({"id": "remotestart", "data": GAME_LIST[gameIdx].nes, "room_id": roomID.value, "player_index": parseInt(playerIndex.value, 10)}));inputTimer + //break; case "pong": // TODO: Change name use one session log("Recv pong. Start webrtc"); @@ -51,7 +61,7 @@ conn.onmessage = e => { function sendPing() { // TODO: format the package with time - conn.send(JSON.stringify({"id": "pingpong", "data": "pingpong"})); + //conn.send(JSON.stringify({"id": "pingpong", "data": "pingpong"})); } function startWebRTC() { diff --git a/ws.go b/ws.go new file mode 100644 index 00000000..e971c410 --- /dev/null +++ b/ws.go @@ -0,0 +1,93 @@ +package main + +import ( + "encoding/json" + "log" + "time" + + "github.com/giongto35/cloud-game/webrtc" + "github.com/gorilla/websocket" +) + +type Client struct { + conn *websocket.Conn + wsoverlord *websocket.Conn + peerconnection *webrtc.WebRTC + + // sendCallback is callback based on packetID + sendCallback map[string]func(req WSPacket) + // recvCallback is callback when receive based on ID of the packet + recvCallback map[string]func(req WSPacket) +} + +type WSPacket struct { + ID string `json:"id"` + Data string `json:"data"` + + RoomID string `json:"room_id"` + PlayerIndex int `json:"player_index"` + + TargetHostID string `json:"target_id"` + PacketID string +} + +func NewClient(conn *websocket.Conn, webrtc *webrtc.WebRTC) *Client { + sendCallback := map[string]func(WSPacket){} + recvCallback := map[string]func(WSPacket){} + return &Client{ + conn: conn, + peerconnection: webrtc, + sendCallback: sendCallback, + recvCallback: recvCallback, + } +} + +// syncSend sends a packet and trigger callback when the packet comes back +func (c *Client) syncSend(packet WSPacket, callback func(msg WSPacket)) { + data, err := json.Marshal(packet) + if err != nil { + return + } + + c.conn.WriteMessage(0, data) + c.sendCallback[packet.PacketID] = callback +} + +// syncReceive receive and response back +func (c *Client) syncReceive(id string, f func(request WSPacket) (response WSPacket)) { + c.recvCallback[id] = func(request WSPacket) { + packet := f(request) + + resp, err := json.Marshal(packet) + if err != nil { + log.Println("[!] json marshal error:", err) + } + c.conn.SetWriteDeadline(time.Now().Add(writeWait)) + c.conn.WriteMessage(websocket.TextMessage, resp) + } +} + +func (c *Client) listen() { + for { + _, rawMsg, err := c.conn.ReadMessage() + if err != nil { + log.Println("[!] read:", err) + break + } + wspacket := WSPacket{} + err = json.Unmarshal(rawMsg, &wspacket) + if err != nil { + continue + } + + // Check if some async send is waiting for the response based on packetID + if callback, ok := c.sendCallback[wspacket.PacketID]; ok { + callback(wspacket) + delete(c.sendCallback, wspacket.PacketID) + } + // Check if some receiver with the ID is registered + if callback, ok := c.recvCallback[wspacket.ID]; ok { + callback(wspacket) + } + } +} From 6457c0ba335e84db91dd178348236971385c2c2d Mon Sep 17 00:00:00 2001 From: giongto35 Date: Sat, 20 Apr 2019 01:04:26 +0800 Subject: [PATCH 02/10] Clean old code --- main.go | 171 -------------------------------------------------------- 1 file changed, 171 deletions(-) diff --git a/main.go b/main.go index 57160d18..29b0c8bc 100644 --- a/main.go +++ b/main.go @@ -168,11 +168,6 @@ func ws(w http.ResponseWriter, r *http.Request) { var playerIndex int client := NewClient(c, webrtc.NewWebRTC()) - //&Client{ - //conn: c, - ////wsOverlord: createOverlordClient(), - //peerconnection: webrtc.NewWebRTC(), - //} client.syncReceive("initwebrtc", func(req WSPacket) WSPacket { log.Println("Received user SDP") @@ -252,172 +247,6 @@ func ws(w http.ResponseWriter, r *http.Request) { client.listen() } -//func wsold(w http.ResponseWriter, r *http.Request) { -//c, err := upgrader.Upgrade(w, r, nil) -//if err != nil { -//log.Print("[!] WS upgrade:", err) -//return -//} -//defer c.Close() - -//client := &Client{ -//ws: c, -//wsOverlord: createOverlordClient(), -//peerconnection: webrtc.NewWebRTC(), -//} - -//if IsOverlord { - -//} - -//log.Println("New ws connection") - -//// streaming game - -//// start new games and webrtc stuff? -////isDone := false - -//var gameName string -//var roomID string -//var playerIndex int - -//for { -//c.SetReadDeadline(time.Now().Add(readWait)) -//mt, message, err := c.ReadMessage() -//if err != nil { -//log.Println("[!] read:", err) -//break -//} - -//req := WSPacket{} -//res := WSPacket{} - -//err = json.Unmarshal(message, &req) -//if err != nil { -//log.Println("[!] json unmarshal:", err) -//break -//} - -//// SDP connection initializations follows WebRTC convention -//// https://developer.mozilla.org/en-US/docs/Web/API/WebRTC_API/Protocols -//switch req.ID { -////case "ping": -////gameName = req.Data -////roomID = req.RoomID -////playerIndex = req.PlayerIndex -////log.Println("Ping from server with game:", gameName) -////res.ID = "pong" - -//case "initwebrtc": -//log.Println("Received user SDP") -//localSession, err := client.peerconnection.StartClient(req.Data, width, height) -//if err != nil { -//log.Fatalln(err) -//} - -//res.ID = "sdp" -//res.Data = localSession - -//// Master received request Offer -////case "requestOffer": -//// Master ask the host connecting to master to request offer from browser - -////case "receivedRemoteSDP": -////log.Println("Received remote rtc") -//// Relay message to overlord, so overlord will send to the correct host -//// TODO: Include room to data -////res.id = "overlordrelaysdp" -////res.Data = req.Data - -////case "receivedRelaySDPFromOverlord": -////// initwebrtc -////log.Println("Received user SDP") -////localSession, err := client.peerconnection.StartClient(req.Data, width, height) -////if err != nil { -////log.Fatalln(err) -////} - -////res.ID = "sdp" -////res.Data = localSession - -//case "candidate": -//// Unuse code -//hi := pionRTC.ICECandidateInit{} -//err = json.Unmarshal([]byte(req.Data), &hi) -//if err != nil { -//log.Println("[!] Cannot parse candidate: ", err) -//} else { -//// webRTC.AddCandidate(hi) -//} -//res.ID = "candidate" - -//case "remoteStart": -//gameName = req.Data -//roomID = req.RoomID -//log.Println("Received Remote start", gameName, roomID) - -//roomID = startSession(client.peerconnection, gameName, roomID, playerIndex) -//res.ID = "start" -//res.RoomID = roomID - -//case "start": -//gameName = req.Data -//roomID = req.RoomID -//playerIndex = req.PlayerIndex -////log.Println("Ping from server with game:", gameName) -////res.ID = "pong" -//log.Println("Starting game") -//if overlord.isRemoteRoom(roomID) { -//res.ID = "overlordRequestOffer" -//} else { -//roomID = startSession(client.peerconnnection, gameName, roomID, playerIndex) -//} -//res.ID = "start" -//res.RoomID = roomID - -//case "save": -//log.Println("Saving game state") -//res.ID = "save" -//res.Data = "ok" -//if roomID != "" { -//err = rooms[roomID].director.SaveGame() -//if err != nil { -//log.Println("[!] Cannot save game state: ", err) -//res.Data = "error" -//} -//} else { -//res.Data = "error" -//} - -//case "load": -//log.Println("Loading game state") -//res.ID = "load" -//res.Data = "ok" -//if roomID != "" { -//err = rooms[roomID].director.LoadGame() -//if err != nil { -//log.Println("[!] Cannot load game state: ", err) -//res.Data = "error" -//} -//} else { -//res.Data = "error" -//} -//} - -//stRes, err := json.Marshal(res) -//if err != nil { -//log.Println("json marshal:", err) -//} - -//c.SetWriteDeadline(time.Now().Add(writeWait)) -//if strings.HasPrefix(res.ID, "overlord") { -//overlord.ws.WriteMessage(res.Data) -//} -//c.WriteMessage(res.Data) -//err = c.WriteMessage(mt, []byte(stRes)) -//} -//} - // generateRoomID generate a unique room ID containing 16 digits func generateRoomID() string { roomID := strconv.FormatInt(rand.Int63(), 16) From 855c43cea6f076324aa7f6c14c5c59f6a48a9163 Mon Sep 17 00:00:00 2001 From: giongto35 Date: Sat, 20 Apr 2019 02:05:50 +0800 Subject: [PATCH 03/10] Setup overlord --- main.go | 68 +++++++++++++++++++++++++++++++++++++++++++++++++---- overlord.go | 28 ---------------------- ws.go | 9 +++---- 3 files changed, 68 insertions(+), 37 deletions(-) delete mode 100644 overlord.go diff --git a/main.go b/main.go index 29b0c8bc..71f830b0 100644 --- a/main.go +++ b/main.go @@ -63,22 +63,31 @@ func main() { fmt.Println("Use debug version") } if len(os.Args) == 3 { - if os.Args[3] == "overlord" { + if os.Args[2] == "overlord" { IsOverlord = true } + fmt.Println("Running as overlord ") } rand.Seed(time.Now().UTC().UnixNano()) - fmt.Println("http://localhost:8000") rooms = map[string]*Room{} // ignore origin upgrader.CheckOrigin = func(r *http.Request) bool { return true } - http.HandleFunc("/ws", ws) http.HandleFunc("/", getWeb) http.Handle("/static/", http.StripPrefix("/static/", http.FileServer(http.Dir("./static")))) - http.ListenAndServe(":8000", nil) + http.HandleFunc("/ws", ws) + + if !IsOverlord { + fmt.Println("http://localhost:8000") + http.ListenAndServe(":8000", nil) + } else { + fmt.Println("http://localhost:9000") + // Overlord expose one more path for handle overlord connections + http.HandleFunc("/wso", wso) + http.ListenAndServe(":9000", nil) + } } func getWeb(w http.ResponseWriter, r *http.Request) { @@ -155,7 +164,9 @@ func startSession(webRTC *webrtc.WebRTC, gameName string, roomID string, playerI return roomID } -func ws(w http.ResponseWriter, r *http.Request) { +// If it's overlord, handle overlord connection (from host to overlord) +func wso(w http.ResponseWriter, r *http.Request) { + fmt.Println("Connected") c, err := upgrader.Upgrade(w, r, nil) if err != nil { log.Print("[!] WS upgrade:", err) @@ -163,6 +174,37 @@ func ws(w http.ResponseWriter, r *http.Request) { } defer c.Close() + client := NewClient(c, webrtc.NewWebRTC()) + + client.syncReceive("ping", func(req WSPacket) WSPacket { + log.Println("received Ping, sending Pong") + return WSPacket{ + ID: "pong", + } + }) + client.listen() +} + +const overlordHost = "ws://localhost:9000/wso" + +func createOverlordClient() (*websocket.Conn, error) { + c, _, err := websocket.DefaultDialer.Dial(overlordHost, nil) + if err != nil { + log.Fatal("dial:", err) + return nil, err + } + + return c, nil +} + +// Handle normal traffic (from browser to host) +func ws(w http.ResponseWriter, r *http.Request) { + c, err := upgrader.Upgrade(w, r, nil) + if err != nil { + log.Print("[!] WS upgrade:", err) + return + } + defer c.Close() var gameName string var roomID string var playerIndex int @@ -244,6 +286,22 @@ func ws(w http.ResponseWriter, r *http.Request) { return res }) + // Create connection to overlord + if !IsOverlord { + oc, err := createOverlordClient() + if err != nil { + log.Println("Cannot connect to overlord") + } + oclient := NewClient(oc, webrtc.NewWebRTC()) + oclient.syncSend(WSPacket{ + ID: "ping", + }, + func(resp WSPacket) { + log.Println("pong") + }, + ) + } + client.listen() } diff --git a/overlord.go b/overlord.go deleted file mode 100644 index 72b7f2a3..00000000 --- a/overlord.go +++ /dev/null @@ -1,28 +0,0 @@ -package main - -import ( - "github.com/gorilla/websocket" -) - -const overlordHost = "http://localhost:9000" - -type Overlord struct { - ws websocket.Conn -} - -//func createOverlordClient() websocket.Conn { -//signal.Notify(interrupt, os.Interrupt) - -////u := url.URL{Scheme: "ws", Host: *addr, Path: "/echo"} -////log.Printf("connecting to %s", u.String()) - -//c, _, err := websocket.DefaultDialer.Dial(overlordHost, nil) -//if err != nil { -//log.Fatal("dial:", err) -//} -//overlord := &Overlord{ -//ws: c, -//} - -//return overlord -//} diff --git a/ws.go b/ws.go index e971c410..a529047e 100644 --- a/ws.go +++ b/ws.go @@ -10,8 +10,8 @@ import ( ) type Client struct { - conn *websocket.Conn - wsoverlord *websocket.Conn + conn *websocket.Conn + peerconnection *webrtc.WebRTC // sendCallback is callback based on packetID @@ -35,7 +35,8 @@ func NewClient(conn *websocket.Conn, webrtc *webrtc.WebRTC) *Client { sendCallback := map[string]func(WSPacket){} recvCallback := map[string]func(WSPacket){} return &Client{ - conn: conn, + conn: conn, + peerconnection: webrtc, sendCallback: sendCallback, recvCallback: recvCallback, @@ -49,7 +50,7 @@ func (c *Client) syncSend(packet WSPacket, callback func(msg WSPacket)) { return } - c.conn.WriteMessage(0, data) + c.conn.WriteMessage(websocket.TextMessage, data) c.sendCallback[packet.PacketID] = callback } From 623a2adb1a13239cd86367292ba7f373c1e30fdb Mon Sep 17 00:00:00 2001 From: giongto35 Date: Sat, 20 Apr 2019 04:24:13 +0800 Subject: [PATCH 04/10] Register server IP --- main.go | 190 +++++++++++++++++++++++++++++----------------------- overlord.go | 57 ++++++++++++++++ ws.go | 35 ++++++++-- 3 files changed, 191 insertions(+), 91 deletions(-) create mode 100644 overlord.go diff --git a/main.go b/main.go index 71f830b0..96fff91b 100644 --- a/main.go +++ b/main.go @@ -145,13 +145,14 @@ func isRoomRunning(roomID string) bool { } // startSession handles one session call -func startSession(webRTC *webrtc.WebRTC, gameName string, roomID string, playerIndex int) string { +func startSession(webRTC *webrtc.WebRTC, gameName string, roomID string, playerIndex int) (rRoomID string, isNewRoom bool) { cleanSession(webRTC) // If the roomID is empty, // or the roomID doesn't have any running sessions (room was closed) // we spawn a new room if roomID == "" || !isRoomRunning(roomID) { roomID = initRoom(roomID, gameName) + isNewRoom = true } // TODO: Might have race condition @@ -161,40 +162,14 @@ func startSession(webRTC *webrtc.WebRTC, gameName string, roomID string, playerI webRTC.AttachRoomID(roomID) go startWebRTCSession(room, webRTC, playerIndex) - return roomID + return roomID, false } -// If it's overlord, handle overlord connection (from host to overlord) -func wso(w http.ResponseWriter, r *http.Request) { - fmt.Println("Connected") - c, err := upgrader.Upgrade(w, r, nil) - if err != nil { - log.Print("[!] WS upgrade:", err) - return - } - defer c.Close() - - client := NewClient(c, webrtc.NewWebRTC()) - - client.syncReceive("ping", func(req WSPacket) WSPacket { - log.Println("received Ping, sending Pong") - return WSPacket{ - ID: "pong", - } - }) - client.listen() -} - -const overlordHost = "ws://localhost:9000/wso" - -func createOverlordClient() (*websocket.Conn, error) { - c, _, err := websocket.DefaultDialer.Dial(overlordHost, nil) - if err != nil { - log.Fatal("dial:", err) - return nil, err - } - - return c, nil +// Session represents a session connected from the browser to the current server +type Session struct { + client *Client + oclient *Client + ServerID string } // Handle normal traffic (from browser to host) @@ -209,11 +184,22 @@ func ws(w http.ResponseWriter, r *http.Request) { var roomID string var playerIndex int + var oclient *Client + // Create connection to overlord client := NewClient(c, webrtc.NewWebRTC()) - client.syncReceive("initwebrtc", func(req WSPacket) WSPacket { + wssession := Session{ + client: client, + // The server session is maintaining + } + + if !IsOverlord { + wssession.NewOverlordClient() + } + + client.syncReceive("initwebrtc", func(resp WSPacket) WSPacket { log.Println("Received user SDP") - localSession, err := client.peerconnection.StartClient(req.Data, width, height) + localSession, err := client.peerconnection.StartClient(resp.Data, width, height) if err != nil { log.Fatalln(err) } @@ -224,84 +210,75 @@ func ws(w http.ResponseWriter, r *http.Request) { } }) - client.syncReceive("save", func(req WSPacket) (res WSPacket) { + client.syncReceive("save", func(resp WSPacket) (req WSPacket) { log.Println("Saving game state") - res.ID = "save" - res.Data = "ok" + req.ID = "save" + req.Data = "ok" if roomID != "" { err = rooms[roomID].director.SaveGame() if err != nil { log.Println("[!] Cannot save game state: ", err) - res.Data = "error" + req.Data = "error" } } else { - res.Data = "error" + req.Data = "error" } - return res + return req }) - client.syncReceive("load", func(req WSPacket) (res WSPacket) { + client.syncReceive("load", func(resp WSPacket) (req WSPacket) { log.Println("Loading game state") - res.ID = "load" - res.Data = "ok" + req.ID = "load" + req.Data = "ok" if roomID != "" { err = rooms[roomID].director.LoadGame() if err != nil { log.Println("[!] Cannot load game state: ", err) - res.Data = "error" + req.Data = "error" } } else { - res.Data = "error" + req.Data = "error" } - return res + return req }) - client.syncReceive("start", func(req WSPacket) (res WSPacket) { - gameName = req.Data - roomID = req.RoomID - playerIndex = req.PlayerIndex + client.syncReceive("start", func(resp WSPacket) (req WSPacket) { + gameName = resp.Data + roomID = resp.RoomID + playerIndex = resp.PlayerIndex + isNewRoom := false //log.Println("Ping from server with game:", gameName) //res.ID = "pong" log.Println("Starting game") - roomID = startSession(client.peerconnection, gameName, roomID, playerIndex) - res.ID = "start" - res.RoomID = roomID + roomID, isNewRoom = startSession(client.peerconnection, gameName, roomID, playerIndex) + if isNewRoom { + oclient.send(WSPacket{ + ID: "RegisterRoom", + Data: roomID, + }) + } + req.ID = "start" + req.RoomID = roomID - return res + return req }) - client.syncReceive("candidate", func(req WSPacket) (res WSPacket) { + client.syncReceive("candidate", func(resp WSPacket) (req WSPacket) { // Unuse code hi := pionRTC.ICECandidateInit{} - err = json.Unmarshal([]byte(req.Data), &hi) + err = json.Unmarshal([]byte(resp.Data), &hi) if err != nil { log.Println("[!] Cannot parse candidate: ", err) } else { // webRTC.AddCandidate(hi) } - res.ID = "candidate" + req.ID = "candidate" - return res + return req }) - // Create connection to overlord - if !IsOverlord { - oc, err := createOverlordClient() - if err != nil { - log.Println("Cannot connect to overlord") - } - oclient := NewClient(oc, webrtc.NewWebRTC()) - oclient.syncSend(WSPacket{ - ID: "ping", - }, - func(resp WSPacket) { - log.Println("pong") - }, - ) - } - client.listen() } @@ -398,12 +375,57 @@ func removeSession(w *webrtc.WebRTC, room *Room) { } } -//func (o *Overlord) isRemoteRoom(roomID string) bool { -//err := c.WriteMessage(websocket.TextMessage, []byte(stRes)) -//o.ws.WriteMessage() -//_, message, err := o.ws.ReadMessage() -//if message == "isRemoteRoom" { -//return true -//} -//return false -//} +func GetServerIDOfRoom(oc Client, roomID string) chan string { + res := make(chan string) + + oc.syncSend(WSPacket{ + ID: "getRoom", + }, func(resp WSPacket) { + res <- resp.Data + }) + return res +} + +const overlordHost = "ws://localhost:9000/wso" + +func createOverlordConnection() (*websocket.Conn, error) { + c, _, err := websocket.DefaultDialer.Dial(overlordHost, nil) + if err != nil { + log.Fatal("dial:", err) + return nil, err + } + + return c, nil +} + +func (s *Session) NewOverlordClient() { + oc, err := createOverlordConnection() + if err != nil { + log.Println("Cannot connect to overlord") + } + oclient := NewClient(oc, webrtc.NewWebRTC()) + oclient.syncSend( + WSPacket{ + ID: "ping", + }, + func(resp WSPacket) { + log.Println("Received pong full flow") + }, + ) + + // Received from overlord the serverID + oclient.syncReceive( + "serverID", + func(response WSPacket) (request WSPacket) { + // Stick session with serverID got from overlord + log.Println("Received serverID ", response.Data) + s.ServerID = response.Data + + return EmptyPacket + }, + ) + + go oclient.listen() + + return +} diff --git a/overlord.go b/overlord.go new file mode 100644 index 00000000..8ab29130 --- /dev/null +++ b/overlord.go @@ -0,0 +1,57 @@ +package main + +import ( + "fmt" + "log" + "math/rand" + "net/http" + "strconv" + + "github.com/giongto35/cloud-game/webrtc" +) + +var overlordRooms = map[string]string{} + +// servers are the map serverID to server Client +var servers = map[string]Client{} + +// If it's overlord, handle overlord connection (from host to overlord) +func wso(w http.ResponseWriter, r *http.Request) { + fmt.Println("Connected") + c, err := upgrader.Upgrade(w, r, nil) + if err != nil { + log.Print("[!] WS upgrade:", err) + return + } + defer c.Close() + + // register new server + serverID := strconv.Itoa(rand.Int()) + log.Println("A new server connected ", serverID) + + client := NewClient(c, webrtc.NewWebRTC()) + + client.send( + WSPacket{ + ID: "serverID", + Data: serverID, + }, + ) + + client.syncReceive("ping", func(resp WSPacket) WSPacket { + log.Println("received Ping, sending Pong") + return WSPacket{ + ID: "pong", + } + }) + + client.syncReceive("registerRoom", func(resp WSPacket) WSPacket { + log.Println("received registerRoom") + overlordRooms[resp.Data] = serverID + return WSPacket{ + ID: "registerRoom", + } + }) + + client.listen() +} diff --git a/ws.go b/ws.go index a529047e..57898684 100644 --- a/ws.go +++ b/ws.go @@ -3,6 +3,8 @@ package main import ( "encoding/json" "log" + "math/rand" + "strconv" "time" "github.com/giongto35/cloud-game/webrtc" @@ -31,6 +33,8 @@ type WSPacket struct { PacketID string } +var EmptyPacket = WSPacket{} + func NewClient(conn *websocket.Conn, webrtc *webrtc.WebRTC) *Client { sendCallback := map[string]func(WSPacket){} recvCallback := map[string]func(WSPacket){} @@ -43,22 +47,39 @@ func NewClient(conn *websocket.Conn, webrtc *webrtc.WebRTC) *Client { } } -// syncSend sends a packet and trigger callback when the packet comes back -func (c *Client) syncSend(packet WSPacket, callback func(msg WSPacket)) { - data, err := json.Marshal(packet) +// send sends a normal packet +func (c *Client) send(request WSPacket) { + data, err := json.Marshal(request) if err != nil { return } c.conn.WriteMessage(websocket.TextMessage, data) - c.sendCallback[packet.PacketID] = callback +} + +// syncSend sends a packet and trigger callback when the packet comes back +func (c *Client) syncSend(request WSPacket, callback func(response WSPacket)) { + request.PacketID = strconv.Itoa(rand.Int()) + data, err := json.Marshal(request) + if err != nil { + return + } + + c.conn.WriteMessage(websocket.TextMessage, data) + c.sendCallback[request.PacketID] = callback } // syncReceive receive and response back -func (c *Client) syncReceive(id string, f func(request WSPacket) (response WSPacket)) { - c.recvCallback[id] = func(request WSPacket) { - packet := f(request) +func (c *Client) syncReceive(id string, f func(response WSPacket) (request WSPacket)) { + c.recvCallback[id] = func(response WSPacket) { + packet := f(response) + // Add Meta data + packet.PacketID = response.PacketID + // Skip rqeuest if it is EmptyPacket + if packet == EmptyPacket { + return + } resp, err := json.Marshal(packet) if err != nil { log.Println("[!] json marshal error:", err) From 563d5edac288d1b43a4ebb0a62751d108fecd896 Mon Sep 17 00:00:00 2001 From: giongto35 Date: Sat, 20 Apr 2019 04:47:52 +0800 Subject: [PATCH 05/10] Add register room logic --- main.go | 24 ++++++++++++++++++------ overlord.go | 13 ++++++++++--- 2 files changed, 28 insertions(+), 9 deletions(-) diff --git a/main.go b/main.go index 96fff91b..a5d1e2d5 100644 --- a/main.go +++ b/main.go @@ -146,6 +146,7 @@ func isRoomRunning(roomID string) bool { // startSession handles one session call func startSession(webRTC *webrtc.WebRTC, gameName string, roomID string, playerIndex int) (rRoomID string, isNewRoom bool) { + isNewRoom = false cleanSession(webRTC) // If the roomID is empty, // or the roomID doesn't have any running sessions (room was closed) @@ -162,7 +163,7 @@ func startSession(webRTC *webrtc.WebRTC, gameName string, roomID string, playerI webRTC.AttachRoomID(roomID) go startWebRTCSession(room, webRTC, playerIndex) - return roomID, false + return roomID, isNewRoom } // Session represents a session connected from the browser to the current server @@ -184,7 +185,6 @@ func ws(w http.ResponseWriter, r *http.Request) { var roomID string var playerIndex int - var oclient *Client // Create connection to overlord client := NewClient(c, webrtc.NewWebRTC()) @@ -252,10 +252,17 @@ func ws(w http.ResponseWriter, r *http.Request) { //log.Println("Ping from server with game:", gameName) //res.ID = "pong" log.Println("Starting game") + roomServerID := <-GetServerIDOfRoom(wssession.oclient, roomID) + log.Println("Server of RoomID ", roomID, " is ", roomServerID) + if roomServerID != "" && wssession.ServerID != roomServerID { + // TODO: Re -register + return + } + roomID, isNewRoom = startSession(client.peerconnection, gameName, roomID, playerIndex) if isNewRoom { - oclient.send(WSPacket{ - ID: "RegisterRoom", + wssession.oclient.send(WSPacket{ + ID: "registerRoom", Data: roomID, }) } @@ -375,12 +382,14 @@ func removeSession(w *webrtc.WebRTC, room *Room) { } } -func GetServerIDOfRoom(oc Client, roomID string) chan string { +func GetServerIDOfRoom(oc *Client, roomID string) chan string { res := make(chan string) oc.syncSend(WSPacket{ - ID: "getRoom", + ID: "getRoom", + Data: roomID, }, func(resp WSPacket) { + log.Println("GetRoom", resp.Data) res <- resp.Data }) return res @@ -427,5 +436,8 @@ func (s *Session) NewOverlordClient() { go oclient.listen() + s.oclient = oclient + // TODO: return oclient + return } diff --git a/overlord.go b/overlord.go index 8ab29130..647c9c37 100644 --- a/overlord.go +++ b/overlord.go @@ -10,7 +10,7 @@ import ( "github.com/giongto35/cloud-game/webrtc" ) -var overlordRooms = map[string]string{} +var roomToServer = map[string]string{} // servers are the map serverID to server Client var servers = map[string]Client{} @@ -46,12 +46,19 @@ func wso(w http.ResponseWriter, r *http.Request) { }) client.syncReceive("registerRoom", func(resp WSPacket) WSPacket { - log.Println("received registerRoom") - overlordRooms[resp.Data] = serverID + log.Println("Received registerRoom ", resp.Data, serverID) + roomToServer[resp.Data] = serverID return WSPacket{ ID: "registerRoom", } }) + client.syncReceive("getRoom", func(resp WSPacket) WSPacket { + return WSPacket{ + ID: "getRoom", + Data: roomToServer[resp.Data], + } + }) + client.listen() } From 73a6a0fe61fc95696c0cdb681fce3767b117dbd3 Mon Sep 17 00:00:00 2001 From: giongto35 Date: Sat, 20 Apr 2019 10:34:49 +0800 Subject: [PATCH 06/10] Add SyncSend to send and wait --- main.go | 43 +++++++++++++++++++++++-------------------- overlord.go | 7 ++++--- ws.go | 31 +++++++++++++++++-------------- 3 files changed, 44 insertions(+), 37 deletions(-) diff --git a/main.go b/main.go index a5d1e2d5..45ef571a 100644 --- a/main.go +++ b/main.go @@ -188,7 +188,7 @@ func ws(w http.ResponseWriter, r *http.Request) { // Create connection to overlord client := NewClient(c, webrtc.NewWebRTC()) - wssession := Session{ + wssession := &Session{ client: client, // The server session is maintaining } @@ -197,7 +197,7 @@ func ws(w http.ResponseWriter, r *http.Request) { wssession.NewOverlordClient() } - client.syncReceive("initwebrtc", func(resp WSPacket) WSPacket { + client.receive("initwebrtc", func(resp WSPacket) WSPacket { log.Println("Received user SDP") localSession, err := client.peerconnection.StartClient(resp.Data, width, height) if err != nil { @@ -210,7 +210,7 @@ func ws(w http.ResponseWriter, r *http.Request) { } }) - client.syncReceive("save", func(resp WSPacket) (req WSPacket) { + client.receive("save", func(resp WSPacket) (req WSPacket) { log.Println("Saving game state") req.ID = "save" req.Data = "ok" @@ -227,7 +227,7 @@ func ws(w http.ResponseWriter, r *http.Request) { return req }) - client.syncReceive("load", func(resp WSPacket) (req WSPacket) { + client.receive("load", func(resp WSPacket) (req WSPacket) { log.Println("Loading game state") req.ID = "load" req.Data = "ok" @@ -244,7 +244,7 @@ func ws(w http.ResponseWriter, r *http.Request) { return req }) - client.syncReceive("start", func(resp WSPacket) (req WSPacket) { + client.receive("start", func(resp WSPacket) (req WSPacket) { gameName = resp.Data roomID = resp.RoomID playerIndex = resp.PlayerIndex @@ -252,10 +252,11 @@ func ws(w http.ResponseWriter, r *http.Request) { //log.Println("Ping from server with game:", gameName) //res.ID = "pong" log.Println("Starting game") - roomServerID := <-GetServerIDOfRoom(wssession.oclient, roomID) + roomServerID := GetServerIDOfRoom(wssession.oclient, roomID) log.Println("Server of RoomID ", roomID, " is ", roomServerID) if roomServerID != "" && wssession.ServerID != roomServerID { // TODO: Re -register + bridgeConnection(wssession, roomServerID) return } @@ -264,7 +265,7 @@ func ws(w http.ResponseWriter, r *http.Request) { wssession.oclient.send(WSPacket{ ID: "registerRoom", Data: roomID, - }) + }, nil) } req.ID = "start" req.RoomID = roomID @@ -272,7 +273,7 @@ func ws(w http.ResponseWriter, r *http.Request) { return req }) - client.syncReceive("candidate", func(resp WSPacket) (req WSPacket) { + client.receive("candidate", func(resp WSPacket) (req WSPacket) { // Unuse code hi := pionRTC.ICECandidateInit{} err = json.Unmarshal([]byte(resp.Data), &hi) @@ -382,17 +383,19 @@ func removeSession(w *webrtc.WebRTC, room *Room) { } } -func GetServerIDOfRoom(oc *Client, roomID string) chan string { - res := make(chan string) +func GetServerIDOfRoom(oc *Client, roomID string) string { + packet := oc.syncSend( + WSPacket{ + ID: "getRoom", + Data: roomID, + }, + ) - oc.syncSend(WSPacket{ - ID: "getRoom", - Data: roomID, - }, func(resp WSPacket) { - log.Println("GetRoom", resp.Data) - res <- resp.Data - }) - return res + return packet.Data +} + +func bridgeConnection(session *Session, serverID string) { + // Ask client to init } const overlordHost = "ws://localhost:9000/wso" @@ -413,7 +416,7 @@ func (s *Session) NewOverlordClient() { log.Println("Cannot connect to overlord") } oclient := NewClient(oc, webrtc.NewWebRTC()) - oclient.syncSend( + oclient.send( WSPacket{ ID: "ping", }, @@ -423,7 +426,7 @@ func (s *Session) NewOverlordClient() { ) // Received from overlord the serverID - oclient.syncReceive( + oclient.receive( "serverID", func(response WSPacket) (request WSPacket) { // Stick session with serverID got from overlord diff --git a/overlord.go b/overlord.go index 647c9c37..5878ba8d 100644 --- a/overlord.go +++ b/overlord.go @@ -36,16 +36,17 @@ func wso(w http.ResponseWriter, r *http.Request) { ID: "serverID", Data: serverID, }, + nil, ) - client.syncReceive("ping", func(resp WSPacket) WSPacket { + client.receive("ping", func(resp WSPacket) WSPacket { log.Println("received Ping, sending Pong") return WSPacket{ ID: "pong", } }) - client.syncReceive("registerRoom", func(resp WSPacket) WSPacket { + client.receive("registerRoom", func(resp WSPacket) WSPacket { log.Println("Received registerRoom ", resp.Data, serverID) roomToServer[resp.Data] = serverID return WSPacket{ @@ -53,7 +54,7 @@ func wso(w http.ResponseWriter, r *http.Request) { } }) - client.syncReceive("getRoom", func(resp WSPacket) WSPacket { + client.receive("getRoom", func(resp WSPacket) WSPacket { return WSPacket{ ID: "getRoom", Data: roomToServer[resp.Data], diff --git a/ws.go b/ws.go index 57898684..e52046f0 100644 --- a/ws.go +++ b/ws.go @@ -47,18 +47,8 @@ func NewClient(conn *websocket.Conn, webrtc *webrtc.WebRTC) *Client { } } -// send sends a normal packet -func (c *Client) send(request WSPacket) { - data, err := json.Marshal(request) - if err != nil { - return - } - - c.conn.WriteMessage(websocket.TextMessage, data) -} - -// syncSend sends a packet and trigger callback when the packet comes back -func (c *Client) syncSend(request WSPacket, callback func(response WSPacket)) { +// send sends a packet and trigger callback when the packet comes back +func (c *Client) send(request WSPacket, callback func(response WSPacket)) { request.PacketID = strconv.Itoa(rand.Int()) data, err := json.Marshal(request) if err != nil { @@ -66,11 +56,14 @@ func (c *Client) syncSend(request WSPacket, callback func(response WSPacket)) { } c.conn.WriteMessage(websocket.TextMessage, data) + if callback == nil { + return + } c.sendCallback[request.PacketID] = callback } -// syncReceive receive and response back -func (c *Client) syncReceive(id string, f func(response WSPacket) (request WSPacket)) { +// receive receive and response back +func (c *Client) receive(id string, f func(response WSPacket) (request WSPacket)) { c.recvCallback[id] = func(response WSPacket) { packet := f(response) // Add Meta data @@ -89,6 +82,16 @@ func (c *Client) syncReceive(id string, f func(response WSPacket) (request WSPac } } +// syncSend sends a packet and wait for callback till the packet comes back +func (c *Client) syncSend(request WSPacket) (response WSPacket) { + res := make(chan WSPacket) + f := func(resp WSPacket) { + res <- resp + } + c.send(request, f) + return <-res +} + func (c *Client) listen() { for { _, rawMsg, err := c.conn.ReadMessage() From 80ca5bdd97ecc341cb813717e71713ca44ecdcff Mon Sep 17 00:00:00 2001 From: giongto35 Date: Sat, 20 Apr 2019 15:25:11 +0800 Subject: [PATCH 07/10] SDP Exchanging --- main.go | 95 ++++++++++++++++++++++++++++++++++++++++++------- overlord.go | 68 +++++++++++++++++++++++++++++++++-- static/js/ws.js | 12 ++++--- ws.go | 15 ++++---- 4 files changed, 163 insertions(+), 27 deletions(-) diff --git a/main.go b/main.go index 45ef571a..7718e0f5 100644 --- a/main.go +++ b/main.go @@ -168,9 +168,10 @@ func startSession(webRTC *webrtc.WebRTC, gameName string, roomID string, playerI // Session represents a session connected from the browser to the current server type Session struct { - client *Client - oclient *Client - ServerID string + client *Client + oclient *Client + peerconnection *webrtc.WebRTC + ServerID string } // Handle normal traffic (from browser to host) @@ -186,10 +187,11 @@ func ws(w http.ResponseWriter, r *http.Request) { var playerIndex int // Create connection to overlord - client := NewClient(c, webrtc.NewWebRTC()) + client := NewClient(c) wssession := &Session{ - client: client, + client: client, + peerconnection: webrtc.NewWebRTC(), // The server session is maintaining } @@ -199,7 +201,7 @@ func ws(w http.ResponseWriter, r *http.Request) { client.receive("initwebrtc", func(resp WSPacket) WSPacket { log.Println("Received user SDP") - localSession, err := client.peerconnection.StartClient(resp.Data, width, height) + localSession, err := wssession.peerconnection.StartClient(resp.Data, width, height) if err != nil { log.Fatalln(err) } @@ -252,15 +254,15 @@ func ws(w http.ResponseWriter, r *http.Request) { //log.Println("Ping from server with game:", gameName) //res.ID = "pong" log.Println("Starting game") - roomServerID := GetServerIDOfRoom(wssession.oclient, roomID) + roomServerID := getServerIDOfRoom(wssession.oclient, roomID) log.Println("Server of RoomID ", roomID, " is ", roomServerID) if roomServerID != "" && wssession.ServerID != roomServerID { // TODO: Re -register - bridgeConnection(wssession, roomServerID) + go bridgeConnection(wssession, roomServerID, gameName, roomID, playerIndex) return } - roomID, isNewRoom = startSession(client.peerconnection, gameName, roomID, playerIndex) + roomID, isNewRoom = startSession(wssession.peerconnection, gameName, roomID, playerIndex) if isNewRoom { wssession.oclient.send(WSPacket{ ID: "registerRoom", @@ -383,7 +385,7 @@ func removeSession(w *webrtc.WebRTC, room *Room) { } } -func GetServerIDOfRoom(oc *Client, roomID string) string { +func getServerIDOfRoom(oc *Client, roomID string) string { packet := oc.syncSend( WSPacket{ ID: "getRoom", @@ -394,8 +396,37 @@ func GetServerIDOfRoom(oc *Client, roomID string) string { return packet.Data } -func bridgeConnection(session *Session, serverID string) { +func bridgeConnection(session *Session, serverID string, gameName string, roomID string, playerIndex int) { + log.Println("Bridging connection to other Host ", serverID) + client := session.client + oclient := session.oclient // Ask client to init + + log.Println("Requesting offer to browser", serverID) + resp := client.syncSend(WSPacket{ + ID: "requestOffer", + Data: "", + }) + + log.Println("Sending offer to overlord to relay message to target host", resp.TargetHostID) + // Ask overlord to relay SDP packet to serverID + resp.TargetHostID = serverID + remoteTargetSDP := oclient.syncSend(resp) + log.Println("Got back remote host SDP, sending to browser", remoteTargetSDP.Data) + // Send back remote SDP of remote server to browser + client.syncSend(WSPacket{ + ID: "sdp", + Data: remoteTargetSDP.Data, + }) + log.Println("Init session done, start game on target host") + + oclient.syncSend(WSPacket{ + ID: "start", + Data: gameName, + RoomID: roomID, + PlayerIndex: playerIndex, + }) + log.Println("Game is started on remote host") } const overlordHost = "ws://localhost:9000/wso" @@ -415,7 +446,7 @@ func (s *Session) NewOverlordClient() { if err != nil { log.Println("Cannot connect to overlord") } - oclient := NewClient(oc, webrtc.NewWebRTC()) + oclient := NewClient(oc) oclient.send( WSPacket{ ID: "ping", @@ -437,6 +468,46 @@ func (s *Session) NewOverlordClient() { }, ) + // Received from overlord the sdp. This is happens when bridging + // TODO: refactor + oclient.receive( + "initwebrtc", + func(resp WSPacket) (req WSPacket) { + log.Println("Received a sdp request from overlord") + log.Println("Start peerconnection from the sdp") + + localSession, err := s.peerconnection.StartClient(resp.Data, width, height) + if err != nil { + log.Fatalln(err) + } + + return WSPacket{ + ID: "sdp", + Data: localSession, + } + }, + ) + + // Received start from overlord. This is happens when bridging + // TODO: refactor + oclient.receive( + "start", + func(resp WSPacket) (req WSPacket) { + log.Println("Received a start request from overlord") + log.Println("Add the connection to current room on the host") + + roomID, isNewRoom := startSession(s.peerconnection, resp.Data, resp.RoomID, resp.PlayerIndex) + // Bridge always access to old room + // TODO: log warn + if isNewRoom == true { + log.Fatal("Bridge should not spawn new room") + } + + req.ID = "start" + req.RoomID = roomID + return req + }, + ) go oclient.listen() s.oclient = oclient diff --git a/overlord.go b/overlord.go index 5878ba8d..1955ce69 100644 --- a/overlord.go +++ b/overlord.go @@ -13,7 +13,7 @@ import ( var roomToServer = map[string]string{} // servers are the map serverID to server Client -var servers = map[string]Client{} +var servers = map[string]*Client{} // If it's overlord, handle overlord connection (from host to overlord) func wso(w http.ResponseWriter, r *http.Request) { @@ -29,7 +29,14 @@ func wso(w http.ResponseWriter, r *http.Request) { serverID := strconv.Itoa(rand.Int()) log.Println("A new server connected ", serverID) - client := NewClient(c, webrtc.NewWebRTC()) + client := NewClient(c) + servers[serverID] = client + + wssession := &Session{ + client: client, + peerconnection: webrtc.NewWebRTC(), + // The server session is maintaining + } client.send( WSPacket{ @@ -61,5 +68,62 @@ func wso(w http.ResponseWriter, r *http.Request) { } }) + client.receive("initwebrtc", func(resp WSPacket) WSPacket { + log.Println("Received a relay sdp request from a host") + // TODO: Abstract + if resp.TargetHostID != serverID { + log.Println("sending relay sdp to target host", resp.TargetHostID) + // relay SDP to target host and get back sdp + // TODO: Async + sdp := servers[resp.TargetHostID].syncSend( + resp, + ) + + return sdp + } + log.Println("Target host is overlord itself: start peerconnection") + // If the target is in master + // start by its old + localSession, err := wssession.peerconnection.StartClient(resp.Data, width, height) + if err != nil { + log.Fatalln(err) + } + + return WSPacket{ + ID: "sdp", + Data: localSession, + } + }) + + // TODO: use relay ID type + // TODO: Merge sdp and start + client.receive("start", func(resp WSPacket) WSPacket { + log.Println("Received a relay start request from a host") + // TODO: Abstract + if resp.TargetHostID != serverID { + // relay SDP to target host and get back sdp + // TODO: Async + resp := servers[resp.TargetHostID].syncSend( + resp, + ) + + return resp + } + log.Println("Target host is overlord itself: start game") + // If the target is in master + // start by its old + roomID, isNewRoom := startSession(wssession.peerconnection, resp.Data, resp.RoomID, resp.PlayerIndex) + // Bridge always access to old room + // TODO: log warn + if isNewRoom == true { + log.Fatal("Bridge should not spawn new room") + } + + return WSPacket{ + ID: "start", + RoomID: roomID, + } + }) + client.listen() } diff --git a/static/js/ws.js b/static/js/ws.js index b50cd19e..936ab8b9 100644 --- a/static/js/ws.js +++ b/static/js/ws.js @@ -1,4 +1,5 @@ var pc; +var curPacketID = ""; // web socket conn = new WebSocket(`ws://${location.host}/ws`); @@ -27,10 +28,10 @@ conn.onmessage = e => { log("Got remote sdp"); pc.setRemoteDescription(new RTCSessionDescription(JSON.parse(atob(d["data"])))); break; - //case "requestOffer": - //pc.createOffer({offerToReceiveVideo: true, offerToReceiveAudio: false}).then(d => { - //pc.setLocalDescription(d).catch(log); - //}) + case "requestOffer": + curPacketID = d["packet_id"]; + log("Received request offer ", curPacketID) + startWebRTC(); //case "sdpremote": //log("Got remote sdp"); @@ -108,7 +109,8 @@ function startWebRTC() { session = btoa(JSON.stringify(pc.localDescription)); localSessionDescription = session; log("Send SDP to remote peer"); - conn.send(JSON.stringify({"id": "initwebrtc", "data": session})); + // TODO: Fix curPacketID + conn.send(JSON.stringify({"id": "initwebrtc", "data": session, "packet_id": curPacketID})); } else { console.log(JSON.stringify(event.candidate)); } diff --git a/ws.go b/ws.go index e52046f0..4f088f9f 100644 --- a/ws.go +++ b/ws.go @@ -7,15 +7,12 @@ import ( "strconv" "time" - "github.com/giongto35/cloud-game/webrtc" "github.com/gorilla/websocket" ) type Client struct { conn *websocket.Conn - peerconnection *webrtc.WebRTC - // sendCallback is callback based on packetID sendCallback map[string]func(req WSPacket) // recvCallback is callback when receive based on ID of the packet @@ -30,20 +27,19 @@ type WSPacket struct { PlayerIndex int `json:"player_index"` TargetHostID string `json:"target_id"` - PacketID string + PacketID string `json:"packet_id"` } var EmptyPacket = WSPacket{} -func NewClient(conn *websocket.Conn, webrtc *webrtc.WebRTC) *Client { +func NewClient(conn *websocket.Conn) *Client { sendCallback := map[string]func(WSPacket){} recvCallback := map[string]func(WSPacket){} return &Client{ conn: conn, - peerconnection: webrtc, - sendCallback: sendCallback, - recvCallback: recvCallback, + sendCallback: sendCallback, + recvCallback: recvCallback, } } @@ -94,6 +90,7 @@ func (c *Client) syncSend(request WSPacket) (response WSPacket) { func (c *Client) listen() { for { + log.Println("Waiting for message") _, rawMsg, err := c.conn.ReadMessage() if err != nil { log.Println("[!] read:", err) @@ -109,6 +106,8 @@ func (c *Client) listen() { if callback, ok := c.sendCallback[wspacket.PacketID]; ok { callback(wspacket) delete(c.sendCallback, wspacket.PacketID) + // Skip receiveCallback to avoid duplication + continue } // Check if some receiver with the ID is registered if callback, ok := c.recvCallback[wspacket.ID]; ok { From 1c26d128053093fe7d2cde6da3bdce98320aafe3 Mon Sep 17 00:00:00 2001 From: giongto35 Date: Sat, 20 Apr 2019 16:06:44 +0800 Subject: [PATCH 08/10] Fix bug, not use async when sending back sdp to browser --- main.go | 22 +++++++++++++++------- static/js/ws.js | 1 + 2 files changed, 16 insertions(+), 7 deletions(-) diff --git a/main.go b/main.go index 7718e0f5..ea5e2d31 100644 --- a/main.go +++ b/main.go @@ -131,11 +131,13 @@ func initRoom(roomID, gameName string) string { // TODO: If we remove sessions from room anytime a session is closed, we can check if the sessions list is empty or not. func isRoomRunning(roomID string) bool { // If no roomID is registered + fmt.Println("rooms list ", rooms) if _, ok := rooms[roomID]; !ok { return false } // If there is running session + fmt.Println("Running session", len(rooms[roomID].rtcSessions)) for _, s := range rooms[roomID].rtcSessions { if !s.IsClosed() { return true @@ -151,6 +153,7 @@ func startSession(webRTC *webrtc.WebRTC, gameName string, roomID string, playerI // If the roomID is empty, // or the roomID doesn't have any running sessions (room was closed) // we spawn a new room + log.Println("Is Room Running", isRoomRunning(roomID)) if roomID == "" || !isRoomRunning(roomID) { roomID = initRoom(roomID, gameName) isNewRoom = true @@ -412,19 +415,24 @@ func bridgeConnection(session *Session, serverID string, gameName string, roomID // Ask overlord to relay SDP packet to serverID resp.TargetHostID = serverID remoteTargetSDP := oclient.syncSend(resp) - log.Println("Got back remote host SDP, sending to browser", remoteTargetSDP.Data) + log.Println("Got back remote host SDP, sending to browser") // Send back remote SDP of remote server to browser - client.syncSend(WSPacket{ + //client.syncSend(WSPacket{ + //ID: "sdp", + //Data: remoteTargetSDP.Data, + //}) + client.send(WSPacket{ ID: "sdp", Data: remoteTargetSDP.Data, - }) + }, nil) log.Println("Init session done, start game on target host") oclient.syncSend(WSPacket{ - ID: "start", - Data: gameName, - RoomID: roomID, - PlayerIndex: playerIndex, + ID: "start", + Data: gameName, + TargetHostID: serverID, + RoomID: roomID, + PlayerIndex: playerIndex, }) log.Println("Game is started on remote host") } diff --git a/static/js/ws.js b/static/js/ws.js index 936ab8b9..d4f7bc40 100644 --- a/static/js/ws.js +++ b/static/js/ws.js @@ -27,6 +27,7 @@ conn.onmessage = e => { case "sdp": log("Got remote sdp"); pc.setRemoteDescription(new RTCSessionDescription(JSON.parse(atob(d["data"])))); + //conn.send(JSON.stringify({"id": "sdpdon", "packet_id": d["packet_id"]})); break; case "requestOffer": curPacketID = d["packet_id"]; From 6abd977d3ecddca05e3ba22b05cb8077832353fa Mon Sep 17 00:00:00 2001 From: giongto35 Date: Sat, 20 Apr 2019 16:13:04 +0800 Subject: [PATCH 09/10] Remove unnecessary log --- main.go | 3 --- static/js/ws.js | 3 +++ 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/main.go b/main.go index ea5e2d31..48169a9b 100644 --- a/main.go +++ b/main.go @@ -131,13 +131,11 @@ func initRoom(roomID, gameName string) string { // TODO: If we remove sessions from room anytime a session is closed, we can check if the sessions list is empty or not. func isRoomRunning(roomID string) bool { // If no roomID is registered - fmt.Println("rooms list ", rooms) if _, ok := rooms[roomID]; !ok { return false } // If there is running session - fmt.Println("Running session", len(rooms[roomID].rtcSessions)) for _, s := range rooms[roomID].rtcSessions { if !s.IsClosed() { return true @@ -153,7 +151,6 @@ func startSession(webRTC *webrtc.WebRTC, gameName string, roomID string, playerI // If the roomID is empty, // or the roomID doesn't have any running sessions (room was closed) // we spawn a new room - log.Println("Is Room Running", isRoomRunning(roomID)) if roomID == "" || !isRoomRunning(roomID) { roomID = initRoom(roomID, gameName) isNewRoom = true diff --git a/static/js/ws.js b/static/js/ws.js index d4f7bc40..cfb8459f 100644 --- a/static/js/ws.js +++ b/static/js/ws.js @@ -33,6 +33,9 @@ conn.onmessage = e => { curPacketID = d["packet_id"]; log("Received request offer ", curPacketID) startWebRTC(); + //pc.createOffer({offerToReceiveVideo: true, offerToReceiveAudio: false}).then(d => { + //pc.setLocalDescription(d).catch(log); + //}) //case "sdpremote": //log("Got remote sdp"); From 84bf7a1a5333e717c2b7416a714ab1ca0cfe4341 Mon Sep 17 00:00:00 2001 From: giongto35 Date: Sat, 20 Apr 2019 22:03:05 +0800 Subject: [PATCH 10/10] Allow change IP --- main.go | 9 ++++++--- overlord.go | 2 +- 2 files changed, 7 insertions(+), 4 deletions(-) diff --git a/main.go b/main.go index 48169a9b..fec6e9be 100644 --- a/main.go +++ b/main.go @@ -54,6 +54,7 @@ type Room struct { } var rooms = map[string]*Room{} +var port string = "8000" func main() { fmt.Println("Usage: ./game [debug]") @@ -62,12 +63,15 @@ func main() { indexFN = debugIndex fmt.Println("Use debug version") } - if len(os.Args) == 3 { + if len(os.Args) >= 3 { if os.Args[2] == "overlord" { IsOverlord = true } fmt.Println("Running as overlord ") } + if len(os.Args) >= 4 { + port = os.Args[3] + } rand.Seed(time.Now().UTC().UnixNano()) rooms = map[string]*Room{} @@ -81,7 +85,7 @@ func main() { if !IsOverlord { fmt.Println("http://localhost:8000") - http.ListenAndServe(":8000", nil) + http.ListenAndServe(":"+port, nil) } else { fmt.Println("http://localhost:9000") // Overlord expose one more path for handle overlord connections @@ -480,7 +484,6 @@ func (s *Session) NewOverlordClient() { func(resp WSPacket) (req WSPacket) { log.Println("Received a sdp request from overlord") log.Println("Start peerconnection from the sdp") - localSession, err := s.peerconnection.StartClient(resp.Data, width, height) if err != nil { log.Fatalln(err) diff --git a/overlord.go b/overlord.go index 1955ce69..06c90341 100644 --- a/overlord.go +++ b/overlord.go @@ -72,7 +72,7 @@ func wso(w http.ResponseWriter, r *http.Request) { log.Println("Received a relay sdp request from a host") // TODO: Abstract if resp.TargetHostID != serverID { - log.Println("sending relay sdp to target host", resp.TargetHostID) + log.Println("sending relay sdp to target host", resp) // relay SDP to target host and get back sdp // TODO: Async sdp := servers[resp.TargetHostID].syncSend(