mirror of
https://github.com/Pardus-LiderAhenk/ahenk
synced 2024-11-26 01:12:15 +03:00
db connection curser became thread safe
This commit is contained in:
parent
1d4429949f
commit
cbc6b64439
1 changed files with 40 additions and 15 deletions
|
@ -3,7 +3,7 @@
|
||||||
# Author: İsmail BAŞARAN <ismail.basaran@tubitak.gov.tr> <basaran.ismaill@gmail.com>
|
# Author: İsmail BAŞARAN <ismail.basaran@tubitak.gov.tr> <basaran.ismaill@gmail.com>
|
||||||
# Author: Volkan Şahin <volkansah.in> <bm.volkansahin@gmail.com>
|
# Author: Volkan Şahin <volkansah.in> <bm.volkansahin@gmail.com>
|
||||||
import sqlite3
|
import sqlite3
|
||||||
|
import threading
|
||||||
from base.scope import Scope
|
from base.scope import Scope
|
||||||
|
|
||||||
|
|
||||||
|
@ -20,6 +20,9 @@ class AhenkDbService(object):
|
||||||
self.connection = None
|
self.connection = None
|
||||||
self.cursor = None
|
self.cursor = None
|
||||||
|
|
||||||
|
self.lock = threading.Lock()
|
||||||
|
|
||||||
|
|
||||||
# TODO get columns anywhere
|
# TODO get columns anywhere
|
||||||
# TODO scheduler db init get here
|
# TODO scheduler db init get here
|
||||||
|
|
||||||
|
@ -50,19 +53,29 @@ class AhenkDbService(object):
|
||||||
self.logger.error('[AhenkDbService] Database connection error ' + str(e))
|
self.logger.error('[AhenkDbService] Database connection error ' + str(e))
|
||||||
|
|
||||||
def check_and_create_table(self, table_name, cols):
|
def check_and_create_table(self, table_name, cols):
|
||||||
|
|
||||||
|
try:
|
||||||
|
self.lock.acquire(True)
|
||||||
if self.cursor:
|
if self.cursor:
|
||||||
cols = ', '.join([str(x) for x in cols])
|
cols = ', '.join([str(x) for x in cols])
|
||||||
self.cursor.execute('create table if not exists ' + table_name + ' (' + cols + ')')
|
self.cursor.execute('create table if not exists ' + table_name + ' (' + cols + ')')
|
||||||
else:
|
else:
|
||||||
self.logger.warning('[AhenkDbService] Could not create table cursor is None! Table Name : ' + str(table_name))
|
self.logger.warning('[AhenkDbService] Could not create table cursor is None! Table Name : ' + str(table_name))
|
||||||
|
finally:
|
||||||
|
self.lock.release()
|
||||||
|
|
||||||
def drop_table(self, table_name):
|
def drop_table(self, table_name):
|
||||||
|
try:
|
||||||
|
self.lock.acquire(True)
|
||||||
sql = 'DROP TABLE ' + table_name
|
sql = 'DROP TABLE ' + table_name
|
||||||
self.cursor.execute(sql)
|
self.cursor.execute(sql)
|
||||||
self.connection.commit()
|
self.connection.commit()
|
||||||
|
finally:
|
||||||
|
self.lock.release()
|
||||||
|
|
||||||
def update(self, table_name, cols, args, criteria=None):
|
def update(self, table_name, cols, args, criteria=None):
|
||||||
try:
|
try:
|
||||||
|
self.lock.acquire(True)
|
||||||
if self.connection:
|
if self.connection:
|
||||||
if criteria is None:
|
if criteria is None:
|
||||||
cols = ', '.join([str(x) for x in cols])
|
cols = ', '.join([str(x) for x in cols])
|
||||||
|
@ -82,14 +95,20 @@ class AhenkDbService(object):
|
||||||
return None
|
return None
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
self.logger.error('[AhenkDbService] Updating table error ! Table Name : ' + str(table_name) + ' ' + str(e))
|
self.logger.error('[AhenkDbService] Updating table error ! Table Name : ' + str(table_name) + ' ' + str(e))
|
||||||
|
finally:
|
||||||
|
self.lock.release()
|
||||||
|
|
||||||
def delete(self, table_name, criteria):
|
def delete(self, table_name, criteria):
|
||||||
|
try:
|
||||||
|
self.lock.acquire(True)
|
||||||
if self.cursor:
|
if self.cursor:
|
||||||
sql = 'DELETE FROM ' + table_name
|
sql = 'DELETE FROM ' + table_name
|
||||||
if criteria:
|
if criteria:
|
||||||
sql += ' where ' + str(criteria)
|
sql += ' where ' + str(criteria)
|
||||||
self.cursor.execute(sql)
|
self.cursor.execute(sql)
|
||||||
self.connection.commit()
|
self.connection.commit()
|
||||||
|
finally:
|
||||||
|
self.lock.release()
|
||||||
|
|
||||||
def findByProperty(self):
|
def findByProperty(self):
|
||||||
# Not implemented yet
|
# Not implemented yet
|
||||||
|
@ -98,6 +117,7 @@ class AhenkDbService(object):
|
||||||
def select(self, table_name, cols='*', criteria='', orderby=''):
|
def select(self, table_name, cols='*', criteria='', orderby=''):
|
||||||
if self.cursor:
|
if self.cursor:
|
||||||
try:
|
try:
|
||||||
|
self.lock.acquire(True)
|
||||||
if not cols == '*':
|
if not cols == '*':
|
||||||
cols = ', '.join([str(x) for x in cols])
|
cols = ', '.join([str(x) for x in cols])
|
||||||
sql = 'SELECT ' + cols + ' FROM ' + table_name
|
sql = 'SELECT ' + cols + ' FROM ' + table_name
|
||||||
|
@ -113,12 +133,15 @@ class AhenkDbService(object):
|
||||||
return rows
|
return rows
|
||||||
except:
|
except:
|
||||||
raise
|
raise
|
||||||
|
finally:
|
||||||
|
self.lock.release()
|
||||||
else:
|
else:
|
||||||
self.logger.warning('[AhenkDbService] Could not select table cursor is None! Table Name : ' + str(table_name))
|
self.logger.warning('[AhenkDbService] Could not select table cursor is None! Table Name : ' + str(table_name))
|
||||||
|
|
||||||
def select_one_result(self, table_name, col, criteria=''):
|
def select_one_result(self, table_name, col, criteria=''):
|
||||||
if self.cursor:
|
if self.cursor:
|
||||||
try:
|
try:
|
||||||
|
self.lock.acquire(True)
|
||||||
sql = 'SELECT ' + col + ' FROM ' + table_name
|
sql = 'SELECT ' + col + ' FROM ' + table_name
|
||||||
if criteria != '':
|
if criteria != '':
|
||||||
sql += ' where '
|
sql += ' where '
|
||||||
|
@ -131,6 +154,8 @@ class AhenkDbService(object):
|
||||||
return None
|
return None
|
||||||
except:
|
except:
|
||||||
raise
|
raise
|
||||||
|
finally:
|
||||||
|
self.lock.release()
|
||||||
else:
|
else:
|
||||||
self.logger.warning('[AhenkDbService] Could not select table cursor is None! Table Name : ' + str(table_name))
|
self.logger.warning('[AhenkDbService] Could not select table cursor is None! Table Name : ' + str(table_name))
|
||||||
|
|
||||||
|
|
Loading…
Reference in a new issue