Project Aura — 调度引擎 (Scheduler)
架构概览
调度系统分为三层:心跳层 → 决策调度层 → 动作执行层
用户状态变化 / 定时心跳
↓
┌──────────────┐
│ 心跳调度器 │ ← 自适应泊松过程
│ Heartbeat │
└──────┬───────┘
↓ 触发唤醒
┌──────────────┐
│ 双系统调度器 │ ← 快/慢双通道
│ DualSystem │
├──────────────┤
│ 快系统(ms级) │ 规则引擎、阈值判断、状态匹配
│ 慢系统(s级) │ LLM评估、上下文检索、意图生成
└──────┬───────┘
↓ 输出决策
┌──────────────┐
│ 执行调度器 │ ← 频率控制 + 队列管理
│ Executor │
└──────┬───────┘
↓ 发起动作
用户界面
1. 心跳调度器 (Heartbeat Scheduler)
数学模型
心跳间隔由自适应泊松过程决定:
λ(t) = λ₀ · (1 + α · Δt_idle / T_window) · (1 - β · interaction_density)
- λ₀: 基础心跳频率(默认 180秒)
- Δt_idle: 自上次用户交互以来的空闲时间
- T_window: 观察窗口(默认 30分钟)
- α: 空闲加速因子(0.1~0.5)— 用户越久没交互,心跳越快
- β: 交互密度衰减因子(0.3~0.8)— 用户越活跃,心跳越慢
状态机
WAITING → (心跳触发) → CHECKING → (价值>阈值) → DECIDING → (决策完成) → WAITING
↓ 价值<阈值 ↓ "等待更佳时机"
WAITING DEFERRED → (N分钟后) → WAITING
代码实现
import time
import random
import threading
from enum import Enum
from dataclasses import dataclass
from typing import Optional, Callable
class HeartbeatState(Enum):
WAITING = "waiting" # 等待下次心跳
CHECKING = "checking" # 正在检查是否需要唤醒
DECIDING = "deciding" # 正在做决策
DEFERRED = "deferred" # 延迟到更佳时机
ACTIVE = "active" # 正在与用户交互
@dataclass
class HeartbeatConfig:
base_interval: float = 180.0 # 基础间隔(秒)
min_interval: float = 30.0 # 最小间隔
max_interval: float = 600.0 # 最大间隔(10分钟无交互就降到最低频)
idle_accel_factor: float = 0.3 # 空闲加速 α
interact_decay: float = 0.5 # 交互衰减 β
observation_window: float = 1800.0 # 观察窗口(30分钟)
defer_backoff: float = 120.0 # 延迟后重试间隔
@dataclass
class UserContext:
last_interaction_time: float = 0.0 # 上次用户交互时间戳
interaction_count_30m: int = 0 # 30分钟内交互次数
current_state: HeartbeatState = HeartbeatState.WAITING
is_active_session: bool = False # 是否在活跃对话中
device_mode: str = "normal" # normal / meeting / sleep / driving
class HeartbeatScheduler:
"""
自适应心跳调度器:决定何时唤醒决策系统
"""
def __init__(self, config: HeartbeatConfig = None):
self.config = config or HeartbeatConfig()
self.context = UserContext()
self._timer: Optional[threading.Timer] = None
self._on_tick: Optional[Callable] = None
self._running = False
def set_callback(self, callback: Callable):
"""设置心跳触发回调"""
self._on_tick = callback
def compute_interval(self) -> float:
"""计算下一次心跳间隔"""
now = time.time()
idle_seconds = now - self.context.last_interaction_time
# 基础间隔 + 空闲加速
raw_interval = self.config.base_interval * (
1 + self.config.idle_accel_factor * idle_seconds / self.config.observation_window
)
# 交互密度衰减
density = self.context.interaction_count_30m / (self.config.observation_window / 60)
raw_interval *= (1 - self.config.interact_decay * min(density, 1.0))
# 场景修正
if self.context.device_mode == "meeting":
raw_interval *= 3 # 会议中降频
elif self.context.device_mode == "sleep":
raw_interval *= 10 # 睡眠中极低频
elif self.context.device_mode == "driving":
raw_interval *= 5
# 如果在活跃会话中,暂时停用心跳
if self.context.is_active_session:
return self.config.max_interval
# 夹紧到[min, max]
return max(self.config.min_interval, min(raw_interval, self.config.max_interval))
def record_interaction(self):
"""记录用户交互"""
now = time.time()
self.context.last_interaction_time = now
self.context.interaction_count_30m += 1
self.context.current_state = HeartbeatState.ACTIVE
2. 双系统调度器 (DualSystem Scheduler)
核心思路
借鉴丹尼尔·卡尼曼的「思考,快与慢」理论:
- 快系统(System 1):规则引擎 + 轻量分类器,毫秒级响应,处理80%的日常判断
- 慢系统(System 2):LLM驱动的深度评估,秒级响应,处理需要语境理解的复杂场景
决策流程图
输入:用户上下文 + 时间 + 设备状态
│
▼
┌──────────────────────┐
│ 快系统(System 1) │
│ │
│ 1. 规则过滤 │
│ - 时间窗检查 │
│ - 频率控制检查 │
│ - 场景模式检查 │
│ │
│ 2. 状态模式匹配 │
│ - 深夜+长时间空闲 │
│ - 会话结束后未回复 │
│ - 日程事件触发 │
│ │
│ 3. 快速价值评分 │
│ Value = w₁·urgency │
│ + w₂·recency│
│ + w₃·affinity│
└──────────┬───────────┘
│
┌──────┴──────┐
│ │
Value>0.7 Value<=0.7
│ │
▼ ▼
直接执行 ┌──────────────┐
│ 慢系统(System2)│
│ │
│ 1. LLM评估 │
│ 2. 上下文检索 │
│ 3. 意图生成 │
│ 4. 价值重估 │
└──────┬───────┘
│
┌─────┴─────┐
│ │
价值>0.85 价值<=0.85
│ │
▼ ▼
执行决策 推迟或放弃
代码实现
from dataclasses import dataclass, field
from typing import List, Optional, Tuple
from enum import Enum
import time
class Decision(Enum):
IMMEDIATE_CARE = "immediate_care" # 立即关心
GENTLE_REMIND = "gentle_remind" # 温和提醒
DEFER = "defer" # 延后
SILENCE = "silence" # 保持沉默
CONTENT_SUGGEST = "content_suggest" # 内容推荐
CHECK_IN = "check_in" # 问候关心
@dataclass
class DecisionContext:
user_id: str
idle_minutes: float
time_of_day: str # morning/afternoon/evening/night
device_mode: str
last_conversation_topic: Optional[str] = None
user_mood: str = "unknown" # happy/sad/anxious/neutral/unknown
upcoming_events: List[str] = field(default_factory=list)
unread_count: int = 0
@dataclass
class ValueScore:
score: float # 0.0 - 1.0
urgency: float # 紧急程度
relevance: float # 相关性
interruption_cost: float # 打扰成本
confidence: float # 置信度
reasoning: str = "" # 决策理由(用于思维快照)
class FastSystem:
"""
快系统:规则引擎 + 轻量评估
响应时间 < 50ms
"""
# 规则集:优先级从高到低
RULES = [
# (条件函数, 决策结果, 价值分)
("深夜独处", lambda ctx: ctx.time_of_day == "night" and ctx.idle_minutes > 120,
Decision.IMMEDIATE_CARE, 0.85),
("长时间沉默", lambda ctx: ctx.idle_minutes > 360 and ctx.time_of_day != "night",
Decision.GENTLE_REMIND, 0.70),
("会话搁置", lambda ctx: ctx.last_conversation_topic and ctx.idle_minutes > 30,
Decision.CHECK_IN, 0.65),
("会议中", lambda ctx: ctx.device_mode == "meeting",
Decision.SILENCE, 0.0),
("睡眠中", lambda ctx: ctx.device_mode == "sleep",
Decision.SILENCE, 0.0),
]
@classmethod
def evaluate(cls, ctx: DecisionContext) -> Tuple[Optional[Decision], ValueScore]:
"""
快系统评估
返回 (决策, 价值分数)
"""
for rule_name, condition, decision, base_score in cls.RULES:
if condition(ctx):
# 简单价值计算
urgency = min(ctx.idle_minutes / 120, 1.0) # 2小时达到最大紧迫度
interruption = 0.3 if ctx.time_of_day == "night" else 0.1
score = ValueScore(
score=base_score * (1 + urgency * 0.3 - interruption * 0.2),
urgency=urgency,
relevance=base_score,
interruption_cost=interruption,
confidence=0.8, # 规则引擎置信度高
reasoning=f"[Fast] 匹配规则: {rule_name}"
)
return decision, score
return None, ValueScore(score=0.0, urgency=0, relevance=0,
interruption_cost=0, confidence=0)
class SlowSystem:
"""
慢系统:LLM驱动的深度评估
响应时间 500ms - 5s
"""
@staticmethod
def evaluate(ctx: DecisionContext) -> Tuple[Optional[Decision], ValueScore]:
"""
慢系统评估(调用LLM进行语义理解)
评估维度:
1. 用户情感状态推断(从历史消息中提取)
2. 当前场景的语义价值(是否有关联话题)
3. 打扰成本的精细化估计(用户当前可能在做什么)
4. 个性化偏好匹配(用户之前对类似行为的反馈)
返回决策和价值评分
"""
# 这里调用LLM进行深度评估
# prompt示例:
"""
你是一个AI主动交互决策引擎。请评估在当前场景下是否应该主动发起对话。
用户上下文:
- 空闲时间:{idle_minutes}分钟
- 时段:{time_of_day}
- 设备模式:{device_mode}
- 上次话题:{last_conversation_topic}
- 用户情绪:{user_mood}
- 近期事件:{upcoming_events}
请从以下维度评分(0-1):
1. 消息紧迫性:用户需要关注的事情
2. 话题相关性:当前是否有关联话题可以延续
3. 打扰风险:打扰用户的代价
4. 情感价值:这个互动能提供的情感价值
综合评分 > 0.85 时建议主动发起,0.7-0.85 时延后,< 0.7 时保持沉默。
"""
# 模拟返回
return Decision.GENTLE_REMIND, ValueScore(
score=0.78, urgency=0.6, relevance=0.7,
interruption_cost=0.3, confidence=0.65,
reasoning="[Slow] 用户已空闲90分钟,上次话题有延续价值"
)
class DualSystemScheduler:
"""
双系统调度器:协调快慢系统的协作
"""
def __init__(self, fast_system=FastSystem, slow_system=SlowSystem):
self.fast = fast_system
self.slow = slow_system
self.confidence_threshold = 0.7 # 快系统置信度阈值
self.deflection_count = 0 # 连续被跳过次数
def decide(self, ctx: DecisionContext) -> Tuple[Optional[Decision], ValueScore, str]:
"""
主决策入口
返回:(决策, 价值分数, 使用系统标识)
"""
# 1. 先跑快系统
fast_decision, fast_score = self.fast.evaluate(ctx)
# 2. 如果快系统有高置信度结果,直接采纳
if fast_decision and fast_score.confidence >= self.confidence_threshold:
# 但连续被打断时逐渐降低信任
effective_confidence = fast_score.confidence * (0.9 ** self.deflection_count)
if effective_confidence >= self.confidence_threshold:
return fast_decision, fast_score, "fast"
# 3. 否则调慢系统做深度评估
slow_decision, slow_score = self.slow.evaluate(ctx)
# 4. 快慢系统加权融合
combined_score = fast_score.score * 0.3 + slow_score.score * 0.7
# 5. 最终决策
if combined_score >= 0.85:
self.deflection_count = 0
return slow_decision or Decision.IMMEDIATE_CARE, ValueScore(
score=combined_score,
urgency=(fast_score.urgency + slow_score.urgency) / 2,
relevance=(fast_score.relevance + slow_score.relevance) / 2,
interruption_cost=(fast_score.interruption_cost + slow_score.interruption_cost) / 2,
confidence=(fast_score.confidence + slow_score.confidence) / 2,
reasoning=f"[Fusion] Fast({fast_score.score:.2f}) + Slow({slow_score.score:.2f}) = {combined_score:.2f}"
), "fusion"
elif combined_score >= 0.70:
self.deflection_count += 1
return Decision.DEFER, ValueScore(
score=combined_score, urgency=0, relevance=0,
interruption_cost=0, confidence=0,
reasoning="Deferred: 价值不足,等待更佳时机"
), "defer"
else:
self.deflection_count += 1
return Decision.SILENCE, ValueScore(
score=combined_score, urgency=0, relevance=0,
interruption_cost=0, confidence=0,
reasoning="Silence: 当前不值得打扰用户"
), "silence"
3. 频率控制 & 反骚扰 (Frequency Control & Anti-Spam)
核心约束
全局约束:
· 每小时主动发言 ≤ 3 次
· 每天主动发言 ≤ 15 次
· 连续两次主动发言间隔 ≥ 20 分钟
场景约束:
· 用户标记"不喜欢"同一类内容后,7天内不再触发同类消息
· 用户明确说"别打扰"后,24小时内完全静默
· 深夜(23:00-07:00)仅允许紧急事件触发
记忆约束:
· 同一话题主动发起 ≤ 2 次(避免唠叨)
· 每次主动发言后进入"冷静期",时长 = 反馈评分 × 基础冷静期
代码实现
@dataclass
class FrequencyConfig:
max_per_hour: int = 3
max_per_day: int = 15
min_interval_seconds: int = 1200 # 20分钟
silence_after_reject_hours: int = 24
same_topic_max: int = 2
night_silence_start: int = 23 # 23:00
night_silence_end: int = 7 # 07:00
class FrequencyController:
"""
频率控制 + 反骚扰
"""
def __init__(self, config: FrequencyConfig = None):
self.config = config or FrequencyConfig()
self.action_log: List[float] = [] # 时间戳列表
self.reject_log: List[Tuple[str, float]] = [] # (topic, 时间戳)
self.silence_until: float = 0 # 静默截止时间
self.topic_count: dict = {}
def can_act(self, decision: Decision, topic: str = "") -> Tuple[bool, str]:
"""检查是否允许执行主动发言"""
now = time.time()
# 1. 静默期内 / 被用户拒绝
if now < self.silence_until:
return False, "Silence period active"
# 2. 每小时上限
recent_hour = [t for t in self.action_log if now - t < 3600]
if len(recent_hour) >= self.config.max_per_hour:
return False, "Hourly limit reached"
# 3. 每天上限
recent_day = [t for t in self.action_log if now - t < 86400]
if len(recent_day) >= self.config.max_per_day:
return False, "Daily limit reached"
# 4. 最小间隔
if self.action_log and (now - self.action_log[-1]) < self.config.min_interval_seconds:
remaining = self.config.min_interval_seconds - (now - self.action_log[-1])
return False, f"Cooldown: {remaining:.0f}s remaining"
# 5. 同一话题检查
if topic and self.topic_count.get(topic, 0) >= self.config.same_topic_max:
return False, f"Topic '{topic}' already triggered {self.same_topic_max} times"
# 6. 深夜限制
hour = time.localtime().tm_hour
if self.config.night_silence_start <= hour or hour < self.config.night_silence_end:
if decision not in (Decision.IMMEDIATE_CARE,):
return False, "Night mode: only urgent allowed"
return True, "OK"
def record_action(self, decision: Decision, topic: str = "", feedback: float = 0.5):
"""记录一次主动发言"""
now = time.time()
self.action_log.append(now)
if topic:
self.topic_count[topic] = self.topic_count.get(topic, 0) + 1
# 根据反馈动态调整冷静期
cooldown = self.config.min_interval_seconds * (1 + feedback)
def record_rejection(self, topic: str = ""):
"""记录用户拒绝"""
now = time.time()
self.reject_log.append((topic, now))
self.silence_until = now + self.config.silence_after_reject_hours * 3600
4. 完整调度链路示例
# 初始化各组件
heartbeat = HeartbeatScheduler()
fast = FastSystem()
slow = SlowSystem()
scheduler = DualSystemScheduler(fast, slow)
freq = FrequencyController()
# 心跳触发回调
def on_heartbeat_tick():
ctx = DecisionContext(
user_id="user_001",
idle_minutes=90,
time_of_day="night",
device_mode="normal",
last_conversation_topic="工作压力",
user_mood="anxious",
)
decision, score, system = scheduler.decide(ctx)
if decision == Decision.SILENCE or decision == Decision.DEFER:
return # 不打扰
allowed, reason = freq.can_act(decision, topic=ctx.last_conversation_topic)
if not allowed:
return # 频率限制
freq.record_action(decision)
# 执行主动发言(调用生成模块)
execute_action(decision, score, ctx)
# 启动心跳
heartbeat.set_callback(on_heartbeat_tick)
heartbeat.start()
5. 性能指标
| 组件 | 响应时间 | 内存 | 适用场景 |
|---|---|---|---|
| 心跳调度器 | <1ms | 2KB | 决定何时检查 |
| 快系统规则引擎 | 5-50ms | 100KB | 80%的日常决策 |
| 慢系统LLM评估 | 500ms-5s | 动态加载 | 20%需要深度理解的场景 |
| 频率控制器 | <1ms | 10KB | 每次执行前检查 |