-
Notifications
You must be signed in to change notification settings - Fork 1
Expand file tree
/
Copy pathmain.py
More file actions
1384 lines (1235 loc) · 59.3 KB
/
Copy pathmain.py
File metadata and controls
1384 lines (1235 loc) · 59.3 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
"""
Iris Memory - AstrBot 轻量化三合一插件
v3.0 架构:
- 记忆侧(源自 Iris Chat Memory 轻量方案):
L1 消息上下文缓冲 / L2 记忆库(FAISS + SQLite)/ L3 知识图谱(SQLite)
+ 用户/群聊画像 + 梦境离线加工 + 图片解析
- 主动回复侧(源自 Iris Reply 统一决策模型):
chime_in 跟话 / follow_up 跟进 / initiate 发起 / watch 被动评估,
SignalGate 本地零成本门控 + 单次 LLM 统一决策 + ThreadAnchor 记账
- 人格自学习迭代:
脱敏采样 + 风格分析 + 候选生成 + 独立审查 + Revision 审批/回滚
钩子编排(等价于原两插件并存时的兼容性契约):
群消息 → on_message(主动回复门控,设 iris_mode extra)
→ on_all_message(入 L1、图片入队)
门控命中 → 统一决策(llm_generate 直调,不触发钩子)
→ 决策发言:劫持主管线,注入 SPEAK_HINTS(mark_as_temp)
on_llm_request → 记忆侧:清空 contexts,注入 L1/L2/L3/画像(mark_as_temp)
on_llm_response → 记忆侧:bot 回复入 L1 → 主动回复侧:按 iris_mode 记账
after_message_sent → 主动回复侧:入滑动窗口 + 写 ThreadAnchor
initiate 直发(context.send_message)→ 手动记账 + 回填 L1
"""
import asyncio
import hashlib
import re
import time
from typing import Any, Optional
from .iris_memory.config import init_config, Config
from astrbot.api import AstrBotConfig
from astrbot.api.event import filter, AstrMessageEvent
from astrbot.api.message_components import At
from astrbot.api.star import Context, Star, StarTools
from astrbot.core.agent.message import TextPart
from astrbot.core.provider.entities import LLMResponse, ProviderRequest
from .iris_memory.core import (
ComponentManager,
get_logger,
create_components,
initialize_components,
shutdown_components,
handle_user_message,
preprocess_llm_request,
handle_llm_response,
handle_agent_done,
handle_pre_request_cleanup,
handle_initiate_backfill,
set_component_manager,
get_run_log_manager,
)
from .iris_memory.tools import register_llm_tools
from .iris_memory.utils.token_counter import warm_up_encoders_async
from .iris_memory.web import register_all_routes
from .iris_memory.llm import LLMManager
from .iris_memory.llm_modules import FRAMEWORK_REPLY, proactive_reply_module
from .iris_memory.commands import (
get_registry,
execute_command,
L1CommandHandler,
L2CommandHandler,
L3CommandHandler,
ProfileCommandHandler,
AllCommandHandler,
LearningCommandHandler,
EvolutionCommandHandler,
)
from .iris_memory.proactive.admin import AdminCommands
from .iris_memory.proactive.api import (
register_web_apis as register_reply_web_apis,
sync_stats_group_state,
)
from .iris_memory.proactive.config import ConfigManager as ReplyConfigManager
from .iris_memory.proactive.decision import (
INPUT_SAFETY_COOLDOWN_MINUTES,
DecisionCore,
DecisionRequest,
SafetyCleanupResult,
)
from .iris_memory.proactive.perception import (
ContextPackager,
Gatekeeper,
SlidingWindow,
WindowMessage,
)
from .iris_memory.proactive.prompts import SPEAK_HINTS
from .iris_memory.proactive.proactive import ProactiveEngine
from .iris_memory.proactive.signals import SignalGate
from .iris_memory.proactive.state import StateManager
from .iris_memory.proactive.stats import StatsCollector
from .iris_memory.proactive.tickets import DecisionTicketRegistry
from .iris_memory.proactive.time_hint import resolve_datetime_reminder
from .iris_memory.proactive.tools import ToolContext
from .iris_memory.extras import ErrorFriendlyProcessor, MarkdownStripper
logger = get_logger("main")
PLUGIN_NAME = "astrbot_plugin_iris_memory"
# 旧版(v2.x)数据自动迁移开关;v4 删除本常量与 iris_memory/legacy_migration/ 即彻底移除
LEGACY_MIGRATION_ENABLED = True
_IRIS_ACTIVE_TIMEOUT = 120
_UMO_KV_KEY = "iris_reply:group_umo"
_PASSIVE_WATCH_SIGNAL = re.compile(
r"[??]|"
r"我(?:会|将|之后|稍后|晚点|等会|回头)|"
r"(?:下次|之后|稍后|晚点|有结果|有消息)(?:再|会)?(?:告诉|提醒|联系|确认|回复)|"
r"(?:等|等待).{0,16}(?:回复|消息|结果)|"
r"\b(?:i(?:'ll| will)|follow up|let you know|remind|check back|waiting for)\b",
re.IGNORECASE,
)
def _detect_passive_trigger(event: AstrMessageEvent, req, context: Context) -> None:
"""检测 LLM 请求是否为被动触发(sampling/主动回复)
当用户消息不以唤醒前缀开头且未 @机器人 时,LLM 请求可能是由 AstrBot 的
active_reply/sampling 机制触发的。此时标记事件,供后续钩子
判断是否跳过图片解析等高 token 消耗操作。
注:本插件主动回复触发的请求会将 is_at_or_wake_command 置 True,
天然不会被误判为被动触发,无需额外处理 iris_mode。
"""
try:
is_at_or_wake = getattr(event, "is_at_or_wake_command", False)
if not is_at_or_wake:
event.set_extra("iris_passive_trigger", True)
logger.debug(
"检测到被动触发(sampling/主动回复),is_at_or_wake_command 为 False"
)
except Exception as e:
logger.debug(f"被动触发检测异常(不影响正常流程):{e}")
def _is_pure_at_self(event: AstrMessageEvent) -> bool:
"""Return whether the message consists only of an @ mention to the bot."""
try:
messages = event.get_messages()
return (
len(messages) == 1
and isinstance(messages[0], At)
and str(messages[0].qq) == str(event.get_self_id())
)
except Exception:
return False
class IrisMemoryPlugin(Star):
"""AstrBot 轻量化三合一插件主类。"""
def __init__(self, context: Context, config: AstrBotConfig | None = None):
super().__init__(context)
self.context: Context = context
try:
# ── 记忆侧初始化 ──
data_dir = StarTools.get_data_dir()
self.config: Config = init_config(config, data_dir)
logger.info(f"插件数据目录:{data_dir}")
components = create_components(context, self)
self.component_manager: Optional[ComponentManager] = ComponentManager(
components
)
self._llm_manager = self.component_manager.get_component(
"llm_manager", LLMManager
)
set_component_manager(self.component_manager)
from .iris_memory.image.recorder_bridge import init_recorder_bridge
init_recorder_bridge(context)
# ── extras(自 v2 保留的低成本功能) ──
self._error_processor = ErrorFriendlyProcessor(self.config)
self._markdown_stripper = MarkdownStripper(
context=self.context,
config=self.config,
)
# ── 主动回复侧初始化 ──
self._reply_config = ReplyConfigManager(
config if config else context.get_config(),
hidden_get=self.config.get,
)
self._state = StateManager(self._reply_config)
self._tool_ctx = ToolContext()
self._gatekeeper = Gatekeeper(self._reply_config, self._state)
self._sliding_window = SlidingWindow(self._reply_config)
self._context_packager = ContextPackager(
self._reply_config, self_id_get=lambda: self._self_id
)
self._signals = SignalGate(self._reply_config, self._state)
self._decision_core = DecisionCore(
self._reply_config, self._state, self._sliding_window, self._context_packager,
time_hint_get=lambda gid: resolve_datetime_reminder(
self.context, self._group_umo.get(gid),
),
)
self._admin = AdminCommands(self._state)
self._stats = StatsCollector()
self._reply_in_progress: dict[str, float] = {}
self._passive_active: dict[str, float] = {}
self._triggering: dict[str, float] = {}
self._decision_tickets = DecisionTicketRegistry(_IRIS_ACTIVE_TIMEOUT)
self._passive_watch_last: dict[str, float] = {}
self._passive_watch_hash: dict[str, str] = {}
self._follow_pending: set[str] = set()
self._group_umo: dict[str, str] = {}
self._umo_dirty: bool = False
self._self_id: str = ""
self._save_task: asyncio.Task | None = None
self._save_interval = 30
self._proactive = ProactiveEngine(
self.context,
self._reply_config,
self._state,
self._sliding_window,
self._signals,
self._decision_core,
self._stats,
llm_manager=self._llm_manager,
packager=self._context_packager,
umo_get=lambda gid: self._group_umo.get(gid),
is_busy=self._is_busy,
self_id_get=lambda: self._self_id,
save_fn=lambda: self._state.save_dirty(self._kv_save),
on_initiate_sent=self._on_initiate_sent,
text_transform=self._strip_initiate_text,
)
# ── 外部注册动作(依赖就绪后在构造函数末尾统一执行) ──
self._register_llm_tools()
self._register_command_handlers()
self._register_web_api()
self._warn_empty_mention_conflict()
logger.info("Iris Memory 整合插件已加载(等待异步初始化)")
except Exception:
logger.error(
"Iris Memory 插件初始化失败,真实错误如下(框架可能将其掩盖为 "
"'missing 1 required positional argument: config',请以下方堆栈为准):",
exc_info=True,
)
raise
def _warn_empty_mention_conflict(self) -> None:
"""Guide users to disable AstrBot's built-in pure-mention waiter."""
try:
if not self.config.get("extras.pure_at_reply.enable", False):
return
platform_settings = self.context.get_config().get("platform_settings", {})
if platform_settings.get("empty_mention_waiting", True):
logger.warning(
"检测到 AstrBot「只 @ 机器人触发等待」已开启:纯 @ 会先被 "
"AstrBot 内置处理器拦截,Iris 无法接管。若要让纯 @ 直接使用 "
"L1 上下文,请关闭 platform_settings.empty_mention_waiting。"
)
except Exception as exc:
logger.debug(f"检查 AstrBot 纯 @ 配置失败(已忽略):{exc}")
async def _prepare_pure_at_request(self, event: AstrMessageEvent) -> bool:
"""Route a pure @ through the standard agent pipeline using existing L1.
AstrBot rejects a truly empty ProviderRequest before on_llm_request hooks
run. A whitespace-only transport prompt keeps the request alive through
that guard; ProviderRequest.assemble_context() drops it before provider
dispatch, so the model receives no synthetic user content and relies on
the normal Iris L1 injection.
"""
if not _is_pure_at_self(event) or not self.component_manager:
return False
try:
if not self.config.get("extras.pure_at_reply.enable", False):
return False
cfg = self.context.get_config(umo=event.unified_msg_origin)
if cfg.get("platform_settings", {}).get("empty_mention_waiting", True):
return False
from .iris_memory.platform import get_adapter
adapter = get_adapter(event)
buffer = self.component_manager.get_available_component("l1_buffer")
if not buffer or not buffer.get_context(adapter.get_session_id(event), 1):
logger.debug("纯 @ 前没有可注入的 L1 上下文,保持 AstrBot 默认空请求行为")
return False
conv_mgr = getattr(self.context, "conversation_manager", None)
if conv_mgr is None:
logger.warning("纯 @ 接管失败:AstrBot conversation_manager 不可用")
return False
umo = event.unified_msg_origin
conversation_id = await conv_mgr.get_curr_conversation_id(umo)
if not conversation_id:
conversation_id = await conv_mgr.new_conversation(
umo, platform_id=event.get_platform_id()
)
conversation = await conv_mgr.get_conversation(umo, conversation_id)
if conversation is None:
logger.warning("纯 @ 接管失败:无法取得标准会话")
return False
request = event.request_llm(
prompt=" ",
session_id=conversation_id,
contexts=[],
system_prompt="",
conversation=conversation,
)
event.set_extra("provider_request", request)
event.set_extra("iris_pure_at", True)
logger.info("已接管纯 @ 消息,将通过标准管道使用现有 L1 上下文")
return True
except Exception as exc:
logger.error(f"纯 @ 消息接管失败,已回退 AstrBot 默认流程:{exc}", exc_info=True)
return False
# ========================================================================
# 记忆侧注册
# ========================================================================
def _register_llm_tools(self) -> None:
"""统一注册全部 LLM Tool"""
try:
self._registered_llm_tools = register_llm_tools(
context=self.context,
owner_module=self.__class__.__module__,
state=self._state,
tool_context=self._tool_ctx,
)
logger.info(f"已注册 {len(self._registered_llm_tools)} 个 Iris LLM Tool")
except Exception as e:
logger.error(f"注册 Iris LLM Tool 失败:{e}", exc_info=True)
def _register_command_handlers(self) -> None:
"""注册记忆侧指令处理器"""
try:
registry = get_registry()
handlers = [
L1CommandHandler(),
L2CommandHandler(),
L3CommandHandler(),
ProfileCommandHandler(),
AllCommandHandler(),
LearningCommandHandler(),
EvolutionCommandHandler(),
]
for handler in handlers:
registry.register(handler)
logger.info(f"已注册 {len(handlers)} 个记忆指令处理器")
except Exception as e:
logger.error(f"注册记忆指令处理器失败:{e}", exc_info=True)
def _register_web_api(self) -> None:
try:
register_all_routes(self.context)
except Exception as e:
logger.error(f"注册记忆 Web API 失败:{e}", exc_info=True)
# ========================================================================
# 生命周期
# ========================================================================
async def initialize(self) -> None:
# 0. 后台预热 tiktoken 编码器:首次运行需联网下载 BPE 文件,
# 放线程池执行,预热期间 token 计数临时走字符估算
self._encoder_warmup_task = asyncio.create_task(warm_up_encoders_async())
# 让预热协程先设置 pending 标志,再进入其余组件初始化;否则本轮
# 事件循环中更早发生的 token 计数仍可能同步触发首次下载。
await asyncio.sleep(0)
# 1. 记忆组件初始化
try:
await initialize_components(self.component_manager)
except Exception as e:
logger.error(f"记忆组件初始化失败:{e}", exc_info=True)
# 2. 旧版(v2.x)数据自动迁移(独立模块,失败不阻断启动)
if LEGACY_MIGRATION_ENABLED:
try:
from .iris_memory.legacy_migration import migrate_if_needed
await migrate_if_needed(
self.context, self, StarTools.get_data_dir(), self.component_manager
)
except Exception:
logger.error("旧数据迁移失败(不影响插件启动)", exc_info=True)
# 3. 主动回复侧初始化
await self._state.load_all(self._kv_load)
umo_data = await self._kv_load(_UMO_KV_KEY)
if isinstance(umo_data, dict):
self._group_umo = {str(k): str(v) for k, v in umo_data.items()}
# 旧版页面配置(KV overrides)一次性迁移到隐藏参数,迁移后清空避免重复覆盖
config_overrides = await self._kv_load("iris_reply:config_overrides")
migrated = ReplyConfigManager.legacy_overrides_to_hidden(config_overrides)
if migrated:
self.config.update_hidden(migrated)
await self._kv_save("iris_reply:config_overrides", {})
logger.info(f"已将 {len(migrated)} 项旧版主动回复页面配置迁移至隐藏参数")
self._save_task = asyncio.create_task(self._periodic_save())
self._stats.enabled = self._reply_config.stats_enabled
register_reply_web_apis(
context=self.context,
plugin_name=PLUGIN_NAME,
state=self._state,
stats=self._stats,
window=self._sliding_window,
kv_save=self._kv_save,
)
await self._proactive.start()
# 4. 功能重叠插件检测(重复注入/门控警告)
for other in ("astrbot_plugin_iris_chat_memory", "astrbot_plugin_iris_reply"):
try:
if self.context.get_registered_star(other):
logger.warning(
f"检测到插件 {other} 已安装,与本插件功能重叠,"
"建议停用其一,避免记忆重复注入与主动回复重复门控"
)
except Exception:
pass
logger.info("Iris Memory 整合插件异步初始化完成")
async def terminate(self):
"""插件卸载清理"""
logger.info("开始关闭插件组件...")
warmup = getattr(self, "_encoder_warmup_task", None)
if warmup and not warmup.done():
warmup.cancel()
try:
await warmup
except asyncio.CancelledError:
pass
# 主动回复侧
await self._proactive.stop()
if self._save_task and not self._save_task.done():
self._save_task.cancel()
try:
await self._save_task
except asyncio.CancelledError:
pass
await self._state.save_all(self._kv_save)
await self._kv_save(_UMO_KV_KEY, dict(self._group_umo))
self._follow_pending.clear()
self._reply_in_progress.clear()
self._passive_active.clear()
self._triggering.clear()
self._decision_tickets.clear()
self._passive_watch_last.clear()
self._passive_watch_hash.clear()
# 记忆侧
await shutdown_components(self.component_manager)
logger.info("Iris Memory 整合插件已卸载")
# ========================================================================
# 主动回复侧:状态保存与互斥
# ========================================================================
async def _periodic_save(self) -> None:
while True:
await asyncio.sleep(self._save_interval)
try:
await self._state.save_dirty(self._kv_save)
if self._umo_dirty:
self._umo_dirty = False
await self._kv_save(_UMO_KV_KEY, dict(self._group_umo))
self._sliding_window.cleanup(self._state.get_whitelist())
self._cleanup_stale_active()
sync_stats_group_state(self._state, self._stats)
except Exception as e:
logger.warning(f"Iris Reply: periodic save error: {e}")
def _cleanup_stale_active(self) -> None:
now = time.time()
stale_rip = [gid for gid, ts in self._reply_in_progress.items() if now - ts > _IRIS_ACTIVE_TIMEOUT]
for gid in stale_rip:
logger.info(f"Iris Reply: cleaning up stale reply_in_progress for group {gid} (timeout)")
self._reply_in_progress.pop(gid, None)
stale_passive = [gid for gid, ts in self._passive_active.items() if now - ts > _IRIS_ACTIVE_TIMEOUT]
for gid in stale_passive:
logger.info(f"Iris Reply: cleaning up stale passive for group {gid} (timeout)")
self._passive_active.pop(gid, None)
stale_triggering = [gid for gid, ts in self._triggering.items()
if gid not in self._reply_in_progress and now - ts > _IRIS_ACTIVE_TIMEOUT]
for gid in stale_triggering:
logger.info(f"Iris Reply: cleaning up stale triggering for group {gid}")
self._triggering.pop(gid, None)
registry = getattr(self, "_decision_tickets", None)
if registry is not None:
for ticket in registry.cleanup(now):
logger.info(
"Iris Reply: cleaning up stale decision ticket "
f"for group {ticket.group_id}, event={ticket.event_id}"
)
if ticket.group_id not in self._reply_in_progress:
self._triggering.pop(ticket.group_id, None)
@staticmethod
def _decision_event_id(event: AstrMessageEvent) -> str:
"""取得平台消息 ID;缺失时退化为当前事件对象 ID。"""
message_obj = getattr(event, "message_obj", None)
raw_id = getattr(message_obj, "message_id", None) if message_obj else None
if raw_id not in (None, ""):
return f"message:{raw_id}"
return f"event:{id(event)}"
def _release_decision_ticket(self, group_id: str, ticket_id: str = "") -> None:
"""仅由 ticket 所有者释放,避免迟到事件清掉新事件的决策状态。"""
registry = getattr(self, "_decision_tickets", None)
if registry is None:
self._triggering.pop(group_id, None)
return
active = registry.get(group_id)
if active is not None and ticket_id and active.ticket_id != ticket_id:
return
registry.release(group_id, ticket_id)
self._triggering.pop(group_id, None)
def _is_busy(self, group_id: str) -> bool:
return (
group_id in self._reply_in_progress
or group_id in self._triggering
or group_id in self._passive_active
)
async def _kv_save(self, key: str, value: Any) -> None:
await self.put_kv_data(key, value)
async def _kv_load(self, key: str) -> Any:
return await self.get_kv_data(key, None)
def _get_group_id(self, event) -> str | None:
group_id = event.get_group_id()
if not group_id:
event.set_result("无法获取群ID")
return None
return group_id
async def _get_provider_id(self, event, preferred: str = "") -> str | None:
if preferred:
return preferred
try:
return await self.context.get_current_chat_provider_id(
event.unified_msg_origin
)
except Exception:
logger.error("Iris Reply: failed to get provider ID")
return None
def _strip_initiate_text(self, text: str) -> str:
"""initiate 直发消息的 Markdown 去除
直发通路(context.send_message)不触发 on_decorating_result 钩子,
消息始终以纯文本发送到平台,此处补齐与管线消息一致的 Markdown 去除,
避免同群内跟话消息与主动发起消息格式处理不一致。
"""
stripper = self._markdown_stripper
if not stripper or not text:
return text
try:
if not stripper.should_strip(text, use_t2i=False):
return text
return stripper.strip(text)
except Exception as e:
logger.warning(f"initiate 消息 Markdown 去除失败:{e}")
return text
async def _on_initiate_sent(self, group_id: str, text: str) -> None:
"""initiate 直发成功后,把 bot 发言回填进 L1 缓冲
直发通路(context.send_message)不触发任何事件钩子,
若不回填,L1 上下文中将看不到这类发起消息。
"""
if not self.component_manager or not text:
return
try:
await handle_initiate_backfill(group_id, text, self.component_manager)
except Exception as e:
logger.warning(f"initiate 消息回填 L1 失败:{e}")
# ========================================================================
# 主动回复侧:LLM 工具(已迁移至 iris_memory/tools/proactive.py,
# 由 _register_llm_tools() 统一注册)
# ========================================================================
# ========================================================================
# 主动回复侧:管理指令
# ========================================================================
@filter.command_group("iris_reply")
@filter.permission_type(filter.PermissionType.ADMIN)
@filter.event_message_type(filter.EventMessageType.GROUP_MESSAGE)
def iris_reply_group(self):
pass
@iris_reply_group.command("enable")
async def cmd_enable(self, event) -> None:
group_id = self._get_group_id(event)
if not group_id:
return
self._state.add_to_whitelist(group_id)
await self._state.save_dirty(self._kv_save)
event.set_result(f"群 {group_id} 已启用主动回复")
@iris_reply_group.command("disable")
async def cmd_disable(self, event) -> None:
group_id = self._get_group_id(event)
if not group_id:
return
self._state.remove_from_whitelist(group_id)
self._sliding_window.remove_group(group_id)
self._state.remove_group_lock(group_id)
await self._state.save_dirty(self._kv_save)
event.set_result(f"群 {group_id} 已禁用主动回复")
@iris_reply_group.command("status")
async def cmd_status(self, event) -> None:
group_id = self._get_group_id(event)
if not group_id:
return
text = self._admin.get_status(group_id)
event.set_result(text)
@iris_reply_group.command("reset")
async def cmd_reset(self, event) -> None:
group_id = self._get_group_id(event)
if not group_id:
return
msg = self._admin.reset_group(group_id)
self._sliding_window.remove_group(group_id)
await self._state.save_dirty(self._kv_save)
event.set_result(msg)
@iris_reply_group.command("cooldown")
async def cmd_cooldown(self, event, minutes: int = 5) -> None:
group_id = self._get_group_id(event)
if not group_id:
return
msg = self._admin.set_cooldown(group_id, minutes)
await self._state.save_dirty(self._kv_save)
event.set_result(msg)
@iris_reply_group.command("willingness")
async def cmd_willingness(self, event, level: str = "") -> None:
group_id = self._get_group_id(event)
if not group_id:
return
if not level.strip():
current = self._admin.get_willingness(group_id)
event.set_result(f"群 {group_id} 当前回复意愿: {current}\n可选: 低/中/高 (low/medium/high)")
return
msg = self._admin.set_willingness(group_id, level.strip())
await self._state.save_dirty(self._kv_save)
event.set_result(msg)
@iris_reply_group.command("initiate")
async def cmd_initiate(self, event) -> None:
group_id = self._get_group_id(event)
if not group_id:
return
result = await self._proactive.attempt_initiate(group_id, force=True)
event.set_result(f"主动发起: {result}")
# ========================================================================
# AstrBot 钩子
# ========================================================================
@filter.event_message_type(filter.EventMessageType.GROUP_MESSAGE)
async def on_message(self, event) -> None:
"""主动回复消息唤醒:门控 → 标记 → 交由 on_llm_request 决策"""
if await self._prepare_pure_at_request(event):
return
if not self._reply_config.enabled:
return
if not self._gatekeeper.should_process(event):
return
group_id = event.get_group_id()
if not group_id:
return
# 缓存会话标识与自身 ID,供主动发起通路使用
umo = getattr(event, "unified_msg_origin", "")
if umo and self._group_umo.get(group_id) != umo:
self._group_umo[group_id] = umo
self._umo_dirty = True
if not self._self_id:
self._self_id = event.get_self_id() or ""
message_str = event.message_str or ""
sender_id = event.get_sender_id()
sender_name = event.get_sender_name() or sender_id
# 发起后的首次接话:清除 pending,该消息直接获得一次跟进评估资格
pending_reply = self._state.consume_initiate_pending(group_id)
is_followed = bool(sender_id and self._state.match_anchor_user(group_id, sender_id))
score = self._gatekeeper.quality_score(message_str)
if score < self._reply_config.quality_threshold and not is_followed and not pending_reply:
return
message_timestamp = time.time()
self._sliding_window.append(
group_id,
WindowMessage(
sender_id=sender_id,
sender_name=sender_name,
content=message_str,
timestamp=message_timestamp,
),
)
# 每条真实群消息都会使旧主动候选失效,并为该群独立重新预约。
self._proactive.notify_human_message(group_id, message_timestamp)
if event.is_at_or_wake_command:
self._release_decision_ticket(group_id)
self._state.increment_msg_count(group_id)
self._passive_active[group_id] = time.time()
event.set_extra("iris_mode", "passive")
return
if self._is_busy(group_id) or self._proactive.is_initiating(group_id):
logger.debug(f"Iris Reply: reply already in progress for group {group_id}")
return
async with self._state.get_lock(group_id):
motive = self._signals.evaluate_message(group_id, sender_id, message_str)
if not motive and pending_reply:
motive = "follow_up"
if not motive:
return
is_follow_up = motive == "follow_up"
if is_follow_up:
if group_id in self._follow_pending:
logger.debug(f"Iris Reply: follow-up aggregation pending for group {group_id}")
return
self._follow_pending.add(group_id)
try:
await asyncio.sleep(self._reply_config.follow_up_aggregate_window)
finally:
self._follow_pending.discard(group_id)
if self._is_busy(group_id):
return
if not pending_reply and not self._state.get_anchor(group_id).active:
return
if not self._state.can_detect(group_id, follow_up=is_follow_up):
logger.debug(f"Iris Reply: trigger rate-limited for group {group_id}")
return
provider_id = self._reply_config.provider_id
if not provider_id:
provider_id = await self._get_provider_id(event)
if not provider_id:
logger.error(f"Iris Reply: failed to get provider ID for group {group_id}")
return
async with self._state.get_lock(group_id):
if group_id in self._triggering:
logger.debug(f"Iris Reply: trigger already in progress for group {group_id}")
return
event_id = self._decision_event_id(event)
registry = getattr(self, "_decision_tickets", None)
ticket = registry.claim(group_id, event_id) if registry else None
if registry is not None and ticket is None:
logger.debug(
f"Iris Reply: duplicate decision event rejected for group {group_id}, "
f"event={event_id}"
)
return
self._state.record_detect_time(group_id)
self._triggering[group_id] = time.time()
event.set_extra("iris_decision", {
"motive": motive,
"provider_id": provider_id,
"event_id": event_id,
"ticket_id": ticket.ticket_id if ticket else "",
})
event.is_at_or_wake_command = True
event.is_wake = True
if provider_id:
event.set_extra("selected_provider", provider_id)
self._tool_ctx.set_context(group_id)
logger.info(
f"Iris Reply: {motive} candidate activated for group {group_id}, deferred to on_llm_request"
)
@filter.event_message_type(filter.EventMessageType.ALL)
async def on_all_message(self, event: AstrMessageEvent) -> None:
"""记忆侧:全类型消息入 L1 缓冲、图片入队"""
if self.component_manager:
await handle_user_message(event, self.component_manager)
@filter.permission_type(filter.PermissionType.ADMIN)
@filter.command("iris_mem")
async def iris_mem(self, event: AstrMessageEvent) -> None:
if self.component_manager:
result = await execute_command(event)
if result:
yield event.plain_result(result)
@filter.permission_type(filter.PermissionType.ADMIN)
@filter.command("memory")
async def memory_legacy_alias(self, event: AstrMessageEvent) -> None:
"""旧版 /memory 指令指引别名
v2 插件的 /memory 指令在本插件中已更名为 /iris_mem。旧指令输入
此前不会被任何处理器响应(无反馈、无删除),迁移用户容易误以为
清除已生效。此处仅返回迁移指引,不执行任何操作。
"""
yield event.plain_result(
"ℹ️ /memory 指令已迁移到 /iris_mem\n"
"按层清除:/iris_mem l1|l2|l3 clear [--group|--all|@用户]\n"
"全部清除:/iris_mem all clear [--group|--all]\n"
"完整用法:/iris_mem help"
)
@filter.on_llm_request()
async def on_llm_request(self, event: AstrMessageEvent, req: ProviderRequest) -> None:
# 1. 主动回复统一决策(仅处理 _triggering 中的群;可能 stop_event 终止请求)
if await self._handle_reply_decision(event):
return
# 2. 记忆侧:被动触发检测 + 上下文接管 + L1/L2/L3/画像注入
if self.component_manager:
_detect_passive_trigger(event, req, self.context)
await handle_pre_request_cleanup(
event, req, self.context, self.component_manager
)
await preprocess_llm_request(event, req, self.component_manager)
# 3. 主动回复发言提示:决策通过后由 _handle_reply_decision 暂存,
# 在记忆注入之后追加,保证 LLM 先看到上下文、再看到发言指令
hint = event.get_extra("iris_speak_hint")
if hint:
req.extra_user_content_parts.append(TextPart(text=hint).mark_as_temp())
# 所有 AstrBot 主管线请求也必须取得 Governor lease。此处刻意位于
# 查询改写/图片关联解析等预处理之后,避免单并发配置下嵌套申请死锁。
mode = event.get_extra("iris_mode")
if self._llm_manager:
provider_id = event.get_extra("iris_llm_provider_id") or ""
if not provider_id:
provider_id = await self._get_provider_id(event) or ""
module = (
proactive_reply_module(mode)
if mode in ("chime_in", "follow_up", "passive")
else FRAMEWORK_REPLY
)
lease = await self._llm_manager.acquire_framework_lease(
module=module,
provider_id=provider_id,
)
try:
# 先把 lease ID 写入事件,确保后续任一步骤抛错时都能在
# 当前栈内释放,而不必等待 watchdog 回收。
event.set_extra("iris_llm_lease_id", lease.lease_id)
await self._llm_manager.record_framework_attempt(module)
event.set_extra(
"iris_llm_tracking",
{
"module": module,
"provider_id": provider_id,
"started_at": time.time(),
"prompt": getattr(req, "prompt", "") or "",
"priority": lease.priority.name,
"queue_wait_ms": lease.queue_wait_ms,
"in_flight_at_start": lease.in_flight_at_start,
"queue_depth_at_start": lease.queue_depth_at_start,
},
)
except Exception:
event.set_extra("iris_llm_lease_id", None)
await self._llm_manager.release_framework_lease(
lease.lease_id,
provider_id=provider_id,
success=False,
)
raise
async def _handle_reply_decision(self, event: AstrMessageEvent) -> bool:
"""主动回复统一决策执行点。
Returns:
True 表示已调用 event.stop_event(),调用方应立即返回,
不再进行记忆注入。
"""
group_id = event.get_group_id()
if not group_id or group_id not in self._triggering:
return False
info = event.get_extra("iris_decision")
if not info:
self._release_decision_ticket(group_id)
return False
motive = info.get("motive", "")
provider_id = info.get("provider_id", "")
ticket_id = str(info.get("ticket_id") or "")
event_id = str(info.get("event_id") or "")
registry = getattr(self, "_decision_tickets", None)
if registry is not None:
active = registry.get(group_id)
# 旧测试/升级中的遗留事件没有 ticket_id,继续走兼容路径;生产事件
# 一旦携带 ticket,就必须同时匹配 ticket_id + event_id。
if ticket_id and not registry.owns(group_id, ticket_id, event_id):
logger.warning(
f"Iris Reply: stale decision ticket rejected for group {group_id}, "
f"event={event_id}"
)
event.stop_event()
return True
if active is not None and not ticket_id:
logger.warning(
f"Iris Reply: decision event without owner ticket rejected for group {group_id}"
)
event.stop_event()
return True
req = DecisionRequest(group_id=group_id, wake="message", motive=motive)
outcome = await self._decision_core.decide(req, self._llm_manager, provider_id)
if outcome.error or outcome.decision is None:
if outcome.error_kind == "input_content_safety_1026":
cleanup, cooldown = await self._clear_rejected_reply_context(
group_id,
outcome.dynamic_context_sources,
record_skip=True,
)
logger.error(
"Iris Reply: decision input rejected by provider safety filter "
f"for group {group_id} (1026, retryable=false, "
f"dynamic_sources={cleanup.dynamic_source_count}, "
f"window_removed={cleanup.window_removed}, "
f"observation_cleared={cleanup.observation_cleared}, "
f"anchor_cleared={cleanup.anchor_cleared}, "
f"cooldown={cooldown}min): "
f"{outcome.error}"
)
else:
logger.error(f"Iris Reply: decision LLM call failed for group {group_id}: {outcome.error}")
async with self._state.get_lock(group_id):
self._state.record_skip_reply(group_id)
await self._state.save_dirty(self._kv_save)
self._stats.record_decision_error(group_id, motive)
self._release_decision_ticket(group_id, ticket_id)
event.stop_event()
return True
decision = outcome.decision
logger.info(
f"Iris Reply: decision raw for group {group_id} (motive={motive}, "
f"len={len(outcome.raw_text)}): {outcome.raw_text:.500s}"
)
self._stats.record_decision(
group_id, motive,
system_prompt=outcome.system_prompt,
user_prompt=outcome.user_prompt,
response_text=outcome.raw_text,
decision=decision,
duration_ms=outcome.duration_ms,
)
logger.info(
f"Iris Reply: decision parsed for group {group_id}: speak={decision.should_speak}, "
f"drifted={decision.drifted}, watch={decision.watch}, "
f"watch_keywords={decision.watch_keywords}, cooldown={decision.cooldown_minutes}"
)
async with self._state.get_lock(group_id):
if decision.observation:
self._state.set_observation(group_id, decision.observation)
if decision.parse_failed:
logger.warning(f"Iris Reply: decision parse failed for group {group_id}")
async with self._state.get_lock(group_id):
self._state.record_skip_reply(group_id)
await self._state.save_dirty(self._kv_save)
self._release_decision_ticket(group_id, ticket_id)
event.stop_event()
return True