feat(notifications): detect persistent download stalls

This commit is contained in:
Sucukdeluxe
2026-08-22 07:23:54 +02:00
parent 4fd30802e5
commit 311475c93b
7 changed files with 1560 additions and 108 deletions
+55 -8
View File
@@ -80,6 +80,7 @@ import { getLegacyDesktopLogDirectory, migrateLogDirectories, prepareLogDirector
import { normalizeStatisticsLedger, saveStatisticsLedger } from "./statistics-ledger";
import { NotificationOutbox } from "./notification-outbox";
import { sendNotification } from "./notify";
import { DownloadHealthMonitor } from "./download-health-monitor";
function sanitizeSettingsPatch(partial: Partial<AppSettings>): Partial<AppSettings> {
const entries = Object.entries(partial || {}).filter(([, value]) => value !== undefined);
@@ -122,6 +123,12 @@ export class AppController {
private notificationOutbox: NotificationOutbox;
private downloadHealthMonitor: DownloadHealthMonitor;
private downloadHealthTimer: NodeJS.Timeout | null = null;
private downloadHealthEvaluation: Promise<void> | null = null;
private logDirectory = this.storagePaths.baseDir;
private onStateHandler: ((snapshot: UiSnapshot) => void) | null = null;
@@ -168,6 +175,7 @@ export class AppController {
timestamp: event.createdAt
})
});
this.downloadHealthMonitor = new DownloadHealthMonitor(this.storagePaths.notificationHealthFile);
void this.notificationOutbox.drain().catch((error) => {
logger.warn(`Notification-Outbox konnte nicht gestartet werden: ${String(error)}`);
});
@@ -189,9 +197,14 @@ export class AppController {
this.recordHistoryEntry(entry);
}
});
this.manager.on("state", (snapshot: UiSnapshot) => {
this.onStateHandler?.(snapshot);
});
this.manager.on("state", (snapshot: UiSnapshot) => {
this.onStateHandler?.(snapshot);
});
void this.evaluateDownloadHealth();
this.downloadHealthTimer = setInterval(() => {
void this.evaluateDownloadHealth();
}, 15_000);
this.downloadHealthTimer.unref?.();
logger.info(`App gestartet v${APP_VERSION}`);
logger.info(`Log-Datei: ${getLogFilePath()}`);
logAuditEvent("INFO", "App gestartet", {
@@ -1284,12 +1297,46 @@ export class AppController {
return this.manager.getPackageLogPath(packageId) || getPackageLogPath(packageId);
}
public getItemLogPath(itemId: string): string | null {
return this.manager.getItemLogPath(itemId) || getItemLogPath(itemId);
}
public getItemLogPath(itemId: string): string | null {
return this.manager.getItemLogPath(itemId) || getItemLogPath(itemId);
}
private evaluateDownloadHealth(): Promise<void> {
if (this.downloadHealthEvaluation) {
return this.downloadHealthEvaluation;
}
const now = Date.now();
const task = this.downloadHealthMonitor.sample(
this.manager.getDownloadHealthSnapshot(now),
now,
{
stallAfterMs: this.settings.notifyStallAfterSeconds * 1000,
cooldownMs: this.settings.notifyStallCooldownMinutes * 60_000,
notifyOnStall: this.settings.notifyOnDownloadStall && Boolean(String(this.settings.notifyUrl || "").trim()),
notifyOnRecovery: this.settings.notifyOnDownloadRecovery && Boolean(String(this.settings.notifyUrl || "").trim())
},
(event) => this.notificationOutbox.enqueue(event)
).then(() => undefined).catch((error) => {
logger.warn(`Download-Health-Monitor konnte nicht ausgewertet werden: ${String(error)}`);
}).finally(() => {
if (this.downloadHealthEvaluation === task) {
this.downloadHealthEvaluation = null;
}
});
this.downloadHealthEvaluation = task;
return task;
}
public async shutdown(): Promise<void> {
if (this.runtimeStatsTimer) {
if (this.downloadHealthTimer) {
clearInterval(this.downloadHealthTimer);
this.downloadHealthTimer = null;
}
this.manager.suspendDownloadHealthMonitoring?.();
if (this.downloadHealthEvaluation) {
await this.downloadHealthEvaluation;
}
if (this.runtimeStatsTimer) {
clearInterval(this.runtimeStatsTimer);
this.runtimeStatsTimer = null;
}
+497
View File
@@ -0,0 +1,497 @@
import fs from "node:fs";
import path from "node:path";
import type { NotificationEvent } from "./notification-outbox";
export type DownloadHealthStatus =
| "idle"
| "suspended"
| "expected_wait"
| "healthy"
| "suspect_scheduler"
| "suspect_no_data"
| "alerted"
| "recovering";
export type DownloadHealthIncidentType = "scheduler" | "no_data";
export interface DownloadHealthSnapshot {
runActive: boolean;
runFingerprint: string;
queueFingerprint: string;
openItems: number;
openPackages: number;
knownDownloadedBytes: number;
activeTasks: number;
startableItems: number;
lastSchedulerTickAt: number;
downloadProgressSequence: number;
itemCompletionSequence: number;
lastPositiveByteAt: number;
technicalRecoveryCount: number;
paused: boolean;
reconnectUntil: number;
nextRetryAt: number;
providerCooldownUntil: number;
blockedOnDisk: boolean;
blockedOnThrottleUntil: number;
activePhaseDeadlineAt: number;
terminalFailure: boolean;
manualStop: boolean;
shuttingDown: boolean;
currentSpeedBps: number;
}
export interface DownloadHealthState {
version: 1;
status: DownloadHealthStatus;
runFingerprint: string | null;
queueFingerprint: string | null;
suspiciousDurationMs: number;
suspiciousSamples: number;
incidentStartedAt: number;
incidentType: DownloadHealthIncidentType | null;
alertedAt: number;
lastAlertAt: number;
cooldownUntil: number;
recoverySamples: number;
lastSampleAt: number | null;
downloadProgressSequence: number;
itemCompletionSequence: number;
lastPositiveByteAt: number;
technicalRecoveryCount: number;
restartPending: boolean;
restartFreshSamples: number;
}
export interface DownloadHealthOptions {
stallAfterMs: number;
cooldownMs: number;
notifyOnStall: boolean;
notifyOnRecovery: boolean;
}
export interface DownloadHealthEvaluation {
state: DownloadHealthState;
events: NotificationEvent[];
}
const HEALTH_STATUSES = new Set<DownloadHealthStatus>([
"idle",
"suspended",
"expected_wait",
"healthy",
"suspect_scheduler",
"suspect_no_data",
"alerted",
"recovering"
]);
const INCIDENT_TYPES = new Set<DownloadHealthIncidentType>(["scheduler", "no_data"]);
const FINGERPRINT_PATTERN = /^[a-f0-9]{64}$/;
const MIN_SUSPICIOUS_SAMPLES = 3;
const SCHEDULER_STALE_AFTER_MS = 30_000;
const ERROR_EVENT_TTL_MS = 24 * 60 * 60 * 1000;
const SUCCESS_EVENT_TTL_MS = 6 * 60 * 60 * 1000;
function finiteInteger(value: unknown, fallback = 0): number {
const numeric = Number(value);
return Number.isFinite(numeric) ? Math.max(0, Math.floor(numeric)) : fallback;
}
function validFingerprint(value: unknown): string | null {
return typeof value === "string" && FINGERPRINT_PATTERN.test(value) ? value : null;
}
function durationText(durationMs: number): string {
const totalSeconds = Math.max(0, Math.floor(durationMs / 1000));
if (totalSeconds < 60) {
return `${totalSeconds} s`;
}
const minutes = Math.floor(totalSeconds / 60);
const seconds = totalSeconds % 60;
return seconds > 0 ? `${minutes} min ${seconds} s` : `${minutes} min`;
}
function byteText(bytes: number): string {
const value = Math.max(0, finiteInteger(bytes));
if (value < 1024) {
return `${value} B`;
}
const units = ["KB", "MB", "GB", "TB"];
let amount = value / 1024;
let unit = units[0];
for (let index = 1; index < units.length && amount >= 1024; index += 1) {
amount /= 1024;
unit = units[index];
}
return `${amount >= 10 ? amount.toFixed(0) : amount.toFixed(1)} ${unit}`;
}
function incidentEvent(
state: DownloadHealthState,
snapshot: DownloadHealthSnapshot,
now: number
): NotificationEvent {
const durationMs = Math.max(state.suspiciousDurationMs, now - state.incidentStartedAt);
return {
id: `health:stall:${snapshot.runFingerprint.slice(0, 16)}:${state.incidentStartedAt}`,
type: "download_stalled",
priority: "error",
createdAt: now,
expiresAt: now + ERROR_EVENT_TTL_MS,
attempts: 0,
nextAttemptAt: now,
payload: {
title: "Downloadstillstand bestätigt",
description: `Seit ${durationText(durationMs)} wurde kein bestätigter Downloadfortschritt erkannt.`,
color: 0xe67e22,
fields: [
{ name: "Dauer", value: durationText(durationMs), inline: true },
{ name: "Offene Pakete", value: String(snapshot.openPackages), inline: true },
{ name: "Offene Dateien", value: String(snapshot.openItems), inline: true },
{ name: "Bekannte Mindestmenge", value: byteText(snapshot.knownDownloadedBytes), inline: true },
{ name: "Aktiv / startfähig", value: `${snapshot.activeTasks} / ${snapshot.startableItems}`, inline: true },
{ name: "Technische Wiederherstellungen", value: String(snapshot.technicalRecoveryCount), inline: true }
]
}
};
}
function recoveryEvent(
state: DownloadHealthState,
snapshot: DownloadHealthSnapshot,
now: number
): NotificationEvent {
return {
id: `health:recovery:${snapshot.runFingerprint.slice(0, 16)}:${state.alertedAt}`,
type: "download_recovered",
priority: "success",
createdAt: now,
expiresAt: now + SUCCESS_EVENT_TTL_MS,
attempts: 0,
nextAttemptAt: now,
payload: {
title: "Download läuft wieder",
description: "Nach dem bestätigten Stillstand wurde neuer Fortschritt erkannt.",
color: 0x2ecc71,
fields: [
{ name: "Incident-Dauer", value: durationText(Math.max(0, now - state.incidentStartedAt)), inline: true },
{ name: "Aktive Downloads", value: String(snapshot.activeTasks), inline: true },
{ name: "Aktuelle Geschwindigkeit", value: `${byteText(snapshot.currentSpeedBps)}/s`, inline: true }
]
}
};
}
export function createDownloadHealthState(overrides: Partial<DownloadHealthState> = {}): DownloadHealthState {
return {
version: 1,
status: "idle",
runFingerprint: null,
queueFingerprint: null,
suspiciousDurationMs: 0,
suspiciousSamples: 0,
incidentStartedAt: 0,
incidentType: null,
alertedAt: 0,
lastAlertAt: 0,
cooldownUntil: 0,
recoverySamples: 0,
lastSampleAt: null,
downloadProgressSequence: 0,
itemCompletionSequence: 0,
lastPositiveByteAt: 0,
technicalRecoveryCount: 0,
restartPending: false,
restartFreshSamples: 0,
...overrides
};
}
function resetForSnapshot(
state: DownloadHealthState,
snapshot: DownloadHealthSnapshot,
now: number
): DownloadHealthState {
return createDownloadHealthState({
runFingerprint: snapshot.runFingerprint,
queueFingerprint: snapshot.queueFingerprint,
lastAlertAt: state.lastAlertAt,
cooldownUntil: state.cooldownUntil,
lastSampleAt: now,
downloadProgressSequence: finiteInteger(snapshot.downloadProgressSequence),
itemCompletionSequence: finiteInteger(snapshot.itemCompletionSequence),
lastPositiveByteAt: finiteInteger(snapshot.lastPositiveByteAt),
technicalRecoveryCount: finiteInteger(snapshot.technicalRecoveryCount)
});
}
function endIncident(state: DownloadHealthState): DownloadHealthState {
return createDownloadHealthState({
lastAlertAt: state.lastAlertAt,
cooldownUntil: state.cooldownUntil,
downloadProgressSequence: state.downloadProgressSequence,
itemCompletionSequence: state.itemCompletionSequence,
lastPositiveByteAt: state.lastPositiveByteAt,
technicalRecoveryCount: state.technicalRecoveryCount
});
}
function expectedWait(snapshot: DownloadHealthSnapshot, now: number): boolean {
return snapshot.paused
|| snapshot.reconnectUntil > now
|| snapshot.nextRetryAt > now
|| snapshot.providerCooldownUntil > now
|| snapshot.blockedOnDisk
|| snapshot.blockedOnThrottleUntil > now
|| snapshot.activePhaseDeadlineAt > now;
}
function suspicionType(snapshot: DownloadHealthSnapshot, now: number): DownloadHealthIncidentType | null {
if (snapshot.activeTasks === 0 && snapshot.startableItems > 0) {
if (snapshot.lastSchedulerTickAt <= 0 || now - snapshot.lastSchedulerTickAt >= SCHEDULER_STALE_AFTER_MS) {
return "scheduler";
}
return null;
}
return "no_data";
}
export function evaluateDownloadHealth(
previous: DownloadHealthState,
snapshot: DownloadHealthSnapshot,
nowValue: number,
options: DownloadHealthOptions
): DownloadHealthEvaluation {
const now = finiteInteger(nowValue);
let state = createDownloadHealthState(previous);
const events: NotificationEvent[] = [];
const runFingerprint = validFingerprint(snapshot.runFingerprint);
const queueFingerprint = validFingerprint(snapshot.queueFingerprint);
if (state.restartPending && !snapshot.runActive && !snapshot.terminalFailure && !snapshot.manualStop && !snapshot.shuttingDown) {
state.status = "suspended";
state.lastSampleAt = null;
return { state, events };
}
if (!snapshot.runActive || snapshot.openItems <= 0 || snapshot.terminalFailure || snapshot.manualStop || snapshot.shuttingDown || !runFingerprint || !queueFingerprint) {
state.downloadProgressSequence = Math.max(state.downloadProgressSequence, finiteInteger(snapshot.downloadProgressSequence));
state.itemCompletionSequence = Math.max(state.itemCompletionSequence, finiteInteger(snapshot.itemCompletionSequence));
state.lastPositiveByteAt = Math.max(state.lastPositiveByteAt, finiteInteger(snapshot.lastPositiveByteAt));
state.technicalRecoveryCount = Math.max(state.technicalRecoveryCount, finiteInteger(snapshot.technicalRecoveryCount));
return { state: endIncident(state), events };
}
if (state.runFingerprint !== runFingerprint || state.queueFingerprint !== queueFingerprint) {
state = resetForSnapshot(state, snapshot, now);
}
state.runFingerprint = runFingerprint;
state.queueFingerprint = queueFingerprint;
if (state.restartPending) {
state.restartFreshSamples += 1;
state.lastSampleAt = now;
if (state.restartFreshSamples < 2) {
state.downloadProgressSequence = finiteInteger(snapshot.downloadProgressSequence);
state.itemCompletionSequence = finiteInteger(snapshot.itemCompletionSequence);
state.lastPositiveByteAt = finiteInteger(snapshot.lastPositiveByteAt);
state.technicalRecoveryCount = finiteInteger(snapshot.technicalRecoveryCount);
return { state, events };
}
state.restartPending = false;
}
const progressSequence = finiteInteger(snapshot.downloadProgressSequence);
const completionSequence = finiteInteger(snapshot.itemCompletionSequence);
const positiveByteProgress = progressSequence > state.downloadProgressSequence;
const completionProgress = completionSequence > state.itemCompletionSequence;
const progress = positiveByteProgress || completionProgress;
state.downloadProgressSequence = Math.max(state.downloadProgressSequence, progressSequence);
state.itemCompletionSequence = Math.max(state.itemCompletionSequence, completionSequence);
state.lastPositiveByteAt = Math.max(state.lastPositiveByteAt, finiteInteger(snapshot.lastPositiveByteAt));
state.technicalRecoveryCount = Math.max(state.technicalRecoveryCount, finiteInteger(snapshot.technicalRecoveryCount));
if (state.alertedAt > 0 && progress) {
state.recoverySamples = completionProgress ? 2 : state.recoverySamples + 1;
if (state.recoverySamples >= 2) {
if (options.notifyOnRecovery) {
events.push(recoveryEvent(state, snapshot, now));
}
const recovered = resetForSnapshot(state, snapshot, now);
recovered.status = "healthy";
recovered.lastAlertAt = state.lastAlertAt;
recovered.cooldownUntil = state.cooldownUntil;
return { state: recovered, events };
}
state.status = "recovering";
state.lastSampleAt = now;
return { state, events };
}
if (state.alertedAt > 0) {
if (expectedWait(snapshot, now)) {
state.status = "expected_wait";
} else {
state.status = "alerted";
}
state.lastSampleAt = now;
return { state, events };
}
if (expectedWait(snapshot, now)) {
state.status = "expected_wait";
state.lastSampleAt = now;
return { state, events };
}
if (progress) {
state.status = "healthy";
state.suspiciousDurationMs = 0;
state.suspiciousSamples = 0;
state.incidentStartedAt = 0;
state.incidentType = null;
state.recoverySamples = 0;
state.lastSampleAt = now;
return { state, events };
}
const incidentType = suspicionType(snapshot, now);
if (!incidentType) {
state.status = "healthy";
state.suspiciousDurationMs = 0;
state.suspiciousSamples = 0;
state.incidentStartedAt = 0;
state.incidentType = null;
state.lastSampleAt = now;
return { state, events };
}
const elapsed = state.lastSampleAt === null ? 0 : Math.max(0, now - state.lastSampleAt);
if (state.incidentStartedAt <= 0) {
state.incidentStartedAt = now;
}
state.suspiciousDurationMs += elapsed;
state.suspiciousSamples += 1;
state.incidentType = incidentType;
state.status = incidentType === "scheduler" ? "suspect_scheduler" : "suspect_no_data";
state.lastSampleAt = now;
const stallAfterMs = Math.max(0, finiteInteger(options.stallAfterMs, 90_000));
const cooldownMs = Math.max(0, finiteInteger(options.cooldownMs, 600_000));
const confirmed = state.suspiciousDurationMs >= stallAfterMs
&& state.suspiciousSamples >= MIN_SUSPICIOUS_SAMPLES;
if (confirmed && options.notifyOnStall && now >= state.cooldownUntil) {
state.status = "alerted";
state.alertedAt = now;
state.lastAlertAt = now;
state.cooldownUntil = now + cooldownMs;
state.recoverySamples = 0;
events.push(incidentEvent(state, snapshot, now));
}
return { state, events };
}
function persistedState(state: DownloadHealthState): Omit<DownloadHealthState, "restartPending" | "restartFreshSamples"> {
return {
version: 1,
status: state.status,
runFingerprint: state.runFingerprint,
queueFingerprint: state.queueFingerprint,
suspiciousDurationMs: finiteInteger(state.suspiciousDurationMs),
suspiciousSamples: finiteInteger(state.suspiciousSamples),
incidentStartedAt: finiteInteger(state.incidentStartedAt),
incidentType: state.incidentType,
alertedAt: finiteInteger(state.alertedAt),
lastAlertAt: finiteInteger(state.lastAlertAt),
cooldownUntil: finiteInteger(state.cooldownUntil),
recoverySamples: finiteInteger(state.recoverySamples),
lastSampleAt: state.lastSampleAt === null ? null : finiteInteger(state.lastSampleAt),
downloadProgressSequence: finiteInteger(state.downloadProgressSequence),
itemCompletionSequence: finiteInteger(state.itemCompletionSequence),
lastPositiveByteAt: finiteInteger(state.lastPositiveByteAt),
technicalRecoveryCount: finiteInteger(state.technicalRecoveryCount)
};
}
export function saveDownloadHealthState(filePath: string, state: DownloadHealthState): void {
fs.mkdirSync(path.dirname(filePath), { recursive: true });
const tempPath = `${filePath}.tmp`;
try {
fs.writeFileSync(tempPath, JSON.stringify(persistedState(state)), "utf8");
fs.renameSync(tempPath, filePath);
} catch (error) {
try {
fs.rmSync(tempPath, { force: true });
} catch {
}
throw error;
}
}
export function loadDownloadHealthState(filePath: string): DownloadHealthState {
try {
const raw = JSON.parse(fs.readFileSync(filePath, "utf8")) as Record<string, unknown>;
const runFingerprint = raw.runFingerprint === null ? null : validFingerprint(raw.runFingerprint);
const queueFingerprint = raw.queueFingerprint === null ? null : validFingerprint(raw.queueFingerprint);
if (raw.version !== 1 || !HEALTH_STATUSES.has(raw.status as DownloadHealthStatus)) {
return createDownloadHealthState();
}
if ((raw.runFingerprint !== null && !runFingerprint) || (raw.queueFingerprint !== null && !queueFingerprint)) {
return createDownloadHealthState();
}
const incidentType = INCIDENT_TYPES.has(raw.incidentType as DownloadHealthIncidentType)
? raw.incidentType as DownloadHealthIncidentType
: null;
return createDownloadHealthState({
status: raw.status as DownloadHealthStatus,
runFingerprint,
queueFingerprint,
suspiciousDurationMs: finiteInteger(raw.suspiciousDurationMs),
suspiciousSamples: finiteInteger(raw.suspiciousSamples),
incidentStartedAt: finiteInteger(raw.incidentStartedAt),
incidentType,
alertedAt: finiteInteger(raw.alertedAt),
lastAlertAt: finiteInteger(raw.lastAlertAt),
cooldownUntil: finiteInteger(raw.cooldownUntil),
recoverySamples: finiteInteger(raw.recoverySamples),
lastSampleAt: null,
downloadProgressSequence: finiteInteger(raw.downloadProgressSequence),
itemCompletionSequence: finiteInteger(raw.itemCompletionSequence),
lastPositiveByteAt: finiteInteger(raw.lastPositiveByteAt),
technicalRecoveryCount: finiteInteger(raw.technicalRecoveryCount),
restartPending: Boolean(runFingerprint && queueFingerprint),
restartFreshSamples: 0
});
} catch {
return createDownloadHealthState();
}
}
export class DownloadHealthMonitor {
private state: DownloadHealthState;
public constructor(private readonly filePath: string, initialState?: DownloadHealthState) {
this.state = initialState ? createDownloadHealthState(initialState) : loadDownloadHealthState(filePath);
}
public getState(): DownloadHealthState {
return createDownloadHealthState(this.state);
}
public async sample(
snapshot: DownloadHealthSnapshot,
now: number,
options: DownloadHealthOptions,
enqueue: (event: NotificationEvent) => Promise<void>
): Promise<DownloadHealthEvaluation> {
const evaluation = evaluateDownloadHealth(this.state, snapshot, now, options);
for (const event of evaluation.events) {
await enqueue(event);
}
saveDownloadHealthState(this.filePath, evaluation.state);
this.state = evaluation.state;
return evaluation;
}
}
+281 -99
View File
@@ -1,6 +1,7 @@
import fs from "node:fs";
import path from "node:path";
import os from "node:os";
import fs from "node:fs";
import path from "node:path";
import os from "node:os";
import { createHash } from "node:crypto";
import { EventEmitter } from "node:events";
import { v4 as uuidv4 } from "uuid";
import {
@@ -93,6 +94,7 @@ import {
} from "./statistics-ledger";
import { finalizePackageResult } from "./package-telemetry";
import type { NotificationEvent } from "./notification-outbox";
import type { DownloadHealthSnapshot } from "./download-health-monitor";
import {
buildHistoryEntry,
buildPackageDigestEvents,
@@ -118,8 +120,12 @@ type ActiveTask = {
stallRetries?: number;
genericErrorRetries?: number;
unrestrictRetries?: number;
blockedOnDiskWrite?: boolean;
blockedOnDiskSince?: number;
blockedOnDiskWrite?: boolean;
blockedOnDiskSince?: number;
blockedOnThrottleUntil?: number;
phase?: "validating" | "downloading" | "integrity_check";
phaseStartedAt?: number;
phaseDeadlineAt?: number;
};
const DOWNLOAD_ACCOUNT_PROVIDERS: readonly DebridProvider[] = [
@@ -1960,7 +1966,23 @@ export class DownloadManager extends EventEmitter {
private itemCount = 0;
private lastSchedulerHeartbeatAt = 0;
private lastSchedulerHeartbeatAt = 0;
private lastSchedulerTickAt = 0;
private downloadProgressSequence = 0;
private itemCompletionSequence = 0;
private lastPositiveByteAt = 0;
private technicalRecoveryCount = 0;
private healthManualStop = false;
private healthShuttingDown = false;
private healthTerminalFailure = false;
private lastReconnectMarkAt = 0;
@@ -2478,11 +2500,118 @@ export class DownloadManager extends EventEmitter {
return cloneSession(this.session);
}
public getSummary(): DownloadSummary | null {
return this.summary;
}
public isSessionRunning(): boolean {
public getSummary(): DownloadSummary | null {
return this.summary;
}
public getDownloadHealthSnapshot(now = nowMs()): DownloadHealthSnapshot {
const openPackages = new Set<string>();
const queueParts: string[] = [];
const activeTasks: ActiveTask[] = [];
let openItems = 0;
let knownDownloadedBytes = 0;
let currentSpeedBps = 0;
let startableItems = 0;
let nextRetryAt = 0;
let providerCooldownUntil = 0;
for (const itemId of this.runItemIds) {
queueParts.push(itemId);
const item = this.session.items[itemId];
if (!item || isFinishedStatus(item.status)) {
continue;
}
const pkg = this.session.packages[item.packageId];
if (!pkg || pkg.cancelled || !pkg.enabled) {
continue;
}
openItems += 1;
openPackages.add(pkg.id);
knownDownloadedBytes += Math.max(0, Math.floor(Number(item.downloadedBytes) || 0));
currentSpeedBps += Math.max(0, Math.floor(Number(item.speedBps) || 0));
const active = this.activeTasks.get(itemId);
if (active && !active.abortController.signal.aborted) {
activeTasks.push(active);
continue;
}
if (item.status !== "queued" && item.status !== "reconnect_wait") {
continue;
}
const retryAt = this.retryAfterByItem.get(itemId) || 0;
const failureKey = this.getProviderFailureKeyForItem(item);
const cooldownAt = this.providerFailures.get(failureKey)?.cooldownUntil || 0;
if (retryAt > now) {
nextRetryAt = nextRetryAt === 0 ? retryAt : Math.min(nextRetryAt, retryAt);
}
if (cooldownAt > now) {
providerCooldownUntil = providerCooldownUntil === 0
? cooldownAt
: Math.min(providerCooldownUntil, cooldownAt);
}
if (retryAt <= now && cooldownAt <= now) {
startableItems += 1;
}
}
const blockedOnDisk = activeTasks.length > 0 && activeTasks.every((active) => Boolean(active.blockedOnDiskWrite));
const throttleDeadlines = activeTasks.map((active) => active.blockedOnThrottleUntil || 0);
const blockedOnThrottleUntil = throttleDeadlines.length > 0 && throttleDeadlines.every((deadline) => deadline > now)
? Math.min(...throttleDeadlines)
: 0;
const phaseDeadlines = activeTasks.map((active) => active.phaseDeadlineAt || 0);
const activePhaseDeadlineAt = phaseDeadlines.length > 0 && phaseDeadlines.every((deadline) => deadline > now)
? Math.min(...phaseDeadlines)
: 0;
const queueFingerprint = createHash("sha256")
.update([...queueParts].sort().join("\n"))
.digest("hex");
const runParts = [...this.runPackageIds].map((packageId) => {
const generation = Math.max(1, Math.floor(Number(this.session.packages[packageId]?.resultGeneration) || 1));
return `${packageId}:${generation}`;
});
const runFingerprint = createHash("sha256")
.update(`${runParts.sort().join("\n")}|${queueFingerprint}`)
.digest("hex");
return {
runActive: this.session.running && openItems > 0,
runFingerprint,
queueFingerprint,
openItems,
openPackages: openPackages.size,
knownDownloadedBytes,
activeTasks: activeTasks.length,
startableItems,
lastSchedulerTickAt: this.lastSchedulerTickAt,
downloadProgressSequence: this.downloadProgressSequence,
itemCompletionSequence: this.itemCompletionSequence,
lastPositiveByteAt: this.lastPositiveByteAt,
technicalRecoveryCount: this.technicalRecoveryCount,
paused: this.session.paused,
reconnectUntil: this.session.reconnectUntil,
nextRetryAt: activeTasks.length === 0 && startableItems === 0 ? nextRetryAt : 0,
providerCooldownUntil: activeTasks.length === 0 && startableItems === 0 ? providerCooldownUntil : 0,
blockedOnDisk,
blockedOnThrottleUntil,
activePhaseDeadlineAt,
terminalFailure: this.healthTerminalFailure,
manualStop: this.healthManualStop,
shuttingDown: this.healthShuttingDown,
currentSpeedBps: this.session.running && !this.session.paused ? currentSpeedBps : 0
};
}
private beginHealthRun(): void {
this.healthManualStop = false;
this.healthShuttingDown = false;
this.healthTerminalFailure = false;
}
public suspendDownloadHealthMonitoring(): void {
this.healthShuttingDown = true;
}
public isSessionRunning(): boolean {
return this.session.running;
}
@@ -5978,6 +6107,7 @@ export class DownloadManager extends EventEmitter {
}
public async startPackages(packageIds: string[]): Promise<void> {
this.beginHealthRun();
this.ensureUsableDownloadAccount();
const targetSet = new Set(packageIds);
for (const packageId of this.packagePostProcessTasks.keys()) {
@@ -6080,6 +6210,7 @@ export class DownloadManager extends EventEmitter {
}
public async startItems(itemIds: string[]): Promise<void> {
this.beginHealthRun();
this.ensureUsableDownloadAccount();
const targetSet = new Set(itemIds);
@@ -6196,6 +6327,7 @@ export class DownloadManager extends EventEmitter {
if (this.session.running) {
return;
}
this.beginHealthRun();
this.ensureUsableDownloadAccount();
this.schedulerGeneration += 1;
@@ -6347,8 +6479,10 @@ export class DownloadManager extends EventEmitter {
});
}
public stop(options?: { parkForRestart?: boolean }): void {
const parkForRestart = options?.parkForRestart === true;
public stop(options?: { parkForRestart?: boolean }): void {
const parkForRestart = options?.parkForRestart === true;
this.healthManualStop = !parkForRestart;
this.healthShuttingDown = parkForRestart;
const abortReason: "stop" | "shutdown" = parkForRestart ? "shutdown" : "stop";
const keepExtraction = this.settings.autoExtractWhenStopped;
const wasRunning = this.session.running;
@@ -6426,6 +6560,7 @@ export class DownloadManager extends EventEmitter {
}
public prepareForShutdown(): void {
this.healthShuttingDown = true;
logger.info(`Shutdown-Vorbereitung gestartet: active=${this.activeTasks.size}, running=${this.session.running}, paused=${this.session.paused}`);
this.updateStatisticsActivity(nowMs());
this.rotationListenerActive = false;
@@ -6922,9 +7057,14 @@ export class DownloadManager extends EventEmitter {
private lastSpeedPruneAt = 0;
private recordSpeed(bytes: number, packageId: string = ""): void {
const now = nowMs();
if (bytes > 0 && this.consecutiveReconnects > 0) {
private recordSpeed(bytes: number, packageId: string = ""): void {
const now = nowMs();
if (!Number.isFinite(bytes) || bytes <= 0) {
return;
}
this.downloadProgressSequence += 1;
this.lastPositiveByteAt = now;
if (bytes > 0 && this.consecutiveReconnects > 0) {
this.consecutiveReconnects = 0;
}
const bucket = now - (now % 120);
@@ -6955,8 +7095,9 @@ export class DownloadManager extends EventEmitter {
this.statisticsDirty = true;
this.statisticsUrgent = true;
}
if (status === "completed" && previous !== "completed") {
this.sessionCompletedFiles += 1;
if (status === "completed" && previous !== "completed") {
this.itemCompletionSequence += 1;
this.sessionCompletedFiles += 1;
this.settings.totalCompletedFilesAllTime = Math.max(0, Number(this.settings.totalCompletedFilesAllTime || 0)) + 1;
this.invalidateStatsCache();
}
@@ -9108,6 +9249,7 @@ export class DownloadManager extends EventEmitter {
try {
while (this.session.running && this.schedulerGeneration === myGeneration) {
const now = nowMs();
this.lastSchedulerTickAt = now;
this.updateStatisticsActivity(now);
if (now - this.lastSchedulerHeartbeatAt >= 60000) {
this.lastSchedulerHeartbeatAt = now;
@@ -9246,8 +9388,9 @@ export class DownloadManager extends EventEmitter {
return;
}
logger.warn(`Globaler Download-Stall erkannt (${Math.floor((now - this.lastGlobalProgressAt) / 1000)}s ohne Fortschritt), ${stalledCount} Task(s) neu starten, diskBlocked=${diskBlockedCount}`);
for (const active of this.activeTasks.values()) {
logger.warn(`Globaler Download-Stall erkannt (${Math.floor((now - this.lastGlobalProgressAt) / 1000)}s ohne Fortschritt), ${stalledCount} Task(s) neu starten, diskBlocked=${diskBlockedCount}`);
this.technicalRecoveryCount += 1;
for (const active of this.activeTasks.values()) {
if (active.abortController.signal.aborted) {
continue;
}
@@ -9576,10 +9719,14 @@ export class DownloadManager extends EventEmitter {
abortController: new AbortController(),
abortReason: "none",
resumable: true,
nonResumableCounted: false,
blockedOnDiskWrite: false,
blockedOnDiskSince: 0
};
nonResumableCounted: false,
blockedOnDiskWrite: false,
blockedOnDiskSince: 0,
blockedOnThrottleUntil: 0,
phase: "validating",
phaseStartedAt: item.updatedAt,
phaseDeadlineAt: item.updatedAt + getUnrestrictTimeoutMs() + 15_000
};
this.activeTasks.set(itemId, active);
this.notePacedStartForItem(item, nowMs());
this.emitState();
@@ -9701,6 +9848,9 @@ export class DownloadManager extends EventEmitter {
this.settings,
item.url
);
active.phase = "validating";
active.phaseStartedAt = nowMs();
active.phaseDeadlineAt = active.phaseStartedAt + unrestrictTimeoutMs;
const unrestrictTimeoutSignal = AbortSignal.timeout(unrestrictTimeoutMs);
const unrestrictedSignal = AbortSignal.any([active.abortController.signal, unrestrictTimeoutSignal]);
let unrestricted;
@@ -9805,6 +9955,9 @@ export class DownloadManager extends EventEmitter {
throw error;
}
item.status = "downloading";
active.phase = "downloading";
active.phaseStartedAt = nowMs();
active.phaseDeadlineAt = 0;
const pLabel = unrestricted.providerLabel;
item.fullStatus = "Starte...";
item.updatedAt = nowMs();
@@ -9855,10 +10008,18 @@ export class DownloadManager extends EventEmitter {
item.status = "integrity_check";
item.fullStatus = "CRC-Check läuft";
item.updatedAt = nowMs();
this.emitState();
const integrityStartedAt = nowMs();
const validation = await validateFileAgainstManifest(item.targetPath, pkg.outputDir);
this.emitState();
const integrityStartedAt = nowMs();
const integrityBytes = Math.max(0, Number(item.downloadedBytes || item.totalBytes || 0));
const integrityBudgetMs = Math.max(
120_000,
Math.min(30 * 60 * 1000, 60_000 + Math.ceil(integrityBytes / (25 * 1024 * 1024)) * 1000)
);
active.phase = "integrity_check";
active.phaseStartedAt = integrityStartedAt;
active.phaseDeadlineAt = integrityStartedAt + integrityBudgetMs;
const validation = await validateFileAgainstManifest(item.targetPath, pkg.outputDir);
if (active.abortController.signal.aborted) {
throw new Error(`aborted:${active.abortReason}`);
}
@@ -9881,9 +10042,12 @@ export class DownloadManager extends EventEmitter {
item.progressPercent = 0;
item.downloadedBytes = 0;
item.totalBytes = mergeKnownTotalBytes(item.totalBytes, unrestricted.fileSize);
this.emitState();
await sleep(300);
continue;
this.emitState();
await sleep(300);
active.phase = "downloading";
active.phaseStartedAt = nowMs();
active.phaseDeadlineAt = 0;
continue;
}
throw new Error(`Integritätsprüfung fehlgeschlagen (${validation.message})`);
}
@@ -11178,7 +11342,7 @@ export class DownloadManager extends EventEmitter {
}
const buffer = Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk.buffer, chunk.byteOffset, chunk.byteLength);
await this.applySpeedLimit(buffer.length, windowBytes, windowStarted, active.abortController.signal);
await this.applySpeedLimit(buffer.length, windowBytes, windowStarted, active);
if (active.abortController.signal.aborted) {
throw new Error(`aborted:${active.abortReason}`);
}
@@ -12241,45 +12405,53 @@ export class DownloadManager extends EventEmitter {
return 0;
}
private async applyGlobalSpeedLimit(chunkBytes: number, bytesPerSecond: number, signal?: AbortSignal): Promise<void> {
const task = this.globalSpeedLimitQueue
private async applyGlobalSpeedLimit(chunkBytes: number, bytesPerSecond: number, active?: ActiveTask): Promise<void> {
const signal = active?.abortController.signal;
const task = this.globalSpeedLimitQueue
.catch(() => undefined)
.then(async () => {
if (signal?.aborted) {
throw new Error("aborted:speed_limit");
}
const now = nowMs();
const waitMs = Math.max(0, this.globalSpeedLimitNextAt - now);
if (waitMs > 0) {
await new Promise<void>((resolve, reject) => {
let timer: NodeJS.Timeout | null = setTimeout(() => {
timer = null;
if (signal) {
signal.removeEventListener("abort", onAbort);
}
resolve();
}, waitMs);
const onAbort = (): void => {
if (timer) {
clearTimeout(timer);
timer = null;
}
signal?.removeEventListener("abort", onAbort);
reject(new Error("aborted:speed_limit"));
};
if (signal) {
if (signal.aborted) {
onAbort();
return;
}
signal.addEventListener("abort", onAbort, { once: true });
}
});
}
if (signal?.aborted) {
const now = nowMs();
const waitMs = Math.max(0, this.globalSpeedLimitNextAt - now);
if (waitMs > 0) {
if (active) {
active.blockedOnThrottleUntil = now + waitMs;
}
try {
await new Promise<void>((resolve, reject) => {
let timer: NodeJS.Timeout | null = setTimeout(() => {
timer = null;
signal?.removeEventListener("abort", onAbort);
resolve();
}, waitMs);
const onAbort = (): void => {
if (timer) {
clearTimeout(timer);
timer = null;
}
signal?.removeEventListener("abort", onAbort);
reject(new Error("aborted:speed_limit"));
};
if (signal) {
if (signal.aborted) {
onAbort();
return;
}
signal.addEventListener("abort", onAbort, { once: true });
}
});
} finally {
if (active) {
active.blockedOnThrottleUntil = 0;
}
}
}
if (signal?.aborted) {
throw new Error("aborted:speed_limit");
}
@@ -12292,7 +12464,8 @@ export class DownloadManager extends EventEmitter {
await task;
}
private async applySpeedLimit(chunkBytes: number, localWindowBytes: number, localWindowStarted: number, signal?: AbortSignal): Promise<void> {
private async applySpeedLimit(chunkBytes: number, localWindowBytes: number, localWindowStarted: number, active?: ActiveTask): Promise<void> {
const signal = active?.abortController.signal;
const limitKbps = this.getEffectiveSpeedLimitKbps();
if (limitKbps <= 0) {
return;
@@ -12304,40 +12477,48 @@ export class DownloadManager extends EventEmitter {
const projected = localWindowBytes + chunkBytes;
const allowed = bytesPerSecond * elapsed;
if (projected > allowed) {
const sleepMs = Math.ceil(((projected - allowed) / bytesPerSecond) * 1000);
if (sleepMs > 0) {
await new Promise<void>((resolve, reject) => {
let timer: NodeJS.Timeout | null = setTimeout(() => {
timer = null;
if (signal) {
signal.removeEventListener("abort", onAbort);
}
resolve();
}, Math.min(300, sleepMs));
const onAbort = (): void => {
if (timer) {
clearTimeout(timer);
timer = null;
}
signal?.removeEventListener("abort", onAbort);
reject(new Error("aborted:speed_limit"));
};
if (signal) {
if (signal.aborted) {
onAbort();
return;
}
signal.addEventListener("abort", onAbort, { once: true });
}
});
const sleepMs = Math.ceil(((projected - allowed) / bytesPerSecond) * 1000);
if (sleepMs > 0) {
const boundedSleepMs = Math.min(300, sleepMs);
if (active) {
active.blockedOnThrottleUntil = nowMs() + boundedSleepMs;
}
try {
await new Promise<void>((resolve, reject) => {
let timer: NodeJS.Timeout | null = setTimeout(() => {
timer = null;
signal?.removeEventListener("abort", onAbort);
resolve();
}, boundedSleepMs);
const onAbort = (): void => {
if (timer) {
clearTimeout(timer);
timer = null;
}
signal?.removeEventListener("abort", onAbort);
reject(new Error("aborted:speed_limit"));
};
if (signal) {
if (signal.aborted) {
onAbort();
return;
}
signal.addEventListener("abort", onAbort, { once: true });
}
});
} finally {
if (active) {
active.blockedOnThrottleUntil = 0;
}
}
}
}
return;
}
await this.applyGlobalSpeedLimit(chunkBytes, bytesPerSecond, signal);
await this.applyGlobalSpeedLimit(chunkBytes, bytesPerSecond, active);
}
private async findReadyArchiveSets(pkg: PackageEntry): Promise<Set<string>> {
@@ -13882,9 +14063,10 @@ export class DownloadManager extends EventEmitter {
this.session.runStartedAt = 0;
const total = this.runItemIds.size;
const outcomes = Array.from(this.runOutcomes.values());
const success = outcomes.filter((status) => status === "completed").length;
const failed = outcomes.filter((status) => status === "failed").length;
const cancelled = outcomes.filter((status) => status === "cancelled").length;
const success = outcomes.filter((status) => status === "completed").length;
const failed = outcomes.filter((status) => status === "failed").length;
const cancelled = outcomes.filter((status) => status === "cancelled").length;
this.healthTerminalFailure = failed > 0;
const extracted = this.runCompletedPackages.size;
const duration = runStartedAt > 0 ? Math.max(1, Math.floor((completedAt - runStartedAt) / 1000)) : 1;
const avgSpeed = Math.floor(this.session.totalDownloadedBytes / duration);
+3 -1
View File
@@ -702,6 +702,7 @@ export interface StoragePaths {
historyFile: string;
statisticsFile: string;
notificationOutboxFile: string;
notificationHealthFile: string;
}
export function createStoragePaths(baseDir: string): StoragePaths {
@@ -711,7 +712,8 @@ export function createStoragePaths(baseDir: string): StoragePaths {
sessionFile: path.join(baseDir, "rd_session_state.json"),
historyFile: path.join(baseDir, "rd_history.json"),
statisticsFile: path.join(baseDir, "rd_statistics.json"),
notificationOutboxFile: path.join(baseDir, "rd_notification_outbox.json")
notificationOutboxFile: path.join(baseDir, "rd_notification_outbox.json"),
notificationHealthFile: path.join(baseDir, "rd_notification_health.json")
};
}