-
Notifications
You must be signed in to change notification settings - Fork 52
Expand file tree
/
Copy pathgrpc_echo_server.py
More file actions
87 lines (72 loc) · 2.96 KB
/
Copy pathgrpc_echo_server.py
File metadata and controls
87 lines (72 loc) · 2.96 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
import time
from concurrent import futures
import logging
import threading
import grpc
from grpc_health.v1 import health
from grpc_health.v1 import health_pb2
from grpc_health.v1 import health_pb2_grpc
from grpc_reflection.v1alpha import reflection
import echo_pb2
import echo_pb2_grpc
def _toggle_health(health_servicer: health.HealthServicer, service: str):
next_status = health_pb2.HealthCheckResponse.SERVING
while True:
if next_status == health_pb2.HealthCheckResponse.SERVING:
next_status = health_pb2.HealthCheckResponse.NOT_SERVING
else:
next_status = health_pb2.HealthCheckResponse.SERVING
health_servicer.set(service, next_status)
time.sleep(5)
def _configure_health_server(server: grpc.Server):
health_servicer = health.HealthServicer(
experimental_non_blocking=True,
experimental_thread_pool=futures.ThreadPoolExecutor(max_workers=10),
)
health_servicer.set("", health_pb2.HealthCheckResponse.SERVING)
health_servicer.set("Echo", health_pb2.HealthCheckResponse.SERVING)
health_pb2_grpc.add_HealthServicer_to_server(health_servicer, server)
# Use a daemon thread to toggle health status
toggle_health_status_thread = threading.Thread(
target=_toggle_health,
args=(health_servicer, "helloworld.Greeter"),
daemon=True,
)
toggle_health_status_thread.start()
class Echo(echo_pb2_grpc.EchoServicer):
def EchoMessage(self, request, context):
time.sleep(0.001)
return echo_pb2.EchoReply(message=f"{request.message}")
def EchoMessageStreaming(self, request_iterator, context):
for request in request_iterator:
yield echo_pb2.EchoReply(message=f"{request.message}")
def serve():
port = "8001"
server = grpc.server(futures.ThreadPoolExecutor(max_workers=10))
echo_pb2_grpc.add_EchoServicer_to_server(Echo(), server)
server.add_insecure_port("[::]:" + port)
_configure_health_server(server)
reflection.enable_server_reflection((
echo_pb2.DESCRIPTOR.services_by_name['Echo'].full_name,
reflection.SERVICE_NAME,
), server)
server.start()
print("Server started, listening on " + port)
server.wait_for_termination()
if __name__ == "__main__":
logging.basicConfig()
serve()