fix(nextcloud-files): WR-03 hoechstens vier gleichzeitige Nextcloud-Aufrufe je Zugang
- die Aufrufsperre vergibt je Zugangsschluessel vier Plaetze bis zur Antwort (Kopfzeilen); weitere Aufrufe warten der Reihe nach und pruefen danach, ob der Zugang inzwischen tot oder der Ursprung gesperrt ist; ein widerrufener Zugang erzeugt so hoechstens vier 401 - Hochladen und Zusammenbau (Antwort erst nach dem ganzen Koerper) belegen keinen Platz - Weboberflaeche: Vorschaubilder laden weiter lazy; schlaegt eines fehl, prueft die Ansicht hoechstens alle 10 s die Verbindung und kehrt bei connectionExpired zur Anmeldung zurueck Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -8,6 +8,8 @@ export const MAX_PAUSE_SECONDS = 60 * 60;
|
||||
export const DEAD_KEY_TTL_MS = 24 * 60 * 60 * 1000;
|
||||
/** Hoechstzahl gemerkter Zugangsschluessel (aelteste fliegen zuerst raus). */
|
||||
export const MAX_DEAD_KEYS = 10_000;
|
||||
/** Hoechstens so viele Aufrufe je Zugangsschluessel warten gleichzeitig auf ihre Antwort (WR-03). */
|
||||
export const MAX_CONCURRENT_PER_KEY = 4;
|
||||
|
||||
/**
|
||||
* Aufrufsperre (quick-261008-mzu, D-O) — ein prozessweites Objekt, das jeder
|
||||
@@ -26,6 +28,12 @@ export const MAX_DEAD_KEYS = 10_000;
|
||||
* tot: spaetere Aufrufe gehen gar nicht erst raus, und laufende Aufrufe
|
||||
* desselben Schluessels werden abgebrochen.
|
||||
*
|
||||
* (c) Hoechstens 4 gleichzeitige Aufrufe je Zugangsschluessel (WR-03): Die
|
||||
* Sperre (b) greift erst, wenn das erste 401 zurueck ist. Ohne Grenze waeren
|
||||
* bis dahin z. B. Dutzende Vorschaubilder unterwegs, und jedes davon zaehlte
|
||||
* bei Nextcloud als Fehlanmeldung der gemeinsamen Server-Adresse. Ein Platz
|
||||
* gilt bis zur Antwort (Kopfzeilen), nicht fuer den ganzen Datenstrom.
|
||||
*
|
||||
* Der Zustand liegt im Arbeitsspeicher; ein Neustart der API hebt ihn auf (die
|
||||
* Nextcloud-Sperre selbst bleibt dort bestehen und wird beim naechsten 429
|
||||
* wieder erkannt).
|
||||
@@ -38,6 +46,7 @@ export class NextcloudCallGate {
|
||||
private readonly pausedUntil = new Map<string, number>();
|
||||
private readonly dead = new Map<string, number>();
|
||||
private readonly controllers = new Map<string, AbortController>();
|
||||
private readonly slots = new Map<string, { active: number; waiting: Array<() => void> }>();
|
||||
|
||||
// --- (a) Sperre je Ursprung --------------------------------------------------
|
||||
|
||||
@@ -119,4 +128,66 @@ export class NextcloudCallGate {
|
||||
}
|
||||
return controller.signal;
|
||||
}
|
||||
|
||||
// --- (c) Gleichzeitige Aufrufe je Zugangsschluessel ---------------------------
|
||||
|
||||
/**
|
||||
* Belegt einen der 4 Plaetze des Schluessels; sind alle belegt, wird in der
|
||||
* Reihenfolge der Anfragen gewartet. Liefert die Freigabe (mehrfacher Aufruf ist
|
||||
* wirkungslos). Bricht `signal` waehrend des Wartens ab, wird mit dessen Grund
|
||||
* abgelehnt.
|
||||
*/
|
||||
async acquireSlot(key: string, signal?: AbortSignal): Promise<() => void> {
|
||||
let slot = this.slots.get(key);
|
||||
if (!slot) {
|
||||
slot = { active: 0, waiting: [] };
|
||||
this.slots.set(key, slot);
|
||||
}
|
||||
if (slot.active < MAX_CONCURRENT_PER_KEY) {
|
||||
slot.active += 1;
|
||||
return this.releaser(key, slot);
|
||||
}
|
||||
const queue = slot;
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
const onAbort = () => {
|
||||
const i = queue.waiting.indexOf(grant);
|
||||
if (i >= 0) queue.waiting.splice(i, 1);
|
||||
reject(signal?.reason ?? new Error('aborted'));
|
||||
};
|
||||
const grant = () => {
|
||||
signal?.removeEventListener('abort', onAbort);
|
||||
resolve();
|
||||
};
|
||||
queue.waiting.push(grant);
|
||||
if (signal) {
|
||||
if (signal.aborted) onAbort();
|
||||
else signal.addEventListener('abort', onAbort, { once: true });
|
||||
}
|
||||
});
|
||||
return this.releaser(key, queue);
|
||||
}
|
||||
|
||||
/** Belegte und wartende Aufrufe eines Schluessels (fuer Tests). */
|
||||
slotUsage(key: string): { active: number; waiting: number } {
|
||||
const slot = this.slots.get(key);
|
||||
return { active: slot?.active ?? 0, waiting: slot?.waiting.length ?? 0 };
|
||||
}
|
||||
|
||||
private releaser(key: string, slot: { active: number; waiting: Array<() => void> }) {
|
||||
let released = false;
|
||||
return () => {
|
||||
if (released) return;
|
||||
released = true;
|
||||
const next = slot.waiting.shift();
|
||||
// Der Platz geht direkt an den naechsten Wartenden (die Zahl der belegten bleibt gleich).
|
||||
if (next) {
|
||||
next();
|
||||
return;
|
||||
}
|
||||
slot.active -= 1;
|
||||
if (slot.active <= 0 && slot.waiting.length === 0 && this.slots.get(key) === slot) {
|
||||
this.slots.delete(key);
|
||||
}
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
@@ -75,6 +75,8 @@ export async function putFile(
|
||||
headersTimeoutMs: DAV_CHUNK_HEADERS_TIMEOUT_MS,
|
||||
bodyTimeoutMs: DAV_CHUNK_BODY_TIMEOUT_MS,
|
||||
signal: opts.signal,
|
||||
// Die Antwort kommt erst nach dem ganzen Koerper; Uploads laufen je Datei nacheinander.
|
||||
unthrottled: true,
|
||||
}),
|
||||
);
|
||||
}
|
||||
@@ -131,6 +133,7 @@ export async function uploadChunk(
|
||||
headersTimeoutMs: DAV_CHUNK_HEADERS_TIMEOUT_MS,
|
||||
bodyTimeoutMs: DAV_CHUNK_BODY_TIMEOUT_MS,
|
||||
signal: opts.signal,
|
||||
unthrottled: true,
|
||||
},
|
||||
),
|
||||
);
|
||||
@@ -166,7 +169,8 @@ export async function uploadAssemble(
|
||||
'MOVE',
|
||||
'/remote.php/dav/uploads/',
|
||||
[uploadId, '.file'],
|
||||
{ headers, headersTimeoutMs: DAV_ASSEMBLE_HEADERS_TIMEOUT_MS },
|
||||
// Der Zusammenbau kann Minuten dauern und darf keinen Platz blockieren (WR-03).
|
||||
{ headers, headersTimeoutMs: DAV_ASSEMBLE_HEADERS_TIMEOUT_MS, unthrottled: true },
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
@@ -48,6 +48,8 @@ export interface DavExtra {
|
||||
headersTimeoutMs?: number;
|
||||
bodyTimeoutMs?: number;
|
||||
signal?: AbortSignal;
|
||||
/** Siehe `NcRequestOptions.unthrottled` (nur Hochladen und Zusammenbau). */
|
||||
unthrottled?: boolean;
|
||||
}
|
||||
|
||||
/** Allgemeiner WebDAV-Aufruf im Benutzerbereich; `segments` sind OHNE die Benutzerkennung. */
|
||||
@@ -74,6 +76,7 @@ export function davRequest(
|
||||
headersTimeoutMs: extra.headersTimeoutMs ?? DAV_SMALL_TIMEOUT_MS,
|
||||
bodyTimeoutMs: extra.bodyTimeoutMs ?? DAV_SMALL_TIMEOUT_MS,
|
||||
signal: extra.signal,
|
||||
unthrottled: extra.unthrottled,
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
@@ -367,6 +367,102 @@ describe('ncRequest — Aufrufsperre', () => {
|
||||
});
|
||||
});
|
||||
|
||||
describe('ncRequest — gleichzeitige Aufrufe je Zugangsschluessel (WR-03)', () => {
|
||||
it('zwanzig gleichzeitige Vorschauen mit widerrufenem Zugang: hoechstens vier erreichen Nextcloud', async () => {
|
||||
const gate = new NextcloudCallGate();
|
||||
let release!: () => void;
|
||||
const hold = new Promise<void>((r) => {
|
||||
release = r;
|
||||
});
|
||||
const { transport, calls } = fakeTransport(async () => {
|
||||
await hold;
|
||||
return reply(401);
|
||||
});
|
||||
const all = Promise.all(
|
||||
Array.from({ length: 20 }, () =>
|
||||
ncRequest(transport, gate, {
|
||||
baseUrl: BASE,
|
||||
prefix: '/index.php/core/preview',
|
||||
query: { fileId: '1' },
|
||||
method: 'GET',
|
||||
authorization: 'Basic eA==',
|
||||
credentialKey: 'key-1',
|
||||
}),
|
||||
),
|
||||
);
|
||||
await new Promise((r) => setTimeout(r, 5));
|
||||
expect(calls).toHaveLength(4);
|
||||
expect(gate.slotUsage('key-1')).toEqual({ active: 4, waiting: 16 });
|
||||
release();
|
||||
const results = await all;
|
||||
expect(calls.length).toBeLessThanOrEqual(4);
|
||||
expect(results.every((r) => !r.ok && r.kind === 'credential-dead')).toBe(true);
|
||||
expect(gate.slotUsage('key-1')).toEqual({ active: 0, waiting: 0 });
|
||||
});
|
||||
|
||||
it('nach den Antworten werden Plaetze frei: alle Aufrufe eines lebenden Zugangs kommen durch', async () => {
|
||||
const gate = new NextcloudCallGate();
|
||||
let inFlight = 0;
|
||||
let peak = 0;
|
||||
const { transport, calls } = fakeTransport(async () => {
|
||||
inFlight += 1;
|
||||
peak = Math.max(peak, inFlight);
|
||||
await new Promise((r) => setTimeout(r, 2));
|
||||
inFlight -= 1;
|
||||
return reply(200, 'ok');
|
||||
});
|
||||
const results = await Promise.all(
|
||||
Array.from({ length: 12 }, () =>
|
||||
ncRequest(transport, gate, {
|
||||
baseUrl: BASE,
|
||||
prefix: '/status.php',
|
||||
method: 'GET',
|
||||
credentialKey: 'key-2',
|
||||
}),
|
||||
),
|
||||
);
|
||||
expect(calls).toHaveLength(12);
|
||||
expect(peak).toBeLessThanOrEqual(4);
|
||||
expect(results.every((r) => r.ok)).toBe(true);
|
||||
});
|
||||
|
||||
it('ein Abbruch waehrend des Wartens gibt sofort aborted zurueck, ohne Aufruf', async () => {
|
||||
const gate = new NextcloudCallGate();
|
||||
const holders = await Promise.all([0, 1, 2, 3].map(() => gate.acquireSlot('key-3')));
|
||||
const { transport, calls } = fakeTransport(async () => reply(200));
|
||||
const controller = new AbortController();
|
||||
const pending = ncRequest(transport, gate, {
|
||||
baseUrl: BASE,
|
||||
prefix: '/status.php',
|
||||
method: 'GET',
|
||||
credentialKey: 'key-3',
|
||||
signal: controller.signal,
|
||||
});
|
||||
controller.abort();
|
||||
expect(await pending).toEqual({ ok: false, kind: 'aborted' });
|
||||
expect(calls).toHaveLength(0);
|
||||
expect(gate.slotUsage('key-3')).toEqual({ active: 4, waiting: 0 });
|
||||
for (const free of holders) free();
|
||||
expect(gate.slotUsage('key-3')).toEqual({ active: 0, waiting: 0 });
|
||||
});
|
||||
|
||||
it('Uploads (unthrottled) belegen keinen Platz', async () => {
|
||||
const gate = new NextcloudCallGate();
|
||||
await Promise.all([0, 1, 2, 3].map(() => gate.acquireSlot('key-4')));
|
||||
const { transport, calls } = fakeTransport(async () => reply(201));
|
||||
const res = await ncRequest(transport, gate, {
|
||||
baseUrl: BASE,
|
||||
prefix: '/remote.php/dav/uploads/',
|
||||
segments: ['anna', 'tessera-x', '00001'],
|
||||
method: 'PUT',
|
||||
credentialKey: 'key-4',
|
||||
unthrottled: true,
|
||||
});
|
||||
expect(res.ok).toBe(true);
|
||||
expect(calls).toHaveLength(1);
|
||||
});
|
||||
});
|
||||
|
||||
describe('readCappedText', () => {
|
||||
it('liest bis zur Grenze', async () => {
|
||||
expect(
|
||||
|
||||
@@ -212,6 +212,12 @@ export interface NcRequestOptions {
|
||||
authorization?: string;
|
||||
/** Gesetzt bei Aufrufen mit gespeichertem App-Passwort (siehe Aufrufsperre). */
|
||||
credentialKey?: string;
|
||||
/**
|
||||
* Nicht unter die Grenze gleichzeitiger Aufrufe je Schluessel (WR-03) fallen: nur fuer
|
||||
* Uebertragungen, deren Antwort erst nach dem ganzen Koerper kommt (PUT eines Stuecks,
|
||||
* Zusammenbau). Sie laufen ohnehin nacheinander je Datei.
|
||||
*/
|
||||
unthrottled?: boolean;
|
||||
/** OCS-Aufruf: setzt `OCS-APIRequest` und `Accept: application/json`. */
|
||||
ocs?: boolean;
|
||||
headersTimeoutMs?: number;
|
||||
@@ -319,6 +325,42 @@ export async function ncRequest(
|
||||
return { ok: false, kind: 'paused', retryAfterSeconds: pause.retryAfterSeconds };
|
||||
}
|
||||
|
||||
// WR-03: hoechstens 4 gleichzeitige Aufrufe je Zugangsschluessel bis zur Antwort. Ein
|
||||
// widerrufener Zugang erzeugt so hoechstens ein paar 401, bevor er als tot gilt.
|
||||
let releaseSlot: (() => void) | undefined;
|
||||
if (opts.credentialKey && !opts.unthrottled) {
|
||||
try {
|
||||
releaseSlot = await gate.acquireSlot(opts.credentialKey, opts.signal);
|
||||
} catch {
|
||||
return { ok: false, kind: 'aborted' };
|
||||
}
|
||||
// Waehrend des Wartens kann der Schluessel gestorben oder der Ursprung gesperrt worden sein.
|
||||
if (gate.isDead(opts.credentialKey)) {
|
||||
releaseSlot();
|
||||
return { ok: false, kind: 'credential-dead' };
|
||||
}
|
||||
const later = gate.isPaused(origin);
|
||||
if (later.paused) {
|
||||
releaseSlot();
|
||||
return { ok: false, kind: 'paused', retryAfterSeconds: later.retryAfterSeconds };
|
||||
}
|
||||
}
|
||||
try {
|
||||
return await sendNcRequest(transport, gate, opts, url, origin);
|
||||
} finally {
|
||||
// Erst NACH der Auswertung (401 -> tot, 429 -> Sperre) freigeben: der naechste Wartende
|
||||
// sieht den neuen Zustand, bevor er sendet.
|
||||
releaseSlot?.();
|
||||
}
|
||||
}
|
||||
|
||||
async function sendNcRequest(
|
||||
transport: NextcloudTransport,
|
||||
gate: NextcloudCallGate,
|
||||
opts: NcRequestOptions,
|
||||
url: string,
|
||||
origin: string,
|
||||
): Promise<NcResult> {
|
||||
const headers: Record<string, string> = { 'user-agent': NC_USER_AGENT };
|
||||
if (opts.ocs) {
|
||||
headers['ocs-apirequest'] = 'true';
|
||||
|
||||
Reference in New Issue
Block a user