diff --git a/CHANGELOG.md b/CHANGELOG.md index 5383285..aa1aa3a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,17 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Added +* `asyncWrites` parameter on the `VectorDatabase` constructor and on `VectorDatabase::open()`, defaulting to `false`, to opt into writing document files in forked child processes. + +### Changed +* Document files are written synchronously by default. `addDocument()` used to fork a child process for every document whenever `ext-pcntl` was loaded, which cost more than the write itself. + +### Fixed +* An async write child ended with `exit(0)`, which ran the shutdown functions and destructors inherited from the parent. In a host application this could close resources the parent still used, such as a shared database connection. The child now ends with `SIGKILL`, so async writes also require `ext-posix`. +* Async write children were reaped only by `save()`, so bulk imports accumulated zombie processes. Finished children are now reaped on every async write. +* A failed async write went unnoticed. It now raises a `RuntimeException` when the child is reaped. + ## [0.4.1] - 2026-09-02 Documentation only. No code changes, no behaviour changes. diff --git a/README.md b/README.md index 780315f..7114cc5 100644 --- a/README.md +++ b/README.md @@ -6,7 +6,7 @@ A pure-PHP vector database implementing **HNSW** (Hierarchical Navigable Small W - PHP 8.2+ - No external PHP extensions required for core functionality -- `ext-pcntl` (optional) — enables asynchronous document writes for lower insert latency +- `ext-pcntl` and `ext-posix` (optional), for the opt-in asynchronous document writes (`asyncWrites: true`) ## Installation @@ -187,10 +187,10 @@ $db = new VectorDatabase( ## Persistence -PHPVector uses a **folder-based** persistence model. Each database lives in its own directory containing separate files for the HNSW graph, the BM25 index, and one file per document. This design has two key advantages: +PHPVector uses a **folder-based** persistence model. Each database lives in its own directory containing separate files for the HNSW graph, the BM25 index, and one file per document. This design keeps memory low and inserts cheap: - **Low memory footprint on load** — only the HNSW graph and BM25 index are loaded into memory. Individual document files (`docs/{n}.bin`) are read lazily, only for the documents that appear in search results. -- **Low insert latency** — document files are written to disk asynchronously in a forked child process (requires `ext-pcntl`), so `addDocument()` returns immediately. +- **Incremental writes**. `addDocument()` writes only its own document file, so inserting does not rewrite the rest of the database. ### Folder layout @@ -208,7 +208,9 @@ PHPVector uses a **folder-based** persistence model. Each database lives in its ### Saving -Pass a `path` to the constructor to enable persistence. Each `addDocument()` call writes the document file to `docs/` (asynchronously when `ext-pcntl` is available). Call `save()` once to flush the HNSW graph and BM25 index — it waits for any outstanding async writes before proceeding. +Pass a `path` to the constructor to enable persistence. Each `addDocument()` call writes the document file to `docs/`. Call `save()` once to flush the HNSW graph and BM25 index. + +Document files are written synchronously by default. Passing `asyncWrites: true` to the constructor or to `open()` writes each one in a forked child process instead (requires `ext-pcntl` and `ext-posix`, otherwise writes stay synchronous), and `save()` waits for the pending writes before flushing the indexes. A fork usually costs more than writing one small file, so only enable it when document writes are slow, for example on network storage. The child ends without running shutdown functions or destructors, so it cannot close resources it shares with the parent, such as a database connection. ```php use PHPVector\Document; diff --git a/src/Persistence/DocumentStore.php b/src/Persistence/DocumentStore.php index cd5d1ef..9201955 100644 --- a/src/Persistence/DocumentStore.php +++ b/src/Persistence/DocumentStore.php @@ -15,8 +15,16 @@ * [4 bytes: textLen] [textLen bytes: utf-8 text] (textLen=0 → null text) * [4 bytes: metaLen] [metaLen bytes: JSON metadata] (metaLen=0 → empty array) * - * Writes can be dispatched to forked child processes (pcntl_fork) to avoid - * blocking the caller. Call waitAll() before reading back or writing index files. + * Writes are synchronous by default. Async writes dispatch each one to a + * forked child process (pcntl_fork) and need both ext-pcntl and ext-posix; + * without them they fall back to synchronous writes. Call waitAll() before + * reading back or writing index files. + * + * An async child ends by sending itself SIGKILL instead of calling exit(), so + * it never runs the shutdown functions and destructors it inherited from the + * parent (which could, for example, close a database connection the parent + * still uses). Finished children are reaped on every async write, so a long + * import does not accumulate zombie processes. * * Every write goes through AtomicFile, so a reader loading `{nodeId}.bin` * lazily always sees a complete record: either the previous version or the new @@ -47,48 +55,77 @@ public function __construct(private readonly string $docsDir) {} /** * Persist a document to disk. * - * When $async is true and pcntl_fork() is available the write is - * dispatched to a child process; the parent returns immediately. - * When unavailable the write is synchronous. + * When $async is true and async writes are supported (see + * supportsAsync()) the write is dispatched to a child process and the + * parent returns immediately. Otherwise the write is synchronous. * * @param int $nodeId * @param string|int $docId Must NOT be null (UUID already assigned by caller). * @param string|null $text * @param array $metadata * @param bool $async + * + * @throws \RuntimeException if a synchronous write fails. */ public function write( int $nodeId, string|int $docId, ?string $text, array $metadata, - bool $async = true, + bool $async = false, ): void { - if ($async && function_exists('pcntl_fork')) { + if ($async && self::supportsAsync()) { + $this->reapFinished(); + $pid = pcntl_fork(); - if ($pid === -1) { - // Fork failed — fall through to synchronous write. - } elseif ($pid === 0) { - // Child: write and exit immediately. - $this->writeSync($nodeId, $docId, $text, $metadata); - exit(0); - } else { - // Parent: record PID keyed by nodeId and return. + if ($pid === 0) { + // Child: never return into the caller's code, whatever happens. + try { + $this->writeSync($nodeId, $docId, $text, $metadata); + } catch (\Throwable $e) { + // The parent notices the missing file when it reaps this child. + error_log(sprintf('PHPVector: async write of node %d failed: %s', $nodeId, $e->getMessage())); + } finally { + posix_kill(getmypid(), SIGKILL); + } + } + + if ($pid > 0) { $this->pendingPids[$nodeId] = $pid; return; } + + // Fork failed: fall through to a synchronous write. } - // Synchronous path (no fork or fork failed). $this->writeSync($nodeId, $docId, $text, $metadata); } + /** + * Whether async writes can run in this process: forking needs ext-pcntl, + * and ending the child without running shutdown code needs ext-posix. + */ + public static function supportsAsync(): bool + { + return function_exists('pcntl_fork') + && function_exists('pcntl_waitpid') + && function_exists('posix_kill'); + } + + /** Number of async writes started but not yet reaped. */ + public function pendingCount(): int + { + return count($this->pendingPids); + } + /** * Block until the async write for a specific node has completed. * * Use this before deleting a node's file so a late child write cannot * recreate {nodeId}.bin after the unlink(). + * + * @throws \RuntimeException if the child finished without writing the file. */ public function waitForNode(int $nodeId): void { @@ -96,25 +133,60 @@ public function waitForNode(int $nodeId): void return; } - if (function_exists('pcntl_waitpid')) { - pcntl_waitpid($this->pendingPids[$nodeId], $status); - } - + pcntl_waitpid($this->pendingPids[$nodeId], $status); unset($this->pendingPids[$nodeId]); + $this->assertWritten($nodeId); } /** * Block until every outstanding async write has completed. * Must be called before index files are written (see VectorDatabase::save()). + * + * @throws \RuntimeException if a child finished without writing its file. */ public function waitAll(): void { - foreach ($this->pendingPids as $pid) { - if (function_exists('pcntl_waitpid')) { - pcntl_waitpid($pid, $status); + $pending = $this->pendingPids; + $this->pendingPids = []; + + foreach ($pending as $pid) { + pcntl_waitpid($pid, $status); + } + foreach (array_keys($pending) as $nodeId) { + $this->assertWritten($nodeId); + } + } + + /** + * Reap children that already exited, without blocking. + * + * @throws \RuntimeException if a child finished without writing its file. + */ + private function reapFinished(): void + { + foreach ($this->pendingPids as $nodeId => $pid) { + if (pcntl_waitpid($pid, $status, WNOHANG) !== 0) { + unset($this->pendingPids[$nodeId]); + $this->assertWritten($nodeId); } } - $this->pendingPids = []; + } + + /** + * Async children cannot report errors through their exit status (they end + * with SIGKILL), so a missing file is the failure signal. + * + * @throws \RuntimeException + */ + private function assertWritten(int $nodeId): void + { + clearstatcache(true, $this->filePath($nodeId)); + if (!is_file($this->filePath($nodeId))) { + throw new \RuntimeException(sprintf( + 'Async write of document file failed: %s', + $this->filePath($nodeId), + )); + } } // ------------------------------------------------------------------ diff --git a/src/VectorDatabase.php b/src/VectorDatabase.php index 10748ed..2322de6 100644 --- a/src/VectorDatabase.php +++ b/src/VectorDatabase.php @@ -57,9 +57,10 @@ * loaded into memory by `open()`; individual `docs/{n}.bin` files are read on * demand when search results are hydrated. * - * Individual document files are written **asynchronously** (via `pcntl_fork`) - * on each `addDocument()` call when the extension is available. `save()` - * waits for all pending writes before flushing the index files. + * Individual document files are written on each `addDocument()` call, + * synchronously by default. With `$asyncWrites` each write runs in a forked + * child process instead (requires ext-pcntl and ext-posix). `save()` waits + * for all pending writes before flushing the index files. * * Quick start * ----------- @@ -117,6 +118,10 @@ final class VectorDatabase /** * @param float $lockTimeout Seconds save() waits for the folder lock before * throwing a LockTimeoutException. + * @param bool $asyncWrites Write document files in forked child processes + * (requires ext-pcntl and ext-posix, otherwise + * writes stay synchronous). Only worth it when a + * single file write is slower than a fork. */ public function __construct( HNSWConfig $hnswConfig = new HNSWConfig(), @@ -125,6 +130,7 @@ public function __construct( private readonly ?string $path = null, private readonly int $overFetchMultiplier = 5, private readonly float $lockTimeout = FileLock::DEFAULT_TIMEOUT, + private readonly bool $asyncWrites = false, ) { if ($overFetchMultiplier < 1) { throw new \InvalidArgumentException('overFetchMultiplier must be at least 1.'); @@ -156,8 +162,8 @@ public function isPersistent(): bool * Add a single document. * * If `$document->id` is null a random UUID v4 is assigned automatically. - * When a folder path is configured the document is written to disk - * asynchronously (pcntl_fork when available, synchronous otherwise). + * When a folder path is configured the document file is written to disk, + * in a forked child process when `$asyncWrites` is enabled. * * @throws \RuntimeException if a document with the same ID already exists. */ @@ -186,7 +192,7 @@ public function addDocument(Document $document): void $this->hnswIndex->insert($document); $this->bm25Index->addDocument($nodeId, $document); - // Persist doc file asynchronously when a path is configured. + // Persist the doc file when a path is configured. if ($this->path !== null) { $this->ensureDocsDir(); $this->getDocumentStore()->write( @@ -194,7 +200,7 @@ public function addDocument(Document $document): void docId: $document->id, text: $document->text, metadata: $document->metadata, - async: true, + async: $this->asyncWrites, ); } } @@ -723,6 +729,7 @@ public function save(): void * * @param float $lockTimeout Seconds to wait for the folder lock before * throwing a LockTimeoutException. + * @param bool $asyncWrites See the constructor. * * @throws LockTimeoutException if another process is writing the folder. * @throws \RuntimeException on I/O failure or distance metric mismatch. @@ -734,6 +741,7 @@ public static function open( TokenizerInterface $tokenizer = new SimpleTokenizer(), int $overFetchMultiplier = 5, float $lockTimeout = FileLock::DEFAULT_TIMEOUT, + bool $asyncWrites = false, ): self { $metaPath = $path . '/meta.json'; if (!file_exists($metaPath)) { @@ -744,7 +752,7 @@ public static function open( $lock->acquireShared($lockTimeout); try { - return self::loadLocked($path, $hnswConfig, $bm25Config, $tokenizer, $overFetchMultiplier, $lockTimeout); + return self::loadLocked($path, $hnswConfig, $bm25Config, $tokenizer, $overFetchMultiplier, $lockTimeout, $asyncWrites); } finally { $lock->release(); } @@ -761,6 +769,7 @@ private static function loadLocked( TokenizerInterface $tokenizer, int $overFetchMultiplier, float $lockTimeout, + bool $asyncWrites, ): self { $metaPath = $path . '/meta.json'; $meta = json_decode(file_get_contents($metaPath), true, 512, JSON_THROW_ON_ERROR); @@ -776,7 +785,7 @@ private static function loadLocked( )); } - $db = new self($hnswConfig, $bm25Config, $tokenizer, $path, $overFetchMultiplier, $lockTimeout); + $db = new self($hnswConfig, $bm25Config, $tokenizer, $path, $overFetchMultiplier, $lockTimeout, $asyncWrites); $db->nextId = (int) $meta['nextId']; $db->docIdToNodeId = $meta['docIdToNodeId']; diff --git a/tests/Persistence/DocumentStoreTest.php b/tests/Persistence/DocumentStoreTest.php new file mode 100644 index 0000000..40202a4 --- /dev/null +++ b/tests/Persistence/DocumentStoreTest.php @@ -0,0 +1,159 @@ +tmpDir = sys_get_temp_dir() . '/phpvtest_store_' . uniqid('', true); + mkdir($this->tmpDir, 0755, true); + } + + protected function tearDown(): void + { + foreach (glob($this->tmpDir . '/{,.}*', GLOB_BRACE) ?: [] as $file) { + if (is_file($file)) { + unlink($file); + } + } + rmdir($this->tmpDir); + } + + private function requireAsync(): void + { + if (!DocumentStore::supportsAsync()) { + self::markTestSkipped('Async writes need ext-pcntl and ext-posix.'); + } + } + + public function testWriteIsSynchronousByDefault(): void + { + $store = new DocumentStore($this->tmpDir); + $store->write(0, 'a', 'hello', ['k' => 'v']); + + self::assertSame(0, $store->pendingCount()); + self::assertSame(['a', 'hello', ['k' => 'v']], $store->read(0)); + } + + public function testAsyncWriteIsReadableAfterWaitAll(): void + { + $this->requireAsync(); + + $store = new DocumentStore($this->tmpDir); + $store->write(0, 'a', 'hello', ['k' => 'v'], async: true); + $store->waitAll(); + + self::assertSame(0, $store->pendingCount()); + self::assertSame(['a', 'hello', ['k' => 'v']], $store->read(0)); + } + + public function testFinishedChildrenAreReapedOnLaterWrites(): void + { + $this->requireAsync(); + + $store = new DocumentStore($this->tmpDir); + $total = 30; + for ($i = 0; $i < $total; $i++) { + $store->write($i, $i, "doc {$i}", [], async: true); + usleep(20_000); + } + + // Every child exits within milliseconds, so most of them must have been + // reaped by the following writes instead of piling up until waitAll(). + self::assertLessThan($total, $store->pendingCount()); + + $store->waitAll(); + self::assertSame([29, 'doc 29', []], $store->read(29)); + } + + public function testFailedAsyncWriteIsReportedByWaitAll(): void + { + $this->requireAsync(); + + // The child logs its failure; keep that out of the test output. + $previousLog = ini_set('error_log', '/dev/null'); + + try { + $store = new DocumentStore($this->tmpDir . '/missing-dir'); + $store->write(0, 'a', 'hello', [], async: true); + + $this->expectException(\RuntimeException::class); + $this->expectExceptionMessage('Async write of document file failed'); + $store->waitAll(); + } finally { + ini_set('error_log', $previousLog === false ? '' : $previousLog); + } + } + + /** + * The child must not run the shutdown functions and destructors it + * inherits from the parent: in a host application they can close shared + * resources such as a database connection. Runs in a separate process so + * that PHPUnit's own shutdown handlers stay out of the picture. + */ + public function testAsyncChildRunsNoShutdownFunctionsOrDestructors(): void + { + $this->requireAsync(); + if (!function_exists('proc_open')) { + self::markTestSkipped('proc_open() is not available.'); + } + + $marker = $this->tmpDir . '/marker.log'; + $script = $this->tmpDir . '/child.php'; + $autoload = dirname(__DIR__, 2) . '/vendor/autoload.php'; + + file_put_contents($script, <<<'CHILD' + marker, 'destruct:' . getmypid() . "\n", FILE_APPEND); + } + } + + $witness = new Witness($marker); + register_shutdown_function(static function () use ($marker): void { + file_put_contents($marker, 'shutdown:' . getmypid() . "\n", FILE_APPEND); + }); + + $store = new PHPVector\Persistence\DocumentStore($docsDir); + $store->write(0, 'a', 'hello', [], async: true); + $store->waitAll(); + + echo getmypid(); + CHILD); + + $descriptors = [1 => ['pipe', 'w'], 2 => ['pipe', 'w']]; + $process = proc_open([PHP_BINARY, $script, $autoload, $this->tmpDir, $marker], $descriptors, $pipes); + if (!is_resource($process)) { + self::markTestSkipped('Could not spawn a child PHP process.'); + } + + $parentPid = (int) stream_get_contents($pipes[1]); + $stderr = (string) stream_get_contents($pipes[2]); + fclose($pipes[1]); + fclose($pipes[2]); + self::assertSame(0, proc_close($process), $stderr); + + self::assertFileExists($this->tmpDir . '/0.bin'); + self::assertSame( + ["shutdown:{$parentPid}", "destruct:{$parentPid}"], + file($marker, FILE_IGNORE_NEW_LINES), + 'Only the parent process may run shutdown functions and destructors.', + ); + } +} diff --git a/tests/PersistenceTest.php b/tests/PersistenceTest.php index 53deeff..2417fc4 100644 --- a/tests/PersistenceTest.php +++ b/tests/PersistenceTest.php @@ -11,6 +11,7 @@ use PHPVector\Document; use PHPVector\HNSW\Config as HNSWConfig; use PHPVector\HybridMode; +use PHPVector\Persistence\DocumentStore; use PHPVector\VectorDatabase; final class PersistenceTest extends TestCase @@ -467,4 +468,47 @@ public function testUpdateDocumentPersistsCorrectly(): void self::assertSame('updated content', $results[0]->document->text); self::assertSame(['version' => 2], $results[0]->document->metadata); } + + // ------------------------------------------------------------------ + // Document file writes + // ------------------------------------------------------------------ + + public function testDocumentFileIsWrittenSynchronouslyByDefault(): void + { + $db = $this->makeDb(); + $db->addDocument(new Document(id: 1, vector: [1.0, 0.0], text: 'on disk right away')); + + // No save() and no wait: the file must already be complete. + self::assertFileExists($this->tmpDir . '/docs/0.bin'); + self::assertSame( + [1, 'on disk right away', []], + (new DocumentStore($this->tmpDir . '/docs'))->read(0), + ); + $db->save(); + } + + public function testAsyncWritesRoundTrip(): void + { + if (!DocumentStore::supportsAsync()) { + self::markTestSkipped('Async writes need ext-pcntl and ext-posix.'); + } + + $db = new VectorDatabase( + hnswConfig: new HNSWConfig(M: 8, efConstruction: 50, efSearch: 20), + tokenizer: new SimpleTokenizer([]), + path: $this->tmpDir, + asyncWrites: true, + ); + for ($i = 0; $i < 10; $i++) { + $db->addDocument(new Document(id: $i, vector: [cos($i), sin($i)], text: "doc {$i}")); + } + $db->deleteDocument(3); + $db->save(); + + $loaded = $this->openDb(); + self::assertSame(9, $loaded->count()); + $top = $loaded->vectorSearch([cos(7), sin(7)], k: 1); + self::assertSame(7, $top[0]->document->id); + self::assertSame('doc 7', $top[0]->document->text); + } }