| name | forecast-runtime |
| description | 基于校正后状态和历史经验的滚动舆情预测引擎。编排 EvoSim 轻量仿真循环, 为 Leader Actor 注入因果经验辅助 LLM 行为决策,通过可选干预 Agent 实时稳定仿真轨迹, 最终由 Report Agent 输出结构化 JSON + 人类可读 Markdown 双格式预测报告。 当需要执行以下任务时使用此 Skill: (1) 在校正后的仿真快照上继续向未来推演 (2) 输出立场分布、极化、关键人物活跃度、传播结构的变化趋势 (3) 基于历史因果经验(原因导致结果)辅助预测并稳定仿真 (4) 生成供人类阅读的预测报告和供下游 Skill 消费的结构化数据 此 Skill 为 EvoPalantir 闭环中的 forecast_runtime 模块, 上游为 calibration_engine 输出的校正状态,下游为 calibration_memory_engine 的经验沉淀。
|
| compatibility | Requires Python 3.10+, EvoSim project, OpenAI-compatible LLM API, and optionally a vector database (Qdrant/Milvus/pgvector). |
| metadata | {"author":"MaYiding","version":"1.0"} |
Forecast Runtime
Overview
在校正后的仿真状态上执行滚动预测。
核心架构决策: 本 Skill 使用外部逐 tick 编排模式,不调用 EvoSim 的 sim.run()。
原因: 需要在每个 tick 的 Leader 行为决策前注入经验到 prompt,这要求拦截 tick 内部步骤,
sim.run() 作为黑盒无法实现此拦截。因此本 Skill 直接调用 EvoSim 的底层原语
(AgentUser.get_feed(), AgentUser.react_to_feed(), FeedPipeline, SnapshotManager 等)。
数据流:
校正后 Snapshot + ActorStates
-> 恢复 EvoSim DB + 初始化组件
-> 外部 tick loop:
获取经验 -> Leader: 经验增强 prompt + LLM 决策
Crowd: 规则决策 -> 推荐器分发 -> 状态更新
-> 指标计算 -> [干预检查] -> [快照]
-> Report Agent (ReACT) 生成报告
-> 输出: JSON + Markdown + DB
Input Parameters
调用时传入 ForecastInput (完整 schema 见 data-schema.md):
| 参数 | 必填 | 默认值 | 说明 |
|---|
event_id | 是 | - | 目标事件 ID |
db_url | 是 | - | 数据库连接串, 如 sqlite:///path/to/sim.db |
project_root | 是 | - | EvoSim 项目根目录, SnapshotManager 需要 |
forecast_horizon_ticks | 是 | - | 预测跨度 (tick 数) |
tick_interval_minutes | 是 | - | 每 tick 代表的真实分钟数 |
llm_config | 是 | - | LLM API 配置 (leader/intervention/report 三模型可分开指定) |
output_dir | 是 | - | 输出目录路径 |
parent_session_id | 否 | 自动取最新 | EvoSim 快照 session ID |
vector_db_config | 否 | None | 向量 DB 配置 |
intervention_enabled | 否 | true | 是否启用干预 Agent |
intervention_check_interval_ticks | 否 | 3 | 干预检查间隔 |
snapshot_interval_ticks | 否 | 5 | 快照保存间隔 |
experience_refresh_strategy | 否 | auto | auto/fixed/on_demand |
experience_cache_ttl_ticks | 否 | 5 | 经验缓存基础 TTL |
recommender_config | 否 | 见 Config 节 | 推荐权重 |
output_to_db | 否 | true | 是否写入数据库 |
Orchestration Workflow
Phase 1: 初始化
import json, os, sys, uuid, random, asyncio
from datetime import datetime
from openai import OpenAI
from scripts.state_loader import StateLoader
from scripts.experience_retriever import (
ExperienceRetriever, VectorDBConfig, SQLConfig,
build_experience_enhanced_prompt, query_user_actions,
query_user_actions_max_id, _ACTION_MAP,
)
from scripts.experience_cache import ExperienceCache, TickMetricsSnapshot
from scripts.metrics_calculator import MetricsCalculator
from scripts.forecast_writer import ForecastWriter, build_summary_from_timeline
evosim_root = forecast_input["project_root"]
evosim_src = os.path.join(evosim_root, "src")
for p in [evosim_root, evosim_src]:
if p not in sys.path:
sys.path.insert(0, p)
loader = StateLoader(db_url=forecast_input["db_url"])
state = loader.load_all(forecast_input["event_id"])
if len(state.actor_states) == 0:
raise ValueError("No actor states found. Ensure calibration_engine has run.")
event = state.event
start_tick = state.start_tick
leader_actors = [a for a in state.actor_states if a.get("role") == "leader"]
crowd_actors = [a for a in state.actor_states if a.get("role") == "crowd"]
import control_flags
control_flags.attack_enabled = False
control_flags.aftercare_enabled = False
control_flags.moderation_enabled = False
from snapshot_manager import create_snapshot_manager
db_path = forecast_input["db_url"].replace("sqlite:///", "")
snapshot_mgr = create_snapshot_manager(forecast_input["project_root"], db_path)
session_id = forecast_input.get("parent_session_id")
if session_id is None:
sessions = snapshot_mgr.list_sessions()
session_id = sessions[0]["session_id"] if sessions else None
if session_id:
snapshot_mgr.restore_from_tick(tick=start_tick, session_id=session_id)
forecast_session_id = snapshot_mgr.create_session(
parent_session_id=session_id, parent_tick=start_tick
)
from database_manager import DatabaseManager
db_mgr = DatabaseManager(db_path, reset_db=False, use_service=False)
from database.database_manager import get_db_manager as get_singleton_db
_singleton = get_singleton_db()
_singleton.set_database_path(db_path)
from user_manager import UserManager
user_mgr = UserManager(
config={"num_users": len(state.actor_states), "temperature": 0.8,
"agent_config_path": "separate"},
db_manager=db_mgr,
restore_existing=True
)
llm_cfg = forecast_input["llm_config"]
api_key_env = llm_cfg.get("api_key_env", "LLM_API_KEY")
openai_client = OpenAI(
api_key=os.environ.get(api_key_env, ""),
base_url=os.environ.get("LLM_API_BASE", llm_cfg.get("api_base", ""))
)
leader_engine = llm_cfg.get("leader_model", "gpt-4")
retriever = ExperienceRetriever(
vector_config=VectorDBConfig(**forecast_input["vector_db_config"])
if forecast_input.get("vector_db_config") else None,
sql_config=SQLConfig(db_url=forecast_input["db_url"])
)
cache = ExperienceCache(
retriever, base_ttl_ticks=forecast_input.get("experience_cache_ttl_ticks", 5)
)
calc = MetricsCalculator()
writer = ForecastWriter(
output_dir=forecast_input["output_dir"],
db_url=forecast_input["db_url"] if forecast_input.get("output_to_db", True) else None
)
import sqlite3
_conn = sqlite3.connect(db_path)
_conn.execute("""CREATE TABLE IF NOT EXISTS actor_states (
actor_id TEXT PRIMARY KEY, event_id TEXT NOT NULL, role TEXT,
stance_posterior TEXT, activity_hazard REAL DEFAULT 0.5,
bridge_score REAL DEFAULT 0.0, topic_mix TEXT,
uncertainty REAL DEFAULT 0.5, last_observed_at TEXT,
state_version INTEGER DEFAULT 0, updated_at_tick INTEGER DEFAULT 0)""")
_conn.execute("""CREATE TABLE IF NOT EXISTS intervention_records (
intervention_id TEXT PRIMARY KEY, event_id TEXT NOT NULL,
tick INTEGER, trigger_reason TEXT, deviation_score REAL,
severity TEXT, adjustments_json TEXT, created_at TEXT)""")
_conn.commit()
_conn.close()
end_tick = start_tick + forecast_input["forecast_horizon_ticks"]
forecast_output = {
"forecast_id": str(uuid.uuid4()),
"event_id": forecast_input["event_id"],
"event_title": event.get("title", forecast_input["event_id"]),
"generated_at": datetime.now().isoformat(),
"forecast_horizon_ticks": forecast_input["forecast_horizon_ticks"],
"tick_interval_minutes": forecast_input["tick_interval_minutes"],
"start_tick": start_tick,
"end_tick": end_tick,
"timeline": [],
"summary": {},
"experience_used": [],
}
intervention_records = []
experience_used_set = {}
prev_phase = event.get("phase", "")
intervention_count = 0
max_interventions = forecast_input["forecast_horizon_ticks"] // 3
Phase 2: Tick 循环
整个 tick 循环包裹在 async 函数中,因为 react_to_feed() 是 async 方法。
外层使用 asyncio.run(run_forecast_loop()) 驱动。
async def run_forecast_loop():
"""async 包装: react_to_feed 是 async 方法, 必须在 async 上下文中调用。"""
for tick in range(start_tick + 1, end_tick + 1):
tick_actions = []
experiences = cache.get_or_fetch(
event=event, current_tick=tick,
actors=leader_actors,
top_k=5
)
for exp in experiences:
if exp.experience_id not in experience_used_set:
experience_used_set[exp.experience_id] = {
"experience_id": exp.experience_id,
"event_type": exp.event_type,
"event_similarity": exp.event_similarity,
"influence_on_forecast": exp.effect_description[:100]
}
for user in user_mgr.users:
actor_state = next((a for a in leader_actors if a["actor_id"] == user.user_id), None)
if actor_state is None:
continue
feed = user.get_feed(experiment_config={}, time_step=tick)
original_persona = user.persona
user.persona = build_experience_enhanced_prompt(
persona=original_persona, feed=feed,
experiences=experiences, actor_state=actor_state
)
_max_id_before = query_user_actions_max_id(user.user_id, db_path)
try:
await user.react_to_feed(openai_client, leader_engine, feed)
except Exception:
pass
user.persona = original_persona
user_tick_actions = query_user_actions(user.user_id, db_path,
after_id=_max_id_before)
tick_actions.extend(user_tick_actions)
for actor_state in crowd_actors:
if random.random() > actor_state.get("activity_hazard", 0.3):
continue
crowd_user = next(
(u for u in user_mgr.users if u.user_id == actor_state["actor_id"]), None
)
crowd_feed = crowd_user.get_feed({}, time_step=tick) if crowd_user else []
if not crowd_feed:
continue
stance = actor_state.get("stance_posterior", {"neutral": 1.0})
dominant_stance = max(stance, key=stance.get)
target = crowd_feed[0] if crowd_feed else None
target_id = target.get("post_id") if isinstance(target, dict) else getattr(target, "post_id", None)
alignment = stance.get("support", 0) - stance.get("oppose", 0)
if alignment > 0.3:
action_type = random.choices(
["comment", "repost", "like"], weights=[0.4, 0.3, 0.3]
)[0]
elif alignment < -0.3:
action_type = random.choices(
["comment", "silence"], weights=[0.3, 0.7]
)[0]
else:
action_type = random.choices(
["like", "silence"], weights=[0.3, 0.7]
)[0]
tick_actions.append({
"actor_id": actor_state["actor_id"],
"action_type": action_type,
"stance_expressed": dominant_stance,
"target_content_id": target_id,
"target_actor_id": None,
})
for actor_list, alpha in [(leader_actors, 0.15), (crowd_actors, 0.08)]:
for actor in actor_list:
actor_actions = [a for a in tick_actions if a["actor_id"] == actor["actor_id"]]
if not actor_actions:
continue
expressed = actor_actions[0].get("stance_expressed", "neutral")
old_stance = actor.get("stance_posterior", {"neutral": 1.0})
new_stance = {}
for k in ["support", "oppose", "neutral", "unclear"]:
one_hot = 1.0 if k == expressed else 0.0
new_stance[k] = (1 - alpha) * old_stance.get(k, 0.0) + alpha * one_hot
total = sum(new_stance.values()) or 1.0
new_stance = {k: v / total for k, v in new_stance.items()}
actor["stance_posterior"] = new_stance
actor["activity_hazard"] = max(0.05, actor.get("activity_hazard", 0.3) * 0.95)
actor["state_version"] = actor.get("state_version", 0) + 1
actor["updated_at_tick"] = tick
all_actors = leader_actors + crowd_actors
metrics = calc.calculate(tick, all_actors, tick_actions)
metrics_snapshot = TickMetricsSnapshot(
tick=tick,
stance_distribution=metrics.stance_distribution,
polarization_index=metrics.polarization_index,
activity_rate=metrics.activity_rate,
has_new_injection=False
)
cache.record_tick_metrics(metrics_snapshot)
current_phase = event.get("phase", "")
if current_phase != prev_phase:
cache.invalidate_on_phase_change()
prev_phase = current_phase
intervention_applied = {"applied": False, "reason": "", "adjustments": {}}
if (forecast_input.get("intervention_enabled", True)
and (tick - start_tick) % forecast_input.get("intervention_check_interval_ticks", 3) == 0
and (tick - start_tick) > 0
and intervention_count < max_interventions):
pass
if tick % forecast_input.get("snapshot_interval_ticks", 5) == 0:
snapshot_mgr.save_tick_snapshot(tick, {
"tick": tick, "timestamp": datetime.now().isoformat(),
"user_count": len(state.actor_states)
})
forecast_output["timeline"].append({
"tick": tick,
"simulated_time": datetime.now().isoformat(),
"metrics": metrics.to_dict(),
"notable_actions": tick_actions[:10],
"intervention_applied": intervention_applied
})
try:
asyncio.run(run_forecast_loop())
except Exception as e:
forecast_output.setdefault("summary", {}).setdefault("risk_signals", []).append(
"simulation_incomplete"
)
action dict 转换说明: EvoSim 的 react_to_feed() 将行为写入 user_actions 表时
使用 action.replace('-', '_') (如 "like_post", "comment_post")。
query_user_actions() 通过 _ACTION_MAP 转换为 MetricsCalculator 统一格式
(如 "like", "comment", "repost", "silence")。
Phase 3: 报告生成
Report Agent 使用 ReACT 模式 (详见 report-guide.md):
report_sections = []
if "stance_trend" not in forecast_output.get("summary", {}):
forecast_output["summary"] = build_summary_from_timeline(
forecast_output["timeline"]
)
forecast_output["summary"].setdefault("risk_signals", []).append(
"summary_from_fallback"
)
experience_used_list = list(experience_used_set.values())
forecast_output["experience_used"] = experience_used_list
Phase 4: 输出
json_path = writer.write_json(forecast_output)
md_path = writer.write_markdown(
forecast_output, report_sections,
experience_used=experience_used_list,
intervention_log=intervention_records
)
if forecast_input.get("output_to_db", True):
writer.write_to_db(forecast_output)
输出文件:
{output_dir}/forecast_output_{event_id}_{timestamp}.json
{output_dir}/forecast_report_{event_id}_{timestamp}.md
Configuration Quick Reference
推荐器权重
ForecastInput.recommender_config 使用 _weight 后缀 key:
similar_feed_weight: 0.5
cross_cutting_weight: 0.2
trending_weight: 0.3
leader_comment_share: 0.7
crowd_comment_share: 0.3
干预 Agent 的 InterventionRecord.adjustments.recommender_weights 使用短 key:
similar_feed, cross_cutting, trending
翻译映射: similar_feed -> similar_feed_weight, cross_cutting -> cross_cutting_weight,
trending -> trending_weight。干预后需将短 key 翻译回 _weight key 更新推荐器配置。
LLM Temperature
Leader Actor: 0.7 # 保留行为多样性
Intervention Agent: 0.3 # 保守/确定
Report Agent: 0.4 # 稳定/可复现
干预阈值
偏差 >= 0.30 或存在趋势异常 -> 触发 LLM 判断 (OR 关系)。
过度干预保护: intervention_count >= forecast_horizon_ticks // 3 时停止。
经验缓存
base_ttl: 5 ticks (高波动减半, 低波动加倍)
立即失效: 注入新内容 / 干预执行 / event.phase 变更 / on_demand 策略
Scripts
References
Error Handling
| 场景 | 处理方式 |
|---|
| actor_states 为空 | Phase 1 抛出 ValueError 终止 |
| EvoSim import 失败 | 捕获 ImportError, 报告缺失模块名 |
| 经验库不可用 | Leader prompt 省略经验 section, 干预仅用硬性指标 |
| LLM 调用失败 | Leader 本 tick 无行动 (try/except), 不中断 |
| tick 循环异常 | 用已收集 timeline 进入 Phase 3, risk_signals 标记 simulation_incomplete |
| Phase 3 summary LLM 失败 | 从 timeline 数据回退构建, risk_signals 标记 summary_from_fallback |
| 过度干预 | 超过 horizon//3 次停止, risk_signals 标记 excessive_intervention |
| 输出写入失败 | 备用路径 {output_dir}/fallback/, 仍失败返回内存结果 |
Environment Variables
LLM_API_KEY # OpenAI 兼容 API key
LLM_API_BASE # API base URL
VECTOR_DB_API_KEY # 向量数据库 API key (如需要)
SQL_DB_URL # 经验库 SQL 连接串