Skip to content
Merged
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
184 changes: 49 additions & 135 deletions src/adapter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,9 +10,9 @@ use tracing::{error, warn};

use crate::acp::{classify_notification, AcpEvent, ContentBlock, SessionPool};
use crate::config::{AttachmentsConfig, LiveStatusConfig, ReactionsConfig, ToolDisplay};
use crate::live_status::StatusEvent as LiveStatusEvent;
use crate::error_display::{format_coded_error, format_user_error};
use crate::format;
use crate::live_status::StatusEvent as LiveStatusEvent;
use crate::markdown::{self, TableMode};
use crate::reactions::StatusReactionController;

Expand Down Expand Up @@ -1510,17 +1510,12 @@ impl AdapterRouter {
} else {
markdown::convert_tables(&final_content, table_mode)
};
// 最終送信も同じ上限で抑える。ここを外すと、streaming で
// 打ち切っても最後に全量が流れ直す。
let chunks = if final_content.is_empty() {
Vec::new()
} else if let Some(collapsed) = collapse_degenerate(&final_content) {
format::split_message(&collapsed, message_limit)
} else {
cap_stream_chunks(
format::split_message(&final_content, message_limit),
MAX_STREAM_MESSAGES,
)
format::split_message(&final_content, message_limit)
};
if let Some(post) = placeholder_post {
if let Some(ref reply_id) = directives.reply_to {
Expand Down Expand Up @@ -1740,37 +1735,6 @@ async fn send_outbound_attachments(
Ok(())
}

/// How many Discord messages one turn may occupy.
///
/// A model that degenerates into a repetition loop produces tens of kilobytes
/// of identical text. Before this cap, every 2000 characters became one more
/// Discord message: on 2026-08-01 a single turn posted 31 messages of
/// `court:\n\n` repeated, 1998 characters each, in 30 seconds. Nothing in the
/// send path bounded the count, so the channel filled until the model stopped.
///
/// The cap is on the *number of messages*, not the content. A long, legitimate
/// answer still gets through — up to this many parts — and anything past it is
/// replaced by one notice saying how much was dropped. Losing the tail of a
/// runaway is the point; losing the tail of a real answer is visible and
/// recoverable, whereas a flooded channel is neither.
pub const MAX_STREAM_MESSAGES: usize = 5;

/// Keep the head, replace the rest with a notice. Returns the input unchanged
/// when it already fits.
fn cap_stream_chunks(chunks: Vec<String>, max: usize) -> Vec<String> {
if max == 0 || chunks.len() <= max {
return chunks;
}
let dropped: usize = chunks.iter().skip(max - 1).map(|c| c.chars().count()).sum();
let dropped_messages = chunks.len() - (max - 1);
let mut capped: Vec<String> = chunks.into_iter().take(max - 1).collect();
capped.push(format!(
"⚠️ 出力が長すぎるため打ち切りました(残り {dropped_messages} 件 / 約 {dropped} 文字)。\n\
同じ内容の繰り返しが続く場合はモデル側の異常です。"
));
capped
}

/// Below this the text is too short to judge; a runaway is never this small.
const DEGENERATE_MIN_CHARS: usize = 4_000;
/// A real answer of this many lines is never built from one or two distinct ones.
Expand All @@ -1780,10 +1744,9 @@ const DEGENERATE_MAX_DISTINCT: usize = 2;

/// Collapse output that is one short line repeated to fill the buffer.
///
/// Capping the number of messages stopped the flood, but the surviving messages
/// were still nothing but `court:` repeated — 2026-08-01 showed four of them
/// before the truncation notice. The reader gains nothing from seeing the same
/// line 2000 times; they need to know *what* repeated and *how much*.
/// A 2026-08-01 incident produced `court:` thousands of times. Full delivery is
/// allowed for varied content, but the reader gains nothing from seeing the same
/// line repeated; they need to know *what* repeated and *how much*.
///
/// The test is deliberately narrow: long output, many lines, and at most a
/// couple of distinct ones. A genuine answer with twenty lines has more variety
Expand Down Expand Up @@ -1831,10 +1794,7 @@ fn split_streaming_display(content: &str, message_limit: usize) -> Vec<String> {
if let Some(collapsed) = collapse_degenerate(content) {
return format::split_message(&collapsed, message_limit);
}
let chunks = cap_stream_chunks(
format::split_message(content, message_limit),
MAX_STREAM_MESSAGES,
);
let chunks = format::split_message(content, message_limit);
if chunks.is_empty() {
vec!["\u{200b}".to_string()]
} else {
Expand Down Expand Up @@ -1987,10 +1947,8 @@ struct StreamingDelivery {
mentions: Vec<String>,
}

/// Plan the final delivery so body edit + mention posts share the
/// `MAX_STREAM_MESSAGES` budget. Without this, a long mention-bearing output
/// splits the mentions unbounded and re-posts 6+ messages in one turn, bypassing
/// the runaway cap the edited path already enforces.
/// Plan the final delivery so mention paragraphs are posted as new messages
/// while the full non-mention body remains editable in place.
fn plan_streaming_delivery(final_content: &str, message_limit: usize) -> StreamingDelivery {
// Degenerate runaway collapses in place; its repetitions must not fan out as
// mention pings.
Expand All @@ -2007,23 +1965,8 @@ fn plan_streaming_delivery(final_content: &str, message_limit: usize) -> Streami
mentions: Vec::new(),
};
}
// Reserve one slot for the body when present so it is never fully crowded
// out, then bound the mention posts to what remains of the cap.
let mention_budget = MAX_STREAM_MESSAGES.saturating_sub(if body.is_empty() { 0 } else { 1 });
let mentions = cap_stream_chunks(
format::split_message(&mentions, message_limit),
mention_budget,
);
let body = if body.is_empty() {
None
} else {
let body_budget = MAX_STREAM_MESSAGES.saturating_sub(mentions.len());
// Whole `split_message` chunks are kept verbatim, so re-joining and
// re-editing yields the same capped count.
Some(
cap_stream_chunks(format::split_message(&body, message_limit), body_budget).join(""),
)
};
let mentions = format::split_message(&mentions, message_limit);
let body = if body.is_empty() { None } else { Some(body) };
StreamingDelivery { body, mentions }
}

Expand Down Expand Up @@ -2253,6 +2196,16 @@ fn compose_display(
mod tests {
use super::*;

fn assert_non_whitespace_preserved(chunks: &[String], expected: &str) {
let actual: String = chunks
.iter()
.flat_map(|chunk| chunk.chars())
.filter(|ch| !ch.is_whitespace())
.collect();
let expected: String = expected.chars().filter(|ch| !ch.is_whitespace()).collect();
assert_eq!(actual, expected);
}

/// Compile-time regression guard: use_streaming() is a required trait method
/// (no default). Any adapter that forgets to implement it will fail to compile.
/// This test documents the contract — see PR #503 / issue #502 for context.
Expand Down Expand Up @@ -2540,55 +2493,44 @@ mod tests {
let text = "court:\nsummary:\n".repeat(2_000);
let collapsed = collapse_degenerate(&text).expect("alternating repetition collapses");
// 反復単位の両方を残す。片方を黙って捨てない。
assert!(collapsed.contains("court:"), "first repeated line must survive");
assert!(collapsed.contains("summary:"), "second repeated line must survive");
assert!(
collapsed.contains("court:"),
"first repeated line must survive"
);
assert!(
collapsed.contains("summary:"),
"second repeated line must survive"
);
}

#[test]
fn a_mention_flood_stays_within_the_cap() {
// 長大な mention 出力が別経路の全件送信で件数上限を迂回しない。
fn a_long_mention_output_is_not_truncated() {
let mut content = String::new();
for i in 0..4000 {
content.push_str(&format!(
"<@123456789> 依頼 {i}: それぞれ異なる内容の指示です。\n\n"
));
}
let plan = plan_streaming_delivery(&content, 2000);
let body_msgs = plan
.body
.as_deref()
.map_or(0, |b| split_streaming_display(b, 2000).len());
let total = body_msgs + plan.mentions.len();
assert!(
total <= MAX_STREAM_MESSAGES,
"mention turn produced {total} messages"
);
// ping は生かす。少なくとも1通の mention が新規postとして残る。
assert!(
plan.mentions.iter().any(|m| m.contains("<@123456789>")),
"the mention must survive so the ping fires"
);
assert!(plan.body.is_none());
assert!(plan.mentions.len() > 5, "long output must not be capped");
assert_non_whitespace_preserved(&plan.mentions, content.trim());
}

#[test]
fn a_body_plus_mention_turn_is_capped_together() {
// 長大な body + 末尾の mention 段落。body の edit と mention の送信が
// 同じ上限を共有する。
fn a_long_body_plus_mention_is_preserved() {
let mut content = String::new();
for i in 0..4000 {
content.push_str(&format!("{i} 行目はそれぞれ異なる説明です。\n"));
}
content.push_str("\n\n<@123456789> レビューお願いします。");
let (expected_body, expected_mentions) = split_off_mention_paragraphs(&content);
let plan = plan_streaming_delivery(&content, 2000);
let body_msgs = plan
.body
.as_deref()
.map_or(0, |b| split_streaming_display(b, 2000).len());
let total = body_msgs + plan.mentions.len();
assert!(total <= MAX_STREAM_MESSAGES, "produced {total} messages");
assert_eq!(plan.body.as_deref(), Some(expected_body.as_str()));
assert_eq!(plan.mentions.concat(), expected_mentions);
assert!(
plan.mentions.iter().any(|m| m.contains("<@123456789>")),
"the mention must survive so the ping fires"
split_streaming_display(plan.body.as_deref().unwrap(), 2000).len() > 5,
"long body must not be capped"
);
}

Expand All @@ -2614,60 +2556,32 @@ mod tests {
}

#[test]
fn a_very_long_varied_turn_is_capped_to_a_few_messages() {
// 畳み込みに当たらない「ただ長い」出力は、件数上限で抑える。
// 2026-08-01 の実障害は繰り返しだったので畳み込み側が受けるが、
// 繰り返しでない暴走も同じだけチャンネルを埋める。
let mut runaway = String::new();
fn a_very_long_varied_turn_is_preserved() {
let mut answer = String::new();
for i in 0..4000 {
runaway.push_str(&format!("{i} 行目はそれぞれ異なる内容の説明です。\n"));
answer.push_str(&format!("{i} 行目はそれぞれ異なる内容の説明です。\n"));
}
let chunks = split_streaming_display(&runaway, 2000);
assert!(
chunks.len() <= MAX_STREAM_MESSAGES,
"runaway turn produced {} messages",
chunks.len()
);
assert!(
chunks.last().unwrap().contains("打ち切りました"),
"truncation must be visible to the reader"
);
let chunks = split_streaming_display(&answer, 2000);
assert!(chunks.len() > 5, "long varied output must not be capped");
assert_non_whitespace_preserved(&chunks, &answer);
}

#[test]
fn a_normal_long_answer_is_not_capped() {
// 上限は「長い回答」を殺すためのものではない。上限内なら素通しする。
fn a_normal_long_answer_is_preserved() {
let text = "本文".repeat(1500); // 2 messages at limit 2000
let chunks = split_streaming_display(&text, 2000);
assert!(chunks.len() > 1, "expected a multi-part answer");
assert!(chunks.len() <= MAX_STREAM_MESSAGES);
assert!(chunks.iter().all(|c| !c.contains("打ち切りました")));
assert_eq!(chunks.concat(), text);
}

#[test]
fn capping_keeps_the_head_of_the_output() {
// 先頭を残す。異常な繰り返しの末尾より、最初の方に意味がある。
fn a_very_long_single_line_is_preserved() {
let mut text = String::from("最初の重要な行\n");
text.push_str(&"x".repeat(60_000));
let chunks = split_streaming_display(&text, 2000);
assert!(chunks[0].starts_with("最初の重要な行"));
assert!(chunks.len() <= MAX_STREAM_MESSAGES);
}

#[test]
fn cap_of_zero_disables_the_limit() {
let chunks = vec!["a".to_string(); 40];
assert_eq!(cap_stream_chunks(chunks.clone(), 0).len(), 40);
}

#[test]
fn the_notice_reports_how_much_was_dropped() {
let chunks: Vec<String> = (0..10).map(|_| "y".repeat(2000)).collect();
let capped = cap_stream_chunks(chunks, 3);
assert_eq!(capped.len(), 3);
let notice = capped.last().unwrap();
assert!(notice.contains("8 件"), "notice was: {notice}");
assert!(notice.contains("16000"), "notice was: {notice}");
assert!(chunks.len() > 5, "long output must not be capped");
assert_non_whitespace_preserved(&chunks, &text);
}

#[test]
Expand Down
Loading