From 7e4e82ebdec197395855c3d0d45db88f4a6ee571 Mon Sep 17 00:00:00 2001 From: Julian Gruber Date: Thu, 6 Nov 2025 08:56:00 +0100 Subject: [PATCH 01/10] piece-retriever: add try different SP on retrieval failure --- piece-retriever/bin/piece-retriever.js | 84 ++++++++++++++----------- piece-retriever/lib/store.js | 57 +++++++---------- piece-retriever/test/store.test.js | 86 ++++++++++++++------------ 3 files changed, 115 insertions(+), 112 deletions(-) diff --git a/piece-retriever/bin/piece-retriever.js b/piece-retriever/bin/piece-retriever.js index 66cbcae0..9ad74329 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, @@ -90,24 +89,39 @@ export default { 'The requested CID was flagged by the Bad Bits Denylist at https://badbits.dwebops.pub', ) - httpAssert( - serviceProviderId, - 404, - `Unsupported Service Provider: ${serviceProviderId}`, - ) + if (retrievalCandidates.length === 0) { + const response = new Response('No retrieval candidate found', { + status: 404, + }) + setContentSecurityPolicy(response) + return response + } + let retrievalCandidate let retrievalResult - 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] + retrievalCandidates.splice(retrievalCandidateIndex) + try { + retrievalResult = await retrieveFile( + ctx, + retrievalCandidate.serviceUrl, + pieceCid, + request, + env.ORIGIN_CACHE_TTL, + { signal: request.signal }, + ) + if (retrievalResult.response.ok) { + break + } + } catch {} + } + + httpAssert(retrievalCandidate, 500, 'should never happen') if (!retrievalResult || retrievalResult.response.status >= 500) { ctx.waitUntil( @@ -117,16 +131,16 @@ 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}` : ''}`, + `Service provider ${retrievalCandidate.serviceProviderId} is unavailable${retrievalResult ? ` at ${retrievalResult.url}` : ''}`, { status: 502, headers: new Headers({ - 'X-Data-Set-ID': dataSetId, + 'X-Data-Set-ID': retrievalCandidate.dataSetId, }), }, ) @@ -145,7 +159,7 @@ export default { egressBytes: 0, requestCountryCode, timestamp: requestTimestamp, - dataSetId, + dataSetId: retrievalCandidate.dataSetId, botName, }), ) @@ -154,7 +168,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 +199,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 +219,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/store.test.js b/piece-retriever/test/store.test.js index 074e7fec..1f068f86 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,20 @@ 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.length).toBe(1) + expect(results[0].serviceProviderId).toBe(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 +103,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 +125,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 +149,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 +169,17 @@ 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.length).toBe(1) + expect(results[0].serviceProviderId).toBe(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 +217,16 @@ 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.length).toBe(2) + expect(results[0].serviceProviderId).toBe(serviceProviderId1) + expect(results[1].serviceProviderId).toBe(serviceProviderId2) }) it('ignores owners that are not approved by Filecoin Warm Storage Service', async () => { @@ -268,13 +266,14 @@ 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({ + expect(results.length).toBe(1) + expect(results[0]).toEqual({ dataSetId: dataSetId1, serviceProviderId: serviceProviderId1.toLowerCase(), serviceUrl: 'https://pdp-provider-1.xyz', @@ -311,7 +310,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 +331,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,13 +351,14 @@ describe('Egress Quota Management', () => { pieceId: 'piece-sufficient', }) - const result = await getStorageProviderAndValidatePayer( + const results = await getRetrievalCandidatesAndValidatePayer( env, payerAddress, pieceCid, true, ) - expect(result).toStrictEqual({ + expect(results.length).toBe(1) + expect(results[0]).toStrictEqual({ dataSetId, serviceProviderId: APPROVED_SERVICE_PROVIDER_ID, serviceUrl: 'https://quota-test-provider.xyz', @@ -523,7 +523,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,13 +545,14 @@ 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({ + expect(results.length).toBe(1) + expect(results[0]).toStrictEqual({ dataSetId, serviceProviderId: APPROVED_SERVICE_PROVIDER_ID, serviceUrl: 'https://quota-test-provider.xyz', @@ -593,13 +594,14 @@ 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({ + expect(results.length).toBe(1) + expect(results[0]).toStrictEqual({ dataSetId, serviceProviderId: APPROVED_SERVICE_PROVIDER_ID, serviceUrl: 'https://quota-test-provider.xyz', @@ -624,13 +626,14 @@ 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({ + expect(results.length).toBe(1) + expect(results[0]).toStrictEqual({ dataSetId, serviceProviderId: APPROVED_SERVICE_PROVIDER_ID, serviceUrl: 'https://quota-test-provider.xyz', @@ -655,13 +658,14 @@ 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({ + expect(results.length).toBe(1) + expect(results[0]).toStrictEqual({ dataSetId, serviceProviderId: APPROVED_SERVICE_PROVIDER_ID, serviceUrl: 'https://quota-test-provider.xyz', From d7437fcee4e77dc01de144b0fe3e24b17a7430fa Mon Sep 17 00:00:00 2001 From: Julian Gruber Date: Thu, 6 Nov 2025 09:28:37 +0100 Subject: [PATCH 02/10] Update piece-retriever/bin/piece-retriever.js MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Miroslav Bajtoš --- piece-retriever/bin/piece-retriever.js | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/piece-retriever/bin/piece-retriever.js b/piece-retriever/bin/piece-retriever.js index 9ad74329..1080ff39 100644 --- a/piece-retriever/bin/piece-retriever.js +++ b/piece-retriever/bin/piece-retriever.js @@ -105,7 +105,7 @@ export default { Math.random() * retrievalCandidates.length, ) retrievalCandidate = retrievalCandidates[retrievalCandidateIndex] - retrievalCandidates.splice(retrievalCandidateIndex) + retrievalCandidates.splice(retrievalCandidateIndex, 1) try { retrievalResult = await retrieveFile( ctx, From 07105aa6ab2f74ffa23ae2b0d7f0078a9ce0ce35 Mon Sep 17 00:00:00 2001 From: Julian Gruber Date: Thu, 6 Nov 2025 09:49:11 +0100 Subject: [PATCH 03/10] improve SPs unavailable error message --- piece-retriever/bin/piece-retriever.js | 4 +++- piece-retriever/test/retriever.test.js | 8 ++++++-- 2 files changed, 9 insertions(+), 3 deletions(-) diff --git a/piece-retriever/bin/piece-retriever.js b/piece-retriever/bin/piece-retriever.js index 1080ff39..70505104 100644 --- a/piece-retriever/bin/piece-retriever.js +++ b/piece-retriever/bin/piece-retriever.js @@ -99,12 +99,14 @@ export default { let retrievalCandidate let retrievalResult + const retrievalAttempts = [] while (retrievalCandidates.length > 0) { const retrievalCandidateIndex = Math.floor( Math.random() * retrievalCandidates.length, ) retrievalCandidate = retrievalCandidates[retrievalCandidateIndex] + retrievalAttempts.push(retrievalCandidate) retrievalCandidates.splice(retrievalCandidateIndex, 1) try { retrievalResult = await retrieveFile( @@ -136,7 +138,7 @@ export default { }), ) const response = new Response( - `Service provider ${retrievalCandidate.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({ diff --git a/piece-retriever/test/retriever.test.js b/piece-retriever/test/retriever.test.js index e9eaa341..50d7d0ba 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()).toBe(`Service provider 2 is unavailable at ${url}`) + 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()).toBe(`Service provider 2 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)) }) }) From 9332577b2eda6783ff6842ad4f7e57d30f941083 Mon Sep 17 00:00:00 2001 From: Julian Gruber Date: Thu, 6 Nov 2025 09:50:55 +0100 Subject: [PATCH 04/10] fix header --- piece-retriever/bin/piece-retriever.js | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/piece-retriever/bin/piece-retriever.js b/piece-retriever/bin/piece-retriever.js index 70505104..4a306288 100644 --- a/piece-retriever/bin/piece-retriever.js +++ b/piece-retriever/bin/piece-retriever.js @@ -142,7 +142,9 @@ export default { { status: 502, headers: new Headers({ - 'X-Data-Set-ID': retrievalCandidate.dataSetId, + 'X-Data-Set-ID': retrievalAttempts + .map((a) => a.dataSetId) + .join(','), }), }, ) From 75ba68c96a4b7d66d170233df772bd9112bf169f Mon Sep 17 00:00:00 2001 From: Julian Gruber Date: Thu, 6 Nov 2025 09:52:35 +0100 Subject: [PATCH 05/10] log --- piece-retriever/bin/piece-retriever.js | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/piece-retriever/bin/piece-retriever.js b/piece-retriever/bin/piece-retriever.js index 4a306288..4908bc28 100644 --- a/piece-retriever/bin/piece-retriever.js +++ b/piece-retriever/bin/piece-retriever.js @@ -120,7 +120,12 @@ export default { if (retrievalResult.response.ok) { break } - } catch {} + } catch { + console.log('Retrieval attempt failed', { + retrievalCandidate, + willRetry: retrievalCandidates.length > 0 + }) + } } httpAssert(retrievalCandidate, 500, 'should never happen') From 7dc2cd74f9f213ec9e0d8aa90731fc71ff185dbe Mon Sep 17 00:00:00 2001 From: Julian Gruber Date: Thu, 6 Nov 2025 09:54:18 +0100 Subject: [PATCH 06/10] log --- piece-retriever/bin/piece-retriever.js | 1 + 1 file changed, 1 insertion(+) diff --git a/piece-retriever/bin/piece-retriever.js b/piece-retriever/bin/piece-retriever.js index 4908bc28..f057454a 100644 --- a/piece-retriever/bin/piece-retriever.js +++ b/piece-retriever/bin/piece-retriever.js @@ -108,6 +108,7 @@ export default { retrievalCandidate = retrievalCandidates[retrievalCandidateIndex] retrievalAttempts.push(retrievalCandidate) retrievalCandidates.splice(retrievalCandidateIndex, 1) + console.log('Attempting retrieval', retrievalCandidate) try { retrievalResult = await retrieveFile( ctx, From 05dedf7a8ec2dd3baf67a39bc4dad91c6c5b01a8 Mon Sep 17 00:00:00 2001 From: Julian Gruber Date: Thu, 6 Nov 2025 09:59:28 +0100 Subject: [PATCH 07/10] refactor --- piece-retriever/bin/piece-retriever.js | 2 +- piece-retriever/test/store.test.js | 119 +++++++++++++------------ 2 files changed, 65 insertions(+), 56 deletions(-) diff --git a/piece-retriever/bin/piece-retriever.js b/piece-retriever/bin/piece-retriever.js index f057454a..8aaaae36 100644 --- a/piece-retriever/bin/piece-retriever.js +++ b/piece-retriever/bin/piece-retriever.js @@ -124,7 +124,7 @@ export default { } catch { console.log('Retrieval attempt failed', { retrievalCandidate, - willRetry: retrievalCandidates.length > 0 + willRetry: retrievalCandidates.length > 0, }) } } diff --git a/piece-retriever/test/store.test.js b/piece-retriever/test/store.test.js index 1f068f86..abc1e589 100644 --- a/piece-retriever/test/store.test.js +++ b/piece-retriever/test/store.test.js @@ -79,8 +79,9 @@ describe('getRetrievalCandidatesAndValidatePayer', () => { pieceCid, true, ) - expect(results.length).toBe(1) - expect(results[0].serviceProviderId).toBe(APPROVED_SERVICE_PROVIDER_ID) + expect(results).toMatchObject([ + { serviceProviderId: APPROVED_SERVICE_PROVIDER_ID }, + ]) }) it('throws error if pieceCid not found', async () => { @@ -175,8 +176,9 @@ describe('getRetrievalCandidatesAndValidatePayer', () => { pieceCid, true, ) - expect(results.length).toBe(1) - expect(results[0].serviceProviderId).toBe(APPROVED_SERVICE_PROVIDER_ID) + expect(results).toMatchObject([ + { serviceProviderId: APPROVED_SERVICE_PROVIDER_ID }, + ]) }) it('returns multiple service providers when they share the same pieceCid', async () => { @@ -224,9 +226,10 @@ describe('getRetrievalCandidatesAndValidatePayer', () => { true, ) - expect(results.length).toBe(2) - expect(results[0].serviceProviderId).toBe(serviceProviderId1) - expect(results[1].serviceProviderId).toBe(serviceProviderId2) + expect(results).toMatchObject([ + { serviceProviderId: serviceProviderId1 }, + { serviceProviderId: serviceProviderId2 }, + ]) }) it('ignores owners that are not approved by Filecoin Warm Storage Service', async () => { @@ -272,14 +275,15 @@ describe('getRetrievalCandidatesAndValidatePayer', () => { pieceCid, true, ) - expect(results.length).toBe(1) - expect(results[0]).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, + }, + ]) }) }) @@ -357,14 +361,15 @@ describe('Egress Quota Management', () => { pieceCid, true, ) - expect(results.length).toBe(1) - expect(results[0]).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 () => { @@ -551,14 +556,15 @@ describe('Egress Quota Management', () => { pieceCid, true, ) - expect(results.length).toBe(1) - expect(results[0]).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, { @@ -600,14 +606,15 @@ describe('Egress Quota Management', () => { pieceCid, false, ) - expect(results.length).toBe(1) - expect(results[0]).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 () => { @@ -632,14 +639,15 @@ describe('Egress Quota Management', () => { pieceCid, false, ) - expect(results.length).toBe(1) - expect(results[0]).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 () => { @@ -664,14 +672,15 @@ describe('Egress Quota Management', () => { pieceCid, false, ) - expect(results.length).toBe(1) - expect(results[0]).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, + }, + ]) }) }) From c2e1d3486479285b32f991826a4ee1c3e4441df2 Mon Sep 17 00:00:00 2001 From: Julian Gruber Date: Fri, 7 Nov 2025 10:09:45 +0100 Subject: [PATCH 08/10] fix 500 --- piece-retriever/bin/piece-retriever.js | 12 +++++------- 1 file changed, 5 insertions(+), 7 deletions(-) diff --git a/piece-retriever/bin/piece-retriever.js b/piece-retriever/bin/piece-retriever.js index 8aaaae36..4a9ff616 100644 --- a/piece-retriever/bin/piece-retriever.js +++ b/piece-retriever/bin/piece-retriever.js @@ -89,13 +89,11 @@ export default { 'The requested CID was flagged by the Bad Bits Denylist at https://badbits.dwebops.pub', ) - if (retrievalCandidates.length === 0) { - const response = new Response('No retrieval candidate found', { - status: 404, - }) - setContentSecurityPolicy(response) - return response - } + httpAssert( + retrievalCandidates.length > 0, + 500, + 'Service provider lookup failed', + ) let retrievalCandidate let retrievalResult From 518df7011fc76854db17afe5064bc0eef2e9d915 Mon Sep 17 00:00:00 2001 From: Julian Gruber Date: Fri, 7 Nov 2025 10:10:44 +0100 Subject: [PATCH 09/10] Update piece-retriever/bin/piece-retriever.js MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Miroslav Bajtoš --- piece-retriever/bin/piece-retriever.js | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/piece-retriever/bin/piece-retriever.js b/piece-retriever/bin/piece-retriever.js index 4a9ff616..e175dba8 100644 --- a/piece-retriever/bin/piece-retriever.js +++ b/piece-retriever/bin/piece-retriever.js @@ -119,8 +119,13 @@ export default { if (retrievalResult.response.ok) { break } + console.log(`Retrieval attempt failed: HTTP ${retrievalResult.response.status}', { + retrievalCandidate, + willRetry: retrievalCandidates.length > 0, + }) + } } catch { - console.log('Retrieval attempt failed', { + console.log('Retrieval attempt failed: ${err.message ?? err}', { retrievalCandidate, willRetry: retrievalCandidates.length > 0, }) From 77a47d91340d6d4ba84ce56254c69fe48427f457 Mon Sep 17 00:00:00 2001 From: Julian Gruber Date: Fri, 7 Nov 2025 10:12:38 +0100 Subject: [PATCH 10/10] fix syntax & type --- piece-retriever/bin/piece-retriever.js | 20 +++++++++++++------- 1 file changed, 13 insertions(+), 7 deletions(-) diff --git a/piece-retriever/bin/piece-retriever.js b/piece-retriever/bin/piece-retriever.js index e175dba8..c197554b 100644 --- a/piece-retriever/bin/piece-retriever.js +++ b/piece-retriever/bin/piece-retriever.js @@ -119,13 +119,19 @@ export default { if (retrievalResult.response.ok) { break } - console.log(`Retrieval attempt failed: HTTP ${retrievalResult.response.status}', { - retrievalCandidate, - willRetry: retrievalCandidates.length > 0, - }) - } - } catch { - console.log('Retrieval attempt failed: ${err.message ?? err}', { + 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, })