Coverage for /pythoncovmergedfiles/medio/medio/usr/local/lib/python3.11/site-packages/sqlalchemy/dialects/postgresql/psycopg.py: 53%

Shortcuts on this page

r m x   toggle line displays

j k   next/prev highlighted chunk

0   (zero) top of page

1   (one) first highlighted chunk

343 statements  

1# dialects/postgresql/psycopg.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# mypy: ignore-errors 

8 

9r""" 

10.. dialect:: postgresql+psycopg 

11 :name: psycopg (a.k.a. psycopg 3) 

12 :dbapi: psycopg 

13 :connectstring: postgresql+psycopg://user:password@host:port/dbname[?key=value&key=value...] 

14 :url: https://pypi.org/project/psycopg/ 

15 

16``psycopg`` is the package and module name for version 3 of the ``psycopg`` 

17database driver, formerly known as ``psycopg2``. This driver is different 

18enough from its ``psycopg2`` predecessor that SQLAlchemy supports it 

19via a totally separate dialect; support for ``psycopg2`` is expected to remain 

20for as long as that package continues to function for modern Python versions, 

21and also remains the default dialect for the ``postgresql://`` dialect 

22series. 

23 

24The SQLAlchemy ``psycopg`` dialect provides both a sync and an async 

25implementation under the same dialect name. The proper version is 

26selected depending on how the engine is created: 

27 

28* calling :func:`_sa.create_engine` with ``postgresql+psycopg://...`` will 

29 automatically select the sync version, e.g.:: 

30 

31 from sqlalchemy import create_engine 

32 

33 sync_engine = create_engine( 

34 "postgresql+psycopg://scott:tiger@localhost/test" 

35 ) 

36 

37* calling :func:`_asyncio.create_async_engine` with 

38 ``postgresql+psycopg://...`` will automatically select the async version, 

39 e.g.:: 

40 

41 from sqlalchemy.ext.asyncio import create_async_engine 

42 

43 asyncio_engine = create_async_engine( 

44 "postgresql+psycopg://scott:tiger@localhost/test" 

45 ) 

46 

47The asyncio version of the dialect may also be specified explicitly using the 

48``psycopg_async`` suffix, as:: 

49 

50 from sqlalchemy.ext.asyncio import create_async_engine 

51 

52 asyncio_engine = create_async_engine( 

53 "postgresql+psycopg_async://scott:tiger@localhost/test" 

54 ) 

55 

56.. seealso:: 

57 

58 :ref:`postgresql_psycopg2` - The SQLAlchemy ``psycopg`` 

59 dialect shares most of its behavior with the ``psycopg2`` dialect. 

60 Further documentation is available there. 

61 

62Using psycopg Connection Pooling 

63-------------------------------- 

64 

65The ``psycopg`` driver provides its own connection pool implementation that 

66may be used in place of SQLAlchemy's pooling functionality. 

67This pool implementation provides support for fixed and dynamic pool sizes 

68(including automatic downsizing for unused connections), connection health 

69pre-checks, and support for both synchronous and asynchronous code 

70environments. 

71 

72Here is an example that uses the sync version of the pool, using 

73``psycopg_pool >= 3.3`` that introduces support for ``close_returns=True``:: 

74 

75 import psycopg_pool 

76 from sqlalchemy import create_engine 

77 from sqlalchemy.pool import NullPool 

78 

79 # Create a psycopg_pool connection pool 

80 my_pool = psycopg_pool.ConnectionPool( 

81 conninfo="postgresql://scott:tiger@localhost/test", 

82 close_returns=True, # Return "closed" active connections to the pool 

83 # ... other pool parameters as desired ... 

84 ) 

85 

86 # Create an engine that uses the connection pool to get a connection 

87 engine = create_engine( 

88 url="postgresql+psycopg://", # Only need the dialect now 

89 poolclass=NullPool, # Disable SQLAlchemy's default connection pool 

90 creator=my_pool.getconn, # Use Psycopg 3 connection pool to obtain connections 

91 ) 

92 

93Similarly an the async example:: 

94 

95 import psycopg_pool 

96 from sqlalchemy.ext.asyncio import create_async_engine 

97 from sqlalchemy.pool import NullPool 

98 

99 

100 async def define_engine(): 

101 # Create a psycopg_pool connection pool 

102 my_pool = psycopg_pool.AsyncConnectionPool( 

103 conninfo="postgresql://scott:tiger@localhost/test", 

104 open=False, # See comment below 

105 close_returns=True, # Return "closed" active connections to the pool 

106 # ... other pool parameters as desired ... 

107 ) 

108 

109 # Must explicitly open AsyncConnectionPool outside constructor 

110 # https://www.psycopg.org/psycopg3/docs/api/pool.html#psycopg_pool.AsyncConnectionPool 

111 await my_pool.open() 

112 

113 # Create an engine that uses the connection pool to get a connection 

114 engine = create_async_engine( 

115 url="postgresql+psycopg://", # Only need the dialect now 

116 poolclass=NullPool, # Disable SQLAlchemy's default connection pool 

117 async_creator=my_pool.getconn, # Use Psycopg 3 connection pool to obtain connections 

118 ) 

119 

120 return engine, my_pool 

121 

122The resulting engine may then be used normally. Internally, Psycopg 3 handles 

123connection pooling:: 

124 

125 with engine.connect() as conn: 

126 print(conn.scalar(text("select 42"))) 

127 

128.. seealso:: 

129 

130 `Connection pools <https://www.psycopg.org/psycopg3/docs/advanced/pool.html>`_ - 

131 the Psycopg 3 documentation for ``psycopg_pool.ConnectionPool``. 

132 

133 `Example for older version of psycopg_pool 

134 <https://github.com/sqlalchemy/sqlalchemy/discussions/12522#discussioncomment-13024666>`_ - 

135 An example about using the ``psycopg_pool<3.3`` that did not have the 

136 ``close_returns``` parameter. 

137 

138Using a different Cursor class 

139------------------------------ 

140 

141One of the differences between ``psycopg`` and the older ``psycopg2`` 

142is how bound parameters are handled: ``psycopg2`` would bind them 

143client side, while ``psycopg`` by default will bind them server side. 

144 

145It's possible to configure ``psycopg`` to do client side binding by 

146specifying the ``cursor_factory`` to be ``ClientCursor`` when creating 

147the engine:: 

148 

149 from psycopg import ClientCursor 

150 

151 client_side_engine = create_engine( 

152 "postgresql+psycopg://...", 

153 connect_args={"cursor_factory": ClientCursor}, 

154 ) 

155 

156Similarly when using an async engine the ``AsyncClientCursor`` can be 

157specified:: 

158 

159 from psycopg import AsyncClientCursor 

160 

161 client_side_engine = create_async_engine( 

162 "postgresql+psycopg://...", 

163 connect_args={"cursor_factory": AsyncClientCursor}, 

164 ) 

165 

166.. seealso:: 

167 

168 `Client-side-binding cursors <https://www.psycopg.org/psycopg3/docs/advanced/cursors.html#client-side-binding-cursors>`_ 

169 

170""" # noqa 

171 

172from __future__ import annotations 

173 

174import collections 

175import logging 

176from types import NoneType 

177from typing import cast 

178from typing import TYPE_CHECKING 

179 

180from . import ranges 

181from ._psycopg_common import _PGDialect_common_psycopg 

182from ._psycopg_common import _PGExecutionContext_common_psycopg 

183from .base import INTERVAL 

184from .base import PGCompiler 

185from .base import PGIdentifierPreparer 

186from .base import REGCONFIG 

187from .json import JSON 

188from .json import JSONB 

189from .json import JSONPathType 

190from .types import CITEXT 

191from ... import util 

192from ...connectors.asyncio import AsyncAdapt_dbapi_connection 

193from ...connectors.asyncio import AsyncAdapt_dbapi_cursor 

194from ...connectors.asyncio import AsyncAdapt_dbapi_module 

195from ...connectors.asyncio import AsyncAdapt_dbapi_ss_cursor 

196from ...sql import sqltypes 

197from ...util.concurrency import await_ 

198 

199if TYPE_CHECKING: 

200 from typing import Iterable 

201 

202 from psycopg import AsyncConnection 

203 

204logger = logging.getLogger("sqlalchemy.dialects.postgresql") 

205 

206 

207class _PGString(sqltypes.String): 

208 render_bind_cast = True 

209 

210 

211class _PGREGCONFIG(REGCONFIG): 

212 render_bind_cast = True 

213 

214 

215class _PGJSON(JSON): 

216 def bind_processor(self, dialect): 

217 """psycopg's bind processor is assembled on the type adapter, 

218 but we still need to wrap the value in a psycopg.Json() object""" 

219 return self._make_bind_processor(None, dialect._psycopg_Json) 

220 

221 

222class _PGJSONB(JSONB): 

223 def bind_processor(self, dialect): 

224 """psycopg's bind processor is assembled on the type adapter, 

225 but we still need to wrap the value in a psycopg.Jsonb() object""" 

226 return self._make_bind_processor(None, dialect._psycopg_Jsonb) 

227 

228 

229class _PGJSONIntIndexType(sqltypes.JSON.JSONIntIndexType): 

230 __visit_name__ = "json_int_index" 

231 

232 render_bind_cast = True 

233 

234 

235class _PGJSONStrIndexType(sqltypes.JSON.JSONStrIndexType): 

236 __visit_name__ = "json_str_index" 

237 

238 render_bind_cast = True 

239 

240 

241class _PGJSONPathType(JSONPathType): 

242 pass 

243 

244 

245class _PGInterval(INTERVAL): 

246 render_bind_cast = True 

247 

248 

249class _PGTimeStamp(sqltypes.DateTime): 

250 render_bind_cast = True 

251 

252 

253class _PGDate(sqltypes.Date): 

254 render_bind_cast = True 

255 

256 

257class _PGTime(sqltypes.Time): 

258 render_bind_cast = True 

259 

260 

261class _PGInteger(sqltypes.Integer): 

262 render_bind_cast = True 

263 

264 

265class _PGSmallInteger(sqltypes.SmallInteger): 

266 render_bind_cast = True 

267 

268 

269class _PGNullType(sqltypes.NullType): 

270 render_bind_cast = True 

271 

272 

273class _PGBigInteger(sqltypes.BigInteger): 

274 render_bind_cast = True 

275 

276 

277class _PGBoolean(sqltypes.Boolean): 

278 render_bind_cast = True 

279 

280 

281class _PsycopgRange(ranges.AbstractSingleRangeImpl): 

282 def bind_processor(self, dialect): 

283 psycopg_Range = cast(PGDialect_psycopg, dialect)._psycopg_Range 

284 

285 def to_range(value): 

286 if isinstance(value, ranges.Range): 

287 value = psycopg_Range( 

288 value.lower, value.upper, value.bounds, value.empty 

289 ) 

290 return value 

291 

292 return to_range 

293 

294 def result_processor(self, dialect, coltype): 

295 def to_range(value): 

296 if value is not None: 

297 value = ranges.Range( 

298 value._lower, 

299 value._upper, 

300 bounds=value._bounds if value._bounds else "[)", 

301 empty=not value._bounds, 

302 ) 

303 return value 

304 

305 return to_range 

306 

307 

308class _PsycopgMultiRange(ranges.AbstractMultiRangeImpl): 

309 def bind_processor(self, dialect): 

310 psycopg_Range = cast(PGDialect_psycopg, dialect)._psycopg_Range 

311 psycopg_Multirange = cast( 

312 PGDialect_psycopg, dialect 

313 )._psycopg_Multirange 

314 

315 def to_range(value): 

316 if isinstance(value, (str, NoneType, psycopg_Multirange)): 

317 return value 

318 

319 return psycopg_Multirange( 

320 [ 

321 psycopg_Range( 

322 element.lower, 

323 element.upper, 

324 element.bounds, 

325 element.empty, 

326 ) 

327 for element in cast("Iterable[ranges.Range]", value) 

328 ] 

329 ) 

330 

331 return to_range 

332 

333 def result_processor(self, dialect, coltype): 

334 def to_range(value): 

335 if value is None: 

336 return None 

337 else: 

338 return ranges.MultiRange( 

339 ranges.Range( 

340 elem._lower, 

341 elem._upper, 

342 bounds=elem._bounds if elem._bounds else "[)", 

343 empty=not elem._bounds, 

344 ) 

345 for elem in value 

346 ) 

347 

348 return to_range 

349 

350 

351class PGExecutionContext_psycopg(_PGExecutionContext_common_psycopg): 

352 pass 

353 

354 

355class PGCompiler_psycopg(PGCompiler): 

356 pass 

357 

358 

359class PGIdentifierPreparer_psycopg(PGIdentifierPreparer): 

360 pass 

361 

362 

363def _log_notices(diagnostic): 

364 logger.info("%s: %s", diagnostic.severity, diagnostic.message_primary) 

365 

366 

367class PGDialect_psycopg(_PGDialect_common_psycopg): 

368 driver = "psycopg" 

369 

370 minimum_dbapi_version = util.VersionInfo((3, 0, 2)) 

371 

372 supports_statement_cache = True 

373 supports_server_side_cursors = True 

374 default_paramstyle = "pyformat" 

375 supports_sane_multi_rowcount = True 

376 

377 supports_native_json_serialization = True 

378 supports_native_json_deserialization = True 

379 dialect_injects_custom_json_deserializer = True 

380 

381 execution_ctx_cls = PGExecutionContext_psycopg 

382 statement_compiler = PGCompiler_psycopg 

383 preparer = PGIdentifierPreparer_psycopg 

384 

385 _has_native_hstore = True 

386 _psycopg_adapters_map = None 

387 

388 colspecs = util.update_copy( 

389 _PGDialect_common_psycopg.colspecs, 

390 { 

391 sqltypes.String: _PGString, 

392 REGCONFIG: _PGREGCONFIG, 

393 JSON: _PGJSON, 

394 CITEXT: CITEXT, 

395 sqltypes.JSON: _PGJSON, 

396 JSONB: _PGJSONB, 

397 sqltypes.JSON.JSONPathType: _PGJSONPathType, 

398 sqltypes.JSON.JSONIntIndexType: _PGJSONIntIndexType, 

399 sqltypes.JSON.JSONStrIndexType: _PGJSONStrIndexType, 

400 sqltypes.Interval: _PGInterval, 

401 INTERVAL: _PGInterval, 

402 sqltypes.Date: _PGDate, 

403 sqltypes.DateTime: _PGTimeStamp, 

404 sqltypes.Time: _PGTime, 

405 sqltypes.Integer: _PGInteger, 

406 sqltypes.SmallInteger: _PGSmallInteger, 

407 sqltypes.BigInteger: _PGBigInteger, 

408 ranges.AbstractSingleRange: _PsycopgRange, 

409 ranges.AbstractMultiRange: _PsycopgMultiRange, 

410 }, 

411 ) 

412 

413 @property 

414 def psycopg_version(self): 

415 """Legacy accessor for :attr:`.Dialect.dbapi_version`. 

416 

417 Retained for backwards compatibility; ``(0, 0)`` is returned when 

418 no version can be determined. 

419 

420 """ 

421 version = self._dbapi_version_or_none 

422 return version if version is not None else (0, 0) 

423 

424 def __init__(self, **kwargs): 

425 super().__init__(**kwargs) 

426 

427 if self.dbapi: 

428 from psycopg.adapt import AdaptersMap 

429 

430 self._psycopg_adapters_map = adapters_map = AdaptersMap( 

431 self.dbapi.adapters 

432 ) 

433 

434 if self._native_inet_types is False: 

435 import psycopg.types.string 

436 

437 adapters_map.register_loader( 

438 "inet", psycopg.types.string.TextLoader 

439 ) 

440 adapters_map.register_loader( 

441 "cidr", psycopg.types.string.TextLoader 

442 ) 

443 

444 if self._json_deserializer: 

445 from psycopg.types.json import set_json_loads 

446 

447 set_json_loads(self._json_deserializer, adapters_map) 

448 

449 if self._json_serializer: 

450 from psycopg.types.json import set_json_dumps 

451 

452 set_json_dumps(self._json_serializer, adapters_map) 

453 

454 def create_connect_args(self, url): 

455 # see https://github.com/psycopg/psycopg/issues/83 

456 cargs, cparams = super().create_connect_args(url) 

457 

458 if self._psycopg_adapters_map: 

459 cparams["context"] = self._psycopg_adapters_map 

460 if self.client_encoding is not None: 

461 cparams["client_encoding"] = self.client_encoding 

462 return cargs, cparams 

463 

464 def _type_info_fetch(self, connection, name): 

465 from psycopg.types import TypeInfo 

466 

467 return TypeInfo.fetch(connection.connection.driver_connection, name) 

468 

469 def initialize(self, connection): 

470 super().initialize(connection) 

471 

472 # PGDialect.initialize() checks server version for <= 8.2 and sets 

473 # this flag to False if so 

474 if not self.insert_returning: 

475 self.insert_executemany_returning = False 

476 

477 # HSTORE can't be registered until we have a connection so that 

478 # we can look up its OID, so we set up this adapter in 

479 # initialize() 

480 if self.use_native_hstore: 

481 info = self._type_info_fetch(connection, "hstore") 

482 self._has_native_hstore = info is not None 

483 if self._has_native_hstore: 

484 from psycopg.types.hstore import register_hstore 

485 

486 # register the adapter for connections made subsequent to 

487 # this one 

488 assert self._psycopg_adapters_map 

489 register_hstore(info, self._psycopg_adapters_map) 

490 

491 # register the adapter for this connection 

492 assert connection.connection 

493 register_hstore(info, connection.connection.driver_connection) 

494 

495 @classmethod 

496 def import_dbapi(cls): 

497 import psycopg 

498 

499 return psycopg 

500 

501 @classmethod 

502 def get_async_dialect_cls(cls, url): 

503 return PGDialectAsync_psycopg 

504 

505 @util.memoized_property 

506 def _isolation_lookup(self): 

507 return { 

508 "READ COMMITTED": self.dbapi.IsolationLevel.READ_COMMITTED, 

509 "READ UNCOMMITTED": self.dbapi.IsolationLevel.READ_UNCOMMITTED, 

510 "REPEATABLE READ": self.dbapi.IsolationLevel.REPEATABLE_READ, 

511 "SERIALIZABLE": self.dbapi.IsolationLevel.SERIALIZABLE, 

512 } 

513 

514 @util.memoized_property 

515 def _psycopg_Json(self): 

516 from psycopg.types import json 

517 

518 return json.Json 

519 

520 @util.memoized_property 

521 def _psycopg_Jsonb(self): 

522 from psycopg.types import json 

523 

524 return json.Jsonb 

525 

526 @util.memoized_property 

527 def _psycopg_TransactionStatus(self): 

528 from psycopg.pq import TransactionStatus 

529 

530 return TransactionStatus 

531 

532 @util.memoized_property 

533 def _psycopg_Range(self): 

534 from psycopg.types.range import Range 

535 

536 return Range 

537 

538 @util.memoized_property 

539 def _psycopg_Multirange(self): 

540 from psycopg.types.multirange import Multirange 

541 

542 return Multirange 

543 

544 def _do_isolation_level(self, connection, autocommit, isolation_level): 

545 connection.autocommit = autocommit 

546 connection.isolation_level = isolation_level 

547 

548 def get_isolation_level(self, dbapi_connection): 

549 status_before = dbapi_connection.info.transaction_status 

550 value = super().get_isolation_level(dbapi_connection) 

551 

552 # don't rely on psycopg providing enum symbols, compare with 

553 # eq/ne 

554 if status_before == self._psycopg_TransactionStatus.IDLE: 

555 dbapi_connection.rollback() 

556 return value 

557 

558 def set_isolation_level(self, dbapi_connection, level): 

559 if level == "AUTOCOMMIT": 

560 self._do_isolation_level( 

561 dbapi_connection, autocommit=True, isolation_level=None 

562 ) 

563 else: 

564 self._do_isolation_level( 

565 dbapi_connection, 

566 autocommit=False, 

567 isolation_level=self._isolation_lookup[level], 

568 ) 

569 

570 def set_readonly(self, connection, value): 

571 connection.read_only = value 

572 

573 def get_readonly(self, connection): 

574 return connection.read_only 

575 

576 def on_connect(self): 

577 def notices(conn): 

578 conn.add_notice_handler(_log_notices) 

579 

580 fns = [notices] 

581 

582 if self.isolation_level is not None: 

583 

584 def on_connect(conn): 

585 self.set_isolation_level(conn, self.isolation_level) 

586 

587 fns.append(on_connect) 

588 

589 # fns always has the notices function 

590 def on_connect(conn): 

591 for fn in fns: 

592 fn(conn) 

593 

594 return on_connect 

595 

596 def is_disconnect(self, e, connection, cursor): 

597 if isinstance(e, self.dbapi.Error) and connection is not None: 

598 if connection.closed or connection.broken: 

599 return True 

600 return False 

601 

602 def _twophase_idle_check(self, dbapi_conn): 

603 # don't rely on psycopg providing enum symbols, compare with eq/ne 

604 return ( 

605 dbapi_conn.info.transaction_status 

606 == self._psycopg_TransactionStatus.IDLE 

607 ) 

608 

609 @util.memoized_property 

610 def _dialect_specific_select_one(self): 

611 return ";" 

612 

613 

614class AsyncAdapt_psycopg_cursor(AsyncAdapt_dbapi_cursor): 

615 __slots__ = () 

616 

617 _awaitable_cursor_close: bool = False 

618 

619 def close(self): 

620 self._rows.clear() 

621 # Normal cursor just call _close() in a non-sync way. 

622 self._cursor._close() 

623 

624 async def _execute_async(self, operation, parameters): 

625 # override to not use mutex, psycopg3 already has mutex 

626 

627 if parameters is None: 

628 result = await self._cursor.execute(operation) 

629 else: 

630 result = await self._cursor.execute(operation, parameters) 

631 

632 # sqlalchemy result is not async, so need to pull all rows here 

633 # (assuming not a server side cursor) 

634 res = self._cursor.pgresult 

635 

636 # don't rely on psycopg providing enum symbols, compare with 

637 # eq/ne 

638 if ( 

639 not self.server_side 

640 and res 

641 and res.status == self._adapt_connection.dbapi.ExecStatus.TUPLES_OK 

642 ): 

643 self._rows = collections.deque(await self._cursor.fetchall()) 

644 return result 

645 

646 async def _executemany_async( 

647 self, 

648 operation, 

649 seq_of_parameters, 

650 ): 

651 # override to not use mutex, psycopg3 already has mutex 

652 return await self._cursor.executemany(operation, seq_of_parameters) 

653 

654 

655class AsyncAdapt_psycopg_ss_cursor( 

656 AsyncAdapt_dbapi_ss_cursor, AsyncAdapt_psycopg_cursor 

657): 

658 __slots__ = ("name",) 

659 

660 name: str 

661 

662 def __init__(self, adapt_connection, name): 

663 self.name = name 

664 super().__init__(adapt_connection) 

665 

666 def _make_new_cursor(self, connection): 

667 return connection.cursor(self.name) 

668 

669 

670class AsyncAdapt_psycopg_connection(AsyncAdapt_dbapi_connection): 

671 _connection: AsyncConnection 

672 __slots__ = () 

673 

674 _cursor_cls = AsyncAdapt_psycopg_cursor 

675 _ss_cursor_cls = AsyncAdapt_psycopg_ss_cursor 

676 

677 def add_notice_handler(self, handler): 

678 self._connection.add_notice_handler(handler) 

679 

680 @property 

681 def info(self): 

682 return self._connection.info 

683 

684 @property 

685 def adapters(self): 

686 return self._connection.adapters 

687 

688 @property 

689 def closed(self): 

690 return self._connection.closed 

691 

692 @property 

693 def broken(self): 

694 return self._connection.broken 

695 

696 @property 

697 def read_only(self): 

698 return self._connection.read_only 

699 

700 @property 

701 def deferrable(self): 

702 return self._connection.deferrable 

703 

704 @property 

705 def autocommit(self): 

706 return self._connection.autocommit 

707 

708 @autocommit.setter 

709 def autocommit(self, value): 

710 self.set_autocommit(value) 

711 

712 def set_autocommit(self, value): 

713 await_(self._connection.set_autocommit(value)) 

714 

715 def set_isolation_level(self, value): 

716 await_(self._connection.set_isolation_level(value)) 

717 

718 def set_read_only(self, value): 

719 await_(self._connection.set_read_only(value)) 

720 

721 def set_deferrable(self, value): 

722 await_(self._connection.set_deferrable(value)) 

723 

724 def cursor(self, name=None, /): 

725 if name: 

726 return AsyncAdapt_psycopg_ss_cursor(self, name) 

727 else: 

728 return AsyncAdapt_psycopg_cursor(self) 

729 

730 def tpc_begin(self, xid): 

731 return await_(self._connection.tpc_begin(xid)) 

732 

733 def tpc_prepare(self): 

734 return await_(self._connection.tpc_prepare()) 

735 

736 def tpc_commit(self, xid=None): 

737 return await_(self._connection.tpc_commit(xid)) 

738 

739 def tpc_rollback(self, xid=None): 

740 return await_(self._connection.tpc_rollback(xid)) 

741 

742 def tpc_recover(self): 

743 return await_(self._connection.tpc_recover()) 

744 

745 

746class PsycopgAdaptDBAPI(AsyncAdapt_dbapi_module): 

747 def __init__(self, psycopg, ExecStatus) -> None: 

748 super().__init__(psycopg) 

749 self.psycopg = psycopg 

750 self.ExecStatus = ExecStatus 

751 

752 for k, v in self.psycopg.__dict__.items(): 

753 if k != "connect": 

754 self.__dict__[k] = v 

755 

756 def connect(self, *arg, **kw): 

757 creator_fn = kw.pop( 

758 "async_creator_fn", self.psycopg.AsyncConnection.connect 

759 ) 

760 return await_( 

761 AsyncAdapt_psycopg_connection.create(self, creator_fn(*arg, **kw)) 

762 ) 

763 

764 

765class PGDialectAsync_psycopg(PGDialect_psycopg): 

766 is_async = True 

767 supports_statement_cache = True 

768 

769 @classmethod 

770 def import_dbapi(cls): 

771 import psycopg 

772 from psycopg.pq import ExecStatus 

773 

774 return PsycopgAdaptDBAPI(psycopg, ExecStatus) 

775 

776 def _type_info_fetch(self, connection, name): 

777 from psycopg.types import TypeInfo 

778 

779 adapted = connection.connection 

780 return await_(TypeInfo.fetch(adapted.driver_connection, name)) 

781 

782 def _do_isolation_level(self, connection, autocommit, isolation_level): 

783 connection.set_autocommit(autocommit) 

784 connection.set_isolation_level(isolation_level) 

785 

786 def _do_autocommit(self, connection, value): 

787 connection.set_autocommit(value) 

788 

789 def set_readonly(self, connection, value): 

790 connection.set_read_only(value) 

791 

792 def set_deferrable(self, connection, value): 

793 connection.set_deferrable(value) 

794 

795 def get_driver_connection(self, connection): 

796 return connection._connection 

797 

798 

799dialect = PGDialect_psycopg 

800dialect_async = PGDialectAsync_psycopg