1
0

6 کامیت‌ها 17d7dd46b5 ... 8c023d13dc

نویسنده SHA1 پیام تاریخ
  MHSanaei 8c023d13dc docs(architecture): fix table padding flagged by oxfmt 12 ساعت پیش
  MHSanaei 823db05966 fix(inbounds): keep the stored client list and enable on inbound save 13 ساعت پیش
  MHSanaei fb7418f7bd fix(mtproto): zero sidecar quotas only for clients whose usage was reset 13 ساعت پیش
  MHSanaei 4210a50cb4 fix(traffic): make a client reset reach every counter enforcing its quota 23 ساعت پیش
  MHSanaei 7c84ca9689 fix(runtime): drop depleted clients by email, not by stale inbound_id 1 روز پیش
  MHSanaei 1110caaa65 fix(inbounds): keep stored client lifecycle and counters on inbound save 1 روز پیش
35فایلهای تغییر یافته به همراه1453 افزوده شده و 229 حذف شده
  1. 3 0
      docs/architecture.md
  2. 10 8
      docs/content/docs/en/reference/api/inbounds.mdx
  3. 1 1
      docs/public/openapi.json
  4. 1 1
      frontend/public/openapi.json
  5. 15 2
      frontend/src/lib/xray/inbound-form-adapter.ts
  6. 1 1
      frontend/src/pages/api-docs/endpoints.ts
  7. 16 13
      frontend/src/pages/inbounds/form/InboundFormModal.tsx
  8. 30 0
      frontend/src/test/inbound-form-modal.test.tsx
  9. 1 0
      internal/database/db.go
  10. 1 0
      internal/database/migrate_data.go
  11. 11 0
      internal/database/model/node_pending_reset.go
  12. 80 0
      internal/web/job/node_reset_replay_test.go
  13. 7 0
      internal/web/job/node_traffic_sync_job.go
  14. 6 0
      internal/web/runtime/remote.go
  15. 33 0
      internal/web/runtime/remote_reset_test.go
  16. 2 1
      internal/web/service/client_renewal_write_validation_test.go
  17. 33 7
      internal/web/service/client_traffic.go
  18. 56 39
      internal/web/service/inbound.go
  19. 22 32
      internal/web/service/inbound_amneziawg.go
  20. 18 0
      internal/web/service/inbound_disable.go
  21. 6 3
      internal/web/service/inbound_hysteria_auth_test.go
  22. 58 59
      internal/web/service/inbound_mtproto.go
  23. 2 1
      internal/web/service/inbound_mtproto_apply_test.go
  24. 15 6
      internal/web/service/inbound_node.go
  25. 145 0
      internal/web/service/inbound_runtime_depleted_filter_test.go
  26. 23 0
      internal/web/service/inbound_settings_commit.go
  27. 21 23
      internal/web/service/inbound_traffic.go
  28. 2 0
      internal/web/service/inbound_traffic_apply.go
  29. 22 32
      internal/web/service/inbound_tuic.go
  30. 170 0
      internal/web/service/inbound_update_stale_form_test.go
  31. 75 0
      internal/web/service/mtproto_fake_test.go
  32. 137 0
      internal/web/service/mtproto_quota_reset_test.go
  33. 3 0
      internal/web/service/node.go
  34. 152 0
      internal/web/service/node_reset_queue.go
  35. 275 0
      internal/web/service/node_reset_undelivered_test.go

+ 3 - 0
docs/architecture.md

@@ -367,6 +367,8 @@ merged with GUID-based baselines to avoid double counting after resets.
 `job/xray_traffic_job.go`, `job/node_traffic_sync_job.go`, `service/inbound_node.go`
 (`SetRemoteTraffic` / `upsertNodeBaseline`), models `xray.ClientTraffic`,
 `model.NodeClientTraffic`, `model.ClientGlobalTraffic` (cross-master totals).
+A client reset is queued per hosting node in `model.NodePendingReset` (`service/node_reset_queue.go`)
+and replayed by the node sync until the node accepts it.
 Periodic resets: `job/periodic_traffic_reset_job.go` (keyed off `Inbound.TrafficReset`).
 
 ### 5.4 Background jobs (cron)
@@ -465,6 +467,7 @@ for AutoMigrate in `internal/database/db.go`.
 | `Host`                          | Subscription host overrides (per inbound) | `Address`, `Port`, `Sni`, `Path`, `Security`, `Fingerprint`, `SortOrder`, visibility/exclusion flags                                                               |
 | `Node`                          | A managed child panel                     | `Guid`, `Address`, `Status`, `TlsVerifyMode`, `PinnedCertSha256`, `ConfigDirty`, version/heartbeat/metric fields                                                   |
 | `NodeClientTraffic`             | Per-node client traffic baseline          | cross-node merge (anti-double-count)                                                                                                                               |
+| `NodePendingReset`              | Client resets a node has not confirmed    | `NodeId`, `Email`, `QueuedAt`; replayed by the node sync, freezes that client's node verdict until delivered                                                       |
 | `NodeClientIp`                  | Per-node client IP attribution            | `NodeGuid`, `Email`, `Ips`                                                                                                                                         |
 | `ClientGlobalTraffic`           | Cross-master usage totals                 | `MasterGuid`, `Email`, `Up`, `Down`                                                                                                                                |
 | `xray.ClientTraffic`            | Per-client counters (`client_traffics`)   | `Email`, `Up`, `Down`, `Total`, `ExpiryTime`, `LastOnline`                                                                                                         |

+ 10 - 8
docs/content/docs/en/reference/api/inbounds.mdx

@@ -60,10 +60,11 @@ _openapi:
         at most once.
       url: '#delete-many-inbounds-in-one-call-processes-the-list-sequentially-failures-are-reported-per-id-and-the-rest-still-proceed-restarts-xray-at-most-once'
     - depth: 2
-      title: Replace an inbound’s configuration. Body shape mirrors /add. Heavy on
-        inbounds with thousands of clients — prefer /setEnable for enable-only
-        flips.
-      url: '#replace-an-inbounds-configuration-body-shape-mirrors-add-heavy-on-inbounds-with-thousands-of-clients--prefer-setenable-for-enable-only-flips'
+      title: 'Replace an inbound’s configuration. Body shape mirrors /add, but the
+        inbound keeps its stored client list and enable flag: settings.clients
+        and enable in the body are ignored. Manage clients through the
+        /panel/api/clients endpoints and toggle the inbound with /setEnable.'
+      url: '#replace-an-inbounds-configuration-body-shape-mirrors-add-but-the-inbound-keeps-its-stored-client-list-and-enable-flag-settingsclients-and-enable-in-the-body-are-ignored-manage-clients-through-the-panelapiclients-endpoints-and-toggle-the-inbound-with-setenable'
     - depth: 2
       title: Toggle only the enable flag without serialising the whole settings JSON.
         Recommended for UI switches on large inbounds.
@@ -151,10 +152,11 @@ _openapi:
           failures are reported per id and the rest still proceed. Restarts xray
           at most once.
         id: delete-many-inbounds-in-one-call-processes-the-list-sequentially-failures-are-reported-per-id-and-the-rest-still-proceed-restarts-xray-at-most-once
-      - content: Replace an inbound’s configuration. Body shape mirrors /add. Heavy on
-          inbounds with thousands of clients — prefer /setEnable for enable-only
-          flips.
-        id: replace-an-inbounds-configuration-body-shape-mirrors-add-heavy-on-inbounds-with-thousands-of-clients--prefer-setenable-for-enable-only-flips
+      - content: 'Replace an inbound’s configuration. Body shape mirrors /add, but the
+          inbound keeps its stored client list and enable flag: settings.clients
+          and enable in the body are ignored. Manage clients through the
+          /panel/api/clients endpoints and toggle the inbound with /setEnable.'
+        id: replace-an-inbounds-configuration-body-shape-mirrors-add-but-the-inbound-keeps-its-stored-client-list-and-enable-flag-settingsclients-and-enable-in-the-body-are-ignored-manage-clients-through-the-panelapiclients-endpoints-and-toggle-the-inbound-with-setenable
       - content: Toggle only the enable flag without serialising the whole settings
           JSON. Recommended for UI switches on large inbounds.
         id: toggle-only-the-enable-flag-without-serialising-the-whole-settings-json-recommended-for-ui-switches-on-large-inbounds

+ 1 - 1
docs/public/openapi.json

@@ -5603,7 +5603,7 @@
         "tags": [
           "Inbounds"
         ],
-        "summary": "Replace an inbound’s configuration. Body shape mirrors /add. Heavy on inbounds with thousands of clients — prefer /setEnable for enable-only flips.",
+        "summary": "Replace an inbound’s configuration. Body shape mirrors /add, but the inbound keeps its stored client list and enable flag: settings.clients and enable in the body are ignored. Manage clients through the /panel/api/clients endpoints and toggle the inbound with /setEnable.",
         "operationId": "post_panel_api_inbounds_update_id",
         "parameters": [
           {

+ 1 - 1
frontend/public/openapi.json

@@ -5603,7 +5603,7 @@
         "tags": [
           "Inbounds"
         ],
-        "summary": "Replace an inbound’s configuration. Body shape mirrors /add. Heavy on inbounds with thousands of clients — prefer /setEnable for enable-only flips.",
+        "summary": "Replace an inbound’s configuration. Body shape mirrors /add, but the inbound keeps its stored client list and enable flag: settings.clients and enable in the body are ignored. Manage clients through the /panel/api/clients endpoints and toggle the inbound with /setEnable.",
         "operationId": "post_panel_api_inbounds_update_id",
         "parameters": [
           {

+ 15 - 2
frontend/src/lib/xray/inbound-form-adapter.ts

@@ -353,9 +353,22 @@ export function dropLegacyOptionalEmpties(
   }
 }
 
-export function formValuesToWirePayload(values: InboundFormValues): WireInboundPayload {
+// An existing inbound's clients change only through the client endpoints, so
+// the edit form neither loads them nor sends them back.
+export function withoutClients(values: InboundFormValues): InboundFormValues {
+  const settings = { ...(values.settings as Record<string, unknown> | undefined) };
+  delete settings.clients;
+  return { ...values, settings } as InboundFormValues;
+}
+
+export function formValuesToWirePayload(
+  values: InboundFormValues,
+  options: { omitClients?: boolean } = {},
+): WireInboundPayload {
   const settingsPruned = (pruneEmpty(values.settings ?? {}) ?? {}) as Record<string, unknown>;
-  if (Array.isArray(settingsPruned.clients)) {
+  if (options.omitClients) {
+    delete settingsPruned.clients;
+  } else if (Array.isArray(settingsPruned.clients)) {
     settingsPruned.clients = normalizeClients(values.protocol, settingsPruned.clients);
   }
   let streamPruned = values.streamSettings

+ 1 - 1
frontend/src/pages/api-docs/endpoints.ts

@@ -325,7 +325,7 @@ export const sections: readonly Section[] = [
         method: 'POST',
         path: '/panel/api/inbounds/update/:id',
         summary:
-          'Replace an inbound’s configuration. Body shape mirrors /add. Heavy on inbounds with thousands of clients — prefer /setEnable for enable-only flips.',
+          'Replace an inbound’s configuration. Body shape mirrors /add, but the inbound keeps its stored client list and enable flag: settings.clients and enable in the body are ignored. Manage clients through the /panel/api/clients endpoints and toggle the inbound with /setEnable.',
         params: [{ name: 'id', in: 'path', type: 'number', desc: 'Inbound ID.' }],
         body: inboundBody,
       },

+ 16 - 13
frontend/src/pages/inbounds/form/InboundFormModal.tsx

@@ -19,7 +19,11 @@ import { Controller, FormProvider, useForm, useWatch } from 'react-hook-form';
 
 import { HttpUtil, NumberFormatter, RandomUtil, SizeFormatter, Wireguard } from '@/utils';
 import type { RealityScanResult } from '@/generated/types';
-import { rawInboundToFormValues, formValuesToWirePayload } from '@/lib/xray/inbound-form-adapter';
+import {
+  rawInboundToFormValues,
+  formValuesToWirePayload,
+  withoutClients,
+} from '@/lib/xray/inbound-form-adapter';
 import { createDefaultInboundSettings } from '@/lib/xray/inbound-defaults';
 import { generateAwgObfuscation } from '@/lib/xray/amneziawg-obfuscation';
 import { composeInboundTag, isAutoInboundTag, type InboundTagInput } from '@/lib/xray/inbound-tag';
@@ -439,7 +443,9 @@ export default function InboundFormModal({
   useEffect(() => {
     if (!open) return;
     const initial =
-      mode === 'edit' && dbInbound ? rawInboundToFormValues(dbInbound) : buildAddModeValues();
+      mode === 'edit' && dbInbound
+        ? withoutClients(rawInboundToFormValues(dbInbound))
+        : buildAddModeValues();
     methods.reset(initial);
     setScanResult(null);
     setActiveTab('basic');
@@ -556,13 +562,8 @@ export default function InboundFormModal({
   }, [mode, methods]);
 
   const saveValues = async () => {
-    /*
-     * getValues() returns the entire form store, including settings.clients and
-     * settings.fallbacks which have no bound field (clients are managed via the
-     * standalone Client modal, not this inbound modal). With shouldUnregister
-     * false those pass-through sub-trees survive from the reset object, so the
-     * update wire payload never silently drops every client on save.
-     */
+    // settings.fallbacks has no bound field; shouldUnregister=false keeps it from
+    // the reset object. An edit sends no clients: the server keeps the stored ones.
     const values = methods.getValues() as InboundFormValues;
     const parsed = InboundFormSchema.safeParse(values);
     if (!parsed.success) {
@@ -577,7 +578,7 @@ export default function InboundFormModal({
     }
     setSaving(true);
     try {
-      const payload = formValuesToWirePayload(parsed.data);
+      const payload = formValuesToWirePayload(parsed.data, { omitClients: mode === 'edit' });
       const url =
         mode === 'edit' && dbInbound
           ? `/panel/api/inbounds/update/${dbInbound.id}`
@@ -615,9 +616,11 @@ export default function InboundFormModal({
 
   const basicTab = (
     <>
-      <FormField name="enable" label={t('enable')} valueProp="checked">
-        <Switch />
-      </FormField>
+      {mode === 'add' && (
+        <FormField name="enable" label={t('enable')} valueProp="checked">
+          <Switch id="inbound-enable" />
+        </FormField>
+      )}
 
       <FormField name="remark" label={t('pages.inbounds.remark')}>
         <Input />

+ 30 - 0
frontend/src/test/inbound-form-modal.test.tsx

@@ -317,4 +317,34 @@ describe('InboundFormModal', () => {
       );
     });
   });
+
+  // Clients and enable change through their own endpoints; the server keeps the
+  // stored ones, so the edit form must neither send nor validate its stale copy.
+  it('edit save neither sends nor validates the clients it loaded', async () => {
+    const post = vi.mocked(HttpUtil.post);
+    post.mockClear();
+    const dbInbound = cloneLikeVlessInbound('example.com:443');
+    const legacy = new DBInbound({
+      ...dbInbound,
+      settings: {
+        ...(dbInbound.settings as Record<string, unknown>),
+        clients: [{ email: 'legacy', id: '' }],
+      },
+    });
+    renderCloneLikeEdit(legacy);
+
+    fireEvent.click(primaryButton());
+
+    await waitFor(() => expect(post).toHaveBeenCalled());
+    const payload = post.mock.calls[0][1] as { settings: string };
+    expect(JSON.parse(payload.settings)).not.toHaveProperty('clients');
+  });
+
+  it('offers the enable switch when adding an inbound but not when editing one', () => {
+    renderModal();
+    expect(document.getElementById('inbound-enable')).not.toBeNull();
+    cleanup();
+    renderCloneLikeEdit(cloneLikeVlessInbound('example.com:443'));
+    expect(document.getElementById('inbound-enable')).toBeNull();
+  });
 });

+ 1 - 0
internal/database/db.go

@@ -83,6 +83,7 @@ func allModels() []any {
 		&model.NodeClientTraffic{},
 		&model.NodeClientIp{},
 		&model.ClientGlobalTraffic{},
+		&model.NodePendingReset{},
 		&model.OutboundSubscription{},
 		&model.SubBalancer{},
 	}

+ 1 - 0
internal/database/migrate_data.go

@@ -56,6 +56,7 @@ func migrationModels() []any {
 		&model.NodeClientTraffic{},
 		&model.NodeClientIp{},
 		&model.ClientGlobalTraffic{},
+		&model.NodePendingReset{},
 		&model.OutboundSubscription{},
 		&model.SubBalancer{},
 	}

+ 11 - 0
internal/database/model/node_pending_reset.go

@@ -0,0 +1,11 @@
+package model
+
+// NodePendingReset is a client traffic reset a hosting node has not confirmed;
+// until it lands the node still counts pre-reset usage, so every sync replays it.
+type NodePendingReset struct {
+	Id     int    `json:"id" gorm:"primaryKey;autoIncrement"`
+	NodeId int    `json:"nodeId" gorm:"uniqueIndex:idx_node_pending_reset,priority:1;not null"`
+	Email  string `json:"email" gorm:"uniqueIndex:idx_node_pending_reset,priority:2;not null"`
+	// QueuedAt (ns) tells a delivery apart from a reset re-queued while it ran.
+	QueuedAt int64 `json:"queuedAt"`
+}

+ 80 - 0
internal/web/job/node_reset_replay_test.go

@@ -0,0 +1,80 @@
+package job
+
+import (
+	"net/http"
+	"net/http/httptest"
+	"path/filepath"
+	"slices"
+	"strconv"
+	"strings"
+	"sync"
+	"testing"
+
+	"github.com/op/go-logging"
+
+	"github.com/mhsanaei/3x-ui/v3/internal/database"
+	"github.com/mhsanaei/3x-ui/v3/internal/database/dbtest"
+	"github.com/mhsanaei/3x-ui/v3/internal/database/model"
+	xuilogger "github.com/mhsanaei/3x-ui/v3/internal/logger"
+	"github.com/mhsanaei/3x-ui/v3/internal/web/runtime"
+	"github.com/mhsanaei/3x-ui/v3/internal/web/service"
+)
+
+// A reset the node missed is replayed by the next sync, ahead of the snapshot
+// fetch so the merge already sees the zeroed counters.
+func TestNodeTrafficSyncReplaysOwedResetBeforeSnapshot(t *testing.T) {
+	xuilogger.InitLogger(logging.ERROR)
+	dbtest.InitDB(t, filepath.Join(t.TempDir(), "x-ui.db"))
+	service.StartTrafficWriter()
+	t.Cleanup(service.StopTrafficWriter)
+	runtime.SetManager(runtime.NewManager(runtime.LocalDeps{APIPort: func() int { return 0 }, SetNeedRestart: func() {}}))
+	t.Cleanup(func() { runtime.SetManager(nil) })
+
+	var mu sync.Mutex
+	var calls []string
+	srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
+		mu.Lock()
+		switch {
+		case strings.Contains(r.URL.Path, "clients/resetTraffic/"):
+			calls = append(calls, "reset:"+r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:])
+		case strings.HasSuffix(r.URL.Path, "inbounds/list"):
+			calls = append(calls, "snapshot")
+		}
+		mu.Unlock()
+		w.Header().Set("Content-Type", "application/json")
+		if strings.HasSuffix(r.URL.Path, "inbounds/list") {
+			_, _ = w.Write([]byte(`{"success":true,"obj":[]}`))
+			return
+		}
+		_, _ = w.Write([]byte(`{"success":true}`))
+	}))
+	t.Cleanup(srv.Close)
+	host, port, _ := strings.Cut(strings.TrimPrefix(srv.URL, "http://"), ":")
+	portNum, _ := strconv.Atoi(port)
+	node := &model.Node{
+		Name: "owes-reset", Scheme: "http", Address: host, Port: portNum, BasePath: "/", ApiToken: "tok",
+		Enable: true, Status: "online", AllowPrivateAddress: true, TlsVerifyMode: "verify",
+	}
+	if err := database.GetDB().Create(node).Error; err != nil {
+		t.Fatalf("create node: %v", err)
+	}
+	if err := database.GetDB().Create(&model.NodePendingReset{NodeId: node.Id, Email: "owed@node", QueuedAt: 1}).Error; err != nil {
+		t.Fatalf("seed pending reset: %v", err)
+	}
+
+	NewNodeTrafficSyncJob().Run()
+
+	mu.Lock()
+	got := slices.Clone(calls)
+	mu.Unlock()
+	if len(got) < 2 || got[0] != "reset:owed@node" || !slices.Contains(got, "snapshot") {
+		t.Fatalf("node calls %v, want the owed reset first, then the snapshot", got)
+	}
+	var left int64
+	if err := database.GetDB().Model(&model.NodePendingReset{}).Count(&left).Error; err != nil {
+		t.Fatalf("count pending: %v", err)
+	}
+	if left != 0 {
+		t.Fatalf("replayed reset still queued (%d rows)", left)
+	}
+}

+ 7 - 0
internal/web/job/node_traffic_sync_job.go

@@ -387,6 +387,13 @@ func (j *NodeTrafficSyncJob) syncOne(mgr *runtime.Manager, n *model.Node, doIpSy
 		}
 	}
 
+	// Before the snapshot, so counters a reset just zeroed are what gets merged.
+	resetCtx, resetCancel := context.WithTimeout(context.Background(), nodeTrafficSyncRequestTimeout)
+	if resetErr := j.inboundService.DeliverNodeResets(resetCtx, n.Id, rt); resetErr != nil {
+		logger.Warningf("node traffic sync: reset delivery to %s failed, retrying next tick: %v", n.Name, resetErr)
+	}
+	resetCancel()
+
 	ctx, cancel := context.WithTimeout(context.Background(), nodeTrafficSyncRequestTimeout)
 	defer cancel()
 

+ 6 - 0
internal/web/runtime/remote.go

@@ -701,6 +701,12 @@ func (r *Remote) ResetClientTraffic(ctx context.Context, _ *model.Inbound, email
 	return err
 }
 
+// ResetClientTraffics zeroes many clients on the node in one request.
+func (r *Remote) ResetClientTraffics(ctx context.Context, emails []string) error {
+	_, err := r.do(ctx, http.MethodPost, "panel/api/clients/bulkResetTraffic", map[string]any{"emails": emails})
+	return err
+}
+
 func (r *Remote) ResetAllTraffics(ctx context.Context) error {
 	_, err := r.do(ctx, http.MethodPost, "panel/api/inbounds/resetAllTraffics", nil)
 	return err

+ 33 - 0
internal/web/runtime/remote_reset_test.go

@@ -0,0 +1,33 @@
+package runtime
+
+import (
+	"context"
+	"encoding/json"
+	"net/http"
+	"net/http/httptest"
+	"slices"
+	"testing"
+)
+
+// The master replays a node's reset backlog through the node's bulk endpoint.
+func TestRemoteResetClientTrafficsPostsEmailsToBulkEndpoint(t *testing.T) {
+	var path string
+	var body struct {
+		Emails []string `json:"emails"`
+	}
+	srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
+		path = r.URL.Path
+		_ = json.NewDecoder(r.Body).Decode(&body)
+		w.Header().Set("Content-Type", "application/json")
+		_, _ = w.Write([]byte(`{"success":true}`))
+	}))
+	t.Cleanup(srv.Close)
+
+	r := NewRemote(nodeForPlainServer(t, srv, "verify", "tok"), nil)
+	if err := r.ResetClientTraffics(context.Background(), []string{"a@x", "b@x"}); err != nil {
+		t.Fatalf("ResetClientTraffics: %v", err)
+	}
+	if path != "/panel/api/clients/bulkResetTraffic" || !slices.Equal(body.Emails, []string{"a@x", "b@x"}) {
+		t.Fatalf("node got %s %v, want /panel/api/clients/bulkResetTraffic [a@x b@x]", path, body.Emails)
+	}
+}

+ 2 - 1
internal/web/service/client_renewal_write_validation_test.go

@@ -48,7 +48,8 @@ func TestClientRenewalWriteValidation(t *testing.T) {
 					}
 					_, _, err = inboundSvc.AddInbound(&update)
 				case "update inbound":
-					_, _, err = inboundSvc.UpdateInbound(&update)
+					// A panel save keeps the stored clients; only a master's push writes them.
+					_, _, err = (&InboundService{FromNodeSync: true}).UpdateInbound(&update)
 				case "add inbound client":
 					client.Email = "invalid-new-client"
 					update.Settings = clientsSettings(t, []model.Client{client})

+ 33 - 7
internal/web/service/client_traffic.go

@@ -74,6 +74,7 @@ func (s *ClientService) BulkResetTraffic(inboundSvc *InboundService, emails []st
 		return 0, err
 	}
 	affected := 0
+	var resetNodes []int
 	err = submitTrafficWrite(func() error {
 		db := database.GetDB()
 		return db.Transaction(func(tx *gorm.DB) error {
@@ -97,12 +98,16 @@ func (s *ClientService) BulkResetTraffic(inboundSvc *InboundService, emails []st
 					return err
 				}
 			}
-			return nil
+			var qErr error
+			resetNodes, qErr = queueNodeResets(tx, cleanEmails)
+			return qErr
 		})
 	})
 	if err != nil {
 		return 0, err
 	}
+	inboundSvc.resetMtprotoClientQuotas(cleanEmails)
+	inboundSvc.deliverNodeResetsNow(resetNodes)
 	// After the zeroing, as in ResetTrafficByEmail: enabling a still-depleted
 	// client first lets the next traffic tick switch it off again.
 	for _, e := range cleanEmails {
@@ -120,18 +125,25 @@ func (s *ClientService) BulkResetTraffic(inboundSvc *InboundService, emails []st
 }
 
 func (s *ClientService) ResetAllClientTraffics(inboundSvc *InboundService, id int) error {
+	var resetNodes []int
+	var resetEmails []string
 	err := submitTrafficWrite(func() error {
-		return s.resetAllClientTrafficsLocked(id)
+		var inner error
+		resetEmails, resetNodes, inner = s.resetAllClientTrafficsLocked(id)
+		return inner
 	})
 	if err == nil {
-		inboundSvc.resetAllMtprotoQuotas()
+		inboundSvc.resetMtprotoClientQuotas(resetEmails)
+		inboundSvc.deliverNodeResetsNow(resetNodes)
 	}
 	return err
 }
 
-func (s *ClientService) resetAllClientTrafficsLocked(id int) error {
+func (s *ClientService) resetAllClientTrafficsLocked(id int) ([]string, []int, error) {
 	db := database.GetDB()
 	now := time.Now().Unix() * 1000
+	var resetNodes []int
+	var reset []string
 
 	if err := db.Transaction(func(tx *gorm.DB) error {
 		// client_traffics.inbound_id is stale: it reflects the inbound the row was
@@ -154,6 +166,7 @@ func (s *ClientService) resetAllClientTrafficsLocked(id int) error {
 		if len(resetEmails) == 0 {
 			return nil
 		}
+		reset = resetEmails
 
 		if err := adjustGroupBaselinesForRemovedTraffic(tx, resetEmails); err != nil {
 			return err
@@ -176,6 +189,10 @@ func (s *ClientService) resetAllClientTrafficsLocked(id int) error {
 				return err
 			}
 		}
+		var qErr error
+		if resetNodes, qErr = queueNodeResets(tx, resetEmails); qErr != nil {
+			return qErr
+		}
 
 		inboundWhereText := "id "
 		if id == -1 {
@@ -190,13 +207,14 @@ func (s *ClientService) resetAllClientTrafficsLocked(id int) error {
 
 		return result.Error
 	}); err != nil {
-		return err
+		return nil, nil, err
 	}
-	return nil
+	return reset, resetNodes, nil
 }
 
 func (s *ClientService) ResetAllTraffics() (bool, error) {
 	var affected int64
+	var resetNodes []int
 	err := submitTrafficWrite(func() error {
 		return database.GetDB().Transaction(func(tx *gorm.DB) error {
 			res := tx.Model(&xray.ClientTraffic{}).
@@ -209,11 +227,19 @@ func (s *ClientService) ResetAllTraffics() (bool, error) {
 			if err := tx.Where("1 = 1").Delete(&model.ClientGlobalTraffic{}).Error; err != nil {
 				return err
 			}
-			return tx.Where("1 = 1").Delete(&model.NodeClientTraffic{}).Error
+			if err := tx.Where("1 = 1").Delete(&model.NodeClientTraffic{}).Error; err != nil {
+				return err
+			}
+			var qErr error
+			resetNodes, qErr = queueNodeResets(tx, nil)
+			return qErr
 		})
 	})
 	if err != nil {
 		return false, err
 	}
+	inbounds := &InboundService{}
+	inbounds.resetAllMtprotoQuotas()
+	inbounds.deliverNodeResetsNow(resetNodes)
 	return affected > 0, nil
 }

+ 56 - 39
internal/web/service/inbound.go

@@ -25,7 +25,6 @@ import (
 	"github.com/mhsanaei/3x-ui/v3/internal/util/common"
 	"github.com/mhsanaei/3x-ui/v3/internal/util/netsafe"
 	wgutil "github.com/mhsanaei/3x-ui/v3/internal/util/wireguard"
-	"github.com/mhsanaei/3x-ui/v3/internal/xray"
 
 	"gorm.io/gorm"
 	"gorm.io/gorm/clause"
@@ -1685,6 +1684,35 @@ func (s *InboundService) SetInboundEnable(id int, enable bool) (bool, error) {
 	return needRestart, nil
 }
 
+func (s *InboundService) validateUpdatedInboundClients(inbound *model.Inbound) error {
+	clients, err := s.GetClients(inbound)
+	if err != nil {
+		return err
+	}
+	if err := validateClientsRenewal(clients); err != nil {
+		return err
+	}
+	for _, client := range clients {
+		switch inbound.Protocol {
+		case model.Hysteria:
+			if client.Auth == "" {
+				return common.NewError("empty client ID")
+			}
+		case model.TUIC:
+			if client.ID == "" {
+				return common.NewError("empty client ID")
+			}
+			if client.Password == "" {
+				return common.NewError("tuic client requires a password")
+			}
+			if client.Email == "" {
+				return common.NewError("empty client email")
+			}
+		}
+	}
+	return nil
+}
+
 func (s *InboundService) UpdateInbound(inbound *model.Inbound) (*model.Inbound, bool, error) {
 	legacyShareAddr := legacyMtprotoShareAddr(inbound)
 	inbound.TrafficResetDay = normalizeTrafficResetDay(inbound.TrafficResetDay)
@@ -1707,34 +1735,6 @@ func (s *InboundService) UpdateInbound(inbound *model.Inbound) (*model.Inbound,
 	}
 	inbound.SubSortIndex = normalizeSubSortIndex(inbound.SubSortIndex)
 
-	clients, err := s.GetClients(inbound)
-	if err != nil {
-		return inbound, false, err
-	}
-	if err := validateClientsRenewal(clients); err != nil {
-		return inbound, false, err
-	}
-	if inbound.Protocol == model.Hysteria {
-		for _, client := range clients {
-			if client.Auth == "" {
-				return inbound, false, common.NewError("empty client ID")
-			}
-		}
-	}
-	if inbound.Protocol == model.TUIC {
-		for _, client := range clients {
-			if client.ID == "" {
-				return inbound, false, common.NewError("empty client ID")
-			}
-			if client.Password == "" {
-				return inbound, false, common.NewError("tuic client requires a password")
-			}
-			if client.Email == "" {
-				return inbound, false, common.NewError("empty client email")
-			}
-		}
-	}
-
 	// Grandfather a row that was already stored incomplete so it stays editable;
 	// only a save that breaks a previously valid TLS block is refused.
 	if !s.FromNodeSync {
@@ -1770,6 +1770,23 @@ func (s *InboundService) UpdateInbound(inbound *model.Inbound) (*model.Inbound,
 	var postCommitApply func()
 
 	txErr := runSerializedTx(func(tx *gorm.DB) error {
+		// Re-read inside the writer: a traffic tick since the read above moved the
+		// counters and may have renewed or disabled clients.
+		stored := &model.Inbound{}
+		if err := tx.First(stored, inbound.Id).Error; err != nil {
+			return err
+		}
+		oldInbound = stored
+		// The form posts back the clients and enable it loaded; both have their
+		// own endpoints, so only a master's push may change them here.
+		if !s.FromNodeSync {
+			inbound.Settings = keepStoredClients(inbound.Settings, stored.Settings)
+			inbound.Enable = stored.Enable
+		}
+		// On the clients actually saved: a protocol switch keeps the stored ones.
+		if err := s.validateUpdatedInboundClients(inbound); err != nil {
+			return err
+		}
 		conflict, cErr := checkPortConflictTx(tx, inbound, inbound.Id)
 		if cErr != nil {
 			return cErr
@@ -2066,16 +2083,16 @@ func (s *InboundService) buildInboundForLocalRuntime(tx *gorm.DB, inbound *model
 		return built, nil
 	}
 
-	var clientStats []xray.ClientTraffic
-	if err := tx.Model(xray.ClientTraffic{}).
-		Where("inbound_id = ?", built.Id).
-		Select("email", "enable").
-		Find(&clientStats).Error; err != nil {
-		return nil, err
+	emails := make([]string, 0, len(clients))
+	for _, client := range clients {
+		if c, ok := client.(map[string]any); ok {
+			email, _ := c["email"].(string)
+			emails = append(emails, email)
+		}
 	}
-	enableMap := make(map[string]bool, len(clientStats))
-	for _, clientTraffic := range clientStats {
-		enableMap[clientTraffic.Email] = clientTraffic.Enable
+	disabled, err := trafficDisabledEmails(tx, emails)
+	if err != nil {
+		return nil, err
 	}
 
 	finalClients := make([]any, 0, len(clients))
@@ -2085,7 +2102,7 @@ func (s *InboundService) buildInboundForLocalRuntime(tx *gorm.DB, inbound *model
 			continue
 		}
 		email, _ := c["email"].(string)
-		if enable, exists := enableMap[email]; exists && !enable {
+		if _, off := disabled[email]; off {
 			continue
 		}
 		if manualEnable, ok := c["enable"].(bool); ok && !manualEnable {

+ 22 - 32
internal/web/service/inbound_amneziawg.go

@@ -14,7 +14,6 @@ import (
 	"github.com/mhsanaei/3x-ui/v3/internal/database/model"
 	"github.com/mhsanaei/3x-ui/v3/internal/logger"
 	wgutil "github.com/mhsanaei/3x-ui/v3/internal/util/wireguard"
-	"github.com/mhsanaei/3x-ui/v3/internal/xray"
 )
 
 // DesiredAmneziaWGInstances derives the AmneziaWG interfaces this panel
@@ -37,47 +36,38 @@ func (s *InboundService) DesiredAmneziaWGInstances() ([]amneziawg.Instance, erro
 		return nil, nil
 	}
 
-	ids := make([]int, 0, len(inbounds))
-	for _, ib := range inbounds {
-		ids = append(ids, ib.Id)
-	}
-	var disabledRows []xray.ClientTraffic
-	err = db.Model(xray.ClientTraffic{}).
-		Where("inbound_id IN ? AND enable = ?", ids, false).
-		Select("inbound_id", "email").
-		Find(&disabledRows).Error
-	if err != nil {
-		return nil, err
-	}
-	disabled := make(map[int]map[string]struct{}, len(disabledRows))
-	for _, row := range disabledRows {
-		if disabled[row.InboundId] == nil {
-			disabled[row.InboundId] = map[string]struct{}{}
-		}
-		disabled[row.InboundId][row.Email] = struct{}{}
-	}
-
 	instances := make([]amneziawg.Instance, 0, len(inbounds))
 	for _, ib := range inbounds {
 		inst, ok := amneziawg.InstanceFromInbound(ib)
 		if !ok {
 			continue
 		}
-		if off := disabled[ib.Id]; len(off) > 0 {
-			kept := make([]amneziawg.Peer, 0, len(inst.Peers))
-			for _, p := range inst.Peers {
-				if _, skip := off[p.Email]; !skip {
-					kept = append(kept, p)
-				}
+		instances = append(instances, inst)
+	}
+	emails := make([]string, 0)
+	for _, inst := range instances {
+		for _, e := range inst.Peers {
+			emails = append(emails, e.Email)
+		}
+	}
+	disabled, err := trafficDisabledEmails(db, emails)
+	if err != nil {
+		return nil, err
+	}
+	served := instances[:0]
+	for _, inst := range instances {
+		kept := make([]amneziawg.Peer, 0, len(inst.Peers))
+		for _, e := range inst.Peers {
+			if _, off := disabled[e.Email]; !off {
+				kept = append(kept, e)
 			}
-			inst.Peers = kept
 		}
-		if len(inst.Peers) == 0 {
-			continue
+		inst.Peers = kept
+		if len(kept) > 0 {
+			served = append(served, inst)
 		}
-		instances = append(instances, inst)
 	}
-	return instances, nil
+	return served, nil
 }
 
 // applyLocalAmneziaWG pushes a single local AmneziaWG inbound's current peer

+ 18 - 0
internal/web/service/inbound_disable.go

@@ -234,3 +234,21 @@ func (s *InboundService) markClientsDisabledInSettings(tx *gorm.DB, inboundID in
 	}
 	return &snapshot, &ib, nil
 }
+
+// trafficDisabledEmails reports which emails have a switched-off stats row. The
+// table is email-keyed and its inbound_id goes stale, so never filter on it.
+func trafficDisabledEmails(db *gorm.DB, emails []string) (map[string]struct{}, error) {
+	disabled := make(map[string]struct{})
+	for _, batch := range chunkStrings(uniqueNonEmptyStrings(emails), sqlInChunk) {
+		var page []string
+		if err := db.Model(xray.ClientTraffic{}).
+			Where("email IN ? AND enable = ?", batch, false).
+			Pluck("email", &page).Error; err != nil {
+			return nil, err
+		}
+		for _, e := range page {
+			disabled[e] = struct{}{}
+		}
+	}
+	return disabled, nil
+}

+ 6 - 3
internal/web/service/inbound_hysteria_auth_test.go

@@ -8,10 +8,12 @@ import (
 	"github.com/mhsanaei/3x-ui/v3/internal/database/model"
 )
 
+// An inbound save keeps the stored clients, so switching to Hysteria is judged
+// on them: ones with no auth would leave an inbound nobody can connect to.
 func TestUpdateInbound_RejectsHysteriaClientWithoutAuth(t *testing.T) {
 	setupConflictDB(t)
 	seedInboundConflict(t, "in-45001-tcp", "0.0.0.0", 45001, model.VLESS,
-		`{"network":"tcp"}`, `{"clients":[]}`)
+		`{"network":"tcp"}`, `{"clients":[{"email":"hysteria@x","enable":true,"password":"not-hysteria-auth"}]}`)
 
 	var existing model.Inbound
 	if err := database.GetDB().Where("tag = ?", "in-45001-tcp").First(&existing).Error; err != nil {
@@ -20,7 +22,7 @@ func TestUpdateInbound_RejectsHysteriaClientWithoutAuth(t *testing.T) {
 
 	update := existing
 	update.Protocol = model.Hysteria
-	update.Settings = `{"clients":[{"email":"hysteria@x","enable":true,"password":"not-hysteria-auth"}]}`
+	update.Settings = `{"clients":[]}`
 
 	svc := &InboundService{}
 	if _, _, err := svc.UpdateInbound(&update); err == nil || !strings.Contains(err.Error(), "empty client ID") {
@@ -54,7 +56,8 @@ func TestUpdateInbound_PreservesHysteriaClientAuth(t *testing.T) {
 	update := existing
 	update.Settings = `{"clients":[{"email":"hysteria@x","enable":true,"password":"` + password + `","auth":"` + wantAuth + `"}]}`
 
-	svc := &InboundService{}
+	// Only a master's push still carries clients through an inbound save.
+	svc := &InboundService{FromNodeSync: true}
 	if _, _, err := svc.UpdateInbound(&update); err != nil {
 		t.Fatalf("UpdateInbound: %v", err)
 	}

+ 58 - 59
internal/web/service/inbound_mtproto.go

@@ -7,7 +7,6 @@ import (
 	"github.com/mhsanaei/3x-ui/v3/internal/database/model"
 	"github.com/mhsanaei/3x-ui/v3/internal/logger"
 	"github.com/mhsanaei/3x-ui/v3/internal/mtproto"
-	"github.com/mhsanaei/3x-ui/v3/internal/xray"
 )
 
 // DesiredMtprotoInstances derives the mtg sidecar configs this panel should be
@@ -32,47 +31,38 @@ func (s *InboundService) DesiredMtprotoInstances() ([]mtproto.Instance, error) {
 		return nil, nil
 	}
 
-	ids := make([]int, 0, len(inbounds))
-	for _, ib := range inbounds {
-		ids = append(ids, ib.Id)
-	}
-	var disabledRows []xray.ClientTraffic
-	err = db.Model(xray.ClientTraffic{}).
-		Where("inbound_id IN ? AND enable = ?", ids, false).
-		Select("inbound_id", "email").
-		Find(&disabledRows).Error
-	if err != nil {
-		return nil, err
-	}
-	disabled := make(map[int]map[string]struct{}, len(disabledRows))
-	for _, row := range disabledRows {
-		if disabled[row.InboundId] == nil {
-			disabled[row.InboundId] = map[string]struct{}{}
-		}
-		disabled[row.InboundId][row.Email] = struct{}{}
-	}
-
 	instances := make([]mtproto.Instance, 0, len(inbounds))
 	for _, ib := range inbounds {
 		inst, ok := mtproto.InstanceFromInbound(ib)
 		if !ok {
 			continue
 		}
-		if off := disabled[ib.Id]; len(off) > 0 {
-			kept := make([]mtproto.SecretEntry, 0, len(inst.Secrets))
-			for _, sec := range inst.Secrets {
-				if _, skip := off[sec.Name]; !skip {
-					kept = append(kept, sec)
-				}
+		instances = append(instances, inst)
+	}
+	emails := make([]string, 0)
+	for _, inst := range instances {
+		for _, e := range inst.Secrets {
+			emails = append(emails, e.Name)
+		}
+	}
+	disabled, err := trafficDisabledEmails(db, emails)
+	if err != nil {
+		return nil, err
+	}
+	served := instances[:0]
+	for _, inst := range instances {
+		kept := make([]mtproto.SecretEntry, 0, len(inst.Secrets))
+		for _, e := range inst.Secrets {
+			if _, off := disabled[e.Name]; !off {
+				kept = append(kept, e)
 			}
-			inst.Secrets = kept
 		}
-		if len(inst.Secrets) == 0 {
-			continue
+		inst.Secrets = kept
+		if len(kept) > 0 {
+			served = append(served, inst)
 		}
-		instances = append(instances, inst)
 	}
-	return instances, nil
+	return served, nil
 }
 
 // applyLocalMtproto pushes a single local mtproto inbound's current client set
@@ -105,16 +95,47 @@ func (s *InboundService) applyLocalMtproto(inboundId int) {
 }
 
 func (s *InboundService) resetMtprotoClientQuota(email string) {
+	s.resetMtprotoClientQuotas([]string{email})
+}
+
+// resetMtprotoClientQuotas zeroes the sidecar's own quota counter for each local
+// MTProto client in emails, or it keeps blocking a client the panel just reset.
+func (s *InboundService) resetMtprotoClientQuotas(emails []string) {
 	mgr := mtproto.GetManager()
-	if !mgr.HasRunning() {
+	if !mgr.HasRunning() || len(emails) == 0 {
 		return
 	}
-	id, ok := s.localMtprotoInboundIdForEmail(email)
-	if !ok {
+	var inbounds []*model.Inbound
+	if err := database.GetDB().Model(model.Inbound{}).
+		Where("protocol = ? AND node_id IS NULL", model.MTProto).
+		Find(&inbounds).Error; err != nil {
 		return
 	}
-	s.applyLocalMtproto(id)
-	mgr.ResetQuota(email)
+	want := make(map[string]struct{}, len(emails))
+	for _, e := range emails {
+		want[e] = struct{}{}
+	}
+	var hit []string
+	for _, ib := range inbounds {
+		inst, ok := mtproto.InstanceFromInbound(ib)
+		if !ok {
+			continue
+		}
+		applied := false
+		for _, sec := range inst.Secrets {
+			if _, ok := want[sec.Name]; !ok {
+				continue
+			}
+			if !applied {
+				s.applyLocalMtproto(ib.Id)
+				applied = true
+			}
+			hit = append(hit, sec.Name)
+		}
+	}
+	for _, email := range hit {
+		mgr.ResetQuota(email)
+	}
 }
 
 func (s *InboundService) resetAllMtprotoQuotas() {
@@ -133,25 +154,3 @@ func (s *InboundService) resetAllMtprotoQuotas() {
 		}
 	}
 }
-
-func (s *InboundService) localMtprotoInboundIdForEmail(email string) (int, bool) {
-	db := database.GetDB()
-	var inbounds []*model.Inbound
-	if err := db.Model(model.Inbound{}).
-		Where("protocol = ? AND node_id IS NULL", model.MTProto).
-		Find(&inbounds).Error; err != nil {
-		return 0, false
-	}
-	for _, ib := range inbounds {
-		inst, ok := mtproto.InstanceFromInbound(ib)
-		if !ok {
-			continue
-		}
-		for _, sec := range inst.Secrets {
-			if sec.Name == email {
-				return ib.Id, true
-			}
-		}
-	}
-	return 0, false
-}

+ 2 - 1
internal/web/service/inbound_mtproto_apply_test.go

@@ -84,7 +84,8 @@ func TestUpdateInboundMtprotoUnchangedDoesNotRestart(t *testing.T) {
 		if !strings.Contains(update.Settings, mtprotoTestSecretD) {
 			t.Fatal("fixture must contain the re-keyed secret")
 		}
-		_, needRestart, err := svc.UpdateInbound(&update)
+		// Clients reach an inbound save only as a master's push to its node.
+		_, needRestart, err := (&InboundService{FromNodeSync: true}).UpdateInbound(&update)
 		if err != nil {
 			t.Fatalf("UpdateInbound: %v", err)
 		}

+ 15 - 6
internal/web/service/inbound_node.go

@@ -547,6 +547,11 @@ func (s *InboundService) setRemoteTrafficLocked(nodeID int, snap *runtime.Traffi
 		centralCSByEmail[centralClientStats[i].Email] = &centralClientStats[i]
 	}
 
+	owedResets, err := pendingNodeResetEmails(db, nodeID)
+	if err != nil {
+		return false, err
+	}
+
 	nodeBaselines := make(map[string]nodeTrafficCounter)
 	var baselineRows []model.NodeClientTraffic
 	if err := db.Model(&model.NodeClientTraffic{}).
@@ -925,6 +930,10 @@ func (s *InboundService) setRemoteTrafficLocked(nodeID int, snap *runtime.Traffi
 
 			// Node-wide total, not this inbound's possibly-stale copy (#5274).
 			canon := nodeEmailTotals[cs.Email]
+			// Until the node applies a reset it owes, its verdict rests on the
+			// pre-reset counters: only usage may move for this client.
+			_, owed := owedResets[cs.Email]
+			clientFrozen := lifecycleFrozen || owed
 
 			base, seen := nodeBaselines[cs.Email]
 			var deltaUp, deltaDown int64
@@ -986,18 +995,18 @@ func (s *InboundService) setRemoteTrafficLocked(nodeID int, snap *runtime.Traffi
 
 			existing := centralCSByEmail[cs.Email]
 			if existing != nil {
-				expiryChanged := !lifecycleFrozen && existing.ExpiryTime != mergeActivationExpiry(existing.ExpiryTime, cs.ExpiryTime)
+				expiryChanged := !clientFrozen && existing.ExpiryTime != mergeActivationExpiry(existing.ExpiryTime, cs.ExpiryTime)
 				// Only a real latch to disabled is structural; one-way merge never
 				// re-enables from the node.
-				enableChanged := !lifecycleFrozen && existing.Enable && !cs.Enable &&
+				enableChanged := !clientFrozen && existing.Enable && !cs.Enable &&
 					!nodeDisableIsStale(existing, cs, now, deltaUp, deltaDown)
-				metaChanged := !lifecycleFrozen && (existing.Total != cs.Total || existing.Reset != cs.Reset || existing.ResetWeekday != cs.ResetWeekday)
+				metaChanged := !clientFrozen && (existing.Total != cs.Total || existing.Reset != cs.Reset || existing.ResetWeekday != cs.ResetWeekday)
 				if enableChanged || metaChanged || expiryChanged {
 					structuralChange = true
 				}
 			}
 
-			renewed := !lifecycleFrozen && seen && existing != nil && nodeClientRenewed(existing, cs, canon, base)
+			renewed := !clientFrozen && seen && existing != nil && nodeClientRenewed(existing, cs, canon, base)
 			if renewed {
 				// Reject when the node's own settings still carry the old absolute:
 				// lagging ClientStats after a master shorten mimic a renew (#6228).
@@ -1037,7 +1046,7 @@ func (s *InboundService) setRemoteTrafficLocked(nodeID int, snap *runtime.Traffi
 				existing.ResetWeekday = cs.ResetWeekday
 				existing.ResetCount = cs.ResetCount
 				structuralChange = true
-			} else if lifecycleFrozen {
+			} else if clientFrozen {
 				// Push pending or just landed: only counters may move, the master
 				// keeps expiry/enable/total/reset.
 				if err := tx.Exec(
@@ -1096,7 +1105,7 @@ func (s *InboundService) setRemoteTrafficLocked(nodeID int, snap *runtime.Traffi
 			}
 			// A dip plus a lagging longer expiry mimics nodeClientRenewed and would
 			// undo a master shorten once the freeze lifts (#6228).
-			if lifecycleFrozen && seen && (canon.Up < base.Up || canon.Down < base.Down) {
+			if clientFrozen && seen && (canon.Up < base.Up || canon.Down < base.Down) {
 				continue
 			}
 			if err := s.upsertNodeBaseline(tx, nodeID, cs.Email, canon.Up, canon.Down); err != nil {

+ 145 - 0
internal/web/service/inbound_runtime_depleted_filter_test.go

@@ -0,0 +1,145 @@
+package service
+
+import (
+	"encoding/json"
+	"testing"
+
+	"github.com/mhsanaei/3x-ui/v3/internal/amneziawg"
+	"github.com/mhsanaei/3x-ui/v3/internal/database"
+	"github.com/mhsanaei/3x-ui/v3/internal/database/model"
+	"github.com/mhsanaei/3x-ui/v3/internal/xray"
+)
+
+// seedDepletedOnSibling attaches clients d and h to inbounds a and b with d's
+// depleted traffic row pointing at b, where AddClientStat's upsert leaves it.
+func seedDepletedOnSibling(t *testing.T, proto model.Protocol, port int, settings string) (a, b *model.Inbound) {
+	t.Helper()
+	setupSettingTestDB(t)
+	db := database.GetDB()
+	for i, dst := range []**model.Inbound{&a, &b} {
+		ib := &model.Inbound{Tag: string(proto) + "-sib-" + string(rune('a'+i)), Enable: true, Port: port + i, Protocol: proto, Settings: settings}
+		if err := db.Create(ib).Error; err != nil {
+			t.Fatalf("create inbound: %v", err)
+		}
+		clients, err := (&InboundService{}).GetClients(ib)
+		if err != nil {
+			t.Fatalf("GetClients: %v", err)
+		}
+		if err := (&ClientService{}).SyncInbound(nil, ib.Id, clients); err != nil {
+			t.Fatalf("SyncInbound: %v", err)
+		}
+		*dst = ib
+	}
+	rows := []xray.ClientTraffic{
+		{InboundId: b.Id, Email: "d", Enable: false, Up: 10, Total: 10},
+		{InboundId: b.Id, Email: "h", Enable: true},
+	}
+	if err := db.Create(&rows).Error; err != nil {
+		t.Fatalf("seed client_traffics: %v", err)
+	}
+	return a, b
+}
+
+func requireOnlyHealthy(t *testing.T, site string, emails []string) {
+	t.Helper()
+	if len(emails) != 1 || emails[0] != "h" {
+		t.Fatalf("%s serves %v on the sibling inbound, want only [h]: depleted d is still served", site, emails)
+	}
+}
+
+func TestRuntimeDropsDepletedClientWhoseStatsRowPointsAtSibling(t *testing.T) {
+	const vless = `{"clients":[{"email":"d","id":"11111111-1111-1111-1111-11111111111d","enable":true},` +
+		`{"email":"h","id":"11111111-1111-1111-1111-11111111111e","enable":true}],"decryption":"none"}`
+
+	t.Run("runtime push", func(t *testing.T) {
+		a, _ := seedDepletedOnSibling(t, model.VLESS, 23311, vless)
+		built, err := (&InboundService{}).buildInboundForLocalRuntime(database.GetDB(), a)
+		if err != nil {
+			t.Fatalf("buildInboundForLocalRuntime: %v", err)
+		}
+		clients, err := (&InboundService{}).GetClients(built)
+		if err != nil {
+			t.Fatalf("GetClients: %v", err)
+		}
+		var emails []string
+		for _, c := range clients {
+			emails = append(emails, c.Email)
+		}
+		requireOnlyHealthy(t, "buildInboundForLocalRuntime", emails)
+	})
+
+	t.Run("mtproto sidecar", func(t *testing.T) {
+		a, _ := seedDepletedOnSibling(t, model.MTProto, 23321,
+			`{"clients":[{"email":"d","secret":"`+mtprotoTestSecretA+`","enable":true},`+
+				`{"email":"h","secret":"`+mtprotoTestSecretB+`","enable":true}]}`)
+		instances, err := (&InboundService{}).DesiredMtprotoInstances()
+		if err != nil {
+			t.Fatalf("DesiredMtprotoInstances: %v", err)
+		}
+		for _, inst := range instances {
+			if inst.Id != a.Id {
+				continue
+			}
+			var emails []string
+			for _, sec := range inst.Secrets {
+				emails = append(emails, sec.Name)
+			}
+			requireOnlyHealthy(t, "DesiredMtprotoInstances", emails)
+			return
+		}
+		t.Fatal("sibling mtproto inbound missing from desired instances")
+	})
+
+	t.Run("tuic sidecar", func(t *testing.T) {
+		a, _ := seedDepletedOnSibling(t, model.TUIC, 23331,
+			`{"certificate":"/c.pem","private_key":"/k.pem","clients":[`+
+				`{"id":"11111111-1111-1111-1111-11111111111d","password":"pd","email":"d","enable":true},`+
+				`{"id":"11111111-1111-1111-1111-11111111111e","password":"ph","email":"h","enable":true}]}`)
+		instances, err := (&InboundService{}).DesiredTuicInstances()
+		if err != nil {
+			t.Fatalf("DesiredTuicInstances: %v", err)
+		}
+		for _, inst := range instances {
+			if inst.Id != a.Id {
+				continue
+			}
+			var emails []string
+			for _, c := range inst.Clients {
+				emails = append(emails, c.Email)
+			}
+			requireOnlyHealthy(t, "DesiredTuicInstances", emails)
+			return
+		}
+		t.Fatal("sibling tuic inbound missing from desired instances")
+	})
+
+	t.Run("amneziawg interface", func(t *testing.T) {
+		settings, err := json.Marshal(amneziawg.InboundSettings{
+			Server: &amneziawg.ServerSettings{SubnetIP: "10.8.1.0", SubnetCIDR: 24},
+			Clients: []model.Client{
+				{Email: "d", Enable: true, PublicKey: "pk-d", AllowedIPs: []string{"10.8.1.2/32"}},
+				{Email: "h", Enable: true, PublicKey: "pk-h", AllowedIPs: []string{"10.8.1.3/32"}},
+			},
+		})
+		if err != nil {
+			t.Fatalf("marshal awg settings: %v", err)
+		}
+		a, _ := seedDepletedOnSibling(t, model.AmneziaWG, 23341, string(settings))
+		instances, err := (&InboundService{}).DesiredAmneziaWGInstances()
+		if err != nil {
+			t.Fatalf("DesiredAmneziaWGInstances: %v", err)
+		}
+		for _, inst := range instances {
+			if inst.Id != a.Id {
+				continue
+			}
+			var emails []string
+			for _, p := range inst.Peers {
+				emails = append(emails, p.Email)
+			}
+			requireOnlyHealthy(t, "DesiredAmneziaWGInstances", emails)
+			return
+		}
+		t.Fatal("sibling amneziawg inbound missing from desired instances")
+	})
+}

+ 23 - 0
internal/web/service/inbound_settings_commit.go

@@ -139,3 +139,26 @@ func mergeClientLists(base, ours, current []any) []any {
 	}
 	return out
 }
+
+// keepStoredClients puts the stored client list back into an inbound save's
+// payload: clients change through the client endpoints, never this form.
+func keepStoredClients(payload, stored string) string {
+	var payloadM, storedM map[string]any
+	if json.Unmarshal([]byte(payload), &payloadM) != nil || json.Unmarshal([]byte(stored), &storedM) != nil {
+		return payload
+	}
+	storedClients, has := storedM["clients"]
+	if reflect.DeepEqual(payloadM["clients"], storedClients) {
+		return payload
+	}
+	if has {
+		payloadM["clients"] = storedClients
+	} else {
+		delete(payloadM, "clients")
+	}
+	b, err := json.MarshalIndent(payloadM, "", "  ")
+	if err != nil {
+		return payload
+	}
+	return string(b)
+}

+ 21 - 23
internal/web/service/inbound_traffic.go

@@ -28,14 +28,16 @@ const depletedClientsClause = "reset = 0 and reset_day = 0 and reset_weekday = 0
 func (s *InboundService) AddTraffic(inboundTraffics []*xray.Traffic, clientTraffics []*xray.ClientTraffic) (needRestart bool, clientsDisabled bool, err error) {
 	var disabledNodeIDs []int
 	var remotePlans []trafficInboundUpdatePlan
+	var renewed []string
 	err = submitTrafficWrite(func() error {
 		var inner error
-		needRestart, clientsDisabled, disabledNodeIDs, remotePlans, inner = s.addTrafficLocked(inboundTraffics, clientTraffics)
+		needRestart, clientsDisabled, disabledNodeIDs, remotePlans, renewed, inner = s.addTrafficLocked(inboundTraffics, clientTraffics)
 		return inner
 	})
 	if err != nil {
 		return
 	}
+	s.resetMtprotoClientQuotas(renewed)
 	// Off the serial writer: a hanging node must not stall traffic accounting.
 	needRestart = s.applyTrafficRemotePlans(remotePlans) || needRestart
 	if len(disabledNodeIDs) > 0 {
@@ -44,7 +46,7 @@ func (s *InboundService) AddTraffic(inboundTraffics []*xray.Traffic, clientTraff
 	return
 }
 
-func (s *InboundService) addTrafficLocked(inboundTraffics []*xray.Traffic, clientTraffics []*xray.ClientTraffic) (bool, bool, []int, []trafficInboundUpdatePlan, error) {
+func (s *InboundService) addTrafficLocked(inboundTraffics []*xray.Traffic, clientTraffics []*xray.ClientTraffic) (bool, bool, []int, []trafficInboundUpdatePlan, []string, error) {
 	db := database.GetDB()
 	// Commit durable traffic before best-effort lifecycle maintenance so helper
 	// failures cannot discard usage already reported by Xray.
@@ -54,7 +56,7 @@ func (s *InboundService) addTrafficLocked(inboundTraffics []*xray.Traffic, clien
 		}
 		return s.addClientTraffic(tx, clientTraffics)
 	}); err != nil {
-		return false, false, nil, nil, err
+		return false, false, nil, nil, nil, err
 	}
 
 	var (
@@ -99,10 +101,10 @@ func (s *InboundService) addTrafficLocked(inboundTraffics []*xray.Traffic, clien
 	})
 	if err != nil {
 		logger.Warning("traffic lifecycle maintenance failed after traffic commit:", err)
-		return false, false, nil, nil, nil
+		return false, false, nil, nil, nil, nil
 	}
 	needRestart = needRestart || s.applyTrafficMutationBatch(batch)
-	return needRestart, clientsDisabled, disabledNodeIDs, batch.remotePlans, nil
+	return needRestart, clientsDisabled, disabledNodeIDs, batch.remotePlans, batch.renewedEmails, nil
 }
 
 func (s *InboundService) addInboundTraffic(tx *gorm.DB, traffics []*xray.Traffic) error {
@@ -515,6 +517,7 @@ func (s *InboundService) autoRenewClients(tx *gorm.DB, mutationBatch *trafficMut
 	if err = clearGlobalTraffic(tx, renewedEmails...); err != nil {
 		return false, 0, err
 	}
+	mutationBatch.renewedEmails = append(mutationBatch.renewedEmails, renewedEmails...)
 	for _, clientToAdd := range clientsToAdd {
 		if clientToAdd.inbound.NodeID != nil {
 			mutationBatch.addNode(*clientToAdd.inbound.NodeID)
@@ -648,33 +651,23 @@ func (s *InboundService) ResetClientTrafficByEmail(clientEmail string) error {
 }
 
 func (s *InboundService) ResetClientTraffic(id int, clientEmail string) (needRestart bool, err error) {
-	var resetInbound *model.Inbound
+	var ownNode *int
 	err = submitTrafficWrite(func() error {
 		var inner error
-		needRestart, resetInbound, inner = s.resetClientTrafficLocked(id, clientEmail)
+		needRestart, ownNode, inner = s.resetClientTrafficLocked(id, clientEmail)
 		return inner
 	})
 	if err == nil {
 		s.resetMtprotoClientQuota(clientEmail)
-		if resetInbound != nil && resetInbound.NodeID != nil {
-			// Attempted whatever the node's status: nothing replays a reset, so a
-			// node still serving after being marked offline must get it now.
-			if rt, rterr := s.runtimeFor(resetInbound); rterr != nil {
-				logger.Warning("ResetClientTraffic: runtime lookup failed:", rterr)
-			} else {
-				ctx, cancel := nodePushContext()
-				e := rt.ResetClientTraffic(ctx, resetInbound, clientEmail)
-				cancel()
-				if e != nil {
-					logger.Warning("ResetClientTraffic: remote propagation to", rt.Name(), "failed:", e)
-				}
-			}
+		// Siblings on other nodes are delivered by their own inbound's reset.
+		if ownNode != nil {
+			s.deliverNodeResetsNow([]int{*ownNode})
 		}
 	}
 	return
 }
 
-func (s *InboundService) resetClientTrafficLocked(id int, clientEmail string) (bool, *model.Inbound, error) {
+func (s *InboundService) resetClientTrafficLocked(id int, clientEmail string) (bool, *int, error) {
 	needRestart := false
 	var reenablePlan *trafficLocalApplyPlan
 	var reenableNodeID *int
@@ -747,6 +740,9 @@ func (s *InboundService) resetClientTrafficLocked(id int, clientEmail string) (b
 		if err := tx.Where("email = ?", clientEmail).Delete(&model.NodeClientTraffic{}).Error; err != nil {
 			return err
 		}
+		if _, err := queueNodeResets(tx, []string{clientEmail}); err != nil {
+			return err
+		}
 		if err := tx.Model(model.Inbound{}).
 			Where("id = ?", id).
 			Update("last_traffic_reset_time", now).Error; err != nil {
@@ -775,7 +771,10 @@ func (s *InboundService) resetClientTrafficLocked(id int, clientEmail string) (b
 		}
 	}
 
-	return needRestart, inbound, nil
+	if inbound != nil {
+		return needRestart, inbound.NodeID, nil
+	}
+	return needRestart, nil, nil
 }
 
 func (s *InboundService) ResetAllTraffics() error {
@@ -784,7 +783,6 @@ func (s *InboundService) ResetAllTraffics() error {
 	})
 	if err == nil {
 		s.propagateResetAllTrafficsToNodes()
-		s.resetAllMtprotoQuotas()
 	}
 	return err
 }

+ 2 - 0
internal/web/service/inbound_traffic_apply.go

@@ -30,6 +30,8 @@ type trafficMutationBatch struct {
 	localPlans  []trafficLocalApplyPlan
 	remotePlans []trafficInboundUpdatePlan
 	nodeIDs     map[int]struct{}
+	// renewedEmails get their MTProto sidecar quota zeroed once the tick commits.
+	renewedEmails []string
 }
 
 type trafficInboundUpdatePlan struct{ oldInbound, newInbound model.Inbound }

+ 22 - 32
internal/web/service/inbound_tuic.go

@@ -7,7 +7,6 @@ import (
 	"github.com/mhsanaei/3x-ui/v3/internal/database/model"
 	"github.com/mhsanaei/3x-ui/v3/internal/logger"
 	"github.com/mhsanaei/3x-ui/v3/internal/tuic"
-	"github.com/mhsanaei/3x-ui/v3/internal/xray"
 )
 
 func (s *InboundService) DesiredTuicInstances() ([]tuic.Instance, error) {
@@ -23,47 +22,38 @@ func (s *InboundService) DesiredTuicInstances() ([]tuic.Instance, error) {
 		return nil, nil
 	}
 
-	ids := make([]int, 0, len(inbounds))
-	for _, ib := range inbounds {
-		ids = append(ids, ib.Id)
-	}
-	var disabledRows []xray.ClientTraffic
-	err = db.Model(xray.ClientTraffic{}).
-		Where("inbound_id IN ? AND enable = ?", ids, false).
-		Select("inbound_id", "email").
-		Find(&disabledRows).Error
-	if err != nil {
-		return nil, err
-	}
-	disabled := make(map[int]map[string]struct{}, len(disabledRows))
-	for _, row := range disabledRows {
-		if disabled[row.InboundId] == nil {
-			disabled[row.InboundId] = map[string]struct{}{}
-		}
-		disabled[row.InboundId][row.Email] = struct{}{}
-	}
-
 	instances := make([]tuic.Instance, 0, len(inbounds))
 	for _, ib := range inbounds {
 		inst, ok := tuic.InstanceFromInbound(ib)
 		if !ok {
 			continue
 		}
-		if off := disabled[ib.Id]; len(off) > 0 {
-			kept := make([]tuic.TuicClientSettings, 0, len(inst.Clients))
-			for _, c := range inst.Clients {
-				if _, skip := off[c.Email]; !skip {
-					kept = append(kept, c)
-				}
+		instances = append(instances, inst)
+	}
+	emails := make([]string, 0)
+	for _, inst := range instances {
+		for _, e := range inst.Clients {
+			emails = append(emails, e.Email)
+		}
+	}
+	disabled, err := trafficDisabledEmails(db, emails)
+	if err != nil {
+		return nil, err
+	}
+	served := instances[:0]
+	for _, inst := range instances {
+		kept := make([]tuic.TuicClientSettings, 0, len(inst.Clients))
+		for _, e := range inst.Clients {
+			if _, off := disabled[e.Email]; !off {
+				kept = append(kept, e)
 			}
-			inst.Clients = kept
 		}
-		if len(inst.Clients) == 0 {
-			continue
+		inst.Clients = kept
+		if len(kept) > 0 {
+			served = append(served, inst)
 		}
-		instances = append(instances, inst)
 	}
-	return instances, nil
+	return served, nil
 }
 
 func (s *InboundService) applyLocalTuic(inboundId int) {

+ 170 - 0
internal/web/service/inbound_update_stale_form_test.go

@@ -0,0 +1,170 @@
+package service
+
+import (
+	"testing"
+	"time"
+
+	"github.com/mhsanaei/3x-ui/v3/internal/database"
+	"github.com/mhsanaei/3x-ui/v3/internal/database/model"
+	"github.com/mhsanaei/3x-ui/v3/internal/xray"
+
+	"gorm.io/gorm"
+)
+
+func readClientTraffic(t *testing.T, email string) xray.ClientTraffic {
+	t.Helper()
+	var row xray.ClientTraffic
+	if err := database.GetDB().Where("email = ?", email).First(&row).Error; err != nil {
+		t.Fatalf("read client_traffics %s: %v", email, err)
+	}
+	return row
+}
+
+// The inbound modal posts the clients it loaded on open; a renewal committed
+// while it was open must survive the save in settings, record and stats.
+func TestInboundFormSaveKeepsClientRenewedWhileOpen(t *testing.T) {
+	setupBulkDB(t)
+	ib := seedRenewableNeighbour(t, 23201, nil)
+	form := *ib
+
+	if err := database.GetDB().Transaction(autoRenewTick); err != nil {
+		t.Fatalf("autoRenew: %v", err)
+	}
+	form.Remark = "edited"
+	if _, _, err := (&InboundService{}).UpdateInbound(&form); err != nil {
+		t.Fatalf("UpdateInbound: %v", err)
+	}
+
+	requireNeighbourRenewed(t, ib.Id)
+	now := time.Now().UnixMilli()
+	if row := readClientTraffic(t, "y@stale"); !row.Enable || row.ExpiryTime <= now {
+		t.Fatalf("client_traffics rolled back: enable=%v expiryTime=%d", row.Enable, row.ExpiryTime)
+	}
+	if rec := lookupClientRecord(t, "y@stale"); !rec.Enable || rec.ExpiryTime <= now {
+		t.Fatalf("client record rolled back: enable=%v expiryTime=%d", rec.Enable, rec.ExpiryTime)
+	}
+}
+
+// On a node the master's push is authoritative, lifecycle fields included.
+func TestInboundUpdateFromMasterAppliesClientLifecycle(t *testing.T) {
+	setupBulkDB(t)
+	ib := seedRenewableNeighbour(t, 23202, nil)
+	clients, err := (&InboundService{}).GetClients(ib)
+	if err != nil {
+		t.Fatalf("GetClients: %v", err)
+	}
+	for i := range clients {
+		if clients[i].Email == "x@stale" {
+			clients[i].Enable = false
+		}
+	}
+	push := *ib
+	push.Settings = clientsSettings(t, clients)
+	if _, _, err := (&InboundService{FromNodeSync: true}).UpdateInbound(&push); err != nil {
+		t.Fatalf("UpdateInbound: %v", err)
+	}
+	if x, _ := settingsClient(t, ib.Id, "x@stale"); x.Enable {
+		t.Fatal("master push disabling x@stale was ignored in settings")
+	}
+	if readClientTraffic(t, "x@stale").Enable {
+		t.Fatal("master push disabling x@stale was ignored in client_traffics")
+	}
+}
+
+// Traffic the poll adds after UpdateInbound read the row must not be written
+// back over by the edit.
+func TestInboundUpdateKeepsTrafficAddedMidEdit(t *testing.T) {
+	setupBulkDB(t)
+	ib := seedRenewableNeighbour(t, 23203, nil)
+	form := *ib
+	form.Remark = "edited"
+	commitTickBetweenReadAndWrite(t, func(tx *gorm.DB) error {
+		return (&InboundService{}).addInboundTraffic(tx, []*xray.Traffic{{IsInbound: true, Tag: ib.Tag, Up: 100, Down: 50}})
+	}, func() {
+		if _, _, err := (&InboundService{}).UpdateInbound(&form); err != nil {
+			t.Errorf("UpdateInbound: %v", err)
+		}
+	})
+	saved, err := (&InboundService{}).GetInbound(ib.Id)
+	if err != nil {
+		t.Fatalf("GetInbound: %v", err)
+	}
+	if saved.Up != 100 || saved.Down != 50 || saved.Remark != "edited" {
+		t.Fatalf("inbound after edit: up=%d down=%d remark=%q, want 100/50/edited", saved.Up, saved.Down, saved.Remark)
+	}
+}
+
+func inboundLinksEmail(t *testing.T, inboundId int, email string) bool {
+	t.Helper()
+	var n int64
+	if err := database.GetDB().Table("client_inbounds").
+		Joins("JOIN clients ON clients.id = client_inbounds.client_id").
+		Where("client_inbounds.inbound_id = ? AND clients.email = ?", inboundId, email).
+		Count(&n).Error; err != nil {
+		t.Fatalf("count links: %v", err)
+	}
+	return n > 0
+}
+
+// A client added while the modal was open is not in the list it posts back;
+// saving the inbound must not detach it.
+func TestInboundFormSaveKeepsClientAddedWhileOpen(t *testing.T) {
+	setupBulkDB(t)
+	ib := seedRenewableNeighbour(t, 23204, nil)
+	form := *ib
+
+	if _, err := (&ClientService{}).AddInboundClient(&InboundService{}, &model.Inbound{
+		Id: ib.Id, Settings: clientsSettings(t, []model.Client{{Email: "z@stale", ID: "aaaaaaaa-0000-0000-0000-00000000000c", Enable: true}}),
+	}); err != nil {
+		t.Fatalf("AddInboundClient: %v", err)
+	}
+	form.Remark = "edited"
+	if _, _, err := (&InboundService{}).UpdateInbound(&form); err != nil {
+		t.Fatalf("UpdateInbound: %v", err)
+	}
+	if _, ok := settingsClient(t, ib.Id, "z@stale"); !ok || !inboundLinksEmail(t, ib.Id, "z@stale") {
+		t.Fatalf("client added while the form was open was dropped: in settings=%v linked=%v", ok, inboundLinksEmail(t, ib.Id, "z@stale"))
+	}
+}
+
+// A client deleted while the modal was open is still in the list it posts
+// back; saving must not restore its access.
+func TestInboundFormSaveDoesNotResurrectDeletedClient(t *testing.T) {
+	setupBulkDB(t)
+	ib := seedRenewableNeighbour(t, 23205, nil)
+	form := *ib
+
+	if _, err := (&ClientService{}).DelInboundClientByEmail(&InboundService{}, ib.Id, "x@stale", false, true); err != nil {
+		t.Fatalf("DelInboundClientByEmail: %v", err)
+	}
+	form.Remark = "edited"
+	if _, _, err := (&InboundService{}).UpdateInbound(&form); err != nil {
+		t.Fatalf("UpdateInbound: %v", err)
+	}
+	if _, ok := settingsClient(t, ib.Id, "x@stale"); ok || inboundLinksEmail(t, ib.Id, "x@stale") {
+		t.Fatalf("deleted client came back: in settings=%v linked=%v", ok, inboundLinksEmail(t, ib.Id, "x@stale"))
+	}
+}
+
+// An inbound switched off while the modal was open stays off when the form,
+// which still holds enable=true, is saved.
+func TestInboundFormSaveKeepsEnableToggledWhileOpen(t *testing.T) {
+	setupBulkDB(t)
+	ib := seedRenewableNeighbour(t, 23206, nil)
+	form := *ib
+
+	if _, err := (&InboundService{}).SetInboundEnable(ib.Id, false); err != nil {
+		t.Fatalf("SetInboundEnable: %v", err)
+	}
+	form.Remark = "edited"
+	if _, _, err := (&InboundService{}).UpdateInbound(&form); err != nil {
+		t.Fatalf("UpdateInbound: %v", err)
+	}
+	saved, err := (&InboundService{}).GetInbound(ib.Id)
+	if err != nil {
+		t.Fatalf("GetInbound: %v", err)
+	}
+	if saved.Enable || saved.Remark != "edited" {
+		t.Fatalf("after save: enable=%v remark=%q, want disabled and edited", saved.Enable, saved.Remark)
+	}
+}

+ 75 - 0
internal/web/service/mtproto_fake_test.go

@@ -2,8 +2,13 @@ package service
 
 import (
 	"fmt"
+	"net"
+	"net/http"
+	"net/url"
 	"os"
 	"path/filepath"
+	"regexp"
+	"slices"
 	"strings"
 	"testing"
 	"time"
@@ -46,9 +51,79 @@ func fakeMtgChildMain() {
 		fmt.Fprintf(f, "%d\n", os.Getpid())
 		f.Close()
 	}
+	if logPath := os.Getenv("MTG_FAKE_APILOG"); logPath != "" && len(os.Args) > 2 {
+		go serveFakeMtgAPI(os.Args[len(os.Args)-1], logPath)
+	}
 	select {}
 }
 
+// serveFakeMtgAPI answers the management API on the config's api-bind-to and
+// logs each reset-quota call, so a test sees which sidecar quotas were zeroed.
+func serveFakeMtgAPI(configPath, logPath string) {
+	cfg, err := os.ReadFile(configPath)
+	if err != nil {
+		return
+	}
+	m := regexp.MustCompile(`api-bind-to = "([^"]+)"`).FindSubmatch(cfg)
+	if m == nil {
+		return
+	}
+	ln, err := net.Listen("tcp", string(m[1]))
+	if err != nil {
+		return
+	}
+	appendFakeMtgLog(logPath, "ready")
+	_ = http.Serve(ln, http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
+		if name, ok := strings.CutSuffix(strings.TrimPrefix(r.URL.Path, "/secrets/"), "/reset-quota"); ok && r.Method == http.MethodPost {
+			if unescaped, err := url.PathUnescape(name); err == nil {
+				appendFakeMtgLog(logPath, "reset:"+unescaped)
+			}
+		}
+		_, _ = w.Write([]byte("{}"))
+	}))
+}
+
+func appendFakeMtgLog(path, line string) {
+	if f, err := os.OpenFile(path, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0o644); err == nil {
+		fmt.Fprintln(f, line)
+		f.Close()
+	}
+}
+
+// installFakeMtgAPI is installFakeMtg whose children also serve the management
+// API; it returns the pid file and the API call log.
+func installFakeMtgAPI(t *testing.T) (string, string) {
+	t.Helper()
+	pidFile := installFakeMtg(t)
+	logPath := filepath.Join(filepath.Dir(pidFile), "mtg-api.log")
+	t.Setenv("MTG_FAKE_APILOG", logPath)
+	return pidFile, logPath
+}
+
+func fakeMtgLog(t *testing.T, logPath string) []string {
+	t.Helper()
+	data, err := os.ReadFile(logPath)
+	if os.IsNotExist(err) {
+		return nil
+	}
+	if err != nil {
+		t.Fatalf("read mtg api log: %v", err)
+	}
+	return strings.Fields(string(data))
+}
+
+// waitFakeMtgLog polls until the log holds want, failing on timeout.
+func waitFakeMtgLog(t *testing.T, logPath, want string) {
+	t.Helper()
+	deadline := time.Now().Add(5 * time.Second)
+	for !slices.Contains(fakeMtgLog(t, logPath), want) {
+		if time.Now().After(deadline) {
+			t.Fatalf("mtg api log never recorded %q: %v", want, fakeMtgLog(t, logPath))
+		}
+		time.Sleep(20 * time.Millisecond)
+	}
+}
+
 // installFakeMtg points the mtproto manager at a copy of the running test
 // binary posing as mtg (via the MTG_FAKE_CHILD gate in TestMain) and returns
 // the pid file whose line count equals the number of processes spawned so far.

+ 137 - 0
internal/web/service/mtproto_quota_reset_test.go

@@ -0,0 +1,137 @@
+package service
+
+import (
+	"slices"
+	"strings"
+	"testing"
+	"time"
+
+	"github.com/mhsanaei/3x-ui/v3/internal/database"
+	"github.com/mhsanaei/3x-ui/v3/internal/database/model"
+	"github.com/mhsanaei/3x-ui/v3/internal/mtproto"
+	"github.com/mhsanaei/3x-ui/v3/internal/web/runtime"
+	"github.com/mhsanaei/3x-ui/v3/internal/xray"
+)
+
+// startQuotaSidecar runs a local MTProto inbound for mtga and mtgb under the fake
+// mtg and returns its API log once the sidecar answers.
+func startQuotaSidecar(t *testing.T, port int, mtga model.Client) (*model.Inbound, string) {
+	t.Helper()
+	setupConflictDB(t)
+	pidFile, logPath := installFakeMtgAPI(t)
+	runtime.SetManager(runtime.NewManager(runtime.LocalDeps{APIPort: func() int { return 0 }, SetNeedRestart: func() {}}))
+	t.Cleanup(func() { runtime.SetManager(nil) })
+
+	mtga.Email, mtga.Secret = "mtga", mtprotoTestSecretA
+	clients := []model.Client{mtga, {Email: "mtgb", Secret: mtprotoTestSecretB, Enable: true}}
+	ib := &model.Inbound{Tag: "mt-quota", Enable: true, Port: port, Protocol: model.MTProto, Settings: clientsSettings(t, clients)}
+	if err := database.GetDB().Create(ib).Error; err != nil {
+		t.Fatalf("create inbound: %v", err)
+	}
+	if err := (&ClientService{}).SyncInbound(nil, ib.Id, clients); err != nil {
+		t.Fatalf("SyncInbound: %v", err)
+	}
+	for _, c := range clients {
+		row := xray.ClientTraffic{InboundId: ib.Id, Email: c.Email, Enable: true, Up: 5, Total: c.TotalGB, ExpiryTime: c.ExpiryTime, Reset: c.Reset}
+		if err := database.GetDB().Create(&row).Error; err != nil {
+			t.Fatalf("seed traffic: %v", err)
+		}
+	}
+	// A running sidecar needs a served client, so prime with the healthy set.
+	inst, ok := mtproto.InstanceFromInbound(&model.Inbound{
+		Id: ib.Id, Tag: ib.Tag, Port: port, Protocol: model.MTProto,
+		Settings: clientsSettings(t, []model.Client{{Email: "mtgb", Secret: mtprotoTestSecretB, Enable: true}}),
+	})
+	if !ok {
+		t.Fatal("seed inbound must produce an mtg instance")
+	}
+	if err := mtproto.GetManager().Ensure(inst); err != nil {
+		t.Fatalf("start mtg: %v", err)
+	}
+	t.Cleanup(func() { mtproto.GetManager().Remove(ib.Id) })
+	waitForSpawns(t, pidFile, 1)
+	waitFakeMtgLog(t, logPath, "ready")
+	return ib, logPath
+}
+
+func quotaResets(t *testing.T, logPath string) []string {
+	t.Helper()
+	var out []string
+	for _, line := range fakeMtgLog(t, logPath) {
+		if name, ok := strings.CutPrefix(line, "reset:"); ok {
+			out = append(out, name)
+		}
+	}
+	slices.Sort(out)
+	return out
+}
+
+// Every path that zeroes a client's panel counters must zero the sidecar's own
+// quota counter too, or the sidecar keeps refusing the client.
+func TestPanelResetsZeroSidecarQuota(t *testing.T) {
+	t.Run("bulk reset", func(t *testing.T) {
+		_, logPath := startQuotaSidecar(t, 46201, model.Client{Enable: true})
+		if _, err := (&ClientService{}).BulkResetTraffic(&InboundService{}, []string{"mtga"}); err != nil {
+			t.Fatalf("BulkResetTraffic: %v", err)
+		}
+		if got := quotaResets(t, logPath); !slices.Equal(got, []string{"mtga"}) {
+			t.Fatalf("sidecar quota resets %v, want [mtga]", got)
+		}
+	})
+	t.Run("inbound clients", func(t *testing.T) {
+		ib, logPath := startQuotaSidecar(t, 46207, model.Client{Enable: true})
+		if err := (&ClientService{}).ResetAllClientTraffics(&InboundService{}, ib.Id); err != nil {
+			t.Fatalf("ResetAllClientTraffics: %v", err)
+		}
+		if got := quotaResets(t, logPath); !slices.Equal(got, []string{"mtga", "mtgb"}) {
+			t.Fatalf("sidecar quota resets %v, want [mtga mtgb]", got)
+		}
+	})
+	t.Run("reset all", func(t *testing.T) {
+		_, logPath := startQuotaSidecar(t, 46202, model.Client{Enable: true})
+		if _, err := (&ClientService{}).ResetAllTraffics(); err != nil {
+			t.Fatalf("ResetAllTraffics: %v", err)
+		}
+		if got := quotaResets(t, logPath); !slices.Equal(got, []string{"mtga", "mtgb"}) {
+			t.Fatalf("sidecar quota resets %v, want [mtga mtgb]", got)
+		}
+	})
+	t.Run("auto renew", func(t *testing.T) {
+		expired := time.Now().Add(-time.Hour).UnixMilli()
+		_, logPath := startQuotaSidecar(t, 46203, model.Client{Enable: true, Reset: 30, ExpiryTime: expired})
+		if _, _, err := (&InboundService{}).AddTraffic(nil, nil); err != nil {
+			t.Fatalf("AddTraffic: %v", err)
+		}
+		if got := quotaResets(t, logPath); !slices.Equal(got, []string{"mtga"}) {
+			t.Fatalf("sidecar quota resets %v, want [mtga]", got)
+		}
+	})
+}
+
+// Resetting inbound counters leaves every client's usage in place, so the
+// sidecar's quota counters must stay too or clients get their quota again free.
+func TestInboundResetAllKeepsSidecarQuota(t *testing.T) {
+	_, logPath := startQuotaSidecar(t, 46204, model.Client{Enable: true})
+	if err := (&InboundService{}).ResetAllTraffics(); err != nil {
+		t.Fatalf("ResetAllTraffics: %v", err)
+	}
+	if got := quotaResets(t, logPath); len(got) != 0 {
+		t.Fatalf("inbound reset zeroed sidecar quotas %v, want none", got)
+	}
+}
+
+// Resetting one inbound's clients zeroes only their sidecar quotas, not those
+// of MTProto clients whose usage the reset left in place.
+func TestInboundClientResetKeepsOtherSidecarQuotas(t *testing.T) {
+	_, logPath := startQuotaSidecar(t, 46205, model.Client{Enable: true})
+	other := mkInbound(t, 46206, model.VLESS, clientsSettings(t, []model.Client{{Email: "vless-only", ID: "11111111-1111-1111-1111-1111111111ab", Enable: true}}))
+	if err := (&ClientService{}).SyncInbound(nil, other.Id, []model.Client{{Email: "vless-only", ID: "11111111-1111-1111-1111-1111111111ab", Enable: true}}); err != nil {
+		t.Fatalf("SyncInbound: %v", err)
+	}
+	if err := (&ClientService{}).ResetAllClientTraffics(&InboundService{}, other.Id); err != nil {
+		t.Fatalf("ResetAllClientTraffics: %v", err)
+	}
+	if got := quotaResets(t, logPath); len(got) != 0 {
+		t.Fatalf("resetting another inbound zeroed sidecar quotas %v, want none", got)
+	}
+}

+ 3 - 0
internal/web/service/node.go

@@ -864,6 +864,9 @@ func (s *NodeService) Delete(id int) error {
 		if err := tx.Where("node_id = ?", id).Delete(&model.NodeClientTraffic{}).Error; err != nil {
 			return err
 		}
+		if err := tx.Where("node_id = ?", id).Delete(&model.NodePendingReset{}).Error; err != nil {
+			return err
+		}
 		guids := []string{synthNodeGuid(id)}
 		if guid != "" {
 			guids = append(guids, guid)

+ 152 - 0
internal/web/service/node_reset_queue.go

@@ -0,0 +1,152 @@
+package service
+
+import (
+	"context"
+	"sync"
+	"time"
+
+	"github.com/mhsanaei/3x-ui/v3/internal/database"
+	"github.com/mhsanaei/3x-ui/v3/internal/database/model"
+	"github.com/mhsanaei/3x-ui/v3/internal/logger"
+	"github.com/mhsanaei/3x-ui/v3/internal/web/runtime"
+
+	"gorm.io/gorm"
+	"gorm.io/gorm/clause"
+)
+
+// nodeBulkResetter is a node runtime that can zero many clients in one call.
+type nodeBulkResetter interface {
+	ResetClientTraffics(ctx context.Context, emails []string) error
+}
+
+type nodeEmail struct {
+	NodeId int    `gorm:"column:node_id"`
+	Email  string `gorm:"column:email"`
+}
+
+// queueNodeResets records a reset for every node hosting one of emails (all
+// node-hosted clients when emails is nil) and returns the nodes involved.
+func queueNodeResets(tx *gorm.DB, emails []string) ([]int, error) {
+	base := func() *gorm.DB {
+		return tx.Table("clients").
+			Select("DISTINCT inbounds.node_id AS node_id, clients.email AS email").
+			Joins("JOIN client_inbounds ON client_inbounds.client_id = clients.id").
+			Joins("JOIN inbounds ON inbounds.id = client_inbounds.inbound_id").
+			Where("inbounds.node_id IS NOT NULL")
+	}
+	var pairs []nodeEmail
+	if emails == nil {
+		if err := base().Scan(&pairs).Error; err != nil {
+			return nil, err
+		}
+	} else {
+		for _, batch := range chunkStrings(uniqueNonEmptyStrings(emails), sqlInChunk) {
+			var page []nodeEmail
+			if err := base().Where("clients.email IN ?", batch).Scan(&page).Error; err != nil {
+				return nil, err
+			}
+			pairs = append(pairs, page...)
+		}
+	}
+	if len(pairs) == 0 {
+		return nil, nil
+	}
+	now := time.Now().UnixNano()
+	rows := make([]model.NodePendingReset, 0, len(pairs))
+	nodes := make(map[int]struct{})
+	for _, p := range pairs {
+		rows = append(rows, model.NodePendingReset{NodeId: p.NodeId, Email: p.Email, QueuedAt: now})
+		nodes[p.NodeId] = struct{}{}
+	}
+	if err := tx.Clauses(clause.OnConflict{
+		Columns:   []clause.Column{{Name: "node_id"}, {Name: "email"}},
+		DoUpdates: clause.AssignmentColumns([]string{"queued_at"}),
+	}).CreateInBatches(rows, 200).Error; err != nil {
+		return nil, err
+	}
+	ids := make([]int, 0, len(nodes))
+	for id := range nodes {
+		ids = append(ids, id)
+	}
+	return ids, nil
+}
+
+// pendingNodeResetEmails lists the clients whose reset the node still owes.
+func pendingNodeResetEmails(tx *gorm.DB, nodeID int) (map[string]struct{}, error) {
+	var emails []string
+	if err := tx.Model(&model.NodePendingReset{}).Where("node_id = ?", nodeID).Pluck("email", &emails).Error; err != nil {
+		return nil, err
+	}
+	out := make(map[string]struct{}, len(emails))
+	for _, e := range emails {
+		out[e] = struct{}{}
+	}
+	return out, nil
+}
+
+var nodeResetDeliveryLocks sync.Map
+
+// DeliverNodeResets sends the node every reset it has not confirmed. A row is
+// dropped only after the node accepted it and only if nothing re-queued it since.
+func (s *InboundService) DeliverNodeResets(ctx context.Context, nodeID int, rt runtime.Runtime) error {
+	lock, _ := nodeResetDeliveryLocks.LoadOrStore(nodeID, &sync.Mutex{})
+	lock.(*sync.Mutex).Lock()
+	defer lock.(*sync.Mutex).Unlock()
+	db := database.GetDB()
+	var rows []model.NodePendingReset
+	if err := db.Where("node_id = ?", nodeID).Order("id").Find(&rows).Error; err != nil {
+		return err
+	}
+	if len(rows) == 0 {
+		return nil
+	}
+	bulk, canBulk := rt.(nodeBulkResetter)
+	for start := 0; start < len(rows); start += sqlInChunk {
+		batch := rows[start:min(start+sqlInChunk, len(rows))]
+		emails := make([]string, len(batch))
+		for i := range batch {
+			emails[i] = batch[i].Email
+		}
+		var err error
+		if canBulk && len(batch) > nodeBulkPushThreshold {
+			err = bulk.ResetClientTraffics(ctx, emails)
+		} else {
+			for _, email := range emails {
+				if err = rt.ResetClientTraffic(ctx, nil, email); err != nil {
+					break
+				}
+			}
+		}
+		if err != nil {
+			return err
+		}
+		for i := range batch {
+			if err := db.Where("id = ? AND queued_at = ?", batch[i].Id, batch[i].QueuedAt).
+				Delete(&model.NodePendingReset{}).Error; err != nil {
+				return err
+			}
+		}
+	}
+	return nil
+}
+
+// deliverNodeResetsNow tries each node once right after a reset commits; what
+// fails stays queued for the node sync job.
+func (s *InboundService) deliverNodeResetsNow(nodeIDs []int) {
+	mgr := runtime.GetManager()
+	if mgr == nil || len(nodeIDs) == 0 {
+		return
+	}
+	fanoutInboundResults(nodeIDs, nodeFanoutConcurrency, func(i int) struct{} {
+		rt, err := mgr.RuntimeFor(&nodeIDs[i])
+		if err != nil {
+			return struct{}{}
+		}
+		ctx, cancel := nodePushContext()
+		defer cancel()
+		if err := s.DeliverNodeResets(ctx, nodeIDs[i], rt); err != nil {
+			logger.Warning("reset delivery to", rt.Name(), "deferred to the next sync:", err)
+		}
+		return struct{}{}
+	})
+}

+ 275 - 0
internal/web/service/node_reset_undelivered_test.go

@@ -0,0 +1,275 @@
+package service
+
+import (
+	"context"
+	"errors"
+	"fmt"
+	"slices"
+	"sync"
+	"sync/atomic"
+	"testing"
+	"time"
+
+	"github.com/mhsanaei/3x-ui/v3/internal/database"
+	"github.com/mhsanaei/3x-ui/v3/internal/database/model"
+	"github.com/mhsanaei/3x-ui/v3/internal/xray"
+
+	"gorm.io/gorm"
+)
+
+const (
+	resetLostOn  = `{"clients":[{"email":"reset-lost","totalGB":100,"enable":true}]}`
+	resetLostOff = `{"clients":[{"email":"reset-lost","totalGB":100,"enable":false}]}`
+)
+
+// seedLatchedNodeClient leaves reset-lost depleted and latched off on the
+// master by its node's own usage, as a real node sync does.
+func seedLatchedNodeClient(t *testing.T, svc *InboundService) (*gorm.DB, *model.Inbound) {
+	t.Helper()
+	db := initTrafficTestDB(t)
+	createNodeInboundWithClient(t, db, 1, "n1-in", 41901, "reset-lost")
+	syncNodeWithSettings(t, svc, 1, "n1-in", resetLostOn, xray.ClientTraffic{Email: "reset-lost", Up: 10, Down: 10, Total: 100, Enable: true})
+	syncNodeWithSettings(t, svc, 1, "n1-in", resetLostOff, xray.ClientTraffic{Email: "reset-lost", Up: 60, Down: 60, Total: 100, Enable: false})
+	if got := readTraffic(t, db, "reset-lost"); got.Enable {
+		t.Fatal("setup: the depleted client should be latched off")
+	}
+	var ib model.Inbound
+	if err := db.Where("tag = ?", "n1-in").First(&ib).Error; err != nil {
+		t.Fatalf("load inbound: %v", err)
+	}
+	return db, &ib
+}
+
+// A reset the node never received leaves its old counters, so the node keeps
+// switching the client off; the master must not adopt that verdict.
+func TestNodeResetNotDeliveredDoesNotRedisableClient(t *testing.T) {
+	resets := []struct {
+		name string
+		run  func(svc *InboundService, ib *model.Inbound) error
+	}{
+		{"single", func(svc *InboundService, ib *model.Inbound) error {
+			_, err := svc.ResetClientTraffic(ib.Id, "reset-lost")
+			return err
+		}},
+		{"bulk", func(svc *InboundService, _ *model.Inbound) error {
+			_, err := (&ClientService{}).BulkResetTraffic(svc, []string{"reset-lost"})
+			return err
+		}},
+		{"inbound", func(svc *InboundService, ib *model.Inbound) error {
+			return (&ClientService{}).ResetAllClientTraffics(svc, ib.Id)
+		}},
+		{"all", func(*InboundService, *model.Inbound) error {
+			_, err := (&ClientService{}).ResetAllTraffics()
+			return err
+		}},
+	}
+	for _, reset := range resets {
+		t.Run(reset.name, func(t *testing.T) {
+			svc := &InboundService{}
+			db, ib := seedLatchedNodeClient(t, svc)
+			if err := reset.run(svc, ib); err != nil {
+				t.Fatalf("reset: %v", err)
+			}
+			syncNodeWithSettings(t, svc, 1, "n1-in", resetLostOff, xray.ClientTraffic{Email: "reset-lost", Up: 60, Down: 60, Total: 100, Enable: false})
+			got := readTraffic(t, db, "reset-lost")
+			if !got.Enable || got.Up+got.Down != 0 {
+				t.Fatalf("after reset: enable=%v used=%d, want enabled at 0 — the undelivered reset re-disabled it", got.Enable, got.Up+got.Down)
+			}
+		})
+	}
+}
+
+// resetRecordingRuntime is a node that accepts per-client resets unless failing.
+type resetRecordingRuntime struct {
+	fakeNodeRuntime
+	mu   sync.Mutex
+	fail bool
+	got  []string
+}
+
+func (r *resetRecordingRuntime) ResetClientTraffic(_ context.Context, _ *model.Inbound, email string) error {
+	r.mu.Lock()
+	defer r.mu.Unlock()
+	if r.fail {
+		return errors.New("node unreachable")
+	}
+	r.got = append(r.got, email)
+	return nil
+}
+
+func (r *resetRecordingRuntime) delivered() []string {
+	r.mu.Lock()
+	defer r.mu.Unlock()
+	return slices.Clone(r.got)
+}
+
+func pendingResetEmails(t *testing.T, nodeID int) []string {
+	t.Helper()
+	var emails []string
+	if err := database.GetDB().Model(&model.NodePendingReset{}).Where("node_id = ?", nodeID).
+		Order("email").Pluck("email", &emails).Error; err != nil {
+		t.Fatalf("read pending resets: %v", err)
+	}
+	return emails
+}
+
+func setupRecordingNode(t *testing.T, fail bool) (int, *resetRecordingRuntime, *model.Inbound) {
+	t.Helper()
+	setupBulkDB(t)
+	mgr := useTestRuntimeManager(t)
+	node := &model.Node{Name: "reset-node", Address: "127.0.0.1", Port: 2096, ApiToken: "tok", Enable: true, Status: "online"}
+	if err := database.GetDB().Create(node).Error; err != nil {
+		t.Fatalf("create node: %v", err)
+	}
+	rec := &resetRecordingRuntime{fail: fail}
+	mgr.SetRuntimeOverride(node.Id, rec)
+	ib := nodeInbound(t, node.Id, 41911, []model.Client{{Email: "reset-lost", ID: "11111111-1111-1111-1111-1111111111aa", Enable: true}})
+	if err := (&InboundService{}).AddClientStat(database.GetDB(), ib.Id, &model.Client{Email: "reset-lost", Enable: true}); err != nil {
+		t.Fatalf("AddClientStat: %v", err)
+	}
+	return node.Id, rec, ib
+}
+
+// A reachable node gets the reset right after the master commits it.
+func TestNodeResetDeliveredRightAway(t *testing.T) {
+	resets := []struct {
+		name string
+		run  func(svc *InboundService, ib *model.Inbound) error
+	}{
+		{"single", func(svc *InboundService, ib *model.Inbound) error {
+			_, err := svc.ResetClientTraffic(ib.Id, "reset-lost")
+			return err
+		}},
+		{"bulk", func(svc *InboundService, _ *model.Inbound) error {
+			_, err := (&ClientService{}).BulkResetTraffic(svc, []string{"reset-lost"})
+			return err
+		}},
+		{"inbound", func(svc *InboundService, ib *model.Inbound) error {
+			return (&ClientService{}).ResetAllClientTraffics(svc, ib.Id)
+		}},
+		{"all", func(*InboundService, *model.Inbound) error {
+			_, err := (&ClientService{}).ResetAllTraffics()
+			return err
+		}},
+	}
+	for _, reset := range resets {
+		t.Run(reset.name, func(t *testing.T) {
+			nodeID, rec, ib := setupRecordingNode(t, false)
+			if err := reset.run(&InboundService{}, ib); err != nil {
+				t.Fatalf("reset: %v", err)
+			}
+			if got := rec.delivered(); !slices.Equal(got, []string{"reset-lost"}) {
+				t.Fatalf("node received resets %v, want [reset-lost]", got)
+			}
+			if left := pendingResetEmails(t, nodeID); len(left) != 0 {
+				t.Fatalf("delivered reset still queued: %v", left)
+			}
+		})
+	}
+}
+
+// bulkResetRuntime also takes a batch in one call.
+type bulkResetRuntime struct {
+	resetRecordingRuntime
+	batches [][]string
+}
+
+func (b *bulkResetRuntime) ResetClientTraffics(_ context.Context, emails []string) error {
+	b.mu.Lock()
+	defer b.mu.Unlock()
+	b.batches = append(b.batches, slices.Clone(emails))
+	return nil
+}
+
+// Above the per-client push threshold a backlog goes out as one bulk request,
+// not one round-trip per client.
+func TestNodeResetBacklogUsesBulkRequest(t *testing.T) {
+	setupBulkDB(t)
+	const nodeID = 7
+	rows := make([]model.NodePendingReset, nodeBulkPushThreshold+1)
+	for i := range rows {
+		rows[i] = model.NodePendingReset{NodeId: nodeID, Email: fmt.Sprintf("owed-%02d", i), QueuedAt: 1}
+	}
+	if err := database.GetDB().Create(&rows).Error; err != nil {
+		t.Fatalf("seed pending resets: %v", err)
+	}
+	rt := &bulkResetRuntime{}
+	if err := (&InboundService{}).DeliverNodeResets(context.Background(), nodeID, rt); err != nil {
+		t.Fatalf("DeliverNodeResets: %v", err)
+	}
+	if len(rt.batches) != 1 || len(rt.batches[0]) != len(rows) || len(rt.delivered()) != 0 {
+		t.Fatalf("bulk batches %d (first %d emails), per-client calls %d; want one batch of %d",
+			len(rt.batches), len(rt.batches[0]), len(rt.delivered()), len(rows))
+	}
+	if left := pendingResetEmails(t, nodeID); len(left) != 0 {
+		t.Fatalf("delivered backlog still queued: %d rows", len(left))
+	}
+}
+
+// An unreachable node keeps the reset queued until a later delivery lands.
+func TestNodeResetReplayedAfterFailure(t *testing.T) {
+	nodeID, rec, ib := setupRecordingNode(t, true)
+	if _, err := (&InboundService{}).ResetClientTraffic(ib.Id, "reset-lost"); err != nil {
+		t.Fatalf("ResetClientTraffic: %v", err)
+	}
+	if left := pendingResetEmails(t, nodeID); !slices.Equal(left, []string{"reset-lost"}) {
+		t.Fatalf("pending after failed delivery = %v, want [reset-lost]", left)
+	}
+
+	rec.mu.Lock()
+	rec.fail = false
+	rec.mu.Unlock()
+	if err := (&InboundService{}).DeliverNodeResets(context.Background(), nodeID, rec); err != nil {
+		t.Fatalf("DeliverNodeResets: %v", err)
+	}
+	if got := rec.delivered(); !slices.Equal(got, []string{"reset-lost"}) {
+		t.Fatalf("node received resets %v, want [reset-lost]", got)
+	}
+	if left := pendingResetEmails(t, nodeID); len(left) != 0 {
+		t.Fatalf("delivered reset still queued: %v", left)
+	}
+}
+
+// slowResetRuntime holds each reset until a second one arrives or a short
+// timeout passes, so two unserialized deliveries both reach the node.
+type slowResetRuntime struct {
+	resetRecordingRuntime
+	calls atomic.Int32
+	both  chan struct{}
+}
+
+func (r *slowResetRuntime) ResetClientTraffic(ctx context.Context, ib *model.Inbound, email string) error {
+	if r.calls.Add(1) == 2 {
+		close(r.both)
+	}
+	select {
+	case <-r.both:
+	case <-time.After(300 * time.Millisecond):
+	}
+	return r.resetRecordingRuntime.ResetClientTraffic(ctx, ib, email)
+}
+
+// The sync job and a reset's own delivery can run at once; the node must still
+// get each owed reset once, or usage made in between is wiped a second time.
+func TestConcurrentNodeResetDeliveriesSendOnce(t *testing.T) {
+	setupBulkDB(t)
+	const nodeID = 9
+	if err := database.GetDB().Create(&model.NodePendingReset{NodeId: nodeID, Email: "once", QueuedAt: 1}).Error; err != nil {
+		t.Fatalf("seed pending reset: %v", err)
+	}
+	rt := &slowResetRuntime{both: make(chan struct{})}
+	var wg sync.WaitGroup
+	for range 2 {
+		wg.Add(1)
+		go func() {
+			defer wg.Done()
+			if err := (&InboundService{}).DeliverNodeResets(context.Background(), nodeID, rt); err != nil {
+				t.Errorf("DeliverNodeResets: %v", err)
+			}
+		}()
+	}
+	wg.Wait()
+	if got := rt.delivered(); !slices.Equal(got, []string{"once"}) {
+		t.Fatalf("node received resets %v, want exactly [once]", got)
+	}
+}