diff options
| -rwxr-xr-x | mail/deliver | 170 |
1 files changed, 112 insertions, 58 deletions
diff --git a/mail/deliver b/mail/deliver index 3ae138386..29db6b656 100755 --- a/mail/deliver +++ b/mail/deliver @@ -16,70 +16,47 @@ # along with this program; if not, write to the Free Software # Foundation, Inc., 59 Temple Place - Suite 330, Boston, MA 02111-1307, USA. -#import sys -#sys.stderr = open("/home/mailman/mailman/logs/debug","a+") -#sys.stdout = sys.stderr -"""Partition a mass delivery into a suitable set of subdeliveries.""" +"""Partition a mass delivery into a suitable set of subdeliveries. + +The script takes the following arguments: + + argv[1] path to file containing message + argv[2] sender + argv[3] number of processes over which to distribute the sends + argv[4:] recipients + +The recipients will be distributed into at most argv[3] batches, grouping +together recipients in the same major domain as much as possible. +""" # Heh, heh, heh, this partition reminds me of the knapsack problem ;-) # Ie, the optimal distribution here is NP Complete. -import os +import os, sys + +TRIES = 5 +REFRACT = 15 # Seconds between fork retries. -if not os.fork(): - import string, sys, regsub +## # Debugging - os.fork can run out or resources, double check the provisions. +## import paths +## from Mailman.Utils import StampedLogger +## try: +## sys.stderr = StampedLogger("error", label = 'deliver', +## manual_reprime=1, nofail=0) +## sys.stdout = sys.stderr +## except IOError: +## pass # Oh well - SOL on redirect, errors show thru. + +def main(): + if not forker(): + do_child() + os._exit(0) + +def do_child(): + import string, sys, re import paths from Mailman import mm_cfg - def ContactTransport(sender, recip, text): - cmd = os.path.join(mm_cfg.SCRIPTS_DIR, "contact_transport") - file = os.popen(string.join([mm_cfg.PYTHON,cmd,sender]+recip," "), 'w') - file.write(text) - file.close() domain_info = {} - def GroupByDomain(addr): - "Collect addrs by major subdomain - e.g. the 'python' in python.org." - parts = regsub.split(addr, '[.@]') - key = string.join(parts[-2:]) - if not domain_info.has_key(key): - domain_info[key] = [addr] - else: - domain_info[key].append(addr) - - def BuildGroups(biglist, num_addrs): - biglist.sort(lambda x,y: len(x) < len(y)) - groups = [] - for i in range(spawns-1): - target_size = num_addrs / (spawns - i) - if not len(biglist): - break - newlist = biglist[0] - biglist.remove(biglist[0]) - j = 0 - while len(newlist) < target_size: - if j >= len(biglist): - break - if len(newlist) + len(biglist[j]) > target_size: - j = j + 1 - continue - newlist = newlist + biglist[j] - biglist.remove(biglist[j]) - groups.append(newlist) - num_adders = num_addrs - len(newlist) - lastgroup = [] - for item in biglist: - lastgroup = lastgroup + item - if len(lastgroup): - groups.append(lastgroup) - return groups - - def ContactTransportForEachGroup(sender, groups, text): - if len(groups) == 1: - ContactTransport(sender,groups[0],text) - return - for group in groups: - if not os.fork(): - ContactTransport(sender,group,text) - os._exit(0) sender = sys.argv[2] spawns = eval(sys.argv[3]) @@ -92,7 +69,84 @@ if not os.fork(): spawns = 1 to_list = sys.argv[4:] - map(GroupByDomain, to_list) - final_groups = BuildGroups(domain_info.values(), len(to_list)) + # Group by domain. + for addr in to_list: + parts = re.split('[.@]', addr) + key = string.join(parts[-2:]) + if not domain_info.has_key(key): + domain_info[key] = [addr] + else: + domain_info[key].append(addr) + final_groups = BuildGroups(domain_info.values(), len(to_list), spawns) ContactTransportForEachGroup(sender, final_groups, text) os.unlink(sys.argv[1]) + +def BuildGroups(biglist, num_addrs, spawns): + biglist.sort(lambda x,y: len(x) < len(y)) + groups = [] + for i in range(spawns-1): + target_size = num_addrs / (spawns - i) + if not len(biglist): + break + newlist = biglist[0] + biglist.remove(biglist[0]) + j = 0 + while len(newlist) < target_size: + if j >= len(biglist): + break + if len(newlist) + len(biglist[j]) > target_size: + j = j + 1 + continue + newlist = newlist + biglist[j] + biglist.remove(biglist[j]) + groups.append(newlist) + num_adders = num_addrs - len(newlist) + lastgroup = [] + for item in biglist: + lastgroup = lastgroup + item + if len(lastgroup): + groups.append(lastgroup) + return groups + +def ContactTransport(sender, recip, text): + import string + from Mailman import mm_cfg + cmd = os.path.join(mm_cfg.SCRIPTS_DIR, "contact_transport") + file = os.popen(string.join([mm_cfg.PYTHON,cmd,sender]+recip," "), 'w') + file.write(text) + file.close() + +def ContactTransportForEachGroup(sender, groups, text): + if len(groups) == 1: + ContactTransport(sender,groups[0],text) + return + for group in groups: + if not forker(): + ContactTransport(sender,group,text) + os._exit(0) + +def forker(tries=TRIES, refract=REFRACT): + """Fork, retrying on EGAIN errors with refract secs pause between tries. + + Returns value of os.fork(), or raises the exception for: + (1) non-EAGAIN exception, or + (2) EGAIN exception encountered more than tries times.""" + got = 0 + # Loop until we successfully fork or the number tries is exceeded. + while 1: + try: + got = os.fork() + break + except os.error, val: + import errno, time + if val[0] == errno.EAGAIN and tries > 0: + # Resource temporarily unavailable - give time to recover. + tries = tries - 1 + time.sleep(refract) + else: + # No go - reraise original exception, same stack frame and all. + raise val, None, sys.exc_info()[2] + return got + +if __name__ == "__main__": + main() |
