Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
67 changes: 52 additions & 15 deletions indexer/enrich.py
Original file line number Diff line number Diff line change
Expand Up @@ -199,13 +199,20 @@ def gateway_key() -> str:
return ''


def caption_image(path: Path) -> str:
def vision_model() -> str:
return os.environ.get('ADE_VISION_MODEL', '')


def caption_image(path: Path) -> str | None:
"""Describe an image with a vision model through the gateway's streaming
Responses API. Opt-in via ADE_VISION_MODEL; failures degrade to OCR-only."""
Responses API. Opt-in via ADE_VISION_MODEL. Returns '' when captioning
is off or the image is too large, and None when the request failed, so
the caller can tell a picture with nothing to say from one that was
never described."""
import base64
import json
import urllib.request
model = os.environ.get('ADE_VISION_MODEL', '')
model = vision_model()
key = gateway_key()
if not model or not key or path.stat().st_size > MAX_CAPTION_BYTES:
return ''
Expand Down Expand Up @@ -239,16 +246,19 @@ def caption_image(path: Path) -> str:
text += ev.get('delta', '')
return text.strip()
except Exception as exc:
print(f' caption failed {path}: {str(exc)[:120]}', file=sys.stderr)
return ''
print(f' caption failed {path}: {str(exc)[:200]}', file=sys.stderr)
return None


def extract_image(path: Path) -> str:
def extract_image_status(path: Path) -> tuple[str, bool]:
"""The image's text, and whether it is complete enough to cache: False
when a vision model is configured and the caption request failed, so the
next pass asks again rather than keeping an empty answer for good."""
# Icons and tiny assets are skipped, except in the chat-attachments
# directory where the user attached the image deliberately.
deliberate = path.parent.name == 'chat' and path.parent.parent.name == 'uploads'
if not deliberate and path.stat().st_size < MIN_IMAGE_BYTES:
return ''
return '', True
parts = []
caption = caption_image(path)
if caption:
Expand All @@ -261,7 +271,18 @@ def extract_image(path: Path) -> str:
parts.append(f'Text found in image (OCR):\n{ocr_text}')
except Exception:
pass
return '\n\n'.join(parts)
return '\n\n'.join(parts), caption is not None


def extract_image(path: Path) -> str:
return extract_image_status(path)[0]


def image_cache_usable(text: str) -> bool:
"""A cached image result is reused unless it is empty while a vision
model is configured: a successful caption is never empty, so that entry
was written by a failed request or before captioning was switched on."""
return bool(text.strip()) or not vision_model()


def extract(path: Path) -> str | None:
Expand Down Expand Up @@ -331,8 +352,13 @@ def main() -> int:
# Expensive extractions (PDF, office, OCR, captions) are
# cached by mtime; a reindex pass reuses them untouched.
cache_file = cache_root / rel_dir / (fpath.name + '.txt')
cached = None
if cacheable and cache_file.exists() and cache_file.stat().st_mtime >= fpath.stat().st_mtime:
text = cache_file.read_text(errors='replace')[:MAX_TEXT]
cached = cache_file.read_text(errors='replace')[:MAX_TEXT]
if suffix in IMAGE_SUFFIXES and not image_cache_usable(cached):
cached = None
if cached is not None:
text = cached
if suffix == '.pdf':
# downmark's thin policy OCRs scans itself; a
# thin cache entry without its marker means the
Expand All @@ -350,14 +376,25 @@ def main() -> int:
except OSError as exc:
print(f' cache write failed {cache_file}: {exc}', file=sys.stderr)
else:
text = extract(fpath)
complete = True
if suffix in IMAGE_SUFFIXES:
try:
text, complete = extract_image_status(fpath)
text = text[:MAX_TEXT]
except Exception as exc:
print(f' extract failed {fpath}: {exc}', file=sys.stderr)
text, complete = None, False
else:
text = extract(fpath)
# Images cache even when empty: captioning and OCR are
# expensive and a picture with no text is a real
# answer. Documents do not: an empty result usually
# means the extractor was missing, and caching it hid
# the file from every later pass, including the one
# after the library was installed.
if cacheable and ((text and text.strip()) or suffix in IMAGE_SUFFIXES):
# answer, unless the caption request failed, which
# the next pass retries. Documents do not cache empty:
# an empty result usually means the extractor was
# missing, and caching it hid the file from every
# later pass, including the one after the library was
# installed.
if cacheable and complete and ((text and text.strip()) or suffix in IMAGE_SUFFIXES):
cache_file.parent.mkdir(parents=True, exist_ok=True)
cache_file.write_text(text or '')
n_extracted += 1
Expand Down
13 changes: 8 additions & 5 deletions indexer/warm_image_cache.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@
from pathlib import Path

sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
from enrich import EXCLUDE_DIRS, IMAGE_SUFFIXES, MIN_IMAGE_BYTES, extract_image # noqa: E402
from enrich import EXCLUDE_DIRS, IMAGE_SUFFIXES, MIN_IMAGE_BYTES, extract_image_status, image_cache_usable # noqa: E402


def main() -> int:
Expand All @@ -37,7 +37,8 @@ def main() -> int:
continue
rel = fpath.relative_to(kb_root)
cache_file = cache_root / rel.parent / (fname + '.txt')
if cache_file.exists() and cache_file.stat().st_mtime >= fpath.stat().st_mtime:
if cache_file.exists() and cache_file.stat().st_mtime >= fpath.stat().st_mtime \
and image_cache_usable(cache_file.read_text(errors='replace')):
continue
todo.append((fpath, cache_file))

Expand All @@ -47,9 +48,11 @@ def main() -> int:
def work(item):
nonlocal done
fpath, cache_file = item
text = extract_image(fpath)
cache_file.parent.mkdir(parents=True, exist_ok=True)
cache_file.write_text(text or '')
text, complete = extract_image_status(fpath)
# A failed caption is not cached, so the next run asks again.
if complete:
cache_file.parent.mkdir(parents=True, exist_ok=True)
cache_file.write_text(text or '')
done += 1
if done % 25 == 0:
print(f' {done}/{len(todo)}', flush=True)
Expand Down
17 changes: 15 additions & 2 deletions server/src/indexing.ts
Original file line number Diff line number Diff line change
Expand Up @@ -112,13 +112,25 @@ export function startIndexJob(rel: string): IndexJob {

export function getIndexJob(id: number): IndexJob | null { return jobs.get(id) ?? null }

/** Where indexing reports what went wrong inside a pass that still
* succeeded; the sweep timer points it at the server log. */
let indexLog: (msg: string) => void = m => console.warn(m)

function run(cmd: string, args: string[], timeoutMs = 10 * 60_000): Promise<string> {
// Settings can change the captioning model at runtime.
const vm = effectiveSettings().visionModel
const env = vm ? { ...process.env, ADE_VISION_MODEL: vm } : process.env
return new Promise((resolve, reject) => {
execFile(cmd, args, { timeout: timeoutMs, maxBuffer: 32 * 1024 * 1024, env }, (err, so, se) =>
err ? reject(new Error(`${path.basename(cmd)}: ${se || err.message}`)) : resolve(so))
execFile(cmd, args, { timeout: timeoutMs, maxBuffer: 32 * 1024 * 1024, env }, (err, so, se) => {
if (err) return reject(new Error(`${path.basename(cmd)}: ${se || err.message}`))
// A pass that succeeds can still have failed for single files (a
// caption the vision model refused, say); those lines were dropped
// with the rest of stderr, which hid why an image had no text.
for (const line of String(se ?? '').split('\n')) {
if (/\bfailed\b/.test(line)) indexLog(`index: ${line.trim().slice(0, 300)}`)
}
resolve(so)
})
})
}

Expand Down Expand Up @@ -400,6 +412,7 @@ async function primeSweepState(): Promise<void> {
}

export function startSweepTimer(log: (msg: string) => void): void {
indexLog = log
// Self-rescheduling so a settings change to the interval applies at the
// next cycle without a restart; 0 pauses sweeping but keeps checking.
const tick = async (): Promise<void> => {
Expand Down
79 changes: 79 additions & 0 deletions server/test/captionRetry.test.mjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,79 @@
// An image whose caption request fails is not cached, so the next pass
// asks again; an empty cached result is retried while a vision model is
// configured. Runs the real indexer/enrich.py against a stand-in gateway.
import test from 'node:test'
import assert from 'node:assert/strict'
import fs from 'node:fs'
import os from 'node:os'
import path from 'node:path'
import http from 'node:http'
import { execFile } from 'node:child_process'
import { promisify } from 'node:util'

const run = promisify(execFile)

const ROOT = path.resolve(import.meta.dirname, '..', '..')
const ENRICH = path.join(ROOT, 'indexer', 'enrich.py')
const work = fs.mkdtempSync(path.join(os.tmpdir(), 'caption-'))
const kb = path.join(work, 'kb'), index = path.join(work, 'index'), cache = path.join(work, 'extract')
fs.mkdirSync(kb); fs.mkdirSync(index)
// Over the 8 KB floor below which images are skipped as icons.
fs.copyFileSync(path.join(ROOT, 'web', 'public', 'icon-512.png'), path.join(kb, 'picture.png'))
fs.writeFileSync(path.join(index, 'db.db'), '')
const cacheFile = path.join(cache, 'picture.png.txt')

let mode = 'refuse'
let calls = 0
const gateway = http.createServer((req, res) => {
calls++
req.resume()
if (mode === 'refuse') { res.writeHead(403); res.end('forbidden'); return }
res.writeHead(200, { 'content-type': 'text/event-stream' })
res.end('data: {"type":"response.output_text.delta","delta":"A navy tile with a white letter S."}\n\ndata: [DONE]\n\n')
})
await new Promise(r => gateway.listen(0, '127.0.0.1', r))
const env = {
...process.env,
PW_API_KEY: 'test-key',
PW_GATEWAY_URL: `http://127.0.0.1:${gateway.address().port}`,
ADE_VISION_MODEL: 'test-vision',
}
// Asynchronous: the stand-in gateway runs in this process and has to answer
// while the indexer waits on it.
const enrich = (extra = {}) => run('python3', [ENRICH, '--kb-root', kb, '--index', index, '--extract-cache', cache],
{ env: { ...env, ...extra }, timeout: 60_000 })

test('a refused caption is not cached, and the next pass asks again', async () => {
await enrich()
assert.equal(calls, 1)
assert.equal(fs.existsSync(cacheFile), false)
mode = 'answer'
await enrich()
assert.equal(calls, 2)
assert.match(fs.readFileSync(cacheFile, 'utf8'), /Image description:\nA navy tile with a white letter S\./)
})

test('a good cached caption is reused without asking again', async () => {
await enrich()
assert.equal(calls, 2)
})

test('an empty cached entry is retried while a vision model is configured', async () => {
fs.writeFileSync(cacheFile, '')
const later = new Date(Date.now() + 60_000)
fs.utimesSync(cacheFile, later, later)
await enrich()
assert.equal(calls, 3)
assert.match(fs.readFileSync(cacheFile, 'utf8'), /A navy tile/)
})

test('with captioning off, an empty entry stands', async () => {
fs.writeFileSync(cacheFile, '')
const later = new Date(Date.now() + 120_000)
fs.utimesSync(cacheFile, later, later)
await enrich({ ADE_VISION_MODEL: '' })
assert.equal(calls, 3)
assert.equal(fs.readFileSync(cacheFile, 'utf8'), '')
})

test.after(() => { gateway.close(); fs.rmSync(work, { recursive: true, force: true }) })
Loading