From 043c45af554873d80a671998ac36776f84eb1290 Mon Sep 17 00:00:00 2001 From: ishaanko Date: Thu, 1 Oct 2026 11:25:19 -0700 Subject: [PATCH 1/2] fix(relay): end idle Live Activities after the display window The 5 minute cron now ends armed iOS cards that have heard nothing for the 15 minute display window once their aggregate is empty. A queued contentless end is skipped when the user has live work again, unless the device turned Live Activities off. --- .../AgentActivityPublisher.test.ts | 64 +++++++++++++++++++ .../agentActivity/AgentActivityPublisher.ts | 56 +++++++++++++++- .../src/agentActivity/ApnsDeliveries.test.ts | 41 ++++++++++++ .../relay/src/agentActivity/ApnsDeliveries.ts | 17 ++++- .../src/agentActivity/FcmDeliveries.test.ts | 1 + .../relay/src/agentActivity/LiveActivities.ts | 49 +++++++++++++- .../agentActivity/MobileRegistrations.test.ts | 2 + infra/relay/src/http/Api.test.ts | 1 + infra/relay/src/worker.ts | 10 ++- 9 files changed, 237 insertions(+), 4 deletions(-) diff --git a/infra/relay/src/agentActivity/AgentActivityPublisher.test.ts b/infra/relay/src/agentActivity/AgentActivityPublisher.test.ts index 366e387f1e6f..c4a0cafffe20 100644 --- a/infra/relay/src/agentActivity/AgentActivityPublisher.test.ts +++ b/infra/relay/src/agentActivity/AgentActivityPublisher.test.ts @@ -2,6 +2,7 @@ import type { RelayAgentActivityState, RelayDeliveryResult } from "@t3tools/cont import { describe, expect, it } from "@effect/vitest"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; +import * as TestClock from "effect/testing/TestClock"; import * as AgentActivityRows from "./AgentActivityRows.ts"; import * as EnvironmentLinks from "../environments/EnvironmentLinks.ts"; @@ -58,6 +59,7 @@ function makeLiveActivities( return { register: () => Effect.void, listTargets: () => Effect.succeed([]), + listIdleArmedTargets: () => Effect.succeed([]), markDelivery: () => Effect.void, markStartQueued: () => Effect.void, clearStartQueued: () => Effect.void, @@ -266,6 +268,68 @@ describe("AgentActivityPublisher", () => { }); }); + it.effect("ends an idle card only once nothing is left to show", () => { + const deliveredBefore: Array = []; + const sent: Array[0]> = + []; + let activeStates: ReadonlyArray = [state]; + const endIdle = AgentActivityPublisher.AgentActivityPublisher.pipe( + Effect.flatMap((publisher) => publisher.endIdleLiveActivities), + Effect.provide( + publisherLayer.pipe( + Layer.provide( + Layer.mergeAll( + Layer.succeed( + AgentActivityRows.AgentActivityRows, + makeAgentActivityRows({ listForUser: () => Effect.sync(() => activeStates) }), + ), + Layer.succeed(EnvironmentLinks.EnvironmentLinks, makeEnvironmentLinks()), + Layer.succeed( + LiveActivities.LiveActivities, + makeLiveActivities({ + listIdleArmedTargets: (input) => + Effect.sync(() => { + deliveredBefore.push(input.deliveredBefore); + return [{ user_id: "dev:julius", device_id: "device-1" }]; + }), + listTargets: () => + Effect.succeed([ + { ...target("device-1"), activity_push_token: "activity-token" }, + target("device-2"), + ]), + }), + ), + Layer.succeed( + ApnsDeliveries.ApnsDeliveries, + makeApnsDeliveries({ + sendForTarget: (input) => + Effect.sync(() => { + sent.push(input); + return null; + }), + }), + ), + ), + ), + ), + ), + ); + + return Effect.gen(function* () { + yield* TestClock.adjust("1 hour"); + // Work started again since the last delivery; its own publish owns the card. + yield* endIdle; + expect(deliveredBefore).toEqual(["1970-01-01T00:45:00.000Z"]); + expect(sent).toEqual([]); + + activeStates = []; + yield* endIdle; + expect(sent).toMatchObject([ + { target: { device_id: "device-1" }, aggregate: null, replay: true }, + ]); + }); + }); + it.effect("publishes listed targets through the APNs delivery service", () => { const firstTarget = target("device-1"); const secondTarget = target("device-2"); diff --git a/infra/relay/src/agentActivity/AgentActivityPublisher.ts b/infra/relay/src/agentActivity/AgentActivityPublisher.ts index 61a420dd7858..41656a93414a 100644 --- a/infra/relay/src/agentActivity/AgentActivityPublisher.ts +++ b/infra/relay/src/agentActivity/AgentActivityPublisher.ts @@ -1,4 +1,7 @@ -import { makeAggregateState } from "./agentActivityAggregate.ts"; +import { + makeAggregateState, + TERMINAL_AGENT_ACTIVITY_DISPLAY_TTL_MS, +} from "./agentActivityAggregate.ts"; export { makeAggregateState, TERMINAL_AGENT_ACTIVITY_DISPLAY_TTL_MS, @@ -44,6 +47,11 @@ export class AgentActivityPublisher extends Context.Service< readonly userId: string; readonly deviceId: string; }) => Effect.Effect; + /** Ends armed iOS cards whose finished rows have outlived the display window. Run by the cron. */ + readonly endIdleLiveActivities: Effect.Effect< + void, + LiveActivities.LiveActivityIdleTargetListPersistenceError + >; } >()("t3code-relay/agentActivity/AgentActivityPublisher") {} @@ -143,6 +151,52 @@ export const make = Effect.gen(function* () { replay: true, }); }), + // A card only hears from the relay when an environment publishes or the app + // re-registers, so one left showing Done rows would keep them for hours. + // Cards with live work again are left to their own publish. + endIdleLiveActivities: Effect.gen(function* () { + const now = yield* DateTime.now; + const idleTargets = yield* liveActivities.listIdleArmedTargets({ + deliveredBefore: DateTime.formatIso( + DateTime.subtract(now, { milliseconds: TERMINAL_AGENT_ACTIVITY_DISPLAY_TTL_MS }), + ), + }); + yield* Effect.annotateCurrentSpan({ "relay.live_activities.idle_count": idleTargets.length }); + yield* Effect.forEach( + idleTargets, + (idle) => + Effect.gen(function* () { + const { activeStates, targets } = yield* Effect.all( + { + activeStates: rows.listForUser({ userId: idle.user_id }), + targets: liveActivities.listTargets({ userId: idle.user_id }), + }, + { concurrency: 2 }, + ); + const target = targets.find((row) => row.device_id === idle.device_id); + const aggregate = makeAggregateState({ + activeStates, + terminalState: null, + nowMs: now.epochMilliseconds, + }); + if (target === undefined || aggregate !== null) return; + yield* apnsDeliveries.sendForTarget({ + target, + aggregate: null, + nowMs: now.epochMilliseconds, + replay: true, + }); + }).pipe( + Effect.catch((error) => + Effect.logWarning("failed to end idle live activity", { + deviceId: idle.device_id, + error, + }), + ), + ), + { concurrency: 4, discard: true }, + ); + }).pipe(Effect.withSpan("relay.agent_activity_publisher.end_idle_live_activities")), publish: Effect.fn("relay.agent_activity_publisher.publish")(function* (input) { yield* Effect.annotateCurrentSpan({ "relay.environment_id": input.environmentId, diff --git a/infra/relay/src/agentActivity/ApnsDeliveries.test.ts b/infra/relay/src/agentActivity/ApnsDeliveries.test.ts index e98c2b639598..8a113d80ca28 100644 --- a/infra/relay/src/agentActivity/ApnsDeliveries.test.ts +++ b/infra/relay/src/agentActivity/ApnsDeliveries.test.ts @@ -230,6 +230,7 @@ function makeLayer(input: { Layer.succeed(LiveActivities.LiveActivities, { register: () => Effect.void, listTargets: () => Effect.succeed(input.currentTargets ?? [target]), + listIdleArmedTargets: () => Effect.succeed([]), markStartQueued: (queued) => Effect.sync(() => { input.queuedStarts?.push(queued); @@ -916,6 +917,46 @@ describe("ApnsDeliveries", () => { }).pipe(Effect.provide(makeLayer({ attempts }))); }); + describe("a queued contentless end while the user has live work", () => { + const signedEnd = signApnsDeliveryJob({ + secret: config.apnsDeliveryJobSigningSecret, + payload: makeApnsDeliveryJobPayload({ + kind: "live_activity_end", + userId: target.user_id, + deviceId: target.device_id, + token: "activity-token", + aggregate: null, + createdAt: "1970-01-01T00:00:00.000Z", + expiresAt: "1970-01-01T00:10:00.000Z", + jobId: "job-end-1", + }), + }); + // Returns the recorded attempt reason for one queued end. + const processEnd = (currentTarget: LiveActivities.TargetRow) => { + const attempts: Array = []; + return ApnsDeliveries.ApnsDeliveries.pipe( + Effect.flatMap((deliveries) => deliveries.processSignedJob(signedEnd)), + Effect.map(() => attempts.map((attempt) => attempt.apnsReason)), + Effect.provide(makeLayer({ attempts, currentTargets: [currentTarget] })), + ); + }; + + // The end was decided while nothing ran; the new work owns the card now. + it.effect("is skipped", () => + Effect.gen(function* () { + expect(yield* processEnd(target)).toEqual(["Stale APNs end job skipped."]); + }), + ); + + it.effect("still goes out when the device turned Live Activities off", () => + Effect.gen(function* () { + const reasons = yield* processEnd({ ...target, preferences_json: disabledPreferences }); + // The test relay's fake key fails at JWT signing, after the stale checks. + expect(reasons).toEqual(["Failed to sign APNs JWT for key key-id."]); + }), + ); + }); + it.effect("skips a queued start when the user no longer has live work", () => { const attempts: Array = []; const clearedStarts: Array< diff --git a/infra/relay/src/agentActivity/ApnsDeliveries.ts b/infra/relay/src/agentActivity/ApnsDeliveries.ts index 651f031f0efa..9d528b96eb04 100644 --- a/infra/relay/src/agentActivity/ApnsDeliveries.ts +++ b/infra/relay/src/agentActivity/ApnsDeliveries.ts @@ -569,7 +569,7 @@ export const make = Effect.gen(function* () { ), ), Effect.catchCause((cause) => - Effect.logWarning("live-work recheck failed; allowing queued start", { cause }).pipe( + Effect.logWarning("live-work recheck failed; assuming live work", { cause }).pipe( Effect.as(true), ), ), @@ -767,6 +767,21 @@ export const make = Effect.gen(function* () { }); return staleJobResult({ deviceId: input.target.device_id, kind: input.kind }); } + // A contentless end was decided while nothing was running. Work that + // started since owns the card, so ending it now would strand that work. + // An end for a device that turned Live Activities off still goes out. + if ( + input.kind === "live_activity_end" && + aggregate === null && + parsePreferences(currentTarget.preferences_json)?.liveActivitiesEnabled !== false && + (yield* userStillHasLiveWork(input.target.user_id)) + ) { + yield* attempts.completeSourceJob({ + sourceJobId: input.sourceJobId, + apnsReason: "Stale APNs end job skipped.", + }); + return staleJobResult({ deviceId: input.target.device_id, kind: input.kind }); + } } if ( input.kind === "live_activity_start" && diff --git a/infra/relay/src/agentActivity/FcmDeliveries.test.ts b/infra/relay/src/agentActivity/FcmDeliveries.test.ts index 46254e1e8b43..2e7c394bf446 100644 --- a/infra/relay/src/agentActivity/FcmDeliveries.test.ts +++ b/infra/relay/src/agentActivity/FcmDeliveries.test.ts @@ -127,6 +127,7 @@ function harness() { Layer.succeed(LiveActivities, { register: () => Effect.void, listTargets: () => Effect.sync(() => [current.target]), + listIdleArmedTargets: () => Effect.succeed([]), markDelivery: (input) => Effect.sync(() => { marked.push(input); diff --git a/infra/relay/src/agentActivity/LiveActivities.ts b/infra/relay/src/agentActivity/LiveActivities.ts index 0c88fbb73c39..d3b9bf2445e3 100644 --- a/infra/relay/src/agentActivity/LiveActivities.ts +++ b/infra/relay/src/agentActivity/LiveActivities.ts @@ -13,7 +13,7 @@ import * as Effect from "effect/Effect"; import * as Function from "effect/Function"; import * as Layer from "effect/Layer"; import * as Schema from "effect/Schema"; -import { and, eq, sql } from "drizzle-orm"; +import { and, eq, isNotNull, sql } from "drizzle-orm"; import * as RelayDb from "../db.ts"; import { relayLiveActivities, relayMobileDevices } from "../persistence/schema.ts"; @@ -43,6 +43,18 @@ export class LiveActivityTargetListPersistenceError extends Schema.TaggedError()( + "LiveActivityIdleTargetListPersistenceError", + { + deliveredBefore: Schema.String, + cause: Schema.Defect(), + }, +) { + override get message(): string { + return `Failed to list idle Live Activities delivered before ${this.deliveredBefore}.`; + } +} + export class LiveActivityDeliveryMarkPersistenceError extends Schema.TaggedError()( "LiveActivityDeliveryMarkPersistenceError", { @@ -97,6 +109,12 @@ export class LiveActivities extends Context.Service< readonly listTargets: (input: { readonly userId: string; }) => Effect.Effect, LiveActivityTargetListPersistenceError>; + readonly listIdleArmedTargets: (input: { + readonly deliveredBefore: string; + }) => Effect.Effect< + ReadonlyArray<{ readonly user_id: string; readonly device_id: string }>, + LiveActivityIdleTargetListPersistenceError + >; readonly markDelivery: (input: { readonly userId: string; readonly deviceId: string; @@ -250,6 +268,35 @@ export const make = Effect.gen(function* () { ); }), + // Armed cards that have heard nothing since `deliveredBefore`. Deliveries + // only run when an environment publishes or the app re-registers, so this is + // how the cron finds cards still showing Done rows past their display window. + listIdleArmedTargets: Effect.fn("relay.live_activities.list_idle_armed_targets")( + function* (input) { + return yield* db + .select({ + user_id: relayLiveActivities.userId, + device_id: relayLiveActivities.deviceId, + }) + .from(relayLiveActivities) + .where( + and( + isNotNull(relayLiveActivities.activityPushToken), + sql`coalesce(${relayLiveActivities.lastLiveActivityDeliveryAt}, ${relayLiveActivities.remoteStartedAt}, ${relayLiveActivities.updatedAt}) < ${input.deliveredBefore}`, + ), + ) + .pipe( + Effect.mapError( + (cause) => + new LiveActivityIdleTargetListPersistenceError({ + deliveredBefore: input.deliveredBefore, + cause, + }), + ), + ); + }, + ), + markDelivery: Effect.fn("relay.live_activities.mark_delivery")(function* (input) { yield* Effect.annotateCurrentSpan({ "relay.mobile.device_id": input.deviceId, diff --git a/infra/relay/src/agentActivity/MobileRegistrations.test.ts b/infra/relay/src/agentActivity/MobileRegistrations.test.ts index 457a545f74b0..73952322c863 100644 --- a/infra/relay/src/agentActivity/MobileRegistrations.test.ts +++ b/infra/relay/src/agentActivity/MobileRegistrations.test.ts @@ -67,6 +67,7 @@ function makeLiveActivities( return { register: () => Effect.void, listTargets: () => Effect.succeed([]), + listIdleArmedTargets: () => Effect.succeed([]), markDelivery: () => Effect.void, markStartQueued: () => Effect.void, clearStartQueued: () => Effect.void, @@ -190,6 +191,7 @@ function makeAgentActivityPublisher( return { publish: () => Effect.succeed({ ok: true, deliveries: [] }), replayForLiveActivityRegistration: () => Effect.succeed(null), + endIdleLiveActivities: Effect.void, ...overrides, }; } diff --git a/infra/relay/src/http/Api.test.ts b/infra/relay/src/http/Api.test.ts index 7e8458949a80..bceafcaaa6c8 100644 --- a/infra/relay/src/http/Api.test.ts +++ b/infra/relay/src/http/Api.test.ts @@ -1185,6 +1185,7 @@ describe("relay routing fallback", () => { return { ok: true as const, deliveries: [] }; }), replayForLiveActivityRegistration: () => Effect.succeed(null), + endIdleLiveActivities: Effect.void, }); const signatures = Layer.succeed(EnvironmentPublishSignatures.EnvironmentPublishSignatures, { verify: (input) => diff --git a/infra/relay/src/worker.ts b/infra/relay/src/worker.ts index d4d674d3342b..683e3feb5e72 100644 --- a/infra/relay/src/worker.ts +++ b/infra/relay/src/worker.ts @@ -356,8 +356,16 @@ export const ApiLive = Api.make( : Effect.logWarning("Failed to clean up inactive managed tunnels", { cause }), ), ), + AgentActivityPublisher.AgentActivityPublisher.pipe( + Effect.flatMap((publisher) => publisher.endIdleLiveActivities), + Effect.catchCause((cause) => + Cause.hasInterrupts(cause) + ? Effect.interrupt + : Effect.logWarning("Failed to end idle Live Activities", { cause }), + ), + ), ], - { concurrency: 2, discard: true }, + { concurrency: 3, discard: true }, ).pipe( Effect.withSpan("relay.cron.prune_expired_state"), // Export cron spans to Axiom like HTTP spans; the scope flushes them before the run ends. From 296bddeb17d2cf175fb12e4db99a7ae207cc367c Mon Sep 17 00:00:00 2001 From: ishaanko Date: Thu, 1 Oct 2026 11:33:52 -0700 Subject: [PATCH 2/2] fix(relay): keep a queued end from retiring a card with new finished work --- .../src/agentActivity/ApnsDeliveries.test.ts | 21 ++++++++++-- .../relay/src/agentActivity/ApnsDeliveries.ts | 32 ++++++++++++++++--- 2 files changed, 46 insertions(+), 7 deletions(-) diff --git a/infra/relay/src/agentActivity/ApnsDeliveries.test.ts b/infra/relay/src/agentActivity/ApnsDeliveries.test.ts index 8a113d80ca28..10dcd333666c 100644 --- a/infra/relay/src/agentActivity/ApnsDeliveries.test.ts +++ b/infra/relay/src/agentActivity/ApnsDeliveries.test.ts @@ -932,12 +932,21 @@ describe("ApnsDeliveries", () => { }), }); // Returns the recorded attempt reason for one queued end. - const processEnd = (currentTarget: LiveActivities.TargetRow) => { + const processEnd = ( + currentTarget: LiveActivities.TargetRow, + activityStates?: ReadonlyArray, + ) => { const attempts: Array = []; return ApnsDeliveries.ApnsDeliveries.pipe( Effect.flatMap((deliveries) => deliveries.processSignedJob(signedEnd)), Effect.map(() => attempts.map((attempt) => attempt.apnsReason)), - Effect.provide(makeLayer({ attempts, currentTargets: [currentTarget] })), + Effect.provide( + makeLayer({ + attempts, + currentTargets: [currentTarget], + ...(activityStates ? { activityStates } : {}), + }), + ), ); }; @@ -948,6 +957,14 @@ describe("ApnsDeliveries", () => { }), ); + it.effect("is skipped when that work already finished inside the display window", () => + Effect.gen(function* () { + expect(yield* processEnd(target, [{ ...state, phase: "completed" }])).toEqual([ + "Stale APNs end job skipped.", + ]); + }), + ); + it.effect("still goes out when the device turned Live Activities off", () => Effect.gen(function* () { const reasons = yield* processEnd({ ...target, preferences_json: disabledPreferences }); diff --git a/infra/relay/src/agentActivity/ApnsDeliveries.ts b/infra/relay/src/agentActivity/ApnsDeliveries.ts index 9d528b96eb04..f32752c67f0c 100644 --- a/infra/relay/src/agentActivity/ApnsDeliveries.ts +++ b/infra/relay/src/agentActivity/ApnsDeliveries.ts @@ -24,6 +24,7 @@ import { sanitizeAgentActivityAggregateState, sanitizeApnsNotificationPayload, } from "./agentActivityPayloads.ts"; +import { makeAggregateState } from "./agentActivityAggregate.ts"; import * as Apns from "./ApnsClient.ts"; import { ApnsDeliveryJobLiveActivityAggregateMissing, @@ -569,7 +570,29 @@ export const make = Effect.gen(function* () { ), ), Effect.catchCause((cause) => - Effect.logWarning("live-work recheck failed; assuming live work", { cause }).pipe( + Effect.logWarning("live-work recheck failed; allowing queued start", { cause }).pipe( + Effect.as(true), + ), + ), + ); + }); + + // Contentless ends are decided when the user's aggregate was empty. Work that + // starts after that, even work that already finished, has content to show + // again. Fails closed: a database hiccup keeps the card for the next sweep. + const userHasContentToShow = Effect.fnUntraced(function* (userId: string) { + const now = yield* DateTime.now; + return yield* activityRows.listForUser({ userId }).pipe( + Effect.map( + (activeStates) => + makeAggregateState({ + activeStates, + terminalState: null, + nowMs: now.epochMilliseconds, + }) !== null, + ), + Effect.catchCause((cause) => + Effect.logWarning("content recheck failed; keeping the card", { cause }).pipe( Effect.as(true), ), ), @@ -767,14 +790,13 @@ export const make = Effect.gen(function* () { }); return staleJobResult({ deviceId: input.target.device_id, kind: input.kind }); } - // A contentless end was decided while nothing was running. Work that - // started since owns the card, so ending it now would strand that work. - // An end for a device that turned Live Activities off still goes out. + // A contentless end must not retire a card that newer work now owns. An + // end for a device that turned Live Activities off still goes out. if ( input.kind === "live_activity_end" && aggregate === null && parsePreferences(currentTarget.preferences_json)?.liveActivitiesEnabled !== false && - (yield* userStillHasLiveWork(input.target.user_id)) + (yield* userHasContentToShow(input.target.user_id)) ) { yield* attempts.completeSourceJob({ sourceJobId: input.sourceJobId,