Skip to content
Draft
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
371 changes: 351 additions & 20 deletions bridge/lan_service.py

Large diffs are not rendered by default.

172 changes: 172 additions & 0 deletions bridge/test_lan_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -2200,6 +2200,178 @@ def test_binary_audio_upload_tracks_telemetry_and_requires_stt_or_transcript(sel
self.assertFalse(session.audio.active)
self.assertEqual(0, session.audio.bytes_received)

def test_audio_capture_absolute_lease_exceeds_firmware_max_but_cannot_be_refreshed(self):
session = LanBridgeSession(LanBridgeConfig())
self.assertGreater(session.config.audio_capture_absolute_lease_ms, 12_000)
self.assertLessEqual(session.config.audio_capture_absolute_lease_ms, 15_000)
session.handle_text(
json.dumps({"type": "utterance_start", "seq": 41, "sample_rate": 16000})
)
session.handle_binary(b"\x01\x00\x02\x00")
current = time.perf_counter()
session.audio.started_at_monotonic = current - (
session.config.audio_capture_absolute_lease_ms + 1
) / 1000.0
session.audio.last_activity_at_monotonic = current

expired = session.handle_binary(b"\x03\x00\x04\x00")

self.assertEqual("audio_capture_absolute_lease_expired", expired[0]["code"])
self.assertEqual(4, expired[0]["audio_bytes"])
self.assertFalse(session.audio.active)
self.assertEqual(0, session.audio.bytes_received)
self.assertEqual(1, expired[0]["audio_uploads_expired"])

def test_audio_capture_valid_after_twelve_seconds_with_recent_pcm(self):
session = LanBridgeSession(LanBridgeConfig())
session.handle_text(
json.dumps({"type": "utterance_start", "seq": 42, "sample_rate": 16000})
)
session.handle_binary(b"\x01\x00\x02\x00")
current = time.perf_counter()
session.audio.started_at_monotonic = current - 12.25
session.audio.last_activity_at_monotonic = current
runner_result = SimpleNamespace(
raw_response=json.dumps(
{
"spoken_text": "I heard the complete sentence.",
"mode": "speak",
"earcon": "none",
"emotion": {"arousal": 0.0, "valence": 0.0},
"memory_write": {},
"memory_forget": [],
}
),
command_source="test",
elapsed_ms=1.0,
approx_tokens_per_sec=10.0,
)

with patch("lan_service.run_runner_profile", return_value=runner_result) as runner:
frames = session.handle_text(
json.dumps(
{"type": "utterance_end", "seq": 42, "text": "complete sentence"}
)
)

runner.assert_called_once()
thinking = next(frame for frame in frames if frame.get("type") == "thinking")
self.assertGreaterEqual(thinking["audio_capture_elapsed_ms"], 12_000)
self.assertNotIn("audio_capture_reject_code", thinking)
self.assertFalse(session.audio.active)

def test_audio_capture_inactivity_discards_partial_pcm_and_rejects_late_end(self):
session = LanBridgeSession(LanBridgeConfig())
session.handle_text(
json.dumps({"type": "utterance_start", "seq": 43, "sample_rate": 16000})
)
session.handle_binary(b"\x01\x00\x02\x00")
current = time.perf_counter()
session.audio.last_activity_at_monotonic = current - (
session.config.audio_capture_inactivity_lease_ms + 1
) / 1000.0

with (
patch("lan_service.transcribe_pcm") as stt,
patch("lan_service.run_runner_profile") as runner,
patch("lan_service.synthesize_speech") as tts,
):
expired = session.handle_binary(b"\x03\x00\x04\x00")
late_end = session.handle_text(
json.dumps({"type": "utterance_end", "seq": 43, "text": "stale"})
)

self.assertEqual("audio_capture_inactivity_expired", expired[0]["code"])
self.assertEqual("utterance_end_stale", late_end[0]["code"])
self.assertEqual(
"audio_capture_inactivity_expired",
late_end[0]["audio_capture_reject_code"],
)
stt.assert_not_called()
runner.assert_not_called()
tts.assert_not_called()

def test_heartbeat_expires_inactive_capture_without_dropping_heartbeat(self):
session = LanBridgeSession(LanBridgeConfig())
session.handle_text(json.dumps({"type": "utterance_start", "seq": 48}))
session.handle_binary(b"\x01\x00\x02\x00")
current = time.perf_counter()
session.audio.last_activity_at_monotonic = current - (
session.config.audio_capture_inactivity_lease_ms + 1
) / 1000.0

frames = session.handle_text(
json.dumps({"type": "heartbeat", "power_source": "wall"})
)

self.assertEqual("audio_capture_inactivity_expired", frames[0]["code"])
self.assertEqual("heartbeat", frames[1]["type"])
self.assertFalse(session.audio.active)

def test_utterance_cancel_is_idempotent_and_late_end_runs_no_pipeline(self):
session = LanBridgeSession(LanBridgeConfig())
session.handle_text(
json.dumps({"type": "utterance_start", "seq": 44, "sample_rate": 16000})
)
session.handle_binary(b"\x01\x00\x02\x00")

with (
patch("lan_service.transcribe_pcm") as stt,
patch("lan_service.run_runner_profile") as runner,
patch("lan_service.synthesize_speech") as tts,
):
cancelled = session.handle_text(
json.dumps(
{"type": "utterance_cancel", "seq": 44, "reason": "capture_discontinuity"}
)
)
duplicate = session.handle_text(
json.dumps(
{"type": "utterance_cancel", "seq": 44, "reason": "capture_discontinuity"}
)
)
late_end = session.handle_text(
json.dumps({"type": "utterance_end", "seq": 44, "text": "unfinished"})
)

self.assertTrue(cancelled[0]["utterance_cancelled"])
self.assertFalse(cancelled[0]["utterance_cancel_duplicate"])
self.assertEqual(4, cancelled[0]["audio_bytes"])
self.assertTrue(duplicate[0]["utterance_cancel_duplicate"])
self.assertEqual("utterance_end_stale", late_end[0]["code"])
self.assertEqual("audio_capture_cancelled", late_end[0]["audio_capture_reject_code"])
self.assertFalse(session.audio.active)
self.assertEqual(0, session.audio.bytes_received)
stt.assert_not_called()
runner.assert_not_called()
tts.assert_not_called()

def test_wrong_terminal_sequence_does_not_finalize_active_capture(self):
session = LanBridgeSession(LanBridgeConfig())
session.handle_text(
json.dumps({"type": "utterance_start", "seq": 45, "sample_rate": 16000})
)
session.handle_binary(b"\x01\x00\x02\x00")

mismatch = session.handle_text(
json.dumps({"type": "utterance_end", "seq": 46, "text": "wrong turn"})
)

self.assertEqual("utterance_end_seq_mismatch", mismatch[0]["code"])
self.assertTrue(session.audio.active)
self.assertEqual(4, session.audio.bytes_received)

def test_audio_lease_telemetry_is_exposed_in_diagnostics(self):
session = LanBridgeSession(LanBridgeConfig())
session.handle_text(json.dumps({"type": "utterance_start", "seq": 47}))

diagnostics = session.handle_text(json.dumps({"type": "diagnostics_request"}))[0]

self.assertTrue(diagnostics["audio"]["audio_upload_active"])
self.assertEqual(1, diagnostics["audio"]["audio_uploads_started"])
self.assertEqual(14_500, diagnostics["audio"]["audio_capture_absolute_lease_ms"])
self.assertEqual(4_000, diagnostics["audio"]["audio_capture_inactivity_lease_ms"])

def test_empty_utterance_end_does_not_run_runner(self):
session = LanBridgeSession(LanBridgeConfig(runner_case="greeting"))

Expand Down
143 changes: 125 additions & 18 deletions src/io/BridgeAudioUplink.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -23,13 +23,15 @@ bool BridgeAudioUplink::begin(const BridgeAudioUplinkConfig& config,
return fail("audio_uplink_session_missing");
}
if (config_.sampleRate == 0 || config_.maxAudioBytes == 0 ||
config_.maxChunkBytes == 0 ||
config_.maxChunkBytes == 0 || config_.terminalRetryMs == 0 ||
config_.maxChunkBytes > kBridgeAudioStreamChunkPayloadMax) {
telemetry_.ready = false;
return fail("audio_uplink_bad_config");
}

telemetry_.lastError[0] = '\0';
pendingTerminalSeq_ = 0;
pendingTerminalReason_[0] = '\0';
return true;
}

Expand All @@ -39,6 +41,8 @@ void BridgeAudioUplink::reset() {
telemetry_.ready = ready;
telemetry_.enabled = config_.enabled;
telemetry_.wakeGateRequired = config_.wakeGateRequired;
pendingTerminalSeq_ = 0;
pendingTerminalReason_[0] = '\0';
if (!config_.enabled) {
copyError("audio_uplink_disabled");
}
Expand All @@ -55,7 +59,7 @@ bool BridgeAudioUplink::beginTurn(uint32_t seq, uint32_t nowMs, bool wakeGateOpe
if (session_ == nullptr) {
return fail("audio_uplink_session_missing");
}
if (telemetry_.active) {
if (telemetry_.active || telemetry_.terminalPending) {
return fail("audio_uplink_already_active");
}
if (config_.wakeGateRequired && !wakeGateOpen) {
Expand All @@ -74,6 +78,7 @@ bool BridgeAudioUplink::beginTurn(uint32_t seq, uint32_t nowMs, bool wakeGateOpe
telemetry_.lastSeq = seq;
telemetry_.activeBytes = 0;
telemetry_.activeChunks = 0;
telemetry_.lastTerminal = BridgeAudioTerminalKind::None;
telemetry_.lastError[0] = '\0';
return true;
}
Expand Down Expand Up @@ -131,7 +136,6 @@ bool BridgeAudioUplink::submitPcmBytes(uint32_t seq,
}

bool BridgeAudioUplink::endTurn(uint32_t seq, uint32_t nowMs) {
(void)nowMs;
if (!configured()) {
return fail("audio_uplink_not_ready");
}
Expand All @@ -145,26 +149,67 @@ bool BridgeAudioUplink::endTurn(uint32_t seq, uint32_t nowMs) {
return fail("audio_uplink_seq_mismatch");
}

return requestTerminal(BridgeAudioTerminalKind::End, seq, nowMs, nullptr);
}

void BridgeAudioUplink::abort(uint32_t nowMs, const char* reason) {
if (!telemetry_.active && !telemetry_.terminalPending) {
copyError(reason != nullptr ? reason : "audio_uplink_aborted");
return;
}
requestTerminal(BridgeAudioTerminalKind::Cancel,
telemetry_.lastSeq,
nowMs,
reason != nullptr ? reason : "audio_uplink_aborted");
}

BridgeAudioTerminalServiceResult BridgeAudioUplink::servicePendingTerminal(uint32_t nowMs) {
if (!telemetry_.terminalPending) {
return BridgeAudioTerminalServiceResult::Idle;
}

char frame[kBridgeEndpointControlResponseMax] = {};
if (!writeEndFrame(seq, frame, sizeof(frame)) || !queueText(frame)) {
telemetry_.queueFailures++;
return fail("utterance_end_queue_failed");
const bool encoded = telemetry_.pendingTerminal == BridgeAudioTerminalKind::End
? writeEndFrame(pendingTerminalSeq_, frame, sizeof(frame))
: writeCancelFrame(pendingTerminalSeq_,
pendingTerminalReason_,
frame,
sizeof(frame));
telemetry_.terminalAttempts++;
if (encoded && queueText(frame)) {
const BridgeAudioTerminalKind delivered = telemetry_.pendingTerminal;
telemetry_.terminalPending = false;
telemetry_.pendingTerminal = BridgeAudioTerminalKind::None;
telemetry_.lastTerminal = delivered;
telemetry_.lastSeq = pendingTerminalSeq_;
if (delivered == BridgeAudioTerminalKind::End) {
telemetry_.turnsCompleted++;
telemetry_.lastError[0] = '\0';
} else {
telemetry_.turnsAborted++;
telemetry_.cancelFramesQueued++;
copyError(pendingTerminalReason_);
}
return BridgeAudioTerminalServiceResult::Queued;
}

telemetry_.active = false;
telemetry_.turnsCompleted++;
telemetry_.lastSeq = seq;
telemetry_.lastError[0] = '\0';
return true;
}
telemetry_.queueFailures++;
telemetry_.terminalRetries++;
if (nowMs - telemetry_.terminalRequestedAtMs < config_.terminalRetryMs) {
copyError("utterance_terminal_queue_pending");
return BridgeAudioTerminalServiceResult::Pending;
}

void BridgeAudioUplink::abort(uint32_t nowMs, const char* reason) {
(void)nowMs;
if (telemetry_.active) {
telemetry_.turnsAborted++;
if (session_ != nullptr) {
session_->stop(nowMs);
}
telemetry_.active = false;
copyError(reason != nullptr ? reason : "audio_uplink_aborted");
telemetry_.terminalTimeouts++;
telemetry_.turnsAborted++;
telemetry_.terminalPending = false;
telemetry_.pendingTerminal = BridgeAudioTerminalKind::None;
telemetry_.lastTerminal = BridgeAudioTerminalKind::Cancel;
fail("utterance_terminal_delivery_timeout");
return BridgeAudioTerminalServiceResult::FailedClosed;
}

bool BridgeAudioUplink::configured() const {
Expand Down Expand Up @@ -206,6 +251,68 @@ bool BridgeAudioUplink::writeEndFrame(uint32_t seq, char* out, size_t outSize) c
return written > 0 && static_cast<size_t>(written) < outSize;
}

bool BridgeAudioUplink::writeCancelFrame(uint32_t seq,
const char* reason,
char* out,
size_t outSize) const {
if (out == nullptr || outSize == 0 || reason == nullptr || reason[0] == '\0') {
return false;
}
const int written = snprintf(out,
outSize,
"{\"type\":\"utterance_cancel\",\"seq\":%lu,"
"\"reason\":\"%s\",\"audio_bytes\":%lu,\"chunks\":%lu}",
static_cast<unsigned long>(seq),
reason,
static_cast<unsigned long>(telemetry_.activeBytes),
static_cast<unsigned long>(telemetry_.activeChunks));
return written > 0 && static_cast<size_t>(written) < outSize;
}

bool BridgeAudioUplink::requestTerminal(BridgeAudioTerminalKind kind,
uint32_t seq,
uint32_t nowMs,
const char* reason) {
if (kind == BridgeAudioTerminalKind::None) {
return fail("audio_uplink_bad_terminal");
}

const uint32_t resolvedSeq = seq != 0 ? seq : telemetry_.lastSeq;
const bool alreadyPending = telemetry_.terminalPending;
telemetry_.active = false;
telemetry_.terminalPending = true;
telemetry_.pendingTerminal = kind;
if (!alreadyPending) {
telemetry_.terminalRequestedAtMs = nowMs;
}
pendingTerminalSeq_ = resolvedSeq;
copyTerminalReason(reason != nullptr ? reason : "audio_uplink_aborted");
const BridgeAudioTerminalServiceResult result = servicePendingTerminal(nowMs);
return result == BridgeAudioTerminalServiceResult::Queued ||
result == BridgeAudioTerminalServiceResult::Pending;
}

void BridgeAudioUplink::copyTerminalReason(const char* reason) {
if (reason == nullptr) {
pendingTerminalReason_[0] = '\0';
return;
}
size_t out = 0;
while (reason[out] != '\0' && out < sizeof(pendingTerminalReason_) - 1u) {
const char value = reason[out];
const bool safe = (value >= 'a' && value <= 'z') ||
(value >= 'A' && value <= 'Z') ||
(value >= '0' && value <= '9') ||
value == '_' || value == '-' || value == '.' || value == ':';
pendingTerminalReason_[out] = safe ? value : '_';
++out;
}
pendingTerminalReason_[out] = '\0';
if (out == 0) {
copyTerminalReason("audio_uplink_aborted");
}
}

bool BridgeAudioUplink::fail(const char* reason) {
telemetry_.errors++;
copyError(reason);
Expand Down
Loading