Documentation
¶
Overview ¶
Package messaging 是 server、worker 共同依赖的消息契约层。
本包定义统一的消息信封、命令结构、事件结构、subject 生成和 JetStream stream 配置, 使 server 和 worker 不再各自维护独立的协议定义和 subject 拼接逻辑。
按 QoS 分 lane 的物理 subject 设计:
Server -> Worker: cmd.run - 会话/task 执行命令,保留 session-keyed debounce cmd.control - cancel run 等控制命令,不经过防抖 cmd.interaction - approval resolve、question answer cmd.file - project file restore,request/reply Worker -> Server/UI: run.stream - 高频 SSE 增量(message delta、reasoning delta、tool delta) run.state - 低频关键状态(run.started、artifact.declared、approval/question、terminal)
Index ¶
- Constants
- func DecodeCommandPayload[T any](body *WorkerCommandBody) (T, error)
- func LLMUsageSubject(orgID uint) (string, error)
- func LLMUsageWildcard() string
- func OrgIDFromSkillPackageSubject(subject string) (uint, error)
- func ProjectNotifySubject(orgID, projectID uint) (string, error)
- func ProjectNotifyWildcard() string
- func RunEventStateWildcard() string
- func RunEventStreamWildcard() string
- func RunEventSubject(orgID uint, sessionID string, lane RunEventLane) (string, error)
- func RunEventWildcard() string
- func SessionRunStateConsumer() string
- func SkillPackageUploadedConsumer() string
- func SkillPackageUploadedSubject(orgID uint) (string, error)
- func SkillPackageUploadedWildcard() string
- func StreamConfigs() map[string]nats.StreamConfig
- func StreamNameFromSubject(subject string) string
- func WorkerCommandSubject(orgID, workerID uint, lane Lane) (string, error)
- func WorkerCommandWildcard() string
- func WorkerControlConsumer() string
- func WorkerInteractionConsumer() string
- func WorkerLaneConsumer(orgID, workerID uint, lane Lane) string
- func WorkerOpsStatusSubject(orgID, workerID uint) (string, error)
- func WorkerRunConsumer() string
- type ActorContext
- type ApprovalDecisionPayload
- type ApprovalRequestPayload
- type ApprovalResolveCommandPayload
- type ArtifactPayload
- type AssistantMessageTrigger
- type Attachment
- type CancelRunCommandPayload
- type ChatMessage
- type CommandType
- type Envelope
- type ExecutionTarget
- type GlobalEventPayload
- type GlobalEventType
- type HumanMessageData
- type InputType
- type Lane
- type MemberBrief
- type MessageCreatedData
- type MessageRole
- type MessageType
- type ModelOptions
- type PlanPublishedPayload
- type PluginSnapshot
- type ProjectContext
- type ProjectFileRestoreCommandPayload
- type ProjectFileRestoreResult
- type QuestionAnswerCommandPayload
- type QuestionAnswerPayload
- type QuestionItem
- type QuestionOption
- type QuestionRequestPayload
- type RouteContext
- type RunCommandMetadata
- type RunCommandPayload
- type RunCompletedPayload
- type RunEvent
- type RunEventBody
- type RunEventError
- type RunEventLane
- type RunEventPayload
- type RunEventRecord
- type RunEventType
- type RunMetadataPayload
- type RunResultPayload
- type RuntimeOptions
- type RuntimeTodoItem
- type SenderType
- type SkillChangeType
- type SkillDetailData
- type SkillListItem
- type SkillPackageUploadedEvent
- type TaskInput
- type TaskPolicy
- type TaskType
- type ToolCallPayload
- type ToolCallResultPayload
- type TraceContext
- type UsagePayload
- type WorkTitleUpdatedPayload
- type WorkerCommand
- func NewApprovalResolveCommand(envID string, route RouteContext, payload ApprovalResolveCommandPayload, ...) WorkerCommand
- func NewCancelRunCommand(envID string, route RouteContext, payload CancelRunCommandPayload, ...) WorkerCommand
- func NewProjectFileRestoreCommand(envID string, route RouteContext, payload ProjectFileRestoreCommandPayload) WorkerCommand
- func NewQuestionAnswerCommand(envID string, route RouteContext, payload QuestionAnswerCommandPayload, ...) WorkerCommand
- func NewRunCommand(envID string, route RouteContext, trace TraceContext, ...) WorkerCommand
- type WorkerCommandBody
- type WorkerCommandResult
- type WorkerRunSummary
- type WorkerStatusRequest
- type WorkerStatusSnapshot
- type WorkspaceOptions
Constants ¶
const ( StreamNameWorker = "WORKER_CMD_STREAM" StreamNameSession = "SESSION_RUN_STREAM" StreamNameGlobalNotify = "GLOBAL_NOTIFY_STREAM" StreamNameLLMUsage = "LLM_USAGE_STREAM" StreamNameSkillPackage = "SKILL_PACKAGE_STREAM" )
const MaxWaitingTasks = 100
MaxWaitingTasks 是等待任务摘要返回条数上限。
Variables ¶
This section is empty.
Functions ¶
func DecodeCommandPayload ¶
func DecodeCommandPayload[T any](body *WorkerCommandBody) (T, error)
DecodeCommandPayload 从 WorkerCommandBody.Payload 解码为指定类型。
func LLMUsageSubject ¶ added in v0.3.0
LLMUsageSubject 构建 worker -> server LLM usage 上报 subject。
格式:org.<org_id>.usage.llm
func LLMUsageWildcard ¶ added in v0.3.0
func LLMUsageWildcard() string
LLMUsageWildcard 返回匹配所有 LLM usage 上报的 wildcard subject。
格式:org.*.usage.llm
func OrgIDFromSkillPackageSubject ¶ added in v0.3.7
OrgIDFromSkillPackageSubject obtains the organization scope carried by the subject.
func ProjectNotifySubject ¶ added in v0.1.23
ProjectNotifySubject 构建 project 级全局通知 subject。
格式:org.<org_id>.project.<project_id>.notify
func ProjectNotifyWildcard ¶ added in v0.1.23
func ProjectNotifyWildcard() string
ProjectNotifyWildcard 返回匹配所有 project notify 的 wildcard subject。
格式:org.*.project.*.notify
func RunEventStateWildcard ¶
func RunEventStateWildcard() string
RunEventStateWildcard 返回匹配所有 state lane 事件的 wildcard subject。
func RunEventStreamWildcard ¶
func RunEventStreamWildcard() string
RunEventStreamWildcard 返回匹配所有 stream lane 事件的 wildcard subject。
func RunEventSubject ¶
func RunEventSubject(orgID uint, sessionID string, lane RunEventLane) (string, error)
RunEventSubject 构建 worker -> server 运行事件 subject。
格式:org.<org_id>.session.<session_id>.run.<lane>
func RunEventWildcard ¶
func RunEventWildcard() string
RunEventWildcard 返回匹配所有 run event 的 wildcard subject。
func SessionRunStateConsumer ¶
func SessionRunStateConsumer() string
SessionRunStateConsumer 返回 session run state projector 的持久化消费者名称。 用于消费 run.state 事件,投影更新 session 的当前运行状态。
func SkillPackageUploadedConsumer ¶ added in v0.3.7
func SkillPackageUploadedConsumer() string
SkillPackageUploadedConsumer returns the durable server consumer name.
func SkillPackageUploadedSubject ¶ added in v0.3.7
SkillPackageUploadedSubject constructs the worker -> server Skill package subject.
Format: org.<org_id>.skill.package.uploaded
func SkillPackageUploadedWildcard ¶ added in v0.3.7
func SkillPackageUploadedWildcard() string
SkillPackageUploadedWildcard matches Skill package upload notifications.
func StreamConfigs ¶
func StreamConfigs() map[string]nats.StreamConfig
StreamConfigs 返回所有预配置的 JetStream stream 配置。
WORKER_CMD_STREAM: server -> worker 方向,覆盖所有 worker command subject (cmd.run、cmd.control、cmd.interaction、cmd.file)。
保留 72h,每 subject 最多 10000 条。使用 DiscardOld, 积压时丢弃最旧消息以确保新命令始终可写入。
SESSION_RUN_STREAM: worker -> server/UI 方向,覆盖所有 run event subject (run.stream、run.state)。
保留 24h,每 subject 最多 10000 条。
func StreamNameFromSubject ¶
StreamNameFromSubject 根据 subject 的路径结构判断它属于哪个 stream。 worker command subjects: org.<id>.worker.<id>.cmd.* → WORKER_CMD_STREAM session event subjects: org.<id>.session.<id>.run.* → SESSION_RUN_STREAM
func WorkerCommandSubject ¶
WorkerCommandSubject 构建 server -> worker 命令 subject。
格式:org.<org_id>.worker.<worker_id>.cmd.<lane>
func WorkerCommandWildcard ¶
func WorkerCommandWildcard() string
WorkerCommandWildcard 返回匹配所有 worker command 的 wildcard subject。
格式:org.*.worker.*.cmd.>
func WorkerControlConsumer ¶
func WorkerControlConsumer() string
WorkerControlConsumer 返回 cmd.control lane 的持久化消费者名称。
func WorkerInteractionConsumer ¶
func WorkerInteractionConsumer() string
WorkerInteractionConsumer 返回 cmd.interaction lane 的持久化消费者名称。
func WorkerLaneConsumer ¶ added in v0.1.22
WorkerLaneConsumer 返回按 org/worker/lane 隔离的 worker 持久化消费者名称。
func WorkerOpsStatusSubject ¶ added in v0.3.9
WorkerOpsStatusSubject 构建 server -> worker 运维状态查询 subject。
该 subject 走 Core NATS request/reply,不进入 JetStream 任务队列, 因此不注册到任何 stream,也不参与 WorkerCommandWildcard。
格式:org.<org_id>.worker.<worker_id>.ops.status
func WorkerRunConsumer ¶
func WorkerRunConsumer() string
WorkerRunConsumer 返回 cmd.run lane 的持久化消费者名称。 用于 SubscribeManualDurable,worker 重启后 NATS 从断点续投。
Types ¶
type ActorContext ¶
type ApprovalDecisionPayload ¶
type ApprovalDecisionPayload struct {
RequestID string `json:"request_id"`
Action string `json:"action"` // "approve" | "deny" | "always"
Reason string `json:"reason,omitempty"`
}
ApprovalDecisionPayload 描述审批结果。
type ApprovalRequestPayload ¶
type ApprovalRequestPayload struct {
RequestID string `json:"request_id"`
ToolName string `json:"tool_name"`
ToolCallID string `json:"tool_call_id,omitempty"`
Description string `json:"description"`
Arguments json.RawMessage `json:"arguments,omitempty"`
Metadata map[string]string `json:"metadata,omitempty"`
}
ApprovalRequestPayload 描述需要用户审批的工具调用。
type ApprovalResolveCommandPayload ¶
type ApprovalResolveCommandPayload struct {
Action string `json:"action"` // "approve" | "deny" | "always"
Reason string `json:"reason,omitempty"`
}
ApprovalResolveCommandPayload 是 approval.resolve 命令的 payload。
type ArtifactPayload ¶
type ArtifactPayload struct {
ArtifactID string `json:"artifact_id,omitempty"`
Title string `json:"title,omitempty"`
Filename string `json:"filename,omitempty"`
OriginalName string `json:"original_name,omitempty"`
Description string `json:"description,omitempty"`
MimeType string `json:"mime_type,omitempty"`
ArtifactType string `json:"artifact_type,omitempty"`
FileSize int64 `json:"file_size,omitempty"`
CreatedAt string `json:"created_at,omitempty"`
RelativePath string `json:"relative_path,omitempty"`
PreviousRelativePath string `json:"previous_relative_path,omitempty"`
StorageKey string `json:"storage_key,omitempty"`
StorageURI string `json:"storage_uri,omitempty"`
Sha256 string `json:"sha256,omitempty"`
Source string `json:"source,omitempty"`
Status string `json:"status,omitempty"`
VersionNo int `json:"version_no,omitempty"`
}
ArtifactPayload 引用单次运行产生的产物。
type AssistantMessageTrigger ¶ added in v0.1.23
type AssistantMessageTrigger struct {
SenderType SenderType `json:"sender_type"`
AssistantID *string `json:"assistant_id,omitempty"`
AssistantName string `json:"assistant_name"`
RunID string `json:"run_id,omitempty"`
}
AssistantMessageTrigger 是 sender_type=assistant 时的 message.created payload。
type Attachment ¶
type CancelRunCommandPayload ¶
type CancelRunCommandPayload struct {
RunID string `json:"run_id"`
Reason string `json:"reason,omitempty"`
}
CancelRunCommandPayload 是 run.cancel 命令的 payload。
type ChatMessage ¶
type ChatMessage struct {
ID string `json:"id,omitempty"`
Role MessageRole `json:"role"`
Content string `json:"content"`
SenderName string `json:"sender_name,omitempty"`
}
type CommandType ¶
type CommandType string
CommandType 表示 server 发给 worker 的命令类型。
const ( // CommandTypeRun 请求 Worker 执行 Agent run。 CommandTypeRun CommandType = "agent.run" // CommandTypeCancel 请求 Worker 取消正在运行的 Agent run。 CommandTypeCancel CommandType = "run.cancel" // CommandTypeApprovalResolve 发送审批决策给 Worker。 CommandTypeApprovalResolve CommandType = "approval.resolve" // CommandTypeQuestionAnswer 发送问题答案给 Worker。 CommandTypeQuestionAnswer CommandType = "question.answer" // CommandTypeProjectFileRestore 请求 Worker 恢复项目文件历史版本。 CommandTypeProjectFileRestore CommandType = "project.file.restore" )
type Envelope ¶
type Envelope[T any] struct { ID string `json:"id"` Type MessageType `json:"type"` CreatedAt time.Time `json:"created_at"` Trace TraceContext `json:"trace"` Route RouteContext `json:"route"` Body T `json:"body"` Metadata json.RawMessage `json:"metadata,omitempty"` }
Envelope 是通用消息信封,用于所有 MQ topic 上的消息传输。
type ExecutionTarget ¶
type ExecutionTarget struct {
// AssistantID 是 leros_digital_assistant.id,自增主键,用于 llm_history 关联。
AssistantID uint `json:"assistant_id,omitempty"`
// AssistantPublicID 是 leros_digital_assistant.public_id,用于 worker 侧日志展示。
AssistantPublicID string `json:"assistant_public_id,omitempty"`
AssistantName string `json:"assistant_name,omitempty"`
AssistantDesc string `json:"assistant_desc,omitempty"`
SystemPrompt string `json:"system_prompt,omitempty"`
Skills []string `json:"skills,omitempty"`
Tools []string `json:"tools,omitempty"`
}
type GlobalEventPayload ¶ added in v0.1.23
type GlobalEventPayload struct {
Type GlobalEventType `json:"type"`
ProjectID uint `json:"project_id"`
SessionID string `json:"session_id"`
Seq uint64 `json:"seq"`
Timestamp int64 `json:"timestamp"`
Data json.RawMessage `json:"data"`
}
GlobalEventPayload 是 GlobalEvents SSE 端点下发的统一信封。
Seq 字段由消费侧从 NATS JetStream message metadata 填充, 发布方无需设置。
type GlobalEventType ¶ added in v0.1.23
type GlobalEventType string
GlobalEventType 标识全局通知事件的类型。
const GlobalEventMessageCreated GlobalEventType = "message.created"
GlobalEventMessageCreated 在用户消息或 AI 队友回复持久化成功后触发, 通知群聊成员有新消息到达。
const GlobalEventWorkTitleUpdated GlobalEventType = "work.title.updated"
GlobalEventWorkTitleUpdated 在项目/任务/会话标题自动更新后触发, 通知项目成员刷新工作台标题展示。
type HumanMessageData ¶ added in v0.1.23
type HumanMessageData struct {
SenderType SenderType `json:"sender_type"`
SenderUin *uint `json:"sender_uin,omitempty"`
SenderName string `json:"sender_name"`
Content string `json:"content"`
MessageType string `json:"message_type"`
Sequence int64 `json:"sequence"`
RunID string `json:"run_id,omitempty"`
CreatedAt string `json:"created_at"`
}
HumanMessageData 是 sender_type=human 时的 message.created payload。 携带完整消息数据,前端可直接渲染,无需再调 GetSessionMessages。
type MemberBrief ¶ added in v0.1.27
type MemberBrief struct {
MemberID uint `json:"member_id"`
MemberType string `json:"member_type"` // user / assistant
MemberRole string `json:"member_role"` // owner / admin / member / viewer
Name string `json:"name"`
IsDefault bool `json:"is_default,omitempty"`
IsCurrentExec bool `json:"is_current_exec,omitempty"` // marks the assistant executing this run
IsCurrentUser bool `json:"is_current_user,omitempty"` // marks the user who initiated this run
}
MemberBrief is a lightweight project member snapshot.
type MessageCreatedData ¶ added in v0.1.23
type MessageCreatedData struct {
MessageID uint `json:"message_id"`
Sequence int64 `json:"sequence"`
SenderType SenderType `json:"sender_type"`
SenderUin *uint `json:"sender_uin,omitempty"`
SenderName string `json:"sender_name"`
AssistantID *string `json:"assistant_id,omitempty"`
AssistantName string `json:"assistant_name,omitempty"`
Content string `json:"content"`
RunID string `json:"run_id,omitempty"`
// 前端 human 消息渲染所需字段(assistant 事件不设置,omitempty 保证不出现)
MessageType string `json:"message_type,omitempty"`
Attachments []types.MessageAttachment `json:"attachments,omitempty"`
Metadata *types.ObjectMetadata `json:"metadata,omitempty"`
}
MessageCreatedData 是 message.created 事件的 payload。
human 场景 Content 为完整消息内容(前端可直接渲染); assistant 场景(run.started 时发布)Content 为空,仅通知"AI 开始回复", 前端收到后用 session_id + assistant_id 订阅 SessionEvents 获取流式输出。
type MessageRole ¶
type MessageRole string
const ( MessageRoleUser MessageRole = "user" MessageRoleAssistant MessageRole = "assistant" MessageRoleSystem MessageRole = "system" MessageRoleTool MessageRole = "tool" )
type MessageType ¶
type MessageType string
MessageType 表示消息的顶层类型标识。
const ( // MessageTypeWorkerCommand 统一 server -> worker 命令。 MessageTypeWorkerCommand MessageType = "worker.command" // MessageTypeRunEvent 统一 worker -> server/UI 运行事件。 MessageTypeRunEvent MessageType = "run.event" )
type ModelOptions ¶
type ModelOptions struct {
ModelID uint `json:"model_id,omitempty"`
Provider string `json:"provider,omitempty"`
Model string `json:"model,omitempty"`
BaseURL string `json:"base_url,omitempty"`
BaseURLHasV1 bool `json:"base_url_has_v1,omitempty"`
APIKey string `json:"api_key,omitempty"`
// Vision 表示该模型是否支持图片(多模态)输入。
Vision bool `json:"vision,omitempty"`
// Temperature 默认采样温度;0 表示未配置,走 provider 默认。
Temperature float64 `json:"temperature,omitempty"`
// MaxTokens 默认最大输出 token 数;0 表示未配置,走 provider 默认。
MaxTokens int `json:"max_tokens,omitempty"`
TopP *float64 `json:"top_p,omitempty"`
FrequencyPenalty *float64 `json:"frequency_penalty,omitempty"`
PresencePenalty *float64 `json:"presence_penalty,omitempty"`
ContextLimit int `json:"context_limit,omitempty"`
OutputLimit int `json:"output_limit,omitempty"`
}
type PlanPublishedPayload ¶ added in v0.1.19
type PlanPublishedPayload struct {
FileID string `json:"file_id"`
Directive string `json:"directive"`
SummaryLines int `json:"summary_lines"`
TotalLines int `json:"total_lines"`
StorageKey string `json:"storage_key,omitempty"`
StorageURI string `json:"storage_uri,omitempty"`
Filename string `json:"filename,omitempty"`
OriginalName string `json:"original_name,omitempty"`
MimeType string `json:"mime_type,omitempty"`
FileSize int64 `json:"file_size,omitempty"`
Sha256 string `json:"sha256,omitempty"`
}
PlanPublishedPayload carries the uploaded plan file info and :::plan directive.
type PluginSnapshot ¶ added in v0.3.7
type PluginSnapshot struct {
PluginID string `json:"plugin_id"`
Code string `json:"code"`
Kind string `json:"kind"`
Revision int `json:"revision"`
Definition json.RawMessage `json:"definition"`
}
PluginSnapshot is the immutable plugin revision selected when a run is published. Definition is the immutable plugin configuration selected for this run.
type ProjectContext ¶ added in v0.1.27
type ProjectContext struct {
Name string `json:"name,omitempty"`
Description string `json:"description,omitempty"`
Objective string `json:"objective,omitempty"`
Members []MemberBrief `json:"members,omitempty"`
}
ProjectContext carries project business context to the worker.
type ProjectFileRestoreCommandPayload ¶ added in v0.2.5
type ProjectFileRestoreCommandPayload struct {
ProjectPublicID string `json:"project_public_id"`
RelativePath string `json:"relative_path"`
Branch string `json:"branch,omitempty"`
DownloadURL string `json:"download_url"`
AuthorName string `json:"author_name,omitempty"`
AuthorEmail string `json:"author_email,omitempty"`
}
ProjectFileRestoreCommandPayload 是 project.file.restore 命令的 payload。
type ProjectFileRestoreResult ¶ added in v0.2.5
type ProjectFileRestoreResult struct {
Success bool `json:"success"`
RelativePath string `json:"relative_path,omitempty"`
CommitSHA string `json:"commit_sha,omitempty"`
Error string `json:"error,omitempty"`
}
ProjectFileRestoreResult 是 Worker 完成项目文件恢复后的同步响应。
type QuestionAnswerCommandPayload ¶
type QuestionAnswerCommandPayload struct {
Answers [][]string `json:"answers"`
}
QuestionAnswerCommandPayload 是 question.answer 命令的 payload。
type QuestionAnswerPayload ¶
type QuestionAnswerPayload struct {
RequestID string `json:"request_id"`
Answers [][]string `json:"answers"`
}
QuestionAnswerPayload 描述用户对问题的回答。
type QuestionItem ¶
type QuestionItem struct {
Question string `json:"question"`
Header string `json:"header,omitempty"`
Options []QuestionOption `json:"options"`
MultiSelect bool `json:"multiple"`
Custom bool `json:"custom"`
}
QuestionItem 是问题请求中的单个问题。
type QuestionOption ¶
type QuestionOption struct {
Label string `json:"label"`
Description string `json:"description,omitempty"`
}
QuestionOption 是问题的选项。
type QuestionRequestPayload ¶
type QuestionRequestPayload struct {
RequestID string `json:"request_id"`
SessionID string `json:"session_id,omitempty"`
Questions []QuestionItem `json:"questions"`
ToolCallID string `json:"tool_call_id,omitempty"`
MessageID string `json:"message_id,omitempty"`
InteractionType string `json:"interaction_type,omitempty"`
Metadata map[string]string `json:"metadata,omitempty"`
}
QuestionRequestPayload 描述引擎向用户提出的澄清问题。
type RouteContext ¶
type RouteContext struct {
OrgID uint `json:"org_id"`
SessionID string `json:"session_id,omitempty"` // SessionID 是 leros_session.public_id,用于消息路由和租户隔离。
// WorkerID 是 leros_worker_deployment.worker_id(worker 主键),用于 worker 消息分发。
WorkerID uint `json:"worker_id,omitempty"`
// WorkerPublicID 是 leros_worker_deployment.public_id,用于日志展示和对外追溯。
WorkerPublicID string `json:"worker_public_id,omitempty"`
// AssistantID 是 leros_digital_assistant.id(assistant 主键),用于内部关联和 llm_history 记录。
AssistantID uint `json:"assistant_id,omitempty"`
// AssistantPublicID 是 leros_digital_assistant.public_id,用于对外标识和 SSE 过滤。
AssistantPublicID string `json:"assistant_public_id,omitempty"`
ClientIP string `json:"client_ip,omitempty"`
}
RouteContext 携带消息路由信息,用于消息投递和租户隔离。
type RunCommandMetadata ¶
type RunCommandMetadata struct {
SessionID string `json:"session_id,omitempty"`
MessageType string `json:"message_type,omitempty"`
Sequence int64 `json:"sequence,omitempty"`
}
RunCommandMetadata contains typed optional metadata for agent.run commands.
type RunCommandPayload ¶
type RunCommandPayload struct {
TaskType TaskType `json:"task_type"`
ExecutionMode string `json:"execution_mode,omitempty"`
Actor ActorContext `json:"actor"`
Execution ExecutionTarget `json:"execution"`
Workspace WorkspaceOptions `json:"workspace,omitempty"`
Project ProjectContext `json:"project,omitempty"`
Input TaskInput `json:"input"`
Model ModelOptions `json:"model,omitempty"`
Runtime RuntimeOptions `json:"runtime,omitempty"`
Policy TaskPolicy `json:"policy,omitempty"`
Plugins []PluginSnapshot `json:"plugins,omitempty"`
// 业务主键 ID,用于 llm_history 等调用记录关联。
// 以下均为对应表的自增主键(int),区别于其它字段中的 public_id(string)。
//
// ProjectID leros_project.id -> 区别于 Workspace.ProjectID(project public_id)
// SessionID leros_session.id -> 区别于 RouteContext.SessionID(session public_id)
// MessageID leros_session_message.id -> 当前触发 run 的消息主键
// AssistantID leros_digital_assistant.id -> 区别于 Execution.AssistantPublicID(assistant public_id)
// AssistantPublicID leros_digital_assistant.public_id -> 用于 worker 侧展示和对外追溯
// Uin leros_user.id -> 区别于 ActorContext.UserID(fmt.Sprintf("%d", uin))
ProjectID uint `json:"project_id"`
SessionID uint `json:"session_id"`
MessageID uint `json:"message_id"`
AssistantID uint `json:"assistant_id"`
AssistantPublicID string `json:"assistant_public_id,omitempty"`
Uin uint `json:"uin"`
// NotAfter Worker 最晚允许开始时间(RFC3339,UTC)。超过该时间 Worker 应拒绝执行。
// 需要持久化进 worker inbox,崩溃恢复时仍能生效。
NotAfter string `json:"not_after,omitempty"`
}
RunCommandPayload 是 agent.run 命令的 payload。
type RunCompletedPayload ¶
type RunCompletedPayload struct {
Status string `json:"status"`
Result RunResultPayload `json:"result"`
Artifacts []ArtifactPayload `json:"artifacts,omitempty"`
Usage *UsagePayload `json:"usage,omitempty"`
Events []RunEventRecord `json:"events,omitempty"`
StartedAt string `json:"started_at,omitempty"`
CompletedAt string `json:"completed_at,omitempty"`
Metadata *RunMetadataPayload `json:"metadata,omitempty"`
}
RunCompletedPayload 归档完整的成功运行。
type RunEventBody ¶
type RunEventBody struct {
Seq int64 `json:"seq"`
Event RunEventType `json:"event"`
Payload RunEventPayload `json:"payload"`
ReplyToMessageIDs []string `json:"reply_to_message_ids,omitempty"`
// MemberCommandIDs / MemberRunIDs 标识同一 Session 合并批次中的所有源命令。
MemberCommandIDs []string `json:"member_command_ids,omitempty"`
MemberRunIDs []string `json:"member_run_ids,omitempty"`
// RunCompleted 仅在终端事件(run.completed/failed/cancelled)时填充。
RunCompleted *RunCompletedPayload `json:"run_completed,omitempty"`
// Error 仅在 run.failed 时填充。
Error *RunEventError `json:"error,omitempty"`
// AssistantPKID 是 leros_digital_assistant.id(主键),用于 worker 侧持久化关联(如 llm_history)。
AssistantPKID uint `json:"assistant_pk_id,omitempty"`
// AssistantID 是 leros_digital_assistant.public_id,用于 server/UI 侧展示和 SSE 过滤。
AssistantID string `json:"assistant_id,omitempty"`
}
RunEventBody 是单个运行事件的 payload。
type RunEventError ¶
RunEventError 描述流执行中的终止或可恢复错误。
type RunEventLane ¶
type RunEventLane string
RunEventLane 表示 run event 发送到哪个 lane。
const ( // RunEventLaneStream 是高频 SSE 增量 lane。 RunEventLaneStream RunEventLane = "run.stream" // RunEventLaneState 是低频关键状态 lane。 RunEventLaneState RunEventLane = "run.state" )
func ClassifyRunEvent ¶
func ClassifyRunEvent(eventType RunEventType) RunEventLane
ClassifyRunEvent 根据事件类型返回对应的 lane。 高频增量事件进入 stream lane,关键状态和终端事件进入 state lane。
type RunEventPayload ¶
type RunEventPayload struct {
MessageID string `json:"message_id,omitempty"`
Role MessageRole `json:"role,omitempty"`
Content string `json:"content,omitempty"`
Usage *UsagePayload `json:"usage,omitempty"`
ToolCall *ToolCallPayload `json:"tool_call,omitempty"`
ToolResult *ToolCallResultPayload `json:"tool_result,omitempty"`
Todos []RuntimeTodoItem `json:"todos,omitempty"`
Artifact *ArtifactPayload `json:"artifact,omitempty"`
ApprovalRequest *ApprovalRequestPayload `json:"approval_request,omitempty"`
ApprovalDecision *ApprovalDecisionPayload `json:"approval_decision,omitempty"`
QuestionRequest *QuestionRequestPayload `json:"question_request,omitempty"`
QuestionAnswer *QuestionAnswerPayload `json:"question_answer,omitempty"`
PlanPublished *PlanPublishedPayload `json:"plan_published,omitempty"`
WorkTitle *WorkTitleUpdatedPayload `json:"work_title,omitempty"`
}
RunEventPayload 携带流事件的特定内容。
type RunEventRecord ¶
type RunEventRecord struct {
Seq int64 `json:"seq,omitempty"`
LastSeq int64 `json:"last_seq,omitempty"`
Type string `json:"type"`
Timestamp int64 `json:"timestamp,omitempty"`
Payload json.RawMessage `json:"payload,omitempty"`
AssistantID string `json:"assistant_id,omitempty"`
}
RunEventRecord 是归一化、已归档的运行时事件。
type RunEventType ¶
type RunEventType string
RunEventType 表示 worker 运行事件的类型。
const ( // RunEventRunStarted 表示 run 已开始。 RunEventRunStarted RunEventType = "run.started" // RunEventRunCompleted 表示 run 成功完成。 RunEventRunCompleted RunEventType = "run.completed" // RunEventRunFailed 表示 run 失败。 RunEventRunFailed RunEventType = "run.failed" // RunEventRunCancelled 表示 run 被取消。 RunEventRunCancelled RunEventType = "run.cancelled" // RunEventArtifactDeclared 表示生成了产物。 RunEventArtifactDeclared RunEventType = "artifact.declared" // RunEventApprovalRequested 表示引擎需要用户审批。 RunEventApprovalRequested RunEventType = "approval.requested" // RunEventApprovalResolved 表示审批已被解决。 RunEventApprovalResolved RunEventType = "approval.resolved" // RunEventQuestionAsked 表示引擎正在向用户提问。 RunEventQuestionAsked RunEventType = "question.asked" // RunEventQuestionAnswered 表示问题已被回答。 RunEventQuestionAnswered RunEventType = "question.answered" // RunEventWorkTitleUpdated 表示项目/任务标题已由 LLM 自动生成。 RunEventWorkTitleUpdated RunEventType = "work.title.updated" // RunEventPlanPublished 表示计划文件已上传到对象存储。 RunEventPlanPublished RunEventType = "plan.published" // RunEventMessageDelta 表示助手文本增量输出。 RunEventMessageDelta RunEventType = "message.delta" // RunEventReasoningDelta 表示推理文本增量输出。 RunEventReasoningDelta RunEventType = "reasoning.delta" // RunEventMessageCompleted 表示最终助手消息已生成。 RunEventMessageCompleted RunEventType = "message.completed" // RunEventToolCallStarted 表示工具调用开始。 RunEventToolCallStarted RunEventType = "tool_call.started" // RunEventToolCallFinished 表示工具调用结束。 RunEventToolCallFinished RunEventType = "tool_call.finished" // RunEventTodoSnapshot 表示完整运行时 todo 列表可用。 RunEventTodoSnapshot RunEventType = "todo.snapshot" // RunEventTodoUpdated 表示运行时 todo 列表已更新。 RunEventTodoUpdated RunEventType = "todo.updated" )
type RunMetadataPayload ¶
type RunMetadataPayload struct {
Runtime string `json:"runtime,omitempty"`
WorkDir string `json:"work_dir,omitempty"`
ProviderID string `json:"provider_id,omitempty"`
SessionID string `json:"session_id,omitempty"`
Phase string `json:"phase,omitempty"`
Resume bool `json:"resume,omitempty"`
}
RunMetadataPayload contains typed run metadata while preserving the wire JSON object.
type RunResultPayload ¶
type RunResultPayload struct {
Message string `json:"message,omitempty"`
}
RunResultPayload 描述 run.completed 中归档的最终助手结果。
type RuntimeOptions ¶
type RuntimeTodoItem ¶
type RuntimeTodoItem struct {
ID string `json:"id"`
Title string `json:"title"`
Status string `json:"status"`
Priority string `json:"priority,omitempty"`
}
RuntimeTodoItem 描述一个运行时本地规划步骤。
type SenderType ¶ added in v0.1.23
type SenderType string
SenderType 区分消息发送方是真人员工还是 AI 队友。
const ( SenderTypeHuman SenderType = "human" // 真人员工发言 SenderTypeAssistant SenderType = "assistant" // AI 队友回复 )
type SkillChangeType ¶ added in v0.3.7
type SkillChangeType string
SkillChangeType describes a publishable worker-side Skill change.
const ( SkillChangeCreated SkillChangeType = "created" SkillChangeUpdated SkillChangeType = "updated" )
type SkillDetailData ¶
type SkillDetailData struct {
SkillID string `json:"skill_id,omitempty"`
Name string `json:"name"`
Description string `json:"description"`
Category string `json:"category"`
Source string `json:"source"`
Trust string `json:"trust"`
Version string `json:"version"`
SkillMD string `json:"skill_md"`
Tags []string `json:"tags"`
Files []string `json:"files"`
}
SkillDetailData 表示已安装 skill 的完整详情,包括 SKILL.md 内容。
type SkillListItem ¶
type SkillListItem struct {
Name string `json:"name"`
Description string `json:"description"`
Category string `json:"category"`
Source string `json:"source"`
Trust string `json:"trust"`
}
SkillListItem 表示已安装的 skill。
type SkillPackageUploadedEvent ¶ added in v0.3.7
type SkillPackageUploadedEvent struct {
EventID string `json:"event_id"`
WorkerID uint `json:"worker_id"`
RunID string `json:"run_id"`
ProjectID uint `json:"project_id,omitempty"`
ActorUIN uint `json:"actor_uin"`
SkillCode string `json:"skill_code"`
ChangeType SkillChangeType `json:"change_type"`
StorageURI string `json:"storage_uri"`
SHA256 string `json:"sha256"`
FileSize int64 `json:"file_size"`
Filename string `json:"filename"`
MimeType string `json:"mime_type"`
}
SkillPackageUploadedEvent reports a Skill package already uploaded through the shared presigned-storage flow. FileUpload identity remains server-owned.
type TaskInput ¶
type TaskInput struct {
Type InputType `json:"type"`
Scene string `json:"scene,omitempty"`
OutputFormat string `json:"output_format,omitempty"`
Messages []ChatMessage `json:"messages,omitempty"`
Attachments []Attachment `json:"attachments,omitempty"`
}
type TaskPolicy ¶
type ToolCallPayload ¶
type ToolCallPayload struct {
ToolCallID string `json:"tool_call_id"`
Name string `json:"name"`
Arguments json.RawMessage `json:"arguments,omitempty"`
}
ToolCallPayload 是工具调用开始和参数事件的标准负载。
type ToolCallResultPayload ¶
type ToolCallResultPayload struct {
ToolCallID string `json:"tool_call_id"`
Name string `json:"name,omitempty"`
Result json.RawMessage `json:"result,omitempty"`
Error string `json:"error,omitempty"`
IsError bool `json:"is_error"`
ElapsedMS int64 `json:"elapsed_ms,omitempty"`
}
ToolCallResultPayload 是工具调用终止事件的标准负载。
type TraceContext ¶
type TraceContext struct {
TraceID string `json:"trace_id"`
RequestID string `json:"request_id,omitempty"`
ReqID string `json:"req_id,omitempty"`
TaskID string `json:"task_id,omitempty"`
RunID string `json:"run_id,omitempty"`
ParentID string `json:"parent_id,omitempty"`
}
TraceContext 携带跨越 UI、Server、Worker、Runtime 的分布式追踪标识。
type UsagePayload ¶
type UsagePayload struct {
TotalTokens int `json:"total_tokens"`
InputTokens int `json:"input_tokens"`
OutputTokens int `json:"output_tokens"`
CacheInputTokens int `json:"cache_input_tokens"`
CacheOutputTokens int `json:"cache_output_tokens"`
}
UsagePayload 描述模型 token 使用情况。
type WorkTitleUpdatedPayload ¶
type WorkTitleUpdatedPayload struct {
ProjectID string `json:"project_id"`
ProjectName string `json:"project_name"`
TaskID string `json:"task_id,omitempty"`
TaskTitle string `json:"task_title,omitempty"`
SessionID string `json:"session_id"`
SessionTitle string `json:"session_title,omitempty"`
}
WorkTitleUpdatedPayload notifies clients that project/task titles were auto-generated.
type WorkerCommand ¶
type WorkerCommand = Envelope[WorkerCommandBody]
WorkerCommand 是 Server -> Worker 的统一命令消息。
func NewApprovalResolveCommand ¶
func NewApprovalResolveCommand(envID string, route RouteContext, payload ApprovalResolveCommandPayload, requestID string) WorkerCommand
NewApprovalResolveCommand 构造一个 approval.resolve WorkerCommand。
func NewCancelRunCommand ¶
func NewCancelRunCommand(envID string, route RouteContext, payload CancelRunCommandPayload, runID string) WorkerCommand
NewCancelRunCommand 构造一个 run.cancel WorkerCommand。
func NewProjectFileRestoreCommand ¶ added in v0.2.5
func NewProjectFileRestoreCommand(envID string, route RouteContext, payload ProjectFileRestoreCommandPayload) WorkerCommand
NewProjectFileRestoreCommand 构造 project.file.restore WorkerCommand。
func NewQuestionAnswerCommand ¶
func NewQuestionAnswerCommand(envID string, route RouteContext, payload QuestionAnswerCommandPayload, requestID string) WorkerCommand
NewQuestionAnswerCommand 构造一个 question.answer WorkerCommand。
func NewRunCommand ¶
func NewRunCommand( envID string, route RouteContext, trace TraceContext, payload RunCommandPayload, metadata *RunCommandMetadata, ) WorkerCommand
NewRunCommand 构造一个 agent.run WorkerCommand。
type WorkerCommandBody ¶
type WorkerCommandBody struct {
CommandType CommandType `json:"command_type"`
Payload json.RawMessage `json:"payload,omitempty"`
// ReplyTo 是 server Request() 注入的 NATS inbox,worker 通过 core NATS 直接回复。
ReplyTo string `json:"reply_to,omitempty"`
}
WorkerCommandBody 是所有 worker 命令的联合 body。
根据 CommandType 不同,Payload 携带不同类型的数据:
- agent.run: RunCommandPayload
- run.cancel: CancelRunCommandPayload
- approval.resolve: ApprovalResolveCommandPayload
- question.answer: QuestionAnswerCommandPayload
- project.file.restore: ProjectFileRestoreCommandPayload
type WorkerCommandResult ¶
type WorkerCommandResult struct {
Success bool `json:"success"`
Action string `json:"action"`
Message string `json:"message,omitempty"`
Error string `json:"error,omitempty"`
Data any `json:"data,omitempty"`
}
WorkerCommandResult 是 Worker -> Server 的同步响应(用于 skill request/reply)。
type WorkerRunSummary ¶ added in v0.3.9
type WorkerRunSummary struct {
RunID string `json:"run_id,omitempty"`
TaskID string `json:"task_id,omitempty"`
SessionID string `json:"session_id,omitempty"`
CommandID string `json:"command_id,omitempty"`
StreamSeq uint64 `json:"stream_seq,omitempty"`
// Status 表示该任务当前所处阶段:running 或 waiting。
Status string `json:"status,omitempty"`
// CreatedAt / UpdatedAt / StartedAt 均为 Unix 秒;StartedAt 零值表示尚未开始。
CreatedAt int64 `json:"created_at,omitempty"`
UpdatedAt int64 `json:"updated_at,omitempty"`
StartedAt int64 `json:"started_at,omitempty"`
}
WorkerRunSummary 是一次运行/等待中任务的轻量摘要。
只包含定位与生命周期字段,不携带 prompt、模型配置、环境变量或原始命令, 避免运维接口泄露敏感信息。
type WorkerStatusRequest ¶ added in v0.3.9
WorkerStatusRequest 是 Server -> Worker 的运维状态查询请求。
该查询使用 Core NATS request/reply(临时 inbox),不写入 JetStream, 避免运维轮询污染任务队列。Worker 依据请求中的 org/worker 归属回答本地运行快照。
type WorkerStatusSnapshot ¶ added in v0.3.9
type WorkerStatusSnapshot struct {
// OrgID / WorkerID 标识生成该快照的 Worker,用于 Server 校验回复归属。
OrgID uint `json:"org_id"`
WorkerID uint `json:"worker_id"`
MaxConcurrency int `json:"max_concurrency"`
// RunningCount 已启动但尚未结束的 Runtime 数量;交互等待中的 Runtime 也计入。
RunningCount int `json:"running_count"`
// WaitingCount 已接收但尚未开始 Runtime 执行的本地 Run 数量。
WaitingCount int `json:"waiting_count"`
// DebounceWaitingCount 仍处于 debounce 窗口内的 Run 数量。
DebounceWaitingCount int `json:"debounce_waiting_count"`
// CoordinatorWaitingCount 已进入 Coordinator 队列或正在等待计算槽的 Run 数量。
CoordinatorWaitingCount int `json:"coordinator_waiting_count"`
// AdmissionWaitingCount 被 Worker 准入 semaphore 实际阻塞的消息数量。
AdmissionWaitingCount int `json:"admission_waiting_count"`
// AcceptedCount Worker 已持久化并由当前进程拥有的命令数量,不含尚未准入的消息。
AcceptedCount int `json:"accepted_count"`
// ComputeBusyCount 当前占用计算槽的 Runtime 数量。
ComputeBusyCount int `json:"compute_busy_count"`
// InteractionWaitingCount 当前处于审批/问答等交互等待的 Runtime 数量。
InteractionWaitingCount int `json:"interaction_waiting_count"`
// InboxPendingCount 持久化 inbox 中处于 pending 状态的记录数。
InboxPendingCount int `json:"inbox_pending_count"`
// InboxProcessingCount 持久化 inbox 中处于 processing 状态的记录数。
InboxProcessingCount int `json:"inbox_processing_count"`
// SnapshotAt 快照生成时间(Unix 秒)。
SnapshotAt int64 `json:"snapshot_at"`
// Degraded 表示快照有局部数据源不可用;Errors 只返回稳定错误码,不暴露底层细节。
Degraded bool `json:"degraded,omitempty"`
Errors []string `json:"errors,omitempty"`
// RunningTasks 正在执行的任务摘要。
RunningTasks []WorkerRunSummary `json:"running_tasks,omitempty"`
// WaitingTasks 等待执行的任务摘要;最多返回 maxWaitingTasks 条。
WaitingTasks []WorkerRunSummary `json:"waiting_tasks,omitempty"`
// WaitingTruncated 表示等待任务摘要因数量上限被截断。
WaitingTruncated bool `json:"waiting_truncated,omitempty"`
}
WorkerStatusSnapshot 是 Worker 返回的本地运行状态快照。
waiting_count 仅统计 Worker 已接收但尚未开始执行的本地任务 (经 debounce 合并后进入 per-session 队列或等待计算槽的批次); NATS 尚未投递到 Worker 的积压不计入该字段。