codekingpro/portable-devtools
114k
1# ext/asyncio/session.py
2# Copyright (C) 2020-2026 the SQLAlchemy authors and contributors
3# <see AUTHORS file>
4#
5# This module is part of SQLAlchemy and is released under
6# the MIT License: https://www.opensource.org/licenses/mit-license.php
7from __future__ import annotations
8
9import asyncio
10from typing import Any
11from typing import Awaitable
12from typing import Callable
13from typing import cast
14from typing import Dict
15from typing import Generic
16from typing import Iterable
17from typing import Iterator
18from typing import NoReturn
19from typing import Optional
20from typing import overload
21from typing import Sequence
22from typing import Tuple
23from typing import Type
24from typing import TYPE_CHECKING
25from typing import TypeVar
26from typing import Union
27
28from . import engine
29from .base import ReversibleProxy
30from .base import StartableContext
31from .result import _ensure_sync_result
32from .result import AsyncResult
33from .result import AsyncScalarResult
34from ... import util
35from ...orm import close_all_sessions as _sync_close_all_sessions
36from ...orm import object_session
37from ...orm import Session
38from ...orm import SessionTransaction
39from ...orm import state as _instance_state
40from ...util.concurrency import greenlet_spawn
41from ...util.typing import Concatenate
42from ...util.typing import ParamSpec
43
44
45if TYPE_CHECKING:
46 from .engine import AsyncConnection
47 from .engine import AsyncEngine
48 from ...engine import Connection
49 from ...engine import Engine
50 from ...engine import Result
51 from ...engine import Row
52 from ...engine import RowMapping
53 from ...engine import ScalarResult
54 from ...engine.interfaces import _CoreAnyExecuteParams
55 from ...engine.interfaces import CoreExecuteOptionsParameter
56 from ...event import dispatcher
57 from ...orm._typing import _IdentityKeyType
58 from ...orm._typing import _O
59 from ...orm._typing import OrmExecuteOptionsParameter
60 from ...orm.identity import IdentityMap
61 from ...orm.interfaces import ORMOption
62 from ...orm.session import _BindArguments
63 from ...orm.session import _EntityBindKey
64 from ...orm.session import _PKIdentityArgument
65 from ...orm.session import _SessionBind
66 from ...orm.session import _SessionBindKey
67 from ...sql._typing import _InfoType
68 from ...sql.base import Executable
69 from ...sql.elements import ClauseElement
70 from ...sql.selectable import ForUpdateParameter
71 from ...sql.selectable import TypedReturnsRows
72
73_AsyncSessionBind = Union["AsyncEngine", "AsyncConnection"]
74
75_P = ParamSpec("_P")
76_T = TypeVar("_T", bound=Any)
77
78
79_EXECUTE_OPTIONS = util.immutabledict({"prebuffer_rows": True})
80_STREAM_OPTIONS = util.immutabledict({"stream_results": True})
81
82
83class AsyncAttrs:
84 """Mixin class which provides an awaitable accessor for all attributes.
85
86 E.g.::
87
88 from __future__ import annotations
89
90 from typing import List
91
92 from sqlalchemy import ForeignKey
93 from sqlalchemy import func
94 from sqlalchemy.ext.asyncio import AsyncAttrs
95 from sqlalchemy.orm import DeclarativeBase
96 from sqlalchemy.orm import Mapped
97 from sqlalchemy.orm import mapped_column
98 from sqlalchemy.orm import relationship
99
100
101 class Base(AsyncAttrs, DeclarativeBase):
102 pass
103
104
105 class A(Base):
106 __tablename__ = "a"
107
108 id: Mapped[int] = mapped_column(primary_key=True)
109 data: Mapped[str]
110 bs: Mapped[List[B]] = relationship()
111
112
113 class B(Base):
114 __tablename__ = "b"
115 id: Mapped[int] = mapped_column(primary_key=True)
116 a_id: Mapped[int] = mapped_column(ForeignKey("a.id"))
117 data: Mapped[str]
118
119 In the above example, the :class:`_asyncio.AsyncAttrs` mixin is applied to
120 the declarative ``Base`` class where it takes effect for all subclasses.
121 This mixin adds a single new attribute
122 :attr:`_asyncio.AsyncAttrs.awaitable_attrs` to all classes, which will
123 yield the value of any attribute as an awaitable. This allows attributes
124 which may be subject to lazy loading or deferred / unexpiry loading to be
125 accessed such that IO can still be emitted::
126
127 a1 = (await async_session.scalars(select(A).where(A.id == 5))).one()
128
129 # use the lazy loader on ``a1.bs`` via the ``.awaitable_attrs``
130 # interface, so that it may be awaited
131 for b1 in await a1.awaitable_attrs.bs:
132 print(b1)
133
134 The :attr:`_asyncio.AsyncAttrs.awaitable_attrs` performs a call against the
135 attribute that is approximately equivalent to using the
136 :meth:`_asyncio.AsyncSession.run_sync` method, e.g.::
137
138 for b1 in await async_session.run_sync(lambda sess: a1.bs):
139 print(b1)
140
141 .. versionadded:: 2.0.13
142
143 .. seealso::
144
145 :ref:`asyncio_orm_avoid_lazyloads`
146
147 """
148
149 class _AsyncAttrGetitem:
150 __slots__ = "_instance"
151
152 def __init__(self, _instance: Any):
153 self._instance = _instance
154
155 def __getattr__(self, name: str) -> Awaitable[Any]:
156 return greenlet_spawn(getattr, self._instance, name)
157
158 @property
159 def awaitable_attrs(self) -> AsyncAttrs._AsyncAttrGetitem:
160 """provide a namespace of all attributes on this object wrapped
161 as awaitables.
162
163 e.g.::
164
165
166 a1 = (await async_session.scalars(select(A).where(A.id == 5))).one()
167
168 some_attribute = await a1.awaitable_attrs.some_deferred_attribute
169 some_collection = await a1.awaitable_attrs.some_collection
170
171 """ # noqa: E501
172
173 return AsyncAttrs._AsyncAttrGetitem(self)
174
175
176@util.create_proxy_methods(
177 Session,
178 ":class:`_orm.Session`",
179 ":class:`_asyncio.AsyncSession`",
180 classmethods=["object_session", "identity_key"],
181 methods=[
182 "__contains__",
183 "__iter__",
184 "add",
185 "add_all",
186 "expire",
187 "expire_all",
188 "expunge",
189 "expunge_all",
190 "is_modified",
191 "in_transaction",
192 "in_nested_transaction",
193 ],
194 attributes=[
195 "dirty",
196 "deleted",
197 "new",
198 "identity_map",
199 "is_active",
200 "autoflush",
201 "no_autoflush",
202 "info",
203 ],
204)
205class AsyncSession(ReversibleProxy[Session]):
206 """Asyncio version of :class:`_orm.Session`.
207
208 The :class:`_asyncio.AsyncSession` is a proxy for a traditional
209 :class:`_orm.Session` instance.
210
211 The :class:`_asyncio.AsyncSession` is **not safe for use in concurrent
212 tasks.**. See :ref:`session_faq_threadsafe` for background.
213
214 .. versionadded:: 1.4
215
216 To use an :class:`_asyncio.AsyncSession` with custom :class:`_orm.Session`
217 implementations, see the
218 :paramref:`_asyncio.AsyncSession.sync_session_class` parameter.
219
220
221 """
222
223 _is_asyncio = True
224
225 dispatch: dispatcher[Session]
226
227 def __init__(
228 self,
229 bind: Optional[_AsyncSessionBind] = None,
230 *,
231 binds: Optional[Dict[_SessionBindKey, _AsyncSessionBind]] = None,
232 sync_session_class: Optional[Type[Session]] = None,
233 **kw: Any,
234 ):
235 r"""Construct a new :class:`_asyncio.AsyncSession`.
236
237 All parameters other than ``sync_session_class`` are passed to the
238 ``sync_session_class`` callable directly to instantiate a new
239 :class:`_orm.Session`. Refer to :meth:`_orm.Session.__init__` for
240 parameter documentation.
241
242 :param sync_session_class:
243 A :class:`_orm.Session` subclass or other callable which will be used
244 to construct the :class:`_orm.Session` which will be proxied. This
245 parameter may be used to provide custom :class:`_orm.Session`
246 subclasses. Defaults to the
247 :attr:`_asyncio.AsyncSession.sync_session_class` class-level
248 attribute.
249
250 .. versionadded:: 1.4.24
251
252 """
253 sync_bind = sync_binds = None
254
255 if bind:
256 self.bind = bind
257 sync_bind = engine._get_sync_engine_or_connection(bind)
258
259 if binds:
260 self.binds = binds
261 sync_binds = {
262 key: engine._get_sync_engine_or_connection(b)
263 for key, b in binds.items()
264 }
265
266 if sync_session_class:
267 self.sync_session_class = sync_session_class
268
269 self.sync_session = self._proxied = self._assign_proxied(
270 self.sync_session_class(bind=sync_bind, binds=sync_binds, **kw)
271 )
272
273 sync_session_class: Type[Session] = Session
274 """The class or callable that provides the
275 underlying :class:`_orm.Session` instance for a particular
276 :class:`_asyncio.AsyncSession`.
277
278 At the class level, this attribute is the default value for the
279 :paramref:`_asyncio.AsyncSession.sync_session_class` parameter. Custom
280 subclasses of :class:`_asyncio.AsyncSession` can override this.
281
282 At the instance level, this attribute indicates the current class or
283 callable that was used to provide the :class:`_orm.Session` instance for
284 this :class:`_asyncio.AsyncSession` instance.
285
286 .. versionadded:: 1.4.24
287
288 """
289
290 sync_session: Session
291 """Reference to the underlying :class:`_orm.Session` this
292 :class:`_asyncio.AsyncSession` proxies requests towards.
293
294 This instance can be used as an event target.
295
296 .. seealso::
297
298 :ref:`asyncio_events`
299
300 """
301
302 @classmethod
303 def _no_async_engine_events(cls) -> NoReturn:
304 raise NotImplementedError(
305 "asynchronous events are not implemented at this time. Apply "
306 "synchronous listeners to the AsyncSession.sync_session."
307 )
308
309 async def refresh(
310 self,
311 instance: object,
312 attribute_names: Optional[Iterable[str]] = None,
313 with_for_update: ForUpdateParameter = None,
314 ) -> None:
315 """Expire and refresh the attributes on the given instance.
316
317 A query will be issued to the database and all attributes will be
318 refreshed with their current database value.
319
320 This is the async version of the :meth:`_orm.Session.refresh` method.
321 See that method for a complete description of all options.
322
323 .. seealso::
324
325 :meth:`_orm.Session.refresh` - main documentation for refresh
326
327 """
328
329 await greenlet_spawn(
330 self.sync_session.refresh,
331 instance,
332 attribute_names=attribute_names,
333 with_for_update=with_for_update,
334 )
335
336 async def run_sync(
337 self,
338 fn: Callable[Concatenate[Session, _P], _T],
339 *arg: _P.args,
340 **kw: _P.kwargs,
341 ) -> _T:
342 '''Invoke the given synchronous (i.e. not async) callable,
343 passing a synchronous-style :class:`_orm.Session` as the first
344 argument.
345
346 This method allows traditional synchronous SQLAlchemy functions to
347 run within the context of an asyncio application.
348
349 E.g.::
350
351 def some_business_method(session: Session, param: str) -> str:
352 """A synchronous function that does not require awaiting
353
354 :param session: a SQLAlchemy Session, used synchronously
355
356 :return: an optional return value is supported
357
358 """
359 session.add(MyObject(param=param))
360 session.flush()
361 return "success"
362
363
364 async def do_something_async(async_engine: AsyncEngine) -> None:
365 """an async function that uses awaiting"""
366
367 with AsyncSession(async_engine) as async_session:
368 # run some_business_method() with a sync-style
369 # Session, proxied into an awaitable
370 return_code = await async_session.run_sync(
371 some_business_method, param="param1"
372 )
373 print(return_code)
374
375 This method maintains the asyncio event loop all the way through
376 to the database connection by running the given callable in a
377 specially instrumented greenlet.
378
379 .. tip::
380
381 The provided callable is invoked inline within the asyncio event
382 loop, and will block on traditional IO calls. IO within this
383 callable should only call into SQLAlchemy's asyncio database
384 APIs which will be properly adapted to the greenlet context.
385
386 .. seealso::
387
388 :class:`.AsyncAttrs` - a mixin for ORM mapped classes that provides
389 a similar feature more succinctly on a per-attribute basis
390
391 :meth:`.AsyncConnection.run_sync`
392
393 :ref:`session_run_sync`
394 ''' # noqa: E501
395
396 return await greenlet_spawn(
397 fn, self.sync_session, *arg, _require_await=False, **kw
398 )
399
400 @overload
401 async def execute(
402 self,
403 statement: TypedReturnsRows[_T],
404 params: Optional[_CoreAnyExecuteParams] = None,
405 *,
406 execution_options: OrmExecuteOptionsParameter = util.EMPTY_DICT,
407 bind_arguments: Optional[_BindArguments] = None,
408 _parent_execute_state: Optional[Any] = None,
409 _add_event: Optional[Any] = None,
410 ) -> Result[_T]: ...
411
412 @overload
413 async def execute(
414 self,
415 statement: Executable,
416 params: Optional[_CoreAnyExecuteParams] = None,
417 *,
418 execution_options: OrmExecuteOptionsParameter = util.EMPTY_DICT,
419 bind_arguments: Optional[_BindArguments] = None,
420 _parent_execute_state: Optional[Any] = None,
421 _add_event: Optional[Any] = None,
422 ) -> Result[Any]: ...
423
424 async def execute(
425 self,
426 statement: Executable,
427 params: Optional[_CoreAnyExecuteParams] = None,
428 *,
429 execution_options: OrmExecuteOptionsParameter = util.EMPTY_DICT,
430 bind_arguments: Optional[_BindArguments] = None,
431 **kw: Any,
432 ) -> Result[Any]:
433 """Execute a statement and return a buffered
434 :class:`_engine.Result` object.
435
436 .. seealso::
437
438 :meth:`_orm.Session.execute` - main documentation for execute
439
440 """
441
442 if execution_options:
443 execution_options = util.immutabledict(execution_options).union(
444 _EXECUTE_OPTIONS
445 )
446 else:
447 execution_options = _EXECUTE_OPTIONS
448
449 result = await greenlet_spawn(
450 self.sync_session.execute,
451 statement,
452 params=params,
453 execution_options=execution_options,
454 bind_arguments=bind_arguments,
455 **kw,
456 )
457 return await _ensure_sync_result(result, self.execute)
458
459 @overload
460 async def scalar(
461 self,
462 statement: TypedReturnsRows[Tuple[_T]],
463 params: Optional[_CoreAnyExecuteParams] = None,
464 *,
465 execution_options: OrmExecuteOptionsParameter = util.EMPTY_DICT,
466 bind_arguments: Optional[_BindArguments] = None,
467 **kw: Any,
468 ) -> Optional[_T]: ...
469
470 @overload
471 async def scalar(
472 self,
473 statement: Executable,
474 params: Optional[_CoreAnyExecuteParams] = None,
475 *,
476 execution_options: OrmExecuteOptionsParameter = util.EMPTY_DICT,
477 bind_arguments: Optional[_BindArguments] = None,
478 **kw: Any,
479 ) -> Any: ...
480
481 async def scalar(
482 self,
483 statement: Executable,
484 params: Optional[_CoreAnyExecuteParams] = None,
485 *,
486 execution_options: OrmExecuteOptionsParameter = util.EMPTY_DICT,
487 bind_arguments: Optional[_BindArguments] = None,
488 **kw: Any,
489 ) -> Any:
490 """Execute a statement and return a scalar result.
491
492 .. seealso::
493
494 :meth:`_orm.Session.scalar` - main documentation for scalar
495
496 """
497
498 if execution_options:
499 execution_options = util.immutabledict(execution_options).union(
500 _EXECUTE_OPTIONS
501 )
502 else:
503 execution_options = _EXECUTE_OPTIONS
504
505 return await greenlet_spawn(
506 self.sync_session.scalar,
507 statement,
508 params=params,
509 execution_options=execution_options,
510 bind_arguments=bind_arguments,
511 **kw,
512 )
513
514 @overload
515 async def scalars(
516 self,
517 statement: TypedReturnsRows[Tuple[_T]],
518 params: Optional[_CoreAnyExecuteParams] = None,
519 *,
520 execution_options: OrmExecuteOptionsParameter = util.EMPTY_DICT,
521 bind_arguments: Optional[_BindArguments] = None,
522 **kw: Any,
523 ) -> ScalarResult[_T]: ...
524
525 @overload
526 async def scalars(
527 self,
528 statement: Executable,
529 params: Optional[_CoreAnyExecuteParams] = None,
530 *,
531 execution_options: OrmExecuteOptionsParameter = util.EMPTY_DICT,
532 bind_arguments: Optional[_BindArguments] = None,
533 **kw: Any,
534 ) -> ScalarResult[Any]: ...
535
536 async def scalars(
537 self,
538 statement: Executable,
539 params: Optional[_CoreAnyExecuteParams] = None,
540 *,
541 execution_options: OrmExecuteOptionsParameter = util.EMPTY_DICT,
542 bind_arguments: Optional[_BindArguments] = None,
543 **kw: Any,
544 ) -> ScalarResult[Any]:
545 """Execute a statement and return scalar results.
546
547 :return: a :class:`_result.ScalarResult` object
548
549 .. versionadded:: 1.4.24 Added :meth:`_asyncio.AsyncSession.scalars`
550
551 .. versionadded:: 1.4.26 Added
552 :meth:`_asyncio.async_scoped_session.scalars`
553
554 .. seealso::
555
556 :meth:`_orm.Session.scalars` - main documentation for scalars
557
558 :meth:`_asyncio.AsyncSession.stream_scalars` - streaming version
559
560 """
561
562 result = await self.execute(
563 statement,
564 params=params,
565 execution_options=execution_options,
566 bind_arguments=bind_arguments,
567 **kw,
568 )
569 return result.scalars()
570
571 async def get(
572 self,
573 entity: _EntityBindKey[_O],
574 ident: _PKIdentityArgument,
575 *,
576 options: Optional[Sequence[ORMOption]] = None,
577 populate_existing: bool = False,
578 with_for_update: ForUpdateParameter = None,
579 identity_token: Optional[Any] = None,
580 execution_options: OrmExecuteOptionsParameter = util.EMPTY_DICT,
581 ) -> Union[_O, None]:
582 """Return an instance based on the given primary key identifier,
583 or ``None`` if not found.
584
585 .. seealso::
586
587 :meth:`_orm.Session.get` - main documentation for get
588
589
590 """
591
592 return await greenlet_spawn(
593 cast("Callable[..., _O]", self.sync_session.get),
594 entity,
595 ident,
596 options=options,
597 populate_existing=populate_existing,
598 with_for_update=with_for_update,
599 identity_token=identity_token,
600 execution_options=execution_options,
601 )
602
603 async def get_one(
604 self,
605 entity: _EntityBindKey[_O],
606 ident: _PKIdentityArgument,
607 *,
608 options: Optional[Sequence[ORMOption]] = None,
609 populate_existing: bool = False,
610 with_for_update: ForUpdateParameter = None,
611 identity_token: Optional[Any] = None,
612 execution_options: OrmExecuteOptionsParameter = util.EMPTY_DICT,
613 ) -> _O:
614 """Return an instance based on the given primary key identifier,
615 or raise an exception if not found.
616
617 Raises :class:`_exc.NoResultFound` if the query selects no rows.
618
619 ..versionadded: 2.0.22
620
621 .. seealso::
622
623 :meth:`_orm.Session.get_one` - main documentation for get_one
624
625 """
626
627 return await greenlet_spawn(
628 cast("Callable[..., _O]", self.sync_session.get_one),
629 entity,
630 ident,
631 options=options,
632 populate_existing=populate_existing,
633 with_for_update=with_for_update,
634 identity_token=identity_token,
635 execution_options=execution_options,
636 )
637
638 @overload
639 async def stream(
640 self,
641 statement: TypedReturnsRows[_T],
642 params: Optional[_CoreAnyExecuteParams] = None,
643 *,
644 execution_options: OrmExecuteOptionsParameter = util.EMPTY_DICT,
645 bind_arguments: Optional[_BindArguments] = None,
646 **kw: Any,
647 ) -> AsyncResult[_T]: ...
648
649 @overload
650 async def stream(
651 self,
652 statement: Executable,
653 params: Optional[_CoreAnyExecuteParams] = None,
654 *,
655 execution_options: OrmExecuteOptionsParameter = util.EMPTY_DICT,
656 bind_arguments: Optional[_BindArguments] = None,
657 **kw: Any,
658 ) -> AsyncResult[Any]: ...
659
660 async def stream(
661 self,
662 statement: Executable,
663 params: Optional[_CoreAnyExecuteParams] = None,
664 *,
665 execution_options: OrmExecuteOptionsParameter = util.EMPTY_DICT,
666 bind_arguments: Optional[_BindArguments] = None,
667 **kw: Any,
668 ) -> AsyncResult[Any]:
669 """Execute a statement and return a streaming
670 :class:`_asyncio.AsyncResult` object.
671
672 """
673
674 if execution_options:
675 execution_options = util.immutabledict(execution_options).union(
676 _STREAM_OPTIONS
677 )
678 else:
679 execution_options = _STREAM_OPTIONS
680
681 result = await greenlet_spawn(
682 self.sync_session.execute,
683 statement,
684 params=params,
685 execution_options=execution_options,
686 bind_arguments=bind_arguments,
687 **kw,
688 )
689 return AsyncResult(result)
690
691 @overload
692 async def stream_scalars(
693 self,
694 statement: TypedReturnsRows[Tuple[_T]],
695 params: Optional[_CoreAnyExecuteParams] = None,
696 *,
697 execution_options: OrmExecuteOptionsParameter = util.EMPTY_DICT,
698 bind_arguments: Optional[_BindArguments] = None,
699 **kw: Any,
700 ) -> AsyncScalarResult[_T]: ...
701
702 @overload
703 async def stream_scalars(
704 self,
705 statement: Executable,
706 params: Optional[_CoreAnyExecuteParams] = None,
707 *,
708 execution_options: OrmExecuteOptionsParameter = util.EMPTY_DICT,
709 bind_arguments: Optional[_BindArguments] = None,
710 **kw: Any,
711 ) -> AsyncScalarResult[Any]: ...
712
713 async def stream_scalars(
714 self,
715 statement: Executable,
716 params: Optional[_CoreAnyExecuteParams] = None,
717 *,
718 execution_options: OrmExecuteOptionsParameter = util.EMPTY_DICT,
719 bind_arguments: Optional[_BindArguments] = None,
720 **kw: Any,
721 ) -> AsyncScalarResult[Any]:
722 """Execute a statement and return a stream of scalar results.
723
724 :return: an :class:`_asyncio.AsyncScalarResult` object
725
726 .. versionadded:: 1.4.24
727
728 .. seealso::
729
730 :meth:`_orm.Session.scalars` - main documentation for scalars
731
732 :meth:`_asyncio.AsyncSession.scalars` - non streaming version
733
734 """
735
736 result = await self.stream(
737 statement,
738 params=params,
739 execution_options=execution_options,
740 bind_arguments=bind_arguments,
741 **kw,
742 )
743 return result.scalars()
744
745 async def delete(self, instance: object) -> None:
746 """Mark an instance as deleted.
747
748 The database delete operation occurs upon ``flush()``.
749
750 As this operation may need to cascade along unloaded relationships,
751 it is awaitable to allow for those queries to take place.
752
753 .. seealso::
754
755 :meth:`_orm.Session.delete` - main documentation for delete
756
757 """
758 await greenlet_spawn(self.sync_session.delete, instance)
759
760 async def merge(
761 self,
762 instance: _O,
763 *,
764 load: bool = True,
765 options: Optional[Sequence[ORMOption]] = None,
766 ) -> _O:
767 """Copy the state of a given instance into a corresponding instance
768 within this :class:`_asyncio.AsyncSession`.
769
770 .. seealso::
771
772 :meth:`_orm.Session.merge` - main documentation for merge
773
774 """
775 return await greenlet_spawn(
776 self.sync_session.merge, instance, load=load, options=options
777 )
778
779 async def flush(self, objects: Optional[Sequence[Any]] = None) -> None:
780 """Flush all the object changes to the database.
781
782 .. seealso::
783
784 :meth:`_orm.Session.flush` - main documentation for flush
785
786 """
787 await greenlet_spawn(self.sync_session.flush, objects=objects)
788
789 def get_transaction(self) -> Optional[AsyncSessionTransaction]:
790 """Return the current root transaction in progress, if any.
791
792 :return: an :class:`_asyncio.AsyncSessionTransaction` object, or
793 ``None``.
794
795 .. versionadded:: 1.4.18
796
797 """
798 trans = self.sync_session.get_transaction()
799 if trans is not None:
800 return AsyncSessionTransaction._retrieve_proxy_for_target(
801 trans, async_session=self
802 )
803 else:
804 return None
805
806 def get_nested_transaction(self) -> Optional[AsyncSessionTransaction]:
807 """Return the current nested transaction in progress, if any.
808
809 :return: an :class:`_asyncio.AsyncSessionTransaction` object, or
810 ``None``.
811
812 .. versionadded:: 1.4.18
813
814 """
815
816 trans = self.sync_session.get_nested_transaction()
817 if trans is not None:
818 return AsyncSessionTransaction._retrieve_proxy_for_target(
819 trans, async_session=self
820 )
821 else:
822 return None
823
824 def get_bind(
825 self,
826 mapper: Optional[_EntityBindKey[_O]] = None,
827 clause: Optional[ClauseElement] = None,
828 bind: Optional[_SessionBind] = None,
829 **kw: Any,
830 ) -> Union[Engine, Connection]:
831 """Return a "bind" to which the synchronous proxied :class:`_orm.Session`
832 is bound.
833
834 Unlike the :meth:`_orm.Session.get_bind` method, this method is
835 currently **not** used by this :class:`.AsyncSession` in any way
836 in order to resolve engines for requests.
837
838 .. note::
839
840 This method proxies directly to the :meth:`_orm.Session.get_bind`
841 method, however is currently **not** useful as an override target,
842 in contrast to that of the :meth:`_orm.Session.get_bind` method.
843 The example below illustrates how to implement custom
844 :meth:`_orm.Session.get_bind` schemes that work with
845 :class:`.AsyncSession` and :class:`.AsyncEngine`.
846
847 The pattern introduced at :ref:`session_custom_partitioning`
848 illustrates how to apply a custom bind-lookup scheme to a
849 :class:`_orm.Session` given a set of :class:`_engine.Engine` objects.
850 To apply a corresponding :meth:`_orm.Session.get_bind` implementation
851 for use with a :class:`.AsyncSession` and :class:`.AsyncEngine`
852 objects, continue to subclass :class:`_orm.Session` and apply it to
853 :class:`.AsyncSession` using
854 :paramref:`.AsyncSession.sync_session_class`. The inner method must
855 continue to return :class:`_engine.Engine` instances, which can be
856 acquired from a :class:`_asyncio.AsyncEngine` using the
857 :attr:`_asyncio.AsyncEngine.sync_engine` attribute::
858
859 # using example from "Custom Vertical Partitioning"
860
861
862 import random
863
864 from sqlalchemy.ext.asyncio import AsyncSession
865 from sqlalchemy.ext.asyncio import create_async_engine
866 from sqlalchemy.ext.asyncio import async_sessionmaker
867 from sqlalchemy.orm import Session
868
869 # construct async engines w/ async drivers
870 engines = {
871 "leader": create_async_engine("sqlite+aiosqlite:///leader.db"),
872 "other": create_async_engine("sqlite+aiosqlite:///other.db"),
873 "follower1": create_async_engine("sqlite+aiosqlite:///follower1.db"),
874 "follower2": create_async_engine("sqlite+aiosqlite:///follower2.db"),
875 }
876
877
878 class RoutingSession(Session):
879 def get_bind(self, mapper=None, clause=None, **kw):
880 # within get_bind(), return sync engines
881 if mapper and issubclass(mapper.class_, MyOtherClass):
882 return engines["other"].sync_engine
883 elif self._flushing or isinstance(clause, (Update, Delete)):
884 return engines["leader"].sync_engine
885 else:
886 return engines[
887 random.choice(["follower1", "follower2"])
888 ].sync_engine
889
890
891 # apply to AsyncSession using sync_session_class
892 AsyncSessionMaker = async_sessionmaker(sync_session_class=RoutingSession)
893
894 The :meth:`_orm.Session.get_bind` method is called in a non-asyncio,
895 implicitly non-blocking context in the same manner as ORM event hooks
896 and functions that are invoked via :meth:`.AsyncSession.run_sync`, so
897 routines that wish to run SQL commands inside of
898 :meth:`_orm.Session.get_bind` can continue to do so using
899 blocking-style code, which will be translated to implicitly async calls
900 at the point of invoking IO on the database drivers.
901
902 """ # noqa: E501
903
904 return self.sync_session.get_bind(
905 mapper=mapper, clause=clause, bind=bind, **kw
906 )
907
908 async def connection(
909 self,
910 bind_arguments: Optional[_BindArguments] = None,
911 execution_options: Optional[CoreExecuteOptionsParameter] = None,
912 **kw: Any,
913 ) -> AsyncConnection:
914 r"""Return a :class:`_asyncio.AsyncConnection` object corresponding to
915 this :class:`.Session` object's transactional state.
916
917 This method may also be used to establish execution options for the
918 database connection used by the current transaction.
919
920 .. versionadded:: 1.4.24 Added \**kw arguments which are passed
921 through to the underlying :meth:`_orm.Session.connection` method.
922
923 .. seealso::
924
925 :meth:`_orm.Session.connection` - main documentation for
926 "connection"
927
928 """
929
930 sync_connection = await greenlet_spawn(
931 self.sync_session.connection,
932 bind_arguments=bind_arguments,
933 execution_options=execution_options,
934 **kw,
935 )
936 return engine.AsyncConnection._retrieve_proxy_for_target(
937 sync_connection
938 )
939
940 def begin(self) -> AsyncSessionTransaction:
941 """Return an :class:`_asyncio.AsyncSessionTransaction` object.
942
943 The underlying :class:`_orm.Session` will perform the
944 "begin" action when the :class:`_asyncio.AsyncSessionTransaction`
945 object is entered::
946
947 async with async_session.begin():
948 ... # ORM transaction is begun
949
950 Note that database IO will not normally occur when the session-level
951 transaction is begun, as database transactions begin on an
952 on-demand basis. However, the begin block is async to accommodate
953 for a :meth:`_orm.SessionEvents.after_transaction_create`
954 event hook that may perform IO.
955
956 For a general description of ORM begin, see
957 :meth:`_orm.Session.begin`.
958
959 """
960
961 return AsyncSessionTransaction(self)
962
963 def begin_nested(self) -> AsyncSessionTransaction:
964 """Return an :class:`_asyncio.AsyncSessionTransaction` object
965 which will begin a "nested" transaction, e.g. SAVEPOINT.
966
967 Behavior is the same as that of :meth:`_asyncio.AsyncSession.begin`.
968
969 For a general description of ORM begin nested, see
970 :meth:`_orm.Session.begin_nested`.
971
972 .. seealso::
973
974 :ref:`aiosqlite_serializable` - special workarounds required
975 with the SQLite asyncio driver in order for SAVEPOINT to work
976 correctly.
977
978 """
979
980 return AsyncSessionTransaction(self, nested=True)
981
982 async def rollback(self) -> None:
983 """Rollback the current transaction in progress.
984
985 .. seealso::
986
987 :meth:`_orm.Session.rollback` - main documentation for
988 "rollback"
989 """
990 await greenlet_spawn(self.sync_session.rollback)
991
992 async def commit(self) -> None:
993 """Commit the current transaction in progress.
994
995 .. seealso::
996
997 :meth:`_orm.Session.commit` - main documentation for
998 "commit"
999 """
1000 await greenlet_spawn(self.sync_session.commit)
1001
1002 async def close(self) -> None:
1003 """Close out the transactional resources and ORM objects used by this
1004 :class:`_asyncio.AsyncSession`.
1005
1006 .. seealso::
1007
1008 :meth:`_orm.Session.close` - main documentation for
1009 "close"
1010
1011 :ref:`session_closing` - detail on the semantics of
1012 :meth:`_asyncio.AsyncSession.close` and
1013 :meth:`_asyncio.AsyncSession.reset`.
1014
1015 """
1016 await greenlet_spawn(self.sync_session.close)
1017
1018 async def reset(self) -> None:
1019 """Close out the transactional resources and ORM objects used by this
1020 :class:`_orm.Session`, resetting the session to its initial state.
1021
1022 .. versionadded:: 2.0.22
1023
1024 .. seealso::
1025
1026 :meth:`_orm.Session.reset` - main documentation for
1027 "reset"
1028
1029 :ref:`session_closing` - detail on the semantics of
1030 :meth:`_asyncio.AsyncSession.close` and
1031 :meth:`_asyncio.AsyncSession.reset`.
1032
1033 """
1034 await greenlet_spawn(self.sync_session.reset)
1035
1036 async def aclose(self) -> None:
1037 """A synonym for :meth:`_asyncio.AsyncSession.close`.
1038
1039 The :meth:`_asyncio.AsyncSession.aclose` name is specifically
1040 to support the Python standard library ``@contextlib.aclosing``
1041 context manager function.
1042
1043 .. versionadded:: 2.0.20
1044
1045 """
1046 await self.close()
1047
1048 async def invalidate(self) -> None:
1049 """Close this Session, using connection invalidation.
1050
1051 For a complete description, see :meth:`_orm.Session.invalidate`.
1052 """
1053 await greenlet_spawn(self.sync_session.invalidate)
1054
1055 @classmethod
1056 @util.deprecated(
1057 "2.0",
1058 "The :meth:`.AsyncSession.close_all` method is deprecated and will be "
1059 "removed in a future release. Please refer to "
1060 ":func:`_asyncio.close_all_sessions`.",
1061 )
1062 async def close_all(cls) -> None:
1063 """Close all :class:`_asyncio.AsyncSession` sessions."""
1064 await close_all_sessions()
1065
1066 async def __aenter__(self: _AS) -> _AS:
1067 return self
1068
1069 async def __aexit__(self, type_: Any, value: Any, traceback: Any) -> None:
1070 task = asyncio.create_task(self.close())
1071 await asyncio.shield(task)
1072
1073 def _maker_context_manager(self: _AS) -> _AsyncSessionContextManager[_AS]:
1074 return _AsyncSessionContextManager(self)
1075
1076 # START PROXY METHODS AsyncSession
1077
1078 # code within this block is **programmatically,
1079 # statically generated** by tools/generate_proxy_methods.py
1080
1081 def __contains__(self, instance: object) -> bool:
1082 r"""Return True if the instance is associated with this session.
1083
1084 .. container:: class_bases
1085
1086 Proxied for the :class:`_orm.Session` class on
1087 behalf of the :class:`_asyncio.AsyncSession` class.
1088
1089 The instance may be pending or persistent within the Session for a
1090 result of True.
1091
1092
1093 """ # noqa: E501
1094
1095 return self._proxied.__contains__(instance)
1096
1097 def __iter__(self) -> Iterator[object]:
1098 r"""Iterate over all pending or persistent instances within this
1099 Session.
1100
1101 .. container:: class_bases
1102
1103 Proxied for the :class:`_orm.Session` class on
1104 behalf of the :class:`_asyncio.AsyncSession` class.
1105
1106
1107 """ # noqa: E501
1108
1109 return self._proxied.__iter__()
1110
1111 def add(self, instance: object, _warn: bool = True) -> None:
1112 r"""Place an object into this :class:`_orm.Session`.
1113
1114 .. container:: class_bases
1115
1116 Proxied for the :class:`_orm.Session` class on
1117 behalf of the :class:`_asyncio.AsyncSession` class.
1118
1119 Objects that are in the :term:`transient` state when passed to the
1120 :meth:`_orm.Session.add` method will move to the
1121 :term:`pending` state, until the next flush, at which point they
1122 will move to the :term:`persistent` state.
1123
1124 Objects that are in the :term:`detached` state when passed to the
1125 :meth:`_orm.Session.add` method will move to the :term:`persistent`
1126 state directly.
1127
1128 If the transaction used by the :class:`_orm.Session` is rolled back,
1129 objects which were transient when they were passed to
1130 :meth:`_orm.Session.add` will be moved back to the
1131 :term:`transient` state, and will no longer be present within this
1132 :class:`_orm.Session`.
1133
1134 .. seealso::
1135
1136 :meth:`_orm.Session.add_all`
1137
1138 :ref:`session_adding` - at :ref:`session_basics`
1139
1140
1141 """ # noqa: E501
1142
1143 return self._proxied.add(instance, _warn=_warn)
1144
1145 def add_all(self, instances: Iterable[object]) -> None:
1146 r"""Add the given collection of instances to this :class:`_orm.Session`.
1147
1148 .. container:: class_bases
1149
1150 Proxied for the :class:`_orm.Session` class on
1151 behalf of the :class:`_asyncio.AsyncSession` class.
1152
1153 See the documentation for :meth:`_orm.Session.add` for a general
1154 behavioral description.
1155
1156 .. seealso::
1157
1158 :meth:`_orm.Session.add`
1159
1160 :ref:`session_adding` - at :ref:`session_basics`
1161
1162
1163 """ # noqa: E501
1164
1165 return self._proxied.add_all(instances)
1166
1167 def expire(
1168 self, instance: object, attribute_names: Optional[Iterable[str]] = None
1169 ) -> None:
1170 r"""Expire the attributes on an instance.
1171
1172 .. container:: class_bases
1173
1174 Proxied for the :class:`_orm.Session` class on
1175 behalf of the :class:`_asyncio.AsyncSession` class.
1176
1177 Marks the attributes of an instance as out of date. When an expired
1178 attribute is next accessed, a query will be issued to the
1179 :class:`.Session` object's current transactional context in order to
1180 load all expired attributes for the given instance. Note that
1181 a highly isolated transaction will return the same values as were
1182 previously read in that same transaction, regardless of changes
1183 in database state outside of that transaction.
1184
1185 To expire all objects in the :class:`.Session` simultaneously,
1186 use :meth:`Session.expire_all`.
1187
1188 The :class:`.Session` object's default behavior is to
1189 expire all state whenever the :meth:`Session.rollback`
1190 or :meth:`Session.commit` methods are called, so that new
1191 state can be loaded for the new transaction. For this reason,
1192 calling :meth:`Session.expire` only makes sense for the specific
1193 case that a non-ORM SQL statement was emitted in the current
1194 transaction.
1195
1196 :param instance: The instance to be refreshed.
1197 :param attribute_names: optional list of string attribute names
1198 indicating a subset of attributes to be expired.
1199
1200 .. seealso::
