Loading scripts/run_tests_locally-context.sh +10 −1 Original line number Diff line number Diff line Loading @@ -47,8 +47,17 @@ docker run --name crdb -d --network=tfs-br --ip 172.254.254.10 -p 26257:26257 -p cockroachdb/cockroach:latest-v22.2 start-single-node docker run --name nats -d --network=tfs-br --ip 172.254.254.11 -p 4222:4222 -p 8222:8222 \ nats:2.9 --http_port 8222 --user tfs --pass tfs123 echo echo "Waiting for initialization..." sleep 10 echo "-----------------------------" #docker logs -f crdb 2>&1 | grep --max-count=1 'finished creating default user "tfs"' while ! docker logs crdb 2>&1 | grep -q 'finished creating default user \"tfs\"'; do sleep 1; done docker logs crdb #docker logs -f nats 2>&1 | grep --max-count=1 'Server is ready' while ! docker logs nats 2>&1 | grep -q 'Server is ready'; do sleep 1; done docker logs nats #sleep 10 docker ps -a echo Loading src/common/Constants.py +2 −2 Original line number Diff line number Diff line Loading @@ -12,7 +12,7 @@ # See the License for the specific language governing permissions and # limitations under the License. import logging, uuid import logging #, uuid from enum import Enum # Default logging level Loading @@ -21,7 +21,7 @@ DEFAULT_LOG_LEVEL = logging.WARNING # Default gRPC server settings DEFAULT_GRPC_BIND_ADDRESS = '0.0.0.0' DEFAULT_GRPC_MAX_WORKERS = 200 DEFAULT_GRPC_GRACE_PERIOD = 60 DEFAULT_GRPC_GRACE_PERIOD = 10 # Default HTTP server settings DEFAULT_HTTP_BIND_ADDRESS = '0.0.0.0' Loading src/common/message_broker/backend/nats/NatsBackend.py +3 −1 Original line number Diff line number Diff line Loading @@ -39,11 +39,13 @@ class NatsBackend(_Backend): def consume(self, topic_names : Set[str], consume_timeout : float) -> Iterator[Tuple[str, str]]: out_queue = queue.Queue[Message]() unsubscribe = threading.Event() tasks = [] for topic_name in topic_names: self._nats_backend_thread.subscribe(topic_name, consume_timeout, out_queue, unsubscribe) tasks.append(self._nats_backend_thread.subscribe(topic_name, consume_timeout, out_queue, unsubscribe)) while not self._terminate.is_set(): try: yield out_queue.get(block=True, timeout=consume_timeout) except queue.Empty: continue unsubscribe.set() for task in tasks: task.cancel() src/common/message_broker/backend/nats/NatsBackendThread.py +17 −3 Original line number Diff line number Diff line Loading @@ -13,6 +13,7 @@ # limitations under the License. import asyncio, nats, nats.errors, queue, threading from typing import List from common.message_broker.Message import Message class NatsBackendThread(threading.Thread): Loading @@ -20,16 +21,23 @@ class NatsBackendThread(threading.Thread): self._nats_uri = nats_uri self._event_loop = asyncio.get_event_loop() self._terminate = asyncio.Event() self._tasks_terminated = asyncio.Event() self._publish_queue = asyncio.Queue[Message]() self._tasks : List[asyncio.Task] = list() super().__init__() def terminate(self) -> None: self._terminate.set() for task in self._tasks: task.cancel() self._tasks_terminated.set() async def _run_publisher(self) -> None: client = await nats.connect(servers=[self._nats_uri]) while not self._terminate.is_set(): try: message : Message = await self._publish_queue.get() except asyncio.CancelledError: break await client.publish(message.topic, message.content.encode('UTF-8')) await client.drain() Loading @@ -46,6 +54,8 @@ class NatsBackendThread(threading.Thread): message = await subscription.next_msg(timeout) except nats.errors.TimeoutError: continue except asyncio.CancelledError: break out_queue.put(Message(message.subject, message.data.decode('UTF-8'))) await subscription.unsubscribe() await client.drain() Loading @@ -53,9 +63,13 @@ class NatsBackendThread(threading.Thread): def subscribe( self, topic_name : str, timeout : float, out_queue : queue.Queue[Message], unsubscribe : threading.Event ) -> None: self._event_loop.create_task(self._run_subscriber(topic_name, timeout, out_queue, unsubscribe)) task = self._event_loop.create_task(self._run_subscriber(topic_name, timeout, out_queue, unsubscribe)) self._tasks.append(task) def run(self) -> None: asyncio.set_event_loop(self._event_loop) self._event_loop.create_task(self._run_publisher()) task = self._event_loop.create_task(self._run_publisher()) self._tasks.append(task) self._event_loop.run_until_complete(self._terminate.wait()) self._tasks.remove(task) self._event_loop.run_until_complete(self._tasks_terminated.wait()) src/context/.gitlab-ci.yml +2 −2 Original line number Diff line number Diff line Loading @@ -67,9 +67,9 @@ unit test context: docker run --name nats -d --network=teraflowbridge -p 4222:4222 -p 8222:8222 nats:2.9 --http_port 8222 --user tfs --pass tfs123 - echo "Waiting for initialization..." - docker logs -f crdb 2>&1 | grep -m 1 'finished creating default database "tfs_test"' - while ! docker logs crdb 2>&1 | grep -q 'finished creating default user \"tfs\"'; do sleep 1; done - docker logs crdb - docker logs -f nats 2>&1 | grep -m 1 'Server is ready' - while ! docker logs nats 2>&1 | grep -q 'Server is ready'; do sleep 1; done - docker logs nats - docker ps -a - CRDB_ADDRESS=$(docker inspect crdb --format "{{.NetworkSettings.Networks.teraflowbridge.IPAddress}}") Loading Loading
scripts/run_tests_locally-context.sh +10 −1 Original line number Diff line number Diff line Loading @@ -47,8 +47,17 @@ docker run --name crdb -d --network=tfs-br --ip 172.254.254.10 -p 26257:26257 -p cockroachdb/cockroach:latest-v22.2 start-single-node docker run --name nats -d --network=tfs-br --ip 172.254.254.11 -p 4222:4222 -p 8222:8222 \ nats:2.9 --http_port 8222 --user tfs --pass tfs123 echo echo "Waiting for initialization..." sleep 10 echo "-----------------------------" #docker logs -f crdb 2>&1 | grep --max-count=1 'finished creating default user "tfs"' while ! docker logs crdb 2>&1 | grep -q 'finished creating default user \"tfs\"'; do sleep 1; done docker logs crdb #docker logs -f nats 2>&1 | grep --max-count=1 'Server is ready' while ! docker logs nats 2>&1 | grep -q 'Server is ready'; do sleep 1; done docker logs nats #sleep 10 docker ps -a echo Loading
src/common/Constants.py +2 −2 Original line number Diff line number Diff line Loading @@ -12,7 +12,7 @@ # See the License for the specific language governing permissions and # limitations under the License. import logging, uuid import logging #, uuid from enum import Enum # Default logging level Loading @@ -21,7 +21,7 @@ DEFAULT_LOG_LEVEL = logging.WARNING # Default gRPC server settings DEFAULT_GRPC_BIND_ADDRESS = '0.0.0.0' DEFAULT_GRPC_MAX_WORKERS = 200 DEFAULT_GRPC_GRACE_PERIOD = 60 DEFAULT_GRPC_GRACE_PERIOD = 10 # Default HTTP server settings DEFAULT_HTTP_BIND_ADDRESS = '0.0.0.0' Loading
src/common/message_broker/backend/nats/NatsBackend.py +3 −1 Original line number Diff line number Diff line Loading @@ -39,11 +39,13 @@ class NatsBackend(_Backend): def consume(self, topic_names : Set[str], consume_timeout : float) -> Iterator[Tuple[str, str]]: out_queue = queue.Queue[Message]() unsubscribe = threading.Event() tasks = [] for topic_name in topic_names: self._nats_backend_thread.subscribe(topic_name, consume_timeout, out_queue, unsubscribe) tasks.append(self._nats_backend_thread.subscribe(topic_name, consume_timeout, out_queue, unsubscribe)) while not self._terminate.is_set(): try: yield out_queue.get(block=True, timeout=consume_timeout) except queue.Empty: continue unsubscribe.set() for task in tasks: task.cancel()
src/common/message_broker/backend/nats/NatsBackendThread.py +17 −3 Original line number Diff line number Diff line Loading @@ -13,6 +13,7 @@ # limitations under the License. import asyncio, nats, nats.errors, queue, threading from typing import List from common.message_broker.Message import Message class NatsBackendThread(threading.Thread): Loading @@ -20,16 +21,23 @@ class NatsBackendThread(threading.Thread): self._nats_uri = nats_uri self._event_loop = asyncio.get_event_loop() self._terminate = asyncio.Event() self._tasks_terminated = asyncio.Event() self._publish_queue = asyncio.Queue[Message]() self._tasks : List[asyncio.Task] = list() super().__init__() def terminate(self) -> None: self._terminate.set() for task in self._tasks: task.cancel() self._tasks_terminated.set() async def _run_publisher(self) -> None: client = await nats.connect(servers=[self._nats_uri]) while not self._terminate.is_set(): try: message : Message = await self._publish_queue.get() except asyncio.CancelledError: break await client.publish(message.topic, message.content.encode('UTF-8')) await client.drain() Loading @@ -46,6 +54,8 @@ class NatsBackendThread(threading.Thread): message = await subscription.next_msg(timeout) except nats.errors.TimeoutError: continue except asyncio.CancelledError: break out_queue.put(Message(message.subject, message.data.decode('UTF-8'))) await subscription.unsubscribe() await client.drain() Loading @@ -53,9 +63,13 @@ class NatsBackendThread(threading.Thread): def subscribe( self, topic_name : str, timeout : float, out_queue : queue.Queue[Message], unsubscribe : threading.Event ) -> None: self._event_loop.create_task(self._run_subscriber(topic_name, timeout, out_queue, unsubscribe)) task = self._event_loop.create_task(self._run_subscriber(topic_name, timeout, out_queue, unsubscribe)) self._tasks.append(task) def run(self) -> None: asyncio.set_event_loop(self._event_loop) self._event_loop.create_task(self._run_publisher()) task = self._event_loop.create_task(self._run_publisher()) self._tasks.append(task) self._event_loop.run_until_complete(self._terminate.wait()) self._tasks.remove(task) self._event_loop.run_until_complete(self._tasks_terminated.wait())
src/context/.gitlab-ci.yml +2 −2 Original line number Diff line number Diff line Loading @@ -67,9 +67,9 @@ unit test context: docker run --name nats -d --network=teraflowbridge -p 4222:4222 -p 8222:8222 nats:2.9 --http_port 8222 --user tfs --pass tfs123 - echo "Waiting for initialization..." - docker logs -f crdb 2>&1 | grep -m 1 'finished creating default database "tfs_test"' - while ! docker logs crdb 2>&1 | grep -q 'finished creating default user \"tfs\"'; do sleep 1; done - docker logs crdb - docker logs -f nats 2>&1 | grep -m 1 'Server is ready' - while ! docker logs nats 2>&1 | grep -q 'Server is ready'; do sleep 1; done - docker logs nats - docker ps -a - CRDB_ADDRESS=$(docker inspect crdb --format "{{.NetworkSettings.Networks.teraflowbridge.IPAddress}}") Loading