Skip to content
5 changes: 2 additions & 3 deletions apps/backend/src/data-retention/data-retention.module.ts
Original file line number Diff line number Diff line change
@@ -1,13 +1,12 @@
import { Module } from "@nestjs/common";
import { TypeOrmModule } from "@nestjs/typeorm";
import { DataRetentionBaseline } from "../database/entities/data-retention-baseline.entity.js";
import { StorageProvider } from "../database/entities/storage-provider.entity.js";
import { PdpSubgraphModule } from "../pdp-subgraph/pdp-subgraph.module.js";
import { WalletSdkModule } from "../wallet-sdk/wallet-sdk.module.js";
import { ProvidersModule } from "../providers/providers.module.js";
import { DataRetentionService } from "./data-retention.service.js";

@Module({
imports: [WalletSdkModule, PdpSubgraphModule, TypeOrmModule.forFeature([DataRetentionBaseline, StorageProvider])],
imports: [ProvidersModule, PdpSubgraphModule, TypeOrmModule.forFeature([DataRetentionBaseline])],
providers: [DataRetentionService],
exports: [DataRetentionService],
})
Expand Down
108 changes: 52 additions & 56 deletions apps/backend/src/data-retention/data-retention.service.spec.ts

Large diffs are not rendered by default.

21 changes: 6 additions & 15 deletions apps/backend/src/data-retention/data-retention.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,18 +3,17 @@ import { ConfigService } from "@nestjs/config";
import { InjectRepository } from "@nestjs/typeorm";
import { InjectMetric } from "@willsoto/nestjs-prometheus";
import { Counter, Gauge } from "prom-client";
import { Raw, Repository } from "typeorm";
import { Repository } from "typeorm";
import { ClickhouseService } from "../clickhouse/clickhouse.service.js";
import { toStructuredError } from "../common/logging.js";
import { isSpBlocked } from "../common/sp-blocklist.js";
import type { Network } from "../common/types.js";
import { IConfig, INetworksConfig } from "../config/index.js";
import { DataRetentionBaseline } from "../database/entities/data-retention-baseline.entity.js";
import { StorageProvider } from "../database/entities/storage-provider.entity.js";
import { buildCheckMetricLabels, CheckMetricLabels } from "../metrics-prometheus/check-metric-labels.js";
import { PDPSubgraphService } from "../pdp-subgraph/pdp-subgraph.service.js";
import { type ProviderDataSetResponse, type SubgraphMeta } from "../pdp-subgraph/types.js";
import { WalletSdkService } from "../wallet-sdk/wallet-sdk.service.js";
import { StorageProviderRepository } from "../providers/repositories/storage-provider.repository.js";
import { type PDPProviderEx } from "../wallet-sdk/wallet-sdk.types.js";

/**
Expand Down Expand Up @@ -53,12 +52,10 @@ export class DataRetentionService {

constructor(
private readonly configService: ConfigService<IConfig, true>,
private readonly walletSdkService: WalletSdkService,
private readonly storageProviderRepository: StorageProviderRepository,
private readonly pdpSubgraphService: PDPSubgraphService,
@InjectRepository(DataRetentionBaseline)
private readonly baselineRepository: Repository<DataRetentionBaseline>,
@InjectRepository(StorageProvider)
private readonly storageProviderRepository: Repository<StorageProvider>,
@InjectMetric("dataSetChallengeStatus")
private readonly dataSetChallengeStatusCounter: Counter,
@InjectMetric("pdp_provider_estimated_overdue_periods")
Expand Down Expand Up @@ -101,7 +98,7 @@ export class DataRetentionService {
// outer catch (which now preserves error type) rethrows it as a dependency failure.
throw new DataRetentionDependencyError("Failed to fetch PDP subgraph meta", { cause: error });
}
const allProviderInfos = this.walletSdkService.getTestingProviders(network);
const allProviderInfos = await this.storageProviderRepository.findTestingProviders(network);
const spBlocklists = this.configService.get("networks", { infer: true })[network];
const providerInfos = allProviderInfos?.filter((p) => !isSpBlocked(spBlocklists, p.serviceProvider, p.id));

Expand Down Expand Up @@ -272,15 +269,9 @@ export class DataRetentionService {
staleProviderCount: staleAddresses.length,
});

let staleProviders: StorageProvider[] = [];
let staleProviders: Awaited<ReturnType<StorageProviderRepository["findByAddressesCaseInsensitive"]>> = [];
try {
staleProviders = await this.storageProviderRepository.find({
where: {
network,
address: Raw((alias) => `LOWER(${alias}) IN (:...addresses)`, { addresses: staleAddresses }),
},
select: ["address", "providerId", "name", "isApproved"],
});
staleProviders = await this.storageProviderRepository.findByAddressesCaseInsensitive(staleAddresses, network);
} catch (error) {
// Bail entirely on DB failure to protect metric baselines
this.logger.error({
Expand Down
Original file line number Diff line number Diff line change
@@ -1,10 +1,11 @@
import { Module } from "@nestjs/common";
import { MetricsPrometheusModule } from "../metrics-prometheus/metrics-prometheus.module.js";
import { ProvidersModule } from "../providers/providers.module.js";
import { WalletSdkModule } from "../wallet-sdk/wallet-sdk.module.js";
import { DataSetLifecycleService } from "./data-set-lifecycle.service.js";

@Module({
imports: [WalletSdkModule, MetricsPrometheusModule],
imports: [WalletSdkModule, ProvidersModule, MetricsPrometheusModule],
providers: [DataSetLifecycleService],
exports: [DataSetLifecycleService],
})
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { DataSetLifecycleCheckMetrics } from "../metrics-prometheus/check-metrics.service.js";
import type { StorageProviderRepository } from "../providers/repositories/storage-provider.repository.js";
import { WalletSdkService } from "../wallet-sdk/wallet-sdk.service.js";
import { DataSetLifecycleService } from "./data-set-lifecycle.service.js";

Expand Down Expand Up @@ -39,10 +40,13 @@ const mockProviderInfo = {
};

const mockWalletSdkService = {
getProviderInfo: vi.fn(() => mockProviderInfo),
getSynapseClient: vi.fn(() => mockClient),
} as unknown as WalletSdkService;

const mockStorageProviderRepository = {
findByAddress: vi.fn(() => mockProviderInfo),
} as unknown as StorageProviderRepository;

const mockMetrics = {
observeCheckDuration: vi.fn(),
recordStatus: vi.fn(),
Expand Down Expand Up @@ -96,7 +100,7 @@ describe("DataSetLifecycleService", () => {

beforeEach(() => {
vi.clearAllMocks();
service = new DataSetLifecycleService(mockWalletSdkService, mockMetrics);
service = new DataSetLifecycleService(mockWalletSdkService, mockStorageProviderRepository, mockMetrics);
});

afterEach(() => {
Expand Down Expand Up @@ -331,7 +335,7 @@ describe("DataSetLifecycleService", () => {
// ─── Shared pre-flight guards ─────────────────────────────────────────────

it("throws when provider is not found in registry", async () => {
vi.mocked(mockWalletSdkService.getProviderInfo).mockReturnValueOnce(undefined);
vi.mocked(mockStorageProviderRepository.findByAddress).mockResolvedValueOnce(undefined);

await expect(service.runLifecycleCheck("0xunknown", "calibration", {})).rejects.toThrow("not found in registry");
});
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import { type ProviderJobContext, toStructuredError } from "../common/logging.js
import type { Network } from "../common/types.js";
import { buildCheckMetricLabels, classifyFailureStatus } from "../metrics-prometheus/check-metric-labels.js";
import { DataSetLifecycleCheckMetrics } from "../metrics-prometheus/check-metrics.service.js";
import { StorageProviderRepository } from "../providers/repositories/storage-provider.repository.js";
import type { SynapseViemClient } from "../wallet-sdk/wallet-sdk.service.js";
import { WalletSdkService } from "../wallet-sdk/wallet-sdk.service.js";
import type { PDPProviderEx } from "../wallet-sdk/wallet-sdk.types.js";
Expand Down Expand Up @@ -50,6 +51,7 @@ export class DataSetLifecycleService {

constructor(
private readonly walletSdkService: WalletSdkService,
private readonly storageProviderRepository: StorageProviderRepository,
private readonly lifecycleCheckMetrics: DataSetLifecycleCheckMetrics,
) {}

Expand All @@ -76,7 +78,7 @@ export class DataSetLifecycleService {
signal?: AbortSignal,
jobContext?: ProviderJobContext,
): Promise<void> {
const providerInfo = this.walletSdkService.getProviderInfo(spAddress, network);
const providerInfo = await this.storageProviderRepository.findByAddress(spAddress, network);
if (!providerInfo) {
throw new Error(`Provider ${spAddress} not found in registry`);
}
Expand Down
Original file line number Diff line number Diff line change
@@ -1,9 +1,10 @@
import { Module } from "@nestjs/common";
import { ProvidersModule } from "../providers/providers.module.js";
import { WalletSdkModule } from "../wallet-sdk/wallet-sdk.module.js";
import { DatasetLivenessService } from "./dataset-liveness.service.js";

@Module({
imports: [WalletSdkModule],
imports: [WalletSdkModule, ProvidersModule],
providers: [DatasetLivenessService],
exports: [DatasetLivenessService],
})
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import { Test, TestingModule } from "@nestjs/testing";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { StorageProviderRepository } from "../providers/repositories/storage-provider.repository.js";
import { WalletSdkService } from "../wallet-sdk/wallet-sdk.service.js";
import { DatasetLivenessService } from "./dataset-liveness.service.js";

Expand Down Expand Up @@ -29,19 +30,25 @@ describe("DatasetLivenessService", () => {
validateDataSet: vi.fn().mockResolvedValue(undefined),
};
const mockWalletSdkService = {
getProviderInfo: vi.fn().mockReturnValue({
id: 101n,
pdp: { serviceURL: "https://sp.example" },
}),
getWalletServices: vi.fn().mockReturnValue({
warmStorageService: mockWarmStorageService,
}),
getSynapseClient: vi.fn().mockReturnValue({ chain: { id: 314 } }),
};
const mockStorageProviderRepository = {
findByAddress: vi.fn().mockResolvedValue({
id: 101n,
pdp: { serviceURL: "https://sp.example" },
}),
};

beforeEach(async () => {
const module: TestingModule = await Test.createTestingModule({
providers: [DatasetLivenessService, { provide: WalletSdkService, useValue: mockWalletSdkService }],
providers: [
DatasetLivenessService,
{ provide: WalletSdkService, useValue: mockWalletSdkService },
{ provide: StorageProviderRepository, useValue: mockStorageProviderRepository },
],
}).compile();
service = module.get<DatasetLivenessService>(DatasetLivenessService);
fetchMock = vi.fn().mockResolvedValue(new Response("ok", { status: 200 }));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import { readContract } from "viem/actions";
import { awaitWithAbort } from "../common/abort-utils.js";
import { toStructuredError } from "../common/logging.js";
import type { Network } from "../common/types.js";
import { StorageProviderRepository } from "../providers/repositories/storage-provider.repository.js";
import { WalletSdkService } from "../wallet-sdk/wallet-sdk.service.js";

const PDP_LIVENESS_PROBE_TIMEOUT_MS = 10_000;
Expand All @@ -27,7 +28,10 @@ const PDP_LIVENESS_PROBE_TIMEOUT_MS = 10_000;
export class DatasetLivenessService {
private readonly logger = new Logger(DatasetLivenessService.name);

constructor(private readonly walletSdkService: WalletSdkService) {}
constructor(
private readonly walletSdkService: WalletSdkService,
private readonly storageProviderRepository: StorageProviderRepository,
) {}

async isDataSetLive(
providerAddress: string,
Expand Down Expand Up @@ -108,7 +112,7 @@ export class DatasetLivenessService {
signal?: AbortSignal,
): Promise<boolean> {
signal?.throwIfAborted();
const providerInfo = this.walletSdkService.getProviderInfo(providerAddress, network);
const providerInfo = await this.storageProviderRepository.findByAddress(providerAddress, network);
if (!providerInfo) {
throw new Error(`Provider ${providerAddress} not found in registry`);
}
Expand Down
5 changes: 3 additions & 2 deletions apps/backend/src/deal/deal.module.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,20 +3,21 @@ import { TypeOrmModule } from "@nestjs/typeorm";
import { DatabaseModule } from "../database/database.module.js";
import { Deal } from "../database/entities/deal.entity.js";
import { Retrieval } from "../database/entities/retrieval.entity.js";
import { StorageProvider } from "../database/entities/storage-provider.entity.js";
import { DataSourceModule } from "../dataSource/dataSource.module.js";
import { DatasetLivenessModule } from "../dataset-liveness/dataset-liveness.module.js";
import { DealAddonsModule } from "../deal-addons/deal-addons.module.js";
import { ProvidersModule } from "../providers/providers.module.js";
import { RetrievalAddonsModule } from "../retrieval-addons/retrieval-addons.module.js";
import { WalletSdkModule } from "../wallet-sdk/wallet-sdk.module.js";
import { DealService } from "./deal.service.js";

@Module({
imports: [
DatabaseModule,
TypeOrmModule.forFeature([Deal, Retrieval, StorageProvider]),
TypeOrmModule.forFeature([Deal, Retrieval]),
DataSourceModule,
WalletSdkModule,
ProvidersModule,
DealAddonsModule,
RetrievalAddonsModule,
DatasetLivenessModule,
Expand Down
Loading