package main

import (
	"bytes"
	"compress/gzip"
	"context"
	"crypto/sha256"
	"encoding/base64"
	"encoding/json"
	"errors"
	"flag"
	"fmt"
	"io"
	"log"
	"math"
	"net/http"
	neturl "net/url"
	"os"
	"os/signal"
	"path/filepath"
	"regexp"
	"sort"
	"strconv"
	"strings"
	"sync"
	"sync/atomic"
	"syscall"
	"time"

	"github.com/gorilla/websocket"
)

const (
	tPing  = 1
	tPong  = 2
	tStart = 3
	tJSON  = 10

	wsWriteTimeout     = 1200 * time.Millisecond
	wsPongWriteTimeout = 120 * time.Millisecond
)

type numberItem struct {
	ID          string
	Phone       string
	Active      string
	Access      string
	DeviceID    string
	Ind         string
	Token       string
	Timer       string
	OrderLine   string
	Fingerprint string
}

type callonState struct {
	ToID        int
	Post        float64
	Salon       float64
	Person      float64
	Data        string
	TimeOT      string
	TimeDO      string
	Access      string
	Ind         string
	ID          string
	Login       string
	Fingerprint string
}

type clientState struct {
	Callon *callonState
}

type wsClient struct {
	conn   *websocket.Conn
	mu     sync.Mutex
	alive  atomic.Bool
	remote string
}

func (c *wsClient) writeBinary(b []byte) bool {
	return c.writeBinaryWithTimeout(b, wsWriteTimeout)
}

func (c *wsClient) writePong() bool {
	return c.writeBinaryWithTimeout([]byte{tPong}, wsPongWriteTimeout)
}

func (c *wsClient) writeBinaryWithTimeout(b []byte, timeout time.Duration) bool {
	c.mu.Lock()
	defer c.mu.Unlock()
	_ = c.conn.SetWriteDeadline(time.Now().Add(timeout))
	if err := c.conn.WriteMessage(websocket.BinaryMessage, b); err != nil {
		_ = c.conn.SetWriteDeadline(time.Time{})
		return false
	}
	_ = c.conn.SetWriteDeadline(time.Time{})
	return true
}

func (c *wsClient) writeJSON(v any) bool {
	body, err := json.Marshal(v)
	if err != nil {
		return false
	}
	buf := make([]byte, 1+len(body))
	buf[0] = tJSON
	copy(buf[1:], body)
	return c.writeBinary(buf)
}

func (c *wsClient) close(code int, reason string) {
	_ = c.conn.WriteControl(websocket.CloseMessage, websocket.FormatCloseMessage(code, reason), time.Now().Add(500*time.Millisecond))
	_ = c.conn.Close()
}

type semaphore struct{ ch chan struct{} }

func newSemaphore(n int) *semaphore {
	if n < 1 {
		n = 1
	}
	return &semaphore{ch: make(chan struct{}, n)}
}

func (s *semaphore) run(ctx context.Context, fn func() error) error {
	select {
	case s.ch <- struct{}{}:
	case <-ctx.Done():
		return ctx.Err()
	}
	defer func() { <-s.ch }()
	return fn()
}

func (s *semaphore) tryAcquire() bool {
	select {
	case s.ch <- struct{}{}:
		return true
	default:
		return false
	}
}

func (s *semaphore) release() {
	<-s.ch
}

type h2Resp struct {
	OK     bool
	Status int
	Data   map[string]any
}

type server struct {
	lineID int
	port   int

	upgrader websocket.Upgrader
	httpSrv  *http.Server

	h2Origin string
	httpc    *http.Client

	clientsMu sync.RWMutex
	clients   map[*wsClient]*clientState
	toIdx     map[int]map[*wsClient]struct{}

	banMu    sync.Mutex
	ban      map[string]time.Time
	banOrder []string
	banMax   int
	banTTL   time.Duration
	banSweep time.Duration

	dedupMu   sync.Mutex
	allOrders map[string]struct{}
	orderQ    []string
	orderMax  int

	myID int64

	numbersMu sync.RWMutex
	numbers   []numberItem
	warmMu    sync.Mutex
	warmNums  map[string]struct{}
	exclMu    sync.Mutex
	exclNums  map[string]struct{}

	serviceOn      atomic.Bool
	warmupOn       atomic.Bool
	ordersSeeded   atomic.Bool
	lastRyadomTime atomic.Value

	newSem   *semaphore
	phoneSem *semaphore

	refreshNumbersMs time.Duration
	newWindowMs      time.Duration

	ordersPath string
	ordersLog  bool
}

func main() {
	scriptName := filepath.Base(os.Args[0])
	var lineFlag int
	flag.IntVar(&lineFlag, "line", 0, "line id")
	flag.Parse()

	line := lineFlag
	if line == 0 {
		// Support: pm2 start ./saryagash12 -- 12
		if flag.NArg() > 0 {
			if n, err := strconv.Atoi(strings.TrimSpace(flag.Arg(0))); err == nil {
				line = n
			}
		}
	}
	if line == 0 {
		line = parseLineFromName(scriptName)
	}
	if line == 0 {
		// Optional fallback for service managers.
		line = envInt("LINE_ID", 0)
	}
	if line <= 0 {
		log.Fatal("line id is required: executable name must contain digits (e.g. astana4) or use -line=4")
	}
	log.Printf("runner=%s line_id=%d", scriptName, line)

	s := newServer(line)
	if err := s.start(); err != nil {
		log.Fatalf("start failed: %v", err)
	}
}

func parseLineFromName(name string) int {
	re := regexp.MustCompile(`\d+`)
	m := re.FindString(name)
	if m == "" {
		return 0
	}
	n, _ := strconv.Atoi(m)
	return n
}

func newServer(line int) *server {
	maxInflight := envInt("H2_MAX_INFLIGHT", 1000)
	tr := &http.Transport{
		MaxIdleConns:        maxInflight * 2,
		MaxIdleConnsPerHost: maxInflight,
		MaxConnsPerHost:     maxInflight,
		IdleConnTimeout:     60 * time.Second,
		ForceAttemptHTTP2:   true,
	}
	return &server{
		lineID: line,
		port:   4000 + line,
		upgrader: websocket.Upgrader{
			ReadBufferSize:  4096,
			WriteBufferSize: 4096,
			CheckOrigin:     func(r *http.Request) bool { return true },
		},
		h2Origin:         envStr("H2_ORIGIN", "https://icl-gw-cf.euce1.indriverapp.com"),
		httpc:            &http.Client{Transport: tr, Timeout: 6 * time.Second},
		clients:          map[*wsClient]*clientState{},
		toIdx:            map[int]map[*wsClient]struct{}{},
		ban:              map[string]time.Time{},
		banOrder:         []string{},
		banMax:           max(100, envInt("BAN_MAX", 5000)),
		banTTL:           time.Duration(max(60_000, envInt("BAN_TTL_MS", 7*24*60*60*1000))) * time.Millisecond,
		banSweep:         time.Duration(max(10_000, envInt("BAN_SWEEP_MS", 60_000))) * time.Millisecond,
		allOrders:        map[string]struct{}{},
		orderQ:           []string{},
		orderMax:         200,
		warmNums:         map[string]struct{}{},
		exclNums:         map[string]struct{}{},
		newSem:           newSemaphore(envInt("NEW_CONCURRENCY", 500)),
		phoneSem:         newSemaphore(envInt("PHONE_CONCURRENCY", 100)),
		refreshNumbersMs: time.Duration(envInt("REFRESH_NUMBERS_MS", 30_000)) * time.Millisecond,
		newWindowMs:      time.Duration(envInt("NEW_WINDOW_MS", 5_200)) * time.Millisecond,
		ordersPath:       envStr("H2_ORDERS_PATH", "/api/v2/orders"),
		ordersLog:        envBool("LOG_ORDERS", true),
	}
}

func (s *server) start() error {
	mux := http.NewServeMux()
	mux.HandleFunc("/", s.handleWS)
	s.httpSrv = &http.Server{
		Addr:              fmt.Sprintf(":%d", s.port),
		Handler:           mux,
		ReadHeaderTimeout: 5 * time.Second,
	}

	ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
	defer stop()

	go s.runMaintenance(ctx)
	go s.runScheduler(ctx)

	go func() {
		<-ctx.Done()
		_ = s.httpSrv.Shutdown(context.Background())
	}()

	log.Printf("[WS %d] listening on %d", s.lineID, s.port)
	err := s.httpSrv.ListenAndServe()
	if err != nil && !errors.Is(err, http.ErrServerClosed) {
		return err
	}
	return nil
}

func (s *server) runMaintenance(ctx context.Context) {
	pingTick := time.NewTicker(15 * time.Second)
	banTick := time.NewTicker(s.banSweep)
	refreshTick := time.NewTicker(s.refreshNumbersMs)
	defer pingTick.Stop()
	defer banTick.Stop()
	defer refreshTick.Stop()

	for {
		select {
		case <-ctx.Done():
			s.closeAllClients()
			return
		case <-pingTick.C:
			s.heartbeat()
		case <-banTick.C:
			s.pruneBan()
		case <-refreshTick.C:
			if s.serviceOn.Load() || s.hasActiveCallonAndMyID() {
				_ = s.refreshNumbers(context.Background())
				s.maybeStartStopService()
			}
		}
	}
}

func (s *server) runScheduler(ctx context.Context) {
	var idx int
	var lastIdleLog time.Time
	for {
		if ctx.Err() != nil {
			return
		}
		if !s.serviceOn.Load() {
			time.Sleep(100 * time.Millisecond)
			continue
		}

		myid := atomic.LoadInt64(&s.myID)
		nums := s.getNumbers()
		activeNums := len(nums)
		selectedKey := ""

		if len(nums) > 0 && myid > 0 {
			item := nums[idx%len(nums)]
			idx++
			selectedKey = s.numberKey(item)
			// Avoid unbounded goroutine growth when the newSem queue is saturated.
			if s.newSem.tryAcquire() {
				go func(it numberItem) {
					defer s.newSem.release()
					s.fireOrders(context.Background(), it)
				}(item)
			}
		}

		delay := 50 * time.Millisecond
		if activeNums > 0 {
			d := s.newWindowMs / time.Duration(activeNums)
			if d < time.Millisecond {
				d = time.Millisecond
			}
			delay = d
		}
		if selectedKey == "" {
			now := time.Now()
			if lastIdleLog.IsZero() || now.Sub(lastIdleLog) >= time.Second {
				s.logOrdersf("[ORDERS line=%d] idle numbers=%d delay_ms=%d myid=%d", s.lineID, activeNums, delay.Milliseconds(), myid)
				lastIdleLog = now
			}
		}
		time.Sleep(delay)
	}
}

func (s *server) handleWS(w http.ResponseWriter, r *http.Request) {
	conn, err := s.upgrader.Upgrade(w, r, nil)
	if err != nil {
		s.logOrdersf("[WS line=%d] upgrade_err remote=%s err=%v", s.lineID, strings.TrimSpace(r.RemoteAddr), err)
		return
	}
	remote := strings.TrimSpace(r.RemoteAddr)
	if remote == "" {
		remote = "unknown"
	}
	c := &wsClient{conn: conn, remote: remote}
	c.alive.Store(true)
	conn.SetReadLimit(256 * 1024)
	conn.SetPongHandler(func(string) error {
		c.alive.Store(true)
		return nil
	})

	s.clientsMu.Lock()
	s.clients[c] = &clientState{}
	totalClients := len(s.clients)
	s.clientsMu.Unlock()
	s.logOrdersf("[WS line=%d] connect client=%p remote=%s ua=%q total=%d", s.lineID, c, c.remote, strings.TrimSpace(r.UserAgent()), totalClients)

	_ = c.writeBinary([]byte{tStart})
	s.maybeStartStopService()

	defer func() {
		removedToID := s.cleanupClient(c)
		if removedToID > 0 {
			s.notifyToIDCount(removedToID)
		}
		s.logOrdersf("[WS line=%d] disconnect client=%p remote=%s removed_toid=%d", s.lineID, c, c.remote, removedToID)
		s.maybeStartStopService()
	}()

	for {
		mt, data, err := conn.ReadMessage()
		if err != nil {
			s.logOrdersf("[WS line=%d] read_err client=%p remote=%s err=%v", s.lineID, c, c.remote, err)
			return
		}
		if mt != websocket.BinaryMessage || len(data) < 1 {
			continue
		}
		typ := data[0]
		if typ == tPing {
			_ = c.writePong()
			continue
		}
		if typ != tJSON || len(data) < 2 {
			continue
		}
		var obj map[string]any
		if err := json.Unmarshal(data[1:], &obj); err != nil {
			s.logOrdersf("[WS line=%d] json_unmarshal_err client=%p remote=%s err=%v", s.lineID, c, c.remote, err)
			continue
		}
		s.logOrdersf("[WS line=%d] recv client=%p remote=%s type=%q message=%q", s.lineID, c, c.remote, str(obj["type"]), str(obj["message"]))
		s.handleMessage(c, obj)
	}
}

func (s *server) handleMessage(c *wsClient, obj map[string]any) {
	if str(obj["message"]) == "callon" {
		s.handleCallon(c, obj)
		return
	}
	switch str(obj["type"]) {
	case "ban":
		if name := str(obj["name"]); name != "" {
			s.banAdd(name)
		}
	case "activeryadom":
		lat, lng := str(obj["lat"]), str(obj["lng"])
		if lat == "" || lng == "" {
			_ = c.writeJSON(map[string]any{"type": "nolatlng"})
			return
		}
		go s.activeryadom(c, str(obj["access"]), str(obj["ind"]), lat, lng)
	case "stopryadom":
		go s.stopryadom(c, str(obj["access"]), str(obj["ind"]))
	case "getryadom":
		lat, lng := str(obj["lat"]), str(obj["lng"])
		if lat == "" || lng == "" {
			_ = c.writeJSON(map[string]any{"type": "nolatlng"})
			return
		}
		go s.getryadom(c, str(obj["access"]), str(obj["ind"]), lat, lng, str(obj["id"]), str(obj["fingerprint"]))
	case "access":
		ack := false
		s.clientsMu.Lock()
		st := s.clients[c]
		if st != nil && st.Callon != nil {
			st.Callon.Access = str(obj["access"])
			st.Callon.Ind = str(obj["ind"])
			ack = true
		}
		s.clientsMu.Unlock()
		if ack {
			_ = c.writeJSON(map[string]any{"type": "access"})
			s.logOrdersf("[WS line=%d] access_ack client=%p remote=%s", s.lineID, c, c.remote)
		} else {
			s.logOrdersf("[WS line=%d] access_reject client=%p remote=%s reason=no_callon", s.lineID, c, c.remote)
			s.kickAndCleanup(c, 4000, "sent")
			return
		}
	default:
		s.logOrdersf("[WS line=%d] unknown_type client=%p remote=%s type=%q", s.lineID, c, c.remote, str(obj["type"]))
	}
}

func (s *server) handleCallon(c *wsClient, obj map[string]any) {
	toid := idInt(obj["toid"])
	prevToID := 0
	prevAccess := ""
	prevInd := ""
	callonLineID, hasCallonLineID := callonLineIDFromPayload(obj)

	s.clientsMu.RLock()
	st := s.clients[c]
	if st != nil && st.Callon != nil {
		prevToID = st.Callon.ToID
		prevAccess = st.Callon.Access
		prevInd = st.Callon.Ind
	}
	s.clientsMu.RUnlock()

	callonAccess := str(obj["access"])
	callonInd := str(obj["ind"])
	if strings.TrimSpace(callonAccess) != "" && strings.TrimSpace(callonInd) != "" {
		_ = s.stopryadomRequest(callonAccess, callonInd)
	} else if prevToID > 0 {
		_ = s.stopryadomRequest(prevAccess, prevInd)
	}

	s.clientsMu.Lock()
	st = s.clients[c]
	if st != nil && st.Callon != nil {
		prevToID = st.Callon.ToID
		s.indexRemoveLocked(prevToID, c)
	}
	s.clientsMu.Unlock()
	s.logOrdersf("[WS line=%d] callon_req client=%p remote=%s toid=%d prev_toid=%d login=%q myid=%d callon_line=%d has_callon_line=%t", s.lineID, c, c.remote, toid, prevToID, str(obj["login"]), int64(int(num(obj["myid"]))), callonLineID, hasCallonLineID)
	if prevToID > 0 && prevToID != toid {
		s.notifyToIDCount(prevToID)
	}
	if hasCallonLineID && callonLineID != s.lineID {
		_ = c.writeJSON(map[string]any{
			"type":        "line_mismatch",
			"line":        callonLineID,
			"server_line": s.lineID,
		})
		s.logOrdersf("[WS line=%d] callon_reject client=%p remote=%s reason=line_mismatch callon_line=%d", s.lineID, c, c.remote, callonLineID)
		s.kickAndCleanup(c, 4000, "line_mismatch")
		return
	}

	if toid <= 0 {
		_ = c.writeJSON(map[string]any{"type": "notoid"})
		s.logOrdersf("[WS line=%d] callon_reject client=%p remote=%s reason=notoid", s.lineID, c, c.remote)
		s.kickAndCleanup(c, 4000, "sent")
		return
	}

	actives := s.fetchactive(str(obj["token"]))
	if actives != "ok" {
		_ = c.writeJSON(map[string]any{"type": "activate"})
		s.logOrdersf("[WS line=%d] callon_reject client=%p remote=%s reason=activate status=%q", s.lineID, c, c.remote, actives)
		s.kickAndCleanup(c, 4000, "sent")
		return
	}

	if atomic.LoadInt64(&s.myID) == 0 {
		atomic.CompareAndSwapInt64(&s.myID, 0, int64(int(num(obj["myid"]))))
	}

	cs := &callonState{
		ToID:        toid,
		Post:        num(obj["post"]),
		Salon:       num(obj["salon"]),
		Person:      num(obj["person"]),
		Data:        str(obj["data"]),
		TimeOT:      str(obj["timeot"]),
		TimeDO:      str(obj["timedo"]),
		Access:      str(obj["access"]),
		Ind:         str(obj["ind"]),
		ID:          str(obj["id"]),
		Login:       str(obj["login"]),
		Fingerprint: str(obj["fingerprint"]),
	}

	s.clientsMu.Lock()
	st = s.clients[c]
	if st != nil {
		st.Callon = cs
		s.indexAddLocked(cs.ToID, c)
	}
	s.clientsMu.Unlock()
	if cs.ToID > 0 {
		s.notifyToIDCount(cs.ToID)
	}
	s.logOrdersf("[WS line=%d] callon_ok client=%p remote=%s toid=%d login=%q toid_count=%d", s.lineID, c, c.remote, cs.ToID, cs.Login, s.toIDCount(cs.ToID))

	s.maybeStartStopService()
	_ = c.writeJSON(map[string]any{"type": "callon_ok", "ts": time.Now().UnixMilli()})
	s.kickWarmupAsync()
}

func callonLineIDFromPayload(obj map[string]any) (int, bool) {
	for _, k := range []string{"line", "line_id", "lineid"} {
		v, ok := obj[k]
		if !ok {
			continue
		}
		return idInt(v), true
	}
	return 0, false
}

func (s *server) activeryadom(c *wsClient, access, ind, lat, lng string) {
	// lat = fixedCoord7(lat)
	// lng = fixedCoord7(lng)
	path := fmt.Sprintf("/api/v2/dispatch/activate?latitude=%s&longitude=%s&market=google", queryEscape(lat), queryEscape(lng))
	res := s.h2Measured(context.Background(), "PUT", path, nil, map[string]string{
		"x-latitude":      lat,
		"x-longitude":     lng,
		"accept-language": "ru_RU",
		"x-os-type":       "android",
		"x-app-flavor":    "indriver",
		"authorization":   "Bearer " + access,
		"mode":            "driver",
		"x-app":           "android 5.155.0",
		"x-fingerprint":   fmt.Sprintf(`{"shield_session_id":"%s"}`, ind),
	}, 5*time.Second)
	dataRaw, dataMarshalErr := json.Marshal(res.Data)
	if dataMarshalErr == nil {
		s.logOrdersf("[RYADOM line=%d] activeryadom status=%d ok=%t lat=%s lng=%s data=%s", s.lineID, res.Status, res.OK, lat, lng, string(dataRaw))
	} else {
		s.logOrdersf("[RYADOM line=%d] activeryadom status=%d ok=%t lat=%s lng=%s marshal_err=%v data=%v", s.lineID, res.Status, res.OK, lat, lng, dataMarshalErr, res.Data)
	}
	logWs := func(respType string, sent bool) {
		if dataMarshalErr == nil {
			s.logOrdersf("[RYADOM line=%d] activeryadom ws_type=%s status=%d sent=%t data=%s", s.lineID, respType, res.Status, sent, string(dataRaw))
		} else {
			s.logOrdersf("[RYADOM line=%d] activeryadom ws_type=%s status=%d sent=%t marshal_err=%v data=%v", s.lineID, respType, res.Status, sent, dataMarshalErr, res.Data)
		}
	}
	if res.Status == http.StatusUnauthorized || str(deepGet(res.Data, "meta", "code")) == "401" {
		sent := c.writeJSON(map[string]any{"type": "401", "data": res.Data})
		logWs("401", sent)
		return
	}
	sent := c.writeJSON(map[string]any{"type": "activeryadom", "data": res.Data})
	logWs("activeryadom", sent)
}

func (s *server) stopryadom(c *wsClient, access, ind string) {
	res := s.stopryadomRequest(access, ind)
	_ = c.writeJSON(map[string]any{"type": "stopryadom", "data": res.Data})
}

func (s *server) stopryadomRequest(access, ind string) h2Resp {
	if strings.TrimSpace(access) == "" || strings.TrimSpace(ind) == "" {
		return h2Resp{OK: false, Status: 0, Data: map[string]any{}}
	}
	return s.h2Measured(context.Background(), "DELETE", "/api/v2/dispatch/deactivate", nil, map[string]string{
		"accept-language": "ru_RU",
		"x-os-type":       "android",
		"x-app-flavor":    "indriver",
		"authorization":   "Bearer " + access,
		"mode":            "driver",
		"x-app":           "android 5.155.0",
		"x-fingerprint":   fmt.Sprintf(`{"shield_session_id":"%s"}`, ind),
	}, 5*time.Second)
}

func (s *server) getryadom(c *wsClient, access, ind, lat, lng, id, fingerprint string) {
	// lat = fixedCoord7(lat)
	// lng = fixedCoord7(lng)
	//log.Printf("[GETRYADOM line=%d] lat=%s lng=%s id=%s", s.lineID, lat, lng, id)
	path := fmt.Sprintf("/api/v2/dispatch/state?latitude=%s&longitude=%s", queryEscape(lat), queryEscape(lng))
	res := s.h2Measured(context.Background(), "GET", path, nil, map[string]string{
		"x-latitude":      lat,
		"x-longitude":     lng,
		"accept-language": "ru_RU",
		"x-os-type":       "android",
		"x-app-flavor":    "indriver",
		"authorization":   "Bearer " + access,
		"mode":            "driver",
		"x-app":           "android 5.155.0",
		"x-fingerprint":   fmt.Sprintf(`{"shield_session_id":"%s"}`, ind),
	}, 5*time.Second)
	if raw, err := json.Marshal(res.Data); err == nil {
		s.logOrdersf("[RYADOM line=%d] getryadom status=%d ok=%t req_id=%s lat=%s lng=%s data=%s", s.lineID, res.Status, res.OK, id, lat, lng, string(raw))
	} else {
		s.logOrdersf("[RYADOM line=%d] getryadom status=%d ok=%t req_id=%s lat=%s lng=%s marshal_err=%v data=%v", s.lineID, res.Status, res.OK, id, lat, lng, err, res.Data)
	}
	if res.Status == http.StatusUnauthorized || str(deepGet(res.Data, "meta", "code")) == "401" {
		_ = c.writeJSON(map[string]any{"type": "401", "data": res.Data})
		return
	}

	userOrder, _ := asMap(res.Data["user_order"])
	if len(userOrder) == 0 {
		_ = c.writeJSON(map[string]any{"type": "getryadom"})
		return
	}
	order, _ := asMap(userOrder["order"])
	data, _ := asMap(order["data"])
	to, _ := asMap(data["to"])
	city, _ := asMap(to["city"])
	toCityID := idInt(city["id"])
	createdAt := str(order["created_at"])
	orderID := str(order["id"])
	if toCityID == 0 || createdAt == "" || orderID == "" {
		_ = c.writeJSON(map[string]any{"type": "getryadom"})
		return
	}

	uniq := orderUniq(userOrder)
	if !s.addOrderDedup(uniq) {
		_ = s.skip(access, ind, orderID)
		_ = c.writeJSON(map[string]any{"type": "getryadom"})
		return
	}

	// Time-order check disabled temporarily.
	// if prev, ok := s.lastRyadomTime.Load().(string); ok && prev != "" {
	// 	tPrev, errPrev := time.Parse(time.RFC3339, prev)
	// 	tCur, errCur := time.Parse(time.RFC3339, createdAt)
	// 	if errPrev == nil && errCur == nil && tCur.Before(tPrev) {
	// 		_ = c.writeJSON(map[string]any{"type": "getryadom"})
	// 		return
	// 	}
	// }
	// s.lastRyadomTime.Store(createdAt)

	payloadType := "collection2"
	if s.isOrderUserBanned(userOrder) {
		payloadType = "ban"
	}
	payload := map[string]any{
		"type":        payloadType,
		"data":        userOrder,
		"phone":       "",
		"access":      access,
		"ind":         ind,
		"fingerprint": fingerprint,
	}
	filters := buildFiltersFromOrder(userOrder)
	sent := s.sendToToIDWithFiltersAndKick(toCityID, filters, payload, c)
	if sent == 0 && payloadType == "collection2" {
		_ = s.skip(access, ind, orderID)
		_ = c.writeJSON(map[string]any{"type": "getryadom"})
	}
}

func (s *server) fireOrders(ctx context.Context, item numberItem) {
	myid := atomic.LoadInt64(&s.myID)
	if myid == 0 {
		return
	}
	numKey := s.numberKey(item)
	res := s.h2Measured(ctx, "POST", s.ordersPath, map[string]any{"from": map[string]any{"city_id": myid}}, s.baseOrderHeaders(item), 5*time.Second)
	collection, _ := asSlice(res.Data["collection"])
	if len(collection) == 0 {
		// No orders in this tick: do not seed, just wait for the next cycle.
		s.logOrdersf("[ORDERS line=%d] empty_collection number=%s", s.lineID, numKey)
		return
	}
	first, _ := asMap(collection[0])
	cursorID := str(first["cursor_id"])
	if cursorID == "" {
		excluded := s.excludeNum(item)
		s.logOrdersf("[ORDERS line=%d] no_cursor_id number=%s collection=%d excluded=%t", s.lineID, numKey, len(collection), excluded)
		return
	}

	// First non-empty /orders response after service start:
	// prefill dedup with the full collection once, then process only unique/new.
	if s.ordersSeeded.CompareAndSwap(false, true) {
		added := s.seedOrderDedup(numKey, collection, 0)
		s.logOrdersf("[ORDERS line=%d] dedup_seed number=%s added=%d total=%d", s.lineID, numKey, added, len(collection))
		return
	}

	s.processOrdersResponse(numKey, collection)
}

func (s *server) seedOrderDedup(numKey string, coll []any, limit int) int {
	n := len(coll)
	if limit > 0 && n > limit {
		n = limit
	}
	added := 0
	for i := 0; i < n; i++ {
		item, ok := asMap(coll[i])
		if !ok {
			continue
		}
		uniq := orderUniq(item)
		if s.addOrderDedup(uniq) {
			added++
			//s.logOrdersf("[ORDERS line=%d] number=%s uniq_name=%s", s.lineID, numKey, orderUniqName(item))
		}
	}
	return added
}

func (s *server) processOrdersResponse(numKey string, coll []any) (int, int, bool) {
	processed := 0
	dedupAdded := 0
	brokeOnDup := false
	for _, it := range coll {
		item, ok := asMap(it)
		if !ok {
			continue
		}
		processed++
		order, _ := asMap(item["order"])
		data, _ := asMap(order["data"])
		to, _ := asMap(data["to"])
		city, _ := asMap(to["city"])
		toCityID := idInt(city["id"])
		createdAt := str(order["created_at"])
		if toCityID == 0 || createdAt == "" {
			continue
		}
		uniq := orderUniq(item)
		if !s.addOrderDedup(uniq) {
			// Collection is ordered with newest items first.
			// Once we hit an already-seen order, the tail is typically old too.
			//s.logOrdersf("[ORDERS line=%d] number=%s uniq_name=%s", s.lineID, numKey, orderUniqName(item))
			brokeOnDup = true
			break
		}
		dedupAdded++
		//s.logOrdersf("[ORDERS line=%d] number=%s uniq_name=%s", s.lineID, numKey, orderUniqName(item))

		payloadType := "collection"
		if s.isOrderUserBanned(item) {
			payloadType = "ban"
		}
		payload := map[string]any{
			"type":  payloadType,
			"data":  item,
			"phone": "",
		}
		filters := buildFiltersFromOrder(item)
		_ = s.sendToToIDWithFiltersAndKick(toCityID, filters, payload, nil)
	}
	return processed, dedupAdded, brokeOnDup
}

func (s *server) sendToToIDWithFiltersAndKick(toid int, filters orderFilters, payload map[string]any, source *wsClient) int {
	type targetClient struct {
		c  *wsClient
		cs callonState
	}

	s.clientsMu.RLock()
	target := s.toIdx[toid]
	if len(target) == 0 {
		s.clientsMu.RUnlock()
		s.logOrdersf("[ROUTE line=%d] toid=%d type=%q targets=0", s.lineID, toid, str(payload["type"]))
		return 0
	}
	list := make([]targetClient, 0, len(target))
	for c := range target {
		st := s.clients[c]
		if st == nil || st.Callon == nil {
			continue
		}
		// Copy callon state while lock is held to avoid repeated lock/unlock
		// and to work with a stable per-client snapshot in this send cycle.
		list = append(list, targetClient{c: c, cs: *st.Callon})
	}
	s.clientsMu.RUnlock()
	if len(list) == 0 {
		s.logOrdersf("[ROUTE line=%d] toid=%d type=%q targets=0 snapshot_empty=1", s.lineID, toid, str(payload["type"]))
		return 0
	}

	sent := 0
	isRyadom := str(payload["type"]) == "collection2"
	payloadType := str(payload["type"])
	var phoneCache string
	sourceGetRyadomSent := false
	filterMissCount := 0
	banCount := 0
	phoneErrCount := 0
	noPhoneCount := 0

	for _, item := range list {
		c := item.c
		cs := item.cs

		if str(payload["type"]) == "ban" {
			_ = c.writeJSON(map[string]any{"type": "collection", "data": payload["data"], "phone": "ban"})
			banCount++
			continue
		}

		filterMiss := (filters.PostGte != nil && *filters.PostGte < cs.Post) ||
			(filters.SalonEq != nil && *filters.SalonEq < cs.Salon) ||
			(filters.PersonEq != nil && *filters.PersonEq < cs.Person) ||
			(cs.Data != "" && filters.DateEq != "" && cs.Data != filters.DateEq) ||
			(filters.TimeAt != "" && !clientCoversTime(cs.TimeOT, cs.TimeDO, filters.TimeAt))
		if filterMiss {
			filterMissCount++
			_ = c.writeJSON(payload)
			if isRyadom && source != nil && c == source {
				_ = c.writeJSON(map[string]any{"type": "getryadom"})
				sourceGetRyadomSent = true
			}
			continue
		}

		if phoneCache == "" {
			orderData, _ := asMap(payload["data"])
			order, _ := asMap(orderData["order"])
			orderID := str(order["id"])
			var code string
			if isRyadom {
				code = s.getphoneRyadom(str(payload["access"]), str(payload["ind"]), orderID, chooseNonEmpty(str(payload["fingerprint"]), cs.Fingerprint))
			} else {
				code = s.getphone(cs.Access, cs.Ind, orderID, cs.Fingerprint)
			}
			//phoneCache = normalizePhoneCache(code)
			phoneCache = code
			if code == "402" || code == "401" || code == "500" || code == "445" {
				_ = c.writeJSON(map[string]any{"type": code, "data": payload["data"], "phone": ""})
				sent++
				phoneErrCount++
				s.kickAndCleanup(c, 4000, "sent")
				phoneCache = ""
				continue
			}
			if code == "404" {
				// No phone available: still send collection payload so client is not left waiting silently.
				noPhoneCount++
				_ = c.writeJSON(payload)
				if isRyadom && source != nil && c == source {
					_ = c.writeJSON(map[string]any{"type": "getryadom"})
					sourceGetRyadomSent = true
				}
				continue
			}
		} else if phoneCache == "404" {
			// Reuse cached 404 behavior for the rest of matched clients in this cycle.
			noPhoneCount++
			_ = c.writeJSON(payload)
			if isRyadom && source != nil && c == source {
				_ = c.writeJSON(map[string]any{"type": "getryadom"})
				sourceGetRyadomSent = true
			}
			continue
		}

		callType := "call"
		if isRyadom {
			callType = "call2"
		}
		_ = c.writeJSON(map[string]any{"type": callType, "data": payload["data"], "phone": phoneCache})
		sent++
		s.kickAndCleanup(c, 4000, "sent")
	}

	if isRyadom && sent == 0 && source != nil && !sourceGetRyadomSent {
		_ = source.writeJSON(map[string]any{"type": "getryadom"})
	}
	s.logOrdersf("[ROUTE line=%d] toid=%d type=%q targets=%d sent=%d filter_miss=%d ban=%d phone_err=%d no_phone=%d ryadom=%t", s.lineID, toid, payloadType, len(list), sent, filterMissCount, banCount, phoneErrCount, noPhoneCount, isRyadom)
	return sent
}

func (s *server) refreshNumbers(ctx context.Context) error {
	url := fmt.Sprintf("http://63.250.60.80/api/server.php?zapros=number&line_id=%d", s.lineID)
	const timeout = 8 * time.Second
	s.logOrdersf("[ORDERS line=%d] refresh_request url=%s", s.lineID, url)
	j, err := s.safeFetchJSON(ctx, url, timeout)
	if err != nil {
		// Retry once on timeout/slow network
		if ctx.Err() == nil {
			time.Sleep(500 * time.Millisecond)
			j, err = s.safeFetchJSON(ctx, url, timeout)
		}
		if err != nil {
			s.logOrdersf("[ORDERS line=%d] refresh_err err=%v", s.lineID, err)
			return err
		}
	}
	if raw, marshalErr := json.Marshal(j); marshalErr == nil {
		s.logOrdersf("[ORDERS line=%d] refresh_response body=%s", s.lineID, string(raw))
	} else {
		s.logOrdersf("[ORDERS line=%d] refresh_response_marshal_err err=%v body=%v", s.lineID, marshalErr, j)
	}
	if str(j["status"]) != "success" {
		s.logOrdersf("[ORDERS line=%d] refresh_bad_status status=%q", s.lineID, str(j["status"]))
		return fmt.Errorf("api status=%s", str(j["status"]))
	}
	arr, _ := asSlice(j["number"])
	apiTotal := len(arr)
	excludedSkipped := 0
	out := make([]numberItem, 0, len(arr))
	for _, it := range arr {
		m, ok := asMap(it)
		if !ok {
			continue
		}
		n := numberItem{
			ID:          str(m["id"]),
			Phone:       str(m["phone"]),
			Active:      str(m["active"]),
			Access:      str(m["access_token"]),
			DeviceID:    str(m["device_id"]),
			Ind:         str(m["ind"]),
			Token:       str(m["token"]),
			Timer:       str(m["timer"]),
			OrderLine:   str(m["order_line_id"]),
			Fingerprint: str(m["fingerprint"]),
		}
		if n.Access != "" && n.Ind != "" {
			if s.isNumExcludedKey(s.numberKey(n)) {
				excludedSkipped++
				continue
			}
			out = append(out, n)
		}
	}
	s.numbersMu.Lock()
	s.numbers = out
	s.numbersMu.Unlock()
	s.logOrdersf("[ORDERS line=%d] refresh api_total=%d active=%d excluded_skipped=%d", s.lineID, apiTotal, len(out), excludedSkipped)
	return nil
}

func (s *server) fetchactive(token string) string {
	url := fmt.Sprintf("http://63.250.60.80/api/server.php?zapros=active&line_id=%s", queryEscape(token))
	for i := 0; i < 2; i++ {
		j, err := s.safeFetchAny(context.Background(), url, 5*time.Second)
		if err == nil {
			if arr, ok := j.([]any); ok && len(arr) > 0 {
				if m, ok := asMap(arr[0]); ok {
					v := str(m["lines2"])
					if v != "" {
						return v
					}
				}
			}
		}
		time.Sleep(150 * time.Millisecond)
	}
	return "bad"
}

func (s *server) h2MeasuredH(ctx context.Context, method, path string, body any, headers http.Header, timeout time.Duration) h2Resp {
	start := time.Now()

	// ⚠️ нужна новая версия h2Request, которая умеет http.Header (см. ниже)
	data, status, err := s.h2RequestH(ctx, method, path, body, headers, timeout)

	s.logHTTPLatency(method, path, status, err, time.Since(start))

	if err != nil {
		return h2Resp{OK: false, Status: status, Data: map[string]any{}}
	}
	return h2Resp{OK: status >= 200 && status < 300, Status: status, Data: data}
}

func (s *server) h2RequestH(ctx context.Context, method, path string, body any, headers http.Header, timeout time.Duration) (map[string]any, int, error) {
	u := strings.TrimRight(s.h2Origin, "/") + path

	var rd io.Reader
	if body != nil {
		b, err := json.Marshal(body)
		if err != nil {
			return nil, 0, err
		}
		rd = bytes.NewReader(b)
	}

	rctx, cancel := context.WithTimeout(ctx, timeout)
	defer cancel()

	req, err := http.NewRequestWithContext(rctx, method, u, rd)
	if err != nil {
		return nil, 0, err
	}

	// ✅ переносим все заголовки с сохранением мульти-значений
	if headers != nil {
		req.Header = headers.Clone()
	}

	if body != nil {
		// чтобы не затереть заранее заданный content-type
		if req.Header.Get("content-type") == "" {
			req.Header.Set("content-type", "application/json")
		}
	}
	//fmt.Println("x-app-signature values:", req.Header.Values("x-app-signature"))

	resp, err := s.httpc.Do(req)
	if err != nil {
		return nil, 0, err
	}
	defer resp.Body.Close()

	br := io.Reader(resp.Body)
	if strings.EqualFold(resp.Header.Get("content-encoding"), "gzip") {
		gz, gzErr := gzip.NewReader(resp.Body)
		if gzErr == nil {
			defer gz.Close()
			br = gz
		}
	}

	raw, _ := io.ReadAll(br)
	raw = bytes.TrimPrefix(raw, []byte{0xEF, 0xBB, 0xBF})
	raw = bytes.TrimSpace(raw)

	if len(raw) == 0 {
		return map[string]any{}, resp.StatusCode, nil
	}

	var out map[string]any
	if err := json.Unmarshal(raw, &out); err != nil {
		return map[string]any{}, resp.StatusCode, nil
	}

	return out, resp.StatusCode, nil
}

func (s *server) getphone3(access, ind, id, fingerprint string) string {
	var out string

	_ = s.phoneSem.run(context.Background(), func() error {
		canonical := "X-App=android 5.155.0&order_id=" + id + "&shield_session_id=" + ind
		sig := xAppSignatureFallback(requestHashFromCanonical(canonical), -9998)
		canonical2 := "X-App=android 5.155.0&shield_session_id=" + ind
		sig2 := xAppSignatureFallback(requestHashFromCanonical(canonical2), -9998)

		h := make(http.Header)
		h.Set("x-latitude", "51.1389885")
		h.Set("x-longitude", "71.3743228")
		h.Set("accept-language", "ru_RU")
		h.Set("x-os-type", "android")

		h.Add("x-app-signature", sig)
		h.Add("x-app-signature", sig2)

		h.Set("x-device-fingerprint", fingerprint)
		h.Set("x-app-flavor", "indriver")
		h.Set("accept-encoding", "gzip")
		h.Set("authorization", "Bearer "+access)
		h.Set("mode", "driver")
		h.Set("x-app", "android 5.155.0")
		h.Set("x-fingerprint", fmt.Sprintf(`{"shield_session_id":"%s"}`, ind))

		res := s.h2MeasuredH(
			context.Background(),
			"GET",
			"/api/v1/order/"+id+"/contact",
			nil,
			h,
			5*time.Second,
		)

		if ph := str(res.Data["phone"]); ph != "" {
			out = ph
			return nil
		}

		switch str(deepGet(res.Data, "meta", "code")) {
		case "402", "401", "404", "445":
			out = str(deepGet(res.Data, "meta", "code"))
		default:
			out = "500"
		}
		return nil
	})

	return out
}

func (s *server) getphone(access, ind, id, fingerprint string) string {
	var out string
	_ = s.phoneSem.run(context.Background(), func() error {
		canonical := "X-App=android 5.155.0&order_id=" + id + "&shield_session_id=" + ind
		sig := xAppSignatureFallback(requestHashFromCanonical(canonical), -9999)
		res := s.h2Measured(context.Background(), "GET", "/api/v1/order/"+id+"/contact", nil, map[string]string{
			"x-latitude":           "51.1389885",
			"x-longitude":          "71.3743228",
			"accept-language":      "ru_RU",
			"x-os-type":            "android",
			"x-app-signature":      sig,
			"x-device-fingerprint": fingerprint,
			"x-app-flavor":         "indriver",
			"accept-encoding":      "gzip",
			"authorization":        "Bearer " + access,
			"mode":                 "driver",
			"x-app":                "android 5.155.0",
			"x-fingerprint":        fmt.Sprintf(`{"shield_session_id":"%s"}`, ind),
		}, 5*time.Second)
		if ph := str(res.Data["phone"]); ph != "" {
			out = ph
			return nil
		}
		switch str(deepGet(res.Data, "meta", "code")) {
		case "402", "401", "404", "445":
			out = str(deepGet(res.Data, "meta", "code"))
		default:
			out = "500"
		}
		return nil
	})
	return out
}

func (s *server) getphoneRyadom(access, ind, id, fingerprint string) string {
	var out string
	_ = s.phoneSem.run(context.Background(), func() error {
		canonical := "X-App=android 5.155.0&order_id=" + id + "&shield_session_id=" + ind
		sig := xAppSignatureFallback(requestHashFromCanonical(canonical), -9999)
		res := s.h2Measured(context.Background(), "GET", "/api/v2/orders/"+id+"/search/phone", nil, map[string]string{
			"x-latitude":           "51.1389885",
			"x-longitude":          "71.3743228",
			"accept-language":      "ru_RU",
			"x-os-type":            "android",
			"x-app-signature":      sig,
			"x-device-fingerprint": fingerprint,
			"x-app-flavor":         "indriver",
			"accept-encoding":      "gzip",
			"authorization":        "Bearer " + access,
			"mode":                 "driver",
			"x-app":                "android 5.155.0",
			"x-fingerprint":        fmt.Sprintf(`{"shield_session_id":"%s"}`, ind),
		}, 5*time.Second)
		if ph := str(res.Data["phone"]); ph != "" {
			out = ph
			_ = s.ryadomfalse(access, ind, id)
			return nil
		}
		switch str(deepGet(res.Data, "meta", "code")) {
		case "402", "401", "404", "445":
			out = str(deepGet(res.Data, "meta", "code"))
		default:
			out = "500"
		}
		return nil
	})
	return out
}

func (s *server) ryadomfalse(access, ind, id string) error {
	_ = s.h2Measured(context.Background(), "POST", "/api/v2/orders/"+id+"/search/confirm?confirmed=true", nil, map[string]string{
		"accept-language": "ru_RU",
		"x-os-type":       "android",
		"x-app-flavor":    "indriver",
		"authorization":   "Bearer " + access,
		"mode":            "driver",
		"x-app":           "android 5.155.0",
		"x-fingerprint":   fmt.Sprintf(`{"shield_session_id":"%s"}`, ind),
	}, 5*time.Second)
	return nil
}

func (s *server) skip(access, ind, id string) error {
	_ = s.h2Measured(context.Background(), "POST", "/api/v2/orders/"+id+"/search/skip", nil, map[string]string{
		"accept-language": "ru_RU",
		"x-os-type":       "android",
		"x-app-flavor":    "indriver",
		"authorization":   "Bearer " + access,
		"mode":            "driver",
		"x-app":           "android 5.155.0",
		"x-fingerprint":   fmt.Sprintf(`{"shield_session_id":"%s"}`, ind),
	}, 5*time.Second)
	return nil
}

func (s *server) h2Measured(ctx context.Context, method, path string, body any, headers map[string]string, timeout time.Duration) h2Resp {
	start := time.Now()
	data, status, err := s.h2Request(ctx, method, path, body, headers, timeout)
	s.logHTTPLatency(method, path, status, err, time.Since(start))
	if err != nil {
		return h2Resp{OK: false, Status: status, Data: map[string]any{}}
	}
	return h2Resp{OK: status >= 200 && status < 300, Status: status, Data: data}
}

func (s *server) h2Request(ctx context.Context, method, path string, body any, headers map[string]string, timeout time.Duration) (map[string]any, int, error) {
	u := strings.TrimRight(s.h2Origin, "/") + path
	var rd io.Reader
	if body != nil {
		b, err := json.Marshal(body)
		if err != nil {
			return nil, 0, err
		}
		rd = bytes.NewReader(b)
	}
	rctx, cancel := context.WithTimeout(ctx, timeout)
	defer cancel()
	req, err := http.NewRequestWithContext(rctx, method, u, rd)
	if err != nil {
		return nil, 0, err
	}
	if body != nil {
		req.Header.Set("content-type", "application/json")
	}
	for k, v := range headers {
		req.Header.Set(k, v)
	}
	resp, err := s.httpc.Do(req)
	if err != nil {
		return nil, 0, err
	}
	defer resp.Body.Close()

	br := io.Reader(resp.Body)
	if strings.EqualFold(resp.Header.Get("content-encoding"), "gzip") {
		gz, gzErr := gzip.NewReader(resp.Body)
		if gzErr == nil {
			defer gz.Close()
			br = gz
		}
	}

	raw, _ := io.ReadAll(br)
	raw = bytes.TrimPrefix(raw, []byte{0xEF, 0xBB, 0xBF})
	raw = bytes.TrimSpace(raw)
	if len(raw) == 0 {
		return map[string]any{}, resp.StatusCode, nil
	}
	var out map[string]any
	if err := json.Unmarshal(raw, &out); err != nil {
		return map[string]any{}, resp.StatusCode, nil
	}
	return out, resp.StatusCode, nil
}

func (s *server) safeFetchJSON(ctx context.Context, url string, timeout time.Duration) (map[string]any, error) {
	v, err := s.safeFetchAny(ctx, url, timeout)
	if err != nil {
		return nil, err
	}
	m, ok := asMap(v)
	if !ok {
		return nil, errors.New("not json object")
	}
	return m, nil
}

func (s *server) safeFetchAny(ctx context.Context, url string, timeout time.Duration) (any, error) {
	rctx, cancel := context.WithTimeout(ctx, timeout)
	defer cancel()
	req, _ := http.NewRequestWithContext(rctx, http.MethodGet, url, nil)
	req.Header.Set("accept", "application/json,text/plain,*/*")
	resp, err := s.httpc.Do(req)
	if err != nil {
		return nil, err
	}
	defer resp.Body.Close()
	b, _ := io.ReadAll(resp.Body)
	b = bytes.TrimPrefix(b, []byte{0xEF, 0xBB, 0xBF})
	b = bytes.TrimSpace(b)
	var out any
	if err := json.Unmarshal(b, &out); err != nil {
		return nil, err
	}
	return out, nil
}

func (s *server) maybeStartStopService() {
	myid := atomic.LoadInt64(&s.myID)
	activeCallonCount := s.activeCallonCount()
	hasActiveCallonAndMyID := activeCallonCount > 0 && myid > 0
	numbersCount := len(s.getNumbers())
	hasNumbers := numbersCount > 0

	if hasActiveCallonAndMyID && !hasNumbers {
		s.logOrdersf("[SVC line=%d] nonumbers_kick active_callon=%d myid=%d numbers=%d", s.lineID, activeCallonCount, myid, numbersCount)
		s.kickClientsNoNumbers()
		myid = atomic.LoadInt64(&s.myID)
		activeCallonCount = s.activeCallonCount()
		hasActiveCallonAndMyID = activeCallonCount > 0 && myid > 0
		numbersCount = len(s.getNumbers())
		hasNumbers = numbersCount > 0
	}

	if hasActiveCallonAndMyID && hasNumbers {
		if s.serviceOn.CompareAndSwap(false, true) {
			log.Printf("[SVC line=%d] start active_callon=%d myid=%d numbers=%d", s.lineID, activeCallonCount, myid, numbersCount)
		}
		return
	}
	s.logOrdersf("[SVC line=%d] no_start active_callon=%d myid=%d numbers=%d service_on=%t", s.lineID, activeCallonCount, myid, numbersCount, s.serviceOn.Load())
	if s.serviceOn.CompareAndSwap(true, false) {
		atomic.StoreInt64(&s.myID, 0)
		s.warmupOn.Store(false)
		s.ordersSeeded.Store(false)
		s.clearWarmNums()
		s.clearExcludedNums()
		log.Printf("[SVC line=%d] stop", s.lineID)
	}
}

func (s *server) hasActiveCallonAndMyID() bool {
	s.clientsMu.RLock()
	hasActiveCallon := len(s.toIdx) > 0
	s.clientsMu.RUnlock()
	return hasActiveCallon && atomic.LoadInt64(&s.myID) > 0
}

func (s *server) activeCallonCount() int {
	s.clientsMu.RLock()
	defer s.clientsMu.RUnlock()
	return len(s.toIdx)
}

func (s *server) kickWarmupAsync() {
	if !s.warmupOn.CompareAndSwap(false, true) {
		return
	}
	go func() {
		defer s.warmupOn.Store(false)
		ctx := context.Background()
		errRefresh := s.refreshNumbers(ctx)
		nums := s.getNumbers()
		if errRefresh != nil && len(nums) == 0 {
			return
		}
		if len(nums) == 0 {
			s.logOrdersf("[SVC line=%d] nonumbers_kick warmup=1", s.lineID)
			s.kickClientsNoNumbers()
			return
		}
		s.maybeStartStopService()
	}()
}

func (s *server) kickClientsNoNumbers() {
	msg := map[string]any{"type": "nonumbers", "ts": time.Now().UnixMilli()}
	s.broadcast(msg)
	s.closeAllClients()
}

func (s *server) heartbeat() {
	s.clientsMu.RLock()
	list := make([]*wsClient, 0, len(s.clients))
	for c := range s.clients {
		list = append(list, c)
	}
	s.clientsMu.RUnlock()
	dead := 0
	for _, c := range list {
		if !c.alive.Load() {
			dead++
			s.kickAndCleanup(c, 4001, "dead")
			continue
		}
		c.alive.Store(false)
		_ = c.conn.WriteControl(websocket.PingMessage, nil, time.Now().Add(500*time.Millisecond))
	}
	if len(list) > 0 || dead > 0 {
		s.logOrdersf("[WS line=%d] heartbeat clients=%d dead=%d", s.lineID, len(list), dead)
	}
}

func (s *server) broadcast(msg map[string]any) {
	s.clientsMu.RLock()
	list := make([]*wsClient, 0, len(s.clients))
	for c := range s.clients {
		list = append(list, c)
	}
	s.clientsMu.RUnlock()
	for _, c := range list {
		_ = c.writeJSON(msg)
	}
}

func (s *server) closeAllClients() {
	s.clientsMu.RLock()
	list := make([]*wsClient, 0, len(s.clients))
	for c := range s.clients {
		list = append(list, c)
	}
	s.clientsMu.RUnlock()
	for _, c := range list {
		s.kickAndCleanup(c, 4000, "shutdown")
	}
}

func (s *server) kickAndCleanup(c *wsClient, code int, reason string) {
	removedToID := s.cleanupClient(c)
	if removedToID > 0 {
		s.notifyToIDCount(removedToID)
	}
	s.logOrdersf("[WS line=%d] kick client=%p remote=%s code=%d reason=%s removed_toid=%d", s.lineID, c, c.remote, code, reason, removedToID)
	c.close(code, reason)
}

func (s *server) toIDCount(toid int) int {
	if toid <= 0 {
		return 0
	}
	s.clientsMu.RLock()
	defer s.clientsMu.RUnlock()
	return len(s.toIdx[toid])
}

func (s *server) cleanupClient(c *wsClient) int {
	removedToID := 0
	s.clientsMu.Lock()
	defer s.clientsMu.Unlock()
	st := s.clients[c]
	if st != nil && st.Callon != nil {
		removedToID = st.Callon.ToID
		s.indexRemoveLocked(removedToID, c)
	}
	delete(s.clients, c)
	return removedToID
}

func (s *server) notifyToIDCount(toid int) {
	if toid <= 0 {
		return
	}
	s.clientsMu.RLock()
	set := s.toIdx[toid]
	if len(set) == 0 {
		s.clientsMu.RUnlock()
		return
	}
	rawCount := len(set)
	list := make([]*wsClient, 0, rawCount)
	logins := make([]string, 0, rawCount)
	missingLogin := 0
	for c := range set {
		list = append(list, c)
		st := s.clients[c]
		if st == nil || st.Callon == nil {
			missingLogin++
			continue
		}
		login := strings.TrimSpace(st.Callon.Login)
		if login == "" {
			missingLogin++
			continue
		}
		logins = append(logins, login)
	}
	s.clientsMu.RUnlock()
	sort.Strings(logins)
	count := len(logins)
	if missingLogin > 0 {
		s.logOrdersf("[WS line=%d] toid_count_missing_login toid=%d raw=%d logins=%d missing=%d", s.lineID, toid, rawCount, count, missingLogin)
	}

	msg := map[string]any{"type": "toid_count", "toid": toid, "count": count, "logins": logins}
	for _, c := range list {
		_ = c.writeJSON(msg)
	}
}

func (s *server) indexAddLocked(toid int, c *wsClient) {
	set := s.toIdx[toid]
	if set == nil {
		set = map[*wsClient]struct{}{}
		s.toIdx[toid] = set
	}
	set[c] = struct{}{}
}

func (s *server) indexRemoveLocked(toid int, c *wsClient) {
	set := s.toIdx[toid]
	if set == nil {
		return
	}
	delete(set, c)
	if len(set) == 0 {
		delete(s.toIdx, toid)
	}
}

func (s *server) banAdd(v string) {
	key := strings.ToLower(strings.TrimSpace(v))
	if key == "" {
		return
	}
	s.banMu.Lock()
	defer s.banMu.Unlock()
	delete(s.ban, key)
	s.ban[key] = time.Now()
	s.banOrder = append(s.banOrder, key)
	s.pruneBanLocked()
}

func (s *server) banHas(v string) bool {
	key := strings.ToLower(strings.TrimSpace(v))
	if key == "" {
		return false
	}
	s.banMu.Lock()
	defer s.banMu.Unlock()
	ts, ok := s.ban[key]
	if !ok {
		return false
	}
	if time.Since(ts) > s.banTTL {
		delete(s.ban, key)
		return false
	}
	return true
}

func (s *server) pruneBan() {
	s.banMu.Lock()
	defer s.banMu.Unlock()
	s.pruneBanLocked()
}

func (s *server) pruneBanLocked() {
	now := time.Now()
	for k, ts := range s.ban {
		if now.Sub(ts) > s.banTTL {
			delete(s.ban, k)
		}
	}
	for len(s.ban) > s.banMax && len(s.banOrder) > 0 {
		k := s.banOrder[0]
		s.banOrder = s.banOrder[1:]
		delete(s.ban, k)
	}
}

func (s *server) isOrderUserBanned(order map[string]any) bool {
	user, _ := asMap(order["user"])
	first := strings.ToLower(strings.TrimSpace(str(user["first_name"])))
	avatar := strings.ToLower(strings.TrimSpace(str(user["avatar"])))
	if first == "" && avatar == "" {
		return false
	}

	now := time.Now()
	s.banMu.Lock()
	defer s.banMu.Unlock()

	check := func(key string) bool {
		if key == "" {
			return false
		}
		ts, ok := s.ban[key]
		if !ok {
			return false
		}
		if now.Sub(ts) > s.banTTL {
			delete(s.ban, key)
			return false
		}
		return true
	}

	return check(first) || check(avatar)
}

func (s *server) addOrderDedup(id string) bool {
	if id == "" {
		return false
	}
	s.dedupMu.Lock()
	defer s.dedupMu.Unlock()
	if _, ok := s.allOrders[id]; ok {
		return false
	}
	s.allOrders[id] = struct{}{}
	s.orderQ = append(s.orderQ, id)
	for len(s.orderQ) > s.orderMax {
		del := s.orderQ[0]
		s.orderQ = s.orderQ[1:]
		delete(s.allOrders, del)
	}
	return true
}

func (s *server) clearOrderDedup() {
	s.dedupMu.Lock()
	s.allOrders = map[string]struct{}{}
	s.orderQ = nil
	s.dedupMu.Unlock()
}

func (s *server) baseOrderHeaders(item numberItem) map[string]string {
	return map[string]string{
		"accept-language": "ru_RU",
		"x-os-type":       "android",
		"x-app-flavor":    "indriver",
		"x-app":           "android 5.155.0",
		"x-fingerprint":   fmt.Sprintf(`{"shield_session_id":"%s"}`, item.Ind),
		"authorization":   "Bearer " + item.Access,
		"mode":            "driver",
		"accept-encoding": "gzip",
	}
}

func (s *server) getNumbers() []numberItem {
	s.numbersMu.RLock()
	defer s.numbersMu.RUnlock()
	out := make([]numberItem, len(s.numbers))
	copy(out, s.numbers)
	return out
}

func (s *server) numberKey(it numberItem) string {
	if it.ID != "" {
		return "id:" + it.ID
	}
	return "acc:" + it.Access + "|ind:" + it.Ind
}

func (s *server) markNumWarm(it numberItem) {
	k := s.numberKey(it)
	if k == "" {
		return
	}
	s.warmMu.Lock()
	_, existed := s.warmNums[k]
	s.warmNums[k] = struct{}{}
	warmCount := len(s.warmNums)
	s.warmMu.Unlock()
	if !existed {
		s.logWarmNumberAdded(k, warmCount)
	}
}

func (s *server) isNumWarm(it numberItem) bool {
	k := s.numberKey(it)
	if k == "" {
		return false
	}
	s.warmMu.Lock()
	_, ok := s.warmNums[k]
	s.warmMu.Unlock()
	return ok
}

func (s *server) clearWarmNums() {
	s.warmMu.Lock()
	s.warmNums = map[string]struct{}{}
	s.warmMu.Unlock()
}

func (s *server) isNumExcludedKey(key string) bool {
	if key == "" {
		return false
	}
	s.exclMu.Lock()
	_, ok := s.exclNums[key]
	s.exclMu.Unlock()
	return ok
}

func (s *server) excludeNum(it numberItem) bool {
	key := s.numberKey(it)
	if key == "" {
		return false
	}

	s.exclMu.Lock()
	_, already := s.exclNums[key]
	s.exclNums[key] = struct{}{}
	s.exclMu.Unlock()

	s.numbersMu.Lock()
	out := make([]numberItem, 0, len(s.numbers))
	removed := false
	for _, n := range s.numbers {
		if s.numberKey(n) == key {
			removed = true
			continue
		}
		out = append(out, n)
	}
	s.numbers = out
	s.numbersMu.Unlock()

	s.warmMu.Lock()
	delete(s.warmNums, key)
	s.warmMu.Unlock()

	return removed || !already
}

func (s *server) clearExcludedNums() {
	s.exclMu.Lock()
	s.exclNums = map[string]struct{}{}
	s.exclMu.Unlock()
}

type orderFilters struct {
	PostGte  *float64
	SalonEq  *float64
	PersonEq *float64
	DateEq   string
	TimeAt   string
}

func buildFiltersFromOrder(it map[string]any) orderFilters {
	var f orderFilters
	order, _ := asMap(it["order"])
	data, _ := asMap(order["data"])
	date, _ := asMap(data["date"])
	if boolVal(date["is_detailed"]) {
		t := str(date["time"])
		if tt, err := time.Parse(time.RFC3339, t); err == nil {
			tt = tt.Add(5 * time.Hour)
			f.DateEq = tt.Format("2006-01-02")
			f.TimeAt = tt.Format("15:04")
		}
	}
	typeObj, _ := asMap(data["type"])
	priceObj, _ := asMap(data["price"])
	price := num(priceObj["amount"]) / 100.0
	switch int(num(typeObj["id"])) {
	case 1:
		f.PersonEq = &price
	case 2:
		f.PostGte = &price
	case 3:
		f.SalonEq = &price
	}
	return f
}

func clientCoversTime(timeOT, timeDO, target string) bool {
	t, ok1 := hhmmToMin(target)
	a, ok2 := hhmmToMin(timeOT)
	b, ok3 := hhmmToMin(timeDO)
	if !ok1 || !ok2 || !ok3 {
		return false
	}
	return t >= a && t <= b
}

func hhmmToMin(v string) (int, bool) {
	parts := strings.Split(v, ":")
	if len(parts) != 2 {
		return 0, false
	}
	h, err1 := strconv.Atoi(parts[0])
	m, err2 := strconv.Atoi(parts[1])
	if err1 != nil || err2 != nil || h < 0 || h > 23 || m < 0 || m > 59 {
		return 0, false
	}
	return h*60 + m, true
}

func orderUniq(item map[string]any) string {
	user, _ := asMap(item["user"])
	order, _ := asMap(item["order"])
	data, _ := asMap(order["data"])
	to, _ := asMap(data["to"])
	city, _ := asMap(to["city"])
	created := str(order["created_at"])
	name := str(user["avatar"])
	if name == "" {
		name = str(user["first_name"])
	}
	return name + created + str(city["id"])
}

func orderUniqName(item map[string]any) string {
	user, _ := asMap(item["user"])
	name := str(user["avatar"])
	if name == "" {
		name = str(user["first_name"])
	}
	return name
}

func requestHashFromCanonical(canonical string) string {
	sum := sha256.Sum256([]byte(canonical))
	return base64.StdEncoding.EncodeToString(sum[:])
}

func xAppSignatureFallback(requestHash string, errorCode int) string {
	payload, _ := json.Marshal(map[string]string{
		"request_hash": requestHash,
		"error_code":   strconv.Itoa(errorCode),
	})
	return base64.StdEncoding.EncodeToString(payload)
}

func deepGet(m map[string]any, path ...string) any {
	var cur any = m
	for _, p := range path {
		mm, ok := cur.(map[string]any)
		if !ok {
			return nil
		}
		cur = mm[p]
	}
	return cur
}

func asMap(v any) (map[string]any, bool) {
	m, ok := v.(map[string]any)
	return m, ok
}

func asSlice(v any) ([]any, bool) {
	a, ok := v.([]any)
	return a, ok
}

func str(v any) string {
	switch t := v.(type) {
	case string:
		return t
	case json.Number:
		return t.String()
	case float64:
		return strconv.FormatFloat(t, 'f', -1, 64)
	case int, int64, int32:
		return fmt.Sprintf("%v", t)
	case nil:
		return ""
	default:
		return fmt.Sprintf("%v", t)
	}
}

func num(v any) float64 {
	switch t := v.(type) {
	case float64:
		return t
	case int:
		return float64(t)
	case int64:
		return float64(t)
	case string:
		f, _ := strconv.ParseFloat(strings.TrimSpace(t), 64)
		return f
	case json.Number:
		f, _ := t.Float64()
		return f
	default:
		return 0
	}
}

func idInt(v any) int {
	switch t := v.(type) {
	case int:
		return t
	case int64:
		if t > int64(math.MaxInt) || t < int64(math.MinInt) {
			return 0
		}
		return int(t)
	case int32:
		return int(t)
	case float64:
		if math.IsNaN(t) || math.IsInf(t, 0) || math.Trunc(t) != t {
			return 0
		}
		if t > float64(math.MaxInt) || t < float64(math.MinInt) {
			return 0
		}
		return int(t)
	case json.Number:
		if i, err := t.Int64(); err == nil {
			if i > int64(math.MaxInt) || i < int64(math.MinInt) {
				return 0
			}
			return int(i)
		}
		f, err := t.Float64()
		if err != nil || math.IsNaN(f) || math.IsInf(f, 0) || math.Trunc(f) != f {
			return 0
		}
		if f > float64(math.MaxInt) || f < float64(math.MinInt) {
			return 0
		}
		return int(f)
	case string:
		s := strings.TrimSpace(t)
		if s == "" {
			return 0
		}
		if i, err := strconv.Atoi(s); err == nil {
			return i
		}
		f, err := strconv.ParseFloat(s, 64)
		if err != nil || math.IsNaN(f) || math.IsInf(f, 0) || math.Trunc(f) != f {
			return 0
		}
		if f > float64(math.MaxInt) || f < float64(math.MinInt) {
			return 0
		}
		return int(f)
	default:
		return 0
	}
}

func boolVal(v any) bool {
	switch t := v.(type) {
	case bool:
		return t
	case string:
		return strings.EqualFold(t, "true") || t == "1"
	case float64:
		return t != 0
	default:
		return false
	}
}

func queryEscape(s string) string {
	// Keep space encoded as %20 for compatibility with previous behavior.
	return strings.ReplaceAll(neturl.QueryEscape(s), "+", "%20")
}

func chooseNonEmpty(v ...string) string {
	for _, s := range v {
		if strings.TrimSpace(s) != "" {
			return s
		}
	}
	return ""
}

func fixedCoord7(v string) string {
	s := strings.TrimSpace(v)
	if s == "" {
		return ""
	}
	f, err := strconv.ParseFloat(strings.ReplaceAll(s, ",", "."), 64)
	if err != nil {
		return s
	}
	return strconv.FormatFloat(f, 'f', 7, 64)
}

func normalizePhoneCache(v string) string {
	raw := strings.TrimSpace(v)
	if raw == "" {
		return ""
	}

	var digits strings.Builder
	digits.Grow(len(raw))
	for _, r := range raw {
		if r >= '0' && r <= '9' {
			digits.WriteRune(r)
		}
	}

	// Keep non-phone values (e.g. 401/404/500) unchanged.
	if digits.Len() < 10 {
		return raw
	}
	return "+" + digits.String()
}

func (s *server) classifyLatencyPath(path string) string {
	switch {
	case path == s.ordersPath:
		return "orders"
	case strings.Contains(path, "/api/v1/order/") && strings.HasSuffix(path, "/contact"):
		return "phone"
	case strings.Contains(path, "/api/v2/orders/") && strings.HasSuffix(path, "/search/phone"):
		return "phone"
	default:
		return ""
	}
}

func (s *server) logHTTPLatency(method, path string, status int, err error, elapsed time.Duration) {
	// latency logs disabled
}

func (s *server) logWarmNumberAdded(key string, warmCount int) {
	// warm-number logs disabled
}

func (s *server) logOrdersf(format string, args ...any) {
	if !s.ordersLog {
		return
	}
	log.Printf(format, args...)
}

func envStr(k, d string) string {
	v := strings.TrimSpace(os.Getenv(k))
	if v == "" {
		return d
	}
	return v
}

func envInt(k string, d int) int {
	v := strings.TrimSpace(os.Getenv(k))
	if v == "" {
		return d
	}
	n, err := strconv.Atoi(v)
	if err != nil {
		return d
	}
	return n
}

func envBool(k string, d bool) bool {
	v := strings.TrimSpace(os.Getenv(k))
	if v == "" {
		return d
	}
	return boolVal(v)
}

func max(a, b int) int {
	if a > b {
		return a
	}
	return b
}
