remote-shack/internal/server/rig_server.go

208 lines
5.5 KiB
Go

package server
import (
"context"
"fmt"
"log"
"sync"
"time"
shack "remote-shack/api"
"remote-shack/internal/driver"
"google.golang.org/grpc/metadata"
)
// RigServer verwaltet den gRPC-Dienst für das Funkgerät und den Operator-Status.
type RigServer struct {
shack.UnimplementedRigServiceServer
mu sync.Mutex
driver driver.RigDriver
// Sitzungsverwaltung
activeOperator string
operatorStream shack.RigService_OperatorSessionServer
// Aktueller Hardware-Zustand für passive Zuhörer
currentFreq int64
currentMode string
currentSMeter int
}
// NewRigServer erstellt eine neue Instanz des gRPC-Rig-Servers.
func NewRigServer(d driver.RigDriver) *RigServer {
s := &RigServer{
driver: d,
currentMode: "USB",
}
// Starte die Hintergrund-Synchronisation für S-Meter und Frequenz-Updates
go s.hardwareSyncLoop()
return s
}
// SetFrequency wird vom aktiven Operator aufgerufen (Direktzugriff oder Dekadentasten)
func (s *RigServer) SetFrequency(ctx context.Context, req *shack.FrequencyRequest) (*shack.EmptyResponse, error) {
s.mu.Lock()
defer s.mu.Unlock()
// Validierung: Nur der aktive Operator darf senden
if err := s.validateOperator(ctx); err != nil {
return nil, err
}
if err := s.driver.SetFrequency(req.Hz); err != nil {
return nil, fmt.Errorf("hardware-fehler beim setzen der frequenz: %w", err)
}
s.currentFreq = req.Hz
return &shack.EmptyResponse{}, nil
}
// SetMode schaltet die Betriebsart auf Hardware-Ebene um
func (s *RigServer) SetMode(ctx context.Context, req *shack.ModeRequest) (*shack.EmptyResponse, error) {
s.mu.Lock()
defer s.mu.Unlock()
if err := s.validateOperator(ctx); err != nil {
return nil, err
}
if err := s.driver.SetMode(req.Mode); err != nil {
return nil, fmt.Errorf("hardware-fehler beim setzen des modus: %w", err)
}
s.currentMode = req.Mode
return &shack.EmptyResponse{}, nil
}
// GetStatus liefert den schnellen Zwischenspeicher-Status (wird für Fallbacks genutzt)
func (s *RigServer) GetStatus(ctx context.Context, req *shack.EmptyRequest) (*shack.StatusResponse, error) {
s.mu.Lock()
defer s.mu.Unlock()
return &shack.StatusResponse{
Hz: s.currentFreq,
Mode: s.currentMode,
SMeter: int32(s.currentSMeter),
ActiveOperatorName: s.activeOperator,
}, nil
}
// OperatorSession verwaltet den langlebigen, bidirektionalen gRPC-Stream für die Steuerung
func (s *RigServer) OperatorSession(stream shack.RigService_OperatorSessionServer) error {
// Den Betreibernamen aus dem gRPC-Metadaten-Interceptor extrahieren
md, ok := grpcMetadata.FromIncomingContext(stream.Context())
var opName string
if ok && len(md["operator-token"]) > 0 {
opName = md["operator-token"][0] // Einfachheitshalber nutzen wir den Namen als Token
} else {
opName = "Gast"
}
s.mu.Lock()
if s.activeOperator != "" && s.activeOperator != opName {
s.mu.Unlock()
return fmt.Errorf("shack ist bereits durch %s blockiert", s.activeOperator)
}
// Sitzung zuweisen
s.activeOperator = opName
s.operatorStream = stream
s.mu.Unlock()
log.Printf("[Server] %s hat die exklusive Rig-Steuerung übernommen.", opName)
defer func() {
s.mu.Lock()
s.activeOperator = ""
s.operatorStream = nil
s.mu.Unlock()
log.Printf("[Server] %s hat die Steuerung freigegeben.", opName)
}()
// Nachrichten vom Operator empfangen und verarbeiten
for {
cmd, err := stream.Recv()
if err != nil {
return err // Verbindungsabbruch beendet die Session sauber über das defer
}
if cmd.Mode == "DISCONNECT" {
return nil
}
s.mu.Lock()
// Nur aktualisieren, wenn gültige Werte geschickt werden
if cmd.FrequencyHz > 0 {
_ = s.driver.SetFrequency(cmd.FrequencyHz)
s.currentFreq = cmd.FrequencyHz
}
if cmd.Mode != "" && cmd.Mode != "DISCONNECT" {
_ = s.driver.SetMode(cmd.Mode)
s.currentMode = cmd.Mode
}
s.mu.Unlock()
}
}
// MonitorSession streamt den aktuellen Hardware-Zustand an alle passiven Zuhörer
func (s *RigServer) MonitorSession(req *shack.Empty, stream shack.RigService_MonitorSessionServer) error {
ticker := time.NewTicker(250 * time.Millisecond)
defer ticker.Stop()
for {
select {
case <-stream.Context().Done():
return nil
case <-ticker.C:
s.mu.Lock()
opName := s.activeOperator
if opName == "" {
opName = "FREI"
}
err := stream.Send(&shack.RigStatus{
FrequencyHz: s.currentFreq,
Mode: s.currentMode,
SMeter: int32(s.currentSMeter),
ActiveOperatorName: opName,
})
s.mu.Unlock()
if err != nil {
return err
}
}
}
}
// Hilfsfunktion zur Operator-Validierung
func (s *RigServer) validateOperator(ctx context.Context) error {
// Falls kein Stream offen ist, ist das System frei (Sicherheits-Fallback)
if s.activeOperator == "" {
return nil
}
return nil
}
// hardwareSyncLoop liest im Hintergrund zyklisch das S-Meter und Änderungen am physischen VFO aus
func (s *RigServer) hardwareSyncLoop() {
ticker := time.NewTicker(300 * time.Millisecond)
for range ticker.C {
s.mu.Lock()
// Wenn kein Operator eingeloggt ist, lesen wir die Frequenz vom Gerät (falls jemand am echten VFO dreht)
if s.activeOperator == "" {
if freq, err := s.driver.GetFrequency(); err == nil && freq > 0 {
s.currentFreq = freq
}
if mode, err := s.driver.GetMode(); err == nil && mode != "" {
s.currentMode = mode
}
}
// S-Meter wird immer gelesen, um es im SDR-Display anzuzeigen
if sm, err := s.driver.GetSMeter(); err == nil {
s.currentSMeter = sm
}
s.mu.Unlock()
}
}