codekingpro/portable-devtools
114k
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
