diff --git a/src/ORM/_migrations/AcceptedShareIndex1780862400000.ts b/src/ORM/_migrations/AcceptedShareIndex1780862400000.ts new file mode 100644 index 0000000..fbb6a0d --- /dev/null +++ b/src/ORM/_migrations/AcceptedShareIndex1780862400000.ts @@ -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 { + 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 { + 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"`); + } +} diff --git a/src/ORM/accepted-share/accepted-share.entity.ts b/src/ORM/accepted-share/accepted-share.entity.ts index cb53662..1296f1a 100644 --- a/src/ORM/accepted-share/accepted-share.entity.ts +++ b/src/ORM/accepted-share/accepted-share.entity.ts @@ -3,6 +3,7 @@ import { Column, Entity, Index, PrimaryColumn, PrimaryGeneratedColumn } from 'ty @Entity() @Index('IDX_accepted_share_accounting_lookup', ['address', 'clientName', '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 }) export class AcceptedShareEntity { @PrimaryGeneratedColumn('uuid') @@ -11,6 +12,9 @@ export class AcceptedShareEntity { @PrimaryColumn({ type: 'timestamptz' }) acceptedAt: Date; + @Column({ type: 'bigint', default: () => `nextval('accepted_share_index_seq')` }) + shareIndex: number; + @Column({ length: 8, type: 'varchar' }) protocol: 'sv1' | 'sv2'; diff --git a/src/database.config.ts b/src/database.config.ts index 0a424c3..f9033b1 100644 --- a/src/database.config.ts +++ b/src/database.config.ts @@ -4,6 +4,7 @@ import { DataSourceOptions } from 'typeorm'; import { InitialTimescaleSchema1780859300000 } from './ORM/_migrations/InitialTimescaleSchema1780859300000'; import { ActiveOnlyUserAgentReport1780860200000 } from './ORM/_migrations/ActiveOnlyUserAgentReport1780860200000'; 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 { AcceptedShareEntity } from './ORM/accepted-share/accepted-share.entity'; import { AddressSettingsEntity } from './ORM/address-settings/address-settings.entity'; @@ -26,6 +27,7 @@ export const databaseMigrations = [ InitialTimescaleSchema1780859300000, ActiveOnlyUserAgentReport1780860200000, TimescaleOperationalHardening1780861200000, + AcceptedShareIndex1780862400000, ]; export function createDatabaseOptions(env: NodeJS.ProcessEnv): TypeOrmModuleOptions & DataSourceOptions { diff --git a/test/timescale-redis.integration-spec.ts b/test/timescale-redis.integration-spec.ts index 786204c..f97fd8e 100644 --- a/test/timescale-redis.integration-spec.ts +++ b/test/timescale-redis.integration-spec.ts @@ -91,6 +91,15 @@ describe('TimescaleDB and Redis integration', () => { expect.objectContaining({ proc_name: 'policy_retention' }), 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 () => { @@ -126,9 +135,33 @@ describe('TimescaleDB and Redis integration', () => { isBlockCandidate: false, 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(Number(rows[0].last)).toBeGreaterThan(Number(rows[0].first)); await dataSource.query(`CALL refresh_continuous_aggregate('accepted_share_10m', NULL, NULL)`); const aggregateRows = await dataSource.query(` @@ -138,7 +171,7 @@ describe('TimescaleDB and Redis integration', () => { `, [client.address, client.clientName]); expect(aggregateRows).toEqual(expect.arrayContaining([ - expect.objectContaining({ shares: 64, acceptedCount: 1 }), + expect.objectContaining({ shares: 96, acceptedCount: 2 }), ])); });