✨ feat(scheduler): 新增热点自动采集功能并优化发布路径
- 新增热点自动采集后台线程,支持定时搜索关键词并执行 AI 分析,结果缓存至结构化状态 - 新增热点分析状态管理接口,提供线程安全的 `get_last_analysis` 和 `set_last_analysis` 方法 - 新增热点数据桥接函数 `feed_hotspot_to_engine`,将分析结果注入 TopicEngine 实现热点加权推荐 - 新增热点选题下拉组件,分析完成后自动填充推荐选题,选中后自动写入选题输入框 - 优化 `generate_from_hotspot` 函数,自动获取结构化分析摘要并增强生成上下文 - 新增热点自动采集配置节点,支持通过 `config.json` 管理关键词和采集间隔 ♻️ refactor(queue): 实现智能排期引擎并统一发布路径 - 新增智能排期引擎,基于 `AnalyticsService` 的 `time_weights` 自动计算最优发布时段 - 新增 `PublishQueue.suggest_schedule_time` 和 `auto_schedule_item` 方法,支持时段冲突检测和内容分布控制 - 修改 `generate_to_queue` 函数,新增 `auto_schedule` 和 `auto_approve` 参数,支持自动排期和自动审核 - 重构 `_scheduler_loop` 的自动发布分支,改为调用 `generate_to_queue` 通过队列发布,统一发布路径 - 重构 `auto_publish_once` 函数,移除直接发布逻辑,改为生成内容入队并返回队列信息 - 新增队列时段使用情况查询方法 `get_slot_usage`,支持 UI 热力图展示 📝 docs(openspec): 新增内容排期优化和热点探测优化规范文档 - 新增 `smart-schedule-engine` 规范,定义智能排期引擎的功能需求和场景 - 新增 `unified-publish-path` 规范,定义统一发布路径的改造方案 - 新增 `hotspot-analysis-state` 规范,定义热点分析状态存储的线程安全接口 - 新增 `hotspot-auto-collector` 规范,定义定时热点自动采集的任务流程 - 新增 `hotspot-engine-bridge` 规范,定义热点数据注入 TopicEngine 的桥接机制 - 新增 `hotspot-topic-selector` 规范,定义热点选题下拉组件的交互行为 - 更新 `services-queue`、`services-scheduler` 和 `services-hotspot` 规范,反映功能修改和新增参数 🔧 chore(config): 新增热点自动采集默认配置 - 在 `DEFAULT_CONFIG` 中新增 `hotspot_auto_collect` 配置节点,包含 `enabled`、`keywords` 和 `interval_hours` 字段 - 提供默认关键词列表 `["穿搭", "美妆", "好物"]` 和默认采集间隔 4 小时 🐛 fix(llm): 增强 JSON 解析容错能力 - 新增 `_try_fix_truncated_json` 方法,尝试修复被 token 限制截断的 JSON 输出 - 支持多种截断场景的自动补全,包括字符串值、数组和嵌套对象的截断修复 - 提高 LLM 分析热点等返回 JSON 的函数的稳定性 💄 style(ui): 优化队列管理和热点探测界面 - 在队列生成区域新增自动排期复选框,勾选后隐藏手动排期输入框 - 在日历视图旁新增推荐时段 Markdown 面板,展示各时段权重和建议热力图 - 在热点探测 Tab 新增推荐选题下拉组件,分析完成后动态填充选项 - 在热点探测 Tab 新增热点自动采集控制区域,支持启动、停止和配置采集参数
This commit is contained in:
@@ -533,6 +533,26 @@ class AnalyticsService:
|
||||
advice_parts.append(f" • {p_name}: 权重 {p_info['weight']}分 (出现{p_info['count']}次)")
|
||||
return "\n".join(advice_parts)
|
||||
|
||||
# ========== 时段权重查询 ==========
|
||||
|
||||
_DEFAULT_TIME_WEIGHTS = {
|
||||
"08-11时": {"weight": 70, "count": 0},
|
||||
"12-14时": {"weight": 60, "count": 0},
|
||||
"18-21时": {"weight": 85, "count": 0},
|
||||
"21-24时": {"weight": 75, "count": 0},
|
||||
}
|
||||
|
||||
def get_time_weights(self) -> dict:
|
||||
"""返回各时段权重字典。
|
||||
|
||||
有分析数据时返回 time_weights;无数据时返回默认高流量时段。
|
||||
返回格式: {"18-21时": {"weight": 85, "count": 12}, ...}
|
||||
"""
|
||||
tw = self._weights.get("time_weights", {})
|
||||
if tw:
|
||||
return tw
|
||||
return dict(self._DEFAULT_TIME_WEIGHTS)
|
||||
|
||||
# ========== LLM 深度分析 ==========
|
||||
|
||||
def generate_llm_analysis_prompt(self) -> str:
|
||||
|
||||
@@ -63,6 +63,12 @@ DEFAULT_CONFIG = {
|
||||
"learn_interval": 6,
|
||||
# 内容排期参数
|
||||
"queue_gen_count": 3,
|
||||
# 热点自动采集参数
|
||||
"hotspot_auto_collect": {
|
||||
"enabled": False,
|
||||
"keywords": ["穿搭", "美妆", "好物"],
|
||||
"interval_hours": 4,
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
|
||||
+98
-15
@@ -2,6 +2,7 @@
|
||||
services/hotspot.py
|
||||
热点探测、热点生成、笔记列表缓存(供评论管家主动评论使用)
|
||||
"""
|
||||
import copy
|
||||
import threading
|
||||
import logging
|
||||
|
||||
@@ -14,10 +15,61 @@ from .persona import _resolve_persona
|
||||
|
||||
logger = logging.getLogger("autobot")
|
||||
|
||||
# ---- 共用: 线程安全缓存 ----
|
||||
# 缓存互斥锁,防止并发回调产生竞态(所有缓存共用)
|
||||
_cache_lock = threading.RLock()
|
||||
# 主动评论缓存
|
||||
_cached_proactive_entries: list[dict] = []
|
||||
# 我的笔记评论缓存
|
||||
_cached_my_note_entries: list[dict] = []
|
||||
|
||||
# ==================================================
|
||||
# Tab 2: 热点探测
|
||||
# ==================================================
|
||||
|
||||
# 最近一次 LLM 热点分析的结构化结果(线程安全,复用 _cache_lock)
|
||||
_last_analysis: dict | None = None
|
||||
|
||||
|
||||
def get_last_analysis() -> dict | None:
|
||||
"""线程安全地获取最近一次热点分析结果的深拷贝"""
|
||||
with _cache_lock:
|
||||
if _last_analysis is None:
|
||||
return None
|
||||
return copy.deepcopy(_last_analysis)
|
||||
|
||||
|
||||
def set_last_analysis(data: dict) -> None:
|
||||
"""线程安全地更新热点分析结果(合并 hot_topics / suggestions 并去重)"""
|
||||
global _last_analysis
|
||||
with _cache_lock:
|
||||
if _last_analysis is None:
|
||||
_last_analysis = copy.deepcopy(data)
|
||||
else:
|
||||
# 合并 hot_topics
|
||||
existing_topics = _last_analysis.get("hot_topics", [])
|
||||
new_topics = data.get("hot_topics", [])
|
||||
seen = set(existing_topics)
|
||||
for t in new_topics:
|
||||
if t not in seen:
|
||||
existing_topics.append(t)
|
||||
seen.add(t)
|
||||
_last_analysis["hot_topics"] = existing_topics
|
||||
|
||||
# 合并 suggestions(按 topic 去重)
|
||||
existing_sug = _last_analysis.get("suggestions", [])
|
||||
existing_sug_topics = {s.get("topic", "") for s in existing_sug}
|
||||
for s in data.get("suggestions", []):
|
||||
if s.get("topic", "") not in existing_sug_topics:
|
||||
existing_sug.append(s)
|
||||
existing_sug_topics.add(s.get("topic", ""))
|
||||
_last_analysis["suggestions"] = existing_sug
|
||||
|
||||
# 其他字段以最新为准
|
||||
for key in data:
|
||||
if key not in ("hot_topics", "suggestions"):
|
||||
_last_analysis[key] = data[key]
|
||||
|
||||
|
||||
def search_hotspots(keyword, sort_by, mcp_url):
|
||||
"""搜索小红书热门内容"""
|
||||
@@ -36,21 +88,25 @@ def search_hotspots(keyword, sort_by, mcp_url):
|
||||
|
||||
|
||||
def analyze_and_suggest(model, keyword, search_result):
|
||||
"""AI 分析热点并给出建议"""
|
||||
"""AI 分析热点并给出建议,同时缓存结构化结果"""
|
||||
if not search_result:
|
||||
return "❌ 请先搜索", "", ""
|
||||
return "❌ 请先搜索", "", "", gr.update(choices=[], value=None)
|
||||
api_key, base_url, _ = _get_llm_config()
|
||||
if not api_key:
|
||||
return "❌ 请先配置 LLM 提供商", "", ""
|
||||
return "❌ 请先配置 LLM 提供商", "", "", gr.update(choices=[], value=None)
|
||||
try:
|
||||
svc = LLMService(api_key, base_url, model)
|
||||
analysis = svc.analyze_hotspots(search_result)
|
||||
|
||||
# 缓存结构化分析结果(在渲染 Markdown 之前)
|
||||
set_last_analysis(analysis)
|
||||
|
||||
topics = "\n".join(f"• {t}" for t in analysis.get("hot_topics", []))
|
||||
patterns = "\n".join(f"• {p}" for p in analysis.get("title_patterns", []))
|
||||
suggestions_list = analysis.get("suggestions", [])
|
||||
suggestions = "\n".join(
|
||||
f"**{s['topic']}** - {s['reason']}"
|
||||
for s in analysis.get("suggestions", [])
|
||||
for s in suggestions_list
|
||||
)
|
||||
structure = analysis.get("content_structure", "")
|
||||
|
||||
@@ -60,14 +116,22 @@ def analyze_and_suggest(model, keyword, search_result):
|
||||
f"## 📐 内容结构\n{structure}\n\n"
|
||||
f"## 💡 推荐选题\n{suggestions}"
|
||||
)
|
||||
return "✅ 分析完成", summary, keyword
|
||||
|
||||
# 构建选题下拉选项
|
||||
topic_choices = [s["topic"] for s in suggestions_list if s.get("topic")]
|
||||
dropdown_update = gr.update(
|
||||
choices=topic_choices,
|
||||
value=topic_choices[0] if topic_choices else None,
|
||||
)
|
||||
|
||||
return "✅ 分析完成", summary, keyword, dropdown_update
|
||||
except Exception as e:
|
||||
logger.error("热点分析失败: %s", e)
|
||||
return f"❌ 分析失败: {e}", "", ""
|
||||
return f"❌ 分析失败: {e}", "", "", gr.update(choices=[], value=None)
|
||||
|
||||
|
||||
def generate_from_hotspot(model, topic_from_hotspot, style, search_result, sd_model_name, persona_text):
|
||||
"""基于热点分析生成文案(自动适配 SD 模型,支持人设)"""
|
||||
"""基于热点分析生成文案(自动适配 SD 模型,支持人设,增强分析上下文)"""
|
||||
if not topic_from_hotspot:
|
||||
return "", "", "", "", "❌ 请先选择或输入选题"
|
||||
api_key, base_url, _ = _get_llm_config()
|
||||
@@ -76,10 +140,30 @@ def generate_from_hotspot(model, topic_from_hotspot, style, search_result, sd_mo
|
||||
try:
|
||||
svc = LLMService(api_key, base_url, model)
|
||||
persona = _resolve_persona(persona_text) if persona_text else None
|
||||
|
||||
# 构建增强参考上下文:结构化分析摘要 + 原始搜索片段
|
||||
analysis = get_last_analysis()
|
||||
reference_parts = []
|
||||
if analysis:
|
||||
topics_str = ", ".join(analysis.get("hot_topics", [])[:5])
|
||||
sug_str = "; ".join(
|
||||
s.get("topic", "") for s in analysis.get("suggestions", [])[:5]
|
||||
)
|
||||
structure = analysis.get("content_structure", "")
|
||||
analysis_summary = (
|
||||
f"[热点分析摘要] 热门选题: {topics_str}\n"
|
||||
f"推荐方向: {sug_str}\n"
|
||||
f"内容结构建议: {structure}\n\n"
|
||||
)
|
||||
reference_parts.append(analysis_summary)
|
||||
if search_result:
|
||||
reference_parts.append(search_result)
|
||||
combined_reference = "".join(reference_parts)[:3000]
|
||||
|
||||
data = svc.generate_copy_with_reference(
|
||||
topic=topic_from_hotspot,
|
||||
style=style,
|
||||
reference_notes=search_result[:2000],
|
||||
reference_notes=combined_reference,
|
||||
sd_model_name=sd_model_name,
|
||||
persona=persona,
|
||||
)
|
||||
@@ -95,19 +179,18 @@ def generate_from_hotspot(model, topic_from_hotspot, style, search_result, sd_mo
|
||||
return "", "", "", "", f"❌ 生成失败: {e}"
|
||||
|
||||
|
||||
def feed_hotspot_to_engine(topic_engine) -> list[dict]:
|
||||
"""将缓存的热点分析结果注入 TopicEngine,返回热点加权推荐列表"""
|
||||
data = get_last_analysis()
|
||||
return topic_engine.recommend_topics(hotspot_data=data)
|
||||
|
||||
|
||||
# ==================================================
|
||||
# Tab 3: 评论管家
|
||||
# ==================================================
|
||||
|
||||
# ---- 共用: 笔记列表缓存(线程安全)----
|
||||
|
||||
# 主动评论缓存
|
||||
_cached_proactive_entries: list[dict] = []
|
||||
# 我的笔记评论缓存
|
||||
_cached_my_note_entries: list[dict] = []
|
||||
# 缓存互斥锁,防止并发回调产生竞态
|
||||
_cache_lock = threading.RLock()
|
||||
|
||||
|
||||
def _set_cache(name: str, entries: list):
|
||||
"""线程安全地更新笔记列表缓存"""
|
||||
|
||||
@@ -717,6 +717,11 @@ class LLMService:
|
||||
except json.JSONDecodeError:
|
||||
pass
|
||||
|
||||
# 策略6: 修复截断的 JSON(LLM 输出被 token 限制截断)
|
||||
truncated = self._try_fix_truncated_json(cleaned)
|
||||
if truncated is not None:
|
||||
return truncated
|
||||
|
||||
# 全部失败,打日志并抛出有用的错误信息
|
||||
preview = raw[:500] if len(raw) > 500 else raw
|
||||
logger.error("JSON 解析全部失败,LLM 原始返回: %s", preview)
|
||||
@@ -726,6 +731,71 @@ class LLMService:
|
||||
f"💡 可能原因: 模型不支持 JSON 输出格式,建议更换模型重试"
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def _try_fix_truncated_json(text: str) -> dict | None:
|
||||
"""
|
||||
尝试修复被 token 限制截断的 JSON。
|
||||
|
||||
常见场景:LLM 输出的 content 字段非常长,JSON 在字符串中间被切断,
|
||||
导致缺少闭合引号和大括号。
|
||||
|
||||
策略:从 '{' 开始,逐步尝试在不同位置截断并补全 JSON。
|
||||
"""
|
||||
# 找到 JSON 起始
|
||||
start = text.find('{')
|
||||
if start < 0:
|
||||
return None
|
||||
fragment = text[start:]
|
||||
|
||||
# 快速检查:如果已完整则无需修复
|
||||
try:
|
||||
return json.loads(fragment)
|
||||
except json.JSONDecodeError:
|
||||
pass
|
||||
|
||||
# 从末尾向前找到最后一个完整的 key-value 对的结束位置
|
||||
# 策略A: 尝试直接补全闭合字符
|
||||
for suffix in [
|
||||
'"}', # 被截断的字符串值 + 闭合对象
|
||||
'"]}', # 被截断的数组中字符串 + 闭合数组 + 闭合对象
|
||||
'"}',
|
||||
'" }',
|
||||
'..."}', # 在截断处加省略号
|
||||
'..."\n}',
|
||||
]:
|
||||
try:
|
||||
result = json.loads(fragment + suffix)
|
||||
logger.info("截断 JSON 修复成功 (补全: %s)", repr(suffix))
|
||||
return result
|
||||
except json.JSONDecodeError:
|
||||
continue
|
||||
|
||||
# 策略B: 回退到最后一个完整字段
|
||||
# 找到所有 "key": "value" 或 "key": [...] 的匹配位置
|
||||
# 从后往前尝试在每个逗号处截断
|
||||
for i in range(len(fragment) - 1, max(0, len(fragment) - 2000), -1):
|
||||
if fragment[i] in (',', '\n'):
|
||||
candidate = fragment[:i].rstrip().rstrip(',')
|
||||
# 计算需要补全的括号
|
||||
open_braces = candidate.count('{') - candidate.count('}')
|
||||
open_brackets = candidate.count('[') - candidate.count(']')
|
||||
# 检查是否在字符串内部(简单启发式:奇数个未转义引号)
|
||||
in_string = (candidate.count('"') - candidate.count('\\"')) % 2 == 1
|
||||
closing = ''
|
||||
if in_string:
|
||||
closing += '"'
|
||||
closing += ']' * max(0, open_brackets)
|
||||
closing += '}' * max(0, open_braces)
|
||||
if closing:
|
||||
try:
|
||||
result = json.loads(candidate + closing)
|
||||
logger.info("截断 JSON 修复成功 (回退到位置 %d, 补全: %s)", i, repr(closing))
|
||||
return result
|
||||
except json.JSONDecodeError:
|
||||
continue
|
||||
|
||||
return None
|
||||
|
||||
# ---------- 业务方法 ----------
|
||||
|
||||
def get_models(self) -> list[str]:
|
||||
|
||||
@@ -489,6 +489,151 @@ class PublishQueue:
|
||||
|
||||
return "\n".join(lines)
|
||||
|
||||
# ---------- 智能排期引擎 ----------
|
||||
|
||||
def get_slot_usage(self, days: int = 7) -> dict:
|
||||
"""查询未来 N 天各日期各时段已排期的数量。
|
||||
|
||||
返回: {"2026-02-28": {"18-21时": 1, "08-11时": 2}, ...}
|
||||
"""
|
||||
conn = self._get_conn()
|
||||
try:
|
||||
now = datetime.now()
|
||||
cutoff = (now + timedelta(days=days)).strftime("%Y-%m-%d %H:%M:%S")
|
||||
rows = conn.execute(
|
||||
"SELECT scheduled_time FROM queue "
|
||||
"WHERE status IN (?, ?) AND scheduled_time IS NOT NULL AND scheduled_time >= ? AND scheduled_time <= ?",
|
||||
(STATUS_SCHEDULED, STATUS_APPROVED, now.strftime("%Y-%m-%d %H:%M:%S"), cutoff),
|
||||
).fetchall()
|
||||
|
||||
usage: dict[str, dict[str, int]] = {}
|
||||
for row in rows:
|
||||
st = row["scheduled_time"]
|
||||
if not st:
|
||||
continue
|
||||
try:
|
||||
dt = datetime.strptime(st[:19], "%Y-%m-%d %H:%M:%S")
|
||||
except (ValueError, TypeError):
|
||||
continue
|
||||
date_key = dt.strftime("%Y-%m-%d")
|
||||
hour = dt.hour
|
||||
# 映射到3小时段
|
||||
slot = self._hour_to_slot(hour)
|
||||
usage.setdefault(date_key, {})
|
||||
usage[date_key][slot] = usage[date_key].get(slot, 0) + 1
|
||||
return usage
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
@staticmethod
|
||||
def _hour_to_slot(hour: int) -> str:
|
||||
"""将小时映射到时段标签。"""
|
||||
brackets = [
|
||||
(0, 3, "00-03时"), (3, 6, "03-06时"), (6, 8, "06-08时"),
|
||||
(8, 11, "08-11时"), (11, 12, "11-12时"), (12, 14, "12-14时"),
|
||||
(14, 18, "14-18时"), (18, 21, "18-21时"), (21, 24, "21-24时"),
|
||||
]
|
||||
for lo, hi, label in brackets:
|
||||
if lo <= hour < hi:
|
||||
return label
|
||||
return "21-24时"
|
||||
|
||||
@staticmethod
|
||||
def _slot_to_hour_range(slot: str) -> tuple[int, int]:
|
||||
"""从时段标签提取起止小时 (start, end)。"""
|
||||
import re as _re
|
||||
m = _re.match(r"(\d{2})-(\d{2})时", slot)
|
||||
if m:
|
||||
return int(m.group(1)), int(m.group(2))
|
||||
return 18, 21 # fallback
|
||||
|
||||
def suggest_schedule_time(self, analytics, max_per_slot: int = 2,
|
||||
max_per_day: int = 5) -> str | None:
|
||||
"""基于时段权重和已有排期,计算最优发布时间。
|
||||
|
||||
返回格式: '%Y-%m-%d %H:%M:%S',所有时段满时返回 None。
|
||||
"""
|
||||
import random as _random
|
||||
|
||||
time_weights = analytics.get_time_weights()
|
||||
if not time_weights:
|
||||
return None
|
||||
|
||||
# 按权重降序排列候选时段
|
||||
sorted_slots = sorted(time_weights.items(),
|
||||
key=lambda x: x[1] if isinstance(x[1], (int, float)) else x[1].get("weight", 0),
|
||||
reverse=True)
|
||||
|
||||
usage = self.get_slot_usage(days=7)
|
||||
now = datetime.now()
|
||||
|
||||
for day_offset in range(8): # 今天 + 未来7天
|
||||
target_date = now + timedelta(days=day_offset)
|
||||
date_key = target_date.strftime("%Y-%m-%d")
|
||||
|
||||
# 检查当天总量
|
||||
day_usage = usage.get(date_key, {})
|
||||
day_total = sum(day_usage.values())
|
||||
if day_total >= max_per_day:
|
||||
continue
|
||||
|
||||
for slot_name, slot_info in sorted_slots:
|
||||
slot_count = day_usage.get(slot_name, 0)
|
||||
if slot_count >= max_per_slot:
|
||||
continue
|
||||
|
||||
start_hour, end_hour = self._slot_to_hour_range(slot_name)
|
||||
|
||||
# 如果是今天,跳过已过去的时段
|
||||
if day_offset == 0 and end_hour <= now.hour:
|
||||
continue
|
||||
# 如果是今天且时段正在进行中,起始小时调整为当前 +1
|
||||
effective_start = start_hour
|
||||
if day_offset == 0 and start_hour <= now.hour < end_hour:
|
||||
effective_start = now.hour + 1
|
||||
if effective_start >= end_hour:
|
||||
continue
|
||||
|
||||
# 在时段内随机选一个时间
|
||||
rand_hour = _random.randint(effective_start, end_hour - 1)
|
||||
rand_minute = _random.randint(0, 59)
|
||||
scheduled = target_date.replace(hour=rand_hour, minute=rand_minute,
|
||||
second=0, microsecond=0)
|
||||
return scheduled.strftime("%Y-%m-%d %H:%M:%S")
|
||||
|
||||
return None
|
||||
|
||||
def auto_schedule_item(self, item_id: int, analytics,
|
||||
max_per_slot: int = 2, max_per_day: int = 5) -> bool:
|
||||
"""为指定队列项自动分配排期时间。
|
||||
|
||||
成功返回 True(状态变为 scheduled),无可用时段返回 False。
|
||||
"""
|
||||
item = self.get(item_id)
|
||||
if not item or item["status"] not in (STATUS_DRAFT, STATUS_APPROVED):
|
||||
return False
|
||||
|
||||
scheduled_time = self.suggest_schedule_time(
|
||||
analytics, max_per_slot=max_per_slot, max_per_day=max_per_day,
|
||||
)
|
||||
if not scheduled_time:
|
||||
logger.warning("auto_schedule_item #%d: 未来7天无可用时段", item_id)
|
||||
return False
|
||||
|
||||
# 更新排期时间 + 状态
|
||||
conn = self._get_conn()
|
||||
try:
|
||||
now_str = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
||||
conn.execute(
|
||||
"UPDATE queue SET status = ?, scheduled_time = ?, updated_at = ? WHERE id = ?",
|
||||
(STATUS_SCHEDULED, scheduled_time, now_str, item_id),
|
||||
)
|
||||
conn.commit()
|
||||
logger.info("auto_schedule_item #%d → %s", item_id, scheduled_time)
|
||||
return True
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
|
||||
class QueuePublisher:
|
||||
"""后台队列发布处理器"""
|
||||
|
||||
+33
-7
@@ -3,15 +3,21 @@ services/queue_ops.py
|
||||
发布队列操作:生成入队、状态管理、发布控制
|
||||
"""
|
||||
import os
|
||||
import re
|
||||
import time
|
||||
import random
|
||||
import logging
|
||||
|
||||
from PIL import Image
|
||||
|
||||
from .config_manager import ConfigManager, OUTPUT_DIR
|
||||
from .publish_queue import (
|
||||
PublishQueue, QueuePublisher,
|
||||
STATUS_DRAFT, STATUS_APPROVED, STATUS_SCHEDULED, STATUS_PUBLISHING,
|
||||
STATUS_PUBLISHED, STATUS_FAILED, STATUS_REJECTED, STATUS_LABELS,
|
||||
)
|
||||
from .llm_service import LLMService
|
||||
from .sd_service import SDService
|
||||
from .mcp_client import get_mcp_client
|
||||
from .connection import _get_llm_config
|
||||
from .persona import DEFAULT_TOPICS, DEFAULT_STYLES, _resolve_persona
|
||||
@@ -52,8 +58,13 @@ def _log(msg: str):
|
||||
|
||||
def generate_to_queue(topics_str, sd_url_val, sd_model_name, model, persona_text=None,
|
||||
quality_mode_val=None, face_swap_on=False, count=1,
|
||||
scheduled_time=None):
|
||||
"""批量生成内容 → 加入发布队列(不直接发布)"""
|
||||
scheduled_time=None, auto_schedule=False, auto_approve=False):
|
||||
"""批量生成内容 → 加入发布队列(不直接发布)
|
||||
|
||||
Args:
|
||||
auto_schedule: 为每篇内容自动分配最优排期时间
|
||||
auto_approve: 入队后自动审核通过
|
||||
"""
|
||||
try:
|
||||
topics = [t.strip() for t in topics_str.split(",") if t.strip()] if topics_str else DEFAULT_TOPICS
|
||||
use_weights = cfg.get("use_smart_weights", True) and _analytics.has_weights
|
||||
@@ -82,10 +93,10 @@ def generate_to_queue(topics_str, sd_url_val, sd_model_name, model, persona_text
|
||||
persona = _resolve_persona(persona_text) if persona_text else None
|
||||
|
||||
if use_weights:
|
||||
weight_insights = f"高权重主题: {', '.join(list(analytics._weights.get('topic_weights', {}).keys())[:5])}\n"
|
||||
weight_insights += f"权重摘要: {analytics.weights_summary}"
|
||||
weight_insights = f"高权重主题: {', '.join(list(_analytics._weights.get('topic_weights', {}).keys())[:5])}\n"
|
||||
weight_insights += f"权重摘要: {_analytics.weights_summary}"
|
||||
title_advice = _analytics.get_title_advice()
|
||||
hot_tags = ", ".join(analytics.get_top_tags(8))
|
||||
hot_tags = ", ".join(_analytics.get_top_tags(8))
|
||||
try:
|
||||
data = svc.generate_weighted_copy(topic, style, weight_insights, title_advice, hot_tags, sd_model_name=sd_model_name, persona=persona)
|
||||
except Exception:
|
||||
@@ -150,7 +161,21 @@ def generate_to_queue(topics_str, sd_url_val, sd_model_name, model, persona_text
|
||||
topic=topic, style=style, persona=persona or "",
|
||||
status=STATUS_DRAFT, scheduled_time=scheduled_time,
|
||||
)
|
||||
results.append(f"#{item_id} {title}")
|
||||
# 自动排期
|
||||
sched_msg = ""
|
||||
if auto_schedule and _analytics:
|
||||
ok = _pub_queue.auto_schedule_item(item_id, _analytics)
|
||||
if ok:
|
||||
item = _pub_queue.get(item_id)
|
||||
sched_msg = f" ⏰{item['scheduled_time'][:16]}" if item else ""
|
||||
_log(f"🕐 #{item_id} 自动排期{sched_msg}")
|
||||
|
||||
# 自动审核
|
||||
if auto_approve:
|
||||
_pub_queue.approve(item_id)
|
||||
_log(f"✅ #{item_id} 自动审核通过")
|
||||
|
||||
results.append(f"#{item_id} {title}{sched_msg}")
|
||||
_log(f"📋 已加入队列 #{item_id}: {title}")
|
||||
|
||||
# 多篇间隔
|
||||
@@ -335,13 +360,14 @@ def queue_batch_approve(status_filter):
|
||||
|
||||
def queue_generate_and_refresh(topics_str, sd_url_val, sd_model_name, model,
|
||||
persona_text, quality_mode_val, face_swap_on,
|
||||
gen_count, gen_schedule_time):
|
||||
gen_count, gen_schedule_time, auto_schedule=False):
|
||||
"""生成内容到队列 + 刷新表格"""
|
||||
msg = generate_to_queue(
|
||||
topics_str, sd_url_val, sd_model_name, model,
|
||||
persona_text=persona_text, quality_mode_val=quality_mode_val,
|
||||
face_swap_on=face_swap_on, count=gen_count,
|
||||
scheduled_time=gen_schedule_time.strip() if gen_schedule_time else None,
|
||||
auto_schedule=auto_schedule,
|
||||
)
|
||||
table = _pub_queue.format_queue_table()
|
||||
calendar = _pub_queue.format_calendar(14)
|
||||
|
||||
+112
-130
@@ -514,143 +514,38 @@ def auto_reply_once(max_replies, mcp_url, model, persona_text):
|
||||
|
||||
|
||||
def auto_publish_once(topics_str, mcp_url, sd_url_val, sd_model_name, model, persona_text=None, quality_mode_val=None, face_swap_on=False):
|
||||
"""一键发布:自动生成文案 → 生成图片 → 本地备份 → 发布到小红书(含限额 + 智能权重 + 人设 + 画质)"""
|
||||
"""一键发布:生成内容 → 加入发布队列(自动排期 + 自动审核)。
|
||||
|
||||
实际发布由 QueuePublisher 后台处理器完成。
|
||||
"""
|
||||
try:
|
||||
if _is_in_cooldown():
|
||||
return "⏳ 错误冷却中,请稍后再试"
|
||||
if not _check_daily_limit("publishes"):
|
||||
return f"🚫 今日发布已达上限 ({DAILY_LIMITS['publishes']})"
|
||||
|
||||
topics = [t.strip() for t in topics_str.split(",") if t.strip()] if topics_str else DEFAULT_TOPICS
|
||||
use_weights = cfg.get("use_smart_weights", True) and _analytics.has_weights
|
||||
# 延迟导入避免循环依赖
|
||||
from .queue_ops import generate_to_queue
|
||||
|
||||
if use_weights:
|
||||
# 智能加权选题
|
||||
topic = _analytics.get_weighted_topic(topics)
|
||||
style = _analytics.get_weighted_style(DEFAULT_STYLES)
|
||||
_auto_log_append(f"🧠 [智能] 主题: {topic} | 风格: {style} (加权选择)")
|
||||
else:
|
||||
topic = random.choice(topics)
|
||||
style = random.choice(DEFAULT_STYLES)
|
||||
_auto_log_append(f"📝 主题: {topic} | 风格: {style} (主题池: {len(topics)} 个)")
|
||||
|
||||
# 生成文案
|
||||
api_key, base_url, _ = _get_llm_config()
|
||||
if not api_key:
|
||||
return "❌ LLM 未配置,请先在全局设置中配置提供商"
|
||||
|
||||
svc = LLMService(api_key, base_url, model)
|
||||
# 解析人设(随机/指定)
|
||||
persona = _resolve_persona(persona_text) if persona_text else None
|
||||
if persona:
|
||||
_auto_log_append(f"🎭 人设: {persona[:20]}...")
|
||||
|
||||
if use_weights:
|
||||
# 使用加权文案生成 (携带权重洞察)
|
||||
weight_insights = f"高权重主题: {', '.join(list(analytics._weights.get('topic_weights', {}).keys())[:5])}\n"
|
||||
weight_insights += f"权重摘要: {analytics.weights_summary}"
|
||||
title_advice = _analytics.get_title_advice()
|
||||
hot_tags = ", ".join(analytics.get_top_tags(8))
|
||||
try:
|
||||
data = svc.generate_weighted_copy(topic, style, weight_insights, title_advice, hot_tags, sd_model_name=sd_model_name, persona=persona)
|
||||
_auto_log_append("🧠 使用智能加权文案模板")
|
||||
except Exception as e:
|
||||
logger.warning("加权文案生成失败, 退回普通模式: %s", e)
|
||||
data = svc.generate_copy(topic, style, sd_model_name=sd_model_name, persona=persona)
|
||||
_auto_log_append("⚠️ 加权模板异常, 使用普通模板")
|
||||
else:
|
||||
data = svc.generate_copy(topic, style, sd_model_name=sd_model_name, persona=persona)
|
||||
|
||||
title = (data.get("title", "") or "")[:20]
|
||||
content = data.get("content", "")
|
||||
sd_prompt = data.get("sd_prompt", "")
|
||||
tags = data.get("tags", [])
|
||||
|
||||
# 如果有高权重标签,补充到 tags 中
|
||||
if use_weights:
|
||||
top_tags = _analytics.get_top_tags(5)
|
||||
for t in top_tags:
|
||||
if t not in tags:
|
||||
tags.append(t)
|
||||
tags = tags[:10] # 限制最多10个标签
|
||||
|
||||
if not title:
|
||||
_record_error()
|
||||
return "❌ 文案生成失败:无标题"
|
||||
_auto_log_append(f"📄 文案: {title}")
|
||||
|
||||
# 生成图片
|
||||
if not sd_url_val or not sd_model_name:
|
||||
return "❌ SD WebUI 未连接或未选择模型,请先在全局设置中连接"
|
||||
|
||||
sd_svc = SDService(sd_url_val)
|
||||
# 自动发布也支持换脸
|
||||
face_image = None
|
||||
if face_swap_on:
|
||||
face_image = SDService.load_face_image()
|
||||
if face_image:
|
||||
_auto_log_append("🎭 换脸已启用")
|
||||
else:
|
||||
_auto_log_append("⚠️ 换脸已启用但未找到头像,跳过换脸")
|
||||
images = sd_svc.txt2img(prompt=sd_prompt, model=sd_model_name,
|
||||
face_image=face_image,
|
||||
quality_mode=quality_mode_val or "快速 (约30秒)",
|
||||
persona=persona)
|
||||
if not images:
|
||||
_record_error()
|
||||
return "❌ 图片生成失败:没有返回图片"
|
||||
_auto_log_append(f"🎨 已生成 {len(images)} 张图片")
|
||||
|
||||
# 本地备份(同时用于发布)
|
||||
ts = int(time.time())
|
||||
safe_title = re.sub(r'[\\/*?:"<>|]', "", title)[:20]
|
||||
backup_dir = os.path.join(OUTPUT_DIR, f"{ts}_{safe_title}")
|
||||
os.makedirs(backup_dir, exist_ok=True)
|
||||
|
||||
# 保存文案
|
||||
with open(os.path.join(backup_dir, "文案.txt"), "w", encoding="utf-8") as f:
|
||||
f.write(f"标题: {title}\n风格: {style}\n主题: {topic}\n\n{content}\n\n标签: {', '.join(tags)}\n\nSD Prompt: {sd_prompt}")
|
||||
|
||||
image_paths = []
|
||||
for idx, img in enumerate(images):
|
||||
if isinstance(img, Image.Image):
|
||||
path = os.path.abspath(os.path.join(backup_dir, f"图{idx+1}.jpg"))
|
||||
if img.mode != "RGB":
|
||||
img = img.convert("RGB")
|
||||
img.save(path, format="JPEG", quality=95)
|
||||
image_paths.append(path)
|
||||
|
||||
if not image_paths:
|
||||
return "❌ 图片保存失败"
|
||||
|
||||
_auto_log_append(f"💾 本地已备份至: {backup_dir}")
|
||||
|
||||
# 发布到小红书
|
||||
client = get_mcp_client(mcp_url)
|
||||
result = client.publish_content(
|
||||
title=title, content=content, images=image_paths, tags=tags
|
||||
topics = topics_str if topics_str else ",".join(DEFAULT_TOPICS)
|
||||
msg = generate_to_queue(
|
||||
topics, sd_url_val, sd_model_name, model,
|
||||
persona_text=persona_text, quality_mode_val=quality_mode_val,
|
||||
face_swap_on=face_swap_on, count=1,
|
||||
auto_schedule=True, auto_approve=True,
|
||||
)
|
||||
if "error" in result:
|
||||
_record_error()
|
||||
_auto_log_append(f"❌ 发布失败: {result['error']} (文案已本地保存)")
|
||||
return f"❌ 发布失败: {result['error']}\n💾 文案和图片已备份至: {backup_dir}"
|
||||
_auto_log_append(f"📋 内容已入队: {msg}")
|
||||
|
||||
_increment_stat("publishes")
|
||||
_clear_error_streak()
|
||||
|
||||
# 清理 _temp_publish 中的旧临时文件
|
||||
temp_dir = os.path.join(OUTPUT_DIR, "_temp_publish")
|
||||
# 检查 QueuePublisher 是否在运行
|
||||
try:
|
||||
if os.path.exists(temp_dir):
|
||||
for f in os.listdir(temp_dir):
|
||||
fp = os.path.join(temp_dir, f)
|
||||
if os.path.isfile(fp) and time.time() - os.path.getmtime(fp) > 3600:
|
||||
os.remove(fp)
|
||||
except Exception:
|
||||
from .queue_ops import _queue_publisher
|
||||
if _queue_publisher and not _queue_publisher.is_running:
|
||||
_auto_log_append("⚠️ 队列处理器未启动,内容已入队但需启动处理器以自动发布")
|
||||
return msg + "\n⚠️ 请启动队列处理器以自动发布"
|
||||
except ImportError:
|
||||
pass
|
||||
|
||||
_auto_log_append(f"🚀 发布成功: {title} (今日第{_daily_stats['publishes']}篇)")
|
||||
return f"✅ 发布成功!\n📌 标题: {title}\n💾 备份: {backup_dir}\n📊 今日发布: {_daily_stats['publishes']}/{DAILY_LIMITS['publishes']}\n{result.get('text', '')}"
|
||||
return msg
|
||||
|
||||
except Exception as e:
|
||||
_record_error()
|
||||
@@ -760,13 +655,17 @@ def _scheduler_loop(comment_enabled, publish_enabled, reply_enabled, like_enable
|
||||
_auto_log_append(f"⏰ 下次收藏: {interval // 60} 分钟后")
|
||||
_update_next_display()
|
||||
|
||||
# 自动发布
|
||||
# 自动发布(通过队列)
|
||||
if publish_enabled and now >= next_publish:
|
||||
try:
|
||||
_auto_log_append("--- 🔄 执行自动发布 ---")
|
||||
msg = auto_publish_once(topics, mcp_url, sd_url_val, sd_model_name, model,
|
||||
persona_text=persona_text, quality_mode_val=quality_mode_val,
|
||||
face_swap_on=face_swap_on)
|
||||
_auto_log_append("--- 🔄 执行自动发布(队列模式) ---")
|
||||
from .queue_ops import generate_to_queue
|
||||
msg = generate_to_queue(
|
||||
topics, sd_url_val, sd_model_name, model,
|
||||
persona_text=persona_text, quality_mode_val=quality_mode_val,
|
||||
face_swap_on=face_swap_on, count=1,
|
||||
auto_schedule=True, auto_approve=True,
|
||||
)
|
||||
_auto_log_append(msg)
|
||||
except Exception as e:
|
||||
_auto_log_append(f"❌ 自动发布异常: {e}")
|
||||
@@ -1065,6 +964,89 @@ def stop_learn_scheduler():
|
||||
return "🛑 定时学习已停止"
|
||||
|
||||
|
||||
# ==================================================
|
||||
# 热点自动采集
|
||||
# ==================================================
|
||||
|
||||
_hotspot_collector_running = threading.Event()
|
||||
_hotspot_collector_thread: threading.Thread | None = None
|
||||
|
||||
|
||||
def _hotspot_collector_loop(keywords: list[str], interval_hours: float, mcp_url: str, model: str):
|
||||
"""热点自动采集后台循环:遍历 keywords → 搜索 → LLM 分析 → 写入状态缓存"""
|
||||
from .hotspot import search_hotspots, analyze_and_suggest
|
||||
|
||||
logger.info("热点自动采集已启动, 关键词=%s, 间隔=%s小时", keywords, interval_hours)
|
||||
_auto_log_append(f"🔥 热点自动采集已启动, 每 {interval_hours} 小时采集一次, 关键词: {', '.join(keywords)}")
|
||||
|
||||
while _hotspot_collector_running.is_set():
|
||||
for kw in keywords:
|
||||
if not _hotspot_collector_running.is_set():
|
||||
break
|
||||
try:
|
||||
_auto_log_append(f"🔥 自动采集热点: 搜索「{kw}」...")
|
||||
status, search_result = search_hotspots(kw, "最多点赞", mcp_url)
|
||||
if "❌" in status or not search_result:
|
||||
_auto_log_append(f"⚠️ 热点搜索失败: {status}")
|
||||
continue
|
||||
|
||||
_auto_log_append(f"🔥 自动采集热点: AI 分析「{kw}」...")
|
||||
a_status, _, _, _ = analyze_and_suggest(model, kw, search_result)
|
||||
_auto_log_append(f"🔥 热点采集「{kw}」: {a_status}")
|
||||
|
||||
except Exception as e:
|
||||
_auto_log_append(f"⚠️ 热点采集「{kw}」异常: {e}")
|
||||
|
||||
# 关键词间间隔,避免过快请求
|
||||
for _ in range(30):
|
||||
if not _hotspot_collector_running.is_set():
|
||||
break
|
||||
time.sleep(1)
|
||||
|
||||
# 等待下一轮
|
||||
wait_seconds = int(interval_hours * 3600)
|
||||
_auto_log_append(f"🔥 热点采集完成一轮, {interval_hours}小时后再次采集")
|
||||
for _ in range(int(wait_seconds / 5)):
|
||||
if not _hotspot_collector_running.is_set():
|
||||
break
|
||||
time.sleep(5)
|
||||
|
||||
logger.info("热点自动采集已停止")
|
||||
_auto_log_append("🔥 热点自动采集已停止")
|
||||
|
||||
|
||||
def start_hotspot_collector(keywords_str: str, interval_hours: float, mcp_url: str, model: str):
|
||||
"""启动热点自动采集"""
|
||||
global _hotspot_collector_thread
|
||||
if _hotspot_collector_running.is_set():
|
||||
return "⚠️ 热点自动采集已在运行中"
|
||||
|
||||
keywords = [k.strip() for k in keywords_str.split(",") if k.strip()]
|
||||
if not keywords:
|
||||
return "❌ 请输入至少一个采集关键词"
|
||||
|
||||
api_key, _, _ = _get_llm_config()
|
||||
if not api_key:
|
||||
return "❌ LLM 未配置,请先在全局设置中配置提供商"
|
||||
|
||||
_hotspot_collector_running.set()
|
||||
_hotspot_collector_thread = threading.Thread(
|
||||
target=_hotspot_collector_loop,
|
||||
args=(keywords, interval_hours, mcp_url, model),
|
||||
daemon=True,
|
||||
)
|
||||
_hotspot_collector_thread.start()
|
||||
return f"✅ 热点自动采集已启动 🔥 每 {int(interval_hours)} 小时采集一次, 关键词: {', '.join(keywords)}"
|
||||
|
||||
|
||||
def stop_hotspot_collector():
|
||||
"""停止热点自动采集"""
|
||||
if not _hotspot_collector_running.is_set():
|
||||
return "⚠️ 热点自动采集未在运行"
|
||||
_hotspot_collector_running.clear()
|
||||
return "🛑 热点自动采集已停止"
|
||||
|
||||
|
||||
# ==================================================
|
||||
# Windows 开机自启管理
|
||||
# ==================================================
|
||||
|
||||
Reference in New Issue
Block a user