Compare commits

..

4 Commits

Author SHA1 Message Date
thloyi
f40c548fc9 update amm virtual quotes reserves 2026-07-17 16:56:44 +08:00
thloyi
44eecac087 tx binary with block time 2026-06-05 10:35:49 +08:00
thloyi
9f17ffce61 fix accounts len check 2026-05-28 10:23:56 +08:00
thloyi
e4eaddec4e reademe 2026-05-18 14:22:48 +08:00
14 changed files with 1029 additions and 115 deletions

251
README.md Normal file
View File

@@ -0,0 +1,251 @@
# pump-parser
Solana transaction parser focused on swap, liquidity, migration, platform, MEV, compute budget, and compact binary persistence workflows.
The package works with a normalized `RawTx` representation, parses it into `Tx`, and emits one or more `Swap` records when a supported protocol action is found.
## Features
- Parse Solana RPC / Yellowstone transactions into local `RawTx`.
- Extract swap and liquidity events into `Tx.Swaps`.
- Preserve transaction metadata such as slot, block time, signer, fee, CU limit, CU consumed, token balances, platform fees, and MEV agent hints.
- Encode/decode parsed transactions with the `PTXB` / `PTXS` binary formats.
- Encode/decode raw transactions with the `PRTX` / `PRTS` / `PRBS` binary formats.
- Stream decode large `PTXS` / `PRTS` payloads without loading every transaction into memory.
- Merge `PTXS` batches while remapping address tables.
## Supported Parsers
Default parser initialization enables the hot-path Pump parsers only:
- Pump
- Pump AMM
`EnableAllParsers()` additionally enables:
- Meteora DLMM
- Meteora Pools
- Meteora DAMM v2
- Meteora Bonding Curve
- Orca Whirlpool
- Raydium AMM v4
- Raydium CLMM
- Raydium CPMM
- Raydium LaunchLab
Use `InitParser(WithMeteoraDlmm())` when only Meteora DLMM should be added to the default parser set.
## Installation
```bash
go get github.com/thloyi/pump-parser
```
This module currently declares `go 1.25.1` in `go.mod`.
## Basic Usage
```go
package main
import (
"fmt"
"log"
pump_parser "github.com/thloyi/pump-parser"
)
func main() {
// Enable all known parser programs. Omit this call to keep the default
// Pump + Pump AMM parser set.
pump_parser.EnableAllParsers()
var rawTx *pump_parser.RawTx
// Fill rawTx from RPC, Yellowstone, JSON, or RawTx binary decoding.
tx, err := pump_parser.ParseRawTx(rawTx)
if err != nil {
log.Fatal(err)
}
fmt.Println("tx:", tx.GetTxHash())
for _, swap := range tx.Swaps {
fmt.Printf("%s %s base=%s quote=%s pool=%s user=%s\n",
swap.Program,
swap.Event,
swap.BaseAmount,
swap.QuoteAmount,
swap.Pool,
swap.User,
)
}
}
```
## Swap Result
`ParseRawTx` returns a `Tx`. A transaction can contain zero, one, or multiple parsed swap records in `Tx.Swaps`.
Each `Swap` describes one protocol-level swap, liquidity, migration, or pool operation that the parser can normalize. The most important fields are:
| Field | Meaning |
| --- | --- |
| `Program` | Normalized protocol name, such as `Pump`, `PumpAMM`, `MeteoraDLMM`, `RaydiumV4`, or `OrcaWhirPool`. |
| `Event` | Normalized action, such as `buy`, `sell`, `add_liquidity`, `remove_liquidity`, `create`, `complete`, or `migrate`. |
| `TxIndex` | Approximate execution order across outer and inner instructions. |
| `InstrIdx` / `InnerIdx` | Outer instruction index and inner instruction index where the swap was found. |
| `Pool` | Pool, pair, bonding curve, or market account for the parsed action. |
| `BaseMint` / `QuoteMint` | Base and quote mint accounts after parser normalization. |
| `BaseTokenProgram` / `QuoteTokenProgram` | Token program for each side, useful when Token-2022 is involved. |
| `BaseMintDecimals` / `QuoteMintDecimals` | Decimals used to interpret raw token amounts. |
| `User` | User or effective owner account for the action. If the parsed user is not on-curve, the parser may fall back to the transaction signer. |
| `BaseAmount` / `QuoteAmount` | Actual parsed base-side and quote-side amounts, stored as `decimal.Decimal`. |
| `BaseReserve` / `QuoteReserve` | Pool reserves when the protocol event or accounts expose them. For Pump AMM buy/sell events, `QuoteReserve` is the post-trade effective quote reserve. |
| `RealQuoteReserve` / `VirtualQuoteReserve` | Quote-reserve components. Pump AMM buy/sell events expose the post-trade real reserve and the signed virtual reserve separately. |
| `UserBaseBalance` / `UserQuoteBalance` | User token balances after the transaction when available from token balance metadata. |
| `AfterSOLBalance` | User or signer SOL balance after the transaction. |
| `EntryContract` | Known router / entry contract account when detected. |
| `Mayhem` / `Cashback` | Protocol or platform-specific labels detected from the transaction path. |
Swap amount and slippage fields normalize instruction intent:
| Field | Meaning |
| --- | --- |
| `SwapMode` | `exact_in`, `exact_out`, or empty when unknown. |
| `FixedAmount` / `FixedAmountSide` / `FixedMint` | The user-specified fixed side of the swap. For `exact_in`, this is the input amount. For `exact_out`, this is the target output amount. |
| `LimitAmountType` | `min_out` for `exact_in`, `max_in` for `exact_out`, or empty when unknown. |
| `LimitAmount` / `LimitAmountSide` / `LimitMint` | User-specified limit on the opposite side. |
| `ActualLimitAmount` / `ActualLimitAmountSide` | Actual executed amount on the limited side. |
| `SlippageBps` | Remaining headroom to the user's limit in basis points. See `SLIPPAGE_MAPPING.md` for protocol-specific derivation rules. |
Liquidity, migration, and DLMM-specific fields are populated only for protocols that expose them:
| Field | Meaning |
| --- | --- |
| `Creator` | Creator account when exposed by pool or launch events. |
| `MigrateToPool` / `MigrateTopProgram` | Destination pool and program for migration events. |
| `LpMint` | LP token mint for liquidity operations when available. |
| `ActiveBinId` / `StartBinId` / `EndBinId` | Meteora DLMM bin identifiers. |
| `RemoveBp` | DLMM remove-liquidity basis points. |
| `PositionAccount` | DLMM position account. |
| `FeeAmount` / `LpFeeAmount` | Parsed fee amounts when the protocol exposes fee breakdowns. |
| `FeeSide` / `FeeMint` / `FeeTokenProgram` / `FeeMintDecimals` | Fee side and mint metadata. |
| `ConsumeUnit` | Per-swap compute unit value when available. |
Amounts are normalized into decimals in parser output, but the exact scale depends on the source event and parser path. For persisted binary output, `tx_binary.go` defines the conversion rules and schema version.
## Convert Transactions
RPC transactions can be converted with `FromRpcTransactionWithMeta`:
```go
rawTx, err := pump_parser.FromRpcTransactionWithMeta(
txWithMeta,
blockTime,
slot,
indexWithinBlock,
)
```
Yellowstone transactions can be converted with `ConvertYellowstoneGrpcTransactionToSolanaTransaction`:
```go
rawTx, err := pump_parser.ConvertYellowstoneGrpcTransactionToSolanaTransaction(
update,
createdUnix,
)
```
Both conversion paths accept optional `RawTxConvertOptions`:
```go
rawTx, err := pump_parser.FromRpcTransactionWithMeta(
txWithMeta,
blockTime,
slot,
indexWithinBlock,
pump_parser.RawTxConvertOptions{
ParseLogEvents: true,
},
)
```
`ParseLogEvents` attaches decoded program log events to matching instructions. `IgnoreLogMessages` skips log-message retention.
## Binary Formats
Parsed transaction binary:
- `EncodeTxBinary` / `DecodeTxBinary` for one parsed `Tx`.
- `EncodeTxsBinary` / `DecodeTxsBinary` for a batch of parsed `Tx`.
- `DecodeTxsBinaryReader` for streaming `PTXS` reads.
- `MergeTxsBinaryBytes` and `MergeTxsBinarySourcesToWriter` for merging `PTXS` batches.
- Schema v5 persists `RealQuoteReserve` and signed `VirtualQuoteReserve`; v3/v4 inputs remain readable and default these components to legacy `QuoteReserve` and zero.
Raw transaction binary:
- `EncodeRawTxBinary` / `DecodeRawTxBinary` for one `RawTx`.
- `EncodeRawTxsBinary` / `DecodeRawTxsBinary` for a batch of `RawTx`.
- `EncodeRawTxBlocksBinary` / `DecodeRawTxBlocksBinary` for grouped block data.
- `DecodeRawTxsBinaryReader` for streaming `PRTS` reads.
Format magic values:
- `PTXB`: one parsed `Tx`.
- `PTXS`: parsed `Tx` batch.
- `PRTX`: one raw `RawTx`.
- `PRTS`: raw `RawTx` batch.
- `PRBS`: raw block-grouped `RawTx` batch.
When adding or renaming transaction-facing enum values, update `tx_binary.go` enum tables by appending new values only. Reordering existing values changes persisted numeric IDs.
## Commands
Parse a live transaction by signature:
```bash
TX_HASH=<signature> go run ./cmd/rpc_parse
```
Collect Yellowstone transactions into a `.prbs` raw-block binary file:
```bash
YELLOWSTONE_X_TOKEN=<token> go run ./cmd/collect_yellowstone_rawtx_binary \
-endpoint ams.rpc.orbitflare.com:10000 \
-duration 5m \
-output testdata/rawtx-binary/sample.prbs
```
Measure parsed `Tx` binary size from a `getBlock` JSON payload:
```bash
go run ./cmd/measure_tx_binary_block \
-file /path/to/block.json \
-slot 413539056 \
-swaps-only
```
Analyze raw transaction binary size distribution:
```bash
go run ./cmd/analyze_rawtx_binary_size \
-file testdata/rawtx-binary/sample.prbs
```
## Development
Run the full test suite:
```bash
go test ./...
```
Useful focused checks:
```bash
go test -run 'TestTxBinary|TestRawTxBinary' .
go test -run 'TestMetaoraPoolSwapEventFromLogsUsesMatchingInvocation|TestMetaoraPoolSwapInstructionOccurrenceIncludesInnerInstructions|TestAttachLogEventsToInstructions' .
TX_HASH=<signature> go run ./cmd/rpc_parse
```
For parser regressions, reproduce with the exact transaction signature first, then add a targeted unit test or fixture around the corrected parser path.

View File

@@ -25,6 +25,9 @@ func chainLinkParser(tx *Tx, instruction Instruction, inners InnerInstructions,
}
decode := instruction.Data
if len(decode) < 4 {
return increaseOffset(offset), nil
}
discriminator := binary.LittleEndian.Uint32(decode[0:4])
switch discriminator {
@@ -55,6 +58,9 @@ func chainLinkSubmitParser(instruction Instruction, inners InnerInstructions, of
if storeInstruction.Accounts[0] >= len(tx.rawTx.accountList) || tx.rawTx.accountList[storeInstruction.Accounts[0]] != chainlinkSOLUSDFeedAccount {
return increaseOffset(offset), InstructionIgnoredError
}
if len(storeInstruction.Data) < 8 {
return increaseOffset(offset), InstructionIgnoredError
}
if !bytes.Equal(storeInstruction.Data[0:8], chainlinkSubmitDiscriminator[:]) {
return increaseOffset(offset), InstructionIgnoredError
}

View File

@@ -741,6 +741,9 @@ func metaoraPoolRemoveLiquidity(tx *Tx, instruction Instruction, innerInstructio
}
func metaoraPoolSwap(tx *Tx, instruction Instruction, innerInstructions InnerInstructions, offset [2]uint) ([]Swap, [2]uint, error) {
if len(instruction.Accounts) < 13 {
return nil, increaseOffset(offset), fmt.Errorf("not enough accounts for swap instruction")
}
swapOffset := offset
var args metaoraPoolSwapArgs
if err := agbinary.NewBorshDecoder(instruction.Data[8:]).Decode(&args); err != nil {

View File

@@ -83,7 +83,13 @@ type MetaoraDammInitializePoolEvent struct {
}
func meteoraDammV2InitializePoolParser(tx *Tx, instruction Instruction, innerInstructions InnerInstructions, offset [2]uint) ([]Swap, [2]uint, error) {
if len(instruction.Accounts) < 12 {
requiredAccounts := 12
if bytes.Equal(instruction.Data[:8], meteoraDammV2InitializePoolWithDynamicConfig[:]) {
requiredAccounts = 13
} else if bytes.Equal(instruction.Data[:8], meteoraDammV2InitializeCustomizablePoolDiscriminator[:]) {
requiredAccounts = 11
}
if len(instruction.Accounts) < requiredAccounts {
return nil, increaseOffset(offset), fmt.Errorf("invalid instruction accounts length")
}
var entryContract = tx.rawTx.accountList[tx.rawTx.Transaction.Message.Instructions[offset[0]].ProgramIDIndex]
@@ -439,7 +445,7 @@ func meteoraDammV2AddLiquidityParser(tx *Tx, instruction Instruction, innerInstr
}
func meteoraDammV2RemoveLiquidityParser(tx *Tx, instruction Instruction, innerInstructions InnerInstructions, offset [2]uint) ([]Swap, [2]uint, error) {
if len(instruction.Accounts) < 8 {
if len(instruction.Accounts) < 9 {
return nil, increaseOffset(offset), fmt.Errorf("invalid instruction accounts length")
}
tokenAMint := tx.rawTx.accountList[instruction.Accounts[7]]

View File

@@ -272,6 +272,8 @@ func (tx *Tx) Parser() error {
quoteMint: swap.QuoteMint,
baseReserve: swap.BaseReserve,
quoteReserve: swap.QuoteReserve,
realQuoteReserve: swap.RealQuoteReserve,
virtualQuoteReserve: swap.VirtualQuoteReserve,
}
}
@@ -281,6 +283,8 @@ func (tx *Tx) Parser() error {
if tx.Swaps[i].BaseMint == v.baseMint && tx.Swaps[i].QuoteMint == v.quoteMint {
tx.Swaps[i].BaseReserve = v.baseReserve
tx.Swaps[i].QuoteReserve = v.quoteReserve
tx.Swaps[i].RealQuoteReserve = v.realQuoteReserve
tx.Swaps[i].VirtualQuoteReserve = v.virtualQuoteReserve
} else if tx.Swaps[i].BaseMint == v.quoteMint && tx.Swaps[i].QuoteMint == v.baseMint {
tx.Swaps[i].BaseReserve = v.quoteReserve
tx.Swaps[i].QuoteReserve = v.baseReserve
@@ -301,6 +305,8 @@ type reserveSnapshot struct {
quoteMint solana.PublicKey
baseReserve decimal.Decimal
quoteReserve decimal.Decimal
realQuoteReserve decimal.Decimal
virtualQuoteReserve decimal.Decimal
}
func cloneSwapPrograms(src map[solana.PublicKey]swapParser) map[solana.PublicKey]swapParser {

View File

@@ -173,9 +173,17 @@ func CreateParser(tx *Tx, instr Instruction, innerInstructions InnerInstructions
}
userIndex := 0
if bytes.HasPrefix(instr.Data, pumpCreateV2Discriminator[:]) {
if len(instr.Accounts) < 6 {
return nil, increaseOffset(offset), InstructionIgnoredError
}
userIndex = instr.Accounts[5]
} else if bytes.HasPrefix(instr.Data, pumpCreateDiscriminator[:]) {
if len(instr.Accounts) < 8 {
return nil, increaseOffset(offset), InstructionIgnoredError
}
userIndex = instr.Accounts[7]
} else {
return nil, increaseOffset(offset), InstructionIgnoredError
}
userBase := getAccountBalanceAfterTx(result, userIndex)
userQuote, _ := GetSolAfterTx(result, userIndex)

View File

@@ -9,7 +9,7 @@ import (
"github.com/shopspring/decimal"
)
type ammBuyEvent struct {
type ammBuyEventPrefix struct {
TimeStamp int64
BaseAmountOut uint64
MaxQuoteAmountIn uint64
@@ -45,6 +45,38 @@ type ammBuyEvent struct {
Cashback uint64
}
type ammBuyEvent struct {
ammBuyEventPrefix
BuybackFeeBasisPoints uint64
BuybackFee uint64
VirtualQuoteReserves agbinary.Int128
CanBoost bool
BaseSupply uint64
}
type ammTradeEventBuybackSuffix struct {
BuybackFeeBasisPoints uint64
BuybackFee uint64
}
type ammTradeEventVirtualSuffix struct {
VirtualQuoteReserves agbinary.Int128
CanBoost bool
BaseSupply uint64
}
func decodeAmmBuyEvent(data []byte) (ammBuyEvent, error) {
var event ammBuyEvent
decoder := agbinary.NewBorshDecoder(data)
if err := decoder.Decode(&event.ammBuyEventPrefix); err != nil {
return ammBuyEvent{}, err
}
if err := decodeAmmTradeEventSuffix(decoder, &event.BuybackFeeBasisPoints, &event.BuybackFee, &event.VirtualQuoteReserves, &event.CanBoost, &event.BaseSupply); err != nil {
return ammBuyEvent{}, err
}
return event, nil
}
type ammCreatePoolEvent struct {
TimeStamp int64
Index uint16
@@ -90,7 +122,7 @@ type ammDepositEvent struct {
UserPoolTokenAccount solana.PublicKey
}
type ammSellEvent struct {
type ammSellEventPrefix struct {
Timestamp int64
BaseAmountIn uint64
MinQuoteAmountOut uint64
@@ -119,6 +151,86 @@ type ammSellEvent struct {
Cashback uint64
}
type ammSellEvent struct {
ammSellEventPrefix
BuybackFeeBasisPoints uint64
BuybackFee uint64
VirtualQuoteReserves agbinary.Int128
CanBoost bool
BaseSupply uint64
}
func decodeAmmSellEvent(data []byte) (ammSellEvent, error) {
var event ammSellEvent
decoder := agbinary.NewBorshDecoder(data)
if err := decoder.Decode(&event.ammSellEventPrefix); err != nil {
return ammSellEvent{}, err
}
if err := decodeAmmTradeEventSuffix(decoder, &event.BuybackFeeBasisPoints, &event.BuybackFee, &event.VirtualQuoteReserves, &event.CanBoost, &event.BaseSupply); err != nil {
return ammSellEvent{}, err
}
return event, nil
}
func decodeAmmTradeEventSuffix(
decoder *agbinary.Decoder,
buybackFeeBasisPoints *uint64,
buybackFee *uint64,
virtualQuoteReserves *agbinary.Int128,
canBoost *bool,
baseSupply *uint64,
) error {
remaining := decoder.Remaining()
if remaining == 0 {
return nil
}
if remaining < 16 {
return fmt.Errorf("pump amm trade event buyback suffix truncated: %d bytes", remaining)
}
var buyback ammTradeEventBuybackSuffix
if err := decoder.Decode(&buyback); err != nil {
return err
}
*buybackFeeBasisPoints = buyback.BuybackFeeBasisPoints
*buybackFee = buyback.BuybackFee
remaining = decoder.Remaining()
if remaining == 0 {
return nil
}
if remaining < 25 {
return fmt.Errorf("pump amm trade event virtual reserve suffix truncated: %d bytes", remaining)
}
var virtual ammTradeEventVirtualSuffix
if err := decoder.Decode(&virtual); err != nil {
return err
}
*virtualQuoteReserves = virtual.VirtualQuoteReserves
*canBoost = virtual.CanBoost
*baseSupply = virtual.BaseSupply
return nil
}
func pumpAmmVirtualQuoteReserve(value agbinary.Int128) decimal.Decimal {
return decimal.NewFromBigInt(value.BigInt(), 0)
}
func pumpAmmBuyPostQuoteReserves(event ammBuyEvent) (real, effective decimal.Decimal) {
// The event reserve is the pre-trade real vault balance. Keep the
// parser's existing LP-fee treatment by applying QuoteAmountIn here.
real = decimal.NewFromUint64(event.PoolQuoteTokenReserve).Add(decimal.NewFromUint64(event.QuoteAmountIn))
return real, real.Add(pumpAmmVirtualQuoteReserve(event.VirtualQuoteReserves))
}
func pumpAmmSellPostQuoteReserves(event ammSellEvent) (real, effective decimal.Decimal) {
// The event reserve is the pre-trade real vault balance. Keep the
// parser's existing LP-fee treatment by applying QuoteAmountOut here.
real = decimal.NewFromUint64(event.PoolQuoteTokenReserves).Sub(decimal.NewFromUint64(event.QuoteAmountOut))
return real, real.Add(pumpAmmVirtualQuoteReserve(event.VirtualQuoteReserves))
}
type ammWithdrawEvent struct {
Timestamp int64
LpTokenAmountIn uint64
@@ -182,6 +294,9 @@ func pumpAmmParser(tx *Tx, instruction Instruction, innerInstructions InnerInstr
}
func ammCreatePoolParser(tx *Tx, instruction Instruction, innerInstructions InnerInstructions, offset [2]uint) ([]Swap, [2]uint, error) {
if len(instruction.Accounts) < 15 {
return nil, increaseOffset(offset), InstructionIgnoredError
}
result := tx.rawTx
var entryContract = result.accountList[result.Transaction.Message.Instructions[offset[0]].ProgramIDIndex]
var err error
@@ -246,6 +361,7 @@ func ammCreatePoolParser(tx *Tx, instruction Instruction, innerInstructions Inne
QuoteAmount: decimal.NewFromUint64(createEvent.QuoteAmountIn),
BaseReserve: decimal.NewFromUint64(createEvent.PoolBaseAmount),
QuoteReserve: decimal.NewFromUint64(createEvent.PoolQuoteAmount),
RealQuoteReserve: decimal.NewFromUint64(createEvent.PoolQuoteAmount),
UserBaseBalance: decimal.Decimal{},
UserQuoteBalance: decimal.Decimal{},
EntryContract: entryContract,
@@ -275,6 +391,9 @@ func pumpAmmSwapAmountInfoFromArgs(args PumpSwapArgs) (swapMode SwapMode, fixedA
}
func failedTxAmmBuyParser(tx *Tx, instruction Instruction, innerInstructions InnerInstructions, offset [2]uint) ([]Swap, [2]uint, error) {
if len(instruction.Accounts) < 13 {
return nil, increaseOffset(offset), InstructionIgnoredError
}
if tx.Err == nil || tx.Err.UnKnown != "" {
return nil, increaseOffset(offset), fmt.Errorf("tx pump amm sell failed but error is nil, offset, %d, %d", offset[0], offset[1])
}
@@ -389,6 +508,7 @@ func failedTxAmmBuyParser(tx *Tx, instruction Instruction, innerInstructions Inn
QuoteAmount: decimal.NewFromUint64(quoteAmount),
BaseReserve: baseReserve,
QuoteReserve: quoteReserve,
RealQuoteReserve: quoteReserve,
Mayhem: isMayhemPump(result.accountList[instruction.Accounts[9]]),
UserBaseBalance: userBase,
UserQuoteBalance: userQuote,
@@ -401,6 +521,9 @@ func failedTxAmmBuyParser(tx *Tx, instruction Instruction, innerInstructions Inn
}
func failedTxAmmSellParser(tx *Tx, instruction Instruction, innerInstructions InnerInstructions, offset [2]uint) ([]Swap, [2]uint, error) {
if len(instruction.Accounts) < 13 {
return nil, increaseOffset(offset), InstructionIgnoredError
}
if tx.Err == nil || tx.Err.UnKnown != "" {
return nil, increaseOffset(offset), fmt.Errorf("tx pump amm sell failed but error is nil, offset, %d, %d", offset[0], offset[1])
}
@@ -509,6 +632,7 @@ func failedTxAmmSellParser(tx *Tx, instruction Instruction, innerInstructions In
QuoteAmount: decimal.NewFromUint64(quoteAmount),
BaseReserve: baseReserve,
QuoteReserve: quoteReserve,
RealQuoteReserve: quoteReserve,
Mayhem: isMayhemPump(result.accountList[instruction.Accounts[9]]),
UserBaseBalance: userBase,
UserQuoteBalance: userQuote,
@@ -525,6 +649,9 @@ func ammBuyParser(tx *Tx, instruction Instruction, innerInstructions InnerInstru
var entryContract = result.accountList[result.Transaction.Message.Instructions[offset[0]].ProgramIDIndex]
var err error
var prefixLen = offset[1]
if len(instruction.Accounts) < 13 {
return nil, increaseOffset(offset), InstructionIgnoredError
}
inners, err := getInnerInstructions(innerInstructions, prefixLen)
if err != nil {
return nil, increaseOffset(offset), fmt.Errorf("pumpamm create get inner instructions error: %v, offset: %d, %d", err, offset[0], prefixLen)
@@ -545,7 +672,7 @@ func ammBuyParser(tx *Tx, instruction Instruction, innerInstructions InnerInstru
if innerInstr.ProgramIDIndex == instruction.ProgramIDIndex &&
bytes.Equal(innerInstr.Data[:8], pumpAmmEventDiscriminator[:]) &&
bytes.Equal(innerInstr.Data[8:16], pumpAmmBuyEventDiscriminator[:]) {
err = agbinary.NewBorshDecoder(innerInstr.Data[16:]).Decode(&event)
event, err = decodeAmmBuyEvent(innerInstr.Data[16:])
if offset[1] == 0 {
offset[0] += 1
} else {
@@ -620,6 +747,7 @@ func ammBuyParser(tx *Tx, instruction Instruction, innerInstructions InnerInstru
if event.IxName == "buy" {
quoteAmount = decimal.NewFromUint64(event.QuoteAmountIn)
}
realQuoteReserve, effectiveQuoteReserve := pumpAmmBuyPostQuoteReserves(event)
swap := Swap{
Program: SolProgramPumpAMM,
Event: "buy",
@@ -635,7 +763,9 @@ func ammBuyParser(tx *Tx, instruction Instruction, innerInstructions InnerInstru
BaseAmount: decimal.NewFromUint64(event.BaseAmountOut),
QuoteAmount: quoteAmount,
BaseReserve: decimal.NewFromUint64(event.PoolBaseTokenReserve - event.BaseAmountOut),
QuoteReserve: decimal.NewFromUint64(event.PoolQuoteTokenReserve + event.QuoteAmountIn),
QuoteReserve: effectiveQuoteReserve,
RealQuoteReserve: realQuoteReserve,
VirtualQuoteReserve: pumpAmmVirtualQuoteReserve(event.VirtualQuoteReserves),
Mayhem: isMayhemPump(result.accountList[instruction.Accounts[9]]),
Cashback: isCashbackCoin,
UserBaseBalance: userBase,
@@ -662,6 +792,9 @@ func ammSellParser(tx *Tx, instruction Instruction, innerInstructions InnerInstr
result := tx.rawTx
var entryContract = result.accountList[result.Transaction.Message.Instructions[offset[0]].ProgramIDIndex]
var err error
if len(instruction.Accounts) < 13 {
return nil, increaseOffset(offset), InstructionIgnoredError
}
var prefixLen = offset[1]
inners, err := getInnerInstructions(innerInstructions, prefixLen)
if err != nil {
@@ -684,7 +817,7 @@ func ammSellParser(tx *Tx, instruction Instruction, innerInstructions InnerInstr
if innerInstr.ProgramIDIndex == instruction.ProgramIDIndex &&
bytes.Equal(innerInstr.Data[:8], pumpAmmEventDiscriminator[:]) &&
bytes.Equal(innerInstr.Data[8:16], pumpAmmSellEventDiscriminator[:]) {
err = agbinary.NewBorshDecoder(innerInstr.Data[16:]).Decode(&event)
event, err = decodeAmmSellEvent(innerInstr.Data[16:])
if offset[1] == 0 {
offset[0] += 1
} else {
@@ -755,6 +888,7 @@ func ammSellParser(tx *Tx, instruction Instruction, innerInstructions InnerInstr
userQuote = userQuote.Add(decimal.NewFromUint64(userBalance))
}
isCashbackCoin := event.CashbackFeeBasisPoints > 0 || event.Cashback > 0
realQuoteReserve, effectiveQuoteReserve := pumpAmmSellPostQuoteReserves(event)
swap := Swap{
Program: SolProgramPumpAMM,
Event: "sell",
@@ -770,7 +904,9 @@ func ammSellParser(tx *Tx, instruction Instruction, innerInstructions InnerInstr
BaseAmount: decimal.NewFromUint64(event.BaseAmountIn),
QuoteAmount: decimal.NewFromUint64(event.UserQuoteAmountOut),
BaseReserve: decimal.NewFromUint64(event.PoolBaseTokenReserves + event.BaseAmountIn),
QuoteReserve: decimal.NewFromUint64(event.PoolQuoteTokenReserves - event.QuoteAmountOut),
QuoteReserve: effectiveQuoteReserve,
RealQuoteReserve: realQuoteReserve,
VirtualQuoteReserve: pumpAmmVirtualQuoteReserve(event.VirtualQuoteReserves),
Mayhem: isMayhemPump(result.accountList[instruction.Accounts[9]]),
Cashback: isCashbackCoin,
UserBaseBalance: userBase,
@@ -789,6 +925,9 @@ func depositParse(tx *Tx, instruction Instruction, innerInstructions InnerInstru
result := tx.rawTx
var entryContract = result.accountList[result.Transaction.Message.Instructions[offset[0]].ProgramIDIndex]
var err error
if len(instruction.Accounts) < 11 {
return nil, increaseOffset(offset), InstructionIgnoredError
}
var prefixLen = offset[1]
inners, err := getInnerInstructions(innerInstructions, prefixLen)
if err != nil {
@@ -875,6 +1014,7 @@ func depositParse(tx *Tx, instruction Instruction, innerInstructions InnerInstru
QuoteAmount: decimal.NewFromUint64(event.QuoteAmountIn),
BaseReserve: decimal.NewFromUint64(event.PoolBaseTokenReserves + event.BaseAmountIn),
QuoteReserve: decimal.NewFromUint64(event.PoolQuoteTokenReserves + event.QuoteAmountIn),
RealQuoteReserve: decimal.NewFromUint64(event.PoolQuoteTokenReserves + event.QuoteAmountIn),
//Mayhem: false,
UserBaseBalance: decimal.NewFromUint64(event.UserBaseTokenReserves - event.BaseAmountIn),
UserQuoteBalance: decimal.NewFromUint64(event.UserQuoteTokenReserves - event.QuoteAmountIn),
@@ -887,6 +1027,9 @@ func withdrawParse(tx *Tx, instruction Instruction, innerInstructions InnerInstr
result := tx.rawTx
var entryContract = result.accountList[result.Transaction.Message.Instructions[offset[0]].ProgramIDIndex]
var err error
if len(instruction.Accounts) < 11 {
return nil, increaseOffset(offset), InstructionIgnoredError
}
var prefixLen = offset[1]
inners, err := getInnerInstructions(innerInstructions, prefixLen)
if err != nil {
@@ -973,6 +1116,7 @@ func withdrawParse(tx *Tx, instruction Instruction, innerInstructions InnerInstr
QuoteAmount: decimal.NewFromUint64(event.QuoteAmountOut),
BaseReserve: decimal.NewFromUint64(event.PoolBaseTokenReserves - event.BaseAmountOut),
QuoteReserve: decimal.NewFromUint64(event.PoolQuoteTokenReserves - event.QuoteAmountOut),
RealQuoteReserve: decimal.NewFromUint64(event.PoolQuoteTokenReserves - event.QuoteAmountOut),
//Mayhem: false,
UserBaseBalance: decimal.NewFromUint64(event.UserBaseTokenReserves + event.BaseAmountOut),
UserQuoteBalance: decimal.NewFromUint64(event.UserQuoteTokenReserves + event.QuoteAmountOut),

View File

@@ -1,11 +1,13 @@
package pump_parser
import (
"bytes"
"encoding/base64"
"fmt"
agbinary "github.com/gagliardetto/binary"
"github.com/gagliardetto/solana-go"
"github.com/mr-tron/base58"
"testing"
)
@@ -40,3 +42,128 @@ func TestAmmBuyEvent(t *testing.T) {
fmt.Println(pumpAmmBuyEventDiscriminator)
fmt.Println(pumpGetFeesDiscriminator)
}
func TestDecodeAmmBuyEventSuffixCompatibility(t *testing.T) {
prefix := ammBuyEventPrefix{TimeStamp: 123, PoolQuoteTokenReserve: 456}
legacy, err := decodeAmmBuyEvent(encodeAmmEventParts(t, prefix))
if err != nil {
t.Fatalf("decodeAmmBuyEvent(legacy) error = %v", err)
}
if legacy.TimeStamp != prefix.TimeStamp || legacy.PoolQuoteTokenReserve != prefix.PoolQuoteTokenReserve {
t.Fatalf("legacy prefix mismatch: %+v", legacy.ammBuyEventPrefix)
}
if legacy.BuybackFeeBasisPoints != 0 || legacy.BuybackFee != 0 || legacy.CanBoost || legacy.BaseSupply != 0 || !pumpAmmVirtualQuoteReserve(legacy.VirtualQuoteReserves).IsZero() {
t.Fatalf("legacy suffix is not zero-valued: %+v", legacy)
}
buyback := ammTradeEventBuybackSuffix{BuybackFeeBasisPoints: 1000, BuybackFee: 77}
withBuyback, err := decodeAmmBuyEvent(encodeAmmEventParts(t, prefix, buyback))
if err != nil {
t.Fatalf("decodeAmmBuyEvent(buyback) error = %v", err)
}
if withBuyback.BuybackFeeBasisPoints != buyback.BuybackFeeBasisPoints || withBuyback.BuybackFee != buyback.BuybackFee {
t.Fatalf("buyback suffix mismatch: %+v", withBuyback)
}
if withBuyback.CanBoost || withBuyback.BaseSupply != 0 || !pumpAmmVirtualQuoteReserve(withBuyback.VirtualQuoteReserves).IsZero() {
t.Fatalf("pre-virtual suffix has unexpected virtual fields: %+v", withBuyback)
}
negativeVirtual := agbinary.Int128(agbinary.Uint128{Lo: ^uint64(122), Hi: ^uint64(0)})
virtual := ammTradeEventVirtualSuffix{
VirtualQuoteReserves: negativeVirtual,
CanBoost: true,
BaseSupply: 1_000_000,
}
current, err := decodeAmmBuyEvent(encodeAmmEventParts(t, prefix, buyback, virtual))
if err != nil {
t.Fatalf("decodeAmmBuyEvent(current) error = %v", err)
}
if got := pumpAmmVirtualQuoteReserve(current.VirtualQuoteReserves).String(); got != "-123" {
t.Fatalf("VirtualQuoteReserves = %s, want -123", got)
}
if !current.CanBoost || current.BaseSupply != virtual.BaseSupply {
t.Fatalf("current suffix mismatch: %+v", current)
}
}
func TestDecodeAmmSellEventSuffixCompatibility(t *testing.T) {
prefix := ammSellEventPrefix{Timestamp: 321, PoolQuoteTokenReserves: 654, QuoteAmountOut: 54}
buyback := ammTradeEventBuybackSuffix{BuybackFeeBasisPoints: 900, BuybackFee: 88}
virtual := ammTradeEventVirtualSuffix{
VirtualQuoteReserves: agbinary.Int128(agbinary.Uint128{Lo: 987}),
CanBoost: true,
BaseSupply: 2_000_000,
}
event, err := decodeAmmSellEvent(encodeAmmEventParts(t, prefix, buyback, virtual))
if err != nil {
t.Fatalf("decodeAmmSellEvent() error = %v", err)
}
if event.Timestamp != prefix.Timestamp || event.BuybackFee != buyback.BuybackFee || !event.CanBoost || event.BaseSupply != virtual.BaseSupply {
t.Fatalf("decoded sell event mismatch: %+v", event)
}
if got := pumpAmmVirtualQuoteReserve(event.VirtualQuoteReserves).String(); got != "987" {
t.Fatalf("VirtualQuoteReserves = %s, want 987", got)
}
real, effective := pumpAmmSellPostQuoteReserves(event)
if real.String() != "600" || effective.String() != "1587" {
t.Fatalf("post reserves = real %s effective %s", real, effective)
}
}
func TestDecodeAmmTradeEventRejectsTruncatedSuffix(t *testing.T) {
prefix := encodeAmmEventParts(t, ammBuyEventPrefix{})
for _, suffixLength := range []int{1, 15, 17, 40} {
data := append(append([]byte(nil), prefix...), make([]byte, suffixLength)...)
if _, err := decodeAmmBuyEvent(data); err == nil {
t.Fatalf("decodeAmmBuyEvent() error = nil for %d-byte suffix", suffixLength)
}
}
}
func TestDecodeAmmBuyEventDevnetVirtualReserve(t *testing.T) {
const eventDataBase58 = "6MF1ykxMQW5eFmo1ErWgds8V2wfEyqJPjEsoz4yRoKQCktPze5U297neF7huAQQ9DhxEJqHPwezxUh9mcQEsfdmnwW84YVxcfEzDUduRqc1zty2dLQZcFGzFeFC6fmpX75tkvuMvTN89HekHa9TEikGRuXhtWkeacjNoVqfbxXUp3FQxGSFoyeBaSdAfPMvXJ7zN4ALTDtn3a5g7HxJCPPWeUtYGcJig6RpKrC3DzAQ3x5SiiJiVhK1zdocmkkRXuTezt5L5Svrf9pbM4QkkwgvxrYqfUs5YeaLNqWYtNZfimLtENhQJcfrhQh2svCosCJ2Ux5iyyLZwJu7VLZibD41CwxuF6ngmrhde8BmEegmQ9f3L4qYhqgW5RJSiu1JnKMFqmTckuwGydQfTiyA9RsCtFxiazaAVim9SMARU2fqe9zfAPdRjBetAvTGkBxrERmmzhWVhER6g448mV1E1eeQjFWYRhtR33hSUpNUUP6dUF8CWae2BuJxaKexpiY3pyqUwbH7osvexbff8jpfpZGubP5UWjqhwijtF8tdkoEdAnb2TgWXqXMu4by5VDsEzwbgNK1DbkGbxrcXYwFTNskqoUqoNLPRaAcfnjTMA9cVme62HAj5WvxuzpL9Gpw7ZzSE2LYnVvzTa6PnN2bbpuqJpZ3d"
data, err := base58.Decode(eventDataBase58)
if err != nil {
t.Fatalf("base58.Decode() error = %v", err)
}
if len(data) < 16 || !bytes.Equal(data[:8], pumpAmmEventDiscriminator[:]) || !bytes.Equal(data[8:16], pumpAmmBuyEventDiscriminator[:]) {
t.Fatalf("unexpected event discriminator: len=%d", len(data))
}
event, err := decodeAmmBuyEvent(data[16:])
if err != nil {
t.Fatalf("decodeAmmBuyEvent(devnet) error = %v", err)
}
if event.TimeStamp != 1784271481 || event.PoolQuoteTokenReserve != 2_213_159_021 || event.QuoteAmountIn != 100_000_000 || event.QuoteAmountInWithLpFee != 98_785_184 {
t.Fatalf("devnet buy event core fields mismatch: %+v", event)
}
if event.BuybackFeeBasisPoints != 1000 || event.BuybackFee != 91_851 {
t.Fatalf("devnet buyback fields mismatch: bps=%d fee=%d", event.BuybackFeeBasisPoints, event.BuybackFee)
}
if got := pumpAmmVirtualQuoteReserve(event.VirtualQuoteReserves).String(); got != "583150126" {
t.Fatalf("VirtualQuoteReserves = %s, want 583150126", got)
}
if !event.CanBoost || event.BaseSupply != 1_000_000_000_000_000 {
t.Fatalf("devnet boost metadata mismatch: can_boost=%t base_supply=%d", event.CanBoost, event.BaseSupply)
}
real, effective := pumpAmmBuyPostQuoteReserves(event)
if real.String() != "2313159021" || effective.String() != "2896309147" {
t.Fatalf("post reserves = real %s effective %s", real, effective)
}
}
func encodeAmmEventParts(t *testing.T, parts ...any) []byte {
t.Helper()
var buf bytes.Buffer
encoder := agbinary.NewBorshEncoder(&buf)
for _, part := range parts {
if err := encoder.Encode(part); err != nil {
t.Fatalf("borsh encode %T: %v", part, err)
}
}
return buf.Bytes()
}

View File

@@ -108,38 +108,14 @@ func raydiumClmmAddLiquidityParser(tx *Tx, instruction Instruction, innerInstruc
switch discriminator {
case raydiumClmmIncreaseLiquidityDiscriminator:
accountMin = 12
market = tx.rawTx.accountList[instruction.Accounts[2]]
vault0 = instruction.Accounts[9]
vault1 = instruction.Accounts[10]
case raydiumClmmIncreaseLiquidityV2Discriminator:
accountMin = 15
market = tx.rawTx.accountList[instruction.Accounts[2]]
vault0 = instruction.Accounts[9]
vault1 = instruction.Accounts[10]
//token0 = tx.rawTx.accountList[instruction.Accounts[13]]
//token1 = tx.rawTx.accountList[instruction.Accounts[14]]
case raydiumClmmOpenPositionDiscriminator:
accountMin = 19
market = tx.rawTx.accountList[instruction.Accounts[5]]
vault0 = instruction.Accounts[12]
vault1 = instruction.Accounts[13]
lpToken = tx.rawTx.accountList[instruction.Accounts[2]]
case raydiumClmmOpenPositionV2Discriminator:
accountMin = 22
market = tx.rawTx.accountList[instruction.Accounts[5]]
vault0 = instruction.Accounts[12]
vault1 = instruction.Accounts[13]
lpToken = tx.rawTx.accountList[instruction.Accounts[2]]
//token0 = tx.rawTx.accountList[instruction.Accounts[20]]
//token1 = tx.rawTx.accountList[instruction.Accounts[21]]
case raydiumClmmOpenPositionWithToken22NftDiscriminator:
accountMin = 20
market = tx.rawTx.accountList[instruction.Accounts[4]]
vault0 = instruction.Accounts[11]
vault1 = instruction.Accounts[12]
lpToken = tx.rawTx.accountList[instruction.Accounts[2]]
//token0 = tx.rawTx.accountList[instruction.Accounts[18]]
//token1 = tx.rawTx.accountList[instruction.Accounts[19]]
default:
return nil, increaseOffset(offset), fmt.Errorf("invalid discriminator")
}
@@ -148,6 +124,23 @@ func raydiumClmmAddLiquidityParser(tx *Tx, instruction Instruction, innerInstruc
return nil, increaseOffset(offset), fmt.Errorf("not enough accounts for raydiumClmm add liquidity instruction, offset, %d, %d", offset[0], offset[1])
}
switch discriminator {
case raydiumClmmIncreaseLiquidityDiscriminator, raydiumClmmIncreaseLiquidityV2Discriminator:
market = tx.rawTx.accountList[instruction.Accounts[2]]
vault0 = instruction.Accounts[9]
vault1 = instruction.Accounts[10]
case raydiumClmmOpenPositionDiscriminator, raydiumClmmOpenPositionV2Discriminator:
market = tx.rawTx.accountList[instruction.Accounts[5]]
vault0 = instruction.Accounts[12]
vault1 = instruction.Accounts[13]
lpToken = tx.rawTx.accountList[instruction.Accounts[2]]
case raydiumClmmOpenPositionWithToken22NftDiscriminator:
market = tx.rawTx.accountList[instruction.Accounts[4]]
vault0 = instruction.Accounts[11]
vault1 = instruction.Accounts[12]
lpToken = tx.rawTx.accountList[instruction.Accounts[2]]
}
baseTokenBalance, err := getTokenBalanceAfterTx(tx.rawTx, vault0)
if err != nil {
return nil, increaseOffset(offset), fmt.Errorf("failed to get token0 vault balance after tx: %v", err)
@@ -197,6 +190,8 @@ func raydiumClmmDecreaseLiquidityParser(tx *Tx, instruction Instruction, innerIn
accountMin = 14
} else if discriminator == raydiumClmmDecreaseLiquidityV2Discriminator {
accountMin = 16
} else {
return nil, increaseOffset(offset), fmt.Errorf("invalid discriminator")
}
if len(instruction.Accounts) < accountMin {
return nil, increaseOffset(offset), fmt.Errorf("not enough accounts for decrease liquidity instruction")
@@ -299,23 +294,21 @@ func raydiumClmmSwapParser(tx *Tx, instruction Instruction, innerInstructions In
}
if discriminator == raydiumClmmSwapDiscriminator {
accountMin = 9
pool = tx.rawTx.accountList[instruction.Accounts[2]]
userTokenInAccount = instruction.Accounts[3]
userTokenOutAccount = instruction.Accounts[4]
tokenInVault = instruction.Accounts[5]
tokenOutVault = instruction.Accounts[6]
} else if discriminator == raydiumClmmSwapV2Discriminator {
accountMin = 13
pool = tx.rawTx.accountList[instruction.Accounts[2]]
userTokenInAccount = instruction.Accounts[3]
userTokenOutAccount = instruction.Accounts[4]
tokenInVault = instruction.Accounts[5]
tokenOutVault = instruction.Accounts[6]
} else {
return nil, increaseOffset(offset), fmt.Errorf("invalid discriminator")
}
if len(instruction.Accounts) < accountMin {
return nil, increaseOffset(offset), fmt.Errorf("not enough accounts for swap instruction")
}
pool = tx.rawTx.accountList[instruction.Accounts[2]]
userTokenInAccount = instruction.Accounts[3]
userTokenOutAccount = instruction.Accounts[4]
tokenInVault = instruction.Accounts[5]
tokenOutVault = instruction.Accounts[6]
baseTokenBalance, err := getTokenBalanceAfterTx(tx.rawTx, tokenInVault)
if err != nil {
return nil, increaseOffset(offset), fmt.Errorf("failed to get tokenIn vault balance after tx: %w", err)

View File

@@ -373,6 +373,9 @@ type raydiumLaunchLabSwapArgs struct {
}
func raydiumLaunchLabSwapParser(tx *Tx, instruction Instruction, innerInstructions InnerInstructions, offset [2]uint) ([]Swap, [2]uint, error) {
if len(instruction.Accounts) < 13 {
return nil, increaseOffset(offset), InstructionIgnoredError
}
platformConfig := tx.rawTx.accountList[instruction.Accounts[3]]
var programName string
if platformConfig.Equals(bonkPlatformConfig) {

View File

@@ -15,6 +15,9 @@ func systemParser(tx *Tx, instruction Instruction, _ InnerInstructions, offset [
}
decode := instruction.Data
if len(decode) < 4 {
return increaseOffset(offset), nil
}
discriminator := binary.LittleEndian.Uint32(decode[0:4])
switch discriminator {
@@ -29,6 +32,9 @@ func TransferParser(result *RawTx, instruction Instruction, offset [2]uint, tx *
if len(decodeData) < 8 {
return increaseOffset(offset), nil
}
if len(instruction.Accounts) < 2 || len(result.Transaction.Message.Instructions[offset[0]].Accounts) < 1 {
return increaseOffset(offset), InstructionIgnoredError
}
var lamports uint64 = binary.LittleEndian.Uint64(decodeData)
from := result.accountList[result.Transaction.Message.Instructions[offset[0]].Accounts[0]]

2
tx.go
View File

@@ -45,6 +45,8 @@ type Swap struct {
BaseReserve decimal.Decimal
QuoteReserve decimal.Decimal
RealQuoteReserve decimal.Decimal
VirtualQuoteReserve decimal.Decimal
Mayhem bool
Cashback bool

View File

@@ -8,6 +8,7 @@ import (
"io"
"iter"
"math"
"math/big"
"sort"
"strconv"
@@ -17,7 +18,10 @@ import (
)
const (
txBinarySchemaVersionCurrent uint16 = 3
txBinarySchemaVersionV3 uint16 = 3
txBinarySchemaVersionV4 uint16 = 4
txBinarySchemaVersionV5 uint16 = 5
txBinarySchemaVersionCurrent = txBinarySchemaVersionV5
txBinaryEnumVersionV1 uint16 = 1
txBinarySOLScale int32 = 9
@@ -27,6 +31,10 @@ const (
var txBinaryMagic = [4]byte{'P', 'T', 'X', 'B'}
var txsBinaryMagic = [4]byte{'P', 'T', 'X', 'S'}
func txBinarySchemaVersionSupported(version uint16) bool {
return version >= txBinarySchemaVersionV3 && version <= txBinarySchemaVersionCurrent
}
type TxBinary struct {
SchemaVersion uint16
EnumVersion uint16
@@ -34,6 +42,7 @@ type TxBinary struct {
Signer uint32
Block uint64
BlockIndex uint64
BlockAt int64
TxHash *[64]byte
CuFee uint64
Swaps []SwapBinary
@@ -96,6 +105,9 @@ type SwapBinary struct {
LpMint uint32
AfterSOLBalance float64
RealQuoteReserve uint64
VirtualQuoteReserve [16]byte
}
type TxsBinary struct {
@@ -204,6 +216,7 @@ func newTxBinaryWithAddressTable(tx *Tx, addressTable []solana.PublicKey, addres
EnumVersion: txBinaryEnumVersionV1,
Block: tx.Block,
BlockIndex: tx.BlockIndex,
BlockAt: tx.BlockAt,
CuLimit: tx.CuLimit,
ComputeUnitsConsumed: tx.ComputeUnitsConsumed,
}
@@ -411,6 +424,8 @@ func MergeTxsBinarySourcesToWriterWithOptions(sources []TxsBinaryReaderSource, w
return fmt.Errorf("source[%d].batch[%d].tx[%d]: %w", sourceIndex, batchIndex, txIndex, err)
}
tx.SchemaVersion = plan.schemaVersion
tx.EnumVersion = plan.enumVersion
bodyBytes, err := txBinaryMarshalTxBody(&tx, plan.enumTable)
if err != nil {
reader.Close()
@@ -432,7 +447,7 @@ func (tx *TxBinary) MarshalBinary() ([]byte, error) {
if tx == nil {
return nil, fmt.Errorf("tx binary is nil")
}
if tx.SchemaVersion != txBinarySchemaVersionCurrent {
if !txBinarySchemaVersionSupported(tx.SchemaVersion) {
return nil, fmt.Errorf("unsupported tx binary schema version: %d", tx.SchemaVersion)
}
@@ -460,7 +475,7 @@ func (txs *TxsBinary) MarshalBinary() ([]byte, error) {
if txs == nil {
return nil, fmt.Errorf("txs binary is nil")
}
if txs.SchemaVersion != txBinarySchemaVersionCurrent {
if !txBinarySchemaVersionSupported(txs.SchemaVersion) {
return nil, fmt.Errorf("unsupported tx binary schema version: %d", txs.SchemaVersion)
}
@@ -478,8 +493,11 @@ func (txs *TxsBinary) MarshalBinary() ([]byte, error) {
}
enc.writeUint32(uint32(len(txs.Txs)))
for i := range txs.Txs {
if err := enc.writeTxBinaryBody(&txs.Txs[i], enumTable); err != nil {
return nil, fmt.Errorf("tx[%d], %s: %w", i, base58.Encode(txs.Txs[i].TxHash[:]), err)
tx := txs.Txs[i]
tx.SchemaVersion = txs.SchemaVersion
tx.EnumVersion = txs.EnumVersion
if err := enc.writeTxBinaryBody(&tx, enumTable); err != nil {
return nil, fmt.Errorf("tx[%d], %s: %w", i, base58.Encode(tx.TxHash[:]), err)
}
}
return enc.bytes(), nil
@@ -520,7 +538,7 @@ func (tx *TxBinary) UnmarshalBinary(data []byte) error {
if err != nil {
return err
}
if tx.SchemaVersion != txBinarySchemaVersionCurrent {
if !txBinarySchemaVersionSupported(tx.SchemaVersion) {
return fmt.Errorf("unsupported tx binary schema version: %d", tx.SchemaVersion)
}
@@ -560,7 +578,7 @@ func (txs *TxsBinary) UnmarshalBinary(data []byte) error {
if err != nil {
return err
}
if txs.SchemaVersion != txBinarySchemaVersionCurrent {
if !txBinarySchemaVersionSupported(txs.SchemaVersion) {
return fmt.Errorf("unsupported tx binary schema version: %d", txs.SchemaVersion)
}
@@ -613,6 +631,7 @@ func (tx *TxBinary) ToTx() (*Tx, error) {
Signer: signer,
Block: tx.Block,
BlockIndex: tx.BlockIndex,
BlockAt: tx.BlockAt,
CuFee: decimal.NewFromUint64(tx.CuFee),
CUPrice: decimal.NewFromUint64(tx.CUPrice).Shift(-txBinaryCUPriceScale),
BeforeSolBalance: txBinaryFloat64ToDecimal(tx.BeforeSolBalance, txBinarySOLScale),
@@ -646,7 +665,7 @@ func (tx *TxBinary) ToTx() (*Tx, error) {
if len(tx.Swaps) > 0 {
out.Swaps = make([]Swap, 0, len(tx.Swaps))
for i, swap := range tx.Swaps {
decodedSwap, err := swap.toSwap(tx.AddressTable, i)
decodedSwap, err := swap.toSwap(tx.AddressTable, i, tx.SchemaVersion)
if err != nil {
return nil, err
}
@@ -784,6 +803,12 @@ func newSwapBinary(swap Swap, index int, addressIndex *txBinaryAddressIndex) (Sw
if out.QuoteReserve, err = txBinaryDecimalToFloat64Raw(swap.QuoteReserve, fmt.Sprintf("swap[%d].quote_reserve", index)); err != nil {
return SwapBinary{}, err
}
if out.RealQuoteReserve, err = txBinaryDecimalToUint64(swap.RealQuoteReserve, fmt.Sprintf("swap[%d].real_quote_reserve", index)); err != nil {
return SwapBinary{}, err
}
if out.VirtualQuoteReserve, err = txBinaryDecimalToInt128(swap.VirtualQuoteReserve, fmt.Sprintf("swap[%d].virtual_quote_reserve", index)); err != nil {
return SwapBinary{}, err
}
if out.UserBaseBalance, err = txBinaryDecimalToUint64(swap.UserBaseBalance, fmt.Sprintf("swap[%d].user_base_balance", index)); err != nil {
return SwapBinary{}, err
}
@@ -797,7 +822,7 @@ func newSwapBinary(swap Swap, index int, addressIndex *txBinaryAddressIndex) (Sw
return out, nil
}
func (swap SwapBinary) toSwap(addressTable []solana.PublicKey, index int) (Swap, error) {
func (swap SwapBinary) toSwap(addressTable []solana.PublicKey, index int, schemaVersion uint16) (Swap, error) {
pool, err := txBinaryAddressAt(addressTable, swap.Pool, fmt.Sprintf("swap[%d].pool", index))
if err != nil {
return Swap{}, err
@@ -851,6 +876,20 @@ func (swap SwapBinary) toSwap(addressTable []solana.PublicKey, index int) (Swap,
return Swap{}, err
}
quoteReserve := txBinaryFloat64ToDecimalRaw(swap.QuoteReserve)
realQuoteReserve := quoteReserve
virtualQuoteReserve := decimal.Zero
if schemaVersion >= txBinarySchemaVersionV5 {
realQuoteReserve = decimal.NewFromUint64(swap.RealQuoteReserve)
virtualQuoteReserve = txBinaryInt128ToDecimal(swap.VirtualQuoteReserve)
// A v3/v4 batch upgraded by the streaming merger has no component
// bytes to carry forward. Preserve its legacy meaning after the
// merged output is written as v5.
if realQuoteReserve.IsZero() && virtualQuoteReserve.IsZero() && !quoteReserve.IsZero() {
realQuoteReserve = quoteReserve
}
}
return Swap{
Program: swap.Program,
Event: swap.Event,
@@ -880,7 +919,9 @@ func (swap SwapBinary) toSwap(addressTable []solana.PublicKey, index int) (Swap,
ActualLimitAmountSide: swap.ActualLimitAmountSide,
SlippageBps: decimal.NewFromUint64(swap.SlippageBps),
BaseReserve: txBinaryFloat64ToDecimalRaw(swap.BaseReserve),
QuoteReserve: txBinaryFloat64ToDecimalRaw(swap.QuoteReserve),
QuoteReserve: quoteReserve,
RealQuoteReserve: realQuoteReserve,
VirtualQuoteReserve: virtualQuoteReserve,
Mayhem: swap.Mayhem,
Cashback: swap.Cashback,
UserBaseBalance: decimal.NewFromUint64(swap.UserBaseBalance),
@@ -1058,6 +1099,44 @@ func txBinaryDecimalToUint64(value decimal.Decimal, field string) (uint64, error
return bigInt.Uint64(), nil
}
func txBinaryDecimalToInt128(value decimal.Decimal, field string) ([16]byte, error) {
var out [16]byte
if !value.Equal(value.Truncate(0)) {
return out, fmt.Errorf("%s must be an integer, got %s", field, value.String())
}
integer := value.BigInt()
limit := new(big.Int).Lsh(big.NewInt(1), 127)
minimum := new(big.Int).Neg(new(big.Int).Set(limit))
maximum := new(big.Int).Sub(new(big.Int).Set(limit), big.NewInt(1))
if integer.Cmp(minimum) < 0 || integer.Cmp(maximum) > 0 {
return out, fmt.Errorf("%s overflows int128: %s", field, value.String())
}
unsigned := new(big.Int).Set(integer)
if unsigned.Sign() < 0 {
unsigned.Add(unsigned, new(big.Int).Lsh(big.NewInt(1), 128))
}
var bigEndian [16]byte
unsigned.FillBytes(bigEndian[:])
for i := range out {
out[i] = bigEndian[len(bigEndian)-1-i]
}
return out, nil
}
func txBinaryInt128ToDecimal(raw [16]byte) decimal.Decimal {
var bigEndian [16]byte
for i := range raw {
bigEndian[i] = raw[len(raw)-1-i]
}
integer := new(big.Int).SetBytes(bigEndian[:])
if raw[len(raw)-1]&0x80 != 0 {
integer.Sub(integer, new(big.Int).Lsh(big.NewInt(1), 128))
}
return decimal.NewFromBigInt(integer, 0)
}
func txBinaryScaledDecimalToUint64(value decimal.Decimal, scale int32, field string) (uint64, error) {
return txBinaryDecimalToUint64(value.Shift(scale), field)
}
@@ -1166,6 +1245,9 @@ func (enc *txBinaryEncoder) writeTxBinaryBody(tx *TxBinary, enumTable *txBinaryE
enc.writeUint32(tx.Signer)
enc.writeUint64(tx.Block)
enc.writeUint64(tx.BlockIndex)
if tx.SchemaVersion >= txBinarySchemaVersionV4 {
enc.writeUint64(uint64(tx.BlockAt))
}
enc.writeBool(tx.TxHash != nil)
if tx.TxHash != nil {
enc.writeBytes(tx.TxHash[:])
@@ -1182,7 +1264,7 @@ func (enc *txBinaryEncoder) writeTxBinaryBody(tx *TxBinary, enumTable *txBinaryE
if err := enc.writeMevAgentEntries(tx.MevAgent, enumTable); err != nil {
return err
}
if err := enc.writeSwaps(tx.Swaps, enumTable); err != nil {
if err := enc.writeSwaps(tx.Swaps, enumTable, tx.SchemaVersion); err != nil {
return err
}
return nil
@@ -1214,7 +1296,7 @@ func (enc *txBinaryEncoder) writeMevAgentEntries(entries []MevAgentBinary, enumT
return nil
}
func (enc *txBinaryEncoder) writeSwaps(swaps []SwapBinary, enumTable *txBinaryEnumTable) error {
func (enc *txBinaryEncoder) writeSwaps(swaps []SwapBinary, enumTable *txBinaryEnumTable, schemaVersion uint16) error {
enc.writeUint32(uint32(len(swaps)))
for i, swap := range swaps {
programID, err := enumTable.programs.id(swap.Program)
@@ -1264,6 +1346,10 @@ func (enc *txBinaryEncoder) writeSwaps(swaps []SwapBinary, enumTable *txBinaryEn
enc.writeUint32(swap.MigrateTopProgram)
enc.writeUint32(swap.LpMint)
enc.writeFloat64(swap.AfterSOLBalance)
if schemaVersion >= txBinarySchemaVersionV5 {
enc.writeUint64(swap.RealQuoteReserve)
enc.writeBytes(swap.VirtualQuoteReserve[:])
}
}
return nil
}
@@ -1379,7 +1465,7 @@ func (dec *txBinaryDecoder) readMevAgentEntries(enumTable *txBinaryEnumTable) ([
}
func (dec *txBinaryDecoder) readSwaps(enumTable *txBinaryEnumTable, _ []solana.PublicKey) ([]SwapBinary, error) {
return txBinaryReadSwaps(dec, enumTable)
return txBinaryReadSwaps(dec, enumTable, txBinarySchemaVersionCurrent)
}
func (dec *txBinaryDecoder) readTxBinaryBody(tx *TxBinary, enumTable *txBinaryEnumTable, addressTable []solana.PublicKey) error {
@@ -1474,7 +1560,7 @@ func (dec *txBinaryStreamDecoder) readTxsBinaryHeader() (*txsBinaryHeader, error
if err != nil {
return nil, err
}
if schemaVersion != txBinarySchemaVersionCurrent {
if !txBinarySchemaVersionSupported(schemaVersion) {
return nil, fmt.Errorf("unsupported tx binary schema version: %d", schemaVersion)
}
@@ -1531,7 +1617,7 @@ func (dec *txBinaryStreamDecoder) readTxsBinaryHeaderOrEOF() (*txsBinaryHeader,
if err != nil {
return nil, err
}
if schemaVersion != txBinarySchemaVersionCurrent {
if !txBinarySchemaVersionSupported(schemaVersion) {
return nil, fmt.Errorf("unsupported tx binary schema version: %d", schemaVersion)
}
@@ -1635,7 +1721,7 @@ func txBinaryReadMevAgentEntries(dec txBinaryBodyReader, enumTable *txBinaryEnum
return out, nil
}
func txBinaryReadSwaps(dec txBinaryBodyReader, enumTable *txBinaryEnumTable) ([]SwapBinary, error) {
func txBinaryReadSwaps(dec txBinaryBodyReader, enumTable *txBinaryEnumTable, schemaVersion uint16) ([]SwapBinary, error) {
count, err := dec.readUint32()
if err != nil {
return nil, err
@@ -1782,6 +1868,16 @@ func txBinaryReadSwaps(dec txBinaryBodyReader, enumTable *txBinaryEnumTable) ([]
if swap.AfterSOLBalance, err = dec.readFloat64(); err != nil {
return nil, err
}
if schemaVersion >= txBinarySchemaVersionV5 {
if swap.RealQuoteReserve, err = dec.readUint64(); err != nil {
return nil, err
}
rawVirtualQuoteReserve, err := dec.readN(len(swap.VirtualQuoteReserve))
if err != nil {
return nil, err
}
copy(swap.VirtualQuoteReserve[:], rawVirtualQuoteReserve)
}
out = append(out, swap)
}
return out, nil
@@ -1799,6 +1895,13 @@ func txBinaryReadTxBody(dec txBinaryBodyReader, tx *TxBinary, enumTable *txBinar
if tx.BlockIndex, err = dec.readUint64(); err != nil {
return err
}
if tx.SchemaVersion >= txBinarySchemaVersionV4 {
blockAt, err := dec.readUint64()
if err != nil {
return err
}
tx.BlockAt = int64(blockAt)
}
hasTxHash, err := dec.readBool()
if err != nil {
@@ -1840,7 +1943,7 @@ func txBinaryReadTxBody(dec txBinaryBodyReader, tx *TxBinary, enumTable *txBinar
if tx.MevAgent, err = txBinaryReadMevAgentEntries(dec, enumTable); err != nil {
return err
}
if tx.Swaps, err = txBinaryReadSwaps(dec, enumTable); err != nil {
if tx.Swaps, err = txBinaryReadSwaps(dec, enumTable, tx.SchemaVersion); err != nil {
return err
}
return nil
@@ -1854,7 +1957,9 @@ func txBinaryBuildMergePlan(sources []TxsBinaryReaderSource, opts TxsBinaryMerge
builder := txBinaryAddressTableBuilder{
index: make(map[solana.PublicKey]struct{}),
}
plan := &txsBinaryMergePlan{}
plan := &txsBinaryMergePlan{
schemaVersion: txBinarySchemaVersionCurrent,
}
hasBatch := false
for sourceIndex, source := range sources {
@@ -1896,15 +2001,10 @@ func txBinaryBuildMergePlan(sources []TxsBinaryReaderSource, opts TxsBinaryMerge
}
if !hasBatch {
plan.schemaVersion = header.schemaVersion
plan.enumVersion = header.enumVersion
plan.enumTable = header.enumTable
hasBatch = true
} else {
if header.schemaVersion != plan.schemaVersion {
reader.Close()
return nil, fmt.Errorf("source[%d].batch[%d]: schema version mismatch: got %d want %d", sourceIndex, batchIndex, header.schemaVersion, plan.schemaVersion)
}
if header.enumVersion != plan.enumVersion {
reader.Close()
return nil, fmt.Errorf("source[%d].batch[%d]: enum version mismatch: got %d want %d", sourceIndex, batchIndex, header.enumVersion, plan.enumVersion)

View File

@@ -19,6 +19,7 @@ func TestTxBinaryRoundTrip(t *testing.T) {
Signer: mustPubKey("So11111111111111111111111111111111111111112"),
Block: 123456789,
BlockIndex: 42,
BlockAt: 1710000000,
TxHash: &txHash,
CuFee: decimal.NewFromInt(5000),
CUPrice: decimal.RequireFromString("0.123456"),
@@ -81,6 +82,8 @@ func TestTxBinaryRoundTrip(t *testing.T) {
SlippageBps: decimal.RequireFromString("833.3333"),
BaseReserve: decimal.NewFromInt(5555),
QuoteReserve: decimal.NewFromInt(9999),
RealQuoteReserve: decimal.NewFromInt(10122),
VirtualQuoteReserve: decimal.NewFromInt(-123),
Mayhem: true,
Cashback: false,
UserBaseBalance: decimal.NewFromInt(777),
@@ -118,6 +121,9 @@ func TestTxBinaryRoundTrip(t *testing.T) {
if decoded.BlockIndex != original.BlockIndex {
t.Fatalf("BlockIndex = %d, want %d", decoded.BlockIndex, original.BlockIndex)
}
if decoded.BlockAt != original.BlockAt {
t.Fatalf("BlockAt = %d, want %d", decoded.BlockAt, original.BlockAt)
}
if decoded.TxHash == nil {
t.Fatal("TxHash = nil, want non-nil")
}
@@ -198,6 +204,12 @@ func TestTxBinaryRoundTrip(t *testing.T) {
if !swap.QuoteReserve.Equal(original.Swaps[0].QuoteReserve) {
t.Fatalf("swap.QuoteReserve = %s, want %s", swap.QuoteReserve, original.Swaps[0].QuoteReserve)
}
if !swap.RealQuoteReserve.Equal(original.Swaps[0].RealQuoteReserve) {
t.Fatalf("swap.RealQuoteReserve = %s, want %s", swap.RealQuoteReserve, original.Swaps[0].RealQuoteReserve)
}
if !swap.VirtualQuoteReserve.Equal(original.Swaps[0].VirtualQuoteReserve) {
t.Fatalf("swap.VirtualQuoteReserve = %s, want %s", swap.VirtualQuoteReserve, original.Swaps[0].VirtualQuoteReserve)
}
if !swap.UserBaseBalance.Equal(original.Swaps[0].UserBaseBalance) {
t.Fatalf("swap.UserBaseBalance = %s, want %s", swap.UserBaseBalance, original.Swaps[0].UserBaseBalance)
}
@@ -510,6 +522,7 @@ func TestTxsBinaryRoundTripWithSharedAddressTable(t *testing.T) {
Signer: mustPubKey("So11111111111111111111111111111111111111112"),
Block: 1,
BlockIndex: 1,
BlockAt: 1710000001,
CuFee: decimal.NewFromInt(1000),
CUPrice: decimal.RequireFromString("0.123456"),
BeforeSolBalance: decimal.RequireFromString("1.000000000"),
@@ -560,6 +573,7 @@ func TestTxsBinaryRoundTripWithSharedAddressTable(t *testing.T) {
tx2 := tx1
tx2.Block = 2
tx2.BlockIndex = 2
tx2.BlockAt = 1710000002
tx2.CuFee = decimal.NewFromInt(2000)
tx2.AfterSOLBalance = decimal.RequireFromString("0.700000000")
tx2.Swaps = []Swap{tx1.Swaps[0]}
@@ -581,6 +595,9 @@ func TestTxsBinaryRoundTripWithSharedAddressTable(t *testing.T) {
if decoded[0].Signer != tx1.Signer || decoded[1].Signer != tx2.Signer {
t.Fatalf("decoded signer mismatch")
}
if decoded[0].BlockAt != tx1.BlockAt || decoded[1].BlockAt != tx2.BlockAt {
t.Fatalf("decoded block_at mismatch")
}
if decoded[0].Swaps[0].Pool != tx1.Swaps[0].Pool || decoded[1].Swaps[0].Pool != tx2.Swaps[0].Pool {
t.Fatalf("decoded shared address mismatch")
}
@@ -603,6 +620,7 @@ func TestDecodeTxsBinaryReader(t *testing.T) {
Signer: mustPubKey("So11111111111111111111111111111111111111112"),
Block: 100,
BlockIndex: 7,
BlockAt: 1710000100,
CuFee: decimal.NewFromInt(111),
CUPrice: decimal.RequireFromString("0.123456"),
BeforeSolBalance: decimal.RequireFromString("1.000000000"),
@@ -649,6 +667,7 @@ func TestDecodeTxsBinaryReader(t *testing.T) {
tx2 := tx1
tx2.Block = 101
tx2.BlockIndex = 8
tx2.BlockAt = 1710000101
tx2.CuFee = decimal.NewFromInt(222)
tx2.AfterSOLBalance = decimal.RequireFromString("0.300000000")
tx2.Swaps = []Swap{tx1.Swaps[0]}
@@ -677,6 +696,9 @@ func TestDecodeTxsBinaryReader(t *testing.T) {
if decoded[0].Block != tx1.Block || decoded[1].Block != tx2.Block {
t.Fatalf("decoded block mismatch")
}
if decoded[0].BlockAt != tx1.BlockAt || decoded[1].BlockAt != tx2.BlockAt {
t.Fatalf("decoded block_at mismatch")
}
if decoded[0].Swaps[0].BaseAmount.Cmp(tx1.Swaps[0].BaseAmount) != 0 {
t.Fatalf("decoded tx1 swap base amount = %s, want %s", decoded[0].Swaps[0].BaseAmount, tx1.Swaps[0].BaseAmount)
}
@@ -724,6 +746,7 @@ func TestMergeTxsBinaryBytes(t *testing.T) {
Signer: mustPubKey("So11111111111111111111111111111111111111112"),
Block: 11,
BlockIndex: 1,
BlockAt: 1710000011,
CuFee: decimal.NewFromInt(10),
CUPrice: decimal.RequireFromString("0.000123"),
BeforeSolBalance: decimal.RequireFromString("1.100000000"),
@@ -755,6 +778,7 @@ func TestMergeTxsBinaryBytes(t *testing.T) {
Signer: mustPubKey("SysvarRent111111111111111111111111111111111"),
Block: 12,
BlockIndex: 2,
BlockAt: 1710000012,
CuFee: decimal.NewFromInt(20),
CUPrice: decimal.RequireFromString("0.000456"),
BeforeSolBalance: decimal.RequireFromString("2.200000000"),
@@ -818,6 +842,9 @@ func TestMergeTxsBinaryBytes(t *testing.T) {
if decoded[0].Block != tx1.Block || decoded[1].Block != tx2.Block {
t.Fatalf("decoded block mismatch")
}
if decoded[0].BlockAt != tx1.BlockAt || decoded[1].BlockAt != tx2.BlockAt {
t.Fatalf("decoded block_at mismatch")
}
}
func TestMergeTxsBinarySourcesToWriterWithConcatenatedBatches(t *testing.T) {
@@ -825,6 +852,7 @@ func TestMergeTxsBinarySourcesToWriterWithConcatenatedBatches(t *testing.T) {
Signer: mustPubKey("So11111111111111111111111111111111111111112"),
Block: 21,
BlockIndex: 1,
BlockAt: 1710000021,
CuFee: decimal.NewFromInt(1),
CUPrice: decimal.RequireFromString("0.000001"),
BeforeSolBalance: decimal.RequireFromString("1.000000000"),
@@ -835,9 +863,11 @@ func TestMergeTxsBinarySourcesToWriterWithConcatenatedBatches(t *testing.T) {
tx2 := tx1
tx2.Block = 22
tx2.BlockIndex = 2
tx2.BlockAt = 1710000022
tx2.Signer = mustPubKey("SysvarRent111111111111111111111111111111111")
tx3 := tx1
tx3.Block = 23
tx3.BlockAt = 1710000023
tx3.BlockIndex = 3
tx3.Signer = mustPubKey("ComputeBudget111111111111111111111111111111")
@@ -880,6 +910,9 @@ func TestMergeTxsBinarySourcesToWriterWithConcatenatedBatches(t *testing.T) {
if decoded[0].Block != tx1.Block || decoded[1].Block != tx2.Block || decoded[2].Block != tx3.Block {
t.Fatalf("decoded block order mismatch")
}
if decoded[0].BlockAt != tx1.BlockAt || decoded[1].BlockAt != tx2.BlockAt || decoded[2].BlockAt != tx3.BlockAt {
t.Fatalf("decoded block_at order mismatch")
}
}
func TestMergeTxsBinarySourcesToWriterWithBatchHeaderFuncSkip(t *testing.T) {
@@ -887,6 +920,7 @@ func TestMergeTxsBinarySourcesToWriterWithBatchHeaderFuncSkip(t *testing.T) {
Signer: mustPubKey("So11111111111111111111111111111111111111112"),
Block: 31,
BlockIndex: 1,
BlockAt: 1710000031,
CuFee: decimal.NewFromInt(1),
CUPrice: decimal.RequireFromString("0.000001"),
BeforeSolBalance: decimal.RequireFromString("1.000000000"),
@@ -897,10 +931,12 @@ func TestMergeTxsBinarySourcesToWriterWithBatchHeaderFuncSkip(t *testing.T) {
tx2 := tx1
tx2.Block = 32
tx2.BlockIndex = 2
tx2.BlockAt = 1710000032
tx2.Signer = mustPubKey("SysvarRent111111111111111111111111111111111")
tx3 := tx1
tx3.Block = 33
tx3.BlockIndex = 3
tx3.BlockAt = 1710000033
tx3.Signer = mustPubKey("ComputeBudget111111111111111111111111111111")
batch1, err := EncodeTxsBinary([]Tx{tx1})
@@ -960,15 +996,238 @@ func TestMergeTxsBinarySourcesToWriterWithBatchHeaderFuncSkip(t *testing.T) {
if decoded[0].Block != tx1.Block || decoded[1].Block != tx3.Block {
t.Fatalf("decoded block order mismatch after skip")
}
if decoded[0].BlockAt != tx1.BlockAt || decoded[1].BlockAt != tx3.BlockAt {
t.Fatalf("decoded block_at order mismatch after skip")
}
if source.opens != 2 {
t.Fatalf("source.opens = %d, want 2", source.opens)
}
}
func TestTxBinaryDecodeSchemaV3LeavesBlockAtZero(t *testing.T) {
original := &Tx{
Signer: mustPubKey("So11111111111111111111111111111111111111112"),
Block: 41,
BlockIndex: 1,
BlockAt: 1710000041,
}
encoded := mustEncodeTxBinaryV3(t, original)
decoded, err := DecodeTxBinary(encoded)
if err != nil {
t.Fatalf("DecodeTxBinary(v3) error = %v", err)
}
if decoded.Block != original.Block || decoded.BlockIndex != original.BlockIndex {
t.Fatalf("decoded block mismatch: got (%d,%d), want (%d,%d)", decoded.Block, decoded.BlockIndex, original.Block, original.BlockIndex)
}
if decoded.BlockAt != 0 {
t.Fatalf("BlockAt = %d, want 0 for legacy v3", decoded.BlockAt)
}
}
func TestTxBinaryDecodeSchemaV4DefaultsQuoteReserveComponents(t *testing.T) {
original := &Tx{
Signer: mustPubKey("So11111111111111111111111111111111111111112"),
Block: 42,
BlockIndex: 2,
BlockAt: 1710000042,
Swaps: []Swap{
{
Program: SolProgramPumpAMM,
Event: TxEventBuy,
QuoteReserve: decimal.NewFromInt(123456789),
RealQuoteReserve: decimal.NewFromInt(123000000),
VirtualQuoteReserve: decimal.NewFromInt(456789),
},
},
}
binaryTx, err := NewTxBinary(original)
if err != nil {
t.Fatalf("NewTxBinary() error = %v", err)
}
binaryTx.SchemaVersion = txBinarySchemaVersionV4
encoded, err := binaryTx.MarshalBinary()
if err != nil {
t.Fatalf("MarshalBinary(v4) error = %v", err)
}
decoded, err := DecodeTxBinary(encoded)
if err != nil {
t.Fatalf("DecodeTxBinary(v4) error = %v", err)
}
if decoded.BlockAt != original.BlockAt {
t.Fatalf("BlockAt = %d, want %d", decoded.BlockAt, original.BlockAt)
}
if len(decoded.Swaps) != 1 {
t.Fatalf("Swaps len = %d, want 1", len(decoded.Swaps))
}
swap := decoded.Swaps[0]
if !swap.QuoteReserve.Equal(original.Swaps[0].QuoteReserve) {
t.Fatalf("QuoteReserve = %s, want %s", swap.QuoteReserve, original.Swaps[0].QuoteReserve)
}
if !swap.RealQuoteReserve.Equal(swap.QuoteReserve) {
t.Fatalf("RealQuoteReserve = %s, want legacy QuoteReserve %s", swap.RealQuoteReserve, swap.QuoteReserve)
}
if !swap.VirtualQuoteReserve.IsZero() {
t.Fatalf("VirtualQuoteReserve = %s, want 0", swap.VirtualQuoteReserve)
}
}
func TestTxBinarySignedInt128RoundTrip(t *testing.T) {
values := []string{
"-170141183460469231731687303715884105728",
"-123",
"-1",
"0",
"1",
"123",
"170141183460469231731687303715884105727",
}
for _, value := range values {
t.Run(value, func(t *testing.T) {
want := decimal.RequireFromString(value)
raw, err := txBinaryDecimalToInt128(want, "value")
if err != nil {
t.Fatalf("txBinaryDecimalToInt128() error = %v", err)
}
if got := txBinaryInt128ToDecimal(raw); !got.Equal(want) {
t.Fatalf("round trip = %s, want %s", got, want)
}
})
}
invalid := []string{
"-170141183460469231731687303715884105729",
"170141183460469231731687303715884105728",
"1.5",
}
for _, value := range invalid {
if _, err := txBinaryDecimalToInt128(decimal.RequireFromString(value), "value"); err == nil {
t.Fatalf("txBinaryDecimalToInt128(%s) error = nil", value)
}
}
}
func TestMergeTxsBinaryBytesUpgradesSchemaV3AndPreservesV4BlockAt(t *testing.T) {
legacyTx := Tx{
Signer: mustPubKey("So11111111111111111111111111111111111111112"),
Block: 51,
BlockIndex: 1,
BlockAt: 1710000051,
}
currentTx := Tx{
Signer: mustPubKey("SysvarRent111111111111111111111111111111111"),
Block: 52,
BlockIndex: 2,
BlockAt: 1710000052,
Swaps: []Swap{
{
Program: SolProgramPumpAMM,
Event: TxEventBuy,
QuoteReserve: decimal.NewFromInt(222),
},
},
}
merged, err := MergeTxsBinaryBytes([][]byte{
mustEncodeTxsBinaryV3(t, []Tx{legacyTx}),
mustEncodeTxsBinaryV4(t, []Tx{currentTx}),
})
if err != nil {
t.Fatalf("MergeTxsBinaryBytes(v3,v4) error = %v", err)
}
var mergedBinary TxsBinary
if err := mergedBinary.UnmarshalBinary(merged); err != nil {
t.Fatalf("UnmarshalBinary(merged) error = %v", err)
}
if mergedBinary.SchemaVersion != txBinarySchemaVersionCurrent {
t.Fatalf("merged schema version = %d, want %d", mergedBinary.SchemaVersion, txBinarySchemaVersionCurrent)
}
decoded, err := DecodeTxsBinary(merged)
if err != nil {
t.Fatalf("DecodeTxsBinary(merged) error = %v", err)
}
if len(decoded) != 2 {
t.Fatalf("decoded len = %d, want 2", len(decoded))
}
if decoded[0].BlockAt != 0 {
t.Fatalf("legacy BlockAt = %d, want 0", decoded[0].BlockAt)
}
if decoded[1].BlockAt != currentTx.BlockAt {
t.Fatalf("current BlockAt = %d, want %d", decoded[1].BlockAt, currentTx.BlockAt)
}
if len(decoded[1].Swaps) != 1 || !decoded[1].Swaps[0].RealQuoteReserve.Equal(currentTx.Swaps[0].QuoteReserve) || !decoded[1].Swaps[0].VirtualQuoteReserve.IsZero() {
t.Fatalf("v4 quote reserve components were not preserved: %+v", decoded[1].Swaps)
}
}
func mustPubKey(value string) solana.PublicKey {
return solana.MustPublicKeyFromBase58(value)
}
func mustEncodeTxBinaryV3(t *testing.T, tx *Tx) []byte {
t.Helper()
binaryTx, err := NewTxBinary(tx)
if err != nil {
t.Fatalf("NewTxBinary() error = %v", err)
}
binaryTx.SchemaVersion = txBinarySchemaVersionV3
encoded, err := binaryTx.MarshalBinary()
if err != nil {
t.Fatalf("MarshalBinary(v3) error = %v", err)
}
return encoded
}
func mustEncodeTxsBinary(t *testing.T, txs []Tx) []byte {
t.Helper()
encoded, err := EncodeTxsBinary(txs)
if err != nil {
t.Fatalf("EncodeTxsBinary() error = %v", err)
}
return encoded
}
func mustEncodeTxsBinaryV3(t *testing.T, txs []Tx) []byte {
t.Helper()
binaryTxs, err := NewTxsBinary(txs)
if err != nil {
t.Fatalf("NewTxsBinary() error = %v", err)
}
binaryTxs.SchemaVersion = txBinarySchemaVersionV3
for i := range binaryTxs.Txs {
binaryTxs.Txs[i].SchemaVersion = txBinarySchemaVersionV3
}
encoded, err := binaryTxs.MarshalBinary()
if err != nil {
t.Fatalf("MarshalBinary(v3) error = %v", err)
}
return encoded
}
func mustEncodeTxsBinaryV4(t *testing.T, txs []Tx) []byte {
t.Helper()
binaryTxs, err := NewTxsBinary(txs)
if err != nil {
t.Fatalf("NewTxsBinary() error = %v", err)
}
binaryTxs.SchemaVersion = txBinarySchemaVersionV4
for i := range binaryTxs.Txs {
binaryTxs.Txs[i].SchemaVersion = txBinarySchemaVersionV4
}
encoded, err := binaryTxs.MarshalBinary()
if err != nil {
t.Fatalf("MarshalBinary(v4) error = %v", err)
}
return encoded
}
func mustTxBinary(t *testing.T, data []byte) *TxsBinary {
t.Helper()