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
6 changes: 6 additions & 0 deletions src/apps/cli/AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -142,6 +142,12 @@ cargo check -p openbitfun-cli
cargo test -p openbitfun-cli
```

For streaming `exec` retry, context recovery, and final-event contracts:

```bash
cargo test --locked -p openbitfun-cli --test cli_command_contracts exec_cli_contracts::stream_json_
```

When a CLI change crosses a shared boundary, use the focused command maintained
by that owner: Agent Runtime for port/SDK behavior, the IPC adapter for shared
protocol behavior, Core for turn/tool/persistence behavior, Terminal for
Expand Down
60 changes: 56 additions & 4 deletions src/apps/cli/tests/cli_command_contracts/exec_cli_contracts.rs
Original file line number Diff line number Diff line change
Expand Up @@ -448,6 +448,58 @@ fn stream_json_malformed_sse_retries_then_completes() {
);
}

#[test]
fn stream_json_context_overflow_compresses_before_reissuing_the_model_request() {
let server = MockOpenAiServer::context_overflow_then_immediate();
let environment = CliTestEnvironment::new();
environment.configure_mock_model(server.base_url());
let mut command = environment.std_command();
command.args([
"exec",
"Remember this request and continue after context recovery",
"--output-format",
"stream-json",
]);
let output = command_output_with_timeout(&mut command, std::time::Duration::from_secs(30));
let stdout = stdout(&output);
assert!(output.status.success(), "{}\n{stdout}", stderr(&output));
// One rejected request, one summary request, then the recovered model round.
server.assert_chat_completion_requests(3);
let requests = server.chat_completion_request_bodies();
assert_ne!(
requests[0]["messages"], requests[1]["messages"],
"overflow must enter compression instead of replaying the original request"
);
assert_ne!(
requests[0]["messages"], requests[2]["messages"],
"recovery must send the compressed context"
);
let events = jsonl_events(&stdout);
let compression_started = events
.iter()
.position(|value| {
value["event"]["type"] == "ContextCompressionStarted"
&& value["event"]["trigger"] == "context_overflow_recovery"
})
.expect("overflow should start recovery compression");
let compression_completed = events
.iter()
.position(|value| value["event"]["type"] == "ContextCompressionCompleted")
.expect("recovery compression should complete");
assert!(compression_started < compression_completed);
assert_eq!(
events
.iter()
.filter(|value| is_terminal_event(value))
.count(),
1
);
assert_eq!(
events.last().unwrap()["event"]["type"],
"DialogTurnCompleted"
);
}

#[test]
fn stream_json_provider_http_403_emits_one_error_terminal() {
let server = MockOpenAiServer::http_403("provider authorization denied");
Expand All @@ -461,7 +513,7 @@ fn stream_json_provider_http_403_emits_one_error_terminal() {
"stream-json",
]);
let output = command_output_with_timeout(&mut command, std::time::Duration::from_secs(30));
server.assert_chat_completion_requests(10);
server.assert_chat_completion_requests(1);

let stdout = stdout(&output);
assert!(!output.status.success(), "{stdout}");
Expand Down Expand Up @@ -524,7 +576,7 @@ fn stream_json_provider_and_patch_failures_publish_one_final_classification() {
&output_target,
]);
let output = command_output_with_timeout(&mut command, std::time::Duration::from_secs(30));
server.assert_chat_completion_requests(10);
server.assert_chat_completion_requests(1);

let stdout = stdout(&output);
let stderr = stderr(&output);
Expand Down Expand Up @@ -562,7 +614,7 @@ fn stream_json_provider_and_patch_failures_publish_one_final_classification() {
}

#[test]
fn stream_json_disconnect_then_exhausted_retry_failure_emits_one_error_terminal() {
fn stream_json_disconnect_then_authorization_failure_emits_one_error_terminal() {
let server = MockOpenAiServer::disconnect_then_http_403();
let environment = CliTestEnvironment::new();
environment.configure_mock_model(server.base_url());
Expand All @@ -574,7 +626,7 @@ fn stream_json_disconnect_then_exhausted_retry_failure_emits_one_error_terminal(
"stream-json",
]);
let output = command_output_with_timeout(&mut command, std::time::Duration::from_secs(30));
server.assert_chat_completion_requests(10);
server.assert_chat_completion_requests(2);

let stdout = stdout(&output);
assert!(!output.status.success(), "{stdout}");
Expand Down
18 changes: 18 additions & 0 deletions src/apps/cli/tests/support/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -320,6 +320,7 @@ enum MockModelResponse {
Http403 { reason: String },
DisconnectThenHttp403,
MalformedSseThenImmediate,
ContextOverflowThenImmediate,
}

impl MockOpenAiServer {
Expand Down Expand Up @@ -349,6 +350,10 @@ impl MockOpenAiServer {
Self::spawn(MockModelResponse::MalformedSseThenImmediate)
}

pub(crate) fn context_overflow_then_immediate() -> Self {
Self::spawn(MockModelResponse::ContextOverflowThenImmediate)
}

pub(crate) fn base_url(&self) -> &str {
&self.base_url
}
Expand Down Expand Up @@ -437,6 +442,10 @@ impl MockOpenAiServer {
| MockModelResponse::DisconnectThenHttp403
) || (matches!(response, MockModelResponse::MalformedSseThenImmediate)
&& attempt < 2)
|| (matches!(
response,
MockModelResponse::ContextOverflowThenImmediate
) && attempt < 3)
|| (matches!(response, MockModelResponse::ProductControlLoop)
&& attempt < 5);
if accepts_more_requests {
Expand Down Expand Up @@ -488,6 +497,15 @@ fn serve_model_response(
release_stream: &mpsc::Receiver<()>,
stream_disconnected: &mpsc::Sender<()>,
) {
if matches!(response, MockModelResponse::ContextOverflowThenImmediate) && attempt == 0 {
let body = json!({
"error": {"code": "context_length_exceeded", "message": "Maximum context length exceeded"}
}).to_string();
write!(stream, "HTTP/1.1 400 Bad Request\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}", body.len())
.expect("write context overflow response");
stream.flush().expect("flush context overflow response");
return;
}
if matches!(response, MockModelResponse::DisconnectThenHttp403) && attempt > 0 {
write_http_403(stream, "provider stream remained unavailable")
.expect("write post-disconnect HTTP error");
Expand Down
6 changes: 6 additions & 0 deletions src/crates/assembly/core/src/agentic/execution/AGENTS.md
Original file line number Diff line number Diff line change
@@ -1,5 +1,11 @@
If you modify `stream_processor.rs`, run the stream integration tests before finishing.

For model retry admission and recovery, use:

```bash
cargo test --locked -p openbitfun-core --no-default-features --features agent-runtime,git --lib agentic::execution::round_executor::tests
```

For complete shell constraint checks, use:

```bash
Expand Down
Loading
Loading