Coverage for chalicelib/core/log_tools/elasticsearch.py: 70%

62 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-10 12:56 +0000

1import logging 

2 

3from chalicelib.core.log_tools import log_tools 

4from chalicelib.utils import ssrf 

5from chalicelib.utils.log import sanitize 

6from elasticsearch import Elasticsearch 

7from schemas import schemas 

8 

9logger = logging.getLogger(__name__) 

10 

11IN_TY = "elasticsearch" 

12 

13 

14def get_all(tenant_id): 

15 return log_tools.get_all_by_tenant(tenant_id=tenant_id, integration=IN_TY) 

16 

17 

18def get(project_id): 

19 return log_tools.get(project_id=project_id, integration=IN_TY) 

20 

21 

22def update(tenant_id, project_id, changes): 

23 options = {} 

24 

25 if "host" in changes: 25 ↛ 27line 25 didn't jump to line 27 because the condition on line 25 was always true

26 options["host"] = changes["host"] 

27 if "apiKeyId" in changes: 27 ↛ 29line 27 didn't jump to line 29 because the condition on line 27 was always true

28 options["apiKeyId"] = changes["apiKeyId"] 

29 if "apiKey" in changes: 29 ↛ 31line 29 didn't jump to line 31 because the condition on line 29 was always true

30 options["apiKey"] = changes["apiKey"] 

31 if "indexes" in changes: 31 ↛ 33line 31 didn't jump to line 33 because the condition on line 31 was always true

32 options["indexes"] = changes["indexes"] 

33 if "port" in changes: 33 ↛ 36line 33 didn't jump to line 36 because the condition on line 33 was always true

34 options["port"] = changes["port"] 

35 

36 return log_tools.edit(project_id=project_id, integration=IN_TY, changes=options) 

37 

38 

39def add(tenant_id, project_id, host, api_key_id, api_key, indexes, port): 

40 options = { 

41 "host": host, "apiKeyId": api_key_id, "apiKey": api_key, "indexes": indexes, "port": port 

42 } 

43 return log_tools.add(project_id=project_id, integration=IN_TY, options=options) 

44 

45 

46def delete(tenant_id, project_id): 

47 return log_tools.delete(project_id=project_id, integration=IN_TY) 

48 

49 

50def add_edit(tenant_id, project_id, data: schemas.IntegrationElasticsearchSchema): 

51 s = get(project_id) 

52 if s is not None: 

53 return update(tenant_id=tenant_id, project_id=project_id, 

54 changes={"host": data.host, "apiKeyId": data.api_key_id, "apiKey": data.api_key, 

55 "indexes": data.indexes, "port": data.port}) 

56 else: 

57 return add(tenant_id=tenant_id, project_id=project_id, 

58 host=data.host, api_key=data.api_key, api_key_id=data.api_key_id, 

59 indexes=data.indexes, port=data.port) 

60 

61 

62def __get_es_client(host, port, api_key_id, api_key, use_ssl=False, timeout=15): 

63 scheme = "http" if host.startswith("http") else "https" 

64 host = host.replace("http://", "").replace("https://", "") 

65 # host/port are user-controlled; without this guard es.ping() is a blind-SSRF 

66 # oracle against internal networks and cloud metadata endpoints 

67 if not isinstance(port, int) or not 0 < port <= 65535: 

68 logger.warning(f"blocked Elasticsearch integration with invalid port: {sanitize(str(port))}") 

69 return None 

70 if host.lower() not in ssrf.get_integration_allowed_hosts() \ 70 ↛ 74line 70 didn't jump to line 74 because the condition on line 70 was always true

71 and ("/" in host or "@" in host or not ssrf.is_public_host(host, port)): 

72 logger.warning(f"blocked Elasticsearch integration with non-public host: {sanitize(host)}") 

73 return None 

74 try: 

75 args = { 

76 "hosts": [{"host": host, "port": port, "scheme": scheme}], 

77 "verify_certs": use_ssl, 

78 "request_timeout": timeout, 

79 "api_key": api_key 

80 } 

81 es = Elasticsearch( 

82 **args 

83 ) 

84 r = es.ping() 

85 if not r and not use_ssl: 

86 return __get_es_client(host, port, api_key_id, api_key, use_ssl=True, timeout=timeout) 

87 if not r: 

88 return None 

89 except Exception as err: 

90 logger.error("================exception connecting to ES host:") 

91 logger.exception(err) 

92 return None 

93 return es 

94 

95 

96def ping(tenant_id, data: schemas.IntegrationElasticsearchTestSchema): 

97 es = __get_es_client(data.host, data.port, data.api_key_id, data.api_key, timeout=3) 

98 if es is None: 98 ↛ 100line 98 didn't jump to line 100 because the condition on line 98 was always true

99 return {"state": False} 

100 return {"state": es.ping()}