diff --git a/fluxer_api/src/api/Config.ts b/fluxer_api/src/api/Config.ts index decff2619..79f0cb344 100644 --- a/fluxer_api/src/api/Config.ts +++ b/fluxer_api/src/api/Config.ts @@ -237,6 +237,11 @@ export function buildAPIConfigFromMaster(master: MasterConfig): APIConfig { jetStreamUrl: master.services.nats?.jetstream_url ?? 'nats://127.0.0.1:4223', authToken: master.services.nats?.auth_token ?? '', }, + storageChangeFeed: { + enabled: master.services.api.storage_change_feed?.enabled ?? false, + stream: master.services.api.storage_change_feed?.stream ?? 'STORAGE_CHANGES', + skipBuckets: master.services.api.storage_change_feed?.skip_buckets ?? [s3Buckets.uploads], + }, search: { engine: master.integrations.search?.engine ?? 'elasticsearch', url: master.integrations.search?.url ?? 'http://127.0.0.1:9200', diff --git a/fluxer_api/src/api/app/APILifecycle.ts b/fluxer_api/src/api/app/APILifecycle.ts index 7410f4051..f99b1a6e5 100644 --- a/fluxer_api/src/api/app/APILifecycle.ts +++ b/fluxer_api/src/api/app/APILifecycle.ts @@ -7,6 +7,7 @@ import {hasDatabaseQueryExecutor, setDatabaseQueryExecutor} from '@app/api/datab import {ensurePostgresKvSchema, PostgresKvQueryExecutor} from '@app/api/database/PostgresKvQueryExecutor'; import {GuildDataRepository} from '@app/api/guild/repositories/GuildDataRepository'; import type {ILogger} from '@app/api/ILogger'; +import {shutdownStorageChangeFeed} from '@app/api/infrastructure/StorageServiceFactory'; import {JobLedgerRepository} from '@app/api/jobs/JobLedgerRepository'; import {startAbuseReplicationSubscriber, stopAbuseReplicationSubscriber} from '@app/api/middleware/AbusiveIpAutoBanner'; import {ipBanCache} from '@app/api/middleware/IpBanMiddleware'; @@ -258,6 +259,7 @@ export function createShutdown(config: APIConfig, logger: ILogger): () => Promis logger.error({error}, 'Error draining JetStream worker connection'); } } + await shutdownStorageChangeFeed(); setInjectedWorkerService(undefined); try { await shutdownSearch(); diff --git a/fluxer_api/src/api/config/APIConfig.ts b/fluxer_api/src/api/config/APIConfig.ts index cb09a9b70..996530a99 100644 --- a/fluxer_api/src/api/config/APIConfig.ts +++ b/fluxer_api/src/api/config/APIConfig.ts @@ -94,6 +94,11 @@ export interface APIConfig { jetStreamUrl: string; authToken: string; }; + storageChangeFeed: { + enabled: boolean; + stream: string; + skipBuckets: Array; + }; search: { engine: 'elasticsearch' | 'meilisearch'; url: string; diff --git a/fluxer_api/src/api/infrastructure/ChangeFeedStorageService.test.ts b/fluxer_api/src/api/infrastructure/ChangeFeedStorageService.test.ts new file mode 100644 index 000000000..423818a26 --- /dev/null +++ b/fluxer_api/src/api/infrastructure/ChangeFeedStorageService.test.ts @@ -0,0 +1,222 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import {Config} from '@app/api/Config'; +import {ChangeFeedStorageService} from '@app/api/infrastructure/ChangeFeedStorageService'; +import { + type StorageChangeEvent, + StorageChangeFeed, + type StorageChangeSink, +} from '@app/api/infrastructure/StorageChangeFeed'; +import {StorageService} from '@app/api/infrastructure/StorageService'; +import {createStorageService, shutdownStorageChangeFeed} from '@app/api/infrastructure/StorageServiceFactory'; +import {MockStorageService} from '@app/api/test/mocks/MockStorageService'; +import {describe, expect, it} from 'vitest'; + +class RecordingSink implements StorageChangeSink { + readonly events: Array = []; + + async publish(subject: string, body: string): Promise { + this.events.push({subject, ...(JSON.parse(body) as StorageChangeEvent)}); + } + + async close(): Promise {} +} + +function createHarness(options: {inner?: MockStorageService; sink?: StorageChangeSink; capacity?: number} = {}): { + service: ChangeFeedStorageService; + inner: MockStorageService; + feed: StorageChangeFeed; + sink: RecordingSink; +} { + const inner = options.inner ?? new MockStorageService(); + const sink = new RecordingSink(); + const feed = new StorageChangeFeed(options.sink ?? sink, { + capacity: options.capacity, + maxAttempts: 1, + retryDelayMs: 0, + }); + const service = new ChangeFeedStorageService(inner, feed, [Config.s3.buckets.uploads]); + return {service, inner, feed, sink}; +} + +const bytes = (length: number) => new Uint8Array(length).fill(7); + +describe('ChangeFeedStorageService', () => { + it('emits a put event after an upload lands', async () => { + const {service, inner, feed, sink} = createHarness(); + + await service.uploadObject({ + bucket: Config.s3.buckets.cdn, + key: 'attachments/1/2/cat.png', + body: bytes(5), + contentType: 'image/png', + }); + await feed.idle(); + + expect(await inner.getObjectMetadata(Config.s3.buckets.cdn, 'attachments/1/2/cat.png')).not.toBeNull(); + expect(sink.events).toEqual([ + { + subject: `storage.${Config.s3.buckets.cdn}.put`, + bucket: Config.s3.buckets.cdn, + key: 'attachments/1/2/cat.png', + op: 'put', + size: 5, + etag: null, + contentType: 'image/png', + at: expect.any(String), + }, + ]); + }); + + it('emits delete events for single and batch deletes', async () => { + const {service, feed, sink} = createHarness(); + + await service.deleteObject(Config.s3.buckets.cdn, 'icons/1/a.webp'); + await service.deleteObjects({bucket: Config.s3.buckets.reports, objects: [{Key: 'r/1'}, {Key: 'r/2'}]}); + await service.deleteAvatar({prefix: 'avatars/9', key: 'hash'}); + await feed.idle(); + + expect(sink.events.map(({subject, bucket, key, op}) => ({subject, bucket, key, op}))).toEqual([ + { + subject: `storage.${Config.s3.buckets.cdn}.delete`, + bucket: Config.s3.buckets.cdn, + key: 'icons/1/a.webp', + op: 'delete', + }, + { + subject: `storage.${Config.s3.buckets.reports}.delete`, + bucket: Config.s3.buckets.reports, + key: 'r/1', + op: 'delete', + }, + { + subject: `storage.${Config.s3.buckets.reports}.delete`, + bucket: Config.s3.buckets.reports, + key: 'r/2', + op: 'delete', + }, + { + subject: `storage.${Config.s3.buckets.cdn}.delete`, + bucket: Config.s3.buckets.cdn, + key: 'avatars/9/hash', + op: 'delete', + }, + ]); + }); + + it('emits a put and a delete for a move between durable buckets', async () => { + const {service, inner, feed, sink} = createHarness(); + await inner.uploadObject({bucket: Config.s3.buckets.cdn, key: 'from', body: bytes(3)}); + + await service.moveObject({ + sourceBucket: Config.s3.buckets.cdn, + sourceKey: 'from', + destinationBucket: Config.s3.buckets.reports, + destinationKey: 'to', + }); + await feed.idle(); + + expect(sink.events.map(({bucket, key, op}) => ({bucket, key, op}))).toEqual([ + {bucket: Config.s3.buckets.reports, key: 'to', op: 'put'}, + {bucket: Config.s3.buckets.cdn, key: 'from', op: 'delete'}, + ]); + }); + + it('emits nothing for the uploads bucket but still reports promotions out of it', async () => { + const {service, inner, feed, sink} = createHarness(); + const uploads = Config.s3.buckets.uploads; + + await service.uploadObject({bucket: uploads, key: 'staged', body: bytes(4)}); + await service.copyObject({ + sourceBucket: uploads, + sourceKey: 'staged', + destinationBucket: Config.s3.buckets.cdn, + destinationKey: 'attachments/1/2/staged.bin', + }); + await service.deleteObject(uploads, 'staged'); + await feed.idle(); + + expect(inner.uploadObjectSpy).toHaveBeenCalledTimes(1); + expect(sink.events.map(({bucket, key, op}) => ({bucket, key, op}))).toEqual([ + {bucket: Config.s3.buckets.cdn, key: 'attachments/1/2/staged.bin', op: 'put'}, + ]); + }); + + it('emits nothing when the storage call itself fails', async () => { + const {service, feed, sink} = createHarness({inner: new MockStorageService({shouldFailUpload: true})}); + + await expect(service.uploadObject({bucket: Config.s3.buckets.cdn, key: 'broken', body: bytes(1)})).rejects.toThrow( + 'Mock storage upload failure', + ); + await feed.idle(); + + expect(sink.events).toEqual([]); + }); + + it('does not let a failing publish reach the caller', async () => { + const sink: StorageChangeSink = { + publish: async () => { + throw new Error('nats: timeout'); + }, + close: async () => undefined, + }; + const {service, inner, feed} = createHarness({sink}); + + await expect( + service.uploadObject({bucket: Config.s3.buckets.cdn, key: 'kept', body: bytes(2)}), + ).resolves.toBeUndefined(); + await expect(service.deleteObject(Config.s3.buckets.cdn, 'kept')).resolves.toBeUndefined(); + await feed.idle(); + + expect(inner.deleteObjectSpy).toHaveBeenCalledWith(Config.s3.buckets.cdn, 'kept'); + expect(feed.stats()).toEqual({queued: 0, published: 0, dropped: 0, failed: 2}); + }); + + it('does not wait for a publish that never answers', async () => { + const sink: StorageChangeSink = { + publish: () => new Promise(() => undefined), + close: async () => undefined, + }; + const {service, feed} = createHarness({sink}); + + await service.uploadObject({bucket: Config.s3.buckets.cdn, key: 'first', body: bytes(1)}); + await new Promise((resolve) => setImmediate(resolve)); + await service.uploadObject({bucket: Config.s3.buckets.cdn, key: 'second', body: bytes(1)}); + + expect(feed.stats()).toEqual({queued: 1, published: 0, dropped: 0, failed: 0}); + }); + + it('drops and counts events when the queue is full without failing the storage call', async () => { + const sink: StorageChangeSink = { + publish: () => new Promise(() => undefined), + close: async () => undefined, + }; + const {service, inner, feed} = createHarness({sink, capacity: 1}); + + for (const key of ['a', 'b', 'c']) { + await service.uploadObject({bucket: Config.s3.buckets.cdn, key, body: bytes(1)}); + } + + expect(inner.uploadObjectSpy).toHaveBeenCalledTimes(3); + expect(feed.stats()).toEqual({queued: 1, published: 0, dropped: 2, failed: 0}); + }); +}); + +describe('createStorageService', () => { + it('returns the plain storage service while the change feed is disabled', () => { + expect(Config.storageChangeFeed.enabled).toBe(false); + expect(Config.storageChangeFeed.skipBuckets).toEqual([Config.s3.buckets.uploads]); + expect(createStorageService()).toBeInstanceOf(StorageService); + }); + + it('wraps the storage service once the change feed is enabled', async () => { + const original = Config.storageChangeFeed.enabled; + Config.storageChangeFeed.enabled = true; + try { + expect(createStorageService()).toBeInstanceOf(ChangeFeedStorageService); + } finally { + Config.storageChangeFeed.enabled = original; + await shutdownStorageChangeFeed(); + } + }); +}); diff --git a/fluxer_api/src/api/infrastructure/ChangeFeedStorageService.ts b/fluxer_api/src/api/infrastructure/ChangeFeedStorageService.ts new file mode 100644 index 000000000..87e63f857 --- /dev/null +++ b/fluxer_api/src/api/infrastructure/ChangeFeedStorageService.ts @@ -0,0 +1,162 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import {Config} from '@app/api/Config'; +import type {IStorageService, ProcessedStorageObjectMetadata} from '@app/api/infrastructure/IStorageService'; +import type {StorageChangeFeed, StorageChangeOp} from '@app/api/infrastructure/StorageChangeFeed'; + +type Params = Parameters; +type Result = ReturnType; + +export class ChangeFeedStorageService implements IStorageService { + private readonly inner: IStorageService; + private readonly feed: StorageChangeFeed; + private readonly skipBuckets: ReadonlySet; + + constructor(inner: IStorageService, feed: StorageChangeFeed, skipBuckets: Iterable) { + this.inner = inner; + this.feed = feed; + this.skipBuckets = new Set(skipBuckets); + } + + private emit( + bucket: string, + key: string, + op: StorageChangeOp, + details: {size?: number | null; contentType?: string | null} = {}, + ): void { + if (this.skipBuckets.has(bucket)) return; + this.feed.record({ + bucket, + key, + op, + size: details.size ?? null, + etag: null, + contentType: details.contentType ?? null, + at: new Date().toISOString(), + }); + } + + async uploadObject(params: Params<'uploadObject'>[0]): Promise { + await this.inner.uploadObject(params); + this.emit(params.bucket, params.key, 'put', { + size: params.body instanceof Uint8Array ? params.body.byteLength : null, + contentType: params.contentType, + }); + } + + async uploadObjectFromFile(params: Params<'uploadObjectFromFile'>[0]): Promise { + await this.inner.uploadObjectFromFile(params); + this.emit(params.bucket, params.key, 'put', {size: params.contentLength, contentType: params.contentType}); + } + + async deleteObject(bucket: string, key: string): Promise { + await this.inner.deleteObject(bucket, key); + this.emit(bucket, key, 'delete'); + } + + getObjectMetadata(...args: Params<'getObjectMetadata'>): Result<'getObjectMetadata'> { + return this.inner.getObjectMetadata(...args); + } + + computeObjectSha256(...args: Params<'computeObjectSha256'>): Result<'computeObjectSha256'> { + return this.inner.computeObjectSha256(...args); + } + + readObject(...args: Params<'readObject'>): Result<'readObject'> { + return this.inner.readObject(...args); + } + + streamObject(...args: Params<'streamObject'>): Result<'streamObject'> { + return this.inner.streamObject(...args); + } + + writeObjectToDisk(...args: Params<'writeObjectToDisk'>): Result<'writeObjectToDisk'> { + return this.inner.writeObjectToDisk(...args); + } + + async copyObject(params: Params<'copyObject'>[0]): Promise { + await this.inner.copyObject(params); + this.emit(params.destinationBucket, params.destinationKey, 'put', {contentType: params.newContentType}); + } + + async copyObjectWithMetadataStripping( + params: Params<'copyObjectWithMetadataStripping'>[0], + ): Promise { + const processed = await this.inner.copyObjectWithMetadataStripping(params); + this.emit(params.destinationBucket, params.destinationKey, 'put', { + size: processed?.contentLength, + contentType: processed?.contentType ?? params.contentType, + }); + return processed; + } + + async moveObject(params: Params<'moveObject'>[0]): Promise { + try { + await this.inner.moveObject(params); + } catch (error) { + this.emit(params.destinationBucket, params.destinationKey, 'put', {contentType: params.newContentType}); + throw error; + } + this.emit(params.destinationBucket, params.destinationKey, 'put', {contentType: params.newContentType}); + this.emit(params.sourceBucket, params.sourceKey, 'delete'); + } + + getPresignedDownloadURL(...args: Params<'getPresignedDownloadURL'>): Result<'getPresignedDownloadURL'> { + return this.inner.getPresignedDownloadURL(...args); + } + + getPresignedUploadURL(...args: Params<'getPresignedUploadURL'>): Result<'getPresignedUploadURL'> { + return this.inner.getPresignedUploadURL(...args); + } + + getPresignedUploadPartURL(...args: Params<'getPresignedUploadPartURL'>): Result<'getPresignedUploadPartURL'> { + return this.inner.getPresignedUploadPartURL(...args); + } + + async purgeBucket(bucket: string): Promise { + const objects = await this.inner.listObjects({bucket, prefix: ''}); + await Promise.all(objects.map(({key}) => this.deleteObject(bucket, key))); + } + + async uploadAvatar(params: Params<'uploadAvatar'>[0]): Promise { + await this.inner.uploadAvatar(params); + this.emit(Config.s3.buckets.cdn, `${params.prefix}/${params.key}`, 'put', {size: params.body.byteLength}); + } + + async deleteAvatar(params: Params<'deleteAvatar'>[0]): Promise { + await this.inner.deleteAvatar(params); + this.emit(Config.s3.buckets.cdn, `${params.prefix}/${params.key}`, 'delete'); + } + + listObjects(...args: Params<'listObjects'>): Result<'listObjects'> { + return this.inner.listObjects(...args); + } + + async deleteObjects(params: Params<'deleteObjects'>[0]): Promise { + await this.inner.deleteObjects(params); + for (const object of params.objects) { + this.emit(params.bucket, object.Key, 'delete'); + } + } + + createMultipartUpload(...args: Params<'createMultipartUpload'>): Result<'createMultipartUpload'> { + return this.inner.createMultipartUpload(...args); + } + + uploadPart(...args: Params<'uploadPart'>): Result<'uploadPart'> { + return this.inner.uploadPart(...args); + } + + listParts(...args: Params<'listParts'>): Result<'listParts'> { + return this.inner.listParts(...args); + } + + async completeMultipartUpload(params: Params<'completeMultipartUpload'>[0]): Promise { + await this.inner.completeMultipartUpload(params); + this.emit(params.bucket, params.key, 'put'); + } + + abortMultipartUpload(...args: Params<'abortMultipartUpload'>): Result<'abortMultipartUpload'> { + return this.inner.abortMultipartUpload(...args); + } +} diff --git a/fluxer_api/src/api/infrastructure/StorageChangeFeed.test.ts b/fluxer_api/src/api/infrastructure/StorageChangeFeed.test.ts new file mode 100644 index 000000000..da92c035b --- /dev/null +++ b/fluxer_api/src/api/infrastructure/StorageChangeFeed.test.ts @@ -0,0 +1,158 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import { + type StorageChangeEvent, + StorageChangeFeed, + type StorageChangeSink, +} from '@app/api/infrastructure/StorageChangeFeed'; +import {describe, expect, it} from 'vitest'; + +interface PublishCall { + subject: string; + body: string; + msgID: string; +} + +class ScriptedSink implements StorageChangeSink { + readonly calls: Array = []; + closed = false; + private readonly outcomes: Array<() => Promise>; + + constructor(outcomes: Array<() => Promise> = []) { + this.outcomes = outcomes; + } + + async publish(subject: string, body: string, msgID: string): Promise { + this.calls.push({subject, body, msgID}); + const next = this.outcomes.shift(); + if (next) await next(); + } + + async close(): Promise { + this.closed = true; + } +} + +function change(key: string, overrides: Partial = {}): StorageChangeEvent { + return { + bucket: 'fluxer', + key, + op: 'put', + size: 12, + etag: null, + contentType: 'image/png', + at: '2026-09-14T12:00:00.000Z', + ...overrides, + }; +} + +describe('StorageChangeFeed', () => { + it('publishes metadata to the bucket and operation subject with a dedupe id', async () => { + const sink = new ScriptedSink(); + const feed = new StorageChangeFeed(sink); + + feed.record(change('attachments/1/2/cat.png')); + feed.record(change('avatars/3/abc', {bucket: 'fluxer.static', op: 'delete', etag: '"abc"'})); + await feed.idle(); + + expect(sink.calls.map(({subject, msgID}) => ({subject, msgID}))).toEqual([ + {subject: 'storage.fluxer.put', msgID: 'fluxer:attachments/1/2/cat.png:2026-09-14T12:00:00.000Z'}, + {subject: 'storage.fluxer_static.delete', msgID: 'fluxer.static:avatars/3/abc:"abc"'}, + ]); + expect(JSON.parse(sink.calls[0].body)).toEqual(change('attachments/1/2/cat.png')); + expect(feed.stats()).toEqual({queued: 0, published: 2, dropped: 0, failed: 0}); + }); + + it('does not publish before the caller has moved on', () => { + const sink = new ScriptedSink(); + const feed = new StorageChangeFeed(sink); + + feed.record(change('attachments/1/2/cat.png')); + + expect(sink.calls).toEqual([]); + expect(feed.stats().queued).toBe(1); + }); + + it('drops and counts events once the queue is full', async () => { + const gate = Promise.withResolvers(); + const sink = new ScriptedSink([() => gate.promise]); + const feed = new StorageChangeFeed(sink, {capacity: 2, concurrency: 1}); + + for (const key of ['a', 'b', 'c', 'd', 'e']) { + feed.record(change(key)); + } + + expect(feed.stats()).toEqual({queued: 2, published: 0, dropped: 3, failed: 0}); + gate.resolve(); + await feed.idle(); + expect(sink.calls.map(({msgID}) => msgID.split(':')[1])).toEqual(['a', 'b']); + expect(feed.stats()).toEqual({queued: 0, published: 2, dropped: 3, failed: 0}); + }); + + it('retries a failed publish with the same dedupe id', async () => { + const sink = new ScriptedSink([() => Promise.reject(new Error('no responders'))]); + const feed = new StorageChangeFeed(sink, {retryDelayMs: 0}); + + feed.record(change('attachments/1/2/cat.png')); + await feed.idle(); + + expect(sink.calls).toHaveLength(2); + expect(sink.calls[0].msgID).toBe(sink.calls[1].msgID); + expect(feed.stats()).toEqual({queued: 0, published: 1, dropped: 0, failed: 0}); + }); + + it('counts an event as failed once its attempts run out and keeps going', async () => { + const reject = () => Promise.reject(new Error('stream offline')); + const sink = new ScriptedSink([reject, reject]); + const feed = new StorageChangeFeed(sink, {maxAttempts: 2, retryDelayMs: 0, concurrency: 1}); + + feed.record(change('lost')); + feed.record(change('kept')); + await feed.idle(); + + expect(sink.calls.map(({msgID}) => msgID.split(':')[1])).toEqual(['lost', 'lost', 'kept']); + expect(feed.stats()).toEqual({queued: 0, published: 1, dropped: 0, failed: 1}); + }); + + it('survives a sink that throws synchronously', async () => { + const sink: StorageChangeSink = { + publish: () => { + throw new Error('NATS connection is not established. Call connect() first.'); + }, + close: async () => undefined, + }; + const feed = new StorageChangeFeed(sink, {maxAttempts: 1, retryDelayMs: 0}); + + expect(() => feed.record(change('attachments/1/2/cat.png'))).not.toThrow(); + await feed.idle(); + + expect(feed.stats()).toEqual({queued: 0, published: 0, dropped: 0, failed: 1}); + }); + + it('flushes queued events on close, then closes the sink and drops later events', async () => { + const sink = new ScriptedSink(); + const feed = new StorageChangeFeed(sink); + + feed.record(change('before-close')); + await feed.close(1000); + feed.record(change('after-close')); + + expect(sink.calls.map(({msgID}) => msgID.split(':')[1])).toEqual(['before-close']); + expect(sink.closed).toBe(true); + expect(feed.stats()).toEqual({queued: 0, published: 1, dropped: 1, failed: 0}); + }); + + it('stops waiting for a stuck publish when close runs out of time', async () => { + const sink = new ScriptedSink([() => new Promise(() => undefined)]); + const feed = new StorageChangeFeed(sink, {concurrency: 1}); + + feed.record(change('stuck')); + feed.record(change('waiting')); + await new Promise((resolve) => setImmediate(resolve)); + await feed.close(10); + + expect(sink.calls.map(({msgID}) => msgID.split(':')[1])).toEqual(['stuck']); + expect(sink.closed).toBe(true); + expect(feed.stats()).toEqual({queued: 0, published: 0, dropped: 1, failed: 0}); + }); +}); diff --git a/fluxer_api/src/api/infrastructure/StorageChangeFeed.ts b/fluxer_api/src/api/infrastructure/StorageChangeFeed.ts new file mode 100644 index 000000000..7459b7678 --- /dev/null +++ b/fluxer_api/src/api/infrastructure/StorageChangeFeed.ts @@ -0,0 +1,265 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import {Logger} from '@app/api/Logger'; +import {JetStreamConnectionManager} from '@pkgs/nats/src/JetStreamConnectionManager'; +import {DiscardPolicy, NatsError, nanos, RetentionPolicy, StorageType} from 'nats'; + +export type StorageChangeOp = 'put' | 'delete'; + +export interface StorageChangeEvent { + bucket: string; + key: string; + op: StorageChangeOp; + size: number | null; + etag: string | null; + contentType: string | null; + at: string; +} + +export interface StorageChangeSink { + publish(subject: string, body: string, msgID: string): Promise; + close(): Promise; +} + +export interface StorageChangeFeedStats { + queued: number; + published: number; + dropped: number; + failed: number; +} + +interface StorageChangeFeedOptions { + capacity?: number; + concurrency?: number; + maxAttempts?: number; + retryDelayMs?: number; + reportIntervalMs?: number; +} + +interface QueuedChange { + subject: string; + body: string; + msgID: string; + attempts: number; +} + +const DEFAULT_CAPACITY = 20_000; +const DEFAULT_CONCURRENCY = 64; +const DEFAULT_MAX_ATTEMPTS = 5; +const DEFAULT_RETRY_DELAY_MS = 1000; +const DEFAULT_REPORT_INTERVAL_MS = 60_000; + +const STREAM_SUBJECTS = 'storage.>'; +const STREAM_MAX_AGE_MS = 14 * 24 * 60 * 60 * 1000; +const STREAM_MAX_BYTES = 1024 * 1024 * 1024; +const STREAM_DUPLICATE_WINDOW_MS = 2 * 60 * 1000; +const PUBLISH_TIMEOUT_MS = 5000; +const STREAM_NAME_IN_USE_ERR_CODE = 10058; +const STREAM_NOT_FOUND_ERR_CODE = 10059; + +function delay(ms: number): Promise { + return new Promise((resolve) => { + setTimeout(resolve, ms).unref(); + }); +} + +function subjectToken(bucket: string): string { + return bucket.replace(/[.*>\s]/gu, '_'); +} + +function jsErrorCode(error: unknown): number | null { + if (!(error instanceof NatsError)) { + return null; + } + return error.jsError()?.err_code ?? null; +} + +export class StorageChangeFeed { + private readonly sink: StorageChangeSink; + private readonly capacity: number; + private readonly concurrency: number; + private readonly maxAttempts: number; + private readonly retryDelayMs: number; + private readonly reportIntervalMs: number; + private readonly queue: Array = []; + private draining: Promise | null = null; + private stopped = false; + private published = 0; + private dropped = 0; + private failed = 0; + private lastReportAt = Number.NEGATIVE_INFINITY; + private reportedDropped = 0; + private reportedFailed = 0; + + constructor(sink: StorageChangeSink, options: StorageChangeFeedOptions = {}) { + this.sink = sink; + this.capacity = options.capacity ?? DEFAULT_CAPACITY; + this.concurrency = options.concurrency ?? DEFAULT_CONCURRENCY; + this.maxAttempts = options.maxAttempts ?? DEFAULT_MAX_ATTEMPTS; + this.retryDelayMs = options.retryDelayMs ?? DEFAULT_RETRY_DELAY_MS; + this.reportIntervalMs = options.reportIntervalMs ?? DEFAULT_REPORT_INTERVAL_MS; + } + + record(event: StorageChangeEvent): void { + try { + if (this.stopped || this.queue.length >= this.capacity) { + this.dropped++; + this.report(null); + return; + } + this.queue.push({ + subject: `storage.${subjectToken(event.bucket)}.${event.op}`, + body: JSON.stringify(event), + msgID: `${event.bucket}:${event.key}:${event.etag ?? event.at}`, + attempts: 0, + }); + this.draining ??= this.scheduleDrain(); + } catch (error) { + this.dropped++; + this.report(error); + } + } + + stats(): StorageChangeFeedStats { + return {queued: this.queue.length, published: this.published, dropped: this.dropped, failed: this.failed}; + } + + async idle(): Promise { + while (this.draining !== null) { + await this.draining; + } + } + + async close(timeoutMs: number): Promise { + const deadline = delay(timeoutMs).then(() => false); + while (this.draining !== null) { + const drained = await Promise.race([this.draining.then(() => true), deadline]); + if (!drained) break; + } + this.stopped = true; + this.dropped += this.queue.length; + this.queue.length = 0; + Logger.info(this.stats(), 'Storage change feed stopped'); + try { + await this.sink.close(); + } catch (error) { + Logger.warn({err: error}, 'Storage change feed connection did not close cleanly'); + } + } + + private scheduleDrain(): Promise { + return new Promise((resolve) => setImmediate(resolve)) + .then(() => this.drain()) + .catch((error: unknown) => { + this.report(error); + }) + .finally(() => { + this.draining = null; + if (this.queue.length > 0 && !this.stopped) { + this.draining = this.scheduleDrain(); + } + }); + } + + private async drain(): Promise { + while (this.queue.length > 0 && !this.stopped) { + const batch = this.queue.splice(0, this.concurrency); + const results = await Promise.allSettled( + batch.map(async (change) => { + await this.sink.publish(change.subject, change.body, change.msgID); + }), + ); + const retry: Array = []; + let failure: PromiseRejectedResult | null = null; + for (const [index, result] of results.entries()) { + if (result.status === 'fulfilled') { + this.published++; + continue; + } + failure = result; + const change = batch[index]; + change.attempts++; + if (this.stopped) { + this.dropped++; + } else if (change.attempts >= this.maxAttempts) { + this.failed++; + } else { + retry.push(change); + } + } + if (failure === null) continue; + this.queue.unshift(...retry); + this.report(failure.reason); + await delay(this.retryDelayMs); + } + } + + private report(error: unknown): void { + const now = Date.now(); + if (now - this.lastReportAt < this.reportIntervalMs) return; + if (error === null && this.dropped === this.reportedDropped && this.failed === this.reportedFailed) return; + this.lastReportAt = now; + this.reportedDropped = this.dropped; + this.reportedFailed = this.failed; + Logger.warn({err: error ?? undefined, ...this.stats()}, 'Storage change feed is dropping or failing events'); + } +} + +export class JetStreamStorageChangeSink implements StorageChangeSink { + private readonly connection: JetStreamConnectionManager; + private readonly stream: string; + private ready: Promise | null = null; + + constructor(options: {url: string; token: string; stream: string}) { + this.connection = new JetStreamConnectionManager({ + url: options.url, + token: options.token || undefined, + name: 'storage-change-feed', + }); + this.stream = options.stream; + } + + async publish(subject: string, body: string, msgID: string): Promise { + this.ready ??= this.open().catch((error: unknown) => { + this.ready = null; + throw error; + }); + await this.ready; + if (this.connection.isClosed()) { + this.ready = null; + throw new Error('Storage change feed connection is closed'); + } + await this.connection.getJetStreamClient().publish(subject, body, {msgID, timeout: PUBLISH_TIMEOUT_MS}); + } + + async close(): Promise { + await this.connection.drain(); + } + + private async open(): Promise { + await this.connection.connect(); + const jsm = await this.connection.getJetStreamManager(); + try { + await jsm.streams.info(this.stream); + return; + } catch (error) { + if (jsErrorCode(error) !== STREAM_NOT_FOUND_ERR_CODE) throw error; + } + try { + await jsm.streams.add({ + name: this.stream, + subjects: [STREAM_SUBJECTS], + retention: RetentionPolicy.Limits, + storage: StorageType.File, + max_age: nanos(STREAM_MAX_AGE_MS), + max_bytes: STREAM_MAX_BYTES, + duplicate_window: nanos(STREAM_DUPLICATE_WINDOW_MS), + discard: DiscardPolicy.Old, + num_replicas: 1, + }); + Logger.info({stream: this.stream}, 'Created storage change feed stream'); + } catch (error) { + if (jsErrorCode(error) !== STREAM_NAME_IN_USE_ERR_CODE) throw error; + } + } +} diff --git a/fluxer_api/src/api/infrastructure/StorageServiceFactory.ts b/fluxer_api/src/api/infrastructure/StorageServiceFactory.ts index ee9ddfdbc..45ee7c8db 100644 --- a/fluxer_api/src/api/infrastructure/StorageServiceFactory.ts +++ b/fluxer_api/src/api/infrastructure/StorageServiceFactory.ts @@ -1,16 +1,39 @@ // SPDX-License-Identifier: AGPL-3.0-or-later import {Config} from '@app/api/Config'; +import {ChangeFeedStorageService} from '@app/api/infrastructure/ChangeFeedStorageService'; import type {IStorageService} from '@app/api/infrastructure/IStorageService'; +import {JetStreamStorageChangeSink, StorageChangeFeed} from '@app/api/infrastructure/StorageChangeFeed'; import {StorageService} from '@app/api/infrastructure/StorageService'; +const CHANGE_FEED_SHUTDOWN_TIMEOUT_MS = 5000; + +let changeFeed: StorageChangeFeed | null = null; + +function withChangeFeed(service: IStorageService): IStorageService { + const {enabled, stream, skipBuckets} = Config.storageChangeFeed; + if (!enabled) { + return service; + } + changeFeed ??= new StorageChangeFeed( + new JetStreamStorageChangeSink({url: Config.nats.jetStreamUrl, token: Config.nats.authToken, stream}), + ); + return new ChangeFeedStorageService(service, changeFeed, skipBuckets); +} + export function createStorageService(): IStorageService { - return new StorageService(); + return withChangeFeed(new StorageService()); } export function createDownloadsStorageService(): IStorageService | null { if (!Config.s3Downloads.isOverridden) { return null; } - return new StorageService(Config.s3Downloads.settings); + return withChangeFeed(new StorageService(Config.s3Downloads.settings)); +} + +export async function shutdownStorageChangeFeed(): Promise { + const feed = changeFeed; + changeFeed = null; + await feed?.close(CHANGE_FEED_SHUTDOWN_TIMEOUT_MS); } diff --git a/fluxer_api/src/api/worker/WorkerMain.ts b/fluxer_api/src/api/worker/WorkerMain.ts index 5c02e5f5b..5e61bab03 100644 --- a/fluxer_api/src/api/worker/WorkerMain.ts +++ b/fluxer_api/src/api/worker/WorkerMain.ts @@ -4,6 +4,7 @@ import {Config} from '@app/api/Config'; import {setDatabaseQueryExecutor} from '@app/api/database/CassandraQueryExecution'; import {ensurePostgresKvSchema, PostgresKvQueryExecutor} from '@app/api/database/PostgresKvQueryExecutor'; import type {ISnowflakeService} from '@app/api/infrastructure/ISnowflakeService'; +import {shutdownStorageChangeFeed} from '@app/api/infrastructure/StorageServiceFactory'; import type {InstanceConfigRepository} from '@app/api/instance/InstanceConfigRepository'; import {JobLedgerRepository} from '@app/api/jobs/JobLedgerRepository'; import {Logger} from '@app/api/Logger'; @@ -114,6 +115,7 @@ export async function startWorkerMain(): Promise { ); }); await voiceShutdown; + await cleanupStep('storage change feed', shutdownStorageChangeFeed); await cleanupStep('jetstream', async () => { await jsConnectionManager?.drain(); jsConnectionManager = null; diff --git a/packages/config/src/ConfigLoader.ts b/packages/config/src/ConfigLoader.ts index 0c98a9a0e..7d173bd4d 100644 --- a/packages/config/src/ConfigLoader.ts +++ b/packages/config/src/ConfigLoader.ts @@ -122,6 +122,10 @@ function defaultConfig(): MasterConfig { content_moderation: { nsfw_threshold: 0.7, }, + storage_change_feed: { + enabled: false, + stream: 'STORAGE_CHANGES', + }, }, nats: { core_url: 'nats://127.0.0.1:4222', @@ -453,6 +457,18 @@ function validateApiWorkerConfig(config: MasterConfig): void { } } +function validateStorageChangeFeedConfig(config: MasterConfig): void { + const feed = config.services.api?.storage_change_feed; + if (!feed?.enabled) { + return; + } + if (feed.stream === undefined || !/^[A-Za-z0-9_-]+$/u.test(feed.stream)) { + throw new Error( + 'FLUXER_API_STORAGE_CHANGE_FEED_STREAM must be letters, digits, underscores or hyphens when the storage change feed is enabled', + ); + } +} + function validateCachePurgeConfig(config: MasterConfig): void { const cachePurge = config.integrations.cache_purge; if (cachePurge.adapter !== 'http') { @@ -541,6 +557,7 @@ function normalizeConfig(config: MasterConfig): MasterConfig { validatePostgresConfig(config); validateCaptchaConfig(config); validateApiWorkerConfig(config); + validateStorageChangeFeedConfig(config); validateCachePurgeConfig(config); assertIntegerInRange(config.services.api.max_inflight_requests, 'FLUXER_API_MAX_INFLIGHT_REQUESTS', 1, 100_000); assertIntegerInRange(config.services.api.headers_timeout_ms, 'FLUXER_API_HEADERS_TIMEOUT_MS', 1_000, 3_600_000); diff --git a/packages/config/src/MasterConfig.ts b/packages/config/src/MasterConfig.ts index 9a6a9db3b..96dc67a2b 100644 --- a/packages/config/src/MasterConfig.ts +++ b/packages/config/src/MasterConfig.ts @@ -127,6 +127,11 @@ export interface MasterConfig { batch?: number; }; }; + storage_change_feed?: { + enabled?: boolean; + stream?: string; + skip_buckets?: Array; + }; }; nats?: { core_url?: string; diff --git a/packages/config/src/__tests__/ConfigLoader.test.ts b/packages/config/src/__tests__/ConfigLoader.test.ts index cd17b4060..3b1eb6242 100644 --- a/packages/config/src/__tests__/ConfigLoader.test.ts +++ b/packages/config/src/__tests__/ConfigLoader.test.ts @@ -309,6 +309,34 @@ describe('ConfigLoader', () => { await expect(loadConfig()).rejects.toThrow('FLUXER_API_WORKER_TASK'); }); + test('leaves the storage change feed disabled by default', async () => { + stubMinimalEnv(); + const config = await loadConfig(); + expect(config.services.api.storage_change_feed).toEqual({enabled: false, stream: 'STORAGE_CHANGES'}); + }); + + test('maps the storage change feed environment variables', async () => { + stubMinimalEnv({ + FLUXER_API_STORAGE_CHANGE_FEED_ENABLED: 'true', + FLUXER_API_STORAGE_CHANGE_FEED_STREAM: 'BACKUP_CHANGES', + FLUXER_API_STORAGE_CHANGE_FEED_SKIP_BUCKETS: 'fluxer-uploads, fluxer-harvests', + }); + const config = await loadConfig(); + expect(config.services.api.storage_change_feed).toEqual({ + enabled: true, + stream: 'BACKUP_CHANGES', + skip_buckets: ['fluxer-uploads', 'fluxer-harvests'], + }); + }); + + test('rejects a storage change feed stream name that JetStream cannot use', async () => { + stubMinimalEnv({ + FLUXER_API_STORAGE_CHANGE_FEED_ENABLED: 'true', + FLUXER_API_STORAGE_CHANGE_FEED_STREAM: 'storage.changes', + }); + await expect(loadConfig()).rejects.toThrow('FLUXER_API_STORAGE_CHANGE_FEED_STREAM'); + }); + test('rejects invalid Postgres typed environment values', async () => { stubMinimalEnv({FLUXER_POSTGRES_PORT: 'abc'}); await expect(loadConfig()).rejects.toThrow('FLUXER_POSTGRES_PORT'); diff --git a/packages/config/src/config_loader/EnvironmentOverrides.ts b/packages/config/src/config_loader/EnvironmentOverrides.ts index 044f7f2bc..0c13bfcbb 100644 --- a/packages/config/src/config_loader/EnvironmentOverrides.ts +++ b/packages/config/src/config_loader/EnvironmentOverrides.ts @@ -112,6 +112,15 @@ const NAMED_FLUXER_ENV_OVERRIDES: Record = { path: ['services', 'api', 'worker', 'lane_concurrency_overrides'], parse: parseJsonObject, }, + FLUXER_API_STORAGE_CHANGE_FEED_ENABLED: { + path: ['services', 'api', 'storage_change_feed', 'enabled'], + parse: parseBoolean, + }, + FLUXER_API_STORAGE_CHANGE_FEED_STREAM: {path: ['services', 'api', 'storage_change_feed', 'stream']}, + FLUXER_API_STORAGE_CHANGE_FEED_SKIP_BUCKETS: { + path: ['services', 'api', 'storage_change_feed', 'skip_buckets'], + parse: parseCsv, + }, FLUXER_API_UNFURL_IGNORED_HOSTS: {path: ['services', 'api', 'unfurl_ignored_hosts'], parse: parseCsv}, FLUXER_API_EMBEDS_OEMBED_HTML_ENABLED: { path: ['services', 'api', 'embeds', 'oembed_html_enabled'],