Serialize library sync to prevent Prisma connection pool exhaustion.
This commit is contained in:
parent
38b73f1894
commit
21d0d7868e
3 changed files with 78 additions and 39 deletions
|
|
@ -25,15 +25,10 @@ export class LibrarySyncService {
|
||||||
|
|
||||||
const lastEvent = events[events.length - 1];
|
const lastEvent = events[events.length - 1];
|
||||||
if (lastEvent) {
|
if (lastEvent) {
|
||||||
const count = await prisma.mediaFile.count({
|
// Skip expensive COUNT on every event — heartbeat already reports libraryFileCount.
|
||||||
where: { nodeId, available: true },
|
|
||||||
});
|
|
||||||
await prisma.node.update({
|
await prisma.node.update({
|
||||||
where: { id: nodeId },
|
where: { id: nodeId },
|
||||||
data: {
|
data: { lastRevision: lastEvent.nodeRevision },
|
||||||
lastRevision: lastEvent.nodeRevision,
|
|
||||||
libraryFileCount: count,
|
|
||||||
},
|
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -293,6 +293,14 @@ class NodeConnectionManager {
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private enqueueNodeJob(nodeId: string, work: () => Promise<void>): 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(
|
private async handleLibraryEvent(
|
||||||
payload: LibraryEventPayload,
|
payload: LibraryEventPayload,
|
||||||
socket: WebSocket
|
socket: WebSocket
|
||||||
|
|
@ -300,26 +308,35 @@ class NodeConnectionManager {
|
||||||
const conn = [...this.connections.values()].find((c) => c.socket === socket);
|
const conn = [...this.connections.values()].find((c) => c.socket === socket);
|
||||||
if (!conn || !this.librarySync) return;
|
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
|
// Don't request another full sync while one is already in progress
|
||||||
if (this.syncState.has(conn.nodeId)) {
|
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);
|
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(
|
private async handleFullSync(
|
||||||
|
|
@ -350,10 +367,9 @@ class NodeConnectionManager {
|
||||||
})
|
})
|
||||||
);
|
);
|
||||||
|
|
||||||
const prev = this.syncQueues.get(conn.nodeId) ?? Promise.resolve();
|
this.enqueueNodeJob(conn.nodeId, async () => {
|
||||||
const job = prev
|
if (!this.librarySync) return;
|
||||||
.then(async () => {
|
try {
|
||||||
if (!this.librarySync) return;
|
|
||||||
await this.librarySync.processFullSync(
|
await this.librarySync.processFullSync(
|
||||||
conn.nodeId,
|
conn.nodeId,
|
||||||
payload.nodeRevision,
|
payload.nodeRevision,
|
||||||
|
|
@ -372,12 +388,12 @@ class NodeConnectionManager {
|
||||||
`Full library sync progress for ${conn.nodeId}: batch ${(batchIndex ?? 0) + 1}/${payload.batchCount ?? "?"}`
|
`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);
|
console.error(`Full sync failed for ${conn.nodeId}:`, err);
|
||||||
this.syncState.delete(conn.nodeId);
|
this.syncState.delete(conn.nodeId);
|
||||||
});
|
throw err;
|
||||||
this.syncQueues.set(conn.nodeId, job);
|
}
|
||||||
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
private handlePlaybackSessionAck(payload: CreatePlaybackSessionAckPayload): void {
|
private handlePlaybackSessionAck(payload: CreatePlaybackSessionAckPayload): void {
|
||||||
|
|
|
||||||
|
|
@ -24,11 +24,13 @@ import (
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
protocolVersion = 1
|
protocolVersion = 1
|
||||||
fullSyncBatchSize = 200
|
fullSyncBatchSize = 200
|
||||||
writeDeadlineShort = 15 * time.Second
|
libraryEventBatchMax = 50
|
||||||
writeDeadlineLong = 60 * time.Second
|
libraryEventFlushWait = 200 * time.Millisecond
|
||||||
syncAckTimeout = 120 * time.Second
|
writeDeadlineShort = 15 * time.Second
|
||||||
|
writeDeadlineLong = 60 * time.Second
|
||||||
|
syncAckTimeout = 120 * time.Second
|
||||||
)
|
)
|
||||||
|
|
||||||
type syncAck struct {
|
type syncAck struct {
|
||||||
|
|
@ -52,6 +54,8 @@ type Client struct {
|
||||||
syncAckCh chan syncAck
|
syncAckCh chan syncAck
|
||||||
fullSyncActive bool
|
fullSyncActive bool
|
||||||
pendingEvents []scanner.LibraryEvent
|
pendingEvents []scanner.LibraryEvent
|
||||||
|
eventBatch []scanner.LibraryEvent
|
||||||
|
eventFlushTimer *time.Timer
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewClient(
|
func NewClient(
|
||||||
|
|
@ -551,9 +555,33 @@ func (c *Client) SendLibraryEvents(events []scanner.LibraryEvent) {
|
||||||
c.mu.Unlock()
|
c.mu.Unlock()
|
||||||
return
|
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) {
|
func (c *Client) flushLibraryEvents(events []scanner.LibraryEvent) {
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue