Skip to content
Merged
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
123 changes: 120 additions & 3 deletions infra/relay/src/agentActivity/ApnsClient.test.ts
Original file line number Diff line number Diff line change
@@ -1,13 +1,22 @@
import * as NodeCrypto from "node:crypto";

import { EnvironmentId, ThreadId } from "@t3tools/contracts";
import type { RelayAgentActivityAggregateState } from "@t3tools/contracts/relay";
import { describe, expect, it } from "@effect/vitest";
import * as DateTime from "effect/DateTime";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import { HttpClient } from "effect/unstable/http";
import * as Redacted from "effect/Redacted";
import * as Schema from "effect/Schema";
import * as HttpClient from "effect/unstable/http/HttpClient";
import * as HttpClientError from "effect/unstable/http/HttpClientError";

import type { RelayAgentActivityAggregateState } from "@t3tools/contracts/relay";
import { EnvironmentId, ThreadId } from "@t3tools/contracts";
import type { ApnsCredentials } from "../Config.ts";
import * as ApnsClient from "./ApnsClient.ts";

const isApnsJwtSigningError = Schema.is(ApnsClient.ApnsJwtSigningError);
const isApnsHttpRequestError = Schema.is(ApnsClient.ApnsHttpRequestError);

const TestLayer = ApnsClient.layer.pipe(
Layer.provide(
Layer.succeed(
Expand Down Expand Up @@ -137,4 +146,112 @@ describe("ApnsClient", () => {
});
}).pipe(Effect.provide(TestLayer)),
);

it.effect("preserves JWT signing context and the crypto cause", () =>
Effect.gen(function* () {
const apns = yield* ApnsClient.ApnsClient;
const request = apns.makePushNotificationRequest({
token: "push-token",
notification: {
title: "Thread",
body: "Input: Project",
environmentId: "env",
threadId: "thread",
deepLink: "/threads/env/thread",
},
});
const error = yield* Effect.flip(
apns.sendPushNotificationRequest({
credentials: {
teamId: "team-1",
keyId: "key-1",
privateKey: Redacted.make("not-a-private-key"),
bundleId: "com.t3tools.test",
environment: "sandbox",
},
request,
issuedAtUnixSeconds: 123,
}),
);

expect(isApnsJwtSigningError(error)).toBe(true);
if (!isApnsJwtSigningError(error)) {
return yield* Effect.die("expected APNs JWT signing error");
}
expect(error).toMatchObject({
teamId: "team-1",
keyId: "key-1",
issuedAtUnixSeconds: 123,
cause: expect.any(Error),
message: "Failed to sign APNs JWT for key key-1.",
});
}).pipe(Effect.provide(TestLayer)),
);

it.effect("preserves APNs request context and the HTTP cause", () => {
const httpCause = new Error("network unavailable");
const { privateKey } = NodeCrypto.generateKeyPairSync("ec", {
namedCurve: "prime256v1",
privateKeyEncoding: { type: "pkcs8", format: "pem" },
publicKeyEncoding: { type: "spki", format: "pem" },
});
const credentials = {
teamId: "team-1",
keyId: "key-1",
privateKey: Redacted.make(privateKey),
bundleId: "com.t3tools.test",
environment: "sandbox",
} satisfies ApnsCredentials;
const failingHttpClient = HttpClient.make((request) =>
Effect.fail(
new HttpClientError.HttpClientError({
reason: new HttpClientError.TransportError({ request, cause: httpCause }),
}),
),
);
const layer = ApnsClient.layer.pipe(
Layer.provide(Layer.succeed(HttpClient.HttpClient, failingHttpClient)),
);

return Effect.gen(function* () {
const apns = yield* ApnsClient.ApnsClient;
const request = apns.makePushNotificationRequest({
token: "long-push-token",
notification: {
title: "Thread",
body: "Input: Project",
environmentId: "env",
threadId: "thread",
deepLink: "/threads/env/thread",
},
});
const error = yield* Effect.flip(
apns.sendPushNotificationRequest({
credentials,
request,
issuedAtUnixSeconds: 123,
}),
);

expect(isApnsHttpRequestError(error)).toBe(true);
if (!isApnsHttpRequestError(error)) {
return yield* Effect.die("expected APNs HTTP request error");
}
expect(error).toMatchObject({
requestKind: "push-notification",
event: null,
environment: "sandbox",
bundleId: "com.t3tools.test",
tokenSuffix: "sh-token",
stage: "send",
status: null,
message: "APNs push-notification request failed during send in sandbox.",
});
expect(error.cause).toBeInstanceOf(HttpClientError.HttpClientError);
expect((error.cause as HttpClientError.HttpClientError).reason).toMatchObject({
_tag: "TransportError",
cause: httpCause,
});
}).pipe(Effect.provide(layer));
});
});
143 changes: 113 additions & 30 deletions infra/relay/src/agentActivity/ApnsClient.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,14 +11,17 @@ import * as Schema from "effect/Schema";
import * as Headers from "effect/unstable/http/Headers";
import * as HttpClient from "effect/unstable/http/HttpClient";
import * as HttpClientRequest from "effect/unstable/http/HttpClientRequest";
import type * as RelayConfiguration from "../Config.ts";
import { ApnsEnvironment as ApnsEnvironmentSchema, type ApnsCredentials } from "../Config.ts";
import type { ApnsNotificationPayload } from "./apnsDeliveryJobs.ts";

const LIVE_ACTIVITY_NAME = "AgentActivity";
const STALE_AFTER_SECONDS = 2 * 60;
const DISMISS_AFTER_SECONDS = 5 * 60;

export type ApnsLiveActivityEvent = "start" | "update" | "end";
const ApnsLiveActivityEventSchema = Schema.Literals(["start", "update", "end"]);
export type ApnsLiveActivityEvent = typeof ApnsLiveActivityEventSchema.Type;

const ApnsRequestKindSchema = Schema.Literals(["live-activity", "push-notification"]);

interface ApnsLiveActivityRequest {
readonly token: string;
Expand All @@ -42,49 +45,55 @@ export interface ApnsDeliveryResult {

export class ApnsJwtEncodingError extends Schema.TaggedErrorClass<ApnsJwtEncodingError>()(
"ApnsJwtEncodingError",
{ cause: Schema.Defect() },
{
component: Schema.Literals(["header", "payload"]),
teamId: Schema.String,
keyId: Schema.String,
issuedAtUnixSeconds: Schema.Number,
cause: Schema.Defect(),
},
) {
override get message(): string {
return "Failed to encode APNs JWT.";
return `Failed to encode APNs JWT ${this.component} for key ${this.keyId}.`;
}
}

export class ApnsJwtSigningError extends Schema.TaggedErrorClass<ApnsJwtSigningError>()(
"ApnsJwtSigningError",
{ cause: Schema.Defect() },
) {
override get message(): string {
return "Failed to sign APNs JWT.";
}
}

export class ApnsHttpRequestError extends Schema.TaggedErrorClass<ApnsHttpRequestError>()(
"ApnsHttpRequestError",
{
teamId: Schema.String,
keyId: Schema.String,
issuedAtUnixSeconds: Schema.Number,
cause: Schema.Defect(),
},
) {
override get message(): string {
return "APNs HTTP request failed.";
return `Failed to sign APNs JWT for key ${this.keyId}.`;
}
}

export class ApnsInvalidResponseError extends Schema.TaggedErrorClass<ApnsInvalidResponseError>()(
"ApnsInvalidResponseError",
export class ApnsHttpRequestError extends Schema.TaggedErrorClass<ApnsHttpRequestError>()(
"ApnsHttpRequestError",
{
requestKind: ApnsRequestKindSchema,
event: Schema.NullOr(ApnsLiveActivityEventSchema),
environment: ApnsEnvironmentSchema,
bundleId: Schema.String,
tokenSuffix: Schema.String,
stage: Schema.Literals(["send", "read-response"]),
status: Schema.NullOr(Schema.Number),
cause: Schema.Defect(),
},
) {
override get message(): string {
return "APNs returned an invalid response.";
return `APNs ${this.requestKind} request failed during ${this.stage} in ${this.environment}.`;
}
}

export const ApnsError = Schema.Union([
ApnsJwtEncodingError,
ApnsJwtSigningError,
ApnsHttpRequestError,
ApnsInvalidResponseError,
]);
export type ApnsError = typeof ApnsError.Type;

Expand Down Expand Up @@ -113,18 +122,38 @@ const encodeApnsJwtPayloadJson = Schema.encodeEffect(
);

const makeApnsJwt = Effect.fn("relay.apns.make_jwt")(function* (input: {
readonly teamId: RelayConfiguration.ApnsCredentials["teamId"];
readonly keyId: RelayConfiguration.ApnsCredentials["keyId"];
readonly privateKey: RelayConfiguration.ApnsCredentials["privateKey"];
readonly teamId: ApnsCredentials["teamId"];
readonly keyId: ApnsCredentials["keyId"];
readonly privateKey: ApnsCredentials["privateKey"];
readonly issuedAtUnixSeconds: number;
}) {
const headerJson = yield* encodeApnsJwtHeaderJson({ alg: "ES256", kid: input.keyId }).pipe(
Effect.mapError((cause) => new ApnsJwtEncodingError({ cause })),
Effect.mapError(
(cause) =>
new ApnsJwtEncodingError({
component: "header",
teamId: input.teamId,
keyId: input.keyId,
issuedAtUnixSeconds: input.issuedAtUnixSeconds,
cause,
}),
),
);
const payloadJson = yield* encodeApnsJwtPayloadJson({
iss: input.teamId,
iat: input.issuedAtUnixSeconds,
}).pipe(Effect.mapError((cause) => new ApnsJwtEncodingError({ cause })));
}).pipe(
Effect.mapError(
(cause) =>
new ApnsJwtEncodingError({
component: "payload",
teamId: input.teamId,
keyId: input.keyId,
issuedAtUnixSeconds: input.issuedAtUnixSeconds,
cause,
}),
),
);

const privateKey = Redacted.value(input.privateKey);
const header = Encoding.encodeBase64Url(headerJson);
Expand All @@ -141,7 +170,13 @@ const makeApnsJwt = Effect.fn("relay.apns.make_jwt")(function* (input: {
});
return `${signingInput}.${Encoding.encodeBase64Url(signature)}`;
},
catch: (cause) => new ApnsJwtSigningError({ cause }),
catch: (cause) =>
new ApnsJwtSigningError({
teamId: input.teamId,
keyId: input.keyId,
issuedAtUnixSeconds: input.issuedAtUnixSeconds,
cause,
}),
});
});

Expand Down Expand Up @@ -251,12 +286,12 @@ export class ApnsClient extends Context.Service<
readonly makeLiveActivityRequest: typeof makeLiveActivityRequest;
readonly makePushNotificationRequest: typeof makePushNotificationRequest;
readonly sendLiveActivityRequest: (input: {
readonly credentials: RelayConfiguration.ApnsCredentials;
readonly credentials: ApnsCredentials;
readonly request: ApnsLiveActivityRequest;
readonly issuedAtUnixSeconds: number;
}) => Effect.Effect<ApnsDeliveryResult, ApnsError>;
readonly sendPushNotificationRequest: (input: {
readonly credentials: RelayConfiguration.ApnsCredentials;
readonly credentials: ApnsCredentials;
readonly request: ApnsPushNotificationRequest;
readonly issuedAtUnixSeconds: number;
}) => Effect.Effect<ApnsDeliveryResult, ApnsError>;
Expand Down Expand Up @@ -287,10 +322,34 @@ export const make = Effect.gen(function* () {
}),
HttpClientRequest.bodyJson(input.request.payload),
Effect.flatMap(httpClient.execute),
Effect.mapError((cause) => new ApnsHttpRequestError({ cause })),
Effect.mapError(
(cause) =>
new ApnsHttpRequestError({
requestKind: "live-activity",
event: input.request.event,
environment: input.credentials.environment,
bundleId: input.credentials.bundleId,
tokenSuffix: input.request.token.slice(-8),
stage: "send",
status: null,
cause,
}),
),
);
const responseText = yield* response.text.pipe(
Effect.mapError((cause) => new ApnsHttpRequestError({ cause })),
Effect.mapError(
(cause) =>
new ApnsHttpRequestError({
requestKind: "live-activity",
event: input.request.event,
environment: input.credentials.environment,
bundleId: input.credentials.bundleId,
tokenSuffix: input.request.token.slice(-8),
stage: "read-response",
status: response.status,
cause,
}),
),
);
const reason = apnsReasonFromBody(responseText);
return {
Expand Down Expand Up @@ -323,10 +382,34 @@ export const make = Effect.gen(function* () {
}),
HttpClientRequest.bodyJson(input.request.payload),
Effect.flatMap(httpClient.execute),
Effect.mapError((cause) => new ApnsHttpRequestError({ cause })),
Effect.mapError(
(cause) =>
new ApnsHttpRequestError({
requestKind: "push-notification",
event: null,
environment: input.credentials.environment,
bundleId: input.credentials.bundleId,
tokenSuffix: input.request.token.slice(-8),
stage: "send",
status: null,
cause,
}),
),
);
const responseText = yield* response.text.pipe(
Effect.mapError((cause) => new ApnsHttpRequestError({ cause })),
Effect.mapError(
(cause) =>
new ApnsHttpRequestError({
requestKind: "push-notification",
event: null,
environment: input.credentials.environment,
bundleId: input.credentials.bundleId,
tokenSuffix: input.request.token.slice(-8),
stage: "read-response",
status: response.status,
cause,
}),
),
);
const reason = apnsReasonFromBody(responseText);
return {
Expand Down
Loading