diff --git a/apps/master-api/src/media/sync.ts b/apps/master-api/src/media/sync.ts index c28a5cf..646d1f9 100644 --- a/apps/master-api/src/media/sync.ts +++ b/apps/master-api/src/media/sync.ts @@ -25,15 +25,10 @@ export class LibrarySyncService { const lastEvent = events[events.length - 1]; if (lastEvent) { - const count = await prisma.mediaFile.count({ - where: { nodeId, available: true }, - }); + // Skip expensive COUNT on every event — heartbeat already reports libraryFileCount. await prisma.node.update({ where: { id: nodeId }, - data: { - lastRevision: lastEvent.nodeRevision, - libraryFileCount: count, - }, + data: { lastRevision: lastEvent.nodeRevision }, }); } } diff --git a/apps/master-api/src/websocket/manager.ts b/apps/master-api/src/websocket/manager.ts index 75f309a..ef1b6fe 100644 --- a/apps/master-api/src/websocket/manager.ts +++ b/apps/master-api/src/websocket/manager.ts @@ -293,6 +293,14 @@ class NodeConnectionManager { }); } + private enqueueNodeJob(nodeId: string, work: () => Promise): void { + const prev = this.syncQueues.get(nodeId) ?? Promise.resolve(); + const job = prev.then(work).catch((err) => { + console.error(`Node job failed for ${nodeId}:`, err); + }); + this.syncQueues.set(nodeId, job); + } + private async handleLibraryEvent( payload: LibraryEventPayload, socket: WebSocket @@ -300,26 +308,35 @@ class NodeConnectionManager { const conn = [...this.connections.values()].find((c) => c.socket === socket); if (!conn || !this.librarySync) return; + // ACK immediately — matching is queued so concurrent WS messages + // cannot exhaust the Prisma connection pool during a large scan. + this.send(socket, createMessage("LIBRARY_EVENT_ACK", { accepted: true })); + // Don't request another full sync while one is already in progress if (this.syncState.has(conn.nodeId)) { + this.enqueueNodeJob(conn.nodeId, async () => { + if (!this.librarySync) return; + await this.librarySync.processEvents(conn.nodeId, payload.events); + }); + return; + } + + this.enqueueNodeJob(conn.nodeId, async () => { + if (!this.librarySync) return; + + const node = await prisma.node.findUnique({ where: { id: conn.nodeId } }); + if (!node) return; + + const expectedRevision = node.lastRevision + 1; + const minRevision = Math.min(...payload.events.map((e) => e.nodeRevision)); + + if (payload.events.length > 0 && minRevision > expectedRevision) { + this.send(socket, createMessage("FULL_LIBRARY_SYNC", { full: true })); + return; + } + await this.librarySync.processEvents(conn.nodeId, payload.events); - this.send(socket, createMessage("LIBRARY_EVENT_ACK", { accepted: true })); - return; - } - - const node = await prisma.node.findUnique({ where: { id: conn.nodeId } }); - if (!node) return; - - const expectedRevision = node.lastRevision + 1; - const minRevision = Math.min(...payload.events.map((e) => e.nodeRevision)); - - if (payload.events.length > 0 && minRevision > expectedRevision) { - this.send(socket, createMessage("FULL_LIBRARY_SYNC", { full: true })); - return; - } - - await this.librarySync.processEvents(conn.nodeId, payload.events); - this.send(socket, createMessage("LIBRARY_EVENT_ACK", { accepted: true })); + }); } private async handleFullSync( @@ -350,10 +367,9 @@ class NodeConnectionManager { }) ); - const prev = this.syncQueues.get(conn.nodeId) ?? Promise.resolve(); - const job = prev - .then(async () => { - if (!this.librarySync) return; + this.enqueueNodeJob(conn.nodeId, async () => { + if (!this.librarySync) return; + try { await this.librarySync.processFullSync( conn.nodeId, payload.nodeRevision, @@ -372,12 +388,12 @@ class NodeConnectionManager { `Full library sync progress for ${conn.nodeId}: batch ${(batchIndex ?? 0) + 1}/${payload.batchCount ?? "?"}` ); } - }) - .catch((err) => { + } catch (err) { console.error(`Full sync failed for ${conn.nodeId}:`, err); this.syncState.delete(conn.nodeId); - }); - this.syncQueues.set(conn.nodeId, job); + throw err; + } + }); } private handlePlaybackSessionAck(payload: CreatePlaybackSessionAckPayload): void { diff --git a/node/media-node/internal/control/client.go b/node/media-node/internal/control/client.go index 412b11a..41fcdfc 100644 --- a/node/media-node/internal/control/client.go +++ b/node/media-node/internal/control/client.go @@ -24,11 +24,13 @@ import ( ) const ( - protocolVersion = 1 - fullSyncBatchSize = 200 - writeDeadlineShort = 15 * time.Second - writeDeadlineLong = 60 * time.Second - syncAckTimeout = 120 * time.Second + protocolVersion = 1 + fullSyncBatchSize = 200 + libraryEventBatchMax = 50 + libraryEventFlushWait = 200 * time.Millisecond + writeDeadlineShort = 15 * time.Second + writeDeadlineLong = 60 * time.Second + syncAckTimeout = 120 * time.Second ) type syncAck struct { @@ -52,6 +54,8 @@ type Client struct { syncAckCh chan syncAck fullSyncActive bool pendingEvents []scanner.LibraryEvent + eventBatch []scanner.LibraryEvent + eventFlushTimer *time.Timer } func NewClient( @@ -551,9 +555,33 @@ func (c *Client) SendLibraryEvents(events []scanner.LibraryEvent) { c.mu.Unlock() return } - c.mu.Unlock() - c.flushLibraryEvents(events) + c.eventBatch = append(c.eventBatch, events...) + if len(c.eventBatch) >= libraryEventBatchMax { + batch := c.eventBatch + c.eventBatch = nil + if c.eventFlushTimer != nil { + c.eventFlushTimer.Stop() + c.eventFlushTimer = nil + } + c.mu.Unlock() + c.flushLibraryEvents(batch) + return + } + + if c.eventFlushTimer == nil { + c.eventFlushTimer = time.AfterFunc(libraryEventFlushWait, func() { + c.mu.Lock() + batch := c.eventBatch + c.eventBatch = nil + c.eventFlushTimer = nil + c.mu.Unlock() + if len(batch) > 0 { + c.flushLibraryEvents(batch) + } + }) + } + c.mu.Unlock() } func (c *Client) flushLibraryEvents(events []scanner.LibraryEvent) {