Files
dela 0335d572de
ci / go (push) Waiting to run
ci / go-db (agent) (push) Waiting to run
ci / go-db (config) (push) Waiting to run
ci / go-db (db) (push) Waiting to run
ci / go-db (evidence) (push) Waiting to run
ci / go-db (llmrec) (push) Waiting to run
ci / go-db (server) (push) Waiting to run
detections / detections (push) Waiting to run
web / web (push) Waiting to run
docs / links (push) Canceled after 0s
First Commit
2026-10-09 08:38:16 +08:00

382 lines
12 KiB
Go

package server
import (
"context"
"encoding/json"
"fmt"
"log"
"strconv"
"strings"
"time"
"github.com/Autumn-27/artex/db"
)
// Scheduler drives P3 triggers: on each tick it fires due interval triggers and
// scans for new findings / newly-met goals (any task) to fire event triggers.
// Each fire enqueues a NEW conversation on the agent's per-agent trigger queue
// (StartTriggeredRun): fires for the same agent run one at a time in FIFO order,
// distinct agents still run concurrently. State is persisted (per-trigger last_fire
// + finding watermark + fired-goal set) so a restart resumes without double-firing.
// Triggers only attach to CUSTOM agents.
type Scheduler struct {
s *Server
pg *db.DB
tick time.Duration
}
const (
schedKeyLastFinding = "last_finding_id" // watermark: max finding node id fired for
schedKeyFiredGoals = "fired_goals" // JSON array of goal node ids already fired
schedKeyLastTimeout = "last_timeout_id" // watermark: max task id fired for task timeout
schedKeyLastToolCall = "last_toolcall_id" // watermark: max activity id fired for tool call
schedKeyLastTaskCreate = "last_taskcreate_id" // watermark: max task id fired for task create
)
func newScheduler(s *Server) *Scheduler {
return &Scheduler{s: s, pg: s.m.pg, tick: 5 * time.Second}
}
// Run loops until ctx is done, ticking the scheduler. Started once from server New.
func (sc *Scheduler) Run(ctx context.Context) {
if sc.pg == nil {
return
}
sc.init()
t := time.NewTicker(sc.tick)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-t.C:
sc.s.reconcileConcurrency() // 并发上限:有空位就把排队任务补位启动
sc.step()
}
}
}
// init seeds the watermarks on first run so pre-existing findings/goals don't all
// fire at once — only events created AFTER the scheduler first starts count.
func (sc *Scheduler) init() {
if sc.mustState(schedKeyLastFinding) == "" {
var maxID int64
if evs, err := sc.pg.NewFindingsSince(0); err == nil {
for _, e := range evs {
if e.NodeID > maxID {
maxID = e.NodeID
}
}
}
_ = sc.pg.SetSchedState(schedKeyLastFinding, strconv.FormatInt(maxID, 10))
}
if sc.mustState(schedKeyFiredGoals) == "" {
set := map[int64]bool{}
if evs, err := sc.pg.MetGoals(); err == nil {
for _, e := range evs {
set[e.NodeID] = true
}
}
sc.saveFiredGoalSet(set)
}
if sc.mustState(schedKeyLastTimeout) == "" {
var maxID int64
if evs, err := sc.pg.TimedOutTasksSince(0); err == nil {
for _, e := range evs {
if e.NodeID > maxID {
maxID = e.NodeID
}
}
}
_ = sc.pg.SetSchedState(schedKeyLastTimeout, strconv.FormatInt(maxID, 10))
}
if sc.mustState(schedKeyLastToolCall) == "" {
var maxID int64
if evs, err := sc.pg.NewToolCallsSince(0); err == nil {
for _, e := range evs {
if e.NodeID > maxID {
maxID = e.NodeID
}
}
}
_ = sc.pg.SetSchedState(schedKeyLastToolCall, strconv.FormatInt(maxID, 10))
}
if sc.mustState(schedKeyLastTaskCreate) == "" {
var maxID int64
if evs, err := sc.pg.NewTasksSince(0); err == nil {
for _, e := range evs {
if e.NodeID > maxID {
maxID = e.NodeID
}
}
}
_ = sc.pg.SetSchedState(schedKeyLastTaskCreate, strconv.FormatInt(maxID, 10))
}
}
func (sc *Scheduler) step() {
triggers, err := sc.pg.ListEnabledTriggers()
if err != nil {
return
}
if len(triggers) == 0 {
return
}
sc.fireIntervals(triggers)
sc.fireFindings(triggers)
sc.fireGoals(triggers)
sc.fireTaskTimeouts(triggers)
sc.fireToolCalls(triggers)
sc.fireTaskCreates(triggers)
}
// fireIntervals fires triggers whose interval has elapsed since last_fire.
func (sc *Scheduler) fireIntervals(triggers []*db.AgentTrigger) {
now := time.Now()
for _, tr := range triggers {
if tr.IntervalSec <= 0 {
continue
}
due := tr.LastFire == nil || now.Sub(*tr.LastFire) >= time.Duration(tr.IntervalSec)*time.Second
if !due {
continue
}
_ = sc.pg.TouchTriggerFire(tr.ID)
// [F10/F19 경계 판정 · 두뇌 입력 보존] 이 파일의 트리거 합성부(정기·finding·목표·초과·
// 생성·도구 fire* 메서드)가 만드는 `【本次…触发】…` 컨텍스트 블록과 trunc 의 `…(已截断…)`
// 절단 표시는 번역하지 않는다. StartTriggeredRun 의 message 인자로 들어가
// finalTriggerMessage → ca.Chat 의 에이전트 입력(두뇌)이자 전사에 노출되는 이중 용도
// 문자열이다(BRIEF 경계 #1·F16 동형, conversations.go taskContextHeader 보존 주석 참조).
// title 인자(`定时触发 · …` 등)는 ca.Chat 에 안 가는 표시 전용이나, 본문이 중국어로
// 고정되어 제목만 바꾸면 대화 안에서 언어가 섞이므로(F7·F8·F12) 함께 보존한다.
ctx := "\n\n【本次为定时触发】" + now.Format(" 2006-01-02 15:04:05 MST")
sc.s.StartTriggeredRun(tr.AgentKey, fmt.Sprintf("定时触发 · %s", now.Format("15:04")), tr.IntervalMessage+ctx, 0, false, "", "")
}
}
// fireFindings fires on_finding triggers for findings above the persisted
// watermark (monotonic node id → no double-fire across restarts).
func (sc *Scheduler) fireFindings(triggers []*db.AgentTrigger) {
var want []*db.AgentTrigger
for _, tr := range triggers {
if tr.OnFinding {
want = append(want, tr)
}
}
last, _ := strconv.ParseInt(sc.mustState(schedKeyLastFinding), 10, 64)
events, err := sc.pg.NewFindingsSince(last)
if err != nil || len(events) == 0 {
return
}
// Advance the watermark whether or not any on_finding trigger is active: a
// finding fires only for triggers live at the moment it appears. Otherwise a
// trigger enabled later would replay the entire historical backlog at once.
maxID := last
for _, e := range events {
if e.NodeID > maxID {
maxID = e.NodeID
}
if len(want) == 0 {
continue
}
msgCtx := fmt.Sprintf("\n\n【本次由任务发现 finding 触发】\n发现: [%s/%s] %s",
e.VulnClass, e.Severity, e.Summary)
for _, tr := range want {
sc.s.StartTriggeredRun(tr.AgentKey, fmt.Sprintf("finding 触发 · task#%d", e.TaskID), tr.FindingMessage+msgCtx, e.TaskID, true, e.TaskDesc, e.TaskGoal)
}
}
_ = sc.pg.SetSchedState(schedKeyLastFinding, strconv.FormatInt(maxID, 10))
}
// fireGoals fires on_goal_met triggers for met goals not yet in the fired set.
func (sc *Scheduler) fireGoals(triggers []*db.AgentTrigger) {
var want []*db.AgentTrigger
for _, tr := range triggers {
if tr.OnGoalMet {
want = append(want, tr)
}
}
events, err := sc.pg.MetGoals()
if err != nil || len(events) == 0 {
return
}
// Mark goals as consumed whether or not a trigger is active, so enabling an
// on_goal_met trigger later doesn't replay every already-met goal.
fired := sc.firedGoalSet()
changed := false
for _, e := range events {
if fired[e.NodeID] {
continue
}
fired[e.NodeID] = true
changed = true
if len(want) == 0 {
continue
}
msgCtx := fmt.Sprintf("\n\n【本次由任务完成目标触发】\n达成目标: %s", e.Summary)
for _, tr := range want {
sc.s.StartTriggeredRun(tr.AgentKey, fmt.Sprintf("目标触发 · task#%d", e.TaskID), tr.GoalMessage+msgCtx, e.TaskID, true, e.TaskDesc, e.TaskGoal)
}
}
if changed {
sc.saveFiredGoalSet(fired)
}
}
// fireTaskTimeouts fires on_task_timeout triggers for tasks that newly reached
// status='timeout' above the persisted watermark (task id → no double-fire).
func (sc *Scheduler) fireTaskTimeouts(triggers []*db.AgentTrigger) {
var want []*db.AgentTrigger
for _, tr := range triggers {
if tr.OnTaskTimeout {
want = append(want, tr)
}
}
last, _ := strconv.ParseInt(sc.mustState(schedKeyLastTimeout), 10, 64)
events, err := sc.pg.TimedOutTasksSince(last)
if err != nil || len(events) == 0 {
return
}
// Advance the watermark even with no active trigger — see fireFindings.
maxID := last
for _, e := range events {
if e.NodeID > maxID {
maxID = e.NodeID
}
if len(want) == 0 {
continue
}
msgCtx := "\n\n【本次由任务超时触发】"
for _, tr := range want {
sc.s.StartTriggeredRun(tr.AgentKey, fmt.Sprintf("超时触发 · task#%d", e.TaskID), tr.TaskTimeoutMessage+msgCtx, e.TaskID, true, e.TaskDesc, e.TaskGoal)
}
}
_ = sc.pg.SetSchedState(schedKeyLastTimeout, strconv.FormatInt(maxID, 10))
}
// fireTaskCreates fires on_task_create triggers for tasks newly created above the
// persisted watermark (task id → no double-fire across restarts).
func (sc *Scheduler) fireTaskCreates(triggers []*db.AgentTrigger) {
var want []*db.AgentTrigger
for _, tr := range triggers {
if tr.OnTaskCreate {
want = append(want, tr)
}
}
last, _ := strconv.ParseInt(sc.mustState(schedKeyLastTaskCreate), 10, 64)
events, err := sc.pg.NewTasksSince(last)
if err != nil || len(events) == 0 {
return
}
// Advance the watermark even with no active trigger — see fireFindings.
maxID := last
for _, e := range events {
if e.NodeID > maxID {
maxID = e.NodeID
}
if len(want) == 0 {
continue
}
msgCtx := "\n\n【本次由任务创建触发】"
for _, tr := range want {
sc.s.StartTriggeredRun(tr.AgentKey, fmt.Sprintf("任务创建触发 · task#%d", e.TaskID), tr.TaskCreateMessage+msgCtx, e.TaskID, true, e.TaskDesc, e.TaskGoal)
}
}
_ = sc.pg.SetSchedState(schedKeyLastTaskCreate, strconv.FormatInt(maxID, 10))
}
// fireToolCalls fires on_tool_call triggers for tool calls (tool_result rows) above
// the persisted watermark whose tool name is in the trigger's selected set. The fire
// message carries the task id/desc/goal + tool name + (truncated) input & output.
func (sc *Scheduler) fireToolCalls(triggers []*db.AgentTrigger) {
var want []*db.AgentTrigger
for _, tr := range triggers {
if tr.OnToolCall && len(tr.ToolNames) > 0 {
want = append(want, tr)
}
}
last, _ := strconv.ParseInt(sc.mustState(schedKeyLastToolCall), 10, 64)
events, err := sc.pg.NewToolCallsSince(last)
if err != nil || len(events) == 0 {
return
}
// Advance the watermark even with no active trigger — see fireFindings.
maxID := last
for _, e := range events {
if e.NodeID > maxID {
maxID = e.NodeID
}
if len(want) == 0 {
continue
}
// A failed finding write has no committed finding to report. Keep other
// tool-error triggers available for user-defined automation.
if e.ToolIsErr && e.Tool == "report_finding" {
continue
}
errTag := ""
if e.ToolIsErr {
errTag = "[error] "
}
msgCtx := fmt.Sprintf("\n\n【本次由工具调用触发】\n工具: %s\n入参: %s\n返回: %s%s",
e.Tool, trunc(e.ToolInput, 1500), errTag, trunc(e.ToolOutput, 1500))
for _, tr := range want {
if !containsFold(tr.ToolNames, e.Tool) {
continue
}
sc.s.StartTriggeredRun(tr.AgentKey, fmt.Sprintf("工具触发 · %s · task#%d", e.Tool, e.TaskID), tr.ToolCallMessage+msgCtx, e.TaskID, true, e.TaskDesc, e.TaskGoal)
}
}
_ = sc.pg.SetSchedState(schedKeyLastToolCall, strconv.FormatInt(maxID, 10))
}
// containsFold reports whether name is in set (case-insensitive).
func containsFold(set []string, name string) bool {
for _, s := range set {
if strings.EqualFold(s, name) {
return true
}
}
return false
}
// trunc caps s to max runes, appending an ellipsis + original length when cut.
func trunc(s string, max int) string {
r := []rune(s)
if len(r) <= max {
return s
}
return string(r[:max]) + fmt.Sprintf("…(已截断,共 %d 字)", len(r))
}
func (sc *Scheduler) mustState(key string) string {
v, _ := sc.pg.GetSchedState(key)
return v
}
func (sc *Scheduler) firedGoalSet() map[int64]bool {
out := map[int64]bool{}
raw := sc.mustState(schedKeyFiredGoals)
if raw == "" {
return out
}
var ids []int64
if json.Unmarshal([]byte(raw), &ids) == nil {
for _, id := range ids {
out[id] = true
}
}
return out
}
func (sc *Scheduler) saveFiredGoalSet(set map[int64]bool) {
ids := make([]int64, 0, len(set))
for id := range set {
ids = append(ids, id)
}
b, _ := json.Marshal(ids)
if err := sc.pg.SetSchedState(schedKeyFiredGoals, string(b)); err != nil {
log.Printf("[scheduler] save fired goals failed: %v", err)
}
}