Skip to content

Commit 41d1ddc

Browse files
committed
feat(synapse-core): add piece batcher with message-size addPieces limiter
Park and pull immediately, then coalesce on-chain addPieces using Curio piece-size bounds and the Filecoin 64KiB message budget instead of a fixed 40-piece count.
1 parent 1b8b54e commit 41d1ddc

13 files changed

Lines changed: 1432 additions & 64 deletions

File tree

packages/synapse-core/src/errors/pdp.ts

Lines changed: 40 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,10 @@
11
import type { Hash } from 'viem'
2+
import type { PieceCID } from '../piece/piece-cid.ts'
23
import type { AddPiecesRejected } from '../sp/add-pieces.ts'
34
import type { CreateDataSetRejected } from '../sp/create-dataset.ts'
45
import { SIZE_CONSTANTS } from '../utils/constants.ts'
56
import { decodePDPError } from '../utils/decode-pdp-errors.ts'
7+
import type { MetadataObject } from '../utils/metadata.ts'
68
import { isSynapseError, SynapseError } from './base.ts'
79

810
export class LocationHeaderError extends SynapseError {
@@ -133,6 +135,44 @@ export class AddPiecesError extends SynapseError {
133135
}
134136
}
135137

138+
export class AddPiecesBatchTooLargeError extends SynapseError {
139+
override name: 'AddPiecesBatchTooLargeError' = 'AddPiecesBatchTooLargeError'
140+
readonly pieceCount: number
141+
142+
constructor(pieceCount: number) {
143+
super(`Piece batch of ${pieceCount} does not fit in a single Filecoin message. Split into smaller batches.`)
144+
this.pieceCount = pieceCount
145+
}
146+
147+
static override is(value: unknown): value is AddPiecesBatchTooLargeError {
148+
return isSynapseError(value) && value.name === 'AddPiecesBatchTooLargeError'
149+
}
150+
}
151+
152+
export class AddPiecesFlushError extends SynapseError {
153+
override name: 'AddPiecesFlushError' = 'AddPiecesFlushError'
154+
readonly pieceCid: PieceCID
155+
readonly metadata?: MetadataObject
156+
/** The window that failed. Same extraData / same tx attempt. */
157+
readonly pieces: Array<{ pieceCid: PieceCID; metadata?: MetadataObject }>
158+
159+
constructor(options: {
160+
pieceCid: PieceCID
161+
metadata?: MetadataObject
162+
pieces: Array<{ pieceCid: PieceCID; metadata?: MetadataObject }>
163+
cause: Error
164+
}) {
165+
super('Failed to add pieces.', { cause: options.cause })
166+
this.pieceCid = options.pieceCid
167+
this.metadata = options.metadata
168+
this.pieces = options.pieces
169+
}
170+
171+
static override is(value: unknown): value is AddPiecesFlushError {
172+
return isSynapseError(value) && value.name === 'AddPiecesFlushError'
173+
}
174+
}
175+
136176
export class WaitForAddPiecesError extends SynapseError {
137177
override name: 'WaitForAddPiecesError' = 'WaitForAddPiecesError'
138178

Lines changed: 139 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,139 @@
1+
import { encodeAbiParameters, encodeFunctionData, type Hex, size, toHex, zeroAddress } from 'viem'
2+
import { pdpVerifierAbi } from '../abis/generated.ts'
3+
import { AddPiecesBatchTooLargeError, InvalidUploadSizeError } from '../errors/pdp.ts'
4+
import { AtLeastOnePieceRequiredError } from '../errors/warm-storage.ts'
5+
import type { PieceCID } from '../piece/piece-cid.ts'
6+
import { signAddPiecesAbiParameters } from '../typed-data/sign-add-pieces.ts'
7+
import { signCreateDataSetAbiParameters } from '../typed-data/sign-create-dataset.ts'
8+
import { signcreateDataSetAndAddPiecesAbiParameters } from '../typed-data/sign-create-dataset-add-pieces.ts'
9+
import { SIZE_CONSTANTS } from '../utils/constants.ts'
10+
import { datasetMetadataObjectToEntry, type MetadataObject, pieceMetadataObjectToEntry } from '../utils/metadata.ts'
11+
import type { PdpDataSet } from '../warm-storage/types.ts'
12+
13+
/** Dummy secp256k1 signature used only to size extraData. */
14+
const DUMMY_SIGNATURE = `0x${'00'.repeat(65)}` as Hex
15+
16+
export type LimiterPiece = {
17+
pieceCid: PieceCID
18+
metadata?: MetadataObject
19+
}
20+
21+
export type LimiterOptions =
22+
| {
23+
kind: 'addPieces'
24+
dataSet?: PdpDataSet
25+
pieces: LimiterPiece[]
26+
}
27+
| {
28+
kind: 'createDataSetAndAddPieces'
29+
metadata?: MetadataObject
30+
cdn?: boolean
31+
pieces: LimiterPiece[]
32+
}
33+
34+
/** `true` if `pieces` still fit in one addPieces / createAndAdd operation. */
35+
export type Limiter = (options: LimiterOptions) => boolean
36+
37+
export namespace addPiecesFits {
38+
export type OptionsType = LimiterOptions
39+
export type OutputType = boolean
40+
}
41+
42+
/**
43+
* Whether a candidate piece list fits in one addPieces / createAndAdd message.
44+
*
45+
* Uses estimated encoded-params size (PieceCID bytes + dummy extraData) against
46+
* {@link SIZE_CONSTANTS.MAX_ADD_PIECES_MESSAGE_SIZE} (64 KiB message cap minus
47+
* overhead). Empty `pieces` does not fit.
48+
*
49+
* @param options - {@link addPiecesFits.OptionsType}
50+
* @returns Whether the pieces fit {@link addPiecesFits.OutputType}
51+
*
52+
* @example
53+
* ```ts
54+
* import { addPiecesFits } from '@filoz/synapse-core/sp'
55+
*
56+
* const fits = addPiecesFits({
57+
* kind: 'addPieces',
58+
* dataSet,
59+
* pieces: [{ pieceCid }],
60+
* })
61+
* ```
62+
*/
63+
export function addPiecesFits(options: addPiecesFits.OptionsType): addPiecesFits.OutputType {
64+
if (options.pieces.length < 1) {
65+
return false
66+
}
67+
return estimateAddPiecesCalldataSize(options) <= SIZE_CONSTANTS.MAX_ADD_PIECES_MESSAGE_SIZE
68+
}
69+
70+
/**
71+
* Throw if a PieceCID's encoded raw size is outside Curio's upload bounds
72+
* ({@link SIZE_CONSTANTS.MIN_UPLOAD_SIZE}–{@link SIZE_CONSTANTS.MAX_UPLOAD_SIZE}).
73+
*
74+
* @throws {@link InvalidUploadSizeError}
75+
*/
76+
export function assertPieceCidSize(pieceCid: PieceCID): void {
77+
const pieceSize = pieceCid.size
78+
if (pieceSize < SIZE_CONSTANTS.MIN_UPLOAD_SIZE || pieceSize > SIZE_CONSTANTS.MAX_UPLOAD_SIZE) {
79+
throw new InvalidUploadSizeError(pieceSize)
80+
}
81+
}
82+
83+
/**
84+
* Throw if `pieces` is empty, a PieceCID is outside Curio's size bounds, or the
85+
* list does not fit in one addPieces / createAndAdd message.
86+
*
87+
* @param options - {@link LimiterOptions}
88+
* @throws {@link AtLeastOnePieceRequiredError} when `pieces` is empty
89+
* @throws {@link InvalidUploadSizeError} when a PieceCID size is below {@link SIZE_CONSTANTS.MIN_UPLOAD_SIZE} or above {@link SIZE_CONSTANTS.MAX_UPLOAD_SIZE}
90+
* @throws {@link AddPiecesBatchTooLargeError} when the estimated message exceeds {@link SIZE_CONSTANTS.MAX_ADD_PIECES_MESSAGE_SIZE}
91+
*/
92+
export function assertAddPiecesFit(options: LimiterOptions): void {
93+
if (options.pieces.length < 1) {
94+
throw new AtLeastOnePieceRequiredError()
95+
}
96+
for (const piece of options.pieces) {
97+
assertPieceCidSize(piece.pieceCid)
98+
}
99+
if (!addPiecesFits(options)) {
100+
throw new AddPiecesBatchTooLargeError(options.pieces.length)
101+
}
102+
}
103+
104+
/**
105+
* Estimated on-chain calldata size in bytes for the given piece list.
106+
*/
107+
export function estimateAddPiecesCalldataSize(options: LimiterOptions): number {
108+
const extraData = dummyExtraData(options)
109+
const pieceData = options.pieces.map((piece) => ({ data: toHex(piece.pieceCid.bytes) }))
110+
const calldata = encodeFunctionData({
111+
abi: pdpVerifierAbi,
112+
functionName: 'addPieces',
113+
args: [0n, zeroAddress, pieceData, extraData],
114+
})
115+
return size(calldata)
116+
}
117+
118+
function dummyExtraData(options: LimiterOptions): Hex {
119+
const addPiecesExtraData = dummyAddPiecesExtraData(options.pieces)
120+
if (options.kind === 'addPieces') {
121+
return addPiecesExtraData
122+
}
123+
const createEntries = datasetMetadataObjectToEntry(options.metadata, { cdn: options.cdn ?? false })
124+
const createExtraData = encodeAbiParameters(signCreateDataSetAbiParameters, [
125+
zeroAddress,
126+
0n,
127+
createEntries.map((entry) => entry.key),
128+
createEntries.map((entry) => entry.value),
129+
DUMMY_SIGNATURE,
130+
])
131+
return encodeAbiParameters(signcreateDataSetAndAddPiecesAbiParameters, [createExtraData, addPiecesExtraData])
132+
}
133+
134+
function dummyAddPiecesExtraData(pieces: LimiterPiece[]): Hex {
135+
const metadataKV = pieces.map((piece) => pieceMetadataObjectToEntry(piece.metadata))
136+
const keys = metadataKV.map((entries) => entries.map((entry) => entry.key))
137+
const values = metadataKV.map((entries) => entries.map((entry) => entry.value))
138+
return encodeAbiParameters(signAddPiecesAbiParameters, [0n, keys, values, DUMMY_SIGNATURE])
139+
}

packages/synapse-core/src/sp/add-pieces.ts

Lines changed: 11 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -3,13 +3,15 @@ import type { ToString } from 'multiformats'
33
import { type Account, type Chain, type Client, type Hex, isHex, type Transport } from 'viem'
44
import * as z from 'zod'
55
import { AddPiecesError, LocationHeaderError } from '../errors/index.ts'
6+
import type { AddPiecesBatchTooLargeError, InvalidUploadSizeError } from '../errors/pdp.ts'
67
import { WaitForAddPiecesError, WaitForAddPiecesRejectedError } from '../errors/pdp.ts'
7-
import { AtLeastOnePieceRequiredError, TooManyPiecesError } from '../errors/warm-storage.ts'
8+
import type { AtLeastOnePieceRequiredError } from '../errors/warm-storage.ts'
89
import type { PieceCID } from '../piece/piece-cid.ts'
910
import { signAddPieces } from '../typed-data/sign-add-pieces.ts'
10-
import { RETRY_CONSTANTS, SIZE_CONSTANTS } from '../utils/constants.ts'
11+
import { RETRY_CONSTANTS } from '../utils/constants.ts'
1112
import { type MetadataObject, pieceMetadataObjectToEntry } from '../utils/metadata.ts'
1213
import { zHex, zNumberToBigInt } from '../utils/schemas.ts'
14+
import { assertAddPiecesFit } from './add-pieces-fits.ts'
1315

1416
export namespace addPiecesApiRequest {
1517
export type OptionsType = {
@@ -114,24 +116,12 @@ export namespace addPieces {
114116
}
115117

116118
export type OutputType = addPiecesApiRequest.OutputType
117-
export type ErrorType = addPiecesApiRequest.ErrorType | signAddPieces.ErrorType
118-
}
119-
120-
/**
121-
* Validate the piece count for an addPieces (or createDataSetAndAddPieces) batch,
122-
* failing early instead of reverting on-chain.
123-
*
124-
* @param pieceCount - Number of pieces in the batch
125-
* @throws AtLeastOnePieceRequiredError when not a positive integer
126-
* @throws TooManyPiecesError when above {@link SIZE_CONSTANTS.MAX_ADD_PIECES_BATCH_SIZE}
127-
*/
128-
export function validateAddPiecesBatch(pieceCount: number): void {
129-
if (!Number.isInteger(pieceCount) || pieceCount < 1) {
130-
throw new AtLeastOnePieceRequiredError()
131-
}
132-
if (pieceCount > SIZE_CONSTANTS.MAX_ADD_PIECES_BATCH_SIZE) {
133-
throw new TooManyPiecesError(pieceCount, SIZE_CONSTANTS.MAX_ADD_PIECES_BATCH_SIZE)
134-
}
119+
export type ErrorType =
120+
| addPiecesApiRequest.ErrorType
121+
| signAddPieces.ErrorType
122+
| AtLeastOnePieceRequiredError
123+
| AddPiecesBatchTooLargeError
124+
| InvalidUploadSizeError
135125
}
136126

137127
/**
@@ -148,7 +138,7 @@ export async function addPieces(
148138
client: Client<Transport, Chain, Account>,
149139
options: addPieces.OptionsType
150140
): Promise<addPieces.OutputType> {
151-
validateAddPiecesBatch(options.pieces.length)
141+
assertAddPiecesFit({ kind: 'addPieces', pieces: options.pieces })
152142
const extraData =
153143
options.extraData ??
154144
(await signAddPieces(client, {

packages/synapse-core/src/sp/create-dataset-add-pieces.ts

Lines changed: 17 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -4,16 +4,20 @@ import { type Account, type Address, type Chain, type Client, type Hex, isHex, t
44
import { asChain } from '../chains.ts'
55
import { CreateDataSetError, LocationHeaderError } from '../errors/index.ts'
66
import type {
7+
AddPiecesBatchTooLargeError,
8+
InvalidUploadSizeError,
79
WaitForAddPiecesError,
810
WaitForAddPiecesRejectedError,
911
WaitForCreateDataSetError,
1012
WaitForCreateDataSetRejectedError,
1113
} from '../errors/pdp.ts'
14+
import type { AtLeastOnePieceRequiredError } from '../errors/warm-storage.ts'
1215
import type { PieceCID } from '../piece/piece-cid.ts'
1316
import { signCreateDataSetAndAddPieces } from '../typed-data/sign-create-dataset-add-pieces.ts'
1417
import { RETRY_CONSTANTS } from '../utils/constants.ts'
1518
import { datasetMetadataObjectToEntry, type MetadataObject, pieceMetadataObjectToEntry } from '../utils/metadata.ts'
16-
import { validateAddPiecesBatch, waitForAddPieces } from './add-pieces.ts'
19+
import { waitForAddPieces } from './add-pieces.ts'
20+
import { assertAddPiecesFit } from './add-pieces-fits.ts'
1721
import { waitForCreateDataSet } from './create-dataset.ts'
1822

1923
export namespace createDataSetAndAddPiecesApiRequest {
@@ -129,7 +133,12 @@ export type CreateDataSetAndAddPiecesOptions = {
129133
export namespace createDataSetAndAddPieces {
130134
export type OptionsType = CreateDataSetAndAddPiecesOptions
131135
export type ReturnType = createDataSetAndAddPiecesApiRequest.OutputType
132-
export type ErrorType = createDataSetAndAddPiecesApiRequest.ErrorType | signCreateDataSetAndAddPieces.ErrorType
136+
export type ErrorType =
137+
| createDataSetAndAddPiecesApiRequest.ErrorType
138+
| signCreateDataSetAndAddPieces.ErrorType
139+
| AtLeastOnePieceRequiredError
140+
| AddPiecesBatchTooLargeError
141+
| InvalidUploadSizeError
133142
}
134143

135144
/**
@@ -144,7 +153,12 @@ export async function createDataSetAndAddPieces(
144153
client: Client<Transport, Chain, Account>,
145154
options: CreateDataSetAndAddPiecesOptions
146155
): Promise<createDataSetAndAddPieces.ReturnType> {
147-
validateAddPiecesBatch(options.pieces.length)
156+
assertAddPiecesFit({
157+
kind: 'createDataSetAndAddPieces',
158+
metadata: options.metadata,
159+
cdn: options.cdn,
160+
pieces: options.pieces,
161+
})
148162
const chain = asChain(client.chain)
149163
const extraData =
150164
options.extraData ??

0 commit comments

Comments
 (0)