/usr/lib/python3/dist-packages/aiopg/connection.py is in python3-aiopg 0.7.0-3.
This file is owned by root:root, with mode 0o644.
The actual contents of the file can be viewed below.
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 | import asyncio
import sys
import warnings
import psycopg2
from psycopg2.extensions import (
POLL_OK, POLL_READ, POLL_WRITE, POLL_ERROR)
from psycopg2 import extras
from .cursor import Cursor
__all__ = ('connect',)
TIMEOUT = 60.0
PY_34 = sys.version_info >= (3, 4)
@asyncio.coroutine
def _enable_hstore(conn):
cur = yield from conn.cursor()
yield from cur.execute("""\
SELECT t.oid, typarray
FROM pg_type t JOIN pg_namespace ns
ON typnamespace = ns.oid
WHERE typname = 'hstore';
""")
rv0, rv1 = [], []
for oids in (yield from cur.fetchall()):
rv0.append(oids[0])
rv1.append(oids[1])
cur.close()
return tuple(rv0), tuple(rv1)
@asyncio.coroutine
def connect(dsn=None, *, timeout=TIMEOUT, loop=None,
enable_json=True, enable_hstore=True, echo=False, **kwargs):
"""A factory for connecting to PostgreSQL.
The coroutine accepts all parameters that psycopg2.connect() does
plus optional keyword-only `loop` and `timeout` parameters.
Returns instantiated Connection object.
"""
if loop is None:
loop = asyncio.get_event_loop()
waiter = asyncio.Future(loop=loop)
conn = Connection(dsn, loop, timeout, waiter, bool(echo), **kwargs)
try:
yield from conn._poll(waiter, timeout)
except Exception:
conn.close()
raise
if enable_json:
extras.register_default_json(conn._conn)
if enable_hstore:
oids = yield from _enable_hstore(conn)
if oids is not None:
oid, array_oid = oids
extras.register_hstore(conn._conn, oid=oid, array_oid=array_oid)
return conn
class Connection:
"""Low-level asynchronous interface for wrapped psycopg2 connection.
The Connection instance encapsulates a database session.
Provides support for creating asynchronous cursors.
"""
def __init__(self, dsn, loop, timeout, waiter, echo, **kwargs):
self._loop = loop
self._conn = psycopg2.connect(dsn, async=True, **kwargs)
self._dsn = self._conn.dsn
assert self._conn.isexecuting(), "Is conn async at all???"
self._fileno = self._conn.fileno()
self._timeout = timeout
self._waiter = waiter
self._reading = False
self._writing = False
self._echo = echo
self._ready()
def _ready(self):
if self._waiter is None:
self._fatal_error("Fatal error on aiopg connection: "
"bad state in _ready callback")
return
try:
state = self._conn.poll()
except (psycopg2.Warning, psycopg2.Error) as exc:
if self._reading:
self._loop.remove_reader(self._fileno)
self._reading = False
if self._writing:
self._loop.remove_writer(self._fileno)
self._writing = False
if not self._waiter.cancelled():
self._waiter.set_exception(exc)
else:
if state == POLL_OK:
if self._reading:
self._loop.remove_reader(self._fileno)
self._reading = False
if self._writing:
self._loop.remove_writer(self._fileno)
self._writing = False
if not self._waiter.cancelled():
self._waiter.set_result(None)
elif state == POLL_READ:
if not self._reading:
self._loop.add_reader(self._fileno, self._ready)
self._reading = True
if self._writing:
self._loop.remove_writer(self._fileno)
self._writing = False
elif state == POLL_WRITE:
if self._reading:
self._loop.remove_reader(self._fileno)
self._reading = False
if not self._writing:
self._loop.add_writer(self._fileno, self._ready)
self._writing = True
elif state == POLL_ERROR:
self._fatal_error("Fatal error on aiopg connection: "
"POLL_ERROR from underlying .poll() call")
else:
self._fatal_error("Fatal error on aiopg connection: "
"unknown answer {} from underlying "
".poll() call"
.format(state))
def _fatal_error(self, message):
# Should be called from exception handler only.
self._loop.call_exception_handler({
'message': message,
'connection': self,
})
self.close()
if self._waiter and not self._waiter.done():
self._waiter.set_exception(psycopg2.OperationalError(message))
def _create_waiter(self, func_name):
if self._waiter is not None:
raise RuntimeError('%s() called while another coroutine is '
'already waiting for incoming data' % func_name)
self._waiter = asyncio.Future(loop=self._loop)
return self._waiter
@asyncio.coroutine
def _poll(self, waiter, timeout):
assert waiter is self._waiter, (waiter, self._waiter)
self._ready()
try:
yield from asyncio.wait_for(self._waiter, timeout, loop=self._loop)
finally:
self._waiter = None
def _isexecuting(self):
return self._conn.isexecuting()
@asyncio.coroutine
def cursor(self, name=None, cursor_factory=None,
scrollable=None, withhold=False, timeout=None):
"""A coroutine that returns a new cursor object using the connection.
*cursor_factory* argument can be used to create non-standard
cursors. The argument must be suclass of
`psycopg2.extensions.cursor`.
*name*, *scrollable* and *withhold* parameters are not supported by
psycopg in asynchronous mode.
"""
if timeout is None:
timeout = self._timeout
impl = yield from self._cursor(name=name,
cursor_factory=cursor_factory,
scrollable=scrollable,
withhold=withhold)
return Cursor(self, impl, timeout, self._echo)
@asyncio.coroutine
def _cursor(self, name=None, cursor_factory=None,
scrollable=None, withhold=False):
if cursor_factory is None:
impl = self._conn.cursor(name=name,
scrollable=scrollable, withhold=withhold)
else:
impl = self._conn.cursor(name=name, cursor_factory=cursor_factory,
scrollable=scrollable, withhold=withhold)
return impl
def close(self):
"""Remove the connection from the event_loop and close it."""
# N.B. If connection contains uncommitted transaction the
# transaction will be discarded
if self._reading:
self._loop.remove_reader(self._fileno)
self._reading = False
if self._writing:
self._loop.remove_writer(self._fileno)
self._writing = False
self._conn.close()
if self._waiter is not None and not self._waiter.done():
self._waiter.set_exception(
psycopg2.OperationalError("Connection closed"))
ret = asyncio.Future(loop=self._loop)
ret.set_result(None)
return ret
@property
def closed(self):
"""Connection status.
Read-only attribute reporting whether the database connection is
open (False) or closed (True).
"""
return self._conn.closed
@property
def raw(self):
"""Underlying psycopg connection object, readonly"""
return self._conn
@asyncio.coroutine
def commit(self):
raise psycopg2.ProgrammingError(
"commit cannot be used in asynchronous mode")
@asyncio.coroutine
def rollback(self):
raise psycopg2.ProgrammingError(
"rollback cannot be used in asynchronous mode")
# TPC
@asyncio.coroutine
def xid(self, format_id, gtrid, bqual):
return self._conn.xid(format_id, gtrid, bqual)
@asyncio.coroutine
def tpc_begin(self, xid=None):
raise psycopg2.ProgrammingError(
"tpc_begin cannot be used in asynchronous mode")
@asyncio.coroutine
def tpc_prepare(self):
raise psycopg2.ProgrammingError(
"tpc_prepare cannot be used in asynchronous mode")
@asyncio.coroutine
def tpc_commit(self, xid=None):
raise psycopg2.ProgrammingError(
"tpc_commit cannot be used in asynchronous mode")
@asyncio.coroutine
def tpc_rollback(self, xid=None):
raise psycopg2.ProgrammingError(
"tpc_rollback cannot be used in asynchronous mode")
@asyncio.coroutine
def tpc_recover(self):
raise psycopg2.ProgrammingError(
"tpc_recover cannot be used in asynchronous mode")
@asyncio.coroutine
def cancel(self, timeout=None):
"""Cancel the current database operation."""
waiter = self._create_waiter('cancel')
self._conn.cancel()
if timeout is None:
timeout = self._timeout
yield from self._poll(waiter, timeout)
@asyncio.coroutine
def reset(self):
raise psycopg2.ProgrammingError(
"reset cannot be used in asynchronous mode")
@property
def dsn(self):
"""DSN connection string.
Read-only attribute representing dsn connection string used
for connectint to PostgreSQL server.
"""
return self._dsn
@asyncio.coroutine
def set_session(self, *, isolation_level=None, readonly=None,
deferrable=None, autocommit=None):
raise psycopg2.ProgrammingError(
"set_session cannot be used in asynchronous mode")
@property
def autocommit(self):
"""Autocommit status"""
return self._conn.autocommit
@autocommit.setter
def autocommit(self, val):
"""Autocommit status"""
self._conn.autocommit = val
@property
def isolation_level(self):
"""Transaction isolation level.
The only allowed value is ISOLATION_LEVEL_READ_COMMITTED.
"""
return self._conn.isolation_level
@asyncio.coroutine
def set_isolation_level(self, val):
"""Transaction isolation level.
The only allowed value is ISOLATION_LEVEL_READ_COMMITTED.
"""
self._conn.set_isolation_level(val)
@property
def encoding(self):
"""Client encoding for SQL operations."""
return self._conn.encoding
@asyncio.coroutine
def set_client_encoding(self, val):
self._conn.set_client_encoding(val)
@property
def notices(self):
"""A list of all db messages sent to the client during the session."""
return self._conn.notices
@property
def cursor_factory(self):
"""The default cursor factory used by .cursor()."""
return self._conn.cursor_factory
@asyncio.coroutine
def get_backend_pid(self):
"""Returns the PID of the backend server process."""
return self._conn.get_backend_pid()
@asyncio.coroutine
def get_parameter_status(self, parameter):
"""Look up a current parameter setting of the server."""
return self._conn.get_parameter_status(parameter)
@asyncio.coroutine
def get_transaction_status(self):
"""Return the current session transaction status as an integer."""
return self._conn.get_transaction_status()
@property
def protocol_version(self):
"""A read-only integer representing protocol being used."""
return self._conn.protocol_version
@property
def server_version(self):
"""A read-only integer representing the backend version."""
return self._conn.server_version
@property
def status(self):
"""A read-only integer representing the status of the connection."""
return self._conn.status
@asyncio.coroutine
def lobject(self, *args, **kwargs):
raise psycopg2.ProgrammingError(
"lobject cannot be used in asynchronous mode")
@property
def timeout(self):
"""Return default timeout for connection operations."""
return self._timeout
@property
def echo(self):
"""Return echo mode status."""
return self._echo
if PY_34: # pragma: no branch
def __del__(self):
if not self._conn.closed:
warnings.warn("Unclosed connection {!r}".format(self),
ResourceWarning)
self.close()
|