Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
24 commits
Select commit Hold shift + click to select a range
a9c545a
fix(webui): keep live projection subscription active
serrrfirat Jul 29, 2026
6b9ce44
fix(webui): restore provider-rate text streaming
serrrfirat Jul 29, 2026
5f3ab32
fix(llm): finalize terminal NEAR AI text streams
serrrfirat Jul 29, 2026
da6e6a9
Revert "fix(llm): finalize terminal NEAR AI text streams"
serrrfirat Jul 29, 2026
2586a17
fix(webui): preserve intermediate model phases
serrrfirat Jul 29, 2026
32e0b5c
fix(webui): render active replies with streamdown
serrrfirat Jul 29, 2026
4f9901b
fix(webui): keep one projection subscription per stream
serrrfirat Jul 30, 2026
931d6b8
fix(webui): prevent streamed markdown transition starvation
serrrfirat Jul 30, 2026
9761db3
fix(webui): conflate cumulative text at paint cadence
serrrfirat Jul 30, 2026
b498051
fix(webui): delegate SSE transport lifecycle
serrrfirat Jul 30, 2026
180d9ab
fix(webui): compact intermediate activity phases
serrrfirat Jul 30, 2026
c4c7b98
docs(webui): clarify intermediate phase boundary
serrrfirat Jul 30, 2026
40e33e4
test(webui): emulate packaged SSE transport
serrrfirat Jul 30, 2026
a0a97d8
fix(webui): keep activity runs collapsed
serrrfirat Jul 30, 2026
a5b292a
Merge branch 'main' into codex/fix-stalled-webchat-stream
serrrfirat Jul 30, 2026
88d5eba
fix(webui): address streaming review feedback
serrrfirat Jul 30, 2026
3030f49
ci: retry Railway preview deployment
serrrfirat Jul 30, 2026
c51db64
Merge branch 'main' into codex/fix-stalled-webchat-stream
serrrfirat Jul 30, 2026
b5e00f7
test(webui): follow continuous stream in gate refresh
serrrfirat Jul 30, 2026
3dc16f8
ci: reduce Reborn PR long tails
serrrfirat Jul 30, 2026
3b37d8a
Revert "ci: reduce Reborn PR long tails"
serrrfirat Jul 30, 2026
ca2a6f7
test(ci): satisfy changed-line coverage gate
serrrfirat Jul 30, 2026
a18288f
fix(ci): compare coverage against PR head
serrrfirat Jul 30, 2026
7c8024a
Merge remote-tracking branch 'origin/main' into codex/fix-stalled-web…
serrrfirat Jul 30, 2026
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
3 changes: 2 additions & 1 deletion .github/workflows/reborn-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -750,12 +750,13 @@ jobs:
if: github.event_name == 'pull_request'
env:
BASE_SHA: ${{ github.event.pull_request.base.sha }}
HEAD_SHA: ${{ github.event.pull_request.head.sha }}
run: |
set -euo pipefail
python3 scripts/ci/reborn_changed_coverage.py \
--lcov reborn-integration-merged.lcov \
--base "$BASE_SHA" \
--head "$(git rev-parse HEAD)" \
--head "$HEAD_SHA" \
--threshold 85 \
--json reborn-changed-coverage.json \
| tee -a "$GITHUB_STEP_SUMMARY"
Expand Down
53 changes: 52 additions & 1 deletion crates/ironclaw_host_api/src/product_surface.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ use std::sync::Mutex;

use async_trait::async_trait;
use serde::{Deserialize, Serialize};
use tokio::sync::{Mutex as AsyncMutex, mpsc};

use crate::{
ActivityId, AdapterInstallationId, AgentId, CapabilityId, ChannelInboundClassification,
Expand Down Expand Up @@ -171,11 +172,55 @@ pub struct ProductSurfaceStreamRequest {
}

/// Generic product event stream response.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[derive(Debug, PartialEq, Serialize, Deserialize)]
pub struct ProductSurfaceStreamResponse {
pub events: Vec<serde_json::Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub next_cursor: Option<String>,
/// Optional continuation of this same logical stream.
///
/// This handle is process-local transport state and is intentionally
/// omitted from the serialized response shape.
#[serde(skip)]
pub subscription: Option<ProductSurfaceEventSubscription>,
}

/// One continuous, caller-scoped product event subscription.
///
/// HTTP transports keep this receiver alive for the lifetime of the client
/// connection. This avoids the event-loss window created by repeatedly
/// tearing down and recreating one-event subscriptions.
pub struct ProductSurfaceEventSubscription {
receiver:
Arc<AsyncMutex<mpsc::Receiver<Result<ProductSurfaceStreamResponse, ProductSurfaceError>>>>,
}

impl ProductSurfaceEventSubscription {
pub fn new(
receiver: mpsc::Receiver<Result<ProductSurfaceStreamResponse, ProductSurfaceError>>,
) -> Self {
Self {
receiver: Arc::new(AsyncMutex::new(receiver)),
}
}

pub async fn next(&self) -> Option<Result<ProductSurfaceStreamResponse, ProductSurfaceError>> {
self.receiver.lock().await.recv().await
}
}

impl std::fmt::Debug for ProductSurfaceEventSubscription {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("ProductSurfaceEventSubscription")
.finish_non_exhaustive()
}
}

impl PartialEq for ProductSurfaceEventSubscription {
fn eq(&self, other: &Self) -> bool {
Arc::ptr_eq(&self.receiver, &other.receiver)
}
}

#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
Expand Down Expand Up @@ -574,6 +619,12 @@ mod tests {
use super::*;
use crate::{AgentId, ProjectId, TenantId, UserId};

#[test]
fn product_stream_continuation_is_single_consumer() {
static_assertions::assert_not_impl_any!(ProductSurfaceEventSubscription: Clone);
static_assertions::assert_not_impl_any!(ProductSurfaceStreamResponse: Clone);
}

fn caller() -> ProductSurfaceCaller {
ProductSurfaceCaller::new(
TenantId::new("tenant-alpha").expect("tenant"),
Expand Down
38 changes: 38 additions & 0 deletions crates/ironclaw_host_api/tests/host_api_contract.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1689,3 +1689,41 @@ fn dispatch_error_auth_required_debug_redacts_required_secrets() {
"provider id must not appear in Debug output; got: {debug_with_requirement}"
);
}

#[test]
fn product_stream_continuation_is_process_local_not_wire_state() {
let (_sender, receiver) = tokio::sync::mpsc::channel(1);
let response = ProductSurfaceStreamResponse {
events: vec![json!({"kind": "projection_update"})],
next_cursor: Some("cursor-1".to_string()),
subscription: Some(ProductSurfaceEventSubscription::new(receiver)),
};

let encoded = serde_json::to_value(response).expect("stream response serializes");
assert_eq!(
encoded,
json!({
"events": [{"kind": "projection_update"}],
"next_cursor": "cursor-1"
}),
"the in-process continuation handle must never cross the wire"
);
let decoded: ProductSurfaceStreamResponse =
serde_json::from_value(encoded).expect("wire response deserializes");
assert!(decoded.subscription.is_none());
}

#[test]
fn product_stream_continuation_debug_and_identity_are_stable() {
let (_sender, receiver) = tokio::sync::mpsc::channel(1);
let subscription = ProductSurfaceEventSubscription::new(receiver);
let (_other_sender, other_receiver) = tokio::sync::mpsc::channel(1);
let other_subscription = ProductSurfaceEventSubscription::new(other_receiver);

assert_eq!(subscription, subscription);
assert_ne!(subscription, other_subscription);
assert_eq!(
format!("{subscription:?}"),
"ProductSurfaceEventSubscription { .. }"
);
}
64 changes: 64 additions & 0 deletions crates/ironclaw_llm/src/nearai_chat.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
//! - **API key auth**: When `NEARAI_API_KEY` is set, uses Bearer API key
//! - **Session token auth**: Otherwise, uses `SessionManager` for Bearer session token
//! with automatic renewal on 401 errors
// arch-exempt: large_file, provider-local streaming regression tests require private parser access pending provider adapter decomposition, plan #6175

use std::collections::HashMap;
use std::sync::Arc;
Expand Down Expand Up @@ -1909,6 +1910,69 @@ mod tests {
.await
}

async fn complete_search_tool_streaming_from_sse(
sse_body: &'static str,
) -> ToolCompletionResponse {
use tokio::io::AsyncWriteExt;
use tokio::net::TcpListener;

let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let base_url = format!("http://{}", listener.local_addr().unwrap());
let server_task = tokio::spawn(async move {
let (mut socket, _) = accept_chat_request(&listener).await;
socket
.write_all(
b"HTTP/1.1 200 OK\r\ncontent-type: text/event-stream\r\ncache-control: no-cache\r\nconnection: close\r\n\r\n",
)
.await
.expect("write sse headers");
socket
.write_all(sse_body.as_bytes())
.await
.expect("write sse body");
});

let provider = NearAiChatProvider::new(test_nearai_config(&base_url), test_session())
.expect("provider");
let response = complete_search_tool_streaming(provider)
.await
.expect("streaming completion");
server_task.await.expect("server task");
response
}

#[tokio::test]
async fn complete_with_tools_streaming_infers_tool_use_without_finish_reason() {
let response = complete_search_tool_streaming_from_sse(
r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_search","type":"function","function":{"name":"search","arguments":"{\"query\":\"near ai\"}"}}]},"finish_reason":null}]}

data: [DONE]

"#,
)
.await;

assert_eq!(response.finish_reason, FinishReason::ToolUse);
assert_eq!(response.tool_calls.len(), 1);
assert_eq!(response.tool_calls[0].name, "search");
}

#[tokio::test]
async fn complete_with_tools_streaming_preserves_unknown_finish_reason() {
let response = complete_search_tool_streaming_from_sse(
r#"data: {"choices":[{"delta":{"content":"answer"},"finish_reason":"vendor_specific"}]}

data: [DONE]

"#,
)
.await;

assert_eq!(response.finish_reason, FinishReason::Unknown);
assert_eq!(response.content.as_deref(), Some("answer"));
assert!(response.tool_calls.is_empty());
}

#[tokio::test]
async fn complete_streaming_emits_delta_before_response_completes() {
use tokio::io::AsyncWriteExt;
Expand Down
Loading
Loading