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
|
"""Table handler.
Per-table decision how to create trigger, copy data and apply events.
"""
"""
-- redirect & create table
partition by batch_time
partition by date field
-- sql handling:
cube1 - I/U/D -> partition, insert
cube2 - I/U/D -> partition, del/insert
field remap
name remap
bublin filter
- replay: filter events
- copy: additional where
- add: add trigger args
multimaster
- replay: conflict handling, add fncall to sql queue?
- add: add 'backup' arg to trigger
plain londiste:
- replay: add to sql queue
"""
import sys, skytools
__all__ = ['RowCache', 'BaseHandler', 'parse_handler', 'build_handler', 'load_handlers']
class RowCache:
def __init__(self, table_name):
self.table_name = table_name
self.keys = {}
self.rows = []
def add_row(self, d):
row = [None] * len(self.keys)
for k, v in d.items():
try:
row[self.keys[k]] = v
except KeyError:
i = len(row)
self.keys[k] = i
row.append(v)
row = tuple(row)
self.rows.append(row)
def get_fields(self):
row = [None] * len(self.keys)
for k, i in self.keys.keys():
row[i] = k
return tuple(row)
def apply_rows(self, curs):
fields = self.get_fields()
skytools.magic_insert(curs, self.table_name, self.rows, fields)
class BaseHandler:
handler_name = 'fwd'
def __init__(self, table_name, next, args, log):
self.table_name = table_name
self.next = next
self.args = args
self.log = log
def add(self, trigger_arg_list):
"""Called when table is added.
Can modify trigger args.
"""
if self.next:
self.next.add(trigger_arg_list)
def reset(self):
"""Called before starting to process a batch.
Should clean any pending data.
"""
if self.next:
self.next.reset()
def prepare_batch(self, batch_info, dst_curs):
"""Called on first event for this table in current batch."""
if self.next:
self.next.prepare_batch(batch_info, dst_curs)
def process_event(self, ev, sql_queue_func, arg):
"""Process a event.
Event should be added to sql_queue or executed directly.
"""
if self.next:
self.next.process_event(ev, sql_queue_func, arg)
def finish_batch(self, batch_info, dst_curs):
"""Called when batch finishes."""
if self.next:
self.next.finish_batch(batch_info, dst_curs)
def prepare_copy(self, expr_list, dst_curs):
"""Can change COPY behaviour.
Returns new expr.
"""
if self.next:
self.next.prepare_copy(expr_list, dst_curs)
class TableHandler(BaseHandler):
handler_name = 'londiste'
sql_command = {
'I': "insert into %s %s;",
'U': "update only %s set %s;",
'D': "delete from only %s where %s;",
}
def process_event(self, ev, sql_queue_func, arg):
if len(ev.type) == 1:
# sql event
fqname = skytools.quote_fqident(ev.extra1)
fmt = self.sql_command[ev.type]
sql = fmt % (fqname, ev.data)
else:
# urlenc event
pklist = ev.type[2:].split(',')
row = skytools.db_urldecode(ev.data)
op = ev.type[0]
tbl = ev.extra1
if op == 'I':
sql = skytools.mk_insert_sql(row, tbl, pklist)
elif op == 'U':
sql = skytools.mk_update_sql(row, tbl, pklist)
elif op == 'D':
sql = skytools.mk_delete_sql(row, tbl, pklist)
sql_queue_func(sql, arg)
_handler_map = {
'londiste': TableHandler,
}
def register_handler_module(modname):
"""Import and module and register handlers."""
__import__(modname)
m = sys.modules[modname]
for h in m.__londiste_handlers__:
_handler_map[h.handler_name] = h
def build_handler(tblname, hlist, log):
"""Execute array of handler initializers."""
klist = []
for h in hlist:
if not h:
continue
pos = h.find('(')
if pos >= 0:
if h[-1] != ')':
raise Exception("handler fmt error")
name = h[:pos].strip()
args = h[pos+1 : -1].split(',')
args = [a.strip() for a in args]
else:
name = h
args = []
klass = _handler_map[name]
klist.append( (klass, args) )
# always append default handler
klist.append( (TableHandler, []) )
# link them together
p = None
klist.reverse()
for klass, args in klist:
p = klass(tblname, p, args, log)
return p
def parse_handler(tblname, hstr, log):
"""Parse and execute string of colon-separated handler initializers."""
hlist = hstr.split(':')
return build_handler(tblname, hlist, log)
def load_handlers(cf):
"""Load and register modules from config."""
lst = cf.getlist('handler_modules', [])
for m in lst:
register_handler_module(m)
|