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