diff --git a/docs/plans/2026-10-08-1451-framework-lease-cleanup.md b/docs/plans/2026-10-08-1451-framework-lease-cleanup.md
new file mode 100644
index 000000000..25793c9e4
--- /dev/null
+++ b/docs/plans/2026-10-08-1451-framework-lease-cleanup.md
@@ -0,0 +1,141 @@
+# Database lease cleanup and response registration
+
+## Outcome and boundaries
+
+Correct the lifecycle defects reported on [PR #653](https://github.com/hypervel/components/pull/653): expired logical owners can silently borrow an unowned pool slot, read-side extensions can enter the wrong construction path, and raw transactions can reach another borrower. Remove unused response stream-ID plumbing. Preserve automatic idle release, familiar Laravel database APIs, callback behavior and inexpensive normal queries.
+
+Work in `/home/binaryfire/workspace/contrib/hypervel/components`, branch `fix/database-lease-lifecycle`, created from `0.4`. AI package work remains paused in its separate worktree. This is a separate framework PR to complete before returning to AI implementation.
+
+**Status:** implementation, verification and review are complete, including failed-initial-owner cleanup. Complete both PRs' remaining CI and bot-review pass before merging. Follow monorepo `CLAUDE.md` and this repository's `AGENTS.md`. The peer is `claude-ai`; do not compact the peer during this work. Leave the historical framework plan unchanged; track the paused AI work's dependency in its existing orchestration document.
+
+Fix confirmed failures at their owning boundary without adding automatic ownership transfer, per-query coroutine checks, SQL parsing, retry machinery or unrelated refactoring. Any newly proposed Laravel API break requires owner approval.
+
+## Evidence and review coverage
+
+| Albert's comments | Verified cause | Owning change |
+|---|---|---|
+| [Lease reacquisition](https://github.com/hypervel/components/pull/653#discussion_r4218935407), [resolver cleanup](https://github.com/hypervel/components/pull/653#discussion_r4218946335) | `ConnectionLease::resolvePdo()` borrows after terminal release; its original deferred cleanup has already run, and no new owner is registered. | Give the logical lease an explicit terminal state. |
+| [Read extensions](https://github.com/hypervel/components/pull/653#discussion_r4219011376) | Pool eligibility inspects top-level configuration; explicit `::read` construction can invoke an extension selected from a read record. Rebuilding that result as a PDO lease bypasses the extension or receives a non-PDO connection. | Classify all effective read configurations before enabling leases. |
+| [Raw transactions](https://github.com/hypervel/components/pull/653#discussion_r4219032223) | Pinning and pool rollback use framework counters. Native PDO and SQL-started transactions bypass those counters; the public `inTransaction()` checks only the writer. | Inspect both physical handles at release boundaries and promptly clean dirty sessions. |
+| [Unused stream ID](https://github.com/hypervel/components/pull/653#discussion_r4219075550) | Response registration stores a stream ID which neither cancellation nor cleanup reads. | Remove the property and argument forwarding. |
+
+Diagnostic probes at `/tmp/hypervel-ai-port/AlbertReviewProbeTest.php` and `albert-review-probes.log` reproduce the orphaned borrow, both read-extension failures and raw-transaction leakage on early and terminal release. They assert the observed defects, not the desired behavior; do not copy those assertions into regressions. A two-PDO probe also confirms a read transaction is invisible to the current connection's `inTransaction()` and `hasPinnedSession()`. Permanent tests below must assert the corrected behavior and fail before the fix.
+
+## 1. End logical ownership explicitly
+
+`ConnectionLease` owns the logical connection, its PDO resolvers and its current pool borrow. The logical connection retains transaction bookkeeping, callbacks, query logs and sticky-routing state. Attaching another cleanup callback alone would not define how that state transfers to another execution or make overlapping use safe. End this resource's ownership cleanly instead of silently reopening it after cleanup.
+
+- Add a protected boolean `$ended = false` to `ConnectionLease`.
+- Make the lease's existing `release()` and `discard()` terminal: set `$ended` before settling the current pooled connection. Repeated terminal calls with no held slot remain harmless.
+- Keep the two internal non-terminal paths recoverable: `releaseIfIdle()` calls `$this->pooledConnection?->release()` directly, while `discardAfterFailure()` calls `$this->pooledConnection?->discard()` directly and keeps its exception precedence.
+- In `resolvePdo()`, check `$ended` only inside the branch about to borrow a new slot. Throw `LogicException` with an actionable message explaining that the owning execution ended and the caller must resolve its connection where it will be used. Do not add checks to every query.
+- Failed initial publication must discard the logical owner terminally; fall back to the borrowed wrapper if construction did not return an owner. The coroutine defer, non-coroutine task release/discard and replacement of a retained non-coroutine owner keep their existing terminal calls. Keep connections registered in context until their release/rollback callbacks finish.
+- Release listeners and rollback callbacks may query or reconnect the still-held slot. Set the terminal flag before settlement so a newly awakened consumer cannot reacquire after detachment; do not reject ordinary access to the slot while cleanup still owns it.
+- Ordinary early release, disconnect/reconnect and failed reacquisition remain reusable within a registered live owner. Do not mark every physical discard terminal.
+
+Core branch, within the existing resolver:
+
+```php
+if ($this->pooledConnection === null) {
+ if ($this->ended) {
+ throw new LogicException('This database connection is no longer available because the coroutine or task that resolved it has finished or failed to set it up. Resolve the connection where you use it.');
+ }
+
+ // Existing borrow and attachment.
+}
+```
+
+This prevents the reported orphaned acquisition. It does not claim to intercept an already-running cursor's PDOStatement, an externally retained raw PDO, or every concurrent reference to a connection.
+
+## 2. Classify effective read extensions
+
+Consolidate the constructor eligibility expression into the existing identity-checking method in `DatabasePool`, renamed to `supportsSessionLeases()` to reflect its full responsibility.
+
+1. Preserve name/top-level extension precedence: an extension returned by `ConnectionFactory::getExtension($config, $name->base)` disables leases immediately.
+2. Select `read` only for an explicit `::read` pool with read configuration; otherwise inspect `write`. If no selected role record exists, the top-level classification suffices. Do not call `configForWrite()` on absent write configuration.
+3. Normalize a single associative record to a one-element list; inspect every record in a configured list. Resolve each effective endpoint using the existing `configForRead()` or `configForWrite()` merge/parser. This inspection must not randomly choose one record for the pool lifetime.
+4. For the read role, any effective endpoint with an extension disables leases for the entire pool. Preserve the existing driver/database/prefix identity comparison across selectable records.
+5. Do not check write-record extensions when deciding lease eligibility: base construction invokes extensions from top-level configuration, then builds PDO write endpoints through `createConnection()` / `Connection::resolverFor()`. Preserve this behavior rather than disabling valid leases unnecessarily.
+
+Keep read pool-option consistency, SQLite restrictions, per-physical-connection endpoint selection and reconnect behavior intact. Eligibility is decided at pool construction, not per query or dynamically after borrowing. Whole-connection extensions retain their existing object and early-release behavior; no new extension interface is required.
+
+## 3. Protect and settle physical transactions
+
+### Pinning without changing public transaction semantics
+
+Add public internal `Connection::hasPhysicalTransaction(): bool`, with a Laravel-style title and `@internal` docblock. Its base implementation delegates to the existing `inTransaction()` contract, so custom whole-connection drivers can report their physical state. No `ConnectionInterface` addition is needed: pool owners already use concrete `Connection`.
+
+Override it in `PdoConnection` to inspect both already-open, distinct handles:
+
+```php
+return $this->inTransaction()
+ || ($this->readPdo instanceof PDO
+ && $this->readPdo !== $this->pdo
+ && $this->readPdo->inTransaction());
+```
+
+Append `|| $this->hasPhysicalTransaction()` to `Connection::hasPinnedSession()` after its existing pin, transaction-counter and FK-suppression checks. Do not resolve a closure or open an unused connection just to inspect it. Leave public `inTransaction()` writer semantics unchanged; `RefreshDatabase` uses that method to assess its writer transaction. Keep public `beginTransaction()`, `commit()`, `rollBack()` and transaction counters unchanged. SQL parsing or synthetic counters are unnecessary.
+
+PDO reports transactions started through its native methods or SQL transaction statements; the new inspection automatically pins those sessions across manual and framework-HTTP release. Raw PDO/statement use outside a transaction still needs explicit pinning when it requires session continuity.
+
+### Inspect after callbacks, then settle once
+
+Preserve `PooledConnection::prepareForRelease()`'s tracked `rollBack(0)` and rollback callbacks. In `release()`, inspect the physical holder after preparation, release listeners and logical-lease detachment, including when preparation threw or was cancelled. The physical holder retains the owned PDO objects; inspecting after callbacks catches a transaction left by a listener as well as one left by application code.
+
+| Final physical-status result | Settlement |
+|---|---|
+| Transaction remains | Decide to discard, log the existing unfinished-transaction error, and call `pool->discard($this)`. |
+| Inspection throws or is cancelled | Decide to discard, retain the failure under existing exception precedence, and still settle. |
+| No transaction | Preserve existing release and reuse checks. An invalid slot from a preparation failure can return to idle; it reconnects before any subsequent application use. |
+
+Do not skip physical inspection merely because an earlier cancellation was captured. Commit the discard decision before invoking logging that might throw. Use one settlement decision, not overlapping flags or an extra cleanup state machine. Run settlement exactly once; cancellation outranks ordinary errors and existing primary-failure rules remain intact.
+
+The existing discard path calls `destroyConnection()` → `close()` → disconnect, reports ordinary close failures, propagates cancellation and frees pool accounting in `finally`. Reuse it. Merely marking a dirty slot invalid is insufficient: it could hold transaction locks while sitting idle until another borrow. Conversely, discarding every failed preparation with no physical transaction would change behavior without addressing this defect.
+
+### Disconnect both handles
+
+Extend `PdoConnection::disconnectDriverResources()` to inspect and roll back both already-open distinct PDOs. Attempt cleanup of the second even when inspecting or rolling back the first fails. Invalidate the physical-session memo after successful rollback; mark it unknown on failure. Preserve existing treatment of lost-connection failures as already disconnected, retain the first other failure, and give cancellation precedence. Always call `forgetDriverResources()` in `finally`; retain `Connection::disconnect()`'s transaction-manager cleanup and failure ordering.
+
+This also fixes read-side cleanup during explicit disconnect, reconnect and discard. Do not add a separate rollback API or duplicate the pool's close machinery. The pool-held PDO for shared in-memory SQLite must survive a successful rollback/discard of its wrapper, preserving committed data and schema. A failed cleanup must never make an unknown session reusable.
+
+## 4. Remove unused stream identity
+
+- Remove `ResponseCancellation::$stream`, its constructor argument and the unused `register()` argument. Keep registration indexed by connection and producer coroutine, identity-safe removal, cancellation behavior and static cleanup.
+- Remove `ResponseBridge::send()`'s unused `streamId` argument and forwarding from HTTP, gRPC and WebSocket servers, including both gRPC call sites and `tests/HttpServer/Fixtures/disconnect-server.php`. Update every source and test caller; do not remove stream IDs used by unrelated HTTP/2 clients.
+- Retain coverage for multiple simultaneous producers on one connection, cancellation opt-in, sibling connections and post-production cleanup. Collapse `ResponseBridgeTest::disconnectCancellationOptions` to one opt-in and one opt-out row; record producer labels in `$closed` instead of unused stream IDs.
+- Update the Swoole stream-cancel entry in `docs/todo.md` to say stream identity will be introduced into response registrations with the released event. Keep existing-runtime limitations and the future integration task; no polling, speculative event API or cancellation redesign.
+
+## 5. Documentation and compatibility
+
+Update `src/docs/database.md` in its existing pooling/releasing sections, using Laravel-style prose. Explain the new-acquisition exception after ownership ends, while retaining ordinary builder reuse across early release. Explain automatic pinning of native/SQL transactions and cleanup of unfinished transactions at execution end. Keep manual pinning advice for other session-dependent operations. Do not claim every physical PDO is replaced: shared in-memory SQLite retains its pool-owned PDO.
+
+In the whole-connection ownership paragraph, include differing drivers alongside database names and table prefixes, and make the `DB::extend` sentence cover effective read-record extensions in explicit `::read` pools. Audit the database README and `src/docs/porting-from-laravel.md` against these changes; their existing promise that active transactions remain pinned becomes accurate, so do not add redundant internal-fix entries. These changes do not alter Laravel public signatures or intended successful results. The removed response arguments are unused Hypervel internal integration plumbing.
+
+Remove stale descriptions, dead arguments and superseded method references from the modified surfaces. Preserve unrelated docs and all still-open TODOs. Do not rewrite historical plans or document review history as product behavior.
+
+## 6. Tests and verification
+
+Prefer additions to existing tests and data providers. Make each new regression fail on the original source and pass after its fix; run each changed test file immediately. Do not introduce a separate test architecture or a driver-by-driver copy of the lease suite.
+
+| Existing test area | Required coverage |
+|---|---|
+| `DatabaseConnectionLeaseLifecycleTest::testTerminalCallbacksResolveTheirOwningConnection` | Retain a builder before cleanup; after each existing coroutine/task/failing-listener row, its next acquisition throws and borrowed count stays zero. Preserve callback access during cleanup. |
+| Existing lease recovery/task-cleanup tests | Retained builder survives early release and another borrower; repeated early release/task cleanup; listener failure during reacquisition remains recoverable; disconnect/reconnect still works. Initial publication failure or cancellation leaves its retained connection unable to borrow, while fresh resolution succeeds. Extend the existing task test's final `discardConnections()` section to prove that retained owner cannot borrow again either; listener mocks alone do not prove this state transition. |
+| `DatabaseConnectionLeaseTest::testConfigFirstExtensionsRetainWholeConnectionOwnership` | Preserve top-level case; add single and listed read-record PDO extension cases. Assert returned extension identity, disabled leases and unaffected early release. Retain mixed-identity, read-selection and existing non-PDO whole-connection tests. |
+| `DatabaseConnectionLeaseTest::testSessionDependentScopesPreventEarlyRelease` | Extend rows for native raw transaction, SQL `BEGIN`, and distinct read-PDO transaction; parameterize the connection name. Check manual and HTTP-triggered release retains the slot, then release succeeds after transaction settlement. Retain explicit/FK/tracked transaction rows. |
+| Raw transaction lifecycle coverage in `DatabaseConnectionLeaseLifecycleTest` | Coroutine/task cleanup, including a transaction opened by a release listener. Assert physical transaction ends before the next borrow, uncommitted data disappears, and pre-transaction schema/committed data survives shared-memory SQLite wrapper replacement. |
+| Existing PDO disconnect tests | Two distinct active handles; one rollback failure does not skip the other; cancellation precedence; lazy handles are not resolved; identical handles are not processed twice. Preserve lost-connection and transaction-manager cleanup tests. Use focused rows/assertions, not a Cartesian matrix. |
+| `PooledConnectionTest::testReleaseRollsBackOpenTransactions` | Add a raw-transaction case for whole-connection ownership without a lease; assert rollback and discard after `resetForPool()`. |
+| Existing `PooledConnectionTest` release-failure tests | A throwing physical-status inspection on a non-PDO fixture settles by discard; logging cannot undo the discard decision (reuse the throwing logger from `testReleasePreservesTheFirstOrdinaryCleanupFailure`); cancellation still settles once. No-transaction preparation failure retains existing invalid/requeue behavior. |
+| Existing response bridge/server tests | Multiple producers on one connection all cancel, another connection stays usable, opt-out and completed production remain unaffected. Update removed-argument callers and eliminate now-identical cases. |
+
+Use the existing real SQLite-backed fixtures for ownership tests and existing mocked PDOs for precise failure paths. Also run `bin/run-database-tests.sh` sequentially for `mysql`, `mariadb` and `pgsql` with their matching local `DB_*` configuration. Do not add new service-specific lease classes just to repeat PDO transaction detection. Framework-wide cleanup changes require full verification before completion, using configured service isolation and established dependency skips.
+
+Implementer workflow after owner authorization:
+
+1. Implement coherent changes, tracing callers/callees and retaining the agreed ownership invariants. Investigate non-trivial new findings as a group and reach peer consensus before changing their design; amend this active plan concisely if needed.
+2. After the immediate changed-file tests, run `composer fix` once for these shared lifecycle changes: formatting, PHPStan, parallel tests, Testbench and dogfood checks. Do not duplicate its component checks beforehand or overlap verification commands. Follow AGENTS' targeted-rerun rules on failure, then self-review the whole diff against this plan and Albert's comments.
+3. Request complete code review from `claude-ai`, resolve findings until signoff, and report verification. The reviewer does not rerun checks or edit code. Do not compact the peer.
+4. At the authorized commit/PR boundary, use coherent whole-file commits with detailed bodies, excluding unrelated edits. Verify authorization for publication rather than inferring it from plan creation. After the follow-up PR merges, reply to each of Albert's five code threads with a link to its fix, within the owner's authorized publication scope. Incorporate the merged `0.4` changes into `feature/ai` before resuming AI package work.
+
+Normal query execution gains no per-query ownership checks. The ended-owner branch is tested only before a new borrow; extension classification runs once per pool; physical transaction checks occur at release/cleanup boundaries and create no new PDOs or database round trips for the built-in drivers. If implementation adds meaningful extra hot-path work or measurements become necessary, agree an idle window with the owner before comparative benchmarks.
diff --git a/docs/todo.md b/docs/todo.md
index ab797b282..c8900cc57 100644
--- a/docs/todo.md
+++ b/docs/todo.md
@@ -51,7 +51,7 @@
## HTTP Server
-- Integrate a server-level HTTP/2 stream-cancel event once Swoole exposes it in a supported release, and raise the `ext-swoole` constraint to that release when integrating it. The current `RST_STREAM` path removes the stream without notifying PHP, so an AI response producer waiting on a silent provider cannot stop immediately when its caller cancels that stream. Use the event's connection and stream identities with Hypervel's active-response map to cancel only the affected opted-in producer. Preserve sibling streams, application close callbacks and terminable work. Whole-connection closes continue using existing onClose handling; no native per-response subscription is needed. Keep the AI package usable on existing supported Swoole releases with connection-close handling where available, failed-write cleanup and provider timeouts; document the remaining HTTP/2 limitation. This enhancement is not an AI-port release prerequisite. Use the actual released event name and arguments rather than polling or anticipating an unreleased API.
+- Integrate a server-level HTTP/2 stream-cancel event once Swoole exposes it in a supported release, and raise the `ext-swoole` constraint to that release when integrating it. The current `RST_STREAM` path removes the stream without notifying PHP, so an AI response producer waiting on a silent provider cannot stop immediately when its caller cancels that stream. Add stream identity to Hypervel's active-response registrations when integrating the released event, then use its connection and stream identities to cancel only the affected opted-in producer. Preserve sibling streams, application close callbacks and terminable work. Whole-connection closes continue using existing onClose handling; no native per-response subscription is needed. Keep the AI package usable on existing supported Swoole releases with connection-close handling where available, failed-write cleanup and provider timeouts; document the remaining HTTP/2 limitation. This enhancement is not an AI-port release prerequisite. Use the actual released event name and arguments rather than polling or anticipating an unreleased API.
- Require a Swoole release that resets signal-listener state in forked server workers before releasing Hypervel 0.4. In Swoole 6.2.3 and earlier, a worker forked after the manager calls `Process::signal()` inherits the listener count, so `Coroutine\System::waitSignal()` fails in it. Hypervel's SIGINT shutdown handling registers a manager callback in both server modes, so after a reload, `max_request` recycling or a crash restart, replacement workers stop receiving configured signal handlers and Artisan traps. Once a fixed release is verified, raise the `ext-swoole` constraint and remove the version skip from `ShutdownOnInterruptListenerTest::testReplacementWorkersKeepTheirSignalHandlers()`.
- Remove trailer-stream one-chunk lookahead once the minimum supported Swoole release includes [swoole-src#6124](https://github.com/swoole/swoole-src/pull/6124). Current releases send an empty `END_STREAM` DATA frame before trailer HEADERS when `end()` receives no body after `write()`, so `ResponseBridge` retains the final chunk for `end($chunk)` and delays delivery by one chunk. Once fixed, raise the `ext-swoole` constraint, write every chunk immediately, emit trailers, call bare `end()`, invert the deterministic bridge ordering tests, and add real gRPC incremental-delivery coverage.
diff --git a/src/database/src/Connection.php b/src/database/src/Connection.php
index be13df0f2..6933d1616 100755
--- a/src/database/src/Connection.php
+++ b/src/database/src/Connection.php
@@ -1101,7 +1101,8 @@ public function hasPinnedSession(): bool
{
return $this->sessionPinDepth > 0
|| $this->transactions > 0
- || $this->foreignKeyConstraintSuppressionDepth > 0;
+ || $this->foreignKeyConstraintSuppressionDepth > 0
+ || $this->hasPhysicalTransaction();
}
/**
@@ -1614,6 +1615,16 @@ private function throwUnsupportedTransactionException(): never
*/
abstract public function inTransaction(): bool;
+ /**
+ * Determine whether any owned driver resource has an active transaction.
+ *
+ * @internal
+ */
+ public function hasPhysicalTransaction(): bool
+ {
+ return $this->inTransaction();
+ }
+
/**
* Set the transaction manager instance on the connection.
*/
diff --git a/src/database/src/ConnectionResolver.php b/src/database/src/ConnectionResolver.php
index 24145bcb3..01ca15c68 100755
--- a/src/database/src/ConnectionResolver.php
+++ b/src/database/src/ConnectionResolver.php
@@ -159,7 +159,7 @@ public function connection(UnitEnum|string|null $name = null): ConnectionInterfa
CoroutineContext::forget($leaseContextKey);
unset($this->nonCoroutineConnections[$connectionOwnerName]);
- $this->discardFailedConnection($pooledConnection, $exception);
+ $this->discardFailedConnection($owner ?? $pooledConnection, $exception);
throw $exception;
}
@@ -252,12 +252,12 @@ protected function getLeaseContextKey(string $name): string
}
/**
- * Discard a failed connection while preserving cancellation precedence.
+ * Discard a failed connection owner while preserving cancellation precedence.
*/
- protected function discardFailedConnection(PooledConnection $pooledConnection, Throwable $exception): void
+ protected function discardFailedConnection(ConnectionLease|PooledConnection $owner, Throwable $exception): void
{
try {
- $pooledConnection->discard();
+ $owner->discard();
} catch (CanceledException $cancellation) {
if (! $exception instanceof CanceledException) {
throw $cancellation;
diff --git a/src/database/src/PdoConnection.php b/src/database/src/PdoConnection.php
index 5e1b575cb..08b03cfa5 100755
--- a/src/database/src/PdoConnection.php
+++ b/src/database/src/PdoConnection.php
@@ -767,18 +767,30 @@ protected function forgetDriverResources(): void
protected function disconnectDriverResources(): void
{
$pdo = $this->getRawPdo();
+ $readPdo = $this->getRawReadPdo();
$exception = null;
try {
- if ($pdo instanceof PDO && $pdo->inTransaction()) {
- $pdo->rollBack();
- $this->invalidateSessionState($pdo);
- }
- } catch (Throwable $throwable) {
- $this->markSessionStateUnknown($pdo);
+ foreach ($readPdo === $pdo ? [$pdo] : [$pdo, $readPdo] as $handle) {
+ if (! $handle instanceof PDO) {
+ continue;
+ }
- if (! $this->causedByLostConnection($throwable)) {
- $exception = $throwable;
+ try {
+ if ($handle->inTransaction()) {
+ $handle->rollBack();
+ $this->invalidateSessionState($handle);
+ }
+ } catch (Throwable $throwable) {
+ $this->markSessionStateUnknown($handle);
+
+ if (! $this->causedByLostConnection($throwable)
+ && ($exception === null
+ || ($throwable instanceof CanceledException && ! $exception instanceof CanceledException))
+ ) {
+ $exception = $throwable;
+ }
+ }
}
} finally {
$this->forgetDriverResources();
@@ -880,6 +892,19 @@ public function inTransaction(): bool
return $this->pdo instanceof PDO && $this->pdo->inTransaction();
}
+ /**
+ * Determine whether either open PDO handle has an active transaction.
+ *
+ * @internal
+ */
+ public function hasPhysicalTransaction(): bool
+ {
+ return $this->inTransaction()
+ || ($this->readPdo instanceof PDO
+ && $this->readPdo !== $this->pdo
+ && $this->readPdo->inTransaction());
+ }
+
/**
* Run the statement to start a new transaction.
*/
diff --git a/src/database/src/Pool/ConnectionLease.php b/src/database/src/Pool/ConnectionLease.php
index 690b85436..71a81d830 100644
--- a/src/database/src/Pool/ConnectionLease.php
+++ b/src/database/src/Pool/ConnectionLease.php
@@ -8,6 +8,7 @@
use Hypervel\Context\NonCopyableContext;
use Hypervel\Database\Connectors\ConnectionFactory;
use Hypervel\Database\PdoConnection;
+use LogicException;
use PDO;
use Swoole\Coroutine\CanceledException;
use Throwable;
@@ -21,6 +22,8 @@ class ConnectionLease implements NonCopyableContext
protected ?PooledConnection $pooledConnection;
+ protected bool $ended = false;
+
/** @var Closure(): PDO */
protected readonly Closure $pdoResolver;
@@ -57,6 +60,10 @@ public function __construct(
protected function resolvePdo(bool $read = false): PDO
{
if ($this->pooledConnection === null) {
+ if ($this->ended) {
+ throw new LogicException('This database connection is no longer available because the coroutine or task that resolved it has finished or failed to set it up. Resolve the connection where you use it.');
+ }
+
/** @var PooledConnection $pooledConnection */
$pooledConnection = $this->pool->borrow();
$this->pooledConnection = $pooledConnection;
@@ -113,23 +120,25 @@ public function reconnect(): PdoConnection
public function releaseIfIdle(): void
{
if (! $this->connection->hasPinnedSession()) {
- $this->release();
+ $this->pooledConnection?->release();
}
}
/**
- * Settle the currently held physical session.
+ * End logical ownership and release the currently held physical session.
*/
public function release(): void
{
+ $this->ended = true;
$this->pooledConnection?->release();
}
/**
- * Discard the currently held physical session.
+ * End logical ownership and discard the currently held physical session.
*/
public function discard(): void
{
+ $this->ended = true;
$this->pooledConnection?->discard();
}
@@ -148,7 +157,7 @@ public function detach(): void
protected function discardAfterFailure(Throwable $exception): void
{
try {
- $this->discard();
+ $this->pooledConnection?->discard();
} catch (CanceledException $cancellation) {
if (! $exception instanceof CanceledException) {
throw $cancellation;
diff --git a/src/database/src/Pool/DatabasePool.php b/src/database/src/Pool/DatabasePool.php
index 86216b125..aa3ca6ba6 100644
--- a/src/database/src/Pool/DatabasePool.php
+++ b/src/database/src/Pool/DatabasePool.php
@@ -72,8 +72,7 @@ public function __construct(Container $container, string $name)
}
$this->config = $config;
- $this->usesSessionLeases = $factory->getExtension($config, $connectionName->base) === null
- && $this->hasConsistentLogicalIdentity($factory, $connectionName, $config);
+ $this->usesSessionLeases = $this->supportsSessionLeases($factory, $connectionName, $config);
$poolOptions = Arr::except(
Arr::get($poolConfig, 'pool', []),
@@ -131,23 +130,33 @@ public function usesSessionLeases(): bool
}
/**
- * Determine whether selectable endpoints share the same logical database identity.
+ * Determine whether endpoints can share a logical connection without bypassing extensions.
*/
- protected function hasConsistentLogicalIdentity(ConnectionFactory $factory, ConnectionName $name, array $config): bool
+ protected function supportsSessionLeases(ConnectionFactory $factory, ConnectionName $name, array $config): bool
{
+ if ($factory->getExtension($config, $name->base) !== null) {
+ return false;
+ }
+
$role = $name->isRead() && $factory->hasReadConfig($config) ? 'read' : 'write';
- if (! isset($config[$role][0])) {
+ if (! isset($config[$role])) {
return true;
}
+ $records = isset($config[$role][0]) ? $config[$role] : [$config[$role]];
$identity = null;
- foreach ($config[$role] as $record) {
+ foreach ($records as $record) {
$candidate = array_replace($config, [$role => $record]);
$endpoint = $role === 'read'
? $factory->configForRead($candidate)
: $factory->configForWrite($candidate);
+
+ if ($role === 'read' && $factory->getExtension($endpoint, $name->base) !== null) {
+ return false;
+ }
+
$candidateIdentity = [$endpoint['driver'], $endpoint['database'], $endpoint['prefix']];
if ($identity !== null && $identity !== $candidateIdentity) {
diff --git a/src/database/src/Pool/PooledConnection.php b/src/database/src/Pool/PooledConnection.php
index 8ac998e2d..c6457c9aa 100644
--- a/src/database/src/Pool/PooledConnection.php
+++ b/src/database/src/Pool/PooledConnection.php
@@ -480,8 +480,25 @@ public function release(): void
$this->lease = null;
}
+ $discard = false;
+
+ try {
+ // Callbacks may leave a raw transaction outside the framework counters.
+ if ($this->connection?->hasPhysicalTransaction()) {
+ $discard = true;
+ $this->logger->error('Database transaction was not committed or rolled back before release.');
+ }
+ } catch (CanceledException $transactionCancellation) {
+ $discard = true;
+ $cancellationFailure ??= $transactionCancellation;
+ } catch (Throwable $exception) {
+ $discard = true;
+ $ordinaryFailure ??= $exception;
+ }
+
try {
- if ($cancellationFailure === null
+ if (! $discard
+ && $cancellationFailure === null
&& $this->connection !== null
&& ! $this->connection->isReusable()
) {
@@ -494,10 +511,14 @@ public function release(): void
$ordinaryFailure ??= $exception;
}
- $this->availableForReuse = true;
+ $this->availableForReuse = ! $discard;
try {
- $this->pool->release($this);
+ if ($discard) {
+ $this->pool->discard($this);
+ } else {
+ $this->pool->release($this);
+ }
} catch (CanceledException $releaseCancellation) {
$cancellationFailure ??= $releaseCancellation;
} catch (Throwable $exception) {
diff --git a/src/docs/database.md b/src/docs/database.md
index acd70dfbb..3149eb223 100644
--- a/src/docs/database.md
+++ b/src/docs/database.md
@@ -233,9 +233,9 @@ The `sticky` option is an *optional* value that can be used to allow the immedia
### Connection Pooling
-Hypervel uses connection pools to keep database access efficient within long-lived Swoole workers. When a coroutine resolves a database connection, Hypervel borrows a physical session from the worker's pool. The coroutine keeps its own connection object, including query logs, callbacks and read / write routing state. Query and schema builders retain that object even if its idle session is returned to the pool and another session is borrowed later. When the coroutine ends, Hypervel rolls back unfinished transactions and returns its remaining borrowed sessions.
+Hypervel uses connection pools to keep database access efficient within long-lived Swoole workers. When a coroutine resolves a database connection, Hypervel borrows a physical session from the worker's pool. The coroutine keeps its own connection object, including query logs, callbacks and read / write routing state. Query and schema builders retain that object even if its idle session is returned to the pool and another session is borrowed later. When the coroutine or task ends, Hypervel rolls back unfinished transactions, including those started directly through PDO or SQL, and returns its remaining borrowed sessions.
-A connection and the query builders created from it belong to the coroutine that resolved them. Build queries inside each child coroutine instead of passing builders or connections between coroutines, and do not keep them after their coroutine ends.
+A connection and the query builders created from it belong to the coroutine or task that resolved them. Build queries inside each child coroutine instead of passing builders or connections between coroutines, and do not keep them after the coroutine or task finishes. Connections that support early release remain usable across releases within that execution, but throw an exception if asked to borrow another session after it finishes.
Each connection may define its own `pool` configuration:
@@ -271,7 +271,7 @@ For the full option reference and custom maintenance behavior, see the [pool doc
For a connection with separate read and write hosts, each base pool slot may lazily open one write PDO and one read PDO. It does not open one PDO per configured host. If `max_connections` is `10`, a worker may therefore hold up to roughly 20 server-side database connections for that configured connection once both sides have been used. Size your database server, PgBouncer, PgDog, or other pooler capacity with that in mind. Increase `max_connections` for more concurrent database work per worker, not simply because you configured more read hosts.
-When a read side is configured, explicit `::read` connections use a separate read-side pool built from the merged read configuration, including the base `pool` settings unless the read configuration overrides them. Drivers registered through `DB::extend` still receive the complete connection configuration so they can select their own endpoints; the pool's options come from the merged read configuration. Without a read side, `::read` uses the base pool. Explicit `::write` connections do not create a separate pool, but a coroutine that uses both `mysql` and `mysql::write` at the same time may borrow two slots from the base pool. Most applications do not need these suffixes in normal query paths because Hypervel already routes reads, writes, transactions, and sticky reads automatically.
+When a read side is configured, explicit `::read` connections use a separate read-side pool built from the merged read configuration, including the base `pool` settings unless the read configuration overrides them. Extensions registered through `DB::extend` for the connection name or its top-level driver receive the complete connection configuration so they can select their own endpoints. An extension selected by a read record receives that record's merged configuration. The pool's options come from the merged read configuration. Without a read side, `::read` uses the base pool. Explicit `::write` connections do not create a separate pool, but a coroutine that uses both `mysql` and `mysql::write` at the same time may borrow two slots from the base pool. Most applications do not need these suffixes in normal query paths because Hypervel already routes reads, writes, transactions, and sticky reads automatically.
If your `read` configuration contains a list of connection records, Hypervel chooses a record for each new physical connection and chooses again when reconnecting. All records in an explicit `::read` pool must have the same effective `pool` options, including inherited defaults. Conflicting options throw an exception when the pool is created, since the records share one pool capacity and lifecycle policy. Derived read pools cannot use in-memory SQLite databases.
@@ -312,7 +312,7 @@ $responses = parallel([
Automatic release occurs before sending a request, not on every read of a streamed response. If your stream consumer performs database work between chunks, you may call `releaseIdleConnections()` before waiting for more data. Faked HTTP requests do not release connections. If you provide your own Guzzle client using `setClient`, call `DB::releaseIdleConnections()` yourself before sending requests.
-Active queries, transactions, open cursors and `Schema::withoutForeignKeyConstraints` callbacks retain their sessions automatically. In particular, an HTTP call inside `DB::transaction()` keeps the transaction's connection. Completed `chunk` or `lazy` query batches may release between batches.
+Active queries, transactions, open cursors and `Schema::withoutForeignKeyConstraints` callbacks retain their sessions automatically. This includes transactions started directly through PDO or SQL, on either the read or write connection. In particular, an HTTP call inside `DB::transaction()` keeps the transaction's connection. Completed `chunk` or `lazy` query batches may release between batches.
If your code requires the same physical session across an external call, wrap the entire operation in `DB::withPinnedSession()`. For a named connection, call `withPinnedSession()` on that connection. For example, a PostgreSQL session-level advisory lock must be released on the session that acquired it:
@@ -332,7 +332,7 @@ $response = $connection->withPinnedSession(function () use ($connection, $accoun
Pinning is also needed for temporary tables, retained raw PDOs or statements, and manual session changes spanning a release boundary. It applies to a manual `disableForeignKeyConstraints` / `enableForeignKeyConstraints` pair; prefer the scoped `withoutForeignKeyConstraints` method. Each pin protects its connection until the callback returns or throws, and pins may be nested. Complete any lazy work inside the callback rather than returning it for later execution. Pinning prevents early release but does not prevent replacing a broken connection during the normal lost-connection retry.
-Connections registered through `DB::extend` retain their complete driver object until execution ends; early release leaves them alone. The same applies when selectable write records (or read records in an explicit `::read` pool) specify different database names or table prefixes. PDO drivers registered through `Connection::resolverFor` participate in early release when those values agree.
+Connections registered through `DB::extend` retain their complete driver object until execution ends; early release leaves them alone. This includes extensions selected by read records in an explicit `::read` pool. The same applies when selectable write records (or read records in an explicit `::read` pool) specify different drivers, database names or table prefixes. PDO drivers registered through `Connection::resolverFor` participate in early release when those values agree.
### Configuring Database Session State
diff --git a/src/grpc/src/Server/Server.php b/src/grpc/src/Server/Server.php
index 1f2385867..2eea497b2 100644
--- a/src/grpc/src/Server/Server.php
+++ b/src/grpc/src/Server/Server.php
@@ -135,7 +135,7 @@ public function onRequest(SwooleRequest $swooleRequest, SwooleResponse $swooleRe
}
$emissionStarted = true;
- ResponseBridge::send($response, $swooleResponse, protocol: 'HTTP/2', request: $request ?? null, streamId: $swooleRequest->streamId ?? 0);
+ ResponseBridge::send($response, $swooleResponse, protocol: 'HTTP/2', request: $request ?? null);
} catch (CanceledException $exception) {
$cancelled = true;
@@ -158,7 +158,7 @@ public function onRequest(SwooleRequest $swooleRequest, SwooleResponse $swooleRe
try {
$response = $this->responses->error($this->exceptions->map($throwable));
$emissionStarted = true;
- ResponseBridge::send($response, $swooleResponse, protocol: 'HTTP/2', request: $request ?? null, streamId: $swooleRequest->streamId ?? 0);
+ ResponseBridge::send($response, $swooleResponse, protocol: 'HTTP/2', request: $request ?? null);
} catch (CanceledException $exception) {
$cancelled = true;
diff --git a/src/http-server/src/ResponseBridge.php b/src/http-server/src/ResponseBridge.php
index df1109f61..e5538d80e 100644
--- a/src/http-server/src/ResponseBridge.php
+++ b/src/http-server/src/ResponseBridge.php
@@ -46,7 +46,6 @@ public static function send(
bool $withBody = true,
string $protocol = 'HTTP/1.1',
?Request $request = null,
- int $streamId = 0,
): void {
if ($response instanceof HasTrailers && $response instanceof BinaryFileResponse) {
throw new RuntimeException('Binary file responses cannot emit trailers.');
@@ -84,7 +83,7 @@ public static function send(
throw new CanceledException('The client disconnected before response production.');
}
- $registration = ResponseCancellation::register($swooleResponse->fd, $streamId);
+ $registration = ResponseCancellation::register($swooleResponse->fd);
}
try {
diff --git a/src/http-server/src/Server.php b/src/http-server/src/Server.php
index 5ccd408a7..45869f25c 100644
--- a/src/http-server/src/Server.php
+++ b/src/http-server/src/Server.php
@@ -144,7 +144,6 @@ public function onRequest(SwooleRequest $swooleRequest, SwooleResponse $swooleRe
withBody: ! isset($rawMethod) || $rawMethod !== 'HEAD',
protocol: is_string($protocol) ? $protocol : 'HTTP/1.1',
request: $request ?? null,
- streamId: $swooleRequest->streamId ?? 0,
);
} catch (CanceledException $throwable) {
$cancellation = $throwable;
diff --git a/src/server/src/ResponseCancellation.php b/src/server/src/ResponseCancellation.php
index 8108a54bb..2fed0c657 100644
--- a/src/server/src/ResponseCancellation.php
+++ b/src/server/src/ResponseCancellation.php
@@ -21,7 +21,6 @@ class ResponseCancellation
*/
private function __construct(
protected int $connection,
- protected int $stream,
protected int $coroutine,
) {
}
@@ -29,7 +28,7 @@ private function __construct(
/**
* Register the current coroutine for the duration of response production.
*/
- public static function register(int $connection, int $stream = 0): ?self
+ public static function register(int $connection): ?self
{
$coroutine = Coroutine::id();
@@ -37,7 +36,7 @@ public static function register(int $connection, int $stream = 0): ?self
return null;
}
- return self::$responses[$connection][$coroutine] = new self($connection, $stream, $coroutine);
+ return self::$responses[$connection][$coroutine] = new self($connection, $coroutine);
}
/**
diff --git a/src/websocket-server/src/Server.php b/src/websocket-server/src/Server.php
index fdd1f388d..92ac47ee5 100644
--- a/src/websocket-server/src/Server.php
+++ b/src/websocket-server/src/Server.php
@@ -199,7 +199,7 @@ public function onHandshake(Request $request, SwooleResponse $response): void
}
try {
- ResponseBridge::send($httpResponse, $response, request: $httpRequest, streamId: $request->streamId ?? 0);
+ ResponseBridge::send($httpResponse, $response, request: $httpRequest);
} catch (CanceledException $exception) {
throw $exception;
} catch (Throwable $throwable) {
diff --git a/tests/Database/DatabaseConnectionLeaseLifecycleTest.php b/tests/Database/DatabaseConnectionLeaseLifecycleTest.php
index 0fbaeda60..17d02c840 100644
--- a/tests/Database/DatabaseConnectionLeaseLifecycleTest.php
+++ b/tests/Database/DatabaseConnectionLeaseLifecycleTest.php
@@ -8,9 +8,11 @@
use Hypervel\Contracts\Foundation\Application;
use Hypervel\Coroutine\Coroutine;
use Hypervel\Database\Pool\PoolManager;
+use Hypervel\Database\QueryException;
use Hypervel\Support\Facades\DB;
use Hypervel\Support\Facades\Event;
use Hypervel\Testbench\TestCase;
+use LogicException;
use PHPUnit\Framework\Attributes\DataProvider;
use RuntimeException;
@@ -42,6 +44,7 @@ public function testTerminalCallbacksResolveTheirOwningConnection(bool $coroutin
$resolver = $this->app->make('db.resolver');
$pool = $this->app->make(PoolManager::class)->pool('leases');
$connection = null;
+ $builder = null;
$callbacks = [];
$registeredAfterCleanup = null;
Event::listen(ConnectionReleasing::class, static function (ConnectionReleasing $event) use (&$callbacks, $listenerFails): void {
@@ -52,9 +55,10 @@ public function testTerminalCallbacksResolveTheirOwningConnection(bool $coroutin
throw new RuntimeException('Release listener failed after querying its owner.');
}
});
- $execute = static function () use (&$connection, &$callbacks): void {
+ $execute = static function () use (&$connection, &$builder, &$callbacks): void {
DB::setDefaultConnection('leases');
$connection = DB::connection();
+ $builder = $connection->query()->selectRaw('1 as value');
$connection->beginTransaction();
$connection->afterRollBack(static function () use (&$callbacks): void {
$resolved = DB::connection();
@@ -82,6 +86,16 @@ public function testTerminalCallbacksResolveTheirOwningConnection(bool $coroutin
$this->assertSame(0, $connection->transactionLevel());
$this->assertSame(0, $pool->getBorrowedCount());
$this->assertSame(1, $pool->getIdleCount());
+
+ try {
+ $builder->first();
+ $this->fail('Expected the finished connection owner to reject a new borrow.');
+ } catch (QueryException $exception) {
+ $this->assertInstanceOf(LogicException::class, $exception->getPrevious());
+ $this->assertStringContainsString('coroutine or task that resolved it has finished', $exception->getMessage());
+ }
+
+ $this->assertSame(0, $pool->getBorrowedCount());
} finally {
$resolver->discardConnections();
$pool->close();
@@ -118,15 +132,81 @@ public function testTaskCleanupSettlesTheCurrentLeaseAfterRepeatedEarlyReleases(
$resolver->releaseConnections();
$this->assertSame(1, $pool->getIdleCount());
- $this->assertNotSame($connection, DB::connection('leases'));
+ $replacement = DB::connection('leases');
+ $this->assertNotSame($connection, $replacement);
$resolver->discardConnections();
$this->assertSame(0, $pool->getManagedCount());
+
+ try {
+ $replacement->selectOne('select 1');
+ $this->fail('Expected the discarded connection owner to reject a new borrow.');
+ } catch (QueryException $exception) {
+ $this->assertInstanceOf(LogicException::class, $exception->getPrevious());
+ $this->assertStringContainsString('coroutine or task that resolved it has finished', $exception->getMessage());
+ }
+
+ $this->assertSame(0, $pool->getBorrowedCount());
+ } finally {
+ $resolver->discardConnections();
+ $pool->close();
+ }
+ }
+
+ #[DataProvider('rawTransactionOwners')]
+ public function testTerminalCleanupRollsBackRawTransactions(bool $coroutine, bool $releaseListener): void
+ {
+ $resolver = $this->app->make('db.resolver');
+ $pool = $this->app->make(PoolManager::class)->pool('leases');
+ $pdo = $pool->getSharedInMemorySqlitePdo();
+ $pdo->exec('create table messages (value integer)');
+ $pdo->exec('insert into messages values (1)');
+ $begin = static function () use ($pdo): void {
+ $pdo->beginTransaction();
+ $pdo->exec('insert into messages values (2)');
+ };
+
+ if ($releaseListener) {
+ Event::listen(ConnectionReleasing::class, $begin);
+ }
+
+ $execute = static function () use ($begin, $releaseListener): void {
+ DB::connection('leases')->selectOne('select 1');
+
+ if (! $releaseListener) {
+ $begin();
+ }
+ };
+
+ try {
+ if ($coroutine) {
+ $this->runInCoroutine($execute);
+ } else {
+ $execute();
+ $resolver->releaseConnections();
+ }
+
+ $this->assertFalse($pdo->inTransaction());
+ $this->assertSame(0, $pool->getManagedCount());
+ $this->assertSame([1], DB::connection('leases')->table('messages')->pluck('value')->all());
+ $this->assertSame($pdo, $pool->getSharedInMemorySqlitePdo());
} finally {
$resolver->discardConnections();
$pool->close();
}
}
+ /**
+ * Provide raw transactions started by application code or release callbacks.
+ */
+ public static function rawTransactionOwners(): array
+ {
+ return [
+ 'coroutine' => [true, false],
+ 'task' => [false, false],
+ 'release listener' => [false, true],
+ ];
+ }
+
public function testPurgingAndReplacingAConnectionDoesNotLoseItsPreviousOwner(): void
{
$resolver = $this->app->make('db.resolver');
diff --git a/tests/Database/DatabaseConnectionLeaseTest.php b/tests/Database/DatabaseConnectionLeaseTest.php
index e17f254d4..47ef05ed6 100644
--- a/tests/Database/DatabaseConnectionLeaseTest.php
+++ b/tests/Database/DatabaseConnectionLeaseTest.php
@@ -27,6 +27,7 @@
use Hypervel\Support\Facades\DB;
use Hypervel\Support\Facades\Event;
use Hypervel\Testbench\TestCase;
+use LogicException;
use PDO;
use PHPUnit\Framework\Attributes\DataProvider;
use RuntimeException;
@@ -237,30 +238,47 @@ public function testHttpRetriesAndRedirectsReleaseSessionsAcquiredByRequestCallb
}
#[DataProvider('pinnedScopes')]
- public function testSessionDependentScopesPreventEarlyRelease(string $scope): void
+ public function testSessionDependentScopesPreventEarlyRelease(string $scope, string $name): void
{
- $connection = DB::connection('leases');
- $check = function () use ($connection): void {
+ $connection = DB::connection($name);
+ $pool = $this->app->make(PoolManager::class)->pool($name);
+ $check = function () use ($connection, $pool): void {
$this->assertSame(1, $connection->selectOne('select 1 as value')->value);
DB::releaseIdleConnections();
- $this->assertSame(0, $this->pool()->getIdleCount());
- (new Factory)->setHandler(function (): PromiseInterface {
- $this->assertSame(0, $this->pool()->getIdleCount());
+ $this->assertSame(0, $pool->getIdleCount());
+ (new Factory)->setHandler(function () use ($pool): PromiseInterface {
+ $this->assertSame(0, $pool->getIdleCount());
return Create::promiseFor(new Response(200));
})->get('http://example.test');
};
- match ($scope) {
- 'explicit' => $connection->withPinnedSession(
- fn () => $connection->withPinnedSession($check)
- ),
- 'transaction' => $connection->transaction($check),
- 'foreign keys' => $connection->getSchemaBuilder()->withoutForeignKeyConstraints($check),
- };
+ if (in_array($scope, ['native transaction', 'SQL transaction', 'read transaction'], true)) {
+ $pdo = $scope === 'read transaction' ? $connection->getReadPdo() : $connection->getPdo();
+
+ if ($scope === 'SQL transaction') {
+ $pdo->exec('BEGIN');
+ } else {
+ $pdo->beginTransaction();
+ }
+
+ try {
+ $check();
+ } finally {
+ $pdo->rollBack();
+ }
+ } else {
+ match ($scope) {
+ 'explicit' => $connection->withPinnedSession(
+ fn () => $connection->withPinnedSession($check)
+ ),
+ 'transaction' => $connection->transaction($check),
+ 'foreign keys' => $connection->getSchemaBuilder()->withoutForeignKeyConstraints($check),
+ };
+ }
DB::releaseIdleConnections();
- $this->assertSame(1, $this->pool()->getIdleCount());
+ $this->assertSame(1, $pool->getIdleCount());
}
/**
@@ -268,7 +286,14 @@ public function testSessionDependentScopesPreventEarlyRelease(string $scope): vo
*/
public static function pinnedScopes(): array
{
- return [['explicit'], ['transaction'], ['foreign keys']];
+ return [
+ ['explicit', 'leases'],
+ ['transaction', 'leases'],
+ ['foreign keys', 'leases'],
+ ['native transaction', 'leases'],
+ ['SQL transaction', 'leases'],
+ ['read transaction', 'physical_leases'],
+ ];
}
public function testAnAbandonedCursorReleasesItsPin(): void
@@ -475,25 +500,49 @@ public function testReadWriteRolesAndStickyRoutingSurviveEarlyRelease(): void
$this->assertSame('primary', $write->getConfig('host'));
}
- public function testConfigFirstExtensionsRetainWholeConnectionOwnership(): void
+ #[DataProvider('configFirstExtensions')]
+ public function testConfigFirstExtensionsRetainWholeConnectionOwnership(bool $read, bool $listed): void
{
$constructions = 0;
- DB::extend('physical_leases', static function (array $config) use (&$constructions): PostgresConnection {
+ $created = null;
+
+ if ($read) {
+ $record = ['driver' => 'custom_read'];
+ config(['database.connections.physical_leases.read' => $listed ? [$record] : $record]);
+ }
+
+ DB::extend($read ? 'custom_read' : 'physical_leases', static function (array $config) use (&$constructions, &$created): PostgresConnection {
++$constructions;
- return new PostgresConnection(new PDO('sqlite::memory:'), $config['database'], '', $config);
+ return $created = new PostgresConnection(new PDO('sqlite::memory:'), $config['database'], '', $config);
});
- $connection = DB::connection('physical_leases');
+ $name = $read ? 'physical_leases::read' : 'physical_leases';
+ $connection = DB::connection($name);
$connection->selectOne('select 1');
DB::releaseIdleConnections();
- $this->assertSame(0, $this->app->make(PoolManager::class)->pool('physical_leases')->getIdleCount());
- $this->assertSame($connection, DB::connection('physical_leases'));
+ $pool = $this->app->make(PoolManager::class)->pool($name);
+ $this->assertFalse($pool->usesSessionLeases());
+ $this->assertSame(0, $pool->getIdleCount());
+ $this->assertSame($created, $connection);
+ $this->assertSame($connection, DB::connection($name));
$this->assertSame(2, $connection->selectOne('select 2 as value')->value);
$this->assertSame(1, $constructions);
}
+ /**
+ * Provide named and effective read-driver extensions.
+ */
+ public static function configFirstExtensions(): array
+ {
+ return [
+ 'named PDO extension' => [false, false],
+ 'single read PDO extension' => [true, false],
+ 'listed read PDO extension' => [true, true],
+ ];
+ }
+
public function testSessionCapabilitiesAreResolvedAgainAfterEarlyRelease(): void
{
Connection::resolverFor(
@@ -551,8 +600,10 @@ public function testConnectionListenerFailuresSettleTheLeaseOnce(bool $reacquire
$failure = $cancel ? new CanceledException('listener canceled') : new RuntimeException('listener failed');
$fail = true;
- Event::listen(ConnectionEstablished::class, static function (ConnectionEstablished $event) use ($failure, &$fail): void {
+ $retained = null;
+ Event::listen(ConnectionEstablished::class, static function (ConnectionEstablished $event) use ($failure, &$fail, &$retained): void {
if ($event->connection->getName() === 'leases' && $fail) {
+ $retained = $event->connection;
$fail = false;
throw $failure;
}
@@ -571,6 +622,17 @@ public function testConnectionListenerFailuresSettleTheLeaseOnce(bool $reacquire
}
$this->assertSame(0, $this->pool()->getManagedCount());
+
+ if (! $reacquire) {
+ try {
+ $retained->getPdo();
+ $this->fail('Expected the failed initial owner to reject a new borrow.');
+ } catch (LogicException) {
+ }
+
+ $this->assertSame(0, $this->pool()->getManagedCount());
+ }
+
$connection ??= DB::connection('leases');
$this->assertSame(1, $connection->selectOne('select 1 as value')->value);
DB::releaseIdleConnections();
diff --git a/tests/Database/DatabasePdoConnectionTest.php b/tests/Database/DatabasePdoConnectionTest.php
index 4c1dda55a..d03c1952b 100755
--- a/tests/Database/DatabasePdoConnectionTest.php
+++ b/tests/Database/DatabasePdoConnectionTest.php
@@ -26,6 +26,8 @@
use PHPUnit\Framework\Attributes\DataProvider;
use RuntimeException;
use stdClass;
+use Swoole\Coroutine\CanceledException;
+use Throwable;
use WeakReference;
class DatabasePdoConnectionTest extends TestCase
@@ -999,7 +1001,8 @@ public function testLostPhysicalRollbackTerminallyDetachesTransactionState(): vo
$this->assertNull($connection->getRawReadPdo());
}
- public function testDisconnectExhaustsCleanupAndPreservesThePhysicalFailure(): void
+ #[DataProvider('disconnectFailures')]
+ public function testDisconnectExhaustsCleanupAndPreservesThePhysicalFailure(bool $readCancels): void
{
$physicalFailure = new RuntimeException('physical rollback failure');
$callbackFailure = new RuntimeException('rollback callback failure');
@@ -1010,8 +1013,18 @@ public function testDisconnectExhaustsCleanupAndPreservesThePhysicalFailure(): v
$pdo->expects($this->once())->method('inTransaction')->willReturn(true);
$pdo->expects($this->once())->method('rollBack')->willThrowException($physicalFailure);
+ $readPdo = $this->getMockBuilder(PDOStub::class)->onlyMethods(['inTransaction', 'rollBack'])->getMock();
+ $readPdo->expects($this->once())->method('inTransaction')->willReturn(true);
+ $cancellation = new CanceledException('Read rollback canceled.');
+
+ if ($readCancels) {
+ $readPdo->expects($this->once())->method('rollBack')->willThrowException($cancellation);
+ } else {
+ $readPdo->expects($this->once())->method('rollBack')->willReturn(true);
+ }
+
$connection = $this->getMockConnection([], $pdo);
- $connection->setReadPdo(new PDOStub);
+ $connection->setReadPdo($readPdo);
$manager = new DatabaseTransactionsManager;
$connection->setTransactionManager($manager);
$connection->beginTransaction();
@@ -1025,8 +1038,8 @@ public function testDisconnectExhaustsCleanupAndPreservesThePhysicalFailure(): v
try {
$connection->disconnect();
$this->fail('Expected disconnect cleanup to fail.');
- } catch (RuntimeException $exception) {
- $this->assertSame($physicalFailure, $exception);
+ } catch (Throwable $exception) {
+ $this->assertSame($readCancels ? $cancellation : $physicalFailure, $exception);
}
$this->assertTrue($rollbackCallbackCalled);
@@ -1040,6 +1053,56 @@ public function testDisconnectExhaustsCleanupAndPreservesThePhysicalFailure(): v
$this->assertFalse($connection->isReusable());
}
+ /**
+ * Provide failures while rolling back physical handles.
+ */
+ public static function disconnectFailures(): array
+ {
+ return [
+ 'write rollback fails' => [false],
+ 'read cancellation wins' => [true],
+ ];
+ }
+
+ #[DataProvider('disconnectHandles')]
+ public function testDisconnectCleansOpenHandlesWithoutResolvingLazyOnes(string $read): void
+ {
+ $pdo = $this->getMockBuilder(PDOStub::class)->onlyMethods(['inTransaction', 'rollBack'])->getMock();
+ $pdo->expects($this->once())->method('inTransaction')->willReturn(true);
+ $pdo->expects($this->once())->method('rollBack')->willReturn(true);
+ $connection = new PdoConnection($pdo);
+ $readPdo = null;
+
+ if ($read === 'distinct') {
+ $readPdo = new PDO('sqlite::memory:');
+ $readPdo->beginTransaction();
+ $connection->setReadPdo($readPdo);
+ } elseif ($read === 'same') {
+ $connection->setReadPdo($pdo);
+ } else {
+ $connection->setReadPdo(function (): never {
+ $this->fail('Disconnect must not resolve a lazy read handle.');
+ });
+ }
+
+ $connection->disconnect();
+
+ if ($readPdo !== null) {
+ $this->assertFalse($readPdo->inTransaction());
+ }
+
+ $this->assertNull($connection->getRawPdo());
+ $this->assertNull($connection->getRawReadPdo());
+ }
+
+ /**
+ * Provide open, shared and unresolved read handles.
+ */
+ public static function disconnectHandles(): array
+ {
+ return [['distinct'], ['same'], ['lazy']];
+ }
+
public function testDisconnectTreatsLostPhysicalRollbackFailureAsAlreadyTerminal(): void
{
$pdo = $this->getMockBuilder(PDOStub::class)
@@ -1051,7 +1114,7 @@ public function testDisconnectTreatsLostPhysicalRollbackFailureAsAlreadyTerminal
);
$connection = $this->getMockConnection([], $pdo);
- $connection->setReadPdo(new PDOStub);
+ $connection->setReadPdo(new PDO('sqlite::memory:'));
$manager = new DatabaseTransactionsManager;
$connection->setTransactionManager($manager);
$manager->begin('test', 1);
@@ -1076,7 +1139,7 @@ public function testDisconnectPreservesManagerFailureAfterLostPhysicalRollbackFa
);
$connection = $this->getMockConnection([], $pdo);
- $connection->setReadPdo(new PDOStub);
+ $connection->setReadPdo(new PDO('sqlite::memory:'));
$manager = new DatabaseTransactionsManager;
$connection->setTransactionManager($manager);
$manager->begin('test', 1);
diff --git a/tests/HttpServer/Fixtures/disconnect-server.php b/tests/HttpServer/Fixtures/disconnect-server.php
index 429fad04b..86aaf2e8c 100644
--- a/tests/HttpServer/Fixtures/disconnect-server.php
+++ b/tests/HttpServer/Fixtures/disconnect-server.php
@@ -89,7 +89,7 @@ protected function defaultCallbacks(): array
}
})()))->cancelOnDisconnect();
- ResponseBridge::send($response, $native, streamId: $request->streamId ?? 0);
+ ResponseBridge::send($response, $native);
},
Event::ON_CLOSE => static function (NativeServer $server, int $connection) use (&$proxyAuthorizations): void {
unset($proxyAuthorizations[$connection]);
diff --git a/tests/HttpServer/ResponseBridgeTest.php b/tests/HttpServer/ResponseBridgeTest.php
index 9a1913535..d0ba5acf9 100644
--- a/tests/HttpServer/ResponseBridgeTest.php
+++ b/tests/HttpServer/ResponseBridgeTest.php
@@ -1238,19 +1238,19 @@ public function testStatusFailureThrowsBeforeHeaders(): void
}
#[DataProvider('disconnectCancellationOptions')]
- public function testDisconnectCancelsAllActiveStreamsOnOnlyThatConnection(bool $cancel, array $streams): void
+ public function testDisconnectCancelsAllActiveStreamsOnOnlyThatConnection(bool $cancel): void
{
$ready = new Channel(3);
$continue = new Channel(3);
$closed = [];
- $produce = function (int $connection, int $stream) use ($ready, $continue, $cancel, &$closed): bool {
- $response = new IterableStreamedResponse((function () use ($connection, $stream, $ready, $continue, &$closed): iterable {
+ $produce = function (int $connection, string $producer) use ($ready, $continue, $cancel, &$closed): bool {
+ $response = new IterableStreamedResponse((function () use ($producer, $ready, $continue, &$closed): iterable {
try {
$ready->push(true);
$this->assertTrue($continue->pop(2));
yield 'body';
} finally {
- $closed[] = [$connection, $stream];
+ $closed[] = $producer;
}
})());
$native = $this->mockSwooleResponse();
@@ -1262,7 +1262,7 @@ public function testDisconnectCancelsAllActiveStreamsOnOnlyThatConnection(bool $
}
try {
- ResponseBridge::send($response, $native, streamId: $stream);
+ ResponseBridge::send($response, $native);
return false;
} catch (CanceledException) {
@@ -1272,9 +1272,9 @@ public function testDisconnectCancelsAllActiveStreamsOnOnlyThatConnection(bool $
try {
$results = parallel([
- 'first' => fn (): bool => $produce(10, $streams[0]),
- 'second' => fn (): bool => $produce(10, $streams[1]),
- 'other' => fn (): bool => $produce(20, 1),
+ 'first' => fn (): bool => $produce(10, 'first'),
+ 'second' => fn (): bool => $produce(10, 'second'),
+ 'other' => fn (): bool => $produce(20, 'other'),
'close' => function () use ($ready, $continue, $cancel): void {
for ($index = 0; $index < 3; ++$index) {
$this->assertTrue($ready->pop(1));
@@ -1292,7 +1292,7 @@ public function testDisconnectCancelsAllActiveStreamsOnOnlyThatConnection(bool $
$this->assertSame($cancel, $results['first']);
$this->assertSame($cancel, $results['second']);
$this->assertFalse($results['other']);
- $this->assertEqualsCanonicalizing([[10, $streams[0]], [10, $streams[1]], [20, 1]], $closed);
+ $this->assertEqualsCanonicalizing(['first', 'second', 'other'], $closed);
} finally {
$ready->close();
$continue->close();
@@ -1305,9 +1305,8 @@ public function testDisconnectCancelsAllActiveStreamsOnOnlyThatConnection(bool $
public static function disconnectCancellationOptions(): array
{
return [
- 'default opt-out' => [false, [1, 3]],
- 'HTTP/2 streams' => [true, [1, 3]],
- 'pipelined HTTP/1 responses' => [true, [0, 0]],
+ 'default opt-out' => [false],
+ 'opt-in' => [true],
];
}
diff --git a/tests/Integration/Database/PooledConnectionTest.php b/tests/Integration/Database/PooledConnectionTest.php
index 176bad13f..2bfaf95c3 100644
--- a/tests/Integration/Database/PooledConnectionTest.php
+++ b/tests/Integration/Database/PooledConnectionTest.php
@@ -32,9 +32,11 @@
use InvalidArgumentException;
use Mockery as m;
use PDO;
+use PHPUnit\Framework\Attributes\DataProvider;
use ReflectionProperty;
use RuntimeException;
use Swoole\Coroutine\CanceledException;
+use Throwable;
use WeakReference;
/**
@@ -467,7 +469,8 @@ public function testReleaseResetsConnectionState(): void
$newPooledConnection->release();
}
- public function testReleaseRollsBackOpenTransactions(): void
+ #[DataProvider('transactionOwners')]
+ public function testReleaseRollsBackOpenTransactions(bool $raw): void
{
$pool = new DatabasePool($this->app, 'pool_test');
@@ -481,14 +484,24 @@ public function testReleaseRollsBackOpenTransactions(): void
$table->string('name');
});
- $connection->beginTransaction();
+ $pdo = $connection->getPdo();
+
+ if ($raw) {
+ $pdo->beginTransaction();
+ } else {
+ $connection->beginTransaction();
+ }
+
$connection->table('test_rollback')->insert(['name' => 'should_be_rolled_back']);
- $this->assertSame(1, $connection->transactionLevel());
+ $this->assertSame($raw ? 0 : 1, $connection->transactionLevel());
// Release should roll back
$pooledConnection->release();
+ $this->assertFalse($pdo->inTransaction());
+ $this->assertSame($raw ? 0 : 1, $pool->getManagedCount());
+
// Get a new connection and verify the data was rolled back
/** @var PooledConnection $newPooledConnection */
$newPooledConnection = $pool->borrow();
@@ -500,6 +513,14 @@ public function testReleaseRollsBackOpenTransactions(): void
$newPooledConnection->release();
}
+ /**
+ * Provide framework-managed and native transactions.
+ */
+ public static function transactionOwners(): array
+ {
+ return ['framework' => [false], 'raw PDO' => [true]];
+ }
+
public function testCleanReleasePreservesMatchingPhysicalSessionState(): void
{
$configurator = new PoolSessionConfigurator;
@@ -989,6 +1010,84 @@ public function testReleasePreservesTheFirstOrdinaryCleanupFailure(): void
$pool->close();
}
+ #[DataProvider('releaseCleanupFailures')]
+ public function testRawTransactionIsDiscardedDespiteCleanupFailure(bool $cancel): void
+ {
+ config(['database.connections.pool_test.pool.events' => [ConnectionReleasing::class]]);
+ $failure = $cancel
+ ? new CanceledException('Release listener canceled.')
+ : new RuntimeException('Transaction logging failed.');
+
+ if (! $cancel) {
+ $logger = m::mock(StdoutLoggerInterface::class);
+ $logger->shouldReceive('error')->once()->andThrow($failure);
+ $this->app->instance(StdoutLoggerInterface::class, $logger);
+ }
+
+ $pool = new DatabasePool($this->app, 'pool_test');
+ $pooledConnection = $pool->borrow();
+ $pdo = $pooledConnection->getConnection()->getPdo();
+ Event::listen(ConnectionReleasing::class, static function () use ($pdo, $cancel, $failure): void {
+ $pdo->beginTransaction();
+
+ if ($cancel) {
+ throw $failure;
+ }
+ });
+
+ try {
+ try {
+ $pooledConnection->release();
+ $this->fail('Expected the cleanup failure to propagate.');
+ } catch (Throwable $exception) {
+ $this->assertSame($failure, $exception);
+ }
+
+ $this->assertFalse($pdo->inTransaction());
+ $this->assertSame(0, $pool->getManagedCount());
+ $this->assertSame(0, $pool->getBorrowedCount());
+ } finally {
+ $pool->close();
+ }
+ }
+
+ #[DataProvider('releaseCleanupFailures')]
+ public function testPhysicalTransactionInspectionFailureDiscardsTheConnection(bool $cancel): void
+ {
+ config(['database.connections.neutral' => ['driver' => 'neutral', 'database' => 'app']]);
+ $connection = new NeutralPoolConnection(1, 'app', '', ['driver' => 'neutral']);
+ $this->app->make('db.factory')->extend('neutral', static fn () => $connection);
+ $pool = new DatabasePool($this->app, 'neutral');
+ $pooledConnection = $pool->borrow();
+ $failure = $cancel
+ ? new CanceledException('Transaction inspection canceled.')
+ : new RuntimeException('Transaction inspection failed.');
+ $connection->transactionFailure = $failure;
+
+ try {
+ try {
+ $pooledConnection->release();
+ $this->fail('Expected the inspection failure to propagate.');
+ } catch (Throwable $exception) {
+ $this->assertSame($failure, $exception);
+ }
+
+ $this->assertSame(1, $connection->disconnectCalls);
+ $this->assertSame(0, $pool->getManagedCount());
+ $this->assertSame(0, $pool->getBorrowedCount());
+ } finally {
+ $pool->close();
+ }
+ }
+
+ /**
+ * Provide ordinary failures and coroutine cancellation during cleanup.
+ */
+ public static function releaseCleanupFailures(): array
+ {
+ return ['ordinary error' => [false], 'cancellation' => [true]];
+ }
+
public function testRollbackCancellationStillReturnsTheConnectionAndEscapesExactly(): void
{
$pool = new DatabasePool($this->app, 'pool_test');
@@ -1850,6 +1949,8 @@ public function release(PoolConnection $connection): void
class NeutralPoolConnection extends Connection
{
+ public ?Throwable $transactionFailure = null;
+
public int $pingCalls = 0;
public int $disconnectCalls = 0;
@@ -1899,6 +2000,10 @@ public function ping(): bool
public function inTransaction(): bool
{
+ if ($this->transactionFailure !== null) {
+ throw $this->transactionFailure;
+ }
+
return false;
}