diff --git a/lib/public-stats.js b/lib/public-stats.js index f595c142..1964cac5 100644 --- a/lib/public-stats.js +++ b/lib/public-stats.js @@ -29,6 +29,7 @@ export const updatePublicStats = async ({ createPgClient, committees, allMeasure await updateRetrievalTimings(pgClient, committees) await updateDailyClientRetrievalStats(pgClient, committees, findDealClients) await updateDailyAllocatorRetrievalStats(pgClient, committees, findDealAllocators) + await updateDailyMinerDealsChecked(pgClient, committees) } finally { await pgClient.end() } @@ -394,3 +395,67 @@ function buildPerPartyStats (committees, perDealParty, partyName) { ) return flatStats } + +/** + * @param {pg.Client} pgClient + * @param {Iterable} committees + */ +const updateDailyMinerDealsChecked = async (pgClient, committees) => { + /** @type {Map>} */ + const minerPayloadCids = new Map() + + for (const c of committees) { + const { minerId, cid } = c.retrievalTask + + let payloadCids = minerPayloadCids.get(minerId) + if (!payloadCids) { + payloadCids = new Set() + minerPayloadCids.set(minerId, payloadCids) + } + + payloadCids.add(cid) + } + + // Convert the map to arrays for the query + const flatStats = Array.from(minerPayloadCids.entries()).map( + // eslint-disable-next-line camelcase + ([minerId, payloadCids]) => ({ + miner_id: minerId, + payload_cids: Array.from(payloadCids) + }) + ) + + if (debug.enabled) { + debug( + 'Updating public stats - daily miner deals checked: miners count = %s, total payload CIDs = %s', + flatStats.length, + flatStats.reduce((sum, stat) => sum + stat.payload_cids.length, 0) + ) + } + + if (flatStats.length === 0) { + debug('No miner deals to record') + return + } + + await pgClient.query(` + INSERT INTO daily_miner_deals_checked ( + day, + miner_id, + payload_cids + ) + SELECT + now(), + miner_id, + payload_cids + FROM jsonb_to_recordset($1::jsonb) AS t (miner_id text, payload_cids text[]) + ON CONFLICT(day, miner_id) DO UPDATE SET + payload_cids = array( + SELECT DISTINCT unnest( + array_cat(daily_miner_deals_checked.payload_cids, EXCLUDED.payload_cids) + ) + ) + `, [ + JSON.stringify(flatStats) + ]) +} diff --git a/migrations/026.do.daily-miner-deals-checked.sql b/migrations/026.do.daily-miner-deals-checked.sql new file mode 100644 index 00000000..ec84a9eb --- /dev/null +++ b/migrations/026.do.daily-miner-deals-checked.sql @@ -0,0 +1,6 @@ +CREATE TABLE daily_miner_deals_checked ( + day DATE NOT NULL, + miner_id TEXT NOT NULL, + payload_cids TEXT[] NOT NULL, + PRIMARY KEY (day, miner_id) +); \ No newline at end of file diff --git a/test/public-stats.test.js b/test/public-stats.test.js index 2b61d693..f965c274 100644 --- a/test/public-stats.test.js +++ b/test/public-stats.test.js @@ -31,6 +31,7 @@ describe('public-stats', () => { await pgClient.query('DELETE FROM retrieval_timings') await pgClient.query('DELETE FROM daily_client_retrieval_stats') await pgClient.query('DELETE FROM daily_allocator_retrieval_stats') + await pgClient.query('DELETE FROM daily_miner_deals_checked') // Run all tests inside a transaction to ensure `now()` always returns the same value // See https://dba.stackexchange.com/a/63549/125312 @@ -1170,6 +1171,133 @@ describe('public-stats', () => { }) }) + describe('updateDailyMinerDealsChecked', () => { + it('collects payload CIDs per miner', async () => { + /** @type {Measurement[]} */ + const honestMeasurements = [ + { ...VALID_MEASUREMENT, minerId: 'f1first', cid: 'cidone' }, + { ...VALID_MEASUREMENT, minerId: 'f1first', cid: 'cidtwo' }, + { ...VALID_MEASUREMENT, minerId: 'f1second', cid: 'cidone' }, + { ...VALID_MEASUREMENT, minerId: 'f1second', cid: 'cidthree' } + ] + const allMeasurements = honestMeasurements + const committees = buildEvaluatedCommitteesFromMeasurements(honestMeasurements) + + await updatePublicStats({ + createPgClient, + committees, + allMeasurements, + findDealClients: (_minerId, _cid) => ['f0client'], + findDealAllocators: (_minerId, _cid) => ['f0allocator'] + }) + + const { rows: created } = await pgClient.query( + 'SELECT day::TEXT, miner_id, payload_cids FROM daily_miner_deals_checked ORDER BY miner_id' + ) + assert.deepStrictEqual(created, [ + { day: today, miner_id: 'f1first', payload_cids: ['cidone', 'cidtwo'] }, + { day: today, miner_id: 'f1second', payload_cids: ['cidone', 'cidthree'] } + ]) + }) + + it('handles duplicate CIDs correctly', async () => { + /** @type {Measurement[]} */ + const honestMeasurements = [ + { ...VALID_MEASUREMENT, minerId: 'f1miner', cid: 'cidone' }, + { ...VALID_MEASUREMENT, minerId: 'f1miner', cid: 'cidone' }, // duplicate + { ...VALID_MEASUREMENT, minerId: 'f1miner', cid: 'cidtwo' } + ] + const allMeasurements = honestMeasurements + const committees = buildEvaluatedCommitteesFromMeasurements(honestMeasurements) + + await updatePublicStats({ + createPgClient, + committees, + allMeasurements, + findDealClients: (_minerId, _cid) => ['f0client'], + findDealAllocators: (_minerId, _cid) => ['f0allocator'] + }) + + const { rows: created } = await pgClient.query( + 'SELECT day::TEXT, miner_id, payload_cids FROM daily_miner_deals_checked' + ) + assert.deepStrictEqual(created, [ + { day: today, miner_id: 'f1miner', payload_cids: ['cidone', 'cidtwo'] } + ]) + }) + + it('updates existing records by merging CID arrays', async () => { + // First, create an initial record + { + /** @type {Measurement[]} */ + const honestMeasurements = [ + { ...VALID_MEASUREMENT, minerId: 'f1miner', cid: 'cidone' } + ] + const allMeasurements = honestMeasurements + const committees = buildEvaluatedCommitteesFromMeasurements(honestMeasurements) + + await updatePublicStats({ + createPgClient, + committees, + allMeasurements, + findDealClients: (_minerId, _cid) => ['f0client'], + findDealAllocators: (_minerId, _cid) => ['f0allocator'] + }) + + const { rows: created } = await pgClient.query( + 'SELECT day::TEXT, miner_id, payload_cids FROM daily_miner_deals_checked' + ) + assert.deepStrictEqual(created, [ + { day: today, miner_id: 'f1miner', payload_cids: ['cidone'] } + ]) + } + + // Now update with additional CIDs + { + /** @type {Measurement[]} */ + const honestMeasurements = [ + { ...VALID_MEASUREMENT, minerId: 'f1miner', cid: 'cidone' }, // duplicate - should be ignored + { ...VALID_MEASUREMENT, minerId: 'f1miner', cid: 'cidtwo' } // new CID + ] + const allMeasurements = honestMeasurements + const committees = buildEvaluatedCommitteesFromMeasurements(honestMeasurements) + + await updatePublicStats({ + createPgClient, + committees, + allMeasurements, + findDealClients: (_minerId, _cid) => ['f0client'], + findDealAllocators: (_minerId, _cid) => ['f0allocator'] + }) + + const { rows: updated } = await pgClient.query( + 'SELECT day::TEXT, miner_id, payload_cids FROM daily_miner_deals_checked' + ) + assert.deepStrictEqual(updated, [ + { day: today, miner_id: 'f1miner', payload_cids: ['cidone', 'cidtwo'] } + ]) + } + }) + + it('handles empty committees gracefully', async () => { + const committees = [] + const allMeasurements = [] + + await updatePublicStats({ + createPgClient, + committees, + allMeasurements, + findDealClients: (_minerId, _cid) => ['f0client'], + findDealAllocators: (_minerId, _cid) => ['f0allocator'] + }) + + const { rows: created } = await pgClient.query( + 'SELECT * FROM daily_miner_deals_checked' + ) + assert.deepStrictEqual(created, []) + }) + }) + const getCurrentDate = async () => { const { rows: [{ today }] } = await pgClient.query('SELECT now()::DATE::TEXT as today') return today