feat(api): publish object storage changes to a JetStream feed (#2796)

This commit is contained in:
Hampus
2026-09-15 17:20:26 +02:00
committed by GitHub
parent 83c8e91955
commit b38e7c6433
13 changed files with 905 additions and 2 deletions
+5
View File
@@ -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',
+2
View File
@@ -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();
+5
View File
@@ -94,6 +94,11 @@ export interface APIConfig {
jetStreamUrl: string;
authToken: string;
};
storageChangeFeed: {
enabled: boolean;
stream: string;
skipBuckets: Array<string>;
};
search: {
engine: 'elasticsearch' | 'meilisearch';
url: string;
@@ -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<StorageChangeEvent & {subject: string}> = [];
async publish(subject: string, body: string): Promise<void> {
this.events.push({subject, ...(JSON.parse(body) as StorageChangeEvent)});
}
async close(): Promise<void> {}
}
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<void>(() => 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<void>(() => 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();
}
});
});
@@ -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<K extends keyof IStorageService> = Parameters<IStorageService[K]>;
type Result<K extends keyof IStorageService> = ReturnType<IStorageService[K]>;
export class ChangeFeedStorageService implements IStorageService {
private readonly inner: IStorageService;
private readonly feed: StorageChangeFeed;
private readonly skipBuckets: ReadonlySet<string>;
constructor(inner: IStorageService, feed: StorageChangeFeed, skipBuckets: Iterable<string>) {
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<void> {
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<void> {
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<void> {
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<void> {
await this.inner.copyObject(params);
this.emit(params.destinationBucket, params.destinationKey, 'put', {contentType: params.newContentType});
}
async copyObjectWithMetadataStripping(
params: Params<'copyObjectWithMetadataStripping'>[0],
): Promise<ProcessedStorageObjectMetadata | null> {
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<void> {
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<void> {
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<void> {
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<void> {
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<void> {
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<void> {
await this.inner.completeMultipartUpload(params);
this.emit(params.bucket, params.key, 'put');
}
abortMultipartUpload(...args: Params<'abortMultipartUpload'>): Result<'abortMultipartUpload'> {
return this.inner.abortMultipartUpload(...args);
}
}
@@ -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<PublishCall> = [];
closed = false;
private readonly outcomes: Array<() => Promise<void>>;
constructor(outcomes: Array<() => Promise<void>> = []) {
this.outcomes = outcomes;
}
async publish(subject: string, body: string, msgID: string): Promise<void> {
this.calls.push({subject, body, msgID});
const next = this.outcomes.shift();
if (next) await next();
}
async close(): Promise<void> {
this.closed = true;
}
}
function change(key: string, overrides: Partial<StorageChangeEvent> = {}): 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<void>();
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<void>(() => 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});
});
});
@@ -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<void>;
close(): Promise<void>;
}
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<void> {
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<QueuedChange> = [];
private draining: Promise<void> | 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<void> {
while (this.draining !== null) {
await this.draining;
}
}
async close(timeoutMs: number): Promise<void> {
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<void> {
return new Promise<void>((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<void> {
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<QueuedChange> = [];
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<void> | 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<void> {
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<void> {
await this.connection.drain();
}
private async open(): Promise<void> {
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;
}
}
}
@@ -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<void> {
const feed = changeFeed;
changeFeed = null;
await feed?.close(CHANGE_FEED_SHUTDOWN_TIMEOUT_MS);
}
+2
View File
@@ -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<void> {
);
});
await voiceShutdown;
await cleanupStep('storage change feed', shutdownStorageChangeFeed);
await cleanupStep('jetstream', async () => {
await jsConnectionManager?.drain();
jsConnectionManager = null;
+17
View File
@@ -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);
+5
View File
@@ -127,6 +127,11 @@ export interface MasterConfig {
batch?: number;
};
};
storage_change_feed?: {
enabled?: boolean;
stream?: string;
skip_buckets?: Array<string>;
};
};
nats?: {
core_url?: string;
@@ -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');
@@ -112,6 +112,15 @@ const NAMED_FLUXER_ENV_OVERRIDES: Record<string, NamedEnvOverride> = {
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'],