初始化提交,包含完整的邮件系统代码
This commit is contained in:
74
app/__init__.py
Normal file
74
app/__init__.py
Normal file
@@ -0,0 +1,74 @@
|
||||
import os
|
||||
import logging
|
||||
from flask import Flask
|
||||
from flask_cors import CORS
|
||||
|
||||
# 修改相对导入为绝对导入
|
||||
import sys
|
||||
sys.path.append(os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
|
||||
from config import active_config
|
||||
|
||||
|
||||
def setup_logging(app):
|
||||
"""设置日志"""
|
||||
log_level = getattr(logging, active_config.LOG_LEVEL.upper(), logging.INFO)
|
||||
|
||||
# 确保日志目录存在
|
||||
log_dir = os.path.dirname(active_config.LOG_FILE)
|
||||
if not os.path.exists(log_dir):
|
||||
os.makedirs(log_dir)
|
||||
|
||||
# 配置日志
|
||||
logging.basicConfig(
|
||||
level=log_level,
|
||||
format='%(asctime)s - %(name)s - %(levelname)s - %(message)s',
|
||||
handlers=[
|
||||
logging.FileHandler(active_config.LOG_FILE),
|
||||
logging.StreamHandler()
|
||||
]
|
||||
)
|
||||
|
||||
app.logger.setLevel(log_level)
|
||||
return app
|
||||
|
||||
|
||||
def create_app(config=None):
|
||||
"""创建并配置Flask应用"""
|
||||
app = Flask(__name__)
|
||||
|
||||
# 加载配置
|
||||
app.config.from_object(active_config)
|
||||
|
||||
# 如果提供了自定义配置,加载它
|
||||
if config:
|
||||
app.config.from_object(config)
|
||||
|
||||
# 允许跨域请求
|
||||
CORS(app)
|
||||
|
||||
# 设置日志
|
||||
app = setup_logging(app)
|
||||
|
||||
# 确保存储邮件的目录存在
|
||||
os.makedirs(active_config.MAIL_STORAGE_PATH, exist_ok=True)
|
||||
|
||||
# 初始化数据库
|
||||
from .models import init_db
|
||||
init_db()
|
||||
|
||||
# 注册蓝图
|
||||
from .api import api_bp
|
||||
app.register_blueprint(api_bp)
|
||||
|
||||
# 首页路由
|
||||
@app.route("/")
|
||||
def index():
|
||||
return {
|
||||
"name": "Email System",
|
||||
"version": "1.0.0",
|
||||
"status": "running"
|
||||
}
|
||||
|
||||
app.logger.info('应用初始化完成')
|
||||
|
||||
return app
|
||||
23
app/api/__init__.py
Normal file
23
app/api/__init__.py
Normal file
@@ -0,0 +1,23 @@
|
||||
# API模块初始化文件
|
||||
from flask import Blueprint
|
||||
import logging
|
||||
|
||||
# 创建API蓝图
|
||||
api_bp = Blueprint('api', __name__, url_prefix='/api')
|
||||
|
||||
# 注册默认路由
|
||||
@api_bp.route('/')
|
||||
def index():
|
||||
return {
|
||||
'name': 'Email System API',
|
||||
'version': '1.0.0',
|
||||
'status': 'running'
|
||||
}
|
||||
|
||||
# 导入并合并所有API路由
|
||||
# 为避免可能的文件读取问题,改为从routes.py模块中导入所有路由定义
|
||||
try:
|
||||
from .routes import *
|
||||
except Exception as e:
|
||||
logging.error(f"导入API路由时出错: {str(e)}")
|
||||
raise
|
||||
BIN
app/api/domain_routes.py
Normal file
BIN
app/api/domain_routes.py
Normal file
Binary file not shown.
198
app/api/email_routes.py
Normal file
198
app/api/email_routes.py
Normal file
@@ -0,0 +1,198 @@
|
||||
from flask import request, jsonify, current_app, send_file
|
||||
from io import BytesIO
|
||||
import time
|
||||
|
||||
from . import api_bp
|
||||
from ..models import get_session, Email, Mailbox
|
||||
|
||||
# 获取邮箱的所有邮件
|
||||
@api_bp.route('/mailboxes/<int:mailbox_id>/emails', methods=['GET'])
|
||||
def get_mailbox_emails(mailbox_id):
|
||||
"""获取指定邮箱的所有邮件"""
|
||||
try:
|
||||
page = int(request.args.get('page', 1))
|
||||
limit = int(request.args.get('limit', 50))
|
||||
unread_only = request.args.get('unread_only', 'false').lower() == 'true'
|
||||
offset = (page - 1) * limit
|
||||
|
||||
db = get_session()
|
||||
try:
|
||||
# 检查邮箱是否存在
|
||||
mailbox = db.query(Mailbox).filter_by(id=mailbox_id).first()
|
||||
if not mailbox:
|
||||
return jsonify({'error': '邮箱不存在'}), 404
|
||||
|
||||
# 查询邮件
|
||||
query = db.query(Email).filter(Email.mailbox_id == mailbox_id)
|
||||
|
||||
if unread_only:
|
||||
query = query.filter(Email.read == False)
|
||||
|
||||
# 获取总数
|
||||
total = query.count()
|
||||
|
||||
# 分页获取邮件
|
||||
emails = query.order_by(Email.received_at.desc()) \
|
||||
.limit(limit) \
|
||||
.offset(offset) \
|
||||
.all()
|
||||
|
||||
# 返回结果
|
||||
result = {
|
||||
'total': total,
|
||||
'page': page,
|
||||
'limit': limit,
|
||||
'emails': [email.to_dict() for email in emails]
|
||||
}
|
||||
|
||||
return jsonify(result), 200
|
||||
finally:
|
||||
db.close()
|
||||
except Exception as e:
|
||||
current_app.logger.error(f"获取邮件列表出错: {str(e)}")
|
||||
return jsonify({'error': '获取邮件列表失败', 'details': str(e)}), 500
|
||||
|
||||
# 获取特定邮件详情
|
||||
@api_bp.route('/emails/<int:email_id>', methods=['GET'])
|
||||
def get_email(email_id):
|
||||
"""获取特定邮件的详细信息"""
|
||||
try:
|
||||
mark_as_read = request.args.get('mark_as_read', 'true').lower() == 'true'
|
||||
|
||||
db = get_session()
|
||||
try:
|
||||
email = db.query(Email).filter_by(id=email_id).first()
|
||||
|
||||
if not email:
|
||||
return jsonify({'error': '邮件不存在'}), 404
|
||||
|
||||
# 标记为已读
|
||||
if mark_as_read and not email.read:
|
||||
email.read = True
|
||||
db.commit()
|
||||
|
||||
# 构建详细响应
|
||||
result = email.to_dict()
|
||||
result['body_text'] = email.body_text
|
||||
result['body_html'] = email.body_html
|
||||
|
||||
# 获取附件信息
|
||||
attachments = []
|
||||
for attachment in email.attachments:
|
||||
attachments.append({
|
||||
'id': attachment.id,
|
||||
'filename': attachment.filename,
|
||||
'content_type': attachment.content_type,
|
||||
'size': attachment.size
|
||||
})
|
||||
result['attachments'] = attachments
|
||||
|
||||
return jsonify(result), 200
|
||||
finally:
|
||||
db.close()
|
||||
except Exception as e:
|
||||
current_app.logger.error(f"获取邮件详情出错: {str(e)}")
|
||||
return jsonify({'error': '获取邮件详情失败', 'details': str(e)}), 500
|
||||
|
||||
# 删除邮件
|
||||
@api_bp.route('/emails/<int:email_id>', methods=['DELETE'])
|
||||
def delete_email(email_id):
|
||||
"""删除特定邮件"""
|
||||
try:
|
||||
db = get_session()
|
||||
try:
|
||||
email = db.query(Email).filter_by(id=email_id).first()
|
||||
|
||||
if not email:
|
||||
return jsonify({'error': '邮件不存在'}), 404
|
||||
|
||||
db.delete(email)
|
||||
db.commit()
|
||||
|
||||
return jsonify({'message': '邮件已删除'}), 200
|
||||
except Exception as e:
|
||||
db.rollback()
|
||||
raise
|
||||
finally:
|
||||
db.close()
|
||||
except Exception as e:
|
||||
current_app.logger.error(f"删除邮件出错: {str(e)}")
|
||||
return jsonify({'error': '删除邮件失败', 'details': str(e)}), 500
|
||||
|
||||
# 下载附件
|
||||
@api_bp.route('/attachments/<int:attachment_id>', methods=['GET'])
|
||||
def download_attachment(attachment_id):
|
||||
"""下载特定的附件"""
|
||||
try:
|
||||
from ..models import Attachment
|
||||
|
||||
db = get_session()
|
||||
try:
|
||||
attachment = db.query(Attachment).filter_by(id=attachment_id).first()
|
||||
|
||||
if not attachment:
|
||||
return jsonify({'error': '附件不存在'}), 404
|
||||
|
||||
# 获取附件内容
|
||||
content = attachment.get_content()
|
||||
if not content:
|
||||
return jsonify({'error': '附件内容不可用'}), 404
|
||||
|
||||
# 创建内存文件对象
|
||||
file_obj = BytesIO(content)
|
||||
|
||||
# 返回文件下载响应
|
||||
return send_file(
|
||||
file_obj,
|
||||
mimetype=attachment.content_type,
|
||||
as_attachment=True,
|
||||
download_name=attachment.filename
|
||||
)
|
||||
finally:
|
||||
db.close()
|
||||
except Exception as e:
|
||||
current_app.logger.error(f"下载附件出错: {str(e)}")
|
||||
return jsonify({'error': '下载附件失败', 'details': str(e)}), 500
|
||||
|
||||
# 获取最新邮件 (轮询API)
|
||||
@api_bp.route('/mailboxes/<int:mailbox_id>/poll', methods=['GET'])
|
||||
def poll_new_emails(mailbox_id):
|
||||
"""轮询指定邮箱的新邮件"""
|
||||
try:
|
||||
# 获取上次检查时间
|
||||
last_check = request.args.get('last_check')
|
||||
if last_check:
|
||||
try:
|
||||
last_check_time = float(last_check)
|
||||
except ValueError:
|
||||
return jsonify({'error': '无效的last_check参数'}), 400
|
||||
else:
|
||||
last_check_time = time.time() - 300 # 默认检查最近5分钟
|
||||
|
||||
db = get_session()
|
||||
try:
|
||||
# 检查邮箱是否存在
|
||||
mailbox = db.query(Mailbox).filter_by(id=mailbox_id).first()
|
||||
if not mailbox:
|
||||
return jsonify({'error': '邮箱不存在'}), 404
|
||||
|
||||
# 查询新邮件
|
||||
new_emails = db.query(Email).filter(
|
||||
Email.mailbox_id == mailbox_id,
|
||||
Email.received_at >= time.strftime('%Y-%m-%d %H:%M:%S', time.localtime(last_check_time))
|
||||
).order_by(Email.received_at.desc()).all()
|
||||
|
||||
# 返回结果
|
||||
result = {
|
||||
'mailbox_id': mailbox_id,
|
||||
'count': len(new_emails),
|
||||
'emails': [email.to_dict() for email in new_emails],
|
||||
'timestamp': time.time()
|
||||
}
|
||||
|
||||
return jsonify(result), 200
|
||||
finally:
|
||||
db.close()
|
||||
except Exception as e:
|
||||
current_app.logger.error(f"轮询新邮件出错: {str(e)}")
|
||||
return jsonify({'error': '轮询新邮件失败', 'details': str(e)}), 500
|
||||
206
app/api/mailbox_routes.py
Normal file
206
app/api/mailbox_routes.py
Normal file
@@ -0,0 +1,206 @@
|
||||
from flask import request, jsonify, current_app
|
||||
from sqlalchemy.exc import IntegrityError
|
||||
import random
|
||||
import string
|
||||
|
||||
from . import api_bp
|
||||
from ..models import get_session, Domain, Mailbox
|
||||
|
||||
# 获取所有邮箱
|
||||
@api_bp.route('/mailboxes', methods=['GET'])
|
||||
def get_mailboxes():
|
||||
"""获取所有邮箱列表"""
|
||||
try:
|
||||
page = int(request.args.get('page', 1))
|
||||
limit = int(request.args.get('limit', 50))
|
||||
offset = (page - 1) * limit
|
||||
|
||||
db = get_session()
|
||||
try:
|
||||
# 查询总数
|
||||
total = db.query(Mailbox).count()
|
||||
|
||||
# 获取分页数据
|
||||
mailboxes = db.query(Mailbox).order_by(Mailbox.created_at.desc()) \
|
||||
.limit(limit) \
|
||||
.offset(offset) \
|
||||
.all()
|
||||
|
||||
# 转换为字典列表
|
||||
result = {
|
||||
'total': total,
|
||||
'page': page,
|
||||
'limit': limit,
|
||||
'mailboxes': [mailbox.to_dict() for mailbox in mailboxes]
|
||||
}
|
||||
|
||||
return jsonify(result), 200
|
||||
finally:
|
||||
db.close()
|
||||
except Exception as e:
|
||||
current_app.logger.error(f"获取邮箱列表出错: {str(e)}")
|
||||
return jsonify({'error': '获取邮箱列表失败', 'details': str(e)}), 500
|
||||
|
||||
# 创建邮箱
|
||||
@api_bp.route('/mailboxes', methods=['POST'])
|
||||
def create_mailbox():
|
||||
"""创建新邮箱"""
|
||||
try:
|
||||
data = request.json
|
||||
|
||||
# 验证必要参数
|
||||
if not data or 'domain_id' not in data:
|
||||
return jsonify({'error': '缺少必要参数'}), 400
|
||||
|
||||
db = get_session()
|
||||
try:
|
||||
# 查询域名是否存在
|
||||
domain = db.query(Domain).filter_by(id=data['domain_id'], active=True).first()
|
||||
if not domain:
|
||||
return jsonify({'error': '指定的域名不存在或未激活'}), 404
|
||||
|
||||
# 生成或使用给定地址
|
||||
if 'address' not in data or not data['address']:
|
||||
# 生成随机地址
|
||||
address = ''.join(random.choices(string.ascii_lowercase + string.digits, k=10))
|
||||
else:
|
||||
address = data['address']
|
||||
|
||||
# 创建邮箱
|
||||
mailbox = Mailbox(
|
||||
address=address,
|
||||
domain_id=domain.id,
|
||||
description=data.get('description', ''),
|
||||
active=True
|
||||
)
|
||||
|
||||
db.add(mailbox)
|
||||
db.commit()
|
||||
|
||||
return jsonify({
|
||||
'message': '邮箱创建成功',
|
||||
'mailbox': mailbox.to_dict()
|
||||
}), 201
|
||||
except IntegrityError:
|
||||
db.rollback()
|
||||
return jsonify({'error': '邮箱地址已存在'}), 409
|
||||
except Exception as e:
|
||||
db.rollback()
|
||||
raise
|
||||
finally:
|
||||
db.close()
|
||||
except Exception as e:
|
||||
current_app.logger.error(f"创建邮箱出错: {str(e)}")
|
||||
return jsonify({'error': '创建邮箱失败', 'details': str(e)}), 500
|
||||
|
||||
# 批量创建邮箱
|
||||
@api_bp.route('/mailboxes/batch', methods=['POST'])
|
||||
def batch_create_mailboxes():
|
||||
"""批量创建邮箱"""
|
||||
try:
|
||||
data = request.json
|
||||
|
||||
# 验证必要参数
|
||||
if not data or 'domain_id' not in data or 'count' not in data:
|
||||
return jsonify({'error': '缺少必要参数'}), 400
|
||||
|
||||
domain_id = data['domain_id']
|
||||
count = min(int(data['count']), 100) # 限制最大数量为100
|
||||
prefix = data.get('prefix', '')
|
||||
|
||||
db = get_session()
|
||||
try:
|
||||
# 查询域名是否存在
|
||||
domain = db.query(Domain).filter_by(id=domain_id, active=True).first()
|
||||
if not domain:
|
||||
return jsonify({'error': '指定的域名不存在或未激活'}), 404
|
||||
|
||||
created_mailboxes = []
|
||||
|
||||
# 批量创建
|
||||
for _ in range(count):
|
||||
# 生成随机地址
|
||||
if prefix:
|
||||
address = f"{prefix}{random.randint(1000, 9999)}"
|
||||
else:
|
||||
address = ''.join(random.choices(string.ascii_lowercase + string.digits, k=10))
|
||||
|
||||
# 尝试创建,如果地址已存在则重试
|
||||
retries = 0
|
||||
while retries < 3: # 最多尝试3次
|
||||
try:
|
||||
mailbox = Mailbox(
|
||||
address=address,
|
||||
domain_id=domain.id,
|
||||
active=True
|
||||
)
|
||||
|
||||
db.add(mailbox)
|
||||
db.flush() # 验证但不提交
|
||||
created_mailboxes.append(mailbox)
|
||||
break
|
||||
except IntegrityError:
|
||||
db.rollback()
|
||||
# 地址已存在,重新生成
|
||||
address = ''.join(random.choices(string.ascii_lowercase + string.digits, k=10))
|
||||
retries += 1
|
||||
|
||||
# 提交所有更改
|
||||
db.commit()
|
||||
|
||||
return jsonify({
|
||||
'message': f'成功创建 {len(created_mailboxes)} 个邮箱',
|
||||
'mailboxes': [mailbox.to_dict() for mailbox in created_mailboxes]
|
||||
}), 201
|
||||
except Exception as e:
|
||||
db.rollback()
|
||||
raise
|
||||
finally:
|
||||
db.close()
|
||||
except Exception as e:
|
||||
current_app.logger.error(f"批量创建邮箱出错: {str(e)}")
|
||||
return jsonify({'error': '批量创建邮箱失败', 'details': str(e)}), 500
|
||||
|
||||
# 获取特定邮箱
|
||||
@api_bp.route('/mailboxes/<int:mailbox_id>', methods=['GET'])
|
||||
def get_mailbox(mailbox_id):
|
||||
"""获取指定ID的邮箱信息"""
|
||||
try:
|
||||
db = get_session()
|
||||
try:
|
||||
mailbox = db.query(Mailbox).filter_by(id=mailbox_id).first()
|
||||
|
||||
if not mailbox:
|
||||
return jsonify({'error': '邮箱不存在'}), 404
|
||||
|
||||
return jsonify(mailbox.to_dict()), 200
|
||||
finally:
|
||||
db.close()
|
||||
except Exception as e:
|
||||
current_app.logger.error(f"获取邮箱详情出错: {str(e)}")
|
||||
return jsonify({'error': '获取邮箱详情失败', 'details': str(e)}), 500
|
||||
|
||||
# 删除邮箱
|
||||
@api_bp.route('/mailboxes/<int:mailbox_id>', methods=['DELETE'])
|
||||
def delete_mailbox(mailbox_id):
|
||||
"""删除指定ID的邮箱"""
|
||||
try:
|
||||
db = get_session()
|
||||
try:
|
||||
mailbox = db.query(Mailbox).filter_by(id=mailbox_id).first()
|
||||
|
||||
if not mailbox:
|
||||
return jsonify({'error': '邮箱不存在'}), 404
|
||||
|
||||
db.delete(mailbox)
|
||||
db.commit()
|
||||
|
||||
return jsonify({'message': '邮箱已删除'}), 200
|
||||
except Exception as e:
|
||||
db.rollback()
|
||||
raise
|
||||
finally:
|
||||
db.close()
|
||||
except Exception as e:
|
||||
current_app.logger.error(f"删除邮箱出错: {str(e)}")
|
||||
return jsonify({'error': '删除邮箱失败', 'details': str(e)}), 500
|
||||
341
app/api/routes.py
Normal file
341
app/api/routes.py
Normal file
@@ -0,0 +1,341 @@
|
||||
from flask import Blueprint, request, jsonify, current_app
|
||||
import json
|
||||
from datetime import datetime, timedelta
|
||||
import os
|
||||
import time
|
||||
import psutil
|
||||
import sys
|
||||
import platform
|
||||
from sqlalchemy import func
|
||||
from ..models import get_session, Domain, Mailbox, Email
|
||||
from ..services import get_smtp_server, get_email_processor
|
||||
|
||||
api_bp = Blueprint('api', __name__, url_prefix='/api')
|
||||
|
||||
|
||||
@api_bp.route('/domains', methods=['GET'])
|
||||
def get_domains():
|
||||
"""获取所有可用域名"""
|
||||
db = get_session()
|
||||
try:
|
||||
domains = db.query(Domain).filter_by(active=True).all()
|
||||
return jsonify({
|
||||
'success': True,
|
||||
'domains': [domain.to_dict() for domain in domains]
|
||||
})
|
||||
except Exception as e:
|
||||
current_app.logger.exception(f"获取域名失败: {str(e)}")
|
||||
return jsonify({'success': False, 'error': '获取域名失败'}), 500
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
@api_bp.route('/domains', methods=['POST'])
|
||||
def create_domain():
|
||||
"""创建新域名"""
|
||||
data = request.json
|
||||
if not data or 'name' not in data:
|
||||
return jsonify({'success': False, 'error': '缺少必要字段'}), 400
|
||||
|
||||
db = get_session()
|
||||
try:
|
||||
# 检查域名是否已存在
|
||||
domain_exists = db.query(Domain).filter_by(name=data['name']).first()
|
||||
if domain_exists:
|
||||
return jsonify({'success': False, 'error': '域名已存在'}), 400
|
||||
|
||||
# 创建新域名
|
||||
domain = Domain(
|
||||
name=data['name'],
|
||||
description=data.get('description', ''),
|
||||
active=data.get('active', True)
|
||||
)
|
||||
db.add(domain)
|
||||
db.commit()
|
||||
|
||||
return jsonify({
|
||||
'success': True,
|
||||
'message': '域名创建成功',
|
||||
'domain': domain.to_dict()
|
||||
})
|
||||
except Exception as e:
|
||||
db.rollback()
|
||||
current_app.logger.exception(f"创建域名失败: {str(e)}")
|
||||
return jsonify({'success': False, 'error': '创建域名失败'}), 500
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
@api_bp.route('/mailboxes', methods=['GET'])
|
||||
def get_mailboxes():
|
||||
"""获取所有邮箱"""
|
||||
db = get_session()
|
||||
try:
|
||||
mailboxes = db.query(Mailbox).all()
|
||||
return jsonify({
|
||||
'success': True,
|
||||
'mailboxes': [mailbox.to_dict() for mailbox in mailboxes]
|
||||
})
|
||||
except Exception as e:
|
||||
current_app.logger.exception(f"获取邮箱失败: {str(e)}")
|
||||
return jsonify({'success': False, 'error': '获取邮箱失败'}), 500
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
@api_bp.route('/mailboxes', methods=['POST'])
|
||||
def create_mailbox():
|
||||
"""创建新邮箱"""
|
||||
data = request.json
|
||||
if not data or 'address' not in data or 'domain_id' not in data:
|
||||
return jsonify({'success': False, 'error': '缺少必要字段'}), 400
|
||||
|
||||
db = get_session()
|
||||
try:
|
||||
# 检查域名是否存在
|
||||
domain = db.query(Domain).filter_by(id=data['domain_id'], active=True).first()
|
||||
if not domain:
|
||||
return jsonify({'success': False, 'error': '域名不存在或未激活'}), 400
|
||||
|
||||
# 检查邮箱是否已存在
|
||||
mailbox_exists = db.query(Mailbox).filter_by(
|
||||
address=data['address'], domain_id=data['domain_id']).first()
|
||||
if mailbox_exists:
|
||||
return jsonify({'success': False, 'error': '邮箱已存在'}), 400
|
||||
|
||||
# 创建新邮箱
|
||||
mailbox = Mailbox(
|
||||
address=data['address'],
|
||||
domain_id=data['domain_id'],
|
||||
description=data.get('description', ''),
|
||||
active=data.get('active', True)
|
||||
)
|
||||
db.add(mailbox)
|
||||
db.commit()
|
||||
|
||||
return jsonify({
|
||||
'success': True,
|
||||
'message': '邮箱创建成功',
|
||||
'mailbox': mailbox.to_dict()
|
||||
})
|
||||
except Exception as e:
|
||||
db.rollback()
|
||||
current_app.logger.exception(f"创建邮箱失败: {str(e)}")
|
||||
return jsonify({'success': False, 'error': '创建邮箱失败'}), 500
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
@api_bp.route('/mailboxes/batch', methods=['POST'])
|
||||
def batch_create_mailboxes():
|
||||
"""批量创建邮箱"""
|
||||
data = request.json
|
||||
if not data or 'domain_id' not in data or 'usernames' not in data or not isinstance(data['usernames'], list):
|
||||
return jsonify({'success': False, 'error': '缺少必要字段或格式不正确'}), 400
|
||||
|
||||
domain_id = data['domain_id']
|
||||
usernames = data['usernames']
|
||||
description = data.get('description', '')
|
||||
|
||||
db = get_session()
|
||||
try:
|
||||
# 检查域名是否存在
|
||||
domain = db.query(Domain).filter_by(id=domain_id, active=True).first()
|
||||
if not domain:
|
||||
return jsonify({'success': False, 'error': '域名不存在或未激活'}), 400
|
||||
|
||||
created_mailboxes = []
|
||||
existed_mailboxes = []
|
||||
|
||||
for username in usernames:
|
||||
# 检查邮箱是否已存在
|
||||
mailbox_exists = db.query(Mailbox).filter_by(
|
||||
username=username, domain_id=domain_id).first()
|
||||
if mailbox_exists:
|
||||
existed_mailboxes.append(username)
|
||||
continue
|
||||
|
||||
# 创建新邮箱
|
||||
mailbox = Mailbox(
|
||||
username=username,
|
||||
domain_id=domain_id,
|
||||
description=description,
|
||||
active=True
|
||||
)
|
||||
db.add(mailbox)
|
||||
created_mailboxes.append(username)
|
||||
|
||||
db.commit()
|
||||
|
||||
return jsonify({
|
||||
'success': True,
|
||||
'message': f'成功创建 {len(created_mailboxes)} 个邮箱,{len(existed_mailboxes)} 个已存在',
|
||||
'created': created_mailboxes,
|
||||
'existed': existed_mailboxes
|
||||
})
|
||||
except Exception as e:
|
||||
db.rollback()
|
||||
current_app.logger.exception(f"批量创建邮箱失败: {str(e)}")
|
||||
return jsonify({'success': False, 'error': '批量创建邮箱失败'}), 500
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
@api_bp.route('/mailboxes/<int:mailbox_id>', methods=['GET'])
|
||||
def get_mailbox(mailbox_id):
|
||||
"""获取特定邮箱的信息"""
|
||||
db = get_session()
|
||||
try:
|
||||
mailbox = db.query(Mailbox).filter_by(id=mailbox_id).first()
|
||||
if not mailbox:
|
||||
return jsonify({'success': False, 'error': '邮箱不存在'}), 404
|
||||
|
||||
# 更新最后访问时间
|
||||
mailbox.last_accessed = datetime.utcnow()
|
||||
db.commit()
|
||||
|
||||
return jsonify({
|
||||
'success': True,
|
||||
'mailbox': mailbox.to_dict()
|
||||
})
|
||||
except Exception as e:
|
||||
current_app.logger.exception(f"获取邮箱信息失败: {str(e)}")
|
||||
return jsonify({'success': False, 'error': '获取邮箱信息失败'}), 500
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
@api_bp.route('/mailboxes/<int:mailbox_id>/emails', methods=['GET'])
|
||||
def get_emails(mailbox_id):
|
||||
"""获取邮箱中的所有邮件"""
|
||||
db = get_session()
|
||||
try:
|
||||
mailbox = db.query(Mailbox).filter_by(id=mailbox_id).first()
|
||||
if not mailbox:
|
||||
return jsonify({'success': False, 'error': '邮箱不存在'}), 404
|
||||
|
||||
# 更新最后访问时间
|
||||
mailbox.last_accessed = datetime.utcnow()
|
||||
db.commit()
|
||||
|
||||
emails = db.query(Email).filter_by(mailbox_id=mailbox_id).order_by(Email.received_at.desc()).all()
|
||||
|
||||
return jsonify({
|
||||
'success': True,
|
||||
'emails': [email.to_dict() for email in emails]
|
||||
})
|
||||
except Exception as e:
|
||||
current_app.logger.exception(f"获取邮件失败: {str(e)}")
|
||||
return jsonify({'success': False, 'error': '获取邮件失败'}), 500
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
@api_bp.route('/emails/<int:email_id>', methods=['GET'])
|
||||
def get_email(email_id):
|
||||
"""获取特定邮件的详细内容"""
|
||||
db = get_session()
|
||||
try:
|
||||
email = db.query(Email).filter_by(id=email_id).first()
|
||||
if not email:
|
||||
return jsonify({'success': False, 'error': '邮件不存在'}), 404
|
||||
|
||||
# 标记为已读
|
||||
if not email.read:
|
||||
email.read = True
|
||||
db.commit()
|
||||
|
||||
return jsonify({
|
||||
'success': True,
|
||||
'email': email.to_dict()
|
||||
})
|
||||
except Exception as e:
|
||||
current_app.logger.exception(f"获取邮件详情失败: {str(e)}")
|
||||
return jsonify({'success': False, 'error': '获取邮件详情失败'}), 500
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
@api_bp.route('/emails/<int:email_id>/verification', methods=['GET'])
|
||||
def get_verification_info(email_id):
|
||||
"""获取邮件中的验证信息(链接和验证码)"""
|
||||
db = get_session()
|
||||
try:
|
||||
email = db.query(Email).filter_by(id=email_id).first()
|
||||
if not email:
|
||||
return jsonify({'success': False, 'error': '邮件不存在'}), 404
|
||||
|
||||
verification_links = json.loads(email.verification_links) if email.verification_links else []
|
||||
verification_codes = json.loads(email.verification_codes) if email.verification_codes else []
|
||||
|
||||
return jsonify({
|
||||
'success': True,
|
||||
'email_id': email_id,
|
||||
'verification_links': verification_links,
|
||||
'verification_codes': verification_codes
|
||||
})
|
||||
except Exception as e:
|
||||
current_app.logger.exception(f"获取验证信息失败: {str(e)}")
|
||||
return jsonify({'success': False, 'error': '获取验证信息失败'}), 500
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
@api_bp.route('/status', methods=['GET'])
|
||||
def system_status():
|
||||
"""获取系统状态"""
|
||||
session = get_session()
|
||||
|
||||
# 获取基本统计信息
|
||||
domain_count = session.query(func.count(Domain.id)).scalar()
|
||||
mailbox_count = session.query(func.count(Mailbox.id)).scalar()
|
||||
email_count = session.query(func.count(Email.id)).scalar()
|
||||
|
||||
# 获取最近24小时的邮件数量
|
||||
recent_emails = session.query(func.count(Email.id)).filter(
|
||||
Email.received_at > datetime.now() - timedelta(hours=24)
|
||||
).scalar()
|
||||
|
||||
# 获取系统资源信息
|
||||
cpu_percent = psutil.cpu_percent(interval=0.5)
|
||||
memory = psutil.virtual_memory()
|
||||
disk = psutil.disk_usage('/')
|
||||
|
||||
# 获取服务状态
|
||||
smtp_server = get_smtp_server()
|
||||
email_processor = get_email_processor()
|
||||
|
||||
smtp_status = "running" if smtp_server and smtp_server.controller else "stopped"
|
||||
processor_status = "running" if email_processor and email_processor.is_running else "stopped"
|
||||
|
||||
# 构建响应
|
||||
status = {
|
||||
"system": {
|
||||
"uptime": round(time.time() - psutil.boot_time()),
|
||||
"time": datetime.now().isoformat(),
|
||||
"platform": platform.platform(),
|
||||
"python_version": sys.version
|
||||
},
|
||||
"resources": {
|
||||
"cpu_percent": cpu_percent,
|
||||
"memory_percent": memory.percent,
|
||||
"memory_used": memory.used,
|
||||
"memory_total": memory.total,
|
||||
"disk_percent": disk.percent,
|
||||
"disk_used": disk.used,
|
||||
"disk_total": disk.total
|
||||
},
|
||||
"application": {
|
||||
"domain_count": domain_count,
|
||||
"mailbox_count": mailbox_count,
|
||||
"email_count": email_count,
|
||||
"recent_emails_24h": recent_emails,
|
||||
"storage_path": os.path.abspath("email_data"),
|
||||
"services": {
|
||||
"smtp_server": smtp_status,
|
||||
"email_processor": processor_status
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return jsonify(status)
|
||||
46
app/models/__init__.py
Normal file
46
app/models/__init__.py
Normal file
@@ -0,0 +1,46 @@
|
||||
from sqlalchemy import create_engine
|
||||
from sqlalchemy.ext.declarative import declarative_base
|
||||
from sqlalchemy.orm import sessionmaker, scoped_session
|
||||
import os
|
||||
import sys
|
||||
|
||||
# 修改相对导入为绝对导入
|
||||
sys.path.append(os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))))
|
||||
from config import active_config
|
||||
|
||||
# 创建数据库引擎
|
||||
engine = create_engine(active_config.SQLALCHEMY_DATABASE_URI)
|
||||
|
||||
# 创建会话工厂
|
||||
session_factory = sessionmaker(bind=engine)
|
||||
Session = scoped_session(session_factory)
|
||||
|
||||
# 创建模型基类
|
||||
Base = declarative_base()
|
||||
|
||||
# 获取数据库会话
|
||||
def get_session():
|
||||
"""获取数据库会话"""
|
||||
return Session()
|
||||
|
||||
# 初始化数据库
|
||||
def init_db():
|
||||
"""初始化数据库,创建所有表"""
|
||||
# 导入所有模型以确保它们被注册
|
||||
from .domain import Domain
|
||||
from .mailbox import Mailbox
|
||||
from .email import Email
|
||||
from .attachment import Attachment
|
||||
|
||||
# 创建表
|
||||
Base.metadata.create_all(engine)
|
||||
|
||||
return get_session()
|
||||
|
||||
# 导出模型类
|
||||
from .domain import Domain
|
||||
from .mailbox import Mailbox
|
||||
from .email import Email
|
||||
from .attachment import Attachment
|
||||
|
||||
__all__ = ['Base', 'get_session', 'init_db', 'Domain', 'Mailbox', 'Email', 'Attachment']
|
||||
71
app/models/attachment.py
Normal file
71
app/models/attachment.py
Normal file
@@ -0,0 +1,71 @@
|
||||
from sqlalchemy import Column, Integer, String, DateTime, ForeignKey, LargeBinary
|
||||
from sqlalchemy.orm import relationship
|
||||
from datetime import datetime
|
||||
import os
|
||||
|
||||
from . import Base
|
||||
|
||||
class Attachment(Base):
|
||||
"""附件模型"""
|
||||
__tablename__ = 'attachments'
|
||||
|
||||
id = Column(Integer, primary_key=True)
|
||||
email_id = Column(Integer, ForeignKey('emails.id'), nullable=False, index=True)
|
||||
filename = Column(String(255), nullable=False)
|
||||
content_type = Column(String(100), nullable=True)
|
||||
size = Column(Integer, nullable=False, default=0)
|
||||
storage_path = Column(String(500), nullable=True) # 用于文件系统存储
|
||||
content = Column(LargeBinary, nullable=True) # 用于小型附件的直接存储
|
||||
created_at = Column(DateTime, default=datetime.utcnow)
|
||||
|
||||
# 关联关系
|
||||
email = relationship("Email", back_populates="attachments")
|
||||
|
||||
@property
|
||||
def is_stored_in_fs(self):
|
||||
"""判断附件是否存储在文件系统中"""
|
||||
return bool(self.storage_path and not self.content)
|
||||
|
||||
def save_to_filesystem(self, content, base_path):
|
||||
"""将附件保存到文件系统"""
|
||||
# 确保目录存在
|
||||
os.makedirs(base_path, exist_ok=True)
|
||||
|
||||
# 创建文件路径
|
||||
file_path = os.path.join(
|
||||
base_path,
|
||||
f"{self.email_id}_{self.id}_{self.filename}"
|
||||
)
|
||||
|
||||
# 写入文件
|
||||
with open(file_path, 'wb') as f:
|
||||
f.write(content)
|
||||
|
||||
# 更新对象属性
|
||||
self.storage_path = file_path
|
||||
self.size = len(content)
|
||||
self.content = None # 清空内存中的内容
|
||||
|
||||
return file_path
|
||||
|
||||
def get_content(self, attachments_dir=None):
|
||||
"""获取附件内容,无论是从数据库还是文件系统"""
|
||||
if self.content:
|
||||
return self.content
|
||||
|
||||
if self.storage_path and os.path.exists(self.storage_path):
|
||||
with open(self.storage_path, 'rb') as f:
|
||||
return f.read()
|
||||
|
||||
return None
|
||||
|
||||
def to_dict(self):
|
||||
"""转换为字典,用于API响应"""
|
||||
return {
|
||||
"id": self.id,
|
||||
"email_id": self.email_id,
|
||||
"filename": self.filename,
|
||||
"content_type": self.content_type,
|
||||
"size": self.size,
|
||||
"created_at": self.created_at.isoformat() if self.created_at else None
|
||||
}
|
||||
35
app/models/domain.py
Normal file
35
app/models/domain.py
Normal file
@@ -0,0 +1,35 @@
|
||||
from sqlalchemy import Column, Integer, String, Boolean, DateTime
|
||||
from sqlalchemy.orm import relationship
|
||||
from datetime import datetime
|
||||
|
||||
from . import Base
|
||||
|
||||
|
||||
class Domain(Base):
|
||||
"""邮件域名模型"""
|
||||
__tablename__ = 'domains'
|
||||
|
||||
id = Column(Integer, primary_key=True)
|
||||
name = Column(String(255), unique=True, nullable=False, index=True)
|
||||
description = Column(String(500), nullable=True)
|
||||
active = Column(Boolean, default=True)
|
||||
created_at = Column(DateTime, default=datetime.utcnow)
|
||||
updated_at = Column(DateTime, default=datetime.utcnow, onupdate=datetime.utcnow)
|
||||
|
||||
# 关系
|
||||
mailboxes = relationship("Mailbox", back_populates="domain", cascade="all, delete-orphan")
|
||||
|
||||
def __repr__(self):
|
||||
return f"<Domain {self.name}>"
|
||||
|
||||
def to_dict(self):
|
||||
"""转换为字典,用于API响应"""
|
||||
return {
|
||||
"id": self.id,
|
||||
"name": self.name,
|
||||
"description": self.description,
|
||||
"active": self.active,
|
||||
"created_at": self.created_at.isoformat() if self.created_at else None,
|
||||
"updated_at": self.updated_at.isoformat() if self.updated_at else None,
|
||||
"mailbox_count": len(self.mailboxes) if self.mailboxes else 0
|
||||
}
|
||||
98
app/models/email.py
Normal file
98
app/models/email.py
Normal file
@@ -0,0 +1,98 @@
|
||||
import os
|
||||
import json
|
||||
from sqlalchemy import Column, Integer, String, Text, DateTime, ForeignKey, Boolean, JSON
|
||||
from sqlalchemy.orm import relationship
|
||||
from datetime import datetime
|
||||
import re
|
||||
import sys
|
||||
|
||||
from . import Base
|
||||
import config
|
||||
active_config = config.active_config
|
||||
|
||||
|
||||
class Email(Base):
|
||||
"""电子邮件模型"""
|
||||
__tablename__ = 'emails'
|
||||
|
||||
id = Column(Integer, primary_key=True)
|
||||
mailbox_id = Column(Integer, ForeignKey('mailboxes.id'), nullable=False, index=True)
|
||||
sender = Column(String(255), nullable=False)
|
||||
recipients = Column(String(1000), nullable=False)
|
||||
subject = Column(String(500), nullable=True)
|
||||
body_text = Column(Text, nullable=True)
|
||||
body_html = Column(Text, nullable=True)
|
||||
received_at = Column(DateTime, default=datetime.utcnow)
|
||||
read = Column(Boolean, default=False)
|
||||
headers = Column(JSON, nullable=True)
|
||||
|
||||
# 提取的验证码和链接
|
||||
verification_code = Column(String(100), nullable=True)
|
||||
verification_link = Column(String(1000), nullable=True)
|
||||
|
||||
# 关联关系
|
||||
mailbox = relationship("Mailbox", back_populates="emails")
|
||||
attachments = relationship("Attachment", back_populates="email", cascade="all, delete-orphan")
|
||||
|
||||
def save_raw_email(self, raw_content):
|
||||
"""保存原始邮件内容到文件"""
|
||||
storage_path = active_config.MAIL_STORAGE_PATH
|
||||
mailbox_dir = os.path.join(storage_path, str(self.mailbox_id))
|
||||
os.makedirs(mailbox_dir, exist_ok=True)
|
||||
|
||||
# 保存原始邮件内容
|
||||
file_path = os.path.join(mailbox_dir, f"{self.id}.eml")
|
||||
with open(file_path, 'wb') as f:
|
||||
f.write(raw_content)
|
||||
|
||||
def extract_verification_data(self):
|
||||
"""
|
||||
尝试从邮件内容中提取验证码和验证链接
|
||||
这个方法会在邮件保存时自动调用
|
||||
"""
|
||||
# 合并文本和HTML内容用于搜索
|
||||
content = f"{self.subject} {self.body_text or ''}"
|
||||
|
||||
# 提取可能的验证码(4-8位数字或字母组合)
|
||||
code_patterns = [
|
||||
r'\b[A-Z0-9]{4,8}\b', # 大写字母和数字
|
||||
r'验证码[::]\s*([A-Z0-9]{4,8})', # 中文格式
|
||||
r'验证码是[::]\s*([A-Z0-9]{4,8})', # 中文格式2
|
||||
r'code[::]\s*([A-Z0-9]{4,8})', # 英文格式
|
||||
]
|
||||
|
||||
for pattern in code_patterns:
|
||||
matches = re.findall(pattern, content, re.IGNORECASE)
|
||||
if matches:
|
||||
self.verification_code = matches[0]
|
||||
break
|
||||
|
||||
# 提取验证链接
|
||||
link_patterns = [
|
||||
r'https?://\S+(?:verify|confirm|activate)\S+',
|
||||
r'https?://\S+(?:token|auth|account)\S+',
|
||||
]
|
||||
|
||||
for pattern in link_patterns:
|
||||
matches = re.findall(pattern, content, re.IGNORECASE)
|
||||
if matches:
|
||||
self.verification_link = matches[0]
|
||||
break
|
||||
|
||||
def __repr__(self):
|
||||
return f"<Email {self.id}: {self.subject}>"
|
||||
|
||||
def to_dict(self):
|
||||
"""转换为字典,用于API响应"""
|
||||
return {
|
||||
"id": self.id,
|
||||
"mailbox_id": self.mailbox_id,
|
||||
"sender": self.sender,
|
||||
"recipients": self.recipients,
|
||||
"subject": self.subject,
|
||||
"received_at": self.received_at.isoformat() if self.received_at else None,
|
||||
"read": self.read,
|
||||
"verification_code": self.verification_code,
|
||||
"verification_link": self.verification_link,
|
||||
"has_attachments": len(self.attachments) > 0 if self.attachments else False
|
||||
}
|
||||
50
app/models/mailbox.py
Normal file
50
app/models/mailbox.py
Normal file
@@ -0,0 +1,50 @@
|
||||
from sqlalchemy import Column, Integer, String, Boolean, DateTime, ForeignKey
|
||||
from sqlalchemy.orm import relationship
|
||||
from datetime import datetime
|
||||
import secrets
|
||||
|
||||
from . import Base
|
||||
|
||||
|
||||
class Mailbox(Base):
|
||||
"""邮箱模型"""
|
||||
__tablename__ = 'mailboxes'
|
||||
|
||||
id = Column(Integer, primary_key=True)
|
||||
address = Column(String(255), unique=True, nullable=False, index=True)
|
||||
domain_id = Column(Integer, ForeignKey('domains.id'), nullable=False)
|
||||
password_hash = Column(String(255), nullable=True)
|
||||
description = Column(String(500), nullable=True)
|
||||
active = Column(Boolean, default=True)
|
||||
api_key = Column(String(64), unique=True, default=lambda: secrets.token_hex(16))
|
||||
created_at = Column(DateTime, default=datetime.utcnow)
|
||||
updated_at = Column(DateTime, default=datetime.utcnow, onupdate=datetime.utcnow)
|
||||
last_accessed = Column(DateTime, nullable=True)
|
||||
|
||||
# 关系
|
||||
domain = relationship("Domain", back_populates="mailboxes")
|
||||
emails = relationship("Email", back_populates="mailbox", cascade="all, delete-orphan")
|
||||
|
||||
@property
|
||||
def full_address(self):
|
||||
"""获取完整邮箱地址 (包含域名)"""
|
||||
return f"{self.address}@{self.domain.name}"
|
||||
|
||||
def __repr__(self):
|
||||
return f"<Mailbox {self.full_address}>"
|
||||
|
||||
def to_dict(self):
|
||||
"""转换为字典,用于API响应"""
|
||||
return {
|
||||
"id": self.id,
|
||||
"address": self.address,
|
||||
"domain_id": self.domain_id,
|
||||
"domain_name": self.domain.name if self.domain else None,
|
||||
"full_address": self.full_address,
|
||||
"description": self.description,
|
||||
"active": self.active,
|
||||
"created_at": self.created_at.isoformat() if self.created_at else None,
|
||||
"updated_at": self.updated_at.isoformat() if self.updated_at else None,
|
||||
"last_accessed": self.last_accessed.isoformat() if self.last_accessed else None,
|
||||
"email_count": len(self.emails) if self.emails else 0
|
||||
}
|
||||
44
app/services/__init__.py
Normal file
44
app/services/__init__.py
Normal file
@@ -0,0 +1,44 @@
|
||||
# 服务层初始化文件
|
||||
# 这里将导入所有服务模块以便于统一调用
|
||||
|
||||
from .smtp_server import SMTPServer
|
||||
from .email_processor import EmailProcessor
|
||||
from .mail_store import MailStore
|
||||
|
||||
# 全局服务实例
|
||||
_smtp_server = None
|
||||
_email_processor = None
|
||||
_mail_store = None
|
||||
|
||||
def register_smtp_server(instance):
|
||||
"""注册SMTP服务器实例"""
|
||||
global _smtp_server
|
||||
_smtp_server = instance
|
||||
|
||||
def register_email_processor(instance):
|
||||
"""注册邮件处理器实例"""
|
||||
global _email_processor
|
||||
_email_processor = instance
|
||||
|
||||
def register_mail_store(instance):
|
||||
"""注册邮件存储实例"""
|
||||
global _mail_store
|
||||
_mail_store = instance
|
||||
|
||||
def get_smtp_server():
|
||||
"""获取SMTP服务器实例"""
|
||||
return _smtp_server
|
||||
|
||||
def get_email_processor():
|
||||
"""获取邮件处理器实例"""
|
||||
return _email_processor
|
||||
|
||||
def get_mail_store():
|
||||
"""获取邮件存储实例"""
|
||||
return _mail_store
|
||||
|
||||
__all__ = [
|
||||
'SMTPServer', 'EmailProcessor', 'MailStore',
|
||||
'register_smtp_server', 'register_email_processor', 'register_mail_store',
|
||||
'get_smtp_server', 'get_email_processor', 'get_mail_store'
|
||||
]
|
||||
123
app/services/email_processor.py
Normal file
123
app/services/email_processor.py
Normal file
@@ -0,0 +1,123 @@
|
||||
import logging
|
||||
import re
|
||||
import threading
|
||||
import time
|
||||
from queue import Queue
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
class EmailProcessor:
|
||||
"""邮件处理器,负责处理邮件并提取验证信息"""
|
||||
|
||||
def __init__(self, mail_store):
|
||||
"""
|
||||
初始化邮件处理器
|
||||
|
||||
参数:
|
||||
mail_store: 邮件存储服务实例
|
||||
"""
|
||||
self.mail_store = mail_store
|
||||
self.processing_queue = Queue()
|
||||
self.is_running = False
|
||||
self.worker_thread = None
|
||||
|
||||
def start(self):
|
||||
"""启动邮件处理器"""
|
||||
if self.is_running:
|
||||
logger.warning("邮件处理器已在运行")
|
||||
return False
|
||||
|
||||
self.is_running = True
|
||||
self.worker_thread = threading.Thread(
|
||||
target=self._processing_worker,
|
||||
daemon=True
|
||||
)
|
||||
self.worker_thread.start()
|
||||
logger.info("邮件处理器已启动")
|
||||
return True
|
||||
|
||||
def stop(self):
|
||||
"""停止邮件处理器"""
|
||||
if not self.is_running:
|
||||
logger.warning("邮件处理器未在运行")
|
||||
return False
|
||||
|
||||
self.is_running = False
|
||||
if self.worker_thread:
|
||||
self.worker_thread.join(timeout=5.0)
|
||||
self.worker_thread = None
|
||||
|
||||
logger.info("邮件处理器已停止")
|
||||
return True
|
||||
|
||||
def queue_email_for_processing(self, email_id):
|
||||
"""将邮件添加到处理队列"""
|
||||
self.processing_queue.put(email_id)
|
||||
return True
|
||||
|
||||
def _processing_worker(self):
|
||||
"""处理队列中的邮件的工作线程"""
|
||||
while self.is_running:
|
||||
try:
|
||||
# 获取队列中的邮件,最多等待1秒
|
||||
try:
|
||||
email_id = self.processing_queue.get(timeout=1.0)
|
||||
except:
|
||||
continue
|
||||
|
||||
# 处理邮件
|
||||
self._process_email(email_id)
|
||||
|
||||
# 标记任务完成
|
||||
self.processing_queue.task_done()
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"处理邮件时出错: {str(e)}")
|
||||
|
||||
def _process_email(self, email_id):
|
||||
"""处理单个邮件,提取验证码和链接"""
|
||||
# 从邮件存储获取邮件
|
||||
email_data = self.mail_store.get_email_by_id(email_id, mark_as_read=False)
|
||||
if not email_data:
|
||||
logger.warning(f"找不到ID为 {email_id} 的邮件")
|
||||
return False
|
||||
|
||||
# 提取验证码和链接已经在Email模型的extract_verification_data方法中实现
|
||||
# 这里可以添加更复杂的提取逻辑或后处理
|
||||
|
||||
logger.info(f"邮件 {email_id} 处理完成")
|
||||
return True
|
||||
|
||||
@staticmethod
|
||||
def extract_verification_code(content):
|
||||
"""从内容中提取验证码"""
|
||||
code_patterns = [
|
||||
r'\b[A-Z0-9]{4,8}\b', # 基本验证码格式
|
||||
r'验证码[::]\s*([A-Z0-9]{4,8})',
|
||||
r'验证码是[::]\s*([A-Z0-9]{4,8})',
|
||||
r'code[::]\s*([A-Z0-9]{4,8})',
|
||||
r'码[::]\s*(\d{4,8})' # 纯数字验证码
|
||||
]
|
||||
|
||||
for pattern in code_patterns:
|
||||
matches = re.findall(pattern, content, re.IGNORECASE)
|
||||
if matches:
|
||||
return matches[0]
|
||||
|
||||
return None
|
||||
|
||||
@staticmethod
|
||||
def extract_verification_link(content):
|
||||
"""从内容中提取验证链接"""
|
||||
link_patterns = [
|
||||
r'(https?://\S+(?:verify|confirm|activate)\S+)',
|
||||
r'(https?://\S+(?:token|auth|account)\S+)',
|
||||
r'href\s*=\s*["\']([^"\']+(?:verify|confirm|activate)[^"\']*)["\']'
|
||||
]
|
||||
|
||||
for pattern in link_patterns:
|
||||
matches = re.findall(pattern, content, re.IGNORECASE)
|
||||
if matches:
|
||||
return matches[0]
|
||||
|
||||
return None
|
||||
263
app/services/mail_store.py
Normal file
263
app/services/mail_store.py
Normal file
@@ -0,0 +1,263 @@
|
||||
import logging
|
||||
import os
|
||||
import email
|
||||
from email.policy import default
|
||||
from sqlalchemy.orm import Session
|
||||
from datetime import datetime
|
||||
import re
|
||||
|
||||
from ..models.domain import Domain
|
||||
from ..models.mailbox import Mailbox
|
||||
from ..models.email import Email
|
||||
from ..models.attachment import Attachment
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
class MailStore:
|
||||
"""邮件存储服务,负责保存和检索邮件"""
|
||||
|
||||
def __init__(self, db_session_factory, storage_path=None):
|
||||
"""
|
||||
初始化邮件存储服务
|
||||
|
||||
参数:
|
||||
db_session_factory: 数据库会话工厂函数
|
||||
storage_path: 附件存储路径
|
||||
"""
|
||||
self.db_session_factory = db_session_factory
|
||||
self.storage_path = storage_path or os.path.join(os.getcwd(), 'email_data')
|
||||
|
||||
# 确保存储目录存在
|
||||
if not os.path.exists(self.storage_path):
|
||||
os.makedirs(self.storage_path)
|
||||
|
||||
async def save_email(self, sender, recipient, message, raw_data):
|
||||
"""
|
||||
保存一封电子邮件
|
||||
|
||||
参数:
|
||||
sender: 发件人地址
|
||||
recipient: 收件人地址
|
||||
message: 解析后的邮件对象
|
||||
raw_data: 原始邮件数据
|
||||
|
||||
返回:
|
||||
成功返回邮件ID,失败返回None
|
||||
"""
|
||||
# 从收件人地址中提取用户名和域名
|
||||
try:
|
||||
address, domain_name = recipient.split('@', 1)
|
||||
except ValueError:
|
||||
logger.warning(f"无效的收件人地址格式: {recipient}")
|
||||
return None
|
||||
|
||||
# 获取数据库会话
|
||||
db = self.db_session_factory()
|
||||
|
||||
try:
|
||||
# 检查域名是否存在且活跃
|
||||
domain = db.query(Domain).filter_by(name=domain_name, active=True).first()
|
||||
if not domain:
|
||||
logger.warning(f"不支持的域名: {domain_name}")
|
||||
return None
|
||||
|
||||
# 查找或创建邮箱
|
||||
mailbox = db.query(Mailbox).filter_by(address=address, domain_id=domain.id).first()
|
||||
if not mailbox:
|
||||
# 自动创建新邮箱
|
||||
mailbox = Mailbox(
|
||||
address=address,
|
||||
domain_id=domain.id,
|
||||
active=True
|
||||
)
|
||||
db.add(mailbox)
|
||||
db.flush() # 获取ID但不提交
|
||||
logger.info(f"已为 {recipient} 自动创建邮箱")
|
||||
|
||||
# 提取邮件内容
|
||||
subject = message.get('subject', '')
|
||||
|
||||
# 获取文本和HTML内容
|
||||
body_text = None
|
||||
body_html = None
|
||||
attachments_data = []
|
||||
|
||||
if message.is_multipart():
|
||||
for part in message.walk():
|
||||
content_type = part.get_content_type()
|
||||
content_disposition = part.get_content_disposition()
|
||||
|
||||
# 处理文本内容
|
||||
if content_disposition is None or content_disposition == 'inline':
|
||||
if content_type == 'text/plain' and not body_text:
|
||||
body_text = part.get_content()
|
||||
elif content_type == 'text/html' and not body_html:
|
||||
body_html = part.get_content()
|
||||
|
||||
# 处理附件
|
||||
elif content_disposition == 'attachment':
|
||||
filename = part.get_filename()
|
||||
if filename:
|
||||
content = part.get_payload(decode=True)
|
||||
if content:
|
||||
attachments_data.append({
|
||||
'filename': filename,
|
||||
'content_type': content_type,
|
||||
'data': content,
|
||||
'size': len(content)
|
||||
})
|
||||
else:
|
||||
# 非多部分邮件
|
||||
content_type = message.get_content_type()
|
||||
if content_type == 'text/plain':
|
||||
body_text = message.get_content()
|
||||
elif content_type == 'text/html':
|
||||
body_html = message.get_content()
|
||||
|
||||
# 创建邮件记录
|
||||
email_obj = Email(
|
||||
mailbox_id=mailbox.id,
|
||||
sender=sender,
|
||||
recipients=recipient,
|
||||
subject=subject,
|
||||
body_text=body_text,
|
||||
body_html=body_html,
|
||||
headers={k: v for k, v in message.items()}
|
||||
)
|
||||
|
||||
# 保存邮件
|
||||
db.add(email_obj)
|
||||
db.flush() # 获取ID但不提交
|
||||
|
||||
# 提取验证信息
|
||||
email_obj.extract_verification_data()
|
||||
|
||||
# 保存附件
|
||||
for attachment_data in attachments_data:
|
||||
attachment = Attachment(
|
||||
email_id=email_obj.id,
|
||||
filename=attachment_data['filename'],
|
||||
content_type=attachment_data['content_type'],
|
||||
size=attachment_data['size']
|
||||
)
|
||||
|
||||
db.add(attachment)
|
||||
db.flush()
|
||||
|
||||
# 决定存储位置
|
||||
if attachment_data['size'] > 1024 * 1024: # 大于1MB的存储到文件系统
|
||||
attachments_dir = os.path.join(self.storage_path, 'attachments')
|
||||
attachment.save_to_filesystem(attachment_data['data'], attachments_dir)
|
||||
else:
|
||||
# 小附件直接存储在数据库
|
||||
attachment.content = attachment_data['data']
|
||||
|
||||
# 提交所有更改
|
||||
db.commit()
|
||||
logger.info(f"邮件已成功保存: {sender} -> {recipient}, ID: {email_obj.id}")
|
||||
return email_obj.id
|
||||
|
||||
except Exception as e:
|
||||
db.rollback()
|
||||
logger.error(f"保存邮件时出错: {str(e)}")
|
||||
return None
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
def get_emails_for_mailbox(self, mailbox_id, limit=50, offset=0, unread_only=False):
|
||||
"""获取指定邮箱的邮件列表"""
|
||||
db = self.db_session_factory()
|
||||
try:
|
||||
query = db.query(Email).filter(Email.mailbox_id == mailbox_id)
|
||||
|
||||
if unread_only:
|
||||
query = query.filter(Email.read == False)
|
||||
|
||||
total = query.count()
|
||||
|
||||
emails = query.order_by(Email.received_at.desc()) \
|
||||
.limit(limit) \
|
||||
.offset(offset) \
|
||||
.all()
|
||||
|
||||
return {
|
||||
'total': total,
|
||||
'items': [email.to_dict() for email in emails]
|
||||
}
|
||||
except Exception as e:
|
||||
logger.error(f"获取邮件列表时出错: {str(e)}")
|
||||
return {'total': 0, 'items': []}
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
def get_email_by_id(self, email_id, mark_as_read=True):
|
||||
"""获取指定ID的邮件详情"""
|
||||
db = self.db_session_factory()
|
||||
try:
|
||||
email = db.query(Email).filter(Email.id == email_id).first()
|
||||
|
||||
if not email:
|
||||
return None
|
||||
|
||||
if mark_as_read and not email.read:
|
||||
email.read = True
|
||||
email.last_read = datetime.utcnow()
|
||||
db.commit()
|
||||
|
||||
# 获取附件信息
|
||||
attachments = [attachment.to_dict() for attachment in email.attachments]
|
||||
|
||||
# 构建完整响应
|
||||
result = email.to_dict()
|
||||
result['body_text'] = email.body_text
|
||||
result['body_html'] = email.body_html
|
||||
result['attachments'] = attachments
|
||||
|
||||
return result
|
||||
except Exception as e:
|
||||
db.rollback()
|
||||
logger.error(f"获取邮件详情时出错: {str(e)}")
|
||||
return None
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
def delete_email(self, email_id):
|
||||
"""删除指定ID的邮件"""
|
||||
db = self.db_session_factory()
|
||||
try:
|
||||
email = db.query(Email).filter(Email.id == email_id).first()
|
||||
|
||||
if not email:
|
||||
return False
|
||||
|
||||
db.delete(email)
|
||||
db.commit()
|
||||
return True
|
||||
except Exception as e:
|
||||
db.rollback()
|
||||
logger.error(f"删除邮件时出错: {str(e)}")
|
||||
return False
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
def get_attachment_content(self, attachment_id):
|
||||
"""获取附件内容"""
|
||||
db = self.db_session_factory()
|
||||
try:
|
||||
attachment = db.query(Attachment).filter(Attachment.id == attachment_id).first()
|
||||
|
||||
if not attachment:
|
||||
return None
|
||||
|
||||
content = attachment.get_content()
|
||||
|
||||
return {
|
||||
'content': content,
|
||||
'filename': attachment.filename,
|
||||
'content_type': attachment.content_type
|
||||
}
|
||||
except Exception as e:
|
||||
logger.error(f"获取附件内容时出错: {str(e)}")
|
||||
return None
|
||||
finally:
|
||||
db.close()
|
||||
136
app/services/smtp_server.py
Normal file
136
app/services/smtp_server.py
Normal file
@@ -0,0 +1,136 @@
|
||||
import asyncio
|
||||
import logging
|
||||
import email
|
||||
import platform
|
||||
from email.policy import default
|
||||
from aiosmtpd.controller import Controller
|
||||
from aiosmtpd.smtp import SMTP as SMTPProtocol
|
||||
from aiosmtpd.handlers import Message
|
||||
import os
|
||||
import sys
|
||||
import threading
|
||||
|
||||
from ..models.domain import Domain
|
||||
from ..models.mailbox import Mailbox
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# 检测是否Windows环境
|
||||
IS_WINDOWS = platform.system().lower() == 'windows'
|
||||
|
||||
class EmailHandler(Message):
|
||||
"""处理接收的电子邮件"""
|
||||
|
||||
def __init__(self, mail_store):
|
||||
super().__init__()
|
||||
self.mail_store = mail_store
|
||||
|
||||
def handle_message(self, message):
|
||||
"""处理邮件消息,这是Message类的抽象方法,必须实现"""
|
||||
# 这个方法在异步DATA处理完成后被调用,但我们的邮件处理逻辑已经在handle_DATA中实现
|
||||
# 所以这里只是一个空实现
|
||||
return
|
||||
|
||||
async def handle_DATA(self, server, session, envelope):
|
||||
"""处理接收到的邮件数据"""
|
||||
try:
|
||||
# 获取收件人和发件人
|
||||
peer = session.peer
|
||||
mail_from = envelope.mail_from
|
||||
rcpt_tos = envelope.rcpt_tos
|
||||
|
||||
# 获取原始邮件内容
|
||||
data = envelope.content
|
||||
mail = email.message_from_bytes(data, policy=default)
|
||||
|
||||
# 保存邮件到存储服务
|
||||
for rcpt in rcpt_tos:
|
||||
result = await self.mail_store.save_email(mail_from, rcpt, mail, data)
|
||||
|
||||
# 记录日志
|
||||
if result:
|
||||
logger.info(f"邮件已保存: {mail_from} -> {rcpt}, 主题: {mail.get('Subject')}")
|
||||
else:
|
||||
logger.warning(f"邮件未保存: {mail_from} -> {rcpt}, 可能是无效地址")
|
||||
|
||||
return '250 Message accepted for delivery'
|
||||
except Exception as e:
|
||||
logger.error(f"处理邮件时出错: {str(e)}")
|
||||
return '451 Requested action aborted: error in processing'
|
||||
|
||||
|
||||
# 为Windows环境自定义SMTP控制器
|
||||
if IS_WINDOWS:
|
||||
class WindowsSafeController(Controller):
|
||||
"""Windows环境安全的Controller,跳过连接测试"""
|
||||
def _trigger_server(self):
|
||||
"""Windows环境下跳过SMTP服务器自检连接测试"""
|
||||
# 在Windows环境下,我们跳过自检连接测试
|
||||
logger.info("Windows环境: 跳过SMTP服务器连接自检")
|
||||
return
|
||||
|
||||
|
||||
class SMTPServer:
|
||||
"""SMTP服务器实现"""
|
||||
|
||||
def __init__(self, host='0.0.0.0', port=25, mail_store=None):
|
||||
self.host = host
|
||||
self.port = port
|
||||
self.mail_store = mail_store
|
||||
self.controller = None
|
||||
self.server_thread = None
|
||||
|
||||
def start(self):
|
||||
"""启动SMTP服务器"""
|
||||
if self.controller:
|
||||
logger.warning("SMTP服务器已经在运行")
|
||||
return
|
||||
|
||||
try:
|
||||
handler = EmailHandler(self.mail_store)
|
||||
|
||||
# 根据环境选择适当的Controller
|
||||
if IS_WINDOWS:
|
||||
# Windows环境使用自定义Controller
|
||||
logger.info(f"Windows环境: 使用自定义Controller启动SMTP服务器 {self.host}:{self.port}")
|
||||
self.controller = WindowsSafeController(
|
||||
handler,
|
||||
hostname=self.host,
|
||||
port=self.port
|
||||
)
|
||||
else:
|
||||
# 非Windows环境使用标准Controller
|
||||
self.controller = Controller(
|
||||
handler,
|
||||
hostname=self.host,
|
||||
port=self.port
|
||||
)
|
||||
|
||||
# 在单独的线程中启动服务器
|
||||
self.server_thread = threading.Thread(
|
||||
target=self.controller.start,
|
||||
daemon=True
|
||||
)
|
||||
self.server_thread.start()
|
||||
|
||||
logger.info(f"SMTP服务器已启动在 {self.host}:{self.port}")
|
||||
return True
|
||||
except Exception as e:
|
||||
logger.error(f"启动SMTP服务器失败: {str(e)}")
|
||||
return False
|
||||
|
||||
def stop(self):
|
||||
"""停止SMTP服务器"""
|
||||
if not self.controller:
|
||||
logger.warning("SMTP服务器没有运行")
|
||||
return
|
||||
|
||||
try:
|
||||
self.controller.stop()
|
||||
self.controller = None
|
||||
self.server_thread = None
|
||||
logger.info("SMTP服务器已停止")
|
||||
return True
|
||||
except Exception as e:
|
||||
logger.error(f"停止SMTP服务器失败: {str(e)}")
|
||||
return False
|
||||
0
app/templates/index.html
Normal file
0
app/templates/index.html
Normal file
Reference in New Issue
Block a user