Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
99 changes: 97 additions & 2 deletions libs/executors/garf/executors/entrypoints/grpc_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,18 +16,77 @@

import argparse
import logging
import os
import subprocess
import time
from concurrent import futures

import garf.executors
import grpc
from garf.executors import execution_context, garf_pb2, garf_pb2_grpc, setup
from garf.executors.entrypoints.tracer import initialize_tracer
from garf.executors import (
execution_context,
fetchers,
garf_pb2,
garf_pb2_grpc,
setup,
telemetry,
version,
)
from garf.executors.entrypoints import utils
from garf.executors.entrypoints.tracer import (
initialize_logger,
initialize_meter,
initialize_tracer,
)
from garf.executors.workflows import workflow, workflow_runner
from google.protobuf.json_format import MessageToDict
from grpc_reflection.v1alpha import reflection
from opentelemetry import metrics
from opentelemetry.instrumentation.grpc import GrpcInstrumentorServer
from opentelemetry.instrumentation.logging import LoggingInstrumentor

OTEL_SERVICE_NAME = 'garf'
LoggingInstrumentor().instrument(set_logging_format=False)

server_start_time = time.time()


def _get_server_info(options):
if not (commit_sha := os.getenv('GIT_COMMIT_SHA')):
try:
commit_sha = (
subprocess.check_output(['git', 'rev-parse', '--short', 'HEAD'])
.decode('ascii')
.strip()
)
except Exception:
commit_sha = 'Unknown'

yield metrics.Observation(
value=1,
attributes={
'version_executors': garf.executors.version.__version__,
'version_core': garf.executors.version.core_version,
'version_io': garf.executors.version.io_version,
'git_commit': commit_sha,
'server_type': 'grpc',
},
)


executor_info = telemetry.meter.create_observable_gauge(
'garf_info',
callbacks=[_get_server_info],
unit='',
description='Build info of garf executor',
)


class GarfService(garf_pb2_grpc.GarfService):
def Execute(self, request, context):
telemetry.executor_requested_counter.add(
1, attributes={'executor.source': request.source}
)
query_executor = setup.setup_executor(
request.source, request.context.fetcher_parameters
)
Expand All @@ -41,6 +100,10 @@ def Execute(self, request, context):
return garf_pb2.ExecuteResponse(results=[result])

def ExecuteBatch(self, request, context):
n_queries = len(request.batch)
telemetry.executor_requested_counter.add(
n_queries, attributes={'executor.source': request.source}
)
query_executor = setup.setup_executor(
request.source, request.context.fetcher_parameters
)
Expand Down Expand Up @@ -73,12 +136,36 @@ def ExecuteWorkflow(self, request, context):
execution_workflow = workflow.Workflow(
**MessageToDict(request.workflow, preserving_proto_field_name=True)
)
telemetry.workflow_requested.add(
1, attributes=execution_workflow.attributes
)
runner = workflow_runner.WorkflowRunner(
execution_workflow=execution_workflow
)
results = runner.run()
return garf_pb2.ExecuteWorkflowResponse(results=results)

def GetVersion(self, request, context):
return garf_pb2.GarfVersion(version=version.__version__)

def GetInfo(self, request, context):
return garf_pb2.GarfInfo(
executors_version=version.__version__,
core_version=version.core_version,
io_version=version.io_version,
)

def ListFetchers(self, request, context):
return garf_pb2.ListFetchersResponse(
results=[
garf_pb2.FetcherInfo(name=name, version=fetcher.version)
for name, fetcher in fetchers.get_all_report_fetchers().items()
]
)

def ListExecutors(self, request, context):
return garf_pb2.ListExecutorsResponse(results=setup.available_executors())


if __name__ == '__main__':
parser = argparse.ArgumentParser()
Expand All @@ -88,6 +175,14 @@ def ExecuteWorkflow(self, request, context):
)
args, _ = parser.parse_known_args()
initialize_tracer()
meter = initialize_meter()
logger = utils.init_logging(
loglevel='INFO', logger_type='local', name=OTEL_SERVICE_NAME
)
logger.addHandler(initialize_logger())

grpc_server_instrumentor = GrpcInstrumentorServer()
grpc_server_instrumentor.instrument()
server = grpc.server(
futures.ThreadPoolExecutor(max_workers=args.parallel_threshold)
)
Expand Down
1 change: 1 addition & 0 deletions libs/executors/garf/executors/entrypoints/server.py
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,7 @@ def _get_server_info(options):
'version_core': garf.executors.version.core_version,
'version_io': garf.executors.version.io_version,
'git_commit': commit_sha,
'server_type': 'http',
},
)

Expand Down
85 changes: 48 additions & 37 deletions libs/executors/garf/executors/garf_pb2.py

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading
Loading