-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathsmtpserver.py
More file actions
168 lines (149 loc) · 5.87 KB
/
Copy pathsmtpserver.py
File metadata and controls
168 lines (149 loc) · 5.87 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
#
# Implementation of simple nonblocking smtp server with tornado
# very much of this code is taken from tornado httpserver.py
#
import errno
import logging
import os
import socket
import time
import email
from tornado import ioloop
from tornado import iostream
try:
import fcntl
except ImportError:
if os.name == 'nt':
import win32_support as fcntl
else:
raise
try:
import multiprocessing
except ImportError:
multiprocessing = None
def _cpu_count():
if multiprocessing is not None:
try:
return multiprocessing.cpu_count()
except NotImplementedError:
pass
try:
return os.sysconf("SC_NPROCESSORS_CONF")
except ValueError:
pass
logging.error("Could not detect number of processors; "
"running with one process")
return 1
class SMTPServer(object):
def __init__(self, request_callback, io_loop=None):
self.request_callback = request_callback
self.io_loop = io_loop
self._socket = None
self._started = False
def listen(self, port, address=""):
self.bind(port,address)
self.start(1)
def bind(self, port, address=""):
assert not self._socket
self._socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM, 0)
flags = fcntl.fcntl(self._socket.fileno(), fcntl.F_GETFD)
flags |= fcntl.FD_CLOEXEC
fcntl.fcntl(self._socket.fileno(), fcntl.F_SETFD, flags)
self._socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
self._socket.setblocking(0)
self._socket.bind((address, port))
self._socket.listen(128)
def start(self, num_processes=1):
"""Starts this server in the IOLoop.
By default, we run the server in this process and do not fork any
additional child process.
If num_processes is None or <= 0, we detect the number of cores
available on this machine and fork that number of child
processes. If num_processes is given and > 1, we fork that
specific number of sub-processes.
Since we use processes and not threads, there is no shared memory
between any server code.
"""
assert not self._started
self._started = True
if num_processes is None or num_processes <= 0:
num_processes = _cpu_count()
if num_processes > 1 and ioloop.IOLoop.initialized():
logging.error("Cannot run in multiple processes: IOLoop instance "
"has already been initialized. You cannot call "
"IOLoop.instance() before calling start()")
num_processes = 1
if num_processes > 1:
logging.info("Pre-forking %d server processes", num_processes)
for i in range(num_processes):
if os.fork() == 0:
self.io_loop = ioloop.IOLoop.instance()
self.io_loop.add_handler(
self._socket.fileno(), self._handle_events,
ioloop.IOLoop.READ)
return
os.waitpid(-1, 0)
else:
if not self.io_loop:
self.io_loop = ioloop.IOLoop.instance()
self.io_loop.add_handler(self._socket.fileno(),
self._handle_events,
ioloop.IOLoop.READ)
def stop(self):
self.io_loop.remove_handler(self._socket.fileno())
self._socket.close()
def _handle_events(self, fd, events):
while True:
try:
connection, address = self._socket.accept()
except socket.error, e:
if e[0] in (errno.EWOULDBLOCK, errno.EAGAIN):
return
raise
try:
stream = iostream.IOStream(connection, io_loop=self.io_loop)
SMTPConnection(stream, address, self.request_callback)
except:
logging.error("Error in connection callback", exc_info=True)
class SMTPConnection(object):
def __init__(self, stream, address, request_callback):
self.stream = stream
self.address = address
self.request_callback = request_callback
self.stream.write("220 myserver Tornado Simple Mail Transfer Service Ready\r\n")
self.stream.read_until("\r\n", self._parse_req)
def write(self, chunk):
assert self._request, "Request closed"
if not self.stream.closed():
self.stream.write(chunk, self._on_write_complete)
def finish(self):
if not self.stream.writing():
self.stream.close()
def _parse_req(self, data):
if data.find("HELO") > -1 or data.find("EHLO") > -1:
self.stream.write("250 myserver\r\n")
self.stream.read_until('\r\n', self._parse_req)
elif data.find("MAIL FROM") > -1:
tokens = data.split(':')
self.from_email = tokens[1]
self.stream.write("200 Ok\r\n")
self.stream.read_until("\r\n", self._parse_req)
elif data.find("RCPT TO") > -1:
tokens = data.split(':')
self.to_email = tokens[1]
self.stream.write("200 Ok\r\n")
self.stream.read_until("\r\n", self._parse_req)
elif data.find("DATA") > -1:
self.stream.write("354 End data with <CR><LF>.<CR><LF>\r\n")
self.stream.read_until("\r\n.\r\n", self._parse_msg)
elif data.find("QUIT") > -1:
self.stream.write("221 myserver Service closing transmission channel\r\n")
self.finish()
else:
print "unknown command: " + data
self.stream.read_until('\r\n', self._parse_req)
def _parse_msg(self, data):
msg = email.message_from_string(data)
self.stream.write("250 Ok\r\n")
self.request_callback(msg)
self.stream.read_until('\r\n', self._parse_req)