diff --git a/.github/workflows/pylint.yml b/.github/workflows/pylint.yml index cb529fb1..c7bf41db 100644 --- a/.github/workflows/pylint.yml +++ b/.github/workflows/pylint.yml @@ -7,7 +7,7 @@ jobs: runs-on: ubuntu-latest strategy: matrix: - python-version: ["3.8", "3.9", "3.10", "3.11"] + python-version: ["3.14"] steps: - uses: actions/checkout@v7 - name: Set up Python ${{ matrix.python-version }} diff --git a/CHANGELOG.md b/CHANGELOG.md index 8cd86c87..8fb0c10d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,18 +5,26 @@ ### Breaking Changes 1. [#238](https://github.com/InfluxCommunity/influxdb3-python/pull/238): Drop support for Python 3.9. Python 3.10 or newer is now required. +1. [#239](https://github.com/InfluxCommunity/influxdb3-python/pull/239): Exception classes now function mostly like a simple data-carrying object. + - `InfluxDBPartialWriteException` class constructor will now accept `message` as an argument. + - `InfluxDBPartialWriteError.from_response(cls, response: HTTPResponse):` function was removed. ### Bug Fixes -1. [#243](https://github.com/InfluxCommunity/influxdb3-python/pull/243): Remove stale `influxdb_client` references from v3 docstrings, examples, comments, and logger names. (Closes #242) - -2. [#237](https://github.com/InfluxCommunity/influxdb3-python/pull/237): Makes the writing API simpler and more consistent with other v3 clients: +1. [#237](https://github.com/InfluxCommunity/influxdb3-python/pull/237): Makes the writing API simpler and more consistent with other v3 clients: - Further simplifies the `WriteApi` request path by constructing v2/v3 requests directly through `RestClient`, while preserving existing write behavior. +1. [#239](https://github.com/InfluxCommunity/influxdb3-python/pull/239): + - Only throws `InfluxDBPartialWriteException` when: + - Error response status code is `400`. + - Error response format `{"error":"...","data":[{"error_message":"...","line_number":2,"original_line": "..."}]}` is returned with `data` must be an array + - `accept_partial` is set to `true`. + - Write endpoint must be `api/v3/write_lp`. 1. [#241](https://github.com/InfluxCommunity/influxdb3-python/pull/241): Harden `MultiprocessingWriter` shutdown and error handling: - Replaces assertion-based runtime state validation with explicit exceptions. - Guarantees queue task completion when worker writes fail. - Adds idempotent `close()` with bounded worker shutdown and a configurable `close_timeout`. - Ensures `on_shutdown` is invoked at most once. +1. [#243](https://github.com/InfluxCommunity/influxdb3-python/pull/243): Remove stale `influxdb_client` references from v3 docstrings, examples, comments, and logger names. (Closes #242) ## 0.21.0 [2026-08-27] diff --git a/README.md b/README.md index ee60fe70..7bfe0341 100644 --- a/README.md +++ b/README.md @@ -153,7 +153,7 @@ Users can import data from CSV, JSON, Feather, ORC, Parquet import influxdb_client_3 as InfluxDBClient3 import pandas as pd import numpy as np -from influxdb_client_3 import write_client_options, WritePrecision, WriteOptions, InfluxDBError +from influxdb_client_3 import write_client_options, WritePrecision, WriteOptions, InfluxDBWriteException class BatchingCallback(object): @@ -165,10 +165,10 @@ class BatchingCallback(object): self.write_count += 1 print(f"Written batch: {conf}, data: {data}") - def error(self, conf, data: str, exception: InfluxDBError): + def error(self, conf, data: str, exception: InfluxDBWriteException): print(f"Cannot write batch: {conf}, data: {data} due: {exception}") - def retry(self, conf, data: str, exception: InfluxDBError): + def retry(self, conf, data: str, exception: InfluxDBWriteException): print(f"Retryable error occurs for batch: {conf}, data: {data} retry: {exception}") callback = BatchingCallback() @@ -245,11 +245,11 @@ client.write_dataframe( ### Accept partial writes and inspect failed lines `accept_partial` defaults to `True` and allows partial success when writing through the V3 API endpoint (`use_v2_api=False`) and a batch contains invalid lines. -On partial failure, the client raises `InfluxDBPartialWriteError` with structured `line_errors`. +On partial failure, the client raises `InfluxDBPartialWriteException` with structured `line_errors`. ```python from influxdb_client_3 import InfluxDBClient3 -from influxdb_client_3.exceptions import InfluxDBPartialWriteError +from influxdb_client_3.exceptions import InfluxDBPartialWriteException client = InfluxDBClient3( host="http://localhost:8181", @@ -261,7 +261,7 @@ lp = "home,room=Sunroom temp=96 1735545600\nhome,room=Sunroom temp=\"hi\" 173554 try: client.write(lp) # accept_partial=True by default on V3 API endpoint -except InfluxDBPartialWriteError as e: +except InfluxDBPartialWriteException as e: for line_err in e.line_errors: print(f"line {line_err.line_number} failed: {line_err.error_message} ({line_err.original_line})") ``` diff --git a/examples/advanced/database_transfer.py b/examples/advanced/database_transfer.py index 19d497c0..57594044 100644 --- a/examples/advanced/database_transfer.py +++ b/examples/advanced/database_transfer.py @@ -4,7 +4,8 @@ import os import time -from influxdb_client_3 import InfluxDBClient3, write_client_options, WriteOptions, InfluxDBError +from influxdb_client_3 import InfluxDBClient3, write_client_options, WriteOptions +from influxdb_client_3.exceptions.write_exceptions import InfluxDBWriteException HOST = os.getenv('INFLUXDB_HOST') or 'http://localhost:8181' TOKEN = os.getenv('INFLUXDB_TOKEN') or 'my-token' @@ -16,10 +17,10 @@ class BatchingCallback(object): def success(self, conf, data: str): print(f"Written batch: {conf}, data: {data}") - def error(self, conf, data: str, exception: InfluxDBError): + def error(self, conf, data: str, exception: InfluxDBWriteException): print(f"Cannot write batch: {conf}, data: {data} due: {exception}") - def retry(self, conf, data: str, exception: InfluxDBError): + def retry(self, conf, data: str, exception: InfluxDBWriteException): print(f"Retryable error occurs for batch: {conf}, data: {data} retry: {exception}") diff --git a/examples/advanced/downsample.py b/examples/advanced/downsample.py index 283888f0..6b2c5c23 100644 --- a/examples/advanced/downsample.py +++ b/examples/advanced/downsample.py @@ -9,7 +9,8 @@ import pandas as pd -from influxdb_client_3 import InfluxDBClient3, InfluxDBError, WriteOptions, write_client_options +from influxdb_client_3 import InfluxDBClient3, WriteOptions, write_client_options +from influxdb_client_3.exceptions.write_exceptions import InfluxDBWriteException dir_path = os.path.dirname(os.path.realpath(__file__)) @@ -23,10 +24,10 @@ class BatchingCallback(object): def success(self, conf, data: str): print(f"Written batch: {conf}, data: {data}") - def error(self, conf, data: str, exception: InfluxDBError): + def error(self, conf, data: str, exception: InfluxDBWriteException): print(f"Cannot write batch: {conf}, data: {data} due: {exception}") - def retry(self, conf, data: str, exception: InfluxDBError): + def retry(self, conf, data: str, exception: InfluxDBWriteException): print(f"Retryable error occurs for batch: {conf}, data: {data} retry: {exception}") @@ -74,7 +75,6 @@ def retry(self, conf, data: str, exception: InfluxDBError): database=DATABASE, enable_gzip=True, write_client_options=wco) as prep_client: - # Generating random data for i in range(num_entries): trainer = random.choice(trainers) @@ -129,7 +129,6 @@ def retry(self, conf, data: str, exception: InfluxDBError): token=TOKEN, host=HOST, database=DATABASE, enable_gzip=True, write_client_options=wco) as ds_client: - # downsample data to average number of catches per quarter-hour sql = ("SELECT date_bin('15 minutes', \"time\") as window_start, \n" "AVG(\"num\") as avg\n" diff --git a/examples/jupyter/basic-write-query.ipynb b/examples/jupyter/basic-write-query.ipynb index 87c184a7..ded8b074 100644 --- a/examples/jupyter/basic-write-query.ipynb +++ b/examples/jupyter/basic-write-query.ipynb @@ -20,7 +20,7 @@ "%env INFLUXDB_HOST=http://localhost:8181\n", "%env INFLUXDB_TOKEN=\n", "%env INFLUXDB_DATABASE=my_db\n", - "from influxdb_client_3 import InfluxDBClient3, Point, WritePrecision, InfluxDBError, WriteOptions, write_client_options" + "from influxdb_client_3 import InfluxDBClient3, Point, InfluxDBWriteException, InfluxDB3ClientQueryException" ], "outputs": [], "execution_count": null @@ -32,9 +32,14 @@ "source": "2. Next, setup a basic sensor class for generating test data. This simple class will generate readings using the built-in Influxdb3 Python `Point` class. Using the `Point` class for writes is recommended in basic applications even though the client `write()` method handles other types of data." }, { + "metadata": { + "SqlCellData": { + "variableName$1": "df_sql1" + } + }, "cell_type": "code", - "id": "532b67ce851fe031", - "metadata": {}, + "outputs": [], + "execution_count": null, "source": [ "class Sensor:\n", "\n", @@ -54,20 +59,17 @@ " .time(timestamp)\n", " )\n" ], - "outputs": [], - "execution_count": null + "id": "85b6e79a3388eb30" }, { "cell_type": "markdown", - "id": "40ac0ffee0bcaefa", + "id": "532b67ce851fe031", "metadata": {}, - "source": [ - "3. Create a client instance. Please note that this client is instantiated using default values only. WriteOptions and QueryOptions can also be added at this step. For simplicity of illustration they have been omitted." - ] + "source": "3. Create a client instance. Please note that this client is instantiated using default values only. WriteOptions and QueryOptions can also be added at this step. For simplicity of illustration they have been omitted." }, { "cell_type": "code", - "id": "fefbdd206fedf89f", + "id": "40ac0ffee0bcaefa", "metadata": {}, "source": [ "import logging\n", @@ -92,15 +94,13 @@ }, { "cell_type": "markdown", - "id": "64228a6a65db2dfc", + "id": "fefbdd206fedf89f", "metadata": {}, - "source": [ - "4. Generate data points and write them to the database." - ] + "source": "4. Generate data points and write them to the database." }, { "cell_type": "code", - "id": "1613427c088c1a56", + "id": "64228a6a65db2dfc", "metadata": {}, "source": [ "import random\n", @@ -137,7 +137,7 @@ "try:\n", " client.write(data)\n", " logging.info(f\"Write successful!\")\n", - "except InfluxDBError as e:\n", + "except InfluxDBWriteException as e:\n", " logging.error(\"InfluxDB error: {}\".format(e))\n" ], "outputs": [], @@ -145,15 +145,13 @@ }, { "cell_type": "markdown", - "id": "3422e2abc73a174", + "id": "1613427c088c1a56", "metadata": {}, - "source": [ - "Now query" - ] + "source": "Now query" }, { "cell_type": "code", - "id": "8a2fb254614ccc27", + "id": "3422e2abc73a174", "metadata": {}, "source": [ "sql = f\"SELECT time, location, model, temperature FROM {measurement} WHERE location = 'hall_east' ORDER BY time DESC\"\n", @@ -163,7 +161,7 @@ "try:\n", " table = client.query(query=sql, language=\"sql\", mode=\"pandas\")\n", " print(table)\n", - "except InfluxDBError as e:\n", + "except InfluxDB3ClientQueryException as e:\n", " print(\"InfluxDB error: {}\".format(e))" ], "outputs": [], @@ -171,15 +169,13 @@ }, { "cell_type": "markdown", - "id": "1788e61d03b7f8c5", + "id": "8a2fb254614ccc27", "metadata": {}, - "source": [ - "Sample results plot below." - ] + "source": "Sample results plot below." }, { "cell_type": "code", - "id": "b56ff2f6e600ed01", + "id": "1788e61d03b7f8c5", "metadata": {}, "source": [ "%matplotlib inline\n", diff --git a/examples/query/handle_query_error.py b/examples/query/handle_query_error.py index 2b1bd1bf..ad5abddd 100644 --- a/examples/query/handle_query_error.py +++ b/examples/query/handle_query_error.py @@ -6,7 +6,7 @@ import os from influxdb_client_3 import InfluxDBClient3 -from influxdb_client_3.exceptions import InfluxDB3ClientQueryError +from influxdb_client_3.exceptions import InfluxDB3ClientQueryException def main() -> None: @@ -29,7 +29,7 @@ def main() -> None: try: # Select from a bucket that doesn't exist client.query("Select a from cpu11") - except InfluxDB3ClientQueryError as e: + except InfluxDB3ClientQueryException as e: logging.log(logging.ERROR, e.message) diff --git a/examples/write/batching.py b/examples/write/batching.py index d74de824..d7531584 100644 --- a/examples/write/batching.py +++ b/examples/write/batching.py @@ -11,13 +11,15 @@ from bson import ObjectId import influxdb_client_3 as InfluxDBClient3 -from influxdb_client_3 import write_client_options, WritePrecision, WriteOptions, InfluxDBError +from influxdb_client_3 import write_client_options, WritePrecision, WriteOptions +from influxdb_client_3.exceptions.write_exceptions import InfluxDBWriteException class BatchingCallback(object): """ Prepare callbacks to be used to handle batching states. """ + def __init__(self): self.write_status_msg = None self.write_count = 0 @@ -29,11 +31,11 @@ def success(self, conf, data: str): self.write_count += 1 self.write_status_msg = f"SUCCESS: {self.write_count} writes" - def error(self, conf, data: str, exception: InfluxDBError): + def error(self, conf, data: str, exception: InfluxDBWriteException): print(f"Cannot write batch: {conf}, data: {len(data)} bytes, due_to: {exception}") self.write_status_msg = f"FAILURE - cause: {exception}" - def retry(self, conf, data: str, exception: InfluxDBError): + def retry(self, conf, data: str, exception: InfluxDBWriteException): print(f"Retryable error occurs for batch: {conf}, data: {len(data)} bytes, retry: {exception}") self.retry_count += 1 @@ -42,7 +44,6 @@ def elapsed(self) -> int: def main() -> None: - host = os.getenv('INFLUXDB_HOST') or 'http://localhost:8181' token = os.getenv('INFLUXDB_TOKEN') or 'my-token' database = os.getenv('INFLUXDB_DATABASE') or 'my-db' diff --git a/examples/write/fileimport.py b/examples/write/fileimport.py index bb7b4053..58e531ee 100644 --- a/examples/write/fileimport.py +++ b/examples/write/fileimport.py @@ -11,7 +11,7 @@ import os import influxdb_client_3 as InfluxDBClient3 -from influxdb_client_3 import write_client_options, WriteOptions, InfluxDBError +from influxdb_client_3 import write_client_options, WriteOptions, InfluxDBWriteException dir_path = os.path.dirname(os.path.realpath(__file__)) @@ -28,10 +28,10 @@ def success(self, conf, data: bytes): self.write_count += 1 print(f"Written batch: {conf}, data: {bytes(data)} bytes") - def error(self, conf, data: bytes, exception: InfluxDBError): + def error(self, conf, data: bytes, exception: InfluxDBWriteException): print(f"Cannot write batch: {conf}, data: {data} due: {exception}") - def retry(self, conf, data: bytes, exception: InfluxDBError): + def retry(self, conf, data: bytes, exception: InfluxDBWriteException): print(f"Retryable error occurred for batch: {conf}, data: {bytes(data)} bytes, retry: {exception}") diff --git a/examples/write/handle_http_error.py b/examples/write/handle_http_error.py index 2ef28a6a..e5b054d4 100644 --- a/examples/write/handle_http_error.py +++ b/examples/write/handle_http_error.py @@ -31,7 +31,7 @@ def main() -> None: try: client.write(lp) - except InfluxDBClient3.InfluxDBError as idberr: + except InfluxDBClient3.InfluxDBWriteException as idberr: logging.log(logging.ERROR, 'WRITE ERROR: %s (%s)', idberr.response.status, idberr.message) diff --git a/examples/write/writeoptions.py b/examples/write/writeoptions.py index 77cb9114..3231a237 100644 --- a/examples/write/writeoptions.py +++ b/examples/write/writeoptions.py @@ -8,14 +8,15 @@ import logging import os -from influxdb_client_3 import (exceptions, InfluxDBClient3, Point, +from influxdb_client_3 import (InfluxDBClient3, Point, WriteOptions, WritePrecision, WriteType, write_client_options) +from influxdb_client_3.exceptions import write_exceptions logger = logging.getLogger("writeoptions") # An illustrative callback - see below -def error_callback(conf, data: bytes, exception: exceptions.InfluxDBError): +def error_callback(conf, data: bytes, exception: write_exceptions.InfluxDBWriteException): now = datetime.datetime.now() logger.warning(f"[{now}] an error occurred on latest write: {exception}") logger.warning(f" conf: {conf}") diff --git a/influxdb_client_3/__init__.py b/influxdb_client_3/__init__.py index b5b3c8f2..25c7896c 100644 --- a/influxdb_client_3/__init__.py +++ b/influxdb_client_3/__init__.py @@ -15,8 +15,8 @@ import polars as pl from pyarrow import ArrowException -from influxdb_client_3.exceptions import InfluxDB3ClientQueryError -from influxdb_client_3.exceptions import InfluxDBError +from influxdb_client_3.exceptions import InfluxDB3ClientQueryException +from influxdb_client_3.exceptions.write_exceptions import InfluxDBWriteException from influxdb_client_3.query.query_api import QueryApi as _QueryApi, QueryApiOptionsBuilder from influxdb_client_3.read_file import UploadFile from influxdb_client_3.write_client import WriteOptions, Point @@ -517,7 +517,7 @@ def write_dataframe( :type database: str, optional :param kwargs: Additional arguments to pass to the write API. :raises TypeError: If df is not a pandas or polars DataFrame. - :raises InfluxDBError: If there is an error writing to the database. + :raises InfluxDBWriteException: If there is an error writing to the database. Example: >>> import pandas as pd @@ -554,7 +554,7 @@ def write_dataframe( data_frame_timestamp_timezone=timestamp_timezone, **kwargs ) - except InfluxDBError as e: + except InfluxDBWriteException as e: raise e def write_file(self, file, measurement_name=None, tag_columns=None, timestamp_column='time', database=None, @@ -637,7 +637,7 @@ def query(self, query: str, language: str = "sql", mode: str = "all", database: try: return self._query_api.query(query=query, language=language, mode=mode, database=database, **kwargs) except ArrowException as e: - raise InfluxDB3ClientQueryError(f"Error while executing query: {e}") + raise InfluxDB3ClientQueryException(f"Error while executing query: {e}") def query_dataframe( self, @@ -715,7 +715,7 @@ async def query_async(self, query: str, language: str = "sql", mode: str = "all" database=database, **kwargs) except ArrowException as e: - raise InfluxDB3ClientQueryError(f"Error while executing query: {e}") + raise InfluxDB3ClientQueryException(f"Error while executing query: {e}") def get_server_version(self) -> Optional[str]: """ diff --git a/influxdb_client_3/cli.py b/influxdb_client_3/cli.py index b51e4c6c..21c952d9 100644 --- a/influxdb_client_3/cli.py +++ b/influxdb_client_3/cli.py @@ -14,7 +14,8 @@ INFLUX_TOKEN, InfluxDBClient3, ) -from influxdb_client_3.exceptions import InfluxDB3ClientQueryError, InfluxDBError +from influxdb_client_3.exceptions import InfluxDB3ClientQueryException +from influxdb_client_3.exceptions.write_exceptions import InfluxDBWriteException def _resolve_option( @@ -196,7 +197,7 @@ def _run_query(args, stdout, stderr, env: Optional[Mapping[str, str]] = None) -> else: stdout.write(payload) return 0 - except (InfluxDB3ClientQueryError, InfluxDBError, OSError, pa.ArrowException) as error: + except (InfluxDB3ClientQueryException, InfluxDBWriteException, OSError, pa.ArrowException) as error: _write_error(stderr, str(error)) return 1 diff --git a/influxdb_client_3/exceptions/__init__.py b/influxdb_client_3/exceptions/__init__.py index eec96cb3..6f6e5b6a 100644 --- a/influxdb_client_3/exceptions/__init__.py +++ b/influxdb_client_3/exceptions/__init__.py @@ -1,4 +1,5 @@ # flake8: noqa -from .exceptions import InfluxDB3ClientQueryError, InfluxDBError, InfluxDB3ClientError, InfluxDBPartialWriteError, \ - InfluxDBPartialWriteLineError +from .exceptions import InfluxDB3ClientException, InfluxDB3ClientQueryException +from .write_exceptions import InfluxDBWriteException, InfluxDBPartialWriteException, InfluxDBRestClientException + diff --git a/influxdb_client_3/exceptions/exceptions.py b/influxdb_client_3/exceptions/exceptions.py index 2492fe3f..f25691a8 100644 --- a/influxdb_client_3/exceptions/exceptions.py +++ b/influxdb_client_3/exceptions/exceptions.py @@ -1,16 +1,14 @@ -"""Exceptions utils for InfluxDB.""" +"""Exception utils for InfluxDB.""" +from __future__ import absolute_import -import json import logging -from dataclasses import dataclass -from typing import List, Optional, Tuple from urllib3 import HTTPResponse logger = logging.getLogger('influxdb_client_3.exceptions') -class InfluxDB3ClientError(Exception): +class InfluxDB3ClientException(Exception): """ Exception raised for errors in the InfluxDB client operations. @@ -23,13 +21,13 @@ class InfluxDB3ClientError(Exception): # This error is for all query operations -class InfluxDB3ClientQueryError(InfluxDB3ClientError): +class InfluxDB3ClientQueryException(InfluxDB3ClientException): """ Represents an error that occurs when querying an InfluxDB client. This class is specifically designed to handle errors originating from client queries to an InfluxDB database. It extends the general - `InfluxDBClientError`, allowing more precise identification and + `InfluxDB3ClientException`, allowing more precise identification and handling of query-related issues. :ivar message: Contains the specific error message describing the @@ -42,214 +40,37 @@ def __init__(self, error_message, *args, **kwargs): self.message = error_message -def _is_partial_write_error(error_message) -> bool: - if not isinstance(error_message, str) or not error_message: - return False - normalized = error_message.lower() - return ( - "partial write of line protocol occurred" in normalized or - "parsing failed for write_lp endpoint" in normalized - ) - - -def _parse_partial_write_data_item(item) -> Optional[Tuple[str, int, str]]: - if item is None: - return None - if not isinstance(item, dict): - raise ValueError("array item is not an object") - - error_message = item.get("error_message") - if not isinstance(error_message, str): - raise ValueError("error_message must be string") - if not error_message: - return None - - line_number_raw = item.get("line_number") - if line_number_raw is None: - line_number = 0 - elif isinstance(line_number_raw, int): - line_number = line_number_raw - else: - raise ValueError("line_number must be int") - - original_line_raw = item.get("original_line") - if original_line_raw is None: - original_line = "" - elif isinstance(original_line_raw, str): - original_line = original_line_raw - else: - raise ValueError("original_line must be string") - - return error_message, line_number, original_line - - -def _parse_typed_partial_write_array(data) -> Optional[List[Tuple[str, int, str]]]: - if not isinstance(data, list): - return None - line_errors: List[Tuple[str, int, str]] = [] - try: - for item in data: - parsed = _parse_partial_write_data_item(item) - if parsed is None: - continue - line_errors.append(parsed) - except ValueError: - return None - return line_errors if len(line_errors) > 0 else None - - -def _parse_typed_partial_write_object_or_none(data) -> Optional[Tuple[str, int, str]]: - try: - return _parse_partial_write_data_item(data) - except ValueError: - return None - - -def _format_partial_write_details(line_errors: List[Tuple[str, int, str]]) -> List[str]: - details: List[str] = [] - for error_message, line_number, original_line in line_errors: - if line_number != 0: - if original_line != "": - details.append(f"\tline {line_number}: {error_message} ({original_line})") - else: - details.append(f"\tline {line_number}: {error_message}") - elif error_message: - details.append(f"\t{error_message}") - return details - - -def _parse_partial_write_line_error_info(data) -> Tuple[List[Tuple[str, int, str]], List[str]]: - if data is None: - return [], [] - - typed_array = _parse_typed_partial_write_array(data) - if typed_array is not None: - return typed_array, _format_partial_write_details(typed_array) - - if isinstance(data, list): - details: List[str] = [] - for item in data: - if item is None: - continue - raw = json.dumps(item, separators=(',', ':')) - if raw and raw.lower() != "null": - details.append(raw) - return [], details - - typed_single = _parse_typed_partial_write_object_or_none(data) - if typed_single is not None: - return [typed_single], _format_partial_write_details([typed_single]) - - return [], [] - - -# This error is for all write operations -class InfluxDBError(InfluxDB3ClientError): - """Raised when a server error occurs.""" - - def __init__(self, response: HTTPResponse = None, message: str = None): - """Initialize the InfluxDBError handler.""" - if response is not None: - self.response = response - self.message = self._get_message(response) - self.retry_after = response.getheader('Retry-After') +class InfluxDBRestClientException(InfluxDB3ClientException): + def __init__(self, http_resp: HTTPResponse = None, status: int = None, message: str = None, reason: str = None): + super().__init__() + if http_resp: + self.status = http_resp.status + self.reason = http_resp.reason + self.body = http_resp.data + self.message = message or '' + self.headers = http_resp.getheaders() + self.response = http_resp else: - self.response = None + self.status = status + self.reason = reason + self.body = None + self.headers = None self.message = message or 'no response' - self.retry_after = None - super().__init__(self.message) - - def _get_message(self, response): - if response.data: - def get(d, key): - if not key or d is None: - return d - if not isinstance(d, dict): - return None - return get(d.get(key[0]), key[1:]) - try: - node = json.loads(response.data) - if isinstance(node, dict): - # InfluxDB v3 error format: { "code": "...", "message": "..." } - code = node.get("code") - message = node.get("message") - if message: - return f"{code}: {message}" if code else message - # InfluxDB v3 write error format: - # { - # "error": "...", - # "data": [ { "error_message": "...", "line_number": 2, "original_line": "..." }, ... ] - # } - error_text = node.get("error") - if error_text and _is_partial_write_error(error_text): - _, details = _parse_partial_write_line_error_info(node.get("data")) - if details: - return error_text + ":\n" + "\n".join( - detail if detail.startswith("\t") else f"\t{detail}" - for detail in details - ) - return error_text - if error_text: - return error_text - for key in [['message'], ['data', 'error_message'], ['error']]: - value = get(node, key) - if value is not None: - return value - return response.data - except Exception as e: - logging.debug(f"Cannot parse error response to JSON: {response.data}, {e}") - return response.data - - # Header - for header_key in ["X-Platform-Error-Code", "X-Influx-Error", "X-InfluxDb-Error"]: - header_value = response.getheader(header_key) - if header_value is not None: - return header_value - - # Http Status - return response.reason + self.response = None def getheaders(self): """Helper method to make response headers more accessible.""" return self.response.getheaders() + def __str__(self): + """Get custom error messages for exception.""" + error_message = "({0})\n" \ + "Reason: {1}\n".format(self.status, self.reason) + if self.headers: + error_message += "HTTP response headers: {0}\n".format( + self.headers) + + if self.body: + error_message += "HTTP response body: {0}\n".format(self.body) -@dataclass(frozen=True) -class InfluxDBPartialWriteLineError: - line_number: int - error_message: str - original_line: str - - -class InfluxDBPartialWriteError(InfluxDBError): - """Structured partial-write error with per-line failures.""" - - def __init__(self, response: HTTPResponse, line_errors: List[InfluxDBPartialWriteLineError]): - super().__init__(response=response) - self.line_errors = line_errors - - @classmethod - def from_response(cls, response: HTTPResponse): - if response is None or not response.data: - return None - try: - node = json.loads(response.data) - except Exception: - return None - if not isinstance(node, dict): - return None - error_text = node.get("error") - if not _is_partial_write_error(error_text): - return None - parsed_line_errors, _ = _parse_partial_write_line_error_info(node.get("data")) - if not parsed_line_errors: - return None - line_errors = [ - InfluxDBPartialWriteLineError( - line_number=line_number, - error_message=error_message, - original_line=original_line, - ) - for error_message, line_number, original_line in parsed_line_errors - ] - return cls(response=response, line_errors=line_errors) + return error_message diff --git a/influxdb_client_3/exceptions/write_exceptions.py b/influxdb_client_3/exceptions/write_exceptions.py new file mode 100644 index 00000000..08104616 --- /dev/null +++ b/influxdb_client_3/exceptions/write_exceptions.py @@ -0,0 +1,229 @@ +# coding: utf-8 + +from __future__ import absolute_import + +import json +import logging +from dataclasses import dataclass +from http import HTTPStatus +from json import JSONDecodeError +from typing import Union, Optional, Any, Tuple, List + +from urllib3 import HTTPResponse + +from influxdb_client_3.exceptions.exceptions import InfluxDBRestClientException + +_UTF_8_encoding = 'utf-8' + +logger = logging.getLogger('influxdb_client_3.write_client.write_exceptions') + + +class InfluxDBWriteException(InfluxDBRestClientException): + def __init__(self, + status=None, + reason=None, + http_resp=None, + message=None): + """Initialize with HTTP response.""" + super().__init__(status=status, reason=reason, http_resp=http_resp, message=message) + if http_resp is not None: + self.retry_after = http_resp.getheader('Retry-After') + else: + self.retry_after = None + + +@dataclass(frozen=True) +class InfluxDBPartialWriteLineException: + line_number: Optional[int] + error_message: Optional[str] + original_line: Optional[str] + + +class InfluxDBPartialWriteException(InfluxDBWriteException): + """Structured partial-write error with per-line failures.""" + + def __init__(self, http_resp: HTTPResponse, message: str, line_errors: List[InfluxDBPartialWriteLineException]): + super().__init__(http_resp=http_resp, message=message) + self.line_errors = line_errors + + @classmethod + def handle_partial_write_error(cls, response, root: dict): + reason = root.get("error") or "" + all_typed, line_errors = InfluxDBPartialWriteException.parse_partial_write_line_errors(root.get("data")) + line_error_details = InfluxDBPartialWriteException.format_partial_write_details(root, all_typed, line_errors) + if line_error_details: + details_str = "".join(f"\n\t{detail}" for detail in line_error_details) + reason = f"{reason}:{details_str}" + return cls(response, reason, line_errors) + + @staticmethod + def is_partial_write_error_data_shape(root): + return (isinstance(root, dict) and + root.get("error") is not None and + (isinstance(root.get('data'), list) and + len(root.get('data')) > 0)) + + @staticmethod + def is_partial_write_error(status_code: int, use_v2_api: bool, accept_partial: bool, root: dict) -> bool: + return (status_code == HTTPStatus.BAD_REQUEST and + accept_partial is True and + use_v2_api is False and + InfluxDBPartialWriteException.is_partial_write_error_data_shape(root) + ) + + @staticmethod + def parse_partial_write_line_errors(data: Any) -> Tuple[bool, List[InfluxDBPartialWriteLineException]]: + all_typed = True + line_errors = [] + for item in data: + ok, line_error = InfluxDBPartialWriteException.parse_partial_write_line_error(item) + if not ok: + all_typed = False + continue + line_errors.append(line_error) + + return all_typed, line_errors + + @staticmethod + def parse_partial_write_line_error(item: dict) -> tuple[bool, Optional[InfluxDBPartialWriteLineException]]: + ok = True + if not isinstance(item, dict): + return False, None + + line_number = item.get("line_number") + if line_number is not None and (not isinstance(line_number, int) or isinstance(line_number, bool)): + ok = False + + error_message = item.get("error_message") + if not error_message: + ok = False + + return ok, InfluxDBPartialWriteLineException(line_number, error_message, item.get("original_line")) + + @staticmethod + def format_partial_write_details( + root: dict, all_typed: bool, line_errors: List[InfluxDBPartialWriteLineException] + ) -> List[str]: + if all_typed: + return [ + InfluxDBPartialWriteException.format_line_error(err.line_number, err.error_message, err.original_line) + for err in line_errors] + + return [ + json.dumps(raw, separators=(',', ':')) + for raw in (root.get('data') or []) + if raw is not None and raw != "null" + ] + + @staticmethod + def format_line_error(line_number: Optional[int], error_message: Optional[str], + original_line: Optional[str]) -> str: + if line_number is not None and original_line is not None: + return f"line {line_number}: {error_message} ({original_line})" + if line_number is not None: + return f"line {line_number}: {error_message}" + return f"{error_message}" + + +def translate_write_exception( + exc: InfluxDBRestClientException, + use_v2_api=False, + accept_partial=False, +) -> Union[InfluxDBWriteException, InfluxDBPartialWriteException]: + if exc.status == HTTPStatus.METHOD_NOT_ALLOWED: + return create_unsupported_endpoint_exception(use_v2_api) + + if exc.body is None and exc.status == 0: + return InfluxDBWriteException(status=exc.status, reason=exc.reason, http_resp=exc.response) + + root = parse_json(exc.body) + if (root is None or + root == "" or + isinstance(root, dict) is False or + (isinstance(root, dict) and (not root.get("error") and not root.get("message")))): + message = extract_fallback_reason(exc.response) + return InfluxDBWriteException(status=exc.status, reason=exc.reason, message=message, http_resp=exc.response) + + if InfluxDBPartialWriteException.is_partial_write_error(exc.status, use_v2_api, accept_partial, root): + # InfluxDB 3 Core/Enterprise partial write error format: + # {"error":"...","data":[{"error_message":"...","line_number":2,"original_line": "..."}]} + return InfluxDBPartialWriteException.handle_partial_write_error(exc.response, root) + + return InfluxDBWriteException(status=exc.status, reason=exc.reason, message=get_message(root, exc.response), + http_resp=exc.response) + + +def get_message(root, response): + if root: + try: + if isinstance(root, dict): + # InfluxDB v3 error format: { "code": "...", "message": "..." } + message = root.get("message") + if message: + code = root.get("code") + return f"{code}: {message}" if code else message + + error_text = root.get("error") + data = root.get("data") + if error_text and isinstance(data, dict): + # Core/Enterprise object format: + # {"error":"...","data":{"error_message":"..."}} + return format_object_data_error(root.get("data"), error_text) + return error_text + except Exception as e: + logger.debug("Cannot parse error response to JSON: %s, %s", response.data, e) + return response.data + return None + + +def parse_json(json_str: Optional[Union[str, bytes]]) -> Optional[Any]: + if not json_str: + return None + try: + return json.loads(json_str) + except (JSONDecodeError, TypeError, ValueError) as e: + logger.debug("Can't parse msg from response body %s: %s", json_str, e) + return None + + +def create_unsupported_endpoint_exception(use_v2_api: bool) -> InfluxDBWriteException: + if use_v2_api: + message = ("Server doesn't support the V2 API endpoint (/api/v2/write). " + "Set use_v2_api=False to use the V3 API endpoint.") + else: + message = ("Server doesn't support the V3 API endpoint (/api/v3/write_lp). " + "Set use_v2_api=True to use the V2 API endpoint.") + ex = InfluxDBWriteException(status=0, reason=message) + ex.message = message + ex.args = (message,) + return ex + + +def format_object_data_error(data_node: dict, error: str) -> str: + line_number = data_node.get("line_number") + error_message = data_node.get("error_message") + original_line = data_node.get("original_line") + + if error_message and isinstance(line_number, int) is False: + return f"{error}:\n\t{error_message}" + elif error_message and isinstance(line_number, int) and not original_line: + return f"{error}:\n\tline {line_number}: {error_message}" + elif error_message and original_line: + return f"{error}:\n\tline {line_number}: {error_message} ({original_line})" + + return error + + +def extract_fallback_reason(response) -> str: + # Fallback to header + for header_key in ["X-Platform-Error-Code", "X-Influx-Error", "X-InfluxDb-Error"]: + header_value = response.getheader(header_key) + if header_value is not None: + return header_value + + # Fallback to raw body + if response.data is not None and response.data != "": + return response.data + + # Fallback to http Status + return response.reason diff --git a/influxdb_client_3/write_client/_sync/rest_client.py b/influxdb_client_3/write_client/_sync/rest_client.py index f124b9d0..fa57ed38 100644 --- a/influxdb_client_3/write_client/_sync/rest_client.py +++ b/influxdb_client_3/write_client/_sync/rest_client.py @@ -10,7 +10,7 @@ from typing import Dict from urllib.parse import urlencode -from influxdb_client_3.write_client.write_exceptions import ApiException +from influxdb_client_3.exceptions.exceptions import InfluxDBRestClientException try: import urllib3 @@ -172,10 +172,13 @@ def request(self, method, path, query_params=None, headers=None, ) except urllib3.exceptions.SSLError as e: msg = "{0}\n{1}".format(type(e).__name__, str(e)) - raise ApiException(status=0, reason=msg) + raise InfluxDBRestClientException(status=0, reason=msg) r = RESTResponse(r) - r.data = r.data.decode('utf8') + if r.data is not None and r.data != "": + r.data = r.data.decode('utf8') + else: + r.data = None if self.debug: RestClient.log_response(r.status) @@ -186,7 +189,7 @@ def request(self, method, path, query_params=None, headers=None, RestClient.log_body(r.data, '<<<') if not 200 <= r.status <= 299: - raise ApiException(http_resp=r) + raise InfluxDBRestClientException(http_resp=r) return r diff --git a/influxdb_client_3/write_client/client/util/helpers.py b/influxdb_client_3/write_client/client/util/helpers.py index f5556a82..ada5be43 100644 --- a/influxdb_client_3/write_client/client/util/helpers.py +++ b/influxdb_client_3/write_client/client/util/helpers.py @@ -1,5 +1,4 @@ """Functions to share utility across client classes.""" -from influxdb_client_3.write_client.write_exceptions import ApiException def _is_id(value): @@ -34,17 +33,17 @@ def get_org_query_param(org, client, required_id=False): try: organizations = client.organizations_api().find_organizations(org=_org) if len(organizations) < 1: - from write_client.client.exceptions import InfluxDBError + from influxdb_client_3.exceptions import InfluxDBRestClientException message = f"The client cannot find organization with name: '{_org}' " \ "to determine their ID. Are you using token with sufficient permission?" - raise InfluxDBError(response=None, message=message) + raise InfluxDBRestClientException(response=None, message=message) return organizations[0].id - except ApiException as e: + except InfluxDBRestClientException as e: if e.status == 404: - from write_client.client.exceptions import InfluxDBError + from influxdb_client_3.exceptions import InfluxDBRestClientException message = f"The client cannot find organization with name: '{_org}' " \ "to determine their ID." - raise InfluxDBError(response=None, message=message) + raise InfluxDBRestClientException(response=None, message=message) raise e return _org diff --git a/influxdb_client_3/write_client/client/util/multiprocessing_helper.py b/influxdb_client_3/write_client/client/util/multiprocessing_helper.py index 23509ae9..7978aded 100644 --- a/influxdb_client_3/write_client/client/util/multiprocessing_helper.py +++ b/influxdb_client_3/write_client/client/util/multiprocessing_helper.py @@ -10,7 +10,7 @@ import queue from influxdb_client_3 import write_client_options -from influxdb_client_3.exceptions import InfluxDBError +from influxdb_client_3.exceptions.write_exceptions import InfluxDBWriteException from influxdb_client_3.write_client import WriteOptions, WriteApi from influxdb_client_3.write_client._sync import rest_client @@ -22,12 +22,12 @@ def _success_callback(conf: (str, str, str), data: str): logger.debug(f"Written batch: {conf}, data: {data}") -def _error_callback(conf: (str, str, str), data: str, exception: InfluxDBError): +def _error_callback(conf: (str, str, str), data: str, exception: InfluxDBWriteException): """Unsuccessfully writen batch.""" logger.debug(f"Cannot write batch: {conf}, data: {data} due: {exception}") -def _retry_callback(conf: (str, str, str), data: str, exception: InfluxDBError): +def _retry_callback(conf: (str, str, str), data: str, exception: InfluxDBWriteException): """Retryable error.""" logger.debug(f"Retryable error occurs for batch: {conf}, data: {data} retry: {exception}") @@ -88,7 +88,7 @@ def main(): .. code-block:: python from influxdb_client_3 import WriteOptions - from influxdb_client_3.exceptions import InfluxDBError + from influxdb_client_3.exceptions.write_exceptions import InfluxDBWriteException from influxdb_client_3.write_client.client.util.multiprocessing_helper import MultiprocessingWriter @@ -97,10 +97,10 @@ class BatchingCallback(object): def success(self, conf: (str, str, str), data: str): print(f"Written batch: {conf}, data: {data}") - def error(self, conf: (str, str, str), data: str, exception: InfluxDBError): + def error(self, conf: (str, str, str), data: str, exception: InfluxDBWriteException): print(f"Cannot write batch: {conf}, data: {data} due: {exception}") - def retry(self, conf: (str, str, str), data: str, exception: InfluxDBError): + def retry(self, conf: (str, str, str), data: str, exception: InfluxDBWriteException): print(f"Retryable error occurs for batch: {conf}, data: {data} retry: {exception}") diff --git a/influxdb_client_3/write_client/client/write/retry.py b/influxdb_client_3/write_client/client/write/retry.py index d469b224..54dce27f 100644 --- a/influxdb_client_3/write_client/client/write/retry.py +++ b/influxdb_client_3/write_client/client/write/retry.py @@ -9,7 +9,9 @@ from urllib3 import Retry from urllib3.exceptions import MaxRetryError, ResponseError -from influxdb_client_3.exceptions import InfluxDBError +from influxdb_client_3.exceptions.write_exceptions import translate_write_exception, \ + InfluxDBWriteException +from influxdb_client_3.exceptions.exceptions import InfluxDBRestClientException logger = logging.getLogger('influxdb_client_3.write_client.client.write.retry') @@ -124,14 +126,14 @@ def increment(self, method=None, url=None, response=None, error=None, _pool=None new_retry = super().increment(method, url, response, error, _pool, _stacktrace) if response is not None: - parsed_error = InfluxDBError(response=response) + parsed_error = translate_write_exception(exc=InfluxDBRestClientException(http_resp=response)) elif error is not None: parsed_error = error else: parsed_error = f"Failed request to: {url}" message = f"The retriable error occurred during request. Reason: '{parsed_error}'." - if isinstance(parsed_error, InfluxDBError): + if isinstance(parsed_error, InfluxDBWriteException): message += f" Retry in {parsed_error.retry_after}s." if self.retry_callback: diff --git a/influxdb_client_3/write_client/client/write_api.py b/influxdb_client_3/write_client/client/write_api.py index 09ceec0d..6c31ab96 100644 --- a/influxdb_client_3/write_client/client/write_api.py +++ b/influxdb_client_3/write_client/client/write_api.py @@ -1,6 +1,5 @@ """Collect and write time series data to InfluxDB Cloud or InfluxDB OSS.""" from __future__ import absolute_import - # TODO Remove after this program no longer supports Python 3.8.* from __future__ import annotations @@ -11,7 +10,6 @@ import warnings from collections import defaultdict from enum import Enum -from http import HTTPStatus from multiprocessing.pool import ThreadPool from random import random from time import sleep @@ -23,15 +21,14 @@ from reactivex.scheduler import ThreadPoolScheduler from reactivex.subject import Subject -from influxdb_client_3.exceptions import InfluxDBPartialWriteError +from influxdb_client_3.exceptions.write_exceptions import _UTF_8_encoding, translate_write_exception +from influxdb_client_3.exceptions.exceptions import InfluxDBRestClientException from influxdb_client_3.write_client._sync.rest_client import RestClient -# from influxdb_client_3.write_client.client._base import _HAS_DATACLASS from influxdb_client_3.write_client.client.write.dataframe_serializer import DataframeSerializer from influxdb_client_3.write_client.client.write.point import Point, DEFAULT_WRITE_PRECISION, sanitize_tag_order from influxdb_client_3.write_client.client.write.retry import WritesRetry from influxdb_client_3.write_client.domain import WritePrecision from influxdb_client_3.write_client.domain.write_precision_converter import WritePrecisionConverter -from influxdb_client_3.write_client.write_exceptions import _UTF_8_encoding, ApiException from influxdb_client_3.write_client.write_defaults import ( DEFAULT_WRITE_ACCEPT_PARTIAL as _DEFAULT_WRITE_ACCEPT_PARTIAL, DEFAULT_WRITE_NO_SYNC as _DEFAULT_WRITE_NO_SYNC, @@ -486,8 +483,8 @@ async def post_write_async(self, org, bucket, body, **kwargs): # noqa: E501,D40 http_kwargs.get('_request_timeout'), http_kwargs.get('urlopen_kw', None), ) - except ApiException as e: - raise self._translate_write_exception(e, use_v2_api) + except InfluxDBRestClientException as e: + raise translate_write_exception(e, use_v2_api, accept_partial=accept_partial) def flush(self): """ @@ -681,8 +678,8 @@ def _post_write(self, _async_req, bucket, org, body, precision, no_sync, accept_ def translated_get(timeout=None): try: return original_get(timeout=timeout) - except ApiException as e: - raise self._translate_write_exception(e, use_v2_api) + except InfluxDBRestClientException as e: + raise translate_write_exception(e, use_v2_api, accept_partial) result.get = translated_get return result @@ -692,8 +689,8 @@ def translated_get(timeout=None): body, http_kwargs.get('_request_timeout'), http_kwargs.get('urlopen_kw', None)) - except ApiException as e: - raise self._translate_write_exception(e, use_v2_api) + except InfluxDBRestClientException as e: + raise translate_write_exception(e, use_v2_api, accept_partial) def _build_write_request(self, org, bucket, precision, no_sync, accept_partial, use_v2_api, **kwargs): if org is None: @@ -783,26 +780,6 @@ def _request(self, request, body, _request_timeout=None, urlopen_kw=None): return response_data - def _translate_write_exception(self, exc, use_v2_api): - if use_v2_api and exc.status == HTTPStatus.METHOD_NOT_ALLOWED: - message = ("Server doesn't support the V2 API endpoint (/api/v2/write). " - "Set use_v2_api=False to use the V3 API endpoint.") - ex = ApiException(status=0, reason=message) - ex.message = message - ex.args = (message,) - return ex - if not use_v2_api and exc.status == HTTPStatus.METHOD_NOT_ALLOWED: - message = ("Server doesn't support the V3 API endpoint (/api/v3/write_lp). " - "Set use_v2_api=True to use the V2 API endpoint.") - ex = ApiException(status=0, reason=message) - ex.message = message - ex.args = (message,) - return ex - partial = InfluxDBPartialWriteError.from_response(exc.response) - if partial is not None: - return partial - return exc - def _should_gzip(self, payload: str, enable_gzip: bool = False, gzip_threshold: int = None) -> bool: """ Determines whether gzip compression should be applied to the given payload based diff --git a/influxdb_client_3/write_client/write_exceptions.py b/influxdb_client_3/write_client/write_exceptions.py deleted file mode 100644 index 4ac437c2..00000000 --- a/influxdb_client_3/write_client/write_exceptions.py +++ /dev/null @@ -1,37 +0,0 @@ -# coding: utf-8 - -from __future__ import absolute_import - -from influxdb_client_3.exceptions import InfluxDBError - -_UTF_8_encoding = 'utf-8' - - -class ApiException(InfluxDBError): - - def __init__(self, status=None, reason=None, http_resp=None): - """Initialize with HTTP response.""" - super().__init__(response=http_resp) - if http_resp: - self.status = http_resp.status - self.reason = http_resp.reason - self.body = http_resp.data - self.headers = http_resp.getheaders() - else: - self.status = status - self.reason = reason - self.body = None - self.headers = None - - def __str__(self): - """Get custom error messages for exception.""" - error_message = "({0})\n" \ - "Reason: {1}\n".format(self.status, self.reason) - if self.headers: - error_message += "HTTP response headers: {0}\n".format( - self.headers) - - if self.body: - error_message += "HTTP response body: {0}\n".format(self.body) - - return error_message diff --git a/tests/test_cli.py b/tests/test_cli.py index c07e0717..37a398e6 100644 --- a/tests/test_cli.py +++ b/tests/test_cli.py @@ -6,7 +6,7 @@ import pyarrow as pa from influxdb_client_3.cli import _coerce_timestamps, _format_table, _run_query, build_parser, main -from influxdb_client_3.exceptions import InfluxDB3ClientQueryError +from influxdb_client_3.exceptions import InfluxDB3ClientQueryException def _args(**overrides): @@ -119,7 +119,7 @@ def test_run_query_reads_query_from_file(tmp_path): def test_run_query_writes_json_error_for_query_exception(): args = _args(query="SELECT bad", host="http://localhost:8181", database="db1") - mock_client = _mock_client(side_effect=InfluxDB3ClientQueryError("bad query")) + mock_client = _mock_client(side_effect=InfluxDB3ClientQueryException("bad query")) stdout, stderr = io.StringIO(), io.StringIO() diff --git a/tests/test_influxdb_client_3.py b/tests/test_influxdb_client_3.py index 1b9f2f5b..a8660640 100644 --- a/tests/test_influxdb_client_3.py +++ b/tests/test_influxdb_client_3.py @@ -8,9 +8,9 @@ from influxdb_client_3 import InfluxDBClient3, WritePrecision, DefaultWriteOptions, Point, WriteOptions, WriteType, \ write_client_options -from influxdb_client_3.exceptions import InfluxDB3ClientQueryError +from influxdb_client_3.exceptions import InfluxDB3ClientQueryException +from influxdb_client_3.exceptions.exceptions import InfluxDBRestClientException from influxdb_client_3.write_client.client.write_api import _BatchItemKey -from influxdb_client_3.write_client.write_exceptions import ApiException from tests.util import asyncio_run from tests.util.mocks import ConstantFlightServer, ConstantData, ErrorFlightServer @@ -547,7 +547,7 @@ def test_disable_grpc_compression_default_is_false(self): def test_query_with_arrow_error(self): f = ErrorFlightServer() with InfluxDBClient3(f"http://localhost:{f.port}", "my_org", "my_db", "my_token") as c: - with self.assertRaises(InfluxDB3ClientQueryError) as err: + with self.assertRaises(InfluxDB3ClientQueryException) as err: c.query("SELECT * FROM my_data") self.assertIn("Error while executing query", str(err.exception)) @@ -555,7 +555,7 @@ def test_query_with_arrow_error(self): async def test_async_query_with_arrow_error(self): f = ErrorFlightServer() with InfluxDBClient3(f"http://localhost:{f.port}", "my_org", "my_db", "my_token") as c: - with self.assertRaises(InfluxDB3ClientQueryError) as err: + with self.assertRaises(InfluxDB3ClientQueryException) as err: await c.query_async("SELECT * FROM my_data") self.assertIn("Error while executing query", str(err.exception)) @@ -597,7 +597,7 @@ def test_get_version_fail(self): response_json={"error": "error"}, status=400 ) - with self.assertRaises(ApiException): + with self.assertRaises(InfluxDBRestClientException): InfluxDBClient3( host=f'http://{server.host}:{server.port}', org="ORG", database="DB", token="TOKEN" ).get_server_version() diff --git a/tests/test_influxdb_client_3_integration.py b/tests/test_influxdb_client_3_integration.py index 102f2804..27d842c0 100644 --- a/tests/test_influxdb_client_3_integration.py +++ b/tests/test_influxdb_client_3_integration.py @@ -15,12 +15,11 @@ from urllib3.exceptions import MaxRetryError, TimeoutError as Url3TimeoutError from influxdb_client_3 import InfluxDBClient3, write_client_options, WriteOptions, \ - WriteType, InfluxDB3ClientQueryError -from influxdb_client_3.exceptions import InfluxDBError, InfluxDBPartialWriteError + WriteType, InfluxDB3ClientQueryException +from influxdb_client_3.exceptions.write_exceptions import InfluxDBWriteException, InfluxDBPartialWriteException from influxdb_client_3.write_client import WriteApi from influxdb_client_3.write_client._sync import rest_client from influxdb_client_3.write_client.client.util.multiprocessing_helper import MultiprocessingWriter -from influxdb_client_3.write_client.write_exceptions import ApiException from tests.util import asyncio_run, lp_to_py_object @@ -186,7 +185,7 @@ def test_v3_error(self): )) ) as client: if accept_partial: - with self.assertRaises(InfluxDBPartialWriteError) as err: + with self.assertRaises(InfluxDBPartialWriteException) as err: client.write(lp) self.assertEqual(1, len(err.exception.line_errors)) @@ -195,11 +194,11 @@ def test_v3_error(self): self.assertIn("invalid column type for column 'temp'", line_error.error_message) self.assertIn("home,room=Sunroom", line_error.original_line) else: - with self.assertRaises(ApiException) as err: + with self.assertRaises(InfluxDBWriteException) as err: client.write(lp) self.assertEqual(400, err.exception.status) - self.assertEqual("line protocol parsing error", err.exception.message) + self.assertIn("line protocol parsing error", err.exception.message) body = json.loads(err.exception.body) self.assertEqual("line protocol parsing error", body["error"]) self.assertEqual(2, body["data"]["line_number"]) @@ -226,11 +225,11 @@ def test_v2_error(self): accept_partial=accept_partial )) ) as client: - with self.assertRaises(ApiException) as err: + with self.assertRaises(InfluxDBWriteException) as err: client.write(lp) self.assertEqual(400, err.exception.status) - self.assertNotIsInstance(err.exception, InfluxDBPartialWriteError) + self.assertNotIsInstance(err.exception, InfluxDBPartialWriteException) body = json.loads(err.exception.body) self.assertEqual("invalid", body["code"]) self.assertIn("write buffer error", body["message"]) @@ -238,7 +237,7 @@ def test_v2_error(self): def test_auth_error_token(self): self.client = InfluxDBClient3(host=self.host, database=self.database, token='fake token') test_id = time.time_ns() - with self.assertRaises(InfluxDBError) as err: + with self.assertRaises(InfluxDBWriteException) as err: self.client.write(f"integration_test_python,type=used value=123.0,test_id={test_id}i") self.assertEqual('Authorization header was malformed, the request was not in the form of ' '\'Authorization: \', supported auth-schemes are Bearer, Token and Basic', @@ -247,7 +246,7 @@ def test_auth_error_token(self): def test_auth_error_auth_scheme(self): self.client = InfluxDBClient3(host=self.host, database=self.database, token=self.token, auth_scheme='Any') test_id = time.time_ns() - with self.assertRaises(InfluxDBError) as err: + with self.assertRaises(InfluxDBWriteException) as err: self.client.write(f"integration_test_python,type=used value=123.0,test_id={test_id}i") self.assertEqual('Authorization header was malformed, the request was not in the form of ' '\'Authorization: \', supported auth-schemes are Bearer, Token and Basic', @@ -267,7 +266,7 @@ def success(conf, data): write_success = True write_count += 1 - def error(conf, data, exception: InfluxDBError): + def error(conf, data, exception: InfluxDBWriteException): nonlocal write_error write_error = True @@ -669,7 +668,7 @@ def test_query_timeout(self): query_timeout=1, ) - with self.assertRaisesRegex(InfluxDB3ClientQueryError, ".*Deadline Exceeded.*"): + with self.assertRaisesRegex(InfluxDB3ClientQueryException, ".*Deadline Exceeded.*"): localClient.query("SELECT * FROM data") def test_write_timeout_per_call_override(self): diff --git a/tests/test_write_api.py b/tests/test_write_api.py index 2f3335a0..134fc0f5 100644 --- a/tests/test_write_api.py +++ b/tests/test_write_api.py @@ -1,22 +1,293 @@ import asyncio +import http import json import unittest import uuid +from dataclasses import dataclass, field +from typing import Optional, List from unittest import mock import pytest from urllib3 import response -from urllib3.exceptions import ConnectTimeoutError +from urllib3.exceptions import ConnectTimeoutError, ProtocolError, SSLError -from influxdb_client_3 import InfluxDBClient3, InfluxDBError -from influxdb_client_3.exceptions import InfluxDBPartialWriteError +from influxdb_client_3 import InfluxDBClient3 +from influxdb_client_3.exceptions.write_exceptions import ( + InfluxDBWriteException, + InfluxDBPartialWriteLineException, + InfluxDBPartialWriteException, + translate_write_exception, +) from influxdb_client_3.version import VERSION -from influxdb_client_3.write_client.write_exceptions import ApiException +from influxdb_client_3.write_client.client.write.retry import WritesRetry _package = "influxdb3-python" _sentHeaders = {} +@dataclass +class TestCase: + __test__ = False + name: str + status_code: int + response_body: str + content_type: Optional[str] = None + use_v2_api: bool = False + accept_partial: bool = False + expected_msg: str = "" + expect_partial: bool = False + expected_lines: List[InfluxDBPartialWriteLineException] = field(default_factory=list) + + def __str__(self): + return self.name + + +POINTS = ( + "home,room=Sunroom temp=96 1735545600\n" + 'home,room=Sunroom temp="hi" 1735545610\n' + "home,room=Sunroom temp=88i 1735545620" +) +REJECTED_LINE = 'home,room=Sunroom temp="hi" 1735545610' +REJECTED_LINE_JSON = 'home,room=Sunroom temp=\\"hi\\" 1735545610' +LINE_ERROR = ( + "invalid column type for column 'temp', expected " + "iox::column_type::field::float, got iox::column_type::field::string" +) + +TEST_CASES = [ + TestCase( + name="V3 accept partial with renamed error and non-empty array", + status_code=http.client.BAD_REQUEST, + content_type="application/json", + response_body=( + f'{{"error":"write completed with rejected rows","data":[' + f'{{"error_message":"{LINE_ERROR}","line_number":2,"original_line":"{REJECTED_LINE_JSON}"}}' + f"]}}" + ), + accept_partial=True, + expected_msg=f"write completed with rejected rows:\n\tline 2: {LINE_ERROR} ({REJECTED_LINE})", + expect_partial=True, + expected_lines=[ + InfluxDBPartialWriteLineException( + error_message=LINE_ERROR, + line_number=2, + original_line=REJECTED_LINE, + ) + ], + ), + TestCase( + name="V3 accept partial without content type", + status_code=http.client.BAD_REQUEST, + content_type=None, + response_body=( + f'{{"error":"write completed with rejected rows","data":[' + f'{{"error_message":"{LINE_ERROR}","line_number":2,"original_line":"{REJECTED_LINE_JSON}"}}' + f"]}}" + ), + accept_partial=True, + expected_msg=f"write completed with rejected rows:\n\tline 2: {LINE_ERROR} ({REJECTED_LINE})", + expect_partial=True, + expected_lines=[ + InfluxDBPartialWriteLineException( + error_message=LINE_ERROR, + line_number=2, + original_line=REJECTED_LINE, + ) + ], + ), + TestCase( + name="V3 accept partial with malformed non-empty array", + status_code=http.client.BAD_REQUEST, + content_type="application/json", + response_body=( + f'{{"error":"write completed with rejected rows","data":[' + f'{{"line_number":"invalid","original_line":"{REJECTED_LINE_JSON}"}}' + f"]}}" + ), + accept_partial=True, + expected_msg=f'write completed with rejected rows:\n\t{{"line_number":"invalid","original_line":' + f'"{REJECTED_LINE_JSON}"}}', + expect_partial=True, + expected_lines=[], + ), + TestCase( + name="V3 accept partial with mixed primitive and typed entries", + status_code=http.client.BAD_REQUEST, + content_type="application/json", + response_body=( + f'{{"error":"write completed with rejected rows","data":[' + f'1,{{"error_message":"{LINE_ERROR}","line_number":2,"original_line":"{REJECTED_LINE_JSON}"}}' + f"]}}" + ), + accept_partial=True, + expected_msg=( + f"write completed with rejected rows:\n\t1\n\t" + f'{{"error_message":"{LINE_ERROR}","line_number":2,"original_line":"{REJECTED_LINE_JSON}"}}' + ), + expect_partial=True, + expected_lines=[ + InfluxDBPartialWriteLineException( + error_message=LINE_ERROR, + line_number=2, + original_line=REJECTED_LINE, + ) + ], + ), + TestCase( + name="V3 accept partial with string entries", + status_code=http.client.BAD_REQUEST, + content_type="application/json", + response_body=f'{{"error":"write completed with rejected rows","data":["{REJECTED_LINE_JSON}"]}}', + accept_partial=True, + expected_msg=f'write completed with rejected rows:\n\t"{REJECTED_LINE_JSON}"', + expect_partial=True, + expected_lines=[], + ), + TestCase( + name="V3 accept partial with error message only", + status_code=http.client.BAD_REQUEST, + content_type="application/json", + response_body=f'{{"error":"write completed with rejected rows",' + f'"data":[{{"error_message":"{LINE_ERROR}"}}]}}', + accept_partial=True, + expected_msg=f"write completed with rejected rows:\n\t{LINE_ERROR}", + expect_partial=True, + expected_lines=[ + InfluxDBPartialWriteLineException( + error_message=LINE_ERROR, + line_number=None, + original_line=None, + ) + ], + ), + TestCase( + name="V3 accept partial with line number but no original line", + status_code=http.client.BAD_REQUEST, + content_type="application/json", + response_body=f'{{"error":"write completed with rejected rows","data":[{{"error_message":"{LINE_ERROR}",' + f'"line_number":2}}]}}', + accept_partial=True, + expected_msg=f"write completed with rejected rows:\n\tline 2: {LINE_ERROR}", + expect_partial=True, + expected_lines=[ + InfluxDBPartialWriteLineException( + error_message=LINE_ERROR, + line_number=2, + original_line=None, + ) + ], + ), + TestCase( + name="V3 accept partial with entry missing error message", + status_code=http.client.BAD_REQUEST, + content_type="application/json", + response_body=f'{{"error":"write completed with rejected rows",' + f'"data":[{{"line_number":2,"original_line":"{REJECTED_LINE_JSON}"}}]}}', + accept_partial=True, + expected_msg=f'write completed with rejected rows:\n\t{{"line_number":2,' + f'"original_line":"{REJECTED_LINE_JSON}"}}', + expect_partial=True, + expected_lines=[], + ), + TestCase( + name="V3 accept partial with empty array", + status_code=http.client.BAD_REQUEST, + content_type="application/json", + response_body='{"error":"write failed","data":[]}', + accept_partial=True, + expected_msg="write failed", + expect_partial=False, + ), + TestCase( + name="V3 accept partial with object details remains generic", + status_code=http.client.BAD_REQUEST, + content_type="application/json", + response_body=( + f'{{"error":"line protocol parsing error","data":' + f'{{"error_message":"{LINE_ERROR}","line_number":2,"original_line":"{REJECTED_LINE_JSON}"}}}}' + ), + accept_partial=True, + expected_msg=f"line protocol parsing error:\n\tline 2: {LINE_ERROR} ({REJECTED_LINE})", + expect_partial=False, + ), + TestCase( + name="V3 reject partial with object details", + status_code=http.client.BAD_REQUEST, + content_type="application/json", + response_body=( + f'{{"error":"line protocol parsing error","data":' + f'{{"error_message":"{LINE_ERROR}","line_number":2,"original_line":"{REJECTED_LINE_JSON}"}}}}' + ), + accept_partial=False, + expected_msg=f"line protocol parsing error:\n\tline 2: {LINE_ERROR} ({REJECTED_LINE})", + expect_partial=False, + ), + TestCase( + name="V2 never returns partial write error", + status_code=http.client.BAD_REQUEST, + content_type="application/json", + response_body=( + f'{{"error":"partial write of line protocol occurred","data":[' + f'{{"error_message":"{LINE_ERROR}","line_number":2,"original_line":"{REJECTED_LINE_JSON}"}}' + f"]}}" + ), + use_v2_api=True, + accept_partial=True, + expected_msg="partial write of line protocol occurred", + expect_partial=False, + ), + TestCase( + name="V3 non-400 never returns partial write error", + status_code=http.client.INTERNAL_SERVER_ERROR, + content_type="application/json", + response_body=( + f'{{"error":"partial write of line protocol occurred","data":[' + f'{{"error_message":"{LINE_ERROR}","line_number":2,"original_line":"{REJECTED_LINE_JSON}"}}' + f"]}}" + ), + accept_partial=True, + expected_msg="partial write of line protocol occurred", + expect_partial=False, + ), + TestCase( + name="V3 scalar data remains generic", + status_code=http.client.BAD_REQUEST, + content_type="application/json", + response_body='{"error":"write failed","data":"invalid"}', + accept_partial=True, + expected_msg="write failed", + expect_partial=False, + ), + TestCase( + name="V3 empty object data remains generic", + status_code=http.client.BAD_REQUEST, + content_type="application/json", + response_body='{"error":"write failed","data":{}}', + accept_partial=True, + expected_msg="write failed", + expect_partial=False, + ), + TestCase( + name="V3 null data remains generic", + status_code=http.client.BAD_REQUEST, + content_type="application/json", + response_body='{"error":"write failed","data":null}', + accept_partial=True, + expected_msg="write failed", + expect_partial=False, + ), + TestCase( + name="V3 malformed JSON preserves raw response", + status_code=http.client.BAD_REQUEST, + content_type="application/json", + response_body='{"error":"write failed"', + accept_partial=True, + expected_msg='{"error":"write failed"', + expect_partial=False, + ), +] + + class WriteApiTests(unittest.TestCase): received_timeout_total = None @@ -29,18 +300,22 @@ def mock_urllib3_timeout_request(method, return response.HTTPResponse(status=200, version=4, reason="OK", decode_content=False, request_url=url) - def _test_api_error(self, body): + def _test_api_error(self, body, header=None, accept_partial=None, use_v2_api=None): client = InfluxDBClient3( host='http://localhost:8181', token='my-token', database='my-bucket', org='my-org' ) + if body is not None: + body = body.encode() + client._write_api.rest_client.pool_manager.request \ = mock.Mock(return_value=response.HTTPResponse(status=400, + headers=header or {}, reason='Bad Request', - body=body.encode())) - client._write_api.write(record="data,foo=bar val=3.14") + body=body)) + client._write_api.write(record="data,foo=bar val=3.14", accept_partial=accept_partial, use_v2_api=use_v2_api) def test_default_headers(self): client = InfluxDBClient3( @@ -190,27 +465,27 @@ def test_build_write_request_requires_org_and_bucket(self): def test_api_error_cloud(self): response_body = '{"message": "parsing failed for write_lp endpoint"}' - with self.assertRaises(InfluxDBError) as err: + with self.assertRaises(InfluxDBWriteException) as err: self._test_api_error(response_body) self.assertEqual('parsing failed for write_lp endpoint', err.exception.message) def test_api_error_oss_without_detail(self): response_body = '{"error": "parsing failed for write_lp endpoint"}' - with self.assertRaises(InfluxDBError) as err: + with self.assertRaises(InfluxDBWriteException) as err: self._test_api_error(response_body) self.assertEqual('parsing failed for write_lp endpoint', err.exception.message) def test_api_error_oss_with_detail(self): response_body = ('{"error":"parsing failed for write_lp endpoint","data":{"error_message":"invalid field value ' 'in line protocol for field \'val\' on line 1"}}') - with self.assertRaises(InfluxDBError) as err: + with self.assertRaises(InfluxDBWriteException) as err: self._test_api_error(response_body) self.assertEqual("parsing failed for write_lp endpoint:\n\tinvalid field value in line protocol for field " "'val' on line 1", err.exception.message) def test_api_error_unknown(self): response_body = '{"detail":"no info"}' - with self.assertRaises(InfluxDBError) as err: + with self.assertRaises(InfluxDBWriteException) as err: self._test_api_error(response_body) self.assertEqual(response_body, err.exception.message) @@ -231,6 +506,8 @@ def test_api_error_v3_with_detail(self): "\tline 3: invalid column type for column 'v', expected iox::column_type::field::float, " "got iox::column_type::field::uinteger (***.INF.remote_***)", True, + False, + 1 ), # error_message only (no line_number/original_line) ( @@ -240,6 +517,8 @@ def test_api_error_v3_with_detail(self): "partial write of line protocol occurred:\n" "\tonly error message", True, + False, + 1 ), # non-dict item in data list is skipped ( @@ -247,8 +526,10 @@ def test_api_error_v3_with_detail(self): '{"error":"partial write of line protocol occurred","data":[null,' '{"error_message":"bad line","line_number":2,"original_line":"bad lp"}]}', "partial write of line protocol occurred:\n" - "\tline 2: bad line (bad lp)", + "\t{\"error_message\":\"bad line\",\"line_number\":2,\"original_line\":\"bad lp\"}", True, + False, + 1 ), # details empty -> return error_text ( @@ -256,7 +537,9 @@ def test_api_error_v3_with_detail(self): '{"error":"partial write of line protocol occurred","data":[{"line_number":2}]}', "partial write of line protocol occurred:\n" "\t{\"line_number\":2}", + True, False, + 0 ), # typed parse fails due line_number type -> raw fallback details ( @@ -265,7 +548,9 @@ def test_api_error_v3_with_detail(self): '[{"error_message":"bad line","line_number":"x","original_line":"bad lp"}]}', "partial write of line protocol occurred:\n" "\t{\"error_message\":\"bad line\",\"line_number\":\"x\",\"original_line\":\"bad lp\"}", + True, False, + 0 ), # mixed valid + malformed in array -> raw fallback for whole array ( @@ -275,7 +560,9 @@ def test_api_error_v3_with_detail(self): "partial write of line protocol occurred:\n" "\t{\"error_message\":\"bad line\",\"line_number\":2,\"original_line\":\"bad lp\"}\n" "\t1", + True, False, + 1 ), # data is not a dict when resolving fallback keys ( @@ -283,14 +570,18 @@ def test_api_error_v3_with_detail(self): '{"error":"data not list","data":"oops"}', "data not list", False, + True, + 0 ), # typed object with empty message is dropped ( "empty error_message in object", - '{"error":"partial write of line protocol occurred","data":' + '{"error":"parsing failed for write_lp endpoint","data":' '{"error_message":"","line_number":2,"original_line":"bad lp"}}', - "partial write of line protocol occurred", + "parsing failed for write_lp endpoint", False, + True, + 0 ), # typed array parse fails, raw fallback skips null item ( @@ -299,64 +590,87 @@ def test_api_error_v3_with_detail(self): '[null,{"error_message":123}]}', "partial write of line protocol occurred:\n" "\t{\"error_message\":123}", + True, False, + 1 ), ] - for name, response_body, expected, is_partial in cases: + for name, response_body, expected, is_partial, use_v2_api, expected_line_error_count in cases: with self.subTest(name): - with self.assertRaises(InfluxDBError) as err: - self._test_api_error(response_body) - self.assertEqual(expected, err.exception.message) if is_partial: - self.assertIsInstance(err.exception, InfluxDBPartialWriteError) - self.assertGreaterEqual(len(err.exception.line_errors), 1) + with self.assertRaises(InfluxDBPartialWriteException) as err: + self._test_api_error(body=response_body, accept_partial=is_partial, use_v2_api=use_v2_api) + self.assertIsInstance(err.exception, InfluxDBPartialWriteException) + self.assertGreaterEqual(len(err.exception.line_errors), expected_line_error_count) else: - self.assertNotIsInstance(err.exception, InfluxDBPartialWriteError) + with self.assertRaises(InfluxDBWriteException) as err: + self._test_api_error(body=response_body, accept_partial=is_partial, use_v2_api=use_v2_api) + self.assertEqual(expected, err.exception.message) - def test_api_error_v3_parsing_failed_object_returns_partial_error(self): + def test_api_error_v3_parsing_failed_object_returns_error(self): response_body = ('{"error":"parsing failed for write_lp endpoint","data":' '{"error_message":"invalid field value","line_number":2,"original_line":"m,t=a f=bad"}}') - with self.assertRaises(InfluxDBPartialWriteError) as err: + with self.assertRaises(InfluxDBWriteException) as err: self._test_api_error(response_body) - self.assertEqual(1, len(err.exception.line_errors)) - self.assertEqual(2, err.exception.line_errors[0].line_number) + self.assertEqual('parsing failed for write_lp endpoint:\n\tline 2: invalid field value (m,t=a f=bad)', + err.exception.message) - def test_api_error_v3_partial_write_with_message_only_object_returns_partial_error(self): - response_body = ('{"error":"partial write of line protocol occurred","data":' + def test_api_error_v3_write_with_message_only_object_returns(self): + response_body = ('{"error":"parsing failed for write_lp endpoint","data":' '{"error_message":"only error message"}}') - with self.assertRaises(InfluxDBPartialWriteError) as err: + with self.assertRaises(InfluxDBWriteException) as err: self._test_api_error(response_body) - self.assertEqual(1, len(err.exception.line_errors)) - self.assertEqual(0, err.exception.line_errors[0].line_number) - self.assertEqual("", err.exception.line_errors[0].original_line) + self.assertEqual("parsing failed for write_lp endpoint:\n\tonly error message", err.exception.message) - def test_api_error_v3_partial_write_with_line_number_without_original_line(self): - response_body = ('{"error":"partial write of line protocol occurred","data":' + def test_api_error_v3_write_with_line_number_without_original_line(self): + response_body = ('{"error":"parsing failed for write_lp endpoint","data":' '{"error_message":"invalid field value","line_number":2}}') - with self.assertRaises(InfluxDBPartialWriteError) as err: + with self.assertRaises(InfluxDBWriteException) as err: self._test_api_error(response_body) - self.assertEqual(1, len(err.exception.line_errors)) - self.assertEqual("partial write of line protocol occurred:\n\tline 2: invalid field value", + self.assertEqual("parsing failed for write_lp endpoint:\n\tline 2: invalid field value", err.exception.message) - def test_partial_write_from_response_guards(self): - self.assertIsNone(InfluxDBPartialWriteError.from_response(None)) - - empty_body = response.HTTPResponse(status=400, reason="Bad Request", body=b"") - self.assertIsNone(InfluxDBPartialWriteError.from_response(empty_body)) + def test_api_error_v3_write_with_invalid_line_number(self): + response_body = ('{"error":"parsing failed for write_lp endpoint","data":' + '{"error_message":"bad line","line_number":"aa"}}') + with self.assertRaises(InfluxDBWriteException) as err: + self._test_api_error(response_body) + self.assertEqual("parsing failed for write_lp endpoint:\n\tbad line", + err.exception.message) - invalid_json = response.HTTPResponse(status=400, reason="Bad Request", body=b"{") - self.assertIsNone(InfluxDBPartialWriteError.from_response(invalid_json)) + def test_fallback_header_or_body(self): + for body in ["{err", "[]", "{}"]: + for is_partial_write in [False, True]: + # Fallback to header message + with self.assertRaises(InfluxDBWriteException) as err: + header = {"X-Influx-Error": "not used"} + self._test_api_error( + body=body, + header=header, + accept_partial=is_partial_write, + use_v2_api=False + ) + self.assertEqual(header["X-Influx-Error"], err.exception.message) - non_dict_json = response.HTTPResponse(status=400, reason="Bad Request", body=b"[]") - self.assertIsNone(InfluxDBPartialWriteError.from_response(non_dict_json)) + # Fallback to raw body + with self.assertRaises(InfluxDBWriteException) as err: + self._test_api_error( + body=body, + accept_partial=is_partial_write, + use_v2_api=False + ) + self.assertEqual(body, err.exception.message) - object_without_typed_line_error = response.HTTPResponse( - status=400, - reason="Bad Request", - body=b'{"error":"partial write of line protocol occurred","data":{"error_message":123}}', - ) - self.assertIsNone(InfluxDBPartialWriteError.from_response(object_without_typed_line_error)) + def test_fallback_status_code_msg(self): + for body in ["", None]: + for is_partial_write in [False, True]: + with self.assertRaises(InfluxDBWriteException) as err: + self._test_api_error( + body=body, + accept_partial=is_partial_write, + use_v2_api=False + ) + self.assertEqual('Bad Request', err.exception.message) def test_api_error_headers(self): body = '{"error": "test error"}' @@ -384,7 +698,7 @@ def test_api_error_headers(self): body=body.encode() ) ) - with self.assertRaises(InfluxDBError) as err: + with self.assertRaises(InfluxDBWriteException) as err: client._write_api.write("TEST_BUCKET", "TEST_ORG", "data,foo=bar val=3.14") self.assertEqual(body_dic['error'], err.exception.message) headers = err.exception.getheaders() @@ -466,22 +780,25 @@ def test_post_write_async_translates_exceptions(self): ( "v2 on v3-only backend", True, + False, response.HTTPResponse(status=405, reason="Method Not Allowed", body=b""), - ApiException, + InfluxDBWriteException, "Server doesn't support the V2 API endpoint (/api/v2/write). " "Set use_v2_api=False to use the V3 API endpoint.", ), ( "v3 on v2-only backend", False, + False, response.HTTPResponse(status=405, reason="Method Not Allowed", body=b""), - ApiException, + InfluxDBWriteException, "Server doesn't support the V3 API endpoint (/api/v3/write_lp). " "Set use_v2_api=True to use the V2 API endpoint.", ), ( "v3 partial write response", False, + True, response.HTTPResponse( status=400, reason="Bad Request", @@ -490,11 +807,11 @@ def test_post_write_async_translates_exceptions(self): b'"line_number":2,"original_line":"home,room=Sunroom temp=\\"hi\\" 1735549200"}]}' ), ), - InfluxDBPartialWriteError, + InfluxDBPartialWriteException, None, ), ] - for name, use_v2_api, http_resp, expected_type, expected_message in cases: + for name, use_v2_api, accept_partial, http_resp, expected_type, expected_message in cases: with self.subTest(name): client = InfluxDBClient3( host='http://localhost:8181', @@ -504,14 +821,14 @@ def test_post_write_async_translates_exceptions(self): ) write_api = client._write_api write_api.rest_client.request = mock.Mock( - side_effect=ApiException(http_resp=http_resp) + side_effect=InfluxDBWriteException(http_resp=http_resp) ) result = write_api._post_write( org="TEST_ORG", bucket="TEST_BUCKET", body="home,room=Sunroom temp=96 1735545600", precision='s', - accept_partial=False, + accept_partial=accept_partial, no_sync=False, async_req=True, _async_req=True, @@ -614,7 +931,7 @@ def test_post_write_async_translates_v3_unsupported(self): write_api = client._write_api write_api.rest_client.request = mock.Mock( - side_effect=ApiException( + side_effect=InfluxDBWriteException( http_resp=response.HTTPResponse(status=405, reason="Method Not Allowed", body=b"") ) ) @@ -627,9 +944,149 @@ async def run(): use_v2_api=False, ) - with self.assertRaises(ApiException) as err: + with self.assertRaises(InfluxDBWriteException) as err: asyncio.run(run()) expected = ("Server doesn't support the V3 API endpoint (/api/v3/write_lp). " "Set use_v2_api=True to use the V2 API endpoint.") self.assertEqual(expected, err.exception.message) + + def test_write_error_classification(self): + for tc in TEST_CASES: + with self.subTest(tc.name): + headers = {"Content-Type": tc.content_type} if tc.content_type is not None else {} + + client = InfluxDBClient3( + host="http://localhost:8086", + token="token", + database="database", + ) + client._write_api.rest_client.pool_manager.request = mock.Mock( + return_value=response.HTTPResponse( + status=tc.status_code, + headers=headers, + body=tc.response_body.encode("utf-8"), + ) + ) + + expected_exc = InfluxDBPartialWriteException if tc.expect_partial else InfluxDBWriteException + with self.assertRaises(expected_exc) as cm: + client.write( + record=POINTS, + use_v2_api=tc.use_v2_api, + accept_partial=tc.accept_partial, + ) + + err = cm.exception + self.assertEqual(tc.expected_msg, err.message) + if tc.expect_partial: + self.assertEqual(tc.expected_lines, err.line_errors) + + def test_increment_with_response(self): + callback = mock.MagicMock() + retry = WritesRetry(total=3, retry_callback=callback) + + mock_response = response.HTTPResponse( + body=b'{"code": "unavailable", "message": "service unavailable"}', + status=503, + headers={"Retry-After": "5"} + ) + + new_retry = retry.increment( + method="POST", + url="http://localhost:8086/api/v2/write", + response=mock_response + ) + + self.assertEqual(new_retry.total, 2) + self.assertEqual(len(new_retry.history), 1) + + callback.assert_called_once() + called_arg = callback.call_args[0][0] + self.assertIsInstance(called_arg, InfluxDBWriteException) + self.assertEqual(called_arg.response, mock_response) + + def test_increment_with_error(self): + callback = mock.MagicMock() + retry = WritesRetry(total=3, retry_callback=callback) + + err = ProtocolError("Connection reset by peer") + with self.assertRaises(ProtocolError) as e: + retry.increment( + method="POST", + url="http://localhost:8086/api/v2/write", + error=err + ) + self.assertEqual(str(e.exception), 'Connection reset by peer') + + def test_increment_with_no_response_and_no_error(self): + callback = mock.MagicMock() + retry = WritesRetry(total=3, retry_callback=callback) + + url = "http://localhost:8086/api/v2/write" + new_retry = retry.increment(method="POST", url=url) + + self.assertEqual(new_retry.total, 2) + callback.assert_called_once_with(f'Failed request to: {url}') + + def test_translate_exception_ssl_error(self): + ssl_error = SSLError("certificate verify failed: self-signed certificate") + msg = "{0}\n{1}".format(type(ssl_error).__name__, str(ssl_error)) + + e = translate_write_exception(InfluxDBWriteException(status=0, reason=msg)) + self.assertIsInstance(e, InfluxDBWriteException) + self.assertEqual(e.reason, msg) + self.assertEqual(e.status, 0) + + def test_is_partial_write_error(self): + valid_root = {"error": "partial write of line protocol occurred", "data": [{"line": 1}]} + + # Positive cases + self.assertTrue( + InfluxDBPartialWriteException.is_partial_write_error( + http.HTTPStatus.BAD_REQUEST, + False, + True, + valid_root)) + + # Negative status codes + for status in [200, 204, 401, 403, 404, 500, None]: + with self.subTest(status=status): + self.assertFalse(InfluxDBPartialWriteException.is_partial_write_error( + status, + False, + True, + valid_root)) + + with self.subTest(use_v2_api=True): + self.assertFalse(InfluxDBPartialWriteException.is_partial_write_error( + 400, + True, + True, + valid_root)) + + with self.subTest(accept_partial=False): + self.assertFalse(InfluxDBPartialWriteException.is_partial_write_error( + 400, + False, + False, + valid_root)) + + # Negative root shapes + invalid_roots = [ + {}, + {"error": None, "data": [{"line": 1}]}, + {"data": [{"line": 1}]}, + {"error": "some error"}, + {"error": "some error", "data": []}, + {"error": "some error", "data": None}, + {"error": "some error", "data": "invalid"}, + {"error": "some error", "data": {"line": 1}}, + ] + for root in invalid_roots: + with self.subTest(root=root): + self.assertFalse(InfluxDBPartialWriteException.is_partial_write_error( + 400, + False, + True, + root)) diff --git a/tests/test_write_exceptions.py b/tests/test_write_exceptions.py new file mode 100644 index 00000000..52846967 --- /dev/null +++ b/tests/test_write_exceptions.py @@ -0,0 +1,362 @@ +import http +import json +import unittest + +from influxdb_client_3.exceptions.exceptions import InfluxDBRestClientException +from influxdb_client_3.exceptions.write_exceptions import ( + translate_write_exception, InfluxDBWriteException, InfluxDBPartialWriteException, InfluxDBPartialWriteLineException +) + + +class DummyHttpResponse: + def __init__(self, status=200, reason="OK", data=None, headers=None): + self.status = status + self.reason = reason + self.data = data + self.headers = headers or {} + + def getheaders(self): + return self.headers + + def getheader(self, name, default=None): + return self.headers.get(name, default) + + +class TestWriteException(unittest.TestCase): + + def test_method_not_allowed_v3(self): + exc = InfluxDBRestClientException(status=http.HTTPStatus.METHOD_NOT_ALLOWED, reason="Method Not Allowed") + result = translate_write_exception(exc, use_v2_api=False) + + self.assertIsInstance(result, InfluxDBWriteException) + self.assertEqual(0, result.status) + expected_msg = ( + "Server doesn't support the V3 API endpoint (/api/v3/write_lp). " + "Set use_v2_api=True to use the V2 API endpoint." + ) + self.assertEqual(expected_msg, result.message) + self.assertEqual(expected_msg, result.reason) + self.assertEqual((expected_msg,), result.args) + + def test_method_not_allowed_v2(self): + exc = InfluxDBRestClientException(status=http.HTTPStatus.METHOD_NOT_ALLOWED, reason="Method Not Allowed") + result = translate_write_exception(exc, use_v2_api=True) + + self.assertIsInstance(result, InfluxDBWriteException) + self.assertEqual(0, result.status) + expected_msg = ( + "Server doesn't support the V2 API endpoint (/api/v2/write). " + "Set use_v2_api=False to use the V3 API endpoint." + ) + self.assertEqual(expected_msg, result.message) + self.assertEqual(expected_msg, result.reason) + self.assertEqual((expected_msg,), result.args) + + def test_status_zero_and_body_none(self): + exc = InfluxDBRestClientException(status=0, reason="Connection aborted") + result = translate_write_exception(exc) + + self.assertIsInstance(result, InfluxDBWriteException) + self.assertEqual(0, result.status) + self.assertEqual("Connection aborted", result.reason) + + def test_fallback_to_headers(self): + header_keys = [ + ("X-Platform-Error-Code", "platform_error_code"), + ("X-Influx-Error", "influx_error_header"), + ("X-InfluxDb-Error", "influxdb_error_header"), + ] + for header_key, header_val in header_keys: + with self.subTest(header_key=header_key): + http_resp = DummyHttpResponse( + status=500, + reason="Internal Server Error", + data=b"raw body text", + headers={header_key: header_val}, + ) + exc = InfluxDBRestClientException(http_resp=http_resp) + result = translate_write_exception(exc) + + self.assertIsInstance(result, InfluxDBWriteException) + self.assertEqual(header_val, result.message) + + def test_fallback_header_precedence(self): + http_resp = DummyHttpResponse( + status=500, + reason="Internal Server Error", + data=b"raw body text", + headers={ + "X-Platform-Error-Code": "platform_code", + "X-Influx-Error": "influx_error", + "X-InfluxDb-Error": "influxdb_error", + }, + ) + exc = InfluxDBRestClientException(http_resp=http_resp) + result = translate_write_exception(exc) + + self.assertIsInstance(result, InfluxDBWriteException) + self.assertEqual("platform_code", result.message) + + def test_fallback_to_raw_body_when_no_headers_and_invalid_json(self): + http_resp = DummyHttpResponse( + status=500, + reason="Internal Server Error", + data=b"raw plain text error", + headers={}, + ) + exc = InfluxDBRestClientException(http_resp=http_resp) + result = translate_write_exception(exc) + + self.assertIsInstance(result, InfluxDBWriteException) + self.assertEqual(b"raw plain text error", result.message) + + def test_fallback_to_status_reason_when_no_headers_and_empty_body(self): + for data in [None, ""]: + with self.subTest(data=data): + http_resp = DummyHttpResponse( + status=500, + reason="Internal Server Error", + data=data, + headers={}, + ) + exc = InfluxDBRestClientException(http_resp=http_resp) + result = translate_write_exception(exc) + + self.assertIsInstance(result, InfluxDBWriteException) + self.assertEqual("Internal Server Error", result.message) + + def test_fallback_when_json_is_not_dict_or_has_no_error_or_message(self): + cases = [ + b'""', + b'"just a string"', + b'123', + b'[1, 2, 3]', + b'{}', + b'{"status": "failed"}', + b'{"other_key": 42}', + ] + for body in cases: + with self.subTest(body=body): + http_resp = DummyHttpResponse( + status=500, + reason="Internal Server Error", + data=body, + headers={}, + ) + exc = InfluxDBRestClientException(http_resp=http_resp) + result = translate_write_exception(exc) + + self.assertIsInstance(result, InfluxDBWriteException) + self.assertEqual(body, result.message) + + def test_v3_message_only(self): + body = json.dumps({"message": "table 'cpu' not found"}).encode("utf-8") + http_resp = DummyHttpResponse(status=404, reason="Not Found", data=body) + exc = InfluxDBRestClientException(http_resp=http_resp) + result = translate_write_exception(exc) + + self.assertIsInstance(result, InfluxDBWriteException) + self.assertEqual("table 'cpu' not found", result.message) + + def test_v3_code_and_message(self): + body = json.dumps({"code": "not_found", "message": "table 'cpu' not found"}).encode("utf-8") + http_resp = DummyHttpResponse(status=404, reason="Not Found", data=body) + exc = InfluxDBRestClientException(http_resp=http_resp) + result = translate_write_exception(exc) + + self.assertIsInstance(result, InfluxDBWriteException) + self.assertEqual("not_found: table 'cpu' not found", result.message) + + def test_error_without_data(self): + body = json.dumps({"error": "syntax error on token"}).encode("utf-8") + http_resp = DummyHttpResponse(status=400, reason="Bad Request", data=body) + exc = InfluxDBRestClientException(http_resp=http_resp) + result = translate_write_exception(exc) + + self.assertIsInstance(result, InfluxDBWriteException) + self.assertEqual("syntax error on token", result.message) + + def test_object_data_error_without_line_number(self): + body = json.dumps({ + "error": "write failed", + "data": { + "error_message": "type conflict for field 'temp'" + } + }).encode("utf-8") + http_resp = DummyHttpResponse(status=400, reason="Bad Request", data=body) + exc = InfluxDBRestClientException(http_resp=http_resp) + result = translate_write_exception(exc) + + self.assertIsInstance(result, InfluxDBWriteException) + self.assertEqual("write failed:\n\ttype conflict for field 'temp'", result.message) + + def test_object_data_error_with_line_number_no_original_line(self): + body = json.dumps({ + "error": "write failed", + "data": { + "line_number": 4, + "error_message": "type conflict for field 'temp'" + } + }).encode("utf-8") + http_resp = DummyHttpResponse(status=400, reason="Bad Request", data=body) + exc = InfluxDBRestClientException(http_resp=http_resp) + result = translate_write_exception(exc) + + self.assertIsInstance(result, InfluxDBWriteException) + self.assertEqual("write failed:\n\tline 4: type conflict for field 'temp'", result.message) + + def test_object_data_error_with_line_number_and_original_line(self): + body = json.dumps({ + "error": "write failed", + "data": { + "line_number": 4, + "error_message": "type conflict for field 'temp'", + "original_line": "cpu,tag=1 temp=10 123456789" + } + }).encode("utf-8") + http_resp = DummyHttpResponse(status=400, reason="Bad Request", data=body) + exc = InfluxDBRestClientException(http_resp=http_resp) + result = translate_write_exception(exc) + + self.assertIsInstance(result, InfluxDBWriteException) + self.assertEqual( + "write failed:\n\tline 4: type conflict for field 'temp' (cpu,tag=1 temp=10 123456789)", + result.message + ) + + def test_object_data_error_without_error_message(self): + body = json.dumps({ + "error": "write failed", + "data": { + "line_number": 4 + } + }).encode("utf-8") + http_resp = DummyHttpResponse(status=400, reason="Bad Request", data=body) + exc = InfluxDBRestClientException(http_resp=http_resp) + result = translate_write_exception(exc) + + self.assertIsInstance(result, InfluxDBWriteException) + self.assertEqual("write failed", result.message) + + def test_partial_write_error_all_typed_details(self): + body = json.dumps({ + "error": "partial write of line protocol occurred", + "data": [ + { + "line_number": 1, + "error_message": "type mismatch", + "original_line": "m,t=a f=1" + }, + { + "line_number": 2, + "error_message": "invalid timestamp" + }, + { + "error_message": "general line error" + } + ] + }).encode("utf-8") + http_resp = DummyHttpResponse(status=400, reason="Bad Request", data=body) + exc = InfluxDBRestClientException(http_resp=http_resp) + result = translate_write_exception(exc, use_v2_api=False, accept_partial=True) + + self.assertIsInstance(result, InfluxDBPartialWriteException) + expected_msg = ( + "partial write of line protocol occurred:\n" + "\tline 1: type mismatch (m,t=a f=1)\n" + "\tline 2: invalid timestamp\n" + "\tgeneral line error" + ) + self.assertEqual(expected_msg, result.message) + self.assertEqual(3, len(result.line_errors)) + self.assertEqual( + [ + InfluxDBPartialWriteLineException(1, "type mismatch", "m,t=a f=1"), + InfluxDBPartialWriteLineException(2, "invalid timestamp", None), + InfluxDBPartialWriteLineException(None, "general line error", None), + ], + result.line_errors + ) + self.assertEqual(http_resp, result.response) + + def test_partial_write_error_untyped_details_fallback(self): + body = json.dumps({ + "error": "partial write of line protocol occurred", + "data": [ + { + "line_number": "not_an_int", + "error_message": "type mismatch", + "original_line": "m,t=a f=1" + }, + None, + "null", + {"line_number": 2}, + 123 + ] + }).encode("utf-8") + http_resp = DummyHttpResponse(status=400, reason="Bad Request", data=body) + exc = InfluxDBRestClientException(http_resp=http_resp) + result = translate_write_exception(exc, use_v2_api=False, accept_partial=True) + + self.assertIsInstance(result, InfluxDBPartialWriteException) + expected_msg = ( + "partial write of line protocol occurred:\n" + '\t{"line_number":"not_an_int","error_message":"type mismatch","original_line":"m,t=a f=1"}\n' + '\t{"line_number":2}\n' + '\t123' + ) + self.assertEqual(expected_msg, result.message) + + def test_partial_write_skipped_when_accept_partial_is_false(self): + body = json.dumps({ + "error": "partial write of line protocol occurred", + "data": [ + { + "line_number": 1, + "error_message": "type mismatch", + "original_line": "m,t=a f=1" + } + ] + }).encode("utf-8") + http_resp = DummyHttpResponse(status=400, reason="Bad Request", data=body) + exc = InfluxDBRestClientException(http_resp=http_resp) + result = translate_write_exception(exc, use_v2_api=False, accept_partial=False) + + self.assertIsInstance(result, InfluxDBWriteException) + self.assertEqual("partial write of line protocol occurred", result.message) + + def test_partial_write_skipped_when_use_v2_api_is_true(self): + body = json.dumps({ + "error": "partial write of line protocol occurred", + "data": [ + { + "line_number": 1, + "error_message": "type mismatch", + "original_line": "m,t=a f=1" + } + ] + }).encode("utf-8") + http_resp = DummyHttpResponse(status=400, reason="Bad Request", data=body) + exc = InfluxDBRestClientException(http_resp=http_resp) + result = translate_write_exception(exc, use_v2_api=True, accept_partial=True) + + self.assertIsInstance(result, InfluxDBWriteException) + self.assertEqual("partial write of line protocol occurred", result.message) + + def test_partial_write_skipped_when_status_is_not_400(self): + body = json.dumps({ + "error": "partial write of line protocol occurred", + "data": [ + { + "line_number": 1, + "error_message": "type mismatch", + "original_line": "m,t=a f=1" + } + ] + }).encode("utf-8") + http_resp = DummyHttpResponse(status=500, reason="Internal Server Error", data=body) + exc = InfluxDBRestClientException(http_resp=http_resp) + result = translate_write_exception(exc, use_v2_api=False, accept_partial=True) + + self.assertIsInstance(result, InfluxDBWriteException) + self.assertEqual("partial write of line protocol occurred", result.message) diff --git a/tests/test_write_local_server.py b/tests/test_write_local_server.py index 252ab041..d8efd996 100644 --- a/tests/test_write_local_server.py +++ b/tests/test_write_local_server.py @@ -8,7 +8,7 @@ from urllib3.exceptions import TimeoutError as urllib3_TimeoutError from influxdb_client_3 import InfluxDBClient3, WriteOptions, WritePrecision, write_client_options, WriteType -from influxdb_client_3.write_client.write_exceptions import ApiException +from influxdb_client_3.exceptions.write_exceptions import InfluxDBWriteException class TestWriteLocalServer: @@ -154,9 +154,9 @@ def test_write_with_v3_on_v2_server(self, httpserver: HTTPServer): expected = ("Server doesn't support the V3 API endpoint (/api/v3/write_lp). " "Set use_v2_api=True to use the V2 API endpoint.") - with pytest.raises(ApiException, match=r".*Server doesn't support the V3 API endpoint " - r"\(/api/v3/write_lp\)\. " - r"Set use_v2_api=True to use the V2 API endpoint\.") as err: + with pytest.raises(InfluxDBWriteException, match=r".*Server doesn't support the V3 API endpoint " + r"\(/api/v3/write_lp\)\. " + r"Set use_v2_api=True to use the V2 API endpoint\.") as err: client.write(self.SAMPLE_RECORD) assert err.value.message == expected assert err.value.reason == expected @@ -174,9 +174,9 @@ def test_write_with_v2_on_v3_server(self, httpserver: HTTPServer): expected = ("Server doesn't support the V2 API endpoint (/api/v2/write). " "Set use_v2_api=False to use the V3 API endpoint.") - with pytest.raises(ApiException, match=r".*Server doesn't support the V2 API endpoint " - r"\(/api/v2/write\)\. " - r"Set use_v2_api=False to use the V3 API endpoint\.") as err: + with pytest.raises(InfluxDBWriteException, match=r".*Server doesn't support the V2 API endpoint " + r"\(/api/v2/write\)\. " + r"Set use_v2_api=False to use the V3 API endpoint\.") as err: client.write(self.SAMPLE_RECORD) assert err.value.message == expected assert err.value.reason == expected