一. 相关结构 1. ContextManager 负责管理多个并发 Job(任务)上下文的核心组件
/// Manages contexts for multiple concurrent jobs.
pub struct ContextManager {
/// Active job contexts.
contexts: RwLock<HashMap<Uuid, JobContext>>,
- HashMap<Uuid, JobContext>:以 Uuid(每个 job 唯一标识)为键,映射到 JobContext
- RwLock:读写锁。多个 reader 可并发读(如多 channel 同时查询 job 状态),writer 独占写(注册/清理 job 时)
/// Memory for each job.
memories: RwLock<HashMap<Uuid, Memory>>,
/// Maximum concurrent jobs.
max_jobs: usize,
}
2. JobContext /// Context for a running job.
#[derive(Debug, Clone, Serialize)]
pub struct JobContext {
/// Unique job ID.
pub job_id: Uuid,
/// Current state.
pub state: JobState,
/// User ID that owns this job (for workspace scoping).
pub user_id: String,
/// Channel-specific requester/actor ID, when different from the owner scope.
#[serde(skip_serializing_if = "Option::is_none")]
pub requester_id: Option<String>,
/// Conversation ID if linked to a conversation.
pub conversation_id: Option<Uuid>,
/// Job title.
pub title: String,
/// Job description.
pub description: String,
/// Job category.
pub category: Option<String>,
/// Budget amount (if from marketplace).
pub budget: Option<Decimal>,
/// Budget token (e.g., "NEAR", "USD").
pub budget_token: Option<String>,
/// Our bid amount.
pub bid_amount: Option<Decimal>,
/// Estimated cost to complete.
pub estimated_cost: Option<Decimal>,
/// Estimated time to complete.
pub estimated_duration: Option<Duration>,
/// Actual cost so far.
pub actual_cost: Decimal,
/// Total tokens consumed by LLM calls in this job.
pub total_tokens_used: u64,
/// Maximum tokens allowed per job (0 = unlimited).
pub max_tokens: u64,
/// When the job was created.
pub created_at: DateTime<Utc>,
/// When the job was started.
pub started_at: Option<DateTime<Utc>>,
/// When the job was completed.
pub completed_at: Option<DateTime<Utc>>,
/// Number of repair attempts.
pub repair_attempts: u32,
/// State transition history.
pub transitions: Vec<StateTransition>,
/// Metadata.
pub metadata: serde_json::Value,
/// Extra environment variables to inject into spawned child processes.
///
/// Used by the worker runtime to pass fetched credentials to tools
/// (e.g., shell commands) without mutating the global process environment
/// via `std::env::set_var`, which is unsafe in multi-threaded programs.
///
/// Wrapped in `Arc` for cheap cloning on every tool invocation.
#[serde(skip)]
pub extra_env: Arc<HashMap<String, String>>,
/// Optional HTTP interceptor for trace recording/replay.
///
/// When set, tools that make outgoing HTTP requests should check this
/// interceptor before sending real requests. During recording, the
/// interceptor captures request/response pairs. During replay, it
/// returns pre-recorded responses.
#[serde(skip)]
pub http_interceptor: Option<Arc<dyn HttpInterceptor>>,
/// Stash of full tool outputs keyed by tool_call_id.
///
/// Tool outputs may be truncated before reaching the LLM context window,
/// but subsequent tools (e.g., `json`) may need the full output. This
/// stash stores the complete, unsanitized output so tools can reference
/// previous results by ID via `$tool_call_id` parameter syntax.
///
/// Also used for cross-tool implicit state (keys prefixed with `__`) such
/// as `__routine_last_name` for fallback recovery in routine tool chains.
#[serde(skip)]
pub tool_output_stash: Arc<tokio::sync::RwLock<HashMap<String, String>>>,
/// User's preferred timezone (IANA name, e.g. "America/New_York"). Defaults to "UTC".
pub user_timezone: String,
/// Approval context for tool execution in this job.
///
/// When set, tools check this context before executing to determine
/// if they're allowed to run in autonomous/non-interactive contexts.
#[serde(skip)]
pub approval_context: Option<ApprovalContext>,//工具执行待批准
}3. Memory // Combined memory for a job.
#[derive(Debug, Clone)]
pub struct Memory {
/// Job ID.
pub job_id: Uuid,
/// Conversation history.
pub conversation: ConversationMemory,//会话记录
/// Action history.
pub actions: Vec<ActionRecord>,//tool行为记录
/// Next action sequence number.
next_sequence: u32,
}仔细核对代码后,结论需要修正得更精确:
> 这里分析的是 `src/` 下的 V1 Job 子系统。`JobContext` 只有“部分、非持续”的数据库持久化;`Memory` 整体完全不持久化,只有其中的 `ActionRecord` 被单独写入数据库。
## 总体结论
| 内容 | 进程内保存 | 数据库保存 | 重启后恢复到 `ContextManager` |
|---|---:|---:|---:|
| `JobContext` | 是 | 部分字段 | 否 |
| `JobContext` 运行时字段 | 是 | 否 | 否 |
| `Memory` 整体 | 是 | 否 | 否 |
| `Memory.conversation` | 是,但生产代码基本不用 | 否 | 否 |
| `Memory.actions` | 是 | 每条 `ActionRecord` 单独保存 | 否 |
| Job 实际 LLM 上下文 | Worker 局部变量 | 主要通过事件等其他路径 | 不由 `Memory` 恢复 |
---
## 1. ContextManager 本身不持久化
`ContextManager` 只有两个进程内 `HashMap`:
```rust
pub struct ContextManager {
contexts: RwLock<HashMap<Uuid, JobContext>>,
memories: RwLock<HashMap<Uuid, Memory>>,
max_jobs: usize,
}
```
见 [manager.rs](E:/codes/Rust/projects/ironclaw/src/context/manager.rs:14)。
应用启动时直接创建一个空的 `ContextManager`:
```rust
let context_manager =
Arc::new(ContextManager::new(self.config.agent.max_parallel_jobs));
```
见 [app.rs](E:/codes/Rust/projects/ironclaw/src/app.rs:1320)。
没有发现启动时读取 `agent_jobs` 或 `job_actions` 并重新插入这两个 map 的代码。因此进程重启后:
```text
contexts = {}
memories = {}
```
数据库里的历史 Job 仍然存在,但不会自动重新成为可运行的 `ContextManager` Job。
---
## 2. JobContext 是否持久化
答案是:部分持久化,但不是完整、持续同步。
### 创建时保存一次
标准 Job 创建路径是:
```text
ContextManager 创建 JobContext
→ 设置 metadata/max_tokens/approval_context
→ store.save_job(&ctx)
→ 调度执行
```
见 [scheduler.rs](E:/codes/Rust/projects/ironclaw/src/agent/scheduler.rs:188) 和 [scheduler.rs](E:/codes/Rust/projects/ironclaw/src/agent/scheduler.rs:244)。
PostgreSQL 和 libSQL 都会写入 `agent_jobs`:
- PostgreSQL:[store.rs](E:/codes/Rust/projects/ironclaw/src/history/store.rs:146)
- libSQL:[jobs.rs](E:/codes/Rust/projects/ironclaw/src/db/libsql/jobs.rs:21)
保存的主要字段包括:
- `job_id`
- `conversation_id`
- `user_id`
- `title`、`description`、`category`
- `state`
- 预算和估算字段
- `actual_cost`
- `repair_attempts`
- `max_tokens`
- `total_tokens_used`
- 创建、开始、完成时间
### 很多字段完全不保存
以下字段在数据库读取时都被重建为空值或默认值:
- `requester_id`
- `transitions`
- `metadata`
- `extra_env`
- `http_interceptor`
- `tool_output_stash`
- `user_timezone`
- `approval_context`
例如数据库读取时明确执行:
```rust
transitions: Vec::new(),
metadata: serde_json::Value::Null,
user_timezone: "UTC".to_string(),
approval_context: None,
```
见 [jobs.rs](E:/codes/Rust/projects/ironclaw/src/db/libsql/jobs.rs:105) 和 [store.rs](E:/codes/Rust/projects/ironclaw/src/history/store.rs:293)。
所以数据库重建出来的 `JobContext` 是一个有损投影,并不是原对象的完整恢复。
### 后续修改没有完整同步
全局搜索后,生产代码中的 `save_job(&ctx)` 只出现在初次调度路径。后续通常只调用:
```rust
update_job_status(job_id, status, failure_reason)
```
这个 SQL 只更新:
- `status`
- `failure_reason`
见 [jobs.rs](E:/codes/Rust/projects/ironclaw/src/db/libsql/jobs.rs:147)。
这意味着以下内存修改通常不会同步回 `agent_jobs`:
- `transitions`
- `metadata`
- `repair_attempts`
- `total_tokens_used`
- `actual_cost`
- `started_at`
- `completed_at`
- 后设置的 `category`
例如,Job 是先以 `Pending` 状态保存,然后才在内存中转成 `InProgress`:[scheduler.rs](E:/codes/Rust/projects/ironclaw/src/agent/scheduler.rs:303)。这里没有同步执行 `update_job_status(InProgress)`。
完成时也是先更新内存中的 `completed_at`,随后数据库只更新状态:[job.rs](E:/codes/Rust/projects/ironclaw/src/worker/job.rs:1049)。
因此准确模型是:
```text
JobContext
├─ 内存对象:运行时权威状态,字段较完整
└─ agent_jobs:创建时快照 + 少量状态更新,可能滞后或缺字段
```
### JobContext 的实际用途
它承担 Job 的运行控制信息:
1. 身份与所有权
`job_id`、`user_id`、`requester_id` 用于 Job 定位、多租户隔离和权限检查。
2. 状态机
管理 `Pending → InProgress → Completed/Failed/Stuck` 等转换,见 [state.rs](E:/codes/Rust/projects/ironclaw/src/context/state.rs:299)。
3. 调度和并发限制
`ContextManager` 根据 Job 状态判断哪些 Job 正在占用并发槽。
4. Worker 初始任务描述
Worker 使用 `title` 和 `description` 构造 LLM 的初始 system message,见 [job.rs](E:/codes/Rust/projects/ironclaw/src/worker/job.rs:266)。
5. Token 和预算控制
`max_tokens`、`total_tokens_used`、`actual_cost` 用于运行期限制。
6. 工具执行上下文
`approval_context`、`extra_env`、`http_interceptor` 和 `tool_output_stash` 会影响权限检查、环境变量注入、HTTP 录制以及跨工具结果引用。
7. 卡住检测和自修复
`state`、`started_at`、`transitions`、`repair_attempts` 用于检测 Stuck Job 和执行恢复。
所以 `JobContext` 更接近:
> 单个 Job 的运行控制块和状态机,而不是 Job 的聊天记录。
---
## 3. Memory 是否持久化
答案是:`Memory` 整体不持久化。
它的定义是:
```rust
pub struct Memory {
pub job_id: Uuid,
pub conversation: ConversationMemory,
pub actions: Vec<ActionRecord>,
next_sequence: u32,
}
```
见 [memory.rs](E:/codes/Rust/projects/ironclaw/src/context/memory.rs:164)。
`Memory` 甚至没有实现 `Serialize`/`Deserialize`,数据库接口中也没有:
```rust
save_memory(...)
get_memory(...)
```
创建 Job 时只是向内存 map 插入:
```rust
memories.insert(job_id, Memory::new(job_id));
```
见 [manager.rs](E:/codes/Rust/projects/ironclaw/src/context/manager.rs:102)。
### actions 会单独持久化
Worker 每次完成工具调用后:
1. 在 `Memory` 中创建并保存 `ActionRecord`
2. 克隆该 `ActionRecord`
3. 异步写入 `job_actions`
见 [job.rs](E:/codes/Rust/projects/ironclaw/src/worker/job.rs:727)。
数据库保存内容包括:
- 工具名称
- 输入参数
- 原始输出
- 脱敏输出
- 脱敏警告
- 耗时
- 成本
- 成功状态
- 错误信息
- 执行时间
见 [store.rs](E:/codes/Rust/projects/ironclaw/src/history/store.rs:384)。
但是数据库读取出的 `ActionRecord` 不会重新放进 `Memory.actions`。项目虽然提供 `get_job_actions(job_id)`,却没有发现任何 Memory hydration 路径。
所以:
```text
Memory.actions ──复制每条记录──> job_actions
↑ │
└──── 重启后不会从数据库恢复 ────┘
```
### conversation 实际基本没用
`Memory.conversation` 提供了最多保留 100 条消息的能力,但全局搜索发现:
- `add_message()` 的调用只存在于测试
- 生产 Worker 没有向 `Memory.conversation` 写入消息
- 也没有读取 `Memory.conversation` 来调用 LLM
Worker 真正使用的对话上下文是局部变量:
```rust
let mut reason_ctx = ReasoningContext::new().with_job(&job_ctx.description);
reason_ctx.messages.push(...);
```
见 [job.rs](E:/codes/Rust/projects/ironclaw/src/worker/job.rs:273)。
因此当前代码里:
> `Memory.conversation` 不是 Job 的真实 LLM 上下文,基本属于未接入生产流程的预留或遗留设计。
### Memory 当前真正的用途
实际用途比较窄:
1. 为工具调用生成 Job 内递增的 `sequence`
2. 保存当前进程生命周期内的工具调用记录
3. 统计成功、失败的 Action 数量
4. 获取最后一次 Action
5. Job 失败或卡住时生成 `FallbackDeliverable`
Fallback 包括:
- 是否已经产生部分成果
- 最后一次操作
- 成功/失败操作数量
- Token、成本、耗时
- 修复次数
见 [fallback.rs](E:/codes/Rust/projects/ironclaw/src/context/fallback.rs:58)。
但这里还有一个持久化缺口:Fallback 最终放入 `JobContext.metadata["fallback_deliverable"]`,而 `metadata` 没有保存到 `agent_jobs`。因此它也只在当前进程的 `ContextManager` 中可用。
另外,`total_cost()`、`total_duration()`、`actions_by_tool()` 在生产代码中没有调用;`ActionRecord::with_cost()` 也只在测试中调用。当前真正用于 fallback 的主要是:
- `successful_actions()`
- `failed_actions()`
- `last_action()`
---
## 4. 是否跨 Job 共享
不共享。每创建一个 Job,就按 `job_id` 创建一份独立 Memory:
```text
ContextManager
├─ job_a → JobContext A + Memory A
├─ job_b → JobContext B + Memory B
└─ job_c → JobContext C + Memory C
```
不同 Job 无法通过 `ContextManager::Memory` 共享信息。
真正负责跨 Job、跨会话长期记忆的是数据库中的:
- `memory_documents`
- `memory_chunks`
- `memory_search`
- `memory_write`
那是 workspace memory 子系统,与这里的 `context::Memory` 没有直接关系。
最后可以用一句话概括:
> `JobContext` 是运行控制状态,数据库只保存它的有损快照;`Memory` 是当前进程内的 Job 工具执行历史,本身不持久化、不跨 Job,且其中的 conversation 当前没有接入真实 LLM 上下文。二. 短期长期记忆写入时机 看完所有生产调用路径后,写入时机如下。需要区分:
- 写入 `ContextManager` 内存
- 写入数据库
- 写入 Job 的 `Memory.actions`
## 普通 Scheduler Job
完整时间线:
```text
dispatch_job()
1. 创建并插入 JobContext(Pending)
2. 创建并插入空 Memory
3. 更新 JobContext 的 metadata/max_tokens/approval_context
4. 将 JobContext 快照写入 agent_jobs
5. 内存 JobContext 转为 InProgress
6. Worker 执行
├─ LLM 返回后更新 JobContext.total_tokens_used
└─ 工具执行结束后写 Memory.actions
└─ 再异步写 job_actions
7. 完成/失败/卡住
├─ 更新内存 JobContext
└─ 异步更新数据库 status
```
### 1. JobContext 首次写入
调用:
```rust
context_manager.create_job_for_user(...)
```
内部构造一个初始状态为 `Pending` 的 `JobContext`,然后写入:
```rust
contexts.insert(job_id, context);
```
见 [manager.rs](E:/codes/Rust/projects/ironclaw/src/context/manager.rs:48) 和 [manager.rs](E:/codes/Rust/projects/ironclaw/src/context/manager.rs:87)。
紧接着创建空 Memory:
```rust
memories.insert(job_id, Memory::new(job_id));
```
注意实际顺序是:
```text
先插入 JobContext
释放 contexts 写锁
再插入 Memory
```
因此这两个 map 的写入不是一个原子操作。正常情况下几乎连续完成,但理论上存在很短的窗口:
```text
get_context(job_id) 成功
get_memory(job_id) 仍然 NotFound
```
### 2. 调度前第二次写 JobContext
Scheduler 创建 Job 后,会立即写入:
- `metadata`
- `max_tokens`
- `approval_context`
使用的是原子更新并返回:
```rust
update_context_and_get(job_id, |ctx| {
ctx.metadata = ...;
ctx.max_tokens = ...;
ctx.approval_context = ...;
})
```
见 [scheduler.rs](E:/codes/Rust/projects/ironclaw/src/agent/scheduler.rs:219)。
### 3. 首次数据库写入
上述字段更新完成后,调用:
```rust
store.save_job(&ctx).await
```
见 [scheduler.rs](E:/codes/Rust/projects/ironclaw/src/agent/scheduler.rs:244)。
这是同步等待的数据库写入,目的是先建立 `agent_jobs` 行,确保后续 `job_actions` 和 `llm_calls` 的外键存在。
此时保存到数据库的状态仍然是:
```text
Pending
```
### 4. 调度时写为 InProgress
数据库保存完成后,Scheduler 调用:
```rust
ctx.transition_to(JobState::InProgress, ...)
```
见 [scheduler.rs](E:/codes/Rust/projects/ironclaw/src/agent/scheduler.rs:303)。
这会写入内存 JobContext 的:
- `state = InProgress`
- `started_at`
- `transitions`
但这里没有调用数据库 `update_job_status(InProgress)`。
因此 Job 正在运行时可能出现:
```text
ContextManager:InProgress
agent_jobs:Pending
```
### 5. LLM 调用后写 Token
只有 `respond_with_tools()` 成功并返回 TokenUsage 后,才会调用:
```rust
ctx.add_tokens(total_tokens)
```
见 [job.rs](E:/codes/Rust/projects/ironclaw/src/worker/job.rs:1569)。
它更新:
- `total_tokens_used`
- 并检查 `max_tokens`
注意:
- `select_tools()` 的 Token 没有统计,因为接口没有返回 TokenUsage。
- 更新后的 `total_tokens_used` 没有再次调用 `save_job()`。
- 因此数据库中的 `total_tokens_used` 通常仍是创建时的值。
### 6. 工具执行结束后写 Memory
工具执行流程是:
```text
权限检查
→ 限流检查
→ Hook
→ Job 取消检查
→ 参数校验
→ tool.execute()
→ 成功/失败/超时
→ 写 Memory.actions
```
只有真正进入 `tool.execute()` 并获得以下结果之一后,才写 Memory:
- 工具成功
- 工具返回错误
- 工具执行超时
写入逻辑见 [job.rs](E:/codes/Rust/projects/ironclaw/src/worker/job.rs:727)。
成功时:
```rust
let rec = mem
.create_action(tool_name, safe_params)
.succeed(sanitized_output, raw_output, elapsed);
mem.record_action(rec.clone());
```
失败时:
```rust
let rec = mem
.create_action(...)
.fail(error, elapsed);
mem.record_action(rec.clone());
```
### 不会写 Memory 的失败
以下错误发生在 `tool.execute()` 之前,因此不会生成 `ActionRecord`:
- Tool 不存在
- Approval 拒绝
- Rate limit
- BeforeToolCall Hook 拒绝
- Job 已取消
- 参数校验失败
也就是说,`Memory.actions` 不是“所有工具尝试”,而是:
> 已经到达实际工具执行边界的调用记录。
### 7. Memory 写完后异步写数据库
先完成:
```rust
mem.record_action(...)
```
然后才:
```rust
tokio::spawn(async move {
store.save_action(job_id, &action).await
});
```
见 [job.rs](E:/codes/Rust/projects/ironclaw/src/worker/job.rs:794)。
顺序是:
```text
Memory.actions 写入完成
→ 启动异步 DB 任务
→ Worker 继续处理结果
```
所以:
- Memory 写入是 awaited,完成后才能继续。
- `job_actions` 数据库写入是 fire-and-forget。
- 进程如果在异步任务完成前退出,可能出现 Memory 中有 Action,但 DB 中没有。
并行调用工具时,每个工具完成后分别抢 `memories` 写锁,因此 Action 的 `sequence` 更接近完成顺序,不保证等于工具提交顺序。
---
## Job 进入终态时
### Completed
Job 完成时更新内存:
```rust
ctx.transition_to(JobState::Completed, ...)
```
见 [job.rs](E:/codes/Rust/projects/ironclaw/src/worker/job.rs:1049)。
它会写入:
- `state`
- `completed_at`
- `transitions`
随后异步调用:
```rust
update_job_status(Completed, ...)
```
数据库只更新:
- `status`
- `failure_reason`
不会同步 `completed_at` 和 `transitions`。
### Failed / Stuck
失败或卡住时,顺序有所不同:
```text
1. 先读取 Memory
2. 读取当前 JobContext
3. 用二者生成 FallbackDeliverable
4. 更新 JobContext.state
5. 将 fallback 写入 JobContext.metadata
6. 异步更新数据库 status
```
见 [job.rs](E:/codes/Rust/projects/ironclaw/src/worker/job.rs:1078) 和 [job.rs](E:/codes/Rust/projects/ironclaw/src/worker/job.rs:1136)。
所以 Memory 的主要消费时机就是 Job 失败或卡住时。
不过:
- `fallback_deliverable` 只写入内存 `JobContext.metadata`
- `metadata` 不会通过 `update_job_status()` 保存
- 因此进程重启后 fallback 会丢失
### Cancelled
取消 Job 时:
```text
更新内存 JobContext → Cancelled
异步更新数据库 status → Cancelled
```
见 [scheduler.rs](E:/codes/Rust/projects/ironclaw/src/agent/scheduler.rs:619)。
### Self-repair
自修复会写入内存 JobContext:
- `InProgress → Stuck`
- `Stuck → InProgress`
- `repair_attempts += 1`
- 超过次数后转为 `Failed`
见 [self_repair.rs](E:/codes/Rust/projects/ironclaw/src/agent/self_repair.rs:157)。
这些路径没有同步保存完整 JobContext;部分状态甚至没有对应的数据库更新。
---
## Sandbox Job
Sandbox Job 的时序不同。
### 1. 先注册 ContextManager
容器创建前调用:
```rust
register_sandbox_job(job_id, ...)
```
见 [job.rs](E:/codes/Rust/projects/ironclaw/src/tools/builtin/job.rs:446)。
它直接创建:
```text
JobContext.state = InProgress
JobContext.started_at = now
Memory = 空
```
也就是说 Sandbox Job 没有普通 Job 的 `Pending → InProgress` 调度过程。
### 2. 再异步保存 SandboxJobRecord
注册 ContextManager 后调用:
```rust
self.persist_job(SandboxJobRecord { ... })
```
见 [job.rs](E:/codes/Rust/projects/ironclaw/src/tools/builtin/job.rs:456)。
`persist_job()` 内部使用 `tokio::spawn`,所以这是异步数据库写入,不等待完成。
数据库写的是 `sandbox_jobs`,不是普通 Job 的 `agent_jobs`。
### 3. 容器运行期间
Sandbox Job 状态由两套路径分别更新:
- `ContextManager.JobContext`:供 `list_jobs`、`job_status`、并发限制等查询
- `sandbox_jobs`:持久化容器状态
容器完成事件到达后,Job monitor 把内存状态转为 `Completed` 或 `Failed`,见 [job_monitor.rs](E:/codes/Rust/projects/ironclaw/src/agent/job_monitor.rs:103)。
### 4. Sandbox 的 Memory 基本为空
虽然注册 Sandbox Job 时也会创建一份 `Memory`,但容器内部工具调用不经过:
```rust
Worker::execute_tool_inner()
```
因此不会写入这份 `Memory.actions`。
Sandbox 的执行历史主要进入:
- `sandbox_jobs`
- `job_events`
- 容器事件流
所以 Sandbox Job 对应的 `ContextManager::Memory` 当前基本只是伴随 JobContext 被创建,没有实际 Action 记录用途。
---
## 写入时机汇总
| 时机 | JobContext 内存 | Memory 内存 | 数据库 |
|---|---|---|---|
| 普通 Job 创建 | 插入 `Pending` | 插入空 Memory | 无 |
| 调度前配置 | 写 metadata、Token 上限、权限 | 无 | 随后 `save_job` |
| DB 初次保存 | 无 | 无 | 写 `agent_jobs`,状态通常为 `Pending` |
| 开始调度 | 转为 `InProgress` | 无 | 没有同步 InProgress |
| LLM 返回 | 累加 Token | 无 | 没有同步完整 Job |
| 工具执行前检查失败 | 无 | 不写 | 不写 `job_actions` |
| 工具成功/失败/超时 | 通常不写 | 追加 Action | 随后异步写 `job_actions` |
| Job Completed | 状态、时间、transition | 不写 | 异步只更新 status |
| Job Failed/Stuck | 状态、时间、fallback metadata | 先读取生成 fallback | 异步只更新 status/reason |
| Self-repair | 状态、repair attempts | 不写 | 通常不同步 |
| Sandbox 创建 | 直接插入 `InProgress` | 插入空 Memory | 异步写 `sandbox_jobs` |
| Sandbox 工具执行 | 通常不写 | 不记录容器 Action | 写 `job_events` |
| Sandbox 完成 | 转为 Completed/Failed | 不写 | 单独更新 `sandbox_jobs` |
最后一个重要事实:生产代码没有调用 `ContextManager::remove_job()`;它目前只在测试中使用。因此普通 Job 完成后,`JobContext` 和 `Memory` 不会立即从 map 删除,而是通常保留到进程退出。容器 Job 的 JobContext 由 Host 侧 ExecuteJobTool::execute_sandbox() 注册,ContextManager 实际构造和保存;容器只上报事件,Host 再据此更新状态。
容器和 Host 之间主要通过一个“Host 提供、容器主动调用”的内部 HTTP API 通信,不是共享内存,也不是 Host 直接连接容器。
```text
Docker bridge 网络
┌────────────────┐ HTTP + JSON + Bearer Token ┌──────────────────┐
│ Container │ ───────────────────────────────► │ Host Orchestrator│
│ Worker/Bridge │ │ Axum :50051 │
└────────────────┘ ◄─────────────────────────────── └──────────────────┘
│ HTTP Response / Prompt
│
└── /workspace bind mount 与 Host 共享项目文件
```
## 1. Host 启动内部 API
Host 启动 `OrchestratorApi`,默认端口是 `50051`:
- Windows/macOS:监听 `127.0.0.1:50051`
- Linux:监听 `0.0.0.0:50051`
见 [api.rs](E:/codes/Rust/projects/ironclaw/src/orchestrator/api.rs:146)。
容器通过:
```text
http://host.docker.internal:50051
```
访问 Host。
Docker 配置还显式设置:
```rust
extra_hosts: Some(vec![
"host.docker.internal:host-gateway".to_string()
])
```
见 [job_manager.rs](E:/codes/Rust/projects/ironclaw/src/orchestrator/job_manager.rs:420)。
因此 Linux Docker bridge 下,`host.docker.internal` 会解析到 Host gateway。
## 2. 创建容器时注入连接信息
Host 创建容器前,为每个 Job 生成独立 Token:
```rust
let token = self.token_store.create_token(job_id).await;
```
然后向容器注入三个环境变量:
```text
IRONCLAW_WORKER_TOKEN=<随机 Token>
IRONCLAW_JOB_ID=<Job UUID>
IRONCLAW_ORCHESTRATOR_URL=http://host.docker.internal:50051
```
见 [job_manager.rs](E:/codes/Rust/projects/ironclaw/src/orchestrator/job_manager.rs:352) 和 [job_manager.rs](E:/codes/Rust/projects/ironclaw/src/orchestrator/job_manager.rs:433)。
Token 特性:
- 32 字节随机数,hex 编码
- 每个 Job 独立
- 只允许访问对应 `job_id`
- 仅保存在 Host 内存
- 容器停止后撤销
- 不记录日志、不写数据库
鉴权实现在 [auth.rs](E:/codes/Rust/projects/ironclaw/src/orchestrator/auth.rs:117)。
容器每次请求都会携带:
```http
Authorization: Bearer <IRONCLAW_WORKER_TOKEN>
```
URL 中的 `job_id` 必须与 Token 绑定的 Job 一致。
## 3. HTTP API
Host 暴露以下接口:
| 方向 | Endpoint | 用途 |
|---|---|---|
| Container → Host | `GET /worker/{id}/job` | 获取任务描述 |
| Container → Host | `GET /worker/{id}/credentials` | 获取授权给该 Job 的凭据 |
| Container → Host | `POST /worker/{id}/llm/complete` | 通过 Host 调用 LLM |
| Container → Host | `POST /worker/{id}/llm/complete_with_tools` | 通过 Host 调用带工具的 LLM |
| Container → Host | `POST /worker/{id}/status` | 上报运行状态和迭代次数 |
| Container → Host | `POST /worker/{id}/event` | 上报消息、工具调用和结果事件 |
| Container → Host | `POST /worker/{id}/complete` | 告知 Host Job 已结束 |
| Container → Host | `GET /worker/{id}/prompt` | 轮询 Host 的后续 Prompt |
路由定义见 [api.rs](E:/codes/Rust/projects/ironclaw/src/orchestrator/api.rs:117)。
容器端统一通过 `WorkerHttpClient` 调用,见 [api.rs](E:/codes/Rust/projects/ironclaw/src/worker/api.rs:122)。
---
## 4. 容器启动后如何获取任务
容器里的可执行程序由模式决定:
```text
Worker mode → worker
Claude Code mode → claude-bridge
ACP mode → acp-bridge
```
启动参数包含:
```text
--job-id <uuid>
--orchestrator-url http://host.docker.internal:50051
```
启动后首先调用:
```http
GET /worker/{job_id}/job
```
Host 从 `ContainerJobManager` 的 `ContainerHandle` 中读取任务描述并返回:
```json
{
"title": "Job <uuid>",
"description": "...",
"project_dir": "..."
}
```
见 [api.rs](E:/codes/Rust/projects/ironclaw/src/orchestrator/api.rs:184)。
容器还会调用 `/credentials` 获取该 Job 显式授权的 Secret,并将其注入子进程环境变量。
## 5. 容器如何上报事件
容器将事件封装为:
```rust
pub struct JobEventPayload {
pub event_type: String,
pub data: serde_json::Value,
}
```
然后发送:
```http
POST /worker/{job_id}/event
Content-Type: application/json
Authorization: Bearer ...
{
"event_type": "tool_use",
"data": {
"tool_name": "shell",
"input": "..."
}
}
```
常见事件包括:
- `message`
- `tool_use`
- `tool_result`
- `status`
- `reasoning`
- `result`
Generic Worker 在每次消息和工具调用前后发送事件,见 [container.rs](E:/codes/Rust/projects/ironclaw/src/worker/container.rs:349)。
Claude Code bridge 则读取 `claude --output-format stream-json` 输出的 NDJSON,转换成 `JobEventPayload` 后逐条 POST 给 Host,见 [claude_bridge.rs](E:/codes/Rust/projects/ironclaw/src/worker/claude_bridge.rs:439)。
ACP bridge 同样把 ACP session notification 转换为 HTTP 事件。
## 6. Host 收到事件后做什么
`POST /event` 到达 Host 后执行两条路径:
```text
┌─► job_events 数据库
Container POST /event ───┤
└─► broadcast channel
├─► JobMonitor
└─► Web Gateway SSE
```
### 持久化
Host 异步调用:
```rust
store.save_job_event(job_id, event_type, data)
```
见 [api.rs](E:/codes/Rust/projects/ironclaw/src/orchestrator/api.rs:318)。
这是 fire-and-forget,因此不会阻塞容器继续运行。
### 转换为 AppEvent
Host 将 wire event 转换为强类型事件:
```text
message → AppEvent::JobMessage
tool_use → AppEvent::JobToolUse
tool_result → AppEvent::JobToolResult
result → AppEvent::JobResult
reasoning → AppEvent::JobReasoning
其他 → AppEvent::JobStatus
```
然后发送到容量为 256 的 Tokio `broadcast` channel。
### JobMonitor 更新 JobContext
`JobMonitor` 订阅该 broadcast:
- Assistant message:注入主 Agent 消息队列
- `JobResult`:更新容器 Job 的内存 `JobContext`
- 成功结果:`Completed`
- 失败结果:`Failed`
见 [job_monitor.rs](E:/codes/Rust/projects/ironclaw/src/agent/job_monitor.rs:59)。
这就是:
```text
容器事件 → Host HTTP API → broadcast → JobMonitor → ContextManager.update_context()
```
### WebUI SSE
Web Gateway 也订阅同一个 broadcast,然后按 `user_id` 将事件发送给对应用户的 SSE 连接,见 [main.rs](E:/codes/Rust/projects/ironclaw/src/main.rs:1024)。
---
## 7. Host 如何给容器发送 Prompt
Host 不主动连接容器,也没有 WebSocket 推送。
流程是:
```text
用户/主 Agent
→ job_prompt 工具
→ 写入 Host PromptQueue
→ 容器定期 GET /worker/{id}/prompt
→ Host 从队列 pop_front()
→ 容器加入自己的 ReasoningContext
```
PromptQueue 是:
```rust
Arc<Mutex<HashMap<Uuid, VecDeque<PendingPrompt>>>>
```
写入见 [job.rs](E:/codes/Rust/projects/ironclaw/src/tools/builtin/job.rs:1668),读取接口见 [api.rs](E:/codes/Rust/projects/ironclaw/src/orchestrator/api.rs:468)。
不同模式的轮询方式:
- Generic Worker:每次 LLM 调用前轮询一次
- Claude Code:每 2 秒轮询,拿到 Prompt 后使用 `claude --resume`
- ACP:在 follow-up loop 中轮询,然后发起新的 ACP `prompt()`
所以 Host → Container 是 pull 模式,不是 push 模式。
## 8. 完成和取消
### 容器主动完成
容器通常先发送终态 `result` 事件,再调用:
```http
POST /worker/{id}/complete
```
`/complete` 会:
1. 保存 `CompletionResult`
2. 停止并删除 Docker 容器
3. 撤销 Bearer Token
4. 保留短期 ContainerHandle,供等待方读取结果
见 [api.rs](E:/codes/Rust/projects/ironclaw/src/orchestrator/api.rs:272) 和 [job_manager.rs](E:/codes/Rust/projects/ironclaw/src/orchestrator/job_manager.rs:700)。
### Host 主动取消
Host 不向容器发送“取消 HTTP 消息”,而是直接通过 Docker API:
```text
stop_container
→ remove_container
→ revoke token
```
因此取消属于 Host 对容器生命周期的强制控制。
## 9. 还有两个非 HTTP 通道
除了内部 HTTP,还有:
- 项目目录通过 bind mount 映射为容器内 `/workspace`
- MCP 配置、Claude 认证配置等通过只读 mount 传入
所以整体通信模型是:
```text
控制面:HTTP API + Bearer Token
事件面:HTTP POST → DB + broadcast
Prompt:Host 内存队列 + Container HTTP polling
文件面:Docker bind mount /workspace
生命周期:Host 通过 Docker API 控制
```
另外代码里有一个值得注意的协议不一致:Generic Worker 的成功 `result` 事件发送的是 `{"success": true}`,而 Host `job_event_handler` 读取的是 `data.status`。缺少 `status` 时会默认转换为 `Failed`。Claude Code 和 ACP 会发送 `status`,Generic Worker 这条路径看起来存在终态误判风险。三. workJob 以上就是完整的分析。核心要点总结:
1. **两条路径,一个引擎** — 无论 in-process 还是 Docker 容器,最终都通过 `run_agentic_loop()` + `LoopDelegate` trait 执行,三条 delegate 实现:`JobDelegate`(路径A)、`ContainerDelegate`(路径B Worker模式)、Claude Code bridge(路径B ClaudeCode模式)
2. **路径A 通信靠 mpsc** — Scheduler spawn tokio task,通过 `mpsc::channel<WorkerMessage>` 发送 Start/Stop/UserMessage 信号;Worker 通过 SSE broadcast + DB 持久化向外报告
3. **路径B 通信靠 HTTP + broadcast** — 容器通过 `WorkerHttpClient`(bearer token 认证)与 Orchestrator Internal API 通信;Orchestrator 将事件广播到 `broadcast::Sender`,`JobMonitor` 订阅后通过 mpsc 注入回主 Agent Loop
4. **安全边界清晰** — 容器内无 API key(LLM 调用全走代理)、凭证通过环境变量注入(非全局环境)、token per-job 隔离且 constant-time 比较、容器本身的 Docker 安全加固(cap-drop ALL、非 root、内存限制)## Job 异常处理与超时机制
IronClaw 有 **多层防护**,从内到外依次作用:
---
## 一、多层超时体系
### 第 1 层:单次工具调用超时(60s 默认)
**`src/tools/tool.rs:422-424`** — `Tool` trait 默认值:
```rust
fn execution_timeout(&self) -> Duration {
Duration::from_secs(60) // 每个工具可覆盖
}
```
**`src/worker/job.rs:688-694`** — Worker 中的实际执行:
```rust
let tool_timeout = tool.execution_timeout();
let result = tokio::time::timeout(tool_timeout, async {
tool.execute(effective_params.clone(), &job_ctx).await
}).await;
```
超时后返回 `ToolError::Timeout { name, timeout }`——**不会终止整个 job**,而是作为工具失败结果返回给 LLM,让 LLM 尝试其他方式。
`ToolDispatcher` 路径同样有超时:`src/tools/dispatch.rs:197`
### 第 2 层:Job 整体超时(`AGENT_JOB_TIMEOUT_SECS`,测试默认 30s)
**`src/config/agent.rs:14,59,92-96`**:
```rust
pub job_timeout: Duration, // 来自 AGENT_JOB_TIMEOUT_SECS
// 默认值:测试 30s,生产通过 env 配置
```
**`src/worker/job.rs:288-293`** — Worker 主循环:
```rust
let result = tokio::time::timeout(self.timeout(), async {
self.execution_loop(&mut rx, &reasoning, &mut reason_ctx).await
}).await;
match result {
Ok(Ok(())) => { /* 正常完成 */ }
Ok(Err(e)) => { self.mark_failed(&e.to_string()).await?; } // 执行错误
Err(_) => { // 超时!
self.mark_stuck("Execution timeout").await?; // 标记为 Stuck
}
}
```
**关键:超时不是直接 Failed,而是 Stuck**。Stuck 状态允许后续 self-repair 恢复。
### 第 3 层:容器整体超时(600s 默认)
**`src/worker/container.rs:40,49`**:
```rust
pub timeout: Duration, // 容器总超时
// 默认:Duration::from_secs(600) = 10 分钟
```
**`src/worker/container.rs:180`**:
```rust
let result = tokio::time::timeout(self.config.timeout, async {
// 整个 ContainerDelegate + run_agentic_loop
}).await;
```
超时后向 Orchestrator 上报 `CompletionReport { success: false, message: "Execution timed out" }`。
### 第 4 层:Claude Code Bridge 超时(1800s 默认)
**`src/worker/mod.rs:140`**:
```rust
timeout: std::time::Duration::from_secs(1800), // 30 分钟
```
### 第 5 层:ACP Bridge 超时(`ACP_TIMEOUT_SECS` 环境变量,默认 1800s)
**`src/worker/mod.rs:42-44`**:
```rust
fn acp_bridge_timeout() -> std::time::Duration {
std::time::Duration::from_secs(crate::config::AcpModeConfig::from_env().timeout_secs)
}
```
---
## 二、迭代次数限制
**`src/worker/job.rs:343-349`**:
```rust
let max_iterations = self.context_manager().get_context(self.job_id).await
.ok().and_then(|ctx| ctx.metadata.get("max_iterations").and_then(|v| v.as_u64()))
.unwrap_or(50) as usize;
let max_iterations = max_iterations.min(ironclaw_common::MAX_WORKER_ITERATIONS as usize);
```
服务端硬上限 `MAX_WORKER_ITERATIONS`,且在容器创建时也有 clamp:
```rust
let capped = iters.clamp(1, MAX_WORKER_ITERATIONS); // src/orchestrator/job_manager.rs:453
```
**`src/worker/job.rs:443-446`** — 超出后:
```rust
LoopOutcome::MaxIterations => {
self.mark_failed("Maximum iterations exceeded: job hit the iteration cap").await?;
}
```
---
## 三、Token 预算
**`src/context/state.rs:380-390`** — 每次 LLM 调用后追踪:
```rust
pub fn add_tokens(&mut self, tokens: u64) -> Result<(), TokenBudgetExceeded> {
self.total_tokens_used += tokens;
if self.max_tokens > 0 && self.total_tokens_used > self.max_tokens {
Err(TokenBudgetExceeded { used: self.total_tokens_used, limit: self.max_tokens })
} else {
Ok(())
}
}
```
**`src/worker/job.rs:1561-1573`** — Worker 中检测:
```rust
if total_tokens > 0
&& let Err(err) = self.context_manager()
.update_context(self.job_id, |ctx| ctx.add_tokens(total_tokens))
.await?
{
self.worker.mark_failed(&err.to_string()).await?;
}
```
配置来源 `AGENT_MAX_TOKENS_PER_JOB`(`src/config/agent.rs:39,151-155`),`0` 表示不限制。
---
## 四、状态机与异常流转
**`src/context/state.rs:49-75`** — 合法状态转移:
```
Pending ──────────► InProgress ──────────► Completed ──────────► Submitted ──────────► Accepted
│ │ │ │ │
│ │ │ │ │
├──► Failed │ ├──► Failed ├──► Failed ├──► Failed
├──► Cancelled │ ├──► Stuck │ │
│ ├──► Cancelled │ │
│ │
│ Stuck ◄──┘ │
│ │ │
│ ├──► InProgress (recovery)
│ ├──► Failed
│ └──► Cancelled
```
**`src/context/state.rs:78-80`** — 终态(不能再转移):
```rust
pub fn is_terminal(&self) -> bool {
matches!(self, Self::Accepted | Self::Failed | Self::Cancelled)
}
```
**注意 `Stuck` 不是终态**——`src/worker/job.rs:1411-1413` 明确写了:
```rust
// job has been cancelled, failed, or already completed — but NOT when Stuck,
// because Stuck is recoverable (Stuck -> InProgress via self-repair).
// Stopping on Stuck would prevent recovery from resuming the worker (issue #892).
```
---
## 五、Self-Repair 自愈机制
### 检测周期
**`src/agent/self_repair.rs:469-479`**:
```rust
pub struct RepairTask {
repair: Arc<dyn SelfRepair>,
check_interval: Duration, // 默认 3600s (1小时),来自 SELF_REPAIR_CHECK_INTERVAL_SECS
}
```
### 从何检测 Stuck Job
两种来源:
1. **状态已是 `Stuck`** — 运行中的 job 超时 → `mark_stuck()` 转换
2. **`InProgress` 但运行超时** — `detect_stuck_jobs()` 用 `stuck_threshold` 找到这些 job,主动将其转 `InProgress → Stuck`
**`src/agent/self_repair.rs:148-178`**:
```rust
let just_transitioned = ctx.state == JobState::InProgress;
if just_transitioned {
let reason = "exceeded stuck_threshold";
// 先转 Stuck,再尝试恢复
context_manager.update_context(job_id, |ctx| ctx.mark_stuck(reason)).await;
}
```
### Stuck 时长从何时算起
**关键设计**:stuck_duration 从 **最后一次 Stuck 转换的时间戳** 算起,而非 `started_at`。
**`src/agent/self_repair.rs:190-202`**:
```rust
let stuck_since = ctx.transitions.iter().rev()
.find(|t| t.to == JobState::Stuck)
.map(|t| t.timestamp);
```
这意味着一个跑了 2 小时的 job 刚刚变 Stuck,不会立即被认为超过 5 分钟阈值。
### 修复流程
**`src/agent/self_repair.rs:224-301`** — `repair_stuck_job()`:
```
1. 检查 repair_attempts >= max_repair_attempts?
└─ 是 → 转 Failed + 返回 ManualRequired(不再尝试)
└─ 否 → 继续
2. 调用 ctx.attempt_recovery() ← Stuck → InProgress, repair_attempts += 1
└─ 成功 → RepairResult::Success(恢复成功,worker 继续跑)
└─ 失败 → RepairResult::Retry(下次周期再试)
```
**`src/agent/self_repair.rs:485-503`** — 周期性检测+修复:
```rust
loop {
tokio::time::sleep(self.check_interval).await;
let stuck_jobs = self.repair.detect_stuck_jobs().await;
for job in stuck_jobs {
match self.repair.repair_stuck_job(&job).await {
Ok(RepairResult::Success { .. }) => { /* 恢复成功 */ }
Ok(RepairResult::Retry { .. }) => { /* 下次再试,不通知用户避免骚扰 */ }
Ok(RepairResult::Failed { .. }) => { /* 记 error */ }
Ok(RepairResult::ManualRequired { .. }) => { /* 通知用户 */ }
// ...
}
}
}
```
**`ManulRequired` 不通知用户**——实际只记 warn 日志(`src/agent/self_repair.rs:496-498`)。
---
## 六、其他异常处理
### LLM 调用异常
**`src/worker/job.rs:1512-1586`** — `JobDelegate::call_llm()`:
| 异常类型 | 处理方式 |
| -------------------------------- | ---------------------------------------------- |
| `RateLimited` | 等待 retry_after,连续 10 次 → `mark_failed()` |
| `EmptyResponse` + 之前有文本输出 | 视为完成,`mark_completed()` |
| `EmptyResponse` + 无文本输出 | 重试 |
| `AuthFailed` / `Http` / `Io` 等 | 直接传播错误 |
### 容器内部异常自适应恢复(AutonomousRecovery)
**`src/worker/job.rs:29-32`** — 针对 LLM 的异常行为:
| 模式 | 症状 | 操作 |
| ------------------- | ------------------ | -------------------------------------------------- |
| `ToolModeNudge` | LLM 返回空工具调用 | 注入 `EMPTY_TOOL_COMPLETION_NUDGE`,重试用工具模式 |
| `ForceTextRecovery` | 连续多次空工具调用 | 清空 `available_tools`,强制纯文本响应 |
| `Fail` | 反复产生无效响应 | `LoopOutcome::Failure`,job 标记失败 |
### 重复失败工具调用检测
**`src/agent/agentic_loop.rs:161-208`** — `DuplicateToolCallTracker`:
- 指纹 `(tool_name, canonicalized_args)` 连续相同 + 全部失败
- 第 3 次注入 `DUPLICATE_TOOL_CALL_WARNING` 提示 LLM 换方法
- 第 5 次 `force_text = true` 强制纯文本模式
### 工具执行安全防护
**`src/worker/job.rs:527-584`** — `execute_tool_inner()` 的完整 pipeline:
```
1. 工具查找 → NotFound 错误
2. 参数归一化
3. Approval 双向检查(job-level + worker-level)
4. 速率限制检查
5. BeforeToolCall hook(可能拒绝)
6. 状态检查(是否已 Cancel)
7. 参数安全校验(Safety::validator)
8. 敏感参数 redact
9. 执行 + per-tool timeout
10. 结果记录(action_record → DB)
```
### 容器断联
Worker 内的 HTTP 调用都是 fallible——`WorkerHttpClient` 的每个方法返回 `WorkerError`。如果 Orchestrator 不可达,容器内的 AgenticLoop 会因 LLM 调用失败而报错,最终容器退出时 Orchestrator 检测到容器停止,通过 `ContainerJobManager` 清理。
---
## 七、全部可配置异常参数
| 参数 | 环境变量 | 默认值 | 作用 |
| --------------- | --------------------------------- | ------- | ------------------------------------ |
| Job 超时 | `AGENT_JOB_TIMEOUT_SECS` | 测试30s | 整个 Worker::run() 的 tokio::timeout |
| 工具超时 | trait 方法覆盖 | 60s | 单次 `tool.execute()` 超时 |
| 容器超时 | hardcoded | 600s | `WorkerRuntime::run()` |
| Claude Bridge | hardcoded | 1800s | Claude Code 桥接 |
| ACP Bridge | `ACP_TIMEOUT_SECS` | 1800s | ACP 代理桥接 |
| 最大迭代 | metadata `max_iterations` | 50 | AgenticLoop 迭代上限 |
| 服务端上限 | hardcoded `MAX_WORKER_ITERATIONS` | — | 硬上限,不可超越 |
| Token 预算 | `AGENT_MAX_TOKENS_PER_JOB` | 0(不限) | 单job token 消耗上限 |
| Stuck 阈值 | `AGENT_STUCK_THRESHOLD_SECS` | 300s | InProgress多久算Stuck |
| 修复间隔 | `SELF_REPAIR_CHECK_INTERVAL_SECS` | 3600s | 自愈检查周期 |
| 最大修复次数 | `SELF_REPAIR_MAX_ATTEMPTS` | — | 超过后转ManualRequired→Failed |
| 最大并行 Job | `AGENT_MAX_PARALLEL_JOBS` | 1 | 同进程并行上限 |
| 每用户 Job 上限 | `MAX_JOBS_PER_USER` | 不限 | 用户级并发限制 |