mirror of
https://github.com/fluxerapp/fluxer
synced 2026-10-07 19:22:14 +09:00
fix(api): handle bad manifests, unfurl errors and the apns key (#2508)
This commit is contained in:
@@ -28,6 +28,7 @@ import {
|
|||||||
getUserRepository,
|
getUserRepository,
|
||||||
} from '../middleware/ServiceSingletons';
|
} from '../middleware/ServiceSingletons';
|
||||||
import {torExitListCache} from '../middleware/TorExitListCache';
|
import {torExitListCache} from '../middleware/TorExitListCache';
|
||||||
|
import {ensureApnsSigningKey} from '../push/ApnsPushService';
|
||||||
import {initializeSearch, shutdownSearch} from '../SearchFactory';
|
import {initializeSearch, shutdownSearch} from '../SearchFactory';
|
||||||
import {warmupAdminSearchIndexes} from '../search/SearchWarmup';
|
import {warmupAdminSearchIndexes} from '../search/SearchWarmup';
|
||||||
import {VisionarySlotInitializer} from '../stripe/VisionarySlotInitializer';
|
import {VisionarySlotInitializer} from '../stripe/VisionarySlotInitializer';
|
||||||
@@ -46,6 +47,7 @@ export function createInitializer(config: APIConfig, logger: ILogger): () => Pro
|
|||||||
return async (): Promise<void> => {
|
return async (): Promise<void> => {
|
||||||
try {
|
try {
|
||||||
logger.info('Initializing API service...');
|
logger.info('Initializing API service...');
|
||||||
|
await ensureApnsSigningKey();
|
||||||
const geoipStartupResult = await ensureGeoipDatabaseOnStartup({
|
const geoipStartupResult = await ensureGeoipDatabaseOnStartup({
|
||||||
geoip: config.geoip,
|
geoip: config.geoip,
|
||||||
s3Config: {
|
s3Config: {
|
||||||
|
|||||||
@@ -17,7 +17,7 @@ import {
|
|||||||
StorageObjectRangeNotSatisfiableError,
|
StorageObjectRangeNotSatisfiableError,
|
||||||
} from '../infrastructure/IStorageService';
|
} from '../infrastructure/IStorageService';
|
||||||
import {Logger} from '../Logger';
|
import {Logger} from '../Logger';
|
||||||
import {isJsonRecord, parseJsonUnknown} from '../utils/JsonBoundaryUtils';
|
import {isJsonRecord, parseJsonRecord, parseJsonUnknown} from '../utils/JsonBoundaryUtils';
|
||||||
import {
|
import {
|
||||||
parseDesktopArtifactScope,
|
parseDesktopArtifactScope,
|
||||||
parseDesktopReleaseDescriptor,
|
parseDesktopReleaseDescriptor,
|
||||||
@@ -731,7 +731,7 @@ export class DownloadService {
|
|||||||
|
|
||||||
private async readJsonObjectFromStorage(key: string): Promise<unknown | null> {
|
private async readJsonObjectFromStorage(key: string): Promise<unknown | null> {
|
||||||
const text = await this.readTextFromStorage(key);
|
const text = await this.readTextFromStorage(key);
|
||||||
return text == null ? null : parseJsonUnknown(text);
|
return text == null ? null : parseJsonRecord(text);
|
||||||
}
|
}
|
||||||
|
|
||||||
private async readTextFromStorage(key: string): Promise<string | null> {
|
private async readTextFromStorage(key: string): Promise<string | null> {
|
||||||
|
|||||||
@@ -0,0 +1,71 @@
|
|||||||
|
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||||
|
|
||||||
|
import {Readable} from 'node:stream';
|
||||||
|
import {describe, expect, it} from 'vitest';
|
||||||
|
import type {IStorageService} from '../../infrastructure/IStorageService';
|
||||||
|
import {DownloadService} from '../DownloadService';
|
||||||
|
|
||||||
|
const PREFIX = 'desktop/stable/darwin/x64';
|
||||||
|
const MANIFEST_KEY = `${PREFIX}/manifest.json`;
|
||||||
|
const LISTED_FILENAME = 'Fluxer-1.2.3-mac-universal.dmg';
|
||||||
|
const MANIFEST_FILENAME = 'Fluxer-1.3.0-mac-universal.dmg';
|
||||||
|
|
||||||
|
const LATEST_PARAMS = {
|
||||||
|
channel: 'stable',
|
||||||
|
plat: 'darwin',
|
||||||
|
arch: 'x64',
|
||||||
|
format: 'dmg',
|
||||||
|
} as const;
|
||||||
|
|
||||||
|
function createService(overrides: {manifestBody?: string | null; objectKeys?: Array<string>} = {}) {
|
||||||
|
const objectKeys = overrides.objectKeys ?? [`${PREFIX}/${LISTED_FILENAME}`];
|
||||||
|
const storageService = {
|
||||||
|
streamObject: async (params: {key: string}) => {
|
||||||
|
if (params.key !== MANIFEST_KEY) {
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
const body = overrides.manifestBody;
|
||||||
|
if (body == null) {
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
const buffer = Buffer.from(body, 'utf8');
|
||||||
|
return {body: Readable.from([buffer]), contentLength: buffer.byteLength};
|
||||||
|
},
|
||||||
|
listObjects: async () => objectKeys.map((key) => ({key})),
|
||||||
|
getObjectMetadata: async (_bucket: string, key: string) =>
|
||||||
|
objectKeys.includes(key) ? {contentLength: 1, contentType: 'application/x-apple-diskimage'} : null,
|
||||||
|
} as unknown as IStorageService;
|
||||||
|
return new DownloadService(storageService);
|
||||||
|
}
|
||||||
|
|
||||||
|
describe('desktop manifest parsing', () => {
|
||||||
|
it('falls back to the object listing when the manifest is not valid JSON', async () => {
|
||||||
|
const service = createService({manifestBody: '{not json'});
|
||||||
|
await expect(service.resolveLatestDesktopKey({...LATEST_PARAMS})).resolves.toBe(`${PREFIX}/${LISTED_FILENAME}`);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('returns null rather than throwing when the manifest is malformed and no artifact is listed', async () => {
|
||||||
|
const service = createService({manifestBody: '{not json', objectKeys: []});
|
||||||
|
await expect(service.resolveLatestDesktopKey({...LATEST_PARAMS})).resolves.toBeNull();
|
||||||
|
});
|
||||||
|
|
||||||
|
it('falls back to the object listing when the manifest parses to an array', async () => {
|
||||||
|
const service = createService({manifestBody: '[]'});
|
||||||
|
await expect(service.resolveLatestDesktopKey({...LATEST_PARAMS})).resolves.toBe(`${PREFIX}/${LISTED_FILENAME}`);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('still resolves through a well-formed manifest', async () => {
|
||||||
|
const service = createService({
|
||||||
|
manifestBody: JSON.stringify({
|
||||||
|
channel: 'stable',
|
||||||
|
platform: 'darwin',
|
||||||
|
arch: 'x64',
|
||||||
|
version: '1.3.0',
|
||||||
|
pub_date: '2026-08-17T00:00:00Z',
|
||||||
|
files: {dmg: MANIFEST_FILENAME},
|
||||||
|
}),
|
||||||
|
objectKeys: [`${PREFIX}/${LISTED_FILENAME}`, `${PREFIX}/${MANIFEST_FILENAME}`],
|
||||||
|
});
|
||||||
|
await expect(service.resolveLatestDesktopKey({...LATEST_PARAMS})).resolves.toBe(`${PREFIX}/${MANIFEST_FILENAME}`);
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -1,5 +1,6 @@
|
|||||||
// SPDX-License-Identifier: AGPL-3.0-or-later
|
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||||
|
|
||||||
|
import {BadGatewayError} from '@fluxer/errors/src/domains/core/BadGatewayError';
|
||||||
import type {GifResponse} from '@fluxer/schema/src/domains/gif/GifSchemas';
|
import type {GifResponse} from '@fluxer/schema/src/domains/gif/GifSchemas';
|
||||||
import {describe, expect, it, vi} from 'vitest';
|
import {describe, expect, it, vi} from 'vitest';
|
||||||
import {GifService} from '../gif/GifService';
|
import {GifService} from '../gif/GifService';
|
||||||
@@ -8,7 +9,7 @@ import type {IMediaService} from '../infrastructure/IMediaService';
|
|||||||
import type {IUnfurlerService} from '../infrastructure/IUnfurlerService';
|
import type {IUnfurlerService} from '../infrastructure/IUnfurlerService';
|
||||||
import {resolveFavoriteGifEntry} from './FavoriteGifResolver';
|
import {resolveFavoriteGifEntry} from './FavoriteGifResolver';
|
||||||
|
|
||||||
function createProvider(gif: GifResponse): IGifProvider {
|
function createProvider(gif: GifResponse, slug: string | null = gif.slug): IGifProvider {
|
||||||
return {
|
return {
|
||||||
meta: {
|
meta: {
|
||||||
name: 'klipy',
|
name: 'klipy',
|
||||||
@@ -23,7 +24,7 @@ function createProvider(gif: GifResponse): IGifProvider {
|
|||||||
suggest: async () => [],
|
suggest: async () => [],
|
||||||
resolveByUrl: async () => gif,
|
resolveByUrl: async () => gif,
|
||||||
buildShareUrl: (slug) => `https://klipy.com/gifs/${slug}`,
|
buildShareUrl: (slug) => `https://klipy.com/gifs/${slug}`,
|
||||||
extractSlugFromUrl: () => gif.slug,
|
extractSlugFromUrl: () => slug,
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -84,4 +85,78 @@ describe('resolveFavoriteGifEntry', () => {
|
|||||||
expect(mediaService.getMetadata).not.toHaveBeenCalled();
|
expect(mediaService.getMetadata).not.toHaveBeenCalled();
|
||||||
expect(unfurlerService.unfurl).not.toHaveBeenCalled();
|
expect(unfurlerService.unfurl).not.toHaveBeenCalled();
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it('degrades only the failing URL when its unfurl stage rejects', async () => {
|
||||||
|
const gif: GifResponse = {
|
||||||
|
id: 'unused',
|
||||||
|
slug: 'unused',
|
||||||
|
provider: 'klipy',
|
||||||
|
title: 'Unused',
|
||||||
|
url: 'https://klipy.com/gifs/unused',
|
||||||
|
src: 'https://static.klipy.example/unused.gif',
|
||||||
|
proxy_src: 'https://media.example/unused.gif',
|
||||||
|
width: 1,
|
||||||
|
height: 1,
|
||||||
|
media: {},
|
||||||
|
placeholder: null,
|
||||||
|
};
|
||||||
|
const brokenUrl = 'https://example.com/broken';
|
||||||
|
const workingUrl = 'https://example.com/working';
|
||||||
|
const mediaService = {
|
||||||
|
getMetadata: vi.fn(async () => null),
|
||||||
|
getExternalMediaProxyURL: vi.fn((url: string) => `https://media.example/proxy?url=${encodeURIComponent(url)}`),
|
||||||
|
} as unknown as IMediaService;
|
||||||
|
const unfurlerService = {
|
||||||
|
unfurl: vi.fn(async (url: string) => {
|
||||||
|
if (url === brokenUrl) {
|
||||||
|
throw new BadGatewayError({message: '[nats-unfurl] service error: upstream exploded'});
|
||||||
|
}
|
||||||
|
return [
|
||||||
|
{
|
||||||
|
image: {
|
||||||
|
url: 'https://cdn.example/working.gif',
|
||||||
|
proxy_url: 'https://media.example/working.gif',
|
||||||
|
content_type: 'image/gif',
|
||||||
|
width: 200,
|
||||||
|
height: 100,
|
||||||
|
placeholder: null,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
];
|
||||||
|
}),
|
||||||
|
} as unknown as IUnfurlerService;
|
||||||
|
const gifService = new GifService(createProvider(gif, null));
|
||||||
|
|
||||||
|
const entries = await Promise.all(
|
||||||
|
[brokenUrl, workingUrl].map((url) =>
|
||||||
|
resolveFavoriteGifEntry({url, locale: 'en-US', country: 'US', gifService, mediaService, unfurlerService}),
|
||||||
|
),
|
||||||
|
);
|
||||||
|
|
||||||
|
expect(entries[0]).toEqual({
|
||||||
|
url: brokenUrl,
|
||||||
|
proxy_url: `https://media.example/proxy?url=${encodeURIComponent(brokenUrl)}`,
|
||||||
|
width: 0,
|
||||||
|
height: 0,
|
||||||
|
media: {},
|
||||||
|
content_type: '',
|
||||||
|
placeholder: null,
|
||||||
|
});
|
||||||
|
expect(entries[1]).toEqual({
|
||||||
|
url: workingUrl,
|
||||||
|
proxy_url: 'https://media.example/working.gif',
|
||||||
|
width: 200,
|
||||||
|
height: 100,
|
||||||
|
media: {
|
||||||
|
gif: {
|
||||||
|
src: 'https://cdn.example/working.gif',
|
||||||
|
proxy_src: 'https://media.example/working.gif',
|
||||||
|
width: 200,
|
||||||
|
height: 100,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
content_type: 'image/gif',
|
||||||
|
placeholder: null,
|
||||||
|
});
|
||||||
|
});
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -4,7 +4,7 @@ import {Logger} from '@fluxer/logger/src/Logger';
|
|||||||
import type {ResolvedGifEntrySchema} from '@fluxer/schema/src/domains/gif/FavoriteGifSchemas';
|
import type {ResolvedGifEntrySchema} from '@fluxer/schema/src/domains/gif/FavoriteGifSchemas';
|
||||||
import {inferFormatContentType, PREVIEW_FORMAT_PRIORITY} from '@fluxer/schema/src/domains/gif/GifMediaFormatKeys';
|
import {inferFormatContentType, PREVIEW_FORMAT_PRIORITY} from '@fluxer/schema/src/domains/gif/GifMediaFormatKeys';
|
||||||
import type {GifMediaFormat, GifResponse} from '@fluxer/schema/src/domains/gif/GifSchemas';
|
import type {GifMediaFormat, GifResponse} from '@fluxer/schema/src/domains/gif/GifSchemas';
|
||||||
import type {EmbedMediaResponse} from '@fluxer/schema/src/domains/message/EmbedSchemas';
|
import type {EmbedMediaResponse, MessageEmbedResponse} from '@fluxer/schema/src/domains/message/EmbedSchemas';
|
||||||
import {tryExtractGifProviderSlug} from '../gif/GifProviderUtils';
|
import {tryExtractGifProviderSlug} from '../gif/GifProviderUtils';
|
||||||
import type {GifService} from '../gif/GifService';
|
import type {GifService} from '../gif/GifService';
|
||||||
import type {IGifProvider} from '../gif/IGifProvider';
|
import type {IGifProvider} from '../gif/IGifProvider';
|
||||||
@@ -130,7 +130,10 @@ async function resolveUnfurledFavoriteGifEntry({
|
|||||||
unfurlerService: IUnfurlerService;
|
unfurlerService: IUnfurlerService;
|
||||||
mediaService: IMediaService;
|
mediaService: IMediaService;
|
||||||
}): Promise<ResolvedGifEntrySchema | null> {
|
}): Promise<ResolvedGifEntrySchema | null> {
|
||||||
const embeds = await unfurlerService.unfurl(url, 'allow');
|
const embeds = await unfurlerService.unfurl(url, 'allow').catch((error: unknown) => {
|
||||||
|
logger.warn({error, url}, 'Failed to unfurl favorite GIF URL');
|
||||||
|
return [] as Array<MessageEmbedResponse>;
|
||||||
|
});
|
||||||
for (const embed of embeds) {
|
for (const embed of embeds) {
|
||||||
const media = [embed.video, embed.image, embed.thumbnail].find((candidate) =>
|
const media = [embed.video, embed.image, embed.thumbnail].find((candidate) =>
|
||||||
isRenderableMediaType(candidate?.content_type),
|
isRenderableMediaType(candidate?.content_type),
|
||||||
|
|||||||
@@ -1,5 +1,8 @@
|
|||||||
// SPDX-License-Identifier: AGPL-3.0-or-later
|
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||||
|
|
||||||
|
import {BadGatewayError} from '@fluxer/errors/src/domains/core/BadGatewayError';
|
||||||
|
import {GatewayTimeoutError} from '@fluxer/errors/src/domains/core/GatewayTimeoutError';
|
||||||
|
import {ServiceUnavailableError} from '@fluxer/errors/src/domains/core/ServiceUnavailableError';
|
||||||
import type {INatsConnectionManager} from '@pkgs/nats/src/INatsConnectionManager';
|
import type {INatsConnectionManager} from '@pkgs/nats/src/INatsConnectionManager';
|
||||||
import {type NatsConnection, StringCodec} from 'nats';
|
import {type NatsConnection, StringCodec} from 'nats';
|
||||||
import {describe, expect, it} from 'vitest';
|
import {describe, expect, it} from 'vitest';
|
||||||
@@ -11,12 +14,23 @@ interface FakeRequest {
|
|||||||
timeout: number | undefined;
|
timeout: number | undefined;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const RESOLVED_REPLY = JSON.stringify({Resolved: {embeds: [], cache_ttl_seconds: null}});
|
||||||
|
|
||||||
|
function natsErrorWithCode(code: string): Error {
|
||||||
|
return Object.assign(new Error('nats request failed'), {code});
|
||||||
|
}
|
||||||
|
|
||||||
class FakeNatsConnectionManager implements INatsConnectionManager {
|
class FakeNatsConnectionManager implements INatsConnectionManager {
|
||||||
private readonly codec = StringCodec();
|
private readonly codec = StringCodec();
|
||||||
private closed = true;
|
private closed = true;
|
||||||
readonly requests: Array<FakeRequest> = [];
|
readonly requests: Array<FakeRequest> = [];
|
||||||
connectCalls = 0;
|
connectCalls = 0;
|
||||||
|
|
||||||
|
constructor(
|
||||||
|
private readonly replyText: string = RESOLVED_REPLY,
|
||||||
|
private readonly requestError: Error | null = null,
|
||||||
|
) {}
|
||||||
|
|
||||||
async connect(): Promise<void> {
|
async connect(): Promise<void> {
|
||||||
this.connectCalls += 1;
|
this.connectCalls += 1;
|
||||||
this.closed = false;
|
this.closed = false;
|
||||||
@@ -33,8 +47,11 @@ class FakeNatsConnectionManager implements INatsConnectionManager {
|
|||||||
body: JSON.parse(this.codec.decode(data)) as Record<string, unknown>,
|
body: JSON.parse(this.codec.decode(data)) as Record<string, unknown>,
|
||||||
timeout: options?.timeout,
|
timeout: options?.timeout,
|
||||||
});
|
});
|
||||||
|
if (this.requestError) {
|
||||||
|
throw this.requestError;
|
||||||
|
}
|
||||||
return {
|
return {
|
||||||
data: this.codec.encode(JSON.stringify({Resolved: {embeds: [], cache_ttl_seconds: null}})),
|
data: this.codec.encode(this.replyText),
|
||||||
};
|
};
|
||||||
},
|
},
|
||||||
} as unknown as NatsConnection;
|
} as unknown as NatsConnection;
|
||||||
@@ -79,4 +96,39 @@ describe('NatsUnfurlerService', () => {
|
|||||||
expect(manager.requests[0]?.body).not.toHaveProperty('media_proxy_endpoint');
|
expect(manager.requests[0]?.body).not.toHaveProperty('media_proxy_endpoint');
|
||||||
expect(manager.requests[0]?.body).not.toHaveProperty('media_proxy_secret_key');
|
expect(manager.requests[0]?.body).not.toHaveProperty('media_proxy_secret_key');
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it('rejects with a bad gateway error when the unfurl service reports a failure', async () => {
|
||||||
|
const manager = new FakeNatsConnectionManager(JSON.stringify({Failed: {message: 'upstream exploded'}}));
|
||||||
|
const service = new NatsUnfurlerService(manager);
|
||||||
|
|
||||||
|
await expect(service.unfurlWithCachePolicy('https://example.com')).rejects.toBeInstanceOf(BadGatewayError);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('rejects with a bad gateway error when the reply payload is unreadable', async () => {
|
||||||
|
const manager = new FakeNatsConnectionManager('not json');
|
||||||
|
const service = new NatsUnfurlerService(manager);
|
||||||
|
|
||||||
|
await expect(service.unfurlWithCachePolicy('https://example.com')).rejects.toBeInstanceOf(BadGatewayError);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('rejects with a gateway timeout error when the request times out', async () => {
|
||||||
|
const manager = new FakeNatsConnectionManager(RESOLVED_REPLY, natsErrorWithCode('TIMEOUT'));
|
||||||
|
const service = new NatsUnfurlerService(manager);
|
||||||
|
|
||||||
|
await expect(service.unfurlWithCachePolicy('https://example.com')).rejects.toBeInstanceOf(GatewayTimeoutError);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('rejects with a service unavailable error when no responders answer', async () => {
|
||||||
|
const manager = new FakeNatsConnectionManager(RESOLVED_REPLY, natsErrorWithCode('503'));
|
||||||
|
const service = new NatsUnfurlerService(manager);
|
||||||
|
|
||||||
|
await expect(service.unfurlWithCachePolicy('https://example.com')).rejects.toBeInstanceOf(ServiceUnavailableError);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('rejects with a service unavailable error when the shard rejects the request', async () => {
|
||||||
|
const manager = new FakeNatsConnectionManager(JSON.stringify({error: 'overloaded'}));
|
||||||
|
const service = new NatsUnfurlerService(manager);
|
||||||
|
|
||||||
|
await expect(service.unfurlWithCachePolicy('https://example.com')).rejects.toBeInstanceOf(ServiceUnavailableError);
|
||||||
|
});
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -1,5 +1,9 @@
|
|||||||
// SPDX-License-Identifier: AGPL-3.0-or-later
|
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||||
|
|
||||||
|
import {BadGatewayError} from '@fluxer/errors/src/domains/core/BadGatewayError';
|
||||||
|
import {GatewayTimeoutError} from '@fluxer/errors/src/domains/core/GatewayTimeoutError';
|
||||||
|
import {ServiceUnavailableError} from '@fluxer/errors/src/domains/core/ServiceUnavailableError';
|
||||||
|
import {FluxerError} from '@fluxer/errors/src/FluxerError';
|
||||||
import type {MessageEmbedResponse} from '@fluxer/schema/src/domains/message/EmbedSchemas';
|
import type {MessageEmbedResponse} from '@fluxer/schema/src/domains/message/EmbedSchemas';
|
||||||
import type {INatsConnectionManager} from '@pkgs/nats/src/INatsConnectionManager';
|
import type {INatsConnectionManager} from '@pkgs/nats/src/INatsConnectionManager';
|
||||||
import {StringCodec} from 'nats';
|
import {StringCodec} from 'nats';
|
||||||
@@ -12,6 +16,8 @@ import {throwForSvcErrorReply} from './SvcErrorReply';
|
|||||||
const NATS_UNFURL_SUBJECT = 'svc.unfurl';
|
const NATS_UNFURL_SUBJECT = 'svc.unfurl';
|
||||||
const NATS_UNFURL_TIMEOUT_MS = 12000;
|
const NATS_UNFURL_TIMEOUT_MS = 12000;
|
||||||
const NATS_UNFURL_CACHE_ONLY_TIMEOUT_MS = 1000;
|
const NATS_UNFURL_CACHE_ONLY_TIMEOUT_MS = 1000;
|
||||||
|
const NATS_NO_RESPONDERS_CODE = '503';
|
||||||
|
const NATS_TIMEOUT_CODE = 'TIMEOUT';
|
||||||
|
|
||||||
interface NatsUnfurlRequest {
|
interface NatsUnfurlRequest {
|
||||||
op: 'Unfurl';
|
op: 'Unfurl';
|
||||||
@@ -54,6 +60,23 @@ function isNatsUnfurlResponse(value: unknown): value is NatsUnfurlResponse {
|
|||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function mapUnfurlTransportError(error: unknown): Error {
|
||||||
|
if (error instanceof FluxerError) {
|
||||||
|
return error;
|
||||||
|
}
|
||||||
|
if (!(error instanceof Error)) {
|
||||||
|
return new BadGatewayError({message: '[nats-unfurl] request failed'});
|
||||||
|
}
|
||||||
|
const code = 'code' in error && typeof error.code === 'string' ? error.code : null;
|
||||||
|
if (code === NATS_NO_RESPONDERS_CODE || error.message === 'NO_RESPONDERS' || error.name === 'NoRespondersError') {
|
||||||
|
return new ServiceUnavailableError({message: '[nats-unfurl] no unfurl service is answering'});
|
||||||
|
}
|
||||||
|
if (code === NATS_TIMEOUT_CODE || error.message === 'TIMEOUT' || error.name === 'TimeoutError') {
|
||||||
|
return new GatewayTimeoutError({message: '[nats-unfurl] unfurl service did not answer in time'});
|
||||||
|
}
|
||||||
|
return new BadGatewayError({message: '[nats-unfurl] request failed'});
|
||||||
|
}
|
||||||
|
|
||||||
export class NatsUnfurlerService extends IUnfurlerService {
|
export class NatsUnfurlerService extends IUnfurlerService {
|
||||||
private readonly connectionManager: INatsConnectionManager;
|
private readonly connectionManager: INatsConnectionManager;
|
||||||
private readonly codec = StringCodec();
|
private readonly codec = StringCodec();
|
||||||
@@ -94,7 +117,7 @@ export class NatsUnfurlerService extends IUnfurlerService {
|
|||||||
const response = parseJsonWithGuard(responseText, isNatsUnfurlResponse);
|
const response = parseJsonWithGuard(responseText, isNatsUnfurlResponse);
|
||||||
if (!response) {
|
if (!response) {
|
||||||
throwForSvcErrorReply('nats-unfurl', parseJsonRecord(responseText));
|
throwForSvcErrorReply('nats-unfurl', parseJsonRecord(responseText));
|
||||||
throw new Error(`[nats-unfurl] invalid response payload: ${responseText}`);
|
throw new BadGatewayError({message: '[nats-unfurl] invalid response payload'});
|
||||||
}
|
}
|
||||||
if ('Resolved' in response) {
|
if ('Resolved' in response) {
|
||||||
return {
|
return {
|
||||||
@@ -103,16 +126,16 @@ export class NatsUnfurlerService extends IUnfurlerService {
|
|||||||
};
|
};
|
||||||
}
|
}
|
||||||
if ('Failed' in response) {
|
if ('Failed' in response) {
|
||||||
throw new Error(`[nats-unfurl] service error: ${response.Failed.message}`);
|
throw new BadGatewayError({message: `[nats-unfurl] service error: ${response.Failed.message}`});
|
||||||
}
|
}
|
||||||
throw new Error(`[nats-unfurl] unexpected response variant: ${responseText}`);
|
throw new BadGatewayError({message: '[nats-unfurl] unexpected response variant'});
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
if (options.signal?.aborted) {
|
if (options.signal?.aborted) {
|
||||||
Logger.warn({url}, '[nats-unfurl] request aborted');
|
Logger.warn({url}, '[nats-unfurl] request aborted');
|
||||||
} else {
|
} else {
|
||||||
Logger.error({error, url}, '[nats-unfurl] failed to unfurl URL');
|
Logger.error({error, url}, '[nats-unfurl] failed to unfurl URL');
|
||||||
}
|
}
|
||||||
throw error;
|
throw mapUnfurlTransportError(error);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -45,6 +45,7 @@ interface Http2Response {
|
|||||||
}
|
}
|
||||||
|
|
||||||
const providerTokenCache = new Map<string, CachedProviderToken>();
|
const providerTokenCache = new Map<string, CachedProviderToken>();
|
||||||
|
const apnsSigningKeys = new Map<string, ReturnType<typeof importPKCS8>>();
|
||||||
const apnsSessions = new Map<string, ClientHttp2Session>();
|
const apnsSessions = new Map<string, ClientHttp2Session>();
|
||||||
|
|
||||||
export async function sendApnsPush(params: SendApnsPushParams): Promise<SendApnsPushResult> {
|
export async function sendApnsPush(params: SendApnsPushParams): Promise<SendApnsPushResult> {
|
||||||
@@ -105,7 +106,7 @@ async function getProviderToken(): Promise<string | null> {
|
|||||||
if (cached && cached.expiresAtSeconds > nowSeconds) {
|
if (cached && cached.expiresAtSeconds > nowSeconds) {
|
||||||
return cached.token;
|
return cached.token;
|
||||||
}
|
}
|
||||||
const key = await importPKCS8(privateKey, 'ES256');
|
const key = await apnsSigningKey(privateKey);
|
||||||
const token = await new SignJWT({})
|
const token = await new SignJWT({})
|
||||||
.setProtectedHeader({alg: 'ES256', kid: cfg.keyId})
|
.setProtectedHeader({alg: 'ES256', kid: cfg.keyId})
|
||||||
.setIssuer(cfg.teamId)
|
.setIssuer(cfg.teamId)
|
||||||
@@ -125,6 +126,27 @@ async function resolveApnsPrivateKey(): Promise<string | null> {
|
|||||||
return normalizePem(pem);
|
return normalizePem(pem);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export async function ensureApnsSigningKey(): Promise<void> {
|
||||||
|
if (!Config.push.apns.enabled) return;
|
||||||
|
const privateKey = await resolveApnsPrivateKey();
|
||||||
|
if (!privateKey) return;
|
||||||
|
await apnsSigningKey(privateKey);
|
||||||
|
Logger.info('APNs signing key loaded');
|
||||||
|
}
|
||||||
|
|
||||||
|
async function apnsSigningKey(privateKey: string): Promise<Awaited<ReturnType<typeof importPKCS8>>> {
|
||||||
|
const cached = apnsSigningKeys.get(privateKey);
|
||||||
|
if (cached) return await cached;
|
||||||
|
const pending = importPKCS8(privateKey, 'ES256');
|
||||||
|
apnsSigningKeys.set(privateKey, pending);
|
||||||
|
try {
|
||||||
|
return await pending;
|
||||||
|
} catch (error) {
|
||||||
|
apnsSigningKeys.delete(privateKey);
|
||||||
|
throw error;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
function normalizePem(value: string): string {
|
function normalizePem(value: string): string {
|
||||||
return value.replaceAll('\\n', '\n');
|
return value.replaceAll('\\n', '\n');
|
||||||
}
|
}
|
||||||
@@ -316,4 +338,5 @@ export const ApnsPushServiceTestHooks = {
|
|||||||
buildApnsPayload,
|
buildApnsPayload,
|
||||||
buildApnsHeaders,
|
buildApnsHeaders,
|
||||||
isPermanentApnsFailure,
|
isPermanentApnsFailure,
|
||||||
|
apnsSigningKey,
|
||||||
};
|
};
|
||||||
|
|||||||
@@ -1,7 +1,17 @@
|
|||||||
// SPDX-License-Identifier: AGPL-3.0-or-later
|
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||||
|
|
||||||
|
import {generateKeyPairSync} from 'node:crypto';
|
||||||
import {describe, expect, it} from 'vitest';
|
import {describe, expect, it} from 'vitest';
|
||||||
import {ApnsPushServiceTestHooks} from '../ApnsPushService';
|
import {ApnsPushServiceTestHooks, ensureApnsSigningKey} from '../ApnsPushService';
|
||||||
|
|
||||||
|
function pkcs8Pem(der: Buffer): string {
|
||||||
|
const base64 = der.toString('base64');
|
||||||
|
const lines: Array<string> = [];
|
||||||
|
for (let index = 0; index < base64.length; index += 64) {
|
||||||
|
lines.push(base64.slice(index, index + 64));
|
||||||
|
}
|
||||||
|
return `-----BEGIN PRIVATE KEY-----\n${lines.join('\n')}\n-----END PRIVATE KEY-----\n`;
|
||||||
|
}
|
||||||
|
|
||||||
describe('ApnsPushService', () => {
|
describe('ApnsPushService', () => {
|
||||||
it('builds modern APNs alert payloads for chat messages', () => {
|
it('builds modern APNs alert payloads for chat messages', () => {
|
||||||
@@ -132,6 +142,21 @@ describe('ApnsPushService', () => {
|
|||||||
expect(payload.aps).not.toHaveProperty('mutable-content');
|
expect(payload.aps).not.toHaveProperty('mutable-content');
|
||||||
expect(payload.author_avatar_url).toBe('https://cdn.example/avatar.png');
|
expect(payload.author_avatar_url).toBe('https://cdn.example/avatar.png');
|
||||||
});
|
});
|
||||||
|
it('imports the APNs signing key once per PEM and rejects a truncated one every time', async () => {
|
||||||
|
const {privateKey} = generateKeyPairSync('ec', {namedCurve: 'prime256v1'});
|
||||||
|
const der = privateKey.export({type: 'pkcs8', format: 'der'});
|
||||||
|
const pem = privateKey.export({type: 'pkcs8', format: 'pem'}).toString();
|
||||||
|
const truncated = pkcs8Pem(der.subarray(0, der.length - 8));
|
||||||
|
const first = await ApnsPushServiceTestHooks.apnsSigningKey(pem);
|
||||||
|
expect(await ApnsPushServiceTestHooks.apnsSigningKey(pem)).toBe(first);
|
||||||
|
await expect(ApnsPushServiceTestHooks.apnsSigningKey(truncated)).rejects.toThrow();
|
||||||
|
await expect(ApnsPushServiceTestHooks.apnsSigningKey(truncated)).rejects.toThrow();
|
||||||
|
});
|
||||||
|
|
||||||
|
it('loads no APNs signing key at startup while APNs is disabled', async () => {
|
||||||
|
await expect(ensureApnsSigningKey()).resolves.toBeUndefined();
|
||||||
|
});
|
||||||
|
|
||||||
it('marks only permanent APNs token failures as subscription deletion signals', () => {
|
it('marks only permanent APNs token failures as subscription deletion signals', () => {
|
||||||
expect(ApnsPushServiceTestHooks.isPermanentApnsFailure(410, 'Unregistered')).toBe(true);
|
expect(ApnsPushServiceTestHooks.isPermanentApnsFailure(410, 'Unregistered')).toBe(true);
|
||||||
expect(ApnsPushServiceTestHooks.isPermanentApnsFailure(400, 'BadDeviceToken')).toBe(true);
|
expect(ApnsPushServiceTestHooks.isPermanentApnsFailure(400, 'BadDeviceToken')).toBe(true);
|
||||||
|
|||||||
Reference in New Issue
Block a user