Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions crates/tw-api/msg-codes.txt
Original file line number Diff line number Diff line change
Expand Up @@ -288,6 +288,7 @@ gw.plugin.failed
gw.plugin.file_changed
gw.plugin.manifest
gw.plugin.memory_limit
gw.plugin.model_not_allowed
gw.plugin.not_located
gw.plugin.nothing_to_try
gw.plugin.output_limit
Expand Down
20 changes: 13 additions & 7 deletions crates/tw-api/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -678,9 +678,10 @@ pub const MSG_CODES: &str = include_str!("../msg-codes.txt");
///
/// **33 起有脚本插件**:`/plugins` 一组端点(列表、试编、装、改、换源码、看改动、批准、
/// 排顺序、删、试跑、日志),事件多了 [`Event::PluginFailed`](插件在请求上出错,或者
/// 文件变了、加载不了而停用),[`RequestDetail`] 多了 `plugins`(每一次运行)和
/// `request_after_plugins`(插件改过的请求体),[`HistoryRow`] 多了 `plugin_changed`。
/// 装、换源码、批准三个端点不给网页调:要在系统的确认框里点头。照 32 写的界面看不到插件。
/// 文件变了、加载不了而停用),[`RequestDetail`] 多了 `plugins`(每一次运行,带着跑在
/// 尝试链的第几跳)和 `request_after_plugins`(插件改过的请求体),[`HistoryRow`] 多了
/// `plugin_changed`。装、换源码、批准三个端点不给网页调:要在系统的确认框里点头。照 32
/// 写的界面看不到插件。
pub const CONTROL_API_VERSION: u32 = 33;

#[derive(Debug, Clone, Serialize, Deserialize)]
Expand Down Expand Up @@ -3696,10 +3697,12 @@ pub struct RequestDetail {
pub row: HistoryRow,
/// 客户端发来的原样
pub request_body: Option<BodyView>,
/// 插件改过之后、发往上游的那一份。**只有插件改了请求才有**
/// 插件改过之后、发往上游的那一份:最后发出去的那一跳收到的(回答的那一家收到的就是
/// 它)。**只有插件改了那一跳的请求才有**
pub request_after_plugins: Option<BodyView>,
pub response_body: Option<BodyView>,
/// 插件在这个请求上的每一次运行,按先后(请求钩子在前,回答钩子在后)
/// 插件在这个请求上的每一次运行,按先后:每一跳的请求钩子,回答那一跳的回答钩子。
/// 按 [`PluginRunView::attempt`] 对着尝试链分组
pub plugins: Vec<PluginRunView>,
/// 这个请求还在跑。**记录在结局到了才落库**,这时的 `row` 是到目前为止
/// 知道的那些:开始时的身份和上游,响应头到了就有状态码,路由走完就有
Expand Down Expand Up @@ -4990,9 +4993,9 @@ impl SettingValue {
pub struct PluginScope {
/// 客户端应用:`claude-code`、`codex`……(请求记录上的 `client_hint`)
pub clients: Vec<String>,
/// 客户端要的模型
/// 发给上游的模型:路由规则改了名的,按改名之后的
pub models: Vec<String>,
/// 服务回答的上游。**只管回答那一段**:改请求时还没选上游
/// 发往的上游。**请求和回答都按它**:请求钩子排在路由之后,每发往一个上游跑一次
pub upstreams: Vec<String>,
}

Expand Down Expand Up @@ -5220,6 +5223,9 @@ pub struct PluginRunView {
/// 当时的名字。**插件写的字**
pub plugin_name: String,
pub hook: PluginHook,
/// 跑在尝试链上的第几跳(从 0 起,对着 [`RoutingView::attempts`])。请求钩子每发往一个
/// 上游跑一次,故障转移换了上游就多一组;回答钩子跑在回答的那一跳上
pub attempt: u32,
pub outcome: PluginOutcome,
/// 出错、拒绝的原因
pub error: Option<Msg>,
Expand Down
4 changes: 2 additions & 2 deletions crates/tw-config/src/plugins.rs
Original file line number Diff line number Diff line change
Expand Up @@ -83,10 +83,10 @@ pub struct PluginScope {
/// 客户端应用:`claude-code`、`codex`……
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub clients: Vec<String>,
/// 客户端要的模型
/// 发给上游的模型:路由规则改了名的,按改名之后的
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub models: Vec<String>,
/// 服务回答的上游。只管回答那一段:改请求时还没选上游
/// 发往的上游。请求和回答都按它:请求钩子每发往一个上游跑一次
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub upstreams: Vec<String>,
}
Expand Down
8 changes: 4 additions & 4 deletions crates/tw-config/tests/manual/schema.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1605,17 +1605,17 @@ pub fn sections() -> Vec<Section> {
Kind::Strs,
Def::Is("[]"),
t(
"Models the client asks for, as model ids or globs (`claude-*`). `[]`: every model.",
"客户端请求的模型,写模型 ID 或通配(`claude-*`)。`[]`:所有模型。",
"Models sent to the upstream, as model ids or globs (`claude-*`). When a routing rule renames the model, the new name is the one that matches. `[]`: every model.",
"发给上游的模型,写模型 ID 或通配(`claude-*`)。路由规则改了模型名的,按改名之后的匹配。`[]`:所有模型。",
),
),
row(
"upstreams",
Kind::Strs,
Def::Is("[]"),
t(
"Upstreams whose answers the plugin handles, by name or glob. It applies to answers only: a request is changed before an upstream is chosen. `[]`: every upstream.",
"插件处理哪些上游的回答,写名字或通配。只作用于回答:请求在选定上游之前就已改写。`[]`:所有上游。",
"Upstreams the plugin handles, by name or glob, for requests and answers alike. `[]`: every upstream.",
"插件处理哪些上游,写名字或通配,请求和回答都按它。`[]`:所有上游。",
),
),
],
Expand Down
7 changes: 7 additions & 0 deletions crates/tw-control/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1056,6 +1056,13 @@ async fn request_detail(
.map_err(records)?
.into_iter()
.map(|r| tw_api::PluginRunView {
// 第几跳记在 `detail` 里(数据面每一次运行都写)
attempt: r
.detail
.as_deref()
.and_then(|d| serde_json::from_str::<serde_json::Value>(d).ok())
.and_then(|d| d.get("attempt").and_then(serde_json::Value::as_u64))
.unwrap_or(0) as u32,
plugin_id: r.plugin_id,
plugin_name: r.plugin_name,
hook: r.hook,
Expand Down
5 changes: 4 additions & 1 deletion crates/tw-control/src/plugins.rs
Original file line number Diff line number Diff line change
Expand Up @@ -889,7 +889,8 @@ fn refused(why: Msg) -> tw_api::PluginTrialResult {
/// 试跑本身在数据面那一侧(视图、写回、占位符都在 [`tw_gateway::plugin::trial`])。
///
/// 存下来的回答是上游的原话:回答它的那一家说什么格式,看服务它的那一跳转换过没有,
/// 和会话记录读回答是同一个办法
/// 和会话记录读回答是同一个办法。插件的 `ctx` 按这一行的路由给:回答它的那一家,和发给
/// 那一家的模型名
async fn run_trial(
s: &ControlState,
active: &Active,
Expand All @@ -916,6 +917,8 @@ async fn run_trial(
query: None,
body,
client: row.client_hint.as_deref(),
upstream: &row.provider,
sent_model: &row.sent_model,
}),
reply
.as_deref()
Expand Down
20 changes: 16 additions & 4 deletions crates/tw-control/tests/plugins.rs
Original file line number Diff line number Diff line change
Expand Up @@ -805,10 +805,17 @@ async fn a_request_shows_its_plugin_runs_and_the_body_after_them() {
let g = b.store.lock().await;
g.db().insert(&row(1, 1_000)).unwrap();
g.db().insert(&row(2, 2_000)).unwrap();
g.record_plugin_run(&run_row(1, 1_000, tw_api::PluginOutcome::Changed));
// 故障转移过一次:第 0 跳、第 1 跳各跑一次请求钩子,回答钩子跑在回答的第 1 跳上
let mut first = run_row(1, 1_000, tw_api::PluginOutcome::Changed);
first.detail = Some(r#"{"attempt":0,"changed":["system"]}"#.into());
g.record_plugin_run(&first);
let mut second = run_row(1, 1_100, tw_api::PluginOutcome::Changed);
second.detail = Some(r#"{"attempt":1,"changed":["system"]}"#.into());
g.record_plugin_run(&second);
let mut reply = run_row(1, 1_500, tw_api::PluginOutcome::Error);
reply.hook = tw_api::PluginHook::Reply;
reply.error = Some(tw_types::msg!("gw.plugin.failed" => "The plugin failed."));
reply.detail = Some(r#"{"attempt":1,"text_calls":1}"#.into());
g.record_plugin_run(&reply);
g.record_plugin_run(&run_row(2, 2_000, tw_api::PluginOutcome::Unchanged));
g.record_body(
Expand All @@ -824,12 +831,17 @@ async fn a_request_shows_its_plugin_runs_and_the_body_after_them() {
let (st, d) = call(&b.app, "GET", "/request/1", None).await;
assert_eq!(st, StatusCode::OK, "{d}");
let runs = d["plugins"].as_array().unwrap();
assert_eq!(runs.len(), 2);
assert_eq!(runs.len(), 3);
assert_eq!(runs[0]["hook"], "request");
assert_eq!(runs[0]["outcome"], "changed");
assert_eq!(runs[0]["cpu_us"], 120);
assert_eq!(runs[1]["hook"], "reply");
assert_eq!(runs[1]["error"]["code"], "gw.plugin.failed");
let attempts: Vec<u64> = runs
.iter()
.map(|r| r["attempt"].as_u64().unwrap())
.collect();
assert_eq!(attempts, [0, 1, 1]);
assert_eq!(runs[2]["hook"], "reply");
assert_eq!(runs[2]["error"]["code"], "gw.plugin.failed");
assert_eq!(d["row"]["plugin_changed"], true);
let after = d["request_after_plugins"]["text"].as_str().unwrap();
assert!(after.contains("today is Friday"), "{after}");
Expand Down
7 changes: 4 additions & 3 deletions crates/tw-gateway/src/bodies.rs
Original file line number Diff line number Diff line change
Expand Up @@ -56,9 +56,10 @@ pub enum BodyKind {
Request,
Response,
/// 插件改过之后的请求体(`Request` 存的是客户端发来的那一份)。**只有插件真的改了
/// 才存**,挨着 `Request` 放。交来的是要发出去的那一份(插件交回的占位符已经换回
/// 原值,见 [`crate::plugin::request`]),带着这个请求的 [`Redaction`]:落盘前和别的
/// 正文一样换掉、打码([`BodyRecord::for_disk`])
/// 才存**,挨着 `Request` 放。请求钩子每一跳跑一次,存的是最后发出去的那一跳收到的
/// 那一份 —— 回答的那一家收到的就是它(客户端那种格式、转换之前,插件交回的占位符
/// 已经换回原值,见 [`crate::plugin::request`]),带着那一跳的 [`Redaction`]:落盘前和
/// 别的正文一样换掉、打码([`BodyRecord::for_disk`])
AfterPlugins,
}

Expand Down
73 changes: 73 additions & 0 deletions crates/tw-gateway/src/guard.rs
Original file line number Diff line number Diff line change
Expand Up @@ -130,6 +130,21 @@ pub fn replace(
(bytes::Bytes::from(r.text), r.ledger)
}

/// `after` 里 `before` 没有的那些值:插件写进请求里的(见 [`crate::plugin::request`])。
///
/// 按规则和打过码的样子比:同一个值在两份里打出来的码一样。客户端原话里就有的值,开头
/// 那一遍已经报过了,插件改过的那一份里再出现不再报一次。
pub fn more_found(before: &[Finding], after: Vec<Finding>) -> Vec<Finding> {
after
.into_iter()
.filter(|f| {
!before
.iter()
.any(|b| b.rule == f.rule && b.masked == f.masked)
})
.collect()
}

/// 找到的东西写成事件里的样子。
pub fn items(found: &[Finding]) -> Vec<tw_api::SecretItem> {
found
Expand Down Expand Up @@ -194,6 +209,64 @@ pub fn screen(
refusal
}

/// 插件改过的请求再看一遍:**只看插件加进来的。**
///
/// 客户端的原话在开头已经看过([`screen`]),该报的报了、该拒的拒了;插件改过的那一份
/// 要是整个再报一遍,同一处藏匿字符、同一条命中会在安全日志里出现两次。所以两份都扫,
/// 原话里就有的那几处减掉,剩下的照 [`screen`] 的规矩报、下结论 —— 拦截档下原话里命中
/// 「拦」的请求走不到这一步,这一遍拒不拒只看插件加进来的。
pub fn screen_more(
bus: &tw_observe::EventBus,
id: u64,
provider: &str,
s: &Screen,
before: &tw_dialect::ir::Request,
after: &tw_dialect::ir::Request,
) -> Option<tw_types::Msg> {
let hidden = if s.hidden_mode.detects() {
let was = tw_guard::hidden::scan_request(before, &s.hidden);
let mut now = tw_guard::hidden::scan_request(after, &s.hidden);
// 一种藏法在一个地方合成一条:插件往同一处又藏了几个,那一条就变了,整条再报
now.retain(|n| {
!was.iter().any(|w| {
w.kind == n.kind
&& w.in_tool_result == n.in_tool_result
&& w.example == n.example
&& w.revealed == n.revealed
&& n.count <= w.count
})
});
now
} else {
Vec::new()
};
let mut refusal = hidden_found(bus, id, provider, s.hidden_mode, &hidden);
if s.content_mode.detects() && !s.content.is_empty() {
let mut was = s.content.scan_request(before);
let mut hits = s.content.scan_request(after);
// 一处一处地减:原话里有一处,插件那一版里同样的一处就不是新的。位置不比 ——
// 插件在前面加了字,后面的位置都挪了
hits.retain(|h| {
match was.iter().position(|w| {
w.rule == h.rule
&& w.custom == h.custom
&& w.action == h.action
&& w.snippet == h.snippet
&& w.in_tool_result == h.in_tool_result
}) {
Some(i) => {
was.swap_remove(i);
false
}
None => true,
}
});
let refused = content_matched(bus, id, provider, s.content_mode, &hits);
refusal = refusal.or(refused);
}
refusal
}

/// 没法按消息结构读的正文(解不开的 WebSocket 帧):**只查藏匿字符** —— 它在任何
/// 地方都没有正当用途;内容规则按整段原文查的话,系统提示里的话也会被当成调用方的。
pub fn screen_text(
Expand Down
4 changes: 3 additions & 1 deletion crates/tw-gateway/src/plugin/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,9 @@
//! 什么都没改时一个字节都不动。
//! - [`bridge`]:插件永远看不到真的密钥(不变式 I5)。进插件之前按出站脱敏的规则把
//! 认得出的密钥换成占位符,出来之后换回去;**不看脱敏开在哪一档**。
//! - [`request`]:请求钩子。一个客户端请求只跑一次(I8),排在内容审查和路由之前(I7)。
//! - [`request`]:请求钩子。排在路由之后,**每发往一个上游跑一次**(契约附录二的 I7、
//! I8):按这一次的客户端、发出去的模型和上游挑插件,从客户端的原话起改;换上游从
//! 原话重来,同一家重发不重跑。改过的请求再过一遍内容审查,然后才转换格式、脱敏。
//! - [`reply`]:回答钩子。排在格式转换之后、工具调用审查和输出长度之前(I7)——
//! 这两道防护看的就是插件改过的那一版。
//! - [`pool`]:插件调用都是阻塞的、吃 CPU 的,放在专用线程池上跑,不占 tokio 的线程。
Expand Down
21 changes: 16 additions & 5 deletions crates/tw-gateway/src/plugin/reply/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -61,11 +61,15 @@ pub struct Call {
pub struct ReplyCtx<'a> {
pub dialect: Dialect,
pub client: Option<&'a str>,
/// 客户端要的模型
/// 发给回答它的那一家的模型名:路由规则、请求钩子改过的是改过之后的
pub model: &'a str,
/// 客户端要的模型
pub requested_model: &'a str,
/// 回答它的那一家
pub upstream: &'a str,
pub request_id: u64,
/// 回答它的那一跳是尝试链上的第几跳。记在每一次运行的 `detail` 里
pub attempt: usize,
}

/// 一个插件在这次回答里的状态。
Expand Down Expand Up @@ -116,6 +120,7 @@ pub struct Chain {
bridge: Bridge,
dialect: Dialect,
request_id: u64,
attempt: usize,
recorded: bool,
}

Expand Down Expand Up @@ -165,7 +170,7 @@ impl Chain {
ctx: &ReplyCtx<'_>,
) -> Result<Option<Chain>, GatewayError> {
let mut stages = Vec::new();
// 跑不了的插件在请求钩子那一步已经按 `on_error` 处理过了:这里只有能跑的
// 跑不了的插件在回答它的那一次发出去之前已经按 `on_error` 处理过了:这里只有能跑的
for a in set.for_reply(ctx.client, ctx.model, ctx.upstream) {
let Some(host) = a.ready().cloned() else {
continue;
Expand All @@ -174,8 +179,9 @@ impl Chain {
let c = super::request::ctx(
ctx.client,
ctx.model,
ctx.requested_model,
ctx.dialect,
Some(ctx.upstream),
ctx.upstream,
&a.settings,
);
let made = state
Expand All @@ -195,7 +201,7 @@ impl Chain {
outcome: PluginOutcome::Error,
error: Some(why.clone()),
cpu_us: 0,
detail: None,
detail: Some(json!({ "attempt": ctx.attempt })),
};
state.plugin_ran(ctx.request_id, &a, run, Vec::new());
if a.on_error == OnError::Reject {
Expand All @@ -207,6 +213,7 @@ impl Chain {
bridge,
dialect: ctx.dialect,
request_id: ctx.request_id,
attempt: ctx.attempt,
recorded: false,
};
started.finish();
Expand All @@ -225,6 +232,7 @@ impl Chain {
bridge,
dialect: ctx.dialect,
request_id: ctx.request_id,
attempt: ctx.attempt,
recorded: false,
}))
}
Expand All @@ -244,8 +252,9 @@ impl Chain {
let c = super::request::ctx(
ctx.client,
ctx.model,
ctx.requested_model,
ctx.dialect,
Some(ctx.upstream),
ctx.upstream,
settings,
);
let instance = pool
Expand All @@ -261,6 +270,7 @@ impl Chain {
bridge: Bridge::new(Arc::new(tw_guard::redact::rules::RuleSet::none())),
dialect: ctx.dialect,
request_id: ctx.request_id,
attempt: ctx.attempt,
recorded: false,
}))
}
Expand Down Expand Up @@ -638,6 +648,7 @@ impl Chain {
error: s.error.clone(),
cpu_us: s.cpu.as_micros().min(u64::MAX as u128) as u64,
detail: Some(json!({
"attempt": self.attempt,
"text_calls": c.text_calls,
"text_changed": c.text_changed,
"tool_calls": c.tool_calls,
Expand Down
2 changes: 2 additions & 0 deletions crates/tw-gateway/src/plugin/reply/tests/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -51,8 +51,10 @@ async fn chain_with(
dialect,
client: None,
model: "m",
requested_model: "m",
upstream: "u",
request_id: 1,
attempt: 0,
},
)
.await
Expand Down
Loading
Loading