208 lines
5.5 KiB
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()
|
|
}
|
|
}
|