diff --git a/docs/de/platform/knowledge/documents.md b/docs/de/platform/knowledge/documents.md index 15fbb9a256..ad97420093 100644 --- a/docs/de/platform/knowledge/documents.md +++ b/docs/de/platform/knowledge/documents.md @@ -79,7 +79,7 @@ Das Löschen eines Ordners löscht jede Datei und jeden Unterordner darin endgü **Neu indexieren** (Zeilenmenü) lässt die Pipeline erneut über die gespeicherte Datei laufen — der richtige Zug nach einem Indexierungsfehler oder wenn ein Dokument **Neuindexierung nötig** zeigt. **Löschen** entfernt das Dokument und seine indexierten Chunks; die Bestätigung sagt es unumwunden — die Aktion lässt sich nicht rückgängig machen. Dieselbe Datei erneut hochzuladen bringt den Inhalt als frisches Dokument zurück. Ein gelenktes Dokument lässt sich nicht mehr löschen, sobald irgendeine Version freigegeben wurde — im Review, freigegeben oder mit offenem nächsten Entwurf zeigt der Menüeintrag stattdessen **Geschütztes gelenktes Dokument**, und ein Ordner mit so einem Datensatz verweigert das Ordner-Löschen genauso. Der freigegebene Stand ist ein aufbewahrtes Dokument; genau dafür gibt es den Lebenszyklus. -Jedes Dokument zeigt einen Status: **In Warteschlange** (wartet — eine ausgelastete Organisation indexiert einige Dateien gleichzeitig, der Rest reiht sich ein), **Wird indexiert**, **Indexiert**, **Fehlgeschlagen** oder **Nicht unterstützt** (ein Altformat wie `.doc`/`.ppt`/`.xls`, das sich problemlos speichern und herunterladen lässt, aber keinen Text-Extraktor hat und daher nie für die Suche indexiert wird). Ein durch ein Zeitlimit oder einen Backend-Neustart unterbrochener Indexierungsvorgang erholt sich innerhalb weniger Minuten von selbst — er wird wiederholt oder als **Fehlgeschlagen** mit Wiederholen-Option markiert, nie steckengelassen. Wenn deine Organisation ein Speicher-Kontingent pro Nutzer durchsetzt, zählen fehlgeschlagene und nicht unterstützte Dateien weiterhin dagegen, bis sie gelöscht werden — Platz schaffen heißt also, nicht mehr benötigte Dateien zu entfernen. +Jedes Dokument zeigt einen Status: **In Warteschlange** (wartet — eine ausgelastete Organisation indexiert einige Dateien gleichzeitig, der Rest reiht sich ein), **Wird indexiert**, **Indexiert**, **Fehlgeschlagen** oder **Nicht unterstützt** (ein Altformat wie `.doc`/`.ppt`/`.xls` oder ein Bild wie `.png`/`.jpg` — lässt sich problemlos speichern und herunterladen, hat aber keinen Text-Extraktor und wird daher nie für die Suche indexiert). Ein durch ein Zeitlimit oder einen Backend-Neustart unterbrochener Indexierungsvorgang erholt sich innerhalb weniger Minuten von selbst — er wird wiederholt oder als **Fehlgeschlagen** mit Wiederholen-Option markiert, nie steckengelassen. Wenn deine Organisation ein Speicher-Kontingent pro Nutzer durchsetzt, zählen fehlgeschlagene und nicht unterstützte Dateien weiterhin dagegen, bis sie gelöscht werden — Platz schaffen heißt also, nicht mehr benötigte Dateien zu entfernen. Ein Klick auf ein Dokument öffnet die Vorschau, mit einer Seitenleiste für Größe, Quelle, RAG-Status, Teams, hochladende Person und Änderungsdatum — der schnellste Weg zu prüfen, worauf ein Zitat wirklich zeigt. diff --git a/docs/en/platform/knowledge/documents.md b/docs/en/platform/knowledge/documents.md index 9e21779766..beaddf37b8 100644 --- a/docs/en/platform/knowledge/documents.md +++ b/docs/en/platform/knowledge/documents.md @@ -79,7 +79,7 @@ Deleting a folder permanently deletes every file and subfolder inside it. Deleti **Reindex** (row menu) re-runs the pipeline on the stored file — the right move after an indexing failure or when a document shows **Needs reindex**. **Delete** removes the document and its indexed chunks; the confirmation says it plainly — the action cannot be undone. Re-uploading the same file brings the content back as a fresh document. A controlled record stops being deletable the moment any of its versions is approved — in review, approved, or drafting the next revision, the menu entry reads **Protected controlled record** instead, and a folder holding such a record refuses folder deletion the same way. The approved snapshot is a retained record; that is the point of the lifecycle. -Each document shows a status: **Queued** (waiting its turn — a busy organization indexes a few files at a time and the rest queue), **Indexing**, **Indexed**, **Failed**, or **Unsupported** (a legacy format such as `.doc`/`.ppt`/`.xls` that stores and downloads fine but has no text extractor, so it is never indexed for search). An indexing job interrupted by a timeout or a backend restart recovers on its own within a few minutes — it is retried or marked **Failed** with a retry option, never left stuck. If your organization enforces a per-user storage quota, failed and unsupported files still count against it until deleted, so freeing space means removing files you no longer need. +Each document shows a status: **Queued** (waiting its turn — a busy organization indexes a few files at a time and the rest queue), **Indexing**, **Indexed**, **Failed**, or **Unsupported** (a legacy format such as `.doc`/`.ppt`/`.xls`, or an image such as `.png`/`.jpg` — it stores and downloads fine but has no text extractor, so it is never indexed for search). An indexing job interrupted by a timeout or a backend restart recovers on its own within a few minutes — it is retried or marked **Failed** with a retry option, never left stuck. If your organization enforces a per-user storage quota, failed and unsupported files still count against it until deleted, so freeing space means removing files you no longer need. Clicking a document opens the preview, with a sidebar showing size, source, RAG status, teams, uploader, and modification date — the fastest way to check what a citation actually points at. diff --git a/docs/fr/platform/knowledge/documents.md b/docs/fr/platform/knowledge/documents.md index 0bd5b598e6..c3beb4e0bf 100644 --- a/docs/fr/platform/knowledge/documents.md +++ b/docs/fr/platform/knowledge/documents.md @@ -79,7 +79,7 @@ Supprimer un dossier supprime définitivement chaque fichier et sous-dossier qu **Réindexer** (menu de la ligne) refait passer le pipeline sur le fichier stocké — le bon geste après un échec d’indexation ou quand un document affiche **Réindexation nécessaire**. **Supprimer** retire le document et ses fragments indexés ; la confirmation le dit sans détour — l’action est irréversible. Retéléverser le même fichier ramène le contenu sous la forme d’un nouveau document. Un document maîtrisé cesse d’être supprimable dès qu’une de ses versions est approuvée — en relecture, approuvé ou avec le brouillon suivant ouvert, l’entrée du menu affiche **Document maîtrisé protégé**, et un dossier qui en contient un refuse la suppression du dossier de la même façon. L’instantané approuvé est un enregistrement conservé ; c’est précisément le rôle du cycle de vie. -Chaque document affiche un statut : **En file** (en attente — une organisation chargée indexe quelques fichiers à la fois et le reste patiente), **Indexation**, **Indexé**, **Échoué** ou **Non pris en charge** (un ancien format comme `.doc`/`.ppt`/`.xls` qui se stocke et se télécharge sans souci mais n’a pas d’extracteur de texte, donc jamais indexé pour la recherche). Une indexation interrompue par un délai dépassé ou un redémarrage du backend se rétablit d’elle-même en quelques minutes — elle est relancée ou marquée **Échoué** avec une option de reprise, jamais laissée bloquée. Si ton organisation applique un quota de stockage par utilisateur, les fichiers échoués et non pris en charge comptent toujours dedans jusqu’à leur suppression : libérer de l’espace revient donc à retirer les fichiers dont tu n’as plus besoin. +Chaque document affiche un statut : **En file** (en attente — une organisation chargée indexe quelques fichiers à la fois et le reste patiente), **Indexation**, **Indexé**, **Échoué** ou **Non pris en charge** (un ancien format comme `.doc`/`.ppt`/`.xls`, ou une image comme `.png`/`.jpg` — ça se stocke et se télécharge sans souci mais n’a pas d’extracteur de texte, donc jamais indexé pour la recherche). Une indexation interrompue par un délai dépassé ou un redémarrage du backend se rétablit d’elle-même en quelques minutes — elle est relancée ou marquée **Échoué** avec une option de reprise, jamais laissée bloquée. Si ton organisation applique un quota de stockage par utilisateur, les fichiers échoués et non pris en charge comptent toujours dedans jusqu’à leur suppression : libérer de l’espace revient donc à retirer les fichiers dont tu n’as plus besoin. Cliquer sur un document ouvre l’aperçu, avec un panneau latéral qui montre la taille, la source, le statut RAG, les équipes, l’auteur du téléversement et la date de modification — le moyen le plus rapide de vérifier ce que vise réellement une citation. diff --git a/services/platform/backend/core/lib/file_io.test.ts b/services/platform/backend/core/lib/file_io.test.ts index 3b0e7d0a81..966bf36151 100644 --- a/services/platform/backend/core/lib/file_io.test.ts +++ b/services/platform/backend/core/lib/file_io.test.ts @@ -1,6 +1,8 @@ // @vitest-environment node import { + chmodSync, + mkdirSync, mkdtempSync, readFileSync, rmSync, @@ -12,7 +14,10 @@ import path from 'node:path'; import { afterEach, beforeEach, describe, expect, it } from 'vitest'; -import { atomicWriteSecret } from './file_io'; +import { atomicWriteSecret, readJsonFile } from './file_io'; + +/** Root bypasses file permissions, so the EACCES lane cannot be produced. */ +const IS_ROOT = typeof process.getuid === 'function' && process.getuid() === 0; let dir: string; let prevUmask: number; @@ -73,3 +78,52 @@ describe('atomicWriteSecret', () => { expect(remaining).toEqual([path.basename(target)]); }); }); + +describe('readJsonFile', () => { + const parse = (content: string): unknown => JSON.parse(content); + + it('reads and hashes a well-formed file', async () => { + const target = path.join(dir, 'config.json'); + writeFileSync(target, '{"a":1}'); + const result = await readJsonFile(target, 1024, parse); + expect(result.ok).toBe(true); + if (result.ok) expect(result.data).toEqual({ a: 1 }); + }); + + it('labels a genuinely missing file not_found', async () => { + const result = await readJsonFile(path.join(dir, 'nope.json'), 1024, parse); + expect(result).toMatchObject({ ok: false, error: 'not_found' }); + }); + + it('labels a path through a regular file not_found (ENOTDIR)', async () => { + const file = path.join(dir, 'file.txt'); + writeFileSync(file, 'x'); + const result = await readJsonFile( + path.join(file, 'config.json'), + 1024, + parse, + ); + expect(result).toMatchObject({ ok: false, error: 'not_found' }); + }); + + // Regression: every stat() failure used to read as `not_found`, so a + // mis-permissioned config volume silently downgraded governance policies + // to their defaults. Only ENOENT/ENOTDIR are "absent"; EACCES is a fault. + it.skipIf(IS_ROOT)( + 'labels a present-but-unreadable file inaccessible, not not_found', + async () => { + const locked = path.join(dir, 'locked'); + mkdirSync(locked); + const target = path.join(locked, 'config.json'); + writeFileSync(target, '{"a":1}'); + chmodSync(locked, 0o000); + try { + const result = await readJsonFile(target, 1024, parse); + expect(result).toMatchObject({ ok: false, error: 'inaccessible' }); + if (!result.ok) expect(result.message).toMatch(/EACCES|permission/i); + } finally { + chmodSync(locked, 0o700); + } + }, + ); +}); diff --git a/services/platform/backend/core/lib/file_io.ts b/services/platform/backend/core/lib/file_io.ts index d434483675..5149f884b9 100644 --- a/services/platform/backend/core/lib/file_io.ts +++ b/services/platform/backend/core/lib/file_io.ts @@ -378,11 +378,23 @@ export async function readJsonFile( let fileStat; try { fileStat = await stat(filePath); - } catch { + } catch (err) { + // Only a genuinely missing file is `not_found` (a component that is not + // a directory means the same thing). A permission or I/O failure is + // `inaccessible`: callers treat `not_found` as "use the defaults", and a + // mis-permissioned config volume must surface, never read as absent. + const code = errnoCode(err); + if (code === 'ENOENT' || code === 'ENOTDIR') { + return { + ok: false, + error: 'not_found', + message: `File not found: ${path.basename(filePath)}`, + }; + } return { ok: false, - error: 'not_found', - message: `File not found: ${path.basename(filePath)}`, + error: 'inaccessible', + message: `Failed to stat file: ${path.basename(filePath)} — ${err instanceof Error ? err.message : String(err)}`, }; } @@ -403,16 +415,9 @@ export async function readJsonFile( await fd.close(); } } catch (err) { - const code = err instanceof Error && 'code' in err ? err.code : undefined; - const errorType = - code === 'ENOENT' - ? 'not_found' - : code === 'EACCES' || code === 'EPERM' - ? 'inaccessible' - : 'inaccessible'; return { ok: false, - error: errorType, + error: errnoCode(err) === 'ENOENT' ? 'not_found' : 'inaccessible', message: `Failed to read file: ${path.basename(filePath)} — ${err instanceof Error ? err.message : String(err)}`, }; } diff --git a/services/platform/backend/core/lib/knowledge/extraction/router.test.ts b/services/platform/backend/core/lib/knowledge/extraction/router.test.ts index d74987e546..3543837a48 100644 --- a/services/platform/backend/core/lib/knowledge/extraction/router.test.ts +++ b/services/platform/backend/core/lib/knowledge/extraction/router.test.ts @@ -1,7 +1,12 @@ import JSZip from 'jszip'; import { describe, expect, it } from 'vitest'; -import { ALL_SUPPORTED_EXTENSIONS, extractText, isSupported } from './router'; +import { + ALL_SUPPORTED_EXTENSIONS, + extractText, + isImageFile, + isSupported, +} from './router'; const enc = (s: string): Uint8Array => new TextEncoder().encode(s); @@ -100,3 +105,26 @@ describe('ALL_SUPPORTED_EXTENSIONS', () => { } }); }); + +describe('isImageFile', () => { + it('recognizes every image extension the router routes to the vision extractor', () => { + for (const f of [ + 'photo.png', + 'PHOTO.JPG', + 'scan.jpeg', + 'anim.gif', + 'pic.webp', + 'raw.bmp', + 'scan.tiff', + 'scan.tif', + ]) { + expect(isImageFile(f)).toBe(true); + } + }); + + it('is false for every non-image format, supported or not', () => { + for (const f of ['doc.pdf', 'doc.docx', 'notes.md', 'data.xlsx', 'x.exe']) { + expect(isImageFile(f)).toBe(false); + } + }); +}); diff --git a/services/platform/backend/core/lib/knowledge/extraction/router.ts b/services/platform/backend/core/lib/knowledge/extraction/router.ts index 52274cdad6..39f68e3c48 100644 --- a/services/platform/backend/core/lib/knowledge/extraction/router.ts +++ b/services/platform/backend/core/lib/knowledge/extraction/router.ts @@ -36,6 +36,17 @@ export function isSupported(filename: string): boolean { return ALL_SUPPORTED_EXTENSIONS.has(extname(filename).toLowerCase()); } +/** + * Does this file route to the IMAGE extractor? Image extraction is entirely + * vision-backed (`extractTextFromImageBytes` yields '' without a + * `VisionClient`), so a caller with no vision lane can decide up front that + * the file has nothing it can index — instead of downloading, extracting + * nothing, and reporting a failure. + */ +export function isImageFile(filename: string): boolean { + return SUPPORTED_IMAGE_EXTENSIONS.has(extname(filename).toLowerCase()); +} + export interface ExtractTextOptions { visionClient?: VisionClient | null; processImages?: boolean; diff --git a/services/platform/backend/core/lib/sops.test.ts b/services/platform/backend/core/lib/sops.test.ts new file mode 100644 index 0000000000..6490c0b606 --- /dev/null +++ b/services/platform/backend/core/lib/sops.test.ts @@ -0,0 +1,316 @@ +// @vitest-environment node + +/** + * The `sops` shell-out is reached from HTTP handlers (object-store resolves, + * credential saves) on a single-threaded event loop, so the contract under + * test is: the child runs OFF the loop, a wedged child is killed at the + * timeout, and concurrent cold reads of one file spawn one child. A fake + * `sops` on PATH makes each of those observable and deterministic; the real + * binary (when installed) proves the encrypt → decrypt round trip. + */ + +import { execFileSync } from 'node:child_process'; +import { + chmodSync, + mkdirSync, + mkdtempSync, + readFileSync, + rmSync, + utimesSync, + writeFileSync, +} from 'node:fs'; +import { tmpdir } from 'node:os'; +import path from 'node:path'; + +import { + afterAll, + afterEach, + beforeAll, + beforeEach, + describe, + expect, + it, +} from 'vitest'; + +import { deriveAgePublicKey } from './age_keygen'; +import { + decryptSecretsFile, + encryptJsonWithSops, + EncryptedFileWithoutKeyError, + invalidateSecretsCache, +} from './sops'; + +/** A throwaway age identity minted for this suite — it protects nothing. */ +const TEST_AGE_KEY = + 'AGE-SECRET-KEY-1V7RF0SP7WHE2LNTTHLYL2ART4WRD09HCP3VU5M4X5NCNUS6MQ9WQ2UDVCK'; + +const ORIGINAL_PATH = process.env.PATH ?? ''; +const ORIGINAL_ENV = { + key: process.env.SOPS_AGE_KEY, + keyFile: process.env.SOPS_AGE_KEY_FILE, +}; + +function realSopsAvailable(): boolean { + try { + execFileSync('sops', ['--version'], { + stdio: 'ignore', + env: { ...process.env, PATH: ORIGINAL_PATH }, + }); + return true; + } catch { + return false; + } +} +const HAS_REAL_SOPS = realSopsAvailable(); + +/** + * The fake `sops`: logs every invocation to FAKE_SOPS_LOG, waits + * FAKE_SOPS_DELAY_MS, then answers `-d` with FAKE_SOPS_DECRYPTED and `-e` + * with an encrypted-shaped document. A single node process, so a kill closes + * its stdio at once (a shell wrapper would leave a `sleep` holding the pipe). + */ +const FAKE_SOPS_SOURCE = `#!/usr/bin/env node +const fs = require('node:fs'); +const args = process.argv.slice(2); +fs.appendFileSync(process.env.FAKE_SOPS_LOG, args.join(' ') + '\\n'); +const delay = Number(process.env.FAKE_SOPS_DELAY_MS || '0'); +setTimeout(() => { + if (process.env.FAKE_SOPS_FAIL) { + process.stderr.write('fake sops: ' + process.env.FAKE_SOPS_FAIL); + process.exitCode = 1; + } else if (args[0] === '-d') { + process.stdout.write(process.env.FAKE_SOPS_DECRYPTED || '{}'); + } else if (args[0] === '-e') { + process.stdout.write(JSON.stringify({ sops: { age: [] }, data: 'ENC[fake]' })); + } else { + process.stderr.write('fake sops: unknown verb ' + args[0]); + process.exitCode = 1; + } +}, delay); +`; + +let root: string; +let fakeBin: string; +let logFile: string; +let fileCounter = 0; + +const ENCRYPTED_SHAPE = JSON.stringify({ + apiKey: 'ENC[AES256_GCM,data:xx,type:str]', + sops: { age: [{ recipient: 'age1test' }], version: '3.9.4' }, +}); + +function writeSecretsFile(content: string): string { + fileCounter += 1; + const file = path.join(root, `secrets-${fileCounter}.json`); + writeFileSync(file, content); + return file; +} + +function invocations(): string[] { + try { + return readFileSync(logFile, 'utf-8').split('\n').filter(Boolean); + } catch { + return []; + } +} + +/** A timer armed alongside an awaited call — records WHEN the loop got to + * run it, which is how a blocked loop shows up as a late tick. */ +function armTick(delayMs: number): { + done: Promise; + firedAt: () => number; +} { + let firedAt = 0; + const done = new Promise((resolve) => { + setTimeout(() => { + firedAt = Date.now(); + resolve(); + }, delayMs); + }); + return { done, firedAt: () => firedAt }; +} + +beforeAll(() => { + root = mkdtempSync(path.join(tmpdir(), 'sops-test-')); + fakeBin = path.join(root, 'bin'); + mkdirSync(fakeBin); + const fake = path.join(fakeBin, 'sops'); + writeFileSync(fake, FAKE_SOPS_SOURCE); + chmodSync(fake, 0o755); + process.env.PATH = `${fakeBin}${path.delimiter}${ORIGINAL_PATH}`; +}); + +afterAll(() => { + process.env.PATH = ORIGINAL_PATH; + rmSync(root, { recursive: true, force: true }); +}); + +beforeEach(() => { + logFile = path.join(root, `log-${Date.now()}-${Math.random()}.txt`); + process.env.FAKE_SOPS_LOG = logFile; + process.env.FAKE_SOPS_DELAY_MS = '0'; + process.env.FAKE_SOPS_DECRYPTED = JSON.stringify({ apiKey: 'decrypted' }); + process.env.SOPS_AGE_KEY = TEST_AGE_KEY; + delete process.env.SOPS_AGE_KEY_FILE; +}); + +afterEach(() => { + if (ORIGINAL_ENV.key === undefined) delete process.env.SOPS_AGE_KEY; + else process.env.SOPS_AGE_KEY = ORIGINAL_ENV.key; + if (ORIGINAL_ENV.keyFile === undefined) delete process.env.SOPS_AGE_KEY_FILE; + else process.env.SOPS_AGE_KEY_FILE = ORIGINAL_ENV.keyFile; +}); + +describe('decryptSecretsFile', () => { + it('runs sops off the event loop — a timer fires while the decrypt is in flight', async () => { + process.env.FAKE_SOPS_DELAY_MS = '600'; + const file = writeSecretsFile(ENCRYPTED_SHAPE); + const startedAt = Date.now(); + const tick = armTick(20); + + const data = await decryptSecretsFile(file); + await tick.done; + + expect(data).toEqual({ apiKey: 'decrypted' }); + // A synchronous spawn would have held the loop for the whole 600ms and + // the 20ms timer would only fire afterwards. + expect(tick.firedAt() - startedAt).toBeLessThan(300); + expect(Date.now() - startedAt).toBeGreaterThanOrEqual(550); + }); + + it('kills a hung sops at the timeout and says so', async () => { + process.env.FAKE_SOPS_DELAY_MS = '60000'; + const file = writeSecretsFile(ENCRYPTED_SHAPE); + const startedAt = Date.now(); + + await expect(decryptSecretsFile(file, { timeoutMs: 300 })).rejects.toThrow( + /timed out after 300ms/, + ); + + expect(Date.now() - startedAt).toBeLessThan(5_000); + expect(invocations()).toHaveLength(1); + }); + + it('shares one sops process across concurrent cold reads of the same file', async () => { + process.env.FAKE_SOPS_DELAY_MS = '200'; + const file = writeSecretsFile(ENCRYPTED_SHAPE); + + const results = await Promise.all( + Array.from({ length: 5 }, () => decryptSecretsFile(file)), + ); + + for (const result of results) { + expect(result).toEqual({ apiKey: 'decrypted' }); + } + expect(invocations().filter((line) => line.startsWith('-d '))).toHaveLength( + 1, + ); + }); + + it('serves the cache until the file changes, and re-reads after invalidation', async () => { + const file = writeSecretsFile(ENCRYPTED_SHAPE); + expect(await decryptSecretsFile(file)).toEqual({ apiKey: 'decrypted' }); + expect(await decryptSecretsFile(file)).toEqual({ apiKey: 'decrypted' }); + expect(invocations()).toHaveLength(1); + + // Same content, new mtime → a re-read. + const later = new Date(Date.now() + 5_000); + utimesSync(file, later, later); + process.env.FAKE_SOPS_DECRYPTED = JSON.stringify({ apiKey: 'rotated' }); + expect(await decryptSecretsFile(file)).toEqual({ apiKey: 'rotated' }); + expect(invocations()).toHaveLength(2); + + invalidateSecretsCache(file); + expect(await decryptSecretsFile(file)).toEqual({ apiKey: 'rotated' }); + expect(invocations()).toHaveLength(3); + }); + + it('reads a plaintext file without spawning sops', async () => { + const file = writeSecretsFile(JSON.stringify({ apiKey: 'plain' })); + expect(await decryptSecretsFile(file)).toEqual({ apiKey: 'plain' }); + expect(invocations()).toHaveLength(0); + }); + + it('refuses an encrypted file when no age key is configured', async () => { + delete process.env.SOPS_AGE_KEY; + const file = writeSecretsFile(ENCRYPTED_SHAPE); + await expect(decryptSecretsFile(file)).rejects.toBeInstanceOf( + EncryptedFileWithoutKeyError, + ); + expect(invocations()).toHaveLength(0); + }); + + it('surfaces a failing sops with its stderr', async () => { + process.env.FAKE_SOPS_FAIL = 'no identity matched any of the recipients'; + try { + const file = writeSecretsFile(ENCRYPTED_SHAPE); + await expect(decryptSecretsFile(file)).rejects.toThrow( + /no identity matched any of the recipients/, + ); + } finally { + delete process.env.FAKE_SOPS_FAIL; + } + }); +}); + +describe('encryptJsonWithSops', () => { + it('runs sops off the event loop and addresses every configured recipient', async () => { + process.env.FAKE_SOPS_DELAY_MS = '300'; + const startedAt = Date.now(); + const tick = armTick(20); + + const encrypted = await encryptJsonWithSops('{"apiKey":"v"}'); + await tick.done; + + // oxlint-disable-next-line typescript/no-unsafe-type-assertion -- the fake's answer shape + const parsed = JSON.parse(encrypted) as { sops?: unknown }; + expect(parsed.sops).toBeDefined(); + expect(tick.firedAt() - startedAt).toBeLessThan(200); + const call = invocations()[0] ?? ''; + expect(call.startsWith('-e ')).toBe(true); + expect(call).toContain(`--age ${deriveAgePublicKey(TEST_AGE_KEY)}`); + }); + + it('kills a hung sops at the timeout', async () => { + process.env.FAKE_SOPS_DELAY_MS = '60000'; + const startedAt = Date.now(); + await expect( + encryptJsonWithSops('{"apiKey":"v"}', { timeoutMs: 300 }), + ).rejects.toThrow(/timed out after 300ms/); + expect(Date.now() - startedAt).toBeLessThan(5_000); + }); + + it('refuses to encrypt without an age key', async () => { + delete process.env.SOPS_AGE_KEY; + await expect(encryptJsonWithSops('{}')).rejects.toThrow( + /No age secret key/, + ); + expect(invocations()).toHaveLength(0); + }); +}); + +describe.skipIf(!HAS_REAL_SOPS)('real sops round trip', () => { + beforeEach(() => { + // The real binary, not the fake. + process.env.PATH = ORIGINAL_PATH; + }); + afterEach(() => { + process.env.PATH = `${fakeBin}${path.delimiter}${ORIGINAL_PATH}`; + }); + + it('encrypts to the configured age recipient and decrypts it back', async () => { + const encrypted = await encryptJsonWithSops( + JSON.stringify({ apiKey: 'round-trip', nested: { n: 1 } }), + ); + // oxlint-disable-next-line typescript/no-unsafe-type-assertion -- sops output shape + const shape = JSON.parse(encrypted) as { sops?: unknown; apiKey?: string }; + expect(shape.sops).toBeDefined(); + expect(shape.apiKey).toMatch(/^ENC\[/); + + const file = writeSecretsFile(encrypted); + expect(await decryptSecretsFile(file)).toEqual({ + apiKey: 'round-trip', + nested: { n: 1 }, + }); + }, 20_000); +}); diff --git a/services/platform/backend/core/lib/sops.ts b/services/platform/backend/core/lib/sops.ts index 590faca790..985661b377 100644 --- a/services/platform/backend/core/lib/sops.ts +++ b/services/platform/backend/core/lib/sops.ts @@ -21,11 +21,18 @@ * for plaintext), so toggling env vars between calls can't return a * mismatched cached value — the cache key is the file's content, not the * current env. + * + * The `sops` child process is ALWAYS spawned asynchronously. Both verbs are + * reached from HTTP handlers (every per-org object-store resolve decrypts on + * cache expiry; every credential save encrypts), and the backend is one + * single-threaded event loop: a synchronous spawn would stall every in-flight + * request for the whole decrypt — up to the timeout when sops hangs. The + * timeout kills a wedged sops so no request waits on it forever, and + * concurrent decrypts of one file share a single child process. */ -import { execFileSync } from 'node:child_process'; -import { mkdtempSync, rmdirSync, unlinkSync, writeFileSync } from 'node:fs'; -import { readFile, stat } from 'node:fs/promises'; +import { execFile } from 'node:child_process'; +import { mkdtemp, readFile, rm, stat, writeFile } from 'node:fs/promises'; import { tmpdir } from 'node:os'; import path from 'node:path'; @@ -38,8 +45,24 @@ interface CacheEntry { const cache = new Map(); +/** Decrypts in flight, keyed by file path — concurrent callers on a cold + * cache share one child process instead of each spawning `sops`. */ +const inflight = new Map>>(); + let plaintextWarnEmitted = false; +/** Upper bound on one `sops` invocation. Secrets files are a few KB, so a + * healthy sops answers in well under a second; anything near this bound is + * a wedged binary (a hung key-file read, a stuck KMS call) and the caller + * gets an error instead of an indefinitely pending request. */ +export const SOPS_TIMEOUT_MS = 10_000; + +export interface SopsOptions { + /** Kill the `sops` child after this many milliseconds (default + * `SOPS_TIMEOUT_MS`). */ + timeoutMs?: number; +} + export class EncryptedFileWithoutKeyError extends Error { constructor(filePath: string) { super( @@ -92,12 +115,56 @@ export function invalidateSecretsCache(filePath: string): void { cache.delete(filePath); } +/** + * Run `sops` with `args` off the event loop and resolve with its stdout. + * Rejects with a legible error on a non-zero exit (stderr attached) and on + * the timeout (the child is killed, and the error says so rather than + * surfacing a bare ETIMEDOUT). + */ +function runSops(args: readonly string[], timeoutMs: number): Promise { + return new Promise((resolve, reject) => { + execFile( + 'sops', + [...args], + { + encoding: 'utf-8', + timeout: timeoutMs, + killSignal: 'SIGKILL', + // Secrets files are a few KB; the default 1MB would already do, but a + // silent truncation on an unusually large file must never happen. + maxBuffer: 8 * 1024 * 1024, + env: { ...process.env }, + }, + (err, stdout, stderr) => { + if (err === null) { + resolve(stdout); + return; + } + const detail = stderr.trim().length > 0 ? stderr.trim() : err.message; + const timedOut = + // oxlint-disable-next-line typescript/no-unsafe-type-assertion -- execFile's error carries `killed`/`signal` (Node's ExecFileException) + (err as NodeJS.ErrnoException & { killed?: boolean }).killed === true; + reject( + Object.assign( + new Error( + timedOut + ? `sops ${args[0]} timed out after ${timeoutMs}ms (killed)` + : `sops ${args[0]} failed: ${detail}`, + ), + { cause: err }, + ), + ); + }, + ); + }); +} + /** * Encrypt a JSON plaintext string with SOPS, addressed to every configured * age recipient — so any key present in `SOPS_AGE_KEY_FILE` can decrypt the - * result later (the rotation primitive). Returns the encrypted SOPS-JSON - * string; throws if no age key is configured, or if `sops` itself fails. - * Callers that want to fall back to plaintext mode must check + * result later (the rotation primitive). Resolves with the encrypted + * SOPS-JSON string; rejects if no age key is configured, or if `sops` itself + * fails. Callers that want to fall back to plaintext mode must check * `hasSopsKey()` themselves before calling this. * * The plaintext is written to a 0o600 file inside a 0o700 mkdtemp @@ -105,7 +172,10 @@ export function invalidateSecretsCache(filePath: string): void { * Cleanup failures are logged rather than swallowed — a leftover temp file * holds a plaintext secret until the OS reaps it. */ -export function encryptJsonWithSops(plaintext: string): string { +export async function encryptJsonWithSops( + plaintext: string, + options: SopsOptions = {}, +): Promise { const recipients = resolveAgeRecipients(); if (recipients.length === 0) { throw new Error( @@ -113,12 +183,11 @@ export function encryptJsonWithSops(plaintext: string): string { 'SOPS_AGE_KEY_FILE (path) in .env, or unset both to use plaintext mode.', ); } - const tmpDir = mkdtempSync(path.join(tmpdir(), 'sops-')); + const tmpDir = await mkdtemp(path.join(tmpdir(), 'sops-')); const tmpFile = path.join(tmpDir, 'plain.json'); try { - writeFileSync(tmpFile, plaintext, { encoding: 'utf-8', mode: 0o600 }); - return execFileSync( - 'sops', + await writeFile(tmpFile, plaintext, { encoding: 'utf-8', mode: 0o600 }); + return await runSops( [ '-e', '--input-type', @@ -129,13 +198,10 @@ export function encryptJsonWithSops(plaintext: string): string { recipients.join(','), tmpFile, ], - { encoding: 'utf-8', timeout: 10_000, stdio: ['pipe', 'pipe', 'pipe'] }, + options.timeoutMs ?? SOPS_TIMEOUT_MS, ); } catch (err) { const message = err instanceof Error ? err.message : String(err); - // Object.assign bolts `cause` onto the Error: convex/tsconfig.json's - // "lib" predates the ES2022 two-argument Error constructor overload, - // even though the runtime itself supports it. throw Object.assign( new Error( `Failed to encrypt secrets with SOPS: ${message}. ` + @@ -145,41 +211,42 @@ export function encryptJsonWithSops(plaintext: string): string { ); } finally { try { - unlinkSync(tmpFile); + await rm(tmpDir, { recursive: true, force: true }); } catch (err) { - // oxlint-disable-next-line typescript/no-unsafe-type-assertion -- Node.js errors always carry .code - if ((err as NodeJS.ErrnoException).code !== 'ENOENT') { - console.warn( - `[sops] failed to remove temp plaintext ${tmpFile}: ${ - err instanceof Error ? err.message : String(err) - }`, - ); - } - } - try { - rmdirSync(tmpDir); - } catch (err) { - // oxlint-disable-next-line typescript/no-unsafe-type-assertion -- Node.js errors always carry .code - if ((err as NodeJS.ErrnoException).code !== 'ENOENT') { - console.warn( - `[sops] failed to remove temp dir ${tmpDir}: ${ - err instanceof Error ? err.message : String(err) - }`, - ); - } + console.warn( + `[sops] failed to remove temp plaintext dir ${tmpDir}: ${ + err instanceof Error ? err.message : String(err) + }`, + ); } } } export async function decryptSecretsFile( filePath: string, + options: SopsOptions = {}, ): Promise> { const fileStat = await stat(filePath); const cached = cache.get(filePath); if (cached && cached.mtimeMs === fileStat.mtimeMs) { return cached.data; } + const pending = inflight.get(filePath); + if (pending) return pending; + const load = loadSecretsFile(filePath, fileStat.mtimeMs, options).finally( + () => { + inflight.delete(filePath); + }, + ); + inflight.set(filePath, load); + return load; +} +async function loadSecretsFile( + filePath: string, + mtimeMs: number, + options: SopsOptions, +): Promise> { const raw = await readFile(filePath, 'utf-8'); let parsed: unknown; @@ -187,9 +254,6 @@ export async function decryptSecretsFile( parsed = JSON.parse(raw); } catch (err) { const message = err instanceof Error ? err.message : String(err); - // Object.assign bolts `cause` onto the Error: convex/tsconfig.json's - // "lib" predates the ES2022 two-argument Error constructor overload, - // even though the runtime itself supports it. throw Object.assign( new Error( `Failed to parse secrets file ${filePath} as JSON: ${message}.`, @@ -205,17 +269,12 @@ export async function decryptSecretsFile( } let stdout: string; try { - stdout = execFileSync('sops', ['-d', '--output-type', 'json', filePath], { - encoding: 'utf-8', - timeout: 10_000, - stdio: ['pipe', 'pipe', 'pipe'], - env: { ...process.env }, - }); + stdout = await runSops( + ['-d', '--output-type', 'json', filePath], + options.timeoutMs ?? SOPS_TIMEOUT_MS, + ); } catch (err) { const message = err instanceof Error ? err.message : String(err); - // Object.assign bolts `cause` onto the Error: convex/tsconfig.json's - // "lib" predates the ES2022 two-argument Error constructor overload, - // even though the runtime itself supports it. throw Object.assign( new Error( `Failed to decrypt secrets file ${filePath}: ${message}. ` + @@ -241,6 +300,6 @@ export async function decryptSecretsFile( data = parsed as Record; } - cache.set(filePath, { data, mtimeMs: fileStat.mtimeMs }); + cache.set(filePath, { data, mtimeMs }); return data; } diff --git a/services/platform/backend/core/lib/storage/blob_access.ts b/services/platform/backend/core/lib/storage/blob_access.ts index facb951efc..e08f947288 100644 --- a/services/platform/backend/core/lib/storage/blob_access.ts +++ b/services/platform/backend/core/lib/storage/blob_access.ts @@ -2,8 +2,10 @@ /** * Backend-aware blob access — the single seam every org-owned blob operation - * routes through so a blob transparently lives in Convex `_storage` (deployment - * default) OR the org's own S3 bucket (bring-your-own object storage). + * routes through. New blobs always land in S3-compatible storage (the org's own + * bucket, else the deployment default's — see `object_store.ts`, which fails + * closed when neither is configured); the Convex-id branches below only read, + * serve, or delete a LEGACY reference that predates the cutover. * * # The blob reference * @@ -49,22 +51,18 @@ export type { BlobRef } from './blob_ref'; export { encodeS3Ref, parseBlobRef, isS3Ref } from './blob_ref'; /** - * Store bytes for an org and return the stored reference. Routes to the org's - * S3 bucket when configured, else Convex `_storage`. Action ctx required (S3 - * signing needs node). + * Store bytes for an org and return the stored reference — always an `s3:` + * ref into the org's resolved bucket; an unconfigured store throws at this + * door instead of failing deeper in the lane. The `ctx` parameter is kept for + * the reused 0.4 call shape. */ export async function putBlob( - ctx: ActionCtx, + _ctx: ActionCtx, orgSlug: string, bytes: Uint8Array, contentType: string, ): Promise { const store = await resolveOrgObjectStore(orgSlug); - if (store.backend === 'convex') { - // oxlint-disable-next-line typescript/no-unsafe-type-assertion -- a Uint8Array is a valid BlobPart at runtime (TS 5.7 ArrayBufferLike variance) - const blob = new Blob([bytes as BlobPart], { type: contentType }); - return await ctx.storage.store(blob); - } const key = buildObjectKey(store, orgSlug); await s3PutObject(store, key, bytes, contentType); return encodeS3Ref(key); @@ -145,21 +143,17 @@ export async function getBlobUrl( } /** - * Upload handoff for the client. Convex: `generateUploadUrl` (the client POSTs - * and learns the id from the response). S3: a presigned PUT plus the ref the - * client will bind (the key is known up front). The caller returns `{ url, - * method, s3Ref }` to the browser; when `s3Ref` is present the client PUTs to - * `url` then binds `s3Ref`, else it POSTs and binds the returned storage id. + * Upload handoff for the client: a presigned PUT plus the ref the client will + * bind (the key is known up front). The caller returns `{ url, method, s3Ref }` + * to the browser, which PUTs to `url` then binds `s3Ref`. The `method` union + * is the reused 0.4 wire shape; 0.5 only ever answers `PUT`. */ export async function generateBlobUpload( - ctx: ActionCtx, + _ctx: ActionCtx, orgSlug: string, opts: { contentType?: string } = {}, ): Promise<{ url: string; method: 'POST' | 'PUT'; s3Ref?: string }> { const store = await resolveOrgObjectStore(orgSlug); - if (store.backend === 'convex') { - return { url: await ctx.storage.generateUploadUrl(), method: 'POST' }; - } const key = buildObjectKey(store, orgSlug); const url = await s3PresignPutUrl(store, key, { contentType: opts.contentType, @@ -169,41 +163,29 @@ export async function generateBlobUpload( export interface ReplacementBlobUploadHandoff { url: string; - method: 'POST' | 'PUT'; - backend: 'convex' | 's3'; + method: 'PUT'; + backend: 's3'; uploadContentType: string; uploadExpiresAt: number; - stagingRef?: BlobRef; - finalRef?: BlobRef; + stagingRef: BlobRef; + finalRef: BlobRef; } -const CONVEX_UPLOAD_TTL_MS = 60 * 60 * 1000; - /** * Mint a replacement-specific upload capability. * * S3 receives two keys: the browser can write only the staging key, while the - * final key is reserved for a create-only server PUT after attestation. Convex - * storage is already immutable, so ownership is proven by an intent nonce in - * the stored content type. + * final key is reserved for a create-only server PUT after attestation. The + * `intentNonce` is kept for the reused call shape. */ export async function generateReplacementBlobUpload( - ctx: ActionCtx, + _ctx: ActionCtx, orgSlug: string, - intentNonce: string, + _intentNonce: string, contentType?: string, ): Promise { const store = await resolveOrgObjectStore(orgSlug); const baseContentType = contentType?.trim() || 'application/octet-stream'; - if (store.backend === 'convex') { - return { - url: await ctx.storage.generateUploadUrl(), - method: 'POST', - backend: 'convex', - uploadContentType: `${baseContentType}; tale-intent=${intentNonce}`, - uploadExpiresAt: Date.now() + CONVEX_UPLOAD_TTL_MS, - }; - } const stagingKey = buildObjectKey(store, orgSlug); const finalKey = buildObjectKey(store, orgSlug); @@ -258,11 +240,5 @@ async function requireS3(orgSlug: string, key: string): Promise { `s3 blob key is outside org '${orgSlug}' namespace; refusing`, ); } - const store = await resolveOrgObjectStore(orgSlug); - if (store.backend !== 's3') { - throw new Error( - `blob references org '${orgSlug}' S3 storage, but no S3 store is configured`, - ); - } - return store; + return resolveOrgObjectStore(orgSlug); } diff --git a/services/platform/backend/core/lib/storage/object_store.test.ts b/services/platform/backend/core/lib/storage/object_store.test.ts index 291cd03417..0f8d6cfbe0 100644 --- a/services/platform/backend/core/lib/storage/object_store.test.ts +++ b/services/platform/backend/core/lib/storage/object_store.test.ts @@ -1,7 +1,14 @@ -import { describe, expect, it } from 'vitest'; +import { mkdirSync, mkdtempSync, rmSync, writeFileSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import path from 'node:path'; + +import { afterEach, beforeEach, describe, expect, it } from 'vitest'; import { buildS3ObjectStore, + clearOrgObjectStoreCache, + ObjectStoreUnconfiguredError, + resolveOrgObjectStore, s3PresignGetUrl, s3PresignPutUrl, } from './object_store'; @@ -98,3 +105,106 @@ describe('s3PresignGetUrl — attachment forcing', () => { expect(url.searchParams.get('X-Amz-Signature')).toBeTruthy(); }); }); + +describe('resolveOrgObjectStore — fail-closed resolution', () => { + const previousConfigDir = process.env.TALE_CONFIG_DIR; + let configRoot: string; + + const connectionJson = (bucket: string): string => + JSON.stringify({ + region: 'us-east-1', + endpoint: 'http://minio.internal:9000', + forcePathStyle: true, + bucket, + }); + const SECRETS = JSON.stringify({ + accessKeyId: 'test-access', + secretAccessKey: 'test-secret', + }); + + function writeTree( + slug: string, + connection: string, + secrets: string | null = SECRETS, + ): void { + const dir = path.join(configRoot, slug, 'object-storage'); + mkdirSync(dir, { recursive: true }); + writeFileSync(path.join(dir, 'connection.json'), connection); + if (secrets !== null) { + writeFileSync(path.join(dir, 'connection.secrets.json'), secrets); + } + } + + beforeEach(() => { + configRoot = mkdtempSync(path.join(tmpdir(), 'object-store-test-')); + process.env.TALE_CONFIG_DIR = configRoot; + clearOrgObjectStoreCache(); + }); + + afterEach(() => { + if (previousConfigDir === undefined) delete process.env.TALE_CONFIG_DIR; + else process.env.TALE_CONFIG_DIR = previousConfigDir; + clearOrgObjectStoreCache(); + rmSync(configRoot, { recursive: true, force: true }); + }); + + // Regression: an org with no connection AND no default tree used to resolve + // to the retired Convex backend, whose `ctx.storage.store` no longer exists + // — the misconfiguration surfaced as a TypeError deep inside a blob lane. + it('throws ObjectStoreUnconfiguredError when neither the org nor the default tree is configured', async () => { + await expect(resolveOrgObjectStore('acme')).rejects.toBeInstanceOf( + ObjectStoreUnconfiguredError, + ); + }); + + it('serves the deployment default tree to an org without its own connection', async () => { + writeTree('default', connectionJson('default-blobs')); + const store = await resolveOrgObjectStore('acme'); + expect(store.backend).toBe('s3'); + expect(store.config.bucket).toBe('default-blobs'); + }); + + it("prefers the org's own connection over the default tree", async () => { + writeTree('default', connectionJson('default-blobs')); + writeTree('acme', connectionJson('acme-own-bucket')); + const store = await resolveOrgObjectStore('acme'); + expect(store.config.bucket).toBe('acme-own-bucket'); + }); + + // Regression: a broken default tree was swallowed (`.catch(() => null)`) and + // fell through to the dead fallback; it must surface as ITS OWN error. + it('surfaces a corrupt default connection as its own error, never a fallback', async () => { + writeTree('default', '{"region":"us-east-1"}'); + const failure = await resolveOrgObjectStore('acme').then( + () => null, + (err: unknown) => err, + ); + expect(failure).toBeInstanceOf(Error); + expect(failure).not.toBeInstanceOf(ObjectStoreUnconfiguredError); + expect(String(failure)).toMatch(/Invalid object-storage connection config/); + }); + + it('surfaces missing default credentials instead of signing with none', async () => { + writeTree('default', connectionJson('default-blobs'), null); + await expect(resolveOrgObjectStore('acme')).rejects.toThrow( + /credentials missing/, + ); + }); + + it('caches a resolution until the cache is cleared', async () => { + writeTree('default', connectionJson('default-blobs')); + expect((await resolveOrgObjectStore('acme')).config.bucket).toBe( + 'default-blobs', + ); + writeTree('acme', connectionJson('acme-own-bucket')); + // Still the cached default within the TTL… + expect((await resolveOrgObjectStore('acme')).config.bucket).toBe( + 'default-blobs', + ); + // …and the org's own bucket once a config write clears the cache. + clearOrgObjectStoreCache(); + expect((await resolveOrgObjectStore('acme')).config.bucket).toBe( + 'acme-own-bucket', + ); + }); +}); diff --git a/services/platform/backend/core/lib/storage/object_store.ts b/services/platform/backend/core/lib/storage/object_store.ts index e8276e7697..a0408ca192 100644 --- a/services/platform/backend/core/lib/storage/object_store.ts +++ b/services/platform/backend/core/lib/storage/object_store.ts @@ -4,22 +4,25 @@ * Per-organization object-store resolution + S3 verbs. * * The SINGLE per-org routing entry point for file blobs, the object-storage - * analogue of `getKnowledgePoolForOrg` for the RAG corpus. `resolveOrgObjectStore - * (orgSlug)` returns either the deployment default (Convex `_storage` — today's - * behaviour, zero regression) or, when the org has configured - * `/object-storage/connection.json`, an S3 backend addressing the org's own - * bucket. Callers that hold an `ActionCtx` use `blob_access.ts`, which routes - * `ctx.storage.*` vs. these S3 verbs off the resolved backend; this module owns - * only the resolution + the raw S3 requests. + * analogue of `getKnowledgePoolForOrg` for the RAG corpus. S3-compatible + * storage is THE blob backend (Convex `_storage` died with the component): + * `resolveOrgObjectStore(orgSlug)` returns the org's own bucket when + * `/object-storage/connection.json` is configured, else the deployment + * default — the `default` config tree's connection — and otherwise FAILS + * CLOSED with `ObjectStoreUnconfiguredError`. There is no fallback store: a + * missing or broken connection is an operator-visible error at the door, never + * a silently different backend deep in a blob lane. Callers that hold an + * `ActionCtx` use `blob_access.ts`; this module owns only the resolution + the + * raw S3 requests. * * S3 requests are signed with `aws4fetch` (a few-KB SigV4 signer) rather than - * `@aws-sdk/client-s3` — the AWS SDK is large and would risk the Convex module - * push-size cap. Works against any S3-compatible store (AWS S3, MinIO, R2, - * Wasabi) via `endpoint` + `forcePathStyle`. + * `@aws-sdk/client-s3`. Works against any S3-compatible store (AWS S3, MinIO, + * R2, Wasabi) via `endpoint` + `forcePathStyle`. * * TENANT ISOLATION: the store is keyed strictly by `orgSlug`; a per-org bucket is * NEVER addressed for another org. Resolution is fail-closed — a present but - * broken per-org config throws rather than silently using the shared default. + * broken config (the org's OR the default tree's) throws rather than silently + * using anything else. */ import { randomUUID } from 'node:crypto'; @@ -32,28 +35,35 @@ import { type ObjectStorageConnectionSecrets, } from '../../object_storage/file_utils'; -/** Deployment default: blobs live in Convex `_storage` (per-org logical scope). */ -export interface ConvexObjectStore { - backend: 'convex'; -} - -/** Per-org bring-your-own S3-compatible bucket (physical isolation). */ +/** An S3-compatible bucket: the org's own (physical isolation) or the + * deployment default's. */ export interface S3ObjectStore { backend: 's3'; client: AwsClient; config: ObjectStorageConnectionFile; } -export type ResolvedObjectStore = ConvexObjectStore | S3ObjectStore; +/** Neither the org nor the deployment default tree has an object-storage + * connection — uploads are refused until an operator configures one. */ +export class ObjectStoreUnconfiguredError extends Error { + constructor() { + super( + 'No object storage configured: neither this org nor the deployment ' + + 'default tree has an object-storage/connection.json', + ); + this.name = 'ObjectStoreUnconfiguredError'; + } +} -const CONVEX_STORE: ConvexObjectStore = { backend: 'convex' }; +/** The config tree whose connection serves every org without its own. */ +const DEFAULT_TREE_SLUG = 'default'; // Short-TTL resolution cache, mirroring `knowledge_db.ts` ORG_URL_TTL_MS: a // config change (admin edits the org's bucket) takes effect within the TTL // without a restart, and the hot path avoids a disk read + SOPS decrypt per blob. const ORG_STORE_TTL_MS = 15_000; interface CacheEntry { - store: ResolvedObjectStore; + store: S3ObjectStore; expires: number; } const orgStoreCache = new Map(); @@ -61,35 +71,40 @@ const orgStoreCache = new Map(); /** * Resolve an org's object store: its own S3 bucket when * `/object-storage/connection.json` is configured, else the deployment - * default (Convex `_storage`). Cached with a short TTL. Throws (fail-closed) if a - * present per-org config is invalid or its credentials can't be decrypted. + * default tree's connection. Cached with a short TTL. Throws (fail-closed) when + * a present config — the org's or the default's — is invalid or its credentials + * can't be decrypted, and `ObjectStoreUnconfiguredError` when neither exists. */ export async function resolveOrgObjectStore( orgSlug: string, -): Promise { +): Promise { const now = Date.now(); const cached = orgStoreCache.get(orgSlug); if (cached && cached.expires > now) { return cached.store; } - // The org's own BYO connection wins; without one, the DEPLOYMENT DEFAULT - // tree's connection serves every org (the 0.5 posture — Convex `_storage` - // dies at cutover). A 0.4 deployment ships no `default` tree, so its - // fallback order is unchanged (org → Convex). + const own = await readOrgObjectStorageConnection(orgSlug); + // A broken default tree is a real misconfiguration and must surface as its + // own error — swallowing it here would hide "undecryptable credentials" + // behind a generic "unconfigured" (or, historically, a dead fallback store). const resolved = - (await readOrgObjectStorageConnection(orgSlug)) ?? - (await readOrgObjectStorageConnection('default').catch(() => null)); - let store: ResolvedObjectStore; + own ?? + (orgSlug === DEFAULT_TREE_SLUG + ? null + : await readOrgObjectStorageConnection(DEFAULT_TREE_SLUG)); if (resolved === null) { - store = CONVEX_STORE; - } else { - store = buildS3ObjectStore(resolved.connection, resolved.secrets); - console.info(`Resolved per-org S3 object store for org '${orgSlug}'`); + throw new ObjectStoreUnconfiguredError(); } + const store = buildS3ObjectStore(resolved.connection, resolved.secrets); orgStoreCache.set(orgSlug, { store, expires: now + ORG_STORE_TTL_MS }); return store; } +/** Drop every cached resolution (test hook + config-write invalidation). */ +export function clearOrgObjectStoreCache(): void { + orgStoreCache.clear(); +} + /** * Build an `S3ObjectStore` from a connection + credentials WITHOUT touching disk * — the shared constructor for both `resolveOrgObjectStore` (which reads the diff --git a/services/platform/backend/core/node_only/sandbox/helpers/session_client.test.ts b/services/platform/backend/core/node_only/sandbox/helpers/session_client.test.ts index ab061973c0..7ee68aefb9 100644 --- a/services/platform/backend/core/node_only/sandbox/helpers/session_client.test.ts +++ b/services/platform/backend/core/node_only/sandbox/helpers/session_client.test.ts @@ -142,6 +142,85 @@ describe('drainSessionExecResilient', () => { }, 15_000); }); +describe('drainSessionExecResilient — a lost first POST', () => { + // Regression: every retry used to ATTACH. When the initial exec POST never + // reached the spawner (network drop, 503 mid-roll), the exec did not exist, + // so each attach answered `exec not found` and the whole budget burned + // against a healthy session. + test('re-creates the exec when the spawner answers not-found with nothing consumed', async () => { + const calls: { url: string; method: string }[] = []; + let n = 0; + // oxlint-disable-next-line typescript-eslint/no-explicit-any + globalThis.fetch = (async (url: any, init?: any) => { + calls.push({ + url: String(url), + // oxlint-disable-next-line typescript-eslint/no-unsafe-member-access + method: String(init?.method ?? 'GET'), + }); + n += 1; + // The POST is lost before any response. + if (n === 1) throw new TypeError('fetch failed'); + if (String(url).includes('/attach')) { + return sseResponse([ + `event: error\ndata: ${JSON.stringify({ message: 'exec exec-3 not found' })}\n\n`, + ]); + } + return sseResponse([ + `event: stdout\ndata: {"text":"OK","seq":1}\n\n`, + RESULT_OK, + ]); + // oxlint-disable-next-line typescript-eslint/no-explicit-any + }) as any; + + const stdout: string[] = []; + const result = await drainSessionExecResilient( + 'ses-3', + { execId: 'exec-3', command: ['x'], timeoutMs: 1_000 }, + new AbortController().signal, + { onStdout: (t) => stdout.push(t) }, + ); + + expect(result.status).toBe('completed'); + expect(stdout).toEqual(['OK']); + // POST (lost) → attach (not found) → POST again (created) — not a fifth + // identical attach. + expect(calls.map((c) => c.method)).toEqual(['POST', 'GET', 'POST']); + expect(calls[1]?.url).toContain('/attach'); + expect(calls[2]?.url).not.toContain('/attach'); + }, 15_000); + + test('never re-creates an exec that already produced output', async () => { + const methods: string[] = []; + let n = 0; + // oxlint-disable-next-line typescript-eslint/no-explicit-any + globalThis.fetch = (async (_url: any, init?: any) => { + // oxlint-disable-next-line typescript-eslint/no-unsafe-member-access + methods.push(String(init?.method ?? 'GET')); + n += 1; + // Progress (seq 2), then the stream drops without a result… + if (n === 1) { + return sseResponse([`event: stdout\ndata: {"text":"AB","seq":2}\n\n`]); + } + // …and the spawner has since lost the exec. Re-POSTing here would run + // the turn twice; the drain must fail instead. + return sseResponse([ + `event: error\ndata: ${JSON.stringify({ message: 'exec exec-4 not found' })}\n\n`, + ]); + // oxlint-disable-next-line typescript-eslint/no-explicit-any + }) as any; + + await expect( + drainSessionExecResilient( + 'ses-4', + { execId: 'exec-4', command: ['x'], timeoutMs: 1_000 }, + new AbortController().signal, + {}, + ), + ).rejects.toThrow(/not found/); + expect(methods.filter((m) => m === 'POST')).toHaveLength(1); + }, 15_000); +}); + describe('chunkStageFiles', () => { const stageFile = (path: string, contentBytes: number): SessionStageFile => ({ path, diff --git a/services/platform/backend/core/node_only/sandbox/helpers/session_client.ts b/services/platform/backend/core/node_only/sandbox/helpers/session_client.ts index ff58b68107..d82d910579 100644 --- a/services/platform/backend/core/node_only/sandbox/helpers/session_client.ts +++ b/services/platform/backend/core/node_only/sandbox/helpers/session_client.ts @@ -42,6 +42,17 @@ export class SessionNotFoundError extends Error { } } +/** The spawner answered an attach with its "exec not found" error + * event: the session is alive but knows no such exec. Distinguished so the + * resilient drain can tell "the exec was never created" (re-POST it) from a + * transient stream drop (re-attach). */ +export class ExecNotFoundError extends Error { + constructor(execId: string) { + super(`sandbox session exec ${execId} not found on the spawner`); + this.name = 'ExecNotFoundError'; + } +} + /** Spawner already owns a live session under this id (HTTP 409 on create). * With deterministic per-(org,user) ids this means an orphan the platform no * longer tracks (e.g. a destroy that raced provisioning) — callers reap it @@ -856,19 +867,30 @@ export async function drainSessionExecResilient( const startWithAttach = opts.resumeSinceSeq !== undefined; let attempt = 0; let seqAtAttemptStart = cursor.lastSeq; + // Whether the next attempt (re)POSTs the exec instead of attaching. A fresh + // turn POSTs once and every retry attaches — EXCEPT when the spawner has just + // answered that it knows no such exec and this drain has consumed nothing + // from it: then the original POST never landed (a network drop or a 503 + // mid-roll swallowed it), and attaching again would fail identically until + // the whole budget was gone. Re-POSTing creates it. Never after progress — + // an exec that produced output existed, and re-running it would execute + // the turn twice — and never on a resume, whose caller owns that recovery. + let recreate = !startWithAttach; for (;;) { try { - return startWithAttach || attempt > 0 - ? await sessionAttachExec( - sessionId, - body.execId, - cursor.lastSeq, - signal, - callbacks, - cursor, - body.timeoutMs, - ) - : await sessionExec(sessionId, body, signal, callbacks, cursor); + if (recreate) { + recreate = false; + return await sessionExec(sessionId, body, signal, callbacks, cursor); + } + return await sessionAttachExec( + sessionId, + body.execId, + cursor.lastSeq, + signal, + callbacks, + cursor, + body.timeoutMs, + ); } catch (err) { if (signal.aborted) throw err; // A 404 means the session is gone, not a transient drop — retrying can't @@ -889,8 +911,12 @@ export async function drainSessionExecResilient( seqAtAttemptStart = cursor.lastSeq; attempt += 1; if (attempt > MAX_RECONNECT_ATTEMPTS) throw err; + recreate = + !startWithAttach && + cursor.lastSeq === 0 && + err instanceof ExecNotFoundError; console.warn( - `[session_client] exec ${body.execId} stream dropped (attempt ${attempt}/${MAX_RECONNECT_ATTEMPTS}, sinceSeq=${cursor.lastSeq}); re-attaching:`, + `[session_client] exec ${body.execId} stream dropped (attempt ${attempt}/${MAX_RECONNECT_ATTEMPTS}, sinceSeq=${cursor.lastSeq}); ${recreate ? 're-creating the exec (the spawner does not know it)' : 're-attaching'}:`, err instanceof Error ? err.message : String(err), ); await new Promise((r) => @@ -936,7 +962,13 @@ async function consumeExecSse( if (parsed) result = parsed; } else if (event === 'error') { const parsed = parseData<{ message?: string }>(data); - throw new Error(parsed?.message ?? 'sandbox session exec stream error'); + const message = parsed?.message ?? 'sandbox session exec stream error'; + // The spawner's attach grammar for an unknown exec (session-routes.ts): + // `exec not found`. + if (message === `exec ${execId} not found`) { + throw new ExecNotFoundError(execId); + } + throw new Error(message); } }; for (;;) { diff --git a/services/platform/backend/core/node_only/sandbox/render_fetch.test.ts b/services/platform/backend/core/node_only/sandbox/render_fetch.test.ts index 8fd6f94b96..3a8bebf49e 100644 --- a/services/platform/backend/core/node_only/sandbox/render_fetch.test.ts +++ b/services/platform/backend/core/node_only/sandbox/render_fetch.test.ts @@ -1,6 +1,18 @@ -import { describe, expect, it } from 'vitest'; +import { execFile } from 'node:child_process'; +import { + mkdirSync, + mkdtempSync, + readFileSync, + rmSync, + writeFileSync, +} from 'node:fs'; +import { tmpdir } from 'node:os'; +import path from 'node:path'; +import { promisify } from 'node:util'; -import { parseRenderResults } from './render_fetch'; +import { afterEach, beforeEach, describe, expect, it } from 'vitest'; + +import { parseRenderResults, RENDER_WORKER_SOURCE } from './render_fetch'; /** * The worker↔engine protocol, pinned: what the crawl engine does with a page @@ -78,3 +90,137 @@ describe('parseRenderResults', () => { } }); }); + +/** + * The staged worker, run for real under node against a fake `playwright-core` + * that answers every navigation with the same multibyte page. What is pinned: + * the output file the host reads back is bounded in BYTES, decided before a + * page is admitted — a page that does not fit is handed back for the next + * batch instead of being written past the cap. + */ +const FAKE_PLAYWRIGHT = ` +const chars = Number(process.env.FAKE_HTML_CHARS || '100'); +// 'é' is one UTF-16 code unit but two UTF-8 bytes. +const html = '' + 'é'.repeat(chars) + ''; +function makePage() { + let current = ''; + return { + async goto(url) { current = url; return { status: () => 200 }; }, + async waitForLoadState() {}, + async evaluate() { return 42; }, + url() { return current; }, + async content() { return html; }, + async close() {}, + }; +} +module.exports = { + chromium: { + async launch() { + return { + async newContext() { return { async newPage() { return makePage(); } }; }, + async close() {}, + }; + }, + }, +}; +`; + +const NODE_BIN = path.basename(process.execPath).startsWith('node') + ? process.execPath + : 'node'; +const execFileAsync = promisify(execFile); + +describe('render worker — output budget in bytes', () => { + let root: string; + let agent: string; + + beforeEach(() => { + root = mkdtempSync(path.join(tmpdir(), 'render-worker-')); + agent = path.join(root, 'agent'); + const fakeDir = path.join(agent, 'code', 'node_modules', 'playwright-core'); + mkdirSync(fakeDir, { recursive: true }); + writeFileSync( + path.join(fakeDir, 'package.json'), + JSON.stringify({ name: 'playwright-core', main: 'index.js' }), + ); + writeFileSync(path.join(fakeDir, 'index.js'), FAKE_PLAYWRIGHT); + // The worker's paths are fixed to /agent inside the sandbox; point them + // at the temp root here. + writeFileSync( + path.join(agent, 'code', 'render.mjs'), + RENDER_WORKER_SOURCE.replaceAll("'/agent/", `'${agent}/`), + ); + }); + + afterEach(() => { + rmSync(root, { recursive: true, force: true }); + }); + + async function runWorker( + urls: readonly string[], + caps: { maxHtmlBytes: number; maxTotalBytes: number }, + htmlChars: number, + ): Promise<{ bytes: number; results: Map }> { + writeFileSync( + path.join(agent, 'code', 'urls.json'), + JSON.stringify({ + urls, + perPageTimeoutMs: 10, + idleTimeoutMs: 10, + softBudgetMs: 60_000, + ...caps, + }), + ); + await execFileAsync(NODE_BIN, [path.join(agent, 'code', 'render.mjs')], { + env: { ...process.env, FAKE_HTML_CHARS: String(htmlChars) }, + timeout: 25_000, + }); + const raw = readFileSync(path.join(agent, 'output', 'pages.json')); + const payload: unknown = JSON.parse(raw.toString('utf8')); + return { + bytes: raw.byteLength, + results: parseRenderResults(payload, urls), + }; + } + + // Regression: the batch total was `html.length` summed AFTER storing each + // page and checked only before the NEXT one, so pages.json could exceed the + // host's read cap (and by more with multibyte text) — the host then saw no + // output and the crawl retried the same batch forever. + it('keeps pages.json under maxTotalBytes and hands back the page that would not fit', async () => { + const urls = ['a', 'b', 'c', 'd'].map((p) => `https://site.example/${p}`); + // 1700 chars = 3400 UTF-8 bytes per page: two fit under 10 000, the third + // would not. + const { bytes, results } = await runWorker( + urls, + { maxHtmlBytes: 4_000, maxTotalBytes: 10_000 }, + 1_700, + ); + expect(bytes).toBeLessThanOrEqual(10_000); + expect(results.get(urls[0] ?? '')).toMatchObject({ + kind: 'ok', + status: 200, + }); + expect(results.get(urls[1] ?? '')).toMatchObject({ + kind: 'ok', + status: 200, + }); + // Not written past the cap, not charged as a failure: due next batch. + expect(results.get(urls[2] ?? '')).toEqual({ kind: 'not_attempted' }); + expect(results.get(urls[3] ?? '')).toEqual({ kind: 'not_attempted' }); + }, 30_000); + + it('applies the per-page bound in bytes, not UTF-16 code units', async () => { + const urls = ['https://site.example/big']; + // 1700 code units pass a 3000 "length" check but are 3400 bytes. + const { results } = await runWorker( + urls, + { maxHtmlBytes: 3_000, maxTotalBytes: 100_000 }, + 1_700, + ); + expect(results.get(urls[0] ?? '')).toEqual({ + kind: 'failed', + reason: 'rendered HTML exceeds the per-page bound', + }); + }, 30_000); +}); diff --git a/services/platform/backend/core/node_only/sandbox/render_fetch.ts b/services/platform/backend/core/node_only/sandbox/render_fetch.ts index f29314371b..c8298c1cc5 100644 --- a/services/platform/backend/core/node_only/sandbox/render_fetch.ts +++ b/services/platform/backend/core/node_only/sandbox/render_fetch.ts @@ -44,8 +44,11 @@ const RENDER_IDLE_TIMEOUT_MS = 5_000; /** The worker stops STARTING pages this far before the exec hard kill, so * it always exits cleanly with its partial results on disk. */ const WORKER_EXIT_MARGIN_MS = 20_000; -/** Rendered-DOM caps: per page, and total per batch — `sessionReadFile` - * serves at most 20MB, so the output file must stay safely under it. */ +/** Rendered-DOM caps in BYTES: per page, and for the whole output file — + * `sessionReadFile` serves at most 20MB, and an output the host cannot read + * back fails the batch deterministically (the crawl would retry the same + * batch forever), so the worker bounds the serialized file before it admits + * each page. */ const RENDER_MAX_HTML_BYTES = 6 * 1024 * 1024; const RENDER_MAX_TOTAL_BYTES = 15 * 1024 * 1024; @@ -245,7 +248,7 @@ export async function renderUrlsInSandbox( * aliases like `convex`) and private/link-local IP literals are refused, on * the original URL and again on the post-redirect landing host. */ -const RENDER_WORKER_SOURCE = ` +export const RENDER_WORKER_SOURCE = ` import { createRequire } from 'node:module'; import { mkdirSync, readFileSync, writeFileSync } from 'node:fs'; @@ -302,13 +305,23 @@ const startedAt = Date.now(); const records = new Map(); for (const url of urls) records.set(url, { url, attempted: false }); mkdirSync('/agent/output', { recursive: true }); +// Rewrites the whole output after every page and returns its size in BYTES — +// the budget is what the host reads back (UTF-8, JSON-escaped), never a count +// of UTF-16 code units. function flush() { - writeFileSync( - '/agent/output/pages.json', - JSON.stringify({ pages: Array.from(records.values()) }), + const serialized = JSON.stringify({ pages: Array.from(records.values()) }); + writeFileSync('/agent/output/pages.json', serialized); + return Buffer.byteLength(serialized, 'utf8'); +} +// The bytes a page would add to the output: its record serialized with the +// rendered fields, minus the record as it stands (escaping included). +function admissionBytes(record, fields) { + return ( + Buffer.byteLength(JSON.stringify({ ...record, ...fields }), 'utf8') - + Buffer.byteLength(JSON.stringify(record), 'utf8') ); } -flush(); +let fileBytes = flush(); const chromium = loadChromium(); const proxyServer = @@ -322,12 +335,11 @@ if (proxyServer) { } const browser = await chromium.launch(launchOptions); -let totalBytes = 0; +let budgetExhausted = false; try { const context = await browser.newContext(); for (const url of urls) { if (Date.now() - startedAt > softBudgetMs) break; - if (totalBytes > maxTotalBytes) break; const record = records.get(url); record.attempted = true; let hostname = ''; @@ -384,13 +396,20 @@ try { record.error = 'HTTP ' + status + ' at render time'; } else { const html = await page.content(); - if (html.length > maxHtmlBytes) { + if (Buffer.byteLength(html, 'utf8') > maxHtmlBytes) { record.error = 'rendered HTML exceeds the per-page bound'; } else { - record.status = status; - record.finalUrl = finalUrl; - record.html = html; - totalBytes += html.length; + const fields = { status, finalUrl, html }; + // Admit the page only while the output stays under the batch cap — + // decided BEFORE storing, in bytes. A page that does not fit is + // handed back (not attempted) for the next batch, and this batch + // ends here. + if (fileBytes + admissionBytes(record, fields) > maxTotalBytes) { + record.attempted = false; + budgetExhausted = true; + } else { + Object.assign(record, fields); + } } } } catch (error) { @@ -398,7 +417,8 @@ try { } finally { await page.close().catch(() => {}); } - flush(); + fileBytes = flush(); + if (budgetExhausted) break; } } finally { await browser.close().catch(() => {}); diff --git a/services/platform/backend/domains/deployment/service.ts b/services/platform/backend/domains/deployment/service.ts index 14d308f38d..1eebf21665 100644 --- a/services/platform/backend/domains/deployment/service.ts +++ b/services/platform/backend/domains/deployment/service.ts @@ -358,7 +358,7 @@ export async function saveDeploymentSecret( } const content = hasSopsKey() - ? encryptJsonWithSops(prepared.plaintext) + ? await encryptJsonWithSops(prepared.plaintext) : prepared.plaintext; await atomicWriteSecret(secretsPath, content); invalidateSecretsCache(secretsPath); diff --git a/services/platform/backend/domains/knowledge/admin.ts b/services/platform/backend/domains/knowledge/admin.ts index 3f58548254..4e2dc27b27 100644 --- a/services/platform/backend/domains/knowledge/admin.ts +++ b/services/platform/backend/domains/knowledge/admin.ts @@ -212,7 +212,9 @@ export async function writeKnowledgeConnection( await removeFileSafe(secretsPath); } else { const plaintext = serializeSecretsJson({ password: args.password }); - const content = hasSopsKey() ? encryptJsonWithSops(plaintext) : plaintext; + const content = hasSopsKey() + ? await encryptJsonWithSops(plaintext) + : plaintext; await atomicWriteSecret(secretsPath, content); } invalidateSecretsCache(secretsPath); diff --git a/services/platform/backend/domains/knowledge/service.image-status.test.ts b/services/platform/backend/domains/knowledge/service.image-status.test.ts new file mode 100644 index 0000000000..b8af82b96c --- /dev/null +++ b/services/platform/backend/domains/knowledge/service.image-status.test.ts @@ -0,0 +1,76 @@ +/** + * An uploaded image has no text extractor today (the vision seam is retired), + * so its indexing outcome must be the honest, terminal 'unsupported' — never + * 'failed', whose badge offers a retry that can only fail again. + */ + +import type { Sql } from 'postgres'; +import { describe, expect, it } from 'vitest'; + +import { indexUploadedFile } from './service.ts'; + +interface Query { + text: string; + values: unknown[]; +} + +function fakeSql(fileName: string, log: Query[]): Sql { + const tag = (strings: TemplateStringsArray, ...values: unknown[]) => { + const text = strings.join('$'); + log.push({ text, values }); + if (text.includes('FROM app.file_metadata WHERE id')) { + return Promise.resolve([ + { + organizationId: 'org-1', + storageRef: 's3:org-1/blob-1', + fileName, + contentType: 'application/octet-stream', + documentId: null, + skipRagIndexing: null, + }, + ]); + } + if (text.includes('FROM "organization"')) { + return Promise.resolve([{ slug: 'acme' }]); + } + if (text.includes('UPDATE app.file_metadata')) { + return Promise.resolve([{ orgId: 'org-1' }]); + } + return Promise.resolve([]); + }; + // oxlint-disable-next-line typescript/no-unsafe-type-assertion -- test double + return Object.assign(tag, { unsafe: (t: string) => t }) as unknown as Sql; +} + +const statusWrites = (log: Query[]): unknown[][] => + log + .filter((q) => q.text.includes('UPDATE app.file_metadata')) + .map((q) => q.values); + +describe('indexUploadedFile — images', () => { + it("marks an image 'unsupported' up front, never 'failed'", async () => { + const log: Query[] = []; + await indexUploadedFile(fakeSql('photo.png', log), 'file-1'); + + const writes = statusWrites(log); + expect(writes).toHaveLength(1); + expect(writes[0]).toContain('unsupported'); + expect(JSON.stringify(writes[0])).toMatch(/vision/i); + expect(JSON.stringify(writes)).not.toContain('failed'); + // Decided from the name alone: no 'running' pass, no blob fetch. + expect(JSON.stringify(writes)).not.toContain('Extracting text'); + }); + + it('still sends a text document down the extraction path', async () => { + const log: Query[] = []; + // Past the status write this reaches the (unfaked) object store; the + // outcome of that is not under test here — only that the image branch + // did not swallow a non-image. + await indexUploadedFile(fakeSql('notes.txt', log), 'file-2').catch( + () => undefined, + ); + const first = statusWrites(log)[0]; + expect(first).toContain('running'); + expect(first).toContain('Extracting text…'); + }); +}); diff --git a/services/platform/backend/domains/knowledge/service.ts b/services/platform/backend/domains/knowledge/service.ts index 5a7cfaac15..964043ef40 100644 --- a/services/platform/backend/domains/knowledge/service.ts +++ b/services/platform/backend/domains/knowledge/service.ts @@ -27,6 +27,7 @@ import { } from '../../core/knowledge/search.ts'; import { extractText, + isImageFile, isSupported, } from '../../core/lib/knowledge/extraction/router.ts'; import { parseBlobRef } from '../../core/lib/storage/blob_ref.ts'; @@ -394,6 +395,22 @@ export async function indexUploadedFile( }); return; } + // Images route to the vision extractor, and the vision seam is retired + // (`extraction/vision_client.ts`): `extractText` below is called with no + // vision client, so an image can only ever yield '' — which used to land as + // 'failed — Indexing skipped (empty)', a badge that invites the user to + // retry a capability that does not exist. An image with nothing to index is + // the honest, terminal 'unsupported', decided before any bytes are fetched. + // Drop this branch when the vision lane returns and a client is wired in. + if (isImageFile(file.fileName)) { + await writeRagStatus(sql, fileId, { + ragStatus: 'unsupported', + ragError: + `Images cannot be indexed for search: no vision (OCR) model lane is ` + + `available to read "${file.fileName}".`, + }); + return; + } await writeRagStatus(sql, fileId, { ragStatus: 'running', diff --git a/services/platform/backend/domains/object_storage/bootstrap.ts b/services/platform/backend/domains/object_storage/bootstrap.ts index ceecf32cf1..939a48618c 100644 --- a/services/platform/backend/domains/object_storage/bootstrap.ts +++ b/services/platform/backend/domains/object_storage/bootstrap.ts @@ -150,7 +150,7 @@ export async function ensureDefaultObjectStore( }); await atomicWriteSecret( secretsPath, - hasSopsKey() ? encryptJsonWithSops(plaintext) : plaintext, + hasSopsKey() ? await encryptJsonWithSops(plaintext) : plaintext, ); invalidateSecretsCache(secretsPath); clearObjectStoreCache(); diff --git a/services/platform/backend/domains/object_storage/service.ts b/services/platform/backend/domains/object_storage/service.ts index 47c1a417bc..6e322df495 100644 --- a/services/platform/backend/domains/object_storage/service.ts +++ b/services/platform/backend/domains/object_storage/service.ts @@ -180,7 +180,9 @@ export async function writeConnection( accessKeyId: args.accessKeyId, secretAccessKey: args.secretAccessKey, }); - const content = hasSopsKey() ? encryptJsonWithSops(plaintext) : plaintext; + const content = hasSopsKey() + ? await encryptJsonWithSops(plaintext) + : plaintext; await atomicWriteSecret(secretsPath, content); invalidateSecretsCache(secretsPath); } diff --git a/services/platform/backend/domains/onedrive/service.claim-heartbeat.test.ts b/services/platform/backend/domains/onedrive/service.claim-heartbeat.test.ts new file mode 100644 index 0000000000..1c4cffca60 --- /dev/null +++ b/services/platform/backend/domains/onedrive/service.claim-heartbeat.test.ts @@ -0,0 +1,156 @@ +/** + * The per-config sync claim: a 'running' stamp older than the stale window is + * treated as a crashed worker and re-claimed. A LIVE run must therefore keep + * its stamp fresh for as long as it runs — otherwise a sync that merely takes + * longer than the window is claimed again by the next cron tick and runs + * twice concurrently. + */ + +import type { Sql } from 'postgres'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; + +import { runSyncConfigJobWith, type SyncProviderAdapter } from './service.ts'; + +interface Query { + text: string; + values: unknown[]; +} + +const CONFIG_ROW = { + id: 'cfg-1', + organizationId: 'org-1', + userId: 'user-1', + itemType: 'folder', + itemId: 'folder-1', + itemName: 'Reports', + itemPath: null, + targetBucket: 'documents', + storagePrefix: null, + teamId: null, + status: 'active', + lastSyncAt: null, + lastSyncStatus: null, + errorMessage: null, +}; + +function fakeSql(log: Query[]): Sql { + const tag = (strings: TemplateStringsArray, ...values: unknown[]) => { + const text = strings.join('$'); + log.push({ text, values }); + // The claim fence. + if (text.includes("last_sync_status = 'running', updated_at_ms")) { + return Promise.resolve([{ id: 'cfg-1' }]); + } + // getSyncConfigRow. + if (text.includes('WHERE id = $ LIMIT 1')) { + return Promise.resolve([CONFIG_ROW]); + } + return Promise.resolve([]); + }; + // oxlint-disable-next-line typescript/no-unsafe-type-assertion -- test double + return Object.assign(tag, { unsafe: (t: string) => t }) as unknown as Sql; +} + +const heartbeats = (log: Query[]): number => + log.filter( + (q) => + q.text.includes('SET updated_at_ms = $') && + q.text.includes("last_sync_status = 'running'"), + ).length; + +const finalStamps = (log: Query[]): number => + log.filter((q) => q.text.includes('last_sync_at_ms = $')).length; + +/** An adapter whose token resolve blocks until the test releases it — the + * "long sync" — then fails so the run ends without touching the reconcile + * lane. */ +function blockingAdapter(gate: Promise): SyncProviderAdapter { + return { + displayName: 'Fake Drive', + sourceProvider: 'fake', + configTable: 'app.onedrive_sync_configs', + configJobName: 'onedrive.sync_config', + singletonPrefix: 'fake-sync-', + metadataItemIdKeys: [], + resolveToken: async () => { + await gate; + return { success: false }; + }, + listFolderContents: () => Promise.resolve({ success: false }), + getFileMetadata: () => Promise.resolve({ success: false }), + buildDownloadUrl: () => '', + runImport: () => + Promise.resolve({ + results: [], + successCount: 0, + skippedCount: 0, + failedCount: 0, + }), + // oxlint-disable-next-line typescript/no-unsafe-type-assertion -- only the fields the claim path touches are exercised + } as unknown as SyncProviderAdapter; +} + +describe('runSyncConfigJobWith — claim heartbeat', () => { + beforeEach(() => { + vi.useFakeTimers(); + }); + afterEach(() => { + vi.useRealTimers(); + }); + + it('renews the claim stamp while the sync runs and stops when it ends', async () => { + const log: Query[] = []; + let release: () => void = () => undefined; + const gate = new Promise((resolve) => { + release = resolve; + }); + + const run = runSyncConfigJobWith( + fakeSql(log), + blockingAdapter(gate), + { organizationId: 'org-1', configId: 'cfg-1' }, + { heartbeatMs: 1_000 }, + ); + + // Three heartbeat periods into a still-running sync: three renewals, no + // outcome stamp yet. + await vi.advanceTimersByTimeAsync(3_500); + expect(heartbeats(log)).toBe(3); + expect(finalStamps(log)).toBe(0); + const renewal = log.find((q) => q.text.includes('SET updated_at_ms = $')); + expect(renewal?.values).toContain('cfg-1'); + expect(renewal?.values).toContain('org-1'); + + release(); + await run; + expect(finalStamps(log)).toBe(1); + + // The interval is cleared with the run: no renewals after the outcome. + await vi.advanceTimersByTimeAsync(5_000); + expect(heartbeats(log)).toBe(3); + }); + + it('does not heartbeat when the claim was not won', async () => { + const log: Query[] = []; + const sql = fakeSql(log); + // A losing claim answers no row. + const losing = Object.assign( + (strings: TemplateStringsArray, ...values: unknown[]) => { + const text = strings.join('$'); + log.push({ text, values }); + return Promise.resolve([]); + }, + { unsafe: (t: string) => t }, + ); + void sql; + await runSyncConfigJobWith( + // oxlint-disable-next-line typescript/no-unsafe-type-assertion -- test double + losing as unknown as Sql, + blockingAdapter(Promise.resolve()), + { organizationId: 'org-1', configId: 'cfg-1' }, + { heartbeatMs: 1_000 }, + ); + await vi.advanceTimersByTimeAsync(5_000); + expect(heartbeats(log)).toBe(0); + }); +}); diff --git a/services/platform/backend/domains/onedrive/service.ts b/services/platform/backend/domains/onedrive/service.ts index 4eab77e077..9642995ef1 100644 --- a/services/platform/backend/domains/onedrive/service.ts +++ b/services/platform/backend/domains/onedrive/service.ts @@ -1371,7 +1371,24 @@ export async function syncOneConfigWith( // ------------------------------------------------------------------- engine /** A run older than this may be re-claimed (crashed worker recovery). */ -const SYNC_CLAIM_STALE_MS = 30 * 60 * 1000; +export const SYNC_CLAIM_STALE_MS = 30 * 60 * 1000; +/** A live run refreshes its claim this often — well inside the stale window, + * so only a run whose process died (no heartbeat) ever reads as stale. */ +export const SYNC_CLAIM_HEARTBEAT_MS = 5 * 60 * 1000; + +/** Refresh a live run's claim stamp; a no-op once the run has stamped its + * outcome (the row is no longer 'running'). */ +async function renewSyncClaim( + sql: Sql, + table: string, + payload: { organizationId: string; configId: string }, +): Promise { + await sql` + UPDATE ${sql.unsafe(table)} SET updated_at_ms = ${Date.now()} + WHERE id = ${payload.configId} AND org_id = ${payload.organizationId} + AND last_sync_status = 'running' + `; +} /** * Cron scan: one per-config job per syncable config. `error` configs are @@ -1400,11 +1417,16 @@ export async function runSyncScanWith( return rows.length; } -/** One per-config sync job: claim, reconcile, stamp the outcome. */ +/** + * One per-config sync job: claim, reconcile, stamp the outcome. The claim is + * kept fresh by a heartbeat for as long as the run is alive, so the stale + * window below only ever re-admits a run whose worker actually died. + */ export async function runSyncConfigJobWith( sql: Sql, adapter: SyncProviderAdapter, payload: { organizationId: string; configId: string }, + opts: { heartbeatMs?: number } = {}, ): Promise { // Claim fence: a second job for the same config no-ops while a fresh run // is in flight; a stale 'running' stamp (crashed worker) is reclaimable. @@ -1425,6 +1447,23 @@ export async function runSyncConfigJobWith( ); if (!config) return; + // Heartbeat: the fence treats a 'running' stamp older than + // SYNC_CLAIM_STALE_MS as a crashed worker. Nothing used to refresh the + // stamp during a run, so a folder sync that merely took longer than that + // (large folder, slow tenant) was re-claimed by the next cron tick and ran + // twice concurrently — racing createDocument into duplicate documents. + const heartbeat = setInterval(() => { + renewSyncClaim(sql, adapter.configTable, payload).catch( + (error: unknown) => { + console.warn( + `[${adapter.displayName} sync] claim heartbeat failed for config ${payload.configId}:`, + error instanceof Error ? error.message : error, + ); + }, + ); + }, opts.heartbeatMs ?? SYNC_CLAIM_HEARTBEAT_MS); + heartbeat.unref(); + try { const result = await syncOneConfigWith(sql, adapter, config); if (result.sourceDeleted === true) { @@ -1458,6 +1497,8 @@ export async function runSyncConfigJobWith( lastSyncStatus: 'error', errorMessage: error instanceof Error ? error.message : String(error), }); + } finally { + clearInterval(heartbeat); } } diff --git a/services/platform/backend/domains/tts/routes.ts b/services/platform/backend/domains/tts/routes.ts index de915f4397..43ea5b9ce2 100644 --- a/services/platform/backend/domains/tts/routes.ts +++ b/services/platform/backend/domains/tts/routes.ts @@ -11,6 +11,7 @@ import { createCtxShim } from '../../lib/ctx-shim.ts'; import { resolveObjectStore, s3PresignGetUrl } from '../../lib/object-store.ts'; import { resolveOrgSlug } from '../../lib/org-config.ts'; import { chatShimHandlers } from '../chat/shim.ts'; +import { loadOwnedThread } from '../chat/threads.ts'; import { getChunkForServe, getMessageChunks, @@ -95,11 +96,19 @@ export function createTtsRoutes(deps: { sql: Sql; auth: Auth }): Hono { return c.json({ chunks }); }); - // The info dialog's per-message voice spend (the 0.4 shape). + // The info dialog's per-message voice spend (the 0.4 shape). Same gate as + // the listing and the audio door: the caller must own the thread. app.get('/messages/:messageId/usage', async (c) => { const threadId = c.req.query('threadId') ?? ''; if (threadId === '') return c.json({ error: 'threadId required' }, 400); const { organizationId, userId } = caller(c); + const thread = await loadOwnedThread( + deps.sql, + organizationId, + userId, + threadId, + ); + if (!thread) return c.json(null); const rows = await deps.sql< { provider: string; @@ -115,11 +124,9 @@ export function createTtsRoutes(deps: { sql: Sql; auth: Auth }): Hono { sum(coalesce(c.cost_estimate_cents, 0))::float8 AS "costCents", count(*)::float8 AS "chunkCount" FROM app.tts_audio_chunks c - JOIN app.threads t ON t.id = c.thread_id WHERE c.message_id = ${c.req.param('messageId')} AND c.thread_id = ${threadId} AND c.org_id = ${organizationId} - AND t.user_id = ${userId} AND c.status = 'ready' GROUP BY c.provider_name, c.model_id, c.voice `; @@ -207,12 +214,13 @@ export function createTtsRoutes(deps: { sql: Sql; auth: Auth }): Hono { } }); - /** Stream one ready chunk's audio. Membership is the session gate; the + /** Stream one ready chunk's audio. The gate is ownership of the chunk's + * thread (the same `loadOwnedThread` the listing and usage doors use); the * bytes are fetched server-side so no replayable URL ever reaches the * client. */ app.get('/audio/:chunkId', async (c) => { const chunk = await getChunkForServe(deps.sql, { - organizationId: c.get('orgId'), + ...caller(c), chunkId: c.req.param('chunkId'), }); if (!chunk) return c.json({ error: 'not found' }, 404); diff --git a/services/platform/backend/domains/tts/service.access-gate.test.ts b/services/platform/backend/domains/tts/service.access-gate.test.ts new file mode 100644 index 0000000000..147dc32899 --- /dev/null +++ b/services/platform/backend/domains/tts/service.access-gate.test.ts @@ -0,0 +1,161 @@ +/** + * The two TTS doors that used to diverge from the thread-ownership gate: + * synthesis accepted any messageId under the caller's own thread (a squatter + * could mint the global `(message_id, chunk_index)` reservation for someone + * else's message), and the audio serve checked org membership only. Both now + * go through the same `loadOwnedThread` rule as the listing and usage doors. + */ + +import type { Sql } from 'postgres'; +import { describe, expect, it } from 'vitest'; + +import { getChunkForServe, synthesizeChunk, TtsError } from './service.ts'; + +type Answer = (text: string, values: unknown[]) => unknown[]; + +/** A `sql` stand-in dispatching on the query text; `begin` runs the callback + * against the same stand-in so audit writes inside a transaction are seen. */ +function fakeSql(answer: Answer, log: { text: string; values: unknown[] }[]) { + const tag = (strings: TemplateStringsArray, ...values: unknown[]) => { + const text = strings.join('$'); + log.push({ text, values }); + return Promise.resolve(answer(text, values)); + }; + const api = { + unsafe: (text: string) => text, + json: (value: unknown) => value, + begin: (fn: (tx: unknown) => Promise) => fn(sql), + }; + const sql = Object.assign(tag, api); + // oxlint-disable-next-line typescript/no-unsafe-type-assertion -- test double + return sql as unknown as Sql; +} + +const ORG = 'org-1'; +const OWNER = 'user-owner'; +const STRANGER = 'user-stranger'; +const THREAD = 'thr-1'; + +/** Answers the thread-ownership read for `owner` only, the message-in-thread + * read per `messageInThread`, and the audit chain writes. */ +function answers(opts: { owner: string; messageInThread: boolean }): Answer { + return (text, values) => { + if (text.includes('FROM app.threads t')) { + return values.includes(opts.owner) ? [{ id: THREAD }] : []; + } + if (text.includes('FROM app.messages')) { + return opts.messageInThread ? [{ one: 1 }] : []; + } + if (text.includes('FROM app.audit_chain_heads')) { + return [{ lastHash: '', lastTs: 0 }]; + } + if (text.includes('INSERT INTO app.audit_logs')) { + return [{ id: 'audit-1' }]; + } + if (text.includes('FROM app.tts_audio_chunks')) { + return [ + { + threadId: THREAD, + storageRef: 's3:org-1/chunk', + status: 'ready', + format: 'mp3', + }, + ]; + } + return []; + }; +} + +describe('synthesizeChunk — the message must belong to the thread', () => { + it('refuses a messageId that is not a row of the (owned) thread, and audits it', async () => { + const log: { text: string; values: unknown[] }[] = []; + const sql = fakeSql(answers({ owner: OWNER, messageInThread: false }), log); + + const outcome = await synthesizeChunk(sql, { + organizationId: ORG, + userId: OWNER, + messageId: 'msg-from-someone-elses-thread', + threadId: THREAD, + index: 0, + text: 'Hello.', + locale: 'en', + }).then( + () => null, + (err: unknown) => err, + ); + + expect(outcome).toBeInstanceOf(TtsError); + if (outcome instanceof TtsError) { + expect(outcome.code).toBe('FORBIDDEN'); + expect(outcome.status).toBe(403); + } + const audit = log.find((q) => + q.text.includes('INSERT INTO app.audit_logs'), + ); + expect(audit?.values).toContain('tts.synthesize_denied'); + expect(JSON.stringify(audit?.values)).toContain('message_not_in_thread'); + // Refused BEFORE the model resolve / reservation: no chunk row is ever + // minted under the squatter's thread. + expect(log.some((q) => q.text.includes('FROM app.thread_metadata'))).toBe( + false, + ); + expect( + log.some((q) => q.text.includes('INSERT INTO app.tts_audio_chunks')), + ).toBe(false); + }); + + it('lets a message of the owned thread through the gate', async () => { + const log: { text: string; values: unknown[] }[] = []; + const sql = fakeSql(answers({ owner: OWNER, messageInThread: true }), log); + + const outcome = await synthesizeChunk(sql, { + organizationId: ORG, + userId: OWNER, + messageId: 'msg-1', + threadId: THREAD, + index: 0, + text: 'Hello.', + locale: 'en', + }).then( + () => null, + (err: unknown) => err, + ); + + // Whatever the (unfaked) model resolver does afterwards, the gate itself + // passed: the next read is the thread's agent binding. + expect(outcome instanceof TtsError && outcome.code === 'FORBIDDEN').toBe( + false, + ); + expect(log.some((q) => q.text.includes('FROM app.thread_metadata'))).toBe( + true, + ); + }); +}); + +describe('getChunkForServe — ownership gate, not org membership', () => { + it('serves a ready chunk to the owner of its thread', async () => { + const log: { text: string; values: unknown[] }[] = []; + const sql = fakeSql(answers({ owner: OWNER, messageInThread: true }), log); + expect( + await getChunkForServe(sql, { + organizationId: ORG, + userId: OWNER, + chunkId: 'chunk-1', + }), + ).toEqual({ storageRef: 's3:org-1/chunk', contentType: 'audio/mpeg' }); + }); + + it('answers null to another member of the org who holds the chunk id', async () => { + const log: { text: string; values: unknown[] }[] = []; + const sql = fakeSql(answers({ owner: OWNER, messageInThread: true }), log); + expect( + await getChunkForServe(sql, { + organizationId: ORG, + userId: STRANGER, + chunkId: 'chunk-1', + }), + ).toBeNull(); + // The gate consulted is the shared thread-ownership read. + expect(log.some((q) => q.text.includes('FROM app.threads t'))).toBe(true); + }); +}); diff --git a/services/platform/backend/domains/tts/service.ts b/services/platform/backend/domains/tts/service.ts index baf1c2dc52..5da06f9ee3 100644 --- a/services/platform/backend/domains/tts/service.ts +++ b/services/platform/backend/domains/tts/service.ts @@ -699,6 +699,25 @@ async function markChunkReadyAndRecordUsage( // ------------------------------------------------------------- synthesize +/** + * Is `messageId` a row of `threadId` (in this org)? Ownership of the thread is + * checked by the caller (`loadOwnedThread`); this closes the remaining gap — + * a message id is not secret (shared threads reveal them), so the thread it + * is synthesized under must be the one it actually belongs to. + */ +async function messageBelongsToThread( + sql: Sql, + args: { organizationId: string; threadId: string; messageId: string }, +): Promise { + const rows = await sql<{ one: number }[]>` + SELECT 1 AS one FROM app.messages + WHERE id = ${args.messageId} AND thread_id = ${args.threadId} + AND org_id = ${args.organizationId} + LIMIT 1 + `; + return rows.length > 0; +} + export interface SynthesizeResult { status: 'ready' | 'in-flight' | 'failed'; errorCode?: string; @@ -736,6 +755,34 @@ export async function synthesizeChunk( if (!thread) { throw new TtsError('FORBIDDEN', 'This conversation does not exist.', 403); } + if (!(await messageBelongsToThread(sql, args))) { + // The caller owns the thread but names a message that is not in it — a + // messageId learnt from another (shared, foreign) conversation. The + // `(message_id, chunk_index)` reservation is global, so admitting this + // would mint the row under the squatter's thread and lock the message's + // real owner out of playback with a 403 and an audit row blaming THEM. + await sql.begin(async (tx) => { + await createAuditLog(tx, { + organizationId: args.organizationId, + actorId: args.userId, + actorType: 'user', + action: 'tts.synthesize_denied', + category: 'security', + resourceType: 'tts_audio_chunk', + metadata: { + reason: 'message_not_in_thread', + requestedMessageId: args.messageId, + requestedThreadId: args.threadId, + }, + status: 'denied', + }); + }); + throw new TtsError( + 'FORBIDDEN', + 'This message does not belong to the conversation.', + 403, + ); + } const meta = await sql<{ agentSlug: string | null }[]>` SELECT agent_slug AS "agentSlug" FROM app.thread_metadata WHERE thread_id = ${args.threadId} LIMIT 1 @@ -952,22 +999,39 @@ export async function getMessageChunks( }); } -/** The audio-serve read: membership is the route's (session) gate; this - * enforces the chunk's own org and readiness. */ +/** + * The audio-serve read. A chunk is spoken content of a PRIVATE conversation, + * so the gate is the one every other TTS door uses — ownership of the chunk's + * thread via `loadOwnedThread` — never org membership alone (a chunk id can + * leak through logs or referers, and unguessability is not authorization). + * Enforces readiness too; `null` covers missing, foreign, and not-ready alike. + */ export async function getChunkForServe( sql: Sql, - args: { organizationId: string; chunkId: string }, + args: { organizationId: string; userId: string; chunkId: string }, ): Promise<{ storageRef: string; contentType: string } | null> { const rows = await sql< - { storageRef: string | null; status: string; format: string | null }[] + { + threadId: string; + storageRef: string | null; + status: string; + format: string | null; + }[] >` - SELECT storage_ref AS "storageRef", status, format + SELECT thread_id AS "threadId", storage_ref AS "storageRef", status, format FROM app.tts_audio_chunks WHERE id = ${args.chunkId} AND org_id = ${args.organizationId} LIMIT 1 `; const row = rows[0]; if (!row || row.status !== 'ready' || row.storageRef === null) return null; + const thread = await loadOwnedThread( + sql, + args.organizationId, + args.userId, + row.threadId, + ); + if (!thread) return null; const format = row.format; const mime = format !== null diff --git a/services/platform/backend/integration-check.ts b/services/platform/backend/integration-check.ts index 6ea25e1085..83d188a03f 100644 --- a/services/platform/backend/integration-check.ts +++ b/services/platform/backend/integration-check.ts @@ -5791,6 +5791,75 @@ async function checkKnowledge( `indexed=${indexed} (status=${statusRows[0]?.status}${statusRows[0]?.error ? `, err=${statusRows[0].error.slice(0, 80)}` : ''}), hits=${search.success ? search.data.hits.length : 'ERR'}, searchHit=${searchRaw.includes('verdigris')}, fetchHit=${fetchRaw.includes('zeppelin ledger')}, documentHints=${ragHints[0]?.count ?? '0'} (want >= 2)`, ); + // An image has no text extractor today (the vision seam is retired), so + // its indexing outcome is the honest, terminal 'unsupported' — never + // 'failed' with a retry affordance that can only fail again. + const PNG_1X1 = Buffer.from( + 'iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAYAAAAfFcSJAAAADUlEQVR42mNkYPhfDwAChwGA60e6kgAAAABJRU5ErkJggg==', + 'base64', + ); + const imageHandoff = z + .object({ storageRef: z.string(), uploadUrl: z.string().url() }) + .safeParse( + await ( + await send('POST', `/api/app/files/upload-handoff?orgId=${orgId}`, { + contentType: 'image/png', + size: PNG_1X1.byteLength, + }) + ).json(), + ); + let imageStatus: + | { status: string | null; error: string | null } + | undefined; + if (imageHandoff.success) { + await fetch(imageHandoff.data.uploadUrl, { + method: 'PUT', + headers: { 'content-type': 'image/png' }, + body: PNG_1X1, + }); + const imageRegistered = z.object({ fileId: z.string() }).safeParse( + await ( + await send('POST', `/api/app/files/register?orgId=${orgId}`, { + storageRef: imageHandoff.data.storageRef, + fileName: 'diagram.png', + contentType: 'image/png', + }) + ).json(), + ); + const imageFileId = imageRegistered.success + ? imageRegistered.data.fileId + : ''; + await send('POST', `/api/app/documents/from-upload?orgId=${orgId}`, { + fileId: imageFileId, + fileName: 'diagram.png', + }); + await waitFor(async () => { + const rows = await sql<{ status: string | null }[]>` + SELECT rag_status AS status FROM app.file_metadata + WHERE id = ${imageFileId} + `; + const status = rows[0]?.status; + return ( + status !== null && + status !== undefined && + status !== 'queued' && + status !== 'running' + ); + }, 20_000); + const rows = await sql<{ status: string | null; error: string | null }[]>` + SELECT rag_status AS status, rag_error AS error FROM app.file_metadata + WHERE id = ${imageFileId} + `; + imageStatus = rows[0]; + } + record( + "knowledge image upload lands 'unsupported', never 'failed'", + imageHandoff.success && + imageStatus?.status === 'unsupported' && + (imageStatus.error ?? '').includes('vision'), + `handoff=${imageHandoff.success}, status=${imageStatus?.status ?? 'none'} (want unsupported), error=${(imageStatus?.error ?? '').slice(0, 90)}`, + ); + // A document larger than one slice (64 chunks) reaches `completed` only // after EVERY slice lands — the regression lock for the port that // stamped completed after slice one and left the corpus 'processing' @@ -19901,6 +19970,23 @@ async function checkTts( ).json(), ); const threadId = created.success ? created.data.id : ''; + // The synthesized message must be a real row of the thread: a messageId + // is not secret (shared threads reveal them), so synthesis names the + // thread it belongs to and a foreign one is refused (probed below). + const insertMessage = async (inThread: string): Promise => { + const rows = await sql<{ id: string }[]>` + INSERT INTO app.messages ( + thread_id, org_id, "order", step_order, role, text, status, + created_at_ms + ) VALUES ( + ${inThread}, ${orgId}, 0, 0, 'assistant', 'Hello voice world.', + 'complete', ${Date.now()} + ) + RETURNING id + `; + return rows[0]?.id ?? ''; + }; + const messageId = await insertMessage(threadId); const capability = z .looseObject({ @@ -19919,7 +20005,7 @@ async function checkTts( .safeParse( await ( await send(`/api/app/tts/synthesize?orgId=${orgId}`, { - messageId: 'itest-tts-msg', + messageId, threadId, index: 0, text, @@ -19934,7 +20020,7 @@ async function checkTts( .safeParse( await ( await send(`/api/app/tts/synthesize?orgId=${orgId}`, { - messageId: 'itest-tts-msg', + messageId, threadId, index: 0, text, @@ -19958,7 +20044,7 @@ async function checkTts( .safeParse( await ( await send( - `/api/app/tts/messages/itest-tts-msg/chunks?orgId=${orgId}&threadId=${threadId}`, + `/api/app/tts/messages/${messageId}/chunks?orgId=${orgId}&threadId=${threadId}`, ) ).json(), ); @@ -19989,7 +20075,7 @@ async function checkTts( .safeParse( await ( await send(`/api/app/tts/synthesize?orgId=${orgId}`, { - messageId: 'itest-tts-msg', + messageId, threadId, index: 1, text: 'Second chunk.', @@ -20003,7 +20089,7 @@ async function checkTts( .safeParse( await ( await send(`/api/app/tts/synthesize?orgId=${orgId}`, { - messageId: 'itest-tts-msg', + messageId, threadId, index: 1, text: 'Second chunk.', @@ -20055,7 +20141,7 @@ async function checkTts( // Guards: out-of-range index, foreign thread. const badIndex = await send(`/api/app/tts/synthesize?orgId=${orgId}`, { - messageId: 'itest-tts-msg', + messageId, threadId, index: 9_999, text: 'nope', @@ -20068,9 +20154,77 @@ async function checkTts( text: 'nope', locale: 'en', }); + // A message of ANOTHER (even owned) thread under this thread: refused + // and audited — it used to mint the global (message, index) reservation + // under the wrong thread and lock the message's real thread out. + const otherThread = z.object({ id: z.string() }).safeParse( + await ( + await send(`/api/app/chat/threads?orgId=${orgId}`, { + title: 'Other voice thread', + }) + ).json(), + ); + const otherMessageId = await insertMessage( + otherThread.success ? otherThread.data.id : '', + ); + const foreignMessage = await send( + `/api/app/tts/synthesize?orgId=${orgId}`, + { + messageId: otherMessageId, + threadId, + index: 0, + text: 'nope', + locale: 'en', + }, + ); + const deniedAudit = await sql<{ count: string }[]>` + SELECT count(*)::text AS count FROM app.audit_logs + WHERE org_id = ${orgId} AND action = 'tts.synthesize_denied' + AND metadata->>'reason' = 'message_not_in_thread' + `; + const otherChunks = await sql<{ count: string }[]>` + SELECT count(*)::text AS count FROM app.tts_audio_chunks + WHERE message_id = ${otherMessageId} + `; + // The audio door is gated like every other TTS door — ownership of the + // chunk's thread — so another member holding a leaked chunk id gets 404. + const strangerSignUp = await fetch(`${base}/api/auth/sign-up/email`, { + method: 'POST', + headers: { 'content-type': 'application/json', origin: base }, + body: JSON.stringify({ + email: `itest-tts-stranger-${Date.now()}@example.com`, + password: 'itest-password-1', + name: 'TTS stranger', + }), + }); + const strangerUser = z + .object({ user: z.object({ id: z.string() }) }) + .safeParse(await strangerSignUp.json()); + await sql` + INSERT INTO "member" ("id", "organizationId", "userId", "role", + "createdAt") + VALUES (gen_random_uuid(), ${orgId}, + ${strangerUser.success ? strangerUser.data.user.id : ''}, + 'member', ${new Date()}) + `; + const strangerCookie = cookieHeaderFrom(strangerSignUp); + const strangerAudio = await fetch( + `${base}/api/app/tts/audio/${chunk0?.chunkId ?? 'missing'}?orgId=${orgId}`, + { headers: { cookie: strangerCookie, origin: base } }, + ); + const strangerChunks = z + .object({ chunks: z.array(z.unknown()) }) + .safeParse( + await ( + await fetch( + `${base}/api/app/tts/messages/${messageId}/chunks?orgId=${orgId}&threadId=${threadId}`, + { headers: { cookie: strangerCookie, origin: base } }, + ) + ).json(), + ); record( - 'tts on pg (reserve/synthesize/serve, ledger, voice-mode cascade)', + 'tts on pg (reserve/synthesize/serve, ledger, voice-mode cascade, thread-scoped doors)', capability.success && capability.data.available && capability.data.modelId === 'gpt-4o-mini-tts' && @@ -20109,8 +20263,14 @@ async function checkTts( !modeVeto.data.enabled && modeVeto.data.source === 'org_policy' && badIndex.status === 400 && - foreignThread.status === 403, - `cap=${capability.success ? `${capability.data.available}/${capability.data.modelId}/${capability.data.voice}` : 'ERR'}, synth=${synth.success ? synth.data.status : 'ERR'} calls=${callsAfterFirst} (want 1) cacheHit=${again.success ? again.data.status : 'ERR'}/calls=${callsAfterSecond} (want 1), chunks=${chunks.success ? chunks.data.chunks.length : 'ERR'} c0=${chunk0?.status}/${chunk0?.voice}/${chunk0?.format}, audio=${audio.status}:${audioBytes.length}B type=${audio.headers.get('content-type')}, ledger chars=${ledger[0]?.characterCount} (want ${text.length}) cost=${ledger[0]?.cost}, fail=${failed.success ? `${failed.data.status}/${failed.data.errorCode}` : 'ERR'} retry=${retried.success ? retried.data.status : 'ERR'}, mode=${modeDefault.success ? modeDefault.data.source : 'ERR'}→${modePref.success ? `${modePref.data.enabled}/${modePref.data.source}` : 'ERR'}→${modeThread.success ? `${modeThread.data.enabled}/${modeThread.data.source}` : 'ERR'}→veto=${modeVeto.success ? `${modeVeto.data.enabled}/${modeVeto.data.source}` : 'ERR'}, badIndex=${badIndex.status} (want 400) foreign=${foreignThread.status} (want 403)`, + foreignThread.status === 403 && + foreignMessage.status === 403 && + Number(deniedAudit[0]?.count ?? '0') >= 1 && + Number(otherChunks[0]?.count ?? '0') === 0 && + strangerAudio.status === 404 && + strangerChunks.success && + strangerChunks.data.chunks.length === 0, + `cap=${capability.success ? `${capability.data.available}/${capability.data.modelId}/${capability.data.voice}` : 'ERR'}, synth=${synth.success ? synth.data.status : 'ERR'} calls=${callsAfterFirst} (want 1) cacheHit=${again.success ? again.data.status : 'ERR'}/calls=${callsAfterSecond} (want 1), chunks=${chunks.success ? chunks.data.chunks.length : 'ERR'} c0=${chunk0?.status}/${chunk0?.voice}/${chunk0?.format}, audio=${audio.status}:${audioBytes.length}B type=${audio.headers.get('content-type')}, ledger chars=${ledger[0]?.characterCount} (want ${text.length}) cost=${ledger[0]?.cost}, fail=${failed.success ? `${failed.data.status}/${failed.data.errorCode}` : 'ERR'} retry=${retried.success ? retried.data.status : 'ERR'}, mode=${modeDefault.success ? modeDefault.data.source : 'ERR'}→${modePref.success ? `${modePref.data.enabled}/${modePref.data.source}` : 'ERR'}→${modeThread.success ? `${modeThread.data.enabled}/${modeThread.data.source}` : 'ERR'}→veto=${modeVeto.success ? `${modeVeto.data.enabled}/${modeVeto.data.source}` : 'ERR'}, badIndex=${badIndex.status} (want 400) foreignThread=${foreignThread.status} (want 403) foreignMessage=${foreignMessage.status} (want 403, audited=${deniedAudit[0]?.count ?? '0'}, squatRows=${otherChunks[0]?.count ?? '?'} want 0) strangerAudio=${strangerAudio.status} (want 404) strangerChunks=${strangerChunks.success ? strangerChunks.data.chunks.length : 'ERR'} (want 0)`, ); } finally { ttsServer.close(); @@ -21304,6 +21464,98 @@ async function checkOneDriveSync( hugeBodyPulls === 0, `browse=${bigBrowse.status}/${bigBody.success ? `${bigBody.data.items?.length}/truncated=${bigBody.data.truncated}` : 'ERR'} (want 250/false), huge=${hugeResult.success ? `${hugeResult.data.failedCount}fail ${hugeRow?.status}: ${hugeRow?.error}` : `PARSE-ERR ${hugeImport.status}`}, bodyPulls=${hugeBodyPulls} (want 0)`, ); + + // 10. A live run keeps its claim fresh. The claim fence re-admits a + // 'running' stamp older than the stale window as a crashed worker; + // a sync that merely outlives the window (large folder, slow tenant) + // used to be claimed again by the next tick and run twice at once. + seed({ id: 'folder-hb', name: 'HbFolder', folder: true }); + seed({ + id: 'f-hb', + name: 'hb.txt', + parent: 'folder-hb', + content: 'hb v1', + hash: 'h-hb-v1', + mime: 'text/plain', + }); + await post('/import', { + importType: 'sync', + items: [ + { + id: 'f-hb', + name: 'hb.txt', + size: 5, + relativePath: 'HbFolder/hb.txt', + selectedParentId: 'folder-hb', + selectedParentName: 'HbFolder', + selectedParentPath: 'HbFolder', + }, + ], + }); + await muteRagJobs(); + const hbConfig = await configByItem('folder-hb'); + const hbConfigId = hbConfig?.id ?? ''; + let listCalls = 0; + let releaseList: () => void = () => undefined; + const listGate = new Promise((resolve) => { + releaseList = resolve; + }); + const slowAdapter = { + ...onedrive.ONEDRIVE_SYNC_ADAPTER, + listFolderContents: async ( + args: Parameters< + typeof onedrive.ONEDRIVE_SYNC_ADAPTER.listFolderContents + >[0], + ) => { + listCalls += 1; + await listGate; + return onedrive.ONEDRIVE_SYNC_ADAPTER.listFolderContents(args); + }, + }; + const hbPayload = { organizationId: orgId, configId: hbConfigId }; + const longRun = onedrive.runSyncConfigJobWith(sql, slowAdapter, hbPayload, { + heartbeatMs: 100, + }); + const claimLanded = await waitFor( + () => Promise.resolve(listCalls === 1), + 5_000, + ); + // Age the claim past the stale window, as a long run's stamp would be… + await sql` + UPDATE app.onedrive_sync_configs + SET updated_at_ms = ${Date.now() - 31 * 60 * 1000} + WHERE id = ${hbConfigId} + `; + // …give the heartbeat a few periods to refresh it… + await sleep(400); + const stampRows = await sql<{ updatedAt: number }[]>` + SELECT updated_at_ms::float8 AS "updatedAt" + FROM app.onedrive_sync_configs WHERE id = ${hbConfigId} + `; + const stampAgeMs = Date.now() - (stampRows[0]?.updatedAt ?? 0); + // …then the next tick's job for the same config must no-op. + const secondRun = onedrive.runSyncConfigJobWith( + sql, + slowAdapter, + hbPayload, + { heartbeatMs: 100 }, + ); + await sleep(300); + const listCallsAfterSecond = listCalls; + releaseList(); + await Promise.all([longRun, secondRun]); + const hbAfter = await configByItem('folder-hb'); + const hbDocs = await docsByExternalId('f-hb'); + record( + 'onedrive live sync keeps its claim fresh (no concurrent re-claim)', + hbConfig?.status === 'active' && + claimLanded && + stampAgeMs < 60_000 && + listCallsAfterSecond === 1 && + hbAfter?.lastSyncStatus === 'success' && + hbDocs.length === 1, + `claimed=${claimLanded}, stampAge=${Math.round(stampAgeMs / 1000)}s (want fresh), listCalls=${listCallsAfterSecond} (want 1: the second job no-ops), final=${hbAfter?.lastSyncStatus}, docs=${hbDocs.length} (want 1)`, + ); } finally { globalThis.fetch = realFetch; if (savedEnv.tenant === undefined) { diff --git a/services/platform/backend/lib/object-store.ts b/services/platform/backend/lib/object-store.ts index eee4787697..dba2130d8e 100644 --- a/services/platform/backend/lib/object-store.ts +++ b/services/platform/backend/lib/object-store.ts @@ -1,6 +1,8 @@ import { buildObjectKey, - buildS3ObjectStore, + clearOrgObjectStoreCache, + ObjectStoreUnconfiguredError, + resolveOrgObjectStore, s3DeleteObject, s3HeadObject, s3PresignGetUrl, @@ -8,7 +10,6 @@ import { s3PutObject, type S3ObjectStore, } from '../core/lib/storage/object_store.ts'; -import { readOrgObjectStorageConnection } from '../core/object_storage/file_utils.ts'; /** * 0.5 object-store resolution — S3-compatible storage is THE blob backend @@ -20,48 +21,23 @@ import { readOrgObjectStorageConnection } from '../core/object_storage/file_util * (compose ships MinIO + a seeded connection at cutover), else * 3. fail closed: uploads are refused until storage is configured. * - * The S3 mechanics (aws4fetch signing, presign lanes, key scheme) are reused - * from `convex/lib/storage/object_store.ts` unchanged, so BYO-org configs - * written under 0.4 keep working verbatim. + * ONE resolver and ONE cache serve both this lane and the reused blob-access + * lane (`core/lib/storage/blob_access.ts`): `resolveObjectStore` IS + * `resolveOrgObjectStore`, so a broken default tree fails identically at + * every door and a config write invalidates every cached resolution at once. */ -export class ObjectStoreUnconfiguredError extends Error { - constructor() { - super( - 'No object storage configured: neither this org nor the deployment ' + - 'default tree has an object-storage/connection.json', - ); - this.name = 'ObjectStoreUnconfiguredError'; - } -} - -const STORE_TTL_MS = 15_000; -const storeCache = new Map(); +export { ObjectStoreUnconfiguredError }; -/** Test hook. */ +/** Test hook + config-write invalidation: drop every cached resolution. */ export function clearObjectStoreCache(): void { - storeCache.clear(); + clearOrgObjectStoreCache(); } export async function resolveObjectStore( orgSlug: string, ): Promise { - const cached = storeCache.get(orgSlug); - if (cached && cached.expires > Date.now()) { - return cached.store; - } - const own = await readOrgObjectStorageConnection(orgSlug); - const resolved = - own ?? - (orgSlug === 'default' - ? null - : await readOrgObjectStorageConnection('default')); - if (!resolved) { - throw new ObjectStoreUnconfiguredError(); - } - const store = buildS3ObjectStore(resolved.connection, resolved.secrets); - storeCache.set(orgSlug, { store, expires: Date.now() + STORE_TTL_MS }); - return store; + return resolveOrgObjectStore(orgSlug); } export {