Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
48 changes: 0 additions & 48 deletions ocp/data/blockchain.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,13 +2,11 @@ package data

import (
"context"
"crypto/ed25519"

"github.com/mr-tron/base58"

"github.com/code-payments/ocp-server/database/query"
"github.com/code-payments/ocp-server/metrics"
"github.com/code-payments/ocp-server/ocp/config"
"github.com/code-payments/ocp-server/solana"
"github.com/code-payments/ocp-server/solana/token"
)
Expand All @@ -24,15 +22,12 @@ type BlockchainData interface {
GetBlockchainAccountDataAfterBlock(ctx context.Context, account string, slot uint64) ([]byte, uint64, error)
GetBlockchainBalance(ctx context.Context, account string, commitment solana.Commitment) (uint64, uint64, error)
GetBlockchainBlock(ctx context.Context, slot uint64) (*solana.Block, error)
GetBlockchainBlockSignatures(ctx context.Context, slot uint64) ([]string, error)
GetBlockchainBlocksWithLimit(ctx context.Context, start uint64, limit uint64) ([]uint64, error)
GetBlockchainHistory(ctx context.Context, account string, commitment solana.Commitment, opts ...query.Option) ([]*solana.TransactionSignature, error)
GetBlockchainMinimumBalanceForRentExemption(ctx context.Context, size uint64) (uint64, error)
GetBlockchainLatestBlockhash(ctx context.Context) (solana.Blockhash, error)
GetBlockchainSignatureStatuses(ctx context.Context, signatures []solana.Signature) ([]*solana.SignatureStatus, error)
GetBlockchainSlot(ctx context.Context, commitment solana.Commitment) (uint64, error)
GetBlockchainTokenAccountInfo(ctx context.Context, account, mint string, commitment solana.Commitment) (*token.Account, error)
GetBlockchainTokenAccountsByOwner(ctx context.Context, account string) ([]ed25519.PublicKey, error)
GetBlockchainTransaction(ctx context.Context, sig string, commitment solana.Commitment) (*solana.ConfirmedTransaction, error)
GetBlockchainTransactionTokenBalances(ctx context.Context, sig string) (*solana.TransactionTokenBalances, error)
GetBlockchainFilteredProgramAccounts(ctx context.Context, program string, offset uint, filterValue []byte) ([]solana.ProgramAccount, uint64, error)
Expand Down Expand Up @@ -130,22 +125,6 @@ func (dp *BlockchainProvider) GetBlockchainTokenAccountInfo(ctx context.Context,
}
return res, err
}
func (dp *BlockchainProvider) GetBlockchainTokenAccountsByOwner(ctx context.Context, account string) ([]ed25519.PublicKey, error) {
tracer := metrics.TraceMethodCall(ctx, blockchainProviderMetricsName, "GetBlockchainTokenAccountsByOwner")
defer tracer.End()

accountId, err := base58.Decode(account)
if err != nil {
return nil, err
}

res, err := dp.sc.GetTokenAccountsByOwner(accountId, config.CoreMintPublicKeyBytes)

if err != nil {
tracer.OnError(err)
}
return res, err
}
func (dp *BlockchainProvider) GetBlockchainSlot(ctx context.Context, commitment solana.Commitment) (uint64, error) {
tracer := metrics.TraceMethodCall(ctx, blockchainProviderMetricsName, "GetBlockchainSlot")
defer tracer.End()
Expand All @@ -158,21 +137,6 @@ func (dp *BlockchainProvider) GetBlockchainSlot(ctx context.Context, commitment
return res, err
}

func (dp *BlockchainProvider) GetBlockchainBlocksWithLimit(ctx context.Context, start uint64, limit uint64) ([]uint64, error) {
tracer := metrics.TraceMethodCall(ctx, blockchainProviderMetricsName, "GetBlockchainBlocksWithLimit")
defer tracer.End()

// TODO: this call is deprecated, remove it
// https://docs.solana.com/developing/clients/jsonrpc-api#getconfirmedblockswithlimit

res, err := dp.sc.GetConfirmedBlocksWithLimit(start, limit)

if err != nil {
tracer.OnError(err)
}
return res, err
}

func (dp *BlockchainProvider) GetBlockchainBlock(ctx context.Context, slot uint64) (*solana.Block, error) {
tracer := metrics.TraceMethodCall(ctx, blockchainProviderMetricsName, "GetBlockchainBlock")
defer tracer.End()
Expand All @@ -185,18 +149,6 @@ func (dp *BlockchainProvider) GetBlockchainBlock(ctx context.Context, slot uint6
return res, err
}

func (dp *BlockchainProvider) GetBlockchainBlockSignatures(ctx context.Context, slot uint64) ([]string, error) {
tracer := metrics.TraceMethodCall(ctx, blockchainProviderMetricsName, "GetBlockchainBlockSignatures")
defer tracer.End()

res, err := dp.sc.GetBlockSignatures(slot)

if err != nil {
tracer.OnError(err)
}
return res, err
}

func (dp *BlockchainProvider) GetBlockchainHistory(ctx context.Context, account string, commitment solana.Commitment, opts ...query.Option) ([]*solana.TransactionSignature, error) {
tracer := metrics.TraceMethodCall(ctx, blockchainProviderMetricsName, "GetBlockchainHistory")
defer tracer.End()
Expand Down
7 changes: 6 additions & 1 deletion ocp/data/transaction/transaction.go
Original file line number Diff line number Diff line change
Expand Up @@ -92,11 +92,16 @@ func FromConfirmedTransaction(tx *solana.ConfirmedTransaction) (*Record, error)
return nil, errors.New("unsupported transaction version")
}

data, err := tx.Transaction.Marshal()
if err != nil {
return nil, err
}

sig := tx.Transaction.Signature()
res := &Record{
Signature: base58.Encode(sig),
Slot: tx.Slot,
Data: tx.Transaction.Marshal(),
Data: data,
HasErrors: tx.Err != nil,
ConfirmationState: ConfirmationFinalized,
CreatedAt: time.Now(),
Expand Down
11 changes: 8 additions & 3 deletions ocp/rpc/transaction/errors.go
Original file line number Diff line number Diff line change
Expand Up @@ -162,26 +162,31 @@ func toInvalidTxnSignatureErrorDetails(
actionId uint32,
txn solana.Transaction,
signature *commonpb.Signature,
) *transactionpb.ErrorDetails {
) (*transactionpb.ErrorDetails, error) {
// Clear out all signatures, so clients have no way of submitting this transaction
var emptySig solana.Signature
for i := range txn.Signatures {
copy(txn.Signatures[i][:], emptySig[:])
}

marshalledTxn, err := txn.Marshal()
if err != nil {
return nil, err
}

return &transactionpb.ErrorDetails{
Type: &transactionpb.ErrorDetails_InvalidSignature{
InvalidSignature: &transactionpb.InvalidSignatureErrorDetails{
ActionId: actionId,
ExpectedBlob: &transactionpb.InvalidSignatureErrorDetails_ExpectedTransaction{
ExpectedTransaction: &commonpb.Transaction{
Value: txn.Marshal(),
Value: marshalledTxn,
},
},
ProvidedSignature: signature,
},
},
}
}, nil
}

func toInvalidVirtualIxnSignatureErrorDetails(
Expand Down
38 changes: 32 additions & 6 deletions ocp/rpc/transaction/stateful_swap.go
Original file line number Diff line number Diff line change
Expand Up @@ -593,7 +593,11 @@ func (s *transactionServer) handleReserveStatefulSwap(

txn.SetBlockhash(selectedNonce.Blockhash)

marshalledTxnMessage := txn.Message.Marshal()
marshalledTxnMessage, err := txn.Message.Marshal()
if err != nil {
log.With(zap.Error(err)).Warn("failure marshalling transaction message")
return handleStatefulSwapError(streamer, err)
}

//
// Section: Server parameters
Expand Down Expand Up @@ -657,10 +661,15 @@ func (s *transactionServer) handleReserveStatefulSwap(
marshalledTxnMessage,
protoSignature.Value,
) {
errorDetails, err := toInvalidTxnSignatureErrorDetails(0, txn, protoSignature)
if err != nil {
log.With(zap.Error(err)).Warn("failure creating error details")
return handleStatefulSwapError(streamer, err)
}
return handleStatefulSwapStructuredError(
streamer,
transactionpb.StatefulSwapResponse_Error_SIGNATURE_ERROR,
toInvalidTxnSignatureErrorDetails(0, txn, protoSignature),
errorDetails,
)
}

Expand All @@ -683,7 +692,11 @@ func (s *transactionServer) handleReserveStatefulSwap(
return handleStatefulSwapError(streamer, err)
}

marshalledTxn := txn.Marshal()
marshalledTxn, err := txn.Marshal()
if err != nil {
log.With(zap.Error(err)).Warn("failure marshalling transaction")
return handleStatefulSwapError(streamer, err)
}

txnSignature := base58.Encode(txn.Signature())

Expand Down Expand Up @@ -1022,7 +1035,11 @@ func (s *transactionServer) handleStablecoinStatefulSwap(

txn.SetBlockhash(selectedNonce.Blockhash)

marshalledTxnMessage := txn.Message.Marshal()
marshalledTxnMessage, err := txn.Message.Marshal()
if err != nil {
log.With(zap.Error(err)).Warn("failure marshalling transaction message")
return handleStatefulSwapError(streamer, err)
}

//
// Section: Server parameters
Expand Down Expand Up @@ -1082,10 +1099,15 @@ func (s *transactionServer) handleStablecoinStatefulSwap(
marshalledTxnMessage,
protoSignature.Value,
) {
errorDetails, err := toInvalidTxnSignatureErrorDetails(0, txn, protoSignature)
if err != nil {
log.With(zap.Error(err)).Warn("failure creating error details")
return handleStatefulSwapError(streamer, err)
}
return handleStatefulSwapStructuredError(
streamer,
transactionpb.StatefulSwapResponse_Error_SIGNATURE_ERROR,
toInvalidTxnSignatureErrorDetails(0, txn, protoSignature),
errorDetails,
)
}

Expand All @@ -1101,7 +1123,11 @@ func (s *transactionServer) handleStablecoinStatefulSwap(
return handleStatefulSwapError(streamer, err)
}

marshalledTxn := txn.Marshal()
marshalledTxn, err := txn.Marshal()
if err != nil {
log.With(zap.Error(err)).Warn("failure marshalling transaction")
return handleStatefulSwapError(streamer, err)
}

txnSignature := base58.Encode(txn.Signature())

Expand Down
13 changes: 11 additions & 2 deletions ocp/rpc/transaction/stateless_swap.go
Original file line number Diff line number Diff line change
Expand Up @@ -274,7 +274,11 @@ func (s *transactionServer) handleStablecoinStatelessSwap(
}
txn.SetBlockhash(blockhash)

marshalledTxnMessage := txn.Message.Marshal()
marshalledTxnMessage, err := txn.Message.Marshal()
if err != nil {
log.With(zap.Error(err)).Warn("failure marshalling transaction message")
return handleStatelessSwapError(streamer, err)
}

//
// Section: Server parameters
Expand Down Expand Up @@ -322,10 +326,15 @@ func (s *transactionServer) handleStablecoinStatelessSwap(
marshalledTxnMessage,
protoSignature.Value,
) {
errorDetails, err := toInvalidTxnSignatureErrorDetails(0, txn, protoSignature)
if err != nil {
log.With(zap.Error(err)).Warn("failure creating error details")
return handleStatelessSwapError(streamer, err)
}
return handleStatelessSwapStructuredError(
streamer,
transactionpb.StatelessSwapResponse_Error_SIGNATURE_ERROR,
toInvalidTxnSignatureErrorDetails(0, txn, protoSignature),
errorDetails,
)
}

Expand Down
7 changes: 6 additions & 1 deletion ocp/worker/currency/feeburner/worker.go
Original file line number Diff line number Diff line change
Expand Up @@ -149,7 +149,12 @@ func (p *runtime) packBurnBatches(targets []*burnTarget, maxBurnsPerBatch int) [

candidate := append(current, target)
txn := p.makeBurnTransaction(candidate)
if len(txn.Marshal()) > solana.MaxTransactionSize {
marshalledTxn, err := txn.Marshal()
if err != nil {
p.log.With(zap.Error(err), zap.String("mint", target.mint)).Warn("skipping currency with unmarshallable burn transaction")
continue
}
if len(marshalledTxn) > solana.MaxLegacyTransactionSize {
if len(current) == 0 {
p.log.With(zap.String("mint", target.mint)).Warn("skipping currency with oversized burn transaction")
continue
Expand Down
8 changes: 6 additions & 2 deletions ocp/worker/currency/feeburner/worker_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,9 @@ func TestPackBurnBatches(t *testing.T) {

for i, batch := range batches {
txn := p.makeBurnTransaction(batch)
assert.LessOrEqual(t, len(txn.Marshal()), solana.MaxTransactionSize, fmt.Sprintf("batch %d exceeds size limit", i))
marshalledTxn, err := txn.Marshal()
require.NoError(t, err)
assert.LessOrEqual(t, len(marshalledTxn), solana.MaxLegacyTransactionSize, fmt.Sprintf("batch %d exceeds size limit", i))
assert.LessOrEqual(t, len(batch), defaultMaxBurnsPerBatch, fmt.Sprintf("batch %d exceeds max burns", i))
}

Expand All @@ -57,7 +59,9 @@ func TestPackBurnBatches(t *testing.T) {
}
overfilled := append(append([]*burnTarget{}, batches[i]...), batches[i+1][0])
txn := p.makeBurnTransaction(overfilled)
assert.Greater(t, len(txn.Marshal()), solana.MaxTransactionSize, fmt.Sprintf("batch %d is not fully packed", i))
marshalledTxn, err := txn.Marshal()
require.NoError(t, err)
assert.Greater(t, len(marshalledTxn), solana.MaxLegacyTransactionSize, fmt.Sprintf("batch %d is not fully packed", i))
}

assert.Greater(t, len(batches[0]), 1)
Expand Down
6 changes: 5 additions & 1 deletion ocp/worker/currency/launcher/util.go
Original file line number Diff line number Diff line change
Expand Up @@ -757,7 +757,11 @@ func (p *runtime) resizeAndExtendBlockchainAccounts(ctx context.Context, account
ixns...,
)

if len(txn.Marshal()) > solana.MaxTransactionSize {
marshalledTxn, err := txn.Marshal()
if err != nil {
return errors.Wrap(err, "error marshalling transaction")
}
if len(marshalledTxn) > solana.MaxLegacyTransactionSize {
return errors.New("transaction exceeds maximum size")
}

Expand Down
5 changes: 4 additions & 1 deletion ocp/worker/sequencer/worker.go
Original file line number Diff line number Diff line change
Expand Up @@ -261,7 +261,10 @@ func (p *runtime) handlePending(ctx context.Context, record *fulfillment.Record)
record.Signature = pointer.String(base58.Encode(txn.Signature()))
record.Nonce = pointer.String(selectedSolanaNonce.Account.PublicKey().ToBase58())
record.Blockhash = pointer.String(base58.Encode(selectedSolanaNonce.Blockhash[:]))
record.Data = txn.Marshal()
record.Data, err = txn.Marshal()
if err != nil {
return err
}

err = selectedSolanaNonce.MarkReservedWithSignature(ctx, *record.Signature)
if err != nil {
Expand Down
9 changes: 7 additions & 2 deletions ocp/worker/sequencer/worker_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -288,13 +288,16 @@ func (e *workerTestEnv) createAnyFulfillmentInState(t *testing.T, state fulfillm

txn.Sign(fakeCodeAccouht.PrivateKey().ToBytes())

marshalledTxn, err := txn.Marshal()
require.NoError(t, err)

fulfillmentRecord := &fulfillment.Record{
Intent: testutil.NewRandomAccount(t).PublicKey().ToBase58(),
IntentType: intent.OpenAccounts,
ActionId: 3,
ActionType: action.OpenAccount,
FulfillmentType: fulfillment.InitializeLockedTimelockAccount,
Data: txn.Marshal(),
Data: marshalledTxn,
Signature: pointer.String(base58.Encode(txn.Signature())),
Source: "source",
Nonce: pointer.String(fakeNonceAccount.PublicKey().ToBase58()),
Expand Down Expand Up @@ -372,7 +375,9 @@ func (e *workerTestEnv) assertFulfillmentCreatedOnDemand(t *testing.T, id uint64
assert.Equal(t, expectedSignature, *fulfillmentRecord.Signature)
assert.Equal(t, nonceAddress, *fulfillmentRecord.Nonce)
assert.Equal(t, blockhash, *fulfillmentRecord.Blockhash)
assert.Equal(t, expectedTxn.Marshal(), fulfillmentRecord.Data)
expectedData, err := expectedTxn.Marshal()
require.NoError(t, err)
assert.Equal(t, expectedData, fulfillmentRecord.Data)

e.assertNonceState(t, nonceAddress, nonce.StateReserved, expectedSignature, blockhash)
}
Expand Down
5 changes: 4 additions & 1 deletion ocp/worker/swap/util.go
Original file line number Diff line number Diff line change
Expand Up @@ -399,7 +399,10 @@ func (p *runtime) markSwapCancelling(
swapRecord.Nonce = cancelNonce.Account.PublicKey().ToBase58()
swapRecord.Blockhash = base58.Encode(cancelNonce.Blockhash[:])
swapRecord.TransactionSignature = cancelTransactionSignature
swapRecord.TransactionBlob = cancelTxn.Marshal()
swapRecord.TransactionBlob, err = cancelTxn.Marshal()
if err != nil {
return err
}
swapRecord.State = swap.StateCancelling
return p.data.SaveSwap(ctx, swapRecord)
})
Expand Down
Loading
Loading