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.

38 changes: 29 additions & 9 deletions bin/twcore/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -782,7 +782,7 @@ fn cmd_serve(path: &Path, port: Option<u16>, safe: bool, parent: Option<u32>) ->
//
// body 的通道在这里建:**它是唯一同时看得见网关和存储的地方**,
// 而两边各有各的同形结构,是为了不让「观测」挂到「转发」下面。
let (body_tx, body_rx) = tokio::sync::mpsc::channel(tw_gateway::bodies::CHANNEL_CAP);
let (body_tx, body_rx) = tw_gateway::bodies::channel();
let store = build_store(&dir, state.bus.clone(), state.pricing.clone(), body_rx);
if store.is_some() {
state.set_body_sink(body_tx);
Expand Down Expand Up @@ -975,25 +975,45 @@ fn build_store(
"could not read the last request id, so this run may overwrite the oldest records: {e}"
),
}
// 两边的 body 结构在这里对接。**一次移动,不复制** —— `Bytes` 的
// 克隆是引用计数。
let (tx, rx) = tokio::sync::mpsc::channel(tw_gateway::bodies::CHANNEL_CAP);
/*
两边的 body 结构在这里对接,**落盘的不是原文**:脱敏规则认得出的值在这里换掉、
打码(`BodyRecord::for_disk`,见 `tw_gateway::bodies`)。放在阻塞线程上 ——
一份 4 MB 的正文要扫好几遍,占着异步线程的话,同一个线程上的转发都得等它。

交给存储层的这一头**只留一个空位**。等着落盘的正文由网关那一头按字节记账
(`bodies::QUEUED_MAX`),一份正文的额度要等存储层收下它才还回去;这里再开一个
大口子的话,攒在这里的那些就没人管了。
*/
let (tx, rx) = tokio::sync::mpsc::channel(1);
let mut bodies = bodies;
tokio::spawn(async move {
while let Some(b) = bodies.recv().await {
let Ok(disk) = tokio::task::spawn_blocking(move || b.for_disk()).await else {
continue;
};
let tw_gateway::bodies::ForDisk {
id,
at_ms,
kind,
body,
original_len,
held,
} = disk;
let mapped = tw_store::StoredBody {
id: b.id,
at_ms: b.at_ms,
which: match b.kind {
id,
at_ms,
which: match kind {
tw_gateway::bodies::BodyKind::Request => tw_store::Which::Request,
tw_gateway::bodies::BodyKind::Response => tw_store::Which::Response,
},
body: b.body,
original_len: b.original_len,
body,
original_len,
};
if tx.send(mapped).await.is_err() {
return;
}
// 存储层收下了:额度还回去
drop(held);
}
});
Some(tw_store::task::spawn(
Expand Down
2 changes: 2 additions & 0 deletions crates/tw-api/src/ep.rs
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,8 @@ endpoints! {
Fixture: GET "/request/{id}/fixture" [id], () => String, text;
Sessions: GET "/sessions", api::ListQuery => Vec<api::SessionView>;
SessionDetail: GET "/sessions/{id}" [id], () => api::SessionDetail;
/// 一次会话读成一段对话:每一轮新说的话、回答、工具调用和结果(已脱敏)
SessionTranscript: GET "/sessions/{id}/transcript" [id], () => api::Transcript;

// ─────────────────────────────────────────────── 测速、回放、试路由
SpeedQuote: POST "/speed/quote", api::SpeedRunRequest => api::SpeedQuote;
Expand Down
142 changes: 136 additions & 6 deletions crates/tw-api/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -672,9 +672,13 @@ pub const MSG_CODES: &str = include_str!("../msg-codes.txt");
/// `cn-resident-id` 和 `bank-card`([`CardNetwork`]、[`CardPrefix`]),[`SecretKind`] 多了
/// `personal`。照 30 写的界面说不出这两条规则按什么认。
///
/// **32 起有脚本插件**:事件多了 [`Event::PluginFailed`](插件在请求上出错,或者文件变了、
/// 加载不了而停用)。照 31 写的界面不认这个事件。
pub const CONTROL_API_VERSION: u32 = 32;
/// **32 起会话能读成一段对话**:新端点 `GET /sessions/{id}/transcript`([`Transcript`])
/// 从存下来的正文里读出每一轮新说的话、回答、推理、工具调用和结果,读不到的地方逐轮说出来
/// ([`TranscriptGap`])。照 31 写的界面只有每一轮的用量和金额。
///
/// **33 起有脚本插件**:事件多了 [`Event::PluginFailed`](插件在请求上出错,或者文件变了、
/// 加载不了而停用)。照 32 写的界面不认这个事件。
pub const CONTROL_API_VERSION: u32 = 33;

#[derive(Debug, Clone, Serialize, Deserialize)]
#[cfg_attr(feature = "ts", derive(ts_rs::TS))]
Expand Down Expand Up @@ -1723,7 +1727,7 @@ pub struct RetentionView {
/// 正文总共最多占多少字节
pub body_max_bytes: u64,
/// 正文现在实际占了多少。**不是配置,是现状** —— 没有它,
/// 「2 GB 上限」是个用户无从判断松紧的数字
/// 「5 GB 上限」是个用户无从判断松紧的数字
pub body_bytes_now: u64,
}

Expand Down Expand Up @@ -3693,14 +3697,23 @@ pub struct RequestDetail {
pub in_flight: bool,
}

/// 一份正文最多存多少字节:请求和回答一样,4 MiB。更长的只存开头,[`BodyView::truncated`]
/// 说出来。
///
/// **存储层按它截,网关攒回答也按它攒**(`tw_store::blobs::MAX_ONE`、
/// `tw_gateway::bodies::RESPONSE_TAP_MAX`)。两个数放在一处:各写各的话,改了一个,
/// 另一个还停在原地 —— 以前回答只攒 256 KB,比存储层肯收的少十几倍。
pub const BODY_MAX: usize = 4 * 1024 * 1024;

/// 一份存下来的 body。
#[derive(Debug, Clone, Serialize, Deserialize)]
#[cfg_attr(feature = "ts", derive(ts_rs::TS))]
pub struct BodyView {
/// **已脱敏**。这段文字会被复制到 issue 里
/// **已脱敏**。这段文字会被复制到 issue 里。落盘的那一份就是换过、打过码的(脱敏规则
/// 认得出的值不会原样写进磁盘),读出来再打一遍
pub text: String,
/// 原本多长。**截断了要能说出来** —— 不说的话用户会以为请求本身
/// 就长这样
/// 就长这样。没截断的就是存下来的这一份的长度:换掉、打码的那几处和原文差几个字节
pub original_len: usize,
pub truncated: bool,
}
Expand Down Expand Up @@ -4047,6 +4060,123 @@ pub struct SessionDetail {
pub turns: Vec<TurnView>,
}

/// 一次会话读成一段对话(`GET /sessions/{id}/transcript`):每一轮新说的话、回答、推理、
/// 工具调用和工具结果。
///
/// **从存下来的正文里读出来**,不是另记的一份:正文只留几天(`retention.body_days`),
/// 太大的只留开头,没存下来的也有。读不到的地方,那一轮的 `gaps` 说出来。
///
/// **已脱敏**,和请求详情里的正文同一套打码。图片只说类型和大小,从不带数据。
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[cfg_attr(feature = "ts", derive(ts_rs::TS))]
pub struct Transcript {
pub session: String,
/// 第一个读得懂的请求里的系统提示:Anthropic 的 `system`、Responses 的 `instructions`、
/// Gemini 的 `systemInstruction`,Chat 和 Responses 还有开头连着的 system、developer
/// 消息,几段之间空一行。没有是 null
pub system: Option<String>,
/// 和 [`SessionDetail::turns`] 同样的请求,同样的顺序
pub turns: Vec<TranscriptTurn>,
}

/// 对话里的一轮,就是会话里的一个请求。
///
/// 客户端每一轮都把整段历史发上来:请求 i 的消息 = 请求 i-1 的消息 + 上一轮的回答 + 新的
/// 用户消息或工具结果。`input` 只放新的那几条;上一轮的回答已经在上一轮的 `output` 里。
///
/// **不生成回答的调用**(数 token、Responses 的压缩)也在这里占一轮,`input`、`output`
/// 都是空的,也不和前后的请求比对:它们问的是这段对话,不是对话里的一句。
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[cfg_attr(feature = "ts", derive(ts_rs::TS))]
pub struct TranscriptTurn {
/// 请求号,写成十进制的字符串。和 [`TurnView::id`] 是同一条请求
pub id: String,
/// 这个请求带的历史没有接着上一个读得懂的请求:压缩过、改过历史,或者它是一串读不懂
/// 的请求之后第一个读得懂的。这时 `input` 是它的整段历史
pub restart: bool,
/// 系统提示和上一个读得懂的请求不一样了:新的那一份(去掉了的是空串)。没变是 null
pub system_changed: Option<String>,
/// 这个请求里新的消息。上一轮的回答没有完整读出来时(那一轮的 `gaps` 里有 `response_*`),
/// 客户端记下的那条助手消息也在这里:它是那一轮说过什么的记录
pub input: Vec<TranscriptMessage>,
/// 回答,从存下来的响应里读出来的。失败的请求(上游回了错误)没有回答,也不算缺
pub output: Vec<TranscriptPart>,
/// 这一轮哪些地方读不出来
pub gaps: Vec<TranscriptGap>,
}

/// 请求里的一条消息。
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[cfg_attr(feature = "ts", derive(ts_rs::TS))]
pub struct TranscriptMessage {
pub role: TranscriptRole,
pub parts: Vec<TranscriptPart>,
}

slug_enum! {
/// 一条消息是谁说的。
pub enum TranscriptRole {
User = "user",
Assistant = "assistant",
/// 只装着工具结果的消息:Anthropic 全是 `tool_result` 的用户消息、Chat 的 `tool`
/// 消息、Responses 的 `function_call_output`
Tool = "tool",
/// 对话中途的 system、developer 消息
System = "system",
}
}

/// 消息或回答里的一块。
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[cfg_attr(feature = "ts", derive(ts_rs::TS))]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum TranscriptPart {
Text {
text: String,
},
/// 推理。只有签名、或者被打码的推理,`text` 是空串
Thinking {
text: String,
},
/// `input` 是参数的 JSON 文本;自由格式的工具(Codex 的 `apply_patch`)是它的原文
ToolCall {
id: String,
name: String,
input: String,
},
/// `call_id` 和它回应的那个 `tool_call` 的 `id` 是同一个
ToolResult {
call_id: String,
text: String,
is_error: bool,
},
/// 图片。**只有类型和大小,从不带数据**;给的是地址的不知道大小
Image {
media_type: Option<String>,
bytes: Option<u64>,
},
/// 别的块:文件、音频、服务端工具的调用和结果……`label` 是它的类型名,原样
Other {
label: String,
},
}

slug_enum! {
/// 一轮里读不出来的地方。
pub enum TranscriptGap {
/// 请求体没有存下来,或者已经清掉了
RequestMissing = "request_missing",
/// 请求体只存了开头,或者解析不了
RequestTruncated = "request_truncated",
/// 回答没有存下来,或者已经清掉了
ResponseMissing = "response_missing",
/// 回答只存了开头:`output` 是读得出来的那一段
ResponseTruncated = "response_truncated",
/// 回答存下来了,但读不懂
ResponseUnreadable = "response_unreadable",
}
}

// ---------------------------------------------------------------- 请求重放

/// 把存下来的那条请求,原样发给另一个上游。
Expand Down
52 changes: 52 additions & 0 deletions crates/tw-api/src/ts.rs
Original file line number Diff line number Diff line change
Expand Up @@ -282,6 +282,58 @@ mod tests {
assert!(event.contains("answered_model?: string"), "{event}");
}

/// 对话记录:请求号是字符串;可以为空的是 null(必有的字段,不是省掉);块是按 `kind`
/// 分的联合,角色和缺口是字面量
#[test]
fn a_transcript_is_turns_of_messages_and_parts() {
let ts = typescript();
assert!(ts.contains(" SessionTranscript: { req: null; res: Transcript };"));
assert!(ts.contains(
" SessionTranscript: { method: \"GET\", path: \"/sessions/{id}/transcript\", params: [\"id\"], format: \"json\" },"
));
let transcript = decl_of(&ts, "Transcript");
for field in [
"session: string",
"system: string | null",
"turns: Array<TranscriptTurn>",
] {
assert!(transcript.contains(field), "{field}: {transcript}");
}
let turn = decl_of(&ts, "TranscriptTurn");
for field in [
"id: string",
"restart: boolean",
"system_changed: string | null",
"input: Array<TranscriptMessage>",
"output: Array<TranscriptPart>",
"gaps: Array<TranscriptGap>",
] {
assert!(turn.contains(field), "{field}: {turn}");
}
assert_eq!(
decl_of(&ts, "TranscriptMessage"),
"export type TranscriptMessage = { role: TranscriptRole, parts: Array<TranscriptPart>, }"
);
assert_eq!(
decl_of(&ts, "TranscriptPart"),
"export type TranscriptPart = { \"kind\": \"text\", text: string, } \
| { \"kind\": \"thinking\", text: string, } \
| { \"kind\": \"tool_call\", id: string, name: string, input: string, } \
| { \"kind\": \"tool_result\", call_id: string, text: string, is_error: boolean, } \
| { \"kind\": \"image\", media_type: string | null, bytes: number | null, } \
| { \"kind\": \"other\", label: string, }"
);
assert_eq!(
decl_of(&ts, "TranscriptRole"),
"export type TranscriptRole = \"user\" | \"assistant\" | \"tool\" | \"system\""
);
assert_eq!(
decl_of(&ts, "TranscriptGap"),
"export type TranscriptGap = \"request_missing\" | \"request_truncated\" \
| \"response_missing\" | \"response_truncated\" | \"response_unreadable\""
);
}

#[test]
fn every_endpoint_is_in_the_table() {
let ts = typescript();
Expand Down
2 changes: 1 addition & 1 deletion crates/tw-config/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1183,7 +1183,7 @@ mod tests {
assert_eq!(back.retention.body_days, 3);
// 没写的那两个仍然是默认值,不是 0 —— 0 会让 gc 把一切都删掉
assert_eq!(back.retention.row_days, 90);
assert_eq!(back.retention.body_max_bytes, 2 * 1024 * 1024 * 1024);
assert_eq!(back.retention.body_max_bytes, 5 * 1024 * 1024 * 1024);
assert!(
serde_yaml_ng::to_string(&back)
.unwrap()
Expand Down
4 changes: 2 additions & 2 deletions crates/tw-config/src/retention.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ pub struct Retention {
/// 一行记录留几天。它撑着「上个月花了多少」那类问题
#[serde(default = "d_row_days")]
pub row_days: u64,
/// 正文总共最多占多少字节。超了从最旧的整天开始删
/// 正文总共最多占多少字节。超了从最旧的整天开始删。出厂 5 GiB
#[serde(default = "d_body_max_bytes")]
pub body_max_bytes: u64,
}
Expand All @@ -33,7 +33,7 @@ fn d_row_days() -> u64 {
90
}
fn d_body_max_bytes() -> u64 {
2 * 1024 * 1024 * 1024
5 * 1024 * 1024 * 1024
}

impl Default for Retention {
Expand Down
6 changes: 3 additions & 3 deletions crates/tw-config/tests/manual/schema.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1165,10 +1165,10 @@ pub fn sections() -> Vec<Section> {
row(
"body_max_bytes",
Kind::Int,
Def::Is("2147483648"),
Def::Is("5368709120"),
t(
"Upper bound on the bytes bodies may take; beyond it the oldest days go first. The default is 2 GiB.",
"正文最多占用的字节数,超出时从最早的日期开始删除。默认 2 GiB。",
"Upper bound on the bytes bodies may take; beyond it the oldest days go first. The default is 5 GiB.",
"正文最多占用的字节数,超出时从最早的日期开始删除。默认 5 GiB。",
),
),
],
Expand Down
Loading
Loading