Coverage for gws-app/gws/plugin/postgres/storage_provider.py: 98%

48 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-08-24 12:46 +0200

1"""Postgres storage provider.""" 

2 

3from typing import Optional 

4 

5import gws 

6import gws.config.util 

7import gws.lib.sa as sa 

8from sqlalchemy.dialects.postgresql import insert as pg_insert 

9 

10from . import provider 

11 

12gws.ext.new.storageProvider('postgres') 

13 

14 

15TABLE_DDL = """ 

16 CREATE TABLE IF NOT EXISTS {table_name} ( 

17 category TEXT NOT NULL, 

18 name TEXT NOT NULL, 

19 user_uid TEXT, 

20 data TEXT, 

21 created TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP, 

22 updated TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP, 

23 PRIMARY KEY (category, name) 

24 ) 

25""" 

26 

27 

28class Config(gws.Config): 

29 """Postgres storage provider. (added in 8.4)""" 

30 

31 dbUid: Optional[str] 

32 """Database provider uid.""" 

33 tableName: str 

34 """Table name for the storage.""" 

35 

36 

37class Object(gws.StorageProvider): 

38 db: provider.Object 

39 tableName: str 

40 

41 def configure(self): 

42 self.configure_provider() 

43 self.configure_table() 

44 

45 def configure_table(self): 

46 self.tableName = self.cfg('tableName') or self.cfg('_defaultTableName') 

47 if not self.db.has_table(self.tableName): 

48 raise gws.ConfigurationError(f'table {self.tableName!r} not found') 

49 

50 def configure_provider(self): 

51 return gws.config.util.configure_database_provider_for(self) 

52 

53 def list_names(self, category): 

54 with self.db.connect() as conn: 

55 tab = self._table() 

56 rs = conn.fetch_all(tab.select().where(tab.c.category == category).with_only_columns(tab.c.name)) 

57 return sorted(rec['name'] for rec in rs) 

58 

59 def read(self, category, name): 

60 with self.db.connect() as conn: 

61 tab = self._table() 

62 rec = conn.fetch_first(tab.select().where(tab.c.category == category, tab.c.name == name).limit(1)) 

63 if rec: 

64 return gws.StorageRecord(**rec) 

65 

66 def write(self, category, name, data, user_uid): 

67 with self.db.connect() as conn: 

68 tab = self._table() 

69 sql = ( 

70 pg_insert(tab) 

71 .values(category=category, name=name, user_uid=user_uid, data=data) 

72 .on_conflict_do_update( 

73 index_elements=['category', 'name'], 

74 set_=dict( 

75 user_uid=user_uid, 

76 data=data, 

77 updated=sa.func.now(), 

78 ) 

79 ) 

80 ) 

81 conn.exec_commit(sql) 

82 

83 def delete(self, category, name): 

84 with self.db.connect() as conn: 

85 tab = self._table() 

86 sql = tab.delete().where(tab.c.category == category, tab.c.name == name) 

87 conn.exec_commit(sql) 

88 

89 def _table(self): 

90 return self.db.table(self.tableName)