summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rwxr-xr-xmail/deliver170
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()