feat(downloads): download completed Xtream catch-up programmes as TS (#1572)

* feat(epg): copy catch-up programme URLs without changing playback

* feat(downloads): save completed Xtream archive programmes as TS

* fix(epg): let newer archive copy requests supersede pending work

* fix(downloads): protect archive partials and independent submissions

* fix(downloads): verify archive identity through finalization

* fix(downloads): bound archive storage and capture cleanup entries

* fix(downloads): preserve archive ownership across failure paths

* fix(downloads): recover explicitly verified archive completions

* fix(downloads): journal archive promotion before publishing files

* fix(downloads): reset archive proof before an explicit restart

* fix(downloads): preserve archive recovery ownership and interruption

* fix(downloads): verify durable archive identity at resume open

* fix(downloads): fence archive commands during completion commit

* fix(downloads): persist archive ownership throughout its lifecycle

* fix(downloads): protect archive removal and missing-file recovery

* fix(downloads): journal private cleanup captures for recovery

* fix(downloads): journal active archive cleanup before removal

* fix(downloads): clean settled archives before deleting stale rows

* fix(downloads): preserve archive ownership on removal and resubmission

* fix(downloads): recover proven archive completions before retry

* fix(downloads): recover local archives before remote transfer checks

* test(downloads): resolve archive fixture from workspace root

* fix(downloads): distinguish reused archive inodes by creation time

* fix(downloads): bind fresh archive reservations to owned files

* fix(downloads): clean reservations when ownership writes fail

* fix(downloads): commit archive reservation and ownership atomically

* fix(downloads): retain captures until replacement restoration succeeds

* fix(downloads): require durable ownership before cleanup relocation

* fix(downloads): preserve 64-bit archive file identities on Windows

* refactor(release): keep capture fixture constants in their shared module

* fix(downloads): preserve the last link of captured foreign files

* fix(downloads): expose retained archive recovery files

* fix(downloads): keep recovery instructions open while copying
This commit is contained in:
4gray authored and GitHub committed 2026-09-08 20:33:05 +02:00
1 parent a7f3860102
commit bad8a0991e
117 files changed
+7510 -880

No files matched your search

@@ -56,7 +56,7 @@ const completedMigrationStateRule: HandlerRule = [
/** The live title index, already carrying the folding tokenizer. */
const foldedIndexRule: HandlerRule = [
'SELECT sql FROM sqlite_master',
"name = 'content_title_fts'",
{
get: () => ({
sql: "CREATE VIRTUAL TABLE content_title_fts USING fts5(title, tokenize='trigram remove_diacritics 1')",
+6 -29
View File
@@ -1,3 +1,8 @@
import {
DOWNLOADS_TABLE_SQL,
DOWNLOADS_INDEX_STATEMENTS,
ensureDownloadsCatchupSchema,
} from './download-schema';
/**
* Database connection and initialization for IPTVnator
* Uses Drizzle ORM with better-sqlite3
@@ -90,35 +95,6 @@ const TMDB_METADATA_TABLE_SQL = `CREATE TABLE IF NOT EXISTS tmdb_metadata (
fetched_at TEXT DEFAULT (datetime('now'))
)`;
const TMDB_METADATA_INDEX_SQL = `CREATE UNIQUE INDEX IF NOT EXISTS tmdb_metadata_lookup_unique ON tmdb_metadata(media_type, lookup_key, language)`;
const DOWNLOADS_TABLE_SQL = `CREATE TABLE IF NOT EXISTS downloads (
id INTEGER PRIMARY KEY AUTOINCREMENT,
playlist_id TEXT NOT NULL,
xtream_id INTEGER NOT NULL,
content_type TEXT NOT NULL CHECK (content_type IN ('vod', 'episode')),
series_xtream_id INTEGER,
season_number INTEGER,
episode_number INTEGER,
episode_identity_scope TEXT,
title TEXT NOT NULL,
url TEXT NOT NULL,
file_name TEXT,
file_path TEXT,
poster_url TEXT,
request_headers TEXT,
resume_validator TEXT,
metadata_snapshot TEXT,
status TEXT NOT NULL DEFAULT 'queued' CHECK (status IN ('queued', 'downloading', 'paused', 'completed', 'failed', 'canceled')),
bytes_downloaded INTEGER DEFAULT 0,
total_bytes INTEGER,
error_message TEXT,
created_at TEXT DEFAULT (datetime('now')),
updated_at TEXT DEFAULT (datetime('now'))
)`;
const DOWNLOADS_INDEX_STATEMENTS = [
`CREATE UNIQUE INDEX IF NOT EXISTS downloads_xtream_playlist_unique ON downloads(xtream_id, playlist_id, content_type)`,
`CREATE INDEX IF NOT EXISTS downloads_playlist_idx ON downloads(playlist_id)`,
`CREATE INDEX IF NOT EXISTS downloads_status_idx ON downloads(status)`,
];
// No unique index: re-recording the same channel is a normal workflow, and
// playlist_id carries no FK so recordings survive source deletion.
const RECORDINGS_TABLE_SQL = `CREATE TABLE IF NOT EXISTS recordings (
@@ -1165,6 +1141,7 @@ function runMigrations(sqliteDb: Database.Database): void {
cleanupLegacyTmdbSearchCache(sqliteDb);
ensureDownloadsPauseResumeSchema(sqliteDb);
runMigrationStatements(sqliteDb, COLUMN_MIGRATION_STATEMENTS);
ensureDownloadsCatchupSchema(sqliteDb);
// The tokenizer upgrade recreates and rebuilds the index itself, so the
// plain rebuild below would only repeat work it just did.
if (!upgradeContentTitleFtsTokenizer(sqliteDb)) {
@@ -0,0 +1,65 @@
import { execFileSync } from 'node:child_process';
import { createRequire } from 'node:module';
import { resolve } from 'node:path';
import { pathToFileURL } from 'node:url';
it('migrates existing downloads without losing files and separates programme identities', () => {
const electron = createRequire(__filename)('electron') as string;
const moduleUrl = pathToFileURL(
resolve(__dirname, 'download-schema.ts')
).href;
const result = execFileSync(
electron,
[
'--import',
'tsx',
'--eval',
`
const { default: Database } = await import('better-sqlite3');
const { DOWNLOADS_TABLE_SQL, ensureDownloadsCatchupSchema } = await import(${JSON.stringify(moduleUrl)});
const db = new Database(':memory:');
db.pragma('foreign_keys = ON');
db.exec(DOWNLOADS_TABLE_SQL.replace(", 'catchup'", '').replace('programme_start INTEGER NOT NULL DEFAULT 0,', '').replace('catchup TEXT,', ''));
db.exec("CREATE UNIQUE INDEX downloads_xtream_playlist_unique ON downloads(xtream_id, playlist_id, content_type)");
db.prepare("INSERT INTO downloads (id,playlist_id,xtream_id,content_type,title,url,status,file_path,bytes_downloaded,resume_validator,request_headers) VALUES (42,'p',1,'vod','Movie','https://host/movie','paused','/safe/movie.mp4',123,'etag','headers')").run();
ensureDownloadsCatchupSchema(db);
ensureDownloadsCatchupSchema(db);
const insert = db.prepare("INSERT INTO downloads (playlist_id,xtream_id,content_type,programme_start,title,url) VALUES ('p',1,?,?, 'Show','https://host/archive')");
insert.run('catchup', 1000); insert.run('catchup', 2000);
let duplicateArchive = false, duplicateMovie = false;
try { insert.run('catchup', 1000); } catch { duplicateArchive = true; }
try { insert.run('vod', 0); } catch { duplicateMovie = true; }
db.prepare("INSERT INTO downloads (id,playlist_id,xtream_id,content_type,title,url) VALUES (999,'p',99,'catchup','Proof','https://host/archive')").run();
db.prepare("INSERT INTO download_archive_finalizations (download_id,proof) VALUES (999,'{}')").run();
db.prepare('DELETE FROM downloads WHERE id=999').run();
const journalCount = db.prepare('SELECT count(*) AS n FROM download_archive_finalizations').get().n;
process.stdout.write(JSON.stringify({journalCount,row:db.prepare('SELECT * FROM downloads WHERE id=42').get(), count:db.prepare('SELECT count(*) AS n FROM downloads').get().n, duplicateArchive,duplicateMovie}));
`,
],
{
cwd: process.cwd(),
encoding: 'utf8',
env: {
...process.env,
ELECTRON_RUN_AS_NODE: '1',
TSX_TSCONFIG_PATH: resolve(process.cwd(), 'tsconfig.base.json'),
},
}
);
expect(JSON.parse(result)).toMatchObject({
count: 3,
journalCount: 0,
duplicateArchive: true,
duplicateMovie: true,
row: {
id: 42,
status: 'paused',
file_path: '/safe/movie.mp4',
bytes_downloaded: 123,
resume_validator: 'etag',
request_headers: 'headers',
programme_start: 0,
catchup: null,
},
});
});
@@ -0,0 +1,94 @@
import type Database from 'better-sqlite3';
export const DOWNLOADS_TABLE_SQL = `CREATE TABLE IF NOT EXISTS downloads (
id INTEGER PRIMARY KEY AUTOINCREMENT,
playlist_id TEXT NOT NULL,
xtream_id INTEGER NOT NULL,
content_type TEXT NOT NULL CHECK (content_type IN ('vod', 'episode', 'catchup')),
programme_start INTEGER NOT NULL DEFAULT 0,
catchup TEXT,
series_xtream_id INTEGER,
season_number INTEGER,
episode_number INTEGER,
episode_identity_scope TEXT,
title TEXT NOT NULL,
url TEXT NOT NULL,
file_name TEXT,
file_path TEXT,
poster_url TEXT,
request_headers TEXT,
resume_validator TEXT,
metadata_snapshot TEXT,
status TEXT NOT NULL DEFAULT 'queued' CHECK (status IN ('queued', 'downloading', 'paused', 'completed', 'failed', 'canceled')),
bytes_downloaded INTEGER DEFAULT 0,
total_bytes INTEGER,
error_message TEXT,
created_at TEXT DEFAULT (datetime('now')),
updated_at TEXT DEFAULT (datetime('now'))
)`;
export const DOWNLOADS_INDEX_STATEMENTS = [
`CREATE UNIQUE INDEX IF NOT EXISTS downloads_xtream_playlist_unique ON downloads(xtream_id, playlist_id, content_type) WHERE content_type != 'catchup'`,
`CREATE INDEX IF NOT EXISTS downloads_playlist_idx ON downloads(playlist_id)`,
`CREATE INDEX IF NOT EXISTS downloads_status_idx ON downloads(status)`,
];
const ARCHIVE_FINALIZATIONS_SQL = `CREATE TABLE IF NOT EXISTS download_archive_finalizations (
download_id INTEGER PRIMARY KEY REFERENCES downloads(id) ON DELETE CASCADE,
proof TEXT NOT NULL
)`;
const CATCHUP_INDEX = `CREATE UNIQUE INDEX IF NOT EXISTS downloads_catchup_unique
ON downloads(xtream_id, playlist_id, programme_start) WHERE content_type = 'catchup'`;
/** Widen the CHECK transactionally; existing ids, files and resume state survive. */
export function ensureDownloadsCatchupSchema(db: Database.Database): void {
const row = db
.prepare(
"SELECT sql FROM sqlite_master WHERE type = 'table' AND name = 'downloads'"
)
.get() as { sql: string } | undefined;
if (!row) return;
if (row.sql.includes("'catchup'")) {
db.exec(CATCHUP_INDEX);
db.exec(ARCHIVE_FINALIZATIONS_SQL);
return;
}
const columns = [
'id',
'playlist_id',
'xtream_id',
'content_type',
'series_xtream_id',
'season_number',
'episode_number',
'episode_identity_scope',
'title',
'url',
'file_name',
'file_path',
'poster_url',
'request_headers',
'resume_validator',
'metadata_snapshot',
'status',
'bytes_downloaded',
'total_bytes',
'error_message',
'created_at',
'updated_at',
];
db.transaction(() => {
db.exec('DROP INDEX IF EXISTS downloads_xtream_playlist_unique');
db.exec('DROP INDEX IF EXISTS downloads_playlist_idx');
db.exec('DROP INDEX IF EXISTS downloads_status_idx');
db.exec('ALTER TABLE downloads RENAME TO downloads_catchup_legacy');
db.exec(DOWNLOADS_TABLE_SQL);
db.exec(
`INSERT INTO downloads (${columns.join(',')}) SELECT ${columns.join(',')} FROM downloads_catchup_legacy`
);
db.exec('DROP TABLE downloads_catchup_legacy');
for (const statement of DOWNLOADS_INDEX_STATEMENTS) db.exec(statement);
db.exec(CATCHUP_INDEX);
db.exec(ARCHIVE_FINALIZATIONS_SQL);
})();
}
@@ -0,0 +1,84 @@
import type { CatchupDownloadMetadata } from '@iptvnator/shared/interfaces';
import { sql } from 'drizzle-orm';
import {
index,
integer,
sqliteTable,
text,
uniqueIndex,
} from 'drizzle-orm/sqlite-core';
// Downloads table
export const downloads = sqliteTable(
'downloads',
{
id: integer('id').primaryKey({ autoIncrement: true }),
playlistId: text('playlist_id').notNull(),
// Content identifiers
xtreamId: integer('xtream_id').notNull(),
contentType: text('content_type', {
enum: ['vod', 'episode', 'catchup'],
}).notNull(),
programmeStart: integer('programme_start').notNull().default(0),
catchup: text('catchup', {
mode: 'json',
}).$type<CatchupDownloadMetadata>(),
// For episodes: store series info
seriesXtreamId: integer('series_xtream_id'),
seasonNumber: integer('season_number'),
episodeNumber: integer('episode_number'),
episodeIdentityScope: text('episode_identity_scope'),
// Download metadata
title: text('title').notNull(),
url: text('url').notNull(),
fileName: text('file_name'),
filePath: text('file_path'),
posterUrl: text('poster_url'),
requestHeaders: text('request_headers'),
resumeValidator: text('resume_validator'),
metadataSnapshot: text('metadata_snapshot'),
// Download progress
status: text('status', {
enum: [
'queued',
'downloading',
'paused',
'completed',
'failed',
'canceled',
],
})
.notNull()
.default('queued'),
bytesDownloaded: integer('bytes_downloaded').default(0),
totalBytes: integer('total_bytes'),
errorMessage: text('error_message'),
// Timestamps
createdAt: text('created_at').default(sql`CURRENT_TIMESTAMP`),
updatedAt: text('updated_at').default(sql`CURRENT_TIMESTAMP`),
},
(table) => ({
playlistIdx: index('downloads_playlist_idx').on(table.playlistId),
statusIdx: index('downloads_status_idx').on(table.status),
xtreamPlaylistUnique: uniqueIndex('downloads_xtream_playlist_unique')
.on(table.xtreamId, table.playlistId, table.contentType)
.where(sql`${table.contentType} != 'catchup'`),
catchupUnique: uniqueIndex('downloads_catchup_unique')
.on(table.xtreamId, table.playlistId, table.programmeStart)
.where(sql`${table.contentType} = 'catchup'`),
})
);
export type Download = typeof downloads.$inferSelect;
export type NewDownload = typeof downloads.$inferInsert;
// Write-ahead archive promotion proof. It outlives process-local DownloadTask.
export const downloadArchiveFinalizations = sqliteTable(
'download_archive_finalizations',
{
downloadId: integer('download_id')
.primaryKey()
.references(() => downloads.id, { onDelete: 'cascade' }),
proof: text('proof').notNull(),
}
);
+1 -56
View File
@@ -335,62 +335,7 @@ export type NewEpgProgramDb = typeof epgPrograms.$inferInsert;
export type PlaybackPosition = typeof playbackPositions.$inferSelect;
export type NewPlaybackPosition = typeof playbackPositions.$inferInsert;
// Downloads table
export const downloads = sqliteTable(
'downloads',
{
id: integer('id').primaryKey({ autoIncrement: true }),
playlistId: text('playlist_id').notNull(),
// Content identifiers
xtreamId: integer('xtream_id').notNull(),
contentType: text('content_type', {
enum: ['vod', 'episode'],
}).notNull(),
// For episodes: store series info
seriesXtreamId: integer('series_xtream_id'),
seasonNumber: integer('season_number'),
episodeNumber: integer('episode_number'),
episodeIdentityScope: text('episode_identity_scope'),
// Download metadata
title: text('title').notNull(),
url: text('url').notNull(),
fileName: text('file_name'),
filePath: text('file_path'),
posterUrl: text('poster_url'),
requestHeaders: text('request_headers'),
resumeValidator: text('resume_validator'),
metadataSnapshot: text('metadata_snapshot'),
// Download progress
status: text('status', {
enum: [
'queued',
'downloading',
'paused',
'completed',
'failed',
'canceled',
],
})
.notNull()
.default('queued'),
bytesDownloaded: integer('bytes_downloaded').default(0),
totalBytes: integer('total_bytes'),
errorMessage: text('error_message'),
// Timestamps
createdAt: text('created_at').default(sql`CURRENT_TIMESTAMP`),
updatedAt: text('updated_at').default(sql`CURRENT_TIMESTAMP`),
},
(table) => ({
playlistIdx: index('downloads_playlist_idx').on(table.playlistId),
statusIdx: index('downloads_status_idx').on(table.status),
xtreamPlaylistUnique: uniqueIndex(
'downloads_xtream_playlist_unique'
).on(table.xtreamId, table.playlistId, table.contentType),
})
);
export type Download = typeof downloads.$inferSelect;
export type NewDownload = typeof downloads.$inferInsert;
export * from './download-tables';
// Live-TV recordings table. Rows are created by the embedded-MPV recording
// tracker, not by the download queue: a recording has no source URL to