From c207eeac38114c03879e1adb5564dbcc6952bd50 Mon Sep 17 00:00:00 2001 From: aman035 Date: Mon, 24 Aug 2026 14:19:41 +0530 Subject: [PATCH 1/4] fix(chains): reject truncated UniversalTx events instead of defaulting tx_type to GAS --- .../chains/common/txtype_routing_test.go | 86 ++++++++++++ universalClient/chains/evm/event_parser.go | 36 ++++-- .../chains/evm/event_parser_test.go | 31 ++++- universalClient/chains/svm/event_parser.go | 79 ++++++------ .../chains/svm/event_parser_test.go | 119 +++++++++++------ .../chains/svm/universal_tx_golden_test.go | 122 ++++++++++++++++++ 6 files changed, 375 insertions(+), 98 deletions(-) create mode 100644 universalClient/chains/common/txtype_routing_test.go create mode 100644 universalClient/chains/svm/universal_tx_golden_test.go diff --git a/universalClient/chains/common/txtype_routing_test.go b/universalClient/chains/common/txtype_routing_test.go new file mode 100644 index 00000000..65426b8b --- /dev/null +++ b/universalClient/chains/common/txtype_routing_test.go @@ -0,0 +1,86 @@ +package common + +import ( + "encoding/json" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/pushchain/push-chain-node/universalClient/store" + uexecutortypes "github.com/pushchain/push-chain-node/x/uexecutor/types" +) + +// The wire values the gateways emit are 0-indexed (Gas, GasAndPayload, Funds, +// FundsAndPayload) while the chain enum reserves 0 for UNSPECIFIED, so the +// mapping is shifted by one. A decoder that leaves TxType unset therefore does +// not produce "unknown", it produces GAS. +func TestConstructInbound_TxTypeMapping(t *testing.T) { + processor := &EventProcessor{} + + for _, tc := range []struct { + wire uint + want uexecutortypes.TxType + }{ + {0, uexecutortypes.TxType_GAS}, + {1, uexecutortypes.TxType_GAS_AND_PAYLOAD}, + {2, uexecutortypes.TxType_FUNDS}, + {3, uexecutortypes.TxType_FUNDS_AND_PAYLOAD}, + {4, uexecutortypes.TxType_UNSPECIFIED_TX}, + {99, uexecutortypes.TxType_UNSPECIFIED_TX}, + } { + data, err := json.Marshal(UniversalTx{ + SourceChain: "solana:devnet", + Sender: "0xabc", + Recipient: "0xdef", + Amount: "5000000", + TxType: tc.wire, + }) + require.NoError(t, err) + + inbound, err := processor.constructInbound(&store.Event{ + EventID: "sig:0", + EventData: data, + }) + require.NoError(t, err) + assert.Equal(t, tc.want, inbound.TxType, "wire value %d", tc.wire) + } +} + +// A FUNDS transfer must never reach the keeper as GAS. The two dispatch to +// different handlers: GAS mints and autoswaps into the sender UEA, FUNDS +// deposits PRC20 to the recipient, so the same amount lands with a different +// party. This is the end to end assertion the finding asks for. +func TestConstructInbound_FundsNeverBecomesGas(t *testing.T) { + processor := &EventProcessor{} + + data, err := json.Marshal(UniversalTx{ + SourceChain: "solana:devnet", + Sender: "0xabc", + Recipient: "0xdef", + Amount: "5000000", + TxType: 2, // Funds, as the real devnet events carry + }) + require.NoError(t, err) + + inbound, err := processor.constructInbound(&store.Event{ + EventID: "sig:0", + EventData: data, + }) + require.NoError(t, err) + + assert.Equal(t, uexecutortypes.TxType_FUNDS, inbound.TxType) + assert.NotEqual(t, uexecutortypes.TxType_GAS, inbound.TxType, + "a FUNDS transfer routed to GAS credits the sender instead of the recipient") +} + +// An event whose data never made it past the decoder must be refused outright +// rather than defaulted. The parsers now discard such events, so this is the +// backstop if one ever reaches the store. +func TestConstructInbound_RejectsEventWithoutData(t *testing.T) { + processor := &EventProcessor{} + + _, err := processor.constructInbound(&store.Event{EventID: "sig:0"}) + require.Error(t, err) + assert.Contains(t, err.Error(), "event data is missing") +} diff --git a/universalClient/chains/evm/event_parser.go b/universalClient/chains/evm/event_parser.go index 2c338232..adbedc77 100644 --- a/universalClient/chains/evm/event_parser.go +++ b/universalClient/chains/evm/event_parser.go @@ -95,8 +95,15 @@ func parseSendFundsEvent(log *types.Log, chainID string, logger zerolog.Logger) ExpiryBlockHeight: 0, // 0 means no expiry } - // Parse universal tx event data - parseUniversalTxEvent(event, log, chainID, logger) + // Parse universal tx event data. A malformed event is dropped rather than + // stored half-decoded: the zero values it would carry are not neutral. + if err := parseUniversalTxEvent(event, log, chainID, logger); err != nil { + logger.Warn(). + Err(err). + Str("event_id", eventID). + Msg("discarding malformed UniversalTx event") + return nil + } return event } @@ -167,10 +174,13 @@ func parseOutboundObservationEvent(log *types.Log, chainID string, logger zerolo } // parseUniversalTxEvent parses a UniversalTx event from log data. -func parseUniversalTxEvent(event *store.Event, log *types.Log, chainID string, logger zerolog.Logger) { +// +// Static words through txType are required. A truncated log is rejected rather +// than returned half-filled, since the zero values are not neutral: txType 0 is +// GAS, which routes funds to a different account than FUNDS does. +func parseUniversalTxEvent(event *store.Event, log *types.Log, chainID string, logger zerolog.Logger) error { if len(log.Topics) < 3 { - logger.Warn().Msg("not enough indexed fields; nothing to do") - return + return fmt.Errorf("need 3 indexed fields, got %d", len(log.Topics)) } payload := common.UniversalTx{ @@ -181,9 +191,7 @@ func parseUniversalTxEvent(event *store.Event, log *types.Log, chainID string, l } if len(log.Data) < 32*5 { - b, _ := json.Marshal(payload) - event.EventData = b - return + return fmt.Errorf("log data has %d bytes, need at least %d for the static words", len(log.Data), 32*5) } // Parse common static fields: token (Word 0), amount (Word 1) @@ -191,7 +199,7 @@ func parseUniversalTxEvent(event *store.Event, log *types.Log, chainID string, l payload.Amount = new(big.Int).SetBytes(log.Data[1*32 : 2*32]).String() dataOffset := new(big.Int).SetBytes(log.Data[2*32 : 3*32]).Uint64() - parseUniversalTx(event, log, dataOffset, &payload, logger) + return parseUniversalTx(event, log, dataOffset, &payload, logger) } // readDynamicBytes decodes ABI-encoded dynamic bytes at the given absolute offset in data. @@ -278,7 +286,7 @@ UniversalTx Event (V2 - upgraded chains): - signatureData (bytes) — Word 5 (offset) - fromCEA (bool) — Word 6 */ -func parseUniversalTx(event *store.Event, log *types.Log, dataOffset uint64, payload *common.UniversalTx, logger zerolog.Logger) { +func parseUniversalTx(event *store.Event, log *types.Log, dataOffset uint64, payload *common.UniversalTx, logger zerolog.Logger) error { data := log.Data decodePayload(data, dataOffset, payload, logger) @@ -288,10 +296,9 @@ func parseUniversalTx(event *store.Event, log *types.Log, dataOffset uint64, pay payload.RevertFundRecipient = ethcommon.BytesToAddress(w[12:32]).Hex() } - // txType (Word 4) - if w := readWord(data, 4); w != nil { - payload.TxType = uint(new(big.Int).SetBytes(w).Uint64()) - } + // txType (Word 4). Always present: the caller rejects anything shorter than + // five words, which is what stops this being left at 0 and read as GAS. + payload.TxType = uint(new(big.Int).SetBytes(readWord(data, 4)).Uint64()) // signatureData (Word 5 offset) if w := readWord(data, 5); w != nil { @@ -304,4 +311,5 @@ func parseUniversalTx(event *store.Event, log *types.Log, dataOffset uint64, pay } finalizeEvent(event, payload, logger) + return nil } diff --git a/universalClient/chains/evm/event_parser_test.go b/universalClient/chains/evm/event_parser_test.go index d3a71e79..4fb4c540 100644 --- a/universalClient/chains/evm/event_parser_test.go +++ b/universalClient/chains/evm/event_parser_test.go @@ -183,7 +183,10 @@ func TestParseEventData(t *testing.T) { assert.NotNil(t, event.EventData) }) - t.Run("handles missing data gracefully", func(t *testing.T) { + // A log too short to carry txType is discarded rather than stored with the + // field left at 0, which is GAS on the wire and routes to a different + // account than FUNDS does. + t.Run("discards a log with no data", func(t *testing.T) { log := &types.Log{ Topics: []ethcommon.Hash{ ethcommon.HexToHash("0x1234"), @@ -193,9 +196,29 @@ func TestParseEventData(t *testing.T) { Data: []byte{}, // Empty data } - event := ParseEvent(log, EventTypeSendFunds, config.Chain, logger) - // Should still create event but with minimal data - require.NotNil(t, event) + assert.Nil(t, ParseEvent(log, EventTypeSendFunds, config.Chain, logger)) + }) + + t.Run("discards a log one word short of txType", func(t *testing.T) { + log := &types.Log{ + Topics: []ethcommon.Hash{ + ethcommon.HexToHash("0x1234"), + ethcommon.HexToHash("0x000000000000000000000000742d35cc6634c0532925a3b844bc9e7595f0beb7"), + ethcommon.HexToHash("0x000000000000000000000000dac17f958d2ee523a2206206994597c13d831ec7"), + }, + Data: make([]byte, 32*4), // words 0..3 present, txType (word 4) missing + } + + assert.Nil(t, ParseEvent(log, EventTypeSendFunds, config.Chain, logger)) + }) + + t.Run("discards a log without the indexed fields", func(t *testing.T) { + log := &types.Log{ + Topics: []ethcommon.Hash{ethcommon.HexToHash("0x1234")}, + Data: make([]byte, 32*8), + } + + assert.Nil(t, ParseEvent(log, EventTypeSendFunds, config.Chain, logger)) }) } diff --git a/universalClient/chains/svm/event_parser.go b/universalClient/chains/svm/event_parser.go index 4c25625d..9df44617 100644 --- a/universalClient/chains/svm/event_parser.go +++ b/universalClient/chains/svm/event_parser.go @@ -93,8 +93,15 @@ func parseSendFundsEvent(log string, signature string, slot uint64, logIndex uin ExpiryBlockHeight: 0, // Will be set based on confirmation type if needed } - // Parse event data from this log - parseUniversalTxEvent(event, decoded, logIndex, chainID, logger) + // Parse event data from this log. A malformed event is dropped rather than + // stored half-decoded: the zero values it would carry are not neutral. + if err := parseUniversalTxEvent(event, decoded, logIndex, chainID, logger); err != nil { + logger.Warn(). + Err(err). + Str("event_id", eventID). + Msg("discarding malformed UniversalTx event") + return nil + } return event } @@ -201,15 +208,11 @@ func parseOutboundObservationEvent(log string, signature string, slot uint64, lo // parseUniversalTxEvent extracts specific data from a single log event // For TxWithFunds events, it JSON-marshals the decoded fields into event.EventData. -func parseUniversalTxEvent(event *store.Event, decoded []byte, logIndex uint, chainID string, logger zerolog.Logger) { - // Parse the TxWithFunds event +func parseUniversalTxEvent(event *store.Event, decoded []byte, logIndex uint, chainID string, logger zerolog.Logger) error { + // Parse the UniversalTx event payload, err := decodeUniversalTxEvent(decoded, logger) if err != nil { - logger.Warn(). - Err(err). - Uint("log_index", logIndex). - Msg("failed to decode TxWithFunds event") - return + return fmt.Errorf("decode UniversalTx event: %w", err) } // Set source chain and log index @@ -217,13 +220,11 @@ func parseUniversalTxEvent(event *store.Event, decoded []byte, logIndex uint, ch payload.LogIndex = logIndex // Marshal and store into event.EventData - if b, err := json.Marshal(payload); err == nil { - event.EventData = b - } else { - logger.Warn(). - Err(err). - Msg("failed to marshal universal tx payload") + b, err := json.Marshal(payload) + if err != nil { + return fmt.Errorf("marshal universal tx payload: %w", err) } + event.EventData = b // if TxType is 0 or 1, use FAST else use STANDARD if payload.TxType == 0 || payload.TxType == 1 { @@ -231,16 +232,20 @@ func parseUniversalTxEvent(event *store.Event, decoded []byte, logIndex uint, ch } else { event.ConfirmationType = store.ConfirmationStandard } + + return nil } -// decodeUniversalTxEvent decodes a TxWithFunds event +// decodeUniversalTxEvent decodes the gateway's Borsh-encoded UniversalTx event: +// +// sender 32, recipient 20, token 32, amount u64, payload (u32 len + bytes), +// revert_recipient 32, tx_type 1, signature_data (u32 len + bytes), from_cea 1 +// +// Every field through signature_data is required. A truncated event is rejected +// rather than returned half-filled, since the zero values are not neutral: +// tx_type 0 is GAS, which routes funds to a different account than FUNDS does. +// Only from_cea is optional, defaulting to false as its absence cannot misroute. func decodeUniversalTxEvent(data []byte, logger zerolog.Logger) (*common.UniversalTx, error) { - if len(data) < 120 { - logger.Warn(). - Int("data_len", len(data)). - Msg("data might be too short for complete TxWithFunds event") - } - offset := 8 payload := &common.UniversalTx{} @@ -286,19 +291,14 @@ func decodeUniversalTxEvent(data []byte, logger zerolog.Logger) (*common.Univers // Parse data field length (4 bytes) if len(data) < offset+4 { - logger.Warn().Msg("not enough data for data field length") - return payload, nil + return nil, fmt.Errorf("not enough data for data field length") } dataLen := binary.LittleEndian.Uint32(data[offset : offset+4]) offset += 4 // Parse data field if len(data) < offset+int(dataLen) { - logger.Warn(). - Uint32("expected_len", dataLen). - Int("available", len(data)-offset). - Msg("not enough data for data field") - return payload, nil + return nil, fmt.Errorf("data field claims %d bytes, only %d available", dataLen, len(data)-offset) } if dataLen > 0 { dataField := data[offset : offset+int(dataLen)] @@ -308,18 +308,19 @@ func decodeUniversalTxEvent(data []byte, logger zerolog.Logger) (*common.Univers // Parse revert_recipient (Pubkey) if len(data) < offset+32 { - logger.Warn().Msg("not enough data for revert recipient") - return payload, nil + return nil, fmt.Errorf("not enough data for revert recipient") } revertRecipient := solana.PublicKey(data[offset : offset+32]) payload.RevertFundRecipient = revertRecipient.String() offset += 32 // Parse tx_type (TxType enum) + // + // No default. Wire 0 is GAS, which credits the sender UEA via swap rather + // than depositing to the recipient, and also selects fast confirmation. + // Guessing it on a truncated event silently changes where the money goes. if len(data) <= offset { - logger.Warn().Msg("not enough data for tx_type, defaulting to Funds") - payload.TxType = uint(0) - return payload, nil + return nil, fmt.Errorf("not enough data for tx_type") } txType := data[offset] payload.TxType = uint(txType) @@ -327,19 +328,14 @@ func decodeUniversalTxEvent(data []byte, logger zerolog.Logger) (*common.Univers // Parse signature data length (4 bytes) if len(data) < offset+4 { - logger.Warn().Msg("not enough data for signature length") - return payload, nil + return nil, fmt.Errorf("not enough data for signature length") } sigLen := binary.LittleEndian.Uint32(data[offset : offset+4]) offset += 4 remainingBytes := len(data) - offset if int(sigLen) > remainingBytes { - logger.Warn(). - Uint32("expected_len", sigLen). - Int("available", remainingBytes). - Msg("signature data length exceeds available data, skipping") - return payload, nil + return nil, fmt.Errorf("signature data claims %d bytes, only %d available", sigLen, remainingBytes) } if sigLen > 0 { @@ -369,4 +365,3 @@ func decodeUniversalTxEvent(data []byte, logger zerolog.Logger) (*common.Univers return payload, nil } - diff --git a/universalClient/chains/svm/event_parser_test.go b/universalClient/chains/svm/event_parser_test.go index da11ce0c..07371d32 100644 --- a/universalClient/chains/svm/event_parser_test.go +++ b/universalClient/chains/svm/event_parser_test.go @@ -23,18 +23,19 @@ func nopLogger() zerolog.Logger { // parseSendFundsEvent / decodeUniversalTxEvent call. // // Layout (Borsh): -// discriminator 8 bytes -// sender 32 bytes (Pubkey) -// recipient 20 bytes (byte20) -// bridge_token 32 bytes (Pubkey) -// bridge_amount 8 bytes (u64 LE) -// data_len 4 bytes (u32 LE) -// data variable -// revert_recip 32 bytes (Pubkey) -// tx_type 1 byte -// sig_len 4 bytes (u32 LE) -// sig_data variable -// fromCEA 1 byte +// +// discriminator 8 bytes +// sender 32 bytes (Pubkey) +// recipient 20 bytes (byte20) +// bridge_token 32 bytes (Pubkey) +// bridge_amount 8 bytes (u64 LE) +// data_len 4 bytes (u32 LE) +// data variable +// revert_recip 32 bytes (Pubkey) +// tx_type 1 byte +// sig_len 4 bytes (u32 LE) +// sig_data variable +// fromCEA 1 byte func buildSendFundsPayload( sender [32]byte, recipient [20]byte, @@ -116,12 +117,12 @@ func TestBase58ToHex(t *testing.T) { }, { name: "known base58 value", - input: "1", // base58 "1" decodes to a single 0x00 byte + input: "1", // base58 "1" decodes to a single 0x00 byte want: "0x00", }, { name: "known base58 multi-byte", - input: "2g", // base58 "2g" decodes to 0x61 + input: "2g", // base58 "2g" decodes to 0x61 want: "0x61", }, { @@ -345,22 +346,30 @@ func TestParseSendFundsEvent_TruncatedData(t *testing.T) { chainID := "solana:devnet" sig := "truncSig" - t.Run("data too short for sender returns event with nil EventData", func(t *testing.T) { + // A truncated event is discarded rather than stored. Every field it fails to + // reach would otherwise be left at its zero value, and tx_type 0 is GAS, + // which credits the sender UEA instead of depositing to the recipient. + t.Run("data too short for sender is discarded", func(t *testing.T) { // Only discriminator (8 bytes), no sender data := make([]byte, 8) event := ParseEvent(wrapAsLog(data), sig, 1, 0, EventTypeSendFunds, chainID, logger) - require.NotNil(t, event) - // Event is created but parseUniversalTxEvent will fail to decode, - // so EventData may be nil - assert.Equal(t, store.EventTypeInbound, event.Type) + assert.Nil(t, event) }) - t.Run("data truncated after sender still returns event", func(t *testing.T) { + t.Run("data truncated after sender is discarded", func(t *testing.T) { // 8 disc + 32 sender = 40 bytes, missing recipient data := make([]byte, 40) event := ParseEvent(wrapAsLog(data), sig, 1, 0, EventTypeSendFunds, chainID, logger) - require.NotNil(t, event) - assert.Equal(t, store.EventTypeInbound, event.Type) + assert.Nil(t, event) + }) + + t.Run("truncated one byte before tx_type is discarded", func(t *testing.T) { + // Everything through revert_recipient, then nothing. This is the exact + // shape that used to decode as TxType 0 and route to GAS. + data := make([]byte, 136) + binary.LittleEndian.PutUint32(data[100:104], 0) + event := ParseEvent(wrapAsLog(data), sig, 1, 0, EventTypeSendFunds, chainID, logger) + assert.Nil(t, event) }) } @@ -582,42 +591,76 @@ func TestDecodeUniversalTxEvent_PartialData(t *testing.T) { assert.Contains(t, err.Error(), "bridge_amount") }) - t.Run("returns partial result when no data field length", func(t *testing.T) { + // Everything through signature_data is required. Returning a partial result + // leaves tx_type at 0, which is GAS on the wire, so a truncated FUNDS + // transfer would be credited to the sender UEA instead of the recipient. + + t.Run("returns error when no data field length", func(t *testing.T) { // 8 + 32 + 20 + 32 + 8 = 100, no data_len data := make([]byte, 100) binary.LittleEndian.PutUint64(data[92:100], 777) - result, err := decodeUniversalTxEvent(data, logger) - require.NoError(t, err) - assert.Equal(t, "777", result.Amount) + _, err := decodeUniversalTxEvent(data, logger) + require.Error(t, err) + assert.Contains(t, err.Error(), "data field length") }) - t.Run("returns partial result when data field exceeds available bytes", func(t *testing.T) { + t.Run("returns error when data field exceeds available bytes", func(t *testing.T) { // 8 + 32 + 20 + 32 + 8 + 4 = 104 data := make([]byte, 104) binary.LittleEndian.PutUint64(data[92:100], 555) binary.LittleEndian.PutUint32(data[100:104], 999) // claims 999 bytes of payload - result, err := decodeUniversalTxEvent(data, logger) - require.NoError(t, err) - assert.Equal(t, "555", result.Amount) - assert.Empty(t, result.RawPayload) // not enough data, so payload is skipped + _, err := decodeUniversalTxEvent(data, logger) + require.Error(t, err) + assert.Contains(t, err.Error(), "claims 999 bytes") }) - t.Run("returns partial result when missing revert recipient", func(t *testing.T) { + t.Run("returns error when missing revert recipient", func(t *testing.T) { // 8 + 32 + 20 + 32 + 8 + 4(data_len=0) = 104 data := make([]byte, 104) binary.LittleEndian.PutUint32(data[100:104], 0) // 0 length payload - result, err := decodeUniversalTxEvent(data, logger) - require.NoError(t, err) - assert.Empty(t, result.RevertFundRecipient) + _, err := decodeUniversalTxEvent(data, logger) + require.Error(t, err) + assert.Contains(t, err.Error(), "revert recipient") }) - t.Run("returns partial result when missing tx_type", func(t *testing.T) { + t.Run("returns error when missing tx_type rather than defaulting to GAS", func(t *testing.T) { // 8 + 32 + 20 + 32 + 8 + 4(data_len=0) + 32(revert) = 136 data := make([]byte, 136) binary.LittleEndian.PutUint32(data[100:104], 0) + _, err := decodeUniversalTxEvent(data, logger) + require.Error(t, err) + assert.Contains(t, err.Error(), "tx_type") + }) + + t.Run("returns error when missing signature length", func(t *testing.T) { + // 136 + 1(tx_type) = 137, no signature length + data := make([]byte, 137) + binary.LittleEndian.PutUint32(data[100:104], 0) + data[136] = 2 // Funds + _, err := decodeUniversalTxEvent(data, logger) + require.Error(t, err) + assert.Contains(t, err.Error(), "signature length") + }) + + t.Run("returns error when signature data exceeds available bytes", func(t *testing.T) { + data := make([]byte, 141) + binary.LittleEndian.PutUint32(data[100:104], 0) + data[136] = 2 + binary.LittleEndian.PutUint32(data[137:141], 500) + _, err := decodeUniversalTxEvent(data, logger) + require.Error(t, err) + assert.Contains(t, err.Error(), "claims 500 bytes") + }) + + t.Run("from_cea stays optional and defaults to false", func(t *testing.T) { + // 141 bytes: complete through signature_data, no from_cea byte. + data := make([]byte, 141) + binary.LittleEndian.PutUint32(data[100:104], 0) + data[136] = 2 // Funds + binary.LittleEndian.PutUint32(data[137:141], 0) result, err := decodeUniversalTxEvent(data, logger) require.NoError(t, err) - // tx_type defaults to 0 when missing - assert.Equal(t, uint(0), result.TxType) + assert.Equal(t, uint(2), result.TxType) + assert.False(t, result.FromCEA) }) } diff --git a/universalClient/chains/svm/universal_tx_golden_test.go b/universalClient/chains/svm/universal_tx_golden_test.go new file mode 100644 index 00000000..3a9c6e51 --- /dev/null +++ b/universalClient/chains/svm/universal_tx_golden_test.go @@ -0,0 +1,122 @@ +package svm + +import ( + "encoding/base64" + "encoding/hex" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/pushchain/push-chain-node/universalClient/store" +) + +// Real UniversalTx events captured from the deployed devnet gateway +// CFVSincHYbETh2k7w6u1ENEkjbSLtveRCEBupKidw2VS. They pin the decoder to what +// the chain actually emits rather than to a hand-built fixture. +// +// Layout, matching the gateway on pc20-3rd-iteration: +// +// disc 8, sender 32, recipient 20, token 32, amount u64, +// payload (u32 len + bytes), revert_recipient 32, tx_type 1, +// signature_data (u32 len + bytes), from_cea 1 = 142 bytes when both vecs are empty +var devnetUniversalTxEvents = []struct { + name string + hex string + wantAmount string + wantTxType uint + wantFromCEA bool + wantConfirmDep string +}{ + { + name: "3000000 lamports, Funds", + hex: "6c9ad829b5ea1d7c5824d1bda3f79e54416ae3d2ec8d8a7456ae9a3e8a85e2f43e3bd20f25a39e5107a26674effcfbef4ee6cc6e8a00dc54801d83d90000000000000000000000000000000000000000000000000000000000000000c0c62d0000000000000000005824d1bda3f79e54416ae3d2ec8d8a7456ae9a3e8a85e2f43e3bd20f25a39e51020000000001", + wantAmount: "3000000", + wantTxType: 2, + wantFromCEA: true, + wantConfirmDep: store.ConfirmationStandard, + }, + { + name: "8000 lamports, Funds", + hex: "6c9ad829b5ea1d7cdc84c8dd7c695f0ed78f3507fd867827812dcd9ccbba61ca9a16d899b6f5ac665c70c864cf1adfb04a0e107ffa248ba3600eab8dcbcae9e66452fe98abcf0fa51e557fc8d671c7fc6ce83c05f92992a3d2bf1932401f00000000000000000000dc84c8dd7c695f0ed78f3507fd867827812dcd9ccbba61ca9a16d899b6f5ac66020000000001", + wantAmount: "8000", + wantTxType: 2, + wantFromCEA: true, + wantConfirmDep: store.ConfirmationStandard, + }, + { + name: "5000000 lamports, Funds", + hex: "6c9ad829b5ea1d7cdc84c8dd7c695f0ed78f3507fd867827812dcd9ccbba61ca9a16d899b6f5ac665c70c864cf1adfb04a0e107ffa248ba3600eab8d0000000000000000000000000000000000000000000000000000000000000000404b4c000000000000000000dc84c8dd7c695f0ed78f3507fd867827812dcd9ccbba61ca9a16d899b6f5ac66020000000001", + wantAmount: "5000000", + wantTxType: 2, + wantFromCEA: true, + wantConfirmDep: store.ConfirmationStandard, + }, +} + +func TestDecodeUniversalTxEvent_RealDevnetEvents(t *testing.T) { + logger := nopLogger() + + for _, tc := range devnetUniversalTxEvents { + t.Run(tc.name, func(t *testing.T) { + data, err := hex.DecodeString(tc.hex) + require.NoError(t, err) + require.Len(t, data, 142, "captured event is not the deployed layout") + + got, err := decodeUniversalTxEvent(data, logger) + require.NoError(t, err) + + assert.Equal(t, tc.wantAmount, got.Amount) + assert.Equal(t, tc.wantTxType, got.TxType, "tx_type must survive decoding, not be defaulted") + assert.Equal(t, tc.wantFromCEA, got.FromCEA) + assert.NotEmpty(t, got.Sender) + assert.NotEmpty(t, got.Recipient) + assert.NotEmpty(t, got.RevertFundRecipient) + }) + } +} + +// FUNDS must take the slower confirmation path. The old default of 0 selected +// FAST as well as routing to GAS, so a high value transfer lost finality too. +func TestParseSendFundsEvent_RealDevnetEventConfirmation(t *testing.T) { + logger := nopLogger() + + for _, tc := range devnetUniversalTxEvents { + t.Run(tc.name, func(t *testing.T) { + data, err := hex.DecodeString(tc.hex) + require.NoError(t, err) + log := "Program data: " + base64.StdEncoding.EncodeToString(data) + + event := ParseEvent(log, "devnetSig", 1, 0, EventTypeSendFunds, "solana:devnet", logger) + require.NotNil(t, event) + require.NotNil(t, event.EventData) + + assert.Equal(t, store.EventTypeInbound, event.Type) + assert.Equal(t, tc.wantConfirmDep, event.ConfirmationType) + }) + } +} + +// Truncating a real event anywhere past bridge_amount must be rejected. Before +// the fix each of these decoded successfully with TxType left at 0. +func TestDecodeUniversalTxEvent_TruncatedRealEventIsRejected(t *testing.T) { + logger := nopLogger() + + full, err := hex.DecodeString(devnetUniversalTxEvents[0].hex) + require.NoError(t, err) + + // 100 is the end of bridge_amount; 142 is the whole event. from_cea is the + // only optional field, so 141 is the shortest valid length. + for n := 100; n < 141; n++ { + got, err := decodeUniversalTxEvent(full[:n], logger) + require.Error(t, err, "%d-byte truncation was accepted", n) + assert.Nil(t, got) + } + + // The two valid lengths still decode, and both carry the real tx_type. + for _, n := range []int{141, 142} { + got, err := decodeUniversalTxEvent(full[:n], logger) + require.NoError(t, err, "%d-byte event was rejected", n) + assert.Equal(t, uint(2), got.TxType) + } +} From dff53e705f26cd43e69b5bb26dbf25a6488e858d Mon Sep 17 00:00:00 2001 From: aman035 Date: Mon, 24 Aug 2026 14:44:58 +0530 Subject: [PATCH 2/4] fix(common): guard event cleaner lifecycle with a mutex and wait for its goroutine --- .../chains/common/event_cleaner.go | 50 ++++-- .../chains/common/event_cleaner_race_test.go | 147 ++++++++++++++++++ .../chains/common/event_cleaner_test.go | 4 +- universalClient/chains/svm/client_test.go | 2 +- 4 files changed, 186 insertions(+), 17 deletions(-) create mode 100644 universalClient/chains/common/event_cleaner_race_test.go diff --git a/universalClient/chains/common/event_cleaner.go b/universalClient/chains/common/event_cleaner.go index b7843a42..0d89d108 100644 --- a/universalClient/chains/common/event_cleaner.go +++ b/universalClient/chains/common/event_cleaner.go @@ -3,6 +3,7 @@ package common import ( "context" "fmt" + "sync" "time" "github.com/pushchain/push-chain-node/universalClient/db" @@ -23,9 +24,14 @@ type EventCleaner struct { cleanupInterval time.Duration retentionPeriod time.Duration logger zerolog.Logger - ticker *time.Ticker - stopCh chan struct{} - running bool + + // mu guards running and stopCh, which Start and Stop both touch. The + // cleanup goroutine reads neither: it closes over its own copies, so the + // only cross-goroutine state is the channel it selects on. + mu sync.Mutex + running bool + stopCh chan struct{} + wg sync.WaitGroup } // NewEventCleaner creates a new event cleaner for a chain @@ -54,9 +60,16 @@ func NewEventCleaner( // Start begins the periodic cleanup process func (ec *EventCleaner) Start(ctx context.Context) error { + ec.mu.Lock() if ec.running { + ec.mu.Unlock() return fmt.Errorf("event cleaner is already running") } + stopCh := make(chan struct{}) + ec.running = true + ec.stopCh = stopCh + ec.wg.Add(1) + ec.mu.Unlock() ec.logger.Debug(). Str("cleanup_interval", ec.cleanupInterval.String()). @@ -69,21 +82,23 @@ func (ec *EventCleaner) Start(ctx context.Context) error { // Don't fail startup on cleanup error, just log it } - ec.running = true - ec.stopCh = make(chan struct{}) - ec.ticker = time.NewTicker(ec.cleanupInterval) + // The ticker and stop channel are the goroutine's own. Holding them on the + // struct let Stop write the fields while the goroutine was still reading + // them, which is the race this shape removes. + ticker := time.NewTicker(ec.cleanupInterval) go func() { - defer ec.ticker.Stop() + defer ec.wg.Done() + defer ticker.Stop() for { select { case <-ctx.Done(): ec.logger.Debug().Msg("context cancelled, stopping event cleaner") return - case <-ec.stopCh: + case <-stopCh: ec.logger.Debug().Msg("stop signal received, stopping event cleaner") return - case <-ec.ticker.C: + case <-ticker.C: if err := ec.performCleanup(); err != nil { ec.logger.Error().Err(err).Msg("failed to perform scheduled cleanup") } @@ -94,17 +109,24 @@ func (ec *EventCleaner) Start(ctx context.Context) error { return nil } -// Stop gracefully stops the event cleaner. No-op if not running. +// Stop gracefully stops the event cleaner and waits for the cleanup goroutine +// to exit. No-op if not running. +// +// Waiting matters on shutdown: the goroutine runs queries against the chain +// database, and returning before it finishes lets the caller close that +// database underneath an in-flight cleanup. func (ec *EventCleaner) Stop() { + ec.mu.Lock() if !ec.running { + ec.mu.Unlock() return } ec.logger.Debug().Msg("stopping event cleaner") - if ec.ticker != nil { - ec.ticker.Stop() - } - close(ec.stopCh) ec.running = false + close(ec.stopCh) + ec.mu.Unlock() + + ec.wg.Wait() } // performCleanup executes cleanup of terminal events (COMPLETED, REORGED, REVERTED) diff --git a/universalClient/chains/common/event_cleaner_race_test.go b/universalClient/chains/common/event_cleaner_race_test.go new file mode 100644 index 00000000..d28a1c10 --- /dev/null +++ b/universalClient/chains/common/event_cleaner_race_test.go @@ -0,0 +1,147 @@ +package common + +import ( + "context" + "sync" + "testing" + "time" + + "github.com/rs/zerolog" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// Start and Stop race against the cleanup goroutine. Run under -race. +func TestEventCleaner_StartStopUnderRace(t *testing.T) { + for i := 0; i < 20; i++ { + database := newTestCleanerDB(t, nil) + cleaner := NewEventCleaner(database, intPtr(3600), intPtr(0), "test-chain", zerolog.Nop()) + // Fast enough that the goroutine is inside performCleanup while Stop runs. + cleaner.cleanupInterval = time.Millisecond + + require.NoError(t, cleaner.Start(context.Background())) + cleaner.Stop() + } +} + +// Concurrent Stop calls must not double close the channel or return before the +// goroutine has exited. +func TestEventCleaner_ConcurrentStop(t *testing.T) { + database := newTestCleanerDB(t, nil) + cleaner := NewEventCleaner(database, intPtr(3600), intPtr(0), "test-chain", zerolog.Nop()) + cleaner.cleanupInterval = time.Millisecond + + require.NoError(t, cleaner.Start(context.Background())) + + var wg sync.WaitGroup + for i := 0; i < 8; i++ { + wg.Add(1) + go func() { + defer wg.Done() + cleaner.Stop() + }() + } + wg.Wait() + + assert.False(t, cleaner.running) +} + +// Concurrent Start calls must leave exactly one goroutine running. +func TestEventCleaner_ConcurrentStart(t *testing.T) { + database := newTestCleanerDB(t, nil) + cleaner := NewEventCleaner(database, intPtr(3600), intPtr(0), "test-chain", zerolog.Nop()) + cleaner.cleanupInterval = time.Millisecond + + var mu sync.Mutex + started := 0 + + var wg sync.WaitGroup + for i := 0; i < 8; i++ { + wg.Add(1) + go func() { + defer wg.Done() + if err := cleaner.Start(context.Background()); err == nil { + mu.Lock() + started++ + mu.Unlock() + } + }() + } + wg.Wait() + + assert.Equal(t, 1, started, "more than one cleanup goroutine was started") + cleaner.Stop() +} + +// Stop must not return while a cleanup is still in flight, otherwise the caller +// can close the chain database underneath an in-flight query. +// +// Held open with a write transaction so the goroutine is genuinely blocked +// inside performCleanup while Stop is called. Without that, the goroutine exits +// so fast that a Stop which does not wait looks identical to one that does. +func TestEventCleaner_StopWaitsForInFlightCleanup(t *testing.T) { + database := newTestCleanerDB(t, nil) + cleaner := NewEventCleaner(database, intPtr(3600), intPtr(0), "test-chain", zerolog.Nop()) + cleaner.cleanupInterval = time.Millisecond + + // Start first: the initial cleanup is synchronous and would block on the lock. + require.NoError(t, cleaner.Start(context.Background())) + + // Take the write lock so the next ticked cleanup blocks on DELETE. + tx := database.Client().Begin() + require.NoError(t, tx.Error) + require.NoError(t, tx.Exec( + "CREATE TABLE IF NOT EXISTS lock_probe (id INTEGER PRIMARY KEY)").Error) + require.NoError(t, tx.Exec("INSERT INTO lock_probe (id) VALUES (1)").Error) + + time.Sleep(50 * time.Millisecond) // let a tick land and block + + stopped := make(chan struct{}) + go func() { + cleaner.Stop() + close(stopped) + }() + + select { + case <-stopped: + tx.Rollback() + t.Fatal("Stop returned while a cleanup was still in flight") + case <-time.After(200 * time.Millisecond): + } + + tx.Rollback() // release the lock; the cleanup can now finish + + select { + case <-stopped: + case <-time.After(5 * time.Second): + t.Fatal("Stop did not return after the cleanup finished") + } +} + +// Cancelling the context stops the goroutine, and a later Stop is still safe. +func TestEventCleaner_ContextCancelThenStop(t *testing.T) { + database := newTestCleanerDB(t, nil) + cleaner := NewEventCleaner(database, intPtr(3600), intPtr(0), "test-chain", zerolog.Nop()) + cleaner.cleanupInterval = time.Millisecond + + ctx, cancel := context.WithCancel(context.Background()) + require.NoError(t, cleaner.Start(ctx)) + cancel() + time.Sleep(20 * time.Millisecond) + + cleaner.Stop() // must not hang or panic + assert.False(t, cleaner.running) +} + +// Restart after Stop gets a fresh channel rather than reusing the closed one. +func TestEventCleaner_RestartAfterStop(t *testing.T) { + database := newTestCleanerDB(t, nil) + cleaner := NewEventCleaner(database, intPtr(3600), intPtr(0), "test-chain", zerolog.Nop()) + cleaner.cleanupInterval = time.Millisecond + + require.NoError(t, cleaner.Start(context.Background())) + cleaner.Stop() + + require.NoError(t, cleaner.Start(context.Background()), "restart was refused") + cleaner.Stop() +} diff --git a/universalClient/chains/common/event_cleaner_test.go b/universalClient/chains/common/event_cleaner_test.go index bc28d2ea..c9701d86 100644 --- a/universalClient/chains/common/event_cleaner_test.go +++ b/universalClient/chains/common/event_cleaner_test.go @@ -75,8 +75,8 @@ func TestEventCleanerStruct(t *testing.T) { assert.Nil(t, ec.database) assert.Equal(t, time.Duration(0), ec.cleanupInterval) assert.Equal(t, time.Duration(0), ec.retentionPeriod) - assert.Nil(t, ec.ticker) assert.Nil(t, ec.stopCh) + assert.False(t, ec.running) }) } @@ -250,7 +250,7 @@ func TestEventCleanerStart(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) require.NoError(t, cleaner.Start(ctx)) - require.NotNil(t, cleaner.ticker) + require.NotNil(t, cleaner.stopCh) cancel() time.Sleep(100 * time.Millisecond) diff --git a/universalClient/chains/svm/client_test.go b/universalClient/chains/svm/client_test.go index 50f1084f..358095f7 100644 --- a/universalClient/chains/svm/client_test.go +++ b/universalClient/chains/svm/client_test.go @@ -255,7 +255,7 @@ func TestApplyDefaults_GasPriceOverride(t *testing.T) { gasPriceInterval := 60 gasPriceMarkup := 20 chainSpecific := &config.ChainSpecificConfig{ - RPCURLs: []string{"https://rpc.example.com"}, + RPCURLs: []string{"https://rpc.example.com"}, GasPriceIntervalSeconds: &gasPriceInterval, GasPriceMarkupPercent: &gasPriceMarkup, } From b090dd43145175546e54530b8d880b3accb09749 Mon Sep 17 00:00:00 2001 From: aman035 Date: Mon, 24 Aug 2026 14:45:17 +0530 Subject: [PATCH 3/4] revert unrelated gofmt change in client_test --- universalClient/chains/svm/client_test.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/universalClient/chains/svm/client_test.go b/universalClient/chains/svm/client_test.go index 358095f7..50f1084f 100644 --- a/universalClient/chains/svm/client_test.go +++ b/universalClient/chains/svm/client_test.go @@ -255,7 +255,7 @@ func TestApplyDefaults_GasPriceOverride(t *testing.T) { gasPriceInterval := 60 gasPriceMarkup := 20 chainSpecific := &config.ChainSpecificConfig{ - RPCURLs: []string{"https://rpc.example.com"}, + RPCURLs: []string{"https://rpc.example.com"}, GasPriceIntervalSeconds: &gasPriceInterval, GasPriceMarkupPercent: &gasPriceMarkup, } From bc7d968c966edc400452ef595dc0eb423ee6dccc Mon Sep 17 00:00:00 2001 From: aman035 Date: Mon, 24 Aug 2026 16:36:30 +0530 Subject: [PATCH 4/4] test: fold new tests into the existing per-source test files --- .../chains/common/event_cleaner_race_test.go | 147 ------------------ .../chains/common/event_cleaner_test.go | 136 ++++++++++++++++ .../chains/common/event_processor_test.go | 74 +++++++++ .../chains/common/txtype_routing_test.go | 86 ---------- .../chains/svm/event_parser_test.go | 110 +++++++++++++ .../chains/svm/universal_tx_golden_test.go | 122 --------------- 6 files changed, 320 insertions(+), 355 deletions(-) delete mode 100644 universalClient/chains/common/event_cleaner_race_test.go delete mode 100644 universalClient/chains/common/txtype_routing_test.go delete mode 100644 universalClient/chains/svm/universal_tx_golden_test.go diff --git a/universalClient/chains/common/event_cleaner_race_test.go b/universalClient/chains/common/event_cleaner_race_test.go deleted file mode 100644 index d28a1c10..00000000 --- a/universalClient/chains/common/event_cleaner_race_test.go +++ /dev/null @@ -1,147 +0,0 @@ -package common - -import ( - "context" - "sync" - "testing" - "time" - - "github.com/rs/zerolog" - "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/require" -) - -// Start and Stop race against the cleanup goroutine. Run under -race. -func TestEventCleaner_StartStopUnderRace(t *testing.T) { - for i := 0; i < 20; i++ { - database := newTestCleanerDB(t, nil) - cleaner := NewEventCleaner(database, intPtr(3600), intPtr(0), "test-chain", zerolog.Nop()) - // Fast enough that the goroutine is inside performCleanup while Stop runs. - cleaner.cleanupInterval = time.Millisecond - - require.NoError(t, cleaner.Start(context.Background())) - cleaner.Stop() - } -} - -// Concurrent Stop calls must not double close the channel or return before the -// goroutine has exited. -func TestEventCleaner_ConcurrentStop(t *testing.T) { - database := newTestCleanerDB(t, nil) - cleaner := NewEventCleaner(database, intPtr(3600), intPtr(0), "test-chain", zerolog.Nop()) - cleaner.cleanupInterval = time.Millisecond - - require.NoError(t, cleaner.Start(context.Background())) - - var wg sync.WaitGroup - for i := 0; i < 8; i++ { - wg.Add(1) - go func() { - defer wg.Done() - cleaner.Stop() - }() - } - wg.Wait() - - assert.False(t, cleaner.running) -} - -// Concurrent Start calls must leave exactly one goroutine running. -func TestEventCleaner_ConcurrentStart(t *testing.T) { - database := newTestCleanerDB(t, nil) - cleaner := NewEventCleaner(database, intPtr(3600), intPtr(0), "test-chain", zerolog.Nop()) - cleaner.cleanupInterval = time.Millisecond - - var mu sync.Mutex - started := 0 - - var wg sync.WaitGroup - for i := 0; i < 8; i++ { - wg.Add(1) - go func() { - defer wg.Done() - if err := cleaner.Start(context.Background()); err == nil { - mu.Lock() - started++ - mu.Unlock() - } - }() - } - wg.Wait() - - assert.Equal(t, 1, started, "more than one cleanup goroutine was started") - cleaner.Stop() -} - -// Stop must not return while a cleanup is still in flight, otherwise the caller -// can close the chain database underneath an in-flight query. -// -// Held open with a write transaction so the goroutine is genuinely blocked -// inside performCleanup while Stop is called. Without that, the goroutine exits -// so fast that a Stop which does not wait looks identical to one that does. -func TestEventCleaner_StopWaitsForInFlightCleanup(t *testing.T) { - database := newTestCleanerDB(t, nil) - cleaner := NewEventCleaner(database, intPtr(3600), intPtr(0), "test-chain", zerolog.Nop()) - cleaner.cleanupInterval = time.Millisecond - - // Start first: the initial cleanup is synchronous and would block on the lock. - require.NoError(t, cleaner.Start(context.Background())) - - // Take the write lock so the next ticked cleanup blocks on DELETE. - tx := database.Client().Begin() - require.NoError(t, tx.Error) - require.NoError(t, tx.Exec( - "CREATE TABLE IF NOT EXISTS lock_probe (id INTEGER PRIMARY KEY)").Error) - require.NoError(t, tx.Exec("INSERT INTO lock_probe (id) VALUES (1)").Error) - - time.Sleep(50 * time.Millisecond) // let a tick land and block - - stopped := make(chan struct{}) - go func() { - cleaner.Stop() - close(stopped) - }() - - select { - case <-stopped: - tx.Rollback() - t.Fatal("Stop returned while a cleanup was still in flight") - case <-time.After(200 * time.Millisecond): - } - - tx.Rollback() // release the lock; the cleanup can now finish - - select { - case <-stopped: - case <-time.After(5 * time.Second): - t.Fatal("Stop did not return after the cleanup finished") - } -} - -// Cancelling the context stops the goroutine, and a later Stop is still safe. -func TestEventCleaner_ContextCancelThenStop(t *testing.T) { - database := newTestCleanerDB(t, nil) - cleaner := NewEventCleaner(database, intPtr(3600), intPtr(0), "test-chain", zerolog.Nop()) - cleaner.cleanupInterval = time.Millisecond - - ctx, cancel := context.WithCancel(context.Background()) - require.NoError(t, cleaner.Start(ctx)) - cancel() - time.Sleep(20 * time.Millisecond) - - cleaner.Stop() // must not hang or panic - assert.False(t, cleaner.running) -} - -// Restart after Stop gets a fresh channel rather than reusing the closed one. -func TestEventCleaner_RestartAfterStop(t *testing.T) { - database := newTestCleanerDB(t, nil) - cleaner := NewEventCleaner(database, intPtr(3600), intPtr(0), "test-chain", zerolog.Nop()) - cleaner.cleanupInterval = time.Millisecond - - require.NoError(t, cleaner.Start(context.Background())) - cleaner.Stop() - - require.NoError(t, cleaner.Start(context.Background()), "restart was refused") - cleaner.Stop() -} diff --git a/universalClient/chains/common/event_cleaner_test.go b/universalClient/chains/common/event_cleaner_test.go index c9701d86..15eee2be 100644 --- a/universalClient/chains/common/event_cleaner_test.go +++ b/universalClient/chains/common/event_cleaner_test.go @@ -3,6 +3,7 @@ package common import ( "context" "fmt" + "sync" "testing" "time" @@ -337,3 +338,138 @@ func TestEventCleanerStartStopLifecycle(t *testing.T) { time.Sleep(50 * time.Millisecond) }) } + +// Start and Stop race against the cleanup goroutine. Run under -race. +func TestEventCleaner_StartStopUnderRace(t *testing.T) { + for i := 0; i < 20; i++ { + database := newTestCleanerDB(t, nil) + cleaner := NewEventCleaner(database, intPtr(3600), intPtr(0), "test-chain", zerolog.Nop()) + // Fast enough that the goroutine is inside performCleanup while Stop runs. + cleaner.cleanupInterval = time.Millisecond + + require.NoError(t, cleaner.Start(context.Background())) + cleaner.Stop() + } +} + +// Concurrent Stop calls must not double close the channel or return before the +// goroutine has exited. +func TestEventCleaner_ConcurrentStop(t *testing.T) { + database := newTestCleanerDB(t, nil) + cleaner := NewEventCleaner(database, intPtr(3600), intPtr(0), "test-chain", zerolog.Nop()) + cleaner.cleanupInterval = time.Millisecond + + require.NoError(t, cleaner.Start(context.Background())) + + var wg sync.WaitGroup + for i := 0; i < 8; i++ { + wg.Add(1) + go func() { + defer wg.Done() + cleaner.Stop() + }() + } + wg.Wait() + + assert.False(t, cleaner.running) +} + +// Concurrent Start calls must leave exactly one goroutine running. +func TestEventCleaner_ConcurrentStart(t *testing.T) { + database := newTestCleanerDB(t, nil) + cleaner := NewEventCleaner(database, intPtr(3600), intPtr(0), "test-chain", zerolog.Nop()) + cleaner.cleanupInterval = time.Millisecond + + var mu sync.Mutex + started := 0 + + var wg sync.WaitGroup + for i := 0; i < 8; i++ { + wg.Add(1) + go func() { + defer wg.Done() + if err := cleaner.Start(context.Background()); err == nil { + mu.Lock() + started++ + mu.Unlock() + } + }() + } + wg.Wait() + + assert.Equal(t, 1, started, "more than one cleanup goroutine was started") + cleaner.Stop() +} + +// Stop must not return while a cleanup is still in flight, otherwise the caller +// can close the chain database underneath an in-flight query. +// +// Held open with a write transaction so the goroutine is genuinely blocked +// inside performCleanup while Stop is called. Without that, the goroutine exits +// so fast that a Stop which does not wait looks identical to one that does. +func TestEventCleaner_StopWaitsForInFlightCleanup(t *testing.T) { + database := newTestCleanerDB(t, nil) + cleaner := NewEventCleaner(database, intPtr(3600), intPtr(0), "test-chain", zerolog.Nop()) + cleaner.cleanupInterval = time.Millisecond + + // Start first: the initial cleanup is synchronous and would block on the lock. + require.NoError(t, cleaner.Start(context.Background())) + + // Take the write lock so the next ticked cleanup blocks on DELETE. + tx := database.Client().Begin() + require.NoError(t, tx.Error) + require.NoError(t, tx.Exec( + "CREATE TABLE IF NOT EXISTS lock_probe (id INTEGER PRIMARY KEY)").Error) + require.NoError(t, tx.Exec("INSERT INTO lock_probe (id) VALUES (1)").Error) + + time.Sleep(50 * time.Millisecond) // let a tick land and block + + stopped := make(chan struct{}) + go func() { + cleaner.Stop() + close(stopped) + }() + + select { + case <-stopped: + tx.Rollback() + t.Fatal("Stop returned while a cleanup was still in flight") + case <-time.After(200 * time.Millisecond): + } + + tx.Rollback() // release the lock; the cleanup can now finish + + select { + case <-stopped: + case <-time.After(5 * time.Second): + t.Fatal("Stop did not return after the cleanup finished") + } +} + +// Cancelling the context stops the goroutine, and a later Stop is still safe. +func TestEventCleaner_ContextCancelThenStop(t *testing.T) { + database := newTestCleanerDB(t, nil) + cleaner := NewEventCleaner(database, intPtr(3600), intPtr(0), "test-chain", zerolog.Nop()) + cleaner.cleanupInterval = time.Millisecond + + ctx, cancel := context.WithCancel(context.Background()) + require.NoError(t, cleaner.Start(ctx)) + cancel() + time.Sleep(20 * time.Millisecond) + + cleaner.Stop() // must not hang or panic + assert.False(t, cleaner.running) +} + +// Restart after Stop gets a fresh channel rather than reusing the closed one. +func TestEventCleaner_RestartAfterStop(t *testing.T) { + database := newTestCleanerDB(t, nil) + cleaner := NewEventCleaner(database, intPtr(3600), intPtr(0), "test-chain", zerolog.Nop()) + cleaner.cleanupInterval = time.Millisecond + + require.NoError(t, cleaner.Start(context.Background())) + cleaner.Stop() + + require.NoError(t, cleaner.Start(context.Background()), "restart was refused") + cleaner.Stop() +} diff --git a/universalClient/chains/common/event_processor_test.go b/universalClient/chains/common/event_processor_test.go index 4a08c305..bd300c48 100644 --- a/universalClient/chains/common/event_processor_test.go +++ b/universalClient/chains/common/event_processor_test.go @@ -1054,3 +1054,77 @@ func TestProcessConfirmedEventsEnabledFlags(t *testing.T) { assert.Equal(t, store.StatusConfirmed, inboundEvt.Status) }) } + +// The wire values the gateways emit are 0-indexed (Gas, GasAndPayload, Funds, +// FundsAndPayload) while the chain enum reserves 0 for UNSPECIFIED, so the +// mapping is shifted by one. A decoder that leaves TxType unset therefore does +// not produce "unknown", it produces GAS. +func TestConstructInbound_TxTypeMapping(t *testing.T) { + processor := &EventProcessor{} + + for _, tc := range []struct { + wire uint + want uexecutortypes.TxType + }{ + {0, uexecutortypes.TxType_GAS}, + {1, uexecutortypes.TxType_GAS_AND_PAYLOAD}, + {2, uexecutortypes.TxType_FUNDS}, + {3, uexecutortypes.TxType_FUNDS_AND_PAYLOAD}, + {4, uexecutortypes.TxType_UNSPECIFIED_TX}, + {99, uexecutortypes.TxType_UNSPECIFIED_TX}, + } { + data, err := json.Marshal(UniversalTx{ + SourceChain: "solana:devnet", + Sender: "0xabc", + Recipient: "0xdef", + Amount: "5000000", + TxType: tc.wire, + }) + require.NoError(t, err) + + inbound, err := processor.constructInbound(&store.Event{ + EventID: "sig:0", + EventData: data, + }) + require.NoError(t, err) + assert.Equal(t, tc.want, inbound.TxType, "wire value %d", tc.wire) + } +} + +// A FUNDS transfer must never reach the keeper as GAS. The two dispatch to +// different handlers: GAS mints and autoswaps into the sender UEA, FUNDS +// deposits PRC20 to the recipient, so the same amount lands with a different +// party. This is the end to end assertion the finding asks for. +func TestConstructInbound_FundsNeverBecomesGas(t *testing.T) { + processor := &EventProcessor{} + + data, err := json.Marshal(UniversalTx{ + SourceChain: "solana:devnet", + Sender: "0xabc", + Recipient: "0xdef", + Amount: "5000000", + TxType: 2, // Funds, as the real devnet events carry + }) + require.NoError(t, err) + + inbound, err := processor.constructInbound(&store.Event{ + EventID: "sig:0", + EventData: data, + }) + require.NoError(t, err) + + assert.Equal(t, uexecutortypes.TxType_FUNDS, inbound.TxType) + assert.NotEqual(t, uexecutortypes.TxType_GAS, inbound.TxType, + "a FUNDS transfer routed to GAS credits the sender instead of the recipient") +} + +// An event whose data never made it past the decoder must be refused outright +// rather than defaulted. The parsers now discard such events, so this is the +// backstop if one ever reaches the store. +func TestConstructInbound_RejectsEventWithoutData(t *testing.T) { + processor := &EventProcessor{} + + _, err := processor.constructInbound(&store.Event{EventID: "sig:0"}) + require.Error(t, err) + assert.Contains(t, err.Error(), "event data is missing") +} diff --git a/universalClient/chains/common/txtype_routing_test.go b/universalClient/chains/common/txtype_routing_test.go deleted file mode 100644 index 65426b8b..00000000 --- a/universalClient/chains/common/txtype_routing_test.go +++ /dev/null @@ -1,86 +0,0 @@ -package common - -import ( - "encoding/json" - "testing" - - "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/require" - - "github.com/pushchain/push-chain-node/universalClient/store" - uexecutortypes "github.com/pushchain/push-chain-node/x/uexecutor/types" -) - -// The wire values the gateways emit are 0-indexed (Gas, GasAndPayload, Funds, -// FundsAndPayload) while the chain enum reserves 0 for UNSPECIFIED, so the -// mapping is shifted by one. A decoder that leaves TxType unset therefore does -// not produce "unknown", it produces GAS. -func TestConstructInbound_TxTypeMapping(t *testing.T) { - processor := &EventProcessor{} - - for _, tc := range []struct { - wire uint - want uexecutortypes.TxType - }{ - {0, uexecutortypes.TxType_GAS}, - {1, uexecutortypes.TxType_GAS_AND_PAYLOAD}, - {2, uexecutortypes.TxType_FUNDS}, - {3, uexecutortypes.TxType_FUNDS_AND_PAYLOAD}, - {4, uexecutortypes.TxType_UNSPECIFIED_TX}, - {99, uexecutortypes.TxType_UNSPECIFIED_TX}, - } { - data, err := json.Marshal(UniversalTx{ - SourceChain: "solana:devnet", - Sender: "0xabc", - Recipient: "0xdef", - Amount: "5000000", - TxType: tc.wire, - }) - require.NoError(t, err) - - inbound, err := processor.constructInbound(&store.Event{ - EventID: "sig:0", - EventData: data, - }) - require.NoError(t, err) - assert.Equal(t, tc.want, inbound.TxType, "wire value %d", tc.wire) - } -} - -// A FUNDS transfer must never reach the keeper as GAS. The two dispatch to -// different handlers: GAS mints and autoswaps into the sender UEA, FUNDS -// deposits PRC20 to the recipient, so the same amount lands with a different -// party. This is the end to end assertion the finding asks for. -func TestConstructInbound_FundsNeverBecomesGas(t *testing.T) { - processor := &EventProcessor{} - - data, err := json.Marshal(UniversalTx{ - SourceChain: "solana:devnet", - Sender: "0xabc", - Recipient: "0xdef", - Amount: "5000000", - TxType: 2, // Funds, as the real devnet events carry - }) - require.NoError(t, err) - - inbound, err := processor.constructInbound(&store.Event{ - EventID: "sig:0", - EventData: data, - }) - require.NoError(t, err) - - assert.Equal(t, uexecutortypes.TxType_FUNDS, inbound.TxType) - assert.NotEqual(t, uexecutortypes.TxType_GAS, inbound.TxType, - "a FUNDS transfer routed to GAS credits the sender instead of the recipient") -} - -// An event whose data never made it past the decoder must be refused outright -// rather than defaulted. The parsers now discard such events, so this is the -// backstop if one ever reaches the store. -func TestConstructInbound_RejectsEventWithoutData(t *testing.T) { - processor := &EventProcessor{} - - _, err := processor.constructInbound(&store.Event{EventID: "sig:0"}) - require.Error(t, err) - assert.Contains(t, err.Error(), "event data is missing") -} diff --git a/universalClient/chains/svm/event_parser_test.go b/universalClient/chains/svm/event_parser_test.go index 07371d32..e8630f40 100644 --- a/universalClient/chains/svm/event_parser_test.go +++ b/universalClient/chains/svm/event_parser_test.go @@ -664,3 +664,113 @@ func TestDecodeUniversalTxEvent_PartialData(t *testing.T) { assert.False(t, result.FromCEA) }) } + +// Real UniversalTx events captured from the deployed devnet gateway +// CFVSincHYbETh2k7w6u1ENEkjbSLtveRCEBupKidw2VS. They pin the decoder to what +// the chain actually emits rather than to a hand-built fixture. +// +// Layout, matching the gateway on pc20-3rd-iteration: +// +// disc 8, sender 32, recipient 20, token 32, amount u64, +// payload (u32 len + bytes), revert_recipient 32, tx_type 1, +// signature_data (u32 len + bytes), from_cea 1 = 142 bytes when both vecs are empty +var devnetUniversalTxEvents = []struct { + name string + hex string + wantAmount string + wantTxType uint + wantFromCEA bool + wantConfirmDep string +}{ + { + name: "3000000 lamports, Funds", + hex: "6c9ad829b5ea1d7c5824d1bda3f79e54416ae3d2ec8d8a7456ae9a3e8a85e2f43e3bd20f25a39e5107a26674effcfbef4ee6cc6e8a00dc54801d83d90000000000000000000000000000000000000000000000000000000000000000c0c62d0000000000000000005824d1bda3f79e54416ae3d2ec8d8a7456ae9a3e8a85e2f43e3bd20f25a39e51020000000001", + wantAmount: "3000000", + wantTxType: 2, + wantFromCEA: true, + wantConfirmDep: store.ConfirmationStandard, + }, + { + name: "8000 lamports, Funds", + hex: "6c9ad829b5ea1d7cdc84c8dd7c695f0ed78f3507fd867827812dcd9ccbba61ca9a16d899b6f5ac665c70c864cf1adfb04a0e107ffa248ba3600eab8dcbcae9e66452fe98abcf0fa51e557fc8d671c7fc6ce83c05f92992a3d2bf1932401f00000000000000000000dc84c8dd7c695f0ed78f3507fd867827812dcd9ccbba61ca9a16d899b6f5ac66020000000001", + wantAmount: "8000", + wantTxType: 2, + wantFromCEA: true, + wantConfirmDep: store.ConfirmationStandard, + }, + { + name: "5000000 lamports, Funds", + hex: "6c9ad829b5ea1d7cdc84c8dd7c695f0ed78f3507fd867827812dcd9ccbba61ca9a16d899b6f5ac665c70c864cf1adfb04a0e107ffa248ba3600eab8d0000000000000000000000000000000000000000000000000000000000000000404b4c000000000000000000dc84c8dd7c695f0ed78f3507fd867827812dcd9ccbba61ca9a16d899b6f5ac66020000000001", + wantAmount: "5000000", + wantTxType: 2, + wantFromCEA: true, + wantConfirmDep: store.ConfirmationStandard, + }, +} + +func TestDecodeUniversalTxEvent_RealDevnetEvents(t *testing.T) { + logger := nopLogger() + + for _, tc := range devnetUniversalTxEvents { + t.Run(tc.name, func(t *testing.T) { + data, err := hex.DecodeString(tc.hex) + require.NoError(t, err) + require.Len(t, data, 142, "captured event is not the deployed layout") + + got, err := decodeUniversalTxEvent(data, logger) + require.NoError(t, err) + + assert.Equal(t, tc.wantAmount, got.Amount) + assert.Equal(t, tc.wantTxType, got.TxType, "tx_type must survive decoding, not be defaulted") + assert.Equal(t, tc.wantFromCEA, got.FromCEA) + assert.NotEmpty(t, got.Sender) + assert.NotEmpty(t, got.Recipient) + assert.NotEmpty(t, got.RevertFundRecipient) + }) + } +} + +// FUNDS must take the slower confirmation path. The old default of 0 selected +// FAST as well as routing to GAS, so a high value transfer lost finality too. +func TestParseSendFundsEvent_RealDevnetEventConfirmation(t *testing.T) { + logger := nopLogger() + + for _, tc := range devnetUniversalTxEvents { + t.Run(tc.name, func(t *testing.T) { + data, err := hex.DecodeString(tc.hex) + require.NoError(t, err) + log := "Program data: " + base64.StdEncoding.EncodeToString(data) + + event := ParseEvent(log, "devnetSig", 1, 0, EventTypeSendFunds, "solana:devnet", logger) + require.NotNil(t, event) + require.NotNil(t, event.EventData) + + assert.Equal(t, store.EventTypeInbound, event.Type) + assert.Equal(t, tc.wantConfirmDep, event.ConfirmationType) + }) + } +} + +// Truncating a real event anywhere past bridge_amount must be rejected. Before +// the fix each of these decoded successfully with TxType left at 0. +func TestDecodeUniversalTxEvent_TruncatedRealEventIsRejected(t *testing.T) { + logger := nopLogger() + + full, err := hex.DecodeString(devnetUniversalTxEvents[0].hex) + require.NoError(t, err) + + // 100 is the end of bridge_amount; 142 is the whole event. from_cea is the + // only optional field, so 141 is the shortest valid length. + for n := 100; n < 141; n++ { + got, err := decodeUniversalTxEvent(full[:n], logger) + require.Error(t, err, "%d-byte truncation was accepted", n) + assert.Nil(t, got) + } + + // The two valid lengths still decode, and both carry the real tx_type. + for _, n := range []int{141, 142} { + got, err := decodeUniversalTxEvent(full[:n], logger) + require.NoError(t, err, "%d-byte event was rejected", n) + assert.Equal(t, uint(2), got.TxType) + } +} diff --git a/universalClient/chains/svm/universal_tx_golden_test.go b/universalClient/chains/svm/universal_tx_golden_test.go deleted file mode 100644 index 3a9c6e51..00000000 --- a/universalClient/chains/svm/universal_tx_golden_test.go +++ /dev/null @@ -1,122 +0,0 @@ -package svm - -import ( - "encoding/base64" - "encoding/hex" - "testing" - - "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/require" - - "github.com/pushchain/push-chain-node/universalClient/store" -) - -// Real UniversalTx events captured from the deployed devnet gateway -// CFVSincHYbETh2k7w6u1ENEkjbSLtveRCEBupKidw2VS. They pin the decoder to what -// the chain actually emits rather than to a hand-built fixture. -// -// Layout, matching the gateway on pc20-3rd-iteration: -// -// disc 8, sender 32, recipient 20, token 32, amount u64, -// payload (u32 len + bytes), revert_recipient 32, tx_type 1, -// signature_data (u32 len + bytes), from_cea 1 = 142 bytes when both vecs are empty -var devnetUniversalTxEvents = []struct { - name string - hex string - wantAmount string - wantTxType uint - wantFromCEA bool - wantConfirmDep string -}{ - { - name: "3000000 lamports, Funds", - hex: "6c9ad829b5ea1d7c5824d1bda3f79e54416ae3d2ec8d8a7456ae9a3e8a85e2f43e3bd20f25a39e5107a26674effcfbef4ee6cc6e8a00dc54801d83d90000000000000000000000000000000000000000000000000000000000000000c0c62d0000000000000000005824d1bda3f79e54416ae3d2ec8d8a7456ae9a3e8a85e2f43e3bd20f25a39e51020000000001", - wantAmount: "3000000", - wantTxType: 2, - wantFromCEA: true, - wantConfirmDep: store.ConfirmationStandard, - }, - { - name: "8000 lamports, Funds", - hex: "6c9ad829b5ea1d7cdc84c8dd7c695f0ed78f3507fd867827812dcd9ccbba61ca9a16d899b6f5ac665c70c864cf1adfb04a0e107ffa248ba3600eab8dcbcae9e66452fe98abcf0fa51e557fc8d671c7fc6ce83c05f92992a3d2bf1932401f00000000000000000000dc84c8dd7c695f0ed78f3507fd867827812dcd9ccbba61ca9a16d899b6f5ac66020000000001", - wantAmount: "8000", - wantTxType: 2, - wantFromCEA: true, - wantConfirmDep: store.ConfirmationStandard, - }, - { - name: "5000000 lamports, Funds", - hex: "6c9ad829b5ea1d7cdc84c8dd7c695f0ed78f3507fd867827812dcd9ccbba61ca9a16d899b6f5ac665c70c864cf1adfb04a0e107ffa248ba3600eab8d0000000000000000000000000000000000000000000000000000000000000000404b4c000000000000000000dc84c8dd7c695f0ed78f3507fd867827812dcd9ccbba61ca9a16d899b6f5ac66020000000001", - wantAmount: "5000000", - wantTxType: 2, - wantFromCEA: true, - wantConfirmDep: store.ConfirmationStandard, - }, -} - -func TestDecodeUniversalTxEvent_RealDevnetEvents(t *testing.T) { - logger := nopLogger() - - for _, tc := range devnetUniversalTxEvents { - t.Run(tc.name, func(t *testing.T) { - data, err := hex.DecodeString(tc.hex) - require.NoError(t, err) - require.Len(t, data, 142, "captured event is not the deployed layout") - - got, err := decodeUniversalTxEvent(data, logger) - require.NoError(t, err) - - assert.Equal(t, tc.wantAmount, got.Amount) - assert.Equal(t, tc.wantTxType, got.TxType, "tx_type must survive decoding, not be defaulted") - assert.Equal(t, tc.wantFromCEA, got.FromCEA) - assert.NotEmpty(t, got.Sender) - assert.NotEmpty(t, got.Recipient) - assert.NotEmpty(t, got.RevertFundRecipient) - }) - } -} - -// FUNDS must take the slower confirmation path. The old default of 0 selected -// FAST as well as routing to GAS, so a high value transfer lost finality too. -func TestParseSendFundsEvent_RealDevnetEventConfirmation(t *testing.T) { - logger := nopLogger() - - for _, tc := range devnetUniversalTxEvents { - t.Run(tc.name, func(t *testing.T) { - data, err := hex.DecodeString(tc.hex) - require.NoError(t, err) - log := "Program data: " + base64.StdEncoding.EncodeToString(data) - - event := ParseEvent(log, "devnetSig", 1, 0, EventTypeSendFunds, "solana:devnet", logger) - require.NotNil(t, event) - require.NotNil(t, event.EventData) - - assert.Equal(t, store.EventTypeInbound, event.Type) - assert.Equal(t, tc.wantConfirmDep, event.ConfirmationType) - }) - } -} - -// Truncating a real event anywhere past bridge_amount must be rejected. Before -// the fix each of these decoded successfully with TxType left at 0. -func TestDecodeUniversalTxEvent_TruncatedRealEventIsRejected(t *testing.T) { - logger := nopLogger() - - full, err := hex.DecodeString(devnetUniversalTxEvents[0].hex) - require.NoError(t, err) - - // 100 is the end of bridge_amount; 142 is the whole event. from_cea is the - // only optional field, so 141 is the shortest valid length. - for n := 100; n < 141; n++ { - got, err := decodeUniversalTxEvent(full[:n], logger) - require.Error(t, err, "%d-byte truncation was accepted", n) - assert.Nil(t, got) - } - - // The two valid lengths still decode, and both carry the real tx_type. - for _, n := range []int{141, 142} { - got, err := decodeUniversalTxEvent(full[:n], logger) - require.NoError(t, err, "%d-byte event was rejected", n) - assert.Equal(t, uint(2), got.TxType) - } -}