diff --git a/server/service/stream/h264/client.go b/server/service/stream/h264/client.go deleted file mode 100644 index 1d72a0d..0000000 --- a/server/service/stream/h264/client.go +++ /dev/null @@ -1,167 +0,0 @@ -package h264 - -import ( - "encoding/json" - "sync" - - "github.com/gorilla/websocket" - "github.com/pion/webrtc/v4" - log "github.com/sirupsen/logrus" -) - -type Client struct { - ws *websocket.Conn - pc *webrtc.PeerConnection - mutex sync.Mutex -} - -type Message struct { - Event string `json:"event"` - Data string `json:"data"` -} - -// add video track -func (c *Client) addTrack() { - videoTrack, err := webrtc.NewTrackLocalStaticSample( - webrtc.RTPCodecCapability{MimeType: webrtc.MimeTypeH264}, - "video", - "pion", - ) - if err != nil { - log.Errorf("failed to create video track: %s", err) - return - } - - _, err = c.pc.AddTrack(videoTrack) - if err != nil { - log.Errorf("failed to add video track: %s", err) - return - } - - trackMap[c.ws] = videoTrack -} - -// register callback events -func (c *Client) register() { - // new ICE candidate found - c.pc.OnICECandidate(func(candidate *webrtc.ICECandidate) { - if candidate == nil { - return - } - - candidateByte, err := json.Marshal(candidate.ToJSON()) - if err != nil { - log.Errorf("failed to marshal candidate: %s", err) - return - } - - _ = c.sendMessage("candidate", string(candidateByte)) - }) - - // ICE connection state has changed - c.pc.OnICEConnectionStateChange(func(state webrtc.ICEConnectionState) { - if state == webrtc.ICEConnectionStateConnected && !isSending { - // start sending h264 data - go send() - isSending = true - } - - log.Debugf("ice connection state has changed to %s", state.String()) - }) -} - -// read websocket message -func (c *Client) readMessage() { - message := &Message{} - - for { - _, raw, err := c.ws.ReadMessage() - if err != nil { - delete(trackMap, c.ws) - if isSending && len(trackMap) == 0 { - // stop sending when all websocket connections are closed - isSending = false - } - - log.Debugf("failed to read message: %s", err) - return - } - - if err := json.Unmarshal(raw, &message); err != nil { - log.Errorf("failed to unmarshal message: %s", err) - continue - } - - log.Debugf("receive message event: %s", message.Event) - - switch message.Event { - case "offer": - offer := webrtc.SessionDescription{} - if err := json.Unmarshal([]byte(message.Data), &offer); err != nil { - log.Errorf("failed to unmarshal offer message: %s", err) - return - } - - if err := c.pc.SetRemoteDescription(offer); err != nil { - log.Errorf("failed to set remote description: %s", err) - return - } - - answer, answerErr := c.pc.CreateAnswer(nil) - if answerErr != nil { - log.Errorf("failed to create answer: %s", answerErr) - return - } - - if err := c.pc.SetLocalDescription(answer); err != nil { - log.Errorf("failed to set local description: %s", err) - return - } - - answerByte, answerByteErr := json.Marshal(answer) - if answerByteErr != nil { - log.Errorf("failed to marshal answer: %s", answerByteErr) - return - } - - _ = c.sendMessage("answer", string(answerByte)) - - case "candidate": - candidate := webrtc.ICECandidateInit{} - if err := json.Unmarshal([]byte(message.Data), &candidate); err != nil { - log.Errorf("failed to unmarshal candidate message: %s", err) - return - } - - if err := c.pc.AddICECandidate(candidate); err != nil { - log.Errorf("failed to add ICE candidate: %s", err) - return - } - - case "heartbeat": - _ = c.sendMessage("heartbeat", "") - - default: - log.Debugf("unhandled message event: %s", message.Event) - } - } -} - -// send websocket message -func (c *Client) sendMessage(event string, data string) error { - c.mutex.Lock() - defer c.mutex.Unlock() - - message := &Message{ - Event: event, - Data: data, - } - - if err := c.ws.WriteJSON(message); err != nil { - log.Errorf("failed to send message %s: %s", event, err) - return err - } - - log.Debugf("send message %s", message.Event) - return nil -} diff --git a/server/service/stream/h264/h264.go b/server/service/stream/h264/h264.go deleted file mode 100644 index 67358de..0000000 --- a/server/service/stream/h264/h264.go +++ /dev/null @@ -1,80 +0,0 @@ -package h264 - -import ( - "NanoKVM-Server/config" - "net/http" - "sync" - "time" - - "github.com/gin-gonic/gin" - "github.com/gorilla/websocket" - "github.com/pion/webrtc/v4" - log "github.com/sirupsen/logrus" -) - -var ( - upgrader = websocket.Upgrader{ - CheckOrigin: func(r *http.Request) bool { - return true - }, - } - trackMap = make(map[*websocket.Conn]*webrtc.TrackLocalStaticSample) - isSending = false -) - -func Connect(c *gin.Context) { - wsConn, err := upgrader.Upgrade(c.Writer, c.Request, nil) - if err != nil { - log.Errorf("failed to create websocket: %s", err) - return - } - - defer func() { - _ = wsConn.Close() - log.Debugf("h264 websocket disconnected") - }() - - var zeroTime time.Time - _ = wsConn.SetReadDeadline(zeroTime) - - conf := config.GetInstance() - - var iceServers []webrtc.ICEServer - - if conf.Stun != "" && conf.Stun != "disable" { - iceServers = append(iceServers, webrtc.ICEServer{ - URLs: []string{"stun:" + conf.Stun}, - }) - } - - if conf.Turn.TurnAddr != "" && conf.Turn.TurnUser != "" && conf.Turn.TurnCred != "" { - iceServers = append(iceServers, webrtc.ICEServer{ - URLs: []string{"turn:" + conf.Turn.TurnAddr}, - Username: conf.Turn.TurnUser, - Credential: conf.Turn.TurnCred, - }) - } - - peerConn, err := webrtc.NewPeerConnection(webrtc.Configuration{ - ICEServers: iceServers, - }) - if err != nil { - log.Errorf("failed to create PeerConnection: %s", err) - return - } - - defer func() { - _ = peerConn.Close() - log.Debugf("PeerConnection disconnected") - }() - - client := &Client{ - ws: wsConn, - pc: peerConn, - mutex: sync.Mutex{}, - } - - client.addTrack() - client.register() - client.readMessage() -} diff --git a/server/service/stream/h264/sender.go b/server/service/stream/h264/sender.go deleted file mode 100644 index e59f57f..0000000 --- a/server/service/stream/h264/sender.go +++ /dev/null @@ -1,49 +0,0 @@ -package h264 - -import ( - "NanoKVM-Server/common" - "time" - - "github.com/pion/webrtc/v4/pkg/media" - log "github.com/sirupsen/logrus" -) - -func send() { - screen := common.GetScreen() - common.CheckScreen() - - fps := screen.FPS - duration := time.Second / time.Duration(fps) - - ticker := time.NewTicker(duration) - defer ticker.Stop() - - vision := common.GetKvmVision() - for range ticker.C { - if !isSending && len(trackMap) == 0 { - return - } - - data, result := vision.ReadH264(screen.Width, screen.Height, screen.BitRate) - if result < 0 || len(data) == 0 { - continue - } - - sample := media.Sample{ - Data: data, - Duration: duration, - } - - for _, track := range trackMap { - if err := track.WriteSample(sample); err != nil { - log.Errorf("failed to send h264 data: %s", err) - } - } - - if screen.FPS != fps { - fps = screen.FPS - duration = time.Second / time.Duration(fps) - ticker.Reset(duration) - } - } -}