Add accepted share order index

This commit is contained in:
Ben
2026-06-07 22:03:32 -04:00
parent 4f51a4d05c
commit 3e1c303fa0
4 changed files with 96 additions and 2 deletions
@@ -0,0 +1,55 @@
import { MigrationInterface, QueryRunner } from 'typeorm';
export class AcceptedShareIndex1780862400000 implements MigrationInterface {
public name = 'AcceptedShareIndex1780862400000';
public transaction = false;
public async up(queryRunner: QueryRunner): Promise<void> {
await queryRunner.query(`CREATE SEQUENCE IF NOT EXISTS "accepted_share_index_seq"`);
await queryRunner.query(`ALTER TABLE "accepted_share_entity" ADD COLUMN IF NOT EXISTS "shareIndex" bigint`);
await queryRunner.query(`
SELECT setval(
'accepted_share_index_seq',
GREATEST((SELECT COUNT(*) FROM "accepted_share_entity") + 1, 1),
false
)
`);
await queryRunner.query(`
ALTER TABLE "accepted_share_entity"
ALTER COLUMN "shareIndex" SET DEFAULT nextval('accepted_share_index_seq')
`);
await queryRunner.query(`
WITH ordered AS (
SELECT
"id",
"acceptedAt",
row_number() OVER (ORDER BY "acceptedAt", "id") AS "shareIndex"
FROM "accepted_share_entity"
WHERE "shareIndex" IS NULL
)
UPDATE "accepted_share_entity" AS share
SET "shareIndex" = ordered."shareIndex"
FROM ordered
WHERE share."id" = ordered."id"
AND share."acceptedAt" = ordered."acceptedAt"
`);
await queryRunner.query(`
ALTER TABLE "accepted_share_entity"
ALTER COLUMN "shareIndex" SET NOT NULL
`);
await queryRunner.query(`
CREATE INDEX IF NOT EXISTS "IDX_accepted_share_order"
ON "accepted_share_entity" ("shareIndex" DESC)
`);
}
public async down(queryRunner: QueryRunner): Promise<void> {
await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_order"`);
await queryRunner.query(`
ALTER TABLE "accepted_share_entity"
ALTER COLUMN "shareIndex" DROP DEFAULT
`);
await queryRunner.query(`ALTER TABLE "accepted_share_entity" DROP COLUMN IF EXISTS "shareIndex"`);
await queryRunner.query(`DROP SEQUENCE IF EXISTS "accepted_share_index_seq"`);
}
}
@@ -3,6 +3,7 @@ import { Column, Entity, Index, PrimaryColumn, PrimaryGeneratedColumn } from 'ty
@Entity() @Entity()
@Index('IDX_accepted_share_accounting_lookup', ['address', 'clientName', 'acceptedAt']) @Index('IDX_accepted_share_accounting_lookup', ['address', 'clientName', 'acceptedAt'])
@Index('IDX_accepted_share_client_lookup', ['clientId', 'acceptedAt']) @Index('IDX_accepted_share_client_lookup', ['clientId', 'acceptedAt'])
@Index('IDX_accepted_share_order', ['shareIndex'])
@Index('IDX_accepted_share_unique_submission', ['acceptedAt', 'protocol', 'sessionId', 'jobId', 'nonce', 'ntime', 'version', 'extraNonce2'], { unique: true }) @Index('IDX_accepted_share_unique_submission', ['acceptedAt', 'protocol', 'sessionId', 'jobId', 'nonce', 'ntime', 'version', 'extraNonce2'], { unique: true })
export class AcceptedShareEntity { export class AcceptedShareEntity {
@PrimaryGeneratedColumn('uuid') @PrimaryGeneratedColumn('uuid')
@@ -11,6 +12,9 @@ export class AcceptedShareEntity {
@PrimaryColumn({ type: 'timestamptz' }) @PrimaryColumn({ type: 'timestamptz' })
acceptedAt: Date; acceptedAt: Date;
@Column({ type: 'bigint', default: () => `nextval('accepted_share_index_seq')` })
shareIndex: number;
@Column({ length: 8, type: 'varchar' }) @Column({ length: 8, type: 'varchar' })
protocol: 'sv1' | 'sv2'; protocol: 'sv1' | 'sv2';
+2
View File
@@ -4,6 +4,7 @@ import { DataSourceOptions } from 'typeorm';
import { InitialTimescaleSchema1780859300000 } from './ORM/_migrations/InitialTimescaleSchema1780859300000'; import { InitialTimescaleSchema1780859300000 } from './ORM/_migrations/InitialTimescaleSchema1780859300000';
import { ActiveOnlyUserAgentReport1780860200000 } from './ORM/_migrations/ActiveOnlyUserAgentReport1780860200000'; import { ActiveOnlyUserAgentReport1780860200000 } from './ORM/_migrations/ActiveOnlyUserAgentReport1780860200000';
import { TimescaleOperationalHardening1780861200000 } from './ORM/_migrations/TimescaleOperationalHardening1780861200000'; import { TimescaleOperationalHardening1780861200000 } from './ORM/_migrations/TimescaleOperationalHardening1780861200000';
import { AcceptedShareIndex1780862400000 } from './ORM/_migrations/AcceptedShareIndex1780862400000';
import { UserAgentReportView } from './ORM/_views/user-agent-report/user-agent-report.view'; import { UserAgentReportView } from './ORM/_views/user-agent-report/user-agent-report.view';
import { AcceptedShareEntity } from './ORM/accepted-share/accepted-share.entity'; import { AcceptedShareEntity } from './ORM/accepted-share/accepted-share.entity';
import { AddressSettingsEntity } from './ORM/address-settings/address-settings.entity'; import { AddressSettingsEntity } from './ORM/address-settings/address-settings.entity';
@@ -26,6 +27,7 @@ export const databaseMigrations = [
InitialTimescaleSchema1780859300000, InitialTimescaleSchema1780859300000,
ActiveOnlyUserAgentReport1780860200000, ActiveOnlyUserAgentReport1780860200000,
TimescaleOperationalHardening1780861200000, TimescaleOperationalHardening1780861200000,
AcceptedShareIndex1780862400000,
]; ];
export function createDatabaseOptions(env: NodeJS.ProcessEnv): TypeOrmModuleOptions & DataSourceOptions { export function createDatabaseOptions(env: NodeJS.ProcessEnv): TypeOrmModuleOptions & DataSourceOptions {
+35 -2
View File
@@ -91,6 +91,15 @@ describe('TimescaleDB and Redis integration', () => {
expect.objectContaining({ proc_name: 'policy_retention' }), expect.objectContaining({ proc_name: 'policy_retention' }),
expect.objectContaining({ proc_name: 'policy_refresh_continuous_aggregate' }), expect.objectContaining({ proc_name: 'policy_refresh_continuous_aggregate' }),
])); ]));
const shareOrderObjects = await dataSource.query(`
SELECT to_regclass('public.accepted_share_index_seq') AS sequence_name,
to_regclass('public."IDX_accepted_share_order"') AS index_name
`);
expect(shareOrderObjects[0]).toEqual({
sequence_name: 'accepted_share_index_seq',
index_name: '"IDX_accepted_share_order"',
});
}); });
it('should persist accepted shares and refresh the 10 minute aggregate', async () => { it('should persist accepted shares and refresh the 10 minute aggregate', async () => {
@@ -126,9 +135,33 @@ describe('TimescaleDB and Redis integration', () => {
isBlockCandidate: false, isBlockCandidate: false,
blockSubmissionResult: null, blockSubmissionResult: null,
}); });
await service.recordAcceptedShare({
protocol: 'sv1',
acceptedAt: new Date(acceptedAt.getTime() + 1),
address: client.address,
clientName: client.clientName,
sessionId: client.sessionId,
clientId: client.id,
jobId: '2',
jobTemplateId: '1',
blockHeight: 900000,
creditedDifficulty: 32,
submissionDifficulty: 64,
networkDifficulty: 100000,
nonce: 'ed460d92',
ntime: '64b3f3ec',
version: '20000000',
extraNonce2: 'c708000000000001',
isBlockCandidate: false,
blockSubmissionResult: null,
});
const rows = await dataSource.query(`SELECT COUNT(*)::int AS count FROM accepted_share_entity`); const rows = await dataSource.query(`
SELECT COUNT(*)::int AS count, MIN("shareIndex")::bigint AS first, MAX("shareIndex")::bigint AS last
FROM accepted_share_entity
`);
expect(rows[0].count).toBeGreaterThanOrEqual(1); expect(rows[0].count).toBeGreaterThanOrEqual(1);
expect(Number(rows[0].last)).toBeGreaterThan(Number(rows[0].first));
await dataSource.query(`CALL refresh_continuous_aggregate('accepted_share_10m', NULL, NULL)`); await dataSource.query(`CALL refresh_continuous_aggregate('accepted_share_10m', NULL, NULL)`);
const aggregateRows = await dataSource.query(` const aggregateRows = await dataSource.query(`
@@ -138,7 +171,7 @@ describe('TimescaleDB and Redis integration', () => {
`, [client.address, client.clientName]); `, [client.address, client.clientName]);
expect(aggregateRows).toEqual(expect.arrayContaining([ expect(aggregateRows).toEqual(expect.arrayContaining([
expect.objectContaining({ shares: 64, acceptedCount: 1 }), expect.objectContaining({ shares: 96, acceptedCount: 2 }),
])); ]));
}); });