fix(downloads): reconcile legacy episode identities

This commit is contained in:
4gray committed 2026-08-02 11:04:37 +02:00
1 parent 02fc0351d1
commit 87727f14b7
4 files changed
+562 -19

No files matched your search

@@ -0,0 +1,202 @@
import * as schema from '../../database/schema';
import type { DownloadsDatabase } from './download-task';
import { resolveExistingDownloadIdentity } from './download-request-identity';
type DownloadRow = typeof schema.downloads.$inferSelect;
const episodeRequest = {
contentType: 'episode' as const,
episodeNumber: 3,
playlistId: 'playlist-1',
seasonNumber: 2,
seriesXtreamId: 100,
xtreamId: 700,
};
function createDownloadRow(overrides: Partial<DownloadRow> = {}): DownloadRow {
return {
bytesDownloaded: 0,
contentType: 'episode',
createdAt: '2026-08-02 10:00:00',
episodeNumber: 3,
errorMessage: null,
fileName: 'episode.mp4',
filePath: null,
id: 42,
metadataSnapshot: null,
playlistId: 'playlist-1',
posterUrl: null,
requestHeaders: null,
resumeValidator: null,
seasonNumber: 2,
seriesXtreamId: 100,
status: 'canceled',
title: 'Episode 3',
totalBytes: null,
updatedAt: '2026-08-02 10:00:00',
url: 'https://example.test/episode.mp4',
xtreamId: 700,
...overrides,
};
}
function createQueryHarness(
canonicalRows: DownloadRow[],
coordinateRows: DownloadRow[] = []
) {
let queryIndex = 0;
const limit = jest.fn((value: number) => {
const rows = queryIndex === 0 ? canonicalRows : coordinateRows;
queryIndex += 1;
return Promise.resolve(rows.slice(0, value));
});
const where = jest.fn(() => ({ limit }));
const from = jest.fn(() => ({ where }));
const select = jest.fn(() => ({ from }));
return {
db: { select } as unknown as DownloadsDatabase,
limit,
select,
};
}
describe('resolveExistingDownloadIdentity', () => {
it('returns a canonical match without migration', async () => {
const row = createDownloadRow();
const harness = createQueryHarness([row], []);
await expect(
resolveExistingDownloadIdentity(harness.db, episodeRequest)
).resolves.toEqual({
item: row,
kind: 'match',
migrateCanonicalId: false,
});
expect(harness.limit).toHaveBeenNthCalledWith(1, 1);
expect(harness.limit).toHaveBeenNthCalledWith(2, 2);
});
it('returns a legacy coordinate match with canonical migration', async () => {
const row = createDownloadRow({ xtreamId: 77 });
const harness = createQueryHarness([], [row]);
await expect(
resolveExistingDownloadIdentity(harness.db, episodeRequest)
).resolves.toEqual({
item: row,
kind: 'match',
migrateCanonicalId: true,
});
});
it('fails closed when canonical and coordinate lookups find different rows', async () => {
const harness = createQueryHarness(
[createDownloadRow()],
[createDownloadRow({ id: 43, xtreamId: 77 })]
);
await expect(
resolveExistingDownloadIdentity(harness.db, episodeRequest)
).resolves.toEqual({ kind: 'conflict' });
});
it('fails closed when the canonical row has conflicting complete coordinates', async () => {
const harness = createQueryHarness(
[createDownloadRow({ episodeNumber: 4 })],
[]
);
await expect(
resolveExistingDownloadIdentity(harness.db, episodeRequest)
).resolves.toEqual({ kind: 'conflict' });
});
it('fails closed when multiple rows share the legacy coordinates', async () => {
const harness = createQueryHarness([], [
createDownloadRow({ id: 42, xtreamId: 77 }),
createDownloadRow({ id: 43, xtreamId: 78 }),
]);
await expect(
resolveExistingDownloadIdentity(harness.db, episodeRequest)
).resolves.toEqual({ kind: 'conflict' });
});
it('performs only the canonical lookup for VOD', async () => {
const row = createDownloadRow({
contentType: 'vod',
episodeNumber: null,
seasonNumber: null,
seriesXtreamId: null,
});
const harness = createQueryHarness([row]);
await expect(
resolveExistingDownloadIdentity(harness.db, {
contentType: 'vod',
playlistId: episodeRequest.playlistId,
xtreamId: episodeRequest.xtreamId,
})
).resolves.toEqual({
item: row,
kind: 'match',
migrateCanonicalId: false,
});
expect(harness.limit).toHaveBeenCalledTimes(1);
expect(harness.limit).toHaveBeenCalledWith(1);
});
it.each([
['series id', { seriesXtreamId: undefined }],
['season number', { seasonNumber: undefined }],
['episode number', { episodeNumber: undefined }],
['unsafe series id', { seriesXtreamId: Number.MAX_SAFE_INTEGER + 1 }],
['unsafe season number', { seasonNumber: Number.NaN }],
['unsafe episode number', { episodeNumber: 1.5 }],
])(
'performs only the canonical lookup when an episode has an incomplete or unsafe %s',
async (_label, override) => {
const row = createDownloadRow({
episodeNumber: null,
seasonNumber: null,
seriesXtreamId: null,
});
const harness = createQueryHarness([row]);
await expect(
resolveExistingDownloadIdentity(harness.db, {
...episodeRequest,
...override,
})
).resolves.toEqual({
item: row,
kind: 'match',
migrateCanonicalId: false,
});
expect(harness.limit).toHaveBeenCalledTimes(1);
expect(harness.limit).toHaveBeenCalledWith(1);
}
);
it('lets the canonical match win when both lookups find the same row', async () => {
const row = createDownloadRow();
const harness = createQueryHarness([row], [row]);
await expect(
resolveExistingDownloadIdentity(harness.db, episodeRequest)
).resolves.toEqual({
item: row,
kind: 'match',
migrateCanonicalId: false,
});
});
it('returns none when neither lookup finds a row', async () => {
const harness = createQueryHarness([], []);
await expect(
resolveExistingDownloadIdentity(harness.db, episodeRequest)
).resolves.toEqual({ kind: 'none' });
});
});
@@ -0,0 +1,159 @@
import { and, eq } from 'drizzle-orm';
import * as schema from '../../database/schema';
import type { DownloadsDatabase } from './download-task';
type DownloadRow = typeof schema.downloads.$inferSelect;
export const DOWNLOAD_IDENTITY_KIND = {
CONFLICT: 'conflict',
MATCH: 'match',
NONE: 'none',
} as const;
interface DownloadIdentityNone {
kind: typeof DOWNLOAD_IDENTITY_KIND.NONE;
}
interface DownloadIdentityConflict {
kind: typeof DOWNLOAD_IDENTITY_KIND.CONFLICT;
}
interface DownloadIdentityMatch {
item: DownloadRow;
kind: typeof DOWNLOAD_IDENTITY_KIND.MATCH;
migrateCanonicalId: boolean;
}
export type DownloadIdentityResolution =
| DownloadIdentityNone
| DownloadIdentityConflict
| DownloadIdentityMatch;
export interface DownloadIdentityRequest {
contentType: DownloadRow['contentType'];
episodeNumber?: number;
playlistId: string;
seasonNumber?: number;
seriesXtreamId?: number;
xtreamId: number;
}
function isSafeInteger(value: number | undefined): value is number {
return Number.isSafeInteger(value);
}
function rowHasConflictingCoordinates(
row: DownloadRow,
request: Required<
Pick<
DownloadIdentityRequest,
'episodeNumber' | 'seasonNumber' | 'seriesXtreamId'
>
>
): boolean {
if (
!isSafeInteger(row.seriesXtreamId ?? undefined) ||
!isSafeInteger(row.seasonNumber ?? undefined) ||
!isSafeInteger(row.episodeNumber ?? undefined)
) {
return false;
}
return (
row.seriesXtreamId !== request.seriesXtreamId ||
row.seasonNumber !== request.seasonNumber ||
row.episodeNumber !== request.episodeNumber
);
}
export async function resolveExistingDownloadIdentity(
db: DownloadsDatabase,
request: DownloadIdentityRequest
): Promise<DownloadIdentityResolution> {
const canonicalRows = await db
.select()
.from(schema.downloads)
.where(
and(
eq(schema.downloads.playlistId, request.playlistId),
eq(schema.downloads.contentType, request.contentType),
eq(schema.downloads.xtreamId, request.xtreamId)
)
)
.limit(1);
const canonicalRow = canonicalRows[0];
const seriesXtreamId = request.seriesXtreamId;
const seasonNumber = request.seasonNumber;
const episodeNumber = request.episodeNumber;
if (
request.contentType !== 'episode' ||
!isSafeInteger(seriesXtreamId) ||
!isSafeInteger(seasonNumber) ||
!isSafeInteger(episodeNumber)
) {
return canonicalRow
? {
item: canonicalRow,
kind: DOWNLOAD_IDENTITY_KIND.MATCH,
migrateCanonicalId: false,
}
: { kind: DOWNLOAD_IDENTITY_KIND.NONE };
}
const coordinateRows = await db
.select()
.from(schema.downloads)
.where(
and(
eq(schema.downloads.playlistId, request.playlistId),
eq(schema.downloads.contentType, 'episode'),
eq(schema.downloads.seriesXtreamId, seriesXtreamId),
eq(schema.downloads.seasonNumber, seasonNumber),
eq(schema.downloads.episodeNumber, episodeNumber)
)
)
.limit(2);
if (coordinateRows.length > 1) {
return { kind: DOWNLOAD_IDENTITY_KIND.CONFLICT };
}
const coordinateRow = coordinateRows[0];
if (
canonicalRow &&
coordinateRow &&
canonicalRow.id !== coordinateRow.id
) {
return { kind: DOWNLOAD_IDENTITY_KIND.CONFLICT };
}
if (
canonicalRow &&
rowHasConflictingCoordinates(canonicalRow, {
episodeNumber,
seasonNumber,
seriesXtreamId,
})
) {
return { kind: DOWNLOAD_IDENTITY_KIND.CONFLICT };
}
if (canonicalRow) {
return {
item: canonicalRow,
kind: DOWNLOAD_IDENTITY_KIND.MATCH,
migrateCanonicalId: false,
};
}
if (coordinateRow) {
return {
item: coordinateRow,
kind: DOWNLOAD_IDENTITY_KIND.MATCH,
migrateCanonicalId: true,
};
}
return { kind: DOWNLOAD_IDENTITY_KIND.NONE };
}
@@ -1,4 +1,5 @@
import type { DownloadMetadataSnapshot } from '@iptvnator/shared/interfaces';
import type { Download } from '../../database/schema';
import type { DownloadDirectoryAuthorizer } from './download-directory-authorization';
const metadataSnapshot: DownloadMetadataSnapshot = {
@@ -9,14 +10,17 @@ const metadataSnapshot: DownloadMetadataSnapshot = {
};
async function setupStartMetadataRequest(
existing: Record<string, unknown> | undefined
existing: Record<string, unknown> | undefined,
coordinateRows: Record<string, unknown>[] = []
) {
jest.resetModules();
const schema = await import('../../database/schema');
const playlistLimit = jest.fn().mockResolvedValue([{ id: 'playlist-1' }]);
const downloadLimit = jest
.fn()
.mockResolvedValue(existing ? [existing] : []);
const downloadLimit = jest.fn((limit: number) =>
Promise.resolve(
limit === 2 ? coordinateRows : existing ? [existing] : []
)
);
const from = jest.fn((table: unknown) => ({
where: jest.fn(() => ({
limit: table === schema.playlists ? playlistLimit : downloadLimit,
@@ -54,6 +58,7 @@ async function setupStartMetadataRequest(
return {
authorizer,
db,
downloadLimit,
enqueueDownload,
insertValues,
set,
@@ -61,6 +66,35 @@ async function setupStartMetadataRequest(
};
}
function createStartDownloadRow(
overrides: Partial<Download> = {}
): Download {
return {
bytesDownloaded: 0,
contentType: 'episode',
createdAt: '2026-08-02 10:00:00',
episodeNumber: 3,
errorMessage: null,
fileName: 'episode.mp4',
filePath: null,
id: 42,
metadataSnapshot: null,
playlistId: 'playlist-1',
posterUrl: null,
requestHeaders: null,
resumeValidator: null,
seasonNumber: 2,
seriesXtreamId: 100,
status: 'canceled',
title: 'Episode 3',
totalBytes: null,
updatedAt: '2026-08-02 10:00:00',
url: 'https://example.test/episode.mp4',
xtreamId: 77,
...overrides,
};
}
function startPayload(
snapshot?: DownloadMetadataSnapshot,
contentType: 'vod' | 'episode' = 'vod'
@@ -76,6 +110,18 @@ function startPayload(
};
}
function episodeStartPayload() {
return {
...startPayload(undefined, 'episode'),
episodeNumber: 3,
seasonNumber: 2,
seriesXtreamId: 100,
title: 'Episode 3',
url: 'https://example.test/episode.mp4',
xtreamId: 700,
};
}
describe('download request metadata snapshots', () => {
it('persists an encoded snapshot for a new download', async () => {
const request = await setupStartMetadataRequest(undefined);
@@ -115,6 +161,7 @@ describe('download request metadata snapshots', () => {
expect(request.set.mock.calls[0][0]).not.toHaveProperty(
'metadataSnapshot'
);
expect(request.set.mock.calls[0][0]).not.toHaveProperty('xtreamId');
});
it('replaces stored metadata when a restart supplies a snapshot', async () => {
@@ -383,6 +430,137 @@ describe('download request metadata snapshots', () => {
});
});
describe('download request identity resolution', () => {
it.each(['queued', 'downloading', 'paused'] as const)(
'returns the stable duplicate result for an active legacy-coordinate %s row',
async (status) => {
const legacyRow = createStartDownloadRow({ status });
const request = await setupStartMetadataRequest(undefined, [
legacyRow,
]);
await expect(
request.startDownloadRequest(
episodeStartPayload(),
request.authorizer
)
).resolves.toEqual({
error: 'Download already in progress',
id: legacyRow.id,
reason: 'already-in-progress',
success: false,
});
expect(request.db.insert).not.toHaveBeenCalled();
expect(request.db.update).not.toHaveBeenCalled();
expect(request.enqueueDownload).not.toHaveBeenCalled();
}
);
it('keeps the stable duplicate reason for an exact active row', async () => {
const exactRow = createStartDownloadRow({
status: 'queued',
xtreamId: 700,
});
const request = await setupStartMetadataRequest(exactRow, [exactRow]);
await expect(
request.startDownloadRequest(
episodeStartPayload(),
request.authorizer
)
).resolves.toEqual({
error: 'Download already in progress',
id: exactRow.id,
reason: 'already-in-progress',
success: false,
});
expect(request.enqueueDownload).not.toHaveBeenCalled();
});
it.each(['failed', 'canceled', 'completed'] as const)(
'reuses and migrates a legacy-coordinate %s row before restart',
async (status) => {
const legacyRow = createStartDownloadRow({ status });
const request = await setupStartMetadataRequest(undefined, [
legacyRow,
]);
await expect(
request.startDownloadRequest(
episodeStartPayload(),
request.authorizer
)
).resolves.toEqual({ id: legacyRow.id, success: true });
expect(request.set).toHaveBeenCalledTimes(1);
expect(request.set).toHaveBeenCalledWith(
expect.objectContaining({
status: 'queued',
xtreamId: episodeStartPayload().xtreamId,
})
);
expect(request.enqueueDownload).toHaveBeenCalledWith(
expect.objectContaining({ id: legacyRow.id })
);
}
);
it('rejects conflicting canonical and coordinate rows without enqueueing', async () => {
const canonicalRow = createStartDownloadRow({
episodeNumber: null,
seasonNumber: null,
seriesXtreamId: null,
xtreamId: 700,
});
const coordinateRow = createStartDownloadRow({ id: 43 });
const request = await setupStartMetadataRequest(canonicalRow, [
coordinateRow,
]);
await expect(
request.startDownloadRequest(
episodeStartPayload(),
request.authorizer
)
).resolves.toEqual({
error: 'Download identity conflict',
success: false,
});
expect(request.db.insert).not.toHaveBeenCalled();
expect(request.db.update).not.toHaveBeenCalled();
expect(request.enqueueDownload).not.toHaveBeenCalled();
});
it('propagates a rejected canonical-id migration before enqueueing', async () => {
const migrationError = new Error(
'UNIQUE constraint failed: downloads.xtream_id'
);
const legacyRow = createStartDownloadRow({ status: 'canceled' });
const request = await setupStartMetadataRequest(undefined, [legacyRow]);
request.set.mockReturnValueOnce({
where: jest.fn().mockRejectedValue(migrationError),
});
await expect(
request.startDownloadRequest(
episodeStartPayload(),
request.authorizer
)
).rejects.toBe(migrationError);
expect(request.set).toHaveBeenCalledWith(
expect.objectContaining({
xtreamId: episodeStartPayload().xtreamId,
})
);
expect(request.db.insert).not.toHaveBeenCalled();
expect(request.enqueueDownload).not.toHaveBeenCalled();
});
});
describe('download requests resume', () => {
it('enqueues a paused download with stored headers and original target path', async () => {
jest.resetModules();
@@ -1,4 +1,7 @@
import type { DownloadMetadataSnapshot } from '@iptvnator/shared/interfaces';
import type {
DownloadMetadataSnapshot,
ElectronBridgeDownloadStartResult,
} from '@iptvnator/shared/interfaces';
import { and, eq, sql } from 'drizzle-orm';
import { basename, dirname, extname } from 'node:path';
import { getDatabase } from '../../database/connection';
@@ -6,6 +9,7 @@ import * as schema from '../../database/schema';
import { assertRemoteUrlAllowed } from '../url-safety';
import { DownloadDirectoryAuthorizer } from './download-directory-authorization';
import { removePartialDownloadFile } from './download-file-path';
import { resolveExistingDownloadIdentity } from './download-request-identity';
import {
assertDownloadMetadataArtworkDiffersFromStream,
assertDownloadMetadataMatchesContentType,
@@ -117,7 +121,7 @@ export function parseStoredHeaders(
export async function startDownloadRequest(
data: StartDownloadRequest,
authorizer: DownloadDirectoryAuthorizer
): Promise<{ success: boolean; error?: string; id?: number }> {
): Promise<ElectronBridgeDownloadStartResult> {
const encodedMetadataSnapshot =
data.metadataSnapshot === undefined
? undefined
@@ -146,22 +150,18 @@ export async function startDownloadRequest(
.from(schema.playlists)
.where(eq(schema.playlists.id, data.playlistId))
.limit(1);
const existing = await db
.select()
.from(schema.downloads)
.where(
and(
eq(schema.downloads.playlistId, data.playlistId),
eq(schema.downloads.xtreamId, data.xtreamId),
eq(schema.downloads.contentType, data.contentType)
)
)
.limit(1);
const identity = await resolveExistingDownloadIdentity(db, data);
if (identity.kind === 'conflict') {
return {
error: 'Download identity conflict',
success: false,
};
}
const fileName = createFileName(data.title, data.url);
const headers = createHeaders(data.headers);
if (existing.length > 0) {
const item = existing[0];
if (identity.kind === 'match') {
const item = identity.item;
if (normalizedMetadataSnapshot) {
assertDownloadMetadataMatchesContentType(
normalizedMetadataSnapshot,
@@ -180,6 +180,7 @@ export async function startDownloadRequest(
return {
error: 'Download already in progress',
id: item.id,
reason: 'already-in-progress',
success: false,
};
}
@@ -220,6 +221,9 @@ export async function startDownloadRequest(
totalBytes: null,
updatedAt: sql`CURRENT_TIMESTAMP`,
url: data.url,
...(identity.migrateCanonicalId
? { xtreamId: data.xtreamId }
: {}),
})
.where(eq(schema.downloads.id, item.id));
enqueueDownload({