199 lines
5.7 KiB
Go
199 lines
5.7 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"log/slog"
|
|
"net/http"
|
|
"os"
|
|
"os/signal"
|
|
"sync"
|
|
"syscall"
|
|
"time"
|
|
|
|
"yovision/sense/internal/auditrelay"
|
|
"yovision/sense/internal/auth"
|
|
"yovision/sense/internal/config"
|
|
"yovision/sense/internal/controlapi"
|
|
"yovision/sense/internal/metrics"
|
|
"yovision/sense/internal/mtx"
|
|
"yovision/sense/internal/onvif"
|
|
"yovision/sense/internal/orphan"
|
|
"yovision/sense/internal/probe"
|
|
"yovision/sense/internal/reconcile"
|
|
"yovision/sense/internal/store"
|
|
)
|
|
|
|
var version = "dev"
|
|
|
|
func main() {
|
|
logger := slog.New(slog.NewJSONHandler(os.Stdout, nil))
|
|
if err := run(logger); err != nil {
|
|
logger.Error("Sense stopped", "error", err)
|
|
os.Exit(1)
|
|
}
|
|
}
|
|
|
|
func run(logger *slog.Logger) error {
|
|
cfg, err := config.Load()
|
|
if err != nil {
|
|
return fmt.Errorf("load configuration: %w", err)
|
|
}
|
|
instanceID := cfg.InstanceID
|
|
if instanceID == "" {
|
|
instanceID, err = metrics.GenerateInstanceID()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
registry := metrics.New(instanceID, version)
|
|
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
|
|
defer stop()
|
|
|
|
repository, err := store.OpenRepository(ctx, cfg.DatabaseDriver, cfg.DatabaseDSN)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer repository.Close()
|
|
var auditWorker *auditrelay.Worker
|
|
if cfg.AuditRelayEnabled {
|
|
relayStore, ok := repository.(auditrelay.Repository)
|
|
if !ok {
|
|
return errors.New("selected repository does not support audit relay")
|
|
}
|
|
if err := relayStore.AuditRelayReady(ctx); err != nil {
|
|
return err
|
|
}
|
|
secret, err := auditrelay.LoadKey(cfg.AuditRelayKeyFile, cfg.AuditRelayKeyID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
client, err := auditrelay.NewClient(cfg.AuditRelayURL, cfg.AuditRelayKeyID, secret, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
auditWorker, err = auditrelay.NewWorker(relayStore, client, instanceID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
var controlHandler http.Handler
|
|
if cfg.ControlAPIEnabled {
|
|
controlStore, ok := repository.(store.ControlRepository)
|
|
if !ok {
|
|
return errors.New("selected repository does not support Sense Control API")
|
|
}
|
|
authenticator, err := auth.LoadStaticSHA256(cfg.ControlAuthFile)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
cursors, err := controlapi.LoadCursorCodec(cfg.ControlCursorKeyFile)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
controlHandler = controlapi.NewHTTPHandler(controlStore, authenticator, cursors)
|
|
}
|
|
mediaClient, err := mtx.NewClient(cfg.MediaMTXURL, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
credentials := onvif.EnvCredentials{}
|
|
var cameraAdapter onvif.Adapter = onvif.UnavailableAdapter{}
|
|
if cfg.ONVIFMode == "standard" {
|
|
cameraAdapter = onvif.NewHTTPAdapter(credentials, nil, onvif.HTTPOptions{
|
|
RTSPRewriteHost: cfg.RTSPRewriteHost,
|
|
RTSPRewritePort: cfg.RTSPRewritePort,
|
|
StripRTSPQuery: cfg.RTSPStripQuery,
|
|
})
|
|
}
|
|
discovery := onvif.NewRouter(cameraAdapter, credentials)
|
|
reconciler := reconcile.NewWithOptions(repository, discovery, mediaClient, reconcile.Options{
|
|
InstanceID: instanceID, LeaseDuration: cfg.ReconcileLeaseDuration,
|
|
OperationTimeout: cfg.ReconcileOperationTimeout, Metrics: registry,
|
|
})
|
|
checker := probe.New(repository, mediaClient)
|
|
var orphanScanner *orphan.Manager
|
|
if cfg.OrphanScanEnabled {
|
|
orphanStore, ok := repository.(store.OrphanRepository)
|
|
if !ok {
|
|
return errors.New("selected repository does not support orphan scanning")
|
|
}
|
|
orphanScanner = orphan.New(orphanStore, mediaClient, instanceID, registry)
|
|
}
|
|
report := func(err error) {
|
|
// Domain and MediaMTX errors intentionally omit stream URIs and credentials.
|
|
logger.Warn("background convergence error", "error", err)
|
|
}
|
|
var background sync.WaitGroup
|
|
startBackground := func(run func()) {
|
|
background.Add(1)
|
|
go func() {
|
|
defer background.Done()
|
|
run()
|
|
}()
|
|
}
|
|
startBackground(func() { reconciler.Run(ctx, cfg.ReconcileInterval, report) })
|
|
startBackground(func() { checker.Run(ctx, cfg.ProbeInterval, report) })
|
|
if orphanScanner != nil {
|
|
startBackground(func() { orphanScanner.Run(ctx, cfg.OrphanScanInterval, report) })
|
|
}
|
|
if auditWorker != nil {
|
|
startBackground(func() { auditWorker.Run(ctx, cfg.AuditRelayInterval, report) })
|
|
}
|
|
|
|
mux := http.NewServeMux()
|
|
mux.HandleFunc("GET /healthz", func(writer http.ResponseWriter, _ *http.Request) {
|
|
writer.Header().Set("Content-Type", "application/json")
|
|
writer.WriteHeader(http.StatusOK)
|
|
_, _ = writer.Write([]byte(`{"status":"ok"}`))
|
|
})
|
|
mux.HandleFunc("GET /readyz", func(writer http.ResponseWriter, _ *http.Request) {
|
|
writer.Header().Set("Content-Type", "application/json")
|
|
writer.WriteHeader(http.StatusOK)
|
|
_, _ = writer.Write([]byte(`{"status":"ready"}`))
|
|
})
|
|
if cfg.MetricsEnabled {
|
|
mux.Handle("GET /metrics", registry.Handler())
|
|
}
|
|
if cfg.ControlAPIEnabled {
|
|
mux.Handle("/api/v1/", controlHandler)
|
|
}
|
|
|
|
server := &http.Server{
|
|
Addr: cfg.HTTPAddress, Handler: mux,
|
|
ReadHeaderTimeout: 5 * time.Second,
|
|
ReadTimeout: 15 * time.Second,
|
|
WriteTimeout: 15 * time.Second,
|
|
IdleTimeout: 60 * time.Second,
|
|
}
|
|
serverErrors := make(chan error, 1)
|
|
go func() {
|
|
logger.Info("Sense listening", "address", cfg.HTTPAddress, "version", version,
|
|
"instance_id", instanceID,
|
|
"control_api_enabled", cfg.ControlAPIEnabled,
|
|
"audit_relay_enabled", cfg.AuditRelayEnabled)
|
|
serverErrors <- server.ListenAndServe()
|
|
}()
|
|
|
|
select {
|
|
case <-ctx.Done():
|
|
case serverErr := <-serverErrors:
|
|
if !errors.Is(serverErr, http.ErrServerClosed) {
|
|
stop()
|
|
background.Wait()
|
|
return fmt.Errorf("serve HTTP: %w", serverErr)
|
|
}
|
|
}
|
|
stop()
|
|
shutdownContext, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer cancel()
|
|
shutdownErr := server.Shutdown(shutdownContext)
|
|
background.Wait()
|
|
if shutdownErr != nil {
|
|
return fmt.Errorf("shutdown HTTP server: %w", shutdownErr)
|
|
}
|
|
return nil
|
|
}
|