Skip to content
Merged
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
39 changes: 27 additions & 12 deletions docs/serving/speech_api.md
Original file line number Diff line number Diff line change
Expand Up @@ -279,10 +279,15 @@ curl -X POST http://localhost:8091/v1/audio/voices \

## Streaming Text Input (WebSocket)

The `/v1/audio/speech/stream` WebSocket endpoint accepts text incrementally and generates audio per sentence as boundaries are detected.
The `/v1/audio/speech/stream` WebSocket endpoint accepts text incrementally.
By default (`split_granularity=none`) it buffers until `input.done` and
synthesizes that flush as **one** TTS request, which keeps long-form timbre
stable. Set `split_granularity` to `sentence` or `clause` to emit a request
at each detected boundary (lower time-to-first-audio for STT/LLM pipelines).

> Note: text input is always streamed incrementally. Audio output remains sentence-scoped:
> use `stream_audio=false` for one binary frame per sentence, or `stream_audio=true` for one or more PCM chunks per sentence.
> Note: `stream_audio` only changes how **audio bytes** are framed (one WAV/PCM
> payload vs chunked PCM). Text segmentation is controlled separately by
> `split_granularity`.

### WebSocket Protocol

Expand Down Expand Up @@ -314,14 +319,16 @@ upstream LLM) pays the WebSocket handshake once instead of once per utterance.

- The session config is sticky. Send `input.text` again straight after
`session.done` to reuse it, or send another `session.config` first to change
voice, format, or reference audio. A `session.config` sent while text is
still buffered is rejected so no pending input is silently dropped.
- An utterance is the flush unit, not a linguistic one: it is whatever text was
buffered when `input.done` arrived, of any length, synthesized as one request.
`utterance_index` counts those flushes across the connection, so it tells you
which `input.done` a frame belongs to. `sentence_index` counts within one
flush and so pairs with `total_sentences`, which means every utterance reports
`sentence_index: 0` of `total_sentences: 1` (or `0` for an empty buffer).
voice, format, or reference audio. A `session.config` sent in the middle of
an utterance is rejected so no pending input is silently dropped and so a
split utterance cannot end up half in one voice and half in another.
- An utterance is the flush unit, not a linguistic one: it is one `input.done`
cycle. `utterance_index` counts those flushes across the connection, so it
tells you which `input.done` a frame belongs to. `sentence_index` counts the
TTS requests inside one flush and so pairs with `total_sentences`: with the
default `split_granularity=none` that is always `sentence_index: 0` of
`total_sentences: 1` (or `0` for an empty buffer), while `sentence` or
`clause` counts the linguistic units actually synthesized.
- End the connection with `session.close`, or by closing the socket. An idle
connection is still closed after the server's idle timeout, which now also
applies to the gap between utterances.
Expand All @@ -332,7 +339,15 @@ All REST API parameters are supported, plus:

| Parameter | Type | Default | Description |
| ----------- | ------ | --------- | ------------- |
| `stream_audio` | bool | false | Stream one or more PCM chunks for the buffered input over WebSocket |
| `stream_audio` | bool | false | Stream one or more PCM chunks for each TTS request over WebSocket |
| `split_granularity` | string | `"none"` | `"none"`: one request per `input.done`. `"sentence"`: split on `.!?` plus CJK `。!?…`, Indic danda `।॥`, and Arabic `؟`. `"clause"`: also split on `,;,;،؛`. |
| `seed` | integer | null | Forwarded to the speech engine for this session |

ASCII punctuation only ends a unit when whitespace or `input.done` follows it,
and decimals (`3.14`), thousands separators (`1,000`), abbreviations (`Dr.`,
`e.g.`) and initials (`J. R.`) are not treated as boundaries. A punctuation run
and any closing quote or bracket stay with the unit they close, so `Wait...`
and `He said "Hello."` are one request each.

```bash
DELETE /v1/audio/voices/{name}
Expand Down
10 changes: 9 additions & 1 deletion examples/online_serving/text_to_speech/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -629,13 +629,21 @@ Raw PCM streaming requires `stream_format="audio"`, `response_format="pcm"`, and

### Streaming WebSocket

The `/v1/audio/speech/stream` endpoint accepts text incrementally and synthesizes the buffered text as one continuous request on `input.done`:
The `/v1/audio/speech/stream` endpoint accepts text incrementally and, by default, synthesizes the buffered text as one continuous request on `input.done`:

```bash
python qwen3_tts/streaming_speech_client.py --text "Hello world. How are you? I am fine."
python qwen3_tts/streaming_speech_client.py --text "..." --simulate-stt --stt-delay 0.1
```

For per-sentence audio (including Indic danda `।`) before `input.done`:

```bash
python qwen3_tts/streaming_speech_client.py \
--text "नमस्ते। कैसे हो?" \
--split-granularity sentence
```

`input.done` flushes without closing, so repeating `--text` synthesizes several utterances over one connection:

```bash
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,11 @@
--text "Hello world. How are you? I am fine." \
--simulate-stt --stt-delay 0.1

# Opt into per-sentence synthesis (Indic danda, CJK, Latin)
python streaming_speech_client.py \
--text "नमस्ते। कैसे हो?" \
--split-granularity sentence

# Receive JSON sidecar chunks with word-level timestamps
python streaming_speech_client.py \
--text "Hello world. How are you?" \
Expand Down Expand Up @@ -73,8 +78,9 @@
def frame_basename(msg: dict) -> str:
"""Name a file after the utterance and sentence a frame belongs to.

Every utterance is sentence 0, so the utterance index is what keeps the
files of one connection from overwriting each other.
Every utterance uses `sentence_index` for units inside that flush.
With default `split_granularity=none` that index stays 0; sentence mode
uses 0, 1, ... so files from one connection do not overwrite each other.
"""
return f"utterance_{msg['utterance_index']:03d}_sentence_{msg['sentence_index']:03d}"

Expand Down Expand Up @@ -368,6 +374,13 @@ def main():
)
parser.add_argument("--speed", type=float, default=1.0, help="Playback speed (0.25-4.0)")
parser.add_argument("--max-new-tokens", type=int, default=None, help="Max tokens")
parser.add_argument("--seed", type=int, default=None, help="Sampling seed forwarded to the server")
parser.add_argument(
"--split-granularity",
default=None,
choices=["none", "sentence", "clause"],
help="Text split mode (default: server none = one request per input.done)",
)

# Base task options
parser.add_argument("--ref-audio", default=None, help="Reference audio")
Expand Down Expand Up @@ -409,6 +422,8 @@ def main():
"max_new_tokens",
"ref_audio",
"ref_text",
"seed",
"split_granularity",
]:
val = getattr(args, key.replace("-", "_"), None)
if val is not None:
Expand Down
161 changes: 160 additions & 1 deletion tests/entrypoints/openai_api/test_serving_speech_stream.py
Original file line number Diff line number Diff line change
Expand Up @@ -197,7 +197,7 @@ def test_session_config_rejected_while_input_is_buffered(self, mocker: MockerFix

error = ws.receive_json()
assert error["type"] == "error"
assert "while input is buffered" in error["message"]
assert "while an utterance is in progress" in error["message"]

# The buffered text survives the rejected reconfiguration.
ws.send_json({"type": "input.done"})
Expand Down Expand Up @@ -742,6 +742,8 @@ async def mock_generate_pcm_chunks(_generator, _request_id, *, include_sample_ra
config.speaker_embedding = None
config.stream_audio = True
config.word_timestamps = False
config.seed = None
config.non_streaming_mode = None

with pytest.raises(WebSocketDisconnect):
asyncio.run(
Expand All @@ -758,6 +760,163 @@ async def mock_generate_pcm_chunks(_generator, _request_id, *, include_sample_ra
assert websocket.send_json.await_count == 2


class TestWebSocketSentenceSplitting:
def test_sentence_granularity_emits_one_request_per_sentence(self, mocker: MockerFixture):
app, speech_service = _build_test_app(mocker=mocker)

with TestClient(app) as client:
with client.websocket_connect("/v1/audio/speech/stream") as ws:
ws.send_json(
{
"type": "session.config",
"voice": "Vivian",
"split_granularity": "sentence",
}
)
ws.send_json({"type": "input.text", "text": "Hello world. How are you? "})
ws.send_json({"type": "input.done"})

first = ws.receive_json()
assert first["sentence_index"] == 0
assert first["sentence_text"] == "Hello world."
ws.receive_bytes()
assert ws.receive_json()["type"] == "audio.done"

second = ws.receive_json()
assert second["sentence_index"] == 1
assert second["sentence_text"] == "How are you?"
ws.receive_bytes()
assert ws.receive_json()["type"] == "audio.done"

assert ws.receive_json() == {
"type": "session.done",
"utterance_index": 0,
"total_sentences": 2,
}

assert speech_service._generate_audio_bytes.await_count == 2
assert [call.args[0].input for call in speech_service._generate_audio_bytes.await_args_list] == [
"Hello world.",
"How are you?",
]

def test_sentence_granularity_emits_before_input_done(self, mocker: MockerFixture):
app, speech_service = _build_test_app(mocker=mocker)

with TestClient(app) as client:
with client.websocket_connect("/v1/audio/speech/stream") as ws:
ws.send_json(
{
"type": "session.config",
"voice": "Vivian",
"split_granularity": "sentence",
}
)
ws.send_json({"type": "input.text", "text": "First sentence. "})

start = ws.receive_json()
assert start["sentence_text"] == "First sentence."
ws.receive_bytes()
assert ws.receive_json()["type"] == "audio.done"
assert speech_service._generate_audio_bytes.await_count == 1

ws.send_json({"type": "input.done"})
assert ws.receive_json() == {
"type": "session.done",
"utterance_index": 0,
"total_sentences": 1,
}

def test_indic_danda_splits_without_latin_period(self, mocker: MockerFixture):
app, speech_service = _build_test_app(mocker=mocker)

with TestClient(app) as client:
with client.websocket_connect("/v1/audio/speech/stream") as ws:
ws.send_json(
{
"type": "session.config",
"voice": "Vivian",
"language": "Auto",
"split_granularity": "sentence",
}
)
ws.send_json({"type": "input.text", "text": "नमस्ते। कैसे हो?"})
ws.send_json({"type": "input.done"})

first = ws.receive_json()
assert first["sentence_text"] == "नमस्ते।"
ws.receive_bytes()
assert ws.receive_json()["type"] == "audio.done"

second = ws.receive_json()
assert second["sentence_text"] == "कैसे हो?"
ws.receive_bytes()
assert ws.receive_json()["type"] == "audio.done"
assert ws.receive_json()["total_sentences"] == 2

def test_session_config_rejected_after_a_split_unit_was_emitted(self, mocker: MockerFixture):
"""The splitter buffer is empty here, but the utterance is still open."""
app, speech_service = _build_test_app(mocker=mocker)

with TestClient(app) as client:
with client.websocket_connect("/v1/audio/speech/stream") as ws:
ws.send_json(
{
"type": "session.config",
"voice": "Vivian",
"split_granularity": "sentence",
}
)
ws.send_json({"type": "input.text", "text": "First sentence. "})
ws.receive_json()
ws.receive_bytes()
assert ws.receive_json()["type"] == "audio.done"

ws.send_json({"type": "session.config", "voice": "Serena"})
error = ws.receive_json()
assert error["type"] == "error"
assert "while an utterance is in progress" in error["message"]

ws.send_json({"type": "input.text", "text": "Second sentence. "})
assert ws.receive_json()["sentence_index"] == 1
ws.receive_bytes()
assert ws.receive_json()["type"] == "audio.done"

ws.send_json({"type": "input.done"})
assert ws.receive_json() == {
"type": "session.done",
"utterance_index": 0,
"total_sentences": 2,
}

# Reconfiguration is allowed again once the utterance closed.
ws.send_json({"type": "session.config", "voice": "Serena"})
ws.send_json({"type": "input.text", "text": "Third."})
ws.send_json({"type": "input.done"})
ws.receive_json()
ws.receive_bytes()
ws.receive_json()
ws.receive_json()

voices = [call.args[0].voice for call in speech_service._generate_audio_bytes.await_args_list]
assert voices == ["Vivian", "Vivian", "Serena"]

def test_seed_is_forwarded_to_speech_request(self, mocker: MockerFixture):
app, speech_service = _build_test_app(mocker=mocker)

with TestClient(app) as client:
with client.websocket_connect("/v1/audio/speech/stream") as ws:
ws.send_json({"type": "session.config", "voice": "Vivian", "seed": 42})
ws.send_json({"type": "input.text", "text": "Hello."})
ws.send_json({"type": "input.done"})
ws.receive_json()
ws.receive_bytes()
ws.receive_json()
ws.receive_json()

assert speech_service._generate_audio_bytes.await_args_list[0].args[0].seed == 42


class TestGeneratePcmChunksContract:
"""Guard: _generate_pcm_chunks must exist on OmniOpenAIServingSpeech.

Expand Down
Loading
Loading