mirror of
https://github.com/d3vyce/fastapi-toolsets.git
synced 2026-09-19 03:09:56 +00:00
Merge pull request #381 from d3vyce/380-dependsdb-resolved-twice-in-one-request-opens-two-sessions-and-silently-discards-writes
fix: Depends(db) resolved twice in one request opens two sessions and silently discards writes
This commit is contained in:
+1
-1
@@ -54,7 +54,7 @@ db = Database(
|
|||||||
|
|
||||||
## Committing before the response
|
## Committing before the response
|
||||||
|
|
||||||
[`db.install(app)`](../reference/db.md#fastapi_toolsets.db.Database) adds a middleware that commits the request's session when the response starts, after the endpoint returns and before the body is sent. With the middleware installed, the dependency does not commit again.
|
[`db.install(app)`](../reference/db.md#fastapi_toolsets.db.Database) adds a middleware that commits the request's session when the response starts, after the endpoint returns and before the body is sent. The dependency commits only if the middleware did not: when a function-scoped dependency unwinds before the response, or when the response never passes through the middleware. Either way the request is committed exactly once.
|
||||||
|
|
||||||
The request is committed as a single transaction:
|
The request is committed as a single transaction:
|
||||||
|
|
||||||
|
|||||||
+1
-1
@@ -101,7 +101,7 @@ exclude = ["*.md"]
|
|||||||
extend-select = ["E712"]
|
extend-select = ["E712"]
|
||||||
|
|
||||||
[tool.ruff.lint.flake8-bugbear]
|
[tool.ruff.lint.flake8-bugbear]
|
||||||
extend-immutable-calls = ["fastapi.Depends"]
|
extend-immutable-calls = ["fastapi.Depends", "fastapi.Security"]
|
||||||
|
|
||||||
[tool.ruff.lint.per-file-ignores]
|
[tool.ruff.lint.per-file-ignores]
|
||||||
"tests/**" = ["RUF012", "RUF059", "SIM117", "DTZ001", "S110", "BLE001"]
|
"tests/**" = ["RUF012", "RUF059", "SIM117", "DTZ001", "S110", "BLE001"]
|
||||||
|
|||||||
@@ -66,10 +66,8 @@ class _CommitOnResponseMiddleware:
|
|||||||
|
|
||||||
async def send_wrapper(message: Message) -> None:
|
async def send_wrapper(message: Message) -> None:
|
||||||
if message["type"] == "http.response.start":
|
if message["type"] == "http.response.start":
|
||||||
# ``scope["state"]`` is the same dict ``request.state`` writes
|
|
||||||
# to, so this is the session stashed by the dependency.
|
|
||||||
state = scope.get("state")
|
state = scope.get("state")
|
||||||
session = state.get(self.state_attr) if state else None
|
session = state.pop(self.state_attr, None) if state else None
|
||||||
if session is not None and session.in_transaction():
|
if session is not None and session.in_transaction():
|
||||||
await session.commit()
|
await session.commit()
|
||||||
await send(message)
|
await send(message)
|
||||||
@@ -158,7 +156,6 @@ class Database:
|
|||||||
# Private, per-instance state attribute; cannot collide with another
|
# Private, per-instance state attribute; cannot collide with another
|
||||||
# Database or be mismatched against the middleware.
|
# Database or be mismatched against the middleware.
|
||||||
self._state_attr = f"_ft_db_session_{id(self):x}"
|
self._state_attr = f"_ft_db_session_{id(self):x}"
|
||||||
self._middleware_installed = False
|
|
||||||
self._disposed = False
|
self._disposed = False
|
||||||
|
|
||||||
async def _dispose(self) -> None:
|
async def _dispose(self) -> None:
|
||||||
@@ -206,7 +203,6 @@ class Database:
|
|||||||
```
|
```
|
||||||
"""
|
"""
|
||||||
app.add_middleware(_CommitOnResponseMiddleware, state_attr=self._state_attr)
|
app.add_middleware(_CommitOnResponseMiddleware, state_attr=self._state_attr)
|
||||||
self._middleware_installed = True
|
|
||||||
|
|
||||||
inner_lifespan = app.router.lifespan_context
|
inner_lifespan = app.router.lifespan_context
|
||||||
|
|
||||||
@@ -243,10 +239,17 @@ class Database:
|
|||||||
return await UserCrud.get(session, [User.id == user_id])
|
return await UserCrud.get(session, [User.id == user_id])
|
||||||
```
|
```
|
||||||
"""
|
"""
|
||||||
|
borrowed = getattr(request.state, self._state_attr, None)
|
||||||
|
if borrowed is not None:
|
||||||
|
yield borrowed
|
||||||
|
return
|
||||||
async with self._open() as session:
|
async with self._open() as session:
|
||||||
setattr(request.state, self._state_attr, session)
|
setattr(request.state, self._state_attr, session)
|
||||||
yield session
|
yield session
|
||||||
if not self._middleware_installed and session.in_transaction():
|
if (
|
||||||
|
getattr(request.state, self._state_attr, None) is session
|
||||||
|
and session.in_transaction()
|
||||||
|
):
|
||||||
await session.commit()
|
await session.commit()
|
||||||
|
|
||||||
@asynccontextmanager
|
@asynccontextmanager
|
||||||
|
|||||||
+128
-9
@@ -6,7 +6,7 @@ from contextlib import asynccontextmanager
|
|||||||
from unittest.mock import AsyncMock, MagicMock, patch
|
from unittest.mock import AsyncMock, MagicMock, patch
|
||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
from fastapi import Depends, FastAPI
|
from fastapi import Depends, FastAPI, Security
|
||||||
from fastapi.responses import StreamingResponse
|
from fastapi.responses import StreamingResponse
|
||||||
from httpx import ASGITransport, AsyncClient
|
from httpx import ASGITransport, AsyncClient
|
||||||
from pydantic import PostgresDsn
|
from pydantic import PostgresDsn
|
||||||
@@ -287,23 +287,37 @@ class TestDatabaseDependency:
|
|||||||
break
|
break
|
||||||
|
|
||||||
@pytest.mark.anyio
|
@pytest.mark.anyio
|
||||||
async def test_skips_commit_when_middleware_installed(self, engine, session_maker):
|
async def test_second_resolution_borrows_session(self, engine):
|
||||||
"""With ``install()``, the dependency must NOT commit — the middleware owns it.
|
"""A second ``Depends(db)`` in one request reuses the stashed session."""
|
||||||
|
db = Database(engine=engine)
|
||||||
|
request = _make_request()
|
||||||
|
|
||||||
Here no middleware actually runs (we call the dependency directly), so the
|
owner_gen = db(request)
|
||||||
open transaction is rolled back on session close and nothing persists.
|
owner = await anext(owner_gen)
|
||||||
"""
|
borrower_gen = db(request)
|
||||||
|
assert await anext(borrower_gen) is owner
|
||||||
|
|
||||||
|
with pytest.raises(StopAsyncIteration): # teardown runs borrower-first
|
||||||
|
await anext(borrower_gen)
|
||||||
|
assert owner.in_transaction() # the borrower must not close what it borrowed
|
||||||
|
|
||||||
|
with pytest.raises(StopAsyncIteration):
|
||||||
|
await anext(owner_gen)
|
||||||
|
|
||||||
|
@pytest.mark.anyio
|
||||||
|
async def test_commits_when_middleware_did_not_run(self, engine, session_maker):
|
||||||
|
"""``install()`` is per-``Database``, but the commit is per-request."""
|
||||||
db = Database(engine=engine)
|
db = Database(engine=engine)
|
||||||
db.install(FastAPI())
|
db.install(FastAPI())
|
||||||
|
|
||||||
async for session in db(_make_request()):
|
async for session in db(_make_request()):
|
||||||
role = Role(name="mw_owns_commit")
|
role = Role(name="mw_never_ran")
|
||||||
session.add(role)
|
session.add(role)
|
||||||
await session.flush()
|
await session.flush()
|
||||||
|
|
||||||
async with session_maker() as verify:
|
async with session_maker() as verify:
|
||||||
result = await RoleCrud.first(verify, [Role.name == "mw_owns_commit"])
|
result = await RoleCrud.first(verify, [Role.name == "mw_never_ran"])
|
||||||
assert result is None
|
assert result is not None
|
||||||
|
|
||||||
|
|
||||||
class TestDatabaseSession:
|
class TestDatabaseSession:
|
||||||
@@ -1523,6 +1537,55 @@ def _build_app(db: Database) -> FastAPI:
|
|||||||
await session.commit()
|
await session.commit()
|
||||||
return {"id": str(role.id), "name": role.name}
|
return {"id": str(role.id), "name": role.name}
|
||||||
|
|
||||||
|
async def _scoped_writer(
|
||||||
|
body: RoleCreate, session: AsyncSession = Security(db, scopes=["roles:write"])
|
||||||
|
) -> int:
|
||||||
|
# Security scopes give this a different dependency cache key than the
|
||||||
|
# endpoint's plain ``Depends(db)``. Without borrowing it opens a second
|
||||||
|
# session, and whichever one the middleware does not hold is discarded.
|
||||||
|
await RoleCrud.create(session, RoleCreate(name=f"{body.name}_sub"))
|
||||||
|
return id(session)
|
||||||
|
|
||||||
|
@app.post("/roles-two-cache-keys")
|
||||||
|
async def create_via_two_cache_keys(
|
||||||
|
body: RoleCreate,
|
||||||
|
sub_session_id: int = Depends(_scoped_writer),
|
||||||
|
session: AsyncSession = Depends(db),
|
||||||
|
) -> dict:
|
||||||
|
await RoleCrud.create(session, body)
|
||||||
|
return {"same_session": sub_session_id == id(session)}
|
||||||
|
|
||||||
|
async def _fn_writer(
|
||||||
|
body: RoleCreate, session: AsyncSession = Depends(db, scope="function")
|
||||||
|
) -> None:
|
||||||
|
# ``scope="function"`` unwinds before the response is sent, taking the
|
||||||
|
# session with it — so the commit cannot be left to the middleware.
|
||||||
|
await RoleCrud.create(session, RoleCreate(name=f"{body.name}_fn"))
|
||||||
|
|
||||||
|
@app.post("/roles-function-scope")
|
||||||
|
async def create_with_function_scope(
|
||||||
|
body: RoleCreate,
|
||||||
|
boom: bool = False,
|
||||||
|
_: None = Depends(_fn_writer),
|
||||||
|
session: AsyncSession = Depends(db),
|
||||||
|
) -> dict:
|
||||||
|
await RoleCrud.create(session, body)
|
||||||
|
if boom:
|
||||||
|
raise RuntimeError("boom after write")
|
||||||
|
return {"ok": True}
|
||||||
|
|
||||||
|
@app.post("/roles-function-scope-borrower")
|
||||||
|
async def function_scope_borrows(
|
||||||
|
body: RoleCreate,
|
||||||
|
session: AsyncSession = Depends(db),
|
||||||
|
_: None = Depends(_fn_writer),
|
||||||
|
) -> dict:
|
||||||
|
# Flipped order: the request-scoped dependency owns the session and the
|
||||||
|
# function-scoped one borrows it. The borrower unwinds early but must not
|
||||||
|
# commit or close — the commit still belongs to the middleware.
|
||||||
|
await RoleCrud.create(session, body)
|
||||||
|
return {"ok": True}
|
||||||
|
|
||||||
@app.get("/roles-stream/{name}")
|
@app.get("/roles-stream/{name}")
|
||||||
async def stream_role(
|
async def stream_role(
|
||||||
name: str, session: AsyncSession = Depends(db)
|
name: str, session: AsyncSession = Depends(db)
|
||||||
@@ -1619,6 +1682,62 @@ class TestCommitIntegration:
|
|||||||
# The write made before the stream began is durably committed.
|
# The write made before the stream began is durably committed.
|
||||||
assert await _row_exists(session_maker, "streamed_role")
|
assert await _row_exists(session_maker, "streamed_role")
|
||||||
|
|
||||||
|
@pytest.mark.anyio
|
||||||
|
async def test_two_cache_keys_share_one_session(self, engine, session_maker):
|
||||||
|
"""Two resolutions of ``Depends(db)`` in one request must share a session."""
|
||||||
|
app = _build_app(Database(engine=engine))
|
||||||
|
transport = ASGITransport(app=app)
|
||||||
|
async with AsyncClient(transport=transport, base_url="http://test") as client:
|
||||||
|
resp = await client.post("/roles-two-cache-keys", json={"name": "two_keys"})
|
||||||
|
|
||||||
|
assert resp.status_code == 200
|
||||||
|
assert resp.json()["same_session"] is True
|
||||||
|
assert await _row_exists(session_maker, "two_keys")
|
||||||
|
assert await _row_exists(session_maker, "two_keys_sub")
|
||||||
|
|
||||||
|
@pytest.mark.anyio
|
||||||
|
async def test_function_scope_commits_before_response(self, engine, session_maker):
|
||||||
|
"""``scope="function"`` unwinds before response-start, so the dependency
|
||||||
|
commits on its way out instead of leaving it to the middleware."""
|
||||||
|
app = _build_app(Database(engine=engine))
|
||||||
|
transport = ASGITransport(app=app)
|
||||||
|
async with AsyncClient(transport=transport, base_url="http://test") as client:
|
||||||
|
resp = await client.post("/roles-function-scope", json={"name": "fn_scope"})
|
||||||
|
|
||||||
|
assert resp.status_code == 200
|
||||||
|
assert await _row_exists(session_maker, "fn_scope")
|
||||||
|
assert await _row_exists(session_maker, "fn_scope_fn")
|
||||||
|
|
||||||
|
@pytest.mark.anyio
|
||||||
|
async def test_function_scope_borrower_leaves_commit_to_middleware(
|
||||||
|
self, engine, session_maker
|
||||||
|
):
|
||||||
|
"""A function-scoped *borrower* unwinds early but owns nothing."""
|
||||||
|
app = _build_app(Database(engine=engine))
|
||||||
|
transport = ASGITransport(app=app)
|
||||||
|
async with AsyncClient(transport=transport, base_url="http://test") as client:
|
||||||
|
resp = await client.post(
|
||||||
|
"/roles-function-scope-borrower", json={"name": "fn_borrow"}
|
||||||
|
)
|
||||||
|
|
||||||
|
assert resp.status_code == 200
|
||||||
|
assert await _row_exists(session_maker, "fn_borrow")
|
||||||
|
assert await _row_exists(session_maker, "fn_borrow_fn")
|
||||||
|
|
||||||
|
@pytest.mark.anyio
|
||||||
|
async def test_function_scope_error_rolls_back(self, engine, session_maker):
|
||||||
|
"""The early commit must still not fire when the request fails."""
|
||||||
|
app = _build_app(Database(engine=engine))
|
||||||
|
transport = ASGITransport(app=app, raise_app_exceptions=False)
|
||||||
|
async with AsyncClient(transport=transport, base_url="http://test") as client:
|
||||||
|
resp = await client.post(
|
||||||
|
"/roles-function-scope?boom=true", json={"name": "fn_ghost"}
|
||||||
|
)
|
||||||
|
|
||||||
|
assert resp.status_code == 500
|
||||||
|
assert not await _row_exists(session_maker, "fn_ghost")
|
||||||
|
assert not await _row_exists(session_maker, "fn_ghost_fn")
|
||||||
|
|
||||||
@pytest.mark.anyio
|
@pytest.mark.anyio
|
||||||
async def test_multi_write_atomicity(self, engine, session_maker):
|
async def test_multi_write_atomicity(self, engine, session_maker):
|
||||||
"""When the 2nd write fails, the 1st must roll back too (one txn)."""
|
"""When the 2nd write fails, the 1st must roll back too (one txn)."""
|
||||||
|
|||||||
Reference in New Issue
Block a user