initial commit from appmail

This commit is contained in:
Adam Cooke
2017-04-19 13:07:25 +01:00
parent a3eff53792
commit 2fdba0ceb5
474 changed files with 51228 additions and 0 deletions
View File
+52
View File
@@ -0,0 +1,52 @@
require 'logger'
module Postal
class AppLogger < Logger
def self.greylog_notifier
@greylog_notifier ||= Postal.config.logging.greylog ? GELF::Notifier.new(Postal.config.logging.greylog.host, Postal.config.logging.greylog.port) : nil
end
def initialize(log_name, *args)
@log_name = log_name
super(*args)
self.formatter = LogFormatter.new
end
def add(severity, message = nil, progname = nil)
super
if n = self.class.greylog_notifier
begin
if message.nil?
message = block_given? ? yield : progname
end
message = message.to_s.force_encoding('UTF-8').scrub
message_without_ansi = message.gsub(/\e\[([\d\;]+)?m/, '') rescue message
n.notify!(:short_message => message_without_ansi, :log_name => @log_name, :facility => 'postal', :application_name => 'postal', :process_name => ENV['PROC_NAME'], :pid => Process.pid)
rescue => e
# Can't log this to GELF. Soz.
Raven.capture_exception(e)
end
end
true
end
end
class LogFormatter
TIME_FORMAT = "%Y-%m-%dT%H:%M:%S.%3N".freeze
COLORS = [32,34,35,31,32,33]
def call(severity, datetime, progname, msg)
time = datetime.strftime(TIME_FORMAT)
if number = ENV['PROC_NAME']
id = number.split('.').last.to_i
proc_text = "\e[#{COLORS[id % COLORS.size]}m[#{ENV['PROC_NAME']}:#{Process.pid}]\e[0m"
else
proc_text = "[#{Process.pid}]"
end
"#{proc_text} [#{time}] #{severity} -- : #{msg}\n"
end
end
end
+55
View File
@@ -0,0 +1,55 @@
module Postal
class BounceMessage
def initialize(server, message)
@server = server
@message = message
end
def raw_message
mail = Mail.new
mail.to = @message.mail_from
mail.from = "Mail Delivery Service <#{@message.route.description}>"
mail.subject = "Mail Delivery Failed (#{@message.subject})"
mail.text_part = body
mail.attachments['Original Message.eml'] = {:mime_type => 'message/rfc822', :encoding => 'quoted-printable', :content => @message.raw_message}
mail.message_id = "<#{SecureRandom.uuid}@#{Postal.config.dns.return_path}>"
mail.to_s
end
def queue
message = @server.message_db.new_message
message.scope = 'outgoing'
message.rcpt_to = @message.mail_from
message.mail_from = @message.route.description
message.domain_id = @message.domain&.id
message.raw_message = self.raw_message
message.bounce = 1
message.bounce_for_id = @message.id
message.save
message.id
end
def postmaster_address
@server.postmaster_address || "postmaster@#{@message.domain&.name || Postal.config.web.host}"
end
private
def body
<<-BODY.strip_heredoc
This is the mail delivery service responsible for delivering mail to #{@message.route.description}.
The message you've sent cannot be delivered. Your original message is attached to this message.
For further assistance please contact #{postmaster_address}. Please include the details below to help us identify the issue.
Message Token: #{@message.token}@#{@server.token}
Orginal Message ID: #{@message.message_id}
Mail from: #{@message.mail_from}
Rcpt To: #{@message.rcpt_to}
BODY
end
end
end
+162
View File
@@ -0,0 +1,162 @@
require 'yaml'
require 'pathname'
require_relative 'error'
require_relative 'version'
module Postal
def self.host
@host ||= config.web.host || "localhost:5000"
end
def self.protocol
@protocol ||= config.web.protocol || "http"
end
def self.host_with_protocol
@host_with_protocol ||= "#{protocol}://#{host}"
end
def self.app_root
@app_root ||= Pathname.new(File.expand_path('../../../', __FILE__))
end
def self.config
@config ||= begin
require 'hashie/mash'
Hashie::Mash.new(yaml_config)
end
end
def self.config_root
@config_root ||= begin
if __FILE__ =~ /\A\/opt\/postal/
Pathname.new("/opt/postal/config")
elsif ENV['AM_CONFIG_ROOT']
Pathname.new(ENV['AM_CONFIG_ROOT'])
else
Pathname.new(File.expand_path("../../../config", __FILE__))
end
end
end
def self.config_file_path
@config_file_path ||= File.join(config_root, 'postal.yml')
end
def self.yaml_config
@yaml_config ||= File.exist?(config_file_path) ? YAML.load_file(config_file_path) : {}
end
def self.database_url
if config.main_db
"mysql2://#{config.main_db.username}:#{config.main_db.password}@#{config.main_db.host}:#{config.main_db.port}/#{config.main_db.database}?encoding=#{config.main_db.encoding || 'utf8mb4'}"
else
"mysql2://root@localhost/postal"
end
end
def self.logger_for(name)
@loggers ||= {}
@loggers[name.to_sym] ||= begin
require 'postal/app_logger'
if config.logging.stdout || ENV['LOG_TO_STDOUT']
Postal::AppLogger.new(name, STDOUT)
else
Postal::AppLogger.new(name, app_root.join('log', "#{name}.log"), config.logging.max_log_files || 10, (config.logging.max_log_file_size || 20).megabytes)
end
end
end
def self.process_name
@process_name ||= begin
string = "host:#{Socket.gethostname} pid:#{Process.pid}"
string += " procname:#{ENV['PROC_NAME']}" if ENV['PROC_NAME']
string
rescue
"pid:#{Process.pid}"
end
end
def self.locker_name
string = process_name.dup
string += " job:#{Thread.current[:job_id]}" if Thread.current[:job_id]
string
end
def self.smtp_from_name
config.smtp&.from_name || "Postal"
end
def self.smtp_from_address
config.smtp&.from_address || "postal@example.com"
end
def self.smtp_private_key
@smtp_private_key ||= OpenSSL::PKey::RSA.new(File.read(smtp_private_key_path))
end
def self.smtp_private_key_path
config_root.join('smtp.key')
end
def self.smtp_certificate_path
config_root.join('smtp.cert')
end
def self.smtp_certificate_data
@smtp_certificate_data ||= File.read(smtp_certificate_path)
end
def self.smtp_certificates
@smtp_certificates ||= begin
certs = self.smtp_certificate_data.scan(/-----BEGIN CERTIFICATE-----.+?-----END CERTIFICATE-----/m)
certs.map do |c|
OpenSSL::X509::Certificate.new(c)
end.freeze
end
end
def self.lets_encrypt_private_key_path
@lets_encrypt_private_key_path ||= Postal.config_root.join('lets_encrypt.pem')
end
def self.signing_key_path
config_root.join('signing.key')
end
def self.signing_key
@signing_key ||= OpenSSL::PKey::RSA.new(File.read(signing_key_path))
end
def self.amrp_dkim_dns_record
public_key = signing_key.public_key.to_s.gsub(/\-+[A-Z ]+\-+\n/, '').gsub(/\n/, '')
"v=DKIM1; t=s; h=sha256; p=#{public_key};"
end
class ConfigError < Postal::Error
end
def self.check_config!
unless File.exist?(self.config_file_path)
raise ConfigError, "No config found at #{self.config_file_path}"
end
unless File.exist?(self.smtp_private_key_path)
raise ConfigError, "No SMTP private key found at #{self.smtp_private_key_path}"
end
unless File.exist?(self.smtp_certificate_path)
raise ConfigError, "No SMTP certificate found at #{self.smtp_certificate_path}"
end
unless File.exists?(self.lets_encrypt_private_key_path)
raise ConfigError, "No Let's Encrypt private key found at #{self.lets_encrypt_private_key_path}"
end
unless File.exists?(self.signing_key_path)
raise ConfigError, "No signing key found at #{self.signing_key_path}"
end
end
end
+5
View File
@@ -0,0 +1,5 @@
module Postal
module Countries
NAMES = ["Afghanistan", "Aland Islands", "Albania", "Algeria", "American Samoa", "Andorra", "Angola", "Anguilla", "Antarctica", "Antigua and Barbuda", "Argentina", "Armenia", "Aruba", "Australia", "Austria", "Azerbaijan", "Bahamas", "Bahrain", "Bangladesh", "Barbados", "Belarus", "Belgium", "Belize", "Benin", "Bermuda", "Bhutan", "Bolivia, Plurinational State of", "Bosnia and Herzegovina", "Botswana", "Bouvet Island", "Brazil", "British Indian Ocean Territory", "Brunei Darussalam", "Bulgaria", "Burkina Faso", "Burundi", "Cambodia", "Cameroon", "Canada", "Cape Verde", "Cayman Islands", "Central African Republic", "Chad", "Chile", "China", "Christmas Island", "Cocos (keeling) Islands", "Colombia", "Comoros", "Congo", "Congo, the Democratic Republic of the", "Cook Islands", "Costa Rica", "Cote D'ivoire", "Croatia", "Cuba", "Cyprus", "Czech Republic", "Denmark", "Djibouti", "Dominica", "Dominican Republic", "Ecuador", "Egypt", "El Salvador", "Equatorial Guinea", "Eritrea", "Estonia", "Ethiopia", "Falkland Islands (malvinas)", "Faroe Islands", "Fiji", "Finland", "France", "French Guiana", "French Polynesia", "French Southern Territories", "Gabon", "Gambia", "Georgia", "Germany", "Ghana", "Gibraltar", "Greece", "Greenland", "Grenada", "Guadeloupe", "Guam", "Guatemala", "Guernsey", "Guinea", "Guinea-bissau", "Guyana", "Haiti", "Heard Island and Mcdonald Islands", "Holy See (Vatican City State)", "Honduras", "Hong Kong", "Hungary", "Iceland", "India", "Indonesia", "Iran, Islamic Republic of", "Iraq", "Ireland", "Isle of Man", "Israel", "Italy", "Jamaica", "Japan", "Jersey", "Jordan", "Kazakhstan", "Kenya", "Kiribati", "Korea, Democratic People's Republic of", "Korea, Republic of", "Kuwait", "Kyrgyzstan", "Lao People's Democratic Republic", "Latvia", "Lebanon", "Lesotho", "Liberia", "Libyan Arab Jamahiriya", "Liechtenstein", "Lithuania", "Luxembourg", "Macao", "Macedonia, the Former Yugoslav Republic of", "Madagascar", "Malawi", "Malaysia", "Maldives", "Mali", "Malta", "Marshall Islands", "Martinique", "Mauritania", "Mauritius", "Mayotte", "Mexico", "Micronesia, Federated States of", "Moldova, Republic of", "Monaco", "Mongolia", "Montenegro", "Montserrat", "Morocco", "Mozambique", "Myanmar", "Namibia", "Nauru", "Nepal", "Netherlands", "Netherlands Antilles", "New Caledonia", "New Zealand", "Nicaragua", "Niger", "Nigeria", "Niue", "Norfolk Island", "Northern Mariana Islands", "Norway", "Oman", "Pakistan", "Palau", "Palestinian Territory, Occupied", "Panama", "Papua New Guinea", "Paraguay", "Peru", "Philippines", "Pitcairn", "Poland", "Portugal", "Puerto Rico", "Qatar", "Reunion", "Romania", "Russian Federation", "Rwanda", "Saint Barthelemy", "Saint Helena, Ascension and Tristan Da Cunha", "Saint Kitts and Nevis", "Saint Lucia", "Saint Martin", "Saint Pierre and Miquelon", "Saint Vincent and the Grenadines", "Samoa", "San Marino", "Sao Tome and Principe", "Saudi Arabia", "Senegal", "Serbia", "Seychelles", "Sierra Leone", "Singapore", "Slovakia", "Slovenia", "Solomon Islands", "Somalia", "South Africa", "South Georgia and the South Sandwich Islands", "Spain", "Sri Lanka", "Sudan", "Suriname", "Svalbard and Jan Mayen", "Swaziland", "Sweden", "Switzerland", "Syrian Arab Republic", "Taiwan", "Tajikistan", "Tanzania, United Republic of", "Thailand", "Timor-leste", "Togo", "Tokelau", "Tonga", "Trinidad and Tobago", "Tunisia", "Turkey", "Turkmenistan", "Turks and Caicos Islands", "Tuvalu", "Uganda", "Ukraine", "United Arab Emirates", "United Kingdom", "United States", "United States Minor Outlying Islands", "Uruguay", "Uzbekistan", "Vanuatu", "Venezuela, Bolivarian Republic of", "Viet Nam", "Virgin Islands, British", "Virgin Islands, U.S.", "Wallis and Futuna", "Western Sahara", "Yemen", "Zambia", "Zimbabwe"]
end
end
+87
View File
@@ -0,0 +1,87 @@
module Postal
class DKIMHeader
def initialize(domain, message)
if domain && domain.dkim_status == 'OK'
@domain_name = domain.name
@dkim_key = domain.dkim_key
@dkim_identifier = domain.dkim_identifier
else
@domain_name = Postal.config.dns.return_path
@dkim_key = Postal.signing_key
@dkim_identifier = 'postal'
end
@domain = domain
@message = message
@raw_headers, @raw_body = @message.split(/\r?\n\r?\n/, 2)
end
def dkim_header
"DKIM-Signature: v=1;" + dkim_properties + signature
end
private
def headers
@headers ||= @raw_headers.to_s.gsub(/\r?\n\s/, ' ').split(/\r?\n/)
end
def header_names
normalized_headers.map{ |h| h.split(':')[0].strip }
end
def normalized_headers
Array.new.tap do |new_headers|
headers.select { |h| h.match(/^(to|from|date|subject|message-id):/i) }.each do |h|
new_headers << normalize_header(h)
end
end
end
def normalize_header(content)
content.gsub!(/[ \t]+/, ' ') # Tidy whitespace
key, value = content.split(':', 2).map{ |a| a.strip } # Split into key/value and strip whitespace
key.downcase! # Downcase the key
key + ':' + value # Rejoin
end
def normalized_body
@normalized_body ||= begin
content = @raw_body.dup
content.gsub!("\r", '') # Make sure we have no random CRs
content.gsub!("\n", "\r\n") # Convert to CRLF
content.gsub!(/[ \t]+/, ' ') # Tidy whitespace
content.gsub!(" \r\n", "\r\n") # Remove trailing whitespace
content.gsub!(/(\r\n)+\z/, "") # Remove trailing lines
content.gsub!(/\z/, "\r\n") # Add a final newline
content
end
end
def body_hash
@body_hash ||= Base64.encode64(Digest::SHA256.digest(normalized_body)).strip
end
def dkim_properties
String.new.tap do |header|
header << " a=rsa-sha256; c=relaxed/relaxed;"
header << " d=#{@domain_name}; s=#{@dkim_identifier}; t=#{Time.now.utc.to_i};"
header << " bh=#{body_hash}; h=#{header_names.join(':')};"
header << " b="
end
end
def dkim_header_for_signing
"dkim-signature:v=1;" + dkim_properties
end
def signable_header_string
(normalized_headers + [dkim_header_for_signing]).join("\r\n")
end
def signature
Base64.encode64(@dkim_key.sign(OpenSSL::Digest::SHA256.new, signable_header_string)).gsub("\n", '')
end
end
end
+7
View File
@@ -0,0 +1,7 @@
module Postal
module Errors
end
class Error < StandardError
end
end
+55
View File
@@ -0,0 +1,55 @@
module Postal
class Espect
def self.inspect(message, scope = :incoming)
if Postal.config.espect&.hosts
hosts = Postal.config.espect.hosts.dup.shuffle
hosts.each do |host|
result = Postal::HTTP.post("#{host}/inspect", :text_body => Base64.encode64(message), :timeout => 20)
if result[:code] == 200 && json = (JSON.parse(result[:body]) rescue nil)
return EspectResult.new(json, scope)
end
end
nil
end
end
end
class EspectResult
EXCLUSIONS = {
:outgoing => ['NO_RECEIVED', 'NO_RELAYS', 'ALL_TRUSTED', 'FREEMAIL_FORGED_REPLYTO', 'RDNS_DYNAMIC', /^SPF\_/, /^HELO\_/, /DKIM_/, /^RCVD_IN_/],
:incoming => []
}
def initialize(reply, scope)
@reply = reply
@scope = scope
end
def spam_score
@spam_score ||= begin
spam_details.inject(0.0) do |total, detail|
total += detail['score'] || 0.0
end
end
end
def spam_details
@spam_details ||= (@reply['spam_details'] || []).reject do |d|
EXCLUSIONS[@scope].any? do |item|
item == d['code'] || (item.is_a?(Regexp) && item =~ d['code'])
end
end
end
def threat?
@reply['threat'] ? true : false
end
def threat_message
@reply['threat_message']
end
end
end
+168
View File
@@ -0,0 +1,168 @@
require 'stringio'
module Postal
module FastServer
class Client
class ClientWentAway < StandardError; end
class BadRequest < StandardError; end
def initialize(socket, options)
@raw_socket = socket
@options = options
end
def run
Timeout.timeout(15) do
if Postal.config.fast_server.proxy_protocol
# gets without readahead
line = ""
char = nil
while(char != "\n")
char = @raw_socket.read(1)
line << char
end
line.chomp!
if m = line.match(/\APROXY (.+) (.+) (.+) (.+) (.+)\z/)
@remote_ip = m[2]
else
return false
end
end
if self.ssl?
@socket = OpenSSL::SSL::SSLSocket.new(@raw_socket, self.class.ssl_context)
@socket.accept
else
@socket = @raw_socket
end
Timeout::timeout(20) do
# Read the request line
request = @socket.gets.to_s.chomp
# Split the request into its 3 parts
method, path, protocol = request.split(' ', 3)
raise BadRequest unless method && path && protocol
# Create an empty header set
header_set = HTTPHeaderSet.new
# Read each header and populate the header set
loop do
header = @socket.gets
if header.nil?
raise ClientWentAway
elsif header.chomp == ""
break
else
header_set << HTTPHeader.from_string(header.chomp)
end
end
# At this point, one might want to read the request body, but I don't think we need it.
# Build rack request
server_name, server_port = header_set['Host'].try(:value).to_s.split(":", 2)
request = {
"REQUEST_METHOD" => method,
"SCRIPT_NAME" => "",
"PATH_INFO" => path.split('?', 2)[0],
"QUERY_STRING" => path.split('?', 2)[1],
"SERVER_NAME" => server_name || "",
"SERVER_PORT" => server_name || "",
"rack.version" => [1, 3],
"rack.url_scheme" => ssl? ? "https" : "http",
"rack.input" => StringIO.new(""),
"rack.errors" => STDERR,
"rack.multithread" => true,
"rack.multiprocess" => true,
"rack.run_once" => false,
"rack.hijack" => false,
"rack.hijack_io" => false,
"REMOTE_ADDR" => remote_ip,
}
# Add request headers to rack hash
header_set.headers.each do |header|
request["HTTP_" + header.key.gsub('-', '_').upcase] = header.value
end
# Call the rack app and process the result
code, headers, body = Interface.new.call(request)
response = "HTTP/1.1 #{code} #{Rack::Utils::HTTP_STATUS_CODES[code]}\r\n"
headers.each do |k,v|
response << "#{k}:#{v}\r\n"
end
response << "\r\n"
body.each do |data|
response << data
end
@socket.write(response)
end
end
rescue ClientWentAway, Timeout::Error, Errno::ECONNRESET
# We don't really care if a client has disapeared, close the sockets and carry on.
rescue OpenSSL::SSL::SSLError
# Don't worry about SSL negotiation failures, disconnect and carry on
rescue BadRequest
# We couldn't read a proper HTTP request, disconnect the client
rescue => e
Raven.capture_exception(e)
ensure
@socket.close rescue nil
@raw_socket.close rescue nil
end
def ssl?
!!@options[:ssl]
end
def remote_ip
@remote_ip || @raw_socket.peeraddr[3].sub('::ffff:', '')
end
def self.ssl_context(domain_name = nil)
@ssl_certificates ||= {}
unless @ssl_certificates_refreshed && @ssl_certificates_refreshed > Time.now.utc.beginning_of_day
@ssl_certificates_refreshed = Time.now.utc
@ssl_certificates = {}
end
@ssl_certificates[domain_name] ||= OpenSSL::SSL::SSLContext.new.tap do |ssl_context|
if domain_name
if domain = TrackCertificate.active.where(:domain => domain_name).first
ssl_context.cert = domain.certificate_object
ssl_context.extra_chain_cert = domain.intermediaries_array
ssl_context.key = domain.key_object
end
end
if ssl_context.cert.nil?
certs = Postal.ssl_certificates
ssl_context.cert = certs.shift
ssl_context.extra_chain_cert = certs
ssl_context.key = Postal.signing_key
end
ssl_context.ssl_version = "SSLv23"
ssl_context.ciphers = 'EECDH+ECDSA+AESGCM EECDH+aRSA+AESGCM EECDH+ECDSA+SHA384 EECDH+ECDSA+SHA256 EECDH+aRSA+SHA384 EECDH+aRSA+SHA256 EECDH+aRSA+RC4 EECDH EDH+aRSA !aNULL !eNULL !LOW !3DES !MD5 !EXP !PSK !SRP !DSS !RC4 !DH'
ssl_context.options = OpenSSL::SSL::SSLContext::DEFAULT_PARAMS[:options] |
OpenSSL::SSL::OP_NO_SSLv2 |
OpenSSL::SSL::OP_NO_SSLv3 |
OpenSSL::SSL::OP_NO_COMPRESSION |
OpenSSL::SSL::OP_CIPHER_SERVER_PREFERENCE
ssl_context.tmp_ecdh_callback = Proc.new do |*a|
OpenSSL::PKey::EC.new("prime256v1")
end
unless domain_name
ssl_context.servername_cb = Proc.new do |ctx, hostname|
self.ssl_context(hostname)
end
end
end
end
end
end
end
+21
View File
@@ -0,0 +1,21 @@
module Postal
module FastServer
class HTTPHeader
attr_accessor :key, :value
def self.from_string(string)
k, v = string.to_s.split(/\:\s*/, 2)
self.new(k.to_s, v.to_s)
end
def initialize(k, v)
@key = k
@value = v
end
def to_s
@key + ": " + @value
end
end
end
end
+37
View File
@@ -0,0 +1,37 @@
module Postal
module FastServer
class HTTPHeaderSet
attr_accessor :headers
def initialize
@headers = []
end
def self.from_string_array(array)
header_set = self.new
header_set.headers = array.map{|h|HTTPHeader.from_string(h)}
header_set
end
def select(key)
@headers.select{|h|h.key.downcase == key.downcase}
end
def [](key)
@headers.find{|h|h.key.downcase == key.downcase}
end
def []=(key, value)
self.delete(key)
@headers << HTTPHeader.new(key, value)
end
def delete(key)
@headers.delete_if{|h|h.key.downcase == key.downcase}
end
def <<(header)
@headers << header
end
end
end
end
+92
View File
@@ -0,0 +1,92 @@
module Postal
module FastServer
class Interface
# TODO: Make this multithreaded? Thread-safe?
TRACKING_PIXEL = File.read(Rails.root.join('app', 'assets', 'images', 'tracking_pixel.png'))
def get_message_db_from_server_token(token)
if server = ::Server.find_by_token(token)
server.message_db
else
nil
end
end
def call(env)
request = Rack::Request.new(env)
if request.path =~ /\A\/(\.well-known\/.*)/
if certificate = ::TrackCertificate.find_by_verification_path($1)
return [200, {'Content-Length' => certificate.verification_string.bytesize.to_s}, [certificate.verification_string]]
else
return [404, {}, ["Verification not found"]]
end
elsif request.path =~ /\A\/img\/([a-z0-9\-]+)\/([a-z0-9\-]+)/i
server_token = $1
message_token = $2
if message_db = get_message_db_from_server_token(server_token)
begin
message = message_db.message(:token => message_token)
message.create_load(request)
rescue Postal::MessageDB::Message::NotFound
# This message has been removed, we'll just continue to serve the image
rescue => e
# Somethign else went wrong. We don't want to stop the image loading though because
# this is our problem. Log this exception though.
Raven.capture_exception(e)
end
source_image = request.params['src']
if source_image.nil?
headers = {}
headers['Content-Type'] = "image/png"
headers['Content-Length'] = TRACKING_PIXEL.bytesize.to_s
return [200, headers, [TRACKING_PIXEL]]
elsif source_image =~ /\Ahttps?\:\/\//
response = Postal::HTTP.get(source_image, :timeout => 3)
if response[:code] == 200
headers = {}
headers['Content-Type'] = response[:headers]['content-type']&.first
headers['Last-Modified'] = response[:headers]['last-modified']&.first
headers['Cache-Control'] = response[:headers]['cache-control']&.first
headers['Etag'] = response[:headers]['etag']&.first
headers['Content-Length'] = response[:body].bytesize.to_s
return [200, headers, [response[:body]]]
else
return [404, {}, ['Not found']]
end
else
return [400, {}, ['Invalid/missing source image']]
end
else
return [404, {}, ['Invalid Server Token']]
end
end
if request.path =~ /\A\/([a-z0-9\-]+)\/([a-z0-9\-]+)/i
server_token = $1
link_token = $2
if message_db = get_message_db_from_server_token(server_token)
if link = message_db.select(:links, :where => {:token => link_token}, :limit => 1).first
time = Time.now.to_f
message_db.update(:messages, {:clicked => time}, :where => {:id => link['message_id']})
message_db.insert(:clicks, {:message_id => link['message_id'], :link_id => link['id'], :ip_address => request.ip, :user_agent => request.user_agent, :timestamp => time})
SendWebhookJob.queue(:main, :server_id => message_db.server_id, :event => 'MessageLinkClicked', :payload => {:_message => link['message_id'], :url => link['url'], :token => link['token'], :ip_address => request.ip, :user_agent => request.user_agent})
return [307, {'Location' => link['url']}, ["Redirected to: #{link['url']}"]]
else
return [404, {}, ['Link not found']]
end
else
return [404, {}, ['Invalid Server Token']]
end
end
[200, {}, ["Hello."]]
end
end
end
end
+37
View File
@@ -0,0 +1,37 @@
require 'socket'
require 'openssl'
module Postal
module FastServer
class Server
def run
Thread.abort_on_exception = true
TrackCertificate
server_sockets = {
TCPServer.new(Postal.config.fast_server.bind_address, Postal.config.fast_server.ssl_port) => {:ssl => true},
TCPServer.new(Postal.config.fast_server.bind_address, Postal.config.fast_server.port) => {:ssl => false},
}
Postal.logger_for(:fast_server).info("Fast server started listening on HTTP port #{Postal.config.fast_server.port}")
Postal.logger_for(:fast_server).info("Fast server started listening on HTTPS port #{Postal.config.fast_server.ssl_port}")
loop do
client = nil
ios = select(server_sockets.keys, nil, nil, 1)
if ios && server_io = ios[0][0]
begin
client_io = server_io.accept_nonblock
client = Client.new(client_io, server_sockets[server_io])
Thread.new(client) { |t_client| t_client.run }
rescue IO::WaitReadable, Errno::EINTR
# Never mind, guess the client went away
rescue => e
Raven.capture_exception(e)
client_io.close rescue nil
end
end
end
end
end
end
end
+10
View File
@@ -0,0 +1,10 @@
module Postal
module Helpers
def self.strip_name_from_address(address)
return nil if address.nil?
address.gsub(/.*</, '').gsub(/>.*/, '').strip
end
end
end
+94
View File
@@ -0,0 +1,94 @@
require 'net/https'
require 'uri'
module Postal
module HTTP
def self.get(url, options = {})
request(Net::HTTP::Get, url, options)
end
def self.post(url, options = {})
request(Net::HTTP::Post, url, options)
end
def self.request(method, url, options = {})
options[:headers] ||= {}
uri = URI.parse(url)
request = method.new(uri.path.length == 0 ? "/" : uri.path)
options[:headers].each { |k,v| request.add_field k, v }
if options[:username]
request.basic_auth(options[:username], options[:password])
end
if options[:params].is_a?(Hash)
# If params has been provided, sent it them as form encoded values
request.set_form_data(options[:params])
elsif options[:json].is_a?(String)
# If we have a JSON string, set the content type and body to be the JSON
# data
request.add_field 'Content-Type', 'application/json'
request.body = options[:json]
elsif options[:text_body]
# Add a plain text body if we have one
request.body = options[:text_body]
end
if options[:sign]
#signature = EncryptoSigno.sign(Postal.signing_key, request.body.to_s).gsub("\n", '')
#request.add_field 'X-Postal-Signature', signature
end
request['User-Agent'] = options[:user_agent] || "Postal/#{Postal::VERSION}"
connection = Net::HTTP.new(uri.host, uri.port)
if uri.scheme == 'https'
connection.use_ssl = true
connection.verify_mode = OpenSSL::SSL::VERIFY_PEER
ssl = true
else
ssl = false
end
begin
timeout = options[:timeout] || 60
Timeout.timeout(timeout) do
result = connection.request(request)
{
:code => result.code.to_i,
:body => result.body,
:headers => result.to_hash,
:secure => @ssl
}
end
rescue OpenSSL::SSL::SSLError => e
{
:code => -3,
:body => "Invalid SSL certificate",
:headers =>{},
:secure => @ssl
}
rescue SocketError, Errno::ECONNRESET, EOFError, Errno::EINVAL, Errno::ENETUNREACH, Errno::EHOSTUNREACH, Errno::ECONNREFUSED => e
{
:code => -2,
:body => e.message,
:headers => {},
:secure => @ssl
}
rescue Timeout::Error => e
{
:code => -1,
:body => "Timed out after #{timeout}s",
:headers => {},
:secure => @ssl
}
end
end
end
end
+125
View File
@@ -0,0 +1,125 @@
module Postal
class HTTPSender < Sender
def initialize(endpoint, options = {})
@endpoint = endpoint
@options = options
@log_id = Nifty::Utils::RandomString.generate(:length => 8).upcase
end
def send_message(message)
start_time = Time.now
result = SendResult.new
result.log_id = @log_id
request_options = {}
request_options[:sign] = true
request_options[:timeout] = @endpoint.timeout || 5
case @endpoint.encoding
when 'BodyAsJSON'
request_options[:json] = parameters(message, :flat => false).to_json
when 'FormData'
request_options[:params] = parameters(message, :flat => true)
end
log "Sending request to #{@endpoint.url}"
response = Postal::HTTP.post(@endpoint.url, request_options)
result.secure = !!response[:secure]
result.details = "Received a #{response[:code]} from #{@endpoint.url}"
log " -> Received: #{response[:code]}"
if response[:body]
log " -> Body: #{response[:body][0,255]}"
result.output = response[:body].to_s[0, 500].strip
end
if response[:code] >= 200 && response[:code] < 300
# This is considered a success
result.type = 'Sent'
elsif response[:code] >= 500 && response[:code] < 600
# This is temporary. They might fix their server so it should soft fail.
result.type = 'SoftFail'
result.retry = true
elsif response[:code] < 0
# Connection/SSL etc... errors
result.type = 'SoftFail'
result.retry = true
result.connect_error = true
else
# This is permanent. Any other error isn't cool with us.
result.type = 'HardFail'
end
result.time = (Time.now - start_time).to_f.round(2)
result
end
private
def log(text)
Postal.logger_for(:http_sender).info("[#{@log_id}] #{text}")
end
def parameters(message, options = {})
case @endpoint.format
when 'Hash'
hash = {
:id => message.id,
:rcpt_to => message.rcpt_to,
:mail_from => message.mail_from,
:token => message.token,
:subject => message.subject,
:message_id => message.message_id,
:timestamp => message.timestamp.to_f,
:size => message.size,
:spam_status => message.spam_status,
:bounce => message.bounce == 1 ? true : false,
:received_with_ssl => message.received_with_ssl == 1,
:to => message.headers['to']&.last,
:cc => message.headers['cc']&.last,
:from => message.headers['from']&.last,
:date => message.headers['date']&.last,
:in_reply_to => message.headers['in-reply-to']&.last,
:references => message.headers['references']&.last,
:html_body => message.html_body,
:attachment_quantity => message.attachments.size
}
if @endpoint.strip_replies
hash[:plain_body], hash[:replies_from_plain_body] = Postal::ReplySeparator.separate(message.plain_body)
else
hash[:plain_body] = message.plain_body
end
if @endpoint.include_attachments?
if options[:flat]
message.attachments.each_with_index do |a, i|
hash["attachments[#{i}][filename]"] = a.filename
hash["attachments[#{i}][content_type]"] = a.content_type
hash["attachments[#{i}][size]"] = a.body.to_s.bytesize.to_s
hash["attachments[#{i}][data]"] = Base64.encode64(a.body.to_s)
end
else
hash[:attachments] = message.attachments.map do |a|
{
:filename => a.filename,
:content_type => a.mime_type,
:size => a.body.to_s.bytesize,
:data => Base64.encode64(a.body.to_s)
}
end
end
end
hash
when 'RawMessage'
{
:id => message.id,
:message => Base64.encode64(message.raw_message),
:base64 => true,
:size => message.size.to_i
}
else
{}
end
end
end
end
+36
View File
@@ -0,0 +1,36 @@
require 'nifty/utils/random_string'
module Postal
class Job
def initialize(id, params = {})
@id = id
@params = params.with_indifferent_access
end
def id
@id
end
def params
@params || {}
end
def perform
end
def log(text)
Worker.logger.info "[#{@id}] #{text}"
end
def self.queue(queue, params = {})
job_id = Nifty::Utils::RandomString.generate(:length => 10).upcase
job_payload = {'params' => params, 'class_name' => self.name, 'id' => job_id, 'queue' => queue}
Postal::Worker.job_queue(queue).publish(job_payload.to_json, :persistent => false)
job_id
end
def self.perform(params = {})
new(nil, params).perform
end
end
end
+24
View File
@@ -0,0 +1,24 @@
require 'acme-client'
module Postal
module LetsEncrypt
def self.client
@client ||= Acme::Client.new(:private_key => private_key, :endpoint => endpoint)
end
def self.private_key
@private_key ||= OpenSSL::PKey::RSA.new(File.open(Postal.lets_encrypt_private_key_path))
end
def self.endpoint
@endpoint ||= Rails.env.development? ? "https://acme-staging.api.letsencrypt.org" : "https://acme-v01.api.letsencrypt.org/"
end
def self.register_private_key(email_address)
registration = client.register(:contact => "mailto:#{email_address}")
registration.agree_terms
end
end
end
+19
View File
@@ -0,0 +1,19 @@
module Postal
module MessageDB
class Click
def initialize(attributes, link)
@url = link['url']
@ip_address = attributes['ip_address']
@user_agent = attributes['user_agent']
@timestamp = Time.at(attributes['timestamp'])
end
attr_reader :ip_address
attr_reader :user_agent
attr_reader :timestamp
attr_reader :url
end
end
end
+379
View File
@@ -0,0 +1,379 @@
module Postal
module MessageDB
class Database
def initialize(organization_id, server_id)
@organization_id = organization_id
@server_id = server_id
end
attr_reader :organization_id
attr_reader :server_id
#
# Return the server
#
def server
@server ||= Server.find_by_id(@server_id)
end
#
# Return the current schema version
#
def schema_version
@schema_version ||= begin
last_migration = select(:migrations, :order => :version, :direction => 'DESC', :limit => 1).first
last_migration ? last_migration['version'] : 0
rescue Mysql2::Error => e
e.message =~ /doesn\'t exist/ ? 0 : raise
end
end
#
# Return a single message. Accepts an ID or an array of conditions
#
def message(*args)
Message.find_one(self, *args)
end
#
# Return an array or count of messages.
#
def messages(*args)
Message.find(self, *args)
end
def messages_with_pagination(*args)
Message.find_with_pagination(self, *args)
end
#
# Create a new message with the given attributes. This won't be saved to the database
# until it has been 'save'd.
#
def new_message(attributes = {})
Message.new(self, attributes)
end
#
# Return the total size of all stored messages
#
def total_size
query("SELECT SUM(size) AS size FROM `#{database_name}`.`raw_message_sizes`").first['size'] || 0
end
#
# Return the live stats instance
#
def live_stats
@live_stats ||= LiveStats.new(self)
end
#
# Return the statistics instance
#
def statistics
@statistics ||= Statistics.new(self)
end
#
# Return the provisioner instance
#
def provisioner
@provisioner ||= Provisioner.new(self)
end
#
# Return the provisioner instance
#
def suppression_list
@suppression_list ||= SuppressionList.new(self)
end
#
# Return the provisioner instance
#
def webhooks
@webhooks ||= Webhooks.new(self)
end
#
# Return the name for a raw message table for a given date
#
def raw_table_name_for_date(date)
date.strftime("raw-%Y-%m-%d")
end
#
# Insert a new raw message into a table (creating it if needed)
#
def insert_raw_message(data, date = Date.today)
table_name = raw_table_name_for_date(date)
begin
headers, body = data.split(/\r?\n\r?\n/, 2)
headers_id = insert(table_name, :data => headers)
body_id = insert(table_name, :data => body)
rescue Mysql2::Error => e
if e.message =~ /doesn\'t exist/
provisioner.create_raw_table(table_name)
retry
else
raise
end
end
[table_name, headers_id, body_id]
end
#
# Selects entries from the database. Accepts a number of options which can be used
# to manipulate the results.
#
# :where => A hash containing the query
# :order => The name of a field to order by
# :direction => The order that should be applied to ordering (ASC or DESC)
# :fields => An array of fields to select
# :limit => Limit the number of results
# :page => Which page number to return
# :per_page => The number of items per page (defaults to 30)
# :count => Return a count of the results instead of the actual data
#
def select(table, options = {})
sql_query = "SELECT"
if options[:count]
sql_query << " COUNT(id) AS count"
elsif options[:fields]
sql_query << " " + options[:fields].map { |f| "`#{f}`" }.join(', ')
else
sql_query << " *"
end
sql_query << " FROM `#{database_name}`.`#{table}`"
if options[:where] && !options[:where].empty?
sql_query << " " + build_where_string(options[:where], ' AND ')
end
if options[:order]
direction = (options[:direction] || 'ASC').upcase
raise Postal::Error, "Invalid direction #{options[:direction]}" unless ['ASC', 'DESC'].include?(direction)
sql_query << " ORDER BY `#{options[:order]}` #{direction}"
end
if options[:limit]
sql_query << " LIMIT #{options[:limit]}"
end
if options[:offset]
sql_query << " OFFSET #{options[:offset]}"
end
result = query(sql_query)
if options[:count]
result.first['count']
else
result.to_a
end
end
#
# A paginated version of select
#
def select_with_pagination(table, page, options = {})
page = page.to_i
page = 1 if page <= 0
per_page = options.delete(:per_page) || 30
offset = (page - 1) * per_page
result = {}
result[:total] = select(table, options.merge(:count => true))
result[:records] = select(table, options.merge(:limit => per_page, :offset => offset))
result[:per_page] = per_page
result[:total_pages], remainder = result[:total].divmod(per_page)
result[:total_pages] += 1 if remainder > 0
result[:page] = page
result
end
#
# Updates a record in the database. Accepts a table name, the attributes to update
# plus some options which are shown below:
#
# :where => The condition to apply to the query
#
# Will return the total number of affected rows.
#
def update(table, attributes, options = {})
sql_query = "UPDATE `#{database_name}`.`#{table}` SET"
sql_query << " #{hash_to_sql(attributes)}"
if options[:where]
sql_query << " " + build_where_string(options[:where])
end
with_mysql do |mysql|
query_on_connection(mysql, sql_query)
mysql.affected_rows
end
end
#
# Insert a record into a given table. A hash of attributes is also provided.
# Will return the ID of the new item.
#
def insert(table, attributes)
sql_query = "INSERT INTO `#{database_name}`.`#{table}`"
sql_query << " (" + attributes.keys.map { |k| "`#{k}`" }.join(', ') + ")"
sql_query << " VALUES (" + attributes.values.map { |v| escape(v) }.join(', ') + ")"
with_mysql do |mysql|
query_on_connection(mysql, sql_query)
mysql.last_id
end
end
#
# Insert multiple rows at the same time in the same query
#
def insert_multi(table, keys, values)
if values.empty?
nil
else
sql_query = "INSERT INTO `#{database_name}`.`#{table}`"
sql_query << " (" + keys.map { |k| "`#{k}`" }.join(', ') + ")"
sql_query << " VALUES "
sql_query << values.map { |v| "(" + v.map { |v| escape(v) }.join(', ') + ")" }.join(', ')
query(sql_query)
end
end
#
# Deletes a in the database. Accepts a table name, and some options which
# are shown below:
#
# :where => The condition to apply to the query
#
# Will return the total number of affected rows.
#
def delete(table, options = {})
sql_query = "DELETE FROM `#{database_name}`.`#{table}`"
sql_query << " " + build_where_string(options[:where], ' AND ')
with_mysql do |mysql|
query_on_connection(mysql, sql_query)
mysql.affected_rows
end
end
#
# Return the correct database name
#
def database_name
@database_name ||= "#{Postal.config.message_db.prefix}-server-#{@server_id}"
end
#
# Run a query, log it and return the result
#
class ResultForExplainPrinter
attr_reader :columns
attr_reader :rows
def initialize(result)
if result.first
@columns = result.first.keys
@rows = result.map { |row| row.map(&:last) }
else
@columns = []
@rows = []
end
end
end
def stringify_keys(hash)
hash.each_with_object({}) do |(key, value), hash|
hash[key.to_s] = value
end
end
def escape(value)
with_mysql do |mysql|
if value == true
'1'
elsif value == false
'0'
elsif value.nil?
'NULL'
else
if value.to_s.length == 0
'NULL'
else
"'" + mysql.escape(value.to_s) + "'"
end
end
end
end
def query(query)
with_mysql do |mysql|
query_on_connection(mysql, query)
end
end
private
def query_on_connection(connection, query)
start_time = Time.now.to_f
result = connection.query(query)
time = Time.now.to_f - start_time
logger.debug " \e[4;34mMessageDB Query (#{time.round(2)}s) \e[0m \e[33m#{query}\e[0m"
if time > 0.5 && query =~ /\A(SELECT|UPDATE|DELETE) /
id = Nifty::Utils::RandomString.generate(:length => 6).upcase
explain_result = ResultForExplainPrinter.new(connection.query("EXPLAIN #{query}"))
slow_query_logger.info "[#{id}] EXPLAIN #{query}"
for line in ActiveRecord::ConnectionAdapters::MySQL::ExplainPrettyPrinter.new.pp(explain_result, time).split("\n")
slow_query_logger.info "[#{id}] " + line
end
end
result
end
def logger
defined?(Rails) ? Rails.logger : Logger.new(STDOUT)
end
def slow_query_logger
Postal.logger_for(:slow_message_db_queries)
end
def with_mysql(&block)
MessageDB::MySQL.client(&block)
end
def build_where_string(attributes, joiner = ', ')
"WHERE #{hash_to_sql(attributes, joiner)}"
end
def hash_to_sql(hash, joiner = ', ')
hash.map do |key, value|
if value.is_a?(Array) && value.all? { |v| v.is_a?(Fixnum) }
"`#{key}` IN (#{value.join(', ')})"
elsif value.is_a?(Array)
escaped_values = value.map { |v| escape(v) }.join(', ')
"`#{key}` IN (#{escaped_values})"
elsif value.is_a?(Hash)
sql = []
value.each do |operator, value|
case operator
when :less_than
sql << "`#{key}` < #{escape(value)}"
when :greater_than
sql << "`#{key}` > #{escape(value)}"
when :less_than_or_equal_to
sql << "`#{key}` <= #{escape(value)}"
when :greater_than_or_equal_to
sql << "`#{key}` >= #{escape(value)}"
end
end
sql.empty? ? "1=1" : sql.join(joiner)
else
"`#{key}` = #{escape(value)}"
end
end.join(joiner)
end
end
end
end
+71
View File
@@ -0,0 +1,71 @@
module Postal
module MessageDB
class Delivery
def self.create(message, attributes = {})
attributes = message.database.stringify_keys(attributes)
attributes = attributes.merge('message_id' => message.id, 'timestamp' => Time.now.to_f)
id = message.database.insert('deliveries', attributes)
delivery = Delivery.new(message, attributes.merge('id' => id))
delivery.update_statistics
delivery.send_webhooks
delivery
end
def initialize(message, attributes)
@message = message
@attributes = attributes.stringify_keys
end
def method_missing(name, value = nil, &block)
if @attributes.has_key?(name.to_s)
@attributes[name.to_s]
else
nil
end
end
def timestamp
@timestamp ||= @attributes['timestamp'] ? Time.at(@attributes['timestamp']) : nil
end
def update_statistics
if self.status == 'Held'
@message.database.statistics.increment_all(self.timestamp, 'held')
end
if self.status == 'Bounced' || self.status == 'HardFail'
@message.database.statistics.increment_all(self.timestamp, 'bounces')
end
end
def send_webhooks
if self.webhook_event
WebhookRequest.trigger(@message.database.server_id, self.webhook_event, self.webhook_hash)
end
end
def webhook_hash
{
:message => @message.webhook_hash,
:status => self.status,
:details => self.details,
:output => self.output.to_s.force_encoding('UTF-8').scrub,
:sent_with_ssl => self.sent_with_ssl,
:timestamp => @attributes['timestamp'],
:time => self.time
}
end
def webhook_event
@webhook_event ||= case self.status
when 'Sent' then 'MessageSent'
when 'SoftFail' then 'MessageDelayed'
when 'HardFail' then 'MessageDeliveryFailed'
when 'Held' then 'MessageHeld'
end
end
end
end
end
+41
View File
@@ -0,0 +1,41 @@
module Postal
module MessageDB
class LiveStats
def initialize(database)
@database = database
end
#
# Increment the live stats by one for the current minute
#
def increment(type)
time = Time.now.utc
type = @database.escape(type.to_s)
sql_query = "INSERT INTO `#{@database.database_name}`.`live_stats` (type, minute, timestamp, count)"
sql_query << " VALUES (#{type}, #{time.min}, #{time.to_f}, 1)"
sql_query << " ON DUPLICATE KEY UPDATE count = if(timestamp < #{time.to_f - 1800}, 1, count + 1), timestamp = #{time.to_f}"
@database.query(sql_query)
end
#
# Return the total number of messages for the last 60 minutes
#
def total(minutes, options = {})
if minutes > 60
raise Postal::Error, "Live stats can only return data for the last 60 minutes."
end
options[:types] ||= [:incoming, :outgoing]
if options[:types].empty?
raise Postal::Error, "You must provide at least one type to return"
else
time = minutes.minutes.ago.beginning_of_minute.utc.to_f
types = options[:types].map {|t| "#{@database.escape(t.to_s)}"}.join(', ')
result = @database.query("SELECT SUM(count) as count FROM `#{@database.database_name}`.`live_stats` WHERE `type` IN (#{types}) AND timestamp > #{time}").first
result['count'] || 0
end
end
end
end
end
+17
View File
@@ -0,0 +1,17 @@
module Postal
module MessageDB
class Load
def initialize(attributes)
@ip_address = attributes['ip_address']
@user_agent = attributes['user_agent']
@timestamp = Time.at(attributes['timestamp'])
end
attr_reader :ip_address
attr_reader :user_agent
attr_reader :timestamp
end
end
end
+576
View File
@@ -0,0 +1,576 @@
module Postal
module MessageDB
class Message
class NotFound < Postal::Error
end
def self.find_one(database, query)
query = {:id => query.to_i} if query.is_a?(Fixnum)
if message = database.select('messages', :where => query, :limit => 1).first
Message.new(database, message)
else
raise NotFound, "No message found matching provided query #{query}"
end
end
def self.find(database, options = {})
if messages = database.select('messages', options)
if messages.is_a?(Array)
messages.map { |m| Message.new(database, m) }
else
messages
end
else
[]
end
end
def self.find_with_pagination(database, page, options = {})
messages = database.select_with_pagination('messages', page, options)
messages[:records] = messages[:records].map { |m| Message.new(database, m) }
messages
end
attr_reader :database
def initialize(database, attributes)
@database = database
@attributes = attributes
end
#
# Return the server for this message
#
def server
@database.server
end
#
# Return the credential for this message
#
def credential
@credential ||= self.credential_id ? Credential.find_by_id(self.credential_id) : nil
end
#
# Return the route for this message
#
def route
@route ||= self.route_id ? Route.find_by_id(self.route_id) : nil
end
#
# Return the endpoint for this message
#
def endpoint
@endpoint ||= begin
if self.endpoint_type && self.endpoint_id
self.endpoint_type.constantize.find_by_id(self.endpoint_id)
elsif self.route && self.route.mode == 'Endpoint'
self.route.endpoint
end
end
end
#
# Return the credential for this message
#
def domain
@domain ||= self.domain_id ? Domain.find_by_id(self.domain_id) : nil
end
#
# Copy appropriate attributes from the raw message to the message itself
#
def copy_attributes_from_raw_message
if self.raw_message
self.subject = self.headers['subject']&.last
self.message_id = self.headers['message-id']&.last
if self.message_id
self.message_id = self.message_id.gsub(/.*</, '').gsub(/>.*/, '').strip
end
end
end
#
# Return the timestamp for this message
#
def timestamp
@timestamp ||= @attributes['timestamp'] ? Time.at(@attributes['timestamp']) : nil
end
#
# Return the time that the last delivery was attempted
#
def last_delivery_attempt
@last_delivery_attempt ||= @attributes['last_delivery_attempt'] ? Time.at(@attributes['last_delivery_attempt']) : nil
end
#
# Return the hold expiry for this message
#
def hold_expiry
@hold_expiry ||= @attributes['hold_expiry'] ? Time.at(@attributes['hold_expiry']) : nil
end
#
# Has this message been read?
#
def read?
!!(loaded || clicked)
end
#
# Add a delivery attempt for this message
#
def create_delivery(status, options = {})
delivery = Delivery.create(self, options.merge(:status => status))
hold_expiry = status == 'Held' ? 7.days.from_now.to_f : nil
self.update(:status => status, :last_delivery_attempt => delivery.timestamp.to_f, :held => status == 'Held' ? 1 : 0, :hold_expiry => hold_expiry)
delivery
end
#
# Return all deliveries for this object
#
def deliveries
@deliveries ||= begin
@database.select('deliveries', :where => {:message_id => self.id}, :order => :timestamp).map do |hash|
Delivery.new(self, hash)
end
end
end
#
# Return all the clicks for this object
#
def clicks
@clicks ||= begin
clicks = @database.select('clicks', :where => {:message_id => self.id}, :order => :timestamp)
if clicks.empty?
[]
else
links = @database.select('links', :where => {:id => clicks.map { |c| c['link_id'].to_i }}).group_by { |l| l['id'] }
clicks.map do |hash|
Click.new(hash, links[hash['link_id']].first)
end
end
end
end
#
# Return all the loads for this object
#
def loads
@loads ||= begin
loads = @database.select('loads', :where => {:message_id => self.id}, :order => :timestamp)
loads.map do |hash|
Load.new(hash)
end
end
end
#
# Return all activity entries
#
def activity_entries
@activity_entries ||= (deliveries + clicks + loads).sort_by(&:timestamp)
end
#
# Provide access to set and get acceptable attributes
#
def method_missing(name, value = nil, &block)
if @attributes.has_key?(name.to_s)
@attributes[name.to_s]
elsif name.to_s =~ /\=\z/
@attributes[name.to_s.gsub('=', '').to_s] = value
else
nil
end
end
#
# Has this message been persisted to the database yet?
#
def persisted?
!@attributes['id'].nil?
end
#
# Save this message
#
def save
save_raw_message
persisted? ? _update : _create
self
end
#
# Update this message
#
def update(attributes_to_change)
@attributes = @attributes.merge(database.stringify_keys(attributes_to_change))
if persisted?
@database.update('messages', attributes_to_change, :where => {:id => self.id})
else
_create
end
end
#
# Delete the message from the database
#
def delete
if persisted?
@database.delete('messages', :where => {:id => self.id})
end
end
#
# Return the headers
#
def raw_headers
if self.raw_table
@raw_headers ||= @database.select(self.raw_table, :where => {:id => self.raw_headers_id}).first&.send(:[], 'data') || ""
else
""
end
end
#
# Return the full raw message body for this message.
#
def raw_body
if self.raw_table
@raw ||= @database.select(self.raw_table, :where => {:id => self.raw_body_id}).first&.send(:[], 'data') || ""
else
""
end
end
#
# Return the full raw message for this message
#
def raw_message
@raw_message ||= "#{raw_headers}\r\n\r\n#{raw_body}"
end
#
# Set the raw message ready for saving later
#
def raw_message=(raw)
@pending_raw_message = raw.force_encoding('BINARY')
end
#
# Save the raw message to the database as appropriate
#
def save_raw_message
if @pending_raw_message
self.size = @pending_raw_message.bytesize
date = Date.today
table_name, headers_id, body_id = @database.insert_raw_message(@pending_raw_message, date)
self.raw_table = table_name
self.raw_headers_id = headers_id
self.raw_body_id = body_id
@raw = nil
@raw_headers = nil
@headers = nil
@mail = nil
@pending_raw_message = nil
copy_attributes_from_raw_message
@database.query("UPDATE `#{@database.database_name}`.`raw_message_sizes` SET size = size + #{self.size} WHERE table_name = '#{table_name}'")
end
end
#
# Is there a raw message?
#
def raw_message?
!!self.raw_table
end
#
# Return the plain body for this message
#
def plain_body
mail&.plain_body
end
#
# Return the HTML body for this message
#
def html_body
mail&.html_body
end
#
# Return the HTML body with any tracking links
#
def html_body_without_tracking_image
html_body.gsub(/\<p class\=['"]ampimg['"].*?\<\/p\>/, '')
end
#
# Return all attachments for this message
#
def attachments
mail&.attachments || []
end
#
# Return the headers for this message
#
def headers
@headers ||= begin
mail = Mail.new(self.raw_headers)
mail.header.fields.each_with_object({}) do |field, hash|
hash[field.name.downcase] ||= []
hash[field.name.downcase] << field.decoded
end
end
end
#
# Return the recipient domain for this message
#
def recipient_domain
self.rcpt_to ? self.rcpt_to.split('@').last : nil
end
#
# Create a new item in the message queue for this message
#
def add_to_message_queue(options = {})
QueuedMessage.create!(:message => self, :server_id => @database.server_id, :batch_key => self.batch_key, :domain => self.recipient_domain, :route_id => self.route_id, :manual => options[:manual]).id
end
#
# Return a suitable batch key for this message
#
def batch_key
case self.scope
when 'outgoing'
key = "outgoing-"
key += self.recipient_domain.to_s
when 'incoming'
key = "incoming-"
key += "rt:#{self.route_id}-ep:#{self.endpoint_id}-#{self.endpoint_type}"
else
key = nil
end
key
end
#
# Return the queued message
#
def queued_message
@queued_message ||= self.id ? QueuedMessage.where(:message_id => self.id, :server_id => @database.server_id).first : nil
end
#
# Return the spam status
#
def spam_status
return 'NotChecked' unless inspected == 1
spam == 1 ? 'Spam' : 'NotSpam'
end
#
# Has this message been held?
#
def held?
status == 'Held'
end
#
# Does this message have our DKIM header yet?
#
def has_outgoing_headers?
!!(raw_headers =~ /^X\-Postal\-MsgID\:/i)
end
#
# Add dkim header
#
def add_outgoing_headers
headers = []
if self.domain
dkim = Postal::DKIMHeader.new(self.domain, self.raw_message)
headers << dkim.dkim_header
end
headers << "X-Postal-MsgID: #{self.token}"
append_headers(*headers)
end
#
# Append a header to the existing headers
#
def append_headers(*headers)
new_headers = headers.join("\r\n")
new_headers = "#{new_headers}\r\n#{self.raw_headers}"
@database.update(self.raw_table, {:data => new_headers}, :where => {:id => self.raw_headers_id})
@raw_headers = new_headers
@raw_message = nil
@headers = nil
end
#
# Return a suitable
#
def webhook_hash
@webhook_hash ||= {
:id => self.id,
:token => self.token,
:direction => self.scope,
:message_id => self.message_id,
:to => self.rcpt_to,
:from => self.mail_from,
:subject => self.subject,
:timestamp => self.timestamp.to_f,
:spam_status => self.spam_status,
:tag => self.tag
}
end
#
# Mark this message as bounced
#
def bounce!(bounce_message)
create_delivery('Bounced', :details => "We've received a bounce message for this e-mail. See <msg:#{bounce_message.id}> for details.")
SendWebhookJob.queue(:main, :server_id => self.database.server_id, :event => "MessageBounced", :payload => {:_original_message => self.id, :_bounce => bounce_message.id})
end
#
# Should bounces be sent for this message?
#
def send_bounces?
self.bounce != 1 && self.mail_from.present?
end
#
# Add a load for this message
#
def create_load(request)
update('loaded' => Time.now.to_f) if loaded.nil?
database.insert(:loads, {:message_id => self.id, :ip_address => request.ip, :user_agent => request.user_agent, :timestamp => Time.now.to_f})
SendWebhookJob.queue(:main, :server_id => self.database.server_id, :event => 'MessageLoaded', :payload => {:_message => self.id, :ip_address => request.ip, :user_agent => request.user_agent})
end
#
# Create a new link
#
def create_link(url)
hash = Digest::SHA1.hexdigest(url.to_s)
token = Nifty::Utils::RandomString.generate(:length => 8)
database.insert(:links, {:message_id => self.id, :hash => hash, :url => url, :timestamp => Time.now.to_f, :token => token})
token
end
#
# Return a message object that this message is a reply to
#
def original_messages
return nil unless self.bounce == 1
other_message_ids = raw_message.scan(/\X\-Postal\-MsgID\:\s*([a-z0-9]+)/i).flatten
if other_message_ids.empty?
[]
else
database.messages(:where => {:token => other_message_ids})
end
end
#
# Was thsi message sent to a return path?
#
def rcpt_to_return_path?
!!(rcpt_to =~ /\@#{Regexp.escape(Postal.config.dns.custom_return_path_prefix)}\./)
end
#
# Inspect this message
#
def inspect_message
if result = Espect.inspect(self.raw_message, self.scope&.to_sym)
# Update the messages table with the results of our inspection
update(:inspected => 1, :spam_score => result.spam_score, :threat => result.threat?, :threat_details => result.threat_message)
# Add any spam details into the spam checks database
self.database.insert_multi(:spam_checks, [:message_id, :code, :score, :description], result.spam_details.map { |d| [self.id, d['code'], d['score'], d['description']]})
# Return the espect result
result
end
end
#
# Return all spam checks for this message
#
def spam_checks
@spam_checks ||= self.database.select(:spam_checks, :where => {:message_id => self.id})
end
#
# Cancel the hold on this message
#
def cancel_hold
if self.status == 'Held'
create_delivery('HoldCancelled', :details => "The hold on this message has been removed without action.")
end
end
#
# Parse the contents of this message
#
def parse_content
parse_result = Postal::MessageParser.new(self)
if parse_result.actioned?
# Somethign was changed, update the raw message
@database.update(self.raw_table, {:data => parse_result.new_body}, :where => {:id => self.raw_body_id})
@raw = parse_result.new_body
@raw_message = nil
end
update('parsed' => 1, 'tracked_links' => parse_result.tracked_links, 'tracked_images' => parse_result.tracked_images)
end
#
# Has this message been parsed?
#
def parsed?
self.parsed == 1
end
#
# Should this message be parsed?
#
def should_parse?
parsed? == false && headers['x-amp'] != 'skip'
end
private
def _update
@database.update('messages', @attributes.reject {|k,v| k == :id }, :where => {:id => @attributes['id']})
end
def _create
self.timestamp = Time.now.to_f if self.timestamp.blank?
self.status = 'Pending' if self.status.blank?
self.token = Nifty::Utils::RandomString.generate(:length => 12) if self.token.blank?
last_id = @database.insert('messages', @attributes.reject {|k,v| k == :id })
@attributes['id'] = last_id
@database.statistics.increment_all(self.timestamp, self.scope)
Statistic.global.increment!(:total_messages)
Statistic.global.increment!("total_#{self.scope}".to_sym)
add_to_message_queue
end
def mail
# This version of mail is only used for accessing the bodies.
@mail ||= raw_message? ? Mail.new(raw_message) : nil
end
end
end
end
+35
View File
@@ -0,0 +1,35 @@
module Postal
module MessageDB
class Migration
def initialize(database)
@database = database
end
def up
end
def self.run(database, start_from = database.schema_version)
files = Dir[Rails.root.join('lib', 'postal', 'message_db', 'migrations', '*.rb')]
files = files.map { |f| id, name = f.split('/').last.split('_', 2); [id.to_i, name] }.sort_by(&:first)
latest_version = files.last.first
if latest_version > start_from
puts "\e[32mMigrating #{database.database_name} from version #{start_from} => #{files.last.first}\e[0m"
else
puts "Nothing to do."
end
files.each do |version, file|
klass_name = file.gsub(/\.rb\z/, '').camelize
next if start_from >= version
puts "\e[45m++ Migrating #{klass_name} (#{version})\e[0m"
require "postal/message_db/migrations/#{version.to_s.rjust(2, '0')}_#{file}"
klass = Postal::MessageDB::Migrations.const_get(klass_name)
instance = klass.new(database)
instance.up
database.insert(:migrations, :version => version)
end
end
end
end
end
@@ -0,0 +1,16 @@
module Postal
module MessageDB
module Migrations
class CreateMigrations < Postal::MessageDB::Migration
def up
@database.provisioner.create_table(:migrations,
:columns => {
:version => 'int(11) NOT NULL',
},
:primary_key => '`version`'
)
end
end
end
end
end
@@ -0,0 +1,57 @@
module Postal
module MessageDB
module Migrations
class CreateMessages < Postal::MessageDB::Migration
def up
@database.provisioner.create_table(:messages,
:columns => {
:id => 'int(11) NOT NULL AUTO_INCREMENT',
:token => 'varchar(255) DEFAULT NULL',
:scope => 'varchar(10) DEFAULT NULL',
:rcpt_to => 'varchar(255) DEFAULT NULL',
:mail_from => 'varchar(255) DEFAULT NULL',
:subject => 'varchar(255) DEFAULT NULL',
:message_id => 'varchar(255) DEFAULT NULL',
:timestamp => 'decimal(18,6) DEFAULT NULL',
:route_id => 'int(11) DEFAULT NULL',
:domain_id => 'int(11) DEFAULT NULL',
:credential_id => 'int(11) DEFAULT NULL',
:status => 'varchar(255) DEFAULT NULL',
:held => 'tinyint(1) DEFAULT 0',
:size => 'varchar(255) DEFAULT NULL',
:last_delivery_attempt => 'decimal(18,6) DEFAULT NULL',
:raw_table => 'varchar(255) DEFAULT NULL',
:raw_body_id => 'int(11) DEFAULT NULL',
:raw_headers_id => 'int(11) DEFAULT NULL',
:inspected => 'tinyint(1) DEFAULT 0',
:spam => 'tinyint(1) DEFAULT 0',
:spam_score => 'decimal(8,2) DEFAULT 0',
:threat => 'tinyint(1) DEFAULT 0',
:threat_details => 'varchar(255) DEFAULT NULL',
:bounce => 'tinyint(1) DEFAULT 0',
:bounce_for_id => 'int(11) DEFAULT 0',
:tag => 'varchar(255) DEFAULT NULL',
:loaded => 'decimal(18,6) DEFAULT NULL',
:clicked => 'decimal(18,6) DEFAULT NULL',
:received_with_ssl => 'tinyint(1) DEFAULT NULL',
},
:indexes => {
:on_message_id => '`message_id`(8)',
:on_token => '`token`(6)',
:on_bounce_for_id => '`bounce_for_id`',
:on_held => '`held`',
:on_scope_and_status => '`scope`(1), `spam`, `status`(6), `timestamp`',
:on_scope_and_tag => '`scope`(1), `spam`, `tag`(8), `timestamp`',
:on_scope_and_spam => '`scope`(1), `spam`, `timestamp`',
:on_scope_and_thr_status => '`scope`(1), `threat`, `status`(6), `timestamp`',
:on_scope_and_threat => '`scope`(1), `threat`, `timestamp`',
:on_rcpt_to => '`rcpt_to`(12), `timestamp`',
:on_mail_from => '`mail_from`(12), `timestamp`',
:on_raw_table => '`raw_table`(14)',
}
)
end
end
end
end
end
@@ -0,0 +1,26 @@
module Postal
module MessageDB
module Migrations
class CreateDeliveries < Postal::MessageDB::Migration
def up
@database.provisioner.create_table(:deliveries,
:columns => {
:id => 'int(11) NOT NULL AUTO_INCREMENT',
:message_id => 'int(11) DEFAULT NULL',
:status => 'varchar(255) DEFAULT NULL',
:code => 'int(11) DEFAULT NULL',
:output => 'varchar(512) DEFAULT NULL',
:details => 'varchar(512) DEFAULT NULL',
:sent_with_ssl => 'tinyint(1) DEFAULT 0',
:log_id => 'varchar(100) DEFAULT NULL',
:timestamp => 'decimal(18,6) DEFAULT NULL'
},
:indexes => {
:on_message_id => '`message_id`'
}
)
end
end
end
end
end
@@ -0,0 +1,19 @@
module Postal
module MessageDB
module Migrations
class CreateLiveStats < Postal::MessageDB::Migration
def up
@database.provisioner.create_table(:live_stats,
:columns => {
:type => 'varchar(20) NOT NULL',
:minute => 'int(11) NOT NULL',
:count => 'int(11) DEFAULT NULL',
:timestamp => 'decimal(18,6) DEFAULT NULL',
},
:primary_key => '`minute`, `type`(8)'
)
end
end
end
end
end
@@ -0,0 +1,20 @@
module Postal
module MessageDB
module Migrations
class CreateRawMessageSizes < Postal::MessageDB::Migration
def up
@database.provisioner.create_table(:raw_message_sizes,
:columns => {
:id => 'int(11) NOT NULL AUTO_INCREMENT',
:table_name => 'varchar(255) DEFAULT NULL',
:size => 'bigint DEFAULT NULL'
},
:indexes => {
:on_table_name => '`table_name`(14)'
}
)
end
end
end
end
end
@@ -0,0 +1,26 @@
module Postal
module MessageDB
module Migrations
class CreateClicks < Postal::MessageDB::Migration
def up
@database.provisioner.create_table(:clicks,
:columns => {
:id => 'int(11) NOT NULL AUTO_INCREMENT',
:message_id => 'int(11) DEFAULT NULL',
:link_id => 'int(11) DEFAULT NULL',
:ip_address => 'varchar(255) DEFAULT NULL',
:country => 'varchar(255) DEFAULT NULL',
:city => 'varchar(255) DEFAULT NULL',
:user_agent => 'varchar(255) DEFAULT NULL',
:timestamp => 'decimal(18,6) DEFAULT NULL'
},
:indexes => {
:on_message_id => '`message_id`',
:on_link_id => '`link_id`'
}
)
end
end
end
end
end
@@ -0,0 +1,24 @@
module Postal
module MessageDB
module Migrations
class CreateLoads < Postal::MessageDB::Migration
def up
@database.provisioner.create_table(:loads,
:columns => {
:id => 'int(11) NOT NULL AUTO_INCREMENT',
:message_id => 'int(11) DEFAULT NULL',
:ip_address => 'varchar(255) DEFAULT NULL',
:country => 'varchar(255) DEFAULT NULL',
:city => 'varchar(255) DEFAULT NULL',
:user_agent => 'varchar(255) DEFAULT NULL',
:timestamp => 'decimal(18,6) DEFAULT NULL'
},
:indexes => {
:on_message_id => '`message_id`'
}
)
end
end
end
end
end
@@ -0,0 +1,26 @@
module Postal
module MessageDB
module Migrations
class CreateStats < Postal::MessageDB::Migration
def up
[:hourly, :daily, :monthly, :yearly].each do |table_name|
@database.provisioner.create_table("stats_#{table_name}",
:columns => {
:id => 'int(11) NOT NULL AUTO_INCREMENT',
:time => 'int(11) DEFAULT NULL',
:incoming => 'bigint DEFAULT NULL',
:outgoing => 'bigint DEFAULT NULL',
:spam => 'bigint DEFAULT NULL',
:bounces => 'bigint DEFAULT NULL',
:held => 'bigint DEFAULT NULL',
},
:unique_indexes => {
:on_time => '`time`'
}
)
end
end
end
end
end
end
@@ -0,0 +1,24 @@
module Postal
module MessageDB
module Migrations
class CreateLinks < Postal::MessageDB::Migration
def up
@database.provisioner.create_table(:links,
:columns => {
:id => 'int(11) NOT NULL AUTO_INCREMENT',
:message_id => 'int(11) DEFAULT NULL',
:token => 'varchar(255) DEFAULT NULL',
:hash => 'varchar(255) DEFAULT NULL',
:url => 'varchar(255) DEFAULT NULL',
:timestamp => 'decimal(18,6) DEFAULT NULL'
},
:indexes => {
:on_message_id => '`message_id`',
:on_token => '`token`(8)',
}
)
end
end
end
end
end
@@ -0,0 +1,23 @@
module Postal
module MessageDB
module Migrations
class CreateSpamChecks < Postal::MessageDB::Migration
def up
@database.provisioner.create_table(:spam_checks,
:columns => {
:id => 'int(11) NOT NULL AUTO_INCREMENT',
:message_id => 'int(11) DEFAULT NULL',
:score => 'decimal(8,2) DEFAULT NULL',
:code => 'varchar(255) DEFAULT NULL',
:description => 'varchar(255) DEFAULT NULL'
},
:indexes => {
:on_message_id => '`message_id`',
:on_code => '`code`(8)',
}
)
end
end
end
end
end
@@ -0,0 +1,11 @@
module Postal
module MessageDB
module Migrations
class AddTimeToDeliveries < Postal::MessageDB::Migration
def up
@database.query("ALTER TABLE `#{@database.database_name}`.`deliveries` ADD COLUMN `time` decimal(8,2)")
end
end
end
end
end
@@ -0,0 +1,11 @@
module Postal
module MessageDB
module Migrations
class AddHoldExpiry < Postal::MessageDB::Migration
def up
@database.query("ALTER TABLE `#{@database.database_name}`.`messages` ADD COLUMN `hold_expiry` decimal(18,6)")
end
end
end
end
end
@@ -0,0 +1,11 @@
module Postal
module MessageDB
module Migrations
class AddIndexToMessageStatus < Postal::MessageDB::Migration
def up
@database.query("ALTER TABLE `#{@database.database_name}`.`messages` ADD INDEX `on_status` (`status`(8)) USING BTREE")
end
end
end
end
end
@@ -0,0 +1,24 @@
module Postal
module MessageDB
module Migrations
class CreateSuppressions < Postal::MessageDB::Migration
def up
@database.provisioner.create_table(:suppressions,
:columns => {
:id => 'int(11) NOT NULL AUTO_INCREMENT',
:type => 'varchar(255) DEFAULT NULL',
:address => 'varchar(255) DEFAULT NULL',
:reason => 'varchar(255) DEFAULT NULL',
:timestamp => 'decimal(18,6) DEFAULT NULL',
:keep_until => 'decimal(18,6) DEFAULT NULL',
},
:indexes => {
:on_address => '`address`(6)',
:on_keep_until => '`keep_until`',
}
)
end
end
end
end
end
@@ -0,0 +1,28 @@
module Postal
module MessageDB
module Migrations
class CreateWebhookRequests < Postal::MessageDB::Migration
def up
@database.provisioner.create_table(:webhook_requests,
:columns => {
:id => 'int(11) NOT NULL AUTO_INCREMENT',
:uuid => 'varchar(255) DEFAULT NULL',
:event => 'varchar(255) DEFAULT NULL',
:attempt => 'int(11) DEFAULT NULL',
:timestamp => 'decimal(18,6) DEFAULT NULL',
:status_code => 'int(1) DEFAULT NULL',
:body => 'text DEFAULT NULL',
:payload => 'text DEFAULT NULL',
:will_retry => 'tinyint DEFAULT NULL'
},
:indexes => {
:on_uuid => '`uuid`(8)',
:on_event => '`event`(8)',
:on_timestamp => '`timestamp`',
}
)
end
end
end
end
end
@@ -0,0 +1,13 @@
module Postal
module MessageDB
module Migrations
class AddUrlAndHookToWebhooks < Postal::MessageDB::Migration
def up
@database.query("ALTER TABLE `#{@database.database_name}`.`webhook_requests` ADD COLUMN `url` varchar(255)")
@database.query("ALTER TABLE `#{@database.database_name}`.`webhook_requests` ADD COLUMN `webhook_id` int(11)")
@database.query("ALTER TABLE `#{@database.database_name}`.`webhook_requests` ADD INDEX `on_webhook_id` (`webhook_id`) USING BTREE")
end
end
end
end
end
@@ -0,0 +1,13 @@
module Postal
module MessageDB
module Migrations
class AddReplacedLinkCountToMessages < Postal::MessageDB::Migration
def up
@database.query("ALTER TABLE `#{@database.database_name}`.`messages` ADD COLUMN `tracked_links` int(11) DEFAULT 0")
@database.query("ALTER TABLE `#{@database.database_name}`.`messages` ADD COLUMN `tracked_images` int(11) DEFAULT 0")
@database.query("ALTER TABLE `#{@database.database_name}`.`messages` ADD COLUMN `parsed` tinyint DEFAULT 0")
end
end
end
end
end
@@ -0,0 +1,11 @@
module Postal
module MessageDB
module Migrations
class AddEndpointsToMessages < Postal::MessageDB::Migration
def up
@database.query("ALTER TABLE `#{@database.database_name}`.`messages` ADD COLUMN `endpoint_id` int(11), ADD COLUMN `endpoint_type` varchar(255)")
end
end
end
end
end
+35
View File
@@ -0,0 +1,35 @@
module Postal
module MessageDB
module MySQL
# This exists here because it needs to be required when the application loads
# so that it isn't unloaded in development. If it was unloaded in development,
# it would be undesirable as we'd just end up with lots of connections.
def self.new_client
Mysql2::Client.new(:host => Postal.config.message_db.host, :username => Postal.config.message_db.username, :password => Postal.config.message_db.password, :port => Postal.config.message_db.port, :reconnect => true)
end
@free_clients = []
def self.client(&block)
client = @free_clients.shift || self.new_client
return_value = nil
tries = 2
begin
return_value = block.call(client)
rescue Mysql2::Error => e
if e.message =~ /(lost connection|gone away)/i && (tries -= 1) > 0
retry
else
raise
end
ensure
@free_clients << client
end
return_value
end
end
end
end
+183
View File
@@ -0,0 +1,183 @@
module Postal
module MessageDB
class Provisioner
def initialize(database)
@database = database
end
#
# Provisions a new database
#
def provision
drop
create
migrate
end
#
# Migrate this database
#
def migrate(start_from = @database.schema_version)
Postal::MessageDB::Migration.run(@database, start_from)
end
#
# Does a database already exist?
#
def exists?
!!@database.query("SELECT schema_name FROM `information_schema`.`schemata` WHERE schema_name = '#{@database.database_name}'").first
end
#
# Creates a new empty database
#
def create
@database.query("CREATE DATABASE `#{@database.database_name}` CHARSET utf8 COLLATE UTF8_UNICODE_CI;")
true
rescue Mysql2::Error => e
e.message =~ /database exists/ ? false : raise
end
#
# Drops the whole message database
#
def drop
@database.query("DROP DATABASE `#{@database.database_name}`;")
true
rescue Mysql2::Error => e
e.message =~ /doesn\'t exist/ ? false : raise
end
#
# Create a new table
#
def create_table(table_name, options)
@database.query(create_table_query(table_name, options))
end
#
# Drop a table
#
def drop_table(table_name)
@database.query("DROP TABLE `#{@database.database_name}`.`#{table_name}`")
end
#
# Creates a new empty raw message table for the given date. Returns nothing.
#
def create_raw_table(table)
begin
@database.query(create_table_query(table,:columns => {
:id => 'int(11) NOT NULL AUTO_INCREMENT',
:data => 'mediumblob DEFAULT NULL',
:next => 'int(11) DEFAULT NULL'
}
))
@database.query("INSERT INTO `#{@database.database_name}`.`raw_message_sizes` (table_name, size) VALUES ('#{table}', 0)")
rescue Mysql2::Error => e
# Don't worry if the table already exists, another thread has already run this code.
raise unless e.message =~ /already exists/
end
end
#
# Return a list of raw message tables that are older than the given date
#
def raw_tables(max_age = 30)
earliest_date = max_age ? Date.today - max_age : nil
[].tap do |tables|
@database.query("SHOW TABLES FROM `#{@database.database_name}` LIKE 'raw-%'").each do |tbl|
tbl_name = tbl.to_a.first.last
date = Date.parse(tbl_name.gsub(/\Araw\-/, ''))
if earliest_date.nil? || date < earliest_date
tables << tbl_name
end
end
end.sort
end
#
# Tidy all messages
#
def remove_raw_tables_older_than(max_age = 30)
raw_tables(max_age).each do |table|
remove_raw_table(table)
end
end
#
# Remove a raw message table
#
def remove_raw_table(table)
@database.query("UPDATE `#{@database.database_name}`.`messages` SET raw_table = NULL, raw_headers_id = NULL, raw_body_id = NULL, size = NULL WHERE raw_table = '#{table}'")
@database.query("DELETE FROM `#{@database.database_name}`.`raw_message_sizes` WHERE table_name = '#{table}'")
drop_table(table)
end
#
# Remove messages from the messages table that are too old to retain
#
def remove_messages(max_age = 60)
time = (Date.today - max_age.days).to_time.end_of_day
if newest_message_to_remove = @database.select(:messages, :where => {:timestamp => {:less_than_or_equal_to => time.to_f}}, :limit => 1, :order => :id, :direction => 'DESC', :fields => [:id]).first
id = newest_message_to_remove['id']
@database.query("DELETE FROM `#{@database.database_name}`.`clicks` WHERE `message_id` <= #{id}")
@database.query("DELETE FROM `#{@database.database_name}`.`loads` WHERE `message_id` <= #{id}")
@database.query("DELETE FROM `#{@database.database_name}`.`deliveries` WHERE `message_id` <= #{id}")
@database.query("DELETE FROM `#{@database.database_name}`.`spam_checks` WHERE `message_id` <= #{id}")
@database.query("DELETE FROM `#{@database.database_name}`.`messages` WHERE `id` <= #{id}")
end
end
#
# Remove raw message tables in order order until size is under the given size (given in MB)
#
def remove_raw_tables_until_less_than_size(size)
tables = self.raw_tables(nil)
tables_removed = []
until @database.total_size <= size
table = tables.shift
tables_removed << table
remove_raw_table(table)
end
tables_removed
end
private
#
# Build a query to load a table
#
def create_table_query(table_name, options)
String.new.tap do |s|
s << "CREATE TABLE `#{@database.database_name}`.`#{table_name}` ("
s << options[:columns].map do |column_name, column_options|
"`#{column_name}` #{column_options}"
end.join(', ')
if options[:indexes]
s << ", "
s << options[:indexes].map do |index_name, index_options|
"KEY `#{index_name}` (#{index_options}) USING BTREE"
end.join(', ')
end
if options[:unique_indexes]
s << ", "
s << options[:unique_indexes].map do |index_name, index_options|
"UNIQUE KEY `#{index_name}` (#{index_options})"
end.join(', ')
end
if options[:primary_key]
s << ", PRIMARY KEY (#{options[:primary_key]})"
else
s << ", PRIMARY KEY (`id`)"
end
s << ") ENGINE=InnoDB AUTO_INCREMENT=1 DEFAULT CHARSET=utf8;"
end
end
end
end
end
+58
View File
@@ -0,0 +1,58 @@
module Postal
module MessageDB
class Statistics
def initialize(database)
@database = database
end
STATS_GAPS = {:hourly => :hour, :daily => :day, :monthly => :month, :yearly => :year}
COUNTERS = [:incoming, :outgoing, :spam, :bounces, :held]
#
# Increment an appropriate counter
#
def increment_one(type, field, time = Time.now)
time = time.utc
initial_values = COUNTERS.map do |c|
field.to_sym == c ? 1 : 0
end
time_i = time.send("beginning_of_#{STATS_GAPS[type]}").utc.to_i
sql_query = "INSERT INTO `#{@database.database_name}`.`stats_#{type}` (time, #{COUNTERS.join(', ')})"
sql_query << " VALUES (#{time_i}, #{initial_values.join(', ')})"
sql_query << " ON DUPLICATE KEY UPDATE #{field} = #{field} + 1"
@database.query(sql_query)
end
#
# Increment all stats counters
#
def increment_all(time, field)
STATS_GAPS.keys.each do |type|
increment_one(type, field, time)
end
end
#
# Get a statistic (or statistics)
#
def get(type, counters, start_date = Time.now, quantity = 10)
date = start_date.utc
items = quantity.times.each_with_object({}) do |i, hash|
hash[(start_date - i.send(STATS_GAPS[type])).send("beginning_of_#{STATS_GAPS[type]}").utc] = counters.each_with_object({}) do |c, h|
h[c] = 0
end
end
@database.select("stats_#{type}", :where => {:time => items.keys.map(&:to_i)}, :fields => [:time] | counters).each do |data|
time = Time.at(data.delete('time'))
data.each do |key, value|
items[time][key.to_sym] = value
end
end
items.to_a.reverse
end
end
end
end
+38
View File
@@ -0,0 +1,38 @@
module Postal
module MessageDB
class SuppressionList
def initialize(database)
@database = database
end
def add(type, address, options = {})
keep_until = (options[:days] || 30).days.from_now.to_f
if existing = @database.select('suppressions', :where => {:type => type, :address => address}, :limit =>1).first
reason = options[:reason] || existing['reason']
@database.update('suppressions', {:reason => reason, :keep_until => keep_until}, :where => {:id => existing['id']})
else
@database.insert('suppressions', {:type => type, :address => address, :reason => options[:reason], :timestamp => Time.now.to_f, :keep_until => keep_until})
end
true
end
def get(type, address)
@database.select('suppressions', :where => {:type => type, :address => address, :keep_until => {:greater_than_or_equal_to => Time.now.to_f}}, :limit => 1).first
end
def all_with_pagination(page)
@database.select_with_pagination(:suppressions, page, :order => :timestamp, :direction => 'desc')
end
def remove(type, address)
@database.delete('suppressions', :where => {:type => type, :address => address}) > 0
end
def prune
@database.delete('suppressions', :where => {:keep_until => {:less_than => Time.now.to_f}}) || 0
end
end
end
end
+88
View File
@@ -0,0 +1,88 @@
module Postal
module MessageDB
class Webhooks
def initialize(database)
@database = database
end
def record(attributes = {})
@database.insert(:webhook_requests, attributes)
end
def list(page)
result = @database.select_with_pagination(:webhook_requests, page, :order => :timestamp, :direction => 'desc')
result[:records] = result[:records].map { |i| Request.new(i) }
result
end
def find(uuid)
request = @database.select(:webhook_requests, :where => {:uuid => uuid}).first || raise(RequestNotFound, "No request found with UUID '#{uuid}'")
Request.new(request)
end
def prune
if last = @database.select(:webhook_requests, :where => {:timestamp => {:less_than => 10.days.ago.to_f}}, :order => 'timestamp', :direction => 'desc', :limit => 1, :fields => ['id']).first
@database.delete(:webhook_requests, :where => {:id => {:less_than_or_equal_to => last['id']}})
end
end
class RequestNotFound < Postal::Error
end
class Request
def initialize(attributes)
@attributes = attributes
end
def [](name)
@attributes[name.to_s]
end
def timestamp
Time.at(@attributes['timestamp'])
end
def event
@attributes['event']
end
def status_code
@attributes['status_code']
end
def url
@attributes['url']
end
def uuid
@attributes['uuid']
end
def payload
@attributes['payload']
end
def pretty_payload
@pretty_payload ||= begin
json = JSON.parse(self.payload)
JSON.pretty_unparse(json)
end
end
def body
@attributes['body']
end
def attempt
@attributes['attempt']
end
def will_retry?
@attributes['will_retry'] == 1
end
end
end
end
end
+144
View File
@@ -0,0 +1,144 @@
module Postal
class MessageParser
URL_REGEX = /(?<url>(?<protocol>https?)\:\/\/(?<domain>[A-Za-z0-9\-\.]+)(?<path>\/[A-Za-z0-9\/\.\/\+\?\&\-\_\%\=\~\:\;]+)?+)/
def initialize(message)
@message = message
@actioned = false
@tracked_links = 0
@tracked_images = 0
@domain = @message.server.track_domains.where(:domain => @message.domain, :dns_status => "OK").first
if @domain
@parsed_output = generate
end
end
attr_reader :tracked_links
attr_reader :tracked_images
def actioned?
@actioned || @tracked_links > 0 || @tracked_images > 0
end
def new_body
@parsed_output.split("\r\n\r\n", 2)[1]
end
private
def generate
@mail = Mail.new(@message.raw_message)
@original_message = @message.raw_message
if @mail.parts.empty?
if @mail.mime_type
if @mail.mime_type =~ /text\/plain/
@mail.body = parse(@mail.body.decoded.dup, :text)
@mail.content_transfer_encoding = nil
@mail.charset = 'UTF-8'
elsif @mail.mime_type =~ /text\/html/
@mail.body = parse(@mail.body.decoded.dup, :html)
@mail.content_transfer_encoding = nil
@mail.charset = 'UTF-8'
end
end
else
parse_parts(@mail.parts)
end
@mail.to_s
rescue => e
if Rails.env.development?
raise
else
Raven.capture_exception(e)
@actioned = false
@tracked_links = 0
@tracked_images = 0
@original_message
end
end
def parse_parts(parts)
parts.each do |part|
if part.content_type =~ /text\/html/
part.body = parse(part.body.decoded.dup, :html)
part.content_transfer_encoding = nil
part.charset = 'UTF-8'
elsif part.content_type =~ /text\/plain/
part.body = parse(part.body.decoded.dup, :text)
part.content_transfer_encoding = nil
part.charset = 'UTF-8'
elsif part.content_type =~ /multipart\/alternative/
unless part.parts.empty?
parse_parts(part.parts)
end
end
end
end
def parse(part, type = nil)
if @domain.track_clicks?
part = insert_links(part, type)
end
if @domain.track_loads? && type == :html
part = insert_tracking_image(part)
end
part
end
def insert_links(part, type = nil)
if type == :text
part.gsub!(/#{URL_REGEX}/) do
if track_domain?($~[:domain])
@tracked_links += 1
token = @message.create_link($~[:url])
"#{domain}/#{@message.server.token}/#{token}"
else
$&
end
end
end
if type == :html
part.gsub!(/href=([\'\"])(#{URL_REGEX})[\'\"]/) do
if track_domain?($~[:domain])
@tracked_links += 1
token = @message.create_link($~[:url])
"href='#{domain}/#{@message.server.token}/#{token}'"
else
$&
end
end
end
part.gsub!(/(https?)\+notrack\:\/\//) do
@actioned = true
"#{$1}://"
end
part
end
def insert_tracking_image(part)
@tracked_images += 1
container = "<p class='ampimg' style='display:none;visibility:none;margin:0;padding:0;line-height:0;'><img src='#{domain}/img/#{@message.server.token}/#{@message.token}' alt=''></p>"
if part =~ /\<\/body\>/
part.gsub("</body>", "#{container}</body>")
else
part + container
end
end
def domain
"#{@domain.use_ssl? ? 'https' : 'http'}://#{@domain.full_name}"
end
def track_domain?(domain)
!@domain.excluded_click_domains_array.include?(domain)
end
end
end
+32
View File
@@ -0,0 +1,32 @@
module Postal
class MessageRequeuer
def run
Signal.trap("INT") { @running ? @exit = true : Process.exit(0) }
Signal.trap("TERM") { @running ? @exit = true : Process.exit(0) }
log "Running message requeuer..."
loop do
@running = true
QueuedMessage.requeue_all
@running = false
check_exit
sleep 5
end
end
private
def log(text)
Postal.logger_for(:message_requeuer).info text
end
def check_exit
if @exit
log "Exiting"
Process.exit(0)
end
end
end
end
+38
View File
@@ -0,0 +1,38 @@
module Postal
class QueryString
def initialize(string)
@string = string.strip + " "
end
def [](value)
to_hash[value.to_s]
end
def empty?
to_hash.empty?
end
def to_hash
@hash ||= @string.scan(/([a-z]+)\:\s*(?:(\d{2,4}\-\d{2}-\d{2}\s\d{2}\:\d{2})|\"(.*?)\"|(.*?))[\s\z]/).each_with_object({}) do |(key, date, string_with_spaces, value), hash|
if date
actual_value = date
elsif string_with_spaces
actual_value = string_with_spaces
elsif value == "[blank]"
actual_value = nil
else
actual_value = value
end
if hash.keys.include?(key.to_s)
hash[key.to_s] = [hash[key.to_s]].flatten
hash[key.to_s] << actual_value
else
hash[key.to_s] = actual_value
end
end
end
end
end
+25
View File
@@ -0,0 +1,25 @@
require 'postal/config'
require 'bunny'
module Postal
module RabbitMQ
def self.create_connection
conn = Bunny.new(
:host => Postal.config.rabbitmq&.host || 'localhost',
:port => Postal.config.rabbitmq&.port || 5672,
:username => Postal.config.rabbitmq&.username || 'guest',
:password => Postal.config.rabbitmq&.password || 'guest',
:vhost => Postal.config.rabbitmq&.vhost || nil
)
conn.start
conn
end
def self.create_channel
conn = self.create_connection
conn.create_channel
end
end
end
+33
View File
@@ -0,0 +1,33 @@
module Postal
class ReplySeparator
RULES = [
/^-{2,10} $.*/m,
/^\>*\s*----- ?Original Message ?-----.*/m,
/^\>*\s*From\:[^\r\n]*[\r\n]+Sent\:.*/m,
/^\>*\s*From\:[^\r\n]*[\r\n]+Date\:.*/m,
/^\>*\s*-----Urspr.ngliche Nachricht----- .*/m,
/^\>*\s*Le[^\r\n]{10,200}a .crit ?\:\s*$.*/,
/^\>*\s*__________________.*/m,
/^\>*\s*On.{10,200}wrote:\s*$.*/m,
/^\>*\s*Sent from my.*/m,
/^\>*\s*=== Please reply above this line ===.*/m,
/(^\>.*\n?){10,}/,
]
def self.separate(text)
return '' unless text.is_a?(String)
text = text.gsub("\r", '')
stripped = ""
RULES.each do |rule|
text.gsub!(rule) do
stripped = $&.to_s + "\n" + stripped
''
end
end
stripped = stripped.strip
[text.strip, stripped.blank? ? nil : stripped]
end
end
end
+12
View File
@@ -0,0 +1,12 @@
module Postal
class SendResult
attr_accessor :type
attr_accessor :details
attr_accessor :retry
attr_accessor :output
attr_accessor :secure
attr_accessor :connect_error
attr_accessor :log_id
attr_accessor :time
end
end
+12
View File
@@ -0,0 +1,12 @@
module Postal
class Sender
def start
end
def send_message(message)
end
def finish
end
end
end
+246
View File
@@ -0,0 +1,246 @@
require 'resolv'
module Postal
class SMTPSender < Sender
def initialize(domain, source_ip_address, options = {})
@domain = domain
@source_ip_address = source_ip_address
@options = options
@smtp_client = nil
@connection_errors = []
@hostnames = []
@log_id = Nifty::Utils::RandomString.generate(:length => 8).upcase
end
def start
servers.each do |server|
if server.is_a?(SMTPEndpoint)
hostname = server.hostname
port = server.port || 25
ssl_mode = server.ssl_mode
else
hostname = server
port = 25
ssl_mode = 'Auto'
end
@hostnames << hostname
[:aaaa, :a].each do |ip_type|
if @source_ip_address && @source_ip_address.ipv6.blank? && ip_type == :aaaa
# Don't try to use IPv6 if the IP address we're sending from doesn't support it.
next
end
begin
@remote_ip = lookup_ip_address(ip_type, hostname)
if @remote_ip.nil?
if ip_type == :a
# As we can't resolve the last IP, we'll put this
@connection_errors << "Could not resolve #{hostname}"
end
next
end
smtp_client = Net::SMTP.new(@remote_ip, port)
if @source_ip_address
# Set the source IP as appropriate
smtp_client.source_address = ip_type == :aaaa ? @source_ip_address.ipv6 : @source_ip_address.ipv4
end
case ssl_mode
when 'Auto'
smtp_client.enable_starttls_auto(self.class.ssl_context_without_verify)
when 'STARTTLS'
smtp_client.enable_starttls(self.class.ssl_context_with_verify)
when 'TLS'
smtp_client.enable_tls(self.class.ssl_context_with_verify)
else
# Nothing
end
smtp_client.start(@source_ip_address ? @source_ip_address.hostname : "localhost")
log "Connected to #{@remote_ip}:#{port} (#{hostname})"
rescue => e
log "Cannot connect to #{@remote_ip}:#{port} (#{hostname}) (#{e.class}: #{e.message})"
@connection_errors << e.message unless @connection_errors.include?(e.message)
smtp_client.disconnect rescue nil
smtp_client = nil
end
if smtp_client
@smtp_client = smtp_client
return true
end
end
end
@connection_errors
end
def reconnect
log "Reconnecting"
@smtp_client&.finish rescue nil
start
end
def safe_rset
# Something went wrong sending the last email. Reset the connection if possible, else disconnect.
begin
@smtp_client.rset
rescue
# Don't reconnect, this would be rather rude if we don't have any more emails to send.
@smtp_client.finish rescue nil
end
end
def send_message(message, force_rcpt_to = nil)
start_time = Time.now
result = SendResult.new
result.log_id = @log_id
if @smtp_client && !@smtp_client.started?
# For some reason we had an SMTP connection but it's no longer connected.
# Make a new one.
start
end
if @smtp_client
result.secure = @smtp_client.secure_socket?
end
begin
if message.bounce == 1
mail_from = ""
elsif message.domain.return_path_status == 'OK'
mail_from = "#{message.server.token}@#{message.domain.return_path_domain}"
else
mail_from = "#{message.server.token}@#{Postal.config.dns.return_path}"
end
raw_message = "Resent-Sender: #{mail_from}\r\n" + message.raw_message
tries = 0
begin
if @smtp_client.nil?
log "-> No SMTP server available for #{@domain}"
log "-> Hostnames: #{@hostnames.inspect}"
log "-> Errors: #{@connection_errors.inspect}"
result.type = 'SoftFail'
result.retry = true
result.details = "No SMTP servers were available for #{@domain}. Tried #{@hostnames.to_sentence}"
result.output = @connection_errors.join(', ')
result.connect_error = true
return result
else
@smtp_client.rset_errors
rcpt_to = force_rcpt_to || @options[:force_rcpt_to] || message.rcpt_to
log "Sending message #{message.server.id}::#{message.id} to #{rcpt_to}"
smtp_result = @smtp_client.send_message(raw_message, mail_from, [rcpt_to])
end
rescue Errno::ECONNRESET, Errno::EPIPE, OpenSSL::SSL::SSLError
if (tries += 1) < 2
reconnect
retry
else
raise
end
end
result.type = 'Sent'
result.details = "Message for #{rcpt_to} accepted by #{destination_host_description}"
if @smtp_client.source_address
result.details += " (from #{@smtp_client.source_address})"
end
result.output = smtp_result.string
log "Message sent ##{message.id} to #{destination_host_description} for #{rcpt_to}"
rescue Net::SMTPServerBusy, Net::SMTPAuthenticationError, Net::SMTPSyntaxError, Net::SMTPUnknownError, Net::ReadTimeout => e
log "#{e.class}: #{e.message}"
result.type = 'SoftFail'
result.retry = true
result.details = "Temporary SMTP delivery error when sending to #{destination_host_description}"
result.output = e.message
if e.to_s =~ /(\d+) seconds/
result.retry = $1.to_i + 10
elsif e.to_s =~ /(\d+) minutes/
result.retry = ($1.to_i * 60) + 10
end
safe_rset
rescue Net::SMTPFatalError => e
log "#{e.class}: #{e.message}"
result.type = 'HardFail'
result.details = "Permanent SMTP delivery error when sending to #{destination_host_description}"
result.output = e.message
safe_rset
rescue => e
log "#{e.class}: #{e.message}"
Raven.capture_exception(e, :extra => {:log_id => @log_id, :server_id => message.server.id, :message_id => message.id})
result.type = 'SoftFail'
result.retry = true
result.details = "An error occurred while sending the message to #{destination_host_description}"
result.output = e.message
safe_rset
end
result.time = (Time.now - start_time).to_f.round(2)
return result
ensure
end
def finish
log "Finishing up"
@smtp_client&.finish
end
private
def servers
@options[:servers] || @servers ||= begin
mx_servers = []
Resolv::DNS.open do |dns|
dns.timeouts = [10,5]
mx_servers = dns.getresources(@domain, Resolv::DNS::Resource::IN::MX).map { |m| [m.preference.to_i, m.exchange.to_s] }.sort.map{ |m| m[1] }
if mx_servers.empty?
mx_servers = [@domain] # This will be resolved to an A or AAAA record later
end
end
mx_servers
end
end
def log(text)
Postal.logger_for(:smtp_sender).info "[#{@log_id}] #{text}"
end
def destination_host_description
"#{@hostnames.last} (#{@remote_ip})"
end
def lookup_ip_address(type, hostname)
records = []
Resolv::DNS.open do |dns|
dns.timeouts = [10,5]
case type
when :a
records = dns.getresources(hostname, Resolv::DNS::Resource::IN::A)
when :aaaa
records = dns.getresources(hostname, Resolv::DNS::Resource::IN::AAAA)
end
end
records.first&.address&.to_s&.downcase
end
def self.ssl_context_with_verify
@ssl_context_with_verify ||= begin
c = OpenSSL::SSL::SSLContext.new
c.verify_mode = OpenSSL::SSL::VERIFY_PEER
c
end
end
def self.ssl_context_without_verify
@ssl_context_without_verify ||= begin
c = OpenSSL::SSL::SSLContext.new
c.verify_mode = OpenSSL::SSL::VERIFY_NONE
c
end
end
end
end
+435
View File
@@ -0,0 +1,435 @@
require 'nifty/utils/random_string'
module Postal
module SMTPServer
class Client
CRAM_MD5_DIGEST = OpenSSL::Digest.new('md5')
attr_reader :logging_enabled
def initialize(ip_address)
@logging_enabled = true
@ip_address = ip_address
if @ip_address
check_ip_address
@state = :welcome
else
@state = :preauth
end
reset
end
def check_ip_address
if @ip_address && Postal.config.smtp_server.log_exclude_ips && @ip_address =~ Regexp.new(Postal.config.smtp_server.log_exclude_ips)
@logging_enabled = false
end
end
def transaction_reset
@recipients = []
@mail_from = nil
@data = nil
@headers = nil
end
def reset
@credential = nil
transaction_reset
end
def id
@id ||= Nifty::Utils::RandomString.generate(:length => 6).upcase
end
def handle(data)
if @state == :preauth
proxy(data)
else
if @proc
log "\e[32m<= #{data.strip}\e[0m"
@proc.call(data)
else
log "\e[32m<= #{data.strip}\e[0m"
handle_command(data)
end
end
end
def finished?
@finished || false
end
def start_tls?
@start_tls || false
end
def start_tls=(value)
@start_tls = value
end
def handle_command(data)
case data
when /^QUIT/i then quit
when /^STARTTLS/i then starttls
when /^EHLO/i then ehlo(data)
when /^HELO/i then helo(data)
when /^RSET/i then rset
when /^NOOP/i then noop
when /^AUTH PLAIN/i then auth_plain(data)
when /^AUTH LOGIN/i then auth_login(data)
when /^AUTH CRAM-MD5/i then auth_cram_md5(data)
when /^MAIL FROM/i then mail_from(data)
when /^RCPT TO/i then rcpt_to(data)
when /^DATA/i then data(data)
else
'502 Invalid/unsupported command'
end
end
def log(text)
return false unless @logging_enabled
Postal.logger_for(:smtp_server).debug "[#{id}] #{text}"
end
private
def resolve_hostname
@hostname = Resolv.new.getname(@ip_address) rescue @ip_address
end
def proxy(data)
if m = data.match(/\APROXY (.+) (.+) (.+) (.+) (.+)\z/)
@ip_address = m[2]
check_ip_address
@state = :welcome
log "\e[35m Client identified as #{@ip_address}\e[0m"
"220 #{Postal.config.dns.smtp_server_hostname} ESMTP Postal/#{id}"
else
@finished = true
'502 Proxy Error'
end
end
def quit
@finished = true
"221 Closing Connection"
end
def starttls
@start_tls = true
@tls = true
"220 Ready to start TLS"
end
def ehlo(data)
resolve_hostname
@helo_name = data.strip.split(' ', 2)[1]
reset
@state = :welcomed
["250-My capabilities are", @tls ? nil : "250-STARTTLS", "250 AUTH CRAM-MD5 PLAIN LOGIN", ]
end
def helo(data)
resolve_hostname
@helo_name = data.strip.split(' ', 2)[1]
reset
@state = :welcomed
"250 #{Postal.config.dns.smtp_server_hostname}"
end
def rset
reset
@state = :welcomed
'250 OK'
end
def noop
'250 OK'
end
def auth_plain(data)
handler = Proc.new do |data|
@proc = nil
data = Base64.decode64(data)
parts = data.split("\0")
username, password = parts[-2], parts[-1]
unless username && password
next '535 Authenticated failed - protocol error'
end
authenticate(password)
end
data = data.gsub(/AUTH PLAIN ?/i, '')
if data.strip == ''
@proc = handler
'334'
else
handler.call(data)
end
end
def auth_login(data)
password_handler = Proc.new do |data|
@proc = nil
password = Base64.decode64(data)
authenticate(password)
end
username_handler = Proc.new do |data|
@proc = password_handler
'334 UGFzc3dvcmQ6'
end
data = data.gsub!(/AUTH LOGIN ?/i, '')
if data.strip == ''
@proc = username_handler
'334 VXNlcm5hbWU6'
end
end
def authenticate(password)
if @credential = Credential.where(:type => 'SMTP', :key => password).first
@credential.use
"235 Granted for #{@credential.server.organization.permalink}/#{@credential.server.permalink}"
else
"535 Invalid credential"
end
end
def auth_cram_md5(data)
challenge = Digest::SHA1.hexdigest(Time.now.to_i.to_s + rand(100000).to_s)
challenge = "<#{challenge[0,20]}@#{Postal.config.dns.smtp_server_hostname}>"
handler = Proc.new do |data|
@proc = nil
username, password = Base64.decode64(data).split(' ', 2).map{ |a| a.chomp }
org_permlink, server_permalink = username.split(/[\/\_]/, 2)
server = ::Server.includes(:organization).where(:organizations => {:permalink => org_permlink}, :permalink => server_permalink).first
next '535 Denied' if server.nil?
grant = nil
server.credentials.where(:type => 'SMTP').each do |credential|
correct_response = OpenSSL::HMAC.hexdigest(CRAM_MD5_DIGEST, credential.key, challenge)
if password == correct_response
@credential = credential
@credential.use
grant = "235 Granted for #{credential.server.organization.permalink}/#{credential.server.permalink}"
break
end
end
grant || '535 Denied'
end
@proc = handler
"334 " + Base64.encode64(challenge).gsub(/[\r\n]/, '')
end
def mail_from(data)
unless in_state(:welcomed, :mail_from_received)
return '503 EHLO/HELO first please'
end
@state = :mail_from_received
transaction_reset
@mail_from = data.gsub(/MAIL FROM\s*:\s*/i, '').gsub(/.*</, '').gsub(/>.*/, '').strip
'250 OK'
end
def rcpt_to(data)
unless in_state(:mail_from_received, :rcpt_to_received)
return '503 EHLO/HELO and MAIL FROM first please'
end
rcpt_to = data.gsub(/RCPT TO\s*:\s*/i, '').gsub(/.*</, '').gsub(/>.*/, '').strip
uname, domain = rcpt_to.split('@', 2)
uname, tag = uname.split('+', 2)
if domain =~ /\A#{Regexp.escape(Postal.config.dns.custom_return_path_prefix)}\./
# This is a return path
@state = :rcpt_to_received
if server = ::Server.where(:token => uname).first
if server.suspended?
'535 Mail server has been suspended'
else
log "Added bounce on server #{server.id}"
@recipients << [:bounce, rcpt_to, server]
'250 OK'
end
else
'550 Invalid server token'
end
elsif domain == Postal.config.dns.route_domain
# This is an email direct to a route. This isn't actually supported yet.
@state = :rcpt_to_received
if route = Route.where(:token => uname).first
if route.server.suspended?
'535 Mail server has been suspended'
elsif route.mode == 'Reject'
'550 Route does not accept incoming messages'
else
log "Added route #{route.id} to recipients (tag: #{tag.inspect})"
actual_rcpt_to = "#{route.name}" + (tag ? "+#{tag}" : "") + "@#{route.domain.name}"
@recipients << [:route, actual_rcpt_to, route.server, :route => route]
'250 OK'
end
else
'550 Invalid route token'
end
elsif @credential
# This is outgoing mail for an authenticated user
@state = :rcpt_to_received
if @credential.server.suspended?
'535 Mail server has been suspended'
else
log "Added external address '#{rcpt_to}'"
@recipients << [:credential, rcpt_to, @credential.server]
'250 OK'
end
elsif uname && domain && route = Route.find_by_name_and_domain(uname, domain)
# This is incoming mail for a route
@state = :rcpt_to_received
if route.server.suspended?
'535 Mail server has been suspended'
elsif route.mode == 'Reject'
'550 Route does not accept incoming messages'
else
log "Added route #{route.id} to recipients (tag: #{tag.inspect})"
@recipients << [:route, rcpt_to, route.server, :route => route]
'250 OK'
end
else
# This is unaccepted mail
'530 Authentication required'
end
end
def data(data)
unless in_state(:rcpt_to_received)
return '503 HELO/EHLO, MAIL FROM and RCPT TO before sending data'
end
@data = "".force_encoding("BINARY")
@headers = {}
@receiving_headers = true
received_header_content = "from #{@helo_name} (#{@hostname} [#{@ip_address}]) by #{Postal.config.dns.smtp_server_hostname} with SMTP; #{Time.now.rfc2822.to_s}".force_encoding('BINARY')
@data << "Received: #{received_header_content}\r\n"
@headers['received'] = [received_header_content]
handler = Proc.new do |data|
if data == '.'
@logging_enabled = true
@proc = nil
finished
else
data = data.to_s.sub(/\A\.\./, '.')
if @credential && @credential.server.log_smtp_data?
# We want to log if enabled
else
log "Not logging further message data."
@logging_enabled = false
end
if @receiving_headers
if data.blank?
@receiving_headers = false
elsif data.to_s =~ /^\s/
# This is a continuation of a header
if @header_key && @headers[@header_key.downcase] && @headers[@header_key.downcase].last
@headers[@header_key.downcase].last << data.to_s
end
else
@header_key, value = data.split(/\:\s*/, 2)
@headers[@header_key.downcase] ||= []
@headers[@header_key.downcase] << value
end
end
@data << data
@data << "\r\n"
nil
end
end
@proc = handler
'354 Go ahead'
end
def finished
if @data.bytesize > 14.megabytes.to_i
return "552 Message too large (maximum size 14MB)"
end
if @headers['received'].select { |r| r =~ /by #{Postal.config.dns.smtp_server_hostname}/ }.count > 4
return '550 Loop detected'
end
authenticated_domain = nil
if @credential
authenticated_domain = @credential.server.find_authenticated_domain_from_headers(@headers)
if authenticated_domain.nil?
return '530 From/Sender name is not valid'
end
end
@recipients.each do |recipient|
type, rcpt_to, server, options = recipient
case type
when :credential
# Outgoing messages are just inserted
message = server.message_db.new_message
message.rcpt_to = rcpt_to
message.mail_from = @mail_from
message.raw_message = @data
message.received_with_ssl = @tls
message.scope = 'outgoing'
message.domain_id = authenticated_domain&.id
message.credential_id = @credential.id
message.save
when :bounce
if rp_route = server.routes.where(:name => "__returnpath__").first
# If there's a return path route, we can use this to create the message
rp_route.create_messages do |message|
message.rcpt_to = rcpt_to
message.mail_from = @mail_from
message.raw_message = @data
message.received_with_ssl = @tls
end
else
# There's no return path route, we just need to insert the mesage
# without going through the route.
message = server.message_db.new_message
message.rcpt_to = rcpt_to
message.mail_from = @mail_from
message.raw_message = @data
message.received_with_ssl = @tls
message.scope = 'incoming'
message.bounce = 1
message.save
end
when :route
options[:route].create_messages do |message|
message.rcpt_to = rcpt_to
message.mail_from = @mail_from
message.raw_message = @data
message.received_with_ssl = @tls
end
end
end
transaction_reset
'250 OK'
end
def in_state(*states)
states.include?(@state)
end
end
end
end
+322
View File
@@ -0,0 +1,322 @@
require 'ipaddr'
require 'epoll' if RUBY_PLATFORM.include?('linux')
module Postal
module SMTPServer
class Server
def initialize(options = {})
@options = options
@options[:ports] ||= Postal.config.smtp_server.ports
@options[:debug] ||= false
prepare_environment
end
def prepare_environment
$\ = "\r\n"
BasicSocket.do_not_reverse_lookup = true
trap("USR1") do
STDOUT.puts "Received USR1 signal, respawning."
fork do
if ENV['APP_ROOT']
Dir.chdir(ENV['APP_ROOT'])
end
ENV.delete('BUNDLE_GEMFILE')
exec("bundle exec --keep-file-descriptors rake postal:smtp_server", :close_others => false)
end
end
trap("TERM") do
STDOUT.puts "Received TERM signal, shutting down."
unlisten
end
end
def ssl_context
@ssl_context ||= begin
ssl_context = OpenSSL::SSL::SSLContext.new
certs = Postal.ssl_certificates
ssl_context.cert = certs.shift
ssl_context.extra_chain_cert = certs
ssl_context.key = Postal.signing_key
ssl_context.ssl_version = "SSLv23"
ssl_context
end
end
def listen
if ENV['SERVER_FD']
@server = TCPServer.for_fd(ENV['SERVER_FD'].to_i)
else
@server = TCPServer.open('::', @options[:ports].first)
end
@server.autoclose = false
@server.close_on_exec = false
if defined?(Socket::SOL_SOCKET) && defined?(Socket::SO_KEEPALIVE)
@server.setsockopt(Socket::SOL_SOCKET, Socket::SO_KEEPALIVE, true)
end
if defined?(Socket::SOL_TCP) && defined?(Socket::TCP_KEEPIDLE) && defined?(Socket::TCP_KEEPINTVL) && defined?(Socket::TCP_KEEPCNT)
@server.setsockopt(Socket::SOL_TCP, Socket::TCP_KEEPIDLE, 50)
@server.setsockopt(Socket::SOL_TCP, Socket::TCP_KEEPINTVL, 10)
@server.setsockopt(Socket::SOL_TCP, Socket::TCP_KEEPCNT, 5)
end
ENV['SERVER_FD'] = @server.to_i.to_s
end
def unlisten
if @epoll
@epoll.del(@server)
if @epoll.size == 0
Process.exit(0)
end
end
@server.close
end
def kill_parent
Process.kill('TERM', Process.ppid)
end
def run_linux
if ENV['SERVER_FD']
listen
kill_parent
else
listen
end
@epoll = Epoll.create
logger.info "Listening"
@epoll.add(@server, Epoll::IN)
buffers = Hash.new { |h, k| h[k] = String.new.force_encoding('BINARY') }
clients = {}
loop do
evlist = @epoll.wait
evlist.each do |ev|
io = ev.data
if io.is_a?(TCPServer)
begin
new_io = io.accept
if Postal.config.smtp_server.proxy_protocol
client = Client.new(nil)
if Postal.config.smtp_server.log_connect
logger.debug "[#{client.id}] \e[35m Connection opened from #{new_io.remote_address.ip_address}\e[0m"
end
else
client = Client.new(new_io.remote_address.ip_address)
if Postal.config.smtp_server.log_connect
logger.debug "[#{client.id}] \e[35m Connection opened from #{new_io.remote_address.ip_address}\e[0m"
end
client.log "\e[35m Client identified as #{new_io.remote_address.ip_address}\e[0m"
new_io.print("220 #{Postal.config.dns.smtp_server_hostname} ESMTP Postal/#{client.id}")
end
clients[new_io] = client
@epoll.add(new_io, Epoll::IN|Epoll::PRI|Epoll::HUP)
rescue => e
Raven.capture_exception(e, :extra => {:log_id => (client.id rescue nil)})
logger.error "An error occurred while accepting a new client."
logger.error "#{e.class}: #{e.message}"
e.backtrace.each do |line|
logger.error line
end
new_io.close rescue nil
end
else
begin
client = clients[io]
eof = false
begin
case io
when OpenSSL::SSL::SSLSocket
buffers[io] << io.readpartial(10240)
while(io.pending > 0)
buffers[io] << io.readpartial(10240)
end
else
buffers[io] << io.readpartial(10240)
end
rescue EOFError, Errno::ECONNRESET
# Client went away
eof = true
end
while buffers[io].index("\n")
if buffers[io].index("\r\n")
line, buffers[io] = buffers[io].split("\r\n", 2)
else
line, buffers[io] = buffers[io].split("\n", 2)
end
result = client.handle(line)
unless result.nil?
result = [result] unless result.is_a?(Array)
result.compact.each do |line|
client.log "\e[34m=> #{line.strip}\e[0m"
begin
io.write(line.to_s + "\r\n")
io.flush
rescue Errno::ECONNRESET
# Client disconnected before we could write response
eof = true
end
end
end
end
if !eof && client.start_tls?
client.start_tls = false
@epoll.del(io)
clients.delete(io)
buffers.delete(io)
tcp_io = io
io = OpenSSL::SSL::SSLSocket.new(io, ssl_context)
@epoll.add(io, Epoll::IN)
clients[io] = client
io.sync_close = true
begin
io.accept
rescue OpenSSL::SSL::SSLError => e
client.log "SSL Negotiation Failed: #{e.message}"
eof = true
end
end
if client.finished? || eof
client.log "\e[35m Connection closed\e[0m"
@epoll.del(io)
clients.delete(io)
buffers.delete(io)
io.close
if @epoll.size == 0
Process.exit(0)
end
end
rescue => e
client_id = client ? client.id : '------'
Raven.capture_exception(e, :extra => {:log_id => (client.id rescue nil)})
logger.error "[#{client_id}] An error occurred while processing data from a client."
logger.error "[#{client_id}] #{e.class}: #{e.message}"
e.backtrace.each do |line|
logger.error "[#{client_id}] #{line}"
end
# Close all IO and forget this client
@epoll.del(io) rescue nil
clients.delete(io)
buffers.delete(io)
io.close rescue nil
if @epoll.size == 0
Process.exit(0)
end
end
end
end
end
end
def run_non_linux
if ENV['SERVER_FD']
listen
kill_parent
else
listen
end
logger.info "Listening"
Thread.abort_on_exception = true
client_threads = []
loop do
s = nil
begin
until s
l = select([@server], [@server], [@server], 0.5)
s = @server.accept if l
end
rescue IOError
STDERR.puts "Server socket was closed."
break
end
client_threads << Thread.new(s) do |io|
begin
if Postal.config.smtp_server.proxy_protocol
client = Client.new(nil)
if Postal.config.smtp_server.log_connect
logger.debug "[#{client.id}] \e[35m Connection opened from #{io.remote_address.ip_address}\e[0m"
end
else
client = Client.new(io.remote_address.ip_address)
if Postal.config.smtp_server.log_connect
logger.debug "[#{client.id}] \e[35m Connection opened from #{io.remote_address.ip_address}\e[0m"
end
client.log "\e[35m Client identified as #{io.remote_address.ip_address}\e[0m"
io.print("220 #{Postal.config.dns.smtp_server_hostname} ESMTP Postal/#{client.id}")
end
loop do
if received_data = io.gets
if result = client.handle(received_data.chomp)
result = [result] unless result.is_a?(Array)
result.compact.each do |line|
client.log "\e[34m=> #{line.strip}\e[0m"
io.write(line.to_s + "\r\n")
io.flush
end
end
end
if client.start_tls?
client.start_tls = false
tcp_io = io
io = OpenSSL::SSL::SSLSocket.new(io, ssl_context)
io.sync_close = true
begin
io.accept
rescue OpenSSL::SSL::SSLError => e
logger.error "SSL Negotiation Failed: #{e.message}"
io.close rescue nil
tcp_io.close rescue nil
eof = true
end
end
if received_data.nil? || client.finished?
client.log "\e[35m Connection closed\e[0m"
io.close
break
end
end
rescue => e
Raven.capture_exception(e, :extra => {:log_id => (client.id rescue nil)})
logger.error "An error occurred while handling a client."
logger.error "#{e.class}: #{e.message}"
e.backtrace.each do |line|
logger.error line
end
# Close all IO
io.close rescue nil
ensure
client_threads.delete(Thread.current)
end
end
end
client_threads.each{ |t| t.join unless t == Thread.current }
end
def run
if ENV['PID_FILE']
File.open(ENV['PID_FILE'], 'w') { |f| f.write(Process.pid.to_s + "\n") }
end
if Postal.config.smtp_server&.evented
logger.info "Running epoll driven server for Linux host.."
run_linux
else
logger.info "Running thread based compatibility server for non-Linux host."
run_non_linux
end
end
private
def logger
Postal.logger_for(:smtp_server)
end
end
end
end
+3
View File
@@ -0,0 +1,3 @@
module Postal
VERSION = '1.0.0'
end
+180
View File
@@ -0,0 +1,180 @@
module Postal
class Worker
def initialize(queues)
@initial_queues = queues
@active_queues = {}
@process_name = $0
end
def work
@running_job = false
Signal.trap("INT") { @exit = true }
Signal.trap("TERM") { @exit = true }
self.class.job_channel.prefetch(1)
@initial_queues.each { |queue | join_queue(queue) }
exit_checks = 0
loop do
if @exit && @running_job == false
logger.info "Exiting immediately because no job running"
exit 0
elsif @exit
if exit_checks >= 60
logger.info "Job did not finish in a timely manner. Exiting"
exit 0
end
if exit_checks == 0
logger.info "Exit requested but job is running. Waiting for job to finish."
end
sleep 60
exit_checks += 1
else
manage_ip_queues
sleep 1
end
end
end
private
def receive_job(delivery_info, properties, body)
@running_job = true
begin
message = JSON.parse(body) rescue nil
if message && message['class_name']
start_time = Time.now
$0 = "#{@process_name} (running #{message['class_name']})"
Thread.current[:job_id] = message['id']
logger.info "[#{message['id']}] Started processing \e[34m#{message['class_name']}\e[0m job"
begin
klass = message['class_name'].constantize.new(message['id'], message['params'])
klass.perform
rescue => e
Raven.capture_exception(e, :extra => {:job_id => message['id']})
logger.warn "[#{message['id']}] \e[31m#{e.class}: #{e.message}\e[0m"
e.backtrace.each do |line|
logger.warn "[#{message['id']}] " + line
end
ensure
logger.info "[#{message['id']}] Finished processing \e[34m#{message['class_name']}\e[0m job in #{Time.now - start_time}s"
end
end
ensure
Thread.current[:job_id] = nil
$0 = @process_name
self.class.job_channel.ack(delivery_info.delivery_tag)
@running_job = false
if @exit
logger.info "Exiting because a job has ended."
exit 0
end
end
end
def join_queue(queue)
if @active_queues[queue]
logger.info "Attempted to join queue #{queue} but already joined."
else
consumer = self.class.job_queue(queue).subscribe(:manual_ack => true) do |delivery_info, properties, body|
receive_job(delivery_info, properties, body)
end
@active_queues[queue] = consumer
logger.info "Joined \e[32m#{queue}\e[0m queue"
end
end
def leave_queue(queue)
if consumer = @active_queues[queue]
consumer.cancel
@active_queues.delete(queue)
logger.info "Left \e[32m#{queue}\e[0m queue"
else
logger.info "Not joined #{queue} so cannot leave"
end
end
def manage_ip_queues
@ip_queues ||= []
@ip_to_id_mapping ||= {}
@unassigned_ips ||= []
@pairs ||= {}
# Get all IP addresses on the system
current_ip_addresses = Socket.ip_address_list.map(&:ip_address)
# Map them to an actual ID in the database if we can and cache that
needed_ip_ids = []
current_ip_addresses.each do |ip|
need = nil
if id = @ip_to_id_mapping[ip]
# We know this IPs ID, we'll just use that.
need = id
elsif @unassigned_ips.include?(ip)
# We know this IP isn't valid. We don't need to do anything
else
# We need to look this up
if ip_address = IPAddress.where("ipv4 = ? OR ipv6 = ?", ip, ip).first
@pairs[ip_address.ipv4] = ip_address.ipv6
@ip_to_id_mapping[ip] = ip_address.id
need = id
else
@unassigned_ips << ip
end
end
if need
pair = @pairs[ip] || @pairs.key(ip)
if current_ip_addresses.include?(pair)
needed_ip_ids << id
else
logger.info "Host has '#{ip}' but its pair (#{pair}) isn't here. Cannot add now."
end
end
end
# Make an array of needed queue names
# Work out what we need to actually do here
missing_queues = needed_ip_ids - @ip_queues
unwanted_queues = @ip_queues - needed_ip_ids
# Leave the queues we don't want any more
unwanted_queues.each do |id|
leave_queue("outgoing-#{id}")
@ip_queues.delete(id)
ip_addresses_to_clear = []
@ip_to_id_mapping.each do |_ip, _id|
if id == _id
ip_addresses_to_clear << _ip
end
end
ip_addresses_to_clear.each { |ip| @ip_to_id_mapping.delete(ip) }
end
# Join any missing queues
missing_queues.uniq.each do |id|
join_queue("outgoing-#{id}")
@ip_queues << id
end
end
def logger
self.class.logger
end
def self.logger
Postal.logger_for(:worker)
end
def self.job_channel
@channel ||= Postal::RabbitMQ.create_channel
end
def self.job_queue(name)
@job_queues ||= {}
@job_queues[name] ||= begin
job_channel.queue("deliver-jobs-#{name}", :durable => true, :arguments => {'x-message-ttl' => 60000})
end
end
end
end
View File
+48
View File
@@ -0,0 +1,48 @@
# NOTE: only doing this in development as some production environments (Heroku)
# NOTE: are sensitive to local FS writes, and besides -- it's just not proper
# NOTE: to have a dev-mode tool do its thing in production.
if Rails.env.development?
task :set_annotation_options do
# You can override any of these by setting an environment variable of the
# same name.
Annotate.set_defaults(
'routes' => 'false',
'position_in_routes' => 'before',
'position_in_class' => 'before',
'position_in_test' => 'before',
'position_in_fixture' => 'before',
'position_in_factory' => 'before',
'position_in_serializer' => 'before',
'show_foreign_keys' => 'true',
'show_indexes' => 'true',
'simple_indexes' => 'false',
'model_dir' => 'app/models',
'root_dir' => '',
'include_version' => 'false',
'require' => '',
'exclude_tests' => 'false',
'exclude_fixtures' => 'false',
'exclude_factories' => 'false',
'exclude_serializers' => 'false',
'exclude_scaffolds' => 'true',
'exclude_controllers' => 'true',
'exclude_helpers' => 'true',
'ignore_model_sub_dir' => 'false',
'ignore_columns' => nil,
'ignore_routes' => nil,
'ignore_unknown_models' => 'false',
'hide_limit_column_types' => 'integer,boolean',
'skip_on_db_migrate' => 'false',
'format_bare' => 'true',
'format_rdoc' => 'false',
'format_markdown' => 'false',
'sort' => 'false',
'force' => 'false',
'trace' => 'false',
'wrapper_open' => nil,
'wrapper_close' => nil
)
end
Annotate.load_tasks
end
+41
View File
@@ -0,0 +1,41 @@
namespace :postal do
desc "Start the backend job worker"
task :worker => :environment do
Postal::Worker.new([:main]).work
end
desc "Start the cron worker"
task :cron => :environment do
require 'clockwork'
require Rails.root.join('config', 'cron')
trap('TERM') { puts "Exiting..."; Process.exit(0) }
Clockwork.run
end
desc 'Start SMTP Server'
task :smtp_server => :environment do
Postal::SMTPServer::Server.new(:debug => true).run
end
desc 'Start the message requeuer'
task :requeuer => :environment do
Postal::MessageRequeuer.new.run
end
desc 'Run all migrations on message databases'
task :migrate_message_databases => :environment do
Server.all.each do |server|
puts "\e[35m-------------------------------------------------------------------\e[0m"
puts "\e[35m#{server.id}: #{server.name} (#{server.permalink})\e[0m"
puts "\e[35m-------------------------------------------------------------------\e[0m"
server.message_db.provisioner.migrate
end
end
desc 'Start the fast server'#
task :fast_server => :environment do
Postal::FastServer::Server.new.run
end
end