@@ -251,6 +251,100 @@ function layerFor(input: {
251251 ) ;
252252}
253253
254+ /**
255+ * Send starting then running four seconds later and assert both
256+ * `live_activity_update` jobs, with phase running / Working on the second.
257+ *
258+ * @param input - Aggregates and the queued-job sink for assertions
259+ */
260+ function * sendStartingThenRunningLiveActivityUpdates ( input : {
261+ readonly startingAggregate : RelayAgentActivityAggregateState ;
262+ readonly runningAggregate : RelayAgentActivityAggregateState ;
263+ readonly queuedJobs : Array < SignedApnsDeliveryJob > ;
264+ } ) {
265+ const deliveries = yield * ApnsDeliveries . ApnsDeliveries ;
266+ const first = yield * deliveries . sendForTarget ( {
267+ target,
268+ aggregate : input . startingAggregate ,
269+ nowMs : 0 ,
270+ } ) ;
271+ expect ( first ?. kind ) . toBe ( "live_activity_update" ) ;
272+
273+ const second = yield * deliveries . sendForTarget ( {
274+ target : {
275+ ...target ,
276+ last_aggregate_json : JSON . stringify ( input . startingAggregate ) ,
277+ last_live_activity_delivery_at : "1970-01-01T00:00:00.000Z" ,
278+ } ,
279+ aggregate : input . runningAggregate ,
280+ nowMs : 4_000 ,
281+ } ) ;
282+
283+ expect ( second ?. kind ) . toBe ( "live_activity_update" ) ;
284+ expect ( input . queuedJobs ) . toMatchObject ( [
285+ {
286+ payload : {
287+ kind : "live_activity_update" ,
288+ target : {
289+ token : "activity-token" ,
290+ } ,
291+ } ,
292+ } ,
293+ {
294+ payload : {
295+ kind : "live_activity_update" ,
296+ target : {
297+ token : "activity-token" ,
298+ } ,
299+ aggregate : {
300+ activities : [ { phase : "running" , status : "Working" } ] ,
301+ } ,
302+ } ,
303+ } ,
304+ ] ) ;
305+ }
306+
307+ /**
308+ * Regression for starting→running inside the 15s Live Activity throttle:
309+ * both updates must queue, and the second payload is the running aggregate.
310+ *
311+ * @returns Effect that drives both deliveries and asserts the queued jobs
312+ */
313+ function queuesLiveActivityUpdateOnStartingToRunningPhaseChange ( ) {
314+ const attempts : Array < DeliveryAttempts . DeliveryAttemptInput > = [ ] ;
315+ const queuedJobs : Array < SignedApnsDeliveryJob > = [ ] ;
316+ const startingAggregate : RelayAgentActivityAggregateState = {
317+ ...aggregate ,
318+ activities : [
319+ {
320+ ...aggregate . activities [ 0 ] ! ,
321+ phase : "starting" ,
322+ status : "Connecting" ,
323+ } ,
324+ ] ,
325+ } ;
326+ const runningAggregate : RelayAgentActivityAggregateState = {
327+ ...startingAggregate ,
328+ updatedAt : "1970-01-01T00:00:04.000Z" ,
329+ activities : [
330+ {
331+ ...startingAggregate . activities [ 0 ] ! ,
332+ phase : "running" ,
333+ status : "Working" ,
334+ updatedAt : "1970-01-01T00:00:04.000Z" ,
335+ } ,
336+ ] ,
337+ } ;
338+
339+ return Effect . gen (
340+ sendStartingThenRunningLiveActivityUpdates . bind ( undefined , {
341+ startingAggregate,
342+ runningAggregate,
343+ queuedJobs,
344+ } ) ,
345+ ) . pipe ( Effect . provide ( layerFor ( { attempts, queuedJobs } ) ) ) ;
346+ }
347+
254348describe ( "ApnsDeliveries" , ( ) => {
255349 it . effect ( "skips Apple delivery when an Android-only relay disables APNs" , ( ) => {
256350 const attempts : Array < DeliveryAttempts . DeliveryAttemptInput > = [ ] ;
@@ -616,6 +710,226 @@ describe("ApnsDeliveries", () => {
616710 } ,
617711 ) ;
618712
713+ it . effect ( "throttles timestamp-only changes while an activity already awaits input" , ( ) => {
714+ const attempts : Array < DeliveryAttempts . DeliveryAttemptInput > = [ ] ;
715+ const queuedJobs : Array < SignedApnsDeliveryJob > = [ ] ;
716+ const waitingAggregate : RelayAgentActivityAggregateState = {
717+ ...aggregate ,
718+ activities : [
719+ {
720+ ...aggregate . activities [ 0 ] ! ,
721+ phase : "waiting_for_input" ,
722+ status : "Input" ,
723+ } ,
724+ ] ,
725+ } ;
726+ const laterWaitingAggregate : RelayAgentActivityAggregateState = {
727+ ...waitingAggregate ,
728+ updatedAt : "1970-01-01T00:00:04.000Z" ,
729+ activities : [
730+ {
731+ ...waitingAggregate . activities [ 0 ] ! ,
732+ updatedAt : "1970-01-01T00:00:04.000Z" ,
733+ } ,
734+ ] ,
735+ } ;
736+ const waitingAggregateJson = JSON . stringify ( waitingAggregate ) ;
737+
738+ return Effect . gen ( function * ( ) {
739+ const deliveries = yield * ApnsDeliveries . ApnsDeliveries ;
740+ const result = yield * deliveries . sendForTarget ( {
741+ target : {
742+ ...target ,
743+ last_aggregate_json : waitingAggregateJson ,
744+ last_live_activity_delivery_at : "1970-01-01T00:00:04.000Z" ,
745+ } ,
746+ aggregate : laterWaitingAggregate ,
747+ nowMs : 5_000 ,
748+ } ) ;
749+
750+ expect ( result ) . toBeNull ( ) ;
751+ expect ( queuedJobs ) . toEqual ( [ ] ) ;
752+ expect ( attempts ) . toEqual ( [ ] ) ;
753+ } ) . pipe ( Effect . provide ( layerFor ( { attempts, queuedJobs } ) ) ) ;
754+ } ) ;
755+
756+ it . effect ( "throttles ordering-only changes while an activity already awaits input" , ( ) => {
757+ const attempts : Array < DeliveryAttempts . DeliveryAttemptInput > = [ ] ;
758+ const queuedJobs : Array < SignedApnsDeliveryJob > = [ ] ;
759+ const waitingRow = {
760+ ...aggregate . activities [ 0 ] ! ,
761+ phase : "waiting_for_input" as const ,
762+ status : "Input" ,
763+ } ;
764+ const runningRow = {
765+ ...aggregate . activities [ 0 ] ! ,
766+ threadId : "thread-2" as RelayAgentActivityState [ "threadId" ] ,
767+ threadTitle : "Other thread" ,
768+ phase : "running" as const ,
769+ status : "Working" ,
770+ } ;
771+ const waitingAggregate : RelayAgentActivityAggregateState = {
772+ ...aggregate ,
773+ activeCount : 2 ,
774+ activities : [ waitingRow , runningRow ] ,
775+ } ;
776+ const reorderedAggregate : RelayAgentActivityAggregateState = {
777+ ...waitingAggregate ,
778+ activities : [ runningRow , waitingRow ] ,
779+ } ;
780+ const waitingAggregateJson = JSON . stringify ( waitingAggregate ) ;
781+
782+ return Effect . gen ( function * ( ) {
783+ const deliveries = yield * ApnsDeliveries . ApnsDeliveries ;
784+ const result = yield * deliveries . sendForTarget ( {
785+ target : {
786+ ...target ,
787+ last_aggregate_json : waitingAggregateJson ,
788+ last_live_activity_delivery_at : "1970-01-01T00:00:04.000Z" ,
789+ } ,
790+ aggregate : reorderedAggregate ,
791+ nowMs : 5_000 ,
792+ } ) ;
793+
794+ expect ( result ) . toBeNull ( ) ;
795+ expect ( queuedJobs ) . toEqual ( [ ] ) ;
796+ expect ( attempts ) . toEqual ( [ ] ) ;
797+ } ) . pipe ( Effect . provide ( layerFor ( { attempts, queuedJobs } ) ) ) ;
798+ } ) ;
799+
800+ it . effect ( "queues an update when a new thread awaits input while another already does" , ( ) => {
801+ const attempts : Array < DeliveryAttempts . DeliveryAttemptInput > = [ ] ;
802+ const queuedJobs : Array < SignedApnsDeliveryJob > = [ ] ;
803+ const waitingRow = {
804+ ...aggregate . activities [ 0 ] ! ,
805+ phase : "waiting_for_input" as const ,
806+ status : "Input" ,
807+ } ;
808+ const runningRow = {
809+ ...aggregate . activities [ 0 ] ! ,
810+ threadId : "thread-2" as RelayAgentActivityState [ "threadId" ] ,
811+ threadTitle : "Other thread" ,
812+ phase : "running" as const ,
813+ status : "Working" ,
814+ } ;
815+ const newWaitingRow = {
816+ ...aggregate . activities [ 0 ] ! ,
817+ threadId : "thread-3" as RelayAgentActivityState [ "threadId" ] ,
818+ threadTitle : "New waiting thread" ,
819+ phase : "waiting_for_input" as const ,
820+ status : "Input" ,
821+ updatedAt : "1970-01-01T00:00:04.000Z" ,
822+ } ;
823+ const previousAggregate : RelayAgentActivityAggregateState = {
824+ ...aggregate ,
825+ activeCount : 2 ,
826+ activities : [ waitingRow , runningRow ] ,
827+ } ;
828+ const nextAggregate : RelayAgentActivityAggregateState = {
829+ ...previousAggregate ,
830+ updatedAt : "1970-01-01T00:00:04.000Z" ,
831+ activities : [ waitingRow , newWaitingRow ] ,
832+ } ;
833+ const previousAggregateJson = JSON . stringify ( previousAggregate ) ;
834+
835+ return Effect . gen ( function * ( ) {
836+ const deliveries = yield * ApnsDeliveries . ApnsDeliveries ;
837+ const result = yield * deliveries . sendForTarget ( {
838+ target : {
839+ ...target ,
840+ last_aggregate_json : previousAggregateJson ,
841+ last_live_activity_delivery_at : "1970-01-01T00:00:04.000Z" ,
842+ } ,
843+ aggregate : nextAggregate ,
844+ nowMs : 5_000 ,
845+ } ) ;
846+
847+ expect ( result ?. kind ) . toBe ( "live_activity_update" ) ;
848+ expect ( queuedJobs ) . toMatchObject ( [
849+ {
850+ payload : {
851+ kind : "live_activity_update" ,
852+ target : {
853+ token : "activity-token" ,
854+ } ,
855+ aggregate : {
856+ activities : [
857+ { phase : "waiting_for_input" , threadId : "thread" } ,
858+ { phase : "waiting_for_input" , threadId : "thread-3" } ,
859+ ] ,
860+ } ,
861+ } ,
862+ } ,
863+ ] ) ;
864+ } ) . pipe ( Effect . provide ( layerFor ( { attempts, queuedJobs } ) ) ) ;
865+ } ) ;
866+
867+ it . effect (
868+ "queues an update when a waiting thread leaves while another still awaits input" ,
869+ ( ) => {
870+ const attempts : Array < DeliveryAttempts . DeliveryAttemptInput > = [ ] ;
871+ const queuedJobs : Array < SignedApnsDeliveryJob > = [ ] ;
872+ const waitingRow = {
873+ ...aggregate . activities [ 0 ] ! ,
874+ phase : "waiting_for_input" as const ,
875+ status : "Input" ,
876+ } ;
877+ const otherWaitingRow = {
878+ ...aggregate . activities [ 0 ] ! ,
879+ threadId : "thread-2" as RelayAgentActivityState [ "threadId" ] ,
880+ threadTitle : "Other thread" ,
881+ phase : "waiting_for_input" as const ,
882+ status : "Input" ,
883+ } ;
884+ const replacementRunningRow = {
885+ ...aggregate . activities [ 0 ] ! ,
886+ threadId : "thread-3" as RelayAgentActivityState [ "threadId" ] ,
887+ threadTitle : "Replacement thread" ,
888+ phase : "running" as const ,
889+ status : "Working" ,
890+ } ;
891+ const previousAggregate : RelayAgentActivityAggregateState = {
892+ ...aggregate ,
893+ activeCount : 2 ,
894+ activities : [ waitingRow , otherWaitingRow ] ,
895+ } ;
896+ const nextAggregate : RelayAgentActivityAggregateState = {
897+ ...previousAggregate ,
898+ updatedAt : "1970-01-01T00:00:04.000Z" ,
899+ activities : [ waitingRow , replacementRunningRow ] ,
900+ } ;
901+ const previousAggregateJson = JSON . stringify ( previousAggregate ) ;
902+
903+ return Effect . gen ( function * ( ) {
904+ const deliveries = yield * ApnsDeliveries . ApnsDeliveries ;
905+ const result = yield * deliveries . sendForTarget ( {
906+ target : {
907+ ...target ,
908+ last_aggregate_json : previousAggregateJson ,
909+ last_live_activity_delivery_at : "1970-01-01T00:00:04.000Z" ,
910+ } ,
911+ aggregate : nextAggregate ,
912+ nowMs : 5_000 ,
913+ } ) ;
914+
915+ expect ( result ?. kind ) . toBe ( "live_activity_update" ) ;
916+ expect ( queuedJobs ) . toMatchObject ( [
917+ {
918+ payload : {
919+ kind : "live_activity_update" ,
920+ aggregate : {
921+ activities : [
922+ { phase : "waiting_for_input" , threadId : "thread" } ,
923+ { phase : "running" , threadId : "thread-3" } ,
924+ ] ,
925+ } ,
926+ } ,
927+ } ,
928+ ] ) ;
929+ } ) . pipe ( Effect . provide ( layerFor ( { attempts, queuedJobs } ) ) ) ;
930+ } ,
931+ ) ;
932+
619933 it . effect (
620934 "throttles updates for changed aggregates with stable counts and no pending attention" ,
621935 ( ) => {
@@ -645,6 +959,11 @@ describe("ApnsDeliveries", () => {
645959 } ,
646960 ) ;
647961
962+ it . effect (
963+ "queues an update when phase changes from starting to running inside the throttle window" ,
964+ queuesLiveActivityUpdateOnStartingToRunningPhaseChange ,
965+ ) ;
966+
648967 it . effect ( "queues an end for an active Live Activity when Live Activities are disabled" , ( ) => {
649968 const attempts : Array < DeliveryAttempts . DeliveryAttemptInput > = [ ] ;
650969 const queuedJobs : Array < SignedApnsDeliveryJob > = [ ] ;
0 commit comments