experimental DAL rewrite (work in progress)

This commit is contained in:
mdipierro
2012-09-30 08:37:27 -05:00
parent 9ce1cde6ca
commit 7ec3fe3bfe
2 changed files with 85 additions and 89 deletions
+1 -1
View File
@@ -1 +1 @@
Version 2.0.9 (2012-09-30 08:31:56) dev Version 2.0.9 (2012-09-30 08:37:19) dev
+84 -88
View File
@@ -218,11 +218,11 @@ try:
except ImportError: except ImportError:
have_validators = False have_validators = False
logger = logging.getLogger("web2py.dal") LOGGER = logging.getLogger("web2py.dal")
DEFAULT = lambda:0 DEFAULT = lambda:0
sql_locker = threading.RLock() GLOBAL_LOCKER = threading.RLock()
thread_local = threading.local() THREAD_LOCAL = threading.local()
# internal representation of tables with field # internal representation of tables with field
# <table>.<field>, tables and fields may only be [a-zA-Z0-9_] # <table>.<field>, tables and fields may only be [a-zA-Z0-9_]
@@ -263,13 +263,13 @@ if not 'google' in DRIVERS:
from pysqlite2 import dbapi2 as sqlite2 from pysqlite2 import dbapi2 as sqlite2
DRIVERS.append('SQLite(sqlite2)') DRIVERS.append('SQLite(sqlite2)')
except ImportError: except ImportError:
logger.debug('no SQLite drivers pysqlite2.dbapi2') LOGGER.debug('no SQLite drivers pysqlite2.dbapi2')
try: try:
from sqlite3 import dbapi2 as sqlite3 from sqlite3 import dbapi2 as sqlite3
DRIVERS.append('SQLite(sqlite3)') DRIVERS.append('SQLite(sqlite3)')
except ImportError: except ImportError:
logger.debug('no SQLite drivers sqlite3') LOGGER.debug('no SQLite drivers sqlite3')
try: try:
# first try contrib driver, then from site-packages (if installed) # first try contrib driver, then from site-packages (if installed)
@@ -284,13 +284,13 @@ if not 'google' in DRIVERS:
import pymysql import pymysql
DRIVERS.append('MySQL(pymysql)') DRIVERS.append('MySQL(pymysql)')
except ImportError: except ImportError:
logger.debug('no MySQL driver pymysql') LOGGER.debug('no MySQL driver pymysql')
try: try:
import MySQLdb import MySQLdb
DRIVERS.append('MySQL(MySQLdb)') DRIVERS.append('MySQL(MySQLdb)')
except ImportError: except ImportError:
logger.debug('no MySQL driver MySQLDB') LOGGER.debug('no MySQL driver MySQLDB')
try: try:
@@ -298,7 +298,7 @@ if not 'google' in DRIVERS:
from psycopg2.extensions import adapt as psycopg2_adapt from psycopg2.extensions import adapt as psycopg2_adapt
DRIVERS.append('PostgreSQL(psycopg2)') DRIVERS.append('PostgreSQL(psycopg2)')
except ImportError: except ImportError:
logger.debug('no PostgreSQL driver psycopg2') LOGGER.debug('no PostgreSQL driver psycopg2')
try: try:
# first try contrib driver, then from site-packages (if installed) # first try contrib driver, then from site-packages (if installed)
@@ -308,13 +308,13 @@ if not 'google' in DRIVERS:
import pg8000.dbapi as pg8000 import pg8000.dbapi as pg8000
DRIVERS.append('PostgreSQL(pg8000)') DRIVERS.append('PostgreSQL(pg8000)')
except ImportError: except ImportError:
logger.debug('no PostgreSQL driver pg8000') LOGGER.debug('no PostgreSQL driver pg8000')
try: try:
import cx_Oracle import cx_Oracle
DRIVERS.append('Oracle(cx_Oracle)') DRIVERS.append('Oracle(cx_Oracle)')
except ImportError: except ImportError:
logger.debug('no Oracle driver cx_Oracle') LOGGER.debug('no Oracle driver cx_Oracle')
try: try:
import pyodbc import pyodbc
@@ -322,53 +322,53 @@ if not 'google' in DRIVERS:
DRIVERS.append('DB2(pyodbc)') DRIVERS.append('DB2(pyodbc)')
DRIVERS.append('Teradata(pyodbc)') DRIVERS.append('Teradata(pyodbc)')
except ImportError: except ImportError:
logger.debug('no MSSQL/DB2/Teradata driver pyodbc') LOGGER.debug('no MSSQL/DB2/Teradata driver pyodbc')
try: try:
import Sybase import Sybase
DRIVERS.append('Sybase(Sybase)') DRIVERS.append('Sybase(Sybase)')
except ImportError: except ImportError:
logger.debug('no Sybase driver') LOGGER.debug('no Sybase driver')
try: try:
import kinterbasdb import kinterbasdb
DRIVERS.append('Interbase(kinterbasdb)') DRIVERS.append('Interbase(kinterbasdb)')
DRIVERS.append('Firebird(kinterbasdb)') DRIVERS.append('Firebird(kinterbasdb)')
except ImportError: except ImportError:
logger.debug('no Firebird/Interbase driver kinterbasdb') LOGGER.debug('no Firebird/Interbase driver kinterbasdb')
try: try:
import fdb import fdb
DRIVERS.append('Firbird(fdb)') DRIVERS.append('Firbird(fdb)')
except ImportError: except ImportError:
logger.debug('no Firebird driver fdb') LOGGER.debug('no Firebird driver fdb')
##### #####
try: try:
import firebirdsql import firebirdsql
DRIVERS.append('Firebird(firebirdsql)') DRIVERS.append('Firebird(firebirdsql)')
except ImportError: except ImportError:
logger.debug('no Firebird driver firebirdsql') LOGGER.debug('no Firebird driver firebirdsql')
try: try:
import informixdb import informixdb
DRIVERS.append('Informix(informixdb)') DRIVERS.append('Informix(informixdb)')
logger.warning('Informix support is experimental') LOGGER.warning('Informix support is experimental')
except ImportError: except ImportError:
logger.debug('no Informix driver informixdb') LOGGER.debug('no Informix driver informixdb')
try: try:
import sapdb import sapdb
DRIVERS.append('SQL(sapdb)') DRIVERS.append('SQL(sapdb)')
logger.warning('SAPDB support is experimental') LOGGER.warning('SAPDB support is experimental')
except ImportError: except ImportError:
logger.debug('no SAP driver sapdb') LOGGER.debug('no SAP driver sapdb')
try: try:
import cubriddb import cubriddb
DRIVERS.append('Cubrid(cubriddb)') DRIVERS.append('Cubrid(cubriddb)')
logger.warning('Cubrid support is experimental') LOGGER.warning('Cubrid support is experimental')
except ImportError: except ImportError:
logger.debug('no Cubrid driver cubriddb') LOGGER.debug('no Cubrid driver cubriddb')
try: try:
from com.ziclix.python.sql import zxJDBC from com.ziclix.python.sql import zxJDBC
@@ -378,36 +378,36 @@ if not 'google' in DRIVERS:
zxJDBC_sqlite = java.sql.DriverManager zxJDBC_sqlite = java.sql.DriverManager
DRIVERS.append('PostgreSQL(zxJDBC)') DRIVERS.append('PostgreSQL(zxJDBC)')
DRIVERS.append('SQLite(zxJDBC)') DRIVERS.append('SQLite(zxJDBC)')
logger.warning('zxJDBC support is experimental') LOGGER.warning('zxJDBC support is experimental')
is_jdbc = True is_jdbc = True
except ImportError: except ImportError:
logger.debug('no SQLite/PostgreSQL driver zxJDBC') LOGGER.debug('no SQLite/PostgreSQL driver zxJDBC')
is_jdbc = False is_jdbc = False
try: try:
import ingresdbi import ingresdbi
DRIVERS.append('Ingres(ingresdbi)') DRIVERS.append('Ingres(ingresdbi)')
except ImportError: except ImportError:
logger.debug('no Ingres driver ingresdbi') LOGGER.debug('no Ingres driver ingresdbi')
# NOTE could try JDBC....... # NOTE could try JDBC.......
try: try:
import couchdb import couchdb
DRIVERS.append('CouchDB(couchdb)') DRIVERS.append('CouchDB(couchdb)')
except ImportError: except ImportError:
logger.debug('no Couchdb driver couchdb') LOGGER.debug('no Couchdb driver couchdb')
try: try:
import pymongo import pymongo
DRIVERS.append('MongoDB(pymongo)') DRIVERS.append('MongoDB(pymongo)')
except: except:
logger.debug('no MongoDB driver pymongo') LOGGER.debug('no MongoDB driver pymongo')
try: try:
import imaplib import imaplib
DRIVERS.append('IMAP(imaplib)') DRIVERS.append('IMAP(imaplib)')
except: except:
logger.debug('no IMAP driver imaplib') LOGGER.debug('no IMAP driver imaplib')
PLURALIZE_RULES = [ PLURALIZE_RULES = [
(re.compile('child$'), re.compile('child$'), 'children'), (re.compile('child$'), re.compile('child$'), 'children'),
@@ -493,42 +493,48 @@ class ConnectionPool(object):
@staticmethod @staticmethod
def set_folder(folder): def set_folder(folder):
thread_local.folder = folder THREAD_LOCAL.folder = folder
# ## this allows gluon to commit/rollback all dbs in this thread # ## this allows gluon to commit/rollback all dbs in this thread
@staticmethod
def recycle_connection(adapter,action):
if action:
if callable(action):
action(adapter)
else:
getattr(adapter, action)()
# ## if you want pools, recycle this connection
really = True
if adapter.pool_size:
GLOBAL_LOCKER.acquire()
pool = ConnectionPool.pools[adapter.uri]
if len(pool) < adapter.pool_size:
pool.append(adapter.connection)
really = False
GLOBAL_LOCKER.release()
if really:
getattr(adapter, 'close')()
@staticmethod @staticmethod
def close_all_instances(action): def close_all_instances(action):
""" to close cleanly databases in a multithreaded environment """ """ to close cleanly databases in a multithreaded environment """
if hasattr(thread_local, 'instances'): dbs = getattr(THREAD_LOCAL,'db_instances',{}).items()
while thread_local.instances: for singleton_code, db in dbs:
instance = thread_local.instances.pop() try:
if action: adapter = db._adapter
if callable(action): except AttributeError:
action(instance) pass
else: else:
getattr(instance, action)() ConnectionPool.recycle_connection(adapter,action)
# ## if you want pools, recycle this connection del THREAD_LOCAL.db_instances[singleton_code]
really = True
if instance.pool_size:
sql_locker.acquire()
pool = ConnectionPool.pools[instance.uri]
if len(pool) < instance.pool_size:
pool.append(instance.connection)
really = False
sql_locker.release()
if really:
getattr(instance, 'close')()
if instance.db._singleton_code in thread_local.db_instances:
del thread_local.db_instances[instance.db._singleton_code]
if callable(action): if callable(action):
action(None) action(None)
return return
def find_or_make_work_folder(self): def find_or_make_work_folder(self):
""" this actually does not make the folder. it has to be there """ """ this actually does not make the folder. it has to be there """
self.folder = getattr(thread_local,'folder','') self.folder = getattr(THREAD_LOCAL,'folder','')
# Creating the folder if it does not exist # Creating the folder if it does not exist
if False and self.folder and not exists(self.folder): if False and self.folder and not exists(self.folder):
@@ -540,7 +546,7 @@ class ConnectionPool(object):
def reconnect(self): def reconnect(self):
""" allows a thread to re-connect to server or re-pool """ """ allows a thread to re-connect to server or re-pool """
self.close_all_instances(False) ConnectionPool.recycle_connection(self,False) ### WHY?
self.pool_connection(self._connection_function) self.pool_connection(self._connection_function)
self.after_connection() self.after_connection()
@@ -560,12 +566,12 @@ class ConnectionPool(object):
else: else:
uri = self.uri uri = self.uri
while True: while True:
sql_locker.acquire() GLOBAL_LOCKER.acquire()
if not uri in pools: if not uri in pools:
pools[uri] = [] pools[uri] = []
if pools[uri]: if pools[uri]:
self.connection = pools[uri].pop() self.connection = pools[uri].pop()
sql_locker.release() GLOBAL_LOCKER.release()
self.cursor = cursor and self.connection.cursor() self.cursor = cursor and self.connection.cursor()
try: try:
if self.cursor and self.check_active_connection: if self.cursor and self.check_active_connection:
@@ -574,14 +580,10 @@ class ConnectionPool(object):
except: except:
pass pass
else: else:
sql_locker.release() GLOBAL_LOCKER.release()
self.connection = f() self.connection = f()
self.cursor = cursor and self.connection.cursor() self.cursor = cursor and self.connection.cursor()
break break
if not hasattr(thread_local,'instances'):
thread_local.instances = []
thread_local.instances.append(self)
################################################################################### ###################################################################################
# this is a generic adapter that does nothing; all others are derived from this one # this is a generic adapter that does nothing; all others are derived from this one
@@ -1663,7 +1665,7 @@ class BaseAdapter(ConnectionPool):
def log_execute(self, *a, **b): def log_execute(self, *a, **b):
command = a[0] command = a[0]
if self.db._debug: if self.db._debug:
logger.debug('SQL: %s' % command) LOGGER.debug('SQL: %s' % command)
self.db._lastsql = command self.db._lastsql = command
t0 = time.time() t0 = time.time()
ret = self.cursor.execute(*a, **b) ret = self.cursor.execute(*a, **b)
@@ -2983,7 +2985,7 @@ class MSSQLAdapter(BaseAdapter):
if not dsn: if not dsn:
raise SyntaxError, 'DSN required' raise SyntaxError, 'DSN required'
except SyntaxError, e: except SyntaxError, e:
logger.error('NdGpatch error') LOGGER.error('NdGpatch error')
raise e raise e
# was cnxn = 'DSN=%s' % dsn # was cnxn = 'DSN=%s' % dsn
cnxn = dsn cnxn = dsn
@@ -3181,7 +3183,7 @@ class SybaseAdapter(MSSQLAdapter):
if not dsn: if not dsn:
raise SyntaxError, 'DSN required' raise SyntaxError, 'DSN required'
except SyntaxError, e: except SyntaxError, e:
logger.error('NdGpatch error') LOGGER.error('NdGpatch error')
raise e raise e
else: else:
m = self.REGEX_URI.match(uri) m = self.REGEX_URI.match(uri)
@@ -4015,7 +4017,7 @@ class GoogleSQLAdapter(UseDatabaseStoredFile,MySQLAdapter):
self.uri = uri self.uri = uri
self.pool_size = pool_size self.pool_size = pool_size
self.db_codec = db_codec self.db_codec = db_codec
self.folder = folder or pjoin('$HOME',thread_local.folder.split( self.folder = folder or pjoin('$HOME',THREAD_LOCAL.folder.split(
os.sep+'applications'+os.sep,1)[1]) os.sep+'applications'+os.sep,1)[1])
ruri = uri.split("://")[1] ruri = uri.split("://")[1]
m = self.REGEX_URI.match(ruri) m = self.REGEX_URI.match(ruri)
@@ -4585,7 +4587,7 @@ class GoogleDatastoreAdapter(NoSQLAdapter):
setattr(item, field.name, self.represent(value,field.type)) setattr(item, field.name, self.represent(value,field.type))
item.put() item.put()
counter += 1 counter += 1
logger.info(str(counter)) LOGGER.info(str(counter))
return counter return counter
def insert(self,table,fields): def insert(self,table,fields):
@@ -5531,12 +5533,12 @@ class IMAPAdapter(NoSQLAdapter):
else: else:
uri = self.uri uri = self.uri
while True: while True:
sql_locker.acquire() GLOBAL_LOCKER.acquire()
if not uri in pools: if not uri in pools:
pools[uri] = [] pools[uri] = []
if pools[uri]: if pools[uri]:
self.connection = pools[uri].pop() self.connection = pools[uri].pop()
sql_locker.release() GLOBAL_LOCKER.release()
self.cursor = cursor and self.connection.cursor() self.cursor = cursor and self.connection.cursor()
if self.cursor and self.check_active_connection: if self.cursor and self.check_active_connection:
try: try:
@@ -5548,15 +5550,11 @@ class IMAPAdapter(NoSQLAdapter):
self.connection = f() self.connection = f()
break break
else: else:
sql_locker.release() GLOBAL_LOCKER.release()
self.connection = f() self.connection = f()
self.cursor = cursor and self.connection.cursor() self.cursor = cursor and self.connection.cursor()
break break
if not hasattr(thread_local,'instances'):
thread_local.instances = []
thread_local.instances.append(self)
def get_last_message(self, tablename): def get_last_message(self, tablename):
last_message = None last_message = None
# request mailbox list to the server # request mailbox list to the server
@@ -5567,7 +5565,7 @@ class IMAPAdapter(NoSQLAdapter):
result = self.connection.select(self.connection.mailbox_names[tablename]) result = self.connection.select(self.connection.mailbox_names[tablename])
last_message = int(result[1][0]) last_message = int(result[1][0])
except (IndexError, ValueError, TypeError, KeyError), e: except (IndexError, ValueError, TypeError, KeyError), e:
logger.debug("Error retrieving the last mailbox sequence number. %s" % str(e)) LOGGER.debug("Error retrieving the last mailbox sequence number. %s" % str(e))
return last_message return last_message
def get_uid_bounds(self, tablename): def get_uid_bounds(self, tablename):
@@ -5599,7 +5597,7 @@ class IMAPAdapter(NoSQLAdapter):
try: try:
dayname, datestring = date.split(",") dayname, datestring = date.split(",")
except (ValueError): except (ValueError):
logger.debug("Could not parse date text: %s" % date) LOGGER.debug("Could not parse date text: %s" % date)
return None return None
date_list = datestring.strip().split() date_list = datestring.strip().split()
year = int(date_list[2]) year = int(date_list[2])
@@ -5729,7 +5727,7 @@ class IMAPAdapter(NoSQLAdapter):
def create_table(self, *args, **kwargs): def create_table(self, *args, **kwargs):
# not implemented # not implemented
logger.debug("Create table feature is not implemented for %s" % type(self)) LOGGER.debug("Create table feature is not implemented for %s" % type(self))
def _select(self,query,fields,attributes): def _select(self,query,fields,attributes):
""" Search and Fetch records and return web2py """ Search and Fetch records and return web2py
@@ -6057,7 +6055,7 @@ class IMAPAdapter(NoSQLAdapter):
try: try:
pedestal, threshold = self.get_uid_bounds(first.tablename) pedestal, threshold = self.get_uid_bounds(first.tablename)
except TypeError, e: except TypeError, e:
logger.debug("Error requesting uid bounds: %s", str(e)) LOGGER.debug("Error requesting uid bounds: %s", str(e))
return "" return ""
try: try:
lower_limit = int(self.expand(second)) + 1 lower_limit = int(self.expand(second)) + 1
@@ -6085,7 +6083,7 @@ class IMAPAdapter(NoSQLAdapter):
try: try:
pedestal, threshold = self.get_uid_bounds(first.tablename) pedestal, threshold = self.get_uid_bounds(first.tablename)
except TypeError, e: except TypeError, e:
logger.debug("Error requesting uid bounds: %s", str(e)) LOGGER.debug("Error requesting uid bounds: %s", str(e))
return "" return ""
lower_limit = self.expand(second) lower_limit = self.expand(second)
result = "UID %s:%s" % (lower_limit, threshold) result = "UID %s:%s" % (lower_limit, threshold)
@@ -6104,7 +6102,7 @@ class IMAPAdapter(NoSQLAdapter):
try: try:
pedestal, threshold = self.get_uid_bounds(first.tablename) pedestal, threshold = self.get_uid_bounds(first.tablename)
except TypeError, e: except TypeError, e:
logger.debug("Error requesting uid bounds: %s", str(e)) LOGGER.debug("Error requesting uid bounds: %s", str(e))
return "" return ""
try: try:
upper_limit = int(self.expand(second)) - 1 upper_limit = int(self.expand(second)) - 1
@@ -6128,7 +6126,7 @@ class IMAPAdapter(NoSQLAdapter):
try: try:
pedestal, threshold = self.get_uid_bounds(first.tablename) pedestal, threshold = self.get_uid_bounds(first.tablename)
except TypeError, e: except TypeError, e:
logger.debug("Error requesting uid bounds: %s", str(e)) LOGGER.debug("Error requesting uid bounds: %s", str(e))
return "" return ""
upper_limit = int(self.expand(second)) upper_limit = int(self.expand(second))
result = "UID %s:%s" % (pedestal, upper_limit) result = "UID %s:%s" % (pedestal, upper_limit)
@@ -6572,19 +6570,19 @@ class DAL(object):
""" """
def __new__(cls, uri='sqlite://dummy.db', *args, **kwargs): def __new__(cls, uri='sqlite://dummy.db', *args, **kwargs):
if not hasattr(thread_local,'db_instances'): if not hasattr(THREAD_LOCAL,'db_instances'):
thread_local.db_instances = {} THREAD_LOCAL.db_instances = {}
if 'singleton_code' in kwargs: if 'singleton_code' in kwargs:
singleton_code = kwargs['singleton_code'] singleton_code = kwargs['singleton_code']
del kwargs['singleton_code'] del kwargs['singleton_code']
singleton_code = hashlib.md5(repr(uri)).hexdigest() singleton_code = hashlib.md5(repr(uri)).hexdigest()
try: try:
db = thread_local.db_instances[singleton_code] db = THREAD_LOCAL.db_instances[singleton_code]
if args or kwargs: if args or kwargs:
raise RuntimeError, 'Cannot duplicate a Singleton' raise RuntimeError, 'Cannot duplicate a Singleton'
except KeyError: except KeyError:
db = super(DAL, cls).__new__(cls, uri, *args, **kwargs) db = super(DAL, cls).__new__(cls, uri, *args, **kwargs)
thread_local.db_instances[singleton_code] = db THREAD_LOCAL.db_instances[singleton_code] = db
db._singleton_code = singleton_code db._singleton_code = singleton_code
return db return db
@@ -7062,12 +7060,12 @@ def index():
args_get('fake_migrate',self._fake_migrate) args_get('fake_migrate',self._fake_migrate)
polymodel = args_get('polymodel',None) polymodel = args_get('polymodel',None)
try: try:
sql_locker.acquire() GLOBAL_LOCKER.acquire()
self._adapter.create_table(table,migrate=migrate, self._adapter.create_table(table,migrate=migrate,
fake_migrate=fake_migrate, fake_migrate=fake_migrate,
polymodel=polymodel) polymodel=polymodel)
finally: finally:
sql_locker.release() GLOBAL_LOCKER.release()
else: else:
table._dbt = None table._dbt = None
on_define = args_get('on_define',None) on_define = args_get('on_define',None)
@@ -7132,10 +7130,8 @@ def index():
def close(self): def close(self):
adapter = self._adapter adapter = self._adapter
if adapter in thread_local.instances: if self._singleton_code in THREAD_LOCAL.db_instances:
thread_local.instances.remove(adapter) del THREAD_LOCAL.db_instances[self._singleton_code]
if self._singleton_code in thread_local.db_instances:
del thread_local.db_instances[self._singleton_code]
adapter.close() adapter.close()
def executesql(self, query, placeholders=None, as_dict=False, def executesql(self, query, placeholders=None, as_dict=False,