Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
64 changes: 64 additions & 0 deletions infra/relay/src/agentActivity/AgentActivityPublisher.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -266,6 +268,68 @@ describe("AgentActivityPublisher", () => {
});
});

it.effect("ends an idle card only once nothing is left to show", () => {
const deliveredBefore: Array<string> = [];
const sent: Array<Parameters<ApnsDeliveries.ApnsDeliveries["Service"]["sendForTarget"]>[0]> =
[];
let activeStates: ReadonlyArray<RelayAgentActivityState> = [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");
Expand Down
56 changes: 55 additions & 1 deletion infra/relay/src/agentActivity/AgentActivityPublisher.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand Down Expand Up @@ -44,6 +47,11 @@ export class AgentActivityPublisher extends Context.Service<
readonly userId: string;
readonly deviceId: string;
}) => Effect.Effect<RelayDeliveryResult | null, AgentActivityPublishError>;
/** 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") {}

Expand Down Expand Up @@ -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,
Expand Down
58 changes: 58 additions & 0 deletions infra/relay/src/agentActivity/ApnsDeliveries.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down Expand Up @@ -916,6 +917,63 @@ 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,
activityStates?: ReadonlyArray<RelayAgentActivityState>,
) => {
const attempts: Array<DeliveryAttempts.DeliveryAttemptInput> = [];
return ApnsDeliveries.ApnsDeliveries.pipe(
Effect.flatMap((deliveries) => deliveries.processSignedJob(signedEnd)),
Effect.map(() => attempts.map((attempt) => attempt.apnsReason)),
Effect.provide(
makeLayer({
attempts,
currentTargets: [currentTarget],
...(activityStates ? { activityStates } : {}),
}),
),
);
};

// 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("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 });
// 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<DeliveryAttempts.DeliveryAttemptInput> = [];
const clearedStarts: Array<
Expand Down
37 changes: 37 additions & 0 deletions infra/relay/src/agentActivity/ApnsDeliveries.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ import {
sanitizeAgentActivityAggregateState,
sanitizeApnsNotificationPayload,
} from "./agentActivityPayloads.ts";
import { makeAggregateState } from "./agentActivityAggregate.ts";
import * as Apns from "./ApnsClient.ts";
import {
ApnsDeliveryJobLiveActivityAggregateMissing,
Expand Down Expand Up @@ -576,6 +577,28 @@ export const make = Effect.gen(function* () {
);
});

// 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),
),
),
);
});

const stateIdentityIsCurrent = Effect.fnUntraced(function* (input: {
readonly userId: string;
readonly environmentId: string;
Expand Down Expand Up @@ -767,6 +790,20 @@ export const make = Effect.gen(function* () {
});
return staleJobResult({ deviceId: input.target.device_id, kind: input.kind });
}
// 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* userHasContentToShow(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" &&
Expand Down
1 change: 1 addition & 0 deletions infra/relay/src/agentActivity/FcmDeliveries.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
49 changes: 48 additions & 1 deletion infra/relay/src/agentActivity/LiveActivities.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -43,6 +43,18 @@ export class LiveActivityTargetListPersistenceError extends Schema.TaggedError<L
}
}

export class LiveActivityIdleTargetListPersistenceError extends Schema.TaggedError<LiveActivityIdleTargetListPersistenceError>()(
"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>()(
"LiveActivityDeliveryMarkPersistenceError",
{
Expand Down Expand Up @@ -97,6 +109,12 @@ export class LiveActivities extends Context.Service<
readonly listTargets: (input: {
readonly userId: string;
}) => Effect.Effect<ReadonlyArray<TargetRow>, 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;
Expand Down Expand Up @@ -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,
Expand Down
2 changes: 2 additions & 0 deletions infra/relay/src/agentActivity/MobileRegistrations.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -190,6 +191,7 @@ function makeAgentActivityPublisher(
return {
publish: () => Effect.succeed({ ok: true, deliveries: [] }),
replayForLiveActivityRegistration: () => Effect.succeed(null),
endIdleLiveActivities: Effect.void,
...overrides,
};
}
Expand Down
1 change: 1 addition & 0 deletions infra/relay/src/http/Api.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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) =>
Expand Down
Loading
Loading