diff --git a/apps/admin-ui/src/app/nodes/page.tsx b/apps/admin-ui/src/app/nodes/page.tsx index d8ecf29..dd86a72 100644 --- a/apps/admin-ui/src/app/nodes/page.tsx +++ b/apps/admin-ui/src/app/nodes/page.tsx @@ -17,6 +17,8 @@ interface Node { freeStorage: string; activeStreams: number; libraryFileCount: number; + fullScanInterval: string; + fullScanAt: string; revoked: boolean; } @@ -37,6 +39,8 @@ export default function NodesPage() { const [editName, setEditName] = useState(""); const [editLocation, setEditLocation] = useState(""); const [editPublicUrl, setEditPublicUrl] = useState(""); + const [editScanInterval, setEditScanInterval] = useState("off"); + const [editScanAt, setEditScanAt] = useState("03:30"); function loadNodes() { fetch("/api/v1/admin/nodes", { credentials: "include" }) @@ -131,7 +135,7 @@ export default function NodesPage() { async function saveEdit() { if (!editId) return; - await fetch(`/api/v1/admin/nodes/${editId}`, { + const res = await fetch(`/api/v1/admin/nodes/${editId}`, { method: "PATCH", credentials: "include", headers: { "Content-Type": "application/json" }, @@ -139,8 +143,15 @@ export default function NodesPage() { name: editName, location: editLocation, publicStreamUrl: editPublicUrl.trim().replace(/\/$/, ""), + fullScanInterval: editScanInterval.trim() || "off", + fullScanAt: editScanAt.trim() || "03:30", }), }); + if (!res.ok) { + const data = await res.json().catch(() => ({})); + alert(data.error?.message ?? "Opslaan mislukt"); + return; + } setEditId(null); loadNodes(); } @@ -226,8 +237,27 @@ export default function NodesPage() { onChange={(e) => setEditPublicUrl(e.target.value)} placeholder="https://node1.example.com" /> + +
+ + setEditScanAt(e.target.value)} + placeholder="03:30" + />

- Zet dit ook in de node-container als MEDIA_NODE_PUBLIC_URL, anders overschrijft de volgende HELLO dit weer. + Catch-up scan in de tijdzone van de node (Unraid TZ). Standaard 03:30. +

+
+
+ + setEditScanInterval(e.target.value)} + placeholder="off" + /> +

+ Laat op off: toevoegen/verwijderen gaat via realtime watch. Admin-Rescan blijft beschikbaar.

@@ -271,6 +301,8 @@ export default function NodesPage() { setEditName(node.name); setEditLocation(node.location ?? ""); setEditPublicUrl(node.publicStreamUrl); + setEditScanInterval(node.fullScanInterval ?? "off"); + setEditScanAt(node.fullScanAt ?? "03:30"); }} > Edit diff --git a/apps/master-api/prisma/migrations/20260825044500_node_scan_settings/migration.sql b/apps/master-api/prisma/migrations/20260825044500_node_scan_settings/migration.sql new file mode 100644 index 0000000..ea0471e --- /dev/null +++ b/apps/master-api/prisma/migrations/20260825044500_node_scan_settings/migration.sql @@ -0,0 +1,3 @@ +-- AlterTable +ALTER TABLE "nodes" ADD COLUMN "full_scan_interval" TEXT NOT NULL DEFAULT 'off'; +ALTER TABLE "nodes" ADD COLUMN "full_scan_at" TEXT NOT NULL DEFAULT '03:30'; diff --git a/apps/master-api/prisma/schema.prisma b/apps/master-api/prisma/schema.prisma index ee53912..f39fc89 100644 --- a/apps/master-api/prisma/schema.prisma +++ b/apps/master-api/prisma/schema.prisma @@ -99,6 +99,10 @@ model Node { currentBandwidth BigInt @default(0) @map("current_bandwidth") libraryFileCount Int @default(0) @map("library_file_count") lastRevision Int @default(0) @map("last_revision") + /** Node scanner: "off" or duration like "24h" / "6h" */ + fullScanInterval String @default("off") @map("full_scan_interval") + /** Nightly catch-up time HH:MM in node local TZ */ + fullScanAt String @default("03:30") @map("full_scan_at") revoked Boolean @default(false) createdAt DateTime @default(now()) @map("created_at") updatedAt DateTime @updatedAt @map("updated_at") diff --git a/apps/master-api/src/nodes/routes.ts b/apps/master-api/src/nodes/routes.ts index de6e459..7ed5c1c 100644 --- a/apps/master-api/src/nodes/routes.ts +++ b/apps/master-api/src/nodes/routes.ts @@ -146,6 +146,8 @@ export async function registerNodeRoutes(app: FastifyInstance) { location?: string; locationType?: "LOCAL" | "REMOTE"; publicStreamUrl?: string; + fullScanInterval?: string; + fullScanAt?: string; }; let locationId: string | undefined | null = undefined; @@ -162,6 +164,23 @@ export async function registerNodeRoutes(app: FastifyInstance) { } } + if (body.fullScanInterval !== undefined) { + const v = body.fullScanInterval.trim().toLowerCase(); + if (v !== "off" && !/^\d+[smhd]$/.test(v)) { + throw new AppError( + "INVALID_REQUEST", + "fullScanInterval must be 'off' or a duration like 6h / 24h", + 400 + ); + } + } + if (body.fullScanAt !== undefined) { + const v = body.fullScanAt.trim(); + if (v.toLowerCase() !== "off" && !/^\d{1,2}:\d{2}$/.test(v)) { + throw new AppError("INVALID_REQUEST", "fullScanAt must be HH:MM or 'off'", 400); + } + } + const node = await prisma.node.update({ where: { id }, data: { @@ -169,10 +188,22 @@ export async function registerNodeRoutes(app: FastifyInstance) { ...(locationId !== undefined ? { locationId } : {}), ...(body.locationType ? { locationType: body.locationType } : {}), ...(body.publicStreamUrl ? { publicStreamUrl: body.publicStreamUrl } : {}), + ...(body.fullScanInterval !== undefined + ? { fullScanInterval: body.fullScanInterval.trim().toLowerCase() } + : {}), + ...(body.fullScanAt !== undefined ? { fullScanAt: body.fullScanAt.trim() } : {}), }, include: { location: true }, }); + if (body.fullScanInterval !== undefined || body.fullScanAt !== undefined) { + const { nodeConnectionManager } = await import("../websocket/manager"); + nodeConnectionManager.sendConfigUpdate(id, { + fullScanInterval: node.fullScanInterval, + fullScanAt: node.fullScanAt, + }); + } + return { node: serializeNode(node) }; }); @@ -270,6 +301,8 @@ function serializeNode(n: { activeStreams: number; currentBandwidth: bigint; libraryFileCount: number; + fullScanInterval?: string; + fullScanAt?: string; revoked: boolean; }) { return { @@ -288,6 +321,8 @@ function serializeNode(n: { activeStreams: n.activeStreams, currentBandwidth: n.currentBandwidth.toString(), libraryFileCount: n.libraryFileCount, + fullScanInterval: n.fullScanInterval ?? "off", + fullScanAt: n.fullScanAt ?? "03:30", revoked: n.revoked, }; } diff --git a/apps/master-api/src/websocket/manager.ts b/apps/master-api/src/websocket/manager.ts index 09d888d..35c1df1 100644 --- a/apps/master-api/src/websocket/manager.ts +++ b/apps/master-api/src/websocket/manager.ts @@ -148,7 +148,9 @@ class NodeConnectionManager { accepted: false, reason: "Invalid credentials", heartbeatIntervalSeconds: 30, - fullScanIntervalSeconds: 21600, + fullScanIntervalSeconds: 0, + fullScanInterval: "off", + fullScanAt: "03:30", }) ); socket.close(4003, "Authentication failed"); @@ -184,12 +186,25 @@ class NodeConnectionManager { console.log(`Node connected: ${payload.nodeId} (${payload.hostname})`); + const node = await prisma.node.findUnique({ where: { id: payload.nodeId } }); + const fullScanInterval = node?.fullScanInterval || "off"; + const fullScanAt = node?.fullScanAt || "03:30"; + const legacySeconds = + fullScanInterval === "off" + ? 0 + : (() => { + const m = fullScanInterval.match(/^(\d+)h$/i); + return m ? parseInt(m[1], 10) * 3600 : 0; + })(); + this.send( socket, createMessage("HELLO_ACK", { accepted: true, heartbeatIntervalSeconds: this.config?.NODE_HEARTBEAT_INTERVAL_SECONDS ?? 30, - fullScanIntervalSeconds: 21600, + fullScanIntervalSeconds: legacySeconds, + fullScanInterval, + fullScanAt, }) ); } @@ -326,6 +341,13 @@ class NodeConnectionManager { return this.sendToNode(nodeId, "RESCAN", { full }); } + sendConfigUpdate( + nodeId: string, + payload: { fullScanInterval?: string; fullScanAt?: string } + ): boolean { + return this.sendToNode(nodeId, "CONFIG_UPDATE", payload); + } + createPlaybackSession(nodeId: string, payload: CreatePlaybackSessionPayload): boolean { return this.sendToNode(nodeId, "CREATE_PLAYBACK_SESSION", payload); } diff --git a/node/media-node/cmd/media-node/main.go b/node/media-node/cmd/media-node/main.go index 1844ee5..e4324ab 100644 --- a/node/media-node/cmd/media-node/main.go +++ b/node/media-node/cmd/media-node/main.go @@ -124,7 +124,7 @@ func runService() { } }) - ctrlClient = control.NewClient(cfg, creds, store, sc, streamer, version) + ctrlClient = control.NewClient(cfg, creds, store, sc, streamer, version, configPath) streamer.SetSessionEndHandler(func(sessionID, reason string) { ctrlClient.NotifySessionEnded(sessionID, reason) }) diff --git a/node/media-node/internal/config/config.go b/node/media-node/internal/config/config.go index d28e4bf..8b8460f 100644 --- a/node/media-node/internal/config/config.go +++ b/node/media-node/internal/config/config.go @@ -2,6 +2,7 @@ package config import ( "fmt" + "log" "os" "strconv" "strings" @@ -71,10 +72,11 @@ func Load(path string) (*Config, error) { if envErr != nil { return nil, fmt.Errorf("read config: %w (and env incomplete: %v)", err, envErr) } + cfg = normalizeScannerDefaults(applyEnvOverrides(cfg)) if saveErr := Save(path, cfg); saveErr != nil { - return applyDefaults(applyEnvOverrides(cfg)), nil + return cfg, nil } - return applyDefaults(applyEnvOverrides(cfg)), nil + return cfg, nil } return nil, fmt.Errorf("read config: %w", err) } @@ -82,7 +84,22 @@ func Load(path string) (*Config, error) { if err := yaml.Unmarshal(data, &cfg); err != nil { return nil, fmt.Errorf("parse config: %w", err) } - return applyDefaults(applyEnvOverrides(&cfg)), nil + beforeInterval := cfg.Scanner.FullScanInterval + beforeAt := cfg.Scanner.FullScanAt + cfgPtr := normalizeScannerDefaults(applyEnvOverrides(&cfg)) + // Persist migration away from legacy 6h interval when nightly schedule is the new default. + if beforeInterval != cfgPtr.Scanner.FullScanInterval || beforeAt != cfgPtr.Scanner.FullScanAt { + if err := Save(path, cfgPtr); err != nil { + log.Printf("Could not persist scanner defaults: %v", err) + } else { + log.Printf( + "Scanner config updated: interval=%s at=%s", + cfgPtr.Scanner.FullScanInterval, + cfgPtr.Scanner.FullScanAt, + ) + } + } + return cfgPtr, nil } // FromEnv builds config from MEDIA_NODE_* environment variables (Docker/Unraid). @@ -125,14 +142,27 @@ func applyEnvOverrides(cfg *Config) *Config { return cfg } -func applyDefaults(cfg *Config) *Config { +// normalizeScannerDefaults migrates legacy 6h-only configs to watch + nightly 03:30. +func normalizeScannerDefaults(cfg *Config) *Config { if cfg.Stream.Listen == "" { cfg.Stream.Listen = DefaultListen } - if cfg.Scanner.FullScanInterval == "" { - cfg.Scanner.FullScanInterval = "off" + interval := strings.TrimSpace(strings.ToLower(cfg.Scanner.FullScanInterval)) + at := strings.TrimSpace(cfg.Scanner.FullScanAt) + + // Empty or legacy default "6h" without an explicit nightly time → new policy. + if interval == "" || interval == "6h" { + if at == "" { + cfg.Scanner.FullScanInterval = "off" + cfg.Scanner.FullScanAt = "03:30" + } else if interval == "6h" && at != "" { + // Explicit nightly already set (e.g. via admin) — drop redundant 6h. + cfg.Scanner.FullScanInterval = "off" + } else { + cfg.Scanner.FullScanInterval = "off" + } } - if cfg.Scanner.FullScanAt == "" { + if strings.TrimSpace(cfg.Scanner.FullScanAt) == "" { cfg.Scanner.FullScanAt = "03:30" } if cfg.Security.SessionIdleTimeout == "" { diff --git a/node/media-node/internal/control/client.go b/node/media-node/internal/control/client.go index aea8742..cd4fb76 100644 --- a/node/media-node/internal/control/client.go +++ b/node/media-node/internal/control/client.go @@ -36,6 +36,7 @@ type syncAck struct { type Client struct { cfg *config.Config + configPath string creds *config.Credentials store *database.Store scanner *scanner.Scanner @@ -58,17 +59,19 @@ func NewClient( sc *scanner.Scanner, streamer *streaming.Server, version string, + configPath string, ) *Client { return &Client{ - cfg: cfg, - creds: creds, - store: store, - scanner: sc, - streamer: streamer, - version: version, - startTime: time.Now(), - stopCh: make(chan struct{}), - syncAckCh: make(chan syncAck, 16), + cfg: cfg, + configPath: configPath, + creds: creds, + store: store, + scanner: sc, + streamer: streamer, + version: version, + startTime: time.Now(), + stopCh: make(chan struct{}), + syncAckCh: make(chan syncAck, 16), } } @@ -252,8 +255,11 @@ func (c *Client) handleMessage(msgType string, payload json.RawMessage) { switch msgType { case "HELLO_ACK": var ack struct { - Accepted bool `json:"accepted"` - Reason string `json:"reason"` + Accepted bool `json:"accepted"` + Reason string `json:"reason"` + FullScanInterval string `json:"fullScanInterval"` + FullScanAt string `json:"fullScanAt"` + FullScanIntervalSeconds int `json:"fullScanIntervalSeconds"` } _ = json.Unmarshal(payload, &ack) if !ack.Accepted { @@ -261,6 +267,7 @@ func (c *Client) handleMessage(msgType string, payload json.RawMessage) { return } log.Println("Connected to Master") + c.applyScannerSettings(ack.FullScanInterval, ack.FullScanAt, ack.FullScanIntervalSeconds) // Must run async so readLoop can receive FULL_LIBRARY_SYNC_ACK go c.sendFullSync() case "FULL_LIBRARY_SYNC_ACK": @@ -334,12 +341,41 @@ func (c *Client) handleMessage(msgType string, payload json.RawMessage) { case "FULL_LIBRARY_SYNC": go c.sendFullSync() case "CONFIG_UPDATE": - log.Println("CONFIG_UPDATE received (apply on next restart)") + var p struct { + FullScanInterval string `json:"fullScanInterval"` + FullScanAt string `json:"fullScanAt"` + } + if json.Unmarshal(payload, &p) == nil { + c.applyScannerSettings(p.FullScanInterval, p.FullScanAt, 0) + } case "ERROR": log.Printf("Master error: %s", string(payload)) } } +func (c *Client) applyScannerSettings(interval, at string, legacySeconds int) { + if interval == "" && at == "" && legacySeconds > 0 { + interval = fmt.Sprintf("%ds", legacySeconds) + } + if interval == "" && at == "" { + return + } + if interval == "" { + interval = c.cfg.Scanner.FullScanInterval + } + if at == "" { + at = c.cfg.Scanner.FullScanAt + } + c.scanner.UpdateSchedule(interval, at) + c.cfg.Scanner.FullScanInterval = interval + c.cfg.Scanner.FullScanAt = at + if c.configPath != "" { + if err := config.Save(c.configPath, c.cfg); err != nil { + log.Printf("Failed to save scanner settings: %v", err) + } + } +} + func (c *Client) SendLibraryEvents(events []scanner.LibraryEvent) { if len(events) == 0 { return diff --git a/node/media-node/internal/scanner/scanner.go b/node/media-node/internal/scanner/scanner.go index da841c1..d9bc655 100644 --- a/node/media-node/internal/scanner/scanner.go +++ b/node/media-node/internal/scanner/scanner.go @@ -41,16 +41,23 @@ type LibraryEvent struct { type EventHandler func(events []LibraryEvent) type Scanner struct { - cfg *config.Config - store *database.Store - onEvents EventHandler - mu sync.Mutex - scanning bool - status string + cfg *config.Config + store *database.Store + onEvents EventHandler + mu sync.Mutex + scanning bool + status string + resetSchedule chan struct{} } func New(cfg *config.Config, store *database.Store, onEvents EventHandler) *Scanner { - return &Scanner{cfg: cfg, store: store, onEvents: onEvents, status: "idle"} + return &Scanner{ + cfg: cfg, + store: store, + onEvents: onEvents, + status: "idle", + resetSchedule: make(chan struct{}, 2), + } } func (s *Scanner) Status() string { @@ -60,43 +67,92 @@ func (s *Scanner) Status() string { } func (s *Scanner) Start() { - // One inventory pass at boot, then realtime watch + nightly catch-up. + // One inventory pass at boot, then realtime watch + scheduled catch-up. go s.runStartupScan() go s.runNightlyFullScan() go s.runIntervalFullScan() go s.watchFilesystem() } +// UpdateSchedule applies admin/master scan settings and restarts timers. +func (s *Scanner) UpdateSchedule(interval, at string) { + s.mu.Lock() + if interval != "" { + s.cfg.Scanner.FullScanInterval = interval + } + if at != "" { + s.cfg.Scanner.FullScanAt = at + } + s.mu.Unlock() + select { + case s.resetSchedule <- struct{}{}: + default: + } + select { + case s.resetSchedule <- struct{}{}: + default: + } + log.Printf("Scanner schedule updated: interval=%s at=%s", s.cfg.Scanner.FullScanInterval, s.cfg.Scanner.FullScanAt) +} + func (s *Scanner) runStartupScan() { log.Println("Startup library inventory scan...") s.FullScan() } func (s *Scanner) runNightlyFullScan() { - hour, minute, ok := s.cfg.FullScanClock() - if !ok { - log.Println("Nightly full scan disabled") - return - } for { + s.mu.Lock() + hour, minute, ok := s.cfg.FullScanClock() + s.mu.Unlock() + if !ok { + select { + case <-s.resetSchedule: + continue + } + } next := nextLocalClock(time.Now(), hour, minute) log.Printf("Next nightly full scan at %s", next.Format(time.RFC3339)) timer := time.NewTimer(time.Until(next)) - <-timer.C - log.Println("Nightly catch-up full scan starting...") - s.FullScan() + select { + case <-timer.C: + log.Println("Nightly catch-up full scan starting...") + s.FullScan() + case <-s.resetSchedule: + if !timer.Stop() { + select { + case <-timer.C: + default: + } + } + } } } func (s *Scanner) runIntervalFullScan() { - d := s.cfg.FullScanDuration() - if d <= 0 { - return - } - log.Printf("Interval full scan every %s", d) - ticker := time.NewTicker(d) - for range ticker.C { - s.FullScan() + for { + s.mu.Lock() + d := s.cfg.FullScanDuration() + s.mu.Unlock() + if d <= 0 { + select { + case <-s.resetSchedule: + continue + } + } + log.Printf("Interval full scan every %s", d) + timer := time.NewTimer(d) + select { + case <-timer.C: + s.FullScan() + case <-s.resetSchedule: + if !timer.Stop() { + select { + case <-timer.C: + default: + } + } + } } } diff --git a/packages/protocol/src/index.ts b/packages/protocol/src/index.ts index 4c538c9..26c17aa 100644 --- a/packages/protocol/src/index.ts +++ b/packages/protocol/src/index.ts @@ -46,7 +46,12 @@ export interface HelloAckPayload { accepted: boolean; reason?: string; heartbeatIntervalSeconds: number; + /** @deprecated prefer fullScanInterval + fullScanAt */ fullScanIntervalSeconds: number; + /** "off" or Go duration e.g. "6h", "24h" */ + fullScanInterval: string; + /** Nightly catch-up HH:MM, or "off" */ + fullScanAt: string; } export interface HeartbeatPayload extends NodeHeartbeat {} @@ -76,6 +81,11 @@ export interface RescanPayload { full: boolean; } +export interface ConfigUpdatePayload { + fullScanInterval?: string; + fullScanAt?: string; +} + export interface PlaybackSessionEndedPayload { sessionId: string; reason: string;