server

package
v0.3.15 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 7, 2026 License: AGPL-3.0 Imports: 67 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

View Source
var BuildVersion = "dev"

BuildVersion is the backend application version, injected from cmd/artex at startup (which in turn gets it from -ldflags "-X main.version=<tag>"). Defaults to "dev" for local builds. Exposed to the frontend via GET /api/health.

Functions

func RestartRequested

func RestartRequested() <-chan struct{}

RestartRequested 返回一个在"请退出并让守护进程重新拉起我"时关闭的 channel。

func SetBootUpdateState

func SetBootUpdateState(st selfupdate.State)

SetBootUpdateState 由 main 在启动时调用一次。

func StartLogCapture

func StartLogCapture()

StartLogCapture redirects the standard log package through the in-memory sink (still writing to stderr). Call once at startup, as early as possible.

Types

type ActivityDTO

type ActivityDTO struct {
	Seq          int64           `json:"seq"`                 // db Activity.ID
	IntentID     string          `json:"intent_id,omitempty"` // db NodeID
	Worker       string          `json:"worker"`
	TS           string          `json:"ts"` // db CreatedAt
	Kind         string          `json:"kind"`
	Tool         string          `json:"tool,omitempty"`
	ToolUseID    string          `json:"tool_use_id,omitempty"`
	IsError      bool            `json:"is_error"`
	Summary      string          `json:"summary"`
	Detail       string          `json:"detail,omitempty"`
	Metadata     json.RawMessage `json:"metadata,omitempty"`
	SourceTaskID string          `json:"source_task_id,omitempty"`
	Inherited    bool            `json:"inherited,omitempty"`
	MainSeg      *int            `json:"main_seg,omitempty"` // main-agent conversation segment (nil for non-mainagent rows)
	// token usage (set only on kind='result'); used for per-session token totals.
	InputTokens      *int `json:"input_tokens,omitempty"`
	OutputTokens     *int `json:"output_tokens,omitempty"`
	CacheReadTokens  *int `json:"cache_read_tokens,omitempty"`
	CacheWriteTokens *int `json:"cache_write_tokens,omitempty"`
}

type AgentDTO

type AgentDTO struct {
	ID               string `json:"id"`
	Key              string `json:"key"`
	Name             string `json:"name"`
	Description      string `json:"description"`
	Role             string `json:"role"`
	Builtin          bool   `json:"builtin"`
	Enabled          bool   `json:"enabled"`
	MaxTurns         int    `json:"max_turns"`
	RunSecs          int    `json:"run_seconds"`
	WebSearch        bool   `json:"web_search"`
	InteractiveShell bool   `json:"interactive_shell"`
	LLMProfileID     *int64 `json:"llm_profile_id"` // 绑定的 LLM 配置;null=跟随任务/全局
	// P3 触发后处理策略(仅自定义 agent 有意义)。
	TriggerRunMode     string `json:"trigger_run_mode"`
	TriggerMergeMode   string `json:"trigger_merge_mode"`
	TriggerMaxParallel int    `json:"trigger_max_parallel"`
	// binding counts (populated only by the list endpoint) — shown on agent cards.
	McpCount   int `json:"mcp_count"`
	SkillCount int `json:"skill_count"`
	ToolCount  int `json:"tool_count"`
}

type Broadcaster

type Broadcaster struct {
	// contains filtered or unexported fields
}

Broadcaster is a per-task in-process pub/sub for live activity events. The engine publishes each appended activity at its single emit point; SSE handlers subscribe per task. Storage (activity table) and the live stream come from the same Publish call, so they never diverge.

func NewBroadcaster

func NewBroadcaster() *Broadcaster

func (*Broadcaster) Publish

func (b *Broadcaster) Publish(task string, a db.Activity)

Publish fans an activity out to all subscribers of a task. Non-blocking: if a subscriber's buffer is full the event is dropped — the client reconnects with its last seq cursor and catches up the gap from the DB, so liveness never stalls the engine.

func (*Broadcaster) Subscribe

func (b *Broadcaster) Subscribe(task string) (<-chan db.Activity, func())

Subscribe returns a buffered channel of activities for a task plus an unsubscribe func the caller must invoke (defer) to release it.

type ConstraintDTO

type ConstraintDTO struct {
	ID     string `json:"id"`
	Kind   string `json:"kind"` // allow | deny
	Text   string `json:"text"`
	Origin string `json:"origin,omitempty"`
	TS     string `json:"ts,omitempty"`
}

ConstraintDTO is one operation constraint (allow/deny) for the 总览「约束管理」UI.

type CoverageAssetRefDTO

type CoverageAssetRefDTO struct {
	ID           int64  `json:"id"`
	Kind         string `json:"kind"`
	State        string `json:"state"`
	Summary      string `json:"summary"`
	SourceTaskID string `json:"source_task_id,omitempty"`
	Inherited    bool   `json:"inherited,omitempty"`
}

type DeleteTaskOptions

type DeleteTaskOptions struct {
	DeleteAssets     bool `json:"delete_assets"`
	DeleteTraffic    bool `json:"delete_traffic"`
	DeleteFiles      bool `json:"delete_files"`
	DeleteFindings   bool `json:"delete_findings"`
	DeleteLLMRecords bool `json:"delete_llm_records"`
}

DeleteTaskOptions controls cleanup of data stored outside the task's own exploration graph. All options default to false for backward compatibility.

type DeleteTaskResult

type DeleteTaskResult struct {
	Deleted           string `json:"deleted"`
	AssetsDeleted     int64  `json:"assets_deleted"`
	AssetsDetached    int64  `json:"assets_detached"`
	TrafficDeleted    int64  `json:"traffic_deleted"`
	FilesDeleted      bool   `json:"files_deleted"`
	FindingsDeleted   int64  `json:"findings_deleted"`
	LLMRecordsDeleted int64  `json:"llm_records_deleted"`
	CleanupWarning    string `json:"cleanup_warning,omitempty"`
}

DeleteTaskResult makes destructive cleanup auditable to API callers.

type EdgeDTO

type EdgeDTO struct {
	Src string `json:"src"`
	Dst string `json:"dst"`
	Rel string `json:"rel"`
}

type Engine

type Engine struct {
	// contains filtered or unexported fields
}

Engine drives the event-driven exploration loop with real LLM agents (docs §4.3/§4.4): on asset/exploration-graph change (debounced) it wakes the planner, which reads the route, queries assets, judges goals and emits intents; N concurrent work agents claim intents and execute them. There is no simulation mode — an LLM provider is required. The planner/worker can be (re)installed at runtime (LLM configured from the UI); the loops always run but idle until an LLM is set.

func NewEngine

func NewEngine(m *Manager) *Engine

func (*Engine) AbortDelete

func (e *Engine) AbortDelete(taskID string, keepPaused bool)

func (*Engine) ActiveLLMCalls

func (e *Engine) ActiveLLMCalls(taskID string) int64

func (*Engine) BeginDelete

func (e *Engine) BeginDelete(taskID string) bool

BeginDelete installs an execution barrier before task data/files are removed. The temporary pause is not a user pause. The server serializes this transition with lifecycle admission and tells AbortDelete whether the persisted task is paused/queued if cleanup fails.

func (*Engine) BeginLLMCall

func (e *Engine) BeginLLMCall(taskID string)

BeginLLMCall/EndLLMCall track actual provider calls separately from the scheduler's task-operation counter. A task can have live loops while all of them are waiting for a trigger; that state must remain idle in the UI.

func (*Engine) Broadcaster

func (e *Engine) Broadcaster() *Broadcaster

Broadcaster exposes the engine's live activity pub/sub (used by the SSE handler).

func (*Engine) ControlWork

func (e *Engine) ControlWork(ctx context.Context, intentID int64, action string) error

ControlWork requests a user-visible pause or cancellation and waits until the worker has fully stopped writing. Cancellation cleanup is performed by the API handler after this returns; pause state is committed by runWorkerStep itself.

func (*Engine) EndLLMCall

func (e *Engine) EndLLMCall(taskID string)

func (*Engine) IsDeleting

func (e *Engine) IsDeleting(taskID string) bool

func (*Engine) IsPaused

func (e *Engine) IsPaused(taskID string) bool

IsPaused reports whether a task is user-paused.

func (*Engine) KillWork

func (e *Engine) KillWork(intentID int64) error

KillWork cancels the in-flight work running intentID (planner's kill_work tool). The work's agent-core session honors ctx cancellation and aborts promptly.

func (*Engine) LastActivity

func (e *Engine) LastActivity(taskID string) int64

LastActivity returns the unix time of the last planner/worker activity for a task (0 if none yet).

func (*Engine) Pause

func (e *Engine) Pause(taskID string, cause error)

Pause stops a task: marks it paused AND cancels any in-flight planner/worker run for it (a long worker.Execute would otherwise keep going until it finishes).

func (*Engine) Ready

func (e *Engine) Ready() bool

Ready reports whether a global LLM provider is configured (via the readiness predicate wired at startup).

func (*Engine) ReadyFor

func (e *Engine) ReadyFor(t *Task) bool

ReadyFor reports whether a specific task can resolve a planner/worker pair. An explicit task profile chain can be runnable even when no global default provider is configured, so task status must not rely on Ready alone.

func (*Engine) Resume

func (e *Engine) Resume(t *Task)

Resume un-pauses a task and nudges a fresh planning round. The next exec under it gets a fresh (uncancelled) context.

func (*Engine) Run

func (e *Engine) Run(ctx context.Context, t *Task)

Run starts the planner loop + N worker loops for a task. The loops always run but no-op until an LLM is configured (so a task created while idle picks up automatically once LLM is set from the UI).

func (*Engine) SetAgentResolver

func (e *Engine) SetAgentResolver(fn func(t *Task) (*agent.Planner, *agent.Worker))

SetAgentResolver installs a per-task planner/worker resolver (wired by the server). Called once at startup before any task loop runs, so no lock is needed on reads.

func (*Engine) SetAuthoritativeAgentResolver

func (e *Engine) SetAuthoritativeAgentResolver(fn func(t *Task) (*agent.Planner, *agent.Worker))

SetAuthoritativeAgentResolver installs a resolver whose nil result must not fall through to the global provider. Task-level failover chains use this so a fully exhausted chain cannot silently bypass its configured boundary.

func (*Engine) SetReadiness

func (e *Engine) SetReadiness(fn func() bool)

SetReadiness wires the global "an LLM provider is configured" predicate (read by Ready() / the llm_configured indicator). Called once at startup.

func (*Engine) Started

func (e *Engine) Started(taskID string) bool

Started reports whether the engine loops are running for a task.

func (*Engine) SteerWork

func (e *Engine) SteerWork(intentID int64, msg string) error

SteerWork queues a mid-run course-correction for the work running intentID (the planner's steer_work tool). The worker delivers it before its next tool call and re-plans — no kill. Errors if no work is currently running that intent.

func (*Engine) StopTask

func (e *Engine) StopTask(taskID string)

StopTask permanently stops every long-lived goroutine and removes all Engine state for a successfully deleted task. The delete barrier remains installed until cleanup finishes, so no new task operation can race the teardown.

type FindingAssetDTO

type FindingAssetDTO struct {
	ID    string `json:"id"`
	Type  string `json:"type"`
	Label string `json:"label"`
}

FindingAssetDTO is one asset a finding is anchored to, pre-labelled for display.

type FindingDTO

type FindingDTO struct {
	TrafficCount          int                        `json:"traffic_count"`
	EvidenceVersion       int64                      `json:"evidence_version"`
	ReportEvidenceVersion int64                      `json:"report_evidence_version"`
	ReportStale           bool                       `json:"report_stale"`
	TrafficBindings       []db.FindingTrafficBinding `json:"traffic_bindings,omitempty"`

	ID        string `json:"id"`
	FindingID string `json:"finding_id,omitempty"` // standalone findings-table id — the handle for status updates
	VulnClass string `json:"vulnclass"`
	Name      string `json:"name,omitempty"` // 漏洞名称;为空时前端回退展示 vulnclass
	Severity  string `json:"severity"`       // critical | high | medium | low
	Status    string `json:"status"`         // pending | in_progress | confirmed | resolved | fixed | false_positive | ignored | duplicate | risk_accepted
	Summary   string `json:"summary"`
	Evidence  string `json:"evidence"`
	Report    string `json:"report,omitempty"` // 详细报告(Markdown);仅详情接口返回,列表为空

	IntentID        string            `json:"intent_id,omitempty"`
	ParamID         string            `json:"param_id,omitempty"`
	TaskID          string            `json:"task_id,omitempty"`
	TaskDescription string            `json:"task_description,omitempty"`
	SourceTaskID    string            `json:"source_task_id,omitempty"`
	Inherited       bool              `json:"inherited,omitempty"`
	Assets          []FindingAssetDTO `json:"assets,omitempty"`
	TS              string            `json:"ts"`
}

type FindingSeverityDTO added in v0.3.14

type FindingSeverityDTO struct {
	Critical int `json:"critical"`
	High     int `json:"high"`
	Medium   int `json:"medium"`
	Low      int `json:"low"`
}

FindingSeverityDTO 是任务列表里按严重度分档的漏洞计数(严重/高/中/低)。

type GoalDTO

type GoalDTO struct {
	ID        string `json:"id"`
	Text      string `json:"text"`
	VulnClass string `json:"vulnclass,omitempty"`
	State     string `json:"state"`
	Origin    string `json:"origin,omitempty"`
	TS        string `json:"ts"`
}

GoalDTO is a goal node with its payload unpacked into text/vulnclass — the shape the 总览「目标管理」UI works with (vs TaskNodeDTO which carries raw payload JSON).

type LLMPoolMemberStatus

type LLMPoolMemberStatus struct {
	ProfileID string `json:"profile_id"`
	Name      string `json:"name"`
	Model     string `json:"model"`
	Format    string `json:"format"`
	Priority  int    `json:"priority"`
	Active    bool   `json:"active"`   // is the globally-activated profile
	Excluded  bool   `json:"excluded"` // pool_exclude — not a failover target
	// Health: state is "ok" | "degraded" (failing but not tripped) | "tripped".
	State        string `json:"state"`
	Fails        int    `json:"fails"`
	Trips        int    `json:"trips"`
	CooldownSecs int    `json:"cooldown_secs"` // remaining cooling-off seconds; 0 = none
	LastError    string `json:"last_error,omitempty"`
	LastAt       string `json:"last_at,omitempty"`
}

LLMPoolMemberStatus is one chain entry as shown in the UI.

type LLMProfileDTO

type LLMProfileDTO struct {
	ID              string  `json:"id"`
	Name            string  `json:"name"`
	Format          string  `json:"format"`
	BaseURL         string  `json:"base_url,omitempty"`
	Proxy           string  `json:"proxy,omitempty"`
	Model           string  `json:"model"`
	APIKeyHint      string  `json:"api_key_hint,omitempty"`
	RatePerSecond   float64 `json:"rate_per_second"`
	RatePerMinute   float64 `json:"rate_per_minute"`
	ContextWindowK  int     `json:"context_window_k"`
	ThinkingType    string  `json:"thinking_type"`
	ReasoningEffort string  `json:"reasoning_effort"`
	IsDefault       bool    `json:"is_default"`
	// 轮询(故障转移)参数:priority 越大越先被选中(激活配置恒为链首);
	// pool_exclude=true 则不作为故障转移目标,但仍可被 agent/任务显式绑定。
	Priority    int  `json:"priority"`
	PoolExclude bool `json:"pool_exclude"`
	// 收发模式:true=流式(SSE) | false=非流式。没有 omitempty —— false 必须出现在
	// 响应里,否则前端读不到「非流式」,开关会回落成默认的流式。
	Streaming bool `json:"streaming"`
	// 单次回复输出上限(0=不发送,由服务端默认值决定),以及它用哪个请求字段名
	// (”=max_tokens | 'max_completion_tokens',仅 openai 格式有意义)。
	MaxTokens      int    `json:"max_tokens"`
	MaxTokensField string `json:"max_tokens_field"`
	// 自定义会话头名:非空时每次请求带该 HTTP 头,头值=当前会话/意图的 session id。
	// ”=不发送。用于按 session-id 头做提示缓存/粘性路由的网关。
	SessionHeaderKey string `json:"session_header_key"`
	// 本配置对重试的覆盖(建连/空响应/同 provider 安全窗口)。每项 attempts:
	// 0=继承全局策略 | -1=关闭该层重试 | >0=次数;interval_ms: 0=用默认指数退避 |
	// >0=改用该固定毫秒间隔。全 0 = 完全跟随全局,即历史行为。
	Retry db.RetryOverride `json:"retry"`
}

type LogLine

type LogLine struct {
	Seq   int64  `json:"seq"`
	DBID  int64  `json:"db_id,omitempty"` // server_logs.id; 0 for pre-persistence entries
	TS    string `json:"ts"`
	Level string `json:"level"` // info | warn | error
	Tag   string `json:"tag"`   // the leading [tag] (pg / planner / activity / …), if any
	Text  string `json:"text"`
}

LogLine is one captured backend log entry exposed by the /api/logs endpoints.

type Manager

type Manager struct {
	// contains filtered or unexported fields
}

Manager owns the PostgreSQL data source (asset graph + every task's exploration graph + config) and the in-memory set of task handles.

func NewManager

func NewManager(dir, proxyAddr string) (*Manager, error)

NewManager connects to PostgreSQL and, if proxyAddr is non-empty, starts the traffic-recording proxy. PostgreSQL is required (it is the single data source).

func (*Manager) ActiveTask

func (m *Manager) ActiveTask() *Task

ActiveTask returns the currently active task (or nil).

func (*Manager) ApplyTaskAdmission

func (m *Manager) ApplyTaskAdmission(id, expectedStatus, status string, queued bool, mode string, preservePosition bool) error

ApplyTaskAdmission atomically commits the lifecycle fields controlled by the concurrency scheduler. Keeping status, paused and queue metadata in one UPDATE prevents a failed resume from leaving a task half-revived (for example running but still user-paused, or dequeued without an Engine start).

preservePosition applies only when the row is already queued. A repeated admission keeps its FIFO timestamp; a task that was explicitly paused and is now re-queued receives a fresh tail position.

func (*Manager) ApplyTaskPause

func (m *Manager) ApplyTaskPause(id string) error

ApplyTaskPause atomically removes a task from the admission queue and records the user pause. queue_mode is intentionally retained so resuming a never-run bootstrap task still performs goal decomposition, but the next enqueue receives a new queued_at timestamp and therefore moves to the FIFO tail.

func (*Manager) Assets

func (m *Manager) Assets() *pgdb.AssetStore

func (*Manager) Close

func (m *Manager) Close() error

func (*Manager) ConcurrencyLimit

func (m *Manager) ConcurrencyLimit() (enabled bool, limit int)

ConcurrencyLimit returns whether the simultaneous-running-task cap is enabled and its limit (default 5 when enabled but unset). limit is always >=1 when enabled.

func (*Manager) CreateTask

func (m *Manager) CreateTask(description, goal string, llmProfileID *int64, timeoutSeconds, planHeartbeatSeconds int) (*Task, error)

CreateTask creates a task + its exploration and makes it active. timeoutSeconds is the task-level wall-clock budget (0 = 不限时).

func (*Manager) CreateTaskWithOptions

func (m *Manager) CreateTaskWithOptions(description, goal string, opts pgdb.TaskCreateOptions) (*Task, error)

func (*Manager) DeleteCompanyWithAssets

func (m *Manager) DeleteCompanyWithAssets(id int64, deleteAssets bool) (int64, error)

DeleteCompanyWithAssets keeps the database cascade and live task handles in one manager-level critical section. This closes the gap where a task could commit its company scope immediately before registration and miss the post-delete in-memory sweep.

func (*Manager) DeleteTask

func (m *Manager) DeleteTask(id string, opts DeleteTaskOptions) (DeleteTaskResult, error)

DeleteTask removes a task and optionally its related global data. Traffic has no task-id column, so related exchanges are resolved by exact hosts from the task's asset rows. Files are staged before the database operation; traffic is staged while PostgreSQL excludes asset/anchor writers. Both are restored on a database failure and purged only after its commit.

func (*Manager) DeleteTaskCategory

func (m *Manager) DeleteTaskCategory(id int64) (bool, error)

DeleteTaskCategory moves all affected live tasks to the uncategorized bucket.

func (*Manager) DequeueTask

func (m *Manager) DequeueTask(id string, clearMode bool) error

DequeueTask removes the concurrency hold. clearMode=false is used when a user pauses a queued task so a later resume still knows whether bootstrap is needed.

func (*Manager) EnqueueTask

func (m *Manager) EnqueueTask(id, mode string) error

EnqueueTask persists the concurrency hold and syncs the in-memory handle.

func (*Manager) Enrich

func (m *Manager) Enrich() *enrich.Engine

Enrich returns the asset auto-completion engine (may be nil if init failed).

func (*Manager) GlobalProxy

func (m *Manager) GlobalProxy() string

GlobalProxy returns the configured global egress proxy URL (empty = direct).

func (*Manager) HostTools

func (m *Manager) HostTools() []actool.CoreTool

HostTools are runtime host-provided tools added to EVERY agent's base list (via ToolAugment); the tools table then filters them per-agent binding. Currently the traffic tools, gated by the global capture switch: empty when capture is off, so no agent gets traffic_search/traffic_get regardless of binding.

func (*Manager) LLMPoolBindFallback

func (m *Manager) LLMPoolBindFallback() bool

LLMPoolBindFallback reports whether an agent/task that is BOUND to a specific profile still falls back to the chain when that profile fails (默认关:绑定即 独占,失败即失败). Only meaningful while LLMPoolEnabled.

func (*Manager) LLMPoolEnabled

func (m *Manager) LLMPoolEnabled() bool

LLMPoolEnabled reports whether LLM failover ("轮询") is on (默认关; settings.llm_pool_enabled). Read when the provider chain is built (applyLLM), so a change requires a rebuild — putSettings does that.

func (*Manager) LLMRecordEnabled

func (m *Manager) LLMRecordEnabled() bool

LLMRecordEnabled reports whether LLM request/response recording is on (默认关;settings.llm_record). The recorder consults this per call, so the toggle takes effect immediately without rebuilding agents.

func (*Manager) List

func (m *Manager) List() []*Task

func (*Manager) LoadExisting

func (m *Manager) LoadExisting() []*Task

LoadExisting rebuilds in-memory task handles from the PG task registry.

func (*Manager) NoaCompactionEnabled

func (m *Manager) NoaCompactionEnabled() bool

NoaCompactionEnabled reports whether the experimental noa context-compression mechanism is on (默认关;settings.noa_compaction). Read per agent run via the injected resolver, so a toggle takes effect on the next run without rebuild.

func (*Manager) PG

func (m *Manager) PG() *pgdb.DB

func (*Manager) ProxyAddr

func (m *Manager) ProxyAddr() string

ProxyAddr returns the egress proxy address agents route target traffic through:

  • capture ON → the recording MITM proxy (which itself exits via the global proxy when one is set); agents also get its CA (see ProxyCACert).
  • capture OFF → the global egress proxy directly (empty CA — real target certs), or "" when no global proxy is set (direct, no recording).

So the global proxy takes effect in both modes: at the MITM's upstream when capturing, in the agent's own bash env / WebFetch when not.

func (*Manager) ProxyCACert

func (m *Manager) ProxyCACert() string

ProxyCACert returns the CA cert path agents must trust to verify HTTPS through the egress proxy. Non-empty ONLY when traffic capture is on (the MITM re-signs certs): the global proxy used directly (capture off) is a plain forwarder that preserves real target certs, so no custom CA is needed there. Its emptiness is also the worker's "recording off" signal (see workerSystem).

func (*Manager) RenameTaskCategory

func (m *Manager) RenameTaskCategory(id int64, name string) (*pgdb.TaskCategory, error)

RenameTaskCategory persists a category name and refreshes every live task DTO that references it. taskStateMu keeps this ordered with task reassignment.

func (*Manager) ReplaceTaskLLMProfiles

func (m *Manager) ReplaceTaskLLMProfiles(id string, profileIDs []int64, activeProfileID int64) (int64, error)

ReplaceTaskLLMProfiles resets a task's ordered provider chain and mirrors the committed state onto the live task handle. Terminal tasks are editable too — their 主 Agent 对话 keeps running on the chain after the task finishes.

func (*Manager) ResolveTask

func (m *Manager) ResolveTask(id string) *Task

func (*Manager) SetActive

func (m *Manager) SetActive(id string) bool

SetActive switches the active task. Returns false if the id is unknown.

func (*Manager) SetConcurrency

func (m *Manager) SetConcurrency(enabled bool, limit int) error

SetConcurrency persists the running-task concurrency cap. limit<1 is clamped to 1.

func (*Manager) SetGlobalProxy

func (m *Manager) SetGlobalProxy(raw string) error

SetGlobalProxy validates, persists and applies the global egress proxy (http/https/socks5, optional user:pass; empty = direct). It updates the MITM's upstream immediately; callers must rebuild agents (applyLLM) afterwards so the capture-off path (bash env / WebFetch) picks up the change too.

func (*Manager) SetLLMPoolBindFallback

func (m *Manager) SetLLMPoolBindFallback(on bool) error

SetLLMPoolBindFallback persists the bound-profile fallback toggle. Callers rebuild agents (applyLLM) afterwards.

func (*Manager) SetLLMPoolEnabled

func (m *Manager) SetLLMPoolEnabled(on bool) error

SetLLMPoolEnabled persists the failover toggle. Callers rebuild agents (applyLLM) afterwards so it takes effect.

func (*Manager) SetLLMRecordEnabled

func (m *Manager) SetLLMRecordEnabled(on bool) error

SetLLMRecordEnabled persists and applies the LLM-record toggle. Effective at once — no applyLLM needed, since the recorder reads the flag on every call.

func (*Manager) SetNoaCompaction

func (m *Manager) SetNoaCompaction(on bool) error

SetNoaCompaction persists the noa toggle. Effective on the next agent run — the resolver reads it per run, so no rebuild is needed.

func (*Manager) SetTaskCategory

func (m *Manager) SetTaskCategory(taskID string, categoryID *int64) (*pgdb.TaskCategory, error)

SetTaskCategory updates one live task without interrupting its runtime.

func (*Manager) SetTaskPaused

func (m *Manager) SetTaskPaused(id string, paused bool) error

SetTaskPaused persists a task's paused state.

func (*Manager) SetTaskStatus

func (m *Manager) SetTaskStatus(id, status string) error

SetTaskStatus persists a task's lifecycle status (e.g. "done") and reflects it on the in-memory handle so the derived DTO status shows it without a reload.

func (*Manager) SetTaskStatusGuarded

func (m *Manager) SetTaskStatusGuarded(id, status string) (won bool, err error)

SetTaskStatusGuarded sets a TERMINAL status only if the task isn't already terminal (resolves the completed↔timeout race — first terminal writer wins). Reflects the won status on the live handle. won=false means another terminal already stuck.

func (*Manager) SetTasksCategory

func (m *Manager) SetTasksCategory(taskIDs []string, categoryID *int64) (map[string]bool, *pgdb.TaskCategory, error)

SetTasksCategory applies one category change to several tasks at once. The database write and the in-memory refresh share taskStateMu, so a concurrent single-task update cannot interleave and leave a live DTO stale. The returned set holds the ids that were actually moved; callers report the rest as gone.

func (*Manager) SetTrafficEnabled

func (m *Manager) SetTrafficEnabled(on bool) error

SetTrafficEnabled persists and applies the traffic-capture toggle. Callers must rebuild the agents (applyLLM) afterwards so the new proxy/tools/prompt take hold.

func (*Manager) SetWebSearch

func (m *Manager) SetWebSearch(on bool, backend string, braveKey, tavilyKey, proxy *string) error

SetWebSearch persists and applies the web-search settings. braveKey, tavilyKey, and proxy are each left untouched when nil (so toggling the switch doesn't wipe a saved key/proxy; pass a pointer to "" to clear). Callers must rebuild agents (applyLLM) afterwards so the settings take effect.

func (*Manager) SetWorkers

func (m *Manager) SetWorkers(n int) error

SetWorkers persists the concurrent work-agent count. Values <=0 are rejected.

func (*Manager) StampTaskFirstRun

func (m *Manager) StampTaskFirstRun(id string) (int64, error)

StampTaskFirstRun stamps first_run_at + deadline_at on the first real run (idempotent in DB) and mirrors deadline_at on the live handle. Returns the deadline unix (0 = 不限).

func (*Manager) Task

func (m *Manager) Task(id string) (*Task, bool)

func (*Manager) TaskStatus

func (m *Manager) TaskStatus(id string) string

TaskStatus returns a task's current in-memory status (empty if unknown).

func (*Manager) Traffic

func (m *Manager) Traffic() *traffic.Traffic

func (*Manager) TrafficEnabled

func (m *Manager) TrafficEnabled() bool

TrafficEnabled reports whether traffic capture is on (default off). When off, no proxy/traffic tools/prompt are injected into agents (nothing is recorded).

func (*Manager) UpdateTaskMetadata

func (m *Manager) UpdateTaskMetadata(taskID string, patch pgdb.TaskPatch) (*Task, error)

UpdateTaskMetadata changes list-only task metadata without interrupting any planner, main-agent, or worker call.

func (*Manager) WebSearch

func (m *Manager) WebSearch() (on bool, backend, braveKey, tavilyKey, proxy string)

WebSearch returns the current web-search config: whether it is enabled, the backend ("ddgs" | "brave-free" | "tavily"), the Brave API key, the Tavily API key (each empty unless set), and the dedicated egress proxy (empty = direct).

func (*Manager) WebSearchOpts

func (m *Manager) WebSearchOpts() agent.WebSearchOpts

WebSearchOpts returns the config as the agent-package struct the server pushes into each agent. Disabled when off, or when a keyed backend is selected without its key (so a half-configured backend never silently drops the tool at session build).

func (*Manager) Workers

func (m *Manager) Workers() int

Workers returns the configured concurrent work-agent count (default 3). Read per-task at engine.Run, so a change applies to tasks started afterwards.

type Notifier added in v0.3.15

type Notifier struct {
	// contains filtered or unexported fields
}

Notifier 是漏洞推送的投递引擎。

与 Scheduler 并列,作为独立 goroutine 运行(见 server.New)。刻意不复用 Scheduler 的 tick:推送的实时性要求(3 秒)与触发器的业务节奏不同, 且两者的失败互不牵连——推送卡住不该影响 agent 触发。

func (*Notifier) Run added in v0.3.15

func (n *Notifier) Run(ctx context.Context)

Run 循环直到 ctx 结束。由 server.New 启动一次。

type Scheduler

type Scheduler struct {
	// contains filtered or unexported fields
}

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.

func (*Scheduler) Run

func (sc *Scheduler) Run(ctx context.Context)

Run loops until ctx is done, ticking the scheduler. Started once from server New.

type Server

type Server struct {
	// contains filtered or unexported fields
}

Server exposes the ARTEX backend over a JSON HTTP API for the shadcn/ui frontend.

func New

func New(ctx context.Context, m *Manager, skillDir string, dataDir string, keyDir string) *Server

func (*Server) Handler

func (s *Server) Handler() http.Handler

func (*Server) StartTriggeredRun

func (s *Server) StartTriggeredRun(agentKey, title, message string, taskID int64, mergeable bool, taskDesc, taskGoal string)

StartTriggeredRun enqueues a P3 trigger fire for agentKey and pumps the queue. The agent's策略 decides concurrency + merge: serial → run one at a time (optionally merging by task / all / none); parallel → run each fire in its own concurrent conversation up to trigger_max_parallel. Distinct agents always run concurrently.

type Task

type Task struct {
	ID           string `json:"id"`
	ExpID        int64  `json:"exploration_id"`
	Name         string `json:"name"` // 可选任务名称;空=未命名
	CategoryID   *int64 `json:"category_id,omitempty"`
	CategoryName string `json:"category_name,omitempty"`
	PinnedAt     int64  `json:"pinned_at,omitempty"`
	Description  string `json:"description"`
	Goal         string `json:"goal"`
	CreatedAt    int64  `json:"created_at"`
	CompletedAt  int64  `json:"completed_at,omitempty"` // 进入终态的 unix 秒;0=未完成
	Paused       bool   `json:"paused"`
	Queued       bool   `json:"queued"` // 因并发上限被挂起、等待空位自动启动;true=尚未开跑
	// QueuedAt is an internal Unix-nanosecond ordering key. It is deliberately
	// finer than CreatedAt so several tasks enqueued in the same second retain
	// their real FIFO order.
	QueuedAt           int64   `json:"queued_at,omitempty"`
	QueueMode          string  `json:"queue_mode,omitempty"`
	ParentRef          string  `json:"parent_ref,omitempty"`     // 父任务 id(编排 spawn 记录)
	LLMProfileID       *int64  `json:"llm_profile_id,omitempty"` // 指定运行本任务 planner/worker 的 LLM 配置;nil=用全局激活配置
	LLMProfileIDs      []int64 `json:"llm_profile_ids,omitempty"`
	ActiveLLMProfileID *int64  `json:"active_llm_profile_id,omitempty"`
	LLMChainRevision   int64   `json:"-"`
	LLMFailoverState   string  `json:"llm_failover_state,omitempty"`
	LLMFailoverReason  string  `json:"llm_failover_reason,omitempty"`
	SourceTaskIDs      []int64 `json:"source_task_ids,omitempty"`
	CompanyIDs         []int64 `json:"company_ids,omitempty"`
	Status             string  `json:"status"` // persisted lifecycle status (done/failed/timeout 为终态;空/其它则由运行态推导)
	// 任务级超时(见 docs/任务级超时与收尾设计.md)。DeadlineAt/FirstRunAt 为 unix 秒,0=未设/未运行。
	TimeoutSeconds       int                    `json:"timeout_seconds"`
	PlanHeartbeatSeconds int                    `json:"plan_heartbeat_seconds"` // planner 心跳触发间隔(秒)
	CoverageEnabled      bool                   `json:"coverage_enabled"`       // 资产覆盖度功能开关(创建时定,默认开)
	FirstRunAt           int64                  `json:"first_run_at,omitempty"`
	DeadlineAt           int64                  `json:"deadline_at,omitempty"`
	Store                *pgdb.ExplorationStore `json:"-"`
	Guard                *guard.Guard           `json:"-"`
	// contains filtered or unexported fields
}

Task is one engagement: a description + goal + its own exploration store, sharing the process-wide asset store. ID is the PG task id as a string; ExpID is the exploration the task owns.

func (*Task) Notify

func (t *Task) Notify()

Notify signals that the asset/exploration graph changed (debounced consumer wakes the planner). Non-blocking.

func (*Task) NotifyCancelled

func (t *Task) NotifyCancelled(intentID int64, summary, reason string)

NotifyCancelled records that the human deleted intentID (reason = 删除原因), then wakes the planner so the next round spells out which intent was removed and why. summary is the intent's text captured before deletion — needed for hard delete, where the node is gone by the time the planner reads the trigger. Applies to both soft (state='deleted') and hard (physical cascade) delete.

func (*Task) NotifyDone

func (t *Task) NotifyDone(intentID int64)

NotifyDone is Notify plus a hint: a worker just finished intentID and that is what triggered this wake-up. The planner reads the accumulated triggers next round so it can spell out which intent finished (+ its output). Events pile up (debounce) until the round drains them via drainTriggers.

func (*Task) NotifyFinding

func (t *Task) NotifyFinding(intentID int64, summary string)

NotifyFinding records that a worker reported a finding on intentID (summary), then wakes the planner — so the round spells out which intent found what.

func (*Task) NotifyGoal

func (t *Task) NotifyGoal(texts []string)

NotifyGoal records that one OR MORE goals were added in a single set_goals call — by the human via the main agent — then wakes the planner, so the next round spells out "人新增了 N 个目标:…" instead of the planner having to spot new open goals in the overview. One call → one trigger event (set_goals 的一次批量算一条,不逐条刷屏). The event survives an early-returning terminal round (drain happens after the gate), so a set_goals that revives a done task still surfaces it once the task is running.

func (*Task) NotifyGoalDeleted

func (t *Task) NotifyGoalDeleted(text string)

NotifyGoalDeleted records that the human deleted a goal (via 总览的目标管理), then wakes the planner so the next round spells out which goal was removed. The event survives an early-returning terminal round (drain happens after the gate).

func (*Task) NotifyGoalEdited

func (t *Task) NotifyGoalEdited(oldText, newText string)

NotifyGoalEdited records that the human edited a goal (via 总览的目标管理), then wakes the planner so the next round spells out the old→new change. The event survives an early-returning terminal round (drain happens after the gate).

func (*Task) NotifyHint

func (t *Task) NotifyHint(texts []string)

NotifyHint records that one OR MORE hints were added in a single add_hint call — by the human via the main agent, or by cross-task orchestration — then wakes the planner, so the next round is told "人新增了 N 条战略提示:…" and looks at them directly instead of having to spot the new hint folded into the graph overview. One call → one trigger event (a batched add_hint counts as one, not one per hint).

type TaskDTO

type TaskDTO struct {
	ID                 string             `json:"id"`
	ExplorationID      int64              `json:"exploration_id"`
	Name               string             `json:"name"` // 可选任务名称;空=未命名
	CategoryID         *int64             `json:"category_id,omitempty"`
	CategoryName       string             `json:"category_name,omitempty"`
	Pinned             bool               `json:"pinned"`
	PinnedAt           string             `json:"pinned_at,omitempty"`
	Description        string             `json:"description"`
	Goal               string             `json:"goal"`
	Status             string             `json:"status"` // created | running | paused | done | failed
	CreatedAt          string             `json:"created_at"`
	CreatedUnix        int64              `json:"created_unix"`       // created_at as unix seconds (for run-duration calc)
	CompletedAt        string             `json:"completed_at"`       // RFC3339 finish time (done/failed); "" if unfinished
	CompletedUnix      int64              `json:"completed_unix"`     // completed_at as unix seconds (0 if unfinished)
	LastActivity       int64              `json:"last_activity_unix"` // unix seconds of the last activity (0 if none)
	Paused             bool               `json:"paused"`
	Queued             bool               `json:"queued"`
	Tokens             TokenTotalDTO      `json:"tokens"` // whole-task token consumption
	GoalsTotal         int                `json:"goals_total"`
	GoalsMet           int                `json:"goals_met"`
	InFlight           int                `json:"in_flight"`                // 运行中 Worker 数(state=running 的意图)
	Findings           FindingSeverityDTO `json:"findings"`                 // 该任务已登记的漏洞数(findings 表,按严重度分档)
	LLMProfileID       *int64             `json:"llm_profile_id,omitempty"` // LLM profile used for this task; nil = default
	LLMProfileIDs      []int64            `json:"llm_profile_ids"`
	ActiveLLMProfileID *int64             `json:"active_llm_profile_id,omitempty"`
	LLMFailoverState   string             `json:"llm_failover_state"`
	LLMFailoverReason  string             `json:"llm_failover_reason,omitempty"`
	SourceTaskIDs      []string           `json:"source_task_ids"`
	ArchiveBlockedBy   string             `json:"archive_blocked_by_task_id,omitempty"`
	CompanyIDs         []int64            `json:"company_ids"`
	CoverageEnabled    bool               `json:"coverage_enabled"` // 资产覆盖度功能开关(创建时定)
}

---- Task (frontend "Task") ---- created_at as RFC3339, plus a derived status.

type TaskNodeDTO

type TaskNodeDTO struct {
	ID           string `json:"id"`
	Type         string `json:"type"` // db Node.Kind
	Payload      string `json:"payload,omitempty"`
	Priority     int    `json:"priority"`
	State        string `json:"state"`
	Origin       string `json:"origin"`
	TS           string `json:"ts"`
	SourceTaskID string `json:"source_task_id,omitempty"`
	Inherited    bool   `json:"inherited,omitempty"`
	DeleteReason string `json:"delete_reason,omitempty"` // 意图假删除(state='deleted')时的删除原因
}

type TokenTotalDTO

type TokenTotalDTO struct {
	InputTokens      int `json:"input_tokens"`
	OutputTokens     int `json:"output_tokens"`
	CacheReadTokens  int `json:"cache_read_tokens"`
	CacheWriteTokens int `json:"cache_write_tokens"`
}

TokenTotalDTO is a whole-task (all agents) token aggregate.

type TrafficExchangeDTO

type TrafficExchangeDTO struct {
	ID          string `json:"id"`
	TS          string `json:"ts"`
	Host        string `json:"host"`
	Method      string `json:"method"`
	URL         string `json:"url"`
	Status      int    `json:"status"`
	ContentType string `json:"content_type"`
	RespLen     int    `json:"resp_len"`
}

---- TrafficExchange (frontend "TrafficExchange") ---- ts as RFC3339.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL