This commit is contained in:
poiuty 2022-08-07 14:18:21 +03:00
parent fec669ea34
commit 7023da390f
6 changed files with 339 additions and 297 deletions

View file

@ -4,73 +4,66 @@ import (
"fmt" "fmt"
"net/http" "net/http"
"net/url" "net/url"
"os"
"runtime" "runtime"
"strings" "strings"
"sync" "time"
// "time"
) )
var memInfo runtime.MemStats
type Info struct { type Info struct {
ch chan struct{}
room string room string
Server string `json:"server"` Server string `json:"server"`
Proxy string `json:"proxy"` Proxy string `json:"proxy"`
Online string `json:"online"`
Rid int64 `json:"rid"`
Start int64 `json:"start"` Start int64 `json:"start"`
Last int64 `json:"last"` Last int64 `json:"last"`
Income int64 `json:"income"` Income int64 `json:"income"`
Dons int64 `json:"dons"`
Tips int64 `json:"tips"`
} }
type Debug struct { func updateFileRooms() string {
Goroutines int for {
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 checkWorker(room string) bool {
chWorker.RLock()
defer chWorker.RUnlock()
if _, ok := chWorker.Map[room]; ok {
return true
}
return false
}
func listRooms() string {
rooms.Json <- "" rooms.Json <- ""
s := <-rooms.Json s := <-rooms.Json
return s err := os.WriteFile(conf.Conn["start"], []byte(s), 0644)
if err != nil {
fmt.Println(err)
}
time.Sleep(10 * time.Second)
}
} }
func listHandler(w http.ResponseWriter, _ *http.Request) { 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) { func debugHandler(w http.ResponseWriter, _ *http.Request) {
ws.Count <- 0
l := <-ws.Count
runtime.ReadMemStats(&memInfo) 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 { if err == nil {
fmt.Fprint(w, string(j)) fmt.Fprint(w, string(j))
} }
@ -81,37 +74,45 @@ func cmdHandler(w http.ResponseWriter, r *http.Request) {
fmt.Fprint(w, "403") fmt.Fprint(w, "403")
return return
} }
params := r.URL.Query() params := r.URL.Query()
if len(params["room"]) > 0 && len(params["server"]) > 0 && len(params["proxy"]) > 0 { if len(params["room"]) > 0 && len(params["server"]) > 0 && len(params["proxy"]) > 0 {
room := params["room"][0] now := time.Now().Unix()
server := params["server"][0] workerData := Info{
proxy := params["proxy"][0] room: params["room"][0],
if checkWorker(room) { Server: params["server"][0],
fmt.Println("Already track:", room) Proxy: params["proxy"][0],
return Online: "0",
Start: now,
Last: now,
Rid: 0,
Income: 0,
Dons: 0,
Tips: 0,
} }
startRoom(workerData)
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"})
} }
if len(params["exit"]) > 0 { if len(params["exit"]) > 0 {
room := strings.Join(params["exit"], "") rooms.Stop <- strings.Join(params["exit"], "")
if checkWorker(room) {
close(chWorker.Map[room].chQuit) // exit gorutine
removeRoom(room)
} }
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"})
} }

View file

@ -17,6 +17,7 @@ func startConfig() {
conf.Conn = map[string]string{ conf.Conn = map[string]string{
"mysql": "user:passwd@unix(/var/run/mysqld/mysqld.sock)/base?interpolateParams=true", "mysql": "user:passwd@unix(/var/run/mysqld/mysqld.sock)/base?interpolateParams=true",
"click": "tcp://127.0.0.1:9000/base?compress=true&debug=false", "click": "tcp://127.0.0.1:9000/base?compress=true&debug=false",
"start": "/tmp/bongaStart.txt",
} }
// 3proxy // 3proxy

View file

@ -1,89 +1,49 @@
package main package main
import ( import (
"github.com/gorilla/websocket"
"net/http" "net/http"
//"fmt"
"github.com/gorilla/websocket"
) )
func newHub() *Hub { var (
return &Hub{ wsClients = make(map[*websocket.Conn]struct{})
broadcast: make(chan []byte),
register: make(chan *Client),
unregister: make(chan *Client),
clients: make(map[*Client]bool),
}
}
type Hub struct { ws = struct {
clients map[*Client]bool Count chan int
broadcast chan []byte Send chan []byte
register chan *Client Add chan *websocket.Conn
unregister chan *Client }{
Count: make(chan int, 100),
Send: make(chan []byte, 100),
Add: make(chan *websocket.Conn, 100),
} }
)
type Client struct { func broadcast() {
hub *Hub
conn *websocket.Conn
send chan []byte
}
func (h *Hub) run() {
for { for {
select { select {
case client := <-h.register: case conn := <-ws.Add:
h.clients[client] = true wsClients[conn] = struct{}{}
case client := <-h.unregister:
if _, ok := h.clients[client]; ok { case <-ws.Count:
delete(h.clients, client) ws.Count <- len(wsClients)
close(client.send)
} case message := <-ws.Send:
case message := <-h.broadcast: for conn := range wsClients {
//fmt.Println("map channel:", len(h.broadcast), cap(h.broadcast)) if err := conn.WriteMessage(1, message); err != nil {
for client := range h.clients { conn.Close()
select { delete(wsClients, conn)
case client.send <- message:
default:
close(client.send)
delete(h.clients, client)
} }
} }
} }
} }
} }
func (c *Client) writePump() { func wsHandler(w http.ResponseWriter, r *http.Request) {
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) {
conn, err := websocket.Upgrade(w, r, w.Header(), 1024, 1024) conn, err := websocket.Upgrade(w, r, w.Header(), 1024, 1024)
if err != nil { if err != nil {
return return
} }
client := &Client{hub: hub, conn: conn, send: make(chan []byte)} ws.Add <- conn
client.hub.register <- client
go client.readPump()
go client.writePump()
} }

View file

@ -1,6 +1,7 @@
package main package main
import ( import (
"fmt"
"log" "log"
"math/rand" "math/rand"
"net" "net"
@ -17,23 +18,29 @@ import (
type Rooms struct { type Rooms struct {
Count chan int Count chan int
Json chan string Json chan string
Add chan Info Check chan string
Stop chan string
Del chan string Del chan string
Add chan Info
} }
var hub = newHub() var (
var Mysql, Clickhouse *sqlx.DB Mysql, Clickhouse *sqlx.DB
var json = jsoniter.ConfigCompatibleWithStandardLibrary
var save = make(chan saveData, 100) json = jsoniter.ConfigCompatibleWithStandardLibrary
var slog = make(chan saveLog, 100)
var rooms = &Rooms{ save = make(chan saveData, 100)
slog = make(chan saveLog, 100)
rooms = &Rooms{
Count: make(chan int), Count: make(chan int),
Json: make(chan string), Json: make(chan string),
Add: make(chan Info), Check: make(chan string),
Stop: make(chan string),
Del: make(chan string), Del: make(chan string),
Add: make(chan Info),
} }
)
func main() { func main() {
rand.Seed(time.Now().UnixNano()) rand.Seed(time.Now().UnixNano())
@ -43,17 +50,19 @@ func main() {
initMysql() initMysql()
initClickhouse() initClickhouse()
go hub.run()
go mapRooms() go mapRooms()
go announceCount() go announceCount()
go saveDB() go saveDB()
go saveLogs() go saveLogs()
go broadcast()
http.HandleFunc("/bongacams/ws/", hub.wsHandler) http.HandleFunc("/bongacams/ws/", wsHandler)
http.HandleFunc("/bongacams/cmd/", cmdHandler) http.HandleFunc("/bongacams/cmd/", cmdHandler)
http.HandleFunc("/bongacams/list/", listHandler) http.HandleFunc("/bongacams/list/", listHandler)
http.HandleFunc("/bongacams/debug/", debugHandler) http.HandleFunc("/bongacams/debug/", debugHandler)
go fastStart()
const SOCK = "/tmp/bongacams.sock" const SOCK = "/tmp/bongacams.sock"
os.Remove(SOCK) os.Remove(SOCK)
unixListener, err := net.Listen("unix", SOCK) unixListener, err := net.Listen("unix", SOCK)
@ -84,3 +93,40 @@ func initClickhouse() {
func randInt(min int, max int) int { func randInt(min int, max int) int {
return min + rand.Intn(max-min) 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)
}
}

View file

@ -5,10 +5,6 @@ import (
"time" "time"
) )
type tID struct {
Id int64 `db:"id"`
}
type saveData struct { type saveData struct {
Room string Room string
From string From string
@ -29,29 +25,46 @@ type DonatorCache struct {
} }
func getDonId(name string) int64 { func getDonId(name string) int64 {
donator := new(tID) var id int64
err := Mysql.Get(donator, "SELECT id FROM donator WHERE name=?", name) err := Mysql.Get(&id, "SELECT id FROM donator WHERE name=?", name)
if err != nil { if err != nil {
res, _ := Mysql.Exec("INSERT INTO donator (`name`) VALUES (?)", name) 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 result := true
room := new(tID) err := Mysql.Get(&id, "SELECT id FROM room WHERE name=?", name)
err := Mysql.Get(room, "SELECT id FROM room WHERE name=?", name)
if err != nil { if err != nil {
result = false 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() { func saveDB() {
last := time.Now().Unix() hours, _, _ := time.Now().Clock()
bulk := make(map[int]saveData) bulk := make(map[int]saveData)
update := make(map[int64]int64)
data := make(map[string]*DonatorCache) data := make(map[string]*DonatorCache)
index := make(map[string]int64)
index = map[string]int64{"hours": int64(hours), "tokens": getSumTokens(), "last": time.Now().Unix()}
for { for {
select { select {
@ -66,6 +79,81 @@ func saveDB() {
data[m.From] = &DonatorCache{Id: getDonId(m.From), Last: now} 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% if randInt(0, 10000) == 777 { // 0.001%
l := len(data) l := len(data)
for k, v := range data { for k, v := range data {
@ -75,49 +163,11 @@ func saveDB() {
} }
fmt.Println("Clean map:", l, "=>", len(data)) 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() { func saveLogs() {
last := time.Now().Unix()
bulk := make(map[int]saveLog) bulk := make(map[int]saveLog)
for { for {
select { select {
@ -125,16 +175,16 @@ func saveLogs() {
if len(m.Mes) > 0 { if len(m.Mes) > 0 {
num := len(bulk) num := len(bulk)
bulk[num] = m bulk[num] = m
now := time.Now().Unix() if num > 2048 {
if num >= 2047 || now >= last+10 {
tx, err := Mysql.Begin() tx, err := Mysql.Begin()
if err == nil { if err == nil {
st, _ := tx.Prepare("INSERT INTO `logs` (`rid`, `time`, `mes`) VALUES (?, ?, ?)")
for _, v := range bulk { 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() tx.Commit()
st.Close()
} }
last = now
bulk = make(map[int]saveLog) bulk = make(map[int]saveLog)
} }
} }

View file

@ -8,7 +8,6 @@ import (
"net/url" "net/url"
"strings" "strings"
"time" "time"
//"bytes"
) )
var uptime = time.Now().Unix() var uptime = time.Now().Unix()
@ -40,16 +39,6 @@ type DonateResponse struct {
A int64 `json:"a"` 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() { func mapRooms() {
data := make(map[string]*Info) data := make(map[string]*Info)
@ -57,7 +46,7 @@ func mapRooms() {
for { for {
select { select {
case m := <-rooms.Add: 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: case s := <-rooms.Json:
j, err := json.Marshal(data) j, err := json.Marshal(data)
@ -71,7 +60,17 @@ func mapRooms() {
case key := <-rooms.Del: case key := <-rooms.Del:
delete(data, key) 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) time.Sleep(30 * time.Second)
rooms.Count <- 0 rooms.Count <- 0
l := <-rooms.Count l := <-rooms.Count
msg, err := json.Marshal(AnnounceCount{Count: l}) msg, err := json.Marshal(struct {
Count int `json:"count"`
}{Count: l})
if err == nil { if err == nil {
hub.broadcast <- msg ws.Send <- msg
} }
} }
} }
@ -118,10 +119,16 @@ func getAMF(room string) (bool, *AuthResponse) {
return true, v return true, v
} }
func statRoom(chQuit chan struct{}, room, server, proxy string, info *tID, u url.URL) { func xWorker(workerData Info, u url.URL) {
//fmt.Println("Start", room, "server", server, "proxy", proxy) 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 { if !ok {
fmt.Println("exit: no amf parms") fmt.Println("exit: no amf parms")
return return
@ -129,85 +136,62 @@ func statRoom(chQuit chan struct{}, room, server, proxy string, info *tID, u url
Dialer := *websocket.DefaultDialer Dialer := *websocket.DefaultDialer
if _, ok := conf.Proxy[proxy]; ok { if _, ok := conf.Proxy[workerData.Proxy]; ok {
Dialer = websocket.Dialer{ Dialer = websocket.Dialer{
Proxy: http.ProxyURL(&url.URL{ Proxy: http.ProxyURL(&url.URL{
Scheme: "http", // or "https" depending on your proxy Scheme: "http", // or "https" depending on your proxy
Host: conf.Proxy[proxy], Host: conf.Proxy[workerData.Proxy],
Path: "/", Path: "/",
}), }),
HandshakeTimeout: 45 * time.Second, // https://pkg.go.dev/github.com/gorilla/websocket 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) c, _, err := Dialer.Dial(u.String(), nil)
if err != nil { if err != nil {
fmt.Println(err.Error(), room) fmt.Println(err.Error(), workerData.room)
return return
} }
defer func() { defer c.Close()
fmt.Println("defer close", room)
c.Close()
}()
c.SetReadDeadline(time.Now().Add(60 * time.Second)) 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 { 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()) fmt.Println(err.Error())
return return
} }
fmt.Println("read first", room)
_, message, err := c.ReadMessage() _, message, err := c.ReadMessage()
if err != nil { if err != nil {
fmt.Println(err.Error(), room) fmt.Println(err.Error(), workerData.room)
return 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)}
if string(message) == `{"id":1,"result":{"audioAvailable":false,"freeShow":false},"error":null}` { 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 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 { 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 return
} }
fmt.Println("read second", room)
_, message, err = c.ReadMessage() _, message, err = c.ReadMessage()
if err != nil { if err != nil {
fmt.Println(err.Error(), room) fmt.Println(err.Error(), workerData.room)
return 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) quit := make(chan bool)
pid := 3 pid := 3
defer func() { defer func() {
fmt.Println("defer quit", room)
quit <- true quit <- true
}() }()
@ -220,8 +204,8 @@ func statRoom(chQuit chan struct{}, room, server, proxy string, info *tID, u url
return return
case <-ticker.C: case <-ticker.C:
if err = c.WriteMessage(websocket.TextMessage, []byte(fmt.Sprintf(`{"id":%d,"name":"ping"}`, pid))); err != nil { if err = c.WriteMessage(websocket.TextMessage, []byte(fmt.Sprintf(`{"id":%d,"name":"ping"}`, pid))); err != nil {
fmt.Println(err.Error(), room) fmt.Println(err.Error(), workerData.room)
close(chWorker.Map[room].chQuit) rooms.Stop <- workerData.room
return return
} }
pid++ pid++
@ -232,11 +216,12 @@ func statRoom(chQuit chan struct{}, room, server, proxy string, info *tID, u url
for { for {
select { select {
case <-chQuit: case <-workerData.ch:
fmt.Println("Exit room:", room) fmt.Println("Exit room:", workerData.room)
return return
default: default:
}
c.SetReadDeadline(time.Now().Add(30 * time.Minute)) c.SetReadDeadline(time.Now().Add(30 * time.Minute))
_, message, err := c.ReadMessage() _, message, err := c.ReadMessage()
if err != nil { if err != nil {
@ -244,14 +229,14 @@ func statRoom(chQuit chan struct{}, room, server, proxy string, info *tID, u url
return return
} }
now = time.Now().Unix() now := time.Now().Unix()
slog <- saveLog{info.Id, now, string(message)} slog <- saveLog{workerData.Rid, now, string(message)}
m := &ServerResponse{} m := &ServerResponse{}
if err = json.Unmarshal(message, m); err != nil { if err = json.Unmarshal(message, m); err != nil {
fmt.Println(err.Error(), room) fmt.Println(err.Error(), workerData.room)
continue continue
} }
@ -259,12 +244,12 @@ func statRoom(chQuit chan struct{}, room, server, proxy string, info *tID, u url
rooms.Add <- workerData rooms.Add <- workerData
if m.Type == "ServerMessageEvent:PERFORMER_STATUS_CHANGE" && string(m.Body) == `"offline"` { if m.Type == "ServerMessageEvent:PERFORMER_STATUS_CHANGE" && string(m.Body) == `"offline"` {
fmt.Println(m.Type, room) fmt.Println(m.Type, workerData.room)
return return
} }
if m.Type == "ServerMessageEvent:ROOM_CLOSE" { if m.Type == "ServerMessageEvent:ROOM_CLOSE" {
fmt.Println(m.Type, room) fmt.Println(m.Type, workerData.room)
return return
} }
@ -273,7 +258,7 @@ func statRoom(chQuit chan struct{}, room, server, proxy string, info *tID, u url
if err = json.Unmarshal(m.Body, d); err == nil { if err = json.Unmarshal(m.Body, d); err == nil {
//fmt.Println(d.F.Username, "send", d.A, "tokens") //fmt.Println(d.F.Username, "send", d.A, "tokens")
save <- saveData{room, d.F.Username, info.Id, d.A, now} save <- saveData{workerData.room, d.F.Username, workerData.Rid, d.A, now}
workerData.Income += d.A workerData.Income += d.A
rooms.Add <- workerData rooms.Add <- workerData
@ -281,4 +266,3 @@ func statRoom(chQuit chan struct{}, room, server, proxy string, info *tID, u url
} }
} }
} }
}