diff --git a/piece-retriever/bin/piece-retriever.js b/piece-retriever/bin/piece-retriever.js index 66cbcae0..c197554b 100644 --- a/piece-retriever/bin/piece-retriever.js +++ b/piece-retriever/bin/piece-retriever.js @@ -5,7 +5,7 @@ import { measureStreamedEgress, } from '../lib/retrieval.js' import { - getStorageProviderAndValidatePayer, + getRetrievalCandidatesAndValidatePayer, logRetrievalResult, updateDataSetStats, } from '../lib/store.js' @@ -71,18 +71,17 @@ export default { // Timestamp to measure file retrieval performance (from cache and from SP) const fetchStartedAt = performance.now() - const [{ serviceProviderId, serviceUrl, dataSetId }, isBadBit] = - await Promise.all([ - getStorageProviderAndValidatePayer( - env, - payerWalletAddress, - pieceCid, - env.ENFORCE_EGRESS_QUOTA, - ), - env.BAD_BITS_KV.get(`bad-bits:${await getBadBitsEntry(pieceCid)}`, { - type: 'json', - }), - ]) + const [retrievalCandidates, isBadBit] = await Promise.all([ + getRetrievalCandidatesAndValidatePayer( + env, + payerWalletAddress, + pieceCid, + env.ENFORCE_EGRESS_QUOTA, + ), + env.BAD_BITS_KV.get(`bad-bits:${await getBadBitsEntry(pieceCid)}`, { + type: 'json', + }), + ]) httpAssert( !isBadBit, @@ -91,23 +90,55 @@ export default { ) httpAssert( - serviceProviderId, - 404, - `Unsupported Service Provider: ${serviceProviderId}`, + retrievalCandidates.length > 0, + 500, + 'Service provider lookup failed', ) + let retrievalCandidate let retrievalResult + const retrievalAttempts = [] - try { - retrievalResult = await retrieveFile( - ctx, - serviceUrl, - pieceCid, - request, - env.ORIGIN_CACHE_TTL, - { signal: request.signal }, + while (retrievalCandidates.length > 0) { + const retrievalCandidateIndex = Math.floor( + Math.random() * retrievalCandidates.length, ) - } catch {} + retrievalCandidate = retrievalCandidates[retrievalCandidateIndex] + retrievalAttempts.push(retrievalCandidate) + retrievalCandidates.splice(retrievalCandidateIndex, 1) + console.log('Attempting retrieval', retrievalCandidate) + try { + retrievalResult = await retrieveFile( + ctx, + retrievalCandidate.serviceUrl, + pieceCid, + request, + env.ORIGIN_CACHE_TTL, + { signal: request.signal }, + ) + if (retrievalResult.response.ok) { + break + } + console.log( + `Retrieval attempt failed: HTTP ${retrievalResult.response.status}`, + { + retrievalCandidate, + willRetry: retrievalCandidates.length > 0, + }, + ) + } catch (err) { + const msg = + typeof err === 'object' && err !== null && 'message' in err + ? err.message + : String(err) + console.log(`Retrieval attempt failed: ${msg}`, { + retrievalCandidate, + willRetry: retrievalCandidates.length > 0, + }) + } + } + + httpAssert(retrievalCandidate, 500, 'should never happen') if (!retrievalResult || retrievalResult.response.status >= 500) { ctx.waitUntil( @@ -117,16 +148,18 @@ export default { egressBytes: 0, requestCountryCode, timestamp: requestTimestamp, - dataSetId, + dataSetId: retrievalCandidate.dataSetId, botName, }), ) const response = new Response( - `Service provider ${serviceProviderId} is unavailable${retrievalResult ? ` at ${retrievalResult.url}` : ''}`, + `No available service provider found. Attempted: ${retrievalAttempts.map((a) => `ID=${a.serviceProviderId} (Service URL=${a.serviceUrl})`).join(', ')}`, { status: 502, headers: new Headers({ - 'X-Data-Set-ID': dataSetId, + 'X-Data-Set-ID': retrievalAttempts + .map((a) => a.dataSetId) + .join(','), }), }, ) @@ -145,7 +178,7 @@ export default { egressBytes: 0, requestCountryCode, timestamp: requestTimestamp, - dataSetId, + dataSetId: retrievalCandidate.dataSetId, botName, }), ) @@ -154,7 +187,7 @@ export default { retrievalResult.response, ) setContentSecurityPolicy(response) - response.headers.set('X-Data-Set-ID', dataSetId) + response.headers.set('X-Data-Set-ID', retrievalCandidate.dataSetId) response.headers.set( 'Cache-Control', `public, max-age=${env.CLIENT_CACHE_TTL}`, @@ -185,12 +218,12 @@ export default { fetchTtlb: lastByteFetchedAt - fetchStartedAt, workerTtfb: firstByteAt - workerStartedAt, }, - dataSetId, + dataSetId: retrievalCandidate.dataSetId, botName, }) await updateDataSetStats(env, { - dataSetId, + dataSetId: retrievalCandidate.dataSetId, egressBytes, cacheMiss: retrievalResult.cacheMiss, enforceEgressQuota: env.ENFORCE_EGRESS_QUOTA, @@ -205,7 +238,7 @@ export default { headers: retrievalResult.response.headers, }) setContentSecurityPolicy(response) - response.headers.set('X-Data-Set-ID', dataSetId) + response.headers.set('X-Data-Set-ID', retrievalCandidate.dataSetId) response.headers.set( 'Cache-Control', `public, max-age=${env.CLIENT_CACHE_TTL}`, diff --git a/piece-retriever/lib/store.js b/piece-retriever/lib/store.js index 64fabb47..d3dba0b6 100644 --- a/piece-retriever/lib/store.js +++ b/piece-retriever/lib/store.js @@ -86,15 +86,17 @@ export async function logRetrievalResult(env, params) { * @param {string} pieceCid - The piece CID to look up * @param {boolean} [enforceEgressQuota=false] - Whether to enforce egress quota * limits. Default is `false` - * @returns {Promise<{ - * serviceProviderId: string - * serviceUrl: string - * dataSetId: string - * cdnEgressQuota: bigint - * cacheMissEgressQuota: bigint - * }>} + * @returns {Promise< + * { + * serviceProviderId: string + * serviceUrl: string + * dataSetId: string + * cdnEgressQuota: bigint + * cacheMissEgressQuota: bigint + * }[] + * >} */ -export async function getStorageProviderAndValidatePayer( +export async function getRetrievalCandidatesAndValidatePayer( env, payerAddress, pieceCid, @@ -206,29 +208,21 @@ export async function getStorageProviderAndValidatePayer( `Cache miss egress quota exhausted for payer '${payerAddress}' and data set '${withSufficientCDNQuota[0]?.data_set_id}'. Please top up your cache miss egress quota.`, ) - const { - data_set_id: dataSetId, - service_provider_id: serviceProviderId, - service_url: serviceUrl, - cdn_egress_quota: cdnEgressQuota, - cache_miss_egress_quota: cacheMissEgressQuota, - } = pickRandom(withSufficientCacheMissQuota) - - // We need this assertion to supress TypeScript error. The compiler is not able to infer that - // `withCDN.filter()` above returns only rows with `service_url` defined. - httpAssert(serviceUrl, 500, 'should never happen') + const retrievalCandidates = withSufficientCacheMissQuota.map((row) => ({ + dataSetId: row.data_set_id, + serviceProviderId: row.service_provider_id, + // We need this cast to supress a TypeScript error. The compiler is not able to infer that + // `withCDN.filter()` above returns only rows with `service_url` defined. + serviceUrl: /** @type {string} */ (row.service_url), + cdnEgressQuota: BigInt(row.cdn_egress_quota ?? '0'), + cacheMissEgressQuota: BigInt(row.cache_miss_egress_quota ?? '0'), + })) console.log( - `Looked up Data set ID '${dataSetId}' and service provider id '${serviceProviderId}' for piece_cid '${pieceCid}' and payer '${payerAddress}'. Service URL: ${serviceUrl}`, + `Looked up ${retrievalCandidates.length} retrieval candidates for piece_cid '${pieceCid}' and payer '${payerAddress}'`, ) - return { - serviceProviderId, - serviceUrl, - dataSetId, - cdnEgressQuota: BigInt(cdnEgressQuota ?? '0'), - cacheMissEgressQuota: BigInt(cacheMissEgressQuota ?? '0'), - } + return retrievalCandidates } /** @@ -266,12 +260,3 @@ export async function updateDataSetStats( ) .run() } - -/** - * @template T - * @param {T[]} arr - * @returns {T} - */ -function pickRandom(arr) { - return arr[Math.floor(Math.random() * arr.length)] -} diff --git a/piece-retriever/test/retriever.test.js b/piece-retriever/test/retriever.test.js index 8a92719c..4b59391d 100644 --- a/piece-retriever/test/retriever.test.js +++ b/piece-retriever/test/retriever.test.js @@ -933,7 +933,9 @@ describe('piece-retriever.fetch', () => { }) await waitOnExecutionContext(ctx) expect(res.status).toBe(502) - expect(await res.text()).toMatch(/^Service provider \d+ is unavailable at /) + expect(await res.text()).toMatch( + /^No available service provider found. Attempted: ID=/, + ) expect(res.headers.get('X-Data-Set-ID')).toBe(String(dataSetId)) const result = await env.DB.prepare( @@ -1063,7 +1065,9 @@ describe('piece-retriever.fetch', () => { }) await waitOnExecutionContext(ctx) expect(res.status).toBe(502) - expect(await res.text()).toMatch(/^Service provider \d+ is unavailable$/) + expect(await res.text()).toMatch( + /^No available service provider found. Attempted: ID=/, + ) expect(res.headers.get('X-Data-Set-ID')).toBe(String(dataSetId)) }) }) diff --git a/piece-retriever/test/store.test.js b/piece-retriever/test/store.test.js index 074e7fec..abc1e589 100644 --- a/piece-retriever/test/store.test.js +++ b/piece-retriever/test/store.test.js @@ -1,7 +1,7 @@ import { describe, it, beforeAll, expect } from 'vitest' import { logRetrievalResult, - getStorageProviderAndValidatePayer, + getRetrievalCandidatesAndValidatePayer, updateDataSetStats, } from '../lib/store.js' import { env } from 'cloudflare:test' @@ -48,7 +48,7 @@ describe('logRetrievalResult', () => { }) }) -describe('getStorageProviderAndValidatePayer', () => { +describe('getRetrievalCandidatesAndValidatePayer', () => { const APPROVED_SERVICE_PROVIDER_ID = '20' beforeAll(async () => { await withApprovedProvider(env, { @@ -73,19 +73,21 @@ describe('getStorageProviderAndValidatePayer', () => { pieceId: 'piece-1', }) - const result = await getStorageProviderAndValidatePayer( + const results = await getRetrievalCandidatesAndValidatePayer( env, payerAddress, pieceCid, true, ) - expect(result.serviceProviderId).toBe(APPROVED_SERVICE_PROVIDER_ID) + expect(results).toMatchObject([ + { serviceProviderId: APPROVED_SERVICE_PROVIDER_ID }, + ]) }) it('throws error if pieceCid not found', async () => { const payerAddress = '0x1234567890abcdef1234567890abcdef12345678' await expect( - getStorageProviderAndValidatePayer( + getRetrievalCandidatesAndValidatePayer( env, payerAddress, 'nonexistent-cid', @@ -102,7 +104,7 @@ describe('getStorageProviderAndValidatePayer', () => { await withPiece(env, { pieceId: 'piece-no-sp', dataSetId, pieceCid }) await expect( - getStorageProviderAndValidatePayer(env, payerAddress, pieceCid, true), + getRetrievalCandidatesAndValidatePayer(env, payerAddress, pieceCid, true), ).rejects.toThrow(/no associated service provider/) }) @@ -124,7 +126,7 @@ describe('getStorageProviderAndValidatePayer', () => { }) await expect( - getStorageProviderAndValidatePayer(env, payerAddress, pieceCid, true), + getRetrievalCandidatesAndValidatePayer(env, payerAddress, pieceCid, true), ).rejects.toThrow( /There is no Filecoin Warm Storage Service deal for payer/, ) @@ -148,11 +150,11 @@ describe('getStorageProviderAndValidatePayer', () => { }) await expect( - getStorageProviderAndValidatePayer(env, payerAddress, pieceCid, true), + getRetrievalCandidatesAndValidatePayer(env, payerAddress, pieceCid, true), ).rejects.toThrow(/withCDN=false/) }) - it('returns serviceProviderId for approved service provider', async () => { + it('returns serviceProviderId for approved service providers', async () => { const pieceCid = 'cid-approved' const dataSetId = 'data-set-approved' const payerAddress = '0xabcdef1234567890abcdef1234567890abcdef12' @@ -168,16 +170,18 @@ describe('getStorageProviderAndValidatePayer', () => { pieceCid, }) - const result = await getStorageProviderAndValidatePayer( + const results = await getRetrievalCandidatesAndValidatePayer( env, payerAddress, pieceCid, true, ) - expect(result.serviceProviderId).toBe(APPROVED_SERVICE_PROVIDER_ID) + expect(results).toMatchObject([ + { serviceProviderId: APPROVED_SERVICE_PROVIDER_ID }, + ]) }) - it('returns a random service provider when multiple service providers share the same pieceCid', async () => { + it('returns multiple service providers when they share the same pieceCid', async () => { const dataSetId1 = 'data-set-a' const dataSetId2 = 'data-set-b' const pieceCid = 'shared-piece-cid' @@ -215,20 +219,17 @@ describe('getStorageProviderAndValidatePayer', () => { payerAddress, }) - const serviceProviderIdsReturned = new Set() - for (let i = 0; i < 100; i++) { - const result = await getStorageProviderAndValidatePayer( - env, - payerAddress, - pieceCid, - true, - ) - serviceProviderIdsReturned.add(result.serviceProviderId) - if (serviceProviderIdsReturned.size === 2) { - return - } - } - throw new Error('Did not return 2 different SPs') + const results = await getRetrievalCandidatesAndValidatePayer( + env, + payerAddress, + pieceCid, + true, + ) + + expect(results).toMatchObject([ + { serviceProviderId: serviceProviderId1 }, + { serviceProviderId: serviceProviderId2 }, + ]) }) it('ignores owners that are not approved by Filecoin Warm Storage Service', async () => { @@ -268,19 +269,21 @@ describe('getStorageProviderAndValidatePayer', () => { }) // Should return service provider 1 because service provider 2 is not approved - const result = await getStorageProviderAndValidatePayer( + const results = await getRetrievalCandidatesAndValidatePayer( env, payerAddress, pieceCid, true, ) - expect(result).toEqual({ - dataSetId: dataSetId1, - serviceProviderId: serviceProviderId1.toLowerCase(), - serviceUrl: 'https://pdp-provider-1.xyz', - cdnEgressQuota: 100n, - cacheMissEgressQuota: 100n, - }) + expect(results).toMatchObject([ + { + dataSetId: dataSetId1, + serviceProviderId: serviceProviderId1.toLowerCase(), + serviceUrl: 'https://pdp-provider-1.xyz', + cdnEgressQuota: 100n, + cacheMissEgressQuota: 100n, + }, + ]) }) }) @@ -311,7 +314,7 @@ describe('Egress Quota Management', () => { }) await expect( - getStorageProviderAndValidatePayer(env, payerAddress, pieceCid, true), + getRetrievalCandidatesAndValidatePayer(env, payerAddress, pieceCid, true), ).rejects.toThrow(/CDN egress quota exhausted for payer/) }) @@ -332,7 +335,7 @@ describe('Egress Quota Management', () => { }) await expect( - getStorageProviderAndValidatePayer(env, payerAddress, pieceCid, true), + getRetrievalCandidatesAndValidatePayer(env, payerAddress, pieceCid, true), ).rejects.toThrow(/Cache miss egress quota exhausted for payer/) }) @@ -352,19 +355,21 @@ describe('Egress Quota Management', () => { pieceId: 'piece-sufficient', }) - const result = await getStorageProviderAndValidatePayer( + const results = await getRetrievalCandidatesAndValidatePayer( env, payerAddress, pieceCid, true, ) - expect(result).toStrictEqual({ - dataSetId, - serviceProviderId: APPROVED_SERVICE_PROVIDER_ID, - serviceUrl: 'https://quota-test-provider.xyz', - cdnEgressQuota: 1n, - cacheMissEgressQuota: 1n, - }) + expect(results).toMatchObject([ + { + dataSetId, + serviceProviderId: APPROVED_SERVICE_PROVIDER_ID, + serviceUrl: 'https://quota-test-provider.xyz', + cdnEgressQuota: 1n, + cacheMissEgressQuota: 1n, + }, + ]) }) it('correctly decrements CDN quota on cache hit', async () => { @@ -523,7 +528,7 @@ describe('Egress Quota Management', () => { // Should return 402 error when quotas are null await expect( - getStorageProviderAndValidatePayer(env, payerAddress, pieceCid, true), + getRetrievalCandidatesAndValidatePayer(env, payerAddress, pieceCid, true), ).rejects.toThrow(/CDN egress quota exhausted for payer/) }) @@ -545,19 +550,21 @@ describe('Egress Quota Management', () => { }) // Quota of exactly 100 should be sufficient (> 0) - const result = await getStorageProviderAndValidatePayer( + const results = await getRetrievalCandidatesAndValidatePayer( env, payerAddress, pieceCid, true, ) - expect(result).toStrictEqual({ - dataSetId, - serviceProviderId: APPROVED_SERVICE_PROVIDER_ID, - serviceUrl: 'https://quota-test-provider.xyz', - cdnEgressQuota: 100n, - cacheMissEgressQuota: 100n, - }) + expect(results).toMatchObject([ + { + dataSetId, + serviceProviderId: APPROVED_SERVICE_PROVIDER_ID, + serviceUrl: 'https://quota-test-provider.xyz', + cdnEgressQuota: 100n, + cacheMissEgressQuota: 100n, + }, + ]) // Decrement by exact amount should result in 0 await updateDataSetStats(env, { @@ -593,19 +600,21 @@ describe('Egress Quota Management', () => { pieceId: 'piece-no-enforce-cdn-exhausted', }) - const result = await getStorageProviderAndValidatePayer( + const results = await getRetrievalCandidatesAndValidatePayer( env, payerAddress, pieceCid, false, ) - expect(result).toStrictEqual({ - dataSetId, - serviceProviderId: APPROVED_SERVICE_PROVIDER_ID, - serviceUrl: 'https://quota-test-provider.xyz', - cdnEgressQuota: 0n, - cacheMissEgressQuota: 1n, - }) + expect(results).toMatchObject([ + { + dataSetId, + serviceProviderId: APPROVED_SERVICE_PROVIDER_ID, + serviceUrl: 'https://quota-test-provider.xyz', + cdnEgressQuota: 0n, + cacheMissEgressQuota: 1n, + }, + ]) }) it('allows retrieval when quota enforcement is disabled and cache-miss quota is exhausted', async () => { @@ -624,19 +633,21 @@ describe('Egress Quota Management', () => { pieceId: 'piece-no-enforce-cache-miss-exhausted', }) - const result = await getStorageProviderAndValidatePayer( + const results = await getRetrievalCandidatesAndValidatePayer( env, payerAddress, pieceCid, false, ) - expect(result).toStrictEqual({ - dataSetId, - serviceProviderId: APPROVED_SERVICE_PROVIDER_ID, - serviceUrl: 'https://quota-test-provider.xyz', - cdnEgressQuota: 1n, - cacheMissEgressQuota: 0n, - }) + expect(results).toMatchObject([ + { + dataSetId, + serviceProviderId: APPROVED_SERVICE_PROVIDER_ID, + serviceUrl: 'https://quota-test-provider.xyz', + cdnEgressQuota: 1n, + cacheMissEgressQuota: 0n, + }, + ]) }) it('allows retrieval when quota enforcement is disabled and both quotas are exhausted', async () => { @@ -655,19 +666,21 @@ describe('Egress Quota Management', () => { pieceId: 'piece-no-enforce-both-exhausted', }) - const result = await getStorageProviderAndValidatePayer( + const results = await getRetrievalCandidatesAndValidatePayer( env, payerAddress, pieceCid, false, ) - expect(result).toStrictEqual({ - dataSetId, - serviceProviderId: APPROVED_SERVICE_PROVIDER_ID, - serviceUrl: 'https://quota-test-provider.xyz', - cdnEgressQuota: 0n, - cacheMissEgressQuota: 0n, - }) + expect(results).toMatchObject([ + { + dataSetId, + serviceProviderId: APPROVED_SERVICE_PROVIDER_ID, + serviceUrl: 'https://quota-test-provider.xyz', + cdnEgressQuota: 0n, + cacheMissEgressQuota: 0n, + }, + ]) }) })