Skip to content

Commit 67174ff

Browse files
committed
refactor: streamline repair dataset handling and enhance CLI functionality
- Updated the repair CLI to target a single IPFS-enabled dataset, simplifying dataset management. - Removed obsolete dataset grouping logic and refactored related database queries for improved clarity. - Enhanced error handling for missing datasets and improved command responses to include relevant dataset information. - Introduced new utility functions for dataset retrieval and updated existing functions to align with the new dataset structure. - Updated TypeScript types for better type safety and clarity in the codebase.
1 parent 2a339ba commit 67174ff

21 files changed

Lines changed: 365 additions & 691 deletions

packages/repair-cli/AGENTS.md

Lines changed: 25 additions & 39 deletions
Original file line numberDiff line numberDiff line change
@@ -13,45 +13,37 @@ Two databases:
1313

1414
Commands get both via `contextMiddleware` (`middleware.ts`): wallet client from config, indexer URL by `chainId` (314 = mainnet, else calibration).
1515

16-
## Piece groups
16+
## Repair dataset
1717

18-
Pieces are grouped by dataset flags (`withCdn`, `withIpfsIndexing`). Groups are **mutually exclusive**:
18+
Repairs now target a single dataset: IPFS indexing enabled and CDN disabled.
1919

20-
- `both`: `withCdn = true`, `withIpfsIndexing = true`
21-
- `cdn`: `withCdn = true`, `withIpfsIndexing = false`
22-
- `ipfs`: `withCdn = false`, `withIpfsIndexing = true`
23-
- `none`: `withCdn = false`, `withIpfsIndexing = false`
24-
25-
Same CID may appear on multiple datasets in one group; dedupe per group when listing or paginating pieces.
26-
27-
Target datasets are looked up with payer + `EARLY_REPAIR_SOURCE` (`utils.ts`, value `early-repair2`).
20+
All source pieces are deduped globally by CID and repaired into this one target dataset. Target datasets are looked up with target provider + payer + `EARLY_REPAIR_SOURCE` (`utils.ts`) + `withIpfsIndexing = true` + `withCdn = false`.
2821

2922
## Repair pipeline
3023

31-
### 1. `repair create --provider-id <id>`
24+
### 1. `repair create --provider-id <id> --target-provider-id <id>`
3225

3326
`createRepair` (`db/create-repair.ts`):
3427

35-
1. **Target provider** — if `--target-provider-id` is set, **`getRepairProvider`** loads that active provider; otherwise **`selectAlternateRepairProvider`** picks one with tier-matched fallback (endorsed → approved → none). Throws if none found.
36-
2. **`getDataSetsByGroup`** — target provider datasets per group.
37-
3. Insert **`repairs`** row (`repairProviderId`, `targetProviderId`).
38-
4. **`forEachPiecesPage`** — paginated `add_piece` operations (`db/get-pieces.ts`, page size 500).
39-
5. **`getRepairGroups`** — distinct groups from pending `add_piece` ops; saved on the repair row.
40-
6. For each group with pending pieces but no target dataset → insert **`create_dataset`** operation (`pending`).
28+
1. **Target provider****`getRepairProvider`** loads the required `--target-provider-id`. Throws if none is found.
29+
2. Insert **`repairs`** row (`repairProviderId`, `targetProviderId`).
30+
3. **`forEachPiecesPage`** — paginated `add_piece` operations (`db/get-pieces.ts`).
31+
4. **`getRepairDataset`** — find the target IPFS-enabled dataset for the repair wallet.
32+
5. If no target dataset exists → insert one **`create_dataset`** operation (`pending`).
4133

4234
### 2. `repair run <repairId>`
4335

44-
1. Run **`create_dataset`** phase (`pipeline/create-datasets.ts`) via fastq `SP.createDataSet`, then `updateOperation`; returns created dataset IDs indexed by group.
45-
2. Run **`add_piece`** pull phase (`pipeline/pull.ts`): pending ops are fetched in same-group pages and fed into fastq as workers free up (`createPullPiecesWorker` — mock logs CIDs per batch).
36+
1. Run **`create_dataset`** phase (`pipeline/create-datasets.ts`) for the single pending dataset operation `SP.createDataSet`, then `updateOperation`; stores the single IPFS target dataset ID.
37+
2. Run **`add_piece`** pull phase (`pipeline/pull.ts`) via `p-queue`: pending ops are fetched in ID order with bounded queue backpressure.
4638

47-
`--reset` retries `pending` and `failed` `create_dataset` ops only; `add_piece` always runs `pending` (failed pieces are skipped). `--batch-size` (default 50) caps pieces per pull job; each batch is one repair group only.
39+
`--reset` retries `pending` and `failed` `create_dataset` ops and also includes failed `add_piece` ops; otherwise `add_piece` runs only `pending` operations. `--batch-size` caps pieces per pull job.
4840

4941
### Operation types
5042

5143
- `create_dataset`: `pending``committing``completed` | `failed`; data has `serviceUrl`, `payee`.
5244
- `add_piece`: `pending``pulling``committing``completed` | `failed`; data has `cid`, `serviceUrl`, `metadata`, `alternateProviders`.
5345

54-
`add_piece` without alternate providers (other replicas) is created as **`failed`** with error `"No alternate providers found"`. `getProvidersByCid` excludes the source `providerId`.
46+
`add_piece` without alternate providers (other replicas) is created as **`skipped`** with error `"No alternate providers found"`.
5547

5648
## Source layout
5749

@@ -66,10 +58,8 @@ src/
6658
wallet.ts
6759
db/
6860
create-repair.ts # createRepair orchestration
69-
get-repair-groups.ts # source groups that need repair
70-
get-datasets-by-group.ts # target datasets per group
61+
get-repair-dataset.ts # target IPFS-enabled dataset
7162
get-providers-by-cid.ts # alternate providers per CID
72-
select-alternate-repair-provider.ts # automatic target provider selection
7363
get-repair-provider.ts # load explicit target provider by ID
7464
update-operation.ts # patch local operation status/result/error
7565
delete-repair.ts # delete repair and its operations
@@ -80,7 +70,7 @@ src/
8070
local-schema.ts # SQLite repairs/operations
8171
indexer-schema.ts # Postgres early-repair schema
8272
middleware.ts # DB + wallet context
83-
types.ts # Group, PIECE_GROUPS, DB types
73+
types.ts # Shared DB and context types
8474
error.ts # NoAlternateProviderError, RepairCreationError
8575
utils.ts # config, client, metadata helpers
8676
```
@@ -89,25 +79,21 @@ src/
8979

9080
- Indexer helpers take `IndexerQueryOptions` (`indexerDb`, `indexerSchema`).
9181
- Extract DB helpers under `src/db/` (indexer queries and local operation updates).
92-
- Add JSDoc on exported functions/types; inline comments only for non-obvious logic (dedupe, pagination, tier fallback).
93-
- Use `PIECE_GROUPS` instead of `Object.keys` for group iteration.
82+
- Add JSDoc on exported functions/types; inline comments only for non-obvious logic (dedupe, pagination).
83+
- Repairs use one IPFS-enabled target dataset; do not add per-operation dataset grouping.
9484

9585
## Indexer API (`src/db/`)
9686

97-
| Function | Module | Purpose |
98-
| ------------------------------- | ----------------------------------- | ------------------------------------------------------- |
99-
| `getRepairGroups` | `get-repair-groups.ts` | Repair groups from pending local `add_piece` operations |
100-
| `getDataSetsByGroup` | `get-datasets-by-group.ts` | One dataset per group for payer + `EARLY_REPAIR_SOURCE` |
101-
| `getProvidersByCid` | `get-providers-by-cid.ts` | Alternate providers per CID; empty array if none |
102-
| `selectAlternateRepairProvider` | `select-alternate-repair-provider.ts` | Automatic target provider selection |
103-
| `getRepairProvider` | `get-repair-provider.ts` | Load explicit target provider by ID |
104-
| `filterPullPiecesNotInDataset` | `filter-pull-pieces-not-in-dataset.ts` | Exclude CIDs already indexed in a dataset |
105-
| `updateOperation` | `update-operation.ts` | Patch local operation status/result/error |
87+
- `getRepairDataset` (`get-repair-dataset.ts`): single IPFS-enabled dataset for payer + `EARLY_REPAIR_SOURCE`.
88+
- `getProvidersByCid` (`get-providers-by-cid.ts`): alternate providers per CID; empty array if none.
89+
- `getRepairProvider` (`get-repair-provider.ts`): load explicit target provider by ID.
90+
- `filterPullPiecesNotInDataset` (`filter-pull-pieces-not-in-dataset.ts`): exclude CIDs already indexed in a dataset.
91+
- `updateOperation` (`update-operation.ts`): patch local operation status/result/error.
10692

10793
## Local schema
10894

109-
- **`repairs`**: `repairProviderId`, `targetProviderId`, `repairGroups`, `status` (`pending` \| `running` \| `completed` \| `failed`).
110-
- **`operations`**: `type`, `group`, `status`, `data` (JSON), `result`, `error`.
95+
- **`repairs`**: `repairProviderId`, `targetProviderId`, `targetDataSetId`, `status` (`pending` \| `completed` \| `failed`).
96+
- **`operations`**: `type`, `status`, `data` (JSON), `result`, `error`.
11197

11298
## Commands
11399

@@ -118,7 +104,7 @@ src/
118104
| `repair wallet balance` | Wallet FIL/USDFC balances and pay account summary |
119105
| `repair wallet deposit <amount>` | Deposit USDFC to pay account |
120106
| `repair wallet withdraw <amount>` | Withdraw USDFC from pay account |
121-
| `repair repair create --provider-id <id>` | Plan repair; optional `--target-provider-id`; returns `repairId` |
107+
| `repair repair create --provider-id <id> --target-provider-id <id>` | Plan repair; returns `repairId` |
122108
| `repair repair list` | List repairs with operation counts |
123109
| `repair repair delete <repairId>` | Delete a repair and its operations |
124110
| `repair repair run <repairId>` | Execute workers; `--concurrency`, `--batch-size`, `--reset` |

packages/repair-cli/package.json

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -110,10 +110,13 @@
110110
"conf": "^15.1.0",
111111
"drizzle-kit": "^0.31.10",
112112
"drizzle-orm": "catalog:",
113-
"fastq": "^1.20.1",
114113
"incur": "^0.4.6",
114+
"iso-base": "^4.4.0",
115115
"iso-web": "^2.2.1",
116+
"p-all": "^5.0.1",
116117
"p-locate": "^7.0.0",
118+
"p-map": "^7.0.4",
119+
"p-queue": "^9.3.0",
117120
"pg": "^8.21.0",
118121
"terminal-link": "^5.0.0"
119122
},

packages/repair-cli/src/commands/repair.ts

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,7 @@ repair.command('create', {
1515
description: 'Create a new repair',
1616
options: globalOptions.extend({
1717
providerId: z.coerce.bigint().describe('Provider ID to repair'),
18-
targetProviderId: z.coerce.bigint().optional().describe('Target provider ID for repair'),
18+
targetProviderId: z.coerce.bigint().describe('Target provider ID for repair'),
1919
}),
2020
middleware: [contextMiddleware],
2121
run: async (c) => {
@@ -32,6 +32,7 @@ repair.command('create', {
3232
repairId,
3333
})
3434
} catch (error) {
35+
console.error(error)
3536
return c.error({
3637
code: 'REPAIR_FAILED',
3738
message: error instanceof Error ? error.message : 'Failed to repair the dataset',
@@ -65,7 +66,8 @@ repair.command('list', {
6566
repairProviderId: repairWithoutOperations.repairProviderId,
6667
targetProviderId: repairWithoutOperations.targetProviderId,
6768
targetProviderUrl: repairWithoutOperations.targetProviderUrl,
68-
targetDataSets: repairWithoutOperations.targetDataSets,
69+
targetDataSetId: repairWithoutOperations.targetDataSetId,
70+
blockNumber: repairWithoutOperations.blockNumber,
6971
createdAt: new Date(repairWithoutOperations.createdAt).toISOString(),
7072
updatedAt: new Date(repairWithoutOperations.updatedAt).toISOString(),
7173
operations: operations.length,
@@ -130,7 +132,7 @@ repair.command('run', {
130132
}),
131133
options: globalOptions.extend({
132134
concurrency: z.coerce.number().default(4).describe('Concurrency level'),
133-
batchSize: z.coerce.number().default(10).describe('Max add_piece operations per pull batch (same group)'),
135+
batchSize: z.coerce.number().default(10).describe('Max add_piece operations per pull batch'),
134136
reset: z.boolean().default(false).describe('Reset the repair'),
135137
}),
136138
middleware: [contextMiddleware],
@@ -152,7 +154,6 @@ repair.command('run', {
152154
await runCreateDatasetsPhase({
153155
...c.var,
154156
repair,
155-
concurrency,
156157
reset: c.options.reset,
157158
})
158159

packages/repair-cli/src/db/create-repair.ts

Lines changed: 35 additions & 42 deletions
Original file line numberDiff line numberDiff line change
@@ -1,16 +1,14 @@
11
import { eq } from 'drizzle-orm'
2+
import { getBlockNumber } from 'viem/actions'
23
import { NoAlternateProviderError, RepairCreationError } from '../error.ts'
3-
import type { RepairTargetDataSets } from '../local-schema.ts'
44
import type { Context } from '../types.ts'
5-
import { getDataSetsByGroup } from './get-datasets-by-group.ts'
65
import { forEachPiecesPage } from './get-pieces.ts'
7-
import { getRepairGroups } from './get-repair-groups.ts'
6+
import { getRepairDataset } from './get-repair-dataset.ts'
87
import { getRepairProvider } from './get-repair-provider.ts'
9-
import { selectAlternateRepairProvider } from './select-alternate-repair-provider.ts'
108

119
export interface CreateRepairOptions extends Context {
1210
repairProviderId: bigint
13-
targetProviderId?: bigint
11+
targetProviderId: bigint
1412
}
1513

1614
/**
@@ -24,20 +22,16 @@ export async function createRepair(options: CreateRepairOptions): Promise<number
2422
const { indexerDb, localDb, repairProviderId, targetProviderId, client } = options
2523
const localSchema = localDb._.fullSchema
2624
const now = Date.now()
25+
const blockNumber = await getBlockNumber(client)
2726

28-
// Select the target provider
29-
if (targetProviderId != null && targetProviderId === repairProviderId) {
27+
// Load the explicit target provider.
28+
if (targetProviderId === repairProviderId) {
3029
throw new RepairCreationError('Target provider must differ from the provider being repaired')
3130
}
32-
const targetProvider = targetProviderId
33-
? await getRepairProvider({
34-
indexerDb,
35-
providerId: targetProviderId,
36-
})
37-
: await selectAlternateRepairProvider({
38-
indexerDb,
39-
providerId: repairProviderId,
40-
})
31+
const targetProvider = await getRepairProvider({
32+
indexerDb,
33+
providerId: targetProviderId,
34+
})
4135

4236
if (!targetProvider) {
4337
throw new NoAlternateProviderError(targetProviderId)
@@ -50,7 +44,8 @@ export async function createRepair(options: CreateRepairOptions): Promise<number
5044
repairProviderId,
5145
targetProviderId: targetProvider.providerId,
5246
targetProviderUrl: targetProvider.serviceUrl,
53-
targetDataSets: {},
47+
targetDataSetId: null,
48+
blockNumber,
5449
createdAt: now,
5550
updatedAt: now,
5651
})
@@ -59,52 +54,50 @@ export async function createRepair(options: CreateRepairOptions): Promise<number
5954
if (!repair) throw new RepairCreationError()
6055

6156
// Add the pieces to the repair
57+
let hasPendingPieces = false
6258
await forEachPiecesPage(
6359
{
6460
indexerDb,
6561
providerId: repairProviderId,
6662
repairId: repair.id,
63+
blockNumber,
6764
},
6865
async (page) => {
6966
if (page.operations.length > 0) {
67+
hasPendingPieces ||= page.operations.some((operation) => operation.status === 'pending')
7068
await localDb.insert(localSchema.operations).values(page.operations)
7169
}
7270
}
7371
)
7472

75-
// Get the repair groups
76-
const repairGroups = await getRepairGroups({ localDb, repairId: repair.id })
77-
// Get the target datasets for the repair. If no dataset is found, a new one will be created.
78-
const dataSetsByGroup = await getDataSetsByGroup({
73+
if (!hasPendingPieces) {
74+
return repair.id
75+
}
76+
77+
// Get the single IPFS-enabled target dataset for the repair. If none exists, create one before pulling.
78+
const targetDataset = await getRepairDataset({
7979
indexerDb,
8080
providerId: targetProvider.providerId,
8181
payer: client.account.address,
82+
blockNumber,
8283
})
8384

84-
// Each group needs a target dataset before pieces can be added; queue missing datasets first.
85-
const targetDataSets: RepairTargetDataSets = {}
86-
for (const group of repairGroups) {
87-
const dataSet = dataSetsByGroup[group]
88-
if (dataSet) {
89-
targetDataSets[group] = dataSet.dataSetId
90-
} else {
91-
targetDataSets[group] = null
92-
await localDb.insert(localSchema.operations).values({
93-
repairId: repair.id,
94-
type: 'create_dataset',
95-
group,
96-
status: 'pending',
97-
data: {
98-
payee: targetProvider.providerAddress,
99-
},
100-
createdAt: now,
101-
updatedAt: now,
102-
})
103-
}
85+
if (!targetDataset) {
86+
await localDb.insert(localSchema.operations).values({
87+
repairId: repair.id,
88+
type: 'create_dataset',
89+
status: 'pending',
90+
data: {
91+
payee: targetProvider.providerAddress,
92+
},
93+
createdAt: now,
94+
updatedAt: now,
95+
})
10496
}
97+
10598
await localDb
10699
.update(localSchema.repairs)
107-
.set({ targetDataSets, updatedAt: now })
100+
.set({ targetDataSetId: targetDataset?.dataSetId ?? null, updatedAt: now })
108101
.where(eq(localSchema.repairs.id, repair.id))
109102

110103
return repair.id

packages/repair-cli/src/db/get-dataset-for-group.ts

Lines changed: 0 additions & 50 deletions
This file was deleted.

0 commit comments

Comments
 (0)