codekingpro/portable-devtools
115k
1# engine/base.py
2# Copyright (C) 2005-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
7"""Defines :class:`_engine.Connection` and :class:`_engine.Engine`."""
8from __future__ import annotations
9
10import contextlib
11import sys
12import typing
13from typing import Any
14from typing import Callable
15from typing import cast
16from typing import Iterable
17from typing import Iterator
18from typing import List
19from typing import Mapping
20from typing import NoReturn
21from typing import Optional
22from typing import overload
23from typing import Tuple
24from typing import Type
25from typing import TypeVar
26from typing import Union
27
28from .interfaces import BindTyping
29from .interfaces import ConnectionEventsTarget
30from .interfaces import DBAPICursor
31from .interfaces import ExceptionContext
32from .interfaces import ExecuteStyle
33from .interfaces import ExecutionContext
34from .interfaces import IsolationLevel
35from .util import _distill_params_20
36from .util import _distill_raw_params
37from .util import TransactionalContext
38from .. import exc
39from .. import inspection
40from .. import log
41from .. import util
42from ..sql import compiler
43from ..sql import util as sql_util
44
45if typing.TYPE_CHECKING:
46 from . import CursorResult
47 from . import ScalarResult
48 from .interfaces import _AnyExecuteParams
49 from .interfaces import _AnyMultiExecuteParams
50 from .interfaces import _CoreAnyExecuteParams
51 from .interfaces import _CoreMultiExecuteParams
52 from .interfaces import _CoreSingleExecuteParams
53 from .interfaces import _DBAPIAnyExecuteParams
54 from .interfaces import _DBAPISingleExecuteParams
55 from .interfaces import _ExecuteOptions
56 from .interfaces import CompiledCacheType
57 from .interfaces import CoreExecuteOptionsParameter
58 from .interfaces import Dialect
59 from .interfaces import SchemaTranslateMapType
60 from .reflection import Inspector # noqa
61 from .url import URL
62 from ..event import dispatcher
63 from ..log import _EchoFlagType
64 from ..pool import _ConnectionFairy
65 from ..pool import Pool
66 from ..pool import PoolProxiedConnection
67 from ..sql import Executable
68 from ..sql._typing import _InfoType
69 from ..sql.compiler import Compiled
70 from ..sql.ddl import ExecutableDDLElement
71 from ..sql.ddl import InvokeDDLBase
72 from ..sql.functions import FunctionElement
73 from ..sql.schema import DefaultGenerator
74 from ..sql.schema import HasSchemaAttr
75 from ..sql.schema import SchemaVisitable
76 from ..sql.selectable import TypedReturnsRows
77
78
79_T = TypeVar("_T", bound=Any)
80_EMPTY_EXECUTION_OPTS: _ExecuteOptions = util.EMPTY_DICT
81NO_OPTIONS: Mapping[str, Any] = util.EMPTY_DICT
82
83
84class Connection(ConnectionEventsTarget, inspection.Inspectable["Inspector"]):
85 """Provides high-level functionality for a wrapped DB-API connection.
86
87 The :class:`_engine.Connection` object is procured by calling the
88 :meth:`_engine.Engine.connect` method of the :class:`_engine.Engine`
89 object, and provides services for execution of SQL statements as well
90 as transaction control.
91
92 The Connection object is **not** thread-safe. While a Connection can be
93 shared among threads using properly synchronized access, it is still
94 possible that the underlying DBAPI connection may not support shared
95 access between threads. Check the DBAPI documentation for details.
96
97 The Connection object represents a single DBAPI connection checked out
98 from the connection pool. In this state, the connection pool has no
99 affect upon the connection, including its expiration or timeout state.
100 For the connection pool to properly manage connections, connections
101 should be returned to the connection pool (i.e. ``connection.close()``)
102 whenever the connection is not in use.
103
104 .. index::
105 single: thread safety; Connection
106
107 """
108
109 dialect: Dialect
110 dispatch: dispatcher[ConnectionEventsTarget]
111
112 _sqla_logger_namespace = "sqlalchemy.engine.Connection"
113
114 # used by sqlalchemy.engine.util.TransactionalContext
115 _trans_context_manager: Optional[TransactionalContext] = None
116
117 # legacy as of 2.0, should be eventually deprecated and
118 # removed. was used in the "pre_ping" recipe that's been in the docs
119 # a long time
120 should_close_with_result = False
121
122 _dbapi_connection: Optional[PoolProxiedConnection]
123
124 _execution_options: _ExecuteOptions
125
126 _transaction: Optional[RootTransaction]
127 _nested_transaction: Optional[NestedTransaction]
128
129 def __init__(
130 self,
131 engine: Engine,
132 connection: Optional[PoolProxiedConnection] = None,
133 _has_events: Optional[bool] = None,
134 _allow_revalidate: bool = True,
135 _allow_autobegin: bool = True,
136 ):
137 """Construct a new Connection."""
138 self.engine = engine
139 self.dialect = dialect = engine.dialect
140
141 if connection is None:
142 try:
143 self._dbapi_connection = engine.raw_connection()
144 except dialect.loaded_dbapi.Error as err:
145 Connection._handle_dbapi_exception_noconnection(
146 err, dialect, engine
147 )
148 raise
149 else:
150 self._dbapi_connection = connection
151
152 self._transaction = self._nested_transaction = None
153 self.__savepoint_seq = 0
154 self.__in_begin = False
155
156 self.__can_reconnect = _allow_revalidate
157 self._allow_autobegin = _allow_autobegin
158 self._echo = self.engine._should_log_info()
159
160 if _has_events is None:
161 # if _has_events is sent explicitly as False,
162 # then don't join the dispatch of the engine; we don't
163 # want to handle any of the engine's events in that case.
164 self.dispatch = self.dispatch._join(engine.dispatch)
165 self._has_events = _has_events or (
166 _has_events is None and engine._has_events
167 )
168
169 self._execution_options = engine._execution_options
170
171 if self._has_events or self.engine._has_events:
172 self.dispatch.engine_connect(self)
173
174 # this can be assigned differently via
175 # characteristics.LoggingTokenCharacteristic
176 _message_formatter: Any = None
177
178 def _log_info(self, message: str, *arg: Any, **kw: Any) -> None:
179 fmt = self._message_formatter
180
181 if fmt:
182 message = fmt(message)
183
184 if log.STACKLEVEL:
185 kw["stacklevel"] = 1 + log.STACKLEVEL_OFFSET
186
187 self.engine.logger.info(message, *arg, **kw)
188
189 def _log_debug(self, message: str, *arg: Any, **kw: Any) -> None:
190 fmt = self._message_formatter
191
192 if fmt:
193 message = fmt(message)
194
195 if log.STACKLEVEL:
196 kw["stacklevel"] = 1 + log.STACKLEVEL_OFFSET
197
198 self.engine.logger.debug(message, *arg, **kw)
199
200 @property
201 def _schema_translate_map(self) -> Optional[SchemaTranslateMapType]:
202 schema_translate_map: Optional[SchemaTranslateMapType] = (
203 self._execution_options.get("schema_translate_map", None)
204 )
205
206 return schema_translate_map
207
208 def schema_for_object(self, obj: HasSchemaAttr) -> Optional[str]:
209 """Return the schema name for the given schema item taking into
210 account current schema translate map.
211
212 """
213
214 name = obj.schema
215 schema_translate_map: Optional[SchemaTranslateMapType] = (
216 self._execution_options.get("schema_translate_map", None)
217 )
218
219 if (
220 schema_translate_map
221 and name in schema_translate_map
222 and obj._use_schema_map
223 ):
224 return schema_translate_map[name]
225 else:
226 return name
227
228 def __enter__(self) -> Connection:
229 return self
230
231 def __exit__(self, type_: Any, value: Any, traceback: Any) -> None:
232 self.close()
233
234 @overload
235 def execution_options(
236 self,
237 *,
238 compiled_cache: Optional[CompiledCacheType] = ...,
239 logging_token: str = ...,
240 isolation_level: IsolationLevel = ...,
241 no_parameters: bool = False,
242 stream_results: bool = False,
243 max_row_buffer: int = ...,
244 yield_per: int = ...,
245 insertmanyvalues_page_size: int = ...,
246 schema_translate_map: Optional[SchemaTranslateMapType] = ...,
247 preserve_rowcount: bool = False,
248 **opt: Any,
249 ) -> Connection: ...
250
251 @overload
252 def execution_options(self, **opt: Any) -> Connection: ...
253
254 def execution_options(self, **opt: Any) -> Connection:
255 r"""Set non-SQL options for the connection which take effect
256 during execution.
257
258 This method modifies this :class:`_engine.Connection` **in-place**;
259 the return value is the same :class:`_engine.Connection` object
260 upon which the method is called. Note that this is in contrast
261 to the behavior of the ``execution_options`` methods on other
262 objects such as :meth:`_engine.Engine.execution_options` and
263 :meth:`_sql.Executable.execution_options`. The rationale is that many
264 such execution options necessarily modify the state of the base
265 DBAPI connection in any case so there is no feasible means of
266 keeping the effect of such an option localized to a "sub" connection.
267
268 .. versionchanged:: 2.0 The :meth:`_engine.Connection.execution_options`
269 method, in contrast to other objects with this method, modifies
270 the connection in-place without creating copy of it.
271
272 As discussed elsewhere, the :meth:`_engine.Connection.execution_options`
273 method accepts any arbitrary parameters including user defined names.
274 All parameters given are consumable in a number of ways including
275 by using the :meth:`_engine.Connection.get_execution_options` method.
276 See the examples at :meth:`_sql.Executable.execution_options`
277 and :meth:`_engine.Engine.execution_options`.
278
279 The keywords that are currently recognized by SQLAlchemy itself
280 include all those listed under :meth:`.Executable.execution_options`,
281 as well as others that are specific to :class:`_engine.Connection`.
282
283 :param compiled_cache: Available on: :class:`_engine.Connection`,
284 :class:`_engine.Engine`.
285
286 A dictionary where :class:`.Compiled` objects
287 will be cached when the :class:`_engine.Connection`
288 compiles a clause
289 expression into a :class:`.Compiled` object. This dictionary will
290 supersede the statement cache that may be configured on the
291 :class:`_engine.Engine` itself. If set to None, caching
292 is disabled, even if the engine has a configured cache size.
293
294 Note that the ORM makes use of its own "compiled" caches for
295 some operations, including flush operations. The caching
296 used by the ORM internally supersedes a cache dictionary
297 specified here.
298
299 :param logging_token: Available on: :class:`_engine.Connection`,
300 :class:`_engine.Engine`, :class:`_sql.Executable`.
301
302 Adds the specified string token surrounded by brackets in log
303 messages logged by the connection, i.e. the logging that's enabled
304 either via the :paramref:`_sa.create_engine.echo` flag or via the
305 ``logging.getLogger("sqlalchemy.engine")`` logger. This allows a
306 per-connection or per-sub-engine token to be available which is
307 useful for debugging concurrent connection scenarios.
308
309 .. versionadded:: 1.4.0b2
310
311 .. seealso::
312
313 :ref:`dbengine_logging_tokens` - usage example
314
315 :paramref:`_sa.create_engine.logging_name` - adds a name to the
316 name used by the Python logger object itself.
317
318 :param isolation_level: Available on: :class:`_engine.Connection`,
319 :class:`_engine.Engine`.
320
321 Set the transaction isolation level for the lifespan of this
322 :class:`_engine.Connection` object.
323 Valid values include those string
324 values accepted by the :paramref:`_sa.create_engine.isolation_level`
325 parameter passed to :func:`_sa.create_engine`. These levels are
326 semi-database specific; see individual dialect documentation for
327 valid levels.
328
329 The isolation level option applies the isolation level by emitting
330 statements on the DBAPI connection, and **necessarily affects the
331 original Connection object overall**. The isolation level will remain
332 at the given setting until explicitly changed, or when the DBAPI
333 connection itself is :term:`released` to the connection pool, i.e. the
334 :meth:`_engine.Connection.close` method is called, at which time an
335 event handler will emit additional statements on the DBAPI connection
336 in order to revert the isolation level change.
337
338 .. note:: The ``isolation_level`` execution option may only be
339 established before the :meth:`_engine.Connection.begin` method is
340 called, as well as before any SQL statements are emitted which
341 would otherwise trigger "autobegin", or directly after a call to
342 :meth:`_engine.Connection.commit` or
343 :meth:`_engine.Connection.rollback`. A database cannot change the
344 isolation level on a transaction in progress.
345
346 .. note:: The ``isolation_level`` execution option is implicitly
347 reset if the :class:`_engine.Connection` is invalidated, e.g. via
348 the :meth:`_engine.Connection.invalidate` method, or if a
349 disconnection error occurs. The new connection produced after the
350 invalidation will **not** have the selected isolation level
351 re-applied to it automatically.
352
353 .. seealso::
354
355 :ref:`dbapi_autocommit`
356
357 :meth:`_engine.Connection.get_isolation_level`
358 - view current actual level
359
360 :param no_parameters: Available on: :class:`_engine.Connection`,
361 :class:`_sql.Executable`.
362
363 When ``True``, if the final parameter
364 list or dictionary is totally empty, will invoke the
365 statement on the cursor as ``cursor.execute(statement)``,
366 not passing the parameter collection at all.
367 Some DBAPIs such as psycopg2 and mysql-python consider
368 percent signs as significant only when parameters are
369 present; this option allows code to generate SQL
370 containing percent signs (and possibly other characters)
371 that is neutral regarding whether it's executed by the DBAPI
372 or piped into a script that's later invoked by
373 command line tools.
374
375 :param stream_results: Available on: :class:`_engine.Connection`,
376 :class:`_sql.Executable`.
377
378 Indicate to the dialect that results should be "streamed" and not
379 pre-buffered, if possible. For backends such as PostgreSQL, MySQL
380 and MariaDB, this indicates the use of a "server side cursor" as
381 opposed to a client side cursor. Other backends such as that of
382 Oracle Database may already use server side cursors by default.
383
384 The usage of
385 :paramref:`_engine.Connection.execution_options.stream_results` is
386 usually combined with setting a fixed number of rows to to be fetched
387 in batches, to allow for efficient iteration of database rows while
388 at the same time not loading all result rows into memory at once;
389 this can be configured on a :class:`_engine.Result` object using the
390 :meth:`_engine.Result.yield_per` method, after execution has
391 returned a new :class:`_engine.Result`. If
392 :meth:`_engine.Result.yield_per` is not used,
393 the :paramref:`_engine.Connection.execution_options.stream_results`
394 mode of operation will instead use a dynamically sized buffer
395 which buffers sets of rows at a time, growing on each batch
396 based on a fixed growth size up until a limit which may
397 be configured using the
398 :paramref:`_engine.Connection.execution_options.max_row_buffer`
399 parameter.
400
401 When using the ORM to fetch ORM mapped objects from a result,
402 :meth:`_engine.Result.yield_per` should always be used with
403 :paramref:`_engine.Connection.execution_options.stream_results`,
404 so that the ORM does not fetch all rows into new ORM objects at once.
405
406 For typical use, the
407 :paramref:`_engine.Connection.execution_options.yield_per` execution
408 option should be preferred, which sets up both
409 :paramref:`_engine.Connection.execution_options.stream_results` and
410 :meth:`_engine.Result.yield_per` at once. This option is supported
411 both at a core level by :class:`_engine.Connection` as well as by the
412 ORM :class:`_engine.Session`; the latter is described at
413 :ref:`orm_queryguide_yield_per`.
414
415 .. seealso::
416
417 :ref:`engine_stream_results` - background on
418 :paramref:`_engine.Connection.execution_options.stream_results`
419
420 :paramref:`_engine.Connection.execution_options.max_row_buffer`
421
422 :paramref:`_engine.Connection.execution_options.yield_per`
423
424 :ref:`orm_queryguide_yield_per` - in the :ref:`queryguide_toplevel`
425 describing the ORM version of ``yield_per``
426
427 :param max_row_buffer: Available on: :class:`_engine.Connection`,
428 :class:`_sql.Executable`. Sets a maximum
429 buffer size to use when the
430 :paramref:`_engine.Connection.execution_options.stream_results`
431 execution option is used on a backend that supports server side
432 cursors. The default value if not specified is 1000.
433
434 .. seealso::
435
436 :paramref:`_engine.Connection.execution_options.stream_results`
437
438 :ref:`engine_stream_results`
439
440
441 :param yield_per: Available on: :class:`_engine.Connection`,
442 :class:`_sql.Executable`. Integer value applied which will
443 set the :paramref:`_engine.Connection.execution_options.stream_results`
444 execution option and invoke :meth:`_engine.Result.yield_per`
445 automatically at once. Allows equivalent functionality as
446 is present when using this parameter with the ORM.
447
448 .. versionadded:: 1.4.40
449
450 .. seealso::
451
452 :ref:`engine_stream_results` - background and examples
453 on using server side cursors with Core.
454
455 :ref:`orm_queryguide_yield_per` - in the :ref:`queryguide_toplevel`
456 describing the ORM version of ``yield_per``
457
458 :param insertmanyvalues_page_size: Available on: :class:`_engine.Connection`,
459 :class:`_engine.Engine`. Number of rows to format into an
460 INSERT statement when the statement uses "insertmanyvalues" mode,
461 which is a paged form of bulk insert that is used for many backends
462 when using :term:`executemany` execution typically in conjunction
463 with RETURNING. Defaults to 1000. May also be modified on a
464 per-engine basis using the
465 :paramref:`_sa.create_engine.insertmanyvalues_page_size` parameter.
466
467 .. versionadded:: 2.0
468
469 .. seealso::
470
471 :ref:`engine_insertmanyvalues`
472
473 :param schema_translate_map: Available on: :class:`_engine.Connection`,
474 :class:`_engine.Engine`, :class:`_sql.Executable`.
475
476 A dictionary mapping schema names to schema names, that will be
477 applied to the :paramref:`_schema.Table.schema` element of each
478 :class:`_schema.Table`
479 encountered when SQL or DDL expression elements
480 are compiled into strings; the resulting schema name will be
481 converted based on presence in the map of the original name.
482
483 .. seealso::
484
485 :ref:`schema_translating`
486
487 :param preserve_rowcount: Boolean; when True, the ``cursor.rowcount``
488 attribute will be unconditionally memoized within the result and
489 made available via the :attr:`.CursorResult.rowcount` attribute.
490 Normally, this attribute is only preserved for UPDATE and DELETE
491 statements. Using this option, the DBAPIs rowcount value can
492 be accessed for other kinds of statements such as INSERT and SELECT,
493 to the degree that the DBAPI supports these statements. See
494 :attr:`.CursorResult.rowcount` for notes regarding the behavior
495 of this attribute.
496
497 .. versionadded:: 2.0.28
498
499 .. seealso::
500
501 :meth:`_engine.Engine.execution_options`
502
503 :meth:`.Executable.execution_options`
504
505 :meth:`_engine.Connection.get_execution_options`
506
507 :ref:`orm_queryguide_execution_options` - documentation on all
508 ORM-specific execution options
509
510 """ # noqa
511 if self._has_events or self.engine._has_events:
512 self.dispatch.set_connection_execution_options(self, opt)
513 self._execution_options = self._execution_options.union(opt)
514 self.dialect.set_connection_execution_options(self, opt)
515 return self
516
517 def get_execution_options(self) -> _ExecuteOptions:
518 """Get the non-SQL options which will take effect during execution.
519
520 .. versionadded:: 1.3
521
522 .. seealso::
523
524 :meth:`_engine.Connection.execution_options`
525 """
526 return self._execution_options
527
528 @property
529 def _still_open_and_dbapi_connection_is_valid(self) -> bool:
530 pool_proxied_connection = self._dbapi_connection
531 return (
532 pool_proxied_connection is not None
533 and pool_proxied_connection.is_valid
534 )
535
536 @property
537 def closed(self) -> bool:
538 """Return True if this connection is closed."""
539
540 return self._dbapi_connection is None and not self.__can_reconnect
541
542 @property
543 def invalidated(self) -> bool:
544 """Return True if this connection was invalidated.
545
546 This does not indicate whether or not the connection was
547 invalidated at the pool level, however
548
549 """
550
551 # prior to 1.4, "invalid" was stored as a state independent of
552 # "closed", meaning an invalidated connection could be "closed",
553 # the _dbapi_connection would be None and closed=True, yet the
554 # "invalid" flag would stay True. This meant that there were
555 # three separate states (open/valid, closed/valid, closed/invalid)
556 # when there is really no reason for that; a connection that's
557 # "closed" does not need to be "invalid". So the state is now
558 # represented by the two facts alone.
559
560 pool_proxied_connection = self._dbapi_connection
561 return pool_proxied_connection is None and self.__can_reconnect
562
563 @property
564 def connection(self) -> PoolProxiedConnection:
565 """The underlying DB-API connection managed by this Connection.
566
567 This is a SQLAlchemy connection-pool proxied connection
568 which then has the attribute
569 :attr:`_pool._ConnectionFairy.dbapi_connection` that refers to the
570 actual driver connection.
571
572 .. seealso::
573
574
575 :ref:`dbapi_connections`
576
577 """
578
579 if self._dbapi_connection is None:
580 try:
581 return self._revalidate_connection()
582 except (exc.PendingRollbackError, exc.ResourceClosedError):
583 raise
584 except BaseException as e:
585 self._handle_dbapi_exception(e, None, None, None, None)
586 else:
587 return self._dbapi_connection
588
589 def get_isolation_level(self) -> IsolationLevel:
590 """Return the current **actual** isolation level that's present on
591 the database within the scope of this connection.
592
593 This attribute will perform a live SQL operation against the database
594 in order to procure the current isolation level, so the value returned
595 is the actual level on the underlying DBAPI connection regardless of
596 how this state was set. This will be one of the four actual isolation
597 modes ``READ UNCOMMITTED``, ``READ COMMITTED``, ``REPEATABLE READ``,
598 ``SERIALIZABLE``. It will **not** include the ``AUTOCOMMIT`` isolation
599 level setting. Third party dialects may also feature additional
600 isolation level settings.
601
602 .. note:: This method **will not report** on the ``AUTOCOMMIT``
603 isolation level, which is a separate :term:`dbapi` setting that's
604 independent of **actual** isolation level. When ``AUTOCOMMIT`` is
605 in use, the database connection still has a "traditional" isolation
606 mode in effect, that is typically one of the four values
607 ``READ UNCOMMITTED``, ``READ COMMITTED``, ``REPEATABLE READ``,
608 ``SERIALIZABLE``.
609
610 Compare to the :attr:`_engine.Connection.default_isolation_level`
611 accessor which returns the isolation level that is present on the
612 database at initial connection time.
613
614 .. seealso::
615
616 :attr:`_engine.Connection.default_isolation_level`
617 - view default level
618
619 :paramref:`_sa.create_engine.isolation_level`
620 - set per :class:`_engine.Engine` isolation level
621
622 :paramref:`.Connection.execution_options.isolation_level`
623 - set per :class:`_engine.Connection` isolation level
624
625 """
626 dbapi_connection = self.connection.dbapi_connection
627 assert dbapi_connection is not None
628 try:
629 return self.dialect.get_isolation_level(dbapi_connection)
630 except BaseException as e:
631 self._handle_dbapi_exception(e, None, None, None, None)
632
633 @property
634 def default_isolation_level(self) -> Optional[IsolationLevel]:
635 """The initial-connection time isolation level associated with the
636 :class:`_engine.Dialect` in use.
637
638 This value is independent of the
639 :paramref:`.Connection.execution_options.isolation_level` and
640 :paramref:`.Engine.execution_options.isolation_level` execution
641 options, and is determined by the :class:`_engine.Dialect` when the
642 first connection is created, by performing a SQL query against the
643 database for the current isolation level before any additional commands
644 have been emitted.
645
646 Calling this accessor does not invoke any new SQL queries.
647
648 .. seealso::
649
650 :meth:`_engine.Connection.get_isolation_level`
651 - view current actual isolation level
652
653 :paramref:`_sa.create_engine.isolation_level`
654 - set per :class:`_engine.Engine` isolation level
655
656 :paramref:`.Connection.execution_options.isolation_level`
657 - set per :class:`_engine.Connection` isolation level
658
659 """
660 return self.dialect.default_isolation_level
661
662 def _invalid_transaction(self) -> NoReturn:
663 raise exc.PendingRollbackError(
664 "Can't reconnect until invalid %stransaction is rolled "
665 "back. Please rollback() fully before proceeding"
666 % ("savepoint " if self._nested_transaction is not None else ""),
667 code="8s2b",
668 )
669
670 def _revalidate_connection(self) -> PoolProxiedConnection:
671 if self.__can_reconnect and self.invalidated:
672 if self._transaction is not None:
673 self._invalid_transaction()
674 self._dbapi_connection = self.engine.raw_connection()
675 return self._dbapi_connection
676 raise exc.ResourceClosedError("This Connection is closed")
677
678 @property
679 def info(self) -> _InfoType:
680 """Info dictionary associated with the underlying DBAPI connection
681 referred to by this :class:`_engine.Connection`, allowing user-defined
682 data to be associated with the connection.
683
684 The data here will follow along with the DBAPI connection including
685 after it is returned to the connection pool and used again
686 in subsequent instances of :class:`_engine.Connection`.
687
688 """
689
690 return self.connection.info
691
692 def invalidate(self, exception: Optional[BaseException] = None) -> None:
693 """Invalidate the underlying DBAPI connection associated with
694 this :class:`_engine.Connection`.
695
696 An attempt will be made to close the underlying DBAPI connection
697 immediately; however if this operation fails, the error is logged
698 but not raised. The connection is then discarded whether or not
699 close() succeeded.
700
701 Upon the next use (where "use" typically means using the
702 :meth:`_engine.Connection.execute` method or similar),
703 this :class:`_engine.Connection` will attempt to
704 procure a new DBAPI connection using the services of the
705 :class:`_pool.Pool` as a source of connectivity (e.g.
706 a "reconnection").
707
708 If a transaction was in progress (e.g. the
709 :meth:`_engine.Connection.begin` method has been called) when
710 :meth:`_engine.Connection.invalidate` method is called, at the DBAPI
711 level all state associated with this transaction is lost, as
712 the DBAPI connection is closed. The :class:`_engine.Connection`
713 will not allow a reconnection to proceed until the
714 :class:`.Transaction` object is ended, by calling the
715 :meth:`.Transaction.rollback` method; until that point, any attempt at
716 continuing to use the :class:`_engine.Connection` will raise an
717 :class:`~sqlalchemy.exc.InvalidRequestError`.
718 This is to prevent applications from accidentally
719 continuing an ongoing transactional operations despite the
720 fact that the transaction has been lost due to an
721 invalidation.
722
723 The :meth:`_engine.Connection.invalidate` method,
724 just like auto-invalidation,
725 will at the connection pool level invoke the
726 :meth:`_events.PoolEvents.invalidate` event.
727
728 :param exception: an optional ``Exception`` instance that's the
729 reason for the invalidation. is passed along to event handlers
730 and logging functions.
731
732 .. seealso::
733
734 :ref:`pool_connection_invalidation`
735
736 """
737
738 if self.invalidated:
739 return
740
741 if self.closed:
742 raise exc.ResourceClosedError("This Connection is closed")
743
744 if self._still_open_and_dbapi_connection_is_valid:
745 pool_proxied_connection = self._dbapi_connection
746 assert pool_proxied_connection is not None
747 pool_proxied_connection.invalidate(exception)
748
749 self._dbapi_connection = None
750
751 def detach(self) -> None:
752 """Detach the underlying DB-API connection from its connection pool.
753
754 E.g.::
755
756 with engine.connect() as conn:
757 conn.detach()
758 conn.execute(text("SET search_path TO schema1, schema2"))
759
760 # work with connection
761
762 # connection is fully closed (since we used "with:", can
763 # also call .close())
764
765 This :class:`_engine.Connection` instance will remain usable.
766 When closed
767 (or exited from a context manager context as above),
768 the DB-API connection will be literally closed and not
769 returned to its originating pool.
770
771 This method can be used to insulate the rest of an application
772 from a modified state on a connection (such as a transaction
773 isolation level or similar).
774
775 """
776
777 if self.closed:
778 raise exc.ResourceClosedError("This Connection is closed")
779
780 pool_proxied_connection = self._dbapi_connection
781 if pool_proxied_connection is None:
782 raise exc.InvalidRequestError(
783 "Can't detach an invalidated Connection"
784 )
785 pool_proxied_connection.detach()
786
787 def _autobegin(self) -> None:
788 if self._allow_autobegin and not self.__in_begin:
789 self.begin()
790
791 def begin(self) -> RootTransaction:
792 """Begin a transaction prior to autobegin occurring.
793
794 E.g.::
795
796 with engine.connect() as conn:
797 with conn.begin() as trans:
798 conn.execute(table.insert(), {"username": "sandy"})
799
800 The returned object is an instance of :class:`_engine.RootTransaction`.
801 This object represents the "scope" of the transaction,
802 which completes when either the :meth:`_engine.Transaction.rollback`
803 or :meth:`_engine.Transaction.commit` method is called; the object
804 also works as a context manager as illustrated above.
805
806 The :meth:`_engine.Connection.begin` method begins a
807 transaction that normally will be begun in any case when the connection
808 is first used to execute a statement. The reason this method might be
809 used would be to invoke the :meth:`_events.ConnectionEvents.begin`
810 event at a specific time, or to organize code within the scope of a
811 connection checkout in terms of context managed blocks, such as::
812
813 with engine.connect() as conn:
814 with conn.begin():
815 conn.execute(...)
816 conn.execute(...)
817
818 with conn.begin():
819 conn.execute(...)
820 conn.execute(...)
821
822 The above code is not fundamentally any different in its behavior than
823 the following code which does not use
824 :meth:`_engine.Connection.begin`; the below style is known
825 as "commit as you go" style::
826
827 with engine.connect() as conn:
828 conn.execute(...)
829 conn.execute(...)
830 conn.commit()
831
832 conn.execute(...)
833 conn.execute(...)
834 conn.commit()
835
836 From a database point of view, the :meth:`_engine.Connection.begin`
837 method does not emit any SQL or change the state of the underlying
838 DBAPI connection in any way; the Python DBAPI does not have any
839 concept of explicit transaction begin.
840
841 .. seealso::
842
843 :ref:`tutorial_working_with_transactions` - in the
844 :ref:`unified_tutorial`
845
846 :meth:`_engine.Connection.begin_nested` - use a SAVEPOINT
847
848 :meth:`_engine.Connection.begin_twophase` -
849 use a two phase /XID transaction
850
851 :meth:`_engine.Engine.begin` - context manager available from
852 :class:`_engine.Engine`
853
854 """
855 if self._transaction is None:
856 self._transaction = RootTransaction(self)
857 return self._transaction
858 else:
859 raise exc.InvalidRequestError(
860 "This connection has already initialized a SQLAlchemy "
861 "Transaction() object via begin() or autobegin; can't "
862 "call begin() here unless rollback() or commit() "
863 "is called first."
864 )
865
866 def begin_nested(self) -> NestedTransaction:
867 """Begin a nested transaction (i.e. SAVEPOINT) and return a transaction
868 handle that controls the scope of the SAVEPOINT.
869
870 E.g.::
871
872 with engine.begin() as connection:
873 with connection.begin_nested():
874 connection.execute(table.insert(), {"username": "sandy"})
875
876 The returned object is an instance of
877 :class:`_engine.NestedTransaction`, which includes transactional
878 methods :meth:`_engine.NestedTransaction.commit` and
879 :meth:`_engine.NestedTransaction.rollback`; for a nested transaction,
880 these methods correspond to the operations "RELEASE SAVEPOINT <name>"
881 and "ROLLBACK TO SAVEPOINT <name>". The name of the savepoint is local
882 to the :class:`_engine.NestedTransaction` object and is generated
883 automatically. Like any other :class:`_engine.Transaction`, the
884 :class:`_engine.NestedTransaction` may be used as a context manager as
885 illustrated above which will "release" or "rollback" corresponding to
886 if the operation within the block were successful or raised an
887 exception.
888
889 Nested transactions require SAVEPOINT support in the underlying
890 database, else the behavior is undefined. SAVEPOINT is commonly used to
891 run operations within a transaction that may fail, while continuing the
892 outer transaction. E.g.::
893
894 from sqlalchemy import exc
895
896 with engine.begin() as connection:
897 trans = connection.begin_nested()
898 try:
899 connection.execute(table.insert(), {"username": "sandy"})
900 trans.commit()
901 except exc.IntegrityError: # catch for duplicate username
902 trans.rollback() # rollback to savepoint
903
904 # outer transaction continues
905 connection.execute(...)
906
907 If :meth:`_engine.Connection.begin_nested` is called without first
908 calling :meth:`_engine.Connection.begin` or
909 :meth:`_engine.Engine.begin`, the :class:`_engine.Connection` object
910 will "autobegin" the outer transaction first. This outer transaction
911 may be committed using "commit-as-you-go" style, e.g.::
912
913 with engine.connect() as connection: # begin() wasn't called
914
915 with connection.begin_nested(): # will auto-"begin()" first
916 connection.execute(...)
917 # savepoint is released
918
919 connection.execute(...)
920
921 # explicitly commit outer transaction
922 connection.commit()
923
924 # can continue working with connection here
925
926 .. versionchanged:: 2.0
927
928 :meth:`_engine.Connection.begin_nested` will now participate
929 in the connection "autobegin" behavior that is new as of
930 2.0 / "future" style connections in 1.4.
931
932 .. seealso::
933
934 :meth:`_engine.Connection.begin`
935
936 :ref:`session_begin_nested` - ORM support for SAVEPOINT
937
938 """
939 if self._transaction is None:
940 self._autobegin()
941
942 return NestedTransaction(self)
943
944 def begin_twophase(self, xid: Optional[Any] = None) -> TwoPhaseTransaction:
945 """Begin a two-phase or XA transaction and return a transaction
946 handle.
947
948 The returned object is an instance of :class:`.TwoPhaseTransaction`,
949 which in addition to the methods provided by
950 :class:`.Transaction`, also provides a
951 :meth:`~.TwoPhaseTransaction.prepare` method.
952
953 :param xid: the two phase transaction id. If not supplied, a
954 random id will be generated.
955
956 .. seealso::
957
958 :meth:`_engine.Connection.begin`
959
960 :meth:`_engine.Connection.begin_twophase`
961
962 """
963
964 if self._transaction is not None:
965 raise exc.InvalidRequestError(
966 "Cannot start a two phase transaction when a transaction "
967 "is already in progress."
968 )
969 if xid is None:
970 xid = self.engine.dialect.create_xid()
971 return TwoPhaseTransaction(self, xid)
972
973 def commit(self) -> None:
974 """Commit the transaction that is currently in progress.
975
976 This method commits the current transaction if one has been started.
977 If no transaction was started, the method has no effect, assuming
978 the connection is in a non-invalidated state.
979
980 A transaction is begun on a :class:`_engine.Connection` automatically
981 whenever a statement is first executed, or when the
982 :meth:`_engine.Connection.begin` method is called.
983
984 .. note:: The :meth:`_engine.Connection.commit` method only acts upon
985 the primary database transaction that is linked to the
986 :class:`_engine.Connection` object. It does not operate upon a
987 SAVEPOINT that would have been invoked from the
988 :meth:`_engine.Connection.begin_nested` method; for control of a
989 SAVEPOINT, call :meth:`_engine.NestedTransaction.commit` on the
990 :class:`_engine.NestedTransaction` that is returned by the
991 :meth:`_engine.Connection.begin_nested` method itself.
992
993
994 """
995 if self._transaction:
996 self._transaction.commit()
997
998 def rollback(self) -> None:
999 """Roll back the transaction that is currently in progress.
1000
1001 This method rolls back the current transaction if one has been started.
1002 If no transaction was started, the method has no effect. If a
1003 transaction was started and the connection is in an invalidated state,
1004 the transaction is cleared using this method.
1005
1006 A transaction is begun on a :class:`_engine.Connection` automatically
1007 whenever a statement is first executed, or when the
1008 :meth:`_engine.Connection.begin` method is called.
1009
1010 .. note:: The :meth:`_engine.Connection.rollback` method only acts
1011 upon the primary database transaction that is linked to the
1012 :class:`_engine.Connection` object. It does not operate upon a
1013 SAVEPOINT that would have been invoked from the
1014 :meth:`_engine.Connection.begin_nested` method; for control of a
1015 SAVEPOINT, call :meth:`_engine.NestedTransaction.rollback` on the
1016 :class:`_engine.NestedTransaction` that is returned by the
1017 :meth:`_engine.Connection.begin_nested` method itself.
1018
1019
1020 """
1021 if self._transaction:
1022 self._transaction.rollback()
1023
1024 def recover_twophase(self) -> List[Any]:
1025 return self.engine.dialect.do_recover_twophase(self)
1026
1027 def rollback_prepared(self, xid: Any, recover: bool = False) -> None:
1028 self.engine.dialect.do_rollback_twophase(self, xid, recover=recover)
1029
1030 def commit_prepared(self, xid: Any, recover: bool = False) -> None:
1031 self.engine.dialect.do_commit_twophase(self, xid, recover=recover)
1032
1033 def in_transaction(self) -> bool:
1034 """Return True if a transaction is in progress."""
1035 return self._transaction is not None and self._transaction.is_active
1036
1037 def in_nested_transaction(self) -> bool:
1038 """Return True if a transaction is in progress."""
1039 return (
1040 self._nested_transaction is not None
1041 and self._nested_transaction.is_active
1042 )
1043
1044 def _is_autocommit_isolation(self) -> bool:
1045 opt_iso = self._execution_options.get("isolation_level", None)
1046 return bool(
1047 opt_iso == "AUTOCOMMIT"
1048 or (
1049 opt_iso is None
1050 and self.engine.dialect._on_connect_isolation_level
1051 == "AUTOCOMMIT"
1052 )
1053 )
1054
1055 def _get_required_transaction(self) -> RootTransaction:
1056 trans = self._transaction
1057 if trans is None:
1058 raise exc.InvalidRequestError("connection is not in a transaction")
1059 return trans
1060
1061 def _get_required_nested_transaction(self) -> NestedTransaction:
1062 trans = self._nested_transaction
1063 if trans is None:
1064 raise exc.InvalidRequestError(
1065 "connection is not in a nested transaction"
1066 )
1067 return trans
1068
1069 def get_transaction(self) -> Optional[RootTransaction]:
1070 """Return the current root transaction in progress, if any.
1071
1072 .. versionadded:: 1.4
1073
1074 """
1075
1076 return self._transaction
1077
1078 def get_nested_transaction(self) -> Optional[NestedTransaction]:
1079 """Return the current nested transaction in progress, if any.
1080
1081 .. versionadded:: 1.4
1082
1083 """
1084 return self._nested_transaction
1085
1086 def _begin_impl(self, transaction: RootTransaction) -> None:
1087 if self._echo:
1088 if self._is_autocommit_isolation():
1089 self._log_info(
1090 "BEGIN (implicit; DBAPI should not BEGIN due to "
1091 "autocommit mode)"
1092 )
1093 else:
1094 self._log_info("BEGIN (implicit)")
1095
1096 self.__in_begin = True
1097
1098 if self._has_events or self.engine._has_events:
1099 self.dispatch.begin(self)
1100
1101 try:
1102 self.engine.dialect.do_begin(self.connection)
1103 except BaseException as e:
1104 self._handle_dbapi_exception(e, None, None, None, None)
1105 finally:
1106 self.__in_begin = False
1107
1108 def _rollback_impl(self) -> None:
1109 if self._has_events or self.engine._has_events:
1110 self.dispatch.rollback(self)
1111
1112 if self._still_open_and_dbapi_connection_is_valid:
1113 if self._echo:
1114 if self._is_autocommit_isolation():
1115 if self.dialect.skip_autocommit_rollback:
1116 self._log_info(
1117 "ROLLBACK will be skipped by "
1118 "skip_autocommit_rollback"
1119 )
1120 else:
1121 self._log_info(
1122 "ROLLBACK using DBAPI connection.rollback(); "
1123 "set skip_autocommit_rollback to prevent fully"
1124 )
1125 else:
1126 self._log_info("ROLLBACK")
1127 try:
1128 self.engine.dialect.do_rollback(self.connection)
1129 except BaseException as e:
1130 self._handle_dbapi_exception(e, None, None, None, None)
1131
1132 def _commit_impl(self) -> None:
1133 if self._has_events or self.engine._has_events:
1134 self.dispatch.commit(self)
1135
1136 if self._echo:
1137 if self._is_autocommit_isolation():
1138 self._log_info(
1139 "COMMIT using DBAPI connection.commit(), "
1140 "has no effect due to autocommit mode"
1141 )
1142 else:
1143 self._log_info("COMMIT")
1144 try:
1145 self.engine.dialect.do_commit(self.connection)
1146 except BaseException as e:
1147 self._handle_dbapi_exception(e, None, None, None, None)
1148
1149 def _savepoint_impl(self, name: Optional[str] = None) -> str:
1150 if self._has_events or self.engine._has_events:
1151 self.dispatch.savepoint(self, name)
1152
1153 if name is None:
1154 self.__savepoint_seq += 1
1155 name = "sa_savepoint_%s" % self.__savepoint_seq
1156 self.engine.dialect.do_savepoint(self, name)
1157 return name
1158
1159 def _rollback_to_savepoint_impl(self, name: str) -> None:
1160 if self._has_events or self.engine._has_events:
1161 self.dispatch.rollback_savepoint(self, name, None)
1162
1163 if self._still_open_and_dbapi_connection_is_valid:
1164 self.engine.dialect.do_rollback_to_savepoint(self, name)
1165
1166 def _release_savepoint_impl(self, name: str) -> None:
1167 if self._has_events or self.engine._has_events:
1168 self.dispatch.release_savepoint(self, name, None)
1169
1170 self.engine.dialect.do_release_savepoint(self, name)
1171
1172 def _begin_twophase_impl(self, transaction: TwoPhaseTransaction) -> None:
1173 if self._echo:
1174 self._log_info("BEGIN TWOPHASE (implicit)")
1175 if self._has_events or self.engine._has_events:
1176 self.dispatch.begin_twophase(self, transaction.xid)
1177
1178 self.__in_begin = True
1179 try:
1180 self.engine.dialect.do_begin_twophase(self, transaction.xid)
1181 except BaseException as e:
1182 self._handle_dbapi_exception(e, None, None, None, None)
1183 finally:
1184 self.__in_begin = False
1185
1186 def _prepare_twophase_impl(self, xid: Any) -> None:
1187 if self._has_events or self.engine._has_events:
1188 self.dispatch.prepare_twophase(self, xid)
1189
1190 assert isinstance(self._transaction, TwoPhaseTransaction)
1191 try:
1192 self.engine.dialect.do_prepare_twophase(self, xid)
1193 except BaseException as e:
1194 self._handle_dbapi_exception(e, None, None, None, None)
1195
1196 def _rollback_twophase_impl(self, xid: Any, is_prepared: bool) -> None:
1197 if self._has_events or self.engine._has_events:
1198 self.dispatch.rollback_twophase(self, xid, is_prepared)
1199
1200 if self._still_open_and_dbapi_connection_is_valid:
