Skip to content
Open
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
77 changes: 72 additions & 5 deletions openless-all/app/src-tauri/src/remote_server/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -421,9 +421,14 @@ fn send_json<T: Serialize>(value: &T) -> Message {
/// Phones only receive the Core session opened by their own connection. Global capsule /
/// plain-text broadcasts have no owner and would send local PC dictation or another phone's
/// content to all paired connections; they must not be a network egress.
///
/// `wants_partials`: the client asked for live text in its hello frame. The owned session's
/// streaming transcript (`TranscriptDelta`) is then forwarded as `partial` frames; it is the
/// same session-scoped text the desktop capsule shows, and never another session's.
fn backend_event_to_phone(
event: &openless_core::BackendEvent,
remote_session_id: &mut Option<openless_core::SessionId>,
wants_partials: bool,
) -> Vec<String> {
use openless_core::{BackendEventKind, DictationPhase};
if remote_session_id.is_none() || event.session_id != *remote_session_id {
Expand Down Expand Up @@ -468,6 +473,13 @@ fn backend_event_to_phone(
}
messages
}
BackendEventKind::TranscriptDelta(delta) if wants_partials => {
// Same contract as TranscriptAccumulator::apply: replace from `offset` (chars).
vec![serde_json::json!({
"type":"partial", "text":delta.text, "offset":delta.offset, "final":delta.is_final,
})
.to_string()]
}
_ => Vec::new(),
}
}
Expand All @@ -481,9 +493,11 @@ type PendingRemoteStop =
async fn handle_ws(mut socket: WebSocket, state: Arc<WsState>, peer_ip: IpAddr) {
// 1) Handshake: wait for the first hello + PIN frame.
let connection_id = openless_core::SessionId::new();
let mut wants_partials = false;
let authed = match tokio::time::timeout(Duration::from_secs(15), socket.recv()).await {
Ok(Some(Ok(Message::Text(txt)))) => {
let candidate = parse_hello_pin(&txt);
wants_partials = parse_hello_partials(&txt);
match state
.backend
.services()
Expand Down Expand Up @@ -609,7 +623,9 @@ async fn handle_ws(mut socket: WebSocket, state: Arc<WsState>, peer_ip: IpAddr)
continue;
}
last_event_sequence = event.sequence;
for msg in backend_event_to_phone(&event, &mut remote_session_id) {
let messages =
backend_event_to_phone(&event, &mut remote_session_id, wants_partials);
for msg in messages {
if socket.send(Message::Text(msg)).await.is_err() {
break 'connection;
}
Expand Down Expand Up @@ -816,6 +832,16 @@ fn parse_hello_pin(txt: &str) -> openless_core::SecretValue {
openless_core::SecretValue::new(pin)
}

/// Optional `"partials": true` in the hello frame opts this connection into live text.
/// Clients that do not send it (the bundled phone page) keep the previous message set.
fn parse_hello_partials(txt: &str) -> bool {
serde_json::from_str::<serde_json::Value>(txt)
.ok()
.filter(|value| value.get("type").and_then(serde_json::Value::as_str) == Some("hello"))
.and_then(|value| value.get("partials").and_then(serde_json::Value::as_bool))
.unwrap_or(false)
}

fn parse_audio_frame(
frame: &[u8],
) -> Result<(openless_core::SessionId, u64, Vec<u8>), openless_core::BackendError> {
Expand All @@ -831,7 +857,10 @@ mod tests {
RemoteInputConfig, RemoteInputService, SessionId,
};

use super::{apply_remote_control, backend_event_to_phone, parse_audio_frame, parse_hello_pin};
use super::{
apply_remote_control, backend_event_to_phone, parse_audio_frame, parse_hello_partials,
parse_hello_pin,
};

fn backend() -> (
OpenLessBackend,
Expand Down Expand Up @@ -1050,23 +1079,61 @@ mod tests {
}),
};
assert!(
backend_event_to_phone(&result(SessionId::new()), &mut owner).is_empty(),
backend_event_to_phone(&result(SessionId::new()), &mut owner, false).is_empty(),
"其他手机/本机结果不可转发"
);
assert_eq!(owner, Some(session));
let messages = backend_event_to_phone(&result(session), &mut owner);
let messages = backend_event_to_phone(&result(session), &mut owner, false);
assert_eq!(messages.len(), 2);
assert_eq!(
serde_json::from_str::<serde_json::Value>(&messages[1]).unwrap(),
serde_json::json!({"type":"result", "text":"自己的结果"})
);
assert_eq!(owner, None);
assert!(
backend_event_to_phone(&result(session), &mut owner).is_empty(),
backend_event_to_phone(&result(session), &mut owner, false).is_empty(),
"取消/终态之后的迟到结果不可转发"
);
}

#[test]
fn websocket_forwards_live_transcript_only_when_asked_and_only_for_the_owned_session() {
use openless_core::{BackendEvent, BackendEventKind, TranscriptDelta};
let session = SessionId::new();
let mut owner = Some(session);
let delta = |id, text: &str, offset| BackendEvent {
sequence: 1,
session_id: Some(id),
kind: BackendEventKind::TranscriptDelta(TranscriptDelta {
text: text.into(),
offset,
is_final: false,
}),
};
assert!(
backend_event_to_phone(&delta(session, "你好", 0), &mut owner, false).is_empty(),
"未在 hello 中请求实时文字的客户端保持原有消息集"
);
assert!(
backend_event_to_phone(&delta(SessionId::new(), "别人的", 0), &mut owner, true)
.is_empty(),
"其他会话的实时文字不可转发"
);
let messages = backend_event_to_phone(&delta(session, "您好", 1), &mut owner, true);
assert_eq!(messages.len(), 1);
assert_eq!(
serde_json::from_str::<serde_json::Value>(&messages[0]).unwrap(),
serde_json::json!({"type":"partial", "text":"您好", "offset":1, "final":false})
);
assert_eq!(owner, Some(session), "实时文字不结束会话所有权");

assert!(parse_hello_partials(
r#"{"type":"hello","pin":"123456","partials":true}"#
));
assert!(!parse_hello_partials(r#"{"type":"hello","pin":"123456"}"#));
assert!(!parse_hello_partials(r#"{"type":"other","partials":true}"#));
}

#[test]
fn websocket_wire_parser_requires_contract_2_frames() {
let session = SessionId::new();
Expand Down