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`
 `job/xray_traffic_job.go`, `job/node_traffic_sync_job.go`, `service/inbound_node.go`
 (`SetRemoteTraffic` / `upsertNodeBaseline`), models `xray.ClientTraffic`,
 (`SetRemoteTraffic` / `upsertNodeBaseline`), models `xray.ClientTraffic`,
 `model.NodeClientTraffic`, `model.ClientGlobalTraffic` (cross-master totals).
 `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`).
 Periodic resets: `job/periodic_traffic_reset_job.go` (keyed off `Inbound.TrafficReset`).
 
 
 ### 5.4 Background jobs (cron)
 ### 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                                                               |
 | `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                                                   |
 | `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)                                                                                                                               |
 | `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`                                                                                                                                         |
 | `NodeClientIp`                  | Per-node client IP attribution            | `NodeGuid`, `Email`, `Ips`                                                                                                                                         |
 | `ClientGlobalTraffic`           | Cross-master usage totals                 | `MasterGuid`, `Email`, `Up`, `Down`                                                                                                                                |
 | `ClientGlobalTraffic`           | Cross-master usage totals                 | `MasterGuid`, `Email`, `Up`, `Down`                                                                                                                                |
 | `xray.ClientTraffic`            | Per-client counters (`client_traffics`)   | `Email`, `Up`, `Down`, `Total`, `ExpiryTime`, `LastOnline`                                                                                                         |
 | `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.
         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'
       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
     - 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
     - depth: 2
       title: Toggle only the enable flag without serialising the whole settings JSON.
       title: Toggle only the enable flag without serialising the whole settings JSON.
         Recommended for UI switches on large inbounds.
         Recommended for UI switches on large inbounds.
@@ -151,10 +152,11 @@ _openapi:
           failures are reported per id and the rest still proceed. Restarts xray
           failures are reported per id and the rest still proceed. Restarts xray
           at most once.
           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
         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
       - content: Toggle only the enable flag without serialising the whole settings
           JSON. Recommended for UI switches on large inbounds.
           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
         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": [
         "tags": [
           "Inbounds"
           "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",
         "operationId": "post_panel_api_inbounds_update_id",
         "parameters": [
         "parameters": [
           {
           {

+ 1 - 1
frontend/public/openapi.json

@@ -5603,7 +5603,7 @@
         "tags": [
         "tags": [
           "Inbounds"
           "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",
         "operationId": "post_panel_api_inbounds_update_id",
         "parameters": [
         "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>;
   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);
     settingsPruned.clients = normalizeClients(values.protocol, settingsPruned.clients);
   }
   }
   let streamPruned = values.streamSettings
   let streamPruned = values.streamSettings

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

@@ -325,7 +325,7 @@ export const sections: readonly Section[] = [
         method: 'POST',
         method: 'POST',
         path: '/panel/api/inbounds/update/:id',
         path: '/panel/api/inbounds/update/:id',
         summary:
         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.' }],
         params: [{ name: 'id', in: 'path', type: 'number', desc: 'Inbound ID.' }],
         body: inboundBody,
         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 { HttpUtil, NumberFormatter, RandomUtil, SizeFormatter, Wireguard } from '@/utils';
 import type { RealityScanResult } from '@/generated/types';
 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 { createDefaultInboundSettings } from '@/lib/xray/inbound-defaults';
 import { generateAwgObfuscation } from '@/lib/xray/amneziawg-obfuscation';
 import { generateAwgObfuscation } from '@/lib/xray/amneziawg-obfuscation';
 import { composeInboundTag, isAutoInboundTag, type InboundTagInput } from '@/lib/xray/inbound-tag';
 import { composeInboundTag, isAutoInboundTag, type InboundTagInput } from '@/lib/xray/inbound-tag';
@@ -439,7 +443,9 @@ export default function InboundFormModal({
   useEffect(() => {
   useEffect(() => {
     if (!open) return;
     if (!open) return;
     const initial =
     const initial =
-      mode === 'edit' && dbInbound ? rawInboundToFormValues(dbInbound) : buildAddModeValues();
+      mode === 'edit' && dbInbound
+        ? withoutClients(rawInboundToFormValues(dbInbound))
+        : buildAddModeValues();
     methods.reset(initial);
     methods.reset(initial);
     setScanResult(null);
     setScanResult(null);
     setActiveTab('basic');
     setActiveTab('basic');
@@ -556,13 +562,8 @@ export default function InboundFormModal({
   }, [mode, methods]);
   }, [mode, methods]);
 
 
   const saveValues = async () => {
   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 values = methods.getValues() as InboundFormValues;
     const parsed = InboundFormSchema.safeParse(values);
     const parsed = InboundFormSchema.safeParse(values);
     if (!parsed.success) {
     if (!parsed.success) {
@@ -577,7 +578,7 @@ export default function InboundFormModal({
     }
     }
     setSaving(true);
     setSaving(true);
     try {
     try {
-      const payload = formValuesToWirePayload(parsed.data);
+      const payload = formValuesToWirePayload(parsed.data, { omitClients: mode === 'edit' });
       const url =
       const url =
         mode === 'edit' && dbInbound
         mode === 'edit' && dbInbound
           ? `/panel/api/inbounds/update/${dbInbound.id}`
           ? `/panel/api/inbounds/update/${dbInbound.id}`
@@ -615,9 +616,11 @@ export default function InboundFormModal({
 
 
   const basicTab = (
   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')}>
       <FormField name="remark" label={t('pages.inbounds.remark')}>
         <Input />
         <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.NodeClientTraffic{},
 		&model.NodeClientIp{},
 		&model.NodeClientIp{},
 		&model.ClientGlobalTraffic{},
 		&model.ClientGlobalTraffic{},
+		&model.NodePendingReset{},
 		&model.OutboundSubscription{},
 		&model.OutboundSubscription{},
 		&model.SubBalancer{},
 		&model.SubBalancer{},
 	}
 	}

+ 1 - 0
internal/database/migrate_data.go

@@ -56,6 +56,7 @@ func migrationModels() []any {
 		&model.NodeClientTraffic{},
 		&model.NodeClientTraffic{},
 		&model.NodeClientIp{},
 		&model.NodeClientIp{},
 		&model.ClientGlobalTraffic{},
 		&model.ClientGlobalTraffic{},
+		&model.NodePendingReset{},
 		&model.OutboundSubscription{},
 		&model.OutboundSubscription{},
 		&model.SubBalancer{},
 		&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)
 	ctx, cancel := context.WithTimeout(context.Background(), nodeTrafficSyncRequestTimeout)
 	defer cancel()
 	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
 	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 {
 func (r *Remote) ResetAllTraffics(ctx context.Context) error {
 	_, err := r.do(ctx, http.MethodPost, "panel/api/inbounds/resetAllTraffics", nil)
 	_, err := r.do(ctx, http.MethodPost, "panel/api/inbounds/resetAllTraffics", nil)
 	return err
 	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)
 					_, _, err = inboundSvc.AddInbound(&update)
 				case "update inbound":
 				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":
 				case "add inbound client":
 					client.Email = "invalid-new-client"
 					client.Email = "invalid-new-client"
 					update.Settings = clientsSettings(t, []model.Client{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
 		return 0, err
 	}
 	}
 	affected := 0
 	affected := 0
+	var resetNodes []int
 	err = submitTrafficWrite(func() error {
 	err = submitTrafficWrite(func() error {
 		db := database.GetDB()
 		db := database.GetDB()
 		return db.Transaction(func(tx *gorm.DB) error {
 		return db.Transaction(func(tx *gorm.DB) error {
@@ -97,12 +98,16 @@ func (s *ClientService) BulkResetTraffic(inboundSvc *InboundService, emails []st
 					return err
 					return err
 				}
 				}
 			}
 			}
-			return nil
+			var qErr error
+			resetNodes, qErr = queueNodeResets(tx, cleanEmails)
+			return qErr
 		})
 		})
 	})
 	})
 	if err != nil {
 	if err != nil {
 		return 0, err
 		return 0, err
 	}
 	}
+	inboundSvc.resetMtprotoClientQuotas(cleanEmails)
+	inboundSvc.deliverNodeResetsNow(resetNodes)
 	// After the zeroing, as in ResetTrafficByEmail: enabling a still-depleted
 	// After the zeroing, as in ResetTrafficByEmail: enabling a still-depleted
 	// client first lets the next traffic tick switch it off again.
 	// client first lets the next traffic tick switch it off again.
 	for _, e := range cleanEmails {
 	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 {
 func (s *ClientService) ResetAllClientTraffics(inboundSvc *InboundService, id int) error {
+	var resetNodes []int
+	var resetEmails []string
 	err := submitTrafficWrite(func() error {
 	err := submitTrafficWrite(func() error {
-		return s.resetAllClientTrafficsLocked(id)
+		var inner error
+		resetEmails, resetNodes, inner = s.resetAllClientTrafficsLocked(id)
+		return inner
 	})
 	})
 	if err == nil {
 	if err == nil {
-		inboundSvc.resetAllMtprotoQuotas()
+		inboundSvc.resetMtprotoClientQuotas(resetEmails)
+		inboundSvc.deliverNodeResetsNow(resetNodes)
 	}
 	}
 	return err
 	return err
 }
 }
 
 
-func (s *ClientService) resetAllClientTrafficsLocked(id int) error {
+func (s *ClientService) resetAllClientTrafficsLocked(id int) ([]string, []int, error) {
 	db := database.GetDB()
 	db := database.GetDB()
 	now := time.Now().Unix() * 1000
 	now := time.Now().Unix() * 1000
+	var resetNodes []int
+	var reset []string
 
 
 	if err := db.Transaction(func(tx *gorm.DB) error {
 	if err := db.Transaction(func(tx *gorm.DB) error {
 		// client_traffics.inbound_id is stale: it reflects the inbound the row was
 		// 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 {
 		if len(resetEmails) == 0 {
 			return nil
 			return nil
 		}
 		}
+		reset = resetEmails
 
 
 		if err := adjustGroupBaselinesForRemovedTraffic(tx, resetEmails); err != nil {
 		if err := adjustGroupBaselinesForRemovedTraffic(tx, resetEmails); err != nil {
 			return err
 			return err
@@ -176,6 +189,10 @@ func (s *ClientService) resetAllClientTrafficsLocked(id int) error {
 				return err
 				return err
 			}
 			}
 		}
 		}
+		var qErr error
+		if resetNodes, qErr = queueNodeResets(tx, resetEmails); qErr != nil {
+			return qErr
+		}
 
 
 		inboundWhereText := "id "
 		inboundWhereText := "id "
 		if id == -1 {
 		if id == -1 {
@@ -190,13 +207,14 @@ func (s *ClientService) resetAllClientTrafficsLocked(id int) error {
 
 
 		return result.Error
 		return result.Error
 	}); err != nil {
 	}); err != nil {
-		return err
+		return nil, nil, err
 	}
 	}
-	return nil
+	return reset, resetNodes, nil
 }
 }
 
 
 func (s *ClientService) ResetAllTraffics() (bool, error) {
 func (s *ClientService) ResetAllTraffics() (bool, error) {
 	var affected int64
 	var affected int64
+	var resetNodes []int
 	err := submitTrafficWrite(func() error {
 	err := submitTrafficWrite(func() error {
 		return database.GetDB().Transaction(func(tx *gorm.DB) error {
 		return database.GetDB().Transaction(func(tx *gorm.DB) error {
 			res := tx.Model(&xray.ClientTraffic{}).
 			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 {
 			if err := tx.Where("1 = 1").Delete(&model.ClientGlobalTraffic{}).Error; err != nil {
 				return err
 				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 {
 	if err != nil {
 		return false, err
 		return false, err
 	}
 	}
+	inbounds := &InboundService{}
+	inbounds.resetAllMtprotoQuotas()
+	inbounds.deliverNodeResetsNow(resetNodes)
 	return affected > 0, nil
 	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/common"
 	"github.com/mhsanaei/3x-ui/v3/internal/util/netsafe"
 	"github.com/mhsanaei/3x-ui/v3/internal/util/netsafe"
 	wgutil "github.com/mhsanaei/3x-ui/v3/internal/util/wireguard"
 	wgutil "github.com/mhsanaei/3x-ui/v3/internal/util/wireguard"
-	"github.com/mhsanaei/3x-ui/v3/internal/xray"
 
 
 	"gorm.io/gorm"
 	"gorm.io/gorm"
 	"gorm.io/gorm/clause"
 	"gorm.io/gorm/clause"
@@ -1685,6 +1684,35 @@ func (s *InboundService) SetInboundEnable(id int, enable bool) (bool, error) {
 	return needRestart, nil
 	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) {
 func (s *InboundService) UpdateInbound(inbound *model.Inbound) (*model.Inbound, bool, error) {
 	legacyShareAddr := legacyMtprotoShareAddr(inbound)
 	legacyShareAddr := legacyMtprotoShareAddr(inbound)
 	inbound.TrafficResetDay = normalizeTrafficResetDay(inbound.TrafficResetDay)
 	inbound.TrafficResetDay = normalizeTrafficResetDay(inbound.TrafficResetDay)
@@ -1707,34 +1735,6 @@ func (s *InboundService) UpdateInbound(inbound *model.Inbound) (*model.Inbound,
 	}
 	}
 	inbound.SubSortIndex = normalizeSubSortIndex(inbound.SubSortIndex)
 	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;
 	// Grandfather a row that was already stored incomplete so it stays editable;
 	// only a save that breaks a previously valid TLS block is refused.
 	// only a save that breaks a previously valid TLS block is refused.
 	if !s.FromNodeSync {
 	if !s.FromNodeSync {
@@ -1770,6 +1770,23 @@ func (s *InboundService) UpdateInbound(inbound *model.Inbound) (*model.Inbound,
 	var postCommitApply func()
 	var postCommitApply func()
 
 
 	txErr := runSerializedTx(func(tx *gorm.DB) error {
 	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)
 		conflict, cErr := checkPortConflictTx(tx, inbound, inbound.Id)
 		if cErr != nil {
 		if cErr != nil {
 			return cErr
 			return cErr
@@ -2066,16 +2083,16 @@ func (s *InboundService) buildInboundForLocalRuntime(tx *gorm.DB, inbound *model
 		return built, nil
 		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))
 	finalClients := make([]any, 0, len(clients))
@@ -2085,7 +2102,7 @@ func (s *InboundService) buildInboundForLocalRuntime(tx *gorm.DB, inbound *model
 			continue
 			continue
 		}
 		}
 		email, _ := c["email"].(string)
 		email, _ := c["email"].(string)
-		if enable, exists := enableMap[email]; exists && !enable {
+		if _, off := disabled[email]; off {
 			continue
 			continue
 		}
 		}
 		if manualEnable, ok := c["enable"].(bool); ok && !manualEnable {
 		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/database/model"
 	"github.com/mhsanaei/3x-ui/v3/internal/logger"
 	"github.com/mhsanaei/3x-ui/v3/internal/logger"
 	wgutil "github.com/mhsanaei/3x-ui/v3/internal/util/wireguard"
 	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
 // DesiredAmneziaWGInstances derives the AmneziaWG interfaces this panel
@@ -37,47 +36,38 @@ func (s *InboundService) DesiredAmneziaWGInstances() ([]amneziawg.Instance, erro
 		return nil, nil
 		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))
 	instances := make([]amneziawg.Instance, 0, len(inbounds))
 	for _, ib := range inbounds {
 	for _, ib := range inbounds {
 		inst, ok := amneziawg.InstanceFromInbound(ib)
 		inst, ok := amneziawg.InstanceFromInbound(ib)
 		if !ok {
 		if !ok {
 			continue
 			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
 // 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
 	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"
 	"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) {
 func TestUpdateInbound_RejectsHysteriaClientWithoutAuth(t *testing.T) {
 	setupConflictDB(t)
 	setupConflictDB(t)
 	seedInboundConflict(t, "in-45001-tcp", "0.0.0.0", 45001, model.VLESS,
 	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
 	var existing model.Inbound
 	if err := database.GetDB().Where("tag = ?", "in-45001-tcp").First(&existing).Error; err != nil {
 	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 := existing
 	update.Protocol = model.Hysteria
 	update.Protocol = model.Hysteria
-	update.Settings = `{"clients":[{"email":"hysteria@x","enable":true,"password":"not-hysteria-auth"}]}`
+	update.Settings = `{"clients":[]}`
 
 
 	svc := &InboundService{}
 	svc := &InboundService{}
 	if _, _, err := svc.UpdateInbound(&update); err == nil || !strings.Contains(err.Error(), "empty client ID") {
 	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 := existing
 	update.Settings = `{"clients":[{"email":"hysteria@x","enable":true,"password":"` + password + `","auth":"` + wantAuth + `"}]}`
 	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 {
 	if _, _, err := svc.UpdateInbound(&update); err != nil {
 		t.Fatalf("UpdateInbound: %v", err)
 		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/database/model"
 	"github.com/mhsanaei/3x-ui/v3/internal/logger"
 	"github.com/mhsanaei/3x-ui/v3/internal/logger"
 	"github.com/mhsanaei/3x-ui/v3/internal/mtproto"
 	"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
 // DesiredMtprotoInstances derives the mtg sidecar configs this panel should be
@@ -32,47 +31,38 @@ func (s *InboundService) DesiredMtprotoInstances() ([]mtproto.Instance, error) {
 		return nil, nil
 		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))
 	instances := make([]mtproto.Instance, 0, len(inbounds))
 	for _, ib := range inbounds {
 	for _, ib := range inbounds {
 		inst, ok := mtproto.InstanceFromInbound(ib)
 		inst, ok := mtproto.InstanceFromInbound(ib)
 		if !ok {
 		if !ok {
 			continue
 			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
 // 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) {
 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()
 	mgr := mtproto.GetManager()
-	if !mgr.HasRunning() {
+	if !mgr.HasRunning() || len(emails) == 0 {
 		return
 		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
 		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() {
 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) {
 		if !strings.Contains(update.Settings, mtprotoTestSecretD) {
 			t.Fatal("fixture must contain the re-keyed secret")
 			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 {
 		if err != nil {
 			t.Fatalf("UpdateInbound: %v", err)
 			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]
 		centralCSByEmail[centralClientStats[i].Email] = &centralClientStats[i]
 	}
 	}
 
 
+	owedResets, err := pendingNodeResetEmails(db, nodeID)
+	if err != nil {
+		return false, err
+	}
+
 	nodeBaselines := make(map[string]nodeTrafficCounter)
 	nodeBaselines := make(map[string]nodeTrafficCounter)
 	var baselineRows []model.NodeClientTraffic
 	var baselineRows []model.NodeClientTraffic
 	if err := db.Model(&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).
 			// Node-wide total, not this inbound's possibly-stale copy (#5274).
 			canon := nodeEmailTotals[cs.Email]
 			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]
 			base, seen := nodeBaselines[cs.Email]
 			var deltaUp, deltaDown int64
 			var deltaUp, deltaDown int64
@@ -986,18 +995,18 @@ func (s *InboundService) setRemoteTrafficLocked(nodeID int, snap *runtime.Traffi
 
 
 			existing := centralCSByEmail[cs.Email]
 			existing := centralCSByEmail[cs.Email]
 			if existing != nil {
 			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
 				// Only a real latch to disabled is structural; one-way merge never
 				// re-enables from the node.
 				// re-enables from the node.
-				enableChanged := !lifecycleFrozen && existing.Enable && !cs.Enable &&
+				enableChanged := !clientFrozen && existing.Enable && !cs.Enable &&
 					!nodeDisableIsStale(existing, cs, now, deltaUp, deltaDown)
 					!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 {
 				if enableChanged || metaChanged || expiryChanged {
 					structuralChange = true
 					structuralChange = true
 				}
 				}
 			}
 			}
 
 
-			renewed := !lifecycleFrozen && seen && existing != nil && nodeClientRenewed(existing, cs, canon, base)
+			renewed := !clientFrozen && seen && existing != nil && nodeClientRenewed(existing, cs, canon, base)
 			if renewed {
 			if renewed {
 				// Reject when the node's own settings still carry the old absolute:
 				// Reject when the node's own settings still carry the old absolute:
 				// lagging ClientStats after a master shorten mimic a renew (#6228).
 				// 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.ResetWeekday = cs.ResetWeekday
 				existing.ResetCount = cs.ResetCount
 				existing.ResetCount = cs.ResetCount
 				structuralChange = true
 				structuralChange = true
-			} else if lifecycleFrozen {
+			} else if clientFrozen {
 				// Push pending or just landed: only counters may move, the master
 				// Push pending or just landed: only counters may move, the master
 				// keeps expiry/enable/total/reset.
 				// keeps expiry/enable/total/reset.
 				if err := tx.Exec(
 				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
 			// A dip plus a lagging longer expiry mimics nodeClientRenewed and would
 			// undo a master shorten once the freeze lifts (#6228).
 			// 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
 				continue
 			}
 			}
 			if err := s.upsertNodeBaseline(tx, nodeID, cs.Email, canon.Up, canon.Down); err != nil {
 			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
 	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) {
 func (s *InboundService) AddTraffic(inboundTraffics []*xray.Traffic, clientTraffics []*xray.ClientTraffic) (needRestart bool, clientsDisabled bool, err error) {
 	var disabledNodeIDs []int
 	var disabledNodeIDs []int
 	var remotePlans []trafficInboundUpdatePlan
 	var remotePlans []trafficInboundUpdatePlan
+	var renewed []string
 	err = submitTrafficWrite(func() error {
 	err = submitTrafficWrite(func() error {
 		var inner 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
 		return inner
 	})
 	})
 	if err != nil {
 	if err != nil {
 		return
 		return
 	}
 	}
+	s.resetMtprotoClientQuotas(renewed)
 	// Off the serial writer: a hanging node must not stall traffic accounting.
 	// Off the serial writer: a hanging node must not stall traffic accounting.
 	needRestart = s.applyTrafficRemotePlans(remotePlans) || needRestart
 	needRestart = s.applyTrafficRemotePlans(remotePlans) || needRestart
 	if len(disabledNodeIDs) > 0 {
 	if len(disabledNodeIDs) > 0 {
@@ -44,7 +46,7 @@ func (s *InboundService) AddTraffic(inboundTraffics []*xray.Traffic, clientTraff
 	return
 	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()
 	db := database.GetDB()
 	// Commit durable traffic before best-effort lifecycle maintenance so helper
 	// Commit durable traffic before best-effort lifecycle maintenance so helper
 	// failures cannot discard usage already reported by Xray.
 	// 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)
 		return s.addClientTraffic(tx, clientTraffics)
 	}); err != nil {
 	}); err != nil {
-		return false, false, nil, nil, err
+		return false, false, nil, nil, nil, err
 	}
 	}
 
 
 	var (
 	var (
@@ -99,10 +101,10 @@ func (s *InboundService) addTrafficLocked(inboundTraffics []*xray.Traffic, clien
 	})
 	})
 	if err != nil {
 	if err != nil {
 		logger.Warning("traffic lifecycle maintenance failed after traffic commit:", err)
 		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)
 	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 {
 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 {
 	if err = clearGlobalTraffic(tx, renewedEmails...); err != nil {
 		return false, 0, err
 		return false, 0, err
 	}
 	}
+	mutationBatch.renewedEmails = append(mutationBatch.renewedEmails, renewedEmails...)
 	for _, clientToAdd := range clientsToAdd {
 	for _, clientToAdd := range clientsToAdd {
 		if clientToAdd.inbound.NodeID != nil {
 		if clientToAdd.inbound.NodeID != nil {
 			mutationBatch.addNode(*clientToAdd.inbound.NodeID)
 			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) {
 func (s *InboundService) ResetClientTraffic(id int, clientEmail string) (needRestart bool, err error) {
-	var resetInbound *model.Inbound
+	var ownNode *int
 	err = submitTrafficWrite(func() error {
 	err = submitTrafficWrite(func() error {
 		var inner error
 		var inner error
-		needRestart, resetInbound, inner = s.resetClientTrafficLocked(id, clientEmail)
+		needRestart, ownNode, inner = s.resetClientTrafficLocked(id, clientEmail)
 		return inner
 		return inner
 	})
 	})
 	if err == nil {
 	if err == nil {
 		s.resetMtprotoClientQuota(clientEmail)
 		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
 	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
 	needRestart := false
 	var reenablePlan *trafficLocalApplyPlan
 	var reenablePlan *trafficLocalApplyPlan
 	var reenableNodeID *int
 	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 {
 		if err := tx.Where("email = ?", clientEmail).Delete(&model.NodeClientTraffic{}).Error; err != nil {
 			return err
 			return err
 		}
 		}
+		if _, err := queueNodeResets(tx, []string{clientEmail}); err != nil {
+			return err
+		}
 		if err := tx.Model(model.Inbound{}).
 		if err := tx.Model(model.Inbound{}).
 			Where("id = ?", id).
 			Where("id = ?", id).
 			Update("last_traffic_reset_time", now).Error; err != nil {
 			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 {
 func (s *InboundService) ResetAllTraffics() error {
@@ -784,7 +783,6 @@ func (s *InboundService) ResetAllTraffics() error {
 	})
 	})
 	if err == nil {
 	if err == nil {
 		s.propagateResetAllTrafficsToNodes()
 		s.propagateResetAllTrafficsToNodes()
-		s.resetAllMtprotoQuotas()
 	}
 	}
 	return err
 	return err
 }
 }

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

@@ -30,6 +30,8 @@ type trafficMutationBatch struct {
 	localPlans  []trafficLocalApplyPlan
 	localPlans  []trafficLocalApplyPlan
 	remotePlans []trafficInboundUpdatePlan
 	remotePlans []trafficInboundUpdatePlan
 	nodeIDs     map[int]struct{}
 	nodeIDs     map[int]struct{}
+	// renewedEmails get their MTProto sidecar quota zeroed once the tick commits.
+	renewedEmails []string
 }
 }
 
 
 type trafficInboundUpdatePlan struct{ oldInbound, newInbound model.Inbound }
 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/database/model"
 	"github.com/mhsanaei/3x-ui/v3/internal/logger"
 	"github.com/mhsanaei/3x-ui/v3/internal/logger"
 	"github.com/mhsanaei/3x-ui/v3/internal/tuic"
 	"github.com/mhsanaei/3x-ui/v3/internal/tuic"
-	"github.com/mhsanaei/3x-ui/v3/internal/xray"
 )
 )
 
 
 func (s *InboundService) DesiredTuicInstances() ([]tuic.Instance, error) {
 func (s *InboundService) DesiredTuicInstances() ([]tuic.Instance, error) {
@@ -23,47 +22,38 @@ func (s *InboundService) DesiredTuicInstances() ([]tuic.Instance, error) {
 		return nil, nil
 		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))
 	instances := make([]tuic.Instance, 0, len(inbounds))
 	for _, ib := range inbounds {
 	for _, ib := range inbounds {
 		inst, ok := tuic.InstanceFromInbound(ib)
 		inst, ok := tuic.InstanceFromInbound(ib)
 		if !ok {
 		if !ok {
 			continue
 			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) {
 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 (
 import (
 	"fmt"
 	"fmt"
+	"net"
+	"net/http"
+	"net/url"
 	"os"
 	"os"
 	"path/filepath"
 	"path/filepath"
+	"regexp"
+	"slices"
 	"strings"
 	"strings"
 	"testing"
 	"testing"
 	"time"
 	"time"
@@ -46,9 +51,79 @@ func fakeMtgChildMain() {
 		fmt.Fprintf(f, "%d\n", os.Getpid())
 		fmt.Fprintf(f, "%d\n", os.Getpid())
 		f.Close()
 		f.Close()
 	}
 	}
+	if logPath := os.Getenv("MTG_FAKE_APILOG"); logPath != "" && len(os.Args) > 2 {
+		go serveFakeMtgAPI(os.Args[len(os.Args)-1], logPath)
+	}
 	select {}
 	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
 // 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
 // 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.
 // 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 {
 		if err := tx.Where("node_id = ?", id).Delete(&model.NodeClientTraffic{}).Error; err != nil {
 			return err
 			return err
 		}
 		}
+		if err := tx.Where("node_id = ?", id).Delete(&model.NodePendingReset{}).Error; err != nil {
+			return err
+		}
 		guids := []string{synthNodeGuid(id)}
 		guids := []string{synthNodeGuid(id)}
 		if guid != "" {
 		if guid != "" {
 			guids = append(guids, 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)
+	}
+}