diff --git a/.env.sample b/.env.sample index e436493..675a413 100644 --- a/.env.sample +++ b/.env.sample @@ -1,6 +1,7 @@ # required PDS_HOSTNAME="" PDS_BSKY_APP_VIEW_URL="https://api.bsky.app" +PDS_DATA_DIRECTORY="/pds/data" # optional PDS_PRIVACY_POLICY_URL="" diff --git a/.gitignore b/.gitignore index a6115f0..4dfa30d 100644 --- a/.gitignore +++ b/.gitignore @@ -9,3 +9,4 @@ .phpunit.result.cache .phpunit.cache .env +/data diff --git a/app/settings.php b/app/settings.php index dcc9e0b..9bf4d71 100644 --- a/app/settings.php +++ b/app/settings.php @@ -8,10 +8,17 @@ use DI\ContainerBuilder; use Monolog\Level; return function (ContainerBuilder $containerBuilder) { + $getRequiredEnv = static function (string $key) { + $value = $_ENV[$key] ?? null; + if (!is_string($value) || $value === '') { + throw new \RuntimeException("$key is required"); + } + return $value; + }; // Global Settings Object $containerBuilder->addDefinitions([ - SettingsInterface::class => function () { + SettingsInterface::class => function () use ($getRequiredEnv) { return new Settings([ 'displayErrorDetails' => true, // Should be set to false in production 'logError' => false, @@ -23,12 +30,30 @@ return function (ContainerBuilder $containerBuilder) { ], // PDS-specific settings 'pds' => [ - 'hostname' => $_ENV['PDS_HOSTNAME'] ?? throw new \RuntimeException('PDS_HOSTNAME is required'), - 'bskyAppViewUrl' => $_ENV['PDS_BSKY_APP_VIEW_URL'] ?? throw new \RuntimeException('PDS_BSKY_APP_VIEW_URL is required'), + 'hostname' => $getRequiredEnv('PDS_HOSTNAME'), + 'bskyAppViewUrl' => $getRequiredEnv('PDS_BSKY_APP_VIEW_URL'), 'privacyPolicyUrl' => $_ENV['PDS_PRIVACY_POLICY_URL'] ?? null, 'termsOfServiceUrl' => $_ENV['PDS_TERMS_OF_SERVICE_URL'] ?? null, 'email' => $_ENV['PDS_CONTACT_EMAIL_ADDRESS'] ?? null, - ] + ], + // Persistence (SQLite) settings + 'database' => (static function () use ($getRequiredEnv): array { + $dataDir = $getRequiredEnv('PDS_DATA_DIRECTORY'); + $resolve = static function (string $suffix) use ($dataDir): string { + if ($dataDir === ':memory:') { + return ':memory:'; + } + return rtrim($dataDir, '/') . '/' . $suffix; + }; + + return [ + 'accountDb' => $resolve('account.sqlite'), + 'sequencerDb' => $resolve('sequencer.sqlite'), + 'didCacheDb' => $resolve('did_cache.sqlite'), + 'actorStoreDir' => $resolve('actors'), + 'blobstoreDir' => $resolve('blocks'), + ]; + })(), ]); } ]); diff --git a/src/Domain/Account/AccountRepository.php b/src/Domain/Account/AccountRepository.php index 9be8af2..6f7d324 100644 --- a/src/Domain/Account/AccountRepository.php +++ b/src/Domain/Account/AccountRepository.php @@ -25,4 +25,9 @@ interface AccountRepository * @throws AccountNotFoundException */ public function findAccountByEmail(string $email): Account; + + /** + * Persist an account. + */ + public function save(Account $account): void; } diff --git a/src/Domain/Actor/ActorRepository.php b/src/Domain/Actor/ActorRepository.php index 201269e..f018cf4 100644 --- a/src/Domain/Actor/ActorRepository.php +++ b/src/Domain/Actor/ActorRepository.php @@ -20,4 +20,9 @@ interface ActorRepository * @throws ActorNotFoundException */ public function findActorByHandle(string $handle): Actor; + + /** + * Persist an actor. + */ + public function save(Actor $actor): void; } diff --git a/src/Infrastructure/Database/Database.php b/src/Infrastructure/Database/Database.php new file mode 100644 index 0000000..6c44fc1 --- /dev/null +++ b/src/Infrastructure/Database/Database.php @@ -0,0 +1,141 @@ +location = $location; + + if ($location !== ':memory:') { + $dir = dirname($location); + if (!is_dir($dir)) { + if (!mkdir($dir, 0o755, true) && !is_dir($dir)) { + throw new \RuntimeException("Could not create database directory: {$dir}"); + } + } + } + + $dsn = 'sqlite:' . $location; + + $this->pdo = new PDO($dsn, null, null, [ + PDO::ATTR_ERRMODE => PDO::ERRMODE_EXCEPTION, + PDO::ATTR_DEFAULT_FETCH_MODE => PDO::FETCH_ASSOC, + PDO::ATTR_EMULATE_PREPARES => false, + ]); + + $this->pdo->exec('PRAGMA foreign_keys = ON'); + + if ($location !== ':memory:') { + $this->pdo->exec('PRAGMA journal_mode = WAL'); + $this->pdo->exec('PRAGMA synchronous = NORMAL'); + } + } + + public function pdo(): PDO + { + return $this->pdo; + } + + public function getLocation(): string + { + return $this->location; + } + + /** + * Prepare and execute a statement, returning the underlying PDOStatement. + * + * Because the connection runs with ERRMODE_EXCEPTION, prepare()/execute() + * never return false in practice. This wrapper exists primarily to give + * static analysers a non-nullable PDOStatement back. + * + * @param array $params + */ + public function prepared(string $sql, array $params = []): \PDOStatement + { + $stmt = $this->pdo->prepare($sql); + $stmt->execute($params); + return $stmt; + } + + /** + * Fetch a single row as an associative array, or null when no row matches. + * + * @param array $params + * @return array|null + */ + public function fetchOne(string $sql, array $params = []): ?array + { + $stmt = $this->prepared($sql, $params); + $row = $stmt->fetch(PDO::FETCH_ASSOC); + if ($row === false) { + return null; + } + /** @var array $row */ + return $row; + } + + /** + * Fetch all matching rows as a list of associative arrays. + * + * @param array $params + * @return list> + */ + public function fetchAll(string $sql, array $params = []): array + { + $stmt = $this->prepared($sql, $params); + /** @var list> $rows */ + $rows = $stmt->fetchAll(PDO::FETCH_ASSOC); + return $rows; + } + + /** + * Execute an INSERT/UPDATE/DELETE/DDL statement and return the affected + * row count. + * + * @param array $params + */ + public function execute(string $sql, array $params = []): int + { + return $this->prepared($sql, $params)->rowCount(); + } + + /** + * Run $callback inside a transaction; commit on success, rollback on + * any throwable. + * + * @template T + * @param callable(PDO): T $callback + * @return T + */ + public function transaction(callable $callback): mixed + { + $this->pdo->beginTransaction(); + try { + $result = $callback($this->pdo); + $this->pdo->commit(); + return $result; + } catch (\Throwable $e) { + if ($this->pdo->inTransaction()) { + $this->pdo->rollBack(); + } + throw $e; + } + } +} diff --git a/src/Infrastructure/Database/Row.php b/src/Infrastructure/Database/Row.php new file mode 100644 index 0000000..e866688 --- /dev/null +++ b/src/Infrastructure/Database/Row.php @@ -0,0 +1,60 @@ + $row */ + public static function str(array $row, string $key): string + { + $v = $row[$key] ?? null; + assert(is_string($v)); + return $v; + } + + /** @param array $row */ + public static function nstr(array $row, string $key): ?string + { + $v = $row[$key] ?? null; + assert($v === null || is_string($v)); + return $v; + } + + /** @param array $row */ + public static function int(array $row, string $key): int + { + $v = $row[$key] ?? null; + assert(is_int($v) || (is_string($v) && $v !== '' && (string) (int) $v === $v)); + return (int) $v; + } + + /** @param array $row */ + public static function nint(array $row, string $key): ?int + { + $v = $row[$key] ?? null; + if ($v === null) { + return null; + } + assert(is_int($v) || is_string($v)); + return (int) $v; + } + + /** @param array $row */ + public static function bool(array $row, string $key): bool + { + $v = $row[$key] ?? null; + assert(is_int($v) || is_bool($v) || is_string($v)); + return (bool) (int) $v; + } +} diff --git a/src/Infrastructure/Database/Schema/AccountSchema.php b/src/Infrastructure/Database/Schema/AccountSchema.php new file mode 100644 index 0000000..664d06d --- /dev/null +++ b/src/Infrastructure/Database/Schema/AccountSchema.php @@ -0,0 +1,185 @@ +pdo(); + + $pdo->exec( + 'CREATE TABLE IF NOT EXISTS account ( + did TEXT PRIMARY KEY, + email TEXT NOT NULL UNIQUE, + password_scrypt TEXT NOT NULL, + email_confirmed_at TEXT, + invites_disabled INTEGER NOT NULL DEFAULT 0 + )' + ); + + $pdo->exec( + 'CREATE TABLE IF NOT EXISTS actor ( + did TEXT PRIMARY KEY, + handle TEXT UNIQUE, + created_at TEXT NOT NULL, + takedown_ref TEXT, + deactivated_at TEXT, + delete_after TEXT + )' + ); + + $pdo->exec( + 'CREATE TABLE IF NOT EXISTS app_password ( + did TEXT NOT NULL, + name TEXT NOT NULL, + password_scrypt TEXT NOT NULL, + created_at TEXT NOT NULL, + privileged INTEGER NOT NULL DEFAULT 0, + PRIMARY KEY (did, name) + )' + ); + + $pdo->exec( + 'CREATE TABLE IF NOT EXISTS email_token ( + purpose TEXT NOT NULL, + did TEXT NOT NULL, + token TEXT NOT NULL, + requested_at TEXT NOT NULL, + PRIMARY KEY (purpose, did), + UNIQUE (purpose, token) + )' + ); + + $pdo->exec( + 'CREATE TABLE IF NOT EXISTS invite_code ( + code TEXT PRIMARY KEY, + available_uses INTEGER NOT NULL, + disabled INTEGER NOT NULL DEFAULT 0, + for_account TEXT NOT NULL, + created_by TEXT NOT NULL, + created_at TEXT NOT NULL + )' + ); + + $pdo->exec( + 'CREATE TABLE IF NOT EXISTS invite_code_use ( + code TEXT NOT NULL, + used_by TEXT NOT NULL, + used_at TEXT NOT NULL, + PRIMARY KEY (code, used_by, used_at) + )' + ); + + $pdo->exec( + 'CREATE TABLE IF NOT EXISTS refresh_token ( + id TEXT PRIMARY KEY, + did TEXT NOT NULL, + expires_at TEXT NOT NULL, + app_password_name TEXT, + next_id TEXT + )' + ); + $pdo->exec('CREATE INDEX IF NOT EXISTS refresh_token_did_idx ON refresh_token (did)'); + + $pdo->exec( + 'CREATE TABLE IF NOT EXISTS account_device ( + did TEXT NOT NULL, + device_id TEXT NOT NULL, + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL, + PRIMARY KEY (did, device_id) + )' + ); + + $pdo->exec( + 'CREATE TABLE IF NOT EXISTS authorized_client ( + did TEXT NOT NULL, + client_id TEXT NOT NULL, + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL, + data_json TEXT NOT NULL, + PRIMARY KEY (did, client_id) + )' + ); + + $pdo->exec( + 'CREATE TABLE IF NOT EXISTS device ( + id TEXT PRIMARY KEY, + session_id TEXT NOT NULL, + user_agent TEXT, + ip_address TEXT NOT NULL, + last_seen_at TEXT NOT NULL + )' + ); + + $pdo->exec( + 'CREATE TABLE IF NOT EXISTS oauth_token ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + did TEXT NOT NULL, + token_id TEXT NOT NULL UNIQUE, + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL, + expires_at TEXT NOT NULL, + client_id TEXT NOT NULL, + client_auth_json TEXT NOT NULL, + device_id TEXT, + parameters_json TEXT NOT NULL, + details_json TEXT, + code TEXT, + current_refresh_token TEXT, + scope TEXT + )' + ); + $pdo->exec('CREATE INDEX IF NOT EXISTS oauth_token_did_idx ON oauth_token (did)'); + $pdo->exec('CREATE INDEX IF NOT EXISTS oauth_token_code_idx ON oauth_token (code)'); + $pdo->exec( + 'CREATE INDEX IF NOT EXISTS oauth_token_refresh_idx ON oauth_token (current_refresh_token)' + ); + + $pdo->exec( + 'CREATE TABLE IF NOT EXISTS used_refresh_token ( + refresh_token TEXT PRIMARY KEY, + token_id INTEGER NOT NULL + )' + ); + $pdo->exec( + 'CREATE INDEX IF NOT EXISTS used_refresh_token_token_id_idx ON used_refresh_token (token_id)' + ); + + $pdo->exec( + 'CREATE TABLE IF NOT EXISTS authorization_request ( + id TEXT PRIMARY KEY, + did TEXT, + device_id TEXT, + client_id TEXT NOT NULL, + client_auth_json TEXT, + parameters_json TEXT NOT NULL, + expires_at TEXT NOT NULL, + code TEXT + )' + ); + $pdo->exec( + 'CREATE INDEX IF NOT EXISTS authorization_request_code_idx ON authorization_request (code)' + ); + + $pdo->exec( + 'CREATE TABLE IF NOT EXISTS lexicon ( + nsid TEXT PRIMARY KEY, + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL, + last_succeeded_at TEXT, + uri TEXT, + lexicon_json TEXT + )' + ); + } +} diff --git a/src/Infrastructure/Database/Schema/ActorStoreSchema.php b/src/Infrastructure/Database/Schema/ActorStoreSchema.php new file mode 100644 index 0000000..5293e37 --- /dev/null +++ b/src/Infrastructure/Database/Schema/ActorStoreSchema.php @@ -0,0 +1,93 @@ +pdo(); + + $pdo->exec( + 'CREATE TABLE IF NOT EXISTS record ( + uri TEXT PRIMARY KEY, + cid TEXT NOT NULL, + collection TEXT NOT NULL, + rkey TEXT NOT NULL, + repo_rev TEXT NOT NULL, + indexed_at TEXT NOT NULL, + takedown_ref TEXT + )' + ); + $pdo->exec('CREATE INDEX IF NOT EXISTS record_collection_idx ON record (collection)'); + + $pdo->exec( + 'CREATE TABLE IF NOT EXISTS record_blob ( + blob_cid TEXT NOT NULL, + record_uri TEXT NOT NULL, + PRIMARY KEY (blob_cid, record_uri) + )' + ); + $pdo->exec( + 'CREATE INDEX IF NOT EXISTS record_blob_record_uri_idx ON record_blob (record_uri)' + ); + + $pdo->exec( + 'CREATE TABLE IF NOT EXISTS backlink ( + uri TEXT NOT NULL, + path TEXT NOT NULL, + link_to TEXT NOT NULL, + PRIMARY KEY (uri, path, link_to) + )' + ); + $pdo->exec('CREATE INDEX IF NOT EXISTS backlink_link_to_idx ON backlink (link_to)'); + + $pdo->exec( + 'CREATE TABLE IF NOT EXISTS blob ( + cid TEXT PRIMARY KEY, + mime_type TEXT NOT NULL, + size INTEGER NOT NULL, + temp_key TEXT, + created_at TEXT NOT NULL, + takedown_ref TEXT + )' + ); + $pdo->exec('CREATE INDEX IF NOT EXISTS blob_temp_key_idx ON blob (temp_key)'); + + $pdo->exec( + 'CREATE TABLE IF NOT EXISTS repo_block ( + cid TEXT PRIMARY KEY, + repo_rev TEXT NOT NULL, + size INTEGER NOT NULL, + content BLOB NOT NULL + )' + ); + $pdo->exec('CREATE INDEX IF NOT EXISTS repo_block_repo_rev_idx ON repo_block (repo_rev)'); + + $pdo->exec( + 'CREATE TABLE IF NOT EXISTS repo_root ( + did TEXT PRIMARY KEY, + cid TEXT NOT NULL, + rev TEXT NOT NULL, + indexed_at TEXT NOT NULL + )' + ); + + $pdo->exec( + 'CREATE TABLE IF NOT EXISTS account_pref ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + name TEXT NOT NULL, + value_json TEXT NOT NULL + )' + ); + $pdo->exec('CREATE INDEX IF NOT EXISTS account_pref_name_idx ON account_pref (name)'); + } +} diff --git a/src/Infrastructure/Database/Schema/DidCacheSchema.php b/src/Infrastructure/Database/Schema/DidCacheSchema.php new file mode 100644 index 0000000..c6b3ce7 --- /dev/null +++ b/src/Infrastructure/Database/Schema/DidCacheSchema.php @@ -0,0 +1,23 @@ +pdo(); + + $pdo->exec( + 'CREATE TABLE IF NOT EXISTS did_doc ( + did TEXT PRIMARY KEY, + doc_json TEXT NOT NULL, + updated_at TEXT NOT NULL + )' + ); + } +} diff --git a/src/Infrastructure/Database/Schema/SequencerSchema.php b/src/Infrastructure/Database/Schema/SequencerSchema.php new file mode 100644 index 0000000..20f5d10 --- /dev/null +++ b/src/Infrastructure/Database/Schema/SequencerSchema.php @@ -0,0 +1,27 @@ +pdo(); + + $pdo->exec( + 'CREATE TABLE IF NOT EXISTS repo_seq ( + seq INTEGER PRIMARY KEY AUTOINCREMENT, + did TEXT NOT NULL, + event_type TEXT NOT NULL, + event BLOB NOT NULL, + sequenced_at TEXT NOT NULL, + invalidated INTEGER NOT NULL DEFAULT 0 + )' + ); + $pdo->exec('CREATE INDEX IF NOT EXISTS repo_seq_did_idx ON repo_seq (did)'); + } +}