Skip to content
Open
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
15 changes: 6 additions & 9 deletions apps/codex-plus-launcher/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -188,16 +188,13 @@ async fn activate_existing_codex_app(options: &LaunchOptions) -> anyhow::Result<
hooks.start_helper(helper_port).await?;
}
let process_ids = codex_plus_core::watcher::find_codex_processes();
let mut activated = false;
#[cfg(windows)]
{
for process_id in &process_ids {
if codex_plus_core::windows_activate_process_window(*process_id) {
activated = true;
break;
}
}
}
let activated = process_ids
.iter()
.copied()
.any(codex_plus_core::windows_activate_process_window);
#[cfg(not(windows))]
let activated = false;
let injection_ready = if settings.enhancements_enabled {
hooks
.ensure_injection(options.debug_port, helper_port, &app_dir)
Expand Down
266 changes: 214 additions & 52 deletions crates/codex-plus-core/src/bridge.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ use std::time::Duration;

use anyhow::{Context, bail};
use base64::Engine;
use futures_util::stream::FuturesUnordered;
use futures_util::{SinkExt, StreamExt};
use serde_json::{Value, json};
use tokio_tungstenite::connect_async;
Expand All @@ -26,6 +27,69 @@ pub type BridgeHandler = Arc<

static NEXT_MESSAGE_ID: AtomicU64 = AtomicU64::new(100);

/// Bridge 会话按注入目标分代。
///
/// 同一目标再次安装 Bridge 时,旧会话会在下一次消息循环中退出并关闭 socket,
/// 避免多份 CDP 会话同时应答同一个页面请求。不同目标互不影响。
static NEXT_BRIDGE_GENERATION: AtomicU64 = AtomicU64::new(1);
static CURRENT_BRIDGE_GENERATIONS: std::sync::LazyLock<std::sync::Mutex<HashMap<String, u64>>> =
std::sync::LazyLock::new(|| std::sync::Mutex::new(HashMap::new()));

#[derive(Clone)]
struct BridgeGeneration {
target: String,
id: u64,
}

type PendingBridgeCall = Pin<Box<dyn Future<Output = CompletedBridgeCall> + Send>>;

struct CompletedBridgeCall {
request_id: String,
generation: Option<BridgeGeneration>,
response: Result<Value, String>,
}

fn next_bridge_generation(target: &str) -> BridgeGeneration {
BridgeGeneration {
target: target.to_string(),
id: NEXT_BRIDGE_GENERATION.fetch_add(1, Ordering::SeqCst),
}
}

fn publish_bridge_generation(generation: &BridgeGeneration) -> bool {
let mut generations = CURRENT_BRIDGE_GENERATIONS
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
if generations
.get(&generation.target)
.is_some_and(|current| *current > generation.id)
{
return false;
}
generations.insert(generation.target.clone(), generation.id);
true
}

fn bridge_generation_is_current(generation: &BridgeGeneration) -> bool {
CURRENT_BRIDGE_GENERATIONS
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.get(&generation.target)
.is_some_and(|current| *current == generation.id)
}

fn release_bridge_generation(generation: &BridgeGeneration) {
let mut generations = CURRENT_BRIDGE_GENERATIONS
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
if generations
.get(&generation.target)
.is_some_and(|current| *current == generation.id)
{
generations.remove(&generation.target);
}
}

pub fn build_bridge_script(binding_name: &str) -> String {
format!(
r#"
Expand Down Expand Up @@ -192,6 +256,8 @@ pub async fn install_bridge(
) -> anyhow::Result<()> {
let socket = connect_cdp_websocket(websocket_url).await?;
let mut session = CdpSession::new(socket).with_handler(handler);
let generation = next_bridge_generation(websocket_url);
session = session.with_generation(generation.clone());

session.send_command(1, "Runtime.enable", json!({})).await?;
session
Expand Down Expand Up @@ -236,17 +302,51 @@ pub async fn install_bridge(
.await?;
}

session.drain_binding_queue().await?;
if !publish_bridge_generation(&generation) {
let _ = crate::diagnostic_log::append_diagnostic_log(
"bridge.generation_superseded_before_publish",
json!({ "generation": generation.id }),
);
session.close().await;
return Ok(());
}
let _ = crate::diagnostic_log::append_diagnostic_log(
"bridge.generation_published",
json!({ "generation": generation.id }),
);

let mut pending_calls = FuturesUnordered::new();
session.enqueue_binding_calls(&mut pending_calls);
tokio::spawn(async move {
loop {
if session.drain_binding_queue().await.is_err() {
if !bridge_generation_is_current(&generation) {
let _ = crate::diagnostic_log::append_diagnostic_log(
"bridge.generation_superseded",
json!({ "generation": generation.id }),
);
break;
}
match session.next_message().await {
Ok(Some(_)) => {}
Ok(None) | Err(_) => break,

session.enqueue_binding_calls(&mut pending_calls);
tokio::select! {
completed = pending_calls.next(), if !pending_calls.is_empty() => {
let Some(completed) = completed else {
continue;
};
if session.finish_binding_call(completed).await.is_err() {
break;
}
}
message = session.next_message() => {
match message {
Ok(Some(_)) => {}
Ok(None) | Err(_) => break,
}
}
}
}
session.close().await;
release_bridge_generation(&generation);
});

Ok(())
Expand Down Expand Up @@ -308,6 +408,7 @@ struct CdpSession<S> {
responses: HashMap<u64, Value>,
binding_calls: VecDeque<Value>,
handler: Option<BridgeHandler>,
generation: Option<BridgeGeneration>,
}

impl<S> CdpSession<S>
Expand All @@ -324,6 +425,7 @@ where
responses: HashMap::new(),
binding_calls: VecDeque::new(),
handler: None,
generation: None,
}
}

Expand All @@ -332,6 +434,22 @@ where
self
}

fn with_generation(mut self, generation: BridgeGeneration) -> Self {
self.generation = Some(generation);
self
}

fn is_current(&self) -> bool {
self.generation
.as_ref()
.is_none_or(bridge_generation_is_current)
}

async fn close(&mut self) {
let _ = self.socket.send(Message::Close(None)).await;
let _ = self.socket.close().await;
}

async fn send_command(
&mut self,
message_id: u64,
Expand Down Expand Up @@ -421,73 +539,117 @@ where
Ok(Some(value))
}

async fn drain_binding_queue(&mut self) -> anyhow::Result<()> {
fn enqueue_binding_calls(&mut self, pending_calls: &mut FuturesUnordered<PendingBridgeCall>) {
while let Some(message) = self.binding_calls.pop_front() {
self.route_binding_call(message).await?;
self.enqueue_binding_call(message, pending_calls);
}
Ok(())
}

fn route_binding_call(
fn enqueue_binding_call(
&mut self,
message: Value,
) -> Pin<Box<dyn Future<Output = anyhow::Result<()>> + Send + '_>> {
Box::pin(async move {
let Some(handler) = self.handler.clone() else {
return Ok(());
};

let Some(payload_text) = message
.get("params")
.and_then(|params| params.get("payload"))
.and_then(Value::as_str)
else {
return Ok(());
};
pending_calls: &mut FuturesUnordered<PendingBridgeCall>,
) {
let Some(handler) = self.handler.clone() else {
return;
};

let parsed: Value = match serde_json::from_str(payload_text) {
Ok(parsed) => parsed,
Err(error) => {
if let Some(request_id) = extract_string_field(payload_text, "id") {
self.reject_bridge_request(
&request_id,
&format!("failed to parse bridge payload: {error}"),
)
.await?;
}
return Ok(());
}
};
self.route_parsed_binding_call(&handler, parsed).await
})
}
let Some(payload_text) = message
.get("params")
.and_then(|params| params.get("payload"))
.and_then(Value::as_str)
else {
return;
};

async fn route_parsed_binding_call(
&mut self,
handler: &BridgeHandler,
parsed: Value,
) -> anyhow::Result<()> {
let Some(request_id) = parsed.get("id").and_then(Value::as_str) else {
return Ok(());
let parsed: Value = match serde_json::from_str(payload_text) {
Ok(parsed) => parsed,
Err(error) => {
let Some(request_id) = extract_string_field(payload_text, "id") else {
return;
};
self.enqueue_completed_binding_call(
request_id,
Err(format!("failed to parse bridge payload: {error}")),
pending_calls,
);
return;
}
};
let Some(request_id) = parsed.get("id").and_then(Value::as_str).map(str::to_string) else {
return;
};
if !self.is_current() {
let _ = crate::diagnostic_log::append_diagnostic_log(
"bridge.stale_request_dropped",
json!({
"request_id": request_id,
"generation": self.generation.as_ref().map(|generation| generation.id)
}),
);
return;
}
let path = parsed
.get("path")
.and_then(Value::as_str)
.unwrap_or_default()
.to_string();
let payload = parsed.get("payload").cloned().unwrap_or_else(|| json!({}));
let generation = self.generation.clone();

pending_calls.push(Box::pin(async move {
CompletedBridgeCall {
request_id,
generation,
response: handler(path, payload)
.await
.map_err(|error| error.to_string()),
}
}));
}

match handler(path, payload).await {
fn enqueue_completed_binding_call(
&self,
request_id: String,
response: Result<Value, String>,
pending_calls: &mut FuturesUnordered<PendingBridgeCall>,
) {
let generation = self.generation.clone();
pending_calls.push(Box::pin(async move {
CompletedBridgeCall {
request_id,
generation,
response,
}
}));
}

async fn finish_binding_call(&mut self, completed: CompletedBridgeCall) -> anyhow::Result<()> {
if completed
.generation
.as_ref()
.is_some_and(|generation| !bridge_generation_is_current(generation))
{
let _ = crate::diagnostic_log::append_diagnostic_log(
"bridge.stale_response_dropped",
json!({
"request_id": completed.request_id,
"generation": completed.generation.as_ref().map(|generation| generation.id)
}),
);
return Ok(());
}

match completed.response {
Ok(result) => {
self.resolve_bridge_request(request_id, &result).await?;
self.resolve_bridge_request(&completed.request_id, &result)
.await
}
Err(error) => {
self.reject_bridge_request(request_id, &error.to_string())
.await?;
Err(message) => {
self.reject_bridge_request(&completed.request_id, &message)
.await
}
}

Ok(())
}

async fn resolve_bridge_request(
Expand Down
Loading
Loading