Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
engine.py1467 linesDownload Raw Back to asyncio
1# ext/asyncio/engine.py
2# Copyright (C) 2020-2024 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
10import contextlib
11from typing import Any
12from typing import AsyncIterator
13from typing import Callable
14from typing import Dict
15from typing import Generator
16from typing import NoReturn
17from typing import Optional
18from typing import overload
19from typing import Tuple
20from typing import Type
21from typing import TYPE_CHECKING
22from typing import TypeVar
23from typing import Union
24
25from . import exc as async_exc
26from .base import asyncstartablecontext
27from .base import GeneratorStartableContext
28from .base import ProxyComparable
29from .base import StartableContext
30from .result import _ensure_sync_result
31from .result import AsyncResult
32from .result import AsyncScalarResult
33from ... import exc
34from ... import inspection
35from ... import util
36from ...engine import Connection
37from ...engine import create_engine as _create_engine
38from ...engine import create_pool_from_url as _create_pool_from_url
39from ...engine import Engine
40from ...engine.base import NestedTransaction
41from ...engine.base import Transaction
42from ...exc import ArgumentError
43from ...util.concurrency import greenlet_spawn
44from ...util.typing import Concatenate
45from ...util.typing import ParamSpec
46
47if TYPE_CHECKING:
48    from ...engine.cursor import CursorResult
49    from ...engine.interfaces import _CoreAnyExecuteParams
50    from ...engine.interfaces import _CoreSingleExecuteParams
51    from ...engine.interfaces import _DBAPIAnyExecuteParams
52    from ...engine.interfaces import _ExecuteOptions
53    from ...engine.interfaces import CompiledCacheType
54    from ...engine.interfaces import CoreExecuteOptionsParameter
55    from ...engine.interfaces import Dialect
56    from ...engine.interfaces import IsolationLevel
57    from ...engine.interfaces import SchemaTranslateMapType
58    from ...engine.result import ScalarResult
59    from ...engine.url import URL
60    from ...pool import Pool
61    from ...pool import PoolProxiedConnection
62    from ...sql._typing import _InfoType
63    from ...sql.base import Executable
64    from ...sql.selectable import TypedReturnsRows
65
66_P = ParamSpec("_P")
67_T = TypeVar("_T", bound=Any)
68
69
70def create_async_engine(url: Union[str, URL], **kw: Any) -> AsyncEngine:
71    """Create a new async engine instance.
72
73    Arguments passed to :func:`_asyncio.create_async_engine` are mostly
74    identical to those passed to the :func:`_sa.create_engine` function.
75    The specified dialect must be an asyncio-compatible dialect
76    such as :ref:`dialect-postgresql-asyncpg`.
77
78    .. versionadded:: 1.4
79
80    :param async_creator: an async callable which returns a driver-level
81        asyncio connection. If given, the function should take no arguments,
82        and return a new asyncio connection from the underlying asyncio
83        database driver; the connection will be wrapped in the appropriate
84        structures to be used with the :class:`.AsyncEngine`.   Note that the
85        parameters specified in the URL are not applied here, and the creator
86        function should use its own connection parameters.
87
88        This parameter is the asyncio equivalent of the
89        :paramref:`_sa.create_engine.creator` parameter of the
90        :func:`_sa.create_engine` function.
91
92        .. versionadded:: 2.0.16
93
94    """
95
96    if kw.get("server_side_cursors", False):
97        raise async_exc.AsyncMethodRequired(
98            "Can't set server_side_cursors for async engine globally; "
99            "use the connection.stream() method for an async "
100            "streaming result set"
101        )
102    kw["_is_async"] = True
103    async_creator = kw.pop("async_creator", None)
104    if async_creator:
105        if kw.get("creator", None):
106            raise ArgumentError(
107                "Can only specify one of 'async_creator' or 'creator', "
108                "not both."
109            )
110
111        def creator() -> Any:
112            # note that to send adapted arguments like
113            # prepared_statement_cache_size, user would use
114            # "creator" and emulate this form here
115            return sync_engine.dialect.dbapi.connect(  # type: ignore
116                async_creator_fn=async_creator
117            )
118
119        kw["creator"] = creator
120    sync_engine = _create_engine(url, **kw)
121    return AsyncEngine(sync_engine)
122
123
124def async_engine_from_config(
125    configuration: Dict[str, Any], prefix: str = "sqlalchemy.", **kwargs: Any
126) -> AsyncEngine:
127    """Create a new AsyncEngine instance using a configuration dictionary.
128
129    This function is analogous to the :func:`_sa.engine_from_config` function
130    in SQLAlchemy Core, except that the requested dialect must be an
131    asyncio-compatible dialect such as :ref:`dialect-postgresql-asyncpg`.
132    The argument signature of the function is identical to that
133    of :func:`_sa.engine_from_config`.
134
135    .. versionadded:: 1.4.29
136
137    """
138    options = {
139        key[len(prefix) :]: value
140        for key, value in configuration.items()
141        if key.startswith(prefix)
142    }
143    options["_coerce_config"] = True
144    options.update(kwargs)
145    url = options.pop("url")
146    return create_async_engine(url, **options)
147
148
149def create_async_pool_from_url(url: Union[str, URL], **kwargs: Any) -> Pool:
150    """Create a new async engine instance.
151
152    Arguments passed to :func:`_asyncio.create_async_pool_from_url` are mostly
153    identical to those passed to the :func:`_sa.create_pool_from_url` function.
154    The specified dialect must be an asyncio-compatible dialect
155    such as :ref:`dialect-postgresql-asyncpg`.
156
157    .. versionadded:: 2.0.10
158
159    """
160    kwargs["_is_async"] = True
161    return _create_pool_from_url(url, **kwargs)
162
163
164class AsyncConnectable:
165    __slots__ = "_slots_dispatch", "__weakref__"
166
167    @classmethod
168    def _no_async_engine_events(cls) -> NoReturn:
169        raise NotImplementedError(
170            "asynchronous events are not implemented at this time.  Apply "
171            "synchronous listeners to the AsyncEngine.sync_engine or "
172            "AsyncConnection.sync_connection attributes."
173        )
174
175
176@util.create_proxy_methods(
177    Connection,
178    ":class:`_engine.Connection`",
179    ":class:`_asyncio.AsyncConnection`",
180    classmethods=[],
181    methods=[],
182    attributes=[
183        "closed",
184        "invalidated",
185        "dialect",
186        "default_isolation_level",
187    ],
188)
189class AsyncConnection(
190    ProxyComparable[Connection],
191    StartableContext["AsyncConnection"],
192    AsyncConnectable,
193):
194    """An asyncio proxy for a :class:`_engine.Connection`.
195
196    :class:`_asyncio.AsyncConnection` is acquired using the
197    :meth:`_asyncio.AsyncEngine.connect`
198    method of :class:`_asyncio.AsyncEngine`::
199
200        from sqlalchemy.ext.asyncio import create_async_engine
201        engine = create_async_engine("postgresql+asyncpg://user:pass@host/dbname")
202
203        async with engine.connect() as conn:
204            result = await conn.execute(select(table))
205
206    .. versionadded:: 1.4
207
208    """  # noqa
209
210    # AsyncConnection is a thin proxy; no state should be added here
211    # that is not retrievable from the "sync" engine / connection, e.g.
212    # current transaction, info, etc.   It should be possible to
213    # create a new AsyncConnection that matches this one given only the
214    # "sync" elements.
215    __slots__ = (
216        "engine",
217        "sync_engine",
218        "sync_connection",
219    )
220
221    def __init__(
222        self,
223        async_engine: AsyncEngine,
224        sync_connection: Optional[Connection] = None,
225    ):
226        self.engine = async_engine
227        self.sync_engine = async_engine.sync_engine
228        self.sync_connection = self._assign_proxied(sync_connection)
229
230    sync_connection: Optional[Connection]
231    """Reference to the sync-style :class:`_engine.Connection` this
232    :class:`_asyncio.AsyncConnection` proxies requests towards.
233
234    This instance can be used as an event target.
235
236    .. seealso::
237
238        :ref:`asyncio_events`
239
240    """
241
242    sync_engine: Engine
243    """Reference to the sync-style :class:`_engine.Engine` this
244    :class:`_asyncio.AsyncConnection` is associated with via its underlying
245    :class:`_engine.Connection`.
246
247    This instance can be used as an event target.
248
249    .. seealso::
250
251        :ref:`asyncio_events`
252
253    """
254
255    @classmethod
256    def _regenerate_proxy_for_target(
257        cls, target: Connection
258    ) -> AsyncConnection:
259        return AsyncConnection(
260            AsyncEngine._retrieve_proxy_for_target(target.engine), target
261        )
262
263    async def start(
264        self, is_ctxmanager: bool = False  # noqa: U100
265    ) -> AsyncConnection:
266        """Start this :class:`_asyncio.AsyncConnection` object's context
267        outside of using a Python ``with:`` block.
268
269        """
270        if self.sync_connection:
271            raise exc.InvalidRequestError("connection is already started")
272        self.sync_connection = self._assign_proxied(
273            await greenlet_spawn(self.sync_engine.connect)
274        )
275        return self
276
277    @property
278    def connection(self) -> NoReturn:
279        """Not implemented for async; call
280        :meth:`_asyncio.AsyncConnection.get_raw_connection`.
281        """
282        raise exc.InvalidRequestError(
283            "AsyncConnection.connection accessor is not implemented as the "
284            "attribute may need to reconnect on an invalidated connection.  "
285            "Use the get_raw_connection() method."
286        )
287
288    async def get_raw_connection(self) -> PoolProxiedConnection:
289        """Return the pooled DBAPI-level connection in use by this
290        :class:`_asyncio.AsyncConnection`.
291
292        This is a SQLAlchemy connection-pool proxied connection
293        which then has the attribute
294        :attr:`_pool._ConnectionFairy.driver_connection` that refers to the
295        actual driver connection. Its
296        :attr:`_pool._ConnectionFairy.dbapi_connection` refers instead
297        to an :class:`_engine.AdaptedConnection` instance that
298        adapts the driver connection to the DBAPI protocol.
299
300        """
301
302        return await greenlet_spawn(getattr, self._proxied, "connection")
303
304    @util.ro_non_memoized_property
305    def info(self) -> _InfoType:
306        """Return the :attr:`_engine.Connection.info` dictionary of the
307        underlying :class:`_engine.Connection`.
308
309        This dictionary is freely writable for user-defined state to be
310        associated with the database connection.
311
312        This attribute is only available if the :class:`.AsyncConnection` is
313        currently connected.   If the :attr:`.AsyncConnection.closed` attribute
314        is ``True``, then accessing this attribute will raise
315        :class:`.ResourceClosedError`.
316
317        .. versionadded:: 1.4.0b2
318
319        """
320        return self._proxied.info
321
322    @util.ro_non_memoized_property
323    def _proxied(self) -> Connection:
324        if not self.sync_connection:
325            self._raise_for_not_started()
326        return self.sync_connection
327
328    def begin(self) -> AsyncTransaction:
329        """Begin a transaction prior to autobegin occurring."""
330        assert self._proxied
331        return AsyncTransaction(self)
332
333    def begin_nested(self) -> AsyncTransaction:
334        """Begin a nested transaction and return a transaction handle."""
335        assert self._proxied
336        return AsyncTransaction(self, nested=True)
337
338    async def invalidate(
339        self, exception: Optional[BaseException] = None
340    ) -> None:
341        """Invalidate the underlying DBAPI connection associated with
342        this :class:`_engine.Connection`.
343
344        See the method :meth:`_engine.Connection.invalidate` for full
345        detail on this method.
346
347        """
348
349        return await greenlet_spawn(
350            self._proxied.invalidate, exception=exception
351        )
352
353    async def get_isolation_level(self) -> IsolationLevel:
354        return await greenlet_spawn(self._proxied.get_isolation_level)
355
356    def in_transaction(self) -> bool:
357        """Return True if a transaction is in progress."""
358
359        return self._proxied.in_transaction()
360
361    def in_nested_transaction(self) -> bool:
362        """Return True if a transaction is in progress.
363
364        .. versionadded:: 1.4.0b2
365
366        """
367        return self._proxied.in_nested_transaction()
368
369    def get_transaction(self) -> Optional[AsyncTransaction]:
370        """Return an :class:`.AsyncTransaction` representing the current
371        transaction, if any.
372
373        This makes use of the underlying synchronous connection's
374        :meth:`_engine.Connection.get_transaction` method to get the current
375        :class:`_engine.Transaction`, which is then proxied in a new
376        :class:`.AsyncTransaction` object.
377
378        .. versionadded:: 1.4.0b2
379
380        """
381
382        trans = self._proxied.get_transaction()
383        if trans is not None:
384            return AsyncTransaction._retrieve_proxy_for_target(trans)
385        else:
386            return None
387
388    def get_nested_transaction(self) -> Optional[AsyncTransaction]:
389        """Return an :class:`.AsyncTransaction` representing the current
390        nested (savepoint) transaction, if any.
391
392        This makes use of the underlying synchronous connection's
393        :meth:`_engine.Connection.get_nested_transaction` method to get the
394        current :class:`_engine.Transaction`, which is then proxied in a new
395        :class:`.AsyncTransaction` object.
396
397        .. versionadded:: 1.4.0b2
398
399        """
400
401        trans = self._proxied.get_nested_transaction()
402        if trans is not None:
403            return AsyncTransaction._retrieve_proxy_for_target(trans)
404        else:
405            return None
406
407    @overload
408    async def execution_options(
409        self,
410        *,
411        compiled_cache: Optional[CompiledCacheType] = ...,
412        logging_token: str = ...,
413        isolation_level: IsolationLevel = ...,
414        no_parameters: bool = False,
415        stream_results: bool = False,
416        max_row_buffer: int = ...,
417        yield_per: int = ...,
418        insertmanyvalues_page_size: int = ...,
419        schema_translate_map: Optional[SchemaTranslateMapType] = ...,
420        preserve_rowcount: bool = False,
421        **opt: Any,
422    ) -> AsyncConnection: ...
423
424    @overload
425    async def execution_options(self, **opt: Any) -> AsyncConnection: ...
426
427    async def execution_options(self, **opt: Any) -> AsyncConnection:
428        r"""Set non-SQL options for the connection which take effect
429        during execution.
430
431        This returns this :class:`_asyncio.AsyncConnection` object with
432        the new options added.
433
434        See :meth:`_engine.Connection.execution_options` for full details
435        on this method.
436
437        """
438
439        conn = self._proxied
440        c2 = await greenlet_spawn(conn.execution_options, **opt)
441        assert c2 is conn
442        return self
443
444    async def commit(self) -> None:
445        """Commit the transaction that is currently in progress.
446
447        This method commits the current transaction if one has been started.
448        If no transaction was started, the method has no effect, assuming
449        the connection is in a non-invalidated state.
450
451        A transaction is begun on a :class:`_engine.Connection` automatically
452        whenever a statement is first executed, or when the
453        :meth:`_engine.Connection.begin` method is called.
454
455        """
456        await greenlet_spawn(self._proxied.commit)
457
458    async def rollback(self) -> None:
459        """Roll back the transaction that is currently in progress.
460
461        This method rolls back the current transaction if one has been started.
462        If no transaction was started, the method has no effect.  If a
463        transaction was started and the connection is in an invalidated state,
464        the transaction is cleared using this method.
465
466        A transaction is begun on a :class:`_engine.Connection` automatically
467        whenever a statement is first executed, or when the
468        :meth:`_engine.Connection.begin` method is called.
469
470
471        """
472        await greenlet_spawn(self._proxied.rollback)
473
474    async def close(self) -> None:
475        """Close this :class:`_asyncio.AsyncConnection`.
476
477        This has the effect of also rolling back the transaction if one
478        is in place.
479
480        """
481        await greenlet_spawn(self._proxied.close)
482
483    async def aclose(self) -> None:
484        """A synonym for :meth:`_asyncio.AsyncConnection.close`.
485
486        The :meth:`_asyncio.AsyncConnection.aclose` name is specifically
487        to support the Python standard library ``@contextlib.aclosing``
488        context manager function.
489
490        .. versionadded:: 2.0.20
491
492        """
493        await self.close()
494
495    async def exec_driver_sql(
496        self,
497        statement: str,
498        parameters: Optional[_DBAPIAnyExecuteParams] = None,
499        execution_options: Optional[CoreExecuteOptionsParameter] = None,
500    ) -> CursorResult[Any]:
501        r"""Executes a driver-level SQL string and return buffered
502        :class:`_engine.Result`.
503
504        """
505
506        result = await greenlet_spawn(
507            self._proxied.exec_driver_sql,
508            statement,
509            parameters,
510            execution_options,
511            _require_await=True,
512        )
513
514        return await _ensure_sync_result(result, self.exec_driver_sql)
515
516    @overload
517    def stream(
518        self,
519        statement: TypedReturnsRows[_T],
520        parameters: Optional[_CoreAnyExecuteParams] = None,
521        *,
522        execution_options: Optional[CoreExecuteOptionsParameter] = None,
523    ) -> GeneratorStartableContext[AsyncResult[_T]]: ...
524
525    @overload
526    def stream(
527        self,
528        statement: Executable,
529        parameters: Optional[_CoreAnyExecuteParams] = None,
530        *,
531        execution_options: Optional[CoreExecuteOptionsParameter] = None,
532    ) -> GeneratorStartableContext[AsyncResult[Any]]: ...
533
534    @asyncstartablecontext
535    async def stream(
536        self,
537        statement: Executable,
538        parameters: Optional[_CoreAnyExecuteParams] = None,
539        *,
540        execution_options: Optional[CoreExecuteOptionsParameter] = None,
541    ) -> AsyncIterator[AsyncResult[Any]]:
542        """Execute a statement and return an awaitable yielding a
543        :class:`_asyncio.AsyncResult` object.
544
545        E.g.::
546
547            result = await conn.stream(stmt):
548            async for row in result:
549                print(f"{row}")
550
551        The :meth:`.AsyncConnection.stream`
552        method supports optional context manager use against the
553        :class:`.AsyncResult` object, as in::
554
555            async with conn.stream(stmt) as result:
556                async for row in result:
557                    print(f"{row}")
558
559        In the above pattern, the :meth:`.AsyncResult.close` method is
560        invoked unconditionally, even if the iterator is interrupted by an
561        exception throw.   Context manager use remains optional, however,
562        and the function may be called in either an ``async with fn():`` or
563        ``await fn()`` style.
564
565        .. versionadded:: 2.0.0b3 added context manager support
566
567
568        :return: an awaitable object that will yield an
569         :class:`_asyncio.AsyncResult` object.
570
571        .. seealso::
572
573            :meth:`.AsyncConnection.stream_scalars`
574
575        """
576        if not self.dialect.supports_server_side_cursors:
577            raise exc.InvalidRequestError(
578                "Cant use `stream` or `stream_scalars` with the current "
579                "dialect since it does not support server side cursors."
580            )
581
582        result = await greenlet_spawn(
583            self._proxied.execute,
584            statement,
585            parameters,
586            execution_options=util.EMPTY_DICT.merge_with(
587                execution_options, {"stream_results": True}
588            ),
589            _require_await=True,
590        )
591        assert result.context._is_server_side
592        ar = AsyncResult(result)
593        try:
594            yield ar
595        except GeneratorExit:
596            pass
597        else:
598            task = asyncio.create_task(ar.close())
599            await asyncio.shield(task)
600
601    @overload
602    async def execute(
603        self,
604        statement: TypedReturnsRows[_T],
605        parameters: Optional[_CoreAnyExecuteParams] = None,
606        *,
607        execution_options: Optional[CoreExecuteOptionsParameter] = None,
608    ) -> CursorResult[_T]: ...
609
610    @overload
611    async def execute(
612        self,
613        statement: Executable,
614        parameters: Optional[_CoreAnyExecuteParams] = None,
615        *,
616        execution_options: Optional[CoreExecuteOptionsParameter] = None,
617    ) -> CursorResult[Any]: ...
618
619    async def execute(
620        self,
621        statement: Executable,
622        parameters: Optional[_CoreAnyExecuteParams] = None,
623        *,
624        execution_options: Optional[CoreExecuteOptionsParameter] = None,
625    ) -> CursorResult[Any]:
626        r"""Executes a SQL statement construct and return a buffered
627        :class:`_engine.Result`.
628
629        :param object: The statement to be executed.  This is always
630         an object that is in both the :class:`_expression.ClauseElement` and
631         :class:`_expression.Executable` hierarchies, including:
632
633         * :class:`_expression.Select`
634         * :class:`_expression.Insert`, :class:`_expression.Update`,
635           :class:`_expression.Delete`
636         * :class:`_expression.TextClause` and
637           :class:`_expression.TextualSelect`
638         * :class:`_schema.DDL` and objects which inherit from
639           :class:`_schema.ExecutableDDLElement`
640
641        :param parameters: parameters which will be bound into the statement.
642         This may be either a dictionary of parameter names to values,
643         or a mutable sequence (e.g. a list) of dictionaries.  When a
644         list of dictionaries is passed, the underlying statement execution
645         will make use of the DBAPI ``cursor.executemany()`` method.
646         When a single dictionary is passed, the DBAPI ``cursor.execute()``
647         method will be used.
648
649        :param execution_options: optional dictionary of execution options,
650         which will be associated with the statement execution.  This
651         dictionary can provide a subset of the options that are accepted
652         by :meth:`_engine.Connection.execution_options`.
653
654        :return: a :class:`_engine.Result` object.
655
656        """
657        result = await greenlet_spawn(
658            self._proxied.execute,
659            statement,
660            parameters,
661            execution_options=execution_options,
662            _require_await=True,
663        )
664        return await _ensure_sync_result(result, self.execute)
665
666    @overload
667    async def scalar(
668        self,
669        statement: TypedReturnsRows[Tuple[_T]],
670        parameters: Optional[_CoreSingleExecuteParams] = None,
671        *,
672        execution_options: Optional[CoreExecuteOptionsParameter] = None,
673    ) -> Optional[_T]: ...
674
675    @overload
676    async def scalar(
677        self,
678        statement: Executable,
679        parameters: Optional[_CoreSingleExecuteParams] = None,
680        *,
681        execution_options: Optional[CoreExecuteOptionsParameter] = None,
682    ) -> Any: ...
683
684    async def scalar(
685        self,
686        statement: Executable,
687        parameters: Optional[_CoreSingleExecuteParams] = None,
688        *,
689        execution_options: Optional[CoreExecuteOptionsParameter] = None,
690    ) -> Any:
691        r"""Executes a SQL statement construct and returns a scalar object.
692
693        This method is shorthand for invoking the
694        :meth:`_engine.Result.scalar` method after invoking the
695        :meth:`_engine.Connection.execute` method.  Parameters are equivalent.
696
697        :return: a scalar Python value representing the first column of the
698         first row returned.
699
700        """
701        result = await self.execute(
702            statement, parameters, execution_options=execution_options
703        )
704        return result.scalar()
705
706    @overload
707    async def scalars(
708        self,
709        statement: TypedReturnsRows[Tuple[_T]],
710        parameters: Optional[_CoreAnyExecuteParams] = None,
711        *,
712        execution_options: Optional[CoreExecuteOptionsParameter] = None,
713    ) -> ScalarResult[_T]: ...
714
715    @overload
716    async def scalars(
717        self,
718        statement: Executable,
719        parameters: Optional[_CoreAnyExecuteParams] = None,
720        *,
721        execution_options: Optional[CoreExecuteOptionsParameter] = None,
722    ) -> ScalarResult[Any]: ...
723
724    async def scalars(
725        self,
726        statement: Executable,
727        parameters: Optional[_CoreAnyExecuteParams] = None,
728        *,
729        execution_options: Optional[CoreExecuteOptionsParameter] = None,
730    ) -> ScalarResult[Any]:
731        r"""Executes a SQL statement construct and returns a scalar objects.
732
733        This method is shorthand for invoking the
734        :meth:`_engine.Result.scalars` method after invoking the
735        :meth:`_engine.Connection.execute` method.  Parameters are equivalent.
736
737        :return: a :class:`_engine.ScalarResult` object.
738
739        .. versionadded:: 1.4.24
740
741        """
742        result = await self.execute(
743            statement, parameters, execution_options=execution_options
744        )
745        return result.scalars()
746
747    @overload
748    def stream_scalars(
749        self,
750        statement: TypedReturnsRows[Tuple[_T]],
751        parameters: Optional[_CoreSingleExecuteParams] = None,
752        *,
753        execution_options: Optional[CoreExecuteOptionsParameter] = None,
754    ) -> GeneratorStartableContext[AsyncScalarResult[_T]]: ...
755
756    @overload
757    def stream_scalars(
758        self,
759        statement: Executable,
760        parameters: Optional[_CoreSingleExecuteParams] = None,
761        *,
762        execution_options: Optional[CoreExecuteOptionsParameter] = None,
763    ) -> GeneratorStartableContext[AsyncScalarResult[Any]]: ...
764
765    @asyncstartablecontext
766    async def stream_scalars(
767        self,
768        statement: Executable,
769        parameters: Optional[_CoreSingleExecuteParams] = None,
770        *,
771        execution_options: Optional[CoreExecuteOptionsParameter] = None,
772    ) -> AsyncIterator[AsyncScalarResult[Any]]:
773        r"""Execute a statement and return an awaitable yielding a
774        :class:`_asyncio.AsyncScalarResult` object.
775
776        E.g.::
777
778            result = await conn.stream_scalars(stmt)
779            async for scalar in result:
780                print(f"{scalar}")
781
782        This method is shorthand for invoking the
783        :meth:`_engine.AsyncResult.scalars` method after invoking the
784        :meth:`_engine.Connection.stream` method.  Parameters are equivalent.
785
786        The :meth:`.AsyncConnection.stream_scalars`
787        method supports optional context manager use against the
788        :class:`.AsyncScalarResult` object, as in::
789
790            async with conn.stream_scalars(stmt) as result:
791                async for scalar in result:
792                    print(f"{scalar}")
793
794        In the above pattern, the :meth:`.AsyncScalarResult.close` method is
795        invoked unconditionally, even if the iterator is interrupted by an
796        exception throw.  Context manager use remains optional, however,
797        and the function may be called in either an ``async with fn():`` or
798        ``await fn()`` style.
799
800        .. versionadded:: 2.0.0b3 added context manager support
801
802        :return: an awaitable object that will yield an
803         :class:`_asyncio.AsyncScalarResult` object.
804
805        .. versionadded:: 1.4.24
806
807        .. seealso::
808
809            :meth:`.AsyncConnection.stream`
810
811        """
812
813        async with self.stream(
814            statement, parameters, execution_options=execution_options
815        ) as result:
816            yield result.scalars()
817
818    async def run_sync(
819        self,
820        fn: Callable[Concatenate[Connection, _P], _T],
821        *arg: _P.args,
822        **kw: _P.kwargs,
823    ) -> _T:
824        """Invoke the given synchronous (i.e. not async) callable,
825        passing a synchronous-style :class:`_engine.Connection` as the first
826        argument.
827
828        This method allows traditional synchronous SQLAlchemy functions to
829        run within the context of an asyncio application.
830
831        E.g.::
832
833            def do_something_with_core(conn: Connection, arg1: int, arg2: str) -> str:
834                '''A synchronous function that does not require awaiting
835
836                :param conn: a Core SQLAlchemy Connection, used synchronously
837
838                :return: an optional return value is supported
839
840                '''
841                conn.execute(
842                    some_table.insert().values(int_col=arg1, str_col=arg2)
843                )
844                return "success"
845
846
847            async def do_something_async(async_engine: AsyncEngine) -> None:
848                '''an async function that uses awaiting'''
849
850                async with async_engine.begin() as async_conn:
851                    # run do_something_with_core() with a sync-style
852                    # Connection, proxied into an awaitable
853                    return_code = await async_conn.run_sync(do_something_with_core, 5, "strval")
854                    print(return_code)
855
856        This method maintains the asyncio event loop all the way through
857        to the database connection by running the given callable in a
858        specially instrumented greenlet.
859
860        The most rudimentary use of :meth:`.AsyncConnection.run_sync` is to
861        invoke methods such as :meth:`_schema.MetaData.create_all`, given
862        an :class:`.AsyncConnection` that needs to be provided to
863        :meth:`_schema.MetaData.create_all` as a :class:`_engine.Connection`
864        object::
865
866            # run metadata.create_all(conn) with a sync-style Connection,
867            # proxied into an awaitable
868            with async_engine.begin() as conn:
869                await conn.run_sync(metadata.create_all)
870
871        .. note::
872
873            The provided callable is invoked inline within the asyncio event
874            loop, and will block on traditional IO calls.  IO within this
875            callable should only call into SQLAlchemy's asyncio database
876            APIs which will be properly adapted to the greenlet context.
877
878        .. seealso::
879
880            :meth:`.AsyncSession.run_sync`
881
882            :ref:`session_run_sync`
883
884        """  # noqa: E501
885
886        return await greenlet_spawn(
887            fn, self._proxied, *arg, _require_await=False, **kw
888        )
889
890    def __await__(self) -> Generator[Any, None, AsyncConnection]:
891        return self.start().__await__()
892
893    async def __aexit__(self, type_: Any, value: Any, traceback: Any) -> None:
894        task = asyncio.create_task(self.close())
895        await asyncio.shield(task)
896
897    # START PROXY METHODS AsyncConnection
898
899    # code within this block is **programmatically,
900    # statically generated** by tools/generate_proxy_methods.py
901
902    @property
903    def closed(self) -> Any:
904        r"""Return True if this connection is closed.
905
906        .. container:: class_bases
907
908            Proxied for the :class:`_engine.Connection` class
909            on behalf of the :class:`_asyncio.AsyncConnection` class.
910
911        """  # noqa: E501
912
913        return self._proxied.closed
914
915    @property
916    def invalidated(self) -> Any:
917        r"""Return True if this connection was invalidated.
918
919        .. container:: class_bases
920
921            Proxied for the :class:`_engine.Connection` class
922            on behalf of the :class:`_asyncio.AsyncConnection` class.
923
924        This does not indicate whether or not the connection was
925        invalidated at the pool level, however
926
927
928        """  # noqa: E501
929
930        return self._proxied.invalidated
931
932    @property
933    def dialect(self) -> Dialect:
934        r"""Proxy for the :attr:`_engine.Connection.dialect` attribute
935        on behalf of the :class:`_asyncio.AsyncConnection` class.
936
937        """  # noqa: E501
938
939        return self._proxied.dialect
940
941    @dialect.setter
942    def dialect(self, attr: Dialect) -> None:
943        self._proxied.dialect = attr
944
945    @property
946    def default_isolation_level(self) -> Any:
947        r"""The initial-connection time isolation level associated with the
948        :class:`_engine.Dialect` in use.
949
950        .. container:: class_bases
951
952            Proxied for the :class:`_engine.Connection` class
953            on behalf of the :class:`_asyncio.AsyncConnection` class.
954
955        This value is independent of the
956        :paramref:`.Connection.execution_options.isolation_level` and
957        :paramref:`.Engine.execution_options.isolation_level` execution
958        options, and is determined by the :class:`_engine.Dialect` when the
959        first connection is created, by performing a SQL query against the
960        database for the current isolation level before any additional commands
961        have been emitted.
962
963        Calling this accessor does not invoke any new SQL queries.
964
965        .. seealso::
966
967            :meth:`_engine.Connection.get_isolation_level`
968            - view current actual isolation level
969
970            :paramref:`_sa.create_engine.isolation_level`
971            - set per :class:`_engine.Engine` isolation level
972
973            :paramref:`.Connection.execution_options.isolation_level`
974            - set per :class:`_engine.Connection` isolation level
975
976
977        """  # noqa: E501
978
979        return self._proxied.default_isolation_level
980
981    # END PROXY METHODS AsyncConnection
982
983
984@util.create_proxy_methods(
985    Engine,
986    ":class:`_engine.Engine`",
987    ":class:`_asyncio.AsyncEngine`",
988    classmethods=[],
989    methods=[
990        "clear_compiled_cache",
991        "update_execution_options",
992        "get_execution_options",
993    ],
994    attributes=["url", "pool", "dialect", "engine", "name", "driver", "echo"],
995)
996class AsyncEngine(ProxyComparable[Engine], AsyncConnectable):
997    """An asyncio proxy for a :class:`_engine.Engine`.
998
999    :class:`_asyncio.AsyncEngine` is acquired using the
1000    :func:`_asyncio.create_async_engine` function::
1001
1002        from sqlalchemy.ext.asyncio import create_async_engine
1003        engine = create_async_engine("postgresql+asyncpg://user:pass@host/dbname")
1004
1005    .. versionadded:: 1.4
1006
1007    """  # noqa
1008
1009    # AsyncEngine is a thin proxy; no state should be added here
1010    # that is not retrievable from the "sync" engine / connection, e.g.
1011    # current transaction, info, etc.   It should be possible to
1012    # create a new AsyncEngine that matches this one given only the
1013    # "sync" elements.
1014    __slots__ = "sync_engine"
1015
1016    _connection_cls: Type[AsyncConnection] = AsyncConnection
1017
1018    sync_engine: Engine
1019    """Reference to the sync-style :class:`_engine.Engine` this
1020    :class:`_asyncio.AsyncEngine` proxies requests towards.
1021
1022    This instance can be used as an event target.
1023
1024    .. seealso::
1025
1026        :ref:`asyncio_events`
1027    """
1028
1029    def __init__(self, sync_engine: Engine):
1030        if not sync_engine.dialect.is_async:
1031            raise exc.InvalidRequestError(
1032                "The asyncio extension requires an async driver to be used. "
1033                f"The loaded {sync_engine.dialect.driver!r} is not async."
1034            )
1035        self.sync_engine = self._assign_proxied(sync_engine)
1036
1037    @util.ro_non_memoized_property
1038    def _proxied(self) -> Engine:
1039        return self.sync_engine
1040
1041    @classmethod
1042    def _regenerate_proxy_for_target(cls, target: Engine) -> AsyncEngine:
1043        return AsyncEngine(target)
1044
1045    @contextlib.asynccontextmanager
1046    async def begin(self) -> AsyncIterator[AsyncConnection]:
1047        """Return a context manager which when entered will deliver an
1048        :class:`_asyncio.AsyncConnection` with an
1049        :class:`_asyncio.AsyncTransaction` established.
1050
1051        E.g.::
1052
1053            async with async_engine.begin() as conn:
1054                await conn.execute(
1055                    text("insert into table (x, y, z) values (1, 2, 3)")
1056                )
1057                await conn.execute(text("my_special_procedure(5)"))
1058
1059
1060        """
1061        conn = self.connect()
1062
1063        async with conn:
1064            async with conn.begin():
1065                yield conn
1066
1067    def connect(self) -> AsyncConnection:
1068        """Return an :class:`_asyncio.AsyncConnection` object.
1069
1070        The :class:`_asyncio.AsyncConnection` will procure a database
1071        connection from the underlying connection pool when it is entered
1072        as an async context manager::
1073
1074            async with async_engine.connect() as conn:
1075                result = await conn.execute(select(user_table))
1076
1077        The :class:`_asyncio.AsyncConnection` may also be started outside of a
1078        context manager by invoking its :meth:`_asyncio.AsyncConnection.start`
1079        method.
1080
1081        """
1082
1083        return self._connection_cls(self)
1084
1085    async def raw_connection(self) -> PoolProxiedConnection:
1086        """Return a "raw" DBAPI connection from the connection pool.
1087
1088        .. seealso::
1089
1090            :ref:`dbapi_connections`
1091
1092        """
1093        return await greenlet_spawn(self.sync_engine.raw_connection)
1094
1095    @overload
1096    def execution_options(
1097        self,
1098        *,
1099        compiled_cache: Optional[CompiledCacheType] = ...,
1100        logging_token: str = ...,
1101        isolation_level: IsolationLevel = ...,
1102        insertmanyvalues_page_size: int = ...,
1103        schema_translate_map: Optional[SchemaTranslateMapType] = ...,
1104        **opt: Any,
1105    ) -> AsyncEngine: ...
1106
1107    @overload
1108    def execution_options(self, **opt: Any) -> AsyncEngine: ...
1109
1110    def execution_options(self, **opt: Any) -> AsyncEngine:
1111        """Return a new :class:`_asyncio.AsyncEngine` that will provide
1112        :class:`_asyncio.AsyncConnection` objects with the given execution
1113        options.
1114
1115        Proxied from :meth:`_engine.Engine.execution_options`.  See that
1116        method for details.
1117
1118        """
1119
1120        return AsyncEngine(self.sync_engine.execution_options(**opt))
1121
1122    async def dispose(self, close: bool = True) -> None:
1123        """Dispose of the connection pool used by this
1124        :class:`_asyncio.AsyncEngine`.
1125
1126        :param close: if left at its default of ``True``, has the
1127         effect of fully closing all **currently checked in**
1128         database connections.  Connections that are still checked out
1129         will **not** be closed, however they will no longer be associated
1130         with this :class:`_engine.Engine`,
1131         so when they are closed individually, eventually the
1132         :class:`_pool.Pool` which they are associated with will
1133         be garbage collected and they will be closed out fully, if
1134         not already closed on checkin.
1135
1136         If set to ``False``, the previous connection pool is de-referenced,
1137         and otherwise not touched in any way.
1138
1139        .. seealso::
1140
1141            :meth:`_engine.Engine.dispose`
1142
1143        """
1144
1145        await greenlet_spawn(self.sync_engine.dispose, close=close)
1146
1147    # START PROXY METHODS AsyncEngine
1148
1149    # code within this block is **programmatically,
1150    # statically generated** by tools/generate_proxy_methods.py
1151
1152    def clear_compiled_cache(self) -> None:
1153        r"""Clear the compiled cache associated with the dialect.
1154
1155        .. container:: class_bases
1156
1157            Proxied for the :class:`_engine.Engine` class on
1158            behalf of the :class:`_asyncio.AsyncEngine` class.
1159
1160        This applies **only** to the built-in cache that is established
1161        via the :paramref:`_engine.create_engine.query_cache_size` parameter.
1162        It will not impact any dictionary caches that were passed via the
1163        :paramref:`.Connection.execution_options.compiled_cache` parameter.
1164
1165        .. versionadded:: 1.4
1166
1167
1168        """  # noqa: E501
1169
1170        return self._proxied.clear_compiled_cache()
1171
1172    def update_execution_options(self, **opt: Any) -> None:
1173        r"""Update the default execution_options dictionary
1174        of this :class:`_engine.Engine`.
1175
1176        .. container:: class_bases
1177
1178            Proxied for the :class:`_engine.Engine` class on
1179            behalf of the :class:`_asyncio.AsyncEngine` class.
1180
1181        The given keys/values in \**opt are added to the
1182        default execution options that will be used for
1183        all connections.  The initial contents of this dictionary
1184        can be sent via the ``execution_options`` parameter
1185        to :func:`_sa.create_engine`.
1186
1187        .. seealso::
1188
1189            :meth:`_engine.Connection.execution_options`
1190
1191            :meth:`_engine.Engine.execution_options`
1192
1193
1194        """  # noqa: E501
1195
1196        return self._proxied.update_execution_options(**opt)
1197
1198    def get_execution_options(self) -> _ExecuteOptions:
1199        r"""Get the non-SQL options which will take effect during execution.
1200

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

codekingpro/portable-devtools · Team Ai