diff --git a/infra/relay/src/agentActivity/ApnsDeliveries.test.ts b/infra/relay/src/agentActivity/ApnsDeliveries.test.ts index e98c2b639598..a636a7004d44 100644 --- a/infra/relay/src/agentActivity/ApnsDeliveries.test.ts +++ b/infra/relay/src/agentActivity/ApnsDeliveries.test.ts @@ -650,6 +650,76 @@ describe("ApnsDeliveries", () => { }, ); + it.effect( + "queues an update when phase changes from starting to running inside the throttle window", + () => { + const attempts: Array = []; + const queuedJobs: Array = []; + const startingAggregate: RelayAgentActivityAggregateState = { + ...aggregate, + activities: [ + { + ...aggregate.activities[0]!, + phase: "starting", + status: "Connecting", + }, + ], + }; + const runningAggregate: RelayAgentActivityAggregateState = { + ...startingAggregate, + updatedAt: "1970-01-01T00:00:04.000Z", + activities: [ + { + ...startingAggregate.activities[0]!, + phase: "running", + status: "Working", + updatedAt: "1970-01-01T00:00:04.000Z", + }, + ], + }; + + return Effect.gen(function* () { + const deliveries = yield* ApnsDeliveries.ApnsDeliveries; + const first = yield* deliveries.sendForTarget({ + target, + aggregate: startingAggregate, + nowMs: 0, + }); + expect(first?.kind).toBe("live_activity_update"); + + const second = yield* deliveries.sendForTarget({ + target: { + ...target, + last_aggregate_json: JSON.stringify(startingAggregate), + last_live_activity_delivery_at: "1970-01-01T00:00:00.000Z", + }, + aggregate: runningAggregate, + nowMs: 4_000, + }); + + expect(second?.kind).toBe("live_activity_update"); + expect(queuedJobs).toMatchObject([ + { + payload: { + kind: "live_activity_update", + target: { + token: "activity-token", + }, + }, + }, + { + payload: { + kind: "live_activity_update", + target: { + token: "activity-token", + }, + }, + }, + ]); + }).pipe(Effect.provide(makeLayer({ attempts, queuedJobs }))); + }, + ); + it.effect("queues an end for an active Live Activity when Live Activities are disabled", () => { const attempts: Array = []; const queuedJobs: Array = []; diff --git a/infra/relay/src/agentActivity/ApnsDeliveries.ts b/infra/relay/src/agentActivity/ApnsDeliveries.ts index 651f031f0efa..e19f9569c118 100644 --- a/infra/relay/src/agentActivity/ApnsDeliveries.ts +++ b/infra/relay/src/agentActivity/ApnsDeliveries.ts @@ -160,6 +160,19 @@ function aggregateNeedsAttention(aggregate: RelayAgentActivityAggregateState): b ); } +function aggregateHasPhaseChange( + previous: RelayAgentActivityAggregateState, + next: RelayAgentActivityAggregateState, +): boolean { + const previousPhases = new Map( + previous.activities.map((row) => [`${row.environmentId}\0${row.threadId}`, row.phase]), + ); + return next.activities.some((row) => { + const previousPhase = previousPhases.get(`${row.environmentId}\0${row.threadId}`); + return previousPhase !== undefined && previousPhase !== row.phase; + }); +} + function shouldUpdateLiveActivity(input: { readonly previousAggregate: RelayAgentActivityAggregateState | null; readonly nextAggregate: RelayAgentActivityAggregateState; @@ -184,6 +197,11 @@ function shouldUpdateLiveActivity(input: { if (newlyTerminalRows(input.previousAggregate, input.nextAggregate, true).length > 0) { return true; } + // starting→running keeps activeCount at 1 and is not attention/terminal, but + // the lock-screen copy changes (Connecting→Working) and is never republished. + if (aggregateHasPhaseChange(input.previousAggregate, input.nextAggregate)) { + return true; + } const lastDeliveryAtMs = input.lastDeliveryAt === null ? null