-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathreplicator.py
More file actions
2204 lines (1999 loc) · 90.6 KB
/
Copy pathreplicator.py
File metadata and controls
2204 lines (1999 loc) · 90.6 KB
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
# Perforce Defect Tracking Integration Project
# <http://www.ravenbrook.com/project/p4dti/>
#
# REPLICATOR.PY -- P4DTI REPLICATOR
#
# Gareth Rees, Ravenbrook Limited, 2000-08-09
#
#
# 1. INTRODUCTION
#
# This Python module implements the P4DTI replicator: the component of
# the P4DTI that copies data from the defect tracker to Perforce and
# vice versa [RB 2000-08-10], in order to keep the defect tracker state
# consistent with the Perforce state [Requirements, 1] and to provide
# the ability to ask questions involving both the defect tracking system
# and Perforce [Requirements, 5].
#
# The replicator is independent of any particular defect tracker: it
# interacts with the defect tracker through the abstract interfaces
# declared in the dt_interface module and documented in [GDR 2000-10-16,
# 7]. This is to make it possible to integrate Perforce with new defect
# tracking systems [Requirements, 20, 21] and to simplify the design of
# the replicator so that it is modifiable [Requirements, 25], stable
# [Requirements, 27] and maintainable [Requirements, 30].
#
# See [GDR 2000-09-13] for the design of the replicator, its algorithms,
# and the specification of the data it stores in Perforce.
#
# The intended readership of this document is project developers.
#
# This document is not confidential.
import catalog
import dt_interface
import message
import p4
import re
import smtplib
import string
import sys
import time
import stacktrace
import types
# 2. CURSOR WRAPPER FOR LISTS
#
# list_cursor is a class that wraps up a list as a cursor with a
# fetchone() method.
#
# This class is used because the specification of the 'all_issues' and
# 'changed_entities' methods in the defect_tracker class has changed
# since P4DTI 1.1.1 was released. In release 1.1.1 they were documented
# to return lists of issues. Now they are documented to return cursors.
# We still want to support people who wrote code that returned lists, so
# this module examines the results of these methods and wraps lists with
# this class.
class list_cursor:
def __init__(self, list):
self.list = list
def fetchone(self):
if self.list:
result = self.list[0]
self.list = self.list[1:]
return result
else:
return None
# 3. DEFECT TRACKER INTERFACE TO PERFORCE
#
# The replicator attempts to be as symmetric as possible, for simplicity
# of design. It treats Perforce as much as possible in the same way as
# the defect tracker.
#
# It should be possible to develop a dt_perforce module that implements
# a full defect tracker interface to Perforce. However, there will need
# to be some changes to the abstract interface [GDR 2000-10-16, 7]
# because the situation is not 100% symmetric (for example, changelists
# are only replicated in one direction).
#
# We haven't had the time to develop the revised interface and the
# implementation, so for the moment the dt_perforce class is a
# placeholder. Eventually it will be fully functional and take over all
# Perforce operations from the replicator, which can then be simplified
# and made fully symmetric.
class dt_perforce(dt_interface.defect_tracker):
error = "Perforce interface error"
p4 = None
def __init__(self, p4_interface, config):
assert isinstance(p4_interface, p4.p4)
self.p4 = p4_interface
# 4. REPLICATOR
class replicator:
# 4.1. Data
# Configuration module.
config = None
# Defect tracker.
dt = None
# Defect tracker interface to Perforce. This is a placeholder; at
# the moment the implementation is incomplete, but eventually it
# will be there.
dt_p4 = None
# Replicator identifier.
rid = None
# Interface to Perforce.
p4 = None
# The number of columns to format e-mail messages to.
columns = 80
# The replicator's counter on the Perforce server.
counter = None
# Error object for fatal errors raised by the replicator.
error = 'P4DTI Replicator error'
# Map from feature name to whether the defect tracker supports it.
# See [GDR 2000-10-16, 3.5] for names of features.
feature = {}
# 4.2. Initialization
def __init__(self, dt, p4_interface, config):
assert isinstance(dt, dt_interface.defect_tracker)
assert isinstance(p4_interface, p4.p4)
self.job_updates = {}
self.dt = dt
self.config = config
self.rid = config.rid
# Set the log failure hook to send a message to the P4DTI
# administrator.
def log_failed_hook(err, context, r=self):
# Log the failure anyway; you never know if some other
# logger may be able to handle it.
r.config.logger.log(context)
# "Error in P4DTI logger: %s"
r.mail_report(catalog.msg(913, err), [context])
self.config.logger.set_log_failed_hook(log_failed_hook)
# This is incomplete. Eventually there will be a real defect
# tracker interface to Perforce.
self.dt_p4 = dt_perforce(p4_interface, config)
self.p4 = p4_interface
# Replicator ids must match.
if self.rid != self.dt.rid:
# "The replicator's RID ('%s') doesn't match the defect
# tracker's RID ('%s')."
raise self.error, catalog.msg(833, (self.rid, self.dt.rid))
# Make a counter name for this replicator.
if not self.counter:
self.counter = 'P4DTI-%s' % self.rid
# Make a client for the replicator.
self.create_client()
# Initialize the defect tracking system.
self.dt.init()
self.determine_supported_features()
# create_client(). Creates a client if one does not exist, or if
# the existing client is broken.
def create_client(self):
client = self.p4.run('client -o')[0]
try:
self.p4.run('client -i', client)
except p4.error, message:
# "Can't use Perforce client %s."
self.log(927, self.p4.client)
self.log(message)
# "Attempting to make working Perforce client %s."
self.log(928, self.p4.client)
# Strip out all non-essential entries in the client record.
for k in client.keys():
if k[0:4] == 'View' or k in ['LineEnd', 'Options', 'Host', 'code']:
del client[k]
client['Description'] = 'Client created and used by P4DTI replicator.'
self.p4.run('client -i', client)
# determine_supported_features(). Determine which optional features
# are supported by the defect tracker interface.
def determine_supported_features(self):
self.feature = {}
if hasattr(self.dt, 'supports'):
# Query the defect tracker's supports() method if available.
for feature in ['filespecs', 'fixes', 'migrate_issues',
'new_issues', 'new_users']:
self.feature[feature] = self.dt.supports(feature)
else:
# Otherwise check to see if the defect tracker and the
# configuration offer the required interface, as promised in
# [GDR 2000-10-16, 13.4].
self.feature['filespecs'] = 1
self.feature['fixes'] = 1
self.feature['migrate_issues'] = 1
for m in ['new_issue', 'new_issues_start',
'new_issues_end']:
if not hasattr(self.dt, m):
self.feature['migrate_issues'] = 0
if not hasattr(self.config, 'prepare_issue_advanced'):
self.feature['migrate_issues'] = 0
self.feature['new_issues'] = self.feature['migrate_issues']
if not hasattr(self.config, 'translate_jobspec_advanced'):
self.feature['migrate_issues'] = 0
self.feature['new_users'] = hasattr(self.dt, 'add_user')
# check_first_time(). Take a look at the old jobspec. If it has no
# P4DTI fields, then we assume that this is the first time the P4DTI
# has been run. If so, check for the existence of jobs; if there
# are any, don't go ahead with the change to the jobspec but instead
# warn the administrator. See job000219.
def check_first_time(self):
if (self.config.jobspec
and not self.p4.jobspec_has_p4dti_fields(
self.p4.get_jobspec(),
warn = 0)
and self.p4.run('jobs')):
# "You must delete your Perforce jobs before running the
# P4DTI for the first time. See section 5.2.3 of the
# Administrator's Guide."
raise self.error, catalog.msg(914)
# update_and_check_jobspec(). If keep_jobspec is set, check the
# current jobspec against the one we want to install. Otherwise,
# just go ahead and install the jobspec. Advanced configurations
# can turn this off by clearing config.jobspec [IG, 8.6].
#
# Subsequently, the installed jobspec is checked to ensure that we
#can run the P4DTI with it.
def update_and_check_jobspec(self):
if self.config.jobspec and not self.config.keep_jobspec:
self.p4.install_jobspec(self.config.jobspec)
self.check_jobspec()
def check_jobspec(self):
self.p4.check_jobspec(self.config.jobspec)
def extend_jobspec(self, force=0):
self.p4.extend_jobspec(self.config.jobspec, force)
# start_logger(). Has the logger been started? If not, start it.
# (We must be careful not to set the logger counter to 0 more than
# once; this will confuse Perforce [Seiwald 2000-09-11].)
def start_logger(self):
counters = self.p4.run('counters')
for c in counters:
if c.get('counter') == 'logger':
return
self.p4.run('counter logger 0')
# 4.3. Logging
# log(msg, args = ()). Write the message to the replicator's log.
# The msg argument can be a message instance, the number of a
# message in the catalog (with the remaining arguments used to fill
# in the message parameters), or a string.
def log(self, msg, args = ()):
if isinstance(msg, message.message):
self.config.logger.log(msg)
elif isinstance(msg, types.IntType):
self.config.logger.log(catalog.msg(msg, args))
else:
# "%s"
self.config.logger.log(catalog.msg(910, msg.encode('utf8')))
# 4.4. Perforce interface
#
# These methods provide the replicator with an interface to
# Perforce. In time they can be moved to the dt_perforce class
# (section 3).
# all_jobs(). Return a list of all jobs.
def all_jobs(self):
return self.p4.run('jobs')
# changed_entities(). Return a 3-tuple consisting of (a) changed
# jobs, (b) changed changelists, and (c) the last log entry that
# was considered. The changed jobs are those that are due for
# replication by this replicator (that is, the P4DTI-rid field of
# the job matches the replicator id), or new jobs which pass the
# replicate_job_p check. The last log entry will be passed to
# mark_changes_done.
def changed_entities(self):
# Get all entries from the log since the last time we updated
# the counter.
log_entries = self.p4.run('logger -t %s' % self.counter)
jobs = {}
changelists = []
last_log_entry = None # The last entry number in the log.
for e in log_entries:
last_log_entry = int(e['sequence'])
if e['key'] == 'job':
jobname = e['attr']
# Can we account for this log entry on the basis of
# updates we made in the previous poll? If so, ignore
# the entry.
if self.job_updates.get(jobname):
n_updates = self.job_updates[jobname]
self.job_updates[jobname] = n_updates - 1
elif jobname == 'new':
# "Perforce has a job called 'new', which is
# illegal and will stop the P4DTI from working."
raise self.error, catalog.msg(896)
elif not jobs.has_key(jobname):
job = self.job(jobname)
p4dti_rid = job.get('P4DTI-rid', 'None')
if (p4dti_rid == self.rid
or (p4dti_rid == 'None'
and self.config.replicate_job_p(job))):
jobs[jobname] = job
elif e['key'] == 'change':
# Collect new and updated changelists here. A
# changelist can change (using p4 change -f) without any
# related jobs changing, so we need to replicate
# changelists as well as replicating the fixes of
# changed jobs.
change_number = e['attr']
try:
changelist = self.p4.run('change -o %s'
% change_number)[0]
changelists.append(changelist)
except p4.error:
# The changelist might not exist any more: it might
# have been a pending changelist that's been
# renumbered. So don't replicate it. Should it be
# deleted from the defect tracker? GDR 2000-11-02.
pass
self.job_updates = {}
return jobs, changelists, last_log_entry
# mark_changes_done(log_entry). Update the Perforce database to
# record the fact that the replicator has replicated all changes up
# to log_entry.
def mark_changes_done(self, log_entry):
assert log_entry == None or isinstance(log_entry, types.IntType)
# Update counter to last entry number in the log that we've
# replicated. If this is the last entry in the log, it has the
# side-effect of deleting the log (see "p4 help undoc").
if log_entry:
self.p4.run('logger -t %s -c %d'
% (self.counter, log_entry))
# clear_logger(). Clear the logger.
def clear_logger(self):
last_log_entry = self.p4.counter_value('logger')
self.p4.run('logger -t %s -c %s'
% (self.counter, last_log_entry))
# job(jobname). Return the Perforce job with the given name if it
# exists, or an empty job specification (otherwise).
def job(self, jobname):
assert isinstance(jobname, basestring)
jobs = self.p4.run('job -o %s' % jobname)
if len(jobs) != 1 or not jobs[0].has_key('Job'):
# "Expected a job but found %s."
raise self.error, catalog.msg(837, str(jobs))
# Compare job names case-insensitively (see job000313).
elif string.lower(jobs[0]['Job']) != string.lower(jobname):
# "Asked for job '%s' but got job '%s'."
raise self.error, catalog.msg(838, (jobname,
jobs[0]['Job']))
else:
return jobs[0]
# job_filespecs(job). Return a list of filespecs for the given job.
# Each element of the list is a filespec, as a string.
def job_filespecs(self, job):
assert isinstance(job, types.DictType)
# if no P4DTI-filespecs field, do the right thing:
job_filespecs = job.get('P4DTI-filespecs', '')
filespecs = string.split(job_filespecs, '\n')
# Since Perforce text fields are terminated with a newline, the
# last item of the list must be an empty string. Remove it.
if filespecs[-1] != '':
# "P4DTI-filespecs field has value '%s': this should end
# in a newline."
raise self.error, catalog.msg(839, job_filespecs)
return filespecs[:-1]
# job_fixes(job). Return a list of fixes for the given job. Each
# element of the list is a dictionary with keys Change, Client,
# User, Job, and Status.
def job_fixes(self, job):
assert isinstance(job, types.DictType)
return self.p4.run('fixes -j %s' % job['Job'])
# job_format(job). Format a job so that people can read it. Also,
# indent the first line of the job so that it can be included in the
# body of a mail message without being wrapped; see mail().
def job_format(self, job):
def format_item(i):
key, value = i
if '\n' in value:
if value[-1] == '\n':
value = value[0:-1]
value = string.join(string.split(value,'\n'),'\n\t')
return "%s:\n\t%s" % (key, value)
else:
return "%s: %s" % (key, value)
items = job.items()
# Remove special Perforce system fields.
items = filter(lambda i: i[0] not in ['code','specdef'], items)
# Sort into lexical order.
items.sort()
return string.join(map(format_item, items), '\n')
# job_mail_recipients(job). Work out the people associated with the
# job who should receive e-mail when there's a problem with that job
# (namely the job's owner and the last person to edit the job,
# unless either of these is the replicator).
def job_mail_recipients(self, job):
recipients = []
# Owner of the job, if any.
owner = job.get(self.config.job_owner_field, None)
if owner:
owner_address = self.user_email_address(owner)
if owner_address:
# "Job owner"
comment = catalog.msg(925)
recipients.append((comment, owner_address))
# Last person to change the job, if neither replicator nor
# owner.
changer = job.get('P4DTI-user', None)
if changer and changer not in [self.config.p4_user, owner]:
changer_address = self.user_email_address(changer)
if changer_address:
# "Job changer"
comment = catalog.msg(926)
recipients.append((comment, changer_address))
return recipients
# job_modifier(job). Return our best guess at who last modified the
# job.
#
# If the Perforce server supports the 'fix_update' feature, then
# this is easy: just return the P4DTI-user field (for safety in the
# case of race conditions -- user changes interleaved with P4DTI
# changes -- we substitute the job owner if P4DTI-user is the
# replicator).
#
# However, in Perforce servers that don't support this feature, the
# "always" fields in a job don't get modified when a job is fixed.
# This means that the P4DTI-user field may not be accurate, since
# there may have been fixes added later.
#
# So our strategy for finding the owner is as follows:
#
# 1. Is there a fix record, submitted more recently than the job
# has been modified, by someone other than the replicator? If so,
# take the person who submitted the most recent such fix as the
# modifier.
#
# 2. If not, does the P4DTI-user field contain a user other than
# the replicator? If so, take them as the modifier.
#
# 3. If not, take the job owner as the modifier.
#
# Note that this doesn't give a 100% accurate answer (for example,
# if you fix a job and then delete the fix), but it's right in all
# but a few exceptional cases.
#
# See job000133 and job000270 for the motivation.
def job_modifier(self, job):
modifier = job.get('P4DTI-user', self.config.p4_user)
if modifier == self.config.p4_user:
modifier = job.get(self.config.job_owner_field, modifier)
# Perforce 2002.1 updates 'always' fields when someone modifies
# a job by making a fix. So the modifier is accurate.
if self.p4.supports('fix_update'):
return modifier
# Dates in job fields look like 2000/12/31 23:59:59, but dates
# in fixes are seconds since 1970-01-01 00:00:00, so convert the
# job modification time to an integer for comparison.
match = re.match('^(\d{4})/(\d{2})/(\d{2}) '
'(\d{2}):(\d{2}):(\d{2})$',
job.get(self.config.job_date_field,
'1970/01/01 00:00:00'))
if not match:
# "Job '%s' has a date field in the wrong format: %s."
raise self.error, catalog.msg(889, (job['Job'], job))
date = time.mktime(tuple(map(int, match.groups()) + [0,0,-1]))
fixes = self.job_fixes(job)
for f in fixes:
if (int(f['Date']) > date
and f['User'] != self.config.p4_user):
modifier = f['User']
date = int(f['Date'])
return modifier
# Map from Perforce job name to the number of times we've updated
# the job in this poll. Used by changed_entities and updated by
# update_job and replicate_fixes_dt_to_p4.
job_updates = {}
def record_job_update(self, job):
jobname = job['Job']
self.job_updates[jobname] = self.job_updates.get(jobname, 0) + 1
# update_job(job, changes). Update the job in Perforce by applying
# the given changes. Also update the "job" dictionary to reflect
# these changes, and also any changes made by Perforce, such as
# picking up the new jobname (if job['Job'] is 'new').
update_job_re = re.compile('^Job ([^ ]+) (.*)')
def update_job(self, job, changes = {}, force = False):
assert isinstance(job, types.DictType)
assert isinstance(changes, types.DictType)
for key, value in changes.items():
job[key] = value
if force:
command = 'job -i -f'
else:
command = 'job -i'
results = self.p4.run(command, job)
# Check that the results of the 'job -i' command are as
# expected: Perforce should say something like 'Job job012345
# saved.' or 'Job job012345 not changed.' If the jobname was
# 'new', then record the jobname that Perforce gave the new
# job so that we can call setup_for_replication() in
# replicate_many().
if len(results) == 1 and results[0].has_key('data'):
match = self.update_job_re.match(results[0]['data'])
if not match or match.group(1) == 'new':
# "Expected Perforce output of 'job -i' to say 'Job
# jobname ...', but found '%s'."
raise self.error, catalog.msg(897, results[0]['data'])
elif job['Job'] == 'new':
job['Job'] = match.group(1)
elif job['Job'] != match.group(1):
# "Tried to update job '%s', but Perforce replied '%s'."
raise self.error, catalog.msg(899, (job['Job'],
results[0]['data']))
if match.group(2) == 'saved.':
self.record_job_update(job)
else:
# "Unexpected output from Perforce command 'job -i': %s."
raise self.error, catalog.msg(898, results)
# Return the e-mail address of a Perforce user, or None if the
# address can't be found.
def user_email_address(self, user):
assert isinstance(user, basestring)
# Even though "p4 user -o foo" doesn't actually create a user,
# it does fail if the user foo doesn't exist and there are no
# spare licenses. So trap that case. See job000204.
try:
u = self.p4.run('user -o %s' % user)
except:
return None
# Existing users have Access and Update fields in the returned
# structure; non-existing users don't.
if (len(u) == 1 and u[0].has_key('Access')
and u[0].has_key('Update') and u[0].has_key('Email')):
return u[0]['Email']
else:
return None
# 4.5. Entry points
# check_consistency(). Run a consistency check on the two
# databases, reporting any inconsistencies.
def check_consistency(self):
# "Checking consistency for replicator '%s'."
self.log(871, self.rid)
self.check_jobspec()
n = 0 # Number of inconsistencies found.
# Get issues and jobs.
issues_cursor = self.dt.all_issues()
# Support old all_issues specification [GDR 2000-10-16, 13.1].
if not hasattr(issues_cursor, 'fetchone'):
issues_cursor = list_cursor(issues_cursor)
issue_id_to_job = {}
jobs = {}
for j in self.p4.run('jobs -e P4DTI-rid=%s' % self.rid):
jobs[j['Job']] = j
while 1:
issue = issues_cursor.fetchone()
if issue == None:
break
id = issue.id()
jobname = issue.corresponding_id()
# "Checking issue '%s' against job '%s'."
self.log(890, (id, jobname))
# Report if issue has no corresponding job.
if issue.rid() != self.rid:
if self.config.replicate_p(issue):
# "Issue '%s' should be replicated but is not."
self.log(872, id)
n = n + 1
continue
issue_id_to_job[id] = jobname
if not jobs.has_key(jobname):
# "Issue '%s' should be replicated to job '%s' but that
# job either does not exist or is not replicated."
self.log(873, (id, jobname))
n = n + 1
continue
# Get corresponding job.
job = jobs[jobname]
del jobs[jobname]
# Report if mapping is in error.
job_issue_id = job.get('P4DTI-issue-id', 'None')
if job_issue_id != id:
# "Issue '%s' is replicated to job '%s' but that job is
# replicated to issue '%s'."
self.log(874, (id, jobname, job_issue_id))
n = n + 1
# Report if job and issue contents don't match.
changes = self.translate_issue_dt_to_p4(issue, job, 1)
if changes:
# "Job '%s' would need the following set of changes in
# order to match issue '%s': %s."
self.log(875, (jobname, id, str(changes)))
n = n + 1
# Report if filespecs don't match.
n = n + self.check_filespecs(issue, job)
# Report if fixes don't match.
n = n + self.check_fixes(issue, job)
# There should be no remaining jobs, so any left are in error.
for job in jobs.values():
job_issue_id = job.get('P4DTI-issue-id', 'None')
if issue_id_to_job.has_key(job_issue_id):
# "Job '%s' is marked as being replicated to issue '%s'
# but that issue is being replicated to job '%s'."
self.log(881, (job['Job'], job_issue_id,
issue_id_to_job[job_issue_id]))
n = n + 1
else:
# "Job '%s' is marked as being replicated to issue '%s'
# but that issue either doesn't exist or is not being
# replicated by this replicator."
self.log(882, (job['Job'], job_issue_id))
n = n + 1
# Report on success/failure.
if len(issue_id_to_job) == 1:
# "Consistency check completed. 1 issue checked."
self.log(883)
else:
# "Consistency check completed. %d issues checked."
self.log(884, len(issue_id_to_job))
if n == 0:
# "Looks all right to me."
self.log(885)
elif n == 1:
# "1 inconsistency found."
self.log(886)
else:
# "%d inconsistencies found."
self.log(887, n)
# check_filespecs(issue, job). Report if the sets of filespecs
# differ between the issue and the job. Return the number of
# inconsistencies found.
def check_filespecs(self, issue, job):
if not self.feature['filespecs']:
return 0
n = 0 # Number of inconsistencies found.
issuename = issue.readable_name()
jobname = job['Job']
p4_filespecs = self.job_filespecs(job)
dt_filespecs = issue.filespecs()
diffs = self.filespecs_differences(dt_filespecs, p4_filespecs)
for p4_filespec, dt_filespec in diffs:
if p4_filespec and not dt_filespec:
# "Job '%s' has associated filespec '%s' but there is no
# corresponding filespec for issue '%s'."
self.log(876, (jobname, p4_filespec, issuename))
n = n + 1
elif not p4_filespec and dt_filespec:
# "Issue '%s' has associated filespec '%s' but there is
# no corresponding filespec for job '%s'."
self.log(877, (issuename, dt_filespec.name(), jobname))
n = n + 1
else:
# Corresponding filespecs can't differ (since their only
# attribute is their name).
assert 0
return n
# check_fixes(issue, job). Report if the sets of fixes differ
# between the issue and the job. Return the number of
# inconsistencies found.
def check_fixes(self, issue, job):
if not self.feature['fixes']:
return 0
n = 0 # Number of inconsistencies found.
issuename = issue.readable_name()
jobname = job['Job']
p4_fixes = self.job_fixes(job)
dt_fixes = issue.fixes()
diffs = self.fixes_differences(dt_fixes, p4_fixes)
for p4_fix, dt_fix in diffs:
if p4_fix and not dt_fix:
# "Change %s fixes job '%s' but there is no
# corresponding fix for issue '%s'."
self.log(878, (p4_fix['Change'], jobname, issuename))
n = n + 1
elif not p4_fix and dt_fix:
# "Change %d fixes issue '%s' but there is no
# corresponding fix for job '%s'."
self.log(879, (dt_fix.change(), issuename, jobname))
n = n + 1
else:
# "Change %s fixes job '%s' with status '%s', but change
# %d fixes issue '%s' with status '%s'."
self.log(880, (p4_fix['Change'], jobname,
p4_fix['Status'], dt_fix.change(),
issuename, dt_fix.status()))
n = n + 1
return n
# migrate_users() ensures that there is a defect tracker user
# corresponding to each Perforce user.
def migrate_users(self):
if not self.feature['new_users']:
# "Defect tracker '%s' does not support migration of
# Perforce users."
raise self.error, catalog.msg(906, self.config.dt_name)
self.dt.add_replicator_user()
p4_users = self.p4.run("users")
for user in p4_users:
self.dt.add_user(user['User'],
user['Email'],
user['FullName'])
# migrate() migrates all existing Perforce jobs to the DT. Note
# that we can't just call replicate_new_issue_p4_to_dt here because
# that method assumes we have the new jobspec in place (so that we
# can replicate backwards) which we don't until migration is
# finished. Instead, we replicate backwards in a bunch after
# migration succeeds.
def migrate(self, starting_with = None):
if not self.feature['migrate_issues']:
# "Defect tracker '%s' does not support migration of
# Perforce jobs."
raise self.error, catalog.msg(905, self.config.dt_name)
jobs = self.all_jobs()
try:
self.dt.new_issues_start()
for job in jobs:
if starting_with and job['Job'] != starting_with:
continue
else:
starting_with = None
if job.get('P4DTI-rid', 'None') != 'None':
# "Not migrating job '%s' (already replicated)."
self.log(916, job['Job'])
elif not self.config.migrate_p(job):
# "Not migrating job '%s' (migrate_p returned 0)."
self.log(917, job['Job'])
else:
try:
# "Before translating jobspec, job '%s' is %s"
self.log(915, (job['Job'], job))
job = self.config.translate_jobspec_advanced(
self.config, self.dt, self.dt_p4, job)
if not isinstance(job, types.DictType):
# "Expected translate_jobspec to return a
# dictionary, but instead it returned %s."
raise self.error, catalog.msg(924, job)
# "After translating jobspec, job '%s' is %s"
self.log(918, (job['Job'], job))
issue = self.create_issue(job)
except:
# "Migrating job '%s'..."
self.log(921, job['Job'])
raise
# "Migrated job '%s' to issue '%s'."
self.log(892, (job['Job'], issue.readable_name()))
# Replicate filespecs, fixes and changelists.
self.replicate_filespecs_p4_to_dt(issue, job)
self.replicate_fixes_p4_to_dt(issue, job)
finally:
self.dt.new_issues_end()
# "Migration completed."
self.log(895)
# poll(). Poll the defect tracker and Perforce, replicate changes,
# then stop.
def poll(self):
self.check_first_time()
self.update_and_check_jobspec()
self.start_logger()
self.poll_databases()
# refresh_perforce_jobs(). Replicate all issues from the defect
# tracker. Note: does not delete jobs first.
def refresh_perforce_jobs(self):
self.update_and_check_jobspec()
self.replicate_all_dt_to_p4()
self.start_logger()
self.clear_logger()
# carefully_poll_databases(). Poll once, handling exceptions
def carefully_poll_databases(self):
try:
self.poll_databases()
# Reset poll period when the poll was successful.
self.poll_period = self.config.poll_period
except AssertionError:
# Assertions indicate severe bugs in the replicator. It
# might cause serious data corruption if we continue.
# We also want these failures to be reported, and they
# might go unreported if the replicator carried on
# going.
raise
except KeyboardInterrupt:
# Allow people to stop the replicator with Control-C.
raise
except:
self.mail_report(
# "The replicator failed to poll successfully."
catalog.msg(863),
# "The replicator failed to poll successfully,
# because of the following problem:"
[catalog.msg(864)])
# The poll failed; it's likely that it will fail again
# for the same reason the next time we poll. Back off
# exponentially so as not to mail bomb the admin. See
# job000215 and job000135.
self.poll_period = self.poll_period * 2
# prepare_to_run(). Invoked once when run() is called, to preform
# startup tasks.
def prepare_to_run(self):
self.check_first_time()
self.update_and_check_jobspec()
self.start_logger()
self.poll_period = self.config.poll_period
self.mail_startup_message()
# run(). Repeatedly (handling exceptions) poll and replicate
# changes.
def run(self):
self.prepare_to_run()
while 1:
self.carefully_poll_databases()
time.sleep(self.poll_period)
# 4.6. E-mail
# mail(recipients, subject, body). Send e-mail to the given
# recipients (pls the administrator) with the given subject and
# body. The recipients argument is a list of pairs (role, address).
# The body argument is a list of paragraphs. Paragraphs belonging
# to the message.message class will be wrapped to 80 columns.
# Ordinary strings will be left alone. Log the contents of the
# message.
def mail(self, recipients, subject, body):
assert isinstance(recipients, types.ListType)
assert isinstance(subject, message.message)
assert isinstance(body, types.ListType)
# Always e-mail the administrator
recipients.append(('P4DTI administrator',
self.config.administrator_address))
# Build the contents of the RFC822 To: header.
to = string.join(map(lambda r: "%s <%s>" % r, recipients), ', ')
# "Mailing '%s'."
self.log(800, to)
self.log(subject)
map(self.log, body)
# Don't send e-mail if administrator_address or smtp_server is
# None.
if (self.config.administrator_address == None
or self.config.smtp_server == None):
return
smtp = smtplib.SMTP(self.config.smtp_server)
message_paragraphs = [
("From: %s\n"
"To: %s\n"
"Subject: %s"
% (self.config.replicator_address, to, subject)),
# "This is an automatically generated e-mail from the
# Perforce Defect Tracking Integration replicator '%s'."
catalog.msg(865, self.rid),
] + body
def fmt(s, columns = self.columns):
if isinstance(s, message.message):
return s.wrap(columns)
else:
return s.encode('utf8')
message_text = string.join(map(fmt, message_paragraphs), "\n\n")
smtp.sendmail(self.config.replicator_address,
map(lambda r: r[1], recipients),
message_text)
smtp.quit()
# exception_message(exc_info). Return a message object describing
# the given exception, or None if there was no exception. The
# exc_info argument must be the results of calling sys.exc_info().
def exception_message(self, exc_info):
exc_type, exc_value = exc_info[0:2]
if isinstance(exc_value, message.message):
return exc_value
elif exc_type is not None:
# "Error (%s): %s"
return catalog.msg(891, (exc_type, exc_value))
else:
# We're not in the context of an exception, so there's
# nothing to report.
return None
def stacktrace(self, exc_info):
return string.join(apply(stacktrace.format_exception, exc_info),
'')
# mail_report(subject, intro, extra=[], job=None, error=1). Compose
# and send e-mail when something's gone wrong. If a job argument is
# supplied, it's the job to which the mail applies, and is used to
# deduce who to send the e-mail to. If no job argument is supplied,
# then mail is to the administrator (only). Iff error is 1, the
# mail includes an error message and traceback.
def mail_report(self, subject, intro, extra=[], job=None, error=1):
assert isinstance(subject, message.message)
assert isinstance(intro, types.ListType)
assert isinstance(extra, types.ListType)
for m in intro + extra:
assert (isinstance(m, basestring)
or isinstance(m, message.message))
assert job is None or isinstance(job, types.DictType)
if error:
try:
exc_info = sys.exc_info()
msg = self.exception_message(exc_info)
body = intro + [ msg ] + extra + [
# "Here's a full Python traceback:"
catalog.msg(852),
self.stacktrace(exc_info),
]
finally: