mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(dictation): rate-limit session minting, disable auto-submit by default
Addresses Wes's review feedback on PR #1511: 1. Per-pubkey rate limit on POST /transcribe/session (5/min, configurable via BUZZ_TRANSCRIBE_SESSIONS_PER_MINUTE). Each session opens a metered OpenAI Realtime connection on the operator's bill. 2. Auto-submit disabled by default (DEFAULT_AUTO_SUBMIT_PHRASE = ''). The infrastructure for configurable phrases remains in place and can be wired to a user setting later. Also: - Move BUZZ_TRANSCRIPTION_MODEL into Config (consistent with other knobs) - Use a shared reqwest::Client via OnceLock (connection pooling) - Stop dictation on manual send (stopDictationRef wiring)
This commit is contained in:
@@ -7,6 +7,7 @@
|
||||
//! Both endpoints require NIP-98 auth (same as `/events`, `/query`, `/count`).
|
||||
|
||||
use std::sync::Arc;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
use axum::{
|
||||
extract::State,
|
||||
@@ -16,13 +17,17 @@ use axum::{
|
||||
use serde::Serialize;
|
||||
use serde_json::Value;
|
||||
|
||||
use buzz_core::CommunityId;
|
||||
|
||||
use crate::state::AppState;
|
||||
|
||||
use super::api_error;
|
||||
|
||||
const OPENAI_REALTIME_CLIENT_SECRETS_URL: &str =
|
||||
"https://api.openai.com/v1/realtime/client_secrets";
|
||||
const DEFAULT_TRANSCRIPTION_MODEL: &str = "whisper-1";
|
||||
|
||||
/// Rate-limit window for transcription session minting.
|
||||
const TRANSCRIBE_RATE_WINDOW: Duration = Duration::from_secs(60);
|
||||
|
||||
/// Response for `GET /transcribe/status`.
|
||||
#[derive(Serialize)]
|
||||
@@ -47,11 +52,11 @@ pub async fn transcribe_status(
|
||||
State(state): State<Arc<AppState>>,
|
||||
headers: HeaderMap,
|
||||
) -> Result<Json<TranscribeStatus>, (StatusCode, Json<Value>)> {
|
||||
authenticate(&state, &headers, "/transcribe/status", "GET").await?;
|
||||
let (_pubkey, _community) = authenticate(&state, &headers, "/transcribe/status", "GET").await?;
|
||||
|
||||
Ok(Json(TranscribeStatus {
|
||||
configured: state.config.openai_api_key.is_some(),
|
||||
model: transcription_model(),
|
||||
model: state.config.transcription_model.clone(),
|
||||
}))
|
||||
}
|
||||
|
||||
@@ -64,7 +69,18 @@ pub async fn create_transcribe_session(
|
||||
State(state): State<Arc<AppState>>,
|
||||
headers: HeaderMap,
|
||||
) -> Result<Json<TranscribeSession>, (StatusCode, Json<Value>)> {
|
||||
authenticate(&state, &headers, "/transcribe/session", "POST").await?;
|
||||
let (pubkey, community) = authenticate(&state, &headers, "/transcribe/session", "POST").await?;
|
||||
|
||||
// Per-(community, pubkey) rate limit — each session mints a metered OpenAI
|
||||
// Realtime connection on the operator's bill.
|
||||
if transcribe_rate_limited(&state, community, &pubkey) {
|
||||
metrics::counter!("buzz_transcribe_session_rejections_total", "reason" => "rate_limit")
|
||||
.increment(1);
|
||||
return Err(api_error(
|
||||
StatusCode::TOO_MANY_REQUESTS,
|
||||
"transcription session rate limit exceeded — try again shortly",
|
||||
));
|
||||
}
|
||||
|
||||
let api_key = state.config.openai_api_key.as_deref().ok_or_else(|| {
|
||||
api_error(
|
||||
@@ -73,10 +89,9 @@ pub async fn create_transcribe_session(
|
||||
)
|
||||
})?;
|
||||
|
||||
let model = transcription_model();
|
||||
let model = state.config.transcription_model.clone();
|
||||
|
||||
let client = reqwest::Client::new();
|
||||
let response = client
|
||||
let response = openai_client()
|
||||
.post(OPENAI_REALTIME_CLIENT_SECRETS_URL)
|
||||
.header("Authorization", format!("Bearer {api_key}"))
|
||||
.header("Content-Type", "application/json")
|
||||
@@ -95,7 +110,6 @@ pub async fn create_transcribe_session(
|
||||
}
|
||||
}
|
||||
}))
|
||||
.timeout(std::time::Duration::from_secs(10))
|
||||
.send()
|
||||
.await
|
||||
.map_err(|e| {
|
||||
@@ -142,12 +156,13 @@ pub async fn create_transcribe_session(
|
||||
|
||||
/// Authenticate the request using the same NIP-98 / X-Pubkey pattern as the
|
||||
/// bridge endpoints, plus replay detection and relay membership enforcement.
|
||||
/// Returns the authenticated pubkey and resolved community on success.
|
||||
async fn authenticate(
|
||||
state: &AppState,
|
||||
headers: &HeaderMap,
|
||||
path: &str,
|
||||
method: &str,
|
||||
) -> Result<(), (StatusCode, Json<Value>)> {
|
||||
) -> Result<(nostr::PublicKey, CommunityId), (StatusCode, Json<Value>)> {
|
||||
let raw_host = headers
|
||||
.get("host")
|
||||
.and_then(|v| v.to_str().ok())
|
||||
@@ -182,14 +197,37 @@ async fn authenticate(
|
||||
)
|
||||
.await?;
|
||||
|
||||
Ok(())
|
||||
Ok((pubkey, tenant.community()))
|
||||
}
|
||||
|
||||
fn transcription_model() -> String {
|
||||
std::env::var("BUZZ_TRANSCRIPTION_MODEL")
|
||||
.ok()
|
||||
.filter(|s| !s.is_empty())
|
||||
.unwrap_or_else(|| DEFAULT_TRANSCRIPTION_MODEL.to_string())
|
||||
/// Per-(community, pubkey) sliding-window rate limiter for transcription session minting.
|
||||
fn transcribe_rate_limited(state: &AppState, community: CommunityId, pubkey: &nostr::PublicKey) -> bool {
|
||||
let key = (community, pubkey.to_bytes());
|
||||
let now = Instant::now();
|
||||
let limit = state.config.transcribe_sessions_per_minute;
|
||||
let mut entry = state.transcribe_rate_limiter.entry(key).or_insert((0, now));
|
||||
let (count, window_start) = entry.value_mut();
|
||||
if now.duration_since(*window_start) >= TRANSCRIBE_RATE_WINDOW {
|
||||
*count = 1;
|
||||
*window_start = now;
|
||||
return false;
|
||||
}
|
||||
if *count >= limit {
|
||||
return true;
|
||||
}
|
||||
*count += 1;
|
||||
false
|
||||
}
|
||||
|
||||
/// Shared HTTP client for OpenAI requests (connection pooling).
|
||||
fn openai_client() -> &'static reqwest::Client {
|
||||
static CLIENT: std::sync::OnceLock<reqwest::Client> = std::sync::OnceLock::new();
|
||||
CLIENT.get_or_init(|| {
|
||||
reqwest::Client::builder()
|
||||
.timeout(Duration::from_secs(10))
|
||||
.build()
|
||||
.expect("failed to build reqwest client")
|
||||
})
|
||||
}
|
||||
|
||||
fn extract_client_secret(value: &Value) -> Option<String> {
|
||||
|
||||
@@ -141,6 +141,12 @@ pub struct Config {
|
||||
pub media_max_concurrent_uploads_per_pubkey: u32,
|
||||
/// Maximum media upload starts accepted from one pubkey per minute.
|
||||
pub media_uploads_per_minute: u32,
|
||||
/// Maximum transcription sessions one pubkey can mint per minute.
|
||||
/// Each session opens a metered OpenAI Realtime connection on the operator's
|
||||
/// bill, so this bounds cost exposure from a single member.
|
||||
pub transcribe_sessions_per_minute: u32,
|
||||
/// Transcription model to use (default: whisper-1).
|
||||
pub transcription_model: String,
|
||||
|
||||
/// Optional override for ephemeral channel TTL (in seconds).
|
||||
/// When set, any channel created with a TTL tag will use this value instead
|
||||
@@ -443,6 +449,17 @@ impl Config {
|
||||
.filter(|&v| v > 0)
|
||||
.unwrap_or(30);
|
||||
|
||||
let transcribe_sessions_per_minute: u32 =
|
||||
std::env::var("BUZZ_TRANSCRIBE_SESSIONS_PER_MINUTE")
|
||||
.ok()
|
||||
.and_then(|v| v.parse().ok())
|
||||
.filter(|&v| v > 0)
|
||||
.unwrap_or(5);
|
||||
let transcription_model: String = std::env::var("BUZZ_TRANSCRIPTION_MODEL")
|
||||
.ok()
|
||||
.filter(|s| !s.is_empty())
|
||||
.unwrap_or_else(|| "whisper-1".to_string());
|
||||
|
||||
let ephemeral_ttl_override = std::env::var("BUZZ_EPHEMERAL_TTL_OVERRIDE")
|
||||
.ok()
|
||||
.and_then(|v| v.parse::<i32>().ok())
|
||||
@@ -541,6 +558,8 @@ impl Config {
|
||||
media_max_concurrent_uploads,
|
||||
media_max_concurrent_uploads_per_pubkey,
|
||||
media_uploads_per_minute,
|
||||
transcribe_sessions_per_minute,
|
||||
transcription_model,
|
||||
ephemeral_ttl_override,
|
||||
git_repo_path,
|
||||
git_max_pack_bytes,
|
||||
|
||||
@@ -374,6 +374,11 @@ pub struct AppState {
|
||||
/// generate fresh Nostr keys.
|
||||
pub invite_claim_rate_limiter:
|
||||
Arc<moka::sync::Cache<ScopedPubkeyKey, Arc<std::sync::atomic::AtomicU32>>>,
|
||||
/// Per-requester sliding-window rate limiter for transcription session
|
||||
/// minting. Key: (community_id, requester pubkey bytes). Value: (count,
|
||||
/// window_start). Bounds the cost of OpenAI Realtime sessions minted on the
|
||||
/// operator's bill.
|
||||
pub transcribe_rate_limiter: Arc<ScopedRateLimiter>,
|
||||
/// Current in-flight media uploads per (community, uploader pubkey).
|
||||
pub media_uploads_in_flight: Arc<DashMap<ScopedPubkeyKey, u32>>,
|
||||
/// Cache for observer agent-owner authorization (kind 24200).
|
||||
@@ -521,6 +526,7 @@ impl AppState {
|
||||
.time_to_live(crate::api::invites::CLAIM_RATE_WINDOW)
|
||||
.build(),
|
||||
),
|
||||
transcribe_rate_limiter: Arc::new(DashMap::new()),
|
||||
media_uploads_in_flight: Arc::new(DashMap::new()),
|
||||
observer_owner_cache: Arc::new(
|
||||
moka::sync::Cache::builder()
|
||||
|
||||
@@ -74,9 +74,17 @@ test("replaceTrailingTranscribedText_noDoubleSpaceBeforePunctuation", () => {
|
||||
// ── getAutoSubmitMatch ──────────────────────────────────────────────────────
|
||||
|
||||
test("getAutoSubmitMatch_returnsNullWhenPhraseAbsent", () => {
|
||||
assert.equal(
|
||||
getAutoSubmitMatch("hello there", parseAutoSubmitPhrases("submit")),
|
||||
null,
|
||||
);
|
||||
});
|
||||
|
||||
test("getAutoSubmitMatch_returnsNullWhenPhrasesEmpty", () => {
|
||||
// DEFAULT_AUTO_SUBMIT_PHRASE is empty (auto-submit disabled by default).
|
||||
assert.equal(
|
||||
getAutoSubmitMatch(
|
||||
"hello there",
|
||||
"send this message submit",
|
||||
parseAutoSubmitPhrases(DEFAULT_AUTO_SUBMIT_PHRASE),
|
||||
),
|
||||
null,
|
||||
@@ -86,7 +94,7 @@ test("getAutoSubmitMatch_returnsNullWhenPhraseAbsent", () => {
|
||||
test("getAutoSubmitMatch_matchesTrailingPhraseAndStripsIt", () => {
|
||||
const match = getAutoSubmitMatch(
|
||||
"send this message submit",
|
||||
parseAutoSubmitPhrases(DEFAULT_AUTO_SUBMIT_PHRASE),
|
||||
parseAutoSubmitPhrases("submit"),
|
||||
);
|
||||
assert.ok(match);
|
||||
assert.equal(match.matchedPhrase, "submit");
|
||||
@@ -98,7 +106,7 @@ test("getAutoSubmitMatch_ignoresPhraseMidSentence", () => {
|
||||
assert.equal(
|
||||
getAutoSubmitMatch(
|
||||
"submit the form later",
|
||||
parseAutoSubmitPhrases(DEFAULT_AUTO_SUBMIT_PHRASE),
|
||||
parseAutoSubmitPhrases("submit"),
|
||||
),
|
||||
null,
|
||||
);
|
||||
@@ -115,7 +123,7 @@ test("getAutoSubmitMatch_requiresWordBoundaryBeforePhrase", () => {
|
||||
test("getAutoSubmitMatch_toleratesTrailingPunctuation", () => {
|
||||
const match = getAutoSubmitMatch(
|
||||
"ship it submit.",
|
||||
parseAutoSubmitPhrases(DEFAULT_AUTO_SUBMIT_PHRASE),
|
||||
parseAutoSubmitPhrases("submit"),
|
||||
);
|
||||
assert.ok(match);
|
||||
assert.equal(match.textWithoutPhrase, "ship it");
|
||||
|
||||
@@ -1,4 +1,10 @@
|
||||
export const DEFAULT_AUTO_SUBMIT_PHRASE = "submit";
|
||||
/**
|
||||
* Default auto-submit phrase. Empty string disables auto-submit — the user
|
||||
* must manually press Enter/Send after dictation. The infrastructure for
|
||||
* configurable phrases is in place (parseAutoSubmitPhrases, getAutoSubmitMatch)
|
||||
* and can be wired to a user setting when we're ready to ship auto-submit.
|
||||
*/
|
||||
export const DEFAULT_AUTO_SUBMIT_PHRASE = "";
|
||||
|
||||
const TRAILING_PUNCTUATION_REGEX = /[\s"'`.,!?;:)\]}]+$/u;
|
||||
|
||||
|
||||
@@ -267,6 +267,7 @@ function MessageComposerImpl({
|
||||
|
||||
const submitMessageRef = React.useRef<() => void>(() => {});
|
||||
const setEditorContentRef = React.useRef<(text: string) => void>(() => {});
|
||||
const stopDictationRef = React.useRef<() => void>(() => {});
|
||||
const dictation = useComposerDictation({
|
||||
syncContentRef: syncContentRefFromEditorRef,
|
||||
disabledRef,
|
||||
@@ -277,6 +278,7 @@ function MessageComposerImpl({
|
||||
submitMessageRef,
|
||||
draftKey: effectiveDraftKey,
|
||||
});
|
||||
stopDictationRef.current = dictation.stopRecording;
|
||||
const composerScrollRef = React.useRef<HTMLDivElement>(null);
|
||||
// Set after `useLinkEditor` exists below; the editor's link-click handler
|
||||
// delegates through this ref to break the hook ordering cycle (the editor
|
||||
@@ -572,11 +574,8 @@ function MessageComposerImpl({
|
||||
spoileredAttachmentUrls,
|
||||
);
|
||||
|
||||
// NIP-30: attach `["emoji", shortcode, url]` tags for custom emoji in the
|
||||
// edited body, exactly like the send path. Without this an edited message
|
||||
// ships with no emoji tags, so the receiver can't resolve a `:shortcode:`
|
||||
// and renders the literal text. `?? []` preserves edit semantics (a
|
||||
// defined-but-empty media set means "wipe attachments").
|
||||
// NIP-30 emoji tags for the edited body (mirrors the send path).
|
||||
// `?? []` preserves edit semantics (empty = "wipe attachments").
|
||||
const outgoingTags =
|
||||
mergeOutgoingTags(
|
||||
mediaTags,
|
||||
@@ -619,6 +618,8 @@ function MessageComposerImpl({
|
||||
return;
|
||||
}
|
||||
|
||||
stopDictationRef.current(); // stop dictation so late transcripts don't refill
|
||||
|
||||
const capturedThreadContext = onCaptureSendContext?.() ?? null;
|
||||
// If a thread-reply composer reported no reply target at submit time,
|
||||
// bail here rather than discovering the null later after async awaits.
|
||||
@@ -711,10 +712,8 @@ function MessageComposerImpl({
|
||||
);
|
||||
|
||||
// ── Keyboard handling ───────────────────────────────────────────────
|
||||
// Tiptap handles formatting shortcuts (⌘B, ⌘I, etc.) natively.
|
||||
// Plain Enter → submit is now handled inside the Tiptap `submitOnEnter`
|
||||
// extension (fires before ProseMirror's splitBlock). This wrapper only
|
||||
// handles autocomplete arrow/enter keys and Escape for edit mode.
|
||||
// Autocomplete arrow/enter keys and Escape for edit mode only; plain
|
||||
// Enter → submit is in the Tiptap `submitOnEnter` extension.
|
||||
const handleEditorKeyDown = React.useCallback(
|
||||
(event: React.KeyboardEvent<HTMLDivElement>) => {
|
||||
// Let autocomplete handle keys first
|
||||
|
||||
Reference in New Issue
Block a user