diff --git a/fluxer_api/src/api/download/DownloadService.ts b/fluxer_api/src/api/download/DownloadService.ts index 7a008f203..4fd040b97 100644 --- a/fluxer_api/src/api/download/DownloadService.ts +++ b/fluxer_api/src/api/download/DownloadService.ts @@ -11,7 +11,12 @@ import type { DesktopPlatform, } from '@fluxer/schema/src/domains/download/DownloadSchemas'; import {Config} from '../Config'; -import type {IStorageService} from '../infrastructure/IStorageService'; +import { + type IStorageService, + StorageObjectListingOverflowError, + StorageObjectRangeNotSatisfiableError, +} from '../infrastructure/IStorageService'; +import {Logger} from '../Logger'; import {isJsonRecord, parseJsonUnknown} from '../utils/JsonBoundaryUtils'; import { parseDesktopArtifactScope, @@ -45,10 +50,14 @@ function isStorageNotFoundError(error: unknown): boolean { } function isUnsatisfiableRangeError(error: unknown): boolean { + if (error instanceof StorageObjectRangeNotSatisfiableError) { + return true; + } return ( error instanceof S3ServiceException && (error.name === 'InvalidRange' || error.$metadata?.httpStatusCode === 416) ); } +const MAX_DESKTOP_OBJECTS_PER_PREFIX = 10_000; const DESKTOP_BUCKET_PREFIX = 'desktop'; const DESKTOP_TEST_BUCKET_PREFIX = 'desktop-test'; const DOWNLOAD_KEY_ALLOWED_PREFIXES = [`${DESKTOP_BUCKET_PREFIX}/`, `${DESKTOP_TEST_BUCKET_PREFIX}/`]; @@ -349,11 +358,8 @@ export class DownloadService { } const prefix = `${basePrefix}/`; try { - const objects = await this.storageService.listObjects({ - bucket: Config.s3.buckets.downloads, - prefix, - }); - if (!objects || objects.length === 0) { + const objects = await this.listDesktopArtifacts(prefix); + if (objects.length === 0) { return {versions: [], hasMore: false}; } const versionMap = new Map< @@ -791,6 +797,22 @@ export class DownloadService { return this.findLatestFilenameForRequestedArch(params); } + private async listDesktopArtifacts(prefix: string): Promise> { + try { + return await this.storageService.listObjects({ + bucket: Config.s3.buckets.downloads, + prefix, + maxObjects: MAX_DESKTOP_OBJECTS_PER_PREFIX, + }); + } catch (error) { + if (error instanceof StorageObjectListingOverflowError) { + Logger.warn({prefix, maxObjects: error.maxObjects}, 'Desktop artifact prefix outgrew its listing cap'); + return []; + } + throw error; + } + } + private isFilenameCompatibleWithRequestedArch(params: ManifestFilenameResolutionParams): boolean { const parsed = this.parseVersionFromFilename(params.filename, params.channel, params.plat, params.arch); if (!parsed) { @@ -815,11 +837,8 @@ export class DownloadService { return null; } const prefix = `${basePrefix}/`; - const objects = await this.storageService.listObjects({ - bucket: Config.s3.buckets.downloads, - prefix, - }); - if (!objects || objects.length === 0) { + const objects = await this.listDesktopArtifacts(prefix); + if (objects.length === 0) { return null; } let latestFilename: string | null = null; diff --git a/fluxer_api/src/api/infrastructure/IStorageService.ts b/fluxer_api/src/api/infrastructure/IStorageService.ts index 56824e67b..0199580ed 100644 --- a/fluxer_api/src/api/infrastructure/IStorageService.ts +++ b/fluxer_api/src/api/infrastructure/IStorageService.ts @@ -2,6 +2,34 @@ import type {Readable} from 'node:stream'; +export class StorageObjectRangeNotSatisfiableError extends Error { + readonly bucket: string; + readonly key: string; + readonly range: string; + + constructor(bucket: string, key: string, range: string) { + super(`Requested range ${range} is not satisfiable for ${bucket}/${key}`); + this.name = 'StorageObjectRangeNotSatisfiableError'; + this.bucket = bucket; + this.key = key; + this.range = range; + } +} + +export class StorageObjectListingOverflowError extends Error { + readonly bucket: string; + readonly prefix: string; + readonly maxObjects: number; + + constructor(bucket: string, prefix: string, maxObjects: number) { + super(`Object listing exceeds maximum of ${maxObjects} objects for ${bucket}/${prefix}`); + this.name = 'StorageObjectListingOverflowError'; + this.bucket = bucket; + this.prefix = prefix; + this.maxObjects = maxObjects; + } +} + export interface ProcessedStorageObjectMetadata { contentType: string; contentLength: number; @@ -97,7 +125,7 @@ export interface IStorageService { purgeBucket(bucket: string): Promise; uploadAvatar(params: {prefix: string; key: string; body: Uint8Array}): Promise; deleteAvatar(params: {prefix: string; key: string}): Promise; - listObjects(params: {bucket: string; prefix: string}): Promise< + listObjects(params: {bucket: string; prefix: string; maxObjects?: number}): Promise< ReadonlyArray<{ key: string; lastModified?: Date; diff --git a/fluxer_api/src/api/infrastructure/StorageService.test.ts b/fluxer_api/src/api/infrastructure/StorageService.test.ts index 829952da5..a8cf27b23 100644 --- a/fluxer_api/src/api/infrastructure/StorageService.test.ts +++ b/fluxer_api/src/api/infrastructure/StorageService.test.ts @@ -1,9 +1,10 @@ // SPDX-License-Identifier: AGPL-3.0-or-later import fs from 'node:fs'; -import type {Readable} from 'node:stream'; +import {PassThrough, Readable} from 'node:stream'; import {describe, expect, it} from 'vitest'; import {Config} from '../Config'; +import {StorageObjectListingOverflowError} from './IStorageService'; import {StorageService} from './StorageService'; interface CopyObjectTestParams { @@ -210,3 +211,131 @@ describe('provider selection', () => { expect(await service.client.config.region()).not.toBe(Config.s3.region); }); }); + +function bodyStream(chunks: Array): PassThrough { + const stream = new PassThrough(); + const writeChunk = (index: number): void => { + if (stream.destroyed) { + return; + } + const chunk = chunks[index]; + if (chunk === undefined) { + stream.end(); + return; + } + stream.write(chunk); + setImmediate(() => writeChunk(index + 1)); + }; + writeChunk(0); + return stream; +} + +function serveGetObject(service: StorageService, out: {Body: PassThrough; ContentLength?: number}): void { + Object.assign(service, {client: {send: async () => out}}); +} + +describe('StorageService.readObject', () => { + it('refuses an object whose declared length is already over the cap', async () => { + const service = new StorageService(); + const body = bodyStream([Buffer.alloc(64, 1)]); + serveGetObject(service, {Body: body, ContentLength: 64}); + + await expect(service.readObject('fluxer-downloads', 'desktop/stable/manifest.json', 32)).rejects.toThrow( + /exceeds maximum buffer size of 32 bytes \(got 64 bytes\)/u, + ); + expect(body.destroyed).toBe(true); + }); + + it('refuses an object that outgrows the cap mid-stream when no length is declared', async () => { + const service = new StorageService(); + const body = bodyStream([Buffer.alloc(24, 1), Buffer.alloc(24, 2)]); + serveGetObject(service, {Body: body}); + + await expect(service.readObject('fluxer-downloads', 'desktop/stable/manifest.json', 32)).rejects.toThrow( + /exceeds maximum buffer size of 32 bytes \(got at least 48 bytes\)/u, + ); + expect(body.destroyed).toBe(true); + }); + + it('returns every byte of an object that exactly fills the cap', async () => { + const service = new StorageService(); + serveGetObject(service, {Body: bodyStream([Buffer.alloc(16, 7), Buffer.alloc(16, 9)])}); + + await expect(service.readObject('fluxer-downloads', 'desktop/stable/manifest.json', 32)).resolves.toEqual( + new Uint8Array([...Array(16).fill(7), ...Array(16).fill(9)]), + ); + }); + + it('rejects rather than hanging when the object store resets the body mid-transfer', async () => { + const service = new StorageService(); + const reset = new Error('socket hang up'); + const body = new Readable({read() {}}); + Object.assign(service, {client: {send: async () => ({Body: body})}}); + const read = service.readObject('fluxer-downloads', 'desktop/stable/manifest.json', 1024); + body.push(Buffer.alloc(8, 1)); + setImmediate(() => body.destroy(reset)); + + await expect(read).rejects.toThrow(/socket hang up/u); + }, 5000); +}); + +interface ListObjectsPage { + Contents?: Array<{Key: string; LastModified?: Date}>; + IsTruncated?: boolean; + NextContinuationToken?: string; +} + +function serveListPages(service: StorageService, pages: Array): {tokens: Array} { + const tokens: Array = []; + let index = 0; + Object.assign(service, { + client: { + send: async (command: {input: {ContinuationToken?: string}}) => { + tokens.push(command.input.ContinuationToken); + const page = pages[index] ?? {}; + index += 1; + return page; + }, + }, + }); + return {tokens}; +} + +describe('StorageService.listObjects', () => { + it('returns every page of a truncated listing', async () => { + const service = new StorageService(); + const {tokens} = serveListPages(service, [ + { + Contents: [{Key: 'desktop/a.exe'}, {Key: 'desktop/b.exe'}], + IsTruncated: true, + NextContinuationToken: 'page-2', + }, + {Contents: [{Key: 'desktop/c.exe'}], IsTruncated: false}, + ]); + const objects = await service.listObjects({bucket: 'fluxer-downloads', prefix: 'desktop/'}); + expect(objects.map(({key}) => key)).toEqual(['desktop/a.exe', 'desktop/b.exe', 'desktop/c.exe']); + expect(tokens).toEqual([undefined, 'page-2']); + }); + + it('throws a typed overflow error when the prefix outgrows the requested cap', async () => { + const service = new StorageService(); + serveListPages(service, [ + { + Contents: [{Key: 'desktop/a.exe'}, {Key: 'desktop/b.exe'}], + IsTruncated: true, + NextContinuationToken: 'page-2', + }, + ]); + await expect( + service.listObjects({bucket: 'fluxer-downloads', prefix: 'desktop/', maxObjects: 2}), + ).rejects.toBeInstanceOf(StorageObjectListingOverflowError); + }); + + it('throws rather than returning a partial listing when the continuation token is missing', async () => { + const service = new StorageService(); + serveListPages(service, [{Contents: [{Key: 'desktop/a.exe'}], IsTruncated: true}]); + await expect(service.listObjects({bucket: 'fluxer-downloads', prefix: 'desktop/'})).rejects.toThrow( + /continuation token/u, + ); + }); +}); diff --git a/fluxer_api/src/api/infrastructure/StorageService.ts b/fluxer_api/src/api/infrastructure/StorageService.ts index 68a7e8ebc..231f0863a 100644 --- a/fluxer_api/src/api/infrastructure/StorageService.ts +++ b/fluxer_api/src/api/infrastructure/StorageService.ts @@ -32,7 +32,12 @@ import {seconds} from 'itty-time'; import {temporaryFile} from 'tempy'; import {Config} from '../Config'; import {Logger} from '../Logger'; -import type {IStorageService, ProcessedStorageObjectMetadata} from './IStorageService'; +import { + type IStorageService, + type ProcessedStorageObjectMetadata, + StorageObjectListingOverflowError, + StorageObjectRangeNotSatisfiableError, +} from './IStorageService'; import {processMediaFile} from './StorageObjectHelpers'; const STREAM_UPLOAD_PART_BYTES = 8 * 1024 * 1024; @@ -73,18 +78,13 @@ async function streamToUint8Array(body: Readable, maxBytes?: number): Promise maxBytes) { - const remaining = maxBytes - total; - if (remaining > 0) { - chunks.push(buf.subarray(0, remaining)); - total += remaining; - } - break; + body.destroy(); + throw new Error( + `Stream exceeds maximum buffer size of ${maxBytes} bytes (got at least ${total + buf.length} bytes)`, + ); } chunks.push(buf); total += buf.length; - if (maxBytes !== undefined && total >= maxBytes) { - break; - } } const out = new Uint8Array(total); let offset = 0; @@ -97,10 +97,15 @@ async function streamToUint8Array(body: Readable, maxBytes?: number): Promise undefined); + return wrapped; } export class StorageService implements IStorageService { @@ -362,7 +367,12 @@ export class StorageService implements IStorageService { async readObject(bucket: string, key: string, maxBytes?: number): Promise { const out = await this.client.send(new GetObjectCommand({Bucket: bucket, Key: key})); - return streamToUint8Array(extractStreamFromGet(out), maxBytes); + const body = extractStreamFromGet(out); + if (maxBytes !== undefined && out.ContentLength !== undefined && out.ContentLength > maxBytes) { + body.destroy(); + throw new Error(`Stream exceeds maximum buffer size of ${maxBytes} bytes (got ${out.ContentLength} bytes)`); + } + return streamToUint8Array(body, maxBytes); } async streamObject(params: {bucket: string; key: string; range?: string}): Promise<{ @@ -395,6 +405,13 @@ export class StorageService implements IStorageService { if (error instanceof S3ServiceException && (error.name === 'NoSuchKey' || error.name === 'NotFound')) { return null; } + if ( + params.range !== undefined && + error instanceof S3ServiceException && + (error.name === 'InvalidRange' || error.$metadata?.httpStatusCode === 416) + ) { + throw new StorageObjectRangeNotSatisfiableError(params.bucket, params.key, params.range); + } throw error; } } @@ -563,31 +580,49 @@ export class StorageService implements IStorageService { return this.client.send(new HeadObjectCommand({Bucket: params.bucket, Key: params.key})); } - async listObjects(params: {bucket: string; prefix: string}): Promise< + async listObjects(params: {bucket: string; prefix: string; maxObjects?: number}): Promise< ReadonlyArray<{ key: string; lastModified?: Date; }> > { + if (params.maxObjects !== undefined && (!Number.isSafeInteger(params.maxObjects) || params.maxObjects <= 0)) { + throw new RangeError('maxObjects must be a positive safe integer'); + } const result: Array<{ key: string; lastModified?: Date; }> = []; + let listedObjects = 0; let continuationToken: string | undefined; do { + const remaining = params.maxObjects === undefined ? undefined : params.maxObjects - listedObjects; const command = new ListObjectsV2Command({ Bucket: params.bucket, Prefix: params.prefix, ContinuationToken: continuationToken, + MaxKeys: remaining === undefined ? undefined : Math.min(remaining, 1000), }); const response = await this.client.send(command); if (response.Contents) { for (const obj of response.Contents) { + listedObjects += 1; if (obj.Key) { result.push({key: obj.Key, lastModified: obj.LastModified}); } } } + if ( + params.maxObjects !== undefined && + (listedObjects > params.maxObjects || (listedObjects === params.maxObjects && response.IsTruncated)) + ) { + throw new StorageObjectListingOverflowError(params.bucket, params.prefix, params.maxObjects); + } + if (response.IsTruncated && !response.NextContinuationToken) { + throw new Error( + `Truncated object listing omitted its continuation token for ${params.bucket}/${params.prefix}`, + ); + } continuationToken = response.IsTruncated ? response.NextContinuationToken : undefined; } while (continuationToken); return result; diff --git a/fluxer_api/src/api/search/MessageSearchResultReconciler.ts b/fluxer_api/src/api/search/MessageSearchResultReconciler.ts index 0ea6b2828..13fc17dbd 100644 --- a/fluxer_api/src/api/search/MessageSearchResultReconciler.ts +++ b/fluxer_api/src/api/search/MessageSearchResultReconciler.ts @@ -2,12 +2,17 @@ import type {SearchResult} from '@fluxer/schema/src/contracts/search/SearchAdapterTypes'; import type {MessageSearchFilters, SearchableMessage} from '@fluxer/schema/src/contracts/search/SearchDocumentTypes'; -import {createChannelID, createMessageID, type MessageID} from '../BrandedTypes'; +import {type ChannelID, createChannelID, createMessageID, type MessageID} from '../BrandedTypes'; import type {IMessageRepository} from '../channel/repositories/IMessageRepository'; +import {Logger} from '../Logger'; +import type {Message} from '../models/Message'; import type {IMessageSearchService} from './IMessageSearchService'; import {deleteMessageSearchDocuments} from './MessageSearchIndexCleanup'; const RECONCILE_BATCH_SIZE = 250; +const MAX_RECONCILE_PAGES = 40; +const MAX_STALE_DELETE_ABSOLUTE = 250; +const MAX_STALE_DELETE_RATIO = 0.5; interface MessageLookupRepository { readonly messages: Pick; @@ -26,6 +31,7 @@ interface SearchExistingMessagesParams { interface ValidatedHits { validHits: Array; staleMessageIds: Array; + lookupErrorCount: number; } export async function searchExistingMessages({ @@ -47,7 +53,9 @@ export async function searchExistingMessages({ return result; } if (cursor?.length) { - await deleteStaleSearchDocuments(searchService, validated.staleMessageIds); + if (validated.lookupErrorCount === 0) { + await deleteStaleSearchDocuments(searchService, validated.staleMessageIds, result.hits.length); + } return { ...result, hits: validated.validHits, @@ -75,18 +83,24 @@ async function reconcileOffsetSearchResult({ const requestedOffset = (page - 1) * hitsPerPage; const pageHits: Array = []; const staleMessageIds: Array = []; + let lookupErrorCount = 0; + let examinedCount = 0; let validTotal = 0; let rawOffset = 0; let rawPage = 1; - while (true) { + let corpusTotal = 0; + while (rawPage <= MAX_RECONCILE_PAGES) { const result = await searchService.searchMessages(query, filters, { hitsPerPage: RECONCILE_BATCH_SIZE, page: rawPage, }); + corpusTotal = result.total; if (result.hits.length === 0) { break; } const validated = await validateSearchHits(messageRepository, result.hits); + lookupErrorCount += validated.lookupErrorCount; + examinedCount += result.hits.length; staleMessageIds.push(...validated.staleMessageIds); for (const hit of validated.validHits) { if (validTotal >= requestedOffset && pageHits.length < hitsPerPage) { @@ -95,15 +109,20 @@ async function reconcileOffsetSearchResult({ validTotal += 1; } rawOffset += result.hits.length; + if (pageHits.length >= hitsPerPage && rawOffset >= requestedOffset) { + break; + } if (rawOffset >= result.total) { break; } rawPage += 1; } - await deleteStaleSearchDocuments(searchService, staleMessageIds); + if (lookupErrorCount === 0) { + await deleteStaleSearchDocuments(searchService, staleMessageIds, examinedCount); + } return { hits: pageHits, - total: validTotal, + total: Math.max(pageHits.length, corpusTotal - staleMessageIds.length), }; } @@ -113,26 +132,33 @@ async function validateSearchHits( ): Promise { const checked = await Promise.all( hits.map(async (hit) => { + let channelId: ChannelID; let messageId: MessageID; try { - const channelId = createChannelID(BigInt(hit.channelId)); + channelId = createChannelID(BigInt(hit.channelId)); messageId = createMessageID(BigInt(hit.id)); - const message = await messageRepository.messages.getMessage(channelId, messageId); - if (message && message.channelId.toString() === hit.channelId) { - return {hit, staleMessageId: null}; - } - } catch (_error) { - try { - messageId = createMessageID(BigInt(hit.id)); - } catch (_invalidMessageId) { - return {hit: null, staleMessageId: null}; - } + } catch (_invalidId) { + return {hit: null, staleMessageId: null, lookupError: false}; } - return {hit: null, staleMessageId: messageId}; + let message: Message | null; + try { + message = await messageRepository.messages.getMessage(channelId, messageId); + } catch (error) { + Logger.warn( + {error, messageId: hit.id, channelId: hit.channelId}, + 'Search read repair lookup failed; keeping document', + ); + return {hit, staleMessageId: null, lookupError: true}; + } + if (message && message.channelId.toString() === hit.channelId) { + return {hit, staleMessageId: null, lookupError: false}; + } + return {hit: null, staleMessageId: messageId, lookupError: false}; }), ); const validHits: Array = []; const staleMessageIds: Array = []; + let lookupErrorCount = 0; for (const item of checked) { if (item.hit) { validHits.push(item.hit); @@ -140,14 +166,31 @@ async function validateSearchHits( if (item.staleMessageId) { staleMessageIds.push(item.staleMessageId); } + if (item.lookupError) { + lookupErrorCount += 1; + } } - return {validHits, staleMessageIds}; + return {validHits, staleMessageIds, lookupErrorCount}; } async function deleteStaleSearchDocuments( searchService: IMessageSearchService, messageIds: Array, + examinedCount: number, ): Promise { + if (messageIds.length === 0) { + return; + } + if ( + messageIds.length > MAX_STALE_DELETE_ABSOLUTE || + (examinedCount > 0 && messageIds.length / examinedCount > MAX_STALE_DELETE_RATIO) + ) { + Logger.warn( + {staleMessageCount: messageIds.length, examinedCount}, + 'Search read repair delete exceeded safety cap; skipping delete', + ); + return; + } await deleteMessageSearchDocuments(messageIds, { searchService, context: {source: 'message_search_read_repair', staleMessageCount: messageIds.length}, diff --git a/fluxer_api/src/api/test/mocks/MockStorageService.ts b/fluxer_api/src/api/test/mocks/MockStorageService.ts index 04ba3e5c3..63543d291 100644 --- a/fluxer_api/src/api/test/mocks/MockStorageService.ts +++ b/fluxer_api/src/api/test/mocks/MockStorageService.ts @@ -6,7 +6,25 @@ import {Readable} from 'node:stream'; import {S3ServiceException} from '@aws-sdk/client-s3'; import {isSupportedMediaContentType} from '@pkgs/mime_utils/src/ContentTypeUtils'; import {vi} from 'vitest'; -import type {IStorageService, ProcessedStorageObjectMetadata} from '../../infrastructure/IStorageService'; +import {Config} from '../../Config'; +import { + type IStorageService, + type ProcessedStorageObjectMetadata, + StorageObjectListingOverflowError, + StorageObjectRangeNotSatisfiableError, +} from '../../infrastructure/IStorageService'; + +const OBJECT_ID_SEPARATOR = '\u0000'; +const BYTE_RANGE_PATTERN = /^bytes=(\d*)-(\d*)$/u; + +interface StoredObject { + data: Uint8Array; + contentType?: string; + etag: string; + lastModified: Date; +} + +type ParsedByteRange = {start: number; end: number} | 'ignored' | 'unsatisfiable'; interface MockStorageServiceConfig { fileData?: Uint8Array | null; @@ -17,14 +35,40 @@ interface MockStorageServiceConfig { shouldFailCopy?: boolean; } +function computeEtag(data: Uint8Array): string { + return `"${createHash('md5').update(data).digest('hex')}"`; +} + +function parseByteRange(range: string, totalLength: number): ParsedByteRange { + const match = BYTE_RANGE_PATTERN.exec(range.trim()); + if (!match) return 'ignored'; + const rawStart = match[1] ?? ''; + const rawEnd = match[2] ?? ''; + if (rawStart === '' && rawEnd === '') return 'ignored'; + if (rawStart === '') { + const suffixLength = Number(rawEnd); + if (suffixLength === 0 || totalLength === 0) return 'unsatisfiable'; + return {start: Math.max(totalLength - suffixLength, 0), end: totalLength - 1}; + } + const start = Number(rawStart); + if (start >= totalLength) return 'unsatisfiable'; + if (rawEnd === '') return {start, end: totalLength - 1}; + const end = Number(rawEnd); + if (end < start) return 'ignored'; + return {start, end: Math.min(end, totalLength - 1)}; +} + +function noSuchKeyError(key: string): S3ServiceException { + return new S3ServiceException({ + name: 'NoSuchKey', + $fault: 'client', + $metadata: {}, + message: `The specified key does not exist: ${key}`, + }); +} + export class MockStorageService implements IStorageService { - private objects: Map< - string, - { - data: Uint8Array; - contentType?: string; - } - > = new Map(); + private objects: Map = new Map(); private multipartUploads: Map< string, { @@ -77,6 +121,19 @@ export class MockStorageService implements IStorageService { this.config = {...this.config, ...config}; } + private objectId(bucket: string, key: string): string { + return `${bucket}${OBJECT_ID_SEPARATOR}${key}`; + } + + private storeObject(bucket: string, key: string, data: Uint8Array, contentType?: string): void { + this.objects.set(this.objectId(bucket, key), { + data, + contentType, + etag: computeEtag(data), + lastModified: new Date(), + }); + } + async uploadObject(params: { bucket: string; key: string; @@ -89,7 +146,7 @@ export class MockStorageService implements IStorageService { throw new Error('Mock storage upload failure'); } const data = params.body instanceof Uint8Array ? params.body : await this.readableToBuffer(params.body); - this.objects.set(params.key, {data, contentType: params.contentType}); + this.storeObject(params.bucket, params.key, data, params.contentType); } async uploadObjectFromFile(params: { @@ -105,7 +162,7 @@ export class MockStorageService implements IStorageService { throw new Error('Mock storage upload failure'); } const data = await fs.promises.readFile(params.filePath); - this.objects.set(params.key, {data: new Uint8Array(data), contentType: params.contentType}); + this.storeObject(params.bucket, params.key, new Uint8Array(data), params.contentType); } private async readableToBuffer(stream: Readable): Promise { @@ -129,7 +186,7 @@ export class MockStorageService implements IStorageService { throw new Error('Mock storage delete failure'); } this.deletedObjects.push({bucket, key}); - this.objects.delete(key); + this.objects.delete(this.objectId(bucket, key)); } async getObjectMetadata( @@ -138,16 +195,23 @@ export class MockStorageService implements IStorageService { ): Promise<{ contentLength: number; contentType: string; + etag?: string; + lastModified?: Date; } | null> { this.getObjectMetadataSpy(bucket, key); - const obj = this.objects.get(key); + const obj = this.objects.get(this.objectId(bucket, key)); if (!obj) return null; - return {contentLength: obj.data.length, contentType: obj.contentType ?? 'application/octet-stream'}; + return { + contentLength: obj.data.length, + contentType: obj.contentType ?? 'application/octet-stream', + etag: obj.etag, + lastModified: obj.lastModified, + }; } async computeObjectSha256(bucket: string, key: string): Promise { this.computeObjectSha256Spy(bucket, key); - const data = this.config.fileData ?? this.objects.get(key)?.data ?? new Uint8Array(); + const data = this.config.fileData ?? this.objects.get(this.objectId(bucket, key))?.data ?? new Uint8Array(); return createHash('sha256').update(data).digest('hex'); } @@ -164,25 +228,13 @@ export class MockStorageService implements IStorageService { } if (this.config.fileData !== undefined) { if (this.config.fileData === null) { - const error = new S3ServiceException({ - name: 'NoSuchKey', - $fault: 'client', - $metadata: {}, - message: `The specified key does not exist: ${key}`, - }); - throw error; + throw noSuchKeyError(key); } return assertWithinLimit(this.config.fileData); } - const obj = this.objects.get(key); + const obj = this.objects.get(this.objectId(bucket, key)); if (!obj) { - const error = new S3ServiceException({ - name: 'NoSuchKey', - $fault: 'client', - $metadata: {}, - message: `The specified key does not exist: ${key}`, - }); - throw error; + throw noSuchKeyError(key); } return assertWithinLimit(obj.data); } @@ -205,35 +257,41 @@ export class MockStorageService implements IStorageService { if (this.config.fileData === null) { return null; } - const obj = this.objects.get(params.key); + const obj = this.objects.get(this.objectId(params.bucket, params.key)); const data = this.config.fileData ?? obj?.data; if (!data) { return null; } + let slice = data; + let contentRange: string | null = null; + if (params.range !== undefined) { + const parsedRange = parseByteRange(params.range, data.length); + if (parsedRange === 'unsatisfiable') { + throw new StorageObjectRangeNotSatisfiableError(params.bucket, params.key, params.range); + } + if (parsedRange !== 'ignored') { + slice = data.subarray(parsedRange.start, parsedRange.end + 1); + contentRange = `bytes ${parsedRange.start}-${parsedRange.end}/${data.length}`; + } + } return { - body: Readable.from([Buffer.from(data)]), - contentLength: data.length, - contentRange: null, + body: Readable.from([Buffer.from(slice)]), + contentLength: slice.length, + contentRange, contentType: obj?.contentType ?? 'application/octet-stream', cacheControl: null, contentDisposition: null, expires: null, - etag: `"${createHash('md5').update(data).digest('hex')}"`, - lastModified: null, + etag: computeEtag(data), + lastModified: obj?.lastModified ?? null, }; } async writeObjectToDisk(bucket: string, key: string, filePath: string): Promise { this.writeObjectToDiskSpy(bucket, key, filePath); - const data = this.config.fileData ?? this.objects.get(key)?.data; + const data = this.config.fileData ?? this.objects.get(this.objectId(bucket, key))?.data; if (!data) { - const error = new S3ServiceException({ - name: 'NoSuchKey', - $fault: 'client', - $metadata: {}, - message: `The specified key does not exist: ${key}`, - }); - throw error; + throw noSuchKeyError(key); } await fs.promises.writeFile(filePath, data); } @@ -255,12 +313,14 @@ export class MockStorageService implements IStorageService { destinationBucket: params.destinationBucket, destinationKey: params.destinationKey, }); - const sourceObj = this.objects.get(params.sourceKey); + const sourceObj = this.objects.get(this.objectId(params.sourceBucket, params.sourceKey)); if (sourceObj) { - this.objects.set(params.destinationKey, { - data: sourceObj.data, - contentType: params.newContentType ?? sourceObj.contentType, - }); + this.storeObject( + params.destinationBucket, + params.destinationKey, + sourceObj.data, + params.newContentType ?? sourceObj.contentType, + ); } } @@ -283,7 +343,8 @@ export class MockStorageService implements IStorageService { if (!isSupportedMediaContentType(params.contentType)) { return null; } - const data = this.objects.get(params.destinationKey)?.data ?? new Uint8Array(); + const data = + this.objects.get(this.objectId(params.destinationBucket, params.destinationKey))?.data ?? new Uint8Array(); return { contentType: params.contentType, contentLength: data.length, @@ -338,37 +399,63 @@ export class MockStorageService implements IStorageService { return `https://presigned-upload.url/test?partNumber=${params.partNumber}&uploadId=${params.uploadId}`; } - async purgeBucket(_bucket: string): Promise { - this.purgeBucketSpy(_bucket); + async purgeBucket(bucket: string): Promise { + this.purgeBucketSpy(bucket); + const prefix = this.objectId(bucket, ''); + for (const id of [...this.objects.keys()]) { + if (id.startsWith(prefix)) { + this.objects.delete(id); + this.deletedObjects.push({bucket, key: id.slice(prefix.length)}); + } + } } async uploadAvatar(params: {prefix: string; key: string; body: Uint8Array}): Promise { this.uploadAvatarSpy(params); - await this.uploadObject({bucket: 'cdn', key: `${params.prefix}/${params.key}`, body: params.body}); + await this.uploadObject({ + bucket: Config.s3.buckets.cdn, + key: `${params.prefix}/${params.key}`, + body: params.body, + }); } async deleteAvatar(params: {prefix: string; key: string}): Promise { this.deleteAvatarSpy(params); - await this.deleteObject('cdn', `${params.prefix}/${params.key}`); + await this.deleteObject(Config.s3.buckets.cdn, `${params.prefix}/${params.key}`); } - async listObjects(_params: {bucket: string; prefix: string}): Promise< + async listObjects(params: {bucket: string; prefix: string; maxObjects?: number}): Promise< ReadonlyArray<{ key: string; lastModified?: Date; }> > { - this.listObjectsSpy(_params); - return []; + this.listObjectsSpy(params); + if (params.maxObjects !== undefined && (!Number.isSafeInteger(params.maxObjects) || params.maxObjects <= 0)) { + throw new RangeError('maxObjects must be a positive safe integer'); + } + const keyOffset = params.bucket.length + OBJECT_ID_SEPARATOR.length; + const matches = [...this.objects.entries()] + .filter(([id]) => id.startsWith(this.objectId(params.bucket, params.prefix))) + .map(([id, object]) => ({key: id.slice(keyOffset), lastModified: object.lastModified})) + .sort((a, b) => (a.key < b.key ? -1 : a.key > b.key ? 1 : 0)); + if (params.maxObjects !== undefined && matches.length > params.maxObjects) { + throw new StorageObjectListingOverflowError(params.bucket, params.prefix, params.maxObjects); + } + return matches; } - async deleteObjects(_params: { + async deleteObjects(params: { bucket: string; objects: ReadonlyArray<{ Key: string; }>; }): Promise { - this.deleteObjectsSpy(_params); + this.deleteObjectsSpy(params); + for (const object of params.objects) { + this.objects.delete(this.objectId(params.bucket, object.Key)); + this.deletedObjects.push({bucket: params.bucket, key: object.Key}); + } } async createMultipartUpload(params: {bucket: string; key: string; contentType?: string}): Promise<{ @@ -440,7 +527,7 @@ export class MockStorageService implements IStorageService { combined.set(data, offset); offset += data.length; } - this.objects.set(upload.key, {data: combined}); + this.storeObject(upload.bucket, upload.key, combined); this.multipartUploads.delete(params.uploadId); } @@ -465,8 +552,8 @@ export class MockStorageService implements IStorageService { return [...this.copiedObjects]; } - hasObject(_bucket: string, key: string): boolean { - return this.objects.has(key); + hasObject(bucket: string, key: string): boolean { + return this.objects.has(this.objectId(bucket, key)); } reset(): void { diff --git a/fluxer_api/src/api/user/controllers/UserContentController.ts b/fluxer_api/src/api/user/controllers/UserContentController.ts index 835da243f..68aa9b19f 100644 --- a/fluxer_api/src/api/user/controllers/UserContentController.ts +++ b/fluxer_api/src/api/user/controllers/UserContentController.ts @@ -18,6 +18,7 @@ import { } from '@fluxer/schema/src/domains/user/UserRequestSchemas'; import {SavedMessageEntryListResponse} from '@fluxer/schema/src/domains/user/UserResponseSchemas'; import {createChannelID, createMessageID} from '../../BrandedTypes'; +import {StorageObjectRangeNotSatisfiableError} from '../../infrastructure/IStorageService'; import {DefaultUserOnly, LoginRequired} from '../../middleware/AuthMiddleware'; import {RateLimitMiddleware} from '../../middleware/RateLimitMiddleware'; import {OpenAPI} from '../../middleware/ResponseTypeMiddleware'; @@ -290,28 +291,35 @@ export function UserContentController(app: HonoApp) { if (!token) { return ctx.text('Not Found', 404); } - const result = await ctx.get('userContentRequestService').streamHarvestDownload({ - harvestId, - token, - range: ctx.req.header('range') ?? undefined, - storageService: ctx.get('storageService'), - }); - if (!result) { - return ctx.text('Not Found', 404); + try { + const result = await ctx.get('userContentRequestService').streamHarvestDownload({ + harvestId, + token, + range: ctx.req.header('range') ?? undefined, + storageService: ctx.get('storageService'), + }); + if (!result) { + return ctx.text('Not Found', 404); + } + const headers = new Headers(); + headers.set('Content-Type', result.contentType ?? 'application/zip'); + headers.set('Content-Disposition', `attachment; filename="${encodeURIComponent(result.filename)}"`); + headers.set('Content-Length', String(result.contentLength)); + headers.set('Cache-Control', 'private, no-store'); + headers.set('Accept-Ranges', 'bytes'); + if (result.contentRange) { + headers.set('Content-Range', result.contentRange); + } + return new Response(Readable.toWeb(result.body) as ReadableStream, { + status: result.contentRange ? 206 : 200, + headers, + }); + } catch (error) { + if (error instanceof StorageObjectRangeNotSatisfiableError) { + return ctx.text('Range Not Satisfiable', 416); + } + throw error; } - const headers = new Headers(); - headers.set('Content-Type', result.contentType ?? 'application/zip'); - headers.set('Content-Disposition', `attachment; filename="${encodeURIComponent(result.filename)}"`); - headers.set('Content-Length', String(result.contentLength)); - headers.set('Cache-Control', 'private, no-store'); - headers.set('Accept-Ranges', 'bytes'); - if (result.contentRange) { - headers.set('Content-Range', result.contentRange); - } - return new Response(Readable.toWeb(result.body) as ReadableStream, { - status: result.contentRange ? 206 : 200, - headers, - }); }, ); app.get( diff --git a/packages/schema/src/domains/message/MessageRequestSchemas.ts b/packages/schema/src/domains/message/MessageRequestSchemas.ts index 0a4c21e78..8db65f59a 100644 --- a/packages/schema/src/domains/message/MessageRequestSchemas.ts +++ b/packages/schema/src/domains/message/MessageRequestSchemas.ts @@ -217,9 +217,9 @@ export const MessageSearchRequest = z.object({ .number() .int() .min(1) - .max(Number.MAX_SAFE_INTEGER) + .max(400) .default(1) - .describe('Page number for pagination (ignored when cursor is provided)'), + .describe('Page number for pagination (ignored when cursor is provided). Use cursor to page beyond this.'), cursor: z .array(z.string()) .optional()