diff --git a/infra/relay/src/agentActivity/AgentActivityPublisher.ts b/infra/relay/src/agentActivity/AgentActivityPublisher.ts index 61a420dd7858..905a2ea73043 100644 --- a/infra/relay/src/agentActivity/AgentActivityPublisher.ts +++ b/infra/relay/src/agentActivity/AgentActivityPublisher.ts @@ -54,6 +54,7 @@ export const make = Effect.gen(function* () { const apnsDeliveries = yield* ApnsDeliveries.ApnsDeliveries; const fcmDeliveries = yield* FcmDeliveries.FcmDeliveries; + /** Keeps environment channel settings separate from each device's delivery settings. */ const publishForDeliveryUser = Effect.fnUntraced(function* (input: { readonly deliveryUser: EnvironmentLinks.AgentAwarenessDeliveryUserRecord; readonly state: RelayAgentActivityState | null; @@ -91,6 +92,10 @@ export const make = Effect.gen(function* () { apnsDeliveries.sendForTarget({ target, aggregate: liveActivityAggregate, + notificationState: + input.deliveryUser.liveActivitiesEnabled && input.deliveryUser.notificationsEnabled + ? input.state + : null, nowMs: input.nowMs, }), notificationOnlyAggregate === null diff --git a/infra/relay/src/agentActivity/ApnsDeliveries.test.ts b/infra/relay/src/agentActivity/ApnsDeliveries.test.ts index e98c2b639598..23d41dc0f8e9 100644 --- a/infra/relay/src/agentActivity/ApnsDeliveries.test.ts +++ b/infra/relay/src/agentActivity/ApnsDeliveries.test.ts @@ -34,6 +34,9 @@ import * as AgentActivityRows from "./AgentActivityRows.ts"; import * as ApnsDeliveries from "./ApnsDeliveries.ts"; import * as ApnsClient from "./ApnsClient.ts"; import * as ApnsProviderTokens from "./ApnsProviderTokens.ts"; +import * as AgentActivityPublisher from "./AgentActivityPublisher.ts"; +import * as EnvironmentLinks from "../environments/EnvironmentLinks.ts"; +import { FcmDeliveries } from "./FcmDeliveries.ts"; const config = RelayConfiguration.RelayConfiguration.of({ relayIssuer: "https://relay.example.test", @@ -146,6 +149,7 @@ const target: LiveActivities.TargetRow = { last_live_activity_delivery_at: null, }; +/** Shares test persistence and queue services between delivery and publisher tests. */ function makeLayer(input: { readonly attempts: Array; readonly sourceJobClaims?: ReadonlyMap; @@ -179,7 +183,7 @@ function makeLayer(input: { Layer.provide(ApnsClient.layer), Layer.provide(ApnsProviderTokens.layer), Layer.provide(ApnsDeliveryQueue.layer.pipe(Layer.provide(NodeCryptoLayer.layer))), - Layer.provide( + Layer.provideMerge( Layer.mergeAll( Layer.succeed(AgentActivityRows.AgentActivityRows, { upsert: () => Effect.void, @@ -257,6 +261,156 @@ function makeLayer(input: { } describe("ApnsDeliveries", () => { + for (const liveActivitiesEnabled of [false, true]) { + for (const [phase, maxAgeMs] of [ + ["waiting_for_input", 24 * 60 * 60 * 1_000], + ["waiting_for_approval", 24 * 60 * 60 * 1_000], + ["completed", 2 * 60 * 1_000], + ["failed", 2 * 60 * 1_000], + ] as const) { + for (const expired of [false, true]) { + it.effect( + `${expired ? "skips" : "queues"} the published ${phase} alert ${expired ? "after" : "at"} its age limit, Live Activities ${liveActivitiesEnabled ? "unarmed" : "disabled"}`, + () => { + const queuedJobs: SignedApnsDeliveryJob[] = []; + return Effect.gen(function* () { + const deliveries = yield* ApnsDeliveries.ApnsDeliveries; + yield* deliveries.sendForTarget({ + target: { + ...target, + push_token: "push-token", + activity_push_token: liveActivitiesEnabled ? null : target.activity_push_token, + preferences_json: liveActivitiesEnabled + ? enabledPreferences + : disabledPreferences, + }, + aggregate: null, + notificationState: { ...state, phase }, + nowMs: maxAgeMs + (expired ? 1 : 0), + }); + expect( + queuedJobs.filter((job) => job.payload.kind === "push_notification"), + ).toHaveLength(expired ? 0 : 1); + }).pipe(Effect.provide(makeLayer({ attempts: [], queuedJobs }))); + }, + ); + } + } + } + + for (const scenario of ["deleted", "replay", "muted", "event muted"] as const) { + it.effect(`keeps a ${scenario} published input state silent`, () => { + const queuedJobs: SignedApnsDeliveryJob[] = []; + return Effect.gen(function* () { + const deliveries = yield* ApnsDeliveries.ApnsDeliveries; + yield* deliveries.sendForTarget({ + target: { + ...target, + push_token: "push-token", + activity_push_token: null, + preferences_json: JSON.stringify({ + ...JSON.parse(enabledPreferences), + notificationsEnabled: scenario !== "muted", + notifyOnInput: scenario !== "event muted", + }), + }, + aggregate: { + ...aggregate, + activities: [{ ...aggregate.activities[0]!, phase: "waiting_for_input" }], + }, + notificationState: + scenario === "deleted" ? null : { ...state, phase: "waiting_for_input" }, + replay: scenario === "replay", + nowMs: 0, + }); + expect(queuedJobs).toEqual([]); + }).pipe(Effect.provide(makeLayer({ attempts: [], queuedJobs }))); + }); + } + + for (const environment of [ + { name: "both channels", liveActivitiesEnabled: true, notificationsEnabled: true }, + { name: "notifications only", liveActivitiesEnabled: false, notificationsEnabled: true }, + { name: "Live Activities only", liveActivitiesEnabled: true, notificationsEnabled: false }, + ]) { + for (const deviceLiveActivitiesEnabled of [false, true]) { + for (const phase of ["completed", "failed"] as const) { + it.effect( + `${environment.notificationsEnabled ? "queues" : "skips"} the published ${phase} alert with other work, environment ${environment.name}, device Live Activities ${deviceLiveActivitiesEnabled ? "unarmed" : "disabled"}`, + () => { + const queuedJobs: SignedApnsDeliveryJob[] = []; + const finished = { ...state, phase }; + const device = { + ...target, + push_token: "push-token", + activity_push_token: deviceLiveActivitiesEnabled ? null : target.activity_push_token, + preferences_json: deviceLiveActivitiesEnabled + ? enabledPreferences + : disabledPreferences, + }; + const otherWork = Array.from({ length: phase === "completed" ? 1 : 5 }, (_, index) => ({ + ...state, + threadId: `other-${index}` as RelayAgentActivityState["threadId"], + })); + return Effect.gen(function* () { + const publisher = yield* AgentActivityPublisher.AgentActivityPublisher; + yield* publisher.publish({ + environmentId: finished.environmentId, + environmentPublicKey: "key", + threadId: finished.threadId, + state: finished, + }); + expect( + queuedJobs + .filter((job) => job.payload.kind === "push_notification") + .map((job) => job.payload.notification), + ).toMatchObject( + environment.notificationsEnabled ? [{ threadId: finished.threadId, phase }] : [], + ); + }).pipe( + Effect.provide( + AgentActivityPublisher.layer.pipe( + Layer.provide( + makeLayer({ + attempts: [], + queuedJobs, + activityStates: [...otherWork, finished], + currentTargets: [device], + }), + ), + Layer.provide( + Layer.succeed(FcmDeliveries, { + enqueue: () => Effect.succeed(null), + process: () => Effect.void, + }), + ), + Layer.provide( + Layer.succeed(EnvironmentLinks.EnvironmentLinks, { + upsert: () => Effect.void, + listUsersForEnvironment: () => Effect.succeed([device.user_id]), + listDeliveryUsersForEnvironment: () => + Effect.succeed([ + { + userId: device.user_id, + notificationsEnabled: environment.notificationsEnabled, + liveActivitiesEnabled: environment.liveActivitiesEnabled, + }, + ]), + listPublicKeysForEnvironment: () => Effect.succeed([]), + listForUser: () => Effect.succeed([]), + getForUser: () => Effect.succeed(null), + revokeForUser: () => Effect.succeed(false), + }), + ), + ), + ), + ); + }, + ); + } + } + } + it.effect("skips Apple delivery when an Android-only relay disables APNs", () => { const attempts: Array = []; const queuedJobs: Array = []; @@ -567,6 +721,7 @@ describe("ApnsDeliveries", () => { last_live_activity_delivery_at: "1970-01-01T00:00:04.000Z", }, aggregate: waitingAggregate, + notificationState: { ...state, phase: "waiting_for_input" }, nowMs: 5_000, }); @@ -1983,8 +2138,13 @@ describe("fast completion delivery", () => { }; return Effect.gen(function* () { const d = yield* ApnsDeliveries.ApnsDeliveries; - yield* d.sendForTarget({ target: device, aggregate, nowMs: 0 }); - yield* d.sendForTarget({ target: device, aggregate: done, nowMs: 0 }); + yield* d.sendForTarget({ target: device, aggregate, notificationState: state, nowMs: 0 }); + yield* d.sendForTarget({ + target: device, + aggregate: done, + notificationState: { ...state, phase: "completed" }, + nowMs: 0, + }); expect( queuedJobs.some((x) => x.payload.alert !== null && x.payload.alert !== undefined), ).toBe(true); diff --git a/infra/relay/src/agentActivity/ApnsDeliveries.ts b/infra/relay/src/agentActivity/ApnsDeliveries.ts index 651f031f0efa..c3039b9b1884 100644 --- a/infra/relay/src/agentActivity/ApnsDeliveries.ts +++ b/infra/relay/src/agentActivity/ApnsDeliveries.ts @@ -1,5 +1,6 @@ import type { RelayAgentActivityAggregateState, + RelayAgentActivityState, RelayAgentAwarenessPreferences, RelayDeliveryKind, RelayDeliveryResult, @@ -42,6 +43,7 @@ import * as LiveActivities from "./LiveActivities.ts"; import * as RelayConfiguration from "../Config.ts"; import * as ApnsDeliveryQueue from "./ApnsDeliveryQueue.ts"; import { withSpanAttributes } from "../observability.ts"; +import { statusForPhase } from "./agentActivityAggregate.ts"; import { alertForAttentionTransition, @@ -198,22 +200,37 @@ function shouldUpdateLiveActivity(input: { ); } -// Completions replayed long after the fact (server restarts republish every -// recently-finished thread) must not ring the device again. - -function notificationForAggregate(input: { +/** + * Selects the published thread independently of the card's order and row limit. + * An omitted notificationState uses the aggregate; null suppresses the push. + * Existing age limits and device notification settings still apply. + */ +function notificationForDelivery(input: { readonly target: LiveActivities.TargetRow; readonly aggregate: RelayAgentActivityAggregateState | null; + readonly notificationState?: RelayAgentActivityState | null; readonly nowMs: number; }): ApnsNotificationPayload | null { - if (!input.target.push_token || input.aggregate === null) { + if (!input.target.push_token) { return null; } const preferences = parsePreferences(input.target.preferences_json); if (!preferences?.notificationsEnabled) { return null; } - const activity = input.aggregate.activities[0]; + if ( + input.notificationState && + isExpiredAgentActivityState(input.notificationState, input.nowMs) + ) { + return null; + } + const activity = + input.notificationState === undefined + ? input.aggregate?.activities[0] + : input.notificationState && { + ...input.notificationState, + status: statusForPhase(input.notificationState.phase), + }; if (!activity) { return null; } @@ -305,9 +322,11 @@ function chooseLiveActivityDelivery(input: { : "suppressed"; } +/** Falls back to a push only when no Live Activity owns the update, preserving silent replays. */ function chooseDelivery(input: { readonly target: LiveActivities.TargetRow; readonly aggregate: RelayAgentActivityAggregateState | null; + readonly notificationState?: RelayAgentActivityState | null; readonly nowMs: number; readonly replay?: boolean; }): ChosenDelivery | null { @@ -318,7 +337,7 @@ function chooseDelivery(input: { if (liveActivityDelivery) { return liveActivityDelivery; } - const notification = input.replay ? null : notificationForAggregate(input); + const notification = input.replay ? null : notificationForDelivery(input); return notification && input.target.push_token ? { kind: "push_notification", @@ -522,6 +541,7 @@ export class ApnsDeliveries extends Context.Service< readonly sendForTarget: (input: { readonly target: LiveActivities.TargetRow; readonly aggregate: RelayAgentActivityAggregateState | null; + readonly notificationState?: RelayAgentActivityState | null; readonly nowMs: number; readonly replay?: boolean; }) => Effect.Effect; @@ -1103,7 +1123,7 @@ export const make = Effect.gen(function* () { sendPushNotificationForTarget: Effect.fnUntraced(function* (input) { if (!config.apns) return null; const now = yield* DateTime.now; - const notification = notificationForAggregate({ + const notification = notificationForDelivery({ target: input.target, aggregate: input.aggregate, nowMs: now.epochMilliseconds, @@ -1122,12 +1142,7 @@ export const make = Effect.gen(function* () { }), sendForTarget: Effect.fnUntraced(function* (input) { if (!config.apns) return null; - const delivery = chooseDelivery({ - target: input.target, - aggregate: input.aggregate, - nowMs: input.nowMs, - replay: input.replay ?? false, - }); + const delivery = chooseDelivery(input); if (!delivery) { return null; } @@ -1142,13 +1157,7 @@ export const make = Effect.gen(function* () { }); return result; } - const notification = input.replay - ? null - : notificationForAggregate({ - target: input.target, - aggregate: input.aggregate, - nowMs: input.nowMs, - }); + const notification = input.replay ? null : notificationForDelivery(input); // The end event doubles as the "task finished" moment. When a companion // push notification is about to ring the device (below), the activity end // stays silent; otherwise the end itself carries the alert so LA-only