fix(downloads): bound file probe callers

This commit is contained in:
4gray committed 2026-08-02 21:30:46 +02:00
1 parent e8640996e9
commit 3ef5293cf4
4 files changed
+125 -77

No files matched your search

+9 -8
View File
@@ -947,14 +947,15 @@ engine` (restart required) or
comes from an asynchronous main-process filesystem recheck before a
completed-missing row can be reset, so a file restored after the renderer
snapshot is not orphaned or downloaded again. The recheck has a one-second
deadline; timeout or probe failure leaves the row untouched and reports a
failed submission so the season loop can continue. Completed-file list
probes use the same deadline; a timeout releases their coalescing slot and is
reported as missing for that snapshot, so a later refresh performs a new
filesystem check instead of joining the stalled operation. 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
caller deadline that starts before shared-slot acquisition; timeout or probe
failure leaves the row untouched and reports a failed submission so the
season loop can continue. Completed-file list callers use the same deadline
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
loading/empty-state resolution intact but disables starts until another
snapshot succeeds.
- Episode ownership uses normalized `episode.id` as the canonical `xtreamId`
@@ -23,9 +23,7 @@ type DownloadAsyncLstat = (
type DownloadFileAvailabilityProbeResult = boolean | 'unknown';
type DownloadFileAvailabilityProbe = (
filePath: string
) => Promise<DownloadFileAvailabilityProbeResult>;
type DownloadFileAvailabilityProbe = (filePath: string) => Promise<boolean>;
const DEFAULT_MAX_CONCURRENT_FILE_PROBES = 4;
const DEFAULT_FILE_PROBE_TIMEOUT_MS = 1_000;
@@ -39,12 +37,11 @@ function createDownloadFileAvailabilityProbe(
): DownloadFileAvailabilityProbe {
const concurrency = Math.max(1, Math.floor(maxConcurrent));
const pending: Array<() => void> = [];
// Coalesce only active probes. Completed results are discarded so an
// externally removed file is visible on the next list refresh.
const inFlight = new Map<
string,
Promise<DownloadFileAvailabilityProbeResult>
>();
// Coalesce unfinished probes, including work waiting for a slot. Caller
// 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>>();
let active = 0;
const acquire = (): Promise<void> => {
@@ -64,28 +61,14 @@ function createDownloadFileAvailabilityProbe(
}
};
const inspect = async (
filePath: string
): Promise<DownloadFileAvailabilityProbeResult> => {
const inspect = async (filePath: string): Promise<boolean> => {
await acquire();
let timeout: ReturnType<typeof setTimeout> | undefined;
try {
const timedOut = new Promise<'unknown'>((resolve) => {
timeout = setTimeout(
() => resolve('unknown'),
DEFAULT_FILE_PROBE_TIMEOUT_MS
);
});
return await Promise.race([
asyncLstat(filePath)
.then((stats) => stats.isFile() && !stats.isSymbolicLink())
.catch(() => false),
timedOut,
]);
const stats = await asyncLstat(filePath);
return stats.isFile() && !stats.isSymbolicLink();
} catch {
return false;
} finally {
if (timeout !== undefined) {
clearTimeout(timeout);
}
release();
}
};
@@ -119,6 +102,31 @@ function toBoundedDownloadFileAvailability(
return available ? 'available' : 'missing';
}
async function probeDownloadFileAvailabilityWithTimeout(
filePath: string,
timeoutMs: number,
probe: DownloadFileAvailabilityProbe
): Promise<DownloadFileAvailabilityProbeResult> {
const boundedTimeoutMs =
Number.isFinite(timeoutMs) && timeoutMs >= 0
? timeoutMs
: DEFAULT_FILE_PROBE_TIMEOUT_MS;
let timeout: ReturnType<typeof setTimeout> | undefined;
const timedOut = new Promise<'unknown'>((resolve) => {
timeout = setTimeout(() => resolve('unknown'), boundedTimeoutMs);
});
try {
return await Promise.race([probe(filePath), timedOut]);
} catch {
return 'unknown';
} finally {
if (timeout !== undefined) {
clearTimeout(timeout);
}
}
}
export function isAvailableDownloadFile(
filePath: string | null | undefined,
lstat: DownloadLstat = lstatSync
@@ -160,14 +168,18 @@ export async function getDownloadFileAvailabilityAsync(
return 'missing';
}
const available = await probe(download.filePath);
const available = await probeDownloadFileAvailabilityWithTimeout(
download.filePath,
DEFAULT_FILE_PROBE_TIMEOUT_MS,
probe
);
return available === true ? 'available' : 'missing';
}
export async function getDownloadFileAvailabilityWithTimeoutAsync(
download: DownloadFileRow,
timeoutMs = DEFAULT_FILE_PROBE_TIMEOUT_MS,
probe?: DownloadFileAvailabilityProbe
probe: DownloadFileAvailabilityProbe = probeDownloadFileAvailability
): Promise<BoundedDownloadFileAvailability> {
if (download.status !== 'completed') {
return 'not-applicable';
@@ -177,34 +189,13 @@ export async function getDownloadFileAvailabilityWithTimeoutAsync(
return 'missing';
}
if (!probe) {
return toBoundedDownloadFileAvailability(
await probeDownloadFileAvailability(download.filePath)
);
}
const boundedTimeoutMs =
Number.isFinite(timeoutMs) && timeoutMs >= 0
? timeoutMs
: DEFAULT_FILE_PROBE_TIMEOUT_MS;
let timeout: ReturnType<typeof setTimeout> | undefined;
const timedOut = new Promise<'unknown'>((resolve) => {
timeout = setTimeout(() => resolve('unknown'), boundedTimeoutMs);
});
try {
const available = await Promise.race([
probe(download.filePath),
timedOut,
]);
return toBoundedDownloadFileAvailability(available);
} catch {
return 'unknown';
} finally {
if (timeout !== undefined) {
clearTimeout(timeout);
}
}
return toBoundedDownloadFileAvailability(
await probeDownloadFileAvailabilityWithTimeout(
download.filePath,
timeoutMs,
probe
)
);
}
export function decorateDownloadItem<T extends DownloadFileRow>(
@@ -77,7 +77,7 @@ describe('downloads events: file availability', () => {
await expect(response).resolves.toHaveLength(2);
});
it('times out an unresponsive probe and lets the next list refresh recheck the file', async () => {
it('bounds each refresh without duplicating an unresponsive filesystem probe', async () => {
jest.useFakeTimers();
try {
const row = {
@@ -91,10 +91,13 @@ describe('downloads events: file availability', () => {
from: jest.fn(() => ({ orderBy })),
})),
});
let finishFirstProbe!: (
value: ReturnType<typeof regularFile>
) => void;
mockLstat
.mockReturnValueOnce(
new Promise(() => {
// Simulate an unresponsive removable/network mount.
new Promise((resolve) => {
finishFirstProbe = resolve;
})
)
.mockResolvedValueOnce(regularFile());
@@ -110,6 +113,21 @@ describe('downloads events: file availability', () => {
},
]);
const secondTimedOutRefresh =
getHandler('DOWNLOADS_GET_LIST')(null);
await jest.advanceTimersByTimeAsync(1_000);
await expect(secondTimedOutRefresh).resolves.toEqual([
{
...row,
metadataSnapshot: undefined,
fileAvailability: 'missing',
},
]);
expect(mockLstat).toHaveBeenCalledTimes(1);
finishFirstProbe(regularFile());
await jest.advanceTimersByTimeAsync(0);
await expect(
getHandler('DOWNLOADS_GET_LIST')(null)
).resolves.toEqual([
@@ -125,6 +143,43 @@ describe('downloads events: file availability', () => {
}
}, 500);
it('starts the list deadline before a completed-file probe waits for a slot', async () => {
jest.useFakeTimers();
try {
const rows = Array.from({ length: 5 }, (_, index) => ({
filePath: `/downloads/offline/movie-${index}.mp4`,
id: index + 1,
status: 'completed',
}));
const orderBy = jest.fn().mockResolvedValue(rows);
mockGetDatabase.mockResolvedValue({
select: jest.fn(() => ({
from: jest.fn(() => ({ orderBy })),
})),
});
mockLstat.mockImplementation(
() =>
new Promise(() => {
// Keep all four shared filesystem slots occupied.
})
);
const refresh = getHandler('DOWNLOADS_GET_LIST')(null);
await jest.advanceTimersByTimeAsync(1_000);
await expect(refresh).resolves.toEqual(
rows.map((row) => ({
...row,
metadataSnapshot: undefined,
fileAvailability: 'missing',
}))
);
expect(mockLstat).toHaveBeenCalledTimes(4);
} finally {
jest.useRealTimers();
}
}, 500);
it('starts at most four completed-file probes concurrently', async () => {
const rows = Array.from({ length: 6 }, (_, index) => ({
filePath: `/downloads/network/movie-${index}.mp4`,
+9 -8
View File
@@ -96,14 +96,15 @@ variants, contextual buttons, and theme-aware styling.
such a completed row, `DOWNLOADS_START` asynchronously rechecks its retained
path in the main process. A restored file returns stable
`reason: 'already-downloaded'` without mutation; active matches return
`reason: 'already-in-progress'`. The recheck has a one-second deadline;
timeout or probe failure leaves the row untouched and returns a failed
submission, allowing the sequential season loop to continue. Completed-file
probes used by list refreshes have the same deadline. A timed-out probe is
released from the coalescing pool, the snapshot reports the file as missing,
and a later refresh performs a fresh check instead of reusing the stalled
filesystem operation. The coordinator counts both stable duplicate reasons
as skipped. There is no batch IPC,
`reason: 'already-in-progress'`. The recheck has a one-second caller deadline
that starts before shared-slot acquisition; timeout or probe failure leaves
the row untouched and returns a failed submission, allowing the sequential
season loop to continue. Completed-file list callers have the same deadline
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,
parallel transfer, or queue reordering: destination authorization, persisted
header handling, and the backend's one-active-transfer FIFO semantics remain
unchanged.