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
17 changes: 16 additions & 1 deletion crates/tw-api/msg-codes.txt
Original file line number Diff line number Diff line change
Expand Up @@ -152,7 +152,6 @@ control.plugin.not_found
control.plugin.order
control.plugin.reserved_id
control.plugin.trial_changed
control.plugin.trial_unavailable
control.plugin.unreadable
control.plugin.write_failed
control.pricing.broke_off
Expand Down Expand Up @@ -279,17 +278,33 @@ gw.oauth.status
gw.oauth.unreachable
gw.output_limit.cut
gw.output_limit.withheld
gw.plugin.answer_unreadable
gw.plugin.api
gw.plugin.bad_output
gw.plugin.changed
gw.plugin.cpu_limit
gw.plugin.engine
gw.plugin.failed
gw.plugin.file_changed
gw.plugin.manifest
gw.plugin.memory_limit
gw.plugin.not_located
gw.plugin.nothing_to_try
gw.plugin.output_limit
gw.plugin.permission_violation
gw.plugin.reason passthrough
gw.plugin.rejected
gw.plugin.reply_failed
gw.plugin.request_failed
gw.plugin.request_unreadable
gw.plugin.setting_type
gw.plugin.setting_unknown
gw.plugin.syntax
gw.plugin.syntax_at
gw.plugin.threw
gw.plugin.too_large
gw.plugin.trap
gw.plugin.unavailable
gw.plugin.unreadable
gw.probe.aws_token_expired
gw.probe.bedrock_list_denied
Expand Down
101 changes: 80 additions & 21 deletions crates/tw-control/src/plugins.rs
Original file line number Diff line number Diff line change
Expand Up @@ -859,17 +859,22 @@ async fn trial(
(row, request, reply)
};
// 跑不了的插件不试:改过的代码不跑(I9),加载不了的也跑不了
if let Some(b) = active.broken() {
return Ok(Json(refused(match b {
Broken::Changed => msg!(
"control.plugin.trial_changed", plugin = &active.name =>
"The file of plugin `{plugin}` changed and has not been approved, so it cannot be \
tried."
),
Broken::Error(m) => m.clone(),
})));
}
Ok(Json(run_trial(&active, &row, request, reply)))
let host = match &active.state {
tw_gateway::plugin::State::Ready(h) => h.clone(),
tw_gateway::plugin::State::Broken(b) => {
return Ok(Json(refused(match b {
Broken::Changed => msg!(
"control.plugin.trial_changed", plugin = &active.name =>
"The file of plugin `{plugin}` changed and has not been approved, so it cannot \
be tried."
),
Broken::Error(m) => m.clone(),
})));
}
};
Ok(Json(
run_trial(&s, &active, host, &row, request, reply).await,
))
}

fn refused(why: Msg) -> tw_api::PluginTrialResult {
Expand All @@ -881,17 +886,71 @@ fn refused(why: Msg) -> tw_api::PluginTrialResult {
}
}

/// 试跑本身在数据面那一侧(视图、写回都在那里)。**还没接上**:在那之前说一句做不了
fn run_trial(
_active: &Active,
_row: &tw_store::RequestRow,
_request: Option<Vec<u8>>,
_reply: Option<Vec<u8>>,
/// 试跑本身在数据面那一侧(视图、写回、占位符都在 [`tw_gateway::plugin::trial`])。
///
/// 存下来的回答是上游的原话:回答它的那一家说什么格式,看服务它的那一跳转换过没有,
/// 和会话记录读回答是同一个办法
async fn run_trial(
s: &ControlState,
active: &Active,
host: std::sync::Arc<dyn tw_gateway::plugin::PluginHost>,
row: &tw_store::RequestRow,
request: Option<Vec<u8>>,
reply: Option<Vec<u8>>,
) -> tw_api::PluginTrialResult {
refused(msg!(
"control.plugin.trial_unavailable" =>
"Trial runs are not available in this build yet."
))
use tw_gateway::plugin::trial::{self, StoredReply, StoredRequest};
use tw_store::search::text::{client_dialect, dialect_of};
let upstream = row
.translated
.as_deref()
.and_then(|j| serde_json::from_str::<tw_api::TranslatedView>(j).ok())
.map(|t| dialect_of(t.to))
.or_else(|| client_dialect(&row.path));
let t = trial::run(
s.gateway.plugin_pool.clone(),
host,
&active.settings,
s.gateway.runtime().redact.clone(),
request.as_deref().map(|body| StoredRequest {
path: &row.path,
query: None,
body,
client: row.client_hint.as_deref(),
}),
reply
.as_deref()
.zip(upstream)
.map(|(body, upstream)| StoredReply {
body,
upstream,
provider: &row.provider,
}),
)
.await;
let at_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |d| d.as_millis() as u64);
let side = |x: trial::Side| tw_api::TrialSide {
before: x.before,
after: x.after,
outcome: x.outcome,
};
tw_api::PluginTrialResult {
request: t.request.map(side),
reply: t.reply.map(side),
logs: t
.logs
.into_iter()
.map(|(hook, l)| tw_api::PluginLogEntry {
at_ms,
request_id: Some(row.id as u64),
hook,
level: l.level,
text: l.text,
})
.collect(),
error: t.error,
}
}

// ---------------------------------------------------------------- 监听
Expand Down
134 changes: 133 additions & 1 deletion crates/tw-control/tests/plugins.rs
Original file line number Diff line number Diff line change
Expand Up @@ -897,7 +897,8 @@ async fn a_trial_needs_a_known_plugin_and_a_recorded_request() {
)
.await;
assert_eq!(st, StatusCode::OK, "{v}");
assert!(v["error"]["code"].is_string(), "{v}");
// 这一条什么正文都没存下来
assert_eq!(v["error"]["code"], "gw.plugin.nothing_to_try", "{v}");
assert!(v["logs"].as_array().unwrap().is_empty());

// 改过还没批准的代码不试
Expand All @@ -913,6 +914,137 @@ async fn a_trial_needs_a_known_plugin_and_a_recorded_request() {
assert_eq!(v["error"]["code"], "control.plugin.trial_changed");
}

/// 跑得起钩子的引擎:manifest 照假引擎读,请求钩子在系统提示后面补一句,回答钩子把字
/// 换成大写
struct Running;

struct RunningHost(Arc<dyn tw_gateway::plugin::PluginHost>);

impl tw_gateway::plugin::PluginHost for RunningHost {
fn manifest(&self) -> &tw_gateway::plugin::Manifest {
self.0.manifest()
}
fn sha256(&self) -> [u8; 32] {
self.0.sha256()
}
fn on_request(
&self,
mut view: Value,
_ctx: Value,
) -> tw_gateway::plugin::Invocation<tw_gateway::plugin::RequestOutcome> {
let system = view["system"].as_str().unwrap_or_default().to_string();
view["system"] = json!(format!("{system} Today is Friday."));
let mut inv =
tw_gateway::plugin::Invocation::ok(tw_gateway::plugin::RequestOutcome::Changed(view));
inv.logs.push(tw_gateway::plugin::LogLine {
level: tw_api::PluginLogLevel::Info,
text: "added the date".into(),
});
inv
}
fn reply(
&self,
_ctx: Value,
) -> Result<Box<dyn tw_gateway::plugin::ReplyHost>, tw_gateway::plugin::RunError> {
use tw_gateway::plugin::{Invocation, ToolCallOutcome};
Ok(Box::new(tw_gateway::plugin::host::double::Closures {
text: Box::new(|t| Invocation::ok(Some(t.to_uppercase()))),
end: Box::new(|| Invocation::ok(None)),
tool: Box::new(|_| Invocation::ok(ToolCallOutcome::Unchanged)),
}))
}
}

impl tw_gateway::plugin::Engine for Running {
fn load(
&self,
source: &[u8],
) -> Result<Arc<dyn tw_gateway::plugin::PluginHost>, tw_gateway::plugin::LoadError> {
Ok(Arc::new(RunningHost(tw_gateway::plugin::Engine::load(
&FakeEngine,
source,
)?)))
}
}

/// 试跑接到数据面上:记下的请求和回答各跑一遍,前后两份都打着码,日志交回来、不进
/// 插件自己的日志
#[tokio::test]
async fn a_trial_runs_the_plugin_on_the_recorded_request_and_answer() {
let b = bed();
b.gw.set_plugin_engine(Arc::new(Running));
let src = source(
json!({"name": "Both", "api": 1, "permissions": ["system", "reply.text"]}),
&["onRequest", "onReplyText"],
);
let id = b.install(&src, json!({})).await;
let key = "sk-ant-api03-TRIALKEYAAAAAAAAAAAAAAAAAAAA";
let request = json!({
"model": "claude-sonnet-4-5", "max_tokens": 64, "stream": true,
"system": "Be brief.",
"messages": [{"role": "user", "content": format!("my key is {key}")}]
})
.to_string();
let answer = [
json!({"type":"message_start","message":{"id":"m","type":"message","role":"assistant","model":"m","content":[],"usage":{"input_tokens":1,"output_tokens":1}}}),
json!({"type":"content_block_start","index":0,"content_block":{"type":"text","text":""}}),
json!({"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"hello there"}}),
json!({"type":"content_block_stop","index":0}),
json!({"type":"message_delta","delta":{"stop_reason":"end_turn"},"usage":{"output_tokens":2}}),
json!({"type":"message_stop"}),
]
.iter()
.map(|c| format!("event: {}\ndata: {c}\n\n", c["type"].as_str().unwrap()))
.collect::<String>();
{
let g = b.store.lock().await;
g.db().insert(&row(9, 1_000)).unwrap();
g.record_body(
1_000,
9,
tw_store::Which::Request,
request.as_bytes(),
request.len(),
);
g.record_body(
1_000,
9,
tw_store::Which::Response,
answer.as_bytes(),
answer.len(),
);
}
let (st, v) = call(
&b.app,
"POST",
&format!("/plugins/{id}/trial"),
Some(json!({"request_id": 9})),
)
.await;
assert_eq!(st, StatusCode::OK, "{v}");
assert!(v["error"].is_null(), "{v}");
assert_eq!(v["request"]["outcome"], "changed", "{v}");
let after = v["request"]["after"].as_str().unwrap();
assert!(after.contains("Be brief. Today is Friday."), "{after}");
assert_eq!(v["reply"]["outcome"], "changed", "{v}");
let reply: Value = serde_json::from_str(v["reply"]["after"].as_str().unwrap()).unwrap();
assert_eq!(reply["content"][0]["text"], "HELLO THERE");
assert!(
!v.to_string().contains("TRIALKEY"),
"a secret was shown: {v}"
);
let logs = v["logs"].as_array().unwrap();
assert_eq!(logs.len(), 1, "{v}");
assert_eq!(logs[0]["hook"], "request");
assert_eq!(logs[0]["request_id"], 9);
// 试跑不进插件自己的日志和计数
let (_, mine) = call(&b.app, "GET", &format!("/plugins/{id}/logs"), None).await;
assert!(
mine.as_array().is_none_or(|l| l.is_empty()),
"the trial was logged: {mine}"
);
}

/// 远程 core:配置不在默认的地方,插件文件就在那份配置旁边 —— 文件由 core 自己写
#[tokio::test]
async fn plugin_files_live_next_to_the_configuration_wherever_it_is() {
Expand Down
5 changes: 3 additions & 2 deletions crates/tw-gateway/src/bodies.rs
Original file line number Diff line number Diff line change
Expand Up @@ -56,8 +56,9 @@ pub enum BodyKind {
Request,
Response,
/// 插件改过之后的请求体(`Request` 存的是客户端发来的那一份)。**只有插件真的改了
/// 才存**,挨着 `Request` 放。交来的是插件交回的那一份(占位符还没换回密钥),带着这个
/// 请求的 [`Redaction`]:落盘前和别的正文一样换掉、打码([`BodyRecord::for_disk`])
/// 才存**,挨着 `Request` 放。交来的是要发出去的那一份(插件交回的占位符已经换回
/// 原值,见 [`crate::plugin::request`]),带着这个请求的 [`Redaction`]:落盘前和别的
/// 正文一样换掉、打码([`BodyRecord::for_disk`])
AfterPlugins,
}

Expand Down
9 changes: 8 additions & 1 deletion crates/tw-gateway/src/guard.rs
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,13 @@ pub fn ledger_for(body: &[u8]) -> Ledger {
/// 就写着的占位符。之后每一跳都接着这本账换([`replace`]),存下来的那份请求也照它换
/// ([`crate::bodies::Redaction`])。不在拦截档时账本是空的。
pub fn look(mode: Mode, rules: &RuleSet, body: &[u8]) -> (Vec<Finding>, Ledger) {
look_from(mode, rules, body, Ledger::new(Scheme::SECRET))
}

/// [`look`],**接着 `seed` 的账编号**:插件跑过的请求,插件看到的占位符是按客户端原文
/// 编的(见 [`crate::plugin::bridge`]),改过之后的这一份接着那本账编,同一个值还是同一个
/// 号;插件改出来的新值接着往后编。
pub fn look_from(mode: Mode, rules: &RuleSet, body: &[u8], seed: Ledger) -> (Vec<Finding>, Ledger) {
let empty = || Ledger::new(Scheme::SECRET);
if !mode.detects() || rules.is_empty() {
return (Vec::new(), empty());
Expand All @@ -90,7 +97,7 @@ pub fn look(mode: Mode, rules: &RuleSet, body: &[u8]) -> (Vec<Finding>, Ledger)
if !mode.acts() {
return (found, empty());
}
let seed = empty().avoiding(text);
let seed = seed.avoiding(text);
let ledger = if hits.is_empty() {
seed
} else {
Expand Down
Loading
Loading