2121import dns .resolver
2222import dns .exception
2323
24+ from multiprocessing import Process as Proc
25+ from multiprocessing import Queue
26+
2427from pubsublogger import publisher
2528from Helper import Process
2629
2730from pyfaup .faup import Faup
2831
29- ## REGEX TIMEOUT ##
30- import signal
31-
32- def timeout_handler (signum , frame ):
33- raise TimeoutException ()
34-
35- class TimeoutException (Exception ):
36- pass
37-
38-
39- signal .signal (signal .SIGALRM , timeout_handler )
40- max_execution_time = 20
41- ## -- ##
42-
4332sys .path .append (os .path .join (os .environ ['AIL_BIN' ], 'packages' ))
4433import Item
4534
@@ -55,6 +44,7 @@ class TimeoutException(Exception):
5544
5645config_loader = None
5746## -- ##
47+
5848def is_mxdomain_in_cache (mxdomain ):
5949 return r_serv_cache .exists ('mxdomain:{}' .format (mxdomain ))
6050
@@ -120,6 +110,9 @@ def check_mx_record(set_mxdomains, dns_server):
120110 print (e )
121111 return valid_mxdomain
122112
113+ def extract_all_emails (queue , item_content ):
114+ queue .put (re .findall (email_regex , item_content ))
115+
123116if __name__ == "__main__" :
124117 publisher .port = 6380
125118 publisher .channel = "Script"
@@ -135,32 +128,37 @@ def check_mx_record(set_mxdomains, dns_server):
135128 # Numbers of Mails needed to Tags
136129 mail_threshold = 10
137130
131+ max_execution_time = 30
132+
138133 email_regex = "[a-zA-Z0-9._%+-]+@[a-zA-Z0-9.-]+\.[a-zA-Z]{2,6}"
139134
135+ q = Queue ()
136+
140137 while True :
141138 message = p .get_from_set ()
142139
143140 if message is not None :
144141 item_id , score = message .split ()
145142
146143 item_content = Item .get_item_content (item_id )
147- item_date = Item .get_item_date (item_id )
148144
149- #print(item_id)
150-
151- # Get all emails address
152- signal .alarm (30 )
145+ proc = Proc (target = extract_all_emails , args = (q , item_content ))
146+ proc .start ()
153147 try :
154- all_emails = re .findall (email_regex , item_content )
155- except TimeoutException :
156- p .incr_module_timeout_statistic ()
157- err_mess = "Mails: processing timeout: {}" .format (item_id )
158- print (err_mess )
159- publisher .info (err_mess )
160- signal .signal (signal .SIGALRM , timeout_handler )
161- continue
162- finally :
163- signal .alarm (0 )
148+ proc .join (max_execution_time )
149+ if proc .is_alive ():
150+ proc .terminate ()
151+ p .incr_module_timeout_statistic ()
152+ err_mess = "Mails: processing timeout: {}" .format (item_id )
153+ print (err_mess )
154+ publisher .info (err_mess )
155+ continue
156+ else :
157+ all_emails = q .get ()
158+ except KeyboardInterrupt :
159+ print ("Caught KeyboardInterrupt, terminating workers" )
160+ proc .terminate ()
161+ sys .exit (0 )
164162
165163 # filtering duplicate
166164 all_emails = set (all_emails )
@@ -179,6 +177,8 @@ def check_mx_record(set_mxdomains, dns_server):
179177
180178 valid_mx = check_mx_record (set_mxdomains , dns_server )
181179
180+ item_date = Item .get_item_date (item_id )
181+
182182 num_valid_email = 0
183183 for domain_mx in valid_mx :
184184 num_valid_email += len (dict_mxdomains_email [domain_mx ])
0 commit comments