mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(dictation): single-request connect, clear transcribing on completion
1. Merge session + SDP into single POST /transcribe/connect The two-step flow (POST /transcribe/session → POST /transcribe/sdp) stored the OpenAI client secret in an in-memory DashMap, which breaks in HA deployments where the two requests may land on different relay replicas. Replaced with a single POST /transcribe/connect that accepts the SDP offer, mints the OpenAI session, proxies the SDP exchange, and returns the SDP answer + model — all in one request. No server-side session state is needed between requests, so this works correctly across any number of replicas. Removed the transcribe_sessions DashMap from AppState entirely. 2. Clear isTranscribing even when final text is unchanged When a TRANSCRIPT_COMPLETED_EVENT produces no text change (final matches accumulated deltas), the early-return in handleRealtimeEvent skipped the setIsTranscribing(false) call. This left the 'Transcribing…' indicator stuck indefinitely. Fixed by updating the transcribing flag before the text-change early return.
This commit is contained in:
@@ -18,7 +18,6 @@ use axum::{
|
||||
};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use serde_json::Value;
|
||||
use uuid::Uuid;
|
||||
|
||||
use buzz_core::CommunityId;
|
||||
|
||||
@@ -31,10 +30,6 @@ const OPENAI_REALTIME_CLIENT_SECRETS_URL: &str =
|
||||
|
||||
const OPENAI_REALTIME_CALLS_URL: &str = "https://api.openai.com/v1/realtime/calls";
|
||||
|
||||
/// Maximum age of a cached transcription session secret before it's considered
|
||||
/// expired. Matches the `expires_after.seconds` sent to OpenAI.
|
||||
const SESSION_SECRET_TTL: Duration = Duration::from_secs(60);
|
||||
|
||||
/// Rate-limit window for transcription session minting.
|
||||
const TRANSCRIBE_RATE_WINDOW: Duration = Duration::from_secs(60);
|
||||
|
||||
@@ -45,26 +40,17 @@ pub struct TranscribeStatus {
|
||||
model: String,
|
||||
}
|
||||
|
||||
/// Response for `POST /transcribe/session`.
|
||||
#[derive(Serialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub struct TranscribeSession {
|
||||
session_id: String,
|
||||
model: String,
|
||||
}
|
||||
|
||||
/// Request body for `POST /transcribe/sdp`.
|
||||
/// Request body for `POST /transcribe/connect`.
|
||||
#[derive(Deserialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub struct SdpExchangeRequest {
|
||||
session_id: String,
|
||||
pub struct TranscribeConnectRequest {
|
||||
sdp: String,
|
||||
}
|
||||
|
||||
/// Response for `POST /transcribe/sdp`.
|
||||
/// Response for `POST /transcribe/connect`.
|
||||
#[derive(Serialize)]
|
||||
pub struct SdpExchangeResponse {
|
||||
pub struct TranscribeConnectResponse {
|
||||
sdp: String,
|
||||
model: String,
|
||||
}
|
||||
|
||||
/// `GET /transcribe/status` — check if transcription is configured.
|
||||
@@ -105,11 +91,20 @@ pub async fn transcribe_status(
|
||||
/// Each session mints a metered OpenAI Realtime connection on the relay
|
||||
/// operator's bill, so we require the caller to be an actual relay member
|
||||
/// regardless of `BUZZ_REQUIRE_RELAY_MEMBERSHIP`.
|
||||
pub async fn create_transcribe_session(
|
||||
/// `POST /transcribe/connect` — mint an OpenAI session and proxy the SDP exchange.
|
||||
///
|
||||
/// Accepts the client's SDP offer, mints an ephemeral OpenAI Realtime session,
|
||||
/// uses the returned client secret to proxy the SDP offer to OpenAI's
|
||||
/// `/v1/realtime/calls` endpoint, and returns the SDP answer + model. The
|
||||
/// entire flow completes in a single request so no server-side session state
|
||||
/// is needed — this works correctly across multiple relay replicas in HA
|
||||
/// deployments without shared storage.
|
||||
pub async fn transcribe_connect(
|
||||
State(state): State<Arc<AppState>>,
|
||||
headers: HeaderMap,
|
||||
) -> Result<Json<TranscribeSession>, (StatusCode, Json<Value>)> {
|
||||
let (pubkey, community) = authenticate(&state, &headers, "/transcribe/session", "POST").await?;
|
||||
Json(body): Json<TranscribeConnectRequest>,
|
||||
) -> Result<Json<TranscribeConnectResponse>, (StatusCode, Json<Value>)> {
|
||||
let (pubkey, community) = authenticate(&state, &headers, "/transcribe/connect", "POST").await?;
|
||||
|
||||
// Hard membership gate — billable endpoint requires actual relay membership
|
||||
// even on open relays where `enforce_relay_membership` would short-circuit.
|
||||
@@ -134,15 +129,11 @@ pub async fn create_transcribe_session(
|
||||
})?;
|
||||
|
||||
let model = state.config.transcription_model.clone();
|
||||
|
||||
let client = openai_client().map_err(|(status, msg)| api_error(status, msg))?;
|
||||
|
||||
// Build the session payload — VAD configuration depends on the model.
|
||||
// `gpt-realtime-whisper` requires manual audio commit (no turn detection),
|
||||
// while other models (e.g. `whisper-1`) use server-side VAD.
|
||||
// Step 1: Mint an ephemeral client secret from OpenAI.
|
||||
let session_payload = build_session_payload(&model);
|
||||
|
||||
let response = client
|
||||
let session_response = client
|
||||
.post(OPENAI_REALTIME_CLIENT_SECRETS_URL)
|
||||
.header("Authorization", format!("Bearer {api_key}"))
|
||||
.header("Content-Type", "application/json")
|
||||
@@ -157,17 +148,17 @@ pub async fn create_transcribe_session(
|
||||
)
|
||||
})?;
|
||||
|
||||
if !response.status().is_success() {
|
||||
let status = response.status();
|
||||
let body = response.text().await.unwrap_or_default();
|
||||
tracing::error!("OpenAI realtime session error ({status}): {body}");
|
||||
if !session_response.status().is_success() {
|
||||
let status = session_response.status();
|
||||
let resp_body = session_response.text().await.unwrap_or_default();
|
||||
tracing::error!("OpenAI realtime session error ({status}): {resp_body}");
|
||||
return Err(api_error(
|
||||
StatusCode::BAD_GATEWAY,
|
||||
"OpenAI rejected the transcription session request",
|
||||
));
|
||||
}
|
||||
|
||||
let body: Value = response.json().await.map_err(|e| {
|
||||
let session_body: Value = session_response.json().await.map_err(|e| {
|
||||
tracing::error!("OpenAI realtime session response parse error: {e}");
|
||||
api_error(
|
||||
StatusCode::BAD_GATEWAY,
|
||||
@@ -175,62 +166,16 @@ pub async fn create_transcribe_session(
|
||||
)
|
||||
})?;
|
||||
|
||||
let client_secret = extract_client_secret(&body).ok_or_else(|| {
|
||||
tracing::error!("OpenAI realtime session response missing client_secret: {body}");
|
||||
let client_secret = extract_client_secret(&session_body).ok_or_else(|| {
|
||||
tracing::error!("OpenAI realtime session response missing client_secret: {session_body}");
|
||||
api_error(
|
||||
StatusCode::BAD_GATEWAY,
|
||||
"transcription service returned unexpected response",
|
||||
)
|
||||
})?;
|
||||
|
||||
// Store the secret server-side — the client receives only an opaque session
|
||||
// ID and must call `/transcribe/sdp` to complete the WebRTC handshake.
|
||||
let session_id = Uuid::new_v4().to_string();
|
||||
state
|
||||
.transcribe_sessions
|
||||
.insert(session_id.clone(), (client_secret, Instant::now()));
|
||||
|
||||
Ok(Json(TranscribeSession { session_id, model }))
|
||||
}
|
||||
|
||||
/// `POST /transcribe/sdp` — proxy the WebRTC SDP exchange to OpenAI.
|
||||
///
|
||||
/// Accepts the client's SDP offer and the session ID returned by
|
||||
/// `/transcribe/session`. The relay looks up the cached client secret,
|
||||
/// forwards the SDP offer to OpenAI's `/v1/realtime/calls` endpoint, and
|
||||
/// returns the SDP answer. The client never sees the bearer token.
|
||||
pub async fn proxy_sdp_exchange(
|
||||
State(state): State<Arc<AppState>>,
|
||||
headers: HeaderMap,
|
||||
Json(body): Json<SdpExchangeRequest>,
|
||||
) -> Result<Json<SdpExchangeResponse>, (StatusCode, Json<Value>)> {
|
||||
// Authenticate — same NIP-98 requirement as session creation.
|
||||
let pubkey = authenticate(&state, &headers, "/transcribe/sdp", "POST").await?;
|
||||
require_relay_member(&state, &headers, &pubkey).await?;
|
||||
|
||||
// Look up the cached client secret.
|
||||
let (client_secret, created_at) = state
|
||||
.transcribe_sessions
|
||||
.remove(&body.session_id)
|
||||
.map(|(_, v)| v)
|
||||
.ok_or_else(|| {
|
||||
api_error(
|
||||
StatusCode::NOT_FOUND,
|
||||
"transcription session not found or already used",
|
||||
)
|
||||
})?;
|
||||
|
||||
// Reject expired sessions.
|
||||
if created_at.elapsed() > SESSION_SECRET_TTL {
|
||||
return Err(api_error(
|
||||
StatusCode::GONE,
|
||||
"transcription session expired — create a new one",
|
||||
));
|
||||
}
|
||||
|
||||
// Proxy the SDP offer to OpenAI.
|
||||
let client = openai_client().map_err(|(status, msg)| api_error(status, msg))?;
|
||||
let response = client
|
||||
// Step 2: Proxy the SDP offer to OpenAI using the minted secret.
|
||||
let sdp_response = client
|
||||
.post(OPENAI_REALTIME_CALLS_URL)
|
||||
.header("Authorization", format!("Bearer {client_secret}"))
|
||||
.header("Content-Type", "application/sdp")
|
||||
@@ -245,9 +190,9 @@ pub async fn proxy_sdp_exchange(
|
||||
)
|
||||
})?;
|
||||
|
||||
if !response.status().is_success() {
|
||||
let status = response.status();
|
||||
let resp_body = response.text().await.unwrap_or_default();
|
||||
if !sdp_response.status().is_success() {
|
||||
let status = sdp_response.status();
|
||||
let resp_body = sdp_response.text().await.unwrap_or_default();
|
||||
tracing::error!("OpenAI SDP exchange error ({status}): {resp_body}");
|
||||
return Err(api_error(
|
||||
StatusCode::BAD_GATEWAY,
|
||||
@@ -255,7 +200,7 @@ pub async fn proxy_sdp_exchange(
|
||||
));
|
||||
}
|
||||
|
||||
let sdp_answer = response.text().await.map_err(|e| {
|
||||
let sdp_answer = sdp_response.text().await.map_err(|e| {
|
||||
tracing::error!("OpenAI SDP answer read error: {e}");
|
||||
api_error(
|
||||
StatusCode::BAD_GATEWAY,
|
||||
@@ -263,7 +208,10 @@ pub async fn proxy_sdp_exchange(
|
||||
)
|
||||
})?;
|
||||
|
||||
Ok(Json(SdpExchangeResponse { sdp: sdp_answer }))
|
||||
Ok(Json(TranscribeConnectResponse {
|
||||
sdp: sdp_answer,
|
||||
model,
|
||||
}))
|
||||
}
|
||||
|
||||
// ── Helpers ───────────────────────────────────────────────────────────────────
|
||||
|
||||
@@ -86,10 +86,9 @@ pub fn build_router(state: Arc<AppState>) -> Router {
|
||||
get(api::transcribe::transcribe_status),
|
||||
)
|
||||
.route(
|
||||
"/transcribe/session",
|
||||
post(api::transcribe::create_transcribe_session),
|
||||
"/transcribe/connect",
|
||||
post(api::transcribe::transcribe_connect),
|
||||
)
|
||||
.route("/transcribe/sdp", post(api::transcribe::proxy_sdp_exchange))
|
||||
// Huddle audio WebSocket route
|
||||
.route(
|
||||
"/huddle/{channel_id}/audio",
|
||||
|
||||
@@ -532,7 +532,6 @@ impl AppState {
|
||||
.build(),
|
||||
),
|
||||
transcribe_rate_limiter: Arc::new(DashMap::new()),
|
||||
transcribe_sessions: Arc::new(DashMap::new()),
|
||||
media_uploads_in_flight: Arc::new(DashMap::new()),
|
||||
observer_owner_cache: Arc::new(
|
||||
moka::sync::Cache::builder()
|
||||
|
||||
@@ -5,13 +5,9 @@ export interface TranscribeStatus {
|
||||
model: string;
|
||||
}
|
||||
|
||||
export interface TranscribeSession {
|
||||
sessionId: string;
|
||||
model: string;
|
||||
}
|
||||
|
||||
export interface SdpExchangeResponse {
|
||||
export interface TranscribeConnectResponse {
|
||||
sdp: string;
|
||||
model: string;
|
||||
}
|
||||
|
||||
/** NIP-98 event kind for HTTP request authorization. */
|
||||
@@ -53,46 +49,28 @@ export async function getTranscribeStatus(): Promise<TranscribeStatus> {
|
||||
return response.json();
|
||||
}
|
||||
|
||||
export async function createTranscribeSession(): Promise<TranscribeSession> {
|
||||
const baseUrl = await getRelayHttpUrl();
|
||||
const url = `${baseUrl}/transcribe/session`;
|
||||
const response = await fetch(url, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
Authorization: await nip98AuthHeader(url, "POST"),
|
||||
},
|
||||
});
|
||||
if (!response.ok) {
|
||||
const body = await response.text().catch(() => "");
|
||||
throw new Error(
|
||||
`Failed to create transcribe session (${response.status}): ${body}`,
|
||||
);
|
||||
}
|
||||
return response.json();
|
||||
}
|
||||
|
||||
/**
|
||||
* Proxy the WebRTC SDP exchange through the relay. The relay holds the
|
||||
* OpenAI client secret server-side — the desktop client never sees it.
|
||||
* Mint an OpenAI Realtime session and complete the WebRTC SDP exchange in a
|
||||
* single relay round-trip. The relay holds the OpenAI bearer token server-side
|
||||
* — the client never sees it. This also works correctly across multiple relay
|
||||
* replicas since no server-side session state is needed between requests.
|
||||
*/
|
||||
export async function proxySdpExchange(
|
||||
sessionId: string,
|
||||
export async function transcribeConnect(
|
||||
sdp: string,
|
||||
): Promise<SdpExchangeResponse> {
|
||||
): Promise<TranscribeConnectResponse> {
|
||||
const baseUrl = await getRelayHttpUrl();
|
||||
const url = `${baseUrl}/transcribe/sdp`;
|
||||
const url = `${baseUrl}/transcribe/connect`;
|
||||
const response = await fetch(url, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
Authorization: await nip98AuthHeader(url, "POST"),
|
||||
},
|
||||
body: JSON.stringify({ sessionId, sdp }),
|
||||
body: JSON.stringify({ sdp }),
|
||||
});
|
||||
if (!response.ok) {
|
||||
const body = await response.text().catch(() => "");
|
||||
throw new Error(`SDP exchange failed (${response.status}): ${body}`);
|
||||
throw new Error(`Transcribe connect failed (${response.status}): ${body}`);
|
||||
}
|
||||
return response.json();
|
||||
}
|
||||
|
||||
@@ -1,9 +1,6 @@
|
||||
import { useCallback, useEffect, useRef, useState } from "react";
|
||||
import { toast } from "sonner";
|
||||
import {
|
||||
createTranscribeSession,
|
||||
getTranscribeStatus,
|
||||
} from "../api/transcribeSession";
|
||||
import { getTranscribeStatus } from "../api/transcribeSession";
|
||||
import {
|
||||
type AudioBufferCapture,
|
||||
type TranscriptEvent,
|
||||
@@ -199,10 +196,13 @@ export function useRealtimeDictation({
|
||||
const prevText = getTranscriptText(segmentStateRef.current);
|
||||
const merged = mergeTranscriptEvent(segmentStateRef.current, event);
|
||||
|
||||
if (merged === prevText) return;
|
||||
// Always update transcribing state — a completed event must clear the
|
||||
// indicator even when the final text matches the accumulated deltas.
|
||||
const stillTranscribing = event.type !== TRANSCRIPT_COMPLETED_EVENT;
|
||||
setIsTranscribing(stillTranscribing);
|
||||
|
||||
if (merged === prevText) return;
|
||||
onTranscriptTextRef.current(merged);
|
||||
setIsTranscribing(event.type !== TRANSCRIPT_COMPLETED_EVENT);
|
||||
},
|
||||
[],
|
||||
);
|
||||
@@ -247,15 +247,7 @@ export function useRealtimeDictation({
|
||||
}
|
||||
audioCaptureRef.current = audioCapture;
|
||||
|
||||
// 3. Create session via relay
|
||||
const session = await createTranscribeSession();
|
||||
if (isStaleRun()) {
|
||||
closeResources({ audioCapture, stream });
|
||||
return;
|
||||
}
|
||||
manualCommitRef.current = requiresManualCommit(session.model);
|
||||
|
||||
// 4. Set up WebRTC
|
||||
// 3. Set up WebRTC peer connection
|
||||
peerConnection = createPeerConnection();
|
||||
peerConnectionRef.current = peerConnection;
|
||||
const activeStream = stream;
|
||||
@@ -276,23 +268,14 @@ export function useRealtimeDictation({
|
||||
// Flush buffered audio once data channel opens
|
||||
const channelToFlush = dataChannel;
|
||||
const captureToFlush = audioCapture;
|
||||
const useManualCommit = manualCommitRef.current;
|
||||
dataChannel.addEventListener("open", () => {
|
||||
// If the user stopped (or restarted) recording between the SDP
|
||||
// exchange and the channel opening, drop this run's buffered audio.
|
||||
if (isStaleRun()) {
|
||||
captureToFlush.close();
|
||||
return;
|
||||
}
|
||||
flushAudioBuffer(channelToFlush, captureToFlush.chunks);
|
||||
// For manual-commit models, commit the initial buffered audio and
|
||||
// start a periodic commit interval so streaming transcripts flow
|
||||
// during recording (server VAD models commit automatically).
|
||||
if (useManualCommit) {
|
||||
if (manualCommitRef.current) {
|
||||
commitAudioBuffer(channelToFlush);
|
||||
// Commit every 2s to produce streaming transcript segments while
|
||||
// the user is still speaking. Each commit triggers a transcription
|
||||
// of the audio accumulated since the last commit.
|
||||
commitIntervalRef.current = setInterval(() => {
|
||||
if (channelToFlush.readyState === "open") {
|
||||
commitAudioBuffer(channelToFlush);
|
||||
@@ -303,16 +286,14 @@ export function useRealtimeDictation({
|
||||
audioCaptureRef.current = null;
|
||||
});
|
||||
|
||||
// 5. SDP exchange (proxied through the relay — client never sees the
|
||||
// OpenAI bearer token)
|
||||
await connectPeerConnection({
|
||||
peerConnection,
|
||||
sessionId: session.sessionId,
|
||||
});
|
||||
// 4. Session creation + SDP exchange in a single relay round-trip.
|
||||
// No cross-replica state needed — works in HA deployments.
|
||||
const { model } = await connectPeerConnection({ peerConnection });
|
||||
if (isStaleRun()) {
|
||||
closeResources({ audioCapture, dataChannel, peerConnection, stream });
|
||||
return;
|
||||
}
|
||||
manualCommitRef.current = requiresManualCommit(model);
|
||||
} catch (error) {
|
||||
closeResources({ audioCapture, dataChannel, peerConnection, stream });
|
||||
if (!isStaleRun()) {
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import { proxySdpExchange } from "../api/transcribeSession";
|
||||
import { transcribeConnect } from "../api/transcribeSession";
|
||||
import {
|
||||
REALTIME_BUFFER_PROCESSOR_NAME,
|
||||
createWorkletBlobUrl,
|
||||
@@ -28,28 +28,28 @@ export function createPeerConnection(): RTCPeerConnection {
|
||||
}
|
||||
|
||||
/**
|
||||
* Complete the WebRTC SDP exchange via the relay proxy.
|
||||
* Complete the WebRTC SDP exchange via the relay.
|
||||
*
|
||||
* The relay holds the OpenAI client secret server-side — the desktop client
|
||||
* sends its SDP offer to the relay, which forwards it to OpenAI and returns
|
||||
* the SDP answer. This prevents the client from ever seeing the bearer token.
|
||||
* The relay mints the OpenAI session and proxies the SDP exchange in a single
|
||||
* request — the client never sees the bearer token, and no server-side session
|
||||
* state is needed between requests (works across HA relay replicas).
|
||||
*
|
||||
* Returns the model name from the relay so the caller can configure VAD mode.
|
||||
*/
|
||||
export async function connectPeerConnection(options: {
|
||||
peerConnection: RTCPeerConnection;
|
||||
sessionId: string;
|
||||
}): Promise<void> {
|
||||
}): Promise<{ model: string }> {
|
||||
const offer = await options.peerConnection.createOffer();
|
||||
await options.peerConnection.setLocalDescription(offer);
|
||||
|
||||
const { sdp: answerSdp } = await proxySdpExchange(
|
||||
options.sessionId,
|
||||
offer.sdp ?? "",
|
||||
);
|
||||
const { sdp: answerSdp, model } = await transcribeConnect(offer.sdp ?? "");
|
||||
|
||||
await options.peerConnection.setRemoteDescription({
|
||||
type: "answer",
|
||||
sdp: answerSdp,
|
||||
});
|
||||
|
||||
return { model };
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
Reference in New Issue
Block a user