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 Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

50 changes: 47 additions & 3 deletions bin/twcore/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -783,9 +783,18 @@ fn cmd_serve(path: &Path, port: Option<u16>, safe: bool, parent: Option<u32>) ->
// body 的通道在这里建:**它是唯一同时看得见网关和存储的地方**,
// 而两边各有各的同形结构,是为了不让「观测」挂到「转发」下面。
let (body_tx, body_rx) = tw_gateway::bodies::channel();
let store = build_store(&dir, state.bus.clone(), state.pricing.clone(), body_rx);
// 插件在每个请求上的运行记录,和正文同一个道理:网关交出去,存储层落库
let (run_tx, run_rx) = tokio::sync::mpsc::channel(tw_gateway::plugin::RUN_CHANNEL_CAP);
let store = build_store(
&dir,
state.bus.clone(),
state.pricing.clone(),
body_rx,
run_rx,
);
if store.is_some() {
state.set_body_sink(body_tx);
state.set_plugin_sink(run_tx);
}

/*
Expand Down Expand Up @@ -843,6 +852,16 @@ fn cmd_serve(path: &Path, port: Option<u16>, safe: bool, parent: Option<u32>) ->
None
}
};
// 插件目录也盯着:**插件文件被改了,那个插件马上停用**,不等下一次改配置。
// 盯不住时退回到每次换配置时重算哈希,所以同样只说一句
let _plugin_watch = match tw_control::plugins::spawn_watcher(state.clone(), manager.path())
{
Ok(w) => Some(w),
Err(e) => {
tracing::warn!("the plugin directory cannot be watched, so a changed plugin file is noticed only at the next configuration change: {e}");
None
}
};

// 控制面无论如何都要起来 —— **网关挂了的时候,用户最需要的恰恰
// 是能改配置**。安全模式就是「只有这一半」。
Expand Down Expand Up @@ -946,6 +965,7 @@ fn build_store(
// **和网关同一份价格簿**,不是一份副本:改了价目表,下一个结束的请求就按新价算
pricing: tw_pricing::Shared,
bodies: tokio::sync::mpsc::Receiver<tw_gateway::bodies::BodyRecord>,
runs: tokio::sync::mpsc::Receiver<tw_gateway::plugin::RunRecord>,
) -> Option<std::sync::Arc<tokio::sync::Mutex<tw_store::Recorder>>> {
let events = bus.subscribe();
let (db, blobs) = match tw_store::open(dir) {
Expand Down Expand Up @@ -1005,6 +1025,7 @@ fn build_store(
which: match kind {
tw_gateway::bodies::BodyKind::Request => tw_store::Which::Request,
tw_gateway::bodies::BodyKind::Response => tw_store::Which::Response,
tw_gateway::bodies::BodyKind::AfterPlugins => tw_store::Which::AfterPlugins,
},
body,
original_len,
Expand All @@ -1016,13 +1037,36 @@ fn build_store(
drop(held);
}
});
Some(tw_store::task::spawn(
let recorder = tw_store::task::spawn(
// 算完价钱往回报一条 —— 见 `Event::RequestPriced`。这里是唯一
// 同时看得见总线和存储层的地方,所以接线在这儿完成。
tw_store::Recorder::new(db, blobs, pricing).reporting_to(bus),
events,
rx,
))
);
// 插件的运行记录同样在这里对接:网关那边的一次运行,换成存储层的一行
let (run_tx, run_rx) = tokio::sync::mpsc::channel(tw_gateway::plugin::RUN_CHANNEL_CAP);
let mut runs = runs;
tokio::spawn(async move {
while let Some(r) = runs.recv().await {
let row = tw_store::PluginRunRow {
request_id: r.request_id as i64,
at_ms: r.at_ms as i64,
plugin_id: r.run.plugin_id,
plugin_name: r.run.plugin_name,
hook: r.run.hook,
outcome: r.run.outcome,
error: r.run.error,
cpu_us: r.run.cpu_us.min(i64::MAX as u64) as i64,
detail: r.run.detail.map(|d| d.to_string()),
};
if run_tx.send(row).await.is_err() {
return;
}
}
});
tw_store::task::record_plugin_runs(recorder.clone(), run_rx);
Some(recorder)
}

/// 等一个「该退了」。
Expand Down
30 changes: 30 additions & 0 deletions crates/tw-api/msg-codes.txt
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,13 @@ config.failover_range
config.name_collision
config.no_clients
config.output_limit_range
config.plugin.bad_id
config.plugin.blank_pattern
config.plugin.duplicate
config.plugin.file
config.plugin.reserved_id
config.plugin.setting_type
config.plugin.sha256
config.rejected
config.rejected_at
config.remote_port_is_gateway
Expand Down Expand Up @@ -136,6 +143,18 @@ control.no_such_version
control.not_a_chatgpt_account
control.patch.no_entry
control.patch.not_an_entry
control.plugin.bad_id
control.plugin.blank_pattern
control.plugin.file_missing
control.plugin.file_moved_on
control.plugin.id_taken
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
control.pricing.not_a_dataset
control.pricing.save_failed
Expand Down Expand Up @@ -260,7 +279,18 @@ gw.oauth.status
gw.oauth.unreachable
gw.output_limit.cut
gw.output_limit.withheld
gw.plugin.api
gw.plugin.engine
gw.plugin.failed
gw.plugin.file_changed
gw.plugin.manifest
gw.plugin.not_located
gw.plugin.setting_type
gw.plugin.setting_unknown
gw.plugin.syntax
gw.plugin.syntax_at
gw.plugin.too_large
gw.plugin.unreadable
gw.probe.aws_token_expired
gw.probe.bedrock_list_denied
gw.probe.connect
Expand Down
29 changes: 29 additions & 0 deletions crates/tw-api/src/ep.rs
Original file line number Diff line number Diff line change
Expand Up @@ -129,6 +129,35 @@ endpoints! {
DeleteCustomRule: DELETE "/security/{guard}/custom/{name}" [guard, name], api::BaseVersion => api::ConfigWritten;
TestSecurity: POST "/security/{guard}/test" [guard], api::SecurityTestRequest => api::SecurityTestResult;

// ─────────────────────────────────────────────── 脚本插件
//
// **装、换源码、批准三个端点不给网页调**(桌面端的 `call` 白名单里没有它们):
// 这三件事要在系统的确认框里点头,那一步在桌面端的 Rust 里 —— 它自己再编一遍
// 源码,把名字、权限和哈希摆给人看,点了头才发请求。网页里的脚本做不到这件事,
// 就做不成这三件事。
/// 全部插件,按运行的顺序:状态、计数
Plugins: GET "/plugins", () => Vec<api::PluginView>;
/// 编一份源码看看它是什么插件。**什么都不留下**
PluginInspect: POST "/plugins/inspect", api::PluginSource => api::PluginInspection;
/// 装一个:写插件文件和它的底稿,配置里加一条。**网页不能调**
CreatePlugin: POST "/plugins", api::PluginCreate => api::ConfigWritten;
/// 排顺序,也就是运行的顺序
ReorderPlugins: PUT "/plugins/order", api::PluginOrder => api::ConfigWritten;
/// 开关、出错时怎么办、范围、设置
UpdatePlugin: PUT "/plugins/{id}" [id], api::PluginUpdate => api::ConfigWritten;
/// 删掉:配置里那一条、插件文件和底稿
DeletePlugin: DELETE "/plugins/{id}" [id], api::BaseVersion => api::ConfigWritten;
/// 换一份源码,批准的就是新的这一份。**网页不能调**
ReplacePluginSource: PUT "/plugins/{id}/source" [id], api::PluginSourceReplace => api::ConfigWritten;
/// 批准过的那一份和磁盘上现在那一份
PluginSourceDiff: GET "/plugins/{id}/source" [id], () => api::PluginSourceView;
/// 批准磁盘上改过的那个文件。**网页不能调**
ApprovePluginFile: POST "/plugins/{id}/approve" [id], api::PluginApprove => api::ConfigWritten;
/// 拿一条记下的请求试跑。**不连上游**
TrialPlugin: POST "/plugins/{id}/trial" [id], api::PluginTrial => api::PluginTrialResult;
/// 最近的日志,老的在前
PluginLogs: GET "/plugins/{id}/logs" [id], () => Vec<api::PluginLogEntry>;

// ─────────────────────────────────────────────── 账号登录
StartChatgptLogin: POST "/chatgpt/login", api::ChatgptLoginStart => api::ChatgptLogin;
ChatgptLoginStatus: GET "/chatgpt/login/{id}" [id], () => api::ChatgptLoginStatus;
Expand Down
Loading
Loading