mirror of
https://github.com/sipeed/NanoKVM.git
synced 2026-09-11 00:22:56 -05:00
feat: show HDMI capture status on desktop screen
Track capture results for each stream mode, broadcast status updates over WebSocket, and display localized warning/error overlays in the desktop view.
This commit is contained in:
167
server/service/stream/capture_status.go
Normal file
167
server/service/stream/capture_status.go
Normal file
@@ -0,0 +1,167 @@
|
||||
package stream
|
||||
|
||||
import (
|
||||
"sort"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
const (
|
||||
CaptureStatusEvent = "capture-status"
|
||||
CaptureModeDirect = "direct"
|
||||
CaptureModeH264 = "h264"
|
||||
CaptureModeMJPEG = "mjpeg"
|
||||
CaptureSeverityError = "error"
|
||||
CaptureSeverityWarning = "warning"
|
||||
)
|
||||
|
||||
type CaptureStatus struct {
|
||||
Ok bool `json:"ok"`
|
||||
Result int `json:"result"`
|
||||
Message string `json:"message"`
|
||||
Mode string `json:"mode"`
|
||||
Severity string `json:"severity"`
|
||||
UpdatedAt time.Time `json:"updatedAt"`
|
||||
}
|
||||
|
||||
type CaptureStatusSubscriber func(CaptureStatus)
|
||||
|
||||
var defaultStore = newStore(time.Now)
|
||||
|
||||
func SubscribeCaptureStatus(subscriber CaptureStatusSubscriber) func() {
|
||||
return defaultStore.Subscribe(subscriber)
|
||||
}
|
||||
|
||||
func UpdateCaptureStatus(mode string, result int) {
|
||||
defaultStore.UpdateCaptureStatus(mode, result)
|
||||
}
|
||||
|
||||
func LatestCaptureStatuses() []CaptureStatus {
|
||||
return defaultStore.LatestCaptureStatuses()
|
||||
}
|
||||
|
||||
type store struct {
|
||||
mutex sync.Mutex
|
||||
latestByMode map[string]CaptureStatus
|
||||
now func() time.Time
|
||||
subscribers map[int]CaptureStatusSubscriber
|
||||
nextSubscriberID int
|
||||
}
|
||||
|
||||
func newStore(now func() time.Time) *store {
|
||||
return &store{
|
||||
now: now,
|
||||
latestByMode: make(map[string]CaptureStatus),
|
||||
subscribers: make(map[int]CaptureStatusSubscriber),
|
||||
}
|
||||
}
|
||||
|
||||
func (s *store) Subscribe(subscriber CaptureStatusSubscriber) func() {
|
||||
s.mutex.Lock()
|
||||
id := s.nextSubscriberID
|
||||
s.nextSubscriberID++
|
||||
s.subscribers[id] = subscriber
|
||||
s.mutex.Unlock()
|
||||
|
||||
return func() {
|
||||
s.mutex.Lock()
|
||||
delete(s.subscribers, id)
|
||||
s.mutex.Unlock()
|
||||
}
|
||||
}
|
||||
|
||||
func (s *store) UpdateCaptureStatus(mode string, result int) {
|
||||
next := newCaptureStatus(mode, result, s.now())
|
||||
subscribers := s.update(next)
|
||||
for _, subscriber := range subscribers {
|
||||
subscriber(next)
|
||||
}
|
||||
}
|
||||
|
||||
func (s *store) update(next CaptureStatus) []CaptureStatusSubscriber {
|
||||
s.mutex.Lock()
|
||||
defer s.mutex.Unlock()
|
||||
|
||||
last, ok := s.latestByMode[next.Mode]
|
||||
if ok && samePublicStatus(last, next) {
|
||||
return nil
|
||||
}
|
||||
|
||||
s.latestByMode[next.Mode] = next
|
||||
|
||||
subscribers := make([]CaptureStatusSubscriber, 0, len(s.subscribers))
|
||||
for _, subscriber := range s.subscribers {
|
||||
subscribers = append(subscribers, subscriber)
|
||||
}
|
||||
return subscribers
|
||||
}
|
||||
|
||||
func (s *store) LatestCaptureStatus(mode string) (CaptureStatus, bool) {
|
||||
s.mutex.Lock()
|
||||
defer s.mutex.Unlock()
|
||||
|
||||
status, ok := s.latestByMode[mode]
|
||||
return status, ok
|
||||
}
|
||||
|
||||
func (s *store) LatestCaptureStatuses() []CaptureStatus {
|
||||
s.mutex.Lock()
|
||||
defer s.mutex.Unlock()
|
||||
|
||||
modes := make([]string, 0, len(s.latestByMode))
|
||||
for mode := range s.latestByMode {
|
||||
modes = append(modes, mode)
|
||||
}
|
||||
sort.Strings(modes)
|
||||
|
||||
statuses := make([]CaptureStatus, 0, len(modes))
|
||||
for _, mode := range modes {
|
||||
statuses = append(statuses, s.latestByMode[mode])
|
||||
}
|
||||
|
||||
return statuses
|
||||
}
|
||||
|
||||
func samePublicStatus(a CaptureStatus, b CaptureStatus) bool {
|
||||
if a.Ok && b.Ok {
|
||||
return true
|
||||
}
|
||||
|
||||
return a.Ok == b.Ok && a.Result == b.Result && a.Mode == b.Mode
|
||||
}
|
||||
|
||||
func newCaptureStatus(mode string, result int, updatedAt time.Time) CaptureStatus {
|
||||
message, severity := captureResultMessage(result)
|
||||
return CaptureStatus{
|
||||
Ok: result >= 0,
|
||||
Result: result,
|
||||
Message: message,
|
||||
Mode: mode,
|
||||
Severity: severity,
|
||||
UpdatedAt: updatedAt,
|
||||
}
|
||||
}
|
||||
|
||||
func captureResultMessage(result int) (string, string) {
|
||||
switch result {
|
||||
case -7:
|
||||
return "HDMI input resolution error", CaptureSeverityError
|
||||
case -6:
|
||||
return "Unsupported HDMI resolution", CaptureSeverityError
|
||||
case -5:
|
||||
return "Retrieving image", CaptureSeverityWarning
|
||||
case -4:
|
||||
return "Changing image resolution", CaptureSeverityWarning
|
||||
case -3:
|
||||
return "Image buffer full", CaptureSeverityError
|
||||
case -2:
|
||||
return "Encoder error", CaptureSeverityError
|
||||
case -1:
|
||||
return "No image captured", CaptureSeverityError
|
||||
default:
|
||||
if result < 0 {
|
||||
return "Capture failed", CaptureSeverityError
|
||||
}
|
||||
return "", ""
|
||||
}
|
||||
}
|
||||
@@ -69,6 +69,7 @@ func (s *Streamer) run() {
|
||||
}
|
||||
|
||||
data, result := vision.ReadH264(screen.Width, screen.Height, screen.BitRate)
|
||||
stream.UpdateCaptureStatus(stream.CaptureModeDirect, result)
|
||||
if result < 0 || len(data) == 0 {
|
||||
continue
|
||||
}
|
||||
|
||||
@@ -85,6 +85,7 @@ func (s *Streamer) run() {
|
||||
}
|
||||
|
||||
data, result := vision.ReadMjpeg(screen.Width, screen.Height, screen.Quality)
|
||||
stream.UpdateCaptureStatus(stream.CaptureModeMJPEG, result)
|
||||
if result < 0 || result == 5 || len(data) == 0 {
|
||||
continue
|
||||
}
|
||||
|
||||
@@ -72,6 +72,7 @@ func (m *WebRTCManager) sendVideoStream() {
|
||||
}
|
||||
|
||||
data, result := vision.ReadH264(screen.Width, screen.Height, screen.BitRate)
|
||||
stream.UpdateCaptureStatus(stream.CaptureModeH264, result)
|
||||
if result < 0 || len(data) == 0 {
|
||||
continue
|
||||
}
|
||||
|
||||
111
server/service/ws/capture_status.go
Normal file
111
server/service/ws/capture_status.go
Normal file
@@ -0,0 +1,111 @@
|
||||
package ws
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"sort"
|
||||
"sync"
|
||||
|
||||
"NanoKVM-Server/service/stream"
|
||||
|
||||
log "github.com/sirupsen/logrus"
|
||||
)
|
||||
|
||||
var captureBroadcaster = newCaptureStatusBroadcaster(broadcastCaptureStatus)
|
||||
|
||||
func init() {
|
||||
go captureBroadcaster.Run()
|
||||
|
||||
stream.SubscribeCaptureStatus(func(status stream.CaptureStatus) {
|
||||
captureBroadcaster.Enqueue(status)
|
||||
})
|
||||
}
|
||||
|
||||
type captureStatusBroadcaster struct {
|
||||
mutex sync.Mutex
|
||||
pending map[string]stream.CaptureStatus
|
||||
notify chan struct{}
|
||||
broadcast func(stream.CaptureStatus)
|
||||
}
|
||||
|
||||
func newCaptureStatusBroadcaster(broadcast func(stream.CaptureStatus)) *captureStatusBroadcaster {
|
||||
return &captureStatusBroadcaster{
|
||||
pending: make(map[string]stream.CaptureStatus),
|
||||
notify: make(chan struct{}, 1),
|
||||
broadcast: broadcast,
|
||||
}
|
||||
}
|
||||
|
||||
func (b *captureStatusBroadcaster) Enqueue(status stream.CaptureStatus) {
|
||||
b.mutex.Lock()
|
||||
b.pending[status.Mode] = status
|
||||
b.mutex.Unlock()
|
||||
|
||||
select {
|
||||
case b.notify <- struct{}{}:
|
||||
default:
|
||||
}
|
||||
}
|
||||
|
||||
func (b *captureStatusBroadcaster) Run() {
|
||||
for range b.notify {
|
||||
for {
|
||||
statuses := b.takePending()
|
||||
if len(statuses) == 0 {
|
||||
break
|
||||
}
|
||||
|
||||
for _, status := range statuses {
|
||||
b.broadcast(status)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (b *captureStatusBroadcaster) takePending() []stream.CaptureStatus {
|
||||
b.mutex.Lock()
|
||||
defer b.mutex.Unlock()
|
||||
|
||||
if len(b.pending) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
modes := make([]string, 0, len(b.pending))
|
||||
for mode := range b.pending {
|
||||
modes = append(modes, mode)
|
||||
}
|
||||
sort.Strings(modes)
|
||||
|
||||
statuses := make([]stream.CaptureStatus, 0, len(modes))
|
||||
for _, mode := range modes {
|
||||
statuses = append(statuses, b.pending[mode])
|
||||
delete(b.pending, mode)
|
||||
}
|
||||
|
||||
return statuses
|
||||
}
|
||||
|
||||
func sendCaptureStatusSnapshot(client *Client) {
|
||||
for _, status := range stream.LatestCaptureStatuses() {
|
||||
if err := sendCaptureStatus(client, status); err != nil {
|
||||
log.Errorf("failed to send capture status snapshot: %s", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func broadcastCaptureStatus(status stream.CaptureStatus) {
|
||||
for _, client := range GetManager().GetClients() {
|
||||
if err := sendCaptureStatus(client, status); err != nil {
|
||||
log.Errorf("failed to send capture status: %s", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func sendCaptureStatus(client *Client, status stream.CaptureStatus) error {
|
||||
payload, err := json.Marshal(status)
|
||||
if err != nil {
|
||||
log.Errorf("failed to marshal capture status: %s", err)
|
||||
return err
|
||||
}
|
||||
|
||||
return client.Write(stream.CaptureStatusEvent, string(payload))
|
||||
}
|
||||
@@ -37,5 +37,7 @@ func (s *Service) Connect(c *gin.Context) {
|
||||
manager.AddClient(ws, client)
|
||||
defer manager.RemoveClient(ws)
|
||||
|
||||
sendCaptureStatusSnapshot(client)
|
||||
|
||||
client.Start()
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user