Skip to content
This repository was archived by the owner on Jan 7, 2026. It is now read-only.
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
65 changes: 65 additions & 0 deletions lib/public-stats.js
Original file line number Diff line number Diff line change
Expand Up @@ -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()
}
Expand Down Expand Up @@ -394,3 +395,67 @@ function buildPerPartyStats (committees, perDealParty, partyName) {
)
return flatStats
}

/**
* @param {pg.Client} pgClient
* @param {Iterable<Committee>} committees
*/
const updateDailyMinerDealsChecked = async (pgClient, committees) => {
/** @type {Map<string, Set<string>>} */
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)
])
Comment thread
bajtos marked this conversation as resolved.
}
6 changes: 6 additions & 0 deletions migrations/026.do.daily-miner-deals-checked.sql
Original file line number Diff line number Diff line change
@@ -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)
);
128 changes: 128 additions & 0 deletions test/public-stats.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down