Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
session.py1948 linesDownload Raw Back to asyncio
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::

Showing the first 1,200 of 1948 lines. Download the file for the rest.

codekingpro/portable-devtools · Team Ai