Files
hockey_new/agent/main.go
2026-08-19 15:08:39 +03:00

660 lines
17 KiB
Go

package main
import (
"bufio"
"bytes"
"crypto/rand"
"crypto/sha1"
"crypto/tls"
"encoding/base64"
"encoding/binary"
"encoding/json"
"errors"
"flag"
"fmt"
"io"
"net"
"net/http"
"net/url"
"os"
"path/filepath"
"runtime"
"strings"
"sync"
"time"
)
const (
agentVersion = "1.1.0"
protocolVersion = 1
websocketGUID = "258EAFA5-E914-47DA-95CA-C5AB0DC85B11"
heartbeatInterval = 8 * time.Second
)
type Config struct {
ServerURL string `json:"server_url"`
VmixURL string `json:"vmix_url"`
DeviceID string `json:"device_id"`
DeviceSecret string `json:"device_secret"`
DeviceName string `json:"device_name"`
}
type VmixStatus struct {
Connected bool `json:"connected"`
Version string `json:"version,omitempty"`
URL string `json:"url"`
Error string `json:"error,omitempty"`
Response string `json:"response,omitempty"`
}
type Agent struct {
cfg Config
mu sync.Mutex
assignmentID string
matchID string
pairedLogin string
active bool
}
type wsConn struct {
conn net.Conn
r *bufio.Reader
mu sync.Mutex
}
func main() {
cfgPath, err := configPath()
if err != nil {
fatal(err)
}
cfg, err := loadConfig(cfgPath)
if err != nil {
fatal(err)
}
server := flag.String("server", "", "Hockey server URL")
vmix := flag.String("vmix", "", "vMix API URL")
name := flag.String("name", "", "Device display name")
resetID := flag.Bool("reset-device", false, "Generate a new Device ID and secret")
flag.Parse()
changed := false
if *server != "" {
cfg.ServerURL = strings.TrimRight(strings.TrimSpace(*server), "/")
changed = true
}
if *vmix != "" {
cfg.VmixURL = strings.TrimSpace(*vmix)
changed = true
}
if *name != "" {
cfg.DeviceName = strings.TrimSpace(*name)
changed = true
}
if *resetID {
cfg.DeviceID = newUUID()
cfg.DeviceSecret = randomToken(40)
changed = true
}
if changed {
if err := saveConfig(cfgPath, cfg); err != nil {
fatal(err)
}
}
fmt.Println("============================================================")
fmt.Println(" HOCKEY vMix AGENT")
fmt.Printf(" Version: %s | Protocol: %d\n", agentVersion, protocolVersion)
fmt.Printf(" Device ID: %s\n", cfg.DeviceID)
fmt.Printf(" Name: %s\n", cfg.DeviceName)
fmt.Printf(" Server: %s\n", cfg.ServerURL)
fmt.Printf(" vMix API: %s\n", cfg.VmixURL)
fmt.Printf(" Config: %s\n", cfgPath)
fmt.Println("============================================================")
fmt.Println("После подключения откройте web-интерфейс → vMix Agent → Прикрепить.")
fmt.Println()
agent := &Agent{cfg: cfg}
backoff := time.Second
for {
if err := agent.runConnection(); err != nil {
fmt.Printf("[%s] SERVER OFFLINE: %v\n", stamp(), err)
}
time.Sleep(backoff)
if backoff < 15*time.Second {
backoff *= 2
}
if backoff > 15*time.Second {
backoff = 15 * time.Second
}
}
}
func (a *Agent) runConnection() error {
target, err := websocketURL(a.cfg.ServerURL)
if err != nil {
return err
}
fmt.Printf("[%s] Connecting %s ...\n", stamp(), target)
ws, err := dialWebSocket(target, 7*time.Second)
if err != nil {
return err
}
defer ws.Close()
fmt.Printf("[%s] SERVER ONLINE\n", stamp())
hello := map[string]any{
"type": "hello",
"protocol": protocolVersion,
"device_id": a.cfg.DeviceID,
"device_secret": a.cfg.DeviceSecret,
"device_name": a.cfg.DeviceName,
"hostname": hostname(),
"agent_version": agentVersion,
"platform": runtime.GOOS + "/" + runtime.GOARCH,
"vmix": probeVmix(a.cfg.VmixURL, nil),
}
if err := ws.WriteJSON(hello); err != nil {
return err
}
lastHeartbeat := time.Time{}
for {
if time.Since(lastHeartbeat) >= heartbeatInterval {
vmix := probeVmix(a.cfg.VmixURL, nil)
a.printVmix(vmix)
a.mu.Lock()
assignmentID, matchID := a.assignmentID, a.matchID
a.mu.Unlock()
_ = ws.WriteJSON(map[string]any{
"type": "heartbeat",
"device_id": a.cfg.DeviceID,
"assignment_id": assignmentID,
"match_id": matchID,
"vmix": vmix,
"error": vmix.Error,
})
lastHeartbeat = time.Now()
}
_ = ws.conn.SetReadDeadline(time.Now().Add(1200 * time.Millisecond))
payload, err := ws.ReadText()
if err != nil {
if ne, ok := err.(net.Error); ok && ne.Timeout() {
continue
}
return err
}
var msg map[string]any
if err := json.Unmarshal(payload, &msg); err != nil {
continue
}
if err := a.handleMessage(ws, msg); err != nil {
fmt.Printf("[%s] MESSAGE ERROR: %v\n", stamp(), err)
}
}
}
func (a *Agent) handleMessage(ws *wsConn, msg map[string]any) error {
typ := stringValue(msg["type"])
switch typ {
case "hello.ok":
fmt.Printf("[%s] Device registered on server\n", stamp())
case "hello.error":
return errors.New(stringValue(msg["detail"]))
case "pairing.confirmed":
account, _ := msg["account"].(map[string]any)
login := stringValue(account["login"])
active := boolValue(msg["active_for_account"])
a.mu.Lock()
a.pairedLogin = login
a.active = active
a.mu.Unlock()
suffix := ""
if active {
suffix = " (ACTIVE)"
}
fmt.Printf("[%s] PAIRED: %s%s\n", stamp(), login, suffix)
case "pairing.revoked":
a.mu.Lock()
a.pairedLogin = ""
a.active = false
a.assignmentID = ""
a.matchID = ""
a.mu.Unlock()
fmt.Printf("[%s] Pairing revoked\n", stamp())
case "device.activated":
a.mu.Lock()
a.active = true
a.mu.Unlock()
fmt.Printf("[%s] Device activated for broadcast\n", stamp())
case "device.deactivated":
a.mu.Lock()
a.active = false
a.assignmentID = ""
a.matchID = ""
a.mu.Unlock()
fmt.Printf("[%s] Device deactivated; match assignment cleared\n", stamp())
case "match.assign":
assignmentID := stringValue(msg["assignment_id"])
matchID := stringValue(msg["game_id"])
a.mu.Lock()
a.assignmentID = assignmentID
a.matchID = matchID
a.mu.Unlock()
fmt.Printf("[%s] MATCH ASSIGNED: %s | assignment %s\n", stamp(), matchID, short(assignmentID))
vmix := probeVmix(a.cfg.VmixURL, nil)
return ws.WriteJSON(map[string]any{
"type": "match.accepted",
"device_id": a.cfg.DeviceID,
"assignment_id": assignmentID,
"match_id": matchID,
"vmix": vmix,
})
case "vmix.probe":
vmix := probeVmix(a.cfg.VmixURL, nil)
a.printVmix(vmix)
return ws.WriteJSON(map[string]any{
"type": "vmix.status",
"request_id": stringValue(msg["request_id"]),
"vmix": vmix,
"error": vmix.Error,
})
case "vmix.command":
return a.handleVmixCommand(ws, msg)
case "heartbeat.ack":
return nil
case "server.error":
fmt.Printf("[%s] SERVER ERROR: %s\n", stamp(), stringValue(msg["detail"]))
}
return nil
}
func (a *Agent) handleVmixCommand(ws *wsConn, msg map[string]any) error {
requestID := stringValue(msg["request_id"])
incomingAssignment := stringValue(msg["assignment_id"])
incomingMatch := stringValue(msg["match_id"])
a.mu.Lock()
currentAssignment, currentMatch := a.assignmentID, a.matchID
a.mu.Unlock()
if incomingAssignment != "" && incomingAssignment != currentAssignment {
return ws.WriteJSON(map[string]any{"type": "command.ack", "request_id": requestID, "ok": false, "reason": "assignment_mismatch", "assignment_id": currentAssignment, "match_id": currentMatch})
}
if incomingMatch != "" && incomingMatch != currentMatch {
return ws.WriteJSON(map[string]any{"type": "command.ack", "request_id": requestID, "ok": false, "reason": "match_mismatch", "assignment_id": currentAssignment, "match_id": currentMatch})
}
rawCommand, _ := msg["command"].(map[string]any)
function := stringValue(rawCommand["Function"])
if function == "" {
function = stringValue(rawCommand["function"])
}
if function == "" {
return ws.WriteJSON(map[string]any{"type": "command.ack", "request_id": requestID, "ok": false, "reason": "missing_function"})
}
params := url.Values{}
params.Set("Function", function)
for key, value := range rawCommand {
if strings.EqualFold(key, "function") || value == nil {
continue
}
text := stringValue(value)
if text != "" {
params.Set(key, text)
}
}
fmt.Printf("[%s] vMix COMMAND: %s | Input=%s | SelectedName=%s\n", stamp(), function, params.Get("Input"), params.Get("SelectedName"))
vmix := probeVmix(a.cfg.VmixURL, params)
a.printVmix(vmix)
if vmix.Connected {
fmt.Printf("[%s] vMix COMMAND OK: %s | response=%s\n", stamp(), function, short(vmix.Response))
}
return ws.WriteJSON(map[string]any{
"type": "command.ack", "request_id": requestID, "ok": vmix.Connected,
"reason": vmix.Error, "assignment_id": currentAssignment, "match_id": currentMatch,
"vmix": vmix, "error": vmix.Error,
})
}
func (a *Agent) printVmix(status VmixStatus) {
if status.Connected {
suffix := ""
if status.Version != "" {
suffix = " " + status.Version
}
fmt.Printf("[%s] vMix CONNECTED%s\n", stamp(), suffix)
} else if status.Error != "" {
fmt.Printf("[%s] vMix NO CONNECTION: %s\n", stamp(), status.Error)
}
}
func probeVmix(base string, params url.Values) VmixStatus {
result := VmixStatus{URL: base}
target := strings.TrimSpace(base)
if params != nil && len(params) > 0 {
sep := "?"
if strings.Contains(target, "?") {
sep = "&"
}
target += sep + params.Encode()
}
client := &http.Client{Timeout: 3 * time.Second}
resp, err := client.Get(target)
if err != nil {
result.Error = err.Error()
return result
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
result.Error = fmt.Sprintf("HTTP %d", resp.StatusCode)
return result
}
body, _ := io.ReadAll(io.LimitReader(resp.Body, 64*1024))
result.Connected = true
text := string(body)
if params != nil && len(params) > 0 {
responseText := strings.TrimSpace(text)
if len(responseText) > 1000 {
responseText = responseText[:1000]
}
result.Response = responseText
}
if i := strings.Index(text, "<version>"); i >= 0 {
rest := text[i+len("<version>"):]
if j := strings.Index(rest, "</version>"); j >= 0 {
result.Version = strings.TrimSpace(rest[:j])
}
}
return result
}
func websocketURL(server string) (string, error) {
raw := strings.TrimRight(strings.TrimSpace(server), "/")
if raw == "" {
return "", errors.New("server_url is empty")
}
u, err := url.Parse(raw)
if err != nil {
return "", err
}
switch strings.ToLower(u.Scheme) {
case "http":
u.Scheme = "ws"
case "https":
u.Scheme = "wss"
case "ws", "wss":
default:
if !strings.Contains(raw, "://") {
return websocketURL("http://" + raw)
}
return "", fmt.Errorf("unsupported server scheme: %s", u.Scheme)
}
u.Path = strings.TrimRight(u.Path, "/") + "/ws/hockey-agent"
u.RawQuery = ""
return u.String(), nil
}
func dialWebSocket(target string, timeout time.Duration) (*wsConn, error) {
u, err := url.Parse(target)
if err != nil {
return nil, err
}
host := u.Hostname()
port := u.Port()
if port == "" {
if u.Scheme == "wss" {
port = "443"
} else {
port = "80"
}
}
address := net.JoinHostPort(host, port)
d := net.Dialer{Timeout: timeout}
var conn net.Conn
if u.Scheme == "wss" {
conn, err = tls.DialWithDialer(&d, "tcp", address, &tls.Config{ServerName: host, MinVersion: tls.VersionTLS12})
} else {
conn, err = d.Dial("tcp", address)
}
if err != nil {
return nil, err
}
key := randomBase64(16)
path := u.EscapedPath()
if path == "" {
path = "/"
}
if u.RawQuery != "" {
path += "?" + u.RawQuery
}
hostHeader := u.Host
req := fmt.Sprintf("GET %s HTTP/1.1\r\nHost: %s\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Key: %s\r\nSec-WebSocket-Version: 13\r\nUser-Agent: Hockey-vMix-Agent/%s\r\n\r\n", path, hostHeader, key, agentVersion)
if _, err := io.WriteString(conn, req); err != nil {
conn.Close()
return nil, err
}
reader := bufio.NewReader(conn)
statusLine, err := reader.ReadString('\n')
if err != nil {
conn.Close()
return nil, err
}
if !strings.Contains(statusLine, " 101 ") {
conn.Close()
return nil, fmt.Errorf("websocket upgrade failed: %s", strings.TrimSpace(statusLine))
}
headers := map[string]string{}
for {
line, err := reader.ReadString('\n')
if err != nil {
conn.Close()
return nil, err
}
line = strings.TrimRight(line, "\r\n")
if line == "" {
break
}
if i := strings.Index(line, ":"); i > 0 {
headers[strings.ToLower(strings.TrimSpace(line[:i]))] = strings.TrimSpace(line[i+1:])
}
}
sum := sha1.Sum([]byte(key + websocketGUID))
expected := base64.StdEncoding.EncodeToString(sum[:])
if headers["sec-websocket-accept"] != expected {
conn.Close()
return nil, errors.New("invalid Sec-WebSocket-Accept")
}
return &wsConn{conn: conn, r: reader}, nil
}
func (w *wsConn) Close() error { return w.conn.Close() }
func (w *wsConn) WriteJSON(value any) error {
payload, err := json.Marshal(value)
if err != nil {
return err
}
return w.writeFrame(0x1, payload)
}
func (w *wsConn) writeFrame(opcode byte, payload []byte) error {
w.mu.Lock()
defer w.mu.Unlock()
var header bytes.Buffer
header.WriteByte(0x80 | opcode)
n := len(payload)
switch {
case n < 126:
header.WriteByte(0x80 | byte(n))
case n <= 65535:
header.WriteByte(0x80 | 126)
_ = binary.Write(&header, binary.BigEndian, uint16(n))
default:
header.WriteByte(0x80 | 127)
_ = binary.Write(&header, binary.BigEndian, uint64(n))
}
mask := make([]byte, 4)
if _, err := rand.Read(mask); err != nil {
return err
}
header.Write(mask)
masked := make([]byte, n)
for i := range payload {
masked[i] = payload[i] ^ mask[i%4]
}
if _, err := w.conn.Write(header.Bytes()); err != nil {
return err
}
_, err := w.conn.Write(masked)
return err
}
func (w *wsConn) ReadText() ([]byte, error) {
for {
first, err := w.r.ReadByte()
if err != nil {
return nil, err
}
second, err := w.r.ReadByte()
if err != nil {
return nil, err
}
fin := first&0x80 != 0
opcode := first & 0x0F
masked := second&0x80 != 0
length := uint64(second & 0x7F)
if length == 126 {
var n uint16
if err := binary.Read(w.r, binary.BigEndian, &n); err != nil {
return nil, err
}
length = uint64(n)
} else if length == 127 {
if err := binary.Read(w.r, binary.BigEndian, &length); err != nil {
return nil, err
}
}
if length > 16*1024*1024 {
return nil, errors.New("websocket frame too large")
}
var mask [4]byte
if masked {
if _, err := io.ReadFull(w.r, mask[:]); err != nil {
return nil, err
}
}
payload := make([]byte, int(length))
if _, err := io.ReadFull(w.r, payload); err != nil {
return nil, err
}
if masked {
for i := range payload {
payload[i] ^= mask[i%4]
}
}
switch opcode {
case 0x1:
if !fin {
return nil, errors.New("fragmented text frames are not supported")
}
return payload, nil
case 0x8:
return nil, io.EOF
case 0x9:
_ = w.writeFrame(0xA, payload)
case 0xA:
continue
default:
continue
}
}
}
func configPath() (string, error) {
exe, err := os.Executable()
if err != nil {
return "", err
}
return filepath.Join(filepath.Dir(exe), "agent_config.json"), nil
}
func loadConfig(path string) (Config, error) {
cfg := Config{ServerURL: "http://127.0.0.1:8000", VmixURL: "http://127.0.0.1:8088/api/", DeviceID: newUUID(), DeviceSecret: randomToken(40), DeviceName: hostname()}
if data, err := os.ReadFile(path); err == nil {
_ = json.Unmarshal(data, &cfg)
}
if cfg.DeviceID == "" {
cfg.DeviceID = newUUID()
}
if cfg.DeviceSecret == "" {
cfg.DeviceSecret = randomToken(40)
}
if cfg.DeviceName == "" {
cfg.DeviceName = hostname()
}
if cfg.ServerURL == "" {
cfg.ServerURL = "http://127.0.0.1:8000"
}
if cfg.VmixURL == "" {
cfg.VmixURL = "http://127.0.0.1:8088/api/"
}
if err := saveConfig(path, cfg); err != nil {
return cfg, err
}
return cfg, nil
}
func saveConfig(path string, cfg Config) error {
data, err := json.MarshalIndent(cfg, "", " ")
if err != nil {
return err
}
return os.WriteFile(path, data, 0600)
}
func newUUID() string {
b := make([]byte, 16)
_, _ = rand.Read(b)
b[6] = (b[6] & 0x0f) | 0x40
b[8] = (b[8] & 0x3f) | 0x80
return fmt.Sprintf("%08x-%04x-%04x-%04x-%012x", b[0:4], b[4:6], b[6:8], b[8:10], b[10:16])
}
func randomBase64(n int) string {
b := make([]byte, n)
_, _ = rand.Read(b)
return base64.StdEncoding.EncodeToString(b)
}
func randomToken(n int) string {
b := make([]byte, n)
_, _ = rand.Read(b)
return base64.RawURLEncoding.EncodeToString(b)
}
func hostname() string {
h, _ := os.Hostname()
if h == "" {
return "HOCKEY-GFX"
}
return h
}
func stamp() string { return time.Now().Format("15:04:05") }
func short(v string) string {
if len(v) > 12 {
return v[:12] + "..."
}
return v
}
func stringValue(v any) string {
if v == nil {
return ""
}
if s, ok := v.(string); ok {
return s
}
return fmt.Sprint(v)
}
func boolValue(v any) bool { b, _ := v.(bool); return b }
func fatal(err error) { fmt.Fprintln(os.Stderr, "FATAL:", err); os.Exit(1) }