From 940e005e38b918305207f9687ae8be3a027b9aee Mon Sep 17 00:00:00 2001 From: Martin Sander Date: Thu, 6 Aug 2026 00:02:06 -0500 Subject: [PATCH 1/5] qa: add retransmit metro test to qa.testnet --- e2e/internal/qa/client_settlement.go | 45 ++++++++++ e2e/qa_multicast_settlement_test.go | 87 ++++++++++++++++--- e2e/qa_retransmit_only_settlement_test.go | 100 +++++++++++----------- e2e/qa_shred_settlement_test.go | 14 +++ sdk/shreds/go/state.go | 11 ++- 5 files changed, 194 insertions(+), 63 deletions(-) diff --git a/e2e/internal/qa/client_settlement.go b/e2e/internal/qa/client_settlement.go index 64a809ee75..74085cff70 100644 --- a/e2e/internal/qa/client_settlement.go +++ b/e2e/internal/qa/client_settlement.go @@ -190,6 +190,38 @@ func (c *Client) ClosestRetransmitOnlyDevice(ctx context.Context) (*Device, map[ return bestDevice, retransmitOnly, nil } +func (c *Client) ClosestNonRetransmitOnlyDevice(ctx context.Context) (*Device, error) { + retransmitOnly, err := c.RetransmitOnlyExchangeKeys(ctx) + if err != nil { + return nil, err + } + + latencies, err := c.GetLatency(ctx) + if err != nil { + return nil, fmt.Errorf("failed to get latency on host %s: %w", c.Host, err) + } + + var bestDevice *Device + var bestAvg uint64 = math.MaxUint64 + for _, l := range latencies { + if !l.Reachable { + continue + } + device, ok := c.devices[l.DeviceCode] + if !ok || retransmitOnly[device.ExchangePubKey] { + continue + } + if l.AvgLatencyNs < bestAvg { + bestAvg = l.AvgLatencyNs + bestDevice = device + } + } + if bestDevice != nil { + c.log.Debug("Determined closest non-retransmit-only device", "host", c.Host, "deviceCode", bestDevice.Code, "avgLatencyNs", bestAvg) + } + return bestDevice, nil +} + // FeedSeatPrice calls the FeedSeatPrice RPC to query seat pricing for a single // device (by pubkey). Querying by pubkey avoids device-code resolution, which // the CLI refuses when it can't classify the cluster (e.g. a private Solana @@ -769,6 +801,19 @@ func (c *Client) IsSeatProratingEnabled(ctx context.Context) (bool, error) { return cfg.IsProratedServiceEnabled(), nil } +func (c *Client) IsRetransmitOnlyOnboardingEnforced(ctx context.Context) (bool, error) { + programID, err := solana.PublicKeyFromBase58(c.ShredSubscriptionProgramID) + if err != nil { + return false, fmt.Errorf("failed to parse shred subscription program ID %q: %w", c.ShredSubscriptionProgramID, err) + } + + cfg, err := c.shredsClient(programID).FetchProgramConfig(ctx) + if err != nil { + return false, fmt.Errorf("failed to fetch program config on host %s: %w", c.Host, err) + } + return cfg.IsRetransmitOnlyOnboardingEnforced(), nil +} + // IsProgramPaused returns true if the shred-subscription program config has // the paused flag set. While paused, the oracle cannot ack instant seat // allocation requests, which leaves the seat un-withdrawable. diff --git a/e2e/qa_multicast_settlement_test.go b/e2e/qa_multicast_settlement_test.go index 0e36a71a4f..56f96a8019 100644 --- a/e2e/qa_multicast_settlement_test.go +++ b/e2e/qa_multicast_settlement_test.go @@ -6,6 +6,7 @@ import ( "context" "flag" "log/slog" + "strconv" "testing" "github.com/malbeclabs/doublezero/e2e/internal/qa" @@ -14,20 +15,48 @@ import ( var enableSettlementTests = flag.Bool("enable-multicast-settlement-tests", false, "enable multicast settlement tests") -// TestQA_MulticastSettlement pays for a full (leader) multicast seat on the -// closest device and verifies pricing, tunnel-up, and the withdraw/refund -// accounting. It is a thin wrapper over the shared settlement flow in -// runShredSettlement; TestQA_RetransmitOnlySettlement mirrors it with -// retransmit-only device selection and a group-subscription assertion. Any fix -// to the settlement machinery belongs in runShredSettlement so both stay in -// sync. +// TestQA_MulticastSettlement pays for a multicast seat and verifies pricing, +// tunnel-up, and the withdraw/refund accounting. It is a thin wrapper over the +// shared settlement flow in runShredSettlement; TestQA_RetransmitOnlySettlement +// mirrors it with retransmit-only device selection and a group-subscription +// assertion. Any fix to the settlement machinery belongs in runShredSettlement +// so both stay in sync. func TestQA_MulticastSettlement(t *testing.T) { + var onboardingEnforced bool runShredSettlement(t, shredSettlementParams{ - enabled: *enableSettlementTests, - skipReason: "Skipping: --enable-multicast-settlement-tests flag not set", - selectSubtestName: "find_closest_device", - selectDevice: selectClosestDevice, - priceLogMsg: "Found epoch price", + enabled: *enableSettlementTests, + skipReason: "Skipping: --enable-multicast-settlement-tests flag not set", + + preflightSubtestName: "reject_new_seat_outside_retransmit_only_metro", + preflight: func(t *testing.T, ctx context.Context, log *slog.Logger, _ *qa.Test, client *qa.Client) { + var err error + onboardingEnforced, err = client.IsRetransmitOnlyOnboardingEnforced(ctx) + require.NoError(t, err, "failed to read the retransmit-only onboarding flag") + if !onboardingEnforced { + log.Info("Retransmit-only onboarding is off; settling on the closest device") + t.Skip("Skipping: the program config does not enforce retransmit-only onboarding") + } + log.Info("Retransmit-only onboarding is on; a new seat must be rejected outside a retransmit-only metro") + assertNewSeatRejected(t, ctx, log, client) + }, + + selectSubtestName: "select_device", + selectDevice: func(t *testing.T, ctx context.Context, log *slog.Logger, test *qa.Test, client *qa.Client) *qa.Device { + if onboardingEnforced { + return selectRetransmitOnlyDevice(t, ctx, log, test, client) + } + return selectClosestDevice(t, ctx, log, test, client) + }, + + priceLogMsg: "Found epoch price", + + extraSubtestName: "assert_subscribed_groups", + extraAssertion: func(t *testing.T, ctx context.Context, log *slog.Logger, client *qa.Client, device *qa.Device) { + if !onboardingEnforced { + t.Skip("Skipping: the seat is not in a retransmit-only metro, so it carries the leader group too") + } + assertSubscribedGroups(t, ctx, log, client, device) + }, }) } @@ -39,3 +68,37 @@ func selectClosestDevice(t *testing.T, ctx context.Context, log *slog.Logger, _ log.Info("Closest device", "code", device.Code, "pubkey", device.PubKey) return device } + +// assertNewSeatRejected pays for a seat on the closest device outside a +// retransmit-only metro and requires the program to reject it. +func assertNewSeatRejected(t *testing.T, ctx context.Context, log *slog.Logger, client *qa.Client) { + device, err := client.ClosestNonRetransmitOnlyDevice(ctx) + require.NoError(t, err, "failed to find a device outside a retransmit-only metro") + if device == nil { + t.Skip("Skipping: every reachable metro is flagged retransmit-only, so no metro rejects a new seat") + } + + prices, err := client.SeatPrices(ctx, device.PubKey) + require.NoError(t, err, "failed to read the onchain seat price for device %s", device.Code) + require.NotZero(t, prices.InstantAllocationDollars, "onchain instant-allocation price is zero for device %s", device.Code) + amount := strconv.FormatUint(prices.InstantAllocationDollars, 10) + + log.Info("Paying for a seat outside a retransmit-only metro", "device", device.Code, + "metro", device.ExchangeCode, "amount", amount) + err = client.FeedSeatPay(ctx, device.PubKey, amount) + if err == nil { + // The settlement flow's cleanup only withdraws the device it selects, so + // withdraw this seat here or it stays active onchain. + if withdrawErr := client.WithdrawSeatWithRetry(ctx, device.PubKey); withdrawErr != nil { + log.Warn("Cleanup: withdraw of the wrongly admitted seat failed; seat left active onchain", + "device", device.Code, "error", withdrawErr) + } + } + require.Error(t, err, "the program admitted a new seat on device %s in metro %s, which is not retransmit-only", + device.Code, device.ExchangeCode) + require.ErrorContains(t, err, "Retransmit-only onboarding enforced: cannot fund seat", + "the payment failed for another reason than retransmit-only onboarding enforcement") + require.ErrorContains(t, err, "with no tenure", + "the payment failed for another reason than retransmit-only onboarding enforcement") + log.Info("The program rejected the new seat", "device", device.Code, "metro", device.ExchangeCode) +} diff --git a/e2e/qa_retransmit_only_settlement_test.go b/e2e/qa_retransmit_only_settlement_test.go index 1459f62589..19169cc378 100644 --- a/e2e/qa_retransmit_only_settlement_test.go +++ b/e2e/qa_retransmit_only_settlement_test.go @@ -6,12 +6,12 @@ import ( "context" "flag" "log/slog" + "strings" "testing" "time" "github.com/gagliardetto/solana-go" "github.com/malbeclabs/doublezero/e2e/internal/qa" - serviceability "github.com/malbeclabs/doublezero/smartcontract/sdk/go/serviceability" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) @@ -25,20 +25,19 @@ const retransmitSubscribeTimeout = 90 * time.Second var ( enableRetransmitOnlyTests = flag.Bool("enable-retransmit-only-settlement-tests", false, "enable the retransmit-only multicast settlement test") retransmitOnlyDeviceFlag = flag.String("retransmit-only-device", "", "device code or pubkey in a retransmit-only metro (overrides auto-discovery)") - leaderGroupCodeFlag = flag.String("leader-group-code", "", "multicast group code for the leader (full) shred feed, expected to be EXCLUDED from the seat") - retransmitGroupCodeFlag = flag.String("retransmit-group-code", "", "multicast group code for the retransmit shred feed, expected to be subscribed") + retransmitGroupCodesFlag = flag.String("retransmit-group-codes", "", "comma-separated multicast group codes a seat in a retransmit-only metro must subscribe to, and nothing else") retransmitPriceFlag = flag.Uint64("retransmit-price", 10, "expected discounted retransmit-only seat price in whole USDC dollars") ) // TestQA_RetransmitOnlySettlement demonstrates the retransmit-only shred // subscription end to end: a client pays for a seat on a device in a // retransmit-only metro, is charged the discounted price, and ends up -// subscribed to the retransmit multicast group only — not the leader group, so -// leader shreds are excluded. It mirrors TestQA_MulticastSettlement over the -// shared runShredSettlement flow, adding retransmit-only device selection, the -// discounted-price assertion, and the subscribed-groups assertion. The test is -// environment-agnostic: the group codes and the expected price are flags, so -// the same binary validates the testnet QA network and later mainnet. +// subscribed to the retransmit groups only, so leader shreds are excluded. It +// mirrors TestQA_MulticastSettlement over the shared runShredSettlement flow, +// adding retransmit-only device selection, the discounted-price assertion, and +// the subscribed-groups assertion. The test is environment-agnostic: the group +// codes and the expected price are flags, so the same binary validates the +// testnet QA network and later mainnet. func TestQA_RetransmitOnlySettlement(t *testing.T) { runShredSettlement(t, shredSettlementParams{ enabled: *enableRetransmitOnlyTests, @@ -97,8 +96,7 @@ func selectRetransmitOnlyDevice(t *testing.T, ctx context.Context, log *slog.Log // The group codes are required to assert leader-exclusion once a device // exists to test. Enabling the test without them is a misconfiguration. - require.NotEmpty(t, *leaderGroupCodeFlag, "--leader-group-code is required") - require.NotEmpty(t, *retransmitGroupCodeFlag, "--retransmit-group-code is required") + require.NotEmpty(t, *retransmitGroupCodesFlag, "--retransmit-group-codes is required") return device } @@ -129,38 +127,34 @@ func flaggedMetroCodes(test *qa.Test, exchangeKeys map[string]bool) []string { } // assertSubscribedGroups polls the seat's onchain multicast subscription until -// it reflects retransmit-only membership: the retransmit group present and the -// leader group absent (leader shreds excluded). On timeout it logs the seat's -// last-seen subscription state so on-call can tell an oracle bug (retransmit -// group never appeared) from a leader-exclusion bug (leader group leaked in) -// without re-deriving it from a rerun. +// it holds the retransmit groups and nothing else, so the leader group, the +// root-node group and any non-retransmit per-metro group all fall away. On +// timeout it logs the seat's last-seen subscription state so on-call can tell an +// oracle bug (a retransmit group never appeared) from a leader-exclusion bug (a +// group leaked in) without re-deriving it from a rerun. func assertSubscribedGroups(t *testing.T, ctx context.Context, log *slog.Logger, client *qa.Client, _ *qa.Device) { - // Resolve the leader and retransmit multicast groups by code. Neither - // onchain data nor the SDK labels a group as leader/retransmit, so the - // operator supplies the codes per network. - leaderGroup, err := client.GetMulticastGroup(ctx, *leaderGroupCodeFlag) - require.NoError(t, err, "failed to resolve leader group %q", *leaderGroupCodeFlag) - require.NotNil(t, leaderGroup, "leader group %q not found onchain", *leaderGroupCodeFlag) - retransmitGroup, err := client.GetMulticastGroup(ctx, *retransmitGroupCodeFlag) - require.NoError(t, err, "failed to resolve retransmit group %q", *retransmitGroupCodeFlag) - require.NotNil(t, retransmitGroup, "retransmit group %q not found onchain", *retransmitGroupCodeFlag) - - subscribed := func(user *serviceability.User, group solana.PublicKey) bool { - for _, sub := range user.Subscribers { - if solana.PublicKeyFromBytes(sub[:]).Equals(group) { - return true - } + // Neither onchain data nor the SDK labels a group as leader or retransmit, + // so the operator supplies the codes per network. + required := make(map[solana.PublicKey]string) + for _, code := range strings.Split(*retransmitGroupCodesFlag, ",") { + code = strings.TrimSpace(code) + if code == "" { + continue } - return false + group, err := client.GetMulticastGroup(ctx, code) + require.NoError(t, err, "failed to resolve multicast group %q", code) + require.NotNil(t, group, "multicast group %q not found onchain", code) + required[group.PK] = code } + require.NotEmpty(t, required, "no multicast group resolved from --retransmit-group-codes %q", *retransmitGroupCodesFlag) // The oracle converges the seat's onchain subscription asynchronously, so // poll until it reflects retransmit-only membership. Capture the last-seen // state on every poll so the timeout branch can report it. var ( - lastSubscribed []string - lastRetransmitFound bool - lastLeaderFound bool + lastSubscribed []string + lastMissing []string + lastExtra []string ) ok := assert.Eventually(t, func() bool { user, err := client.GetServiceabilityUser(ctx) @@ -168,29 +162,39 @@ func assertSubscribedGroups(t *testing.T, ctx context.Context, log *slog.Logger, log.Info("serviceability user poll error", "error", err) return false } - lastRetransmitFound = subscribed(user, retransmitGroup.PK) - lastLeaderFound = subscribed(user, leaderGroup.PK) + subscribed := make(map[solana.PublicKey]bool, len(user.Subscribers)) subs := make([]string, 0, len(user.Subscribers)) for _, sub := range user.Subscribers { - subs = append(subs, solana.PublicKeyFromBytes(sub[:]).String()) + group := solana.PublicKeyFromBytes(sub[:]) + subscribed[group] = true + subs = append(subs, group.String()) + } + var missing, extra []string + for group, code := range required { + if !subscribed[group] { + missing = append(missing, code) + } + } + for group := range subscribed { + if _, want := required[group]; !want { + extra = append(extra, group.String()) + } } - lastSubscribed = subs - return lastRetransmitFound && !lastLeaderFound + lastSubscribed, lastMissing, lastExtra = subs, missing, extra + return len(missing) == 0 && len(extra) == 0 }, retransmitSubscribeTimeout, 5*time.Second) if !ok { log.Warn("seat did not converge to retransmit-only subscription within timeout", - "retransmit_group", retransmitGroup.PK, - "retransmit_present", lastRetransmitFound, - "leader_group", leaderGroup.PK, - "leader_present", lastLeaderFound, + "missing_groups", lastMissing, + "extra_groups", lastExtra, "subscribed_groups", lastSubscribed, ) } require.True(t, ok, - "seat should subscribe to the retransmit group and NOT the leader group") - log.Info("Subscribed to the retransmit group only", - "retransmit_group", retransmitGroup.PK, - "leader_group_excluded", leaderGroup.PK, + "the seat should subscribe to %s and to nothing else; it is missing %v and wrongly subscribes to %v", + *retransmitGroupCodesFlag, lastMissing, lastExtra) + log.Info("Subscribed to the retransmit groups only", + "retransmit_groups", *retransmitGroupCodesFlag, "subscriber_count", len(lastSubscribed), ) } diff --git a/e2e/qa_shred_settlement_test.go b/e2e/qa_shred_settlement_test.go index 082cb22212..422aee0660 100644 --- a/e2e/qa_shred_settlement_test.go +++ b/e2e/qa_shred_settlement_test.go @@ -40,6 +40,12 @@ type shredSettlementParams struct { enabled bool skipReason string + // preflightSubtestName / preflight, when both set, run as a subtest after the + // reconciler is enabled and before device selection, so a caller can assert + // on a payment the program must reject before this run pays for a seat. + preflightSubtestName string + preflight func(t *testing.T, ctx context.Context, log *slog.Logger, test *qa.Test, client *qa.Client) + // selectSubtestName is the subtest name under which selectDevice runs, e.g. // "find_closest_device" or "select_retransmit_only_device". selectSubtestName string @@ -188,6 +194,14 @@ func runShredSettlement(t *testing.T, p shredSettlementParams) { return } + if p.preflight != nil { + if !t.Run(p.preflightSubtestName, func(t *testing.T) { + p.preflight(t, ctx, log, test, client) + }) { + return + } + } + if !t.Run(p.selectSubtestName, func(t *testing.T) { device = p.selectDevice(t, ctx, log, test, client) }) { diff --git a/sdk/shreds/go/state.go b/sdk/shreds/go/state.go index c1ce79a82f..8e432d7883 100644 --- a/sdk/shreds/go/state.go +++ b/sdk/shreds/go/state.go @@ -69,9 +69,10 @@ type ProgramConfig struct { // Flag bits in ProgramConfig.Flags. Mirrors the onchain bit indices from the // shred-subscription program. const ( - programConfigFlagIsPausedBit = 0 - programConfigFlagIsMigratedBit = 1 - programConfigFlagIsProratedServiceEnabledBit = 2 + programConfigFlagIsPausedBit = 0 + programConfigFlagIsMigratedBit = 1 + programConfigFlagIsProratedServiceEnabledBit = 2 + programConfigFlagRetransmitOnlyOnboardingEnforcedBit = 7 ) // IsPaused returns true if the program is paused. @@ -90,6 +91,10 @@ func (p *ProgramConfig) IsProratedServiceEnabled() bool { return p.Flags&(1< Date: Thu, 6 Aug 2026 00:52:14 -0500 Subject: [PATCH 2/5] remove file with just a helper function and make test pass --- CHANGELOG.md | 2 + e2e/qa_multicast_settlement_test.go | 672 +++++++++++++++++++++- e2e/qa_retransmit_only_settlement_test.go | 200 ------- e2e/qa_shred_settlement_test.go | 590 ------------------- 4 files changed, 663 insertions(+), 801 deletions(-) delete mode 100644 e2e/qa_retransmit_only_settlement_test.go delete mode 100644 e2e/qa_shred_settlement_test.go diff --git a/CHANGELOG.md b/CHANGELOG.md index 87af8fb96b..0ef8f6ea46 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -23,6 +23,8 @@ All notable changes to this project will be documented in this file. - Samples dropped because a partition's onchain account is full are now counted too, under `reason="account_full"` plus `submitter_account_full` on the errors counter. That path reports success to its caller, so a warning was the only trace, and its count was wrong: it reported the whole flushed partition rather than the samples actually lost. (#4145) - A failed submission now retries from the first unwritten sample rather than restarting at the beginning of the flushed partition, and only the unwritten remainder is requeued. Previously a mid-partition error made every subsequent attempt re-send batches that were already onchain, appending those samples a second time and pushing the account toward the sample cap it is measured against. (#4145) - Agent logs now identify the ledger RPC endpoint in use, and the peer count and the stale-program-data warning report state transitions instead of firing on every refresh. New lines name the resolved remote address of each connection (a bad load balancer address behind a hostname was invisible before), and a refresh that finds no peers is now a warning rather than a Debug line. The stale-cache warning also moves off the package-global `slog` onto the agent's own logger, so it is formatted and leveled with everything else. New: `doublezero_device_telemetry_agent_peers` gauge, and `pinger_epoch_fetch` on the errors counter for every exhausted epoch fetch. (#4147) +QA + - Rework existing TestQA_MulticastSettlement and adapt it to the new `FLAG_RETRANSMIT_ONLY_ONBOARDING_ENFORCED_BIT` flag. Test now checks for this flag in the ProgramConfig solana account and depending on if it's on or off tries to assert that no new user can subscribe to a non retransmit-only metro unless that metro has the retransmit-only flag enabled in the MetroHistory account. (#4156) ## [v0.33.0](https://github.com/malbeclabs/doublezero/compare/client/v0.32.0...client/v0.33.0) - 2026-07-31 diff --git a/e2e/qa_multicast_settlement_test.go b/e2e/qa_multicast_settlement_test.go index 56f96a8019..3a4f40d4bf 100644 --- a/e2e/qa_multicast_settlement_test.go +++ b/e2e/qa_multicast_settlement_test.go @@ -7,20 +7,30 @@ import ( "flag" "log/slog" "strconv" + "strings" "testing" + "time" + "github.com/gagliardetto/solana-go" "github.com/malbeclabs/doublezero/e2e/internal/qa" + pb "github.com/malbeclabs/doublezero/e2e/proto/qa/gen/pb-go" + "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) -var enableSettlementTests = flag.Bool("enable-multicast-settlement-tests", false, "enable multicast settlement tests") +const ( + retransmitSubscribeTimeout = 90 * time.Second + balanceSettleTimeout = 30 * time.Second +) + +var ( + enableSettlementTests = flag.Bool("enable-multicast-settlement-tests", false, "enable multicast settlement tests") + retransmitOnlyDeviceFlag = flag.String("retransmit-only-device", "", "device code or pubkey in a retransmit-only metro (overrides auto-discovery)") + retransmitGroupCodesFlag = flag.String("retransmit-group-codes", "", "comma-separated multicast group codes a seat in a retransmit-only metro must subscribe to, and nothing else") + keypairFlag = flag.String("keypair", "$HOME/.config/doublezero/id.json", "path to keypair file for settlement commands") + settlementClientFlag = flag.String("multicast-settlement-client", "", "host of the client to use for settlement tests (overrides random selection)") +) -// TestQA_MulticastSettlement pays for a multicast seat and verifies pricing, -// tunnel-up, and the withdraw/refund accounting. It is a thin wrapper over the -// shared settlement flow in runShredSettlement; TestQA_RetransmitOnlySettlement -// mirrors it with retransmit-only device selection and a group-subscription -// assertion. Any fix to the settlement machinery belongs in runShredSettlement -// so both stay in sync. func TestQA_MulticastSettlement(t *testing.T) { var onboardingEnforced bool runShredSettlement(t, shredSettlementParams{ @@ -60,8 +70,6 @@ func TestQA_MulticastSettlement(t *testing.T) { }) } -// selectClosestDevice picks the reachable device with the lowest latency, -// regardless of metro flags. It backs TestQA_MulticastSettlement. func selectClosestDevice(t *testing.T, ctx context.Context, log *slog.Logger, _ *qa.Test, client *qa.Client) *qa.Device { device, err := client.ClosestDevice(ctx) require.NoError(t, err, "failed to find closest device") @@ -69,8 +77,6 @@ func selectClosestDevice(t *testing.T, ctx context.Context, log *slog.Logger, _ return device } -// assertNewSeatRejected pays for a seat on the closest device outside a -// retransmit-only metro and requires the program to reject it. func assertNewSeatRejected(t *testing.T, ctx context.Context, log *slog.Logger, client *qa.Client) { device, err := client.ClosestNonRetransmitOnlyDevice(ctx) require.NoError(t, err, "failed to find a device outside a retransmit-only metro") @@ -102,3 +108,647 @@ func assertNewSeatRejected(t *testing.T, ctx context.Context, log *slog.Logger, "the payment failed for another reason than retransmit-only onboarding enforcement") log.Info("The program rejected the new seat", "device", device.Code, "metro", device.ExchangeCode) } + +// No flagged metro means the feature is not configured on this network, so the +// test skips. A flagged metro with no reachable device fails instead, because +// the feature would otherwise go unexercised. +func selectRetransmitOnlyDevice(t *testing.T, ctx context.Context, log *slog.Logger, test *qa.Test, client *qa.Client) *qa.Device { + var device *qa.Device + if *retransmitOnlyDeviceFlag != "" { + pinned, ok := test.DeviceByCodeOrPubkey(*retransmitOnlyDeviceFlag) + require.True(t, ok, "pinned device %q not found", *retransmitOnlyDeviceFlag) + retransmitOnly, err := client.RetransmitOnlyExchangeKeys(ctx) + require.NoError(t, err, "failed to read retransmit-only metros") + require.True(t, retransmitOnly[pinned.ExchangePubKey], + "pinned device %s (metro %s) is not in a retransmit-only metro", pinned.Code, pinned.ExchangeCode) + device = pinned + } else { + selected, retransmitOnly, err := client.ClosestRetransmitOnlyDevice(ctx) + require.NoError(t, err, "failed to find a retransmit-only device") + if len(retransmitOnly) == 0 { + t.Skip("Skipping: no metro is flagged retransmit-only (feature not deployed/configured on this network)") + } + require.NotNil(t, selected, + "retransmit-only metros %v are configured but no reachable device matched; the feature under test cannot be exercised", + flaggedMetroCodes(test, retransmitOnly)) + device = selected + } + log.Info("Selected retransmit-only device", "code", device.Code, "pubkey", device.PubKey, "metro", device.ExchangeCode) + + require.NotEmpty(t, *retransmitGroupCodesFlag, "--retransmit-group-codes is required") + return device +} + +func flaggedMetroCodes(test *qa.Test, exchangeKeys map[string]bool) []string { + seen := make(map[string]bool) + codes := make([]string, 0, len(exchangeKeys)) + for _, d := range test.Devices() { + if !exchangeKeys[d.ExchangePubKey] || seen[d.ExchangePubKey] { + continue + } + seen[d.ExchangePubKey] = true + code := d.ExchangeCode + if code == "" { + code = d.ExchangePubKey + } + codes = append(codes, code) + } + for key := range exchangeKeys { + if !seen[key] { + codes = append(codes, key) + } + } + return codes +} + +func assertSubscribedGroups(t *testing.T, ctx context.Context, log *slog.Logger, client *qa.Client, _ *qa.Device) { + // Nothing onchain labels a group as leader or retransmit, so the operator + // names the retransmit groups per network. + required := make(map[solana.PublicKey]string) + for _, code := range strings.Split(*retransmitGroupCodesFlag, ",") { + code = strings.TrimSpace(code) + if code == "" { + continue + } + group, err := client.GetMulticastGroup(ctx, code) + require.NoError(t, err, "failed to resolve multicast group %q", code) + require.NotNil(t, group, "multicast group %q not found onchain", code) + required[group.PK] = code + } + require.NotEmpty(t, required, "no multicast group resolved from --retransmit-group-codes %q", *retransmitGroupCodesFlag) + + // The oracle converges the seat's onchain subscription asynchronously, so + // poll until it reflects retransmit-only membership. + var ( + lastSubscribed []string + lastMissing []string + lastExtra []string + ) + ok := assert.Eventually(t, func() bool { + user, err := client.GetServiceabilityUser(ctx) + if err != nil { + log.Info("serviceability user poll error", "error", err) + return false + } + subscribed := make(map[solana.PublicKey]bool, len(user.Subscribers)) + subs := make([]string, 0, len(user.Subscribers)) + for _, sub := range user.Subscribers { + group := solana.PublicKeyFromBytes(sub[:]) + subscribed[group] = true + subs = append(subs, group.String()) + } + var missing, extra []string + for group, code := range required { + if !subscribed[group] { + missing = append(missing, code) + } + } + for group := range subscribed { + if _, want := required[group]; !want { + extra = append(extra, group.String()) + } + } + lastSubscribed, lastMissing, lastExtra = subs, missing, extra + return len(missing) == 0 && len(extra) == 0 + }, retransmitSubscribeTimeout, 5*time.Second) + if !ok { + log.Warn("seat did not converge to retransmit-only subscription within timeout", + "missing_groups", lastMissing, + "extra_groups", lastExtra, + "subscribed_groups", lastSubscribed, + ) + } + require.True(t, ok, + "the seat should subscribe to %s and to nothing else; it is missing %v and wrongly subscribes to %v", + *retransmitGroupCodesFlag, lastMissing, lastExtra) + log.Info("Subscribed to the retransmit groups only", + "retransmit_groups", *retransmitGroupCodesFlag, + "subscriber_count", len(lastSubscribed), + ) +} + +type shredSettlementParams struct { + enabled bool + skipReason string + + // preflight runs after the reconciler is enabled and before device selection. + preflightSubtestName string + preflight func(t *testing.T, ctx context.Context, log *slog.Logger, test *qa.Test, client *qa.Client) + + selectSubtestName string + // selectDevice may call t.Skip when the feature is not configured. It logs + // its own selection detail. + selectDevice func(t *testing.T, ctx context.Context, log *slog.Logger, test *qa.Test, client *qa.Client) *qa.Device + + priceLogMsg string + + // extraAssertion runs after the tunnel is up and before the seat is withdrawn. + extraSubtestName string + extraAssertion func(t *testing.T, ctx context.Context, log *slog.Logger, client *qa.Client, device *qa.Device) +} + +// runShredSettlement picks a device, pays its seat price, checks the debit and +// the tunnel, then withdraws the seat and checks the refund. +func runShredSettlement(t *testing.T, p shredSettlementParams) { + if !p.enabled { + t.Skip(p.skipReason) + } + + log := newTestLogger(t) + ctx := t.Context() + test, err := qa.NewTest(ctx, log, hostsArg, portArg, networkConfig, nil) + require.NoError(t, err, "failed to create test") + + var client *qa.Client + if *settlementClientFlag != "" { + var ok bool + client, ok = test.ClientByHost(*settlementClientFlag) + require.True(t, ok, "client %q not found in hosts", *settlementClientFlag) + } else { + client = test.RandomClient() + } + if *keypairFlag != "" { + client.Keypair = *keypairFlag + } + log.Info("Selected client", "host", client.Host) + + var device *qa.Device + var amount string + var quoted *pb.DevicePrice + var onchain *qa.SeatPrices + var fundedAmount uint64 + var effectivePrice uint64 + var balanceBeforePay uint64 + var balanceAfterPay uint64 + seatPaid := false + + t.Cleanup(func() { + if seatPaid && device != nil { + // Retry the withdraw: a single-shot withdraw that hit the spurious + // "request in flight" preflight bail (or a transient RPC failure) + // leaves the seat active onchain with an open escrow, poisoning + // every subsequent hourly run. Retrying over a bounded window heals + // the state instead of letting the escrow grow one epoch per run. + // Bound the cleanup so a hung withdraw can't block teardown forever. + cleanupCtx, cancel := context.WithTimeout(context.Background(), 4*time.Minute) + defer cancel() + if withdrawErr := client.WithdrawSeatWithRetry(cleanupCtx, device.PubKey); withdrawErr != nil { + // Warn, not Info: the seat is left active onchain and the escrow + // will grow next run, so this must stand out in the run log. + log.Warn("Cleanup: seat withdraw failed after retries; seat left active onchain", "error", withdrawErr) + } + } + if t.Failed() { + client.DumpDiagnostics(nil) + } + }) + + if !t.Run("ensure_program_unpaused", func(t *testing.T) { + // Migrations pause the program; while paused the oracle cannot ack + // instant seat allocation requests, which would leave the seat + // un-withdrawable and fail the rest of the test with a confusing + // "invalid account data for instruction" rejection. + paused, err := client.IsProgramPaused(ctx) + require.NoError(t, err, "failed to read program-paused flag") + if paused { + t.Skip("Skipping: shred-subscription program is paused (migration in progress)") + } + }) { + return + } + + if !t.Run("ensure_multicast_disconnected", func(t *testing.T) { + // Self-heal a seat left stuck-active onchain by a previous run whose + // withdraw bailed. This is the poisoned state that can't be seen from a + // session status: `shreds pay` on an already-active seat only tops up the + // escrow and never creates a new allocation request, so the seat never + // re-acks and the tunnel never comes up. Detect and withdraw it before + // the session check so the run starts from a clean slate. + healed, err := client.SelfHealStuckSeats(ctx) + require.NoError(t, err, "failed to self-heal stuck-active seats") + if healed > 0 { + log.Info("Self-healed stuck-active seat(s)", "count", healed) + } + + statuses, err := client.GetUserStatuses(ctx) + if err != nil { + log.Info("No active sessions") + return + } + var mcast *pb.Status + for _, s := range statuses { + if s.UserType == "Multicast" && s.SessionStatus != qa.UserStatusDisconnected { + mcast = s + break + } + } + if mcast == nil { + log.Info("No active multicast session") + return + } + log.Info("Active multicast session found, withdrawing", "device", mcast.CurrentDevice, "status", mcast.SessionStatus) + dev, ok := test.Devices()[mcast.CurrentDevice] + require.True(t, ok, "device %q not found in devices map", mcast.CurrentDevice) + err = client.WithdrawSeatWithRetry(ctx, dev.PubKey) + require.NoError(t, err, "failed to withdraw existing seat") + err = client.WaitForMulticastStatusDisconnected(ctx) + require.NoError(t, err, "existing multicast session did not disconnect") + }) { + return + } + + if !t.Run("enable_reconciler", func(t *testing.T) { + err := client.FeedEnable(ctx) + require.NoError(t, err, "failed to enable reconciler") + }) { + return + } + + if p.preflight != nil { + if !t.Run(p.preflightSubtestName, func(t *testing.T) { + p.preflight(t, ctx, log, test, client) + }) { + return + } + } + + if !t.Run(p.selectSubtestName, func(t *testing.T) { + device = p.selectDevice(t, ctx, log, test, client) + }) { + return + } + if device == nil { + // selectDevice skipped its subtest (e.g. no retransmit-only metro is + // configured on this network). t.Run reports a skipped subtest as + // success, so the skip does not stop the parent — skip it explicitly + // here rather than dereferencing a nil device in query_seat_price below. + // A nil device after a non-failed subtest can only mean the selector + // skipped: a failure would have made t.Run return false and returned + // above, and the success path always assigns a device. + t.Skip("Skipping: device selection skipped (feature not configured on this network)") + } + + if !t.Run("query_seat_price", func(t *testing.T) { + prices, err := client.FeedSeatPrice(ctx, device.PubKey) + require.NoError(t, err, "failed to get seat prices") + + // Match by pubkey, not code: querying by --device skips code resolution, + // so the returned rows may not carry a device_code. + for _, pr := range prices { + if pr.DevicePubkey == device.PubKey { + quoted = pr + break + } + } + require.NotNil(t, quoted, "no price found for device %s", device.Code) + require.NotZero(t, quoted.EpochPrice, "epoch price is zero for device %s", device.Code) + log.Info(p.priceLogMsg, "device", device.Code, + "epoch_price", quoted.EpochPrice, + "instant_allocation_price", quoted.GetInstantAllocationPrice(), + "reports_instant_allocation_price", quoted.GetReportsInstantAllocationPrice()) + }) { + return + } + + if !t.Run("query_onchain_seat_price", func(t *testing.T) { + // Read the prices the program itself computes straight off the chain, as + // an oracle for the CLI quote. This snapshot is taken next to the quote so + // the two comparisons below are as close to simultaneous as possible; the + // amount actually funded is re-read just before paying, since the wait for + // the open phase can outlive this read. + var err error + onchain, err = client.SeatPrices(ctx, device.PubKey) + require.NoError(t, err, "failed to compute the onchain seat prices") + require.NotZero(t, onchain.InstantAllocationDollars, "onchain instant-allocation price is zero for device %s", device.Code) + + log.Info("Onchain seat prices", "device", device.Code, + "instant_allocation_price", onchain.InstantAllocationDollars, + "last_settled_epoch", onchain.LastSettledEpoch, + "current_epoch_price", onchain.CurrentEpochDollars, + "current_subscription_epoch", onchain.CurrentSubscriptionEpoch) + }) { + return + } + + // The next two subtests are deliberately not guarded with + // `if !t.Run(...) { return }`: a CLI/chain divergence is a real product bug + // and must fail the run loudly, but the payment below funds from the onchain + // price, so the rest of the settlement flow is still worth running. Both + // compare whole dollars, never prorated micro-USDC: proration is a function + // of the current slot, so an exact comparison against a separately-timed read + // is inherently racy. + t.Run("validate_instant_allocation_price_matches_chain", func(t *testing.T) { + switch { + case !quoted.GetReportsInstantAllocationPrice(): + // QA hosts install doublezero-solana from a version-pinned apt package + // (doublezero_solana_version in malbeclabs/infra + // ansible/inventory/*/group_vars/all.yml, 0.5.10-1 at time of writing), + // so the field only appears once a release carrying it is published and + // the pin bumped. Asserting against an absent field would read 0 and + // fail as "quoted 0, chain 43" — a misleading failure that looks like a + // new bug rather than a rollout gap. + t.Skipf("Skipping: installed doublezero-solana does not report instant_allocation_price (needs a release newer than the pinned 0.5.10-1); chain says %d USDC at last_settled_epoch=%d", + onchain.InstantAllocationDollars, onchain.LastSettledEpoch) + case quoted.InstantAllocationPrice == nil: + // Reported, but null: the CLI could not find the settled-epoch ring + // entry. That is a real condition, not a rollout artifact — the program + // performs the same lookup and would reject the allocation. + t.Fatalf("`shreds price` reported instant_allocation_price as unavailable for device %s, but the chain has a price of %d USDC at last_settled_epoch=%d", + device.Code, onchain.InstantAllocationDollars, onchain.LastSettledEpoch) + default: + require.Equal(t, onchain.InstantAllocationDollars, quoted.GetInstantAllocationPrice(), + "CLI quoted %d USDC for an instant allocation but the program charges %d USDC (last_settled_epoch=%d); `shreds price` and `shreds pay` read different ring entries", + quoted.GetInstantAllocationPrice(), onchain.InstantAllocationDollars, onchain.LastSettledEpoch) + } + }) + + t.Run("validate_epoch_price_matches_chain", func(t *testing.T) { + // The other half of the invariant: epoch_price must keep meaning the + // current-epoch price a recurring subscriber pays. Without this, someone + // "fixing" the divergence by repointing epoch_price at last_settled_epoch + // would go green here while silently breaking recurring subscribers. + if !onchain.HasCurrentEpoch { + // Part-way through UpdatingPrices the ring has not been advanced to the + // current epoch for every metro and device yet, so there is nothing to + // compare against. Transient by design, not a regression. + t.Skipf("Skipping: no onchain price entry yet for current_subscription_epoch=%d (prices are still being updated)", + onchain.CurrentSubscriptionEpoch) + } + require.Equal(t, onchain.CurrentEpochDollars, quoted.EpochPrice, + "CLI quoted epoch_price %d USDC but the chain has %d USDC at current_subscription_epoch=%d; epoch_price must stay the price a recurring subscriber pays next epoch", + quoted.EpochPrice, onchain.CurrentEpochDollars, onchain.CurrentSubscriptionEpoch) + }) + + // Set when wait_for_open_phase times out inside the epoch-tail closed + // window (verified against live chain state), so the parent can skip the + // remaining subtests: the program stays closed until the epoch boundary, + // which the 2-minute wait cannot bridge, so pay/withdraw below could only + // fail and page for a by-design condition. + var epochTailWindow *qa.EpochTailWindow + if !t.Run("wait_for_open_phase", func(t *testing.T) { + // Record where the wait begins: a timeout means the program was closed + // for the entire wait, so the classification below can require the + // whole span — not just the timeout-time slot — to be inside the + // window. Best-effort: on a read failure classification degrades to + // the timeout-time slot only. + waitStartSlot, slotErr := client.CurrentSolanaSlot(ctx) + if slotErr != nil { + log.Warn("Failed to read wait-start slot; epoch-tail classification will use the timeout-time slot only", "error", slotErr) + } + err := client.WaitForOpenForRequests(ctx) + if err != nil { + // For the last grace-period slots of every epoch the shred oracle + // closes the program by design (settle seats, update prices) and + // reopens it just after the epoch boundary. Verify against live + // chain state — onchain grace period, controller phase, and RPC + // epoch schedule — whether this timeout landed in that window; a + // timeout outside it must keep failing exactly as loudly as before. + win, winErr := client.EpochTailClosedWindow(ctx, waitStartSlot) + switch { + case winErr != nil: + log.Warn("Failed to classify epoch-tail closed window; treating timeout as a real failure", "error", winErr) + case win.Benign: + epochTailWindow = &win + t.Skipf("expected epoch-tail closed window: %s", win) + default: + // Give on-call the computed window so a real outage's distance + // from the benign window is visible in the run log. + log.Info("Timeout is not the benign epoch-tail closed window", "window", win.String()) + } + } + require.NoError(t, err, "shred-subscription program did not enter OpenForRequests phase within timeout") + }) { + return + } + if epochTailWindow != nil { + // t.Run reports a skipped subtest as success, so skip the parent + // explicitly to stop the run here. + t.Skipf("expected epoch-tail closed window: %s", epochTailWindow) + } + + if !t.Run("refresh_onchain_seat_price", func(t *testing.T) { + // wait_for_open_phase blocks for up to two minutes, and a settlement + // completing inside that window advances last_settled_epoch — the very + // read the charge is derived from. Funding a price captured before the + // wait would underfund the escrow and reproduce the opaque pay-time + // rejection this test exists to avoid, so the funded amount comes from a + // read taken here, immediately before paying. A rollover between this read + // and the transaction landing is irreducible (any payer races it), but the + // window shrinks from minutes to seconds. + refreshed, err := client.SeatPrices(ctx, device.PubKey) + require.NoError(t, err, "failed to re-read the onchain seat prices before paying") + require.NotZero(t, refreshed.InstantAllocationDollars, "onchain instant-allocation price is zero for device %s", device.Code) + + if refreshed.LastSettledEpoch != onchain.LastSettledEpoch || + refreshed.InstantAllocationDollars != onchain.InstantAllocationDollars { + // Warn, not Info: this means the quote comparisons above were made + // against a snapshot the payment no longer uses, so a failure up there + // should be read in that light. + log.Warn("Onchain seat price moved while waiting for the open-for-requests phase; funding the refreshed price", + "device", device.Code, + "before_price", onchain.InstantAllocationDollars, "before_last_settled_epoch", onchain.LastSettledEpoch, + "after_price", refreshed.InstantAllocationDollars, "after_last_settled_epoch", refreshed.LastSettledEpoch) + } + + onchain = refreshed + amount = strconv.FormatUint(onchain.InstantAllocationDollars, 10) + fundedAmount = onchain.InstantAllocationDollars * 1_000_000 // dollars to USDC raw units (6 decimals) + log.Info("Funding the seat escrow from the onchain price", "device", device.Code, + "amount", amount, "last_settled_epoch", onchain.LastSettledEpoch) + }) { + return + } + + if !t.Run("record_balance_before_pay", func(t *testing.T) { + var err error + balanceBeforePay, err = client.GetUSDCBalance(ctx) + require.NoError(t, err, "failed to get USDC balance before pay") + log.Info("USDC balance before pay", "balance", balanceBeforePay) + }) { + return + } + + if !t.Run("pay_for_seat", func(t *testing.T) { + err := client.FeedSeatPay(ctx, device.PubKey, amount) + require.NoError(t, err, "failed to pay for seat") + seatPaid = true + }) { + return + } + + if !t.Run("validate_balance_after_pay", func(t *testing.T) { + // Poll until the balance reflects the debit. FeedSeatPay returns + // after the tx is submitted, and the RPC balance view can lag the + // confirmed state briefly, so a one-shot read races. + var lastDebit uint64 + require.Eventually(t, func() bool { + bal, err := client.GetUSDCBalance(ctx) + if err != nil { + log.Info("USDC balance poll error", "error", err) + return false + } + balanceAfterPay = bal + lastDebit = balanceBeforePay - bal + return lastDebit == fundedAmount + }, balanceSettleTimeout, 5*time.Second, "USDC balance should decrease by the paid amount") + log.Info("USDC balance after pay", "balance", balanceAfterPay, "debit", lastDebit, "expected_debit", fundedAmount) + }) { + return + } + + if !t.Run("query_effective_seat_price", func(t *testing.T) { + // Built on the onchain price, not the CLI quote: this feeds the + // non-prorating balance assertion below, which must predict what the + // program actually charged. GetEffectiveSeatPrice applies the seat's price + // override on top when one is set. + var err error + effectivePrice, err = client.GetEffectiveSeatPrice(ctx, device.PubKey, onchain.InstantAllocationDollars) + require.NoError(t, err, "failed to get effective seat price") + log.Info("Effective seat price", "effective_usdc", effectivePrice, "funded_usdc", fundedAmount) + }) { + return + } + + if !t.Run("validate_tunnel_up", func(t *testing.T) { + err := client.WaitForMulticastStatusUp(ctx) + require.NoError(t, err, "multicast tunnel did not come up after seat payment") + }) { + return + } + + if !t.Run("validate_device_assignment", func(t *testing.T) { + statuses, err := client.GetUserStatuses(ctx) + require.NoError(t, err, "failed to get user statuses") + mcastStatus := qa.FindMulticastStatus(statuses) + require.NotNil(t, mcastStatus, "no multicast status found after seat payment") + require.Equal(t, device.Code, mcastStatus.CurrentDevice, "tunnel connected to wrong device") + log.Info("Tunnel up and device matches", "device", mcastStatus.CurrentDevice, "dzIP", mcastStatus.DoubleZeroIp) + }) { + return + } + + if p.extraAssertion != nil { + if !t.Run(p.extraSubtestName, func(t *testing.T) { + p.extraAssertion(t, ctx, log, client, device) + }) { + return + } + } + + if !t.Run("withdraw_seat", func(t *testing.T) { + // Withdraw is rejected while this run's instant allocation request is + // in flight (or a stale RPC read claims it is), so retry with endpoint + // rotation rather than waiting on an ack the harness cannot observe + // reliably. + err := client.WithdrawSeatWithRetry(ctx, device.PubKey) + require.NoError(t, err, "failed to withdraw seat") + seatPaid = false + }) { + return + } + + if !t.Run("validate_tunnel_down", func(t *testing.T) { + err := client.WaitForMulticastStatusDisconnected(ctx) + require.NoError(t, err, "tunnel did not come down after seat withdrawal") + }) { + return + } + + t.Run("validate_balance_after_withdraw", func(t *testing.T) { + // Read onchain whether the shred-subscription program has prorated + // service enabled. This lets the test self-adapt across environments + // (testnet has it on, mainnet does not) without needing a CI flag. + proratingEnabled, err := client.IsSeatProratingEnabled(ctx) + require.NoError(t, err, "failed to read prorating flag from program config") + + var balanceAfterWithdraw uint64 + if proratingEnabled { + // Prorating refunds the unused portion of the epoch to the wallet. + // Poll until the refund is reflected (balance strictly greater + // than after-pay). + require.Eventually(t, func() bool { + bal, err := client.GetUSDCBalance(ctx) + if err != nil { + log.Info("USDC balance poll error", "error", err) + return false + } + balanceAfterWithdraw = bal + return bal > balanceAfterPay + }, balanceSettleTimeout, 5*time.Second, + "USDC balance should increase to reflect the prorated refund") + } else { + expectedBalance := balanceBeforePay - effectivePrice + require.Eventually(t, func() bool { + bal, err := client.GetUSDCBalance(ctx) + if err != nil { + log.Info("USDC balance poll error", "error", err) + return false + } + balanceAfterWithdraw = bal + return bal == expectedBalance + }, balanceSettleTimeout, 5*time.Second, + "USDC balance should equal before_pay minus the effective seat price") + } + + refund := balanceAfterWithdraw - balanceAfterPay + + // A seat's payment escrow can carry a balance from an earlier run whose + // withdraw did not complete (e.g. during the reservoir-ack outage on + // devnet). Closing the escrow now refunds that leftover too, so the + // wallet-measured refund exceeds what was paid this run and no longer + // isolates this payment (`retained` would underflow). In that case the + // wallet-delta proration check is not meaningful, so skip it rather than + // fail — the settlement path itself is still covered by the pay/ack/ + // tunnel/withdraw sub-tests above. + if refund > fundedAmount { + log.Warn("skipping wallet-delta proration check: refund exceeds amount paid this run (pre-existing escrow drained)", + "refund", refund, + "paid_amount", fundedAmount, + "before_pay", balanceBeforePay, + "after_pay", balanceAfterPay, + "after_withdraw", balanceAfterWithdraw, + ) + return + } + // Equivalent to balanceBeforePay - balanceAfterWithdraw, but computed from + // the amount paid this run so it cannot underflow given the guard above. + retained := fundedAmount - refund + + log.Info("USDC balance after withdraw", + "balance", balanceAfterWithdraw, + "before_pay", balanceBeforePay, + "after_pay", balanceAfterPay, + "paid_amount", fundedAmount, + "effective_price", effectivePrice, + "refund", refund, + "retained", retained, + "prorating_enabled", proratingEnabled, + ) + + // Accounting invariant: regardless of prorating, the sum of what was + // refunded to the wallet and what the program retained must equal the + // amount debited at pay time. This uses fundedAmount rather than + // effectivePrice because a seat with a zero price override is still + // charged fundedAmount at pay and fully refunded on withdraw. + require.Equal(t, fundedAmount, refund+retained, + "refund + retained must equal the amount paid") + + if !proratingEnabled || effectivePrice == 0 { + return + } + + // With prorating enabled we avoid replicating the onchain formula + // against client-side RPC state (epoch schedule + current epoch reads + // are fragile on DZ ledger). Instead assert the qualitative invariants + // that distinguish a real partial refund from a regression: + // - refund > 0 (prorating actually happened) + // - retained > 0 (the seat was not free for the used portion) + // - retained < effective_price (kept less than a full epoch) + require.Greater(t, refund, uint64(0), + "prorating: refund should be strictly greater than zero") + require.Greater(t, retained, uint64(0), + "prorating: retained should be strictly greater than zero") + require.Less(t, retained, effectivePrice, + "prorating: retained should be strictly less than the effective price") + }) +} diff --git a/e2e/qa_retransmit_only_settlement_test.go b/e2e/qa_retransmit_only_settlement_test.go deleted file mode 100644 index 19169cc378..0000000000 --- a/e2e/qa_retransmit_only_settlement_test.go +++ /dev/null @@ -1,200 +0,0 @@ -//go:build qa - -package e2e - -import ( - "context" - "flag" - "log/slog" - "strings" - "testing" - "time" - - "github.com/gagliardetto/solana-go" - "github.com/malbeclabs/doublezero/e2e/internal/qa" - "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/require" -) - -// retransmitSubscribeTimeout bounds how long we wait for the oracle to converge -// the seat's onchain multicast subscription after the tunnel comes up. The -// tunnel-up wait already covers the subscribe path, so this is a safety buffer -// against the oracle's reconcile cadence. -const retransmitSubscribeTimeout = 90 * time.Second - -var ( - enableRetransmitOnlyTests = flag.Bool("enable-retransmit-only-settlement-tests", false, "enable the retransmit-only multicast settlement test") - retransmitOnlyDeviceFlag = flag.String("retransmit-only-device", "", "device code or pubkey in a retransmit-only metro (overrides auto-discovery)") - retransmitGroupCodesFlag = flag.String("retransmit-group-codes", "", "comma-separated multicast group codes a seat in a retransmit-only metro must subscribe to, and nothing else") - retransmitPriceFlag = flag.Uint64("retransmit-price", 10, "expected discounted retransmit-only seat price in whole USDC dollars") -) - -// TestQA_RetransmitOnlySettlement demonstrates the retransmit-only shred -// subscription end to end: a client pays for a seat on a device in a -// retransmit-only metro, is charged the discounted price, and ends up -// subscribed to the retransmit groups only, so leader shreds are excluded. It -// mirrors TestQA_MulticastSettlement over the shared runShredSettlement flow, -// adding retransmit-only device selection, the discounted-price assertion, and -// the subscribed-groups assertion. The test is environment-agnostic: the group -// codes and the expected price are flags, so the same binary validates the -// testnet QA network and later mainnet. -func TestQA_RetransmitOnlySettlement(t *testing.T) { - runShredSettlement(t, shredSettlementParams{ - enabled: *enableRetransmitOnlyTests, - skipReason: "Skipping: --enable-retransmit-only-settlement-tests flag not set", - - selectSubtestName: "select_retransmit_only_device", - selectDevice: selectRetransmitOnlyDevice, - - priceLogMsg: "Found discounted epoch price", - assertPrice: func(t *testing.T, device *qa.Device, price uint64) { - // The retransmit-only metro is priced at the discount, so the seat - // price must equal the expected retransmit price (default 10 USDC). - // This is the price the program charges, read from chain, not the CLI - // quote — the quote is checked against chain separately. - require.Equal(t, *retransmitPriceFlag, price, - "retransmit-only device %s should be priced at the discounted retransmit price", device.Code) - }, - - extraSubtestName: "assert_subscribed_groups", - extraAssertion: assertSubscribedGroups, - }) -} - -// selectRetransmitOnlyDevice picks the device to settle against for the -// retransmit-only test: the -retransmit-only-device pin when set, otherwise the -// closest reachable device in a retransmit-only metro (auto-discovery). -// -// Auto-discovery distinguishes two outcomes: when no metro is flagged -// retransmit-only the feature is simply not deployed here, so it skips; but when -// metros are flagged and yet none of their devices is reachable, it fails -// (naming the flagged metros) rather than skipping — otherwise the one feature -// this test exists to guard would go silently unexercised on the deployed -// network. -func selectRetransmitOnlyDevice(t *testing.T, ctx context.Context, log *slog.Logger, test *qa.Test, client *qa.Client) *qa.Device { - var device *qa.Device - if *retransmitOnlyDeviceFlag != "" { - pinned, ok := test.DeviceByCodeOrPubkey(*retransmitOnlyDeviceFlag) - require.True(t, ok, "pinned device %q not found", *retransmitOnlyDeviceFlag) - retransmitOnly, err := client.RetransmitOnlyExchangeKeys(ctx) - require.NoError(t, err, "failed to read retransmit-only metros") - require.True(t, retransmitOnly[pinned.ExchangePubKey], - "pinned device %s (metro %s) is not in a retransmit-only metro", pinned.Code, pinned.ExchangeCode) - device = pinned - } else { - selected, retransmitOnly, err := client.ClosestRetransmitOnlyDevice(ctx) - require.NoError(t, err, "failed to find a retransmit-only device") - if len(retransmitOnly) == 0 { - t.Skip("Skipping: no metro is flagged retransmit-only (feature not deployed/configured on this network)") - } - require.NotNil(t, selected, - "retransmit-only metros %v are configured but no reachable device matched; the feature under test cannot be exercised", - flaggedMetroCodes(test, retransmitOnly)) - device = selected - } - log.Info("Selected retransmit-only device", "code", device.Code, "pubkey", device.PubKey, "metro", device.ExchangeCode) - - // The group codes are required to assert leader-exclusion once a device - // exists to test. Enabling the test without them is a misconfiguration. - require.NotEmpty(t, *retransmitGroupCodesFlag, "--retransmit-group-codes is required") - return device -} - -// flaggedMetroCodes resolves the retransmit-only exchange pubkeys to readable -// exchange codes via the devices map, so a failure message can name the metros -// that were flagged but had no reachable device. Falls back to the raw pubkey -// for a flagged metro with no device in the map. -func flaggedMetroCodes(test *qa.Test, exchangeKeys map[string]bool) []string { - seen := make(map[string]bool) - codes := make([]string, 0, len(exchangeKeys)) - for _, d := range test.Devices() { - if !exchangeKeys[d.ExchangePubKey] || seen[d.ExchangePubKey] { - continue - } - seen[d.ExchangePubKey] = true - code := d.ExchangeCode - if code == "" { - code = d.ExchangePubKey - } - codes = append(codes, code) - } - for key := range exchangeKeys { - if !seen[key] { - codes = append(codes, key) - } - } - return codes -} - -// assertSubscribedGroups polls the seat's onchain multicast subscription until -// it holds the retransmit groups and nothing else, so the leader group, the -// root-node group and any non-retransmit per-metro group all fall away. On -// timeout it logs the seat's last-seen subscription state so on-call can tell an -// oracle bug (a retransmit group never appeared) from a leader-exclusion bug (a -// group leaked in) without re-deriving it from a rerun. -func assertSubscribedGroups(t *testing.T, ctx context.Context, log *slog.Logger, client *qa.Client, _ *qa.Device) { - // Neither onchain data nor the SDK labels a group as leader or retransmit, - // so the operator supplies the codes per network. - required := make(map[solana.PublicKey]string) - for _, code := range strings.Split(*retransmitGroupCodesFlag, ",") { - code = strings.TrimSpace(code) - if code == "" { - continue - } - group, err := client.GetMulticastGroup(ctx, code) - require.NoError(t, err, "failed to resolve multicast group %q", code) - require.NotNil(t, group, "multicast group %q not found onchain", code) - required[group.PK] = code - } - require.NotEmpty(t, required, "no multicast group resolved from --retransmit-group-codes %q", *retransmitGroupCodesFlag) - - // The oracle converges the seat's onchain subscription asynchronously, so - // poll until it reflects retransmit-only membership. Capture the last-seen - // state on every poll so the timeout branch can report it. - var ( - lastSubscribed []string - lastMissing []string - lastExtra []string - ) - ok := assert.Eventually(t, func() bool { - user, err := client.GetServiceabilityUser(ctx) - if err != nil { - log.Info("serviceability user poll error", "error", err) - return false - } - subscribed := make(map[solana.PublicKey]bool, len(user.Subscribers)) - subs := make([]string, 0, len(user.Subscribers)) - for _, sub := range user.Subscribers { - group := solana.PublicKeyFromBytes(sub[:]) - subscribed[group] = true - subs = append(subs, group.String()) - } - var missing, extra []string - for group, code := range required { - if !subscribed[group] { - missing = append(missing, code) - } - } - for group := range subscribed { - if _, want := required[group]; !want { - extra = append(extra, group.String()) - } - } - lastSubscribed, lastMissing, lastExtra = subs, missing, extra - return len(missing) == 0 && len(extra) == 0 - }, retransmitSubscribeTimeout, 5*time.Second) - if !ok { - log.Warn("seat did not converge to retransmit-only subscription within timeout", - "missing_groups", lastMissing, - "extra_groups", lastExtra, - "subscribed_groups", lastSubscribed, - ) - } - require.True(t, ok, - "the seat should subscribe to %s and to nothing else; it is missing %v and wrongly subscribes to %v", - *retransmitGroupCodesFlag, lastMissing, lastExtra) - log.Info("Subscribed to the retransmit groups only", - "retransmit_groups", *retransmitGroupCodesFlag, - "subscriber_count", len(lastSubscribed), - ) -} diff --git a/e2e/qa_shred_settlement_test.go b/e2e/qa_shred_settlement_test.go deleted file mode 100644 index 422aee0660..0000000000 --- a/e2e/qa_shred_settlement_test.go +++ /dev/null @@ -1,590 +0,0 @@ -//go:build qa - -package e2e - -import ( - "context" - "flag" - "log/slog" - "strconv" - "testing" - "time" - - "github.com/malbeclabs/doublezero/e2e/internal/qa" - pb "github.com/malbeclabs/doublezero/e2e/proto/qa/gen/pb-go" - "github.com/stretchr/testify/require" -) - -// balanceSettleTimeout bounds how long we wait for a USDC balance change to -// become visible after a settlement transaction is submitted. 30s covers -// the lag between FeedSeatPay/FeedSeatWithdraw returning and the balance -// RPC reflecting the debit/credit. -const balanceSettleTimeout = 30 * time.Second - -var ( - keypairFlag = flag.String("keypair", "$HOME/.config/doublezero/id.json", "path to keypair file for settlement commands") - settlementClientFlag = flag.String("multicast-settlement-client", "", "host of the client to use for settlement tests (overrides random selection)") -) - -// shredSettlementParams parameterizes the shared shred-pay settlement flow that -// backs both TestQA_MulticastSettlement and TestQA_RetransmitOnlySettlement. -// -// The two tests differ only in device selection, the seat-price assertion, and -// an optional post-tunnel-up assertion; everything else (self-heal, epoch-tail -// classification, escrow-drain guard, withdraw-retry, and the balance -// accounting invariants) is identical settlement machinery. That machinery has -// churned repeatedly (#4066, #4069), so it lives here once rather than being -// copied into each test where a fix would have to land twice. -type shredSettlementParams struct { - // enabled gates the whole test; when false it skips with skipReason. - enabled bool - skipReason string - - // preflightSubtestName / preflight, when both set, run as a subtest after the - // reconciler is enabled and before device selection, so a caller can assert - // on a payment the program must reject before this run pays for a seat. - preflightSubtestName string - preflight func(t *testing.T, ctx context.Context, log *slog.Logger, test *qa.Test, client *qa.Client) - - // selectSubtestName is the subtest name under which selectDevice runs, e.g. - // "find_closest_device" or "select_retransmit_only_device". - selectSubtestName string - // selectDevice picks and returns the device to settle against. It runs - // inside selectSubtestName, may call t.Skip (feature not configured) or - // fail, and is expected to log its own selection detail. - selectDevice func(t *testing.T, ctx context.Context, log *slog.Logger, test *qa.Test, client *qa.Client) *qa.Device - - // priceLogMsg is logged after the CLI seat price is queried, e.g. "Found - // epoch price" or "Found discounted epoch price". - priceLogMsg string - // assertPrice, when non-nil, runs an extra assertion on the seat price the - // program will charge (whole USDC dollars, read from chain) inside the - // query_onchain_seat_price subtest. It deliberately does not see the CLI - // quote: what a caller means to pin is what the seat actually costs. - assertPrice func(t *testing.T, device *qa.Device, price uint64) - - // extraSubtestName / extraAssertion, when both set, run as a gated subtest - // after the tunnel is up and the device assignment is validated, before the - // seat is withdrawn. The retransmit-only test uses it to assert the seat's - // multicast group subscription. - extraSubtestName string - extraAssertion func(t *testing.T, ctx context.Context, log *slog.Logger, client *qa.Client, device *qa.Device) -} - -// runShredSettlement drives the full shred-pay settlement flow: pick a device, -// query its seat price from the CLI and cross-check it against the price the -// program will charge (read from chain), wait for the program's -// open-for-requests phase, re-read that price and pay it, verify the debit and -// that the multicast tunnel comes up on the right device, (optionally) assert -// extra invariants, then withdraw and verify the refund accounting. Both -// settlement QA tests are thin wrappers over this. -func runShredSettlement(t *testing.T, p shredSettlementParams) { - if !p.enabled { - t.Skip(p.skipReason) - } - - log := newTestLogger(t) - ctx := t.Context() - test, err := qa.NewTest(ctx, log, hostsArg, portArg, networkConfig, nil) - require.NoError(t, err, "failed to create test") - - var client *qa.Client - if *settlementClientFlag != "" { - var ok bool - client, ok = test.ClientByHost(*settlementClientFlag) - require.True(t, ok, "client %q not found in hosts", *settlementClientFlag) - } else { - client = test.RandomClient() - } - if *keypairFlag != "" { - client.Keypair = *keypairFlag - } - log.Info("Selected client", "host", client.Host) - - // Shared state across subtests. - var device *qa.Device - var amount string - var quoted *pb.DevicePrice - var onchain *qa.SeatPrices - var fundedAmount uint64 - var effectivePrice uint64 - var balanceBeforePay uint64 - var balanceAfterPay uint64 - seatPaid := false - - t.Cleanup(func() { - if seatPaid && device != nil { - // Retry the withdraw: a single-shot withdraw that hit the spurious - // "request in flight" preflight bail (or a transient RPC failure) - // leaves the seat active onchain with an open escrow, poisoning - // every subsequent hourly run. Retrying over a bounded window heals - // the state instead of letting the escrow grow one epoch per run. - // Bound the cleanup so a hung withdraw can't block teardown forever. - cleanupCtx, cancel := context.WithTimeout(context.Background(), 4*time.Minute) - defer cancel() - if withdrawErr := client.WithdrawSeatWithRetry(cleanupCtx, device.PubKey); withdrawErr != nil { - // Warn, not Info: the seat is left active onchain and the escrow - // will grow next run, so this must stand out in the run log. - log.Warn("Cleanup: seat withdraw failed after retries; seat left active onchain", "error", withdrawErr) - } - } - if t.Failed() { - client.DumpDiagnostics(nil) - } - }) - - if !t.Run("ensure_program_unpaused", func(t *testing.T) { - // Migrations pause the program; while paused the oracle cannot ack - // instant seat allocation requests, which would leave the seat - // un-withdrawable and fail the rest of the test with a confusing - // "invalid account data for instruction" rejection. - paused, err := client.IsProgramPaused(ctx) - require.NoError(t, err, "failed to read program-paused flag") - if paused { - t.Skip("Skipping: shred-subscription program is paused (migration in progress)") - } - }) { - return - } - - if !t.Run("ensure_multicast_disconnected", func(t *testing.T) { - // Self-heal a seat left stuck-active onchain by a previous run whose - // withdraw bailed. This is the poisoned state that can't be seen from a - // session status: `shreds pay` on an already-active seat only tops up the - // escrow and never creates a new allocation request, so the seat never - // re-acks and the tunnel never comes up. Detect and withdraw it before - // the session check so the run starts from a clean slate. - healed, err := client.SelfHealStuckSeats(ctx) - require.NoError(t, err, "failed to self-heal stuck-active seats") - if healed > 0 { - log.Info("Self-healed stuck-active seat(s)", "count", healed) - } - - statuses, err := client.GetUserStatuses(ctx) - if err != nil { - log.Info("No active sessions") - return - } - var mcast *pb.Status - for _, s := range statuses { - if s.UserType == "Multicast" && s.SessionStatus != qa.UserStatusDisconnected { - mcast = s - break - } - } - if mcast == nil { - log.Info("No active multicast session") - return - } - log.Info("Active multicast session found, withdrawing", "device", mcast.CurrentDevice, "status", mcast.SessionStatus) - dev, ok := test.Devices()[mcast.CurrentDevice] - require.True(t, ok, "device %q not found in devices map", mcast.CurrentDevice) - err = client.WithdrawSeatWithRetry(ctx, dev.PubKey) - require.NoError(t, err, "failed to withdraw existing seat") - err = client.WaitForMulticastStatusDisconnected(ctx) - require.NoError(t, err, "existing multicast session did not disconnect") - }) { - return - } - - if !t.Run("enable_reconciler", func(t *testing.T) { - err := client.FeedEnable(ctx) - require.NoError(t, err, "failed to enable reconciler") - }) { - return - } - - if p.preflight != nil { - if !t.Run(p.preflightSubtestName, func(t *testing.T) { - p.preflight(t, ctx, log, test, client) - }) { - return - } - } - - if !t.Run(p.selectSubtestName, func(t *testing.T) { - device = p.selectDevice(t, ctx, log, test, client) - }) { - return - } - if device == nil { - // selectDevice skipped its subtest (e.g. no retransmit-only metro is - // configured on this network). t.Run reports a skipped subtest as - // success, so the skip does not stop the parent — skip it explicitly - // here rather than dereferencing a nil device in query_seat_price below. - // A nil device after a non-failed subtest can only mean the selector - // skipped: a failure would have made t.Run return false and returned - // above, and the success path always assigns a device. - t.Skip("Skipping: device selection skipped (feature not configured on this network)") - } - - if !t.Run("query_seat_price", func(t *testing.T) { - prices, err := client.FeedSeatPrice(ctx, device.PubKey) - require.NoError(t, err, "failed to get seat prices") - - // Match by pubkey, not code: querying by --device skips code resolution, - // so the returned rows may not carry a device_code. - for _, pr := range prices { - if pr.DevicePubkey == device.PubKey { - quoted = pr - break - } - } - require.NotNil(t, quoted, "no price found for device %s", device.Code) - require.NotZero(t, quoted.EpochPrice, "epoch price is zero for device %s", device.Code) - log.Info(p.priceLogMsg, "device", device.Code, - "epoch_price", quoted.EpochPrice, - "instant_allocation_price", quoted.GetInstantAllocationPrice(), - "reports_instant_allocation_price", quoted.GetReportsInstantAllocationPrice()) - }) { - return - } - - if !t.Run("query_onchain_seat_price", func(t *testing.T) { - // Read the prices the program itself computes straight off the chain, as - // an oracle for the CLI quote. This snapshot is taken next to the quote so - // the two comparisons below are as close to simultaneous as possible; the - // amount actually funded is re-read just before paying, since the wait for - // the open phase can outlive this read. - var err error - onchain, err = client.SeatPrices(ctx, device.PubKey) - require.NoError(t, err, "failed to compute the onchain seat prices") - require.NotZero(t, onchain.InstantAllocationDollars, "onchain instant-allocation price is zero for device %s", device.Code) - - // The price assertion belongs here, not on the CLI quote: what a caller - // means by "this seat costs the discounted price" is what the program - // charges, and that is the onchain figure. - if p.assertPrice != nil { - p.assertPrice(t, device, onchain.InstantAllocationDollars) - } - log.Info("Onchain seat prices", "device", device.Code, - "instant_allocation_price", onchain.InstantAllocationDollars, - "last_settled_epoch", onchain.LastSettledEpoch, - "current_epoch_price", onchain.CurrentEpochDollars, - "current_subscription_epoch", onchain.CurrentSubscriptionEpoch) - }) { - return - } - - // The next two subtests are deliberately not guarded with - // `if !t.Run(...) { return }`: a CLI/chain divergence is a real product bug - // and must fail the run loudly, but the payment below funds from the onchain - // price, so the rest of the settlement flow is still worth running. Both - // compare whole dollars, never prorated micro-USDC: proration is a function - // of the current slot, so an exact comparison against a separately-timed read - // is inherently racy. - t.Run("validate_instant_allocation_price_matches_chain", func(t *testing.T) { - switch { - case !quoted.GetReportsInstantAllocationPrice(): - // QA hosts install doublezero-solana from a version-pinned apt package - // (doublezero_solana_version in malbeclabs/infra - // ansible/inventory/*/group_vars/all.yml, 0.5.10-1 at time of writing), - // so the field only appears once a release carrying it is published and - // the pin bumped. Asserting against an absent field would read 0 and - // fail as "quoted 0, chain 43" — a misleading failure that looks like a - // new bug rather than a rollout gap. - t.Skipf("Skipping: installed doublezero-solana does not report instant_allocation_price (needs a release newer than the pinned 0.5.10-1); chain says %d USDC at last_settled_epoch=%d", - onchain.InstantAllocationDollars, onchain.LastSettledEpoch) - case quoted.InstantAllocationPrice == nil: - // Reported, but null: the CLI could not find the settled-epoch ring - // entry. That is a real condition, not a rollout artifact — the program - // performs the same lookup and would reject the allocation. - t.Fatalf("`shreds price` reported instant_allocation_price as unavailable for device %s, but the chain has a price of %d USDC at last_settled_epoch=%d", - device.Code, onchain.InstantAllocationDollars, onchain.LastSettledEpoch) - default: - require.Equal(t, onchain.InstantAllocationDollars, quoted.GetInstantAllocationPrice(), - "CLI quoted %d USDC for an instant allocation but the program charges %d USDC (last_settled_epoch=%d); `shreds price` and `shreds pay` read different ring entries", - quoted.GetInstantAllocationPrice(), onchain.InstantAllocationDollars, onchain.LastSettledEpoch) - } - }) - - t.Run("validate_epoch_price_matches_chain", func(t *testing.T) { - // The other half of the invariant: epoch_price must keep meaning the - // current-epoch price a recurring subscriber pays. Without this, someone - // "fixing" the divergence by repointing epoch_price at last_settled_epoch - // would go green here while silently breaking recurring subscribers. - if !onchain.HasCurrentEpoch { - // Part-way through UpdatingPrices the ring has not been advanced to the - // current epoch for every metro and device yet, so there is nothing to - // compare against. Transient by design, not a regression. - t.Skipf("Skipping: no onchain price entry yet for current_subscription_epoch=%d (prices are still being updated)", - onchain.CurrentSubscriptionEpoch) - } - require.Equal(t, onchain.CurrentEpochDollars, quoted.EpochPrice, - "CLI quoted epoch_price %d USDC but the chain has %d USDC at current_subscription_epoch=%d; epoch_price must stay the price a recurring subscriber pays next epoch", - quoted.EpochPrice, onchain.CurrentEpochDollars, onchain.CurrentSubscriptionEpoch) - }) - - // Set when wait_for_open_phase times out inside the epoch-tail closed - // window (verified against live chain state), so the parent can skip the - // remaining subtests: the program stays closed until the epoch boundary, - // which the 2-minute wait cannot bridge, so pay/withdraw below could only - // fail and page for a by-design condition. - var epochTailWindow *qa.EpochTailWindow - if !t.Run("wait_for_open_phase", func(t *testing.T) { - // Record where the wait begins: a timeout means the program was closed - // for the entire wait, so the classification below can require the - // whole span — not just the timeout-time slot — to be inside the - // window. Best-effort: on a read failure classification degrades to - // the timeout-time slot only. - waitStartSlot, slotErr := client.CurrentSolanaSlot(ctx) - if slotErr != nil { - log.Warn("Failed to read wait-start slot; epoch-tail classification will use the timeout-time slot only", "error", slotErr) - } - err := client.WaitForOpenForRequests(ctx) - if err != nil { - // For the last grace-period slots of every epoch the shred oracle - // closes the program by design (settle seats, update prices) and - // reopens it just after the epoch boundary. Verify against live - // chain state — onchain grace period, controller phase, and RPC - // epoch schedule — whether this timeout landed in that window; a - // timeout outside it must keep failing exactly as loudly as before. - win, winErr := client.EpochTailClosedWindow(ctx, waitStartSlot) - switch { - case winErr != nil: - log.Warn("Failed to classify epoch-tail closed window; treating timeout as a real failure", "error", winErr) - case win.Benign: - epochTailWindow = &win - t.Skipf("expected epoch-tail closed window: %s", win) - default: - // Give on-call the computed window so a real outage's distance - // from the benign window is visible in the run log. - log.Info("Timeout is not the benign epoch-tail closed window", "window", win.String()) - } - } - require.NoError(t, err, "shred-subscription program did not enter OpenForRequests phase within timeout") - }) { - return - } - if epochTailWindow != nil { - // t.Run reports a skipped subtest as success, so skip the parent - // explicitly to stop the run here. - t.Skipf("expected epoch-tail closed window: %s", epochTailWindow) - } - - if !t.Run("refresh_onchain_seat_price", func(t *testing.T) { - // wait_for_open_phase blocks for up to two minutes, and a settlement - // completing inside that window advances last_settled_epoch — the very - // read the charge is derived from. Funding a price captured before the - // wait would underfund the escrow and reproduce the opaque pay-time - // rejection this test exists to avoid, so the funded amount comes from a - // read taken here, immediately before paying. A rollover between this read - // and the transaction landing is irreducible (any payer races it), but the - // window shrinks from minutes to seconds. - refreshed, err := client.SeatPrices(ctx, device.PubKey) - require.NoError(t, err, "failed to re-read the onchain seat prices before paying") - require.NotZero(t, refreshed.InstantAllocationDollars, "onchain instant-allocation price is zero for device %s", device.Code) - - if refreshed.LastSettledEpoch != onchain.LastSettledEpoch || - refreshed.InstantAllocationDollars != onchain.InstantAllocationDollars { - // Warn, not Info: this means the quote comparisons above were made - // against a snapshot the payment no longer uses, so a failure up there - // should be read in that light. - log.Warn("Onchain seat price moved while waiting for the open-for-requests phase; funding the refreshed price", - "device", device.Code, - "before_price", onchain.InstantAllocationDollars, "before_last_settled_epoch", onchain.LastSettledEpoch, - "after_price", refreshed.InstantAllocationDollars, "after_last_settled_epoch", refreshed.LastSettledEpoch) - } - - onchain = refreshed - amount = strconv.FormatUint(onchain.InstantAllocationDollars, 10) - fundedAmount = onchain.InstantAllocationDollars * 1_000_000 // dollars to USDC raw units (6 decimals) - log.Info("Funding the seat escrow from the onchain price", "device", device.Code, - "amount", amount, "last_settled_epoch", onchain.LastSettledEpoch) - }) { - return - } - - if !t.Run("record_balance_before_pay", func(t *testing.T) { - var err error - balanceBeforePay, err = client.GetUSDCBalance(ctx) - require.NoError(t, err, "failed to get USDC balance before pay") - log.Info("USDC balance before pay", "balance", balanceBeforePay) - }) { - return - } - - if !t.Run("pay_for_seat", func(t *testing.T) { - err := client.FeedSeatPay(ctx, device.PubKey, amount) - require.NoError(t, err, "failed to pay for seat") - seatPaid = true - }) { - return - } - - if !t.Run("validate_balance_after_pay", func(t *testing.T) { - // Poll until the balance reflects the debit. FeedSeatPay returns - // after the tx is submitted, and the RPC balance view can lag the - // confirmed state briefly, so a one-shot read races. - var lastDebit uint64 - require.Eventually(t, func() bool { - bal, err := client.GetUSDCBalance(ctx) - if err != nil { - log.Info("USDC balance poll error", "error", err) - return false - } - balanceAfterPay = bal - lastDebit = balanceBeforePay - bal - return lastDebit == fundedAmount - }, balanceSettleTimeout, 5*time.Second, "USDC balance should decrease by the paid amount") - log.Info("USDC balance after pay", "balance", balanceAfterPay, "debit", lastDebit, "expected_debit", fundedAmount) - }) { - return - } - - if !t.Run("query_effective_seat_price", func(t *testing.T) { - // Built on the onchain price, not the CLI quote: this feeds the - // non-prorating balance assertion below, which must predict what the - // program actually charged. GetEffectiveSeatPrice applies the seat's price - // override on top when one is set. - var err error - effectivePrice, err = client.GetEffectiveSeatPrice(ctx, device.PubKey, onchain.InstantAllocationDollars) - require.NoError(t, err, "failed to get effective seat price") - log.Info("Effective seat price", "effective_usdc", effectivePrice, "funded_usdc", fundedAmount) - }) { - return - } - - if !t.Run("validate_tunnel_up", func(t *testing.T) { - err := client.WaitForMulticastStatusUp(ctx) - require.NoError(t, err, "multicast tunnel did not come up after seat payment") - }) { - return - } - - if !t.Run("validate_device_assignment", func(t *testing.T) { - statuses, err := client.GetUserStatuses(ctx) - require.NoError(t, err, "failed to get user statuses") - mcastStatus := qa.FindMulticastStatus(statuses) - require.NotNil(t, mcastStatus, "no multicast status found after seat payment") - require.Equal(t, device.Code, mcastStatus.CurrentDevice, "tunnel connected to wrong device") - log.Info("Tunnel up and device matches", "device", mcastStatus.CurrentDevice, "dzIP", mcastStatus.DoubleZeroIp) - }) { - return - } - - if p.extraAssertion != nil { - if !t.Run(p.extraSubtestName, func(t *testing.T) { - p.extraAssertion(t, ctx, log, client, device) - }) { - return - } - } - - if !t.Run("withdraw_seat", func(t *testing.T) { - // Withdraw is rejected while this run's instant allocation request is - // in flight (or a stale RPC read claims it is), so retry with endpoint - // rotation rather than waiting on an ack the harness cannot observe - // reliably. - err := client.WithdrawSeatWithRetry(ctx, device.PubKey) - require.NoError(t, err, "failed to withdraw seat") - seatPaid = false - }) { - return - } - - if !t.Run("validate_tunnel_down", func(t *testing.T) { - err := client.WaitForMulticastStatusDisconnected(ctx) - require.NoError(t, err, "tunnel did not come down after seat withdrawal") - }) { - return - } - - t.Run("validate_balance_after_withdraw", func(t *testing.T) { - // Read onchain whether the shred-subscription program has prorated - // service enabled. This lets the test self-adapt across environments - // (testnet has it on, mainnet does not) without needing a CI flag. - proratingEnabled, err := client.IsSeatProratingEnabled(ctx) - require.NoError(t, err, "failed to read prorating flag from program config") - - var balanceAfterWithdraw uint64 - if proratingEnabled { - // Prorating refunds the unused portion of the epoch to the wallet. - // Poll until the refund is reflected (balance strictly greater - // than after-pay). - require.Eventually(t, func() bool { - bal, err := client.GetUSDCBalance(ctx) - if err != nil { - log.Info("USDC balance poll error", "error", err) - return false - } - balanceAfterWithdraw = bal - return bal > balanceAfterPay - }, balanceSettleTimeout, 5*time.Second, - "USDC balance should increase to reflect the prorated refund") - } else { - expectedBalance := balanceBeforePay - effectivePrice - require.Eventually(t, func() bool { - bal, err := client.GetUSDCBalance(ctx) - if err != nil { - log.Info("USDC balance poll error", "error", err) - return false - } - balanceAfterWithdraw = bal - return bal == expectedBalance - }, balanceSettleTimeout, 5*time.Second, - "USDC balance should equal before_pay minus the effective seat price") - } - - refund := balanceAfterWithdraw - balanceAfterPay - - // A seat's payment escrow can carry a balance from an earlier run whose - // withdraw did not complete (e.g. during the reservoir-ack outage on - // devnet). Closing the escrow now refunds that leftover too, so the - // wallet-measured refund exceeds what was paid this run and no longer - // isolates this payment (`retained` would underflow). In that case the - // wallet-delta proration check is not meaningful, so skip it rather than - // fail — the settlement path itself is still covered by the pay/ack/ - // tunnel/withdraw sub-tests above. - if refund > fundedAmount { - log.Warn("skipping wallet-delta proration check: refund exceeds amount paid this run (pre-existing escrow drained)", - "refund", refund, - "paid_amount", fundedAmount, - "before_pay", balanceBeforePay, - "after_pay", balanceAfterPay, - "after_withdraw", balanceAfterWithdraw, - ) - return - } - // Equivalent to balanceBeforePay - balanceAfterWithdraw, but computed from - // the amount paid this run so it cannot underflow given the guard above. - retained := fundedAmount - refund - - log.Info("USDC balance after withdraw", - "balance", balanceAfterWithdraw, - "before_pay", balanceBeforePay, - "after_pay", balanceAfterPay, - "paid_amount", fundedAmount, - "effective_price", effectivePrice, - "refund", refund, - "retained", retained, - "prorating_enabled", proratingEnabled, - ) - - // Accounting invariant: regardless of prorating, the sum of what was - // refunded to the wallet and what the program retained must equal the - // amount debited at pay time. This uses fundedAmount rather than - // effectivePrice because a seat with a zero price override is still - // charged fundedAmount at pay and fully refunded on withdraw. - require.Equal(t, fundedAmount, refund+retained, - "refund + retained must equal the amount paid") - - if !proratingEnabled || effectivePrice == 0 { - return - } - - // With prorating enabled we avoid replicating the onchain formula - // against client-side RPC state (epoch schedule + current epoch reads - // are fragile on DZ ledger). Instead assert the qualitative invariants - // that distinguish a real partial refund from a regression: - // - refund > 0 (prorating actually happened) - // - retained > 0 (the seat was not free for the used portion) - // - retained < effective_price (kept less than a full epoch) - require.Greater(t, refund, uint64(0), - "prorating: refund should be strictly greater than zero") - require.Greater(t, retained, uint64(0), - "prorating: retained should be strictly greater than zero") - require.Less(t, retained, effectivePrice, - "prorating: retained should be strictly less than the effective price") - }) -} From cb3ec334b63de285e61d1f2c9e3933164ea8c3c3 Mon Sep 17 00:00:00 2001 From: Martin Sander Date: Thu, 6 Aug 2026 10:31:14 -0500 Subject: [PATCH 3/5] qa: address review on the retransmit-only settlement test The rejection payment now runs after wait_for_open_phase. The program is closed for the last grace-period slots of every epoch, so a payment made before the wait fails for the phase and never carries the onboarding message the test asserts on. The settlement flow is inlined into the test. It had one caller after the retransmit-only test went away, and the hook indirection hid the branch. The onboarding flag is read outside every subtest, so a -run filter can no longer skip the read and take the closest-device path in silence. The price assertion is back as --retransmit-price, checked against the price the program charges rather than the CLI quote. TestProgramConfigFlags tables all four program config flag accessors, including bit 7 for retransmit-only onboarding. ClosestRetransmitOnlyDevice and ClosestNonRetransmitOnlyDevice now share one selector taking the membership it wants. --- CHANGELOG.md | 2 +- e2e/internal/qa/client_settlement.go | 38 +++----- e2e/qa_multicast_settlement_test.go | 131 ++++++++++----------------- sdk/shreds/go/state_test.go | 37 ++++++++ 4 files changed, 100 insertions(+), 108 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 0ef8f6ea46..e3b6d1e686 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -23,7 +23,7 @@ All notable changes to this project will be documented in this file. - Samples dropped because a partition's onchain account is full are now counted too, under `reason="account_full"` plus `submitter_account_full` on the errors counter. That path reports success to its caller, so a warning was the only trace, and its count was wrong: it reported the whole flushed partition rather than the samples actually lost. (#4145) - A failed submission now retries from the first unwritten sample rather than restarting at the beginning of the flushed partition, and only the unwritten remainder is requeued. Previously a mid-partition error made every subsequent attempt re-send batches that were already onchain, appending those samples a second time and pushing the account toward the sample cap it is measured against. (#4145) - Agent logs now identify the ledger RPC endpoint in use, and the peer count and the stale-program-data warning report state transitions instead of firing on every refresh. New lines name the resolved remote address of each connection (a bad load balancer address behind a hostname was invisible before), and a refresh that finds no peers is now a warning rather than a Debug line. The stale-cache warning also moves off the package-global `slog` onto the agent's own logger, so it is formatted and leveled with everything else. New: `doublezero_device_telemetry_agent_peers` gauge, and `pinger_epoch_fetch` on the errors counter for every exhausted epoch fetch. (#4147) -QA +- QA - Rework existing TestQA_MulticastSettlement and adapt it to the new `FLAG_RETRANSMIT_ONLY_ONBOARDING_ENFORCED_BIT` flag. Test now checks for this flag in the ProgramConfig solana account and depending on if it's on or off tries to assert that no new user can subscribe to a non retransmit-only metro unless that metro has the retransmit-only flag enabled in the MetroHistory account. (#4156) ## [v0.33.0](https://github.com/malbeclabs/doublezero/compare/client/v0.32.0...client/v0.33.0) - 2026-07-31 diff --git a/e2e/internal/qa/client_settlement.go b/e2e/internal/qa/client_settlement.go index 74085cff70..85d2f891ff 100644 --- a/e2e/internal/qa/client_settlement.go +++ b/e2e/internal/qa/client_settlement.go @@ -164,38 +164,27 @@ func (c *Client) ClosestRetransmitOnlyDevice(ctx context.Context) (*Device, map[ return nil, retransmitOnly, nil } - latencies, err := c.GetLatency(ctx) + device, err := c.closestDeviceInMetros(ctx, retransmitOnly, true) if err != nil { - return nil, retransmitOnly, fmt.Errorf("failed to get latency on host %s: %w", c.Host, err) + return nil, retransmitOnly, err } - - var bestDevice *Device - var bestAvg uint64 = math.MaxUint64 - for _, l := range latencies { - if !l.Reachable { - continue - } - device, ok := c.devices[l.DeviceCode] - if !ok || !retransmitOnly[device.ExchangePubKey] { - continue - } - if l.AvgLatencyNs < bestAvg { - bestAvg = l.AvgLatencyNs - bestDevice = device - } - } - if bestDevice != nil { - c.log.Debug("Determined closest retransmit-only device", "host", c.Host, "deviceCode", bestDevice.Code, "avgLatencyNs", bestAvg) - } - return bestDevice, retransmitOnly, nil + return device, retransmitOnly, nil } +// ClosestNonRetransmitOnlyDevice returns the reachable device with the lowest +// average latency whose metro is not flagged retransmit-only. A nil device means +// every reachable metro is flagged, so no metro is left to reject a new seat. func (c *Client) ClosestNonRetransmitOnlyDevice(ctx context.Context) (*Device, error) { retransmitOnly, err := c.RetransmitOnlyExchangeKeys(ctx) if err != nil { return nil, err } + return c.closestDeviceInMetros(ctx, retransmitOnly, false) +} +// closestDeviceInMetros returns the lowest-latency reachable device whose metro +// membership in exchangeKeys equals want. +func (c *Client) closestDeviceInMetros(ctx context.Context, exchangeKeys map[string]bool, want bool) (*Device, error) { latencies, err := c.GetLatency(ctx) if err != nil { return nil, fmt.Errorf("failed to get latency on host %s: %w", c.Host, err) @@ -208,7 +197,7 @@ func (c *Client) ClosestNonRetransmitOnlyDevice(ctx context.Context) (*Device, e continue } device, ok := c.devices[l.DeviceCode] - if !ok || retransmitOnly[device.ExchangePubKey] { + if !ok || exchangeKeys[device.ExchangePubKey] != want { continue } if l.AvgLatencyNs < bestAvg { @@ -217,7 +206,8 @@ func (c *Client) ClosestNonRetransmitOnlyDevice(ctx context.Context) (*Device, e } } if bestDevice != nil { - c.log.Debug("Determined closest non-retransmit-only device", "host", c.Host, "deviceCode", bestDevice.Code, "avgLatencyNs", bestAvg) + c.log.Debug("Determined closest device", "host", c.Host, "deviceCode", bestDevice.Code, + "avgLatencyNs", bestAvg, "retransmitOnly", want) } return bestDevice, nil } diff --git a/e2e/qa_multicast_settlement_test.go b/e2e/qa_multicast_settlement_test.go index 3a4f40d4bf..4543be6bdf 100644 --- a/e2e/qa_multicast_settlement_test.go +++ b/e2e/qa_multicast_settlement_test.go @@ -27,49 +27,11 @@ var ( enableSettlementTests = flag.Bool("enable-multicast-settlement-tests", false, "enable multicast settlement tests") retransmitOnlyDeviceFlag = flag.String("retransmit-only-device", "", "device code or pubkey in a retransmit-only metro (overrides auto-discovery)") retransmitGroupCodesFlag = flag.String("retransmit-group-codes", "", "comma-separated multicast group codes a seat in a retransmit-only metro must subscribe to, and nothing else") + retransmitPriceFlag = flag.Uint64("retransmit-price", 0, "expected seat price in whole USDC dollars in a retransmit-only metro; 0 asserts nothing") keypairFlag = flag.String("keypair", "$HOME/.config/doublezero/id.json", "path to keypair file for settlement commands") settlementClientFlag = flag.String("multicast-settlement-client", "", "host of the client to use for settlement tests (overrides random selection)") ) -func TestQA_MulticastSettlement(t *testing.T) { - var onboardingEnforced bool - runShredSettlement(t, shredSettlementParams{ - enabled: *enableSettlementTests, - skipReason: "Skipping: --enable-multicast-settlement-tests flag not set", - - preflightSubtestName: "reject_new_seat_outside_retransmit_only_metro", - preflight: func(t *testing.T, ctx context.Context, log *slog.Logger, _ *qa.Test, client *qa.Client) { - var err error - onboardingEnforced, err = client.IsRetransmitOnlyOnboardingEnforced(ctx) - require.NoError(t, err, "failed to read the retransmit-only onboarding flag") - if !onboardingEnforced { - log.Info("Retransmit-only onboarding is off; settling on the closest device") - t.Skip("Skipping: the program config does not enforce retransmit-only onboarding") - } - log.Info("Retransmit-only onboarding is on; a new seat must be rejected outside a retransmit-only metro") - assertNewSeatRejected(t, ctx, log, client) - }, - - selectSubtestName: "select_device", - selectDevice: func(t *testing.T, ctx context.Context, log *slog.Logger, test *qa.Test, client *qa.Client) *qa.Device { - if onboardingEnforced { - return selectRetransmitOnlyDevice(t, ctx, log, test, client) - } - return selectClosestDevice(t, ctx, log, test, client) - }, - - priceLogMsg: "Found epoch price", - - extraSubtestName: "assert_subscribed_groups", - extraAssertion: func(t *testing.T, ctx context.Context, log *slog.Logger, client *qa.Client, device *qa.Device) { - if !onboardingEnforced { - t.Skip("Skipping: the seat is not in a retransmit-only metro, so it carries the leader group too") - } - assertSubscribedGroups(t, ctx, log, client, device) - }, - }) -} - func selectClosestDevice(t *testing.T, ctx context.Context, log *slog.Logger, _ *qa.Test, client *qa.Client) *qa.Device { device, err := client.ClosestDevice(ctx) require.NoError(t, err, "failed to find closest device") @@ -77,6 +39,9 @@ func selectClosestDevice(t *testing.T, ctx context.Context, log *slog.Logger, _ return device } +// assertNewSeatRejected pays for a seat outside a retransmit-only metro and +// requires the program to refuse it. The client holds no seat at this point, so +// the seat has no tenure and onboarding enforcement applies to it. func assertNewSeatRejected(t *testing.T, ctx context.Context, log *slog.Logger, client *qa.Client) { device, err := client.ClosestNonRetransmitOnlyDevice(ctx) require.NoError(t, err, "failed to find a device outside a retransmit-only metro") @@ -227,31 +192,19 @@ func assertSubscribedGroups(t *testing.T, ctx context.Context, log *slog.Logger, ) } -type shredSettlementParams struct { - enabled bool - skipReason string - - // preflight runs after the reconciler is enabled and before device selection. - preflightSubtestName string - preflight func(t *testing.T, ctx context.Context, log *slog.Logger, test *qa.Test, client *qa.Client) - - selectSubtestName string - // selectDevice may call t.Skip when the feature is not configured. It logs - // its own selection detail. - selectDevice func(t *testing.T, ctx context.Context, log *slog.Logger, test *qa.Test, client *qa.Client) *qa.Device - - priceLogMsg string - - // extraAssertion runs after the tunnel is up and before the seat is withdrawn. - extraSubtestName string - extraAssertion func(t *testing.T, ctx context.Context, log *slog.Logger, client *qa.Client, device *qa.Device) -} - -// runShredSettlement picks a device, pays its seat price, checks the debit and -// the tunnel, then withdraws the seat and checks the refund. -func runShredSettlement(t *testing.T, p shredSettlementParams) { - if !p.enabled { - t.Skip(p.skipReason) +// TestQA_MulticastSettlement picks a device, queries its seat price from the CLI +// and cross-checks it against the price the program will charge (read from +// chain), waits for the open-for-requests phase, re-reads that price and pays +// it, checks the debit and the tunnel, then withdraws the seat and checks the +// refund accounting. +// +// The program config decides which shape the run takes. When it enforces +// retransmit-only onboarding, the test also requires the program to reject a new +// seat outside a retransmit-only metro, and it settles inside one. Otherwise it +// settles on the closest device, which is what mainnet does today. +func TestQA_MulticastSettlement(t *testing.T) { + if !*enableSettlementTests { + t.Skip("Skipping: --enable-multicast-settlement-tests flag not set") } log := newTestLogger(t) @@ -272,6 +225,10 @@ func runShredSettlement(t *testing.T, p shredSettlementParams) { } log.Info("Selected client", "host", client.Host) + retransmitOnboardingEnforced, err := client.IsRetransmitOnlyOnboardingEnforced(ctx) + require.NoError(t, err, "failed to read the retransmit-only onboarding flag") + log.Info("Retransmit-only onboarding", "enforced", retransmitOnboardingEnforced) + var device *qa.Device var amount string var quoted *pb.DevicePrice @@ -364,27 +321,18 @@ func runShredSettlement(t *testing.T, p shredSettlementParams) { return } - if p.preflight != nil { - if !t.Run(p.preflightSubtestName, func(t *testing.T) { - p.preflight(t, ctx, log, test, client) - }) { + if !t.Run("select_device", func(t *testing.T) { + if retransmitOnboardingEnforced { + device = selectRetransmitOnlyDevice(t, ctx, log, test, client) return } - } - - if !t.Run(p.selectSubtestName, func(t *testing.T) { - device = p.selectDevice(t, ctx, log, test, client) + device = selectClosestDevice(t, ctx, log, test, client) }) { return } if device == nil { - // selectDevice skipped its subtest (e.g. no retransmit-only metro is - // configured on this network). t.Run reports a skipped subtest as - // success, so the skip does not stop the parent — skip it explicitly - // here rather than dereferencing a nil device in query_seat_price below. - // A nil device after a non-failed subtest can only mean the selector - // skipped: a failure would have made t.Run return false and returned - // above, and the success path always assigns a device. + // The selector skipped its subtest, because no metro is flagged + // retransmit-only on this network. t.Skip("Skipping: device selection skipped (feature not configured on this network)") } @@ -402,7 +350,7 @@ func runShredSettlement(t *testing.T, p shredSettlementParams) { } require.NotNil(t, quoted, "no price found for device %s", device.Code) require.NotZero(t, quoted.EpochPrice, "epoch price is zero for device %s", device.Code) - log.Info(p.priceLogMsg, "device", device.Code, + log.Info("Found epoch price", "device", device.Code, "epoch_price", quoted.EpochPrice, "instant_allocation_price", quoted.GetInstantAllocationPrice(), "reports_instant_allocation_price", quoted.GetReportsInstantAllocationPrice()) @@ -421,6 +369,13 @@ func runShredSettlement(t *testing.T, p shredSettlementParams) { require.NoError(t, err, "failed to compute the onchain seat prices") require.NotZero(t, onchain.InstantAllocationDollars, "onchain instant-allocation price is zero for device %s", device.Code) + // Pin the price the program charges, not the CLI argument. + if retransmitOnboardingEnforced && *retransmitPriceFlag != 0 { + require.Equal(t, *retransmitPriceFlag, onchain.InstantAllocationDollars, + "device %s in retransmit-only metro %s should cost the price --retransmit-price names", + device.Code, device.ExchangeCode) + } + log.Info("Onchain seat prices", "device", device.Code, "instant_allocation_price", onchain.InstantAllocationDollars, "last_settled_epoch", onchain.LastSettledEpoch, @@ -526,6 +481,14 @@ func runShredSettlement(t *testing.T, p shredSettlementParams) { t.Skipf("expected epoch-tail closed window: %s", epochTailWindow) } + if retransmitOnboardingEnforced { + if !t.Run("reject_new_seat_outside_retransmit_only_metro", func(t *testing.T) { + assertNewSeatRejected(t, ctx, log, client) + }) { + return + } + } + if !t.Run("refresh_onchain_seat_price", func(t *testing.T) { // wait_for_open_phase blocks for up to two minutes, and a settlement // completing inside that window advances last_settled_epoch — the very @@ -627,9 +590,11 @@ func runShredSettlement(t *testing.T, p shredSettlementParams) { return } - if p.extraAssertion != nil { - if !t.Run(p.extraSubtestName, func(t *testing.T) { - p.extraAssertion(t, ctx, log, client, device) + // The seat carries the leader group in any other metro, so the group + // assertion only means something once the metro is retransmit-only. + if retransmitOnboardingEnforced { + if !t.Run("assert_subscribed_groups", func(t *testing.T) { + assertSubscribedGroups(t, ctx, log, client, device) }) { return } diff --git a/sdk/shreds/go/state_test.go b/sdk/shreds/go/state_test.go index 88312205f2..1a6f48a9f4 100644 --- a/sdk/shreds/go/state_test.go +++ b/sdk/shreds/go/state_test.go @@ -308,6 +308,43 @@ func TestMetroHistoryFlags(t *testing.T) { } } +// The bit indices mirror the shred-subscription program's ProgramConfig +// constants: FLAG_IS_PAUSED_BIT 0, FLAG_PRORATED_SERVICE_ENABLED_BIT 2 and +// FLAG_RETRANSMIT_ONLY_ONBOARDING_ENFORCED_BIT 7. A wrong index here reads as +// "feature off" rather than as an error, so nothing else catches it. +func TestProgramConfigFlags(t *testing.T) { + cases := []struct { + flags uint64 + paused bool + migrated bool + prorated bool + onboardingEnforced bool + }{ + {flags: 0}, + {flags: 1 << 0, paused: true}, + {flags: 1 << 1, migrated: true}, + {flags: 1 << 2, prorated: true}, + {flags: 1 << 7, onboardingEnforced: true}, + {flags: 0x7c, prorated: true}, + {flags: 0xff, paused: true, migrated: true, prorated: true, onboardingEnforced: true}, + } + for _, tc := range cases { + cfg := ProgramConfig{Flags: tc.flags} + if got := cfg.IsPaused(); got != tc.paused { + t.Errorf("flags=%#b: IsPaused() = %v, want %v", tc.flags, got, tc.paused) + } + if got := cfg.IsMigrated(); got != tc.migrated { + t.Errorf("flags=%#b: IsMigrated() = %v, want %v", tc.flags, got, tc.migrated) + } + if got := cfg.IsProratedServiceEnabled(); got != tc.prorated { + t.Errorf("flags=%#b: IsProratedServiceEnabled() = %v, want %v", tc.flags, got, tc.prorated) + } + if got := cfg.IsRetransmitOnlyOnboardingEnforced(); got != tc.onboardingEnforced { + t.Errorf("flags=%#b: IsRetransmitOnlyOnboardingEnforced() = %v, want %v", tc.flags, got, tc.onboardingEnforced) + } + } +} + func TestDeviceHistoryDeserialization(t *testing.T) { data := make([]byte, unsafe.Sizeof(DeviceHistory{})) // Device key From b32b81511735cf7584df4bb5f783130c053429cf Mon Sep 17 00:00:00 2001 From: Martin Sander Date: Thu, 6 Aug 2026 10:40:03 -0500 Subject: [PATCH 4/5] qa: read the multicast user for the subscribed-groups check A host running -multi-tunnel holds an IBRL user at the same client IP, and GetServiceabilityUser returns whichever account comes first. Reading that one leaves the subscriber list empty and the group check reports missing groups. --- e2e/internal/qa/client.go | 15 +++++++++++++++ e2e/qa_multicast_settlement_test.go | 2 +- 2 files changed, 16 insertions(+), 1 deletion(-) diff --git a/e2e/internal/qa/client.go b/e2e/internal/qa/client.go index 6c3e44d74e..0e344b21eb 100644 --- a/e2e/internal/qa/client.go +++ b/e2e/internal/qa/client.go @@ -507,6 +507,21 @@ func (c *Client) GetServiceabilityUser(ctx context.Context) (*serviceability.Use return nil, fmt.Errorf("serviceability user not found for client IP %s on host %s", publicIP, c.Host) } +func (c *Client) GetMulticastServiceabilityUser(ctx context.Context) (*serviceability.User, error) { + data, err := getProgramDataWithRetry(ctx, c.serviceability) + if err != nil { + return nil, fmt.Errorf("failed to get program data on host %s: %w", c.Host, err) + } + publicIP := c.publicIP.To4().String() + for i := range data.Users { + user := &data.Users[i] + if net.IP(user.ClientIp[:]).String() == publicIP && user.UserType == serviceability.UserTypeMulticast { + return user, nil + } + } + return nil, fmt.Errorf("multicast serviceability user not found for client IP %s on host %s", publicIP, c.Host) +} + func (c *Client) GetOwnerPubkey(ctx context.Context) (solana.PublicKey, error) { user, err := c.GetServiceabilityUser(ctx) if err != nil { diff --git a/e2e/qa_multicast_settlement_test.go b/e2e/qa_multicast_settlement_test.go index 4543be6bdf..9ca811ec95 100644 --- a/e2e/qa_multicast_settlement_test.go +++ b/e2e/qa_multicast_settlement_test.go @@ -150,7 +150,7 @@ func assertSubscribedGroups(t *testing.T, ctx context.Context, log *slog.Logger, lastExtra []string ) ok := assert.Eventually(t, func() bool { - user, err := client.GetServiceabilityUser(ctx) + user, err := client.GetMulticastServiceabilityUser(ctx) if err != nil { log.Info("serviceability user poll error", "error", err) return false From 72e5802bf505d818026b90fc8c80debd22d214df Mon Sep 17 00:00:00 2001 From: Martin Sander Date: Thu, 6 Aug 2026 11:10:20 -0500 Subject: [PATCH 5/5] qa: name the multicast groups by code in the group check A failure printed raw pubkeys, which is what made the first report of this check unreadable. One read of the program data now maps every group pubkey to its code, and the check labels both the extras and the subscribed list from it. That read also replaces the per-code lookups, each of which fetched the whole program data again. The seat-price assertion logs at Info when --retransmit-price is unset, so a run cannot look like it checked the price when it did not. --- e2e/internal/qa/client_multicast.go | 14 ++++++++++ e2e/qa_multicast_settlement_test.go | 41 ++++++++++++++++++++++------- 2 files changed, 45 insertions(+), 10 deletions(-) diff --git a/e2e/internal/qa/client_multicast.go b/e2e/internal/qa/client_multicast.go index f46beb875a..a1e337be30 100644 --- a/e2e/internal/qa/client_multicast.go +++ b/e2e/internal/qa/client_multicast.go @@ -141,6 +141,20 @@ func (c *Client) GetMulticastGroup(ctx context.Context, code string) (*Multicast return nil, nil } +// MulticastGroupCodes maps every multicast group pubkey to its code, so a caller +// holding pubkeys off a user account can name them in a log or a failure. +func (c *Client) MulticastGroupCodes(ctx context.Context) (map[solana.PublicKey]string, error) { + data, err := getProgramDataWithRetry(ctx, c.serviceability) + if err != nil { + return nil, fmt.Errorf("failed to get program data on host %s: %w", c.Host, err) + } + codes := make(map[solana.PublicKey]string, len(data.MulticastGroups)) + for _, group := range data.MulticastGroups { + codes[solana.PublicKeyFromBytes(group.PubKey[:])] = group.Code + } + return codes, nil +} + func (c *Client) CreateMulticastGroup(ctx context.Context, code string, maxBandwidth string) (*MulticastGroup, error) { c.log.Debug("Creating multicast group", "host", c.Host, "code", code, "maxBandwidth", maxBandwidth) resp, err := c.grpcClient.CreateMulticastGroup(ctx, &pb.CreateMulticastGroupRequest{ diff --git a/e2e/qa_multicast_settlement_test.go b/e2e/qa_multicast_settlement_test.go index 9ca811ec95..a462247933 100644 --- a/e2e/qa_multicast_settlement_test.go +++ b/e2e/qa_multicast_settlement_test.go @@ -129,19 +129,34 @@ func flaggedMetroCodes(test *qa.Test, exchangeKeys map[string]bool) []string { func assertSubscribedGroups(t *testing.T, ctx context.Context, log *slog.Logger, client *qa.Client, _ *qa.Device) { // Nothing onchain labels a group as leader or retransmit, so the operator // names the retransmit groups per network. + groupCodes, err := client.MulticastGroupCodes(ctx) + require.NoError(t, err, "failed to read the multicast groups") + keysByCode := make(map[string]solana.PublicKey, len(groupCodes)) + for key, code := range groupCodes { + keysByCode[code] = key + } + required := make(map[solana.PublicKey]string) for _, code := range strings.Split(*retransmitGroupCodesFlag, ",") { code = strings.TrimSpace(code) if code == "" { continue } - group, err := client.GetMulticastGroup(ctx, code) - require.NoError(t, err, "failed to resolve multicast group %q", code) - require.NotNil(t, group, "multicast group %q not found onchain", code) - required[group.PK] = code + key, ok := keysByCode[code] + require.True(t, ok, "multicast group %q not found onchain", code) + required[key] = code } require.NotEmpty(t, required, "no multicast group resolved from --retransmit-group-codes %q", *retransmitGroupCodesFlag) + // A failure names the groups by code. The raw pubkeys are what made the + // first report of this check unreadable. + label := func(group solana.PublicKey) string { + if code, ok := groupCodes[group]; ok { + return code + } + return group.String() + } + // The oracle converges the seat's onchain subscription asynchronously, so // poll until it reflects retransmit-only membership. var ( @@ -160,7 +175,7 @@ func assertSubscribedGroups(t *testing.T, ctx context.Context, log *slog.Logger, for _, sub := range user.Subscribers { group := solana.PublicKeyFromBytes(sub[:]) subscribed[group] = true - subs = append(subs, group.String()) + subs = append(subs, label(group)) } var missing, extra []string for group, code := range required { @@ -170,7 +185,7 @@ func assertSubscribedGroups(t *testing.T, ctx context.Context, log *slog.Logger, } for group := range subscribed { if _, want := required[group]; !want { - extra = append(extra, group.String()) + extra = append(extra, label(group)) } } lastSubscribed, lastMissing, lastExtra = subs, missing, extra @@ -370,10 +385,16 @@ func TestQA_MulticastSettlement(t *testing.T) { require.NotZero(t, onchain.InstantAllocationDollars, "onchain instant-allocation price is zero for device %s", device.Code) // Pin the price the program charges, not the CLI argument. - if retransmitOnboardingEnforced && *retransmitPriceFlag != 0 { - require.Equal(t, *retransmitPriceFlag, onchain.InstantAllocationDollars, - "device %s in retransmit-only metro %s should cost the price --retransmit-price names", - device.Code, device.ExchangeCode) + if retransmitOnboardingEnforced { + if *retransmitPriceFlag == 0 { + log.Info("Not asserting the seat price; --retransmit-price is unset", + "device", device.Code, "metro", device.ExchangeCode, + "price", onchain.InstantAllocationDollars) + } else { + require.Equal(t, *retransmitPriceFlag, onchain.InstantAllocationDollars, + "device %s in retransmit-only metro %s should cost the price --retransmit-price names", + device.Code, device.ExchangeCode) + } } log.Info("Onchain seat prices", "device", device.Code,