diff --git a/libs/portal/shared/data-access/src/index.ts b/libs/portal/shared/data-access/src/index.ts index 4f245dfcb..0887a2743 100644 --- a/libs/portal/shared/data-access/src/index.ts +++ b/libs/portal/shared/data-access/src/index.ts @@ -1,2 +1,3 @@ export * from './lib/collection'; +export * from './lib/downloads'; export * from './lib/multi-source'; diff --git a/libs/portal/shared/data-access/src/lib/downloads/index.ts b/libs/portal/shared/data-access/src/lib/downloads/index.ts new file mode 100644 index 000000000..7b97d720d --- /dev/null +++ b/libs/portal/shared/data-access/src/lib/downloads/index.ts @@ -0,0 +1,2 @@ +export * from './season-download.models'; +export * from './season-download-coordinator.service'; diff --git a/libs/portal/shared/data-access/src/lib/downloads/season-download-coordinator.service.spec.ts b/libs/portal/shared/data-access/src/lib/downloads/season-download-coordinator.service.spec.ts new file mode 100644 index 000000000..96f46eadd --- /dev/null +++ b/libs/portal/shared/data-access/src/lib/downloads/season-download-coordinator.service.spec.ts @@ -0,0 +1,409 @@ +import { signal, type WritableSignal } from '@angular/core'; +import { TestBed } from '@angular/core/testing'; +import { + DownloadsService, + type DownloadItem, + type DownloadStartInput, +} from '@iptvnator/services'; +import type { + ElectronBridgeDownloadStartResult, + XtreamSerieEpisode, +} from '@iptvnator/shared/interfaces'; +import type { EpisodeDownloadIdentity } from '@iptvnator/portal/shared/util'; +import type { + EpisodeDownloadCandidate, + SeasonEpisodeDownloadAdapter, +} from './season-download.models'; +import { SeasonDownloadCoordinator } from './season-download-coordinator.service'; + +interface DownloadsServiceStub { + readonly downloads: WritableSignal; + readonly hasLoadedDownloads: WritableSignal; + readonly isAvailable: WritableSignal; + readonly startDownload: jest.MockedFunction< + DownloadsService['startDownload'] + >; + readonly loadDownloads: jest.MockedFunction< + DownloadsService['loadDownloads'] + >; +} + +interface Deferred { + readonly promise: Promise; + readonly resolve: (value: T) => void; + readonly reject: (reason?: unknown) => void; +} + +const FIRST_IDENTITY: EpisodeDownloadIdentity = { + playlistId: 'playlist-1', + contentType: 'episode', + xtreamId: 101, + seriesXtreamId: 20, + seasonNumber: 1, + episodeNumber: 1, +}; + +describe('SeasonDownloadCoordinator', () => { + let downloadsService: DownloadsServiceStub; + let coordinator: SeasonDownloadCoordinator; + + beforeEach(() => { + downloadsService = { + downloads: signal([]), + hasLoadedDownloads: signal(true), + isAvailable: signal(true), + startDownload: jest.fn().mockResolvedValue({ success: true }), + loadDownloads: jest.fn().mockResolvedValue(undefined), + }; + + TestBed.configureTestingModule({ + providers: [ + SeasonDownloadCoordinator, + { provide: DownloadsService, useValue: downloadsService }, + ], + }); + coordinator = TestBed.inject(SeasonDownloadCoordinator); + }); + + it('reserves one identity synchronously without blocking a distinct episode', async () => { + const firstPrepare = deferred(); + const first = candidate(FIRST_IDENTITY, () => firstPrepare.promise); + const second = candidate(identity(102, 2)); + + const submission = coordinator.enqueueOne(first); + + expect(coordinator.isPending(first.identity)).toBe(true); + expect(coordinator.isEligible(first)).toBe(false); + expect(coordinator.isPending(second.identity)).toBe(false); + expect(coordinator.isEligible(second)).toBe(true); + + firstPrepare.resolve(request(first.identity)); + await expect(submission).resolves.toBe('added'); + }); + + it('starts only one download for two rapid submissions of the same identity', async () => { + const prepare = deferred(); + const first = candidate(FIRST_IDENTITY, () => prepare.promise); + const duplicate = candidate(FIRST_IDENTITY); + + const firstSubmission = coordinator.enqueueOne(first); + const secondSubmission = coordinator.enqueueOne(duplicate); + + await expect(secondSubmission).resolves.toBe('skipped'); + expect(duplicate.prepare).not.toHaveBeenCalled(); + + prepare.resolve(request(first.identity)); + await expect(firstSubmission).resolves.toBe('added'); + expect(downloadsService.startDownload).toHaveBeenCalledTimes(1); + }); + + it('submits a season sequentially in display order', async () => { + const firstStart = deferred(); + const first = candidate(FIRST_IDENTITY); + const second = candidate(identity(102, 2)); + const adapter = adapterFor([first, second]); + downloadsService.startDownload + .mockImplementationOnce(() => firstStart.promise) + .mockResolvedValueOnce({ success: true }); + + const submission = coordinator.enqueueSeason( + [episode(101, 1), episode(102, 2)], + adapter, + '1' + ); + await flushMicrotasks(); + + expect(first.prepare).toHaveBeenCalledTimes(1); + expect(second.prepare).not.toHaveBeenCalled(); + expect(downloadsService.startDownload).toHaveBeenCalledTimes(1); + + firstStart.resolve({ success: true }); + await flushMicrotasks(); + + expect(second.prepare).toHaveBeenCalledTimes(1); + expect(downloadsService.startDownload).toHaveBeenNthCalledWith( + 2, + request(second.identity) + ); + await expect(submission).resolves.toEqual({ + added: 2, + skipped: 0, + failed: 0, + }); + }); + + it('partitions invalid, duplicate, blocked, stable skip, and unsuccessful submissions', async () => { + const stable = candidate(FIRST_IDENTITY); + const duplicate = candidate(FIRST_IDENTITY); + const unsuccessful = candidate(identity(102, 2)); + const paused = candidate(identity(103, 3)); + downloadsService.downloads.set([ + download(paused.identity, { status: 'paused' }), + ]); + downloadsService.startDownload + .mockResolvedValueOnce({ + success: false, + reason: 'already-in-progress', + }) + .mockResolvedValueOnce({ + success: false, + error: 'Backend declined the request', + }); + const adapter = adapterFor([ + null, + stable, + duplicate, + unsuccessful, + paused, + ]); + + await expect( + coordinator.enqueueSeason( + [ + episode(100, 0), + episode(101, 1), + episode(101, 1), + episode(102, 2), + episode(103, 3), + ], + adapter, + '1' + ) + ).resolves.toEqual({ added: 0, skipped: 4, failed: 1 }); + + expect(stable.prepare).toHaveBeenCalledTimes(1); + expect(duplicate.prepare).not.toHaveBeenCalled(); + expect(unsuccessful.prepare).toHaveBeenCalledTimes(1); + expect(paused.prepare).not.toHaveBeenCalled(); + expect(downloadsService.loadDownloads).not.toHaveBeenCalled(); + }); + + it('continues after preparation rejects and refreshes once after successes', async () => { + const warnSpy = jest.spyOn(console, 'warn').mockImplementation(() => { + /* captured */ + }); + const rejected = candidate(FIRST_IDENTITY, () => + Promise.reject( + new Error( + 'Request failed: https://portal.invalid?username=alice&password=hunter2' + ) + ) + ); + const ipcFailure = candidate(identity(102, 2)); + const firstSuccess = candidate(identity(103, 3)); + const secondSuccess = candidate(identity(104, 4)); + downloadsService.startDownload.mockRejectedValueOnce( + new Error( + 'IPC failed: https://portal.invalid?username=alice&password=hunter2' + ) + ); + + try { + await expect( + coordinator.enqueueSeason( + [ + episode(101, 1), + episode(102, 2), + episode(103, 3), + episode(104, 4), + ], + adapterFor([ + rejected, + ipcFailure, + firstSuccess, + secondSuccess, + ]) + ) + ).resolves.toEqual({ added: 2, skipped: 0, failed: 2 }); + + expect(downloadsService.startDownload).toHaveBeenCalledTimes(3); + expect(downloadsService.loadDownloads).toHaveBeenCalledTimes(1); + expect(warnSpy).toHaveBeenCalled(); + expect(loggedText(warnSpy)).not.toContain('hunter2'); + expect(loggedText(warnSpy)).not.toContain('alice'); + } finally { + warnSpy.mockRestore(); + } + }); + + it('blocks eligibility and reservation until downloads have loaded', async () => { + downloadsService.hasLoadedDownloads.set(false); + const item = candidate(FIRST_IDENTITY); + + expect(coordinator.isEligible(item)).toBe(false); + await expect(coordinator.enqueueOne(item)).resolves.toBe('skipped'); + expect(item.prepare).not.toHaveBeenCalled(); + + downloadsService.hasLoadedDownloads.set(true); + downloadsService.isAvailable.set(false); + expect(coordinator.isEligible(item)).toBe(false); + await expect(coordinator.enqueueOne(item)).resolves.toBe('skipped'); + }); + + it('allows a completed missing episode and skips a completed available one', async () => { + const missing = candidate(FIRST_IDENTITY); + const available = candidate(identity(102, 2)); + downloadsService.downloads.set([ + download(missing.identity, { + status: 'completed', + fileAvailability: 'missing', + }), + download(available.identity, { + status: 'completed', + fileAvailability: 'available', + }), + ]); + + expect(coordinator.findDownload(missing.identity)?.fileAvailability).toBe( + 'missing' + ); + expect(coordinator.isEligible(missing)).toBe(true); + expect(coordinator.isEligible(available)).toBe(false); + await expect(coordinator.enqueueOne(available)).resolves.toBe( + 'skipped' + ); + expect(available.prepare).not.toHaveBeenCalled(); + }); + + it('prepares a duplicate identity only once within one batch', async () => { + const first = candidate(FIRST_IDENTITY); + const duplicate = candidate(FIRST_IDENTITY); + + await expect( + coordinator.enqueueSeason( + [episode(101, 1), episode(101, 1)], + adapterFor([first, duplicate]) + ) + ).resolves.toEqual({ added: 1, skipped: 1, failed: 0 }); + + expect(first.prepare).toHaveBeenCalledTimes(1); + expect(duplicate.prepare).not.toHaveBeenCalled(); + expect(downloadsService.startDownload).toHaveBeenCalledTimes(1); + }); + + it('keeps successful keys pending and releases them when refresh rejects', async () => { + const refresh = deferred(); + const refreshStarted = deferred(); + const first = candidate(FIRST_IDENTITY); + const second = candidate(identity(102, 2)); + downloadsService.loadDownloads.mockImplementation(() => { + refreshStarted.resolve(undefined); + return refresh.promise; + }); + + const submission = coordinator.enqueueSeason( + [episode(101, 1), episode(102, 2)], + adapterFor([first, second]) + ); + await refreshStarted.promise; + + expect(coordinator.isPending(first.identity)).toBe(true); + expect(coordinator.isPending(second.identity)).toBe(true); + + refresh.reject(new Error('Refresh failed')); + await expect(submission).rejects.toThrow('Refresh failed'); + expect(coordinator.isPending(first.identity)).toBe(false); + expect(coordinator.isPending(second.identity)).toBe(false); + }); +}); + +function identity( + xtreamId: number, + episodeNumber: number +): EpisodeDownloadIdentity { + return { ...FIRST_IDENTITY, xtreamId, episodeNumber }; +} + +function request(value: EpisodeDownloadIdentity): DownloadStartInput { + return { + playlistId: value.playlistId, + xtreamId: value.xtreamId, + contentType: value.contentType, + title: `Episode ${value.episodeNumber}`, + url: `https://stream.invalid/${value.xtreamId}`, + seriesXtreamId: value.seriesXtreamId, + seasonNumber: value.seasonNumber, + episodeNumber: value.episodeNumber, + }; +} + +function candidate( + value: EpisodeDownloadIdentity, + prepare: () => Promise = () => + Promise.resolve(request(value)) +): EpisodeDownloadCandidate & { readonly prepare: jest.Mock } { + return { identity: value, prepare: jest.fn(prepare) }; +} + +function episode( + xtreamId: number, + episodeNumber: number +): XtreamSerieEpisode { + return { + id: String(xtreamId), + episode_num: episodeNumber, + title: `Episode ${episodeNumber}`, + container_extension: 'mp4', + info: [], + custom_sid: '', + added: '', + season: 1, + direct_source: '', + }; +} + +function adapterFor( + candidates: readonly (EpisodeDownloadCandidate | null)[] +): SeasonEpisodeDownloadAdapter { + let index = 0; + return { + createCandidate: jest.fn(() => candidates[index++] ?? null), + }; +} + +function download( + value: EpisodeDownloadIdentity, + overrides: Partial = {} +): DownloadItem { + return { + id: value.xtreamId, + playlistId: value.playlistId, + xtreamId: value.xtreamId, + contentType: value.contentType, + seriesXtreamId: value.seriesXtreamId, + seasonNumber: value.seasonNumber, + episodeNumber: value.episodeNumber, + title: `Episode ${value.episodeNumber}`, + url: 'https://stream.invalid/download', + status: 'queued', + ...overrides, + }; +} + +function deferred(): Deferred { + let resolve!: (value: T) => void; + let reject!: (reason?: unknown) => void; + const promise = new Promise((resolvePromise, rejectPromise) => { + resolve = resolvePromise; + reject = rejectPromise; + }); + return { promise, resolve, reject }; +} + +async function flushMicrotasks(): Promise { + await Promise.resolve(); + await Promise.resolve(); +} + +function loggedText(spy: jest.SpyInstance): string { + return spy.mock.calls + .flat() + .map((value) => + value instanceof Error + ? `${value.message}\n${value.stack ?? ''}` + : typeof value === 'string' + ? value + : JSON.stringify(value) + ) + .join('\n'); +} diff --git a/libs/portal/shared/data-access/src/lib/downloads/season-download-coordinator.service.ts b/libs/portal/shared/data-access/src/lib/downloads/season-download-coordinator.service.ts new file mode 100644 index 000000000..a7ee6a1d0 --- /dev/null +++ b/libs/portal/shared/data-access/src/lib/downloads/season-download-coordinator.service.ts @@ -0,0 +1,174 @@ +import { inject, Injectable, signal } from '@angular/core'; +import { + DownloadsService, + type DownloadItem, + type DownloadStartInput, +} from '@iptvnator/services'; +import { + ELECTRON_BRIDGE_DOWNLOAD_START_REASONS, + type XtreamSerieEpisode, +} from '@iptvnator/shared/interfaces'; +import { + createEpisodeDownloadIdentityKey, + createLogger, + findEpisodeDownload, + isEpisodeDownloadEligible, + type EpisodeDownloadIdentity, +} from '@iptvnator/portal/shared/util'; +import { + EPISODE_DOWNLOAD_SUBMISSIONS, + type EpisodeDownloadCandidate, + type EpisodeDownloadSubmission, + type SeasonDownloadResult, + type SeasonEpisodeDownloadAdapter, +} from './season-download.models'; + +@Injectable({ providedIn: 'root' }) +export class SeasonDownloadCoordinator { + private readonly downloadsService = inject(DownloadsService); + private readonly logger = createLogger('SeasonDownloadCoordinator'); + private readonly pending = signal>(new Set()); + + isPending(identity: EpisodeDownloadIdentity): boolean { + return this.pending().has(createEpisodeDownloadIdentityKey(identity)); + } + + findDownload(identity: EpisodeDownloadIdentity): DownloadItem | undefined { + return findEpisodeDownload(identity, this.downloadsService.downloads()); + } + + isEligible(candidate: EpisodeDownloadCandidate): boolean { + return ( + this.downloadsService.isAvailable() && + this.downloadsService.hasLoadedDownloads() && + !this.isPending(candidate.identity) && + isEpisodeDownloadEligible(this.findDownload(candidate.identity)) + ); + } + + async enqueueOne( + candidate: EpisodeDownloadCandidate + ): Promise { + if (!this.reserve(candidate)) { + return EPISODE_DOWNLOAD_SUBMISSIONS.Skipped; + } + + const submission = await this.submit(candidate); + if (submission !== EPISODE_DOWNLOAD_SUBMISSIONS.Added) { + this.release(candidate.identity); + return submission; + } + + try { + await this.downloadsService.loadDownloads(); + return submission; + } finally { + this.release(candidate.identity); + } + } + + async enqueueSeason( + episodes: readonly XtreamSerieEpisode[], + adapter: SeasonEpisodeDownloadAdapter, + fallbackSeasonKey?: string + ): Promise { + const result = { added: 0, skipped: 0, failed: 0 }; + const reserved: EpisodeDownloadCandidate[] = []; + + for (const episode of episodes) { + const candidate = this.createCandidate( + adapter, + episode, + fallbackSeasonKey + ); + if (!candidate || !this.reserve(candidate)) { + result.skipped += 1; + continue; + } + reserved.push(candidate); + } + + const accepted: EpisodeDownloadIdentity[] = []; + for (const candidate of reserved) { + const submission = await this.submit(candidate); + result[submission] += 1; + if (submission === EPISODE_DOWNLOAD_SUBMISSIONS.Added) { + accepted.push(candidate.identity); + } else { + this.release(candidate.identity); + } + } + + if (accepted.length > 0) { + try { + await this.downloadsService.loadDownloads(); + } finally { + this.releaseAll(accepted); + } + } + + return result; + } + + private reserve(candidate: EpisodeDownloadCandidate): boolean { + if (!this.isEligible(candidate)) { + return false; + } + + const key = createEpisodeDownloadIdentityKey(candidate.identity); + const next = new Set(this.pending()); + next.add(key); + this.pending.set(next); + return true; + } + + private async submit( + candidate: EpisodeDownloadCandidate + ): Promise { + try { + const request: DownloadStartInput = await candidate.prepare(); + const result = await this.downloadsService.startDownload(request); + if (result.success) { + return EPISODE_DOWNLOAD_SUBMISSIONS.Added; + } + if ( + result.reason === + ELECTRON_BRIDGE_DOWNLOAD_START_REASONS.AlreadyInProgress + ) { + return EPISODE_DOWNLOAD_SUBMISSIONS.Skipped; + } + return EPISODE_DOWNLOAD_SUBMISSIONS.Failed; + } catch { + this.logger.warn('Episode download submission failed'); + return EPISODE_DOWNLOAD_SUBMISSIONS.Failed; + } + } + + private createCandidate( + adapter: SeasonEpisodeDownloadAdapter, + episode: XtreamSerieEpisode, + fallbackSeasonKey: string | undefined + ): EpisodeDownloadCandidate | null { + try { + return adapter.createCandidate(episode, fallbackSeasonKey); + } catch { + this.logger.warn('Episode download candidate creation failed'); + return null; + } + } + + private release(identity: EpisodeDownloadIdentity): void { + const key = createEpisodeDownloadIdentityKey(identity); + const next = new Set(this.pending()); + next.delete(key); + this.pending.set(next); + } + + private releaseAll(identities: readonly EpisodeDownloadIdentity[]): void { + const next = new Set(this.pending()); + for (const identity of identities) { + next.delete(createEpisodeDownloadIdentityKey(identity)); + } + this.pending.set(next); + } +} diff --git a/libs/portal/shared/data-access/src/lib/downloads/season-download.models.ts b/libs/portal/shared/data-access/src/lib/downloads/season-download.models.ts new file mode 100644 index 000000000..e53013969 --- /dev/null +++ b/libs/portal/shared/data-access/src/lib/downloads/season-download.models.ts @@ -0,0 +1,30 @@ +import type { DownloadStartInput } from '@iptvnator/services'; +import type { XtreamSerieEpisode } from '@iptvnator/shared/interfaces'; +import type { EpisodeDownloadIdentity } from '@iptvnator/portal/shared/util'; + +export const EPISODE_DOWNLOAD_SUBMISSIONS = { + Added: 'added', + Skipped: 'skipped', + Failed: 'failed', +} as const; + +export type EpisodeDownloadSubmission = + (typeof EPISODE_DOWNLOAD_SUBMISSIONS)[keyof typeof EPISODE_DOWNLOAD_SUBMISSIONS]; + +export interface EpisodeDownloadCandidate { + readonly identity: EpisodeDownloadIdentity; + readonly prepare: () => Promise; +} + +export interface SeasonEpisodeDownloadAdapter { + createCandidate( + episode: XtreamSerieEpisode, + fallbackSeasonKey: string | undefined + ): EpisodeDownloadCandidate | null; +} + +export interface SeasonDownloadResult { + readonly added: number; + readonly skipped: number; + readonly failed: number; +}