fix: rollback forward mutation on tunnel switch failure

This commit is contained in:
sagit
2026-02-08 02:30:36 +00:00
parent d720b9e00c
commit 7f9c05172b
3 changed files with 355 additions and 6 deletions
@@ -283,6 +283,32 @@ func (h *Handler) listUserTunnelIDs(userID, tunnelID int64) ([]int64, error) {
return out, nil
}
func (h *Handler) listUserTunnelIDsByUser(userID int64) ([]int64, error) {
rows, err := h.repo.DB().Query(`
SELECT id
FROM user_tunnel
WHERE user_id = ?
ORDER BY id ASC
`, userID)
if err != nil {
return nil, err
}
defer rows.Close()
out := make([]int64, 0)
for rows.Next() {
var id int64
if err := rows.Scan(&id); err != nil {
return nil, err
}
out = append(out, id)
}
if err := rows.Err(); err != nil {
return nil, err
}
return out, nil
}
func (h *Handler) syncForwardServices(forward *forwardRecord, method string, allowFallbackAdd bool) error {
if h == nil || forward == nil {
return errors.New("invalid forward sync context")
@@ -342,7 +368,14 @@ func (h *Handler) controlForwardServices(forward *forwardRecord, commandType str
if err != nil {
return err
}
bases := buildForwardServiceBaseCandidates(forward.ID, forward.UserID, userTunnelID, userTunnelIDs)
allUserTunnelIDs, err := h.listUserTunnelIDsByUser(forward.UserID)
if err != nil {
return err
}
candidateTunnelIDs := make([]int64, 0, len(userTunnelIDs)+len(allUserTunnelIDs))
candidateTunnelIDs = append(candidateTunnelIDs, userTunnelIDs...)
candidateTunnelIDs = append(candidateTunnelIDs, allUserTunnelIDs...)
bases := buildForwardServiceBaseCandidates(forward.ID, forward.UserID, userTunnelID, candidateTunnelIDs)
seen := map[int64]struct{}{}
for _, fp := range ports {
if _, ok := seen[fp.NodeID]; ok {
+71 -5
View File
@@ -948,6 +948,11 @@ func (h *Handler) forwardUpdate(w http.ResponseWriter, r *http.Request) {
response.WriteJSON(w, response.Err(-2, err.Error()))
return
}
oldPorts, err := h.listForwardPorts(id)
if err != nil {
response.WriteJSON(w, response.Err(-2, err.Error()))
return
}
tunnelID := asInt64(req["tunnelId"], forward.TunnelID)
if tunnelID <= 0 {
@@ -1000,13 +1005,19 @@ func (h *Handler) forwardUpdate(w http.ResponseWriter, r *http.Request) {
response.WriteJSON(w, response.Err(-2, err.Error()))
return
}
_ = h.replaceForwardPorts(id, tunnelID, port)
if err := h.replaceForwardPorts(id, tunnelID, port); err != nil {
h.rollbackForwardMutation(forward, oldPorts)
response.WriteJSON(w, response.Err(-2, err.Error()))
return
}
updatedForward, err := h.getForwardRecord(id)
if err != nil {
h.rollbackForwardMutation(forward, oldPorts)
response.WriteJSON(w, response.Err(-2, err.Error()))
return
}
if err := h.syncForwardServices(updatedForward, "UpdateService", true); err != nil {
h.rollbackForwardMutation(forward, oldPorts)
response.WriteJSON(w, response.ErrDefault(err.Error()))
return
}
@@ -1291,6 +1302,11 @@ func (h *Handler) forwardBatchChangeTunnel(w http.ResponseWriter, r *http.Reques
fail++
continue
}
oldPorts, listPortsErr := h.listForwardPorts(id)
if listPortsErr != nil {
fail++
continue
}
var port sql.NullInt64
_ = h.repo.DB().QueryRow(`SELECT MIN(port) FROM forward_port WHERE forward_id = ?`, id).Scan(&port)
_, err := h.repo.DB().Exec(`UPDATE forward SET tunnel_id = ?, updated_time = ? WHERE id = ?`, req.TargetTunnelID, time.Now().UnixMilli(), id)
@@ -1305,13 +1321,19 @@ func (h *Handler) forwardBatchChangeTunnel(w http.ResponseWriter, r *http.Reques
if p <= 0 {
p = h.pickTunnelPort(req.TargetTunnelID)
}
_ = h.replaceForwardPorts(id, req.TargetTunnelID, p)
if err := h.replaceForwardPorts(id, req.TargetTunnelID, p); err != nil {
h.rollbackForwardMutation(forward, oldPorts)
fail++
continue
}
updatedForward, fetchErr := h.getForwardRecord(id)
if fetchErr != nil {
h.rollbackForwardMutation(forward, oldPorts)
fail++
continue
}
if err := h.syncForwardServices(updatedForward, "UpdateService", true); err != nil {
h.rollbackForwardMutation(forward, oldPorts)
fail++
continue
}
@@ -2407,14 +2429,58 @@ func (h *Handler) replaceForwardPorts(forwardID, tunnelID int64, port int) error
return err
}
defer func() { _ = tx.Rollback() }()
_, _ = tx.Exec(`DELETE FROM forward_port WHERE forward_id = ?`, forwardID)
entryNodes, _ := h.tunnelEntryNodeIDs(tunnelID)
if _, err := tx.Exec(`DELETE FROM forward_port WHERE forward_id = ?`, forwardID); err != nil {
return err
}
entryNodes, err := h.tunnelEntryNodeIDs(tunnelID)
if err != nil {
return err
}
for _, nodeID := range entryNodes {
_, _ = tx.Exec(`INSERT INTO forward_port(forward_id, node_id, port) VALUES(?, ?, ?)`, forwardID, nodeID, port)
if _, err := tx.Exec(`INSERT INTO forward_port(forward_id, node_id, port) VALUES(?, ?, ?)`, forwardID, nodeID, port); err != nil {
return err
}
}
return tx.Commit()
}
func (h *Handler) replaceForwardPortsWithRecords(forwardID int64, ports []forwardPortRecord) error {
tx, err := h.repo.DB().Begin()
if err != nil {
return err
}
defer func() { _ = tx.Rollback() }()
if _, err := tx.Exec(`DELETE FROM forward_port WHERE forward_id = ?`, forwardID); err != nil {
return err
}
for _, fp := range ports {
if _, err := tx.Exec(`INSERT INTO forward_port(forward_id, node_id, port) VALUES(?, ?, ?)`, forwardID, fp.NodeID, fp.Port); err != nil {
return err
}
}
return tx.Commit()
}
func (h *Handler) rollbackForwardMutation(oldForward *forwardRecord, oldPorts []forwardPortRecord) {
if h == nil || oldForward == nil || h.repo == nil || h.repo.DB() == nil {
return
}
_, _ = h.repo.DB().Exec(`
UPDATE forward
SET user_id = ?, user_name = ?, name = ?, tunnel_id = ?, remote_addr = ?, strategy = ?, status = ?, updated_time = ?
WHERE id = ?
`, oldForward.UserID, oldForward.UserName, oldForward.Name, oldForward.TunnelID, oldForward.RemoteAddr, oldForward.Strategy, oldForward.Status, time.Now().UnixMilli(), oldForward.ID)
if err := h.replaceForwardPortsWithRecords(oldForward.ID, oldPorts); err != nil {
return
}
_ = h.syncForwardServices(oldForward, "UpdateService", true)
}
func (h *Handler) upsertUserTunnel(req map[string]interface{}) error {
userID := asInt64(req["userId"], 0)
tunnelID := asInt64(req["tunnelId"], 0)