"""
Portfolio Monitor Service.
Runs scheduled AI analysis on manual positions and sends notifications.
"""
from __future__ import annotations
import hashlib
import json
import threading
import time
import traceback
from typing import Any, Dict, List, Optional
from app.utils.db import get_db_connection
from app.utils.logger import get_logger
from app.services.analysis import AnalysisService
from app.services.signal_notifier import SignalNotifier
from app.services.kline import KlineService
logger = get_logger(__name__)
DEFAULT_USER_ID = 1
_monitor_thread: Optional[threading.Thread] = None
_stop_event = threading.Event()
# 多语言消息模板
ALERT_MESSAGES = {
'zh-CN': {
'price_above': '🔔 价格突破预警: {symbol} 当前价格 ${current_price:.4f} 已突破 ${threshold:.4f}',
'price_below': '🔔 价格跌破预警: {symbol} 当前价格 ${current_price:.4f} 已跌破 ${threshold:.4f}',
'pnl_above': '🎉 盈利预警: {symbol} 当前盈亏 {pnl_percent:.1f}% 已达到 {threshold:.1f}% 目标',
'pnl_below': '⚠️ 亏损预警: {symbol} 当前盈亏 {pnl_percent:.1f}% 已触及 {threshold:.1f}% 止损线',
'alert_title': '价格/盈亏预警'
},
'en-US': {
'price_above': '🔔 Price Alert: {symbol} current price ${current_price:.4f} has exceeded ${threshold:.4f}',
'price_below': '🔔 Price Alert: {symbol} current price ${current_price:.4f} has dropped below ${threshold:.4f}',
'pnl_above': '🎉 Profit Alert: {symbol} P&L {pnl_percent:.1f}% has reached {threshold:.1f}% target',
'pnl_below': '⚠️ Loss Alert: {symbol} P&L {pnl_percent:.1f}% has hit {threshold:.1f}% stop-loss',
'alert_title': 'Price/P&L Alert'
}
}
def _get_alert_message(alert_type: str, language: str = 'en-US', **kwargs) -> str:
"""Get localized alert message."""
lang = 'zh-CN' if language and language.startswith('zh') else 'en-US'
templates = ALERT_MESSAGES.get(lang, ALERT_MESSAGES['en-US'])
template = templates.get(alert_type, '')
if template:
return template.format(**kwargs)
return ''
def _get_alert_title(language: str = 'en-US') -> str:
"""Get localized alert title."""
lang = 'zh-CN' if language and language.startswith('zh') else 'en-US'
return ALERT_MESSAGES.get(lang, ALERT_MESSAGES['en-US']).get('alert_title', 'Alert')
def _now_ts() -> int:
return int(time.time())
def _safe_json_loads(value, default=None):
"""Safely parse JSON string."""
if default is None:
default = {}
if isinstance(value, dict):
return value
if isinstance(value, list):
return value
if isinstance(value, str) and value.strip():
try:
return json.loads(value)
except Exception:
return default
return default
def _get_positions_for_monitor(position_ids: List[int] = None, user_id: int = None) -> List[Dict[str, Any]]:
"""Get positions, optionally filtered by IDs and user_id."""
try:
kline_service = KlineService()
effective_user_id = user_id if user_id is not None else DEFAULT_USER_ID
with get_db_connection() as db:
cur = db.cursor()
if position_ids:
placeholders = ','.join(['?' for _ in position_ids])
cur.execute(
f"""
SELECT id, market, symbol, name, side, quantity, entry_price
FROM qd_manual_positions
WHERE user_id = ? AND id IN ({placeholders})
""",
[effective_user_id] + list(position_ids)
)
else:
cur.execute(
"""
SELECT id, market, symbol, name, side, quantity, entry_price
FROM qd_manual_positions
WHERE user_id = ?
""",
(effective_user_id,)
)
rows = cur.fetchall() or []
cur.close()
positions = []
for row in rows:
market = row.get('market')
symbol = row.get('symbol')
entry_price = float(row.get('entry_price') or 0)
quantity = float(row.get('quantity') or 0)
side = row.get('side') or 'long'
# Get current price (use realtime price API)
current_price = 0
try:
price_data = kline_service.get_realtime_price(market, symbol)
current_price = float(price_data.get('price') or 0)
except Exception:
pass
# Calculate PnL
if side == 'long':
pnl = (current_price - entry_price) * quantity
else:
pnl = (entry_price - current_price) * quantity
pnl_percent = round(pnl / (entry_price * quantity) * 100, 2) if entry_price * quantity > 0 else 0
positions.append({
'id': row.get('id'),
'market': market,
'symbol': symbol,
'name': row.get('name') or symbol,
'side': side,
'quantity': quantity,
'entry_price': entry_price,
'current_price': current_price,
'pnl': round(pnl, 2),
'pnl_percent': pnl_percent
})
return positions
except Exception as e:
logger.error(f"_get_positions_for_monitor failed: {e}")
return []
def _run_ai_analysis(positions: List[Dict[str, Any]], config: Dict[str, Any]) -> Dict[str, Any]:
"""
Run full multi-agent AI analysis on positions.
Uses the same 13-agent analysis flow as the AI Analysis page.
"""
try:
language = config.get('language', 'en-US')
custom_prompt = config.get('prompt', '')
# Analyze each position using the full agent analysis flow
position_analyses = []
for pos in positions:
market = pos.get('market')
symbol = pos.get('symbol')
name = pos.get('name') or symbol
if not market or not symbol:
continue
try:
logger.info(f"Running multi-agent analysis for {market}:{symbol}")
# Use the full AnalysisService (13-agent flow)
analysis_result = AnalysisService().analyze(
market=market,
symbol=symbol,
language=language,
timeframe='1D'
)
# Extract key information from the analysis
final_decision = analysis_result.get('final_decision', {})
trader_decision = analysis_result.get('trader_decision', {})
overview = analysis_result.get('overview', {})
risk_report = analysis_result.get('risk', {})
position_analysis = {
'market': market,
'symbol': symbol,
'name': name,
'entry_price': pos.get('entry_price'),
'current_price': pos.get('current_price'),
'pnl': pos.get('pnl'),
'pnl_percent': pos.get('pnl_percent'),
'quantity': pos.get('quantity'),
'side': pos.get('side'),
# Multi-agent analysis results
'final_decision': final_decision.get('decision', 'HOLD'),
'confidence': final_decision.get('confidence', 50),
'reasoning': final_decision.get('reasoning', ''),
'trader_decision': trader_decision.get('decision', 'HOLD'),
'trader_reasoning': trader_decision.get('reasoning', ''),
'overview_report': overview.get('report', ''),
'risk_report': risk_report.get('report', ''),
'error': analysis_result.get('error')
}
position_analyses.append(position_analysis)
logger.info(f"Analysis completed for {market}:{symbol}: {final_decision.get('decision', 'N/A')}")
except Exception as e:
logger.error(f"Failed to analyze {market}:{symbol}: {e}")
position_analyses.append({
'market': market,
'symbol': symbol,
'name': name,
'error': str(e)
})
# Build comprehensive report
analysis_report = _build_comprehensive_report(positions, position_analyses, language, custom_prompt)
return {
'success': True,
'analysis': analysis_report,
'position_analyses': position_analyses,
'position_count': len(positions),
'analyzed_count': len([p for p in position_analyses if not p.get('error')]),
'timestamp': _now_ts()
}
except Exception as e:
logger.error(f"_run_ai_analysis failed: {e}")
logger.error(traceback.format_exc())
return {
'success': False,
'error': str(e),
'timestamp': _now_ts()
}
def _build_comprehensive_report(
positions: List[Dict[str, Any]],
position_analyses: List[Dict[str, Any]],
language: str,
custom_prompt: str = ''
) -> str:
"""Build a comprehensive text report (backward compatible)."""
# Use HTML report as the main format
return _build_html_report(positions, position_analyses, language, custom_prompt)
def _build_html_report(
positions: List[Dict[str, Any]],
position_analyses: List[Dict[str, Any]],
language: str,
custom_prompt: str = ''
) -> str:
"""Build a beautiful HTML report with collapsible sections."""
# Calculate portfolio summary
total_cost = sum(float(p.get('entry_price', 0)) * float(p.get('quantity', 0)) for p in positions)
total_pnl = sum(float(p.get('pnl', 0)) for p in positions)
total_pnl_percent = round(total_pnl / total_cost * 100, 2) if total_cost > 0 else 0
total_market_value = sum(float(p.get('current_price', 0)) * float(p.get('quantity', 0)) for p in positions)
# Count recommendations
buy_count = len([p for p in position_analyses if p.get('final_decision') == 'BUY'])
sell_count = len([p for p in position_analyses if p.get('final_decision') == 'SELL'])
hold_count = len([p for p in position_analyses if p.get('final_decision') == 'HOLD'])
is_zh = language.startswith('zh')
# Text translations
texts = {
'title': '投资组合AI分析报告' if is_zh else 'Portfolio AI Analysis Report',
'subtitle': '由 QuantDinger 多智能体分析系统生成' if is_zh else 'Generated by QuantDinger Multi-Agent Analysis System',
'overview': '组合概览' if is_zh else 'Portfolio Overview',
'positions': '持仓数量' if is_zh else 'Positions',
'total_value': '总市值' if is_zh else 'Total Value',
'total_cost': '总成本' if is_zh else 'Total Cost',
'total_pnl': '总盈亏' if is_zh else 'Total P&L',
'ai_recommendations': '🤖 AI智能分析建议' if is_zh else '🤖 AI Recommendations',
'buy': '买入' if is_zh else 'Buy',
'sell': '卖出' if is_zh else 'Sell',
'hold': '持有' if is_zh else 'Hold',
'position_analysis': '📈 各持仓详细分析' if is_zh else '📈 Position Analysis',
'current_price': '当前价格' if is_zh else 'Current',
'entry_price': '买入价' if is_zh else 'Entry',
'pnl': '盈亏' if is_zh else 'P&L',
'quantity': '数量' if is_zh else 'Qty',
'side': '方向' if is_zh else 'Side',
'long': '做多' if is_zh else 'Long',
'short': '做空' if is_zh else 'Short',
'ai_decision': 'AI决策' if is_zh else 'AI Decision',
'confidence': '置信度' if is_zh else 'Confidence',
'reasoning': '分析摘要' if is_zh else 'Summary',
'trader_report': '📋 交易员详细评估' if is_zh else '📋 Trader Analysis',
'risk_report': '⚠️ 风险评估' if is_zh else '⚠️ Risk Assessment',
'overview_report': '📊 市场概览' if is_zh else '📊 Market Overview',
'click_expand': '点击展开详情' if is_zh else 'Click to expand',
'user_focus': '👤 用户关注点' if is_zh else '👤 User Focus',
'generated_at': '报告生成时间' if is_zh else 'Generated at',
'disclaimer': '本报告仅供参考,不构成投资建议。' if is_zh else 'For reference only. Not investment advice.',
'analysis_failed': '分析失败' if is_zh else 'Analysis failed'
}
# CSS Styles
css = '''
'''
# Build HTML
pnl_class = 'positive' if total_pnl >= 0 else 'negative'
pnl_sign = '+' if total_pnl >= 0 else ''
html = f'''
{css}
{texts['overview']}
{texts['positions']}
{len(positions)}
{texts['total_value']}
${total_market_value:,.2f}
{texts['total_cost']}
${total_cost:,.2f}
{texts['total_pnl']}
{pnl_sign}${total_pnl:,.2f}({pnl_sign}{total_pnl_percent:.1f}%)
{texts['ai_recommendations']}
🟢
{buy_count}
{texts['buy']}
🔴
{sell_count}
{texts['sell']}
🟡
{hold_count}
{texts['hold']}
{texts['position_analysis']}
'''
for pa in position_analyses:
symbol = pa.get('symbol', '')
name = pa.get('name', symbol)
market = pa.get('market', '')
if pa.get('error'):
html += f'''
{texts['analysis_failed']}: {pa.get('error')}
'''
continue
decision = pa.get('final_decision', 'HOLD')
decision_lower = decision.lower()
decision_text = texts.get(decision_lower, decision)
confidence = pa.get('confidence', 50)
current_price = pa.get('current_price', 0)
entry_price = pa.get('entry_price', 0)
pnl = pa.get('pnl', 0)
pnl_pct = pa.get('pnl_percent', 0)
quantity = pa.get('quantity', 0)
side = pa.get('side', 'long')
side_text = texts['long'] if side == 'long' else texts['short']
pnl_class = 'positive' if pnl >= 0 else 'negative'
pnl_sign = '+' if pnl >= 0 else ''
reasoning = pa.get('reasoning', '')
trader_reasoning = pa.get('trader_reasoning', '')
overview_report = pa.get('overview_report', '')
risk_report = pa.get('risk_report', '')
html += f'''
{texts['current_price']}
${current_price:.4f}
{texts['entry_price']}
${entry_price:.4f}
{texts['pnl']}
{pnl_sign}${pnl:.2f} ({pnl_sign}{pnl_pct:.1f}%)
{texts['quantity']} / {texts['side']}
{quantity} / {side_text}
'''
# Reasoning summary
if reasoning:
html += f'''
{texts['reasoning']}
{reasoning[:500]}{'...' if len(reasoning) > 500 else ''}
'''
# Generate unique ID for collapsible sections (use symbol hash to avoid special chars)
section_id_base = hashlib.md5(f"{symbol}_{market}".encode()).hexdigest()[:8]
# Collapsible: Trader Analysis
if trader_reasoning:
trader_id = f"trader_{section_id_base}"
html += f'''
'''
# Collapsible: Market Overview
if overview_report:
overview_id = f"overview_{section_id_base}"
html += f'''
'''
# Collapsible: Risk Assessment
if risk_report:
risk_id = f"risk_{section_id_base}"
html += f'''
'''
html += '''
'''
# User focus section
if custom_prompt:
html += f'''
{texts['user_focus']}
{custom_prompt}
'''
# Footer
html += f'''
'''
return html
def _build_telegram_report(
positions: List[Dict[str, Any]],
position_analyses: List[Dict[str, Any]],
language: str,
custom_prompt: str = ''
) -> str:
"""Build a concise report suitable for Telegram (HTML format)."""
# Calculate summary
total_cost = sum(float(p.get('entry_price', 0)) * float(p.get('quantity', 0)) for p in positions)
total_pnl = sum(float(p.get('pnl', 0)) for p in positions)
total_pnl_percent = round(total_pnl / total_cost * 100, 2) if total_cost > 0 else 0
buy_count = len([p for p in position_analyses if p.get('final_decision') == 'BUY'])
sell_count = len([p for p in position_analyses if p.get('final_decision') == 'SELL'])
hold_count = len([p for p in position_analyses if p.get('final_decision') == 'HOLD'])
is_zh = language.startswith('zh')
pnl_sign = '+' if total_pnl >= 0 else ''
if is_zh:
lines = [
"📊 投资组合AI分析报告",
"",
"📈 组合概览",
f"• 持仓: {len(positions)} 个",
f"• 总成本: ${total_cost:,.2f}",
f"• 总盈亏: {pnl_sign}${total_pnl:,.2f} ({pnl_sign}{total_pnl_percent:.1f}%)",
"",
"🤖 AI建议汇总",
f"🟢 买入: {buy_count} | 🔴 卖出: {sell_count} | 🟡 持有: {hold_count}",
"",
"📋 持仓分析"
]
for pa in position_analyses:
if pa.get('error'):
lines.append(f"⚠️ {pa.get('name', pa.get('symbol'))}: 分析失败")
continue
decision = pa.get('final_decision', 'HOLD')
emoji = {'BUY': '🟢', 'SELL': '🔴', 'HOLD': '🟡'}.get(decision, '⚪')
text = {'BUY': '买入', 'SELL': '卖出', 'HOLD': '持有'}.get(decision, '持有')
pnl = pa.get('pnl', 0)
pnl_pct = pa.get('pnl_percent', 0)
pnl_s = '+' if pnl >= 0 else ''
lines.append(f"\n{emoji} {pa.get('name', pa.get('symbol'))} ({pa.get('market')}/{pa.get('symbol')})")
lines.append(f" 💰 ${pa.get('current_price', 0):.2f} | 盈亏: {pnl_s}${pnl:.2f} ({pnl_s}{pnl_pct:.1f}%)")
lines.append(f" 🎯 建议: {text} (置信度 {pa.get('confidence', 50)}%)")
reasoning = pa.get('reasoning', '')
if reasoning:
lines.append(f" 📝 {reasoning[:150]}{'...' if len(reasoning) > 150 else ''}")
if custom_prompt:
lines.extend(["", f"👤 关注点: {custom_prompt}"])
lines.extend([
"",
"─────────────────────",
f"⏰ {time.strftime('%Y-%m-%d %H:%M')}",
"由 QuantDinger 多智能体系统生成"
])
else:
lines = [
"📊 Portfolio AI Analysis Report",
"",
"📈 Overview",
f"• Positions: {len(positions)}",
f"• Total Cost: ${total_cost:,.2f}",
f"• Total P&L: {pnl_sign}${total_pnl:,.2f} ({pnl_sign}{total_pnl_percent:.1f}%)",
"",
"🤖 AI Recommendations",
f"🟢 Buy: {buy_count} | 🔴 Sell: {sell_count} | 🟡 Hold: {hold_count}",
"",
"📋 Position Analysis"
]
for pa in position_analyses:
if pa.get('error'):
lines.append(f"⚠️ {pa.get('name', pa.get('symbol'))}: Analysis failed")
continue
decision = pa.get('final_decision', 'HOLD')
emoji = {'BUY': '🟢', 'SELL': '🔴', 'HOLD': '🟡'}.get(decision, '⚪')
pnl = pa.get('pnl', 0)
pnl_pct = pa.get('pnl_percent', 0)
pnl_s = '+' if pnl >= 0 else ''
lines.append(f"\n{emoji} {pa.get('name', pa.get('symbol'))} ({pa.get('market')}/{pa.get('symbol')})")
lines.append(f" 💰 ${pa.get('current_price', 0):.2f} | P&L: {pnl_s}${pnl:.2f} ({pnl_s}{pnl_pct:.1f}%)")
lines.append(f" 🎯 Rec: {decision} (Conf: {pa.get('confidence', 50)}%)")
reasoning = pa.get('reasoning', '')
if reasoning:
lines.append(f" 📝 {reasoning[:150]}{'...' if len(reasoning) > 150 else ''}")
if custom_prompt:
lines.extend(["", f"👤 Focus: {custom_prompt}"])
lines.extend([
"",
"─────────────────────",
f"⏰ {time.strftime('%Y-%m-%d %H:%M')}",
"Generated by QuantDinger Multi-Agent System"
])
return '\n'.join(lines)
def _send_monitor_notification(
monitor_name: str,
result: Dict[str, Any],
notification_config: Dict[str, Any],
positions: List[Dict[str, Any]] = None,
position_analyses: List[Dict[str, Any]] = None,
language: str = 'en-US',
custom_prompt: str = '',
user_id: int = None
) -> None:
"""Send notification with analysis result using appropriate format for each channel."""
try:
notifier = SignalNotifier()
effective_user_id = user_id if user_id is not None else DEFAULT_USER_ID
channels = notification_config.get('channels', ['browser'])
targets = notification_config.get('targets', {})
title = f"📊 资产监测: {monitor_name}" if language.startswith('zh') else f"📊 Portfolio Monitor: {monitor_name}"
if not result.get('success'):
error_title = f"⚠️ 资产监测失败: {monitor_name}" if language.startswith('zh') else f"⚠️ Monitor Failed: {monitor_name}"
error_msg = f"分析失败: {result.get('error', 'Unknown error')}" if language.startswith('zh') else f"Analysis failed: {result.get('error', 'Unknown error')}"
for channel in channels:
try:
ch = str(channel).strip().lower()
if ch == 'browser':
with get_db_connection() as db:
cur = db.cursor()
cur.execute(
"""
INSERT INTO qd_strategy_notifications
(user_id, strategy_id, symbol, signal_type, channels, title, message, payload_json, created_at)
VALUES (?, NULL, ?, ?, ?, ?, ?, ?, NOW())
""",
(effective_user_id, 'PORTFOLIO', 'ai_monitor', 'browser', error_title, error_msg,
json.dumps(result, ensure_ascii=False))
)
db.commit()
cur.close()
elif ch == 'telegram':
chat_id = targets.get('telegram', '')
if chat_id:
notifier._notify_telegram(chat_id=chat_id, text=f"{error_title}\n\n{error_msg}", parse_mode="HTML")
elif ch == 'email':
to_email = targets.get('email', '')
if to_email:
notifier._notify_email(to_email=to_email, subject=error_title, body_text=error_msg)
except Exception as e:
logger.warning(f"Failed to send error notification to {channel}: {e}")
return
# Generate reports for different channels
html_report = result.get('analysis', '') # This is already HTML from _build_html_report
# Generate Telegram-specific report if we have the data
telegram_report = ''
if positions is not None and position_analyses is not None:
telegram_report = _build_telegram_report(positions, position_analyses, language, custom_prompt)
else:
# Fallback: strip HTML tags for Telegram
import re
telegram_report = re.sub(r'<[^>]+>', '', html_report)
if len(telegram_report) > 4000:
telegram_report = telegram_report[:4000] + '...'
# Send to each channel
for channel in channels:
try:
ch = str(channel).strip().lower()
if ch == 'browser':
# Browser notification uses HTML report
with get_db_connection() as db:
cur = db.cursor()
cur.execute(
"""
INSERT INTO qd_strategy_notifications
(user_id, strategy_id, symbol, signal_type, channels, title, message, payload_json, created_at)
VALUES (?, NULL, ?, ?, ?, ?, ?, ?, NOW())
""",
(effective_user_id, 'PORTFOLIO', 'ai_monitor', 'browser', title, html_report,
json.dumps(result, ensure_ascii=False))
)
db.commit()
cur.close()
elif ch == 'telegram':
chat_id = targets.get('telegram', '')
if chat_id:
# Use Telegram-optimized format
notifier._notify_telegram(
chat_id=chat_id,
text=telegram_report,
parse_mode="HTML"
)
elif ch == 'email':
to_email = targets.get('email', '')
if to_email:
# Email uses full HTML report
notifier._notify_email(
to_email=to_email,
subject=title,
body_text=html_report,
body_html=html_report # Send as HTML email
)
elif ch == 'webhook':
url = targets.get('webhook', '')
if url:
notifier._notify_webhook(
url=url,
payload={
'type': 'portfolio_monitor',
'monitor_name': monitor_name,
'result': result,
'html_report': html_report
}
)
except Exception as e:
logger.warning(f"Failed to send notification to {channel}: {e}")
except Exception as e:
logger.error(f"_send_monitor_notification failed: {e}")
def run_single_monitor(monitor_id: int, override_language: str = None, user_id: int = None) -> Dict[str, Any]:
"""Run a single monitor and return the result.
Args:
monitor_id: The monitor ID to run
override_language: Optional language override (e.g., 'zh-CN', 'en-US')
If provided, will override the language in monitor config
user_id: Optional user ID for user isolation
"""
try:
# Use provided user_id or default
effective_user_id = user_id if user_id is not None else DEFAULT_USER_ID
with get_db_connection() as db:
cur = db.cursor()
cur.execute(
"""
SELECT id, user_id, name, position_ids, monitor_type, config, notification_config
FROM qd_position_monitors
WHERE id = ? AND user_id = ?
""",
(monitor_id, effective_user_id)
)
row = cur.fetchone()
cur.close()
if not row:
return {'success': False, 'error': 'Monitor not found'}
monitor_user_id = int(row.get('user_id') or effective_user_id)
name = row.get('name') or f'Monitor #{monitor_id}'
position_ids = _safe_json_loads(row.get('position_ids'), [])
monitor_type = row.get('monitor_type') or 'ai'
config = _safe_json_loads(row.get('config'), {})
notification_config = _safe_json_loads(row.get('notification_config'), {})
# Override language if provided (from frontend)
if override_language:
config['language'] = override_language
# Get positions for this user
positions = _get_positions_for_monitor(position_ids if position_ids else None, user_id=monitor_user_id)
if not positions:
return {'success': False, 'error': 'No positions to analyze'}
# Run analysis based on type
if monitor_type == 'ai':
result = _run_ai_analysis(positions, config)
else:
# For other types, we can add price_alert, pnl_alert logic later
result = {'success': False, 'error': f'Unsupported monitor type: {monitor_type}'}
# Update monitor record
interval_minutes = int(config.get('interval_minutes') or 60)
with get_db_connection() as db:
cur = db.cursor()
cur.execute(
"""
UPDATE qd_position_monitors
SET last_run_at = NOW(),
next_run_at = NOW() + INTERVAL '%s minutes',
last_result = ?,
run_count = run_count + 1,
updated_at = NOW()
WHERE id = ?
""",
(interval_minutes, json.dumps(result, ensure_ascii=False), monitor_id)
)
db.commit()
cur.close()
# Send notification
if notification_config.get('channels'):
language = config.get('language', 'en-US')
custom_prompt = config.get('prompt', '')
position_analyses = result.get('position_analyses', [])
_send_monitor_notification(
monitor_name=name,
result=result,
notification_config=notification_config,
positions=positions,
position_analyses=position_analyses,
language=language,
custom_prompt=custom_prompt,
user_id=monitor_user_id
)
return result
except Exception as e:
logger.error(f"run_single_monitor failed: {e}")
logger.error(traceback.format_exc())
return {'success': False, 'error': str(e)}
def _check_position_alerts():
"""Check all active alerts and trigger notifications if conditions are met."""
from datetime import datetime, timezone
try:
kline_service = KlineService()
notifier = SignalNotifier()
now = datetime.now(timezone.utc)
with get_db_connection() as db:
cur = db.cursor()
# Get active alerts for all users that haven't been triggered (or can repeat)
cur.execute(
"""
SELECT a.id, a.user_id, a.position_id, a.market, a.symbol, a.alert_type, a.threshold,
a.notification_config, a.is_triggered, a.last_triggered_at, a.repeat_interval,
p.entry_price, p.quantity, p.side, p.name as position_name
FROM qd_position_alerts a
LEFT JOIN qd_manual_positions p ON a.position_id = p.id
WHERE a.is_active = 1
"""
)
alerts = cur.fetchall() or []
cur.close()
for alert in alerts:
try:
alert_id = alert.get('id')
alert_user_id = int(alert.get('user_id') or 1)
alert_type = alert.get('alert_type')
threshold = float(alert.get('threshold') or 0)
market = alert.get('market')
symbol = alert.get('symbol')
is_triggered = bool(alert.get('is_triggered'))
last_triggered_at = alert.get('last_triggered_at') # datetime or None
repeat_interval = int(alert.get('repeat_interval') or 0)
notification_config = _safe_json_loads(alert.get('notification_config'), {})
# Check if we can trigger (not triggered yet, or repeat interval passed)
can_trigger = not is_triggered
if is_triggered and repeat_interval > 0 and last_triggered_at:
# Convert last_triggered_at to timezone-aware if needed
if last_triggered_at.tzinfo is None:
last_triggered_at = last_triggered_at.replace(tzinfo=timezone.utc)
elapsed_seconds = (now - last_triggered_at).total_seconds()
if elapsed_seconds >= repeat_interval:
can_trigger = True
if not can_trigger:
continue
# Get current price (use realtime price API)
current_price = 0
try:
price_data = kline_service.get_realtime_price(market, symbol)
current_price = float(price_data.get('price') or 0)
except Exception:
continue
if current_price <= 0:
continue
triggered = False
alert_message = ""
# Get language from notification_config (saved when alert was created)
alert_language = notification_config.get('language', 'en-US')
if alert_type == 'price_above':
if current_price >= threshold:
triggered = True
alert_message = _get_alert_message(
'price_above', alert_language,
symbol=symbol, current_price=current_price, threshold=threshold
)
elif alert_type == 'price_below':
if current_price <= threshold:
triggered = True
alert_message = _get_alert_message(
'price_below', alert_language,
symbol=symbol, current_price=current_price, threshold=threshold
)
elif alert_type in ('pnl_above', 'pnl_below'):
entry_price = float(alert.get('entry_price') or 0)
quantity = float(alert.get('quantity') or 0)
side = alert.get('side') or 'long'
if entry_price > 0 and quantity > 0:
if side == 'long':
pnl = (current_price - entry_price) * quantity
else:
pnl = (entry_price - current_price) * quantity
pnl_percent = pnl / (entry_price * quantity) * 100
if alert_type == 'pnl_above' and pnl_percent >= threshold:
triggered = True
alert_message = _get_alert_message(
'pnl_above', alert_language,
symbol=symbol, pnl_percent=pnl_percent, threshold=threshold
)
elif alert_type == 'pnl_below' and pnl_percent <= threshold:
triggered = True
alert_message = _get_alert_message(
'pnl_below', alert_language,
symbol=symbol, pnl_percent=pnl_percent, threshold=threshold
)
if triggered:
logger.info(f"Alert #{alert_id} triggered: {alert_message}")
# Update alert status
with get_db_connection() as db:
cur = db.cursor()
cur.execute(
"""
UPDATE qd_position_alerts
SET is_triggered = 1, last_triggered_at = NOW(), trigger_count = trigger_count + 1, updated_at = NOW()
WHERE id = ?
""",
(alert_id,)
)
db.commit()
cur.close()
# Send notification
channels = notification_config.get('channels', ['browser'])
targets = notification_config.get('targets', {})
alert_title = _get_alert_title(alert_language)
for channel in channels:
try:
ch = str(channel).strip().lower()
if ch == 'browser':
with get_db_connection() as db:
cur = db.cursor()
cur.execute(
"""
INSERT INTO qd_strategy_notifications
(user_id, strategy_id, symbol, signal_type, channels, title, message, payload_json, created_at)
VALUES (?, NULL, ?, ?, ?, ?, ?, ?, NOW())
""",
(alert_user_id, symbol, 'price_alert', 'browser', alert_title, alert_message,
json.dumps({'alert_id': alert_id, 'alert_type': alert_type}, ensure_ascii=False))
)
db.commit()
cur.close()
elif ch == 'telegram':
chat_id = targets.get('telegram', '')
if chat_id:
notifier._notify_telegram(chat_id=chat_id, text=alert_message, parse_mode="HTML")
elif ch == 'email':
to_email = targets.get('email', '')
if to_email:
notifier._notify_email(to_email=to_email, subject=alert_title, body_text=alert_message)
except Exception as e:
logger.warning(f"Failed to send alert notification: {e}")
except Exception as e:
logger.warning(f"Error processing alert: {e}")
except Exception as e:
logger.error(f"_check_position_alerts failed: {e}")
def notify_strategy_signal_for_positions(market: str, symbol: str, signal_type: str, signal_detail: str, user_id: int = None):
"""
Called when a strategy signal is triggered.
Check if user has manual positions in this symbol and send notification.
"""
try:
symbol = (symbol or '').strip().upper()
if not symbol:
return
with get_db_connection() as db:
cur = db.cursor()
# Query positions for all users or specific user
if user_id is not None:
cur.execute(
"""
SELECT id, user_id, market, symbol, name, side, quantity, entry_price, group_name
FROM qd_manual_positions
WHERE user_id = ? AND symbol = ?
""",
(user_id, symbol)
)
else:
cur.execute(
"""
SELECT id, user_id, market, symbol, name, side, quantity, entry_price, group_name
FROM qd_manual_positions
WHERE symbol = ?
""",
(symbol,)
)
positions = cur.fetchall() or []
cur.close()
if not positions:
return
# User has positions in this symbol - send notification
notifier = SignalNotifier()
now = _now_ts()
for pos in positions:
pos_user_id = int(pos.get('user_id') or 1)
pos_name = pos.get('name') or symbol
pos_side = pos.get('side') or 'long'
quantity = float(pos.get('quantity') or 0)
entry_price = float(pos.get('entry_price') or 0)
title = f"🔗 策略信号联动: {pos_name}"
message = f"""策略发出 {signal_type} 信号!
标的: {market}/{symbol}
您的持仓: {pos_side.upper()} {quantity} @ {entry_price:.4f}
信号详情:
{signal_detail}
请注意检查您的持仓是否需要调整。"""
# Save browser notification
with get_db_connection() as db:
cur = db.cursor()
cur.execute(
"""
INSERT INTO qd_strategy_notifications
(user_id, strategy_id, symbol, signal_type, channels, title, message, payload_json, created_at)
VALUES (?, NULL, ?, ?, ?, ?, ?, ?, NOW())
""",
(pos_user_id, symbol, 'strategy_linkage', 'browser', title, message,
json.dumps({'signal_type': signal_type}, ensure_ascii=False))
)
db.commit()
cur.close()
logger.info(f"Strategy signal linkage: notified {len(positions)} position(s) for {symbol}")
except Exception as e:
logger.error(f"notify_strategy_signal_for_positions failed: {e}")
def _monitor_loop():
"""Background loop that checks and runs due monitors."""
logger.info("Portfolio monitor background loop started")
while not _stop_event.is_set():
try:
# 1. Check position alerts (price/pnl alerts) for all users
_check_position_alerts()
# 2. Find AI monitors that are due for all users
with get_db_connection() as db:
cur = db.cursor()
cur.execute(
"""
SELECT id, user_id FROM qd_position_monitors
WHERE is_active = 1 AND next_run_at <= NOW()
ORDER BY next_run_at ASC
LIMIT 10
"""
)
rows = cur.fetchall() or []
cur.close()
for row in rows:
if _stop_event.is_set():
break
monitor_id = row.get('id')
monitor_user_id = int(row.get('user_id') or 1)
if monitor_id:
logger.info(f"Running due monitor #{monitor_id} for user #{monitor_user_id}")
try:
run_single_monitor(monitor_id, user_id=monitor_user_id)
except Exception as e:
logger.error(f"Monitor #{monitor_id} execution failed: {e}")
except Exception as e:
logger.error(f"Monitor loop error: {e}")
# Sleep for 30 seconds before next check
_stop_event.wait(30)
logger.info("Portfolio monitor background loop stopped")
def start_monitor_service():
"""Start the background monitor service."""
global _monitor_thread
if _monitor_thread and _monitor_thread.is_alive():
logger.info("Portfolio monitor service already running")
return
_stop_event.clear()
_monitor_thread = threading.Thread(target=_monitor_loop, daemon=True, name="PortfolioMonitor")
_monitor_thread.start()
logger.info("Portfolio monitor service started")
def stop_monitor_service():
"""Stop the background monitor service."""
global _monitor_thread
_stop_event.set()
if _monitor_thread:
_monitor_thread.join(timeout=5)
_monitor_thread = None
logger.info("Portfolio monitor service stopped")