Skip to content
Merged
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
125 changes: 74 additions & 51 deletions weaviate/tests/conftest.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,84 +3,107 @@
# Licensed under a 3-clause BSD style license (see LICENSE)
import json
import os
import time
from contextlib import ExitStack

import pytest
import requests

from datadog_checks.dev import get_here
from datadog_checks.dev._env import get_state, save_state
from datadog_checks.dev.kind import kind_run
from datadog_checks.dev.kube_port_forward import port_forward
from datadog_checks.dev.subprocess import run_command
from datadog_checks.weaviate.check import DEFAULT_LIVENESS_ENDPOINT

from .common import BATCH_OBJECTS, USE_AUTH

HERE = get_here()
opj = os.path.join

NAMESPACE = 'weaviate'
# No Service targets the StatefulSet's metrics port, only the API port (see weaviate_install.yaml /
# weaviate_auth.yaml), so the API endpoint uses Service DNS while metrics fall back to the pod IP.
WEAVIATE_API_ENDPOINT = f'http://weaviate.{NAMESPACE}.svc.cluster.local:80'
POD_IP_STATE = 'weaviate_pod_ip'


def setup_weaviate():
run_command(['kubectl', 'create', 'ns', 'weaviate'])
run_command(['kubectl', 'create', 'ns', 'weaviate'], check=True)

if USE_AUTH:
run_command(['kubectl', 'apply', '-f', opj(HERE, 'kind', 'weaviate_auth.yaml'), '-n', 'weaviate'])
run_command(['kubectl', 'apply', '-f', opj(HERE, 'kind', 'weaviate_auth.yaml'), '-n', 'weaviate'], check=True)
else:
run_command(['kubectl', 'apply', '-f', opj(HERE, 'kind', 'weaviate_install.yaml'), '-n', 'weaviate'])
run_command(
['kubectl', 'apply', '-f', opj(HERE, 'kind', 'weaviate_install.yaml'), '-n', 'weaviate'], check=True
)

# Tries to ensure that the Kubernetes resources are deployed and ready before we do anything else
run_command(['kubectl', 'rollout', 'status', 'statefulset/weaviate', '-n', 'weaviate'])
run_command(['kubectl', 'wait', 'pods', '--all', '-n', 'weaviate', '--for=condition=Ready', '--timeout=600s'])
run_command(['kubectl', 'rollout', 'status', 'statefulset/weaviate', '-n', 'weaviate'], check=True)
run_command(
['kubectl', 'wait', 'pods', '--all', '-n', 'weaviate', '--for=condition=Ready', '--timeout=600s'],
check=True,
)

# `setup_weaviate` only runs once, on the initial `ddev env start`. Later invocations of the
# `dd_environment` fixture (e.g. during `ddev env stop`) run in a fresh process after the cluster
# is torn down, so the pod IP is cached here via `save_state`/`get_state` rather than looked up live.
save_state(POD_IP_STATE, get_weaviate_pod_ip())

make_weaviate_request()
Comment thread
vitkyrka marked this conversation as resolved.


def get_weaviate_pod_ip() -> str:
result = run_command(
['kubectl', 'get', 'pods', '--namespace', NAMESPACE, '--selector', 'app=weaviate', '--output', 'json'],
capture='out',
check=True,
)
pods = json.loads(result.stdout)['items']
if len(pods) != 1 or not pods[0].get('status', {}).get('podIP'):
raise RuntimeError(f'Expected one ready Weaviate pod, found {len(pods)}')
return pods[0]['status']['podIP']


def make_weaviate_request():
# This helps seed some dummy data in to Weaviate to make some metrics available. Run from a
# temporary pod since the host cannot reach the cluster directly.
weaviate_batch_endpoint = f'{WEAVIATE_API_ENDPOINT}/v1/batch/objects'

command = [
'kubectl',
'run',
'weaviate-seed-data',
'--namespace',
NAMESPACE,
'--image=curlimages/curl',
'--restart=Never',
'--attach',
'--rm',
'--quiet',
'--',
'curl',
'-sf',
'-X',
'POST',
weaviate_batch_endpoint,
'-H',
'Content-Type: application/json',
]
if USE_AUTH:
command.extend(['-H', 'Authorization: Bearer test123'])
command.extend(['-d', json.dumps(BATCH_OBJECTS)])

run_command(command, capture='both', check=True)


@pytest.fixture(scope='session')
def dd_environment():
with kind_run(conditions=[setup_weaviate]) as kubeconfig, ExitStack() as stack:
weaviate_host, weaviate_port = stack.enter_context(
port_forward(kubeconfig, 'weaviate', 2112, 'statefulset', 'weaviate')
)
weaviate_host, weaviate_api_port = stack.enter_context(
port_forward(kubeconfig, 'weaviate', 8080, 'statefulset', 'weaviate')
)
with kind_run(conditions=[setup_weaviate]) as kubeconfig:
weaviate_metrics_port = 2112

instance = {
'openmetrics_endpoint': f'http://{weaviate_host}:{weaviate_port}/metrics',
'weaviate_api_endpoint': f'http://{weaviate_host}:{weaviate_api_port}',
'openmetrics_endpoint': f'http://{get_state(POD_IP_STATE)}:{weaviate_metrics_port}/metrics',
'weaviate_api_endpoint': WEAVIATE_API_ENDPOINT,
}
if USE_AUTH:
instance['headers'] = {'Authorization': 'Bearer test123'}

make_weaviate_request(instance)
yield instance


def make_weaviate_request(instance):
# This helps seed some dummy data in to Weaviate to make some metrics available
weaviate_api_endpoint = instance.get('weaviate_api_endpoint')
weaviate_batch_endpoint = f'{weaviate_api_endpoint}/v1/batch/objects'
headers = {'content-type': 'application/json'}

if instance.get('headers'):
headers.update(instance['headers'])

if ready_check(weaviate_api_endpoint, 300):
requests.post(weaviate_batch_endpoint, headers=headers, data=json.dumps(BATCH_OBJECTS))


def ready_check(endpoint, timeout=300):
# Sometimes the API endpoint isn't ready when the cluster is ready. This will try to ensure the
# API is ready for requests before we seed some dummy data.
stop_time = time.time() + timeout
endpoint = f'{endpoint}{DEFAULT_LIVENESS_ENDPOINT}'
while time.time() < stop_time:
try:
response = requests.get(endpoint, timeout=5)
if response.ok:
return True
except requests.RequestException as e:
print(f'Request failed: {e}')

time.sleep(1)
metadata = {'agent_type': 'kubernetes', 'kubernetes': {'kubeconfig': kubeconfig}}

return False
yield instance, metadata
Loading