• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

localstack / localstack / f292942b-ed07-405a-825d-2a00ad2cd753

24 Jan 2025 06:19PM UTC coverage: 86.884% (+0.02%) from 86.86%
f292942b-ed07-405a-825d-2a00ad2cd753

push

circleci

web-flow
StepFunctions: Support for Output Blocks in Choice Rules, Improvments to JSONata Choice Defaults (#12075)

44 of 47 new or added lines in 4 files covered. (93.62%)

81 existing lines in 7 files now uncovered.

61118 of 70344 relevant lines covered (86.88%)

0.87 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

91.3
/localstack-core/localstack/utils/server/tcp_proxy.py
1
import logging
1✔
2
import select
1✔
3
import socket
1✔
4
from concurrent.futures import ThreadPoolExecutor
1✔
5
from typing import Callable
1✔
6

7
from localstack.utils.serving import Server
1✔
8

9
LOG = logging.getLogger(__name__)
1✔
10

11

12
class TCPProxy(Server):
1✔
13
    """
14
    Server based TCP proxy abstraction.
15
    This uses a ThreadPoolExecutor, so the maximum number of parallel connections is limited.
16
    """
17

18
    _target_address: str
1✔
19
    _target_port: int
1✔
20
    _handler: Callable[[bytes], tuple[bytes, bytes]] | None
1✔
21
    _buffer_size: int
1✔
22
    _thread_pool: ThreadPoolExecutor
1✔
23
    _server_socket: socket.socket | None
1✔
24

25
    def __init__(
1✔
26
        self,
27
        target_address: str,
28
        target_port: int,
29
        port: int,
30
        host: str,
31
        handler: Callable[[bytes], tuple[bytes, bytes]] = None,
32
    ) -> None:
33
        super().__init__(port, host)
1✔
34
        self._target_address = target_address
1✔
35
        self._target_port = target_port
1✔
36
        self._handler = handler
1✔
37
        self._buffer_size = 1024
1✔
38
        # thread pool limited to 64 workers for now - can be increased or made configurable if this should not suffice
39
        # for certain use cases
40
        self._thread_pool = ThreadPoolExecutor(thread_name_prefix="tcp-proxy", max_workers=64)
1✔
41
        self._server_socket = None
1✔
42

43
    def _handle_request(self, s_src: socket.socket):
1✔
44
        try:
1✔
45
            s_dst = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
1✔
46
            with s_src as s_src, s_dst as s_dst:
1✔
47
                s_dst.connect((self._target_address, self._target_port))
1✔
48

49
                sockets = [s_src, s_dst]
1✔
50
                while not self._stopped.is_set():
1✔
51
                    s_read, _, _ = select.select(sockets, [], [], 1)
1✔
52

53
                    for s in s_read:
1✔
54
                        data = s.recv(self._buffer_size)
1✔
55
                        if not data:
1✔
56
                            return
1✔
57

58
                        if s == s_src:
1✔
59
                            forward, response = data, None
1✔
60
                            if self._handler:
1✔
61
                                forward, response = self._handler(data)
×
62
                            if forward is not None:
1✔
63
                                s_dst.sendall(forward)
1✔
64
                            elif response is not None:
×
65
                                s_src.sendall(response)
×
66
                                return
×
67
                        elif s == s_dst:
1✔
68
                            s_src.sendall(data)
1✔
69
        except Exception as e:
1✔
70
            LOG.error(
1✔
71
                "Error while handling request from %s to %s:%s: %s",
72
                s_src.getpeername(),
73
                self._target_address,
74
                self._target_port,
75
                e,
76
            )
77

78
    def do_run(self):
1✔
79
        self._server_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
1✔
80
        self._server_socket.bind((self.host, self.port))
1✔
81
        self._server_socket.listen(1)
1✔
82
        self._server_socket.settimeout(10)
1✔
83
        LOG.debug(
1✔
84
            "Starting TCP proxy bound on %s:%s forwarding to %s:%s",
85
            self.host,
86
            self.port,
87
            self._target_address,
88
            self._target_port,
89
        )
90

91
        with self._server_socket:
1✔
92
            while not self._stopped.is_set():
1✔
93
                try:
1✔
94
                    src_socket, _ = self._server_socket.accept()
1✔
95
                    self._thread_pool.submit(self._handle_request, src_socket)
1✔
96
                except socket.timeout:
1✔
UNCOV
97
                    pass
×
98
                except OSError as e:
1✔
99
                    # avoid creating an error message if OSError is thrown due to socket closing
100
                    if not self._stopped.is_set():
1✔
101
                        LOG.warning("Error during during TCPProxy socket accept: %s", e)
×
102

103
    def do_shutdown(self):
1✔
104
        if self._server_socket:
1✔
105
            self._server_socket.shutdown(socket.SHUT_RDWR)
1✔
106
            self._server_socket.close()
1✔
107
        self._thread_pool.shutdown(cancel_futures=True)
1✔
108
        LOG.debug("Shut down TCPProxy on %s:%s", self.host, self.port)
1✔
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc