mirror of
https://github.com/4gray/iptvnator.git
synced 2026-10-10 18:36:15 -08:00
fix(downloads): fail closed before provider prep
This commit is contained in:
1 parent
dc8da931ea
commit
645a37afa1
6 files changed
+204
-16
No files matched your search
@@ -942,8 +942,10 @@ engine` (restart required) or
|
||||
reports added, skipped, and failed counts. Xtream and Stalker adapters remain
|
||||
responsible for provider URLs, headers, and metadata; the backend still runs
|
||||
one active transfer with a FIFO queue. `DOWNLOADS_START` remains the sole
|
||||
start IPC; its
|
||||
stable `reason: 'already-in-progress'` and `reason: 'already-downloaded'`
|
||||
start IPC. A reserved completed-missing match triggers one authoritative
|
||||
preflight refresh before provider preparation, so a restored Stalker file can
|
||||
become a stable skip without a portal request. The IPC's stable
|
||||
`reason: 'already-in-progress'` and `reason: 'already-downloaded'`
|
||||
results are counted as skipped, and no batch IPC is introduced. The latter
|
||||
comes from an asynchronous main-process filesystem recheck before a
|
||||
completed-missing row can be reset, so a file restored after the renderer
|
||||
@@ -954,9 +956,11 @@ engine` (restart required) or
|
||||
and report a timeout as missing for that snapshot. The underlying filesystem
|
||||
operation remains coalesced and charged against the four-probe cap until it
|
||||
settles, so later callers have independent bounded waits without duplicating
|
||||
stalled native work. Episode and season download actions require an
|
||||
authoritative global list. A successful snapshot remains authoritative while
|
||||
a later background refresh is in flight; a latest refresh failure leaves
|
||||
stalled native work. Only `ENOENT` and `ENOTDIR` prove absence; permission,
|
||||
I/O, and other filesystem errors remain unknown and cannot clear a completed
|
||||
row. Episode and season download actions require an authoritative global
|
||||
list. A successful snapshot remains authoritative while a later background
|
||||
refresh is in flight; a latest refresh failure leaves
|
||||
loading/empty-state resolution intact but disables starts until another
|
||||
snapshot succeeds.
|
||||
- Episode ownership uses normalized `episode.id` as the canonical `xtreamId`
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import {
|
||||
createDownloadFileAvailabilityProbe,
|
||||
decorateDownloadItem,
|
||||
getDownloadFileAvailability,
|
||||
getDownloadFileAvailabilityWithTimeoutAsync,
|
||||
@@ -17,6 +18,32 @@ function lstatResult(options: {
|
||||
}
|
||||
|
||||
describe('download file availability', () => {
|
||||
it.each([
|
||||
['EACCES', 'unknown'],
|
||||
['EIO', 'unknown'],
|
||||
['ENOENT', 'missing'],
|
||||
['ENOTDIR', 'missing'],
|
||||
] as const)(
|
||||
'classifies an %s filesystem probe without risking a completed row',
|
||||
async (code, expected) => {
|
||||
const error = Object.assign(new Error(code), { code });
|
||||
const probe = createDownloadFileAvailabilityProbe(async () => {
|
||||
throw error;
|
||||
});
|
||||
|
||||
await expect(
|
||||
getDownloadFileAvailabilityWithTimeoutAsync(
|
||||
{
|
||||
filePath: '/downloads/episode.mp4',
|
||||
status: 'completed',
|
||||
},
|
||||
25,
|
||||
probe
|
||||
)
|
||||
).resolves.toBe(expected);
|
||||
}
|
||||
);
|
||||
|
||||
it('bounds a restored-file probe and returns unknown on timeout', async () => {
|
||||
jest.useFakeTimers();
|
||||
try {
|
||||
|
||||
@@ -23,7 +23,9 @@ type DownloadAsyncLstat = (
|
||||
|
||||
type DownloadFileAvailabilityProbeResult = boolean | 'unknown';
|
||||
|
||||
type DownloadFileAvailabilityProbe = (filePath: string) => Promise<boolean>;
|
||||
type DownloadFileAvailabilityProbe = (
|
||||
filePath: string
|
||||
) => Promise<DownloadFileAvailabilityProbeResult>;
|
||||
|
||||
const DEFAULT_MAX_CONCURRENT_FILE_PROBES = 4;
|
||||
const DEFAULT_FILE_PROBE_TIMEOUT_MS = 1_000;
|
||||
@@ -31,7 +33,15 @@ const DEFAULT_FILE_PROBE_TIMEOUT_MS = 1_000;
|
||||
export type BoundedDownloadFileAvailability =
|
||||
ElectronDownloadFileAvailability | 'unknown';
|
||||
|
||||
function createDownloadFileAvailabilityProbe(
|
||||
function isMissingFileSystemError(error: unknown): boolean {
|
||||
if (!error || typeof error !== 'object' || !('code' in error)) {
|
||||
return false;
|
||||
}
|
||||
const code = (error as { code?: unknown }).code;
|
||||
return code === 'ENOENT' || code === 'ENOTDIR';
|
||||
}
|
||||
|
||||
export function createDownloadFileAvailabilityProbe(
|
||||
asyncLstat: DownloadAsyncLstat = lstat,
|
||||
maxConcurrent = DEFAULT_MAX_CONCURRENT_FILE_PROBES
|
||||
): DownloadFileAvailabilityProbe {
|
||||
@@ -41,7 +51,10 @@ function createDownloadFileAvailabilityProbe(
|
||||
// deadlines never evict this raw operation, so a stalled lstat remains
|
||||
// charged to the concurrency cap instead of being duplicated by refreshes.
|
||||
// Completed results are discarded so later filesystem changes stay visible.
|
||||
const inFlight = new Map<string, Promise<boolean>>();
|
||||
const inFlight = new Map<
|
||||
string,
|
||||
Promise<DownloadFileAvailabilityProbeResult>
|
||||
>();
|
||||
let active = 0;
|
||||
|
||||
const acquire = (): Promise<void> => {
|
||||
@@ -61,13 +74,15 @@ function createDownloadFileAvailabilityProbe(
|
||||
}
|
||||
};
|
||||
|
||||
const inspect = async (filePath: string): Promise<boolean> => {
|
||||
const inspect = async (
|
||||
filePath: string
|
||||
): Promise<DownloadFileAvailabilityProbeResult> => {
|
||||
await acquire();
|
||||
try {
|
||||
const stats = await asyncLstat(filePath);
|
||||
return stats.isFile() && !stats.isSymbolicLink();
|
||||
} catch {
|
||||
return false;
|
||||
} catch (error) {
|
||||
return isMissingFileSystemError(error) ? false : 'unknown';
|
||||
} finally {
|
||||
release();
|
||||
}
|
||||
|
||||
@@ -74,8 +74,11 @@ variants, contextual buttons, and theme-aware styling.
|
||||
the pending-to-queued/downloaded handoff, and the coordinator returns
|
||||
`added`, `skipped`, and `failed` counts. Xtream and Stalker adapters own
|
||||
provider URL, request header, and metadata preparation; the coordinator owns
|
||||
only provider-neutral orchestration. Both providers use normalized
|
||||
`episode.id` as the canonical
|
||||
only provider-neutral orchestration. When a reserved candidate matches a
|
||||
completed-missing row, the coordinator performs one authoritative preflight
|
||||
refresh before any provider preparation. Restored files therefore become
|
||||
stable skips without requiring a Stalker URL/network request. Both providers
|
||||
use normalized `episode.id` as the canonical
|
||||
episode `xtreamId`; Stalker `originalCmd` and `originalId` participate only
|
||||
in URL resolution.
|
||||
The exact `(playlistId, contentType, xtreamId)` identity is authoritative.
|
||||
@@ -104,8 +107,11 @@ variants, contextual buttons, and theme-aware styling.
|
||||
and report a timeout as missing for that snapshot. The underlying filesystem
|
||||
operation remains coalesced and charged against the four-probe cap until it
|
||||
actually settles, so later callers get independent bounded waits without
|
||||
duplicating stalled native work. The coordinator counts both stable duplicate
|
||||
reasons as skipped. There is no batch IPC,
|
||||
duplicating stalled native work. Only `ENOENT` and `ENOTDIR` are authoritative
|
||||
absence; permission, I/O, and other probe errors remain unknown, so
|
||||
`DOWNLOADS_START` leaves the completed row and file path untouched. The
|
||||
coordinator counts both stable duplicate reasons as skipped. There is no
|
||||
batch IPC,
|
||||
parallel transfer, or queue reordering: destination authorization, persisted
|
||||
header handling, and the backend's one-active-transfer FIFO semantics remain
|
||||
unchanged.
|
||||
|
||||
+67
@@ -311,6 +311,73 @@ describe('SeasonDownloadCoordinator', () => {
|
||||
expect(available.prepare).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('refreshes a completed-missing episode before provider preparation', async () => {
|
||||
const restored = candidate(FIRST_IDENTITY, () =>
|
||||
Promise.reject(new Error('portal unavailable'))
|
||||
);
|
||||
downloadsService.downloads.set([
|
||||
download(restored.identity, {
|
||||
status: 'completed',
|
||||
fileAvailability: 'missing',
|
||||
}),
|
||||
]);
|
||||
downloadsService.loadDownloads.mockImplementation(async () => {
|
||||
downloadsService.downloads.set([
|
||||
download(restored.identity, {
|
||||
status: 'completed',
|
||||
fileAvailability: 'available',
|
||||
}),
|
||||
]);
|
||||
});
|
||||
|
||||
await expect(coordinator.enqueueOne(restored)).resolves.toBe('skipped');
|
||||
|
||||
expect(downloadsService.loadDownloads).toHaveBeenCalledTimes(1);
|
||||
expect(restored.prepare).not.toHaveBeenCalled();
|
||||
expect(downloadsService.startDownload).not.toHaveBeenCalled();
|
||||
expect(coordinator.isPending(restored.identity)).toBe(false);
|
||||
});
|
||||
|
||||
it('refreshes completed-missing season candidates once before provider preparation', async () => {
|
||||
const first = candidate(FIRST_IDENTITY, () =>
|
||||
Promise.reject(new Error('portal unavailable'))
|
||||
);
|
||||
const second = candidate(identity(102, 2), () =>
|
||||
Promise.reject(new Error('portal unavailable'))
|
||||
);
|
||||
downloadsService.downloads.set([
|
||||
download(first.identity, {
|
||||
status: 'completed',
|
||||
fileAvailability: 'missing',
|
||||
}),
|
||||
download(second.identity, {
|
||||
status: 'completed',
|
||||
fileAvailability: 'missing',
|
||||
}),
|
||||
]);
|
||||
downloadsService.loadDownloads.mockImplementation(async () => {
|
||||
downloadsService.downloads.set([
|
||||
download(first.identity, {
|
||||
status: 'completed',
|
||||
fileAvailability: 'available',
|
||||
}),
|
||||
download(second.identity, {
|
||||
status: 'completed',
|
||||
fileAvailability: 'available',
|
||||
}),
|
||||
]);
|
||||
});
|
||||
|
||||
await expect(
|
||||
coordinator.enqueueSeason([first, second])
|
||||
).resolves.toEqual({ added: 0, skipped: 2, failed: 0 });
|
||||
|
||||
expect(downloadsService.loadDownloads).toHaveBeenCalledTimes(1);
|
||||
expect(first.prepare).not.toHaveBeenCalled();
|
||||
expect(second.prepare).not.toHaveBeenCalled();
|
||||
expect(downloadsService.startDownload).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('skips an episode when multiple managed rows claim its coordinates', async () => {
|
||||
const item = candidate(FIRST_IDENTITY);
|
||||
downloadsService.downloads.set([
|
||||
|
||||
+70
-1
@@ -20,6 +20,12 @@ import {
|
||||
type SeasonDownloadResult,
|
||||
} from './season-download.models';
|
||||
|
||||
interface EpisodeDownloadPreflight {
|
||||
readonly ready: EpisodeDownloadCandidate[];
|
||||
readonly skipped: EpisodeDownloadCandidate[];
|
||||
readonly failed: EpisodeDownloadCandidate[];
|
||||
}
|
||||
|
||||
@Injectable({ providedIn: 'root' })
|
||||
export class SeasonDownloadCoordinator {
|
||||
private readonly downloadsService = inject(DownloadsService);
|
||||
@@ -56,6 +62,14 @@ export class SeasonDownloadCoordinator {
|
||||
return EPISODE_DOWNLOAD_SUBMISSIONS.Skipped;
|
||||
}
|
||||
|
||||
const preflight = await this.preflightCompletedMissing([candidate]);
|
||||
if (preflight.ready.length === 0) {
|
||||
this.release(candidate.identity);
|
||||
return preflight.skipped.length > 0
|
||||
? EPISODE_DOWNLOAD_SUBMISSIONS.Skipped
|
||||
: EPISODE_DOWNLOAD_SUBMISSIONS.Failed;
|
||||
}
|
||||
|
||||
const submission = await this.submit(candidate);
|
||||
if (submission === EPISODE_DOWNLOAD_SUBMISSIONS.Failed) {
|
||||
this.release(candidate.identity);
|
||||
@@ -84,8 +98,17 @@ export class SeasonDownloadCoordinator {
|
||||
reserved.push(candidate);
|
||||
}
|
||||
|
||||
const preflight = await this.preflightCompletedMissing(reserved);
|
||||
result.skipped += preflight.skipped.length;
|
||||
result.failed += preflight.failed.length;
|
||||
this.releaseAll(
|
||||
[...preflight.skipped, ...preflight.failed].map(
|
||||
({ identity }) => identity
|
||||
)
|
||||
);
|
||||
|
||||
const refreshPending: EpisodeDownloadIdentity[] = [];
|
||||
for (const candidate of reserved) {
|
||||
for (const candidate of preflight.ready) {
|
||||
const submission = await this.submit(candidate);
|
||||
result[submission] += 1;
|
||||
if (submission !== EPISODE_DOWNLOAD_SUBMISSIONS.Failed) {
|
||||
@@ -106,6 +129,52 @@ export class SeasonDownloadCoordinator {
|
||||
return result;
|
||||
}
|
||||
|
||||
private async preflightCompletedMissing(
|
||||
candidates: readonly EpisodeDownloadCandidate[]
|
||||
): Promise<EpisodeDownloadPreflight> {
|
||||
if (
|
||||
!candidates.some(({ identity }) =>
|
||||
this.isCompletedMissing(identity)
|
||||
)
|
||||
) {
|
||||
return { ready: [...candidates], skipped: [], failed: [] };
|
||||
}
|
||||
|
||||
try {
|
||||
await this.downloadsService.loadDownloads();
|
||||
} catch {
|
||||
this.logger.warn('Completed episode file recheck failed');
|
||||
return { ready: [], skipped: [], failed: [...candidates] };
|
||||
}
|
||||
|
||||
if (
|
||||
!this.downloadsService.isAvailable() ||
|
||||
!this.downloadsService.hasAuthoritativeDownloadList() ||
|
||||
!this.downloadsService.hasLoadedDownloads()
|
||||
) {
|
||||
return { ready: [], skipped: [], failed: [...candidates] };
|
||||
}
|
||||
|
||||
const ready: EpisodeDownloadCandidate[] = [];
|
||||
const skipped: EpisodeDownloadCandidate[] = [];
|
||||
for (const candidate of candidates) {
|
||||
const resolution = this.resolveDownload(candidate.identity);
|
||||
(isEpisodeDownloadEligible(resolution) ? ready : skipped).push(
|
||||
candidate
|
||||
);
|
||||
}
|
||||
return { ready, skipped, failed: [] };
|
||||
}
|
||||
|
||||
private isCompletedMissing(identity: EpisodeDownloadIdentity): boolean {
|
||||
const resolution = this.resolveDownload(identity);
|
||||
return (
|
||||
resolution.kind === 'match' &&
|
||||
resolution.download.status === 'completed' &&
|
||||
resolution.download.fileAvailability === 'missing'
|
||||
);
|
||||
}
|
||||
|
||||
private reserve(candidate: EpisodeDownloadCandidate): boolean {
|
||||
if (!this.isEligible(candidate)) {
|
||||
return false;
|
||||
|
||||
Reference in new issue
Block a user