diff --git a/app/cmd.go b/app/cmd.go index f8b6895..94d4822 100644 --- a/app/cmd.go +++ b/app/cmd.go @@ -4,73 +4,66 @@ import ( "fmt" "net/http" "net/url" + "os" "runtime" "strings" - "sync" - // "time" + "time" ) +var memInfo runtime.MemStats + type Info struct { + ch chan struct{} room string Server string `json:"server"` Proxy string `json:"proxy"` + Online string `json:"online"` + Rid int64 `json:"rid"` Start int64 `json:"start"` Last int64 `json:"last"` Income int64 `json:"income"` + Dons int64 `json:"dons"` + Tips int64 `json:"tips"` } -type Debug struct { - Goroutines int - Alloc uint64 - HeapSys uint64 - Uptime int64 -} - -type Worker struct { - chQuit chan struct{} -} - -type Workers struct { - sync.RWMutex - Map map[string]*Worker -} - -var ( - memInfo runtime.MemStats - chWorker = &Workers{Map: make(map[string]*Worker)} -) - -func removeRoom(room string) { - if checkWorker(room) { - chWorker.Lock() - //fmt.Printf("%v remove %v from chWorker.Map \n", time.Now().UnixMilli(), room ) - delete(chWorker.Map, room) - chWorker.Unlock() +func updateFileRooms() string { + for { + rooms.Json <- "" + s := <-rooms.Json + err := os.WriteFile(conf.Conn["start"], []byte(s), 0644) + if err != nil { + fmt.Println(err) + } + time.Sleep(10 * time.Second) } } -func checkWorker(room string) bool { - chWorker.RLock() - defer chWorker.RUnlock() - if _, ok := chWorker.Map[room]; ok { - return true - } - return false -} - -func listRooms() string { - rooms.Json <- "" - s := <-rooms.Json - return s -} - func listHandler(w http.ResponseWriter, _ *http.Request) { - fmt.Fprint(w, listRooms()) + dat, err := os.ReadFile(conf.Conn["start"]) + if err != nil { + fmt.Println(err) + return + } + fmt.Fprint(w, string(dat)) } func debugHandler(w http.ResponseWriter, _ *http.Request) { + ws.Count <- 0 + l := <-ws.Count runtime.ReadMemStats(&memInfo) - j, err := json.Marshal(Debug{runtime.NumGoroutine(), memInfo.Alloc, memInfo.HeapSys, uptime}) + j, err := json.Marshal(struct { + Goroutines int + WebSocket int + Uptime int64 + Alloc uint64 + HeapSys uint64 + }{ + Goroutines: runtime.NumGoroutine(), + Alloc: memInfo.Alloc, + HeapSys: memInfo.HeapSys, + Uptime: uptime, + WebSocket: l, + }) if err == nil { fmt.Fprint(w, string(j)) } @@ -81,37 +74,45 @@ func cmdHandler(w http.ResponseWriter, r *http.Request) { fmt.Fprint(w, "403") return } - params := r.URL.Query() if len(params["room"]) > 0 && len(params["server"]) > 0 && len(params["proxy"]) > 0 { - room := params["room"][0] - server := params["server"][0] - proxy := params["proxy"][0] - if checkWorker(room) { - fmt.Println("Already track:", room) - return + now := time.Now().Unix() + workerData := Info{ + room: params["room"][0], + Server: params["server"][0], + Proxy: params["proxy"][0], + Online: "0", + Start: now, + Last: now, + Rid: 0, + Income: 0, + Dons: 0, + Tips: 0, } - - info, ok := getRoomInfo(room) - if !ok { - fmt.Println("No room in MySQL:", room) - return - } - - chQuit := make(chan struct{}) - - chWorker.Lock() - chWorker.Map[room] = &Worker{chQuit: chQuit} - chWorker.Unlock() - - go statRoom(chQuit, room, server, proxy, info, url.URL{Scheme: "wss", Host: server + ".bcccdn.com", Path: "/websocket"}) - + startRoom(workerData) } if len(params["exit"]) > 0 { - room := strings.Join(params["exit"], "") - if checkWorker(room) { - close(chWorker.Map[room].chQuit) // exit gorutine - removeRoom(room) - } + rooms.Stop <- strings.Join(params["exit"], "") } + fmt.Fprint(w, string("ok")) +} + +func startRoom(workerData Info) { + rooms.Check <- workerData.room + testRoom := <-rooms.Check + if testRoom == workerData.room { + fmt.Println("Already track:", workerData.room) + return + } + + rid, ok := getRoomInfo(workerData.room) + if !ok { + fmt.Println("No room in MySQL:", workerData.room) + return + } + + workerData.Rid = rid + workerData.ch = make(chan struct{}) + + go xWorker(workerData, url.URL{Scheme: "wss", Host: workerData.Server + ".bcccdn.com", Path: "/websocket"}) } diff --git a/app/conf.go b/app/conf.go index 6176b5b..2ac77b3 100644 --- a/app/conf.go +++ b/app/conf.go @@ -17,6 +17,7 @@ func startConfig() { conf.Conn = map[string]string{ "mysql": "user:passwd@unix(/var/run/mysqld/mysqld.sock)/base?interpolateParams=true", "click": "tcp://127.0.0.1:9000/base?compress=true&debug=false", + "start": "/tmp/bongaStart.txt", } // 3proxy diff --git a/app/echo.go b/app/echo.go index bbd46d2..3bfddf1 100644 --- a/app/echo.go +++ b/app/echo.go @@ -1,89 +1,49 @@ package main import ( - "github.com/gorilla/websocket" "net/http" - //"fmt" + + "github.com/gorilla/websocket" ) -func newHub() *Hub { - return &Hub{ - broadcast: make(chan []byte), - register: make(chan *Client), - unregister: make(chan *Client), - clients: make(map[*Client]bool), +var ( + wsClients = make(map[*websocket.Conn]struct{}) + + ws = struct { + Count chan int + Send chan []byte + Add chan *websocket.Conn + }{ + Count: make(chan int, 100), + Send: make(chan []byte, 100), + Add: make(chan *websocket.Conn, 100), } -} +) -type Hub struct { - clients map[*Client]bool - broadcast chan []byte - register chan *Client - unregister chan *Client -} - -type Client struct { - hub *Hub - conn *websocket.Conn - send chan []byte -} - -func (h *Hub) run() { +func broadcast() { for { select { - case client := <-h.register: - h.clients[client] = true - case client := <-h.unregister: - if _, ok := h.clients[client]; ok { - delete(h.clients, client) - close(client.send) - } - case message := <-h.broadcast: - //fmt.Println("map channel:", len(h.broadcast), cap(h.broadcast)) - for client := range h.clients { - select { - case client.send <- message: - default: - close(client.send) - delete(h.clients, client) + case conn := <-ws.Add: + wsClients[conn] = struct{}{} + + case <-ws.Count: + ws.Count <- len(wsClients) + + case message := <-ws.Send: + for conn := range wsClients { + if err := conn.WriteMessage(1, message); err != nil { + conn.Close() + delete(wsClients, conn) } } } } } -func (c *Client) writePump() { - for { - message, ok := <-c.send - if !ok { - // The hub closed the channel. - c.conn.WriteMessage(websocket.CloseMessage, []byte{}) - return - } - c.conn.WriteMessage(1, message) - } - c.conn.Close() -} - -func (c *Client) readPump() { - for { - // Client close connection - _, _, err := c.conn.ReadMessage() - if err != nil { - break - } - } - c.hub.unregister <- c - c.conn.Close() -} - -func (hub *Hub) wsHandler(w http.ResponseWriter, r *http.Request) { +func wsHandler(w http.ResponseWriter, r *http.Request) { conn, err := websocket.Upgrade(w, r, w.Header(), 1024, 1024) if err != nil { return } - client := &Client{hub: hub, conn: conn, send: make(chan []byte)} - client.hub.register <- client - go client.readPump() - go client.writePump() + ws.Add <- conn } diff --git a/app/main.go b/app/main.go index 77aeb20..53cb10e 100644 --- a/app/main.go +++ b/app/main.go @@ -1,6 +1,7 @@ package main import ( + "fmt" "log" "math/rand" "net" @@ -17,23 +18,29 @@ import ( type Rooms struct { Count chan int Json chan string - Add chan Info + Check chan string + Stop chan string Del chan string + Add chan Info } -var hub = newHub() -var Mysql, Clickhouse *sqlx.DB -var json = jsoniter.ConfigCompatibleWithStandardLibrary +var ( + Mysql, Clickhouse *sqlx.DB -var save = make(chan saveData, 100) -var slog = make(chan saveLog, 100) + json = jsoniter.ConfigCompatibleWithStandardLibrary -var rooms = &Rooms{ - Count: make(chan int), - Json: make(chan string), - Add: make(chan Info), - Del: make(chan string), -} + save = make(chan saveData, 100) + slog = make(chan saveLog, 100) + + rooms = &Rooms{ + Count: make(chan int), + Json: make(chan string), + Check: make(chan string), + Stop: make(chan string), + Del: make(chan string), + Add: make(chan Info), + } +) func main() { rand.Seed(time.Now().UnixNano()) @@ -43,17 +50,19 @@ func main() { initMysql() initClickhouse() - go hub.run() go mapRooms() go announceCount() go saveDB() go saveLogs() + go broadcast() - http.HandleFunc("/bongacams/ws/", hub.wsHandler) + http.HandleFunc("/bongacams/ws/", wsHandler) http.HandleFunc("/bongacams/cmd/", cmdHandler) http.HandleFunc("/bongacams/list/", listHandler) http.HandleFunc("/bongacams/debug/", debugHandler) + go fastStart() + const SOCK = "/tmp/bongacams.sock" os.Remove(SOCK) unixListener, err := net.Listen("unix", SOCK) @@ -84,3 +93,40 @@ func initClickhouse() { func randInt(min int, max int) int { return min + rand.Intn(max-min) } + +func fastStart() { + defer func() { + go updateFileRooms() + }() + dat, err := os.ReadFile(conf.Conn["start"]) + if err != nil { + fmt.Println(err) + return + } + list := make(map[string]Info) + if err := json.Unmarshal(dat, &list); err != nil { + fmt.Println(err.Error()) + return + } + now := time.Now().Unix() + for k, v := range list { + if now > v.Last+60*20 { + continue + } + fmt.Println("fastStart:", k, v.Server, v.Proxy) + workerData := Info{ + room: k, + Server: v.Server, + Proxy: v.Proxy, + Online: v.Online, + Start: v.Start, + Last: now, + Rid: v.Rid, + Income: v.Income, + Dons: v.Dons, + Tips: v.Tips, + } + startRoom(workerData) + time.Sleep(2 * time.Second) + } +} diff --git a/app/save.go b/app/save.go index c327406..7a7779c 100644 --- a/app/save.go +++ b/app/save.go @@ -5,10 +5,6 @@ import ( "time" ) -type tID struct { - Id int64 `db:"id"` -} - type saveData struct { Room string From string @@ -29,29 +25,46 @@ type DonatorCache struct { } func getDonId(name string) int64 { - donator := new(tID) - err := Mysql.Get(donator, "SELECT id FROM donator WHERE name=?", name) + var id int64 + err := Mysql.Get(&id, "SELECT id FROM donator WHERE name=?", name) if err != nil { res, _ := Mysql.Exec("INSERT INTO donator (`name`) VALUES (?)", name) - donator.Id, _ = res.LastInsertId() + id, _ = res.LastInsertId() } - return donator.Id + return id } -func getRoomInfo(name string) (*tID, bool) { +func getRoomInfo(name string) (int64, bool) { + var id int64 result := true - room := new(tID) - err := Mysql.Get(room, "SELECT id FROM room WHERE name=?", name) + err := Mysql.Get(&id, "SELECT id FROM room WHERE name=?", name) if err != nil { result = false } - return room, result + return id, result +} + +func getSumTokens() int64 { + r := struct { + Date string + Sum int64 + }{} + err := Clickhouse.Get(&r, "SELECT toStartOfHour(toDateTime(`unix`)) as date, SUM(`token`) as sum FROM `stat` WHERE time = today() GROUP BY date ORDER BY date DESC LIMIT 1") + if err == nil && r.Sum > 0 { + return r.Sum + } + return 0 } func saveDB() { - last := time.Now().Unix() + hours, _, _ := time.Now().Clock() + bulk := make(map[int]saveData) + update := make(map[int64]int64) data := make(map[string]*DonatorCache) + index := make(map[string]int64) + + index = map[string]int64{"hours": int64(hours), "tokens": getSumTokens(), "last": time.Now().Unix()} for { select { @@ -66,6 +79,81 @@ func saveDB() { data[m.From] = &DonatorCache{Id: getDonId(m.From), Last: now} } + num := len(bulk) + + bulk[num] = m + + update[m.Rid] = m.Now + + if num > 512 { + + tx, err := Mysql.Begin() + if err == nil { + st, _ := tx.Prepare("INSERT INTO `stat` (`did`, `rid`, `token`, `time`) VALUES (?, ?, ?, ?)") + for _, v := range bulk { + st.Exec(data[v.From].Id, v.Rid, v.Amount, v.Now) + } + tx.Commit() + st.Close() + } + + tx, err = Mysql.Begin() + if err == nil { + st, _ := tx.Prepare("UPDATE `room` SET `last` = ? WHERE `id` = ?") + for k, v := range update { + st.Exec(v, k) + } + tx.Commit() + st.Close() + } + + tx, err = Clickhouse.Begin() + if err == nil { + st, _ := tx.Prepare("INSERT INTO stat VALUES (?, ?, ?, ?, ?)") + for _, v := range bulk { + st.Exec(uint32(data[v.From].Id), uint32(v.Rid), uint32(v.Amount), time.Unix(v.Now, 0), uint32(v.Now)) + } + tx.Commit() + st.Close() + } + + bulk = make(map[int]saveData) + update = make(map[int64]int64) + } + + if m.Amount > 99 { + msg, err := json.Marshal(struct { + Room string `json:"room"` + Donator string `json:"donator"` + Amount int64 `json:"amount"` + }{ + Room: m.Room, + Donator: m.From, + Amount: m.Amount, + }) + if err == nil { + ws.Send <- msg + } + } + + hours, minutes, seconds := time.Now().Clock() + if int64(hours) == index["hours"] { + index["tokens"] += m.Amount + } else { + index = map[string]int64{"hours": int64(hours), "tokens": 0, "last": 0} + } + + if minutes >= 5 && now > index["last"]+30 { + seconds += minutes * 60 + msg, err := json.Marshal(struct { + Index int64 `json:"index"` + }{Index: index["tokens"] / int64(seconds) * 3600 / 1000 * 5 / 100}) + if err == nil { + ws.Send <- msg + } + index["last"] = now + } + if randInt(0, 10000) == 777 { // 0.001% l := len(data) for k, v := range data { @@ -75,49 +163,11 @@ func saveDB() { } fmt.Println("Clean map:", l, "=>", len(data)) } - - Mysql.Exec("UPDATE `room` SET `last` = ? WHERE `id` = ?", m.Now, m.Rid) - - num := len(bulk) - - bulk[num] = m - - if num >= 999 || now >= last+10 { - tx, err := Mysql.Begin() - if err == nil { - for _, v := range bulk { - tx.Exec("INSERT INTO `stat` (`did`, `rid`, `token`, `time`) VALUES (?, ?, ?, ?)", data[v.From].Id, v.Rid, v.Amount, v.Now) - } - } - tx.Commit() - - tx, err = Clickhouse.Begin() - if err == nil { - st, _ := tx.Prepare("INSERT INTO stat VALUES (?, ?, ?, ?)") - //fmt.Println("G:", err) - for _, v := range bulk { - st.Exec(uint32(data[v.From].Id), uint32(v.Rid), uint32(v.Amount), time.Unix(v.Now, 0)) - //fmt.Println("B:", aaa, sss) - } - tx.Commit() - st.Close() - } - - last = now - bulk = make(map[int]saveData) - } - if m.Amount > 99 { - msg, err := json.Marshal(AnnounceDonate{Room: m.Room, Donator: m.From, Amount: m.Amount}) - if err == nil { - hub.broadcast <- msg - } - } } } } func saveLogs() { - last := time.Now().Unix() bulk := make(map[int]saveLog) for { select { @@ -125,16 +175,16 @@ func saveLogs() { if len(m.Mes) > 0 { num := len(bulk) bulk[num] = m - now := time.Now().Unix() - if num >= 2047 || now >= last+10 { + if num > 2048 { tx, err := Mysql.Begin() if err == nil { + st, _ := tx.Prepare("INSERT INTO `logs` (`rid`, `time`, `mes`) VALUES (?, ?, ?)") for _, v := range bulk { - tx.Exec("INSERT INTO `logs` (`rid`, `time`, `mes`) VALUES (?, ?, ?)", v.Rid, v.Now, v.Mes) + st.Exec(v.Rid, v.Now, v.Mes) } tx.Commit() + st.Close() } - last = now bulk = make(map[int]saveLog) } } diff --git a/app/worker.go b/app/worker.go index f69c330..e577063 100644 --- a/app/worker.go +++ b/app/worker.go @@ -8,7 +8,6 @@ import ( "net/url" "strings" "time" - //"bytes" ) var uptime = time.Now().Unix() @@ -40,16 +39,6 @@ type DonateResponse struct { A int64 `json:"a"` } -type AnnounceCount struct { - Count int `json:"count"` -} - -type AnnounceDonate struct { - Room string `json:"room"` - Donator string `json:"donator"` - Amount int64 `json:"amount"` -} - func mapRooms() { data := make(map[string]*Info) @@ -57,7 +46,7 @@ func mapRooms() { for { select { case m := <-rooms.Add: - data[m.room] = &Info{Server: m.Server, Proxy: m.Proxy, Start: m.Start, Last: m.Last, Income: m.Income} + data[m.room] = &Info{Server: m.Server, Proxy: m.Proxy, Start: m.Start, Last: m.Last, Online: m.Online, Income: m.Income, Dons: m.Dons, Tips: m.Tips, ch: m.ch} case s := <-rooms.Json: j, err := json.Marshal(data) @@ -71,7 +60,17 @@ func mapRooms() { case key := <-rooms.Del: delete(data, key) - removeRoom(key) + + case room := <-rooms.Check: + if _, ok := data[room]; !ok { + room = "" + } + rooms.Check <- room + + case room := <-rooms.Stop: + if _, ok := data[room]; ok { + close(data[room].ch) + } } } } @@ -81,9 +80,11 @@ func announceCount() { time.Sleep(30 * time.Second) rooms.Count <- 0 l := <-rooms.Count - msg, err := json.Marshal(AnnounceCount{Count: l}) + msg, err := json.Marshal(struct { + Count int `json:"count"` + }{Count: l}) if err == nil { - hub.broadcast <- msg + ws.Send <- msg } } } @@ -118,10 +119,16 @@ func getAMF(room string) (bool, *AuthResponse) { return true, v } -func statRoom(chQuit chan struct{}, room, server, proxy string, info *tID, u url.URL) { - //fmt.Println("Start", room, "server", server, "proxy", proxy) +func xWorker(workerData Info, u url.URL) { + fmt.Println("Start", workerData.room, "server", workerData.Server, "proxy", workerData.Proxy) - ok, v := getAMF(room) + rooms.Add <- workerData + + defer func() { + rooms.Del <- workerData.room + }() + + ok, v := getAMF(workerData.room) if !ok { fmt.Println("exit: no amf parms") return @@ -129,85 +136,62 @@ func statRoom(chQuit chan struct{}, room, server, proxy string, info *tID, u url Dialer := *websocket.DefaultDialer - if _, ok := conf.Proxy[proxy]; ok { + if _, ok := conf.Proxy[workerData.Proxy]; ok { Dialer = websocket.Dialer{ Proxy: http.ProxyURL(&url.URL{ Scheme: "http", // or "https" depending on your proxy - Host: conf.Proxy[proxy], + Host: conf.Proxy[workerData.Proxy], Path: "/", }), HandshakeTimeout: 45 * time.Second, // https://pkg.go.dev/github.com/gorilla/websocket } } - now := time.Now().Unix() - workerData := Info{room, server, proxy, now, now, 0} - rooms.Add <- workerData - - defer func() { - fmt.Println("defer remove map", room) - rooms.Del <- room - }() - c, _, err := Dialer.Dial(u.String(), nil) if err != nil { - fmt.Println(err.Error(), room) + fmt.Println(err.Error(), workerData.room) return } - defer func() { - fmt.Println("defer close", room) - c.Close() - }() + defer c.Close() c.SetReadDeadline(time.Now().Add(60 * time.Second)) - fmt.Println("send first", room) - if err = c.WriteMessage(websocket.TextMessage, []byte(fmt.Sprintf(`{"id":%d,"name":"joinRoom","args":["%s",{"username":"%s","displayName":"%s","location":"%s","chathost":"%s","isRu":%t,"isPerformer":false,"hasStream":false,"isLogged":false,"isPayable":false,"showType":"public"},"%s"]}`, 1, v.UserData.Chathost, v.UserData.Username, v.UserData.DisplayName, v.UserData.Location, v.UserData.Chathost, v.UserData.IsRu, v.LocalData.DataKey))); err != nil { fmt.Println(err.Error()) return } - fmt.Println("read first", room) - _, message, err := c.ReadMessage() if err != nil { - fmt.Println(err.Error(), room) + fmt.Println(err.Error(), workerData.room) return } - fmt.Println(room, len(string(message)), string(message)) - - slog <- saveLog{info.Id, now, string(message)} + slog <- saveLog{workerData.Rid, time.Now().Unix(), string(message)} if string(message) == `{"id":1,"result":{"audioAvailable":false,"freeShow":false},"error":null}` { - fmt.Println("room offline, exit", room) + fmt.Println("room offline, exit", workerData.room) return } - fmt.Println("send second", room) - if err = c.WriteMessage(websocket.TextMessage, []byte(fmt.Sprintf(`{"id":%d,"name":"ChatModule.connect","args":["public-chat"]}`, 2))); err != nil { - fmt.Println(err.Error(), room) + fmt.Println(err.Error(), workerData.room) return } - fmt.Println("read second", room) _, message, err = c.ReadMessage() if err != nil { - fmt.Println(err.Error(), room) + fmt.Println(err.Error(), workerData.room) return } - fmt.Println(room, len(string(message)), string(message)) + slog <- saveLog{workerData.Rid, time.Now().Unix(), string(message)} - slog <- saveLog{info.Id, now, string(message)} quit := make(chan bool) pid := 3 defer func() { - fmt.Println("defer quit", room) quit <- true }() @@ -220,8 +204,8 @@ func statRoom(chQuit chan struct{}, room, server, proxy string, info *tID, u url return case <-ticker.C: if err = c.WriteMessage(websocket.TextMessage, []byte(fmt.Sprintf(`{"id":%d,"name":"ping"}`, pid))); err != nil { - fmt.Println(err.Error(), room) - close(chWorker.Map[room].chQuit) + fmt.Println(err.Error(), workerData.room) + rooms.Stop <- workerData.room return } pid++ @@ -232,52 +216,52 @@ func statRoom(chQuit chan struct{}, room, server, proxy string, info *tID, u url for { select { - case <-chQuit: - fmt.Println("Exit room:", room) + case <-workerData.ch: + fmt.Println("Exit room:", workerData.room) return - default: - c.SetReadDeadline(time.Now().Add(30 * time.Minute)) - _, message, err := c.ReadMessage() - if err != nil { - fmt.Println(err.Error()) - return - } + } - now = time.Now().Unix() + c.SetReadDeadline(time.Now().Add(30 * time.Minute)) + _, message, err := c.ReadMessage() + if err != nil { + fmt.Println(err.Error()) + return + } - slog <- saveLog{info.Id, now, string(message)} + now := time.Now().Unix() - m := &ServerResponse{} + slog <- saveLog{workerData.Rid, now, string(message)} - if err = json.Unmarshal(message, m); err != nil { - fmt.Println(err.Error(), room) - continue - } + m := &ServerResponse{} - workerData.Last = now - rooms.Add <- workerData + if err = json.Unmarshal(message, m); err != nil { + fmt.Println(err.Error(), workerData.room) + continue + } - if m.Type == "ServerMessageEvent:PERFORMER_STATUS_CHANGE" && string(m.Body) == `"offline"` { - fmt.Println(m.Type, room) - return - } + workerData.Last = now + rooms.Add <- workerData - if m.Type == "ServerMessageEvent:ROOM_CLOSE" { - fmt.Println(m.Type, room) - return - } + if m.Type == "ServerMessageEvent:PERFORMER_STATUS_CHANGE" && string(m.Body) == `"offline"` { + fmt.Println(m.Type, workerData.room) + return + } - if m.Type == "ServerMessageEvent:INCOMING_TIP" { - d := &DonateResponse{} - if err = json.Unmarshal(m.Body, d); err == nil { - //fmt.Println(d.F.Username, " send ", d.A, "tokens") + if m.Type == "ServerMessageEvent:ROOM_CLOSE" { + fmt.Println(m.Type, workerData.room) + return + } - save <- saveData{room, d.F.Username, info.Id, d.A, now} + if m.Type == "ServerMessageEvent:INCOMING_TIP" { + d := &DonateResponse{} + if err = json.Unmarshal(m.Body, d); err == nil { + //fmt.Println(d.F.Username, "send", d.A, "tokens") - workerData.Income += d.A - rooms.Add <- workerData - } + save <- saveData{workerData.room, d.F.Username, workerData.Rid, d.A, now} + + workerData.Income += d.A + rooms.Add <- workerData } } }