/
idl.py
2041 lines (1753 loc) · 83.1 KB
/
idl.py
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
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
# Copyright (c) 2009, 2010, 2011, 2012, 2013, 2016 Nicira, Inc.
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at:
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
import collections
import functools
import uuid
import ovs.db.data as data
import ovs.db.parser
import ovs.db.schema
import ovs.jsonrpc
import ovs.ovsuuid
import ovs.poller
import ovs.vlog
from ovs.db import custom_index
from ovs.db import error
vlog = ovs.vlog.Vlog("idl")
__pychecker__ = 'no-classattr no-objattrs'
ROW_CREATE = "create"
ROW_UPDATE = "update"
ROW_DELETE = "delete"
OVSDB_UPDATE = 0
OVSDB_UPDATE2 = 1
CLUSTERED = "clustered"
Notice = collections.namedtuple('Notice', ('event', 'row', 'updates'))
Notice.__new__.__defaults__ = (None,) # default updates=None
class Idl(object):
"""Open vSwitch Database Interface Definition Language (OVSDB IDL).
The OVSDB IDL maintains an in-memory replica of a database. It issues RPC
requests to an OVSDB database server and parses the responses, converting
raw JSON into data structures that are easier for clients to digest.
The IDL also assists with issuing database transactions. The client
creates a transaction, manipulates the IDL data structures, and commits or
aborts the transaction. The IDL then composes and issues the necessary
JSON-RPC requests and reports to the client whether the transaction
completed successfully.
The client is allowed to access the following attributes directly, in a
read-only fashion:
- 'tables': This is the 'tables' map in the ovs.db.schema.DbSchema provided
to the Idl constructor. Each ovs.db.schema.TableSchema in the map is
annotated with a new attribute 'rows', which is a dict from a uuid.UUID
to a Row object.
The client may directly read and write the Row objects referenced by the
'rows' map values. Refer to Row for more details.
- 'change_seqno': A number that represents the IDL's state. When the IDL
is updated (by Idl.run()), its value changes. The sequence number can
occasionally change even if the database does not. This happens if the
connection to the database drops and reconnects, which causes the
database contents to be reloaded even if they didn't change. (It could
also happen if the database server sends out a "change" that reflects
what the IDL already thought was in the database. The database server is
not supposed to do that, but bugs could in theory cause it to do so.)
- 'lock_name': The name of the lock configured with Idl.set_lock(), or None
if no lock is configured.
- 'has_lock': True, if the IDL is configured to obtain a lock and owns that
lock, and False otherwise.
Locking and unlocking happens asynchronously from the database client's
point of view, so the information is only useful for optimization
(e.g. if the client doesn't have the lock then there's no point in trying
to write to the database).
- 'is_lock_contended': True, if the IDL is configured to obtain a lock but
the database server has indicated that some other client already owns the
requested lock, and False otherwise.
- 'txn': The ovs.db.idl.Transaction object for the database transaction
currently being constructed, if there is one, or None otherwise.
"""
IDL_S_INITIAL = 0
IDL_S_SERVER_SCHEMA_REQUESTED = 1
IDL_S_SERVER_MONITOR_REQUESTED = 2
IDL_S_DATA_MONITOR_REQUESTED = 3
IDL_S_DATA_MONITOR_COND_REQUESTED = 4
def __init__(self, remote, schema_helper, probe_interval=None,
leader_only=True):
"""Creates and returns a connection to the database named 'db_name' on
'remote', which should be in a form acceptable to
ovs.jsonrpc.session.open(). The connection will maintain an in-memory
replica of the remote database.
'remote' can be comma separated multiple remotes and each remote
should be in a form acceptable to ovs.jsonrpc.session.open().
'schema_helper' should be an instance of the SchemaHelper class which
generates schema for the remote database. The caller may have cut it
down by removing tables or columns that are not of interest. The IDL
will only replicate the tables and columns that remain. The caller may
also add an attribute named 'alert' to selected remaining columns,
setting its value to False; if so, then changes to those columns will
not be considered changes to the database for the purpose of the return
value of Idl.run() and Idl.change_seqno. This is useful for columns
that the IDL's client will write but not read.
As a convenience to users, 'schema' may also be an instance of the
SchemaHelper class.
The IDL uses and modifies 'schema' directly.
If 'leader_only' is set to True (default value) the IDL will only
monitor and transact with the leader of the cluster.
If "probe_interval" is zero it disables the connection keepalive
feature. If non-zero the value will be forced to at least 1000
milliseconds. If None it will just use the default value in OVS.
"""
assert isinstance(schema_helper, SchemaHelper)
schema = schema_helper.get_idl_schema()
self.tables = schema.tables
self.readonly = schema.readonly
self._db = schema
remotes = self._parse_remotes(remote)
self._session = ovs.jsonrpc.Session.open_multiple(remotes,
probe_interval=probe_interval)
self._monitor_request_id = None
self._last_seqno = None
self.change_seqno = 0
self.uuid = uuid.uuid1()
# Server monitor.
self._server_schema_request_id = None
self._server_monitor_request_id = None
self._db_change_aware_request_id = None
self._server_db_name = '_Server'
self._server_db_table = 'Database'
self.server_tables = None
self._server_db = None
self.server_monitor_uuid = uuid.uuid1()
self.leader_only = leader_only
self.cluster_id = None
self._min_index = 0
self.state = self.IDL_S_INITIAL
# Database locking.
self.lock_name = None # Name of lock we need, None if none.
self.has_lock = False # Has db server said we have the lock?
self.is_lock_contended = False # Has db server said we can't get lock?
self._lock_request_id = None # JSON-RPC ID of in-flight lock request.
# Transaction support.
self.txn = None
self._outstanding_txns = {}
for table in schema.tables.values():
for column in table.columns.values():
if not hasattr(column, 'alert'):
column.alert = True
table.need_table = False
table.rows = custom_index.IndexedRows(table)
table.idl = self
table.condition = [True]
table.cond_changed = False
def _parse_remotes(self, remote):
# If remote is -
# "tcp:10.0.0.1:6641,unix:/tmp/db.sock,t,s,tcp:10.0.0.2:6642"
# this function returns
# ["tcp:10.0.0.1:6641", "unix:/tmp/db.sock,t,s", tcp:10.0.0.2:6642"]
remotes = []
for r in remote.split(','):
if remotes and r.find(":") == -1:
remotes[-1] += "," + r
else:
remotes.append(r)
return remotes
def set_cluster_id(self, cluster_id):
"""Set the id of the cluster that this idl must connect to."""
self.cluster_id = cluster_id
if self.state != self.IDL_S_INITIAL:
self.force_reconnect()
def index_create(self, table, name):
"""Create a named multi-column index on a table"""
return self.tables[table].rows.index_create(name)
def index_irange(self, table, name, start, end):
"""Return items in a named index between start/end inclusive"""
return self.tables[table].rows.indexes[name].irange(start, end)
def index_equal(self, table, name, value):
"""Return items in a named index matching a value"""
return self.tables[table].rows.indexes[name].irange(value, value)
def close(self):
"""Closes the connection to the database. The IDL will no longer
update."""
self._session.close()
def run(self):
"""Processes a batch of messages from the database server. Returns
True if the database as seen through the IDL changed, False if it did
not change. The initial fetch of the entire contents of the remote
database is considered to be one kind of change. If the IDL has been
configured to acquire a database lock (with Idl.set_lock()), then
successfully acquiring the lock is also considered to be a change.
This function can return occasional false positives, that is, report
that the database changed even though it didn't. This happens if the
connection to the database drops and reconnects, which causes the
database contents to be reloaded even if they didn't change. (It could
also happen if the database server sends out a "change" that reflects
what we already thought was in the database, but the database server is
not supposed to do that.)
As an alternative to checking the return value, the client may check
for changes in self.change_seqno."""
assert not self.txn
initial_change_seqno = self.change_seqno
self.send_cond_change()
self._session.run()
i = 0
while i < 50:
i += 1
if not self._session.is_connected():
break
seqno = self._session.get_seqno()
if seqno != self._last_seqno:
self._last_seqno = seqno
self.__txn_abort_all()
self.__send_server_schema_request()
if self.lock_name:
self.__send_lock_request()
break
msg = self._session.recv()
if msg is None:
break
if (msg.type == ovs.jsonrpc.Message.T_NOTIFY
and msg.method == "update2"
and len(msg.params) == 2):
# Database contents changed.
self.__parse_update(msg.params[1], OVSDB_UPDATE2)
elif (msg.type == ovs.jsonrpc.Message.T_NOTIFY
and msg.method == "update"
and len(msg.params) == 2):
# Database contents changed.
if msg.params[0] == str(self.server_monitor_uuid):
self.__parse_update(msg.params[1], OVSDB_UPDATE,
tables=self.server_tables)
self.change_seqno = initial_change_seqno
if not self.__check_server_db():
self.force_reconnect()
break
else:
self.__parse_update(msg.params[1], OVSDB_UPDATE)
elif (msg.type == ovs.jsonrpc.Message.T_REPLY
and self._monitor_request_id is not None
and self._monitor_request_id == msg.id):
# Reply to our "monitor" request.
try:
self.change_seqno += 1
self._monitor_request_id = None
self.__clear()
if self.state == self.IDL_S_DATA_MONITOR_COND_REQUESTED:
self.__parse_update(msg.result, OVSDB_UPDATE2)
else:
assert self.state == self.IDL_S_DATA_MONITOR_REQUESTED
self.__parse_update(msg.result, OVSDB_UPDATE)
except error.Error as e:
vlog.err("%s: parse error in received schema: %s"
% (self._session.get_name(), e))
self.__error()
elif (msg.type == ovs.jsonrpc.Message.T_REPLY
and self._server_schema_request_id is not None
and self._server_schema_request_id == msg.id):
# Reply to our "get_schema" of _Server request.
try:
self._server_schema_request_id = None
sh = SchemaHelper(None, msg.result)
sh.register_table(self._server_db_table)
schema = sh.get_idl_schema()
self._server_db = schema
self.server_tables = schema.tables
self.__send_server_monitor_request()
except error.Error as e:
vlog.err("%s: error receiving server schema: %s"
% (self._session.get_name(), e))
if self.cluster_id:
self.__error()
break
else:
self.change_seqno = initial_change_seqno
self.__send_monitor_request()
elif (msg.type == ovs.jsonrpc.Message.T_REPLY
and self._server_monitor_request_id is not None
and self._server_monitor_request_id == msg.id):
# Reply to our "monitor" of _Server request.
try:
self._server_monitor_request_id = None
self.__parse_update(msg.result, OVSDB_UPDATE,
tables=self.server_tables)
self.change_seqno = initial_change_seqno
if self.__check_server_db():
self.__send_monitor_request()
self.__send_db_change_aware()
else:
self.force_reconnect()
break
except error.Error as e:
vlog.err("%s: parse error in received schema: %s"
% (self._session.get_name(), e))
if self.cluster_id:
self.__error()
break
else:
self.change_seqno = initial_change_seqno
self.__send_monitor_request()
elif (msg.type == ovs.jsonrpc.Message.T_REPLY
and self._db_change_aware_request_id is not None
and self._db_change_aware_request_id == msg.id):
# Reply to us notifying the server of our change awarness.
self._db_change_aware_request_id = None
elif (msg.type == ovs.jsonrpc.Message.T_REPLY
and self._lock_request_id is not None
and self._lock_request_id == msg.id):
# Reply to our "lock" request.
self.__parse_lock_reply(msg.result)
elif (msg.type == ovs.jsonrpc.Message.T_NOTIFY
and msg.method == "locked"):
# We got our lock.
self.__parse_lock_notify(msg.params, True)
elif (msg.type == ovs.jsonrpc.Message.T_NOTIFY
and msg.method == "stolen"):
# Someone else stole our lock.
self.__parse_lock_notify(msg.params, False)
elif msg.type == ovs.jsonrpc.Message.T_NOTIFY and msg.id == "echo":
# Reply to our echo request. Ignore it.
pass
elif (msg.type == ovs.jsonrpc.Message.T_ERROR and
self.state == self.IDL_S_DATA_MONITOR_COND_REQUESTED and
self._monitor_request_id == msg.id):
if msg.error == "unknown method":
self.__send_monitor_request()
elif (msg.type == ovs.jsonrpc.Message.T_ERROR and
self._server_schema_request_id is not None and
self._server_schema_request_id == msg.id):
self._server_schema_request_id = None
if self.cluster_id:
self.force_reconnect()
break
else:
self.change_seqno = initial_change_seqno
self.__send_monitor_request()
elif (msg.type in (ovs.jsonrpc.Message.T_ERROR,
ovs.jsonrpc.Message.T_REPLY)
and self.__txn_process_reply(msg)):
# __txn_process_reply() did everything needed.
pass
else:
# This can happen if a transaction is destroyed before we
# receive the reply, so keep the log level low.
vlog.dbg("%s: received unexpected %s message"
% (self._session.get_name(),
ovs.jsonrpc.Message.type_to_string(msg.type)))
return initial_change_seqno != self.change_seqno
def send_cond_change(self):
if not self._session.is_connected():
return
for table in self.tables.values():
if table.cond_changed:
self.__send_cond_change(table, table.condition)
table.cond_changed = False
def cond_change(self, table_name, cond):
"""Sets the condition for 'table_name' to 'cond', which should be a
conditional expression suitable for use directly in the OVSDB
protocol, with the exception that the empty condition []
matches no rows (instead of matching every row). That is, []
is equivalent to [False], not to [True].
"""
table = self.tables.get(table_name)
if not table:
raise error.Error('Unknown table "%s"' % table_name)
if cond == []:
cond = [False]
if table.condition != cond:
table.condition = cond
table.cond_changed = True
def wait(self, poller):
"""Arranges for poller.block() to wake up when self.run() has something
to do or when activity occurs on a transaction on 'self'."""
self._session.wait(poller)
self._session.recv_wait(poller)
def has_ever_connected(self):
"""Returns True, if the IDL successfully connected to the remote
database and retrieved its contents (even if the connection
subsequently dropped and is in the process of reconnecting). If so,
then the IDL contains an atomic snapshot of the database's contents
(but it might be arbitrarily old if the connection dropped).
Returns False if the IDL has never connected or retrieved the
database's contents. If so, the IDL is empty."""
return self.change_seqno != 0
def force_reconnect(self):
"""Forces the IDL to drop its connection to the database and reconnect.
In the meantime, the contents of the IDL will not change."""
self._session.force_reconnect()
def session_name(self):
return self._session.get_name()
def set_lock(self, lock_name):
"""If 'lock_name' is not None, configures the IDL to obtain the named
lock from the database server and to avoid modifying the database when
the lock cannot be acquired (that is, when another client has the same
lock).
If 'lock_name' is None, drops the locking requirement and releases the
lock."""
assert not self.txn
assert not self._outstanding_txns
if self.lock_name and (not lock_name or lock_name != self.lock_name):
# Release previous lock.
self.__send_unlock_request()
self.lock_name = None
self.is_lock_contended = False
if lock_name and not self.lock_name:
# Acquire new lock.
self.lock_name = lock_name
self.__send_lock_request()
def notify(self, event, row, updates=None):
"""Hook for implementing create/update/delete notifications
:param event: The event that was triggered
:type event: ROW_CREATE, ROW_UPDATE, or ROW_DELETE
:param row: The row as it is after the operation has occured
:type row: Row
:param updates: For updates, row with only old values of the changed
columns
:type updates: Row
"""
def __send_cond_change(self, table, cond):
monitor_cond_change = {table.name: [{"where": cond}]}
old_uuid = str(self.uuid)
self.uuid = uuid.uuid1()
params = [old_uuid, str(self.uuid), monitor_cond_change]
msg = ovs.jsonrpc.Message.create_request("monitor_cond_change", params)
self._session.send(msg)
def __clear(self):
changed = False
for table in self.tables.values():
if table.rows:
changed = True
table.rows = custom_index.IndexedRows(table)
if changed:
self.change_seqno += 1
def __update_has_lock(self, new_has_lock):
if new_has_lock and not self.has_lock:
if self._monitor_request_id is None:
self.change_seqno += 1
else:
# We're waiting for a monitor reply, so don't signal that the
# database changed. The monitor reply will increment
# change_seqno anyhow.
pass
self.is_lock_contended = False
self.has_lock = new_has_lock
def __do_send_lock_request(self, method):
self.__update_has_lock(False)
self._lock_request_id = None
if self._session.is_connected():
msg = ovs.jsonrpc.Message.create_request(method, [self.lock_name])
msg_id = msg.id
self._session.send(msg)
else:
msg_id = None
return msg_id
def __send_lock_request(self):
self._lock_request_id = self.__do_send_lock_request("lock")
def __send_unlock_request(self):
self.__do_send_lock_request("unlock")
def __parse_lock_reply(self, result):
self._lock_request_id = None
got_lock = isinstance(result, dict) and result.get("locked") is True
self.__update_has_lock(got_lock)
if not got_lock:
self.is_lock_contended = True
def __parse_lock_notify(self, params, new_has_lock):
if (self.lock_name is not None
and isinstance(params, (list, tuple))
and params
and params[0] == self.lock_name):
self.__update_has_lock(new_has_lock)
if not new_has_lock:
self.is_lock_contended = True
def __send_db_change_aware(self):
msg = ovs.jsonrpc.Message.create_request("set_db_change_aware",
[True])
self._db_change_aware_request_id = msg.id
self._session.send(msg)
def __send_monitor_request(self):
if (self.state in [self.IDL_S_SERVER_MONITOR_REQUESTED,
self.IDL_S_INITIAL]):
self.state = self.IDL_S_DATA_MONITOR_COND_REQUESTED
method = "monitor_cond"
else:
self.state = self.IDL_S_DATA_MONITOR_REQUESTED
method = "monitor"
monitor_requests = {}
for table in self.tables.values():
columns = []
for column in table.columns.keys():
if ((table.name not in self.readonly) or
(table.name in self.readonly) and
(column not in self.readonly[table.name])):
columns.append(column)
monitor_request = {"columns": columns}
if method == "monitor_cond" and table.condition != [True]:
monitor_request["where"] = table.condition
table.cond_change = False
monitor_requests[table.name] = [monitor_request]
msg = ovs.jsonrpc.Message.create_request(
method, [self._db.name, str(self.uuid), monitor_requests])
self._monitor_request_id = msg.id
self._session.send(msg)
def __send_server_schema_request(self):
self.state = self.IDL_S_SERVER_SCHEMA_REQUESTED
msg = ovs.jsonrpc.Message.create_request(
"get_schema", [self._server_db_name, str(self.uuid)])
self._server_schema_request_id = msg.id
self._session.send(msg)
def __send_server_monitor_request(self):
self.state = self.IDL_S_SERVER_MONITOR_REQUESTED
monitor_requests = {}
table = self.server_tables[self._server_db_table]
columns = [column for column in table.columns.keys()]
for column in table.columns.values():
if not hasattr(column, 'alert'):
column.alert = True
table.rows = custom_index.IndexedRows(table)
table.need_table = False
table.idl = self
monitor_request = {"columns": columns}
monitor_requests[table.name] = [monitor_request]
msg = ovs.jsonrpc.Message.create_request(
'monitor', [self._server_db.name,
str(self.server_monitor_uuid),
monitor_requests])
self._server_monitor_request_id = msg.id
self._session.send(msg)
def __parse_update(self, update, version, tables=None):
try:
if not tables:
self.__do_parse_update(update, version, self.tables)
else:
self.__do_parse_update(update, version, tables)
except error.Error as e:
vlog.err("%s: error parsing update: %s"
% (self._session.get_name(), e))
def __do_parse_update(self, table_updates, version, tables):
if not isinstance(table_updates, dict):
raise error.Error("<table-updates> is not an object",
table_updates)
notices = []
for table_name, table_update in table_updates.items():
table = tables.get(table_name)
if not table:
raise error.Error('<table-updates> includes unknown '
'table "%s"' % table_name)
if not isinstance(table_update, dict):
raise error.Error('<table-update> for table "%s" is not '
'an object' % table_name, table_update)
for uuid_string, row_update in table_update.items():
if not ovs.ovsuuid.is_valid_string(uuid_string):
raise error.Error('<table-update> for table "%s" '
'contains bad UUID "%s" as member '
'name' % (table_name, uuid_string),
table_update)
uuid = ovs.ovsuuid.from_string(uuid_string)
if not isinstance(row_update, dict):
raise error.Error('<table-update> for table "%s" '
'contains <row-update> for %s that '
'is not an object'
% (table_name, uuid_string))
if version == OVSDB_UPDATE2:
changes = self.__process_update2(table, uuid, row_update)
if changes:
notices.append(changes)
self.change_seqno += 1
continue
parser = ovs.db.parser.Parser(row_update, "row-update")
old = parser.get_optional("old", [dict])
new = parser.get_optional("new", [dict])
parser.finish()
if not old and not new:
raise error.Error('<row-update> missing "old" and '
'"new" members', row_update)
changes = self.__process_update(table, uuid, old, new)
if changes:
notices.append(changes)
self.change_seqno += 1
for notice in notices:
self.notify(*notice)
def __process_update2(self, table, uuid, row_update):
"""Returns Notice if a column changed, False otherwise."""
row = table.rows.get(uuid)
if "delete" in row_update:
if row:
del table.rows[uuid]
return Notice(ROW_DELETE, row)
else:
# XXX rate-limit
vlog.warn("cannot delete missing row %s from table"
"%s" % (uuid, table.name))
elif "insert" in row_update or "initial" in row_update:
if row:
vlog.warn("cannot add existing row %s from table"
" %s" % (uuid, table.name))
del table.rows[uuid]
row = self.__create_row(table, uuid)
if "insert" in row_update:
row_update = row_update['insert']
else:
row_update = row_update['initial']
self.__add_default(table, row_update)
changed = self.__row_update(table, row, row_update)
table.rows[uuid] = row
if changed:
return Notice(ROW_CREATE, row)
elif "modify" in row_update:
if not row:
raise error.Error('Modify non-existing row')
old_row = self.__apply_diff(table, row, row_update['modify'])
return Notice(ROW_UPDATE, row, Row(self, table, uuid, old_row))
else:
raise error.Error('<row-update> unknown operation',
row_update)
return False
def __process_update(self, table, uuid, old, new):
"""Returns Notice if a column changed, False otherwise."""
row = table.rows.get(uuid)
changed = False
if not new:
# Delete row.
if row:
del table.rows[uuid]
return Notice(ROW_DELETE, row)
else:
# XXX rate-limit
vlog.warn("cannot delete missing row %s from table %s"
% (uuid, table.name))
elif not old:
# Insert row.
op = ROW_CREATE
if not row:
row = self.__create_row(table, uuid)
changed = True
else:
# XXX rate-limit
op = ROW_UPDATE
vlog.warn("cannot add existing row %s to table %s"
% (uuid, table.name))
changed |= self.__row_update(table, row, new)
if op == ROW_CREATE:
table.rows[uuid] = row
if changed:
return Notice(ROW_CREATE, row)
else:
op = ROW_UPDATE
if not row:
row = self.__create_row(table, uuid)
changed = True
op = ROW_CREATE
# XXX rate-limit
vlog.warn("cannot modify missing row %s in table %s"
% (uuid, table.name))
changed |= self.__row_update(table, row, new)
if op == ROW_CREATE:
table.rows[uuid] = row
if changed:
return Notice(op, row, Row.from_json(self, table, uuid, old))
return False
def __check_server_db(self):
"""Returns True if this is a valid server database, False otherwise."""
session_name = self.session_name()
if self._server_db_table not in self.server_tables:
vlog.info("%s: server does not have %s table in its %s database"
% (session_name, self._server_db_table,
self._server_db_name))
return False
rows = self.server_tables[self._server_db_table].rows
database = None
for row in rows.values():
if self.cluster_id:
if self.cluster_id in \
map(lambda x: str(x)[:4], row.cid):
database = row
break
elif row.name == self._db.name:
database = row
break
if not database:
vlog.info("%s: server does not have %s database"
% (session_name, self._db.name))
return False
if database.model == CLUSTERED:
if not database.schema:
vlog.info('%s: clustered database server has not yet joined '
'cluster; trying another server' % session_name)
return False
if not database.connected:
vlog.info('%s: clustered database server is disconnected '
'from cluster; trying another server' % session_name)
return False
if (self.leader_only and
not database.leader):
vlog.info('%s: clustered database server is not cluster '
'leader; trying another server' % session_name)
return False
if database.index:
if database.index[0] < self._min_index:
vlog.warn('%s: clustered database server has stale data; '
'trying another server' % session_name)
return False
self._min_index = database.index[0]
return True
def __column_name(self, column):
if column.type.key.type == ovs.db.types.UuidType:
return ovs.ovsuuid.to_json(column.type.key.type.default)
else:
return column.type.key.type.default
def __add_default(self, table, row_update):
for column in table.columns.values():
if column.name not in row_update:
if ((table.name not in self.readonly) or
(table.name in self.readonly) and
(column.name not in self.readonly[table.name])):
if column.type.n_min != 0 and not column.type.is_map():
row_update[column.name] = self.__column_name(column)
def __apply_diff(self, table, row, row_diff):
old_row = {}
for column_name, datum_diff_json in row_diff.items():
column = table.columns.get(column_name)
if not column:
# XXX rate-limit
vlog.warn("unknown column %s updating table %s"
% (column_name, table.name))
continue
try:
datum_diff = data.Datum.from_json(column.type, datum_diff_json)
except error.Error as e:
# XXX rate-limit
vlog.warn("error parsing column %s in table %s: %s"
% (column_name, table.name, e))
continue
old_row[column_name] = row._data[column_name].copy()
datum = row._data[column_name].diff(datum_diff)
if datum != row._data[column_name]:
row._data[column_name] = datum
return old_row
def __row_update(self, table, row, row_json):
changed = False
for column_name, datum_json in row_json.items():
column = table.columns.get(column_name)
if not column:
# XXX rate-limit
vlog.warn("unknown column %s updating table %s"
% (column_name, table.name))
continue
try:
datum = data.Datum.from_json(column.type, datum_json)
except error.Error as e:
# XXX rate-limit
vlog.warn("error parsing column %s in table %s: %s"
% (column_name, table.name, e))
continue
if datum != row._data[column_name]:
row._data[column_name] = datum
if column.alert:
changed = True
else:
# Didn't really change but the OVSDB monitor protocol always
# includes every value in a row.
pass
return changed
def __create_row(self, table, uuid):
data = {}
for column in table.columns.values():
data[column.name] = ovs.db.data.Datum.default(column.type)
return Row(self, table, uuid, data)
def __error(self):
self._session.force_reconnect()
def __txn_abort_all(self):
while self._outstanding_txns:
txn = self._outstanding_txns.popitem()[1]
txn._status = Transaction.TRY_AGAIN
def __txn_process_reply(self, msg):
txn = self._outstanding_txns.pop(msg.id, None)
if txn:
txn._process_reply(msg)
return True
def _uuid_to_row(atom, base):
if base.ref_table:
return base.ref_table.rows.get(atom)
else:
return atom
def _row_to_uuid(value):
if isinstance(value, Row):
return value.uuid
else:
return value
@functools.total_ordering
class Row(object):
"""A row within an IDL.
The client may access the following attributes directly:
- 'uuid': a uuid.UUID object whose value is the row's database UUID.
- An attribute for each column in the Row's table, named for the column,
whose values are as returned by Datum.to_python() for the column's type.
If some error occurs (e.g. the database server's idea of the column is
different from the IDL's idea), then the attribute values is the
"default" value return by Datum.default() for the column's type. (It is
important to know this because the default value may violate constraints
for the column's type, e.g. the default integer value is 0 even if column
contraints require the column's value to be positive.)
When a transaction is active, column attributes may also be assigned new
values. Committing the transaction will then cause the new value to be
stored into the database.
*NOTE*: In the current implementation, the value of a column is a *copy*
of the value in the database. This means that modifying its value
directly will have no useful effect. For example, the following:
row.mycolumn["a"] = "b" # don't do this
will not change anything in the database, even after commit. To modify
the column, instead assign the modified column value back to the column:
d = row.mycolumn
d["a"] = "b"
row.mycolumn = d
"""
def __init__(self, idl, table, uuid, data):
# All of the explicit references to self.__dict__ below are required
# to set real attributes with invoking self.__getattr__().
self.__dict__["uuid"] = uuid
self.__dict__["_idl"] = idl
self.__dict__["_table"] = table
# _data is the committed data. It takes the following values:
#
# - A dictionary that maps every column name to a Datum, if the row
# exists in the committed form of the database.
#
# - None, if this row is newly inserted within the active transaction
# and thus has no committed form.
self.__dict__["_data"] = data
# _changes describes changes to this row within the active transaction.
# It takes the following values:
#
# - {}, the empty dictionary, if no transaction is active or if the
# row has yet not been changed within this transaction.
#
# - A dictionary that maps a column name to its new Datum, if an
# active transaction changes those columns' values.
#
# - A dictionary that maps every column name to a Datum, if the row
# is newly inserted within the active transaction.
#
# - None, if this transaction deletes this row.
self.__dict__["_changes"] = {}
# _mutations describes changes to this row to be handled via a
# mutate operation on the wire. It takes the following values:
#
# - {}, the empty dictionary, if no transaction is active or if the
# row has yet not been mutated within this transaction.
#
# - A dictionary that contains two keys:
#
# - "_inserts" contains a dictionary that maps column names to
# new keys/key-value pairs that should be inserted into the
# column
# - "_removes" contains a dictionary that maps column names to
# the keys/key-value pairs that should be removed from the
# column
#
# - None, if this transaction deletes this row.
self.__dict__["_mutations"] = {}
# A dictionary whose keys are the names of columns that must be
# verified as prerequisites when the transaction commits. The values
# in the dictionary are all None.
self.__dict__["_prereqs"] = {}
def __lt__(self, other):
if not isinstance(other, Row):
return NotImplemented
return bool(self.__dict__['uuid'] < other.__dict__['uuid'])
def __eq__(self, other):
if not isinstance(other, Row):
return NotImplemented