Skip to content

Commit 88fbe2e

Browse files
committed
Add interruption handling to realtime out
1 parent 3b75440 commit 88fbe2e

2 files changed

Lines changed: 90 additions & 43 deletions

File tree

‎src/simulflow/transport/out.clj‎

Lines changed: 48 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@
22
(:require
33
[clojure.core.async :as a :refer [<!! >!! chan timeout]]
44
[clojure.core.async.flow :as flow]
5-
[simulflow.async :refer [vthread-loop]]
5+
[simulflow.async :as async :refer [vthread-loop]]
66
[simulflow.frame :as frame]
77
[simulflow.schema :as schema]
88
[simulflow.transport.protocols :as tp]
@@ -119,7 +119,8 @@
119119
;; Channels following activity monitor pattern
120120
timer-in-ch (a/chan 1024)
121121
timer-out-ch (a/chan 1024)
122-
audio-write-ch (a/chan 1024)]
122+
audio-write-ch (a/chan 1024)
123+
command-ch (a/chan 1024)]
123124

124125
;; Minimal timer process - just sends timing events (like activity monitor)
125126
(vthread-loop []
@@ -130,28 +131,43 @@
130131

131132
;; Audio writer process - handles only audio I/O side effects
132133
(vthread-loop []
133-
(when-let [audio-command (<!! audio-write-ch)]
134-
(when (= (:command audio-command) :write-audio)
135-
(let [current-time (u/mono-time)
136-
delay-until (:delay-until audio-command 0)
137-
wait-time (max 0 (- delay-until current-time))]
138-
(t/log! {:data {:command :write-audio
139-
:delay-until (:delay-until audio-command 0)
140-
:current-time current-time
141-
:wait-time wait-time}
142-
:level :debug
143-
:sample 0.05
144-
:id :realtime-out})
145-
(when (pos? wait-time)
146-
(<!! (timeout wait-time)))
147-
(>!! chan (:data audio-command))))
148-
(recur)))
134+
;; using :priority true to always prefer command-ch in case both audio-write and command chans have value
135+
(let [[val port] (a/alts!! [command-ch audio-write-ch] :priority true)]
136+
(when val
137+
(cond
138+
(= port command-ch)
139+
(case (:command/kind val)
140+
:command/drain-queue (do
141+
(async/drain-channel! audio-write-ch)
142+
(t/log! {:level :debug :id :transport-out} "Drained audio queue"))
143+
;; unknown command
144+
nil)
145+
146+
(= port audio-write-ch)
147+
(when (= (:command/kind val) :command/write-audio)
148+
(let [current-time (u/mono-time)
149+
delay-until (:delay-until val 0)
150+
wait-time (max 0 (- delay-until current-time))]
151+
(t/log! {:data {:command :write-audio
152+
:delay-until (:delay-until val 0)
153+
:current-time current-time
154+
:wait-time wait-time}
155+
:level :debug
156+
:sample 0.05
157+
:id :realtime-out})
158+
(when (pos? wait-time)
159+
(<!! (timeout wait-time)))
160+
(>!! chan (:data val))))
161+
:else
162+
nil)
163+
(recur))))
149164

150165
;; Return state with minimal setup
151166
(into parsed-params
152-
{::flow/in-ports {:timer-out timer-out-ch}
153-
::flow/out-ports {:timer-in timer-in-ch
154-
:audio-write audio-write-ch}
167+
{::flow/in-ports {::timer-out timer-out-ch}
168+
::flow/out-ports {::timer-in timer-in-ch
169+
::command command-ch
170+
::audio-write audio-write-ch}
155171
;; Initial business logic state (managed in transform)
156172
::speaking? false
157173
::last-send-time 0
@@ -206,23 +222,23 @@
206222
audio-frame (if serializer
207223
(tp/serialize-frame serializer frame)
208224
audio-data)
209-
audio-command {:command :write-audio
225+
audio-command {:command/kind :command/write-audio
210226
:data audio-frame
211227
:delay-until next-send-time
212228
:sample-rate sample-rate}]
213229

214230
[updated-state (-> (frame/send (when should-emit-start? (frame/bot-speech-start true)))
215-
(assoc :audio-write [audio-command]))]))
231+
(assoc ::audio-write [audio-command]))]))
216232

217233
(defn base-realtime-out-transform
218234
[{::keys [now] :as state
219235
:or {now (u/mono-time)}} input-port frame]
220236
(cond
221237
;; Handle incoming audio frames - core business logic moved here
222-
(frame/audio-output-raw? frame)
238+
(and (frame/audio-output-raw? frame) (not (:pipeline/interrupted? state)))
223239
(process-realtime-out-audio-frame state frame now)
224240

225-
(and (= input-port :timer-out)
241+
(and (= input-port ::timer-out)
226242
(:timer/tick frame))
227243
(let [silence-duration (- (:timer/timestamp frame) (::last-send-time state 0))
228244
should-emit-stop? (and (::speaking? state)
@@ -240,8 +256,14 @@
240256
[(assoc state :transport/serializer new-serializer) {}]
241257
[state {}])
242258

259+
(frame/control-interrupt-start? frame)
260+
[(assoc state :pipeline/interrupted? true) {::command [{:command/kind :command/drain-queue}]}]
261+
262+
(frame/control-interrupt-stop? frame)
263+
[(assoc state :pipeline/interrupted? false)]
264+
243265
;; Default case
244-
:else [state {}]))
266+
:else [state]))
245267

246268
(defn realtime-out-fn
247269
"Processor fn that sends audio chunks to output channel in a realtime manner"

‎test/simulflow/transport/out_test.clj‎

Lines changed: 42 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -19,8 +19,8 @@
1919
(is (= current-time (::sut/last-send-time new-state)))
2020
(is (= 1 (count (:sys-out output))))
2121
(is (frame/bot-speech-start? (first (:sys-out output))))
22-
(is (= 1 (count (:audio-write output))))
23-
(is (= 16000 (:sample-rate (first (:audio-write output)))))))
22+
(is (= 1 (count (::sut/audio-write output))))
23+
(is (= 16000 (:sample-rate (first (::sut/audio-write output)))))))
2424

2525
(testing "process-realtime-out-audio-frame with different sample rate"
2626
(let [state {::sut/speaking? false
@@ -29,15 +29,15 @@
2929
frame (frame/audio-output-raw {:audio (byte-array [1 2 3]) :sample-rate 24000})
3030
current-time 2000
3131
[new-state output] (sut/process-realtime-out-audio-frame state frame current-time)
32-
command (first (:audio-write output))]
33-
(is (= (update command :data vec) {:command :write-audio
32+
command (first (::sut/audio-write output))]
33+
(is (= (update command :data vec) {:command/kind :command/write-audio
3434
:data [1, 2, 3]
3535
:delay-until 2000
3636
:sample-rate 24000}))
3737

3838
(is (true? (::sut/speaking? new-state)))
39-
(is (= 1 (count (:audio-write output))))
40-
(is (= 24000 (:sample-rate (first (:audio-write output))))))))
39+
(is (= 1 (count (::sut/audio-write output))))
40+
(is (= 24000 (:sample-rate (first (::sut/audio-write output))))))))
4141

4242
(deftest realtime-out-transform-test
4343
(testing "transform with audio frame"
@@ -48,18 +48,43 @@
4848
(is (true? (::sut/speaking? new-state)))
4949
(is (= 1 (count (:sys-out output))))
5050
(is (frame/bot-speech-start? (first (:sys-out output))))
51-
(is (= 1 (count (:audio-write output))))))
51+
(is (= 1 (count (::sut/audio-write output))))))
5252

5353
(testing "transform with timer tick (stop speaking)"
5454
(let [state {::sut/speaking? true
5555
::sut/last-send-time 1000
5656
:activity-detection/silence-threshold-ms 200}
5757
timer-frame {:timer/tick true :timer/timestamp 1300}
58-
[new-state output] (sut/base-realtime-out-transform state :timer-out timer-frame)]
58+
[new-state output] (sut/base-realtime-out-transform state ::sut/timer-out timer-frame)]
5959

6060
(is (false? (::sut/speaking? new-state)))
6161
(is (= 1 (count (:sys-out output))))
62-
(is (frame/bot-speech-stop? (first (:sys-out output)))))))
62+
(is (frame/bot-speech-stop? (first (:sys-out output))))))
63+
(testing "Handling of interruptions"
64+
(testing "Sends drain-queue command when an interruption starts"
65+
(let [state {::sut/speaking? true
66+
::sut/last-send-time 1000
67+
:activity-detection/silence-threshold-ms 200}
68+
[new-state output] (sut/base-realtime-out-transform state :sys-in (frame/control-interrupt-start true))]
69+
(is (true? (:pipeline/interrupted? new-state)))
70+
(is (= (first (::sut/command output)) {:command/kind :command/drain-queue}))))
71+
72+
(testing "Sends drain-queue command when an interruption starts"
73+
(let [state {::sut/speaking? false
74+
::sut/last-send-time 1000
75+
:pipeline/interrupted? true
76+
:activity-detection/silence-threshold-ms 200}
77+
[new-state output] (sut/base-realtime-out-transform state :sys-in (frame/control-interrupt-stop true))]
78+
(is (false? (:pipeline/interrupted? new-state)))
79+
(is (nil? output))))
80+
81+
(testing "audio-output frames are dropped when the pipeline is interrupted"
82+
(let [state {::sut/speaking? false :audio.out/sending-interval 25 :pipeline/interrupted? true}
83+
frame (frame/audio-output-raw {:audio (byte-array [1 2 3]) :sample-rate 16000})
84+
[new-state output] (sut/base-realtime-out-transform state :in frame)]
85+
86+
(is (= state new-state))
87+
(is (nil? output))))))
6388

6489
(deftest test-realtime-speakers-out-describe
6590
(testing "describe function returns correct structure"
@@ -104,7 +129,7 @@
104129
::sut/last-send-time 1000
105130
:activity-detection/silence-threshold-ms 200}
106131
timer-frame {:timer/tick true :timer/timestamp 1300} ; 300ms silence > 200ms threshold
107-
[new-state output] (sut/base-realtime-out-transform state :timer-out timer-frame)]
132+
[new-state output] (sut/base-realtime-out-transform state ::sut/timer-out timer-frame)]
108133

109134
(is (false? (::sut/speaking? new-state)))
110135
(is (= 1 (count (:sys-out output))))
@@ -142,8 +167,8 @@
142167
frame (frame/audio-output-raw {:audio (byte-array [1 2 3]) :sample-rate 16000})]
143168
(with-redefs [u/mono-time (constantly 1000)]
144169
(let [[_ output] (sut/base-realtime-out-transform state :in frame)
145-
audio-write (first (:audio-write output))]
146-
(is (= :write-audio (:command audio-write)))
170+
audio-write (first (::sut/audio-write output))]
171+
(is (= :command/write-audio (:command/kind audio-write)))
147172
(is (= [99 99 99] (:frame/data (:data audio-write)))))))) ; Should use serialized data
148173

149174
(testing "transform without serializer"
@@ -154,8 +179,8 @@
154179
frame (frame/audio-output-raw {:audio (byte-array [1 2 3]) :sample-rate 16000})]
155180
(with-redefs [u/mono-time (constantly 1000)]
156181
(let [[_ output] (sut/base-realtime-out-transform state :in frame)
157-
audio-write (first (:audio-write output))]
158-
(is (= :write-audio (:command audio-write)))
182+
audio-write (first (::sut/audio-write output))]
183+
(is (= :command/write-audio (:command/kind audio-write)))
159184
(is (= [1 2 3] (vec (:data audio-write)))))))))
160185

161186
(deftest test-realtime-speakers-out-edge-cases
@@ -225,7 +250,7 @@
225250
:timer/timestamp (+ (::sut/last-send-time state2) 1000)}
226251
[state3 output3] (sut/base-realtime-out-transform
227252
(assoc state2 :activity-detection/silence-threshold-ms 500)
228-
:timer-out timer-frame)]
253+
::sut/timer-out timer-frame)]
229254

230255
;; Verify state progression
231256
(is (false? (::sut/speaking? initial-state)))
@@ -249,7 +274,7 @@
249274
::sut/now current-time}
250275
frame (frame/audio-output-raw {:audio (byte-array [1 2 3]) :sample-rate 16000})
251276
[new-state output] (sut/base-realtime-out-transform state :in frame)
252-
audio-write (first (:audio-write output))
277+
audio-write (first (::sut/audio-write output))
253278
[next-state] (sut/base-realtime-out-transform (assoc new-state ::sut/now 1020) :in frame)
254279
[next-state2] (sut/base-realtime-out-transform (assoc next-state ::sut/now 1025) :in frame)]
255280

@@ -278,7 +303,7 @@
278303

279304
;; Timer tick over threshold
280305
timer-frame-over {:timer/tick true :timer/timestamp (+ base-time 350)}
281-
[state-over output-over] (sut/base-realtime-out-transform state :timer-out timer-frame-over)]
306+
[state-over output-over] (sut/base-realtime-out-transform state ::sut/timer-out timer-frame-over)]
282307

283308
;; Under threshold: still speaking
284309
(is (true? (::sut/speaking? state-under)))

0 commit comments

Comments
 (0)