diff --git a/CHANGELOG.md b/CHANGELOG.md index 5178c68..acdeceb 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,46 @@ All notable changes to this project are documented in this file. ## Unreleased +- New transports: SQL (`BaseSQLTransport` + `SQLiteTransport`, + `PostgresTransport`, `MySQLTransport`), NoSQL (`MongoDBTransport`, + `DynamoDBTransport`, `RedisTransport`), message queues + (`BaseQueueTransport` + `KafkaTransport`, `RabbitMQTransport`, + `SQSTransport`, `PubSubTransport`), and cloud-native sinks + (`CloudWatchTransport`, `CloudLoggingTransport`, `AppInsightsTransport`, + `DatadogTransport`, `ElasticsearchTransport`, `NewRelicTransport`) — + full parity with `logquill-js` 0.2.0. All of it sits on a new shared + `BatchingTransport` base that bounds its buffer by both record count and + estimated byte size, swaps the buffer out before sending so a + synchronous re-entrant flush can't double-send, and catches a failing + send rather than propagating it to the caller (logged via Python's + stdlib `logging.getLogger("logquill")`) — a slow or down sink can't + crash the process. Every optional backend driver (`psycopg2-binary`, + `pymysql`, `pymongo`, `boto3`, `redis`, `kafka-python`, `pika`, + `google-cloud-pubsub`, `google-cloud-logging`) is a lazy, injectable + dependency behind a new `pyproject.toml` extra (`postgres`, `mysql`, + `mongodb`, `redis`, `kafka`, `rabbitmq`, `pubsub`, `gcp-logging`, and a + shared `aws` extra for CloudWatch/DynamoDB/SQS, all boto3-backed); a + missing driver raises an actionable `ImportError` rather than a + cryptic one, and every test injects a hand-written fake instead of + requiring a live service. `SQLiteTransport` needs no extra at all + (stdlib `sqlite3`). + + Two deliberate departures from `logquill-js`'s implementation, same + outward behavior: `AppInsightsTransport` posts to Application Insights' + public ingestion endpoint via stdlib `urllib` instead of an Azure SDK + dependency, and `SQSTransport` dispatches its 10-message chunks + sequentially rather than concurrently, since this project's dispatch is + still fully synchronous end to end (true concurrency arrives once a + non-blocking async worker exists). `SyslogTransport` isn't included + here either, matching `logquill-js` 0.2.0, which didn't ship it; it's a + shared follow-up for both packages, not a Python-only gap. + + Also restructured `logquill/transport.py`, `console_transport.py`, + `file_transport.py`, and `http_transport.py` into a new + `logquill/transports/` subpackage (with `sql/`, `nosql/`, `queue/`, and + `cloud/` subpackages) to hold the 17 new transports — a pure move, the + public `from logquill import ...` surface is unchanged. + - Plugin pipeline: `Plugin` base (`before_log`/`after_log`/`on_error`, all optional to override), `ContextPlugin` (merges fixed context into `meta`), `RedactPlugin` (replaces sensitive `meta` values by key, case- diff --git a/README.md b/README.md index 276921e..e4f9ee9 100644 --- a/README.md +++ b/README.md @@ -22,7 +22,7 @@ landed so far. - **Structured by default** — every call carries a `meta` dict, not just a message string - **Cross-language record shape** — identical JSON shape and level names/weights as [`logquill` on npm](https://www.npmjs.com/package/logquill) -- **Pluggable transports** — `ConsoleTransport` (colorized, stderr for errors), `FileTransport` (rotation), `HTTPTransport` (batched); write your own by subclassing `Transport` +- **Pluggable transports** — `ConsoleTransport` (colorized, stderr for errors), `FileTransport` (rotation), `HTTPTransport` (batched), plus SQL/NoSQL/message-queue/cloud-native sinks (see [Transports](#transports)); write your own by subclassing `Transport` - **Pluggable formatters** — `JSONFormatter` out of the box; implement `format(record) -> str` for your own - **Plugin pipeline** — `ContextPlugin`, `RedactPlugin`, `SamplingPlugin` out of the box; a broken plugin can't crash logging - **Zero required runtime dependencies** — stdlib only; `aiohttp` is opt-in, for async HTTP @@ -101,6 +101,152 @@ logger.info("hello") assert sink.records[0]["message"] == "hello" ``` +### SQL, NoSQL, message queue, and cloud-native transports + +Every transport below shares one design: records are **always batched** +(bounded by both count and estimated byte size via a shared +`BatchingTransport` base — never one write per log call), and every +optional backend driver is a **lazy, injectable dependency** — pass a +pre-built client/connection for tests or an alternate setup, or let the +transport construct one itself from the real driver on first use. A +missing driver raises an actionable `ImportError` telling you which +extra to install, the same shape every transport in this list follows. + +`SQLiteTransport` needs no optional dependency at all (stdlib `sqlite3`), +so it's fully runnable as-is: + +```python +from logquill import Logger, SQLiteTransport + +transport = SQLiteTransport(filename="app.db", ensure_schema=True, max_records=100) +logger = Logger("app", transports=[transport]) + +logger.info("user signed up", user_id=42, run_id="run-1") +logger.close() # flushes any buffered rows +``` + +Every other backend follows the same injection shape — here's +`MongoDBTransport` with a hand-rolled fake standing in for a real +`pymongo` collection (the same pattern every transport's own test suite +uses, so you never need a live service to test your own logging setup): + +```python +from logquill import Logger, MongoDBTransport + +class FakeCollection: + def __init__(self): + self.documents = [] + def insert_many(self, documents): + self.documents.extend(documents) + +collection = FakeCollection() +transport = MongoDBTransport(collection=collection, max_records=1) +logger = Logger("app", transports=[transport]) + +logger.info("user signed up", user_id=42) +assert collection.documents[0]["message"] == "user signed up" +``` + +Passing a real `pymongo.Collection` instead of a fake works identically — +`MongoDBTransport(uri="mongodb://localhost:27017", database="app", collection_name="logs")` +builds one lazily via the optional `pymongo` peer dependency. + +**SQL** — `BaseSQLTransport` (a fixed `logs` table: `timestamp`/`level`/ +`logger`/`message`/`meta`, plus `run_id`/`span_id`/`parent_span_id`/ +`trace_id` for upcoming cross-service trace-correlation support). +`ensure_schema=True` is a dev/test convenience only — production +schema/migrations are your responsibility, same as every batching +transport below. + +| Transport | Driver | Extra | +|---|---|---| +| `SQLiteTransport` | stdlib `sqlite3` | *(none)* | +| `PostgresTransport` | `psycopg2-binary` | `pip install logquill[postgres]` | +| `MySQLTransport` | `pymysql` | `pip install logquill[mysql]` | + +**NoSQL** + +| Transport | Driver | Extra | +|---|---|---| +| `MongoDBTransport` | `pymongo` | `pip install logquill[mongodb]` | +| `DynamoDBTransport` | `boto3` | `pip install logquill[aws]` | +| `RedisTransport` | `redis` | `pip install logquill[redis]` | + +`DynamoDBTransport` partitions by `meta["run_id"]` (falling back to +`meta["trace_id"]`, then the logger name) with `timestamp` as the sort +key. `RedisTransport` writes to a Redis Stream via `XADD` — a fast local +buffer/tail, not a durable store. + +**Message queues** — `BaseQueueTransport` (`topic` names the Kafka +topic / RabbitMQ queue / SQS queue URL / GCP Pub/Sub topic path). +Decouples log producers from consumers so a SIEM, an analytics pipeline, +and an alerting system can all fan out from one topic. `SQSTransport` +chunks at the API's 10-message `SendMessageBatch` cap: + +```python +from logquill import Logger, SQSTransport + +class FakeSQSClient: + def __init__(self): + self.calls = [] + def send_message_batch(self, QueueUrl, Entries): + self.calls.append((QueueUrl, Entries)) + +client = FakeSQSClient() +transport = SQSTransport( + topic="https://sqs.us-east-1.amazonaws.com/123456789012/app-logs", + client=client, + max_records=12, +) +logger = Logger("app", transports=[transport]) +for i in range(12): + logger.info(f"event {i}") +# chunked into two send_message_batch calls: 10 messages, then 2 +``` + +| Transport | Driver | Extra | +|---|---|---| +| `KafkaTransport` | `kafka-python` | `pip install logquill[kafka]` | +| `RabbitMQTransport` | `pika` | `pip install logquill[rabbitmq]` | +| `SQSTransport` | `boto3` | `pip install logquill[aws]` | +| `PubSubTransport` | `google-cloud-pubsub` | `pip install logquill[pubsub]` | + +**Cloud-native** — `DatadogTransport`, `ElasticsearchTransport`, and +`AppInsightsTransport` need no client SDK at all: each POSTs directly to +its provider's public ingestion endpoint via stdlib `urllib`, with an +injectable `sender` for tests: + +```python +from logquill import DatadogTransport, Logger + +class FakeSender: + def __init__(self): + self.calls = [] + def __call__(self, url, api_key, batch): + self.calls.append((url, api_key, batch)) + +sender = FakeSender() +transport = DatadogTransport(api_key="dd-api-key", sender=sender, max_records=1) +logger = Logger("app", transports=[transport]) + +logger.info("user signed up", user_id=42) +``` + +| Transport | Mechanism | Extra | +|---|---|---| +| `CloudWatchTransport` | `boto3` | `pip install logquill[aws]` | +| `CloudLoggingTransport` | `google-cloud-logging` | `pip install logquill[gcp-logging]` | +| `AppInsightsTransport` | stdlib `urllib` (public ingestion endpoint) | *(none)* | +| `DatadogTransport` | stdlib `urllib` | *(none)* | +| `ElasticsearchTransport` | stdlib `urllib` (`_bulk` API) | *(none)* | +| `NewRelicTransport` | stdlib `urllib` + `gzip` | *(none)* | + +`NewRelicTransport` gzips every payload, strips `meta["eventType"]` (New +Relic's reserved key), and on a `429` response reads `Retry-After` and +pauses sends until it elapses — dropping (not requeuing) any batch +flushed during that window, since New Relic blocks the rest of that +minute on a rate-limit breach anyway. + ## Plugins Plugins hook into the pipeline around each log call: `before_log(record)` can diff --git a/logquill/__init__.py b/logquill/__init__.py index 3b40671..3854676 100644 --- a/logquill/__init__.py +++ b/logquill/__init__.py @@ -1,31 +1,70 @@ -from logquill.console_transport import ConsoleTransport -from logquill.context_plugin import ContextPlugin -from logquill.file_transport import FileTransport from logquill.formatter import Formatter, JSONFormatter -from logquill.http_transport import HTTPTransport from logquill.levels import Level, parse_level from logquill.logger import Logger -from logquill.plugin import Plugin +from logquill.plugins.context_plugin import ContextPlugin +from logquill.plugins.plugin import Plugin +from logquill.plugins.redact_plugin import RedactPlugin +from logquill.plugins.sampling_plugin import SamplingPlugin from logquill.records import LogRecord -from logquill.redact_plugin import RedactPlugin -from logquill.sampling_plugin import SamplingPlugin -from logquill.transport import CollectingTransport, Transport +from logquill.transports.batching_transport import BatchingTransport +from logquill.transports.cloud.app_insights_transport import AppInsightsTransport +from logquill.transports.cloud.cloud_logging_transport import CloudLoggingTransport +from logquill.transports.cloud.cloudwatch_transport import CloudWatchTransport +from logquill.transports.cloud.datadog_transport import DatadogTransport +from logquill.transports.cloud.elasticsearch_transport import ElasticsearchTransport +from logquill.transports.cloud.new_relic_transport import NewRelicTransport +from logquill.transports.console_transport import ConsoleTransport +from logquill.transports.file_transport import FileTransport +from logquill.transports.http_transport import HTTPTransport +from logquill.transports.nosql.dynamodb_transport import DynamoDBTransport +from logquill.transports.nosql.mongodb_transport import MongoDBTransport +from logquill.transports.nosql.redis_transport import RedisTransport +from logquill.transports.queue.base_queue_transport import BaseQueueTransport +from logquill.transports.queue.kafka_transport import KafkaTransport +from logquill.transports.queue.pubsub_transport import PubSubTransport +from logquill.transports.queue.rabbitmq_transport import RabbitMQTransport +from logquill.transports.queue.sqs_transport import SQSTransport +from logquill.transports.sql.base_sql_transport import BaseSQLTransport, SQLLogRow +from logquill.transports.sql.mysql_transport import MySQLTransport +from logquill.transports.sql.postgres_transport import PostgresTransport +from logquill.transports.sql.sqlite_transport import SQLiteTransport +from logquill.transports.transport import CollectingTransport, Transport __version__ = "0.1.3" __all__ = [ + "AppInsightsTransport", + "BaseQueueTransport", + "BaseSQLTransport", + "BatchingTransport", + "CloudLoggingTransport", + "CloudWatchTransport", "CollectingTransport", "ConsoleTransport", "ContextPlugin", + "DatadogTransport", + "DynamoDBTransport", + "ElasticsearchTransport", "FileTransport", "Formatter", "HTTPTransport", "JSONFormatter", + "KafkaTransport", "Level", "LogRecord", "Logger", + "MongoDBTransport", + "MySQLTransport", + "NewRelicTransport", "Plugin", + "PostgresTransport", + "PubSubTransport", + "RabbitMQTransport", "RedactPlugin", + "RedisTransport", + "SQLLogRow", + "SQLiteTransport", + "SQSTransport", "SamplingPlugin", "Transport", "parse_level", diff --git a/logquill/logger.py b/logquill/logger.py index e40f290..208fdf8 100644 --- a/logquill/logger.py +++ b/logquill/logger.py @@ -4,9 +4,9 @@ from typing import Any from logquill.levels import Level, parse_level -from logquill.plugin import Plugin +from logquill.plugins.plugin import Plugin from logquill.records import LogRecord, create_record -from logquill.transport import Transport +from logquill.transports.transport import Transport class Logger: diff --git a/logquill/plugins/__init__.py b/logquill/plugins/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/logquill/context_plugin.py b/logquill/plugins/context_plugin.py similarity index 92% rename from logquill/context_plugin.py rename to logquill/plugins/context_plugin.py index ddfff18..ee82923 100644 --- a/logquill/context_plugin.py +++ b/logquill/plugins/context_plugin.py @@ -2,7 +2,7 @@ from typing import Any -from logquill.plugin import Plugin +from logquill.plugins.plugin import Plugin from logquill.records import LogRecord diff --git a/logquill/plugin.py b/logquill/plugins/plugin.py similarity index 100% rename from logquill/plugin.py rename to logquill/plugins/plugin.py diff --git a/logquill/redact_plugin.py b/logquill/plugins/redact_plugin.py similarity index 95% rename from logquill/redact_plugin.py rename to logquill/plugins/redact_plugin.py index 3fe7616..108f018 100644 --- a/logquill/redact_plugin.py +++ b/logquill/plugins/redact_plugin.py @@ -2,7 +2,7 @@ from collections.abc import Iterable -from logquill.plugin import Plugin +from logquill.plugins.plugin import Plugin from logquill.records import LogRecord DEFAULT_REDACTED_KEYS = frozenset({"password", "token", "secret", "api_key", "authorization"}) diff --git a/logquill/sampling_plugin.py b/logquill/plugins/sampling_plugin.py similarity index 93% rename from logquill/sampling_plugin.py rename to logquill/plugins/sampling_plugin.py index c46ff19..e0d087c 100644 --- a/logquill/sampling_plugin.py +++ b/logquill/plugins/sampling_plugin.py @@ -3,7 +3,7 @@ import random from typing import Callable -from logquill.plugin import Plugin +from logquill.plugins.plugin import Plugin from logquill.records import LogRecord diff --git a/logquill/transports/__init__.py b/logquill/transports/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/logquill/transports/batching_transport.py b/logquill/transports/batching_transport.py new file mode 100644 index 0000000..ef060b9 --- /dev/null +++ b/logquill/transports/batching_transport.py @@ -0,0 +1,69 @@ +from __future__ import annotations + +import json +import logging +from abc import abstractmethod +from typing import Generic, Sequence, TypeVar, cast + +from logquill.formatter import Formatter +from logquill.records import LogRecord +from logquill.transports.transport import Transport + +T = TypeVar("T") + +_logger = logging.getLogger("logquill") + + +class BatchingTransport(Transport, Generic[T]): + """Shared base for every transport that buffers records and sends them in + batches (SQL, NoSQL, message queue, and cloud-native sinks). + + Bounds the buffer by **both** record count and estimated byte size — + count alone lets a handful of huge `meta` payloads blow past reasonable + memory before a batch triggers. A flush fires as soon as either bound is + hit, checked after every `write()`. `_send_batch` is never called with an + empty batch, and a failing send is caught and logged rather than + propagated — a slow or down sink can never crash the caller's process. + """ + + def __init__( + self, + *, + formatter: Formatter | None = None, + max_records: int = 100, + max_bytes: int = 1_000_000, + ) -> None: + super().__init__(formatter) + self.max_records = max_records + self.max_bytes = max_bytes + self._buffer: list[T] = [] + self._buffer_bytes = 0 + + def _to_item(self, formatted: str, record: LogRecord) -> T: + return cast(T, record) + + def _size_of(self, item: T) -> int: + return len(json.dumps(item, separators=(",", ":"), default=str).encode("utf-8")) + + def write(self, formatted: str, record: LogRecord) -> None: + item = self._to_item(formatted, record) + self._buffer.append(item) + self._buffer_bytes += self._size_of(item) + if len(self._buffer) >= self.max_records or self._buffer_bytes >= self.max_bytes: + self.flush() + + def flush(self) -> None: + if not self._buffer: + return + batch, self._buffer = self._buffer, [] + self._buffer_bytes = 0 + try: + self._send_batch(batch) + except Exception: + _logger.exception("%s: failed to send log batch", type(self).__name__) + + def close(self) -> None: + self.flush() + + @abstractmethod + def _send_batch(self, batch: Sequence[T]) -> None: ... diff --git a/logquill/transports/cloud/__init__.py b/logquill/transports/cloud/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/logquill/transports/cloud/app_insights_transport.py b/logquill/transports/cloud/app_insights_transport.py new file mode 100644 index 0000000..639e482 --- /dev/null +++ b/logquill/transports/cloud/app_insights_transport.py @@ -0,0 +1,76 @@ +from __future__ import annotations + +import json +import urllib.request +from typing import Callable, Sequence + +from logquill.formatter import Formatter +from logquill.records import LogRecord +from logquill.transports.batching_transport import BatchingTransport + +Sender = Callable[[str, Sequence[str]], None] + +_SEVERITY = {"TRACE": 0, "DEBUG": 0, "INFO": 1, "WARN": 2, "ERROR": 3, "FATAL": 4} + + +def _urllib_sender(url: str, batch: Sequence[str]) -> None: + body = "\n".join(batch).encode("utf-8") + request = urllib.request.Request( + url, + data=body, + headers={"Content-Type": "application/x-json-stream"}, + method="POST", + ) + with urllib.request.urlopen(request, timeout=10) as response: # noqa: S310 + response.read() + + +class AppInsightsTransport(BatchingTransport[LogRecord]): + """Ships records to Azure Application Insights as trace telemetry. + + Deliberately diverges from logquill-js's SDK-based approach: Application + Insights has a documented public ingestion endpoint + (`https://dc.services.visualstudio.com/v2/track`), so this posts + newline-delimited telemetry envelopes there directly via stdlib + `urllib` instead of pulling in an Azure SDK dependency. Same outward + behavior — records land as trace telemetry — zero new dependency. + """ + + URL = "https://dc.services.visualstudio.com/v2/track" + + def __init__( + self, + *, + instrumentation_key: str, + formatter: Formatter | None = None, + max_records: int = 100, + max_bytes: int = 1_000_000, + sender: Sender | None = None, + ) -> None: + super().__init__(formatter=formatter, max_records=max_records, max_bytes=max_bytes) + self.instrumentation_key = instrumentation_key + self._sender: Sender = sender or _urllib_sender + + def _envelope(self, record: LogRecord) -> str: + return json.dumps( + { + "name": "Microsoft.ApplicationInsights.Trace", + "time": record["timestamp"], + "iKey": self.instrumentation_key, + "data": { + "baseType": "MessageData", + "baseData": { + "ver": 2, + "message": record["message"], + "severityLevel": _SEVERITY.get(record["level"], 1), + "properties": {k: str(v) for k, v in record["meta"].items()}, + }, + }, + }, + separators=(",", ":"), + default=str, + ) + + def _send_batch(self, batch: Sequence[LogRecord]) -> None: + lines = [self._envelope(record) for record in batch] + self._sender(self.URL, lines) diff --git a/logquill/transports/cloud/cloud_logging_transport.py b/logquill/transports/cloud/cloud_logging_transport.py new file mode 100644 index 0000000..983e06d --- /dev/null +++ b/logquill/transports/cloud/cloud_logging_transport.py @@ -0,0 +1,67 @@ +from __future__ import annotations + +from typing import Any, Protocol, Sequence, cast + +from logquill.formatter import Formatter +from logquill.records import LogRecord +from logquill.transports.batching_transport import BatchingTransport + +_SEVERITY = { + "TRACE": "DEBUG", + "DEBUG": "DEBUG", + "INFO": "INFO", + "WARN": "WARNING", + "ERROR": "ERROR", + "FATAL": "CRITICAL", +} + + +class CloudLoggingClientLike(Protocol): + def log_struct(self, info: dict[str, Any], severity: str) -> None: ... + + +class CloudLoggingTransport(BatchingTransport[LogRecord]): + """Ships records to Google Cloud Logging via `log_struct`, through the + optional `google-cloud-logging` peer dependency, mapping `level` onto + Cloud Logging's `severity` scale. One `log_struct` call per record — + the client library has no native multi-record batch form at this + level. Pass `client` to inject a pre-built Cloud Logging `Logger` (or + a fake, for tests), or `log_name` to let this transport connect + itself.""" + + def __init__( + self, + *, + log_name: str = "logquill", + client: CloudLoggingClientLike | None = None, + formatter: Formatter | None = None, + max_records: int = 100, + max_bytes: int = 1_000_000, + ) -> None: + super().__init__(formatter=formatter, max_records=max_records, max_bytes=max_bytes) + self._log_name = log_name + self._injected = client + self._client: CloudLoggingClientLike | None = None + + def _resolved_client(self) -> CloudLoggingClientLike: + if self._injected is not None: + return self._injected + if self._client is None: + self._client = self._import_client() + return self._client + + def _import_client(self) -> CloudLoggingClientLike: + try: + import google.cloud.logging as cloud_logging # type: ignore # stubs vary by env + except ImportError as exc: + raise ImportError( + "CloudLoggingTransport: install `google-cloud-logging` to " + "use this transport without providing a client — " + "`pip install logquill[gcp-logging]`" + ) from exc + return cast(CloudLoggingClientLike, cloud_logging.Client().logger(self._log_name)) + + def _send_batch(self, batch: Sequence[LogRecord]) -> None: + client = self._resolved_client() + for record in batch: + client.log_struct(dict(record), severity=_SEVERITY.get(record["level"], "DEFAULT")) diff --git a/logquill/transports/cloud/cloudwatch_transport.py b/logquill/transports/cloud/cloudwatch_transport.py new file mode 100644 index 0000000..0aec29a --- /dev/null +++ b/logquill/transports/cloud/cloudwatch_transport.py @@ -0,0 +1,80 @@ +from __future__ import annotations + +from typing import Any, Protocol, Sequence, cast + +from logquill.formatter import Formatter +from logquill.records import LogRecord +from logquill.transports.batching_transport import BatchingTransport + + +class CloudWatchClientLike(Protocol): + def put_log_events( + self, + logGroupName: str, + logStreamName: str, + logEvents: Sequence[dict[str, Any]], # noqa: N803 + ) -> object: ... + + +def _to_millis(iso_timestamp: str) -> int: + from datetime import datetime + + normalized = iso_timestamp[:-1] + "+00:00" if iso_timestamp.endswith("Z") else iso_timestamp + return int(datetime.fromisoformat(normalized).timestamp() * 1000) + + +class CloudWatchTransport(BatchingTransport[LogRecord]): + """Ships batched records to AWS CloudWatch Logs via `put_log_events`, + through the optional `boto3` peer dependency (shared with + `DynamoDBTransport`/`SQSTransport`). Sorts events by timestamp before + each call, since the API requires events within a single + `PutLogEvents` request to be in chronological order. Pass `client` to + inject a pre-built `boto3` `logs` client (or a fake, for tests), or + `region` to let this transport connect itself.""" + + def __init__( + self, + *, + log_group: str, + log_stream: str, + client: CloudWatchClientLike | None = None, + region: str | None = None, + formatter: Formatter | None = None, + max_records: int = 100, + max_bytes: int = 1_000_000, + ) -> None: + super().__init__(formatter=formatter, max_records=max_records, max_bytes=max_bytes) + self.log_group = log_group + self.log_stream = log_stream + self._injected = client + self._region = region + self._client: CloudWatchClientLike | None = None + + def _resolved_client(self) -> CloudWatchClientLike: + if self._injected is not None: + return self._injected + if self._client is None: + self._client = self._import_client() + return self._client + + def _import_client(self) -> CloudWatchClientLike: + try: + import boto3 # type: ignore # stub availability for this optional dep varies by env + except ImportError as exc: + raise ImportError( + "CloudWatchTransport: install `boto3` to use this " + "transport without providing a client — " + "`pip install logquill[aws]`" + ) from exc + return cast(CloudWatchClientLike, boto3.client("logs", region_name=self._region)) + + def _send_batch(self, batch: Sequence[LogRecord]) -> None: + client = self._resolved_client() + events = [ + {"timestamp": _to_millis(record["timestamp"]), "message": self.format(record)} + for record in batch + ] + events.sort(key=lambda event: cast(int, event["timestamp"])) + client.put_log_events( + logGroupName=self.log_group, logStreamName=self.log_stream, logEvents=events + ) diff --git a/logquill/transports/cloud/datadog_transport.py b/logquill/transports/cloud/datadog_transport.py new file mode 100644 index 0000000..abb70ea --- /dev/null +++ b/logquill/transports/cloud/datadog_transport.py @@ -0,0 +1,57 @@ +from __future__ import annotations + +import urllib.error +import urllib.request +from typing import Callable, Sequence + +from logquill.formatter import Formatter +from logquill.records import LogRecord +from logquill.transports.batching_transport import BatchingTransport + +DatadogSender = Callable[[str, str, Sequence[str]], None] + + +def _urllib_sender(url: str, api_key: str, batch: Sequence[str]) -> None: + body = ("[" + ",".join(batch) + "]").encode("utf-8") + request = urllib.request.Request( + url, + data=body, + headers={"Content-Type": "application/json", "DD-API-KEY": api_key}, + method="POST", + ) + try: + with urllib.request.urlopen(request, timeout=10) as response: # noqa: S310 + response.read() + except urllib.error.HTTPError as exc: + raise RuntimeError( + f"DatadogTransport: request to {url} failed with status " + f"{exc.code} — check the API key and site region" + ) from exc + + +class DatadogTransport(BatchingTransport[LogRecord]): + """Batches records and POSTs them as a JSON array to Datadog's Logs + intake API via stdlib `urllib` — no client dependency. `site` + selects the region-specific endpoint (`datadoghq.com`, `datadoghq.eu`, + `us3.datadoghq.com`, `us5.datadoghq.com`, `ap1.datadoghq.com`, ...). + Pass `sender` to inject a fake for tests or an alternate transport.""" + + def __init__( + self, + *, + api_key: str, + site: str = "datadoghq.com", + formatter: Formatter | None = None, + max_records: int = 100, + max_bytes: int = 1_000_000, + sender: DatadogSender | None = None, + ) -> None: + super().__init__(formatter=formatter, max_records=max_records, max_bytes=max_bytes) + self.api_key = api_key + self.site = site + self.url = f"https://http-intake.logs.{site}/api/v2/logs" + self._sender: DatadogSender = sender or _urllib_sender + + def _send_batch(self, batch: Sequence[LogRecord]) -> None: + formatted = [self.format(record) for record in batch] + self._sender(self.url, self.api_key, formatted) diff --git a/logquill/transports/cloud/elasticsearch_transport.py b/logquill/transports/cloud/elasticsearch_transport.py new file mode 100644 index 0000000..ef6a115 --- /dev/null +++ b/logquill/transports/cloud/elasticsearch_transport.py @@ -0,0 +1,59 @@ +from __future__ import annotations + +import json +import urllib.error +import urllib.request +from typing import Callable, Sequence + +from logquill.formatter import Formatter +from logquill.records import LogRecord +from logquill.transports.batching_transport import BatchingTransport + +ElasticsearchSender = Callable[[str, Sequence[str]], None] + + +def _urllib_sender(url: str, batch: Sequence[str]) -> None: + body = ("\n".join(batch) + "\n").encode("utf-8") + request = urllib.request.Request( + url, + data=body, + headers={"Content-Type": "application/x-ndjson"}, + method="POST", + ) + try: + with urllib.request.urlopen(request, timeout=10) as response: # noqa: S310 + response.read() + except urllib.error.HTTPError as exc: + raise RuntimeError( + f"ElasticsearchTransport: request to {url} failed with status " + f"{exc.code} — check the URL and index" + ) from exc + + +class ElasticsearchTransport(BatchingTransport[LogRecord]): + """Batches records and POSTs them to Elasticsearch's `_bulk` API as + newline-delimited action+source pairs via stdlib `urllib` — no client + dependency. Pass `sender` to inject a fake for tests or an alternate + transport.""" + + def __init__( + self, + *, + url: str, + index: str = "logs", + formatter: Formatter | None = None, + max_records: int = 100, + max_bytes: int = 1_000_000, + sender: ElasticsearchSender | None = None, + ) -> None: + super().__init__(formatter=formatter, max_records=max_records, max_bytes=max_bytes) + self.index = index + self.bulk_url = f"{url.rstrip('/')}/_bulk" + self._sender: ElasticsearchSender = sender or _urllib_sender + + def _send_batch(self, batch: Sequence[LogRecord]) -> None: + lines: list[str] = [] + for record in batch: + lines.append(json.dumps({"index": {"_index": self.index}}, separators=(",", ":"))) + lines.append(self.format(record)) + self._sender(self.bulk_url, lines) diff --git a/logquill/transports/cloud/new_relic_transport.py b/logquill/transports/cloud/new_relic_transport.py new file mode 100644 index 0000000..ac3efc0 --- /dev/null +++ b/logquill/transports/cloud/new_relic_transport.py @@ -0,0 +1,136 @@ +from __future__ import annotations + +import gzip +import json +import logging +import time +import urllib.error +import urllib.request +from email.utils import parsedate_to_datetime +from typing import Callable, Dict, Literal, Sequence, TypedDict + +from logquill.formatter import Formatter +from logquill.records import LogRecord +from logquill.transports.batching_transport import BatchingTransport + +NewRelicRegion = Literal["US", "EU"] + +_URLS: dict[NewRelicRegion, str] = { + "US": "https://log-api.newrelic.com/log/v1", + "EU": "https://log-api.eu.newrelic.com/log/v1", +} + +_logger = logging.getLogger("logquill") + + +class NewRelicSenderResult(TypedDict): + ok: bool + status: int + retry_after: str | None + + +NewRelicSender = Callable[[str, Dict[str, str], bytes], NewRelicSenderResult] + + +def _urllib_sender(url: str, headers: dict[str, str], body: bytes) -> NewRelicSenderResult: + request = urllib.request.Request(url, data=body, headers=headers, method="POST") + try: + with urllib.request.urlopen(request, timeout=10) as response: # noqa: S310 + response.read() + return {"ok": True, "status": response.status, "retry_after": None} + except urllib.error.HTTPError as exc: + retry_after = exc.headers.get("Retry-After") if exc.headers else None + return {"ok": False, "status": exc.code, "retry_after": retry_after} + + +def _without_event_type(record: LogRecord) -> LogRecord: + meta = dict(record["meta"]) + meta.pop("eventType", None) + return LogRecord( + timestamp=record["timestamp"], + level=record["level"], + logger=record["logger"], + message=record["message"], + meta=meta, + ) + + +def _resume_timestamp(retry_after: str | None, now: float) -> float: + if retry_after is None: + return now + 60.0 + try: + return now + float(retry_after) + except ValueError: + pass + try: + parsed = parsedate_to_datetime(retry_after) + except (TypeError, ValueError): + return now + 60.0 + return parsed.timestamp() + + +class NewRelicTransport(BatchingTransport[LogRecord]): + """Batches and gzips records before POSTing to New Relic's Log API via + stdlib `urllib` + `gzip` — no client dependency. `region` selects the + US or EU endpoint. Strips `meta["eventType"]` from every record — New + Relic reserves and silently drops that key. On a 429 response, reads + `Retry-After` (integer seconds or an HTTP-date) and pauses further + sends until that time elapses, **dropping** (not requeuing) any batch + flushed during the pause window, since New Relic blocks further sends + for the rest of that minute on a rate-limit breach and retrying + immediately just wastes calls. `clock` is injectable so the backoff is + testable without real waits.""" + + def __init__( + self, + *, + license_key: str, + region: NewRelicRegion = "US", + formatter: Formatter | None = None, + max_records: int = 100, + max_bytes: int = 1_000_000, + sender: NewRelicSender | None = None, + clock: Callable[[], float] | None = None, + ) -> None: + super().__init__(formatter=formatter, max_records=max_records, max_bytes=max_bytes) + self.license_key = license_key + self.region = region + self.url = _URLS[region] + self._sender: NewRelicSender = sender or _urllib_sender + self._clock: Callable[[], float] = clock or time.time + self._paused_until: float | None = None + + def _send_batch(self, batch: Sequence[LogRecord]) -> None: + now = self._clock() + if self._paused_until is not None and now < self._paused_until: + _logger.error( + "NewRelicTransport: sends paused until %s after a 429 " + "rate-limit response — skipping this batch rather than " + "making a doomed request", + self._paused_until, + ) + return + self._paused_until = None + + records = [_without_event_type(record) for record in batch] + body = gzip.compress(json.dumps(records, separators=(",", ":"), default=str).encode()) + headers = { + "Content-Type": "application/json", + "Content-Encoding": "gzip", + "Api-Key": self.license_key, + } + + result = self._sender(self.url, headers, body) + + if result["status"] == 429: + self._paused_until = _resume_timestamp(result["retry_after"], now) + _logger.error( + "NewRelicTransport: received 429 — pausing sends until %s", + self._paused_until, + ) + return + if not result["ok"]: + raise RuntimeError( + f"NewRelicTransport: request to {self.url} failed with " + f"status {result['status']} — check the license key and region" + ) diff --git a/logquill/console_transport.py b/logquill/transports/console_transport.py similarity index 96% rename from logquill/console_transport.py rename to logquill/transports/console_transport.py index 82b3bb3..7d69691 100644 --- a/logquill/console_transport.py +++ b/logquill/transports/console_transport.py @@ -6,7 +6,7 @@ from logquill.formatter import Formatter from logquill.levels import Level, parse_level from logquill.records import LogRecord -from logquill.transport import Transport +from logquill.transports.transport import Transport _COLORS = { Level.TRACE: "\x1b[90m", # gray diff --git a/logquill/file_transport.py b/logquill/transports/file_transport.py similarity index 97% rename from logquill/file_transport.py rename to logquill/transports/file_transport.py index 40d4b69..007ab23 100644 --- a/logquill/file_transport.py +++ b/logquill/transports/file_transport.py @@ -5,7 +5,7 @@ from logquill.formatter import Formatter from logquill.records import LogRecord -from logquill.transport import Transport +from logquill.transports.transport import Transport class FileTransport(Transport): diff --git a/logquill/http_transport.py b/logquill/transports/http_transport.py similarity index 97% rename from logquill/http_transport.py rename to logquill/transports/http_transport.py index e351ca7..f6574bd 100644 --- a/logquill/http_transport.py +++ b/logquill/transports/http_transport.py @@ -5,7 +5,7 @@ from logquill.formatter import Formatter from logquill.records import LogRecord -from logquill.transport import Transport +from logquill.transports.transport import Transport Sender = Callable[[str, Sequence[str]], None] diff --git a/logquill/transports/nosql/__init__.py b/logquill/transports/nosql/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/logquill/transports/nosql/dynamodb_transport.py b/logquill/transports/nosql/dynamodb_transport.py new file mode 100644 index 0000000..205319c --- /dev/null +++ b/logquill/transports/nosql/dynamodb_transport.py @@ -0,0 +1,95 @@ +from __future__ import annotations + +from typing import Any, ContextManager, Protocol, Sequence, cast + +from logquill.formatter import Formatter +from logquill.records import LogRecord +from logquill.transports.batching_transport import BatchingTransport + + +class DynamoBatchWriterLike(Protocol): + def put_item(self, Item: dict[str, Any]) -> object: ... # noqa: N803 — matches boto3's kwarg + + +class DynamoTableLike(Protocol): + def batch_writer(self) -> ContextManager[DynamoBatchWriterLike]: ... + + +def _partition_key(record: LogRecord) -> str: + meta = record["meta"] + run_id = meta.get("run_id") + if isinstance(run_id, str) and run_id: + return run_id + trace_id = meta.get("trace_id") + if isinstance(trace_id, str) and trace_id: + return trace_id + return record["logger"] + + +def _to_item(record: LogRecord) -> dict[str, Any]: + meta = record["meta"] + item: dict[str, Any] = { + "run_id": _partition_key(record), + "timestamp": record["timestamp"], + "level": record["level"], + "logger": record["logger"], + "message": record["message"], + "meta": meta, + } + for key in ("span_id", "parent_span_id", "trace_id"): + value = meta.get(key) + if isinstance(value, str): + item[key] = value + return item + + +class DynamoDBTransport(BatchingTransport[LogRecord]): + """Writes batches to DynamoDB via a `Table.batch_writer()` context + manager — boto3 auto-marshals Python types, auto-chunks to the API's + 25-item `BatchWriteItem` cap, and auto-retries unprocessed items, so no + hand-rolled `AttributeValue` marshalling or chunking is needed here. + Partitions by `meta["run_id"]` (falling back to `meta["trace_id"]`, + then the logger name) with `timestamp` as the sort key. Pass `table` to + inject a pre-built `boto3` `Table` (or a fake, for tests), or + `table_name`/`region` to let this transport connect itself via the + optional `boto3` peer dependency.""" + + def __init__( + self, + *, + table: DynamoTableLike | None = None, + table_name: str = "logs", + region: str | None = None, + formatter: Formatter | None = None, + max_records: int = 100, + max_bytes: int = 1_000_000, + ) -> None: + super().__init__(formatter=formatter, max_records=max_records, max_bytes=max_bytes) + self._injected = table + self._table_name = table_name + self._region = region + self._table: DynamoTableLike | None = None + + def _resolved_table(self) -> DynamoTableLike: + if self._injected is not None: + return self._injected + if self._table is None: + self._table = self._import_table() + return self._table + + def _import_table(self) -> DynamoTableLike: + try: + import boto3 # type: ignore # stub availability for this optional dep varies by env + except ImportError as exc: + raise ImportError( + "DynamoDBTransport: install `boto3` to use this transport " + "without providing a table — `pip install logquill[aws]`" + ) from exc + resource = boto3.resource("dynamodb", region_name=self._region) + return cast(DynamoTableLike, resource.Table(self._table_name)) + + def _send_batch(self, batch: Sequence[LogRecord]) -> None: + table = self._resolved_table() + with table.batch_writer() as writer: + for record in batch: + writer.put_item(Item=_to_item(record)) diff --git a/logquill/transports/nosql/mongodb_transport.py b/logquill/transports/nosql/mongodb_transport.py new file mode 100644 index 0000000..cbaf956 --- /dev/null +++ b/logquill/transports/nosql/mongodb_transport.py @@ -0,0 +1,71 @@ +from __future__ import annotations + +from typing import Any, Protocol, Sequence, cast + +from logquill.formatter import Formatter +from logquill.records import LogRecord +from logquill.transports.batching_transport import BatchingTransport + + +class MongoCollectionLike(Protocol): + def insert_many(self, documents: Sequence[dict[str, Any]]) -> object: ... + + +class MongoClientLike(Protocol): + def close(self) -> None: ... + + +class MongoDBTransport(BatchingTransport[LogRecord]): + """Batches log records into `insert_many()` calls against a MongoDB + collection — records map 1:1 to documents, no JSON-in-a-column + workaround needed. Pass `collection` to inject a pre-built + `pymongo.Collection` (or a fake, for tests), or `uri`/`database`/ + `collection_name` to let this transport connect itself via the + optional `pymongo` peer dependency.""" + + def __init__( + self, + *, + collection: MongoCollectionLike | None = None, + uri: str = "mongodb://localhost:27017", + database: str = "logquill", + collection_name: str = "logs", + formatter: Formatter | None = None, + max_records: int = 100, + max_bytes: int = 1_000_000, + ) -> None: + super().__init__(formatter=formatter, max_records=max_records, max_bytes=max_bytes) + self._injected = collection + self._uri = uri + self._database = database + self._collection_name = collection_name + self._collection: MongoCollectionLike | None = None + self._client: MongoClientLike | None = None + + def _resolved_collection(self) -> MongoCollectionLike: + if self._injected is not None: + return self._injected + if self._collection is None: + self._collection = self._import_collection() + return self._collection + + def _import_collection(self) -> MongoCollectionLike: + try: + import pymongo # type: ignore # stub availability for this optional dep varies by env + except ImportError as exc: + raise ImportError( + "MongoDBTransport: install `pymongo` to use this transport " + "without providing a collection — `pip install logquill[mongodb]`" + ) from exc + client = pymongo.MongoClient(self._uri) + self._client = cast(MongoClientLike, client) + return cast(MongoCollectionLike, client[self._database][self._collection_name]) + + def close(self) -> None: + super().close() + if self._injected is None and self._client is not None: + self._client.close() + + def _send_batch(self, batch: Sequence[LogRecord]) -> None: + collection = self._resolved_collection() + collection.insert_many([dict(record) for record in batch]) diff --git a/logquill/transports/nosql/redis_transport.py b/logquill/transports/nosql/redis_transport.py new file mode 100644 index 0000000..14c7650 --- /dev/null +++ b/logquill/transports/nosql/redis_transport.py @@ -0,0 +1,76 @@ +from __future__ import annotations + +import json +from typing import Any, Protocol, Sequence + +from logquill.formatter import Formatter +from logquill.records import LogRecord +from logquill.transports.batching_transport import BatchingTransport + + +class RedisClientLike(Protocol): + def xadd(self, name: str, fields: dict[str, Any]) -> object: ... + def close(self) -> None: ... + + +class RedisTransport(BatchingTransport[LogRecord]): + """Appends each record to a Redis Stream via `XADD`, through the + optional `redis` peer dependency — a fast local buffer/tail, **not** a + durable store, worth using as such rather than as a replacement for a + SQL/NoSQL sink. One `XADD` per record: Streams has no native + multi-record batch form, unlike `SendMessageBatch`/`BatchWriteItem`. + Pass `client` to inject a pre-built `redis.Redis` (or a fake, for + tests), or `url` to let this transport connect itself.""" + + def __init__( + self, + *, + client: RedisClientLike | None = None, + stream: str = "logs", + url: str = "redis://localhost:6379", + formatter: Formatter | None = None, + max_records: int = 100, + max_bytes: int = 1_000_000, + ) -> None: + super().__init__(formatter=formatter, max_records=max_records, max_bytes=max_bytes) + self._injected = client + self._stream = stream + self._url = url + self._client: RedisClientLike | None = None + + def _resolved_client(self) -> RedisClientLike: + if self._injected is not None: + return self._injected + if self._client is None: + self._client = self._import_client() + return self._client + + def _import_client(self) -> RedisClientLike: + try: + import redis # type: ignore # stub availability for this optional dep varies by env + except ImportError as exc: + raise ImportError( + "RedisTransport: install `redis` to use this transport " + "without providing a client — `pip install logquill[redis]`" + ) from exc + client: RedisClientLike = redis.Redis.from_url(self._url) + return client + + def close(self) -> None: + super().close() + if self._injected is None and self._client is not None: + self._client.close() + + def _send_batch(self, batch: Sequence[LogRecord]) -> None: + client = self._resolved_client() + for record in batch: + client.xadd( + self._stream, + { + "timestamp": record["timestamp"], + "level": record["level"], + "logger": record["logger"], + "message": record["message"], + "meta": json.dumps(record["meta"], separators=(",", ":"), default=str), + }, + ) diff --git a/logquill/transports/queue/__init__.py b/logquill/transports/queue/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/logquill/transports/queue/base_queue_transport.py b/logquill/transports/queue/base_queue_transport.py new file mode 100644 index 0000000..0441e7c --- /dev/null +++ b/logquill/transports/queue/base_queue_transport.py @@ -0,0 +1,34 @@ +from __future__ import annotations + +from abc import abstractmethod +from typing import Sequence + +from logquill.formatter import Formatter +from logquill.records import LogRecord +from logquill.transports.batching_transport import BatchingTransport + + +class BaseQueueTransport(BatchingTransport[LogRecord]): + """Shared base for every message-queue sink: decouples log producers + from consumers so multiple downstream systems (a SIEM, an analytics + pipeline, an alerting system) can fan out from one topic, and buffers + through a downstream outage instead of blocking or dropping. Always + batches — never publishes one message per log call. + """ + + def __init__( + self, + *, + topic: str, + formatter: Formatter | None = None, + max_records: int = 100, + max_bytes: int = 1_000_000, + ) -> None: + super().__init__(formatter=formatter, max_records=max_records, max_bytes=max_bytes) + self.topic = topic + + def _send_batch(self, batch: Sequence[LogRecord]) -> None: + self._publish_batch(batch) + + @abstractmethod + def _publish_batch(self, records: Sequence[LogRecord]) -> None: ... diff --git a/logquill/transports/queue/kafka_transport.py b/logquill/transports/queue/kafka_transport.py new file mode 100644 index 0000000..e50a0c0 --- /dev/null +++ b/logquill/transports/queue/kafka_transport.py @@ -0,0 +1,83 @@ +from __future__ import annotations + +import json +from typing import Protocol, Sequence, cast + +from logquill.formatter import Formatter +from logquill.records import LogRecord +from logquill.transports.queue.base_queue_transport import BaseQueueTransport + + +class KafkaProducerLike(Protocol): + def send(self, topic: str, value: bytes, key: bytes | None = None) -> object: ... + def flush(self) -> None: ... + def close(self) -> None: ... + + +def _partition_key(record: LogRecord) -> bytes | None: + meta = record["meta"] + run_id = meta.get("run_id") + if isinstance(run_id, str) and run_id: + return run_id.encode("utf-8") + trace_id = meta.get("trace_id") + if isinstance(trace_id, str) and trace_id: + return trace_id.encode("utf-8") + return None + + +class KafkaTransport(BaseQueueTransport): + """Publishes batches to Kafka via the optional `kafka-python` peer + dependency (pure Python — no `librdkafka` system dependency, unlike + `confluent-kafka`), keyed by `meta["run_id"]`/`meta["trace_id"]` so one + trace's records stay ordered on the same partition. Pass `producer` to + inject a pre-built `KafkaProducer` (or a fake, for tests), or + `bootstrap_servers` to let this transport connect itself.""" + + def __init__( + self, + *, + topic: str, + producer: KafkaProducerLike | None = None, + bootstrap_servers: str = "localhost:9092", + formatter: Formatter | None = None, + max_records: int = 100, + max_bytes: int = 1_000_000, + ) -> None: + super().__init__( + topic=topic, formatter=formatter, max_records=max_records, max_bytes=max_bytes + ) + self._injected = producer + self._bootstrap_servers = bootstrap_servers + self._producer: KafkaProducerLike | None = None + + def _resolved_producer(self) -> KafkaProducerLike: + if self._injected is not None: + return self._injected + if self._producer is None: + self._producer = self._import_producer() + return self._producer + + def _import_producer(self) -> KafkaProducerLike: + try: + import kafka # type: ignore # stub availability for this optional dep varies by env + except ImportError as exc: + raise ImportError( + "KafkaTransport: install `kafka-python` to use this " + "transport without providing a producer — " + "`pip install logquill[kafka]`" + ) from exc + return cast( + KafkaProducerLike, kafka.KafkaProducer(bootstrap_servers=self._bootstrap_servers) + ) + + def close(self) -> None: + super().close() + if self._injected is None and self._producer is not None: + self._producer.close() + + def _publish_batch(self, records: Sequence[LogRecord]) -> None: + producer = self._resolved_producer() + for record in records: + value = json.dumps(record, separators=(",", ":"), default=str).encode("utf-8") + producer.send(self.topic, value=value, key=_partition_key(record)) + producer.flush() diff --git a/logquill/transports/queue/pubsub_transport.py b/logquill/transports/queue/pubsub_transport.py new file mode 100644 index 0000000..4dc365c --- /dev/null +++ b/logquill/transports/queue/pubsub_transport.py @@ -0,0 +1,70 @@ +from __future__ import annotations + +import json +from typing import Protocol, Sequence, cast + +from logquill.formatter import Formatter +from logquill.records import LogRecord +from logquill.transports.queue.base_queue_transport import BaseQueueTransport + + +class PubSubFutureLike(Protocol): + def result(self) -> object: ... + + +class PubSubTopicLike(Protocol): + def publish(self, topic: str, data: bytes) -> PubSubFutureLike: ... + + +class PubSubTransport(BaseQueueTransport): + """Publishes to GCP Pub/Sub via the optional `google-cloud-pubsub` peer + dependency. `topic` is the fully-qualified topic path (e.g. + `projects//topics/`). Each record's publish future is + waited on before the flush returns, keeping this transport's dispatch + synchronous like every other transport in the project today. Pass + `client` to inject a pre-built `PublisherClient` (or a fake, for + tests).""" + + def __init__( + self, + *, + topic: str, + client: PubSubTopicLike | None = None, + formatter: Formatter | None = None, + max_records: int = 100, + max_bytes: int = 1_000_000, + ) -> None: + super().__init__( + topic=topic, formatter=formatter, max_records=max_records, max_bytes=max_bytes + ) + self._injected = client + self._client: PubSubTopicLike | None = None + + def _resolved_client(self) -> PubSubTopicLike: + if self._injected is not None: + return self._injected + if self._client is None: + self._client = self._import_client() + return self._client + + def _import_client(self) -> PubSubTopicLike: + try: + from google.cloud import pubsub_v1 # type: ignore # stub availability varies by env + except ImportError as exc: + raise ImportError( + "PubSubTransport: install `google-cloud-pubsub` to use " + "this transport without providing a client — " + "`pip install logquill[pubsub]`" + ) from exc + return cast(PubSubTopicLike, pubsub_v1.PublisherClient()) + + def _publish_batch(self, records: Sequence[LogRecord]) -> None: + client = self._resolved_client() + futures = [ + client.publish( + self.topic, data=json.dumps(record, separators=(",", ":"), default=str).encode() + ) + for record in records + ] + for future in futures: + future.result() diff --git a/logquill/transports/queue/rabbitmq_transport.py b/logquill/transports/queue/rabbitmq_transport.py new file mode 100644 index 0000000..83181b2 --- /dev/null +++ b/logquill/transports/queue/rabbitmq_transport.py @@ -0,0 +1,74 @@ +from __future__ import annotations + +import json +from typing import Protocol, Sequence, cast + +from logquill.formatter import Formatter +from logquill.records import LogRecord +from logquill.transports.queue.base_queue_transport import BaseQueueTransport + + +class AMQPChannelLike(Protocol): + def basic_publish(self, exchange: str, routing_key: str, body: bytes) -> object: ... + + +class AMQPConnectionLike(Protocol): + def close(self) -> None: ... + + +class RabbitMQTransport(BaseQueueTransport): + """Publishes to RabbitMQ via the optional `pika` peer dependency. One + `basic_publish` call per record on the default exchange (AMQP has no + native multi-message batch form). Pass `channel` to inject a pre-built + `pika` channel (or a fake, for tests), or `url` to let this transport + connect itself.""" + + def __init__( + self, + *, + topic: str, + channel: AMQPChannelLike | None = None, + url: str = "amqp://guest:guest@localhost:5672/%2F", + formatter: Formatter | None = None, + max_records: int = 100, + max_bytes: int = 1_000_000, + ) -> None: + super().__init__( + topic=topic, formatter=formatter, max_records=max_records, max_bytes=max_bytes + ) + self._injected = channel + self._url = url + self._channel: AMQPChannelLike | None = None + self._connection: AMQPConnectionLike | None = None + + def _resolved_channel(self) -> AMQPChannelLike: + if self._injected is not None: + return self._injected + if self._channel is None: + self._channel = self._import_channel() + return self._channel + + def _import_channel(self) -> AMQPChannelLike: + try: + import pika # type: ignore # stub availability for this optional dep varies by env + except ImportError as exc: + raise ImportError( + "RabbitMQTransport: install `pika` to use this transport " + "without providing a channel — `pip install logquill[rabbitmq]`" + ) from exc + connection = pika.BlockingConnection(pika.URLParameters(self._url)) + self._connection = cast(AMQPConnectionLike, connection) + channel = connection.channel() + channel.queue_declare(queue=self.topic) + return cast(AMQPChannelLike, channel) + + def close(self) -> None: + super().close() + if self._injected is None and self._connection is not None: + self._connection.close() + + def _publish_batch(self, records: Sequence[LogRecord]) -> None: + channel = self._resolved_channel() + for record in records: + body = json.dumps(record, separators=(",", ":"), default=str).encode("utf-8") + channel.basic_publish(exchange="", routing_key=self.topic, body=body) diff --git a/logquill/transports/queue/sqs_transport.py b/logquill/transports/queue/sqs_transport.py new file mode 100644 index 0000000..2652555 --- /dev/null +++ b/logquill/transports/queue/sqs_transport.py @@ -0,0 +1,74 @@ +from __future__ import annotations + +import json +from typing import Protocol, Sequence, cast + +from logquill.formatter import Formatter +from logquill.records import LogRecord +from logquill.transports.queue.base_queue_transport import BaseQueueTransport + +_SQS_BATCH_LIMIT = 10 + + +class SQSClientLike(Protocol): + def send_message_batch(self, QueueUrl: str, Entries: Sequence[dict[str, str]]) -> object: ... # noqa: N803 + + +class SQSTransport(BaseQueueTransport): + """Publishes to SQS via batched `send_message_batch` calls, chunked at + the API's 10-message cap, through the optional `boto3` peer dependency + (shared with `CloudWatchTransport`/`DynamoDBTransport`, since boto3 is + one monolithic package). Chunks are dispatched **sequentially**, not + concurrently — this project's dispatch is still fully synchronous + project-wide (no non-blocking async worker exists yet), so true + concurrent chunk dispatch is deferred until one does, rather than + hand-rolled here with threads. `topic` is the queue URL. Pass `client` + to inject a pre-built `boto3` SQS client (or a fake, for tests), or + `region` to let this transport connect itself.""" + + def __init__( + self, + *, + topic: str, + client: SQSClientLike | None = None, + region: str | None = None, + formatter: Formatter | None = None, + max_records: int = 100, + max_bytes: int = 1_000_000, + ) -> None: + super().__init__( + topic=topic, formatter=formatter, max_records=max_records, max_bytes=max_bytes + ) + self._injected = client + self._region = region + self._client: SQSClientLike | None = None + + def _resolved_client(self) -> SQSClientLike: + if self._injected is not None: + return self._injected + if self._client is None: + self._client = self._import_client() + return self._client + + def _import_client(self) -> SQSClientLike: + try: + import boto3 # type: ignore # stub availability for this optional dep varies by env + except ImportError as exc: + raise ImportError( + "SQSTransport: install `boto3` to use this transport " + "without providing a client — `pip install logquill[aws]`" + ) from exc + return cast(SQSClientLike, boto3.client("sqs", region_name=self._region)) + + def _publish_batch(self, records: Sequence[LogRecord]) -> None: + client = self._resolved_client() + for start in range(0, len(records), _SQS_BATCH_LIMIT): + chunk = records[start : start + _SQS_BATCH_LIMIT] + entries = [ + { + "Id": str(index), + "MessageBody": json.dumps(record, separators=(",", ":"), default=str), + } + for index, record in enumerate(chunk) + ] + client.send_message_batch(QueueUrl=self.topic, Entries=entries) diff --git a/logquill/transports/sql/__init__.py b/logquill/transports/sql/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/logquill/transports/sql/base_sql_transport.py b/logquill/transports/sql/base_sql_transport.py new file mode 100644 index 0000000..194a2a5 --- /dev/null +++ b/logquill/transports/sql/base_sql_transport.py @@ -0,0 +1,103 @@ +from __future__ import annotations + +import json +from abc import abstractmethod +from typing import Any, Sequence, TypedDict + +from logquill.formatter import Formatter +from logquill.records import LogRecord +from logquill.transports.batching_transport import BatchingTransport + + +class SQLLogRow(TypedDict): + """One row of the fixed `logs` table schema every SQL transport writes.""" + + timestamp: str + level: str + logger: str + message: str + meta: str + run_id: str | None + span_id: str | None + parent_span_id: str | None + trace_id: str | None + + +def _meta_str(meta: dict[str, Any], key: str) -> str | None: + value = meta.get(key) + return value if isinstance(value, str) else None + + +class BaseSQLTransport(BatchingTransport[SQLLogRow]): + """Shared base for every SQL sink: a fixed `logs` table schema and + always-batched inserts — never one `INSERT` per log call, that defeats + the point of buffering. + + Schema is never auto-created in production: `ensure_schema=True` is a + dev/test-only convenience, and even then the DDL runs exactly once, + before the first batch send, regardless of how many synchronous + re-entrant flushes happen back to back (the "done" flag is set *before* + the DDL runs, not after, so a second flush triggered while the first is + still in flight can't re-run it). + """ + + def __init__( + self, + *, + formatter: Formatter | None = None, + max_records: int = 100, + max_bytes: int = 1_000_000, + table_name: str = "logs", + ensure_schema: bool = False, + ) -> None: + super().__init__(formatter=formatter, max_records=max_records, max_bytes=max_bytes) + self.table_name = table_name + self.ensure_schema = ensure_schema + self._schema_ensured = False + + def _to_item(self, formatted: str, record: LogRecord) -> SQLLogRow: + meta = record["meta"] + return SQLLogRow( + timestamp=record["timestamp"], + level=record["level"], + logger=record["logger"], + message=record["message"], + meta=json.dumps(meta, separators=(",", ":"), default=str), + run_id=_meta_str(meta, "run_id"), + span_id=_meta_str(meta, "span_id"), + parent_span_id=_meta_str(meta, "parent_span_id"), + trace_id=_meta_str(meta, "trace_id"), + ) + + def _size_of(self, item: SQLLogRow) -> int: + return len(item["message"].encode("utf-8")) + len(item["meta"].encode("utf-8")) + 96 + + def _send_batch(self, batch: Sequence[SQLLogRow]) -> None: + if self.ensure_schema and not self._schema_ensured: + self._schema_ensured = True + self._ensure_table() + self._insert_rows(batch) + + def create_table_sql(self) -> str: + """Dialect-generic `CREATE TABLE IF NOT EXISTS` DDL. Subclasses + override for dialect-correct column types (e.g. Postgres `JSONB`).""" + return ( + f"CREATE TABLE IF NOT EXISTS {self.table_name} (" + "id INTEGER PRIMARY KEY AUTOINCREMENT, " + "timestamp TEXT NOT NULL, " + "level TEXT NOT NULL, " + "logger TEXT NOT NULL, " + "message TEXT NOT NULL, " + "meta TEXT NOT NULL, " + "run_id TEXT, " + "span_id TEXT, " + "parent_span_id TEXT, " + "trace_id TEXT" + ")" + ) + + @abstractmethod + def _ensure_table(self) -> None: ... + + @abstractmethod + def _insert_rows(self, rows: Sequence[SQLLogRow]) -> None: ... diff --git a/logquill/transports/sql/mysql_transport.py b/logquill/transports/sql/mysql_transport.py new file mode 100644 index 0000000..2596006 --- /dev/null +++ b/logquill/transports/sql/mysql_transport.py @@ -0,0 +1,132 @@ +from __future__ import annotations + +from typing import Protocol, Sequence, cast + +from logquill.formatter import Formatter +from logquill.transports.sql.base_sql_transport import BaseSQLTransport, SQLLogRow + + +class MySQLCursorLike(Protocol): + def execute(self, sql: str, parameters: Sequence[object] = ()) -> object: ... + + +class MySQLConnectionLike(Protocol): + def cursor(self) -> MySQLCursorLike: ... + def commit(self) -> None: ... + def close(self) -> None: ... + + +class MySQLTransport(BaseSQLTransport): + """Batches records into one parameterized multi-row `INSERT` per flush, + via the optional `pymysql` peer dependency (pure Python — no C toolchain + needed, unlike `mysqlclient`). Pass `dsn`-equivalent connection kwargs + via a pre-built `connection`, or let this transport connect itself from + `host`/`user`/`password`/`database`.""" + + def __init__( + self, + *, + connection: MySQLConnectionLike | None = None, + host: str = "localhost", + port: int = 3306, + user: str = "root", + password: str = "", + database: str = "logquill", + formatter: Formatter | None = None, + max_records: int = 100, + max_bytes: int = 1_000_000, + table_name: str = "logs", + ensure_schema: bool = False, + ) -> None: + super().__init__( + formatter=formatter, + max_records=max_records, + max_bytes=max_bytes, + table_name=table_name, + ensure_schema=ensure_schema, + ) + self._injected = connection + self._host = host + self._port = port + self._user = user + self._password = password + self._database = database + self._connection: MySQLConnectionLike | None = None + + def create_table_sql(self) -> str: + return ( + f"CREATE TABLE IF NOT EXISTS {self.table_name} (" + "id INT AUTO_INCREMENT PRIMARY KEY, " + "timestamp TEXT NOT NULL, " + "level VARCHAR(16) NOT NULL, " + "logger TEXT NOT NULL, " + "message TEXT NOT NULL, " + "meta JSON NOT NULL, " + "run_id TEXT, " + "span_id TEXT, " + "parent_span_id TEXT, " + "trace_id TEXT" + ")" + ) + + def _resolved_connection(self) -> MySQLConnectionLike: + if self._injected is not None: + return self._injected + if self._connection is None: + self._connection = self._import_connection() + return self._connection + + def _import_connection(self) -> MySQLConnectionLike: + try: + import pymysql # type: ignore # stub availability for this optional dep varies by env + except ImportError as exc: + raise ImportError( + "MySQLTransport: install `pymysql` to use this transport " + "without providing a connection — `pip install logquill[mysql]`" + ) from exc + return cast( + MySQLConnectionLike, + pymysql.connect( + host=self._host, + port=self._port, + user=self._user, + password=self._password, + database=self._database, + ), + ) + + def close(self) -> None: + super().close() + if self._injected is None and self._connection is not None: + self._connection.close() + + def _ensure_table(self) -> None: + connection = self._resolved_connection() + connection.cursor().execute(self.create_table_sql()) + connection.commit() + + def _insert_rows(self, rows: Sequence[SQLLogRow]) -> None: + connection = self._resolved_connection() + values_sql = ", ".join(["(%s, %s, %s, %s, %s, %s, %s, %s, %s)"] * len(rows)) + sql = ( + f"INSERT INTO {self.table_name} " + "(timestamp, level, logger, message, meta, run_id, span_id, parent_span_id, trace_id) " + f"VALUES {values_sql}" + ) + params: list[object] = [] + for row in rows: + params.extend( + [ + row["timestamp"], + row["level"], + row["logger"], + row["message"], + row["meta"], + row["run_id"], + row["span_id"], + row["parent_span_id"], + row["trace_id"], + ] + ) + connection.cursor().execute(sql, params) + connection.commit() diff --git a/logquill/transports/sql/postgres_transport.py b/logquill/transports/sql/postgres_transport.py new file mode 100644 index 0000000..7358b28 --- /dev/null +++ b/logquill/transports/sql/postgres_transport.py @@ -0,0 +1,116 @@ +from __future__ import annotations + +from typing import Protocol, Sequence, cast + +from logquill.formatter import Formatter +from logquill.transports.sql.base_sql_transport import BaseSQLTransport, SQLLogRow + + +class PostgresCursorLike(Protocol): + def execute(self, sql: str, parameters: Sequence[object] = ()) -> object: ... + + +class PostgresConnectionLike(Protocol): + def cursor(self) -> PostgresCursorLike: ... + def commit(self) -> None: ... + def close(self) -> None: ... + + +class PostgresTransport(BaseSQLTransport): + """Batches records into one parameterized multi-row `INSERT` per flush, + via the optional `psycopg2-binary` peer dependency (raw DB-API, not an + ORM — avoids ORM overhead on a pure insert-heavy path). Pass `dsn` to + let this transport connect itself, or inject an already-open + `connection` for tests or an alternate setup.""" + + def __init__( + self, + *, + connection: PostgresConnectionLike | None = None, + dsn: str | None = None, + formatter: Formatter | None = None, + max_records: int = 100, + max_bytes: int = 1_000_000, + table_name: str = "logs", + ensure_schema: bool = False, + ) -> None: + super().__init__( + formatter=formatter, + max_records=max_records, + max_bytes=max_bytes, + table_name=table_name, + ensure_schema=ensure_schema, + ) + self._injected = connection + self._dsn = dsn + self._connection: PostgresConnectionLike | None = None + + def create_table_sql(self) -> str: + return ( + f"CREATE TABLE IF NOT EXISTS {self.table_name} (" + "id SERIAL PRIMARY KEY, " + "timestamp TEXT NOT NULL, " + "level TEXT NOT NULL, " + "logger TEXT NOT NULL, " + "message TEXT NOT NULL, " + "meta JSONB NOT NULL, " + "run_id TEXT, " + "span_id TEXT, " + "parent_span_id TEXT, " + "trace_id TEXT" + ")" + ) + + def _resolved_connection(self) -> PostgresConnectionLike: + if self._injected is not None: + return self._injected + if self._connection is None: + self._connection = self._import_connection() + return self._connection + + def _import_connection(self) -> PostgresConnectionLike: + try: + import psycopg2 # type: ignore # stub availability for this optional dep varies by env + except ImportError as exc: + raise ImportError( + "PostgresTransport: install `psycopg2-binary` to use this " + "transport without providing a connection — " + "`pip install logquill[postgres]`" + ) from exc + return cast(PostgresConnectionLike, psycopg2.connect(self._dsn)) + + def close(self) -> None: + super().close() + if self._injected is None and self._connection is not None: + self._connection.close() + + def _ensure_table(self) -> None: + connection = self._resolved_connection() + connection.cursor().execute(self.create_table_sql()) + connection.commit() + + def _insert_rows(self, rows: Sequence[SQLLogRow]) -> None: + connection = self._resolved_connection() + values_sql = ", ".join(["(%s, %s, %s, %s, %s::jsonb, %s, %s, %s, %s)"] * len(rows)) + sql = ( + f"INSERT INTO {self.table_name} " + "(timestamp, level, logger, message, meta, run_id, span_id, parent_span_id, trace_id) " + f"VALUES {values_sql}" + ) + params: list[object] = [] + for row in rows: + params.extend( + [ + row["timestamp"], + row["level"], + row["logger"], + row["message"], + row["meta"], + row["run_id"], + row["span_id"], + row["parent_span_id"], + row["trace_id"], + ] + ) + connection.cursor().execute(sql, params) + connection.commit() diff --git a/logquill/transports/sql/sqlite_transport.py b/logquill/transports/sql/sqlite_transport.py new file mode 100644 index 0000000..bea301f --- /dev/null +++ b/logquill/transports/sql/sqlite_transport.py @@ -0,0 +1,86 @@ +from __future__ import annotations + +import sqlite3 +from typing import Iterable, Protocol, Sequence + +from logquill.formatter import Formatter +from logquill.transports.sql.base_sql_transport import BaseSQLTransport, SQLLogRow + + +class SQLiteConnectionLike(Protocol): + def execute(self, sql: str, parameters: Sequence[object] = ()) -> object: ... + def executemany(self, sql: str, seq_of_parameters: Iterable[Sequence[object]]) -> object: ... + def commit(self) -> None: ... + def close(self) -> None: ... + + +class SQLiteTransport(BaseSQLTransport): + """Zero-setup SQL sink via stdlib `sqlite3` — no server process, no + optional dependency, works against a file path or `:memory:`. Pass + `connection` to inject an already-open connection for tests or an + alternate setup.""" + + def __init__( + self, + *, + connection: SQLiteConnectionLike | None = None, + filename: str = ":memory:", + formatter: Formatter | None = None, + max_records: int = 100, + max_bytes: int = 1_000_000, + table_name: str = "logs", + ensure_schema: bool = False, + ) -> None: + super().__init__( + formatter=formatter, + max_records=max_records, + max_bytes=max_bytes, + table_name=table_name, + ensure_schema=ensure_schema, + ) + self._injected = connection + self._filename = filename + self._connection: SQLiteConnectionLike | None = None + + def _resolved_connection(self) -> SQLiteConnectionLike: + if self._injected is not None: + return self._injected + if self._connection is None: + self._connection = sqlite3.connect(self._filename) + return self._connection + + def close(self) -> None: + super().close() + if self._injected is None and self._connection is not None: + self._connection.close() + + def _ensure_table(self) -> None: + connection = self._resolved_connection() + connection.execute(self.create_table_sql()) + connection.commit() + + def _insert_rows(self, rows: Sequence[SQLLogRow]) -> None: + connection = self._resolved_connection() + sql = ( + f"INSERT INTO {self.table_name} " + "(timestamp, level, logger, message, meta, run_id, span_id, parent_span_id, trace_id) " + "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)" + ) + connection.executemany( + sql, + [ + ( + row["timestamp"], + row["level"], + row["logger"], + row["message"], + row["meta"], + row["run_id"], + row["span_id"], + row["parent_span_id"], + row["trace_id"], + ) + for row in rows + ], + ) + connection.commit() diff --git a/logquill/transport.py b/logquill/transports/transport.py similarity index 100% rename from logquill/transport.py rename to logquill/transports/transport.py diff --git a/pyproject.toml b/pyproject.toml index b4f0b01..d963fd6 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -39,6 +39,15 @@ dependencies = [] [project.optional-dependencies] http = ["aiohttp>=3.9"] +postgres = ["psycopg2-binary>=2.9"] +mysql = ["pymysql>=1.1"] +mongodb = ["pymongo>=4.6"] +redis = ["redis>=5.0"] +kafka = ["kafka-python>=2.0"] +rabbitmq = ["pika>=1.3"] +pubsub = ["google-cloud-pubsub>=2.21"] +gcp-logging = ["google-cloud-logging>=3.10"] +aws = ["boto3>=1.34"] dev = [ "ruff>=0.6", "mypy>=1.11", diff --git a/tests/test_plugins/__init__.py b/tests/test_plugins/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/tests/test_context_plugin.py b/tests/test_plugins/test_context_plugin.py similarity index 90% rename from tests/test_context_plugin.py rename to tests/test_plugins/test_context_plugin.py index 79e1c33..f7f1baf 100644 --- a/tests/test_context_plugin.py +++ b/tests/test_plugins/test_context_plugin.py @@ -1,5 +1,5 @@ -from logquill.context_plugin import ContextPlugin from logquill.logger import Logger +from logquill.plugins.context_plugin import ContextPlugin def test_injects_fixed_context_into_meta() -> None: diff --git a/tests/test_plugin.py b/tests/test_plugins/test_plugin.py similarity index 96% rename from tests/test_plugin.py rename to tests/test_plugins/test_plugin.py index 3fb6f0d..f7da2a8 100644 --- a/tests/test_plugin.py +++ b/tests/test_plugins/test_plugin.py @@ -1,9 +1,9 @@ from __future__ import annotations from logquill.logger import Logger -from logquill.plugin import Plugin +from logquill.plugins.plugin import Plugin from logquill.records import LogRecord -from logquill.transport import CollectingTransport +from logquill.transports.transport import CollectingTransport class UppercasePlugin(Plugin): diff --git a/tests/test_redact_plugin.py b/tests/test_plugins/test_redact_plugin.py similarity index 94% rename from tests/test_redact_plugin.py rename to tests/test_plugins/test_redact_plugin.py index 006c29e..cce9d4c 100644 --- a/tests/test_redact_plugin.py +++ b/tests/test_plugins/test_redact_plugin.py @@ -1,5 +1,5 @@ from logquill.logger import Logger -from logquill.redact_plugin import RedactPlugin +from logquill.plugins.redact_plugin import RedactPlugin def test_redacts_default_sensitive_keys() -> None: diff --git a/tests/test_sampling_plugin.py b/tests/test_plugins/test_sampling_plugin.py similarity index 93% rename from tests/test_sampling_plugin.py rename to tests/test_plugins/test_sampling_plugin.py index 2fcaee2..4df7456 100644 --- a/tests/test_sampling_plugin.py +++ b/tests/test_plugins/test_sampling_plugin.py @@ -1,7 +1,7 @@ import pytest from logquill.logger import Logger -from logquill.sampling_plugin import SamplingPlugin +from logquill.plugins.sampling_plugin import SamplingPlugin def test_rate_zero_drops_everything() -> None: diff --git a/tests/test_transports/__init__.py b/tests/test_transports/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/tests/test_transports/cloud/__init__.py b/tests/test_transports/cloud/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/tests/test_transports/cloud/test_app_insights_transport.py b/tests/test_transports/cloud/test_app_insights_transport.py new file mode 100644 index 0000000..a0a9235 --- /dev/null +++ b/tests/test_transports/cloud/test_app_insights_transport.py @@ -0,0 +1,38 @@ +from __future__ import annotations + +import json +from typing import Sequence + +from logquill.logger import Logger +from logquill.transports.cloud.app_insights_transport import AppInsightsTransport + + +class FakeSender: + def __init__(self) -> None: + self.calls: list[tuple[str, Sequence[str]]] = [] + + def __call__(self, url: str, batch: Sequence[str]) -> None: + self.calls.append((url, batch)) + + +def test_sends_ndjson_trace_envelopes_with_instrumentation_key() -> None: + sender = FakeSender() + transport = AppInsightsTransport(instrumentation_key="ikey-123", sender=sender, max_records=2) + logger = Logger("app.test", transports=[transport]) + + logger.warn("careful", user_id="u-1") + logger.fatal("boom") + + assert len(sender.calls) == 1 + url, lines = sender.calls[0] + assert url == AppInsightsTransport.URL + assert len(lines) == 2 + + first = json.loads(lines[0]) + assert first["iKey"] == "ikey-123" + assert first["data"]["baseData"]["message"] == "careful" + assert first["data"]["baseData"]["severityLevel"] == 2 # WARN + assert first["data"]["baseData"]["properties"]["user_id"] == "u-1" + + second = json.loads(lines[1]) + assert second["data"]["baseData"]["severityLevel"] == 4 # FATAL diff --git a/tests/test_transports/cloud/test_cloud_logging_transport.py b/tests/test_transports/cloud/test_cloud_logging_transport.py new file mode 100644 index 0000000..eecdb69 --- /dev/null +++ b/tests/test_transports/cloud/test_cloud_logging_transport.py @@ -0,0 +1,41 @@ +from __future__ import annotations + +import logging +from typing import Any + +import pytest + +from logquill.logger import Logger +from logquill.transports.cloud.cloud_logging_transport import CloudLoggingTransport + + +class FakeCloudLoggingClient: + def __init__(self) -> None: + self.calls: list[tuple[dict[str, Any], str]] = [] + + def log_struct(self, info: dict[str, Any], severity: str) -> None: + self.calls.append((info, severity)) + + +def test_maps_level_onto_severity_scale() -> None: + client = FakeCloudLoggingClient() + transport = CloudLoggingTransport(client=client, max_records=1) + logger = Logger("app.test", transports=[transport]) + + logger.warn("careful") + logger.fatal("boom") + + assert client.calls[0][1] == "WARNING" + assert client.calls[1][1] == "CRITICAL" + + +def test_missing_dependency_logs_actionable_install_hint( + caplog: pytest.LogCaptureFixture, +) -> None: + transport = CloudLoggingTransport(max_records=1) + logger = Logger("app.test", transports=[transport]) + + with caplog.at_level(logging.ERROR, logger="logquill"): + logger.info("hello") # must not raise + + assert "pip install logquill[gcp-logging]" in caplog.text diff --git a/tests/test_transports/cloud/test_cloudwatch_transport.py b/tests/test_transports/cloud/test_cloudwatch_transport.py new file mode 100644 index 0000000..c061d2a --- /dev/null +++ b/tests/test_transports/cloud/test_cloudwatch_transport.py @@ -0,0 +1,52 @@ +from __future__ import annotations + +import logging +from typing import Any + +import pytest + +from logquill.logger import Logger +from logquill.transports.cloud.cloudwatch_transport import CloudWatchTransport + + +class FakeCloudWatchClient: + def __init__(self) -> None: + self.calls: list[tuple[str, str, list[dict[str, Any]]]] = [] + + def put_log_events( + self, + logGroupName: str, + logStreamName: str, + logEvents: list[dict[str, Any]], # noqa: N803 + ) -> None: + self.calls.append((logGroupName, logStreamName, logEvents)) + + +def test_puts_events_sorted_by_timestamp() -> None: + client = FakeCloudWatchClient() + transport = CloudWatchTransport( + log_group="my-group", log_stream="my-stream", client=client, max_records=2 + ) + logger = Logger("app.test", transports=[transport]) + + logger.info("first") + logger.info("second") + + assert len(client.calls) == 1 + group, stream, events = client.calls[0] + assert group == "my-group" + assert stream == "my-stream" + assert len(events) == 2 + assert events[0]["timestamp"] <= events[1]["timestamp"] + + +def test_missing_dependency_logs_actionable_install_hint( + caplog: pytest.LogCaptureFixture, +) -> None: + transport = CloudWatchTransport(log_group="g", log_stream="s", max_records=1) + logger = Logger("app.test", transports=[transport]) + + with caplog.at_level(logging.ERROR, logger="logquill"): + logger.info("hello") # must not raise + + assert "pip install logquill[aws]" in caplog.text diff --git a/tests/test_transports/cloud/test_datadog_transport.py b/tests/test_transports/cloud/test_datadog_transport.py new file mode 100644 index 0000000..54bc018 --- /dev/null +++ b/tests/test_transports/cloud/test_datadog_transport.py @@ -0,0 +1,49 @@ +from __future__ import annotations + +import logging +from typing import Sequence + +import pytest + +from logquill.logger import Logger +from logquill.transports.cloud.datadog_transport import DatadogTransport + + +class FakeSender: + def __init__(self) -> None: + self.calls: list[tuple[str, str, Sequence[str]]] = [] + + def __call__(self, url: str, api_key: str, batch: Sequence[str]) -> None: + self.calls.append((url, api_key, batch)) + + +def test_url_is_region_specific_and_batch_is_sent() -> None: + sender = FakeSender() + transport = DatadogTransport( + api_key="key-123", site="datadoghq.eu", sender=sender, max_records=2 + ) + logger = Logger("app.test", transports=[transport]) + + logger.info("first") + logger.info("second") + + assert len(sender.calls) == 1 + url, api_key, batch = sender.calls[0] + assert url == "https://http-intake.logs.datadoghq.eu/api/v2/logs" + assert api_key == "key-123" + assert len(batch) == 2 + + +def test_failing_send_is_caught_and_logged_not_raised( + caplog: pytest.LogCaptureFixture, +) -> None: + def failing_sender(url: str, api_key: str, batch: Sequence[str]) -> None: + raise RuntimeError(f"DatadogTransport: request to {url} failed with status 403") + + transport = DatadogTransport(api_key="bad-key", sender=failing_sender, max_records=1) + logger = Logger("app.test", transports=[transport]) + + with caplog.at_level(logging.ERROR, logger="logquill"): + logger.info("hello") # must not raise + + assert "failed with status 403" in caplog.text diff --git a/tests/test_transports/cloud/test_elasticsearch_transport.py b/tests/test_transports/cloud/test_elasticsearch_transport.py new file mode 100644 index 0000000..9b2b56d --- /dev/null +++ b/tests/test_transports/cloud/test_elasticsearch_transport.py @@ -0,0 +1,34 @@ +from __future__ import annotations + +import json +from typing import Sequence + +from logquill.logger import Logger +from logquill.transports.cloud.elasticsearch_transport import ElasticsearchTransport + + +class FakeSender: + def __init__(self) -> None: + self.calls: list[tuple[str, Sequence[str]]] = [] + + def __call__(self, url: str, batch: Sequence[str]) -> None: + self.calls.append((url, batch)) + + +def test_sends_ndjson_action_and_source_pairs_to_bulk_endpoint() -> None: + sender = FakeSender() + transport = ElasticsearchTransport( + url="http://localhost:9200/", index="app-logs", sender=sender, max_records=1 + ) + logger = Logger("app.test", transports=[transport]) + + logger.info("hello") + + assert len(sender.calls) == 1 + url, lines = sender.calls[0] + assert url == "http://localhost:9200/_bulk" + assert len(lines) == 2 + action = json.loads(lines[0]) + assert action == {"index": {"_index": "app-logs"}} + source = json.loads(lines[1]) + assert source["message"] == "hello" diff --git a/tests/test_transports/cloud/test_new_relic_transport.py b/tests/test_transports/cloud/test_new_relic_transport.py new file mode 100644 index 0000000..ccb9744 --- /dev/null +++ b/tests/test_transports/cloud/test_new_relic_transport.py @@ -0,0 +1,84 @@ +from __future__ import annotations + +import gzip +import json +from typing import Callable + +from logquill.logger import Logger +from logquill.transports.cloud.new_relic_transport import ( + NewRelicSenderResult, + NewRelicTransport, +) + + +class FakeSender: + def __init__(self) -> None: + self.calls: list[tuple[str, dict[str, str], bytes]] = [] + self.results: list[NewRelicSenderResult] = [] + + def __call__(self, url: str, headers: dict[str, str], body: bytes) -> NewRelicSenderResult: + self.calls.append((url, headers, body)) + if self.results: + return self.results.pop(0) + return {"ok": True, "status": 202, "retry_after": None} + + +def _clock(values: list[float]) -> Callable[[], float]: + def clock() -> float: + return values.pop(0) + + return clock + + +def test_region_url_gzip_body_and_headers() -> None: + sender = FakeSender() + transport = NewRelicTransport(license_key="lk-1", region="EU", sender=sender, max_records=1) + logger = Logger("app.test", transports=[transport]) + + logger.info("hello", eventType="Custom", extra="kept") + + assert len(sender.calls) == 1 + url, headers, body = sender.calls[0] + assert url == "https://log-api.eu.newrelic.com/log/v1" + assert headers["Api-Key"] == "lk-1" + assert headers["Content-Encoding"] == "gzip" + + records = json.loads(gzip.decompress(body)) + assert "eventType" not in records[0]["meta"] + assert records[0]["meta"]["extra"] == "kept" + + +def test_429_pauses_sends_and_drops_batches_during_the_window() -> None: + sender = FakeSender() + sender.results = [{"ok": False, "status": 429, "retry_after": "60"}] + clock_values = [1000.0, 1010.0, 1070.0] + transport = NewRelicTransport( + license_key="lk-1", sender=sender, clock=_clock(clock_values), max_records=1 + ) + logger = Logger("app.test", transports=[transport]) + + logger.info("first") # 429 at t=1000, pause until t=1060 + assert len(sender.calls) == 1 + + logger.info("second") # t=1010, still paused -> dropped, no send + assert len(sender.calls) == 1 + + logger.info("third") # t=1070, past pause -> sends again + assert len(sender.calls) == 2 + + +def test_missing_retry_after_defaults_to_sixty_second_pause() -> None: + sender = FakeSender() + sender.results = [{"ok": False, "status": 429, "retry_after": None}] + clock_values = [0.0, 59.0, 61.0] + transport = NewRelicTransport( + license_key="lk-1", sender=sender, clock=_clock(clock_values), max_records=1 + ) + logger = Logger("app.test", transports=[transport]) + + logger.info("first") + assert len(sender.calls) == 1 + logger.info("second") # still within default 60s pause + assert len(sender.calls) == 1 + logger.info("third") # past the 60s pause + assert len(sender.calls) == 2 diff --git a/tests/test_transports/nosql/__init__.py b/tests/test_transports/nosql/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/tests/test_transports/nosql/test_dynamodb_transport.py b/tests/test_transports/nosql/test_dynamodb_transport.py new file mode 100644 index 0000000..14b678c --- /dev/null +++ b/tests/test_transports/nosql/test_dynamodb_transport.py @@ -0,0 +1,69 @@ +from __future__ import annotations + +import logging +from typing import Any + +import pytest + +from logquill.logger import Logger +from logquill.transports.nosql.dynamodb_transport import DynamoDBTransport + + +class FakeDynamoTable: + def __init__(self) -> None: + self.put_items: list[dict[str, Any]] = [] + self.entered = 0 + self.exited = 0 + + def batch_writer(self) -> FakeDynamoTable: + return self + + def __enter__(self) -> FakeDynamoTable: + self.entered += 1 + return self + + def __exit__(self, *args: object) -> bool: + self.exited += 1 + return False + + def put_item(self, Item: dict[str, Any]) -> None: # noqa: N803 + self.put_items.append(Item) + + +def test_partition_key_prefers_run_id_over_trace_id_over_logger() -> None: + table = FakeDynamoTable() + transport = DynamoDBTransport(table=table, max_records=3) + logger = Logger("app.test", transports=[transport]) + + logger.info("has run_id", run_id="run-1", trace_id="trace-1") + logger.info("has trace_id only", trace_id="trace-2") + logger.info("has neither") + + assert table.put_items[0]["run_id"] == "run-1" + assert table.put_items[1]["run_id"] == "trace-2" + assert table.put_items[2]["run_id"] == "app.test" + + +def test_sort_key_is_timestamp_and_batch_writer_used_once_per_flush() -> None: + table = FakeDynamoTable() + transport = DynamoDBTransport(table=table, max_records=2) + logger = Logger("app.test", transports=[transport]) + + logger.info("first") + logger.info("second") + + assert table.entered == 1 + assert table.exited == 1 + assert all("timestamp" in item for item in table.put_items) + + +def test_missing_dependency_logs_actionable_install_hint( + caplog: pytest.LogCaptureFixture, +) -> None: + transport = DynamoDBTransport(max_records=1) + logger = Logger("app.test", transports=[transport]) + + with caplog.at_level(logging.ERROR, logger="logquill"): + logger.info("hello") # must not raise + + assert "pip install logquill[aws]" in caplog.text diff --git a/tests/test_transports/nosql/test_mongodb_transport.py b/tests/test_transports/nosql/test_mongodb_transport.py new file mode 100644 index 0000000..bccb5c1 --- /dev/null +++ b/tests/test_transports/nosql/test_mongodb_transport.py @@ -0,0 +1,42 @@ +from __future__ import annotations + +import logging +from typing import Any, Sequence + +import pytest + +from logquill.logger import Logger +from logquill.transports.nosql.mongodb_transport import MongoDBTransport + + +class FakeCollection: + def __init__(self) -> None: + self.batches: list[list[dict[str, Any]]] = [] + + def insert_many(self, documents: Sequence[dict[str, Any]]) -> None: + self.batches.append(list(documents)) + + +def test_batches_records_as_documents_one_to_one() -> None: + collection = FakeCollection() + transport = MongoDBTransport(collection=collection, max_records=2) + logger = Logger("app.test", transports=[transport]) + + logger.info("first") + logger.info("second") + + assert len(collection.batches) == 1 + assert len(collection.batches[0]) == 2 + assert collection.batches[0][0]["message"] == "first" + + +def test_missing_dependency_logs_actionable_install_hint( + caplog: pytest.LogCaptureFixture, +) -> None: + transport = MongoDBTransport(max_records=1) + logger = Logger("app.test", transports=[transport]) + + with caplog.at_level(logging.ERROR, logger="logquill"): + logger.info("hello") # must not raise + + assert "pip install logquill[mongodb]" in caplog.text diff --git a/tests/test_transports/nosql/test_redis_transport.py b/tests/test_transports/nosql/test_redis_transport.py new file mode 100644 index 0000000..dc5340d --- /dev/null +++ b/tests/test_transports/nosql/test_redis_transport.py @@ -0,0 +1,43 @@ +from __future__ import annotations + +import json +import logging +from typing import Any + +import pytest + +from logquill.logger import Logger +from logquill.transports.nosql.redis_transport import RedisTransport + + +class FakeRedisClient: + def __init__(self) -> None: + self.calls: list[tuple[str, dict[str, Any]]] = [] + + def xadd(self, name: str, fields: dict[str, Any]) -> None: + self.calls.append((name, fields)) + + +def test_one_xadd_call_per_record_in_the_batch() -> None: + client = FakeRedisClient() + transport = RedisTransport(client=client, stream="my-logs", max_records=2) + logger = Logger("app.test", transports=[transport]) + + logger.info("first") + logger.info("second") + + assert len(client.calls) == 2 + assert all(name == "my-logs" for name, _ in client.calls) + assert json.loads(client.calls[0][1]["meta"]) == {} + + +def test_missing_dependency_logs_actionable_install_hint( + caplog: pytest.LogCaptureFixture, +) -> None: + transport = RedisTransport(max_records=1) + logger = Logger("app.test", transports=[transport]) + + with caplog.at_level(logging.ERROR, logger="logquill"): + logger.info("hello") # must not raise + + assert "pip install logquill[redis]" in caplog.text diff --git a/tests/test_transports/queue/__init__.py b/tests/test_transports/queue/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/tests/test_transports/queue/test_kafka_transport.py b/tests/test_transports/queue/test_kafka_transport.py new file mode 100644 index 0000000..56d9db3 --- /dev/null +++ b/tests/test_transports/queue/test_kafka_transport.py @@ -0,0 +1,61 @@ +from __future__ import annotations + +import json +import logging + +import pytest + +from logquill.logger import Logger +from logquill.transports.queue.kafka_transport import KafkaTransport + + +class FakeKafkaProducer: + def __init__(self) -> None: + self.sent: list[tuple[str, bytes, bytes | None]] = [] + self.flush_count = 0 + + def send(self, topic: str, value: bytes, key: bytes | None = None) -> None: + self.sent.append((topic, value, key)) + + def flush(self) -> None: + self.flush_count += 1 + + +def test_publishes_batch_keyed_by_run_id() -> None: + producer = FakeKafkaProducer() + transport = KafkaTransport(topic="app-logs", producer=producer, max_records=2) + logger = Logger("app.test", transports=[transport]) + + logger.info("first", run_id="run-1") + logger.info("second", run_id="run-1") + + assert len(producer.sent) == 2 + assert all(topic == "app-logs" for topic, _, _ in producer.sent) + assert producer.sent[0][2] == b"run-1" + assert producer.flush_count == 1 + + +def test_falls_back_to_trace_id_then_no_key() -> None: + producer = FakeKafkaProducer() + transport = KafkaTransport(topic="app-logs", producer=producer, max_records=1) + logger = Logger("app.test", transports=[transport]) + + logger.info("traced only", trace_id="trace-9") + logger.info("untraced") + + assert producer.sent[0][2] == b"trace-9" + assert producer.sent[1][2] is None + payload = json.loads(producer.sent[1][1]) + assert payload["message"] == "untraced" + + +def test_missing_dependency_logs_actionable_install_hint( + caplog: pytest.LogCaptureFixture, +) -> None: + transport = KafkaTransport(topic="app-logs", max_records=1) + logger = Logger("app.test", transports=[transport]) + + with caplog.at_level(logging.ERROR, logger="logquill"): + logger.info("hello") # must not raise + + assert "pip install logquill[kafka]" in caplog.text diff --git a/tests/test_transports/queue/test_pubsub_transport.py b/tests/test_transports/queue/test_pubsub_transport.py new file mode 100644 index 0000000..422190c --- /dev/null +++ b/tests/test_transports/queue/test_pubsub_transport.py @@ -0,0 +1,55 @@ +from __future__ import annotations + +import logging + +import pytest + +from logquill.logger import Logger +from logquill.transports.queue.pubsub_transport import PubSubTransport + + +class FakeFuture: + def __init__(self) -> None: + self.waited = False + + def result(self) -> str: + self.waited = True + return "message-id" + + +class FakePubSubTopic: + def __init__(self) -> None: + self.published: list[tuple[str, bytes]] = [] + self.futures: list[FakeFuture] = [] + + def publish(self, topic: str, data: bytes) -> FakeFuture: + self.published.append((topic, data)) + future = FakeFuture() + self.futures.append(future) + return future + + +def test_publishes_and_waits_on_every_future() -> None: + topic_path = "projects/my-project/topics/app-logs" + client = FakePubSubTopic() + transport = PubSubTransport(topic=topic_path, client=client, max_records=2) + logger = Logger("app.test", transports=[transport]) + + logger.info("first") + logger.info("second") + + assert len(client.published) == 2 + assert all(topic == topic_path for topic, _ in client.published) + assert all(future.waited for future in client.futures) + + +def test_missing_dependency_logs_actionable_install_hint( + caplog: pytest.LogCaptureFixture, +) -> None: + transport = PubSubTransport(topic="projects/p/topics/t", max_records=1) + logger = Logger("app.test", transports=[transport]) + + with caplog.at_level(logging.ERROR, logger="logquill"): + logger.info("hello") # must not raise + + assert "pip install logquill[pubsub]" in caplog.text diff --git a/tests/test_transports/queue/test_rabbitmq_transport.py b/tests/test_transports/queue/test_rabbitmq_transport.py new file mode 100644 index 0000000..30cbff3 --- /dev/null +++ b/tests/test_transports/queue/test_rabbitmq_transport.py @@ -0,0 +1,41 @@ +from __future__ import annotations + +import logging + +import pytest + +from logquill.logger import Logger +from logquill.transports.queue.rabbitmq_transport import RabbitMQTransport + + +class FakeChannel: + def __init__(self) -> None: + self.published: list[tuple[str, str, bytes]] = [] + + def basic_publish(self, exchange: str, routing_key: str, body: bytes) -> None: + self.published.append((exchange, routing_key, body)) + + +def test_publishes_one_message_per_record_to_the_default_exchange() -> None: + channel = FakeChannel() + transport = RabbitMQTransport(topic="app-logs", channel=channel, max_records=2) + logger = Logger("app.test", transports=[transport]) + + logger.info("first") + logger.info("second") + + assert len(channel.published) == 2 + assert all(exchange == "" for exchange, _, _ in channel.published) + assert all(routing_key == "app-logs" for _, routing_key, _ in channel.published) + + +def test_missing_dependency_logs_actionable_install_hint( + caplog: pytest.LogCaptureFixture, +) -> None: + transport = RabbitMQTransport(topic="app-logs", max_records=1) + logger = Logger("app.test", transports=[transport]) + + with caplog.at_level(logging.ERROR, logger="logquill"): + logger.info("hello") # must not raise + + assert "pip install logquill[rabbitmq]" in caplog.text diff --git a/tests/test_transports/queue/test_sqs_transport.py b/tests/test_transports/queue/test_sqs_transport.py new file mode 100644 index 0000000..f46bdd4 --- /dev/null +++ b/tests/test_transports/queue/test_sqs_transport.py @@ -0,0 +1,60 @@ +from __future__ import annotations + +import logging + +import pytest + +from logquill.logger import Logger +from logquill.transports.queue.sqs_transport import SQSTransport + + +class FakeSQSClient: + def __init__(self) -> None: + self.calls: list[tuple[str, list[dict[str, str]]]] = [] + + def send_message_batch(self, QueueUrl: str, Entries: list[dict[str, str]]) -> None: # noqa: N803 + self.calls.append((QueueUrl, Entries)) + + +def test_chunks_at_ten_messages_per_send_message_batch_call() -> None: + client = FakeSQSClient() + transport = SQSTransport( + topic="https://sqs.us-east-1.amazonaws.com/123/my-queue", + client=client, + max_records=25, + ) + logger = Logger("app.test", transports=[transport]) + + for i in range(25): + logger.info(f"message {i}") + + expected_url = "https://sqs.us-east-1.amazonaws.com/123/my-queue" + assert len(client.calls) == 3 + assert [len(entries) for _, entries in client.calls] == [10, 10, 5] + assert all(queue_url == expected_url for queue_url, _ in client.calls) + + +def test_entry_ids_reset_per_chunk() -> None: + client = FakeSQSClient() + transport = SQSTransport(topic="q", client=client, max_records=15) + logger = Logger("app.test", transports=[transport]) + + for i in range(15): + logger.info(f"message {i}") + + first_chunk_ids = [entry["Id"] for entry in client.calls[0][1]] + second_chunk_ids = [entry["Id"] for entry in client.calls[1][1]] + assert first_chunk_ids == [str(i) for i in range(10)] + assert second_chunk_ids == [str(i) for i in range(5)] + + +def test_missing_dependency_logs_actionable_install_hint( + caplog: pytest.LogCaptureFixture, +) -> None: + transport = SQSTransport(topic="q", max_records=1) + logger = Logger("app.test", transports=[transport]) + + with caplog.at_level(logging.ERROR, logger="logquill"): + logger.info("hello") # must not raise + + assert "pip install logquill[aws]" in caplog.text diff --git a/tests/test_transports/sql/__init__.py b/tests/test_transports/sql/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/tests/test_transports/sql/test_mysql_transport.py b/tests/test_transports/sql/test_mysql_transport.py new file mode 100644 index 0000000..2aac347 --- /dev/null +++ b/tests/test_transports/sql/test_mysql_transport.py @@ -0,0 +1,55 @@ +from __future__ import annotations + +import logging + +import pytest + +from logquill.logger import Logger +from logquill.transports.sql.mysql_transport import MySQLCursorLike, MySQLTransport + + +class FakeCursor: + def __init__(self, sink: list[tuple[str, object]]) -> None: + self._sink = sink + + def execute(self, sql: str, parameters: object = ()) -> None: + self._sink.append((sql, parameters)) + + +class FakeConnection: + def __init__(self) -> None: + self.calls: list[tuple[str, object]] = [] + self.commit_count = 0 + + def cursor(self) -> MySQLCursorLike: + return FakeCursor(self.calls) + + def commit(self) -> None: + self.commit_count += 1 + + +def test_batches_multi_row_insert() -> None: + fake = FakeConnection() + transport = MySQLTransport(connection=fake, max_records=2) + logger = Logger("app.test", transports=[transport]) + + logger.info("first") + logger.info("second") + + assert len(fake.calls) == 1 + sql, params = fake.calls[0] + assert sql.startswith("INSERT INTO logs") + assert len(params) == 18 # 9 columns * 2 rows + assert fake.commit_count == 1 + + +def test_missing_dependency_logs_actionable_install_hint( + caplog: pytest.LogCaptureFixture, +) -> None: + transport = MySQLTransport(max_records=1) + logger = Logger("app.test", transports=[transport]) + + with caplog.at_level(logging.ERROR, logger="logquill"): + logger.info("hello") # must not raise + + assert "pip install logquill[mysql]" in caplog.text diff --git a/tests/test_transports/sql/test_postgres_transport.py b/tests/test_transports/sql/test_postgres_transport.py new file mode 100644 index 0000000..11db140 --- /dev/null +++ b/tests/test_transports/sql/test_postgres_transport.py @@ -0,0 +1,55 @@ +from __future__ import annotations + +import logging + +import pytest + +from logquill.logger import Logger +from logquill.transports.sql.postgres_transport import PostgresCursorLike, PostgresTransport + + +class FakeCursor: + def __init__(self, sink: list[tuple[str, object]]) -> None: + self._sink = sink + + def execute(self, sql: str, parameters: object = ()) -> None: + self._sink.append((sql, parameters)) + + +class FakeConnection: + def __init__(self) -> None: + self.calls: list[tuple[str, object]] = [] + self.commit_count = 0 + + def cursor(self) -> PostgresCursorLike: + return FakeCursor(self.calls) + + def commit(self) -> None: + self.commit_count += 1 + + +def test_batches_multi_row_insert_with_jsonb_cast() -> None: + fake = FakeConnection() + transport = PostgresTransport(connection=fake, max_records=2) + logger = Logger("app.test", transports=[transport]) + + logger.info("first") + logger.info("second") + + assert len(fake.calls) == 1 + sql, params = fake.calls[0] + assert sql.count("::jsonb") == 2 + assert len(params) == 18 # 9 columns * 2 rows + assert fake.commit_count == 1 + + +def test_missing_dependency_logs_actionable_install_hint( + caplog: pytest.LogCaptureFixture, +) -> None: + transport = PostgresTransport(max_records=1) + logger = Logger("app.test", transports=[transport]) + + with caplog.at_level(logging.ERROR, logger="logquill"): + logger.info("hello") # must not raise + + assert "pip install logquill[postgres]" in caplog.text diff --git a/tests/test_transports/sql/test_sqlite_transport.py b/tests/test_transports/sql/test_sqlite_transport.py new file mode 100644 index 0000000..0afcf49 --- /dev/null +++ b/tests/test_transports/sql/test_sqlite_transport.py @@ -0,0 +1,86 @@ +from __future__ import annotations + +from typing import Iterable, Sequence + +from logquill.logger import Logger +from logquill.transports.sql.sqlite_transport import SQLiteConnectionLike, SQLiteTransport + + +class FakeSQLiteConnection: + def __init__(self) -> None: + self.exec_calls: list[str] = [] + self.executemany_calls: list[tuple[str, list[Sequence[object]]]] = [] + self.commit_count = 0 + + def execute(self, sql: str, parameters: Sequence[object] = ()) -> None: + self.exec_calls.append(sql) + + def executemany(self, sql: str, seq_of_parameters: Iterable[Sequence[object]]) -> None: + self.executemany_calls.append((sql, list(seq_of_parameters))) + + def commit(self) -> None: + self.commit_count += 1 + + +def _fake() -> tuple[FakeSQLiteConnection, SQLiteConnectionLike]: + connection = FakeSQLiteConnection() + return connection, connection + + +def test_batches_inserts_and_never_creates_schema_by_default() -> None: + connection, injected = _fake() + transport = SQLiteTransport(connection=injected, max_records=2) + logger = Logger("app.test", transports=[transport]) + + logger.info("first") + assert connection.executemany_calls == [] + logger.info("second") + + assert len(connection.executemany_calls) == 1 + assert len(connection.executemany_calls[0][1]) == 2 + assert connection.exec_calls == [] + + +def test_ensure_schema_runs_ddl_exactly_once_across_reentrant_flushes() -> None: + connection, injected = _fake() + transport = SQLiteTransport(connection=injected, max_records=1, ensure_schema=True) + logger = Logger("app.test", transports=[transport]) + + logger.info("first") + logger.info("second") + + assert len(connection.exec_calls) == 1 + assert "CREATE TABLE IF NOT EXISTS logs" in connection.exec_calls[0] + + +def test_close_flushes_partial_batch() -> None: + connection, injected = _fake() + transport = SQLiteTransport(connection=injected, max_records=100) + logger = Logger("app.test", transports=[transport]) + + logger.info("only one") + logger.close() + + assert len(connection.executemany_calls) == 1 + assert connection.commit_count == 1 + + +def test_row_carries_run_id_and_trace_id_from_meta() -> None: + connection, injected = _fake() + transport = SQLiteTransport(connection=injected, max_records=1) + logger = Logger("app.test", transports=[transport]) + + logger.info("traced", run_id="run-1", trace_id="trace-1") + + _, params = connection.executemany_calls[0] + row = params[0] + assert row[5] == "run-1" # run_id column + assert row[8] == "trace-1" # trace_id column + + +def test_default_construction_uses_real_stdlib_sqlite3_in_memory() -> None: + transport = SQLiteTransport(max_records=1, ensure_schema=True) + logger = Logger("app.test", transports=[transport]) + + logger.info("hello") # must not raise — exercises the real sqlite3 path + logger.close() diff --git a/tests/test_transports/test_batching_transport.py b/tests/test_transports/test_batching_transport.py new file mode 100644 index 0000000..e6b914b --- /dev/null +++ b/tests/test_transports/test_batching_transport.py @@ -0,0 +1,74 @@ +from __future__ import annotations + +from typing import Sequence + +from logquill.levels import Level +from logquill.records import LogRecord, create_record +from logquill.transports.batching_transport import BatchingTransport + + +class CollectingBatchTransport(BatchingTransport[LogRecord]): + """Minimal concrete BatchingTransport for testing the shared base's + buffering/flush logic in isolation from any real backend.""" + + def __init__(self, *, max_records: int = 100, max_bytes: int = 1_000_000) -> None: + super().__init__(max_records=max_records, max_bytes=max_bytes) + self.sent_batches: list[list[LogRecord]] = [] + self.fail_next = False + + def _send_batch(self, batch: Sequence[LogRecord]) -> None: + if self.fail_next: + self.fail_next = False + raise RuntimeError("boom") + self.sent_batches.append(list(batch)) + + +def _record(message: str = "hello") -> LogRecord: + return create_record(level=Level.INFO, logger="app.test", message=message, meta={}) + + +def test_flushes_exactly_at_max_records() -> None: + transport = CollectingBatchTransport(max_records=3, max_bytes=10_000_000) + for _ in range(2): + transport.write("x", _record()) + assert transport.sent_batches == [] + + transport.write("x", _record()) + assert len(transport.sent_batches) == 1 + assert len(transport.sent_batches[0]) == 3 + + +def test_max_bytes_flushes_immediately_even_with_huge_max_records() -> None: + transport = CollectingBatchTransport(max_records=1000, max_bytes=1) + transport.write("x", _record()) + assert len(transport.sent_batches) == 1 + + +def test_close_flushes_a_partial_batch() -> None: + transport = CollectingBatchTransport(max_records=100, max_bytes=10_000_000) + transport.write("x", _record("only one")) + transport.close() + assert len(transport.sent_batches) == 1 + assert transport.sent_batches[0][0]["message"] == "only one" + + +def test_close_on_empty_buffer_sends_nothing() -> None: + transport = CollectingBatchTransport() + transport.close() + assert transport.sent_batches == [] + + +def test_failed_send_is_caught_and_does_not_raise() -> None: + transport = CollectingBatchTransport(max_records=1) + transport.fail_next = True + transport.write("x", _record()) # must not raise + assert transport.sent_batches == [] + + +def test_buffer_is_cleared_before_send_avoiding_reentrant_double_send() -> None: + transport = CollectingBatchTransport(max_records=1) + transport.write("x", _record("first")) + transport.write("x", _record("second")) + assert len(transport.sent_batches) == 2 + assert transport.sent_batches[0][0]["message"] == "first" + assert transport.sent_batches[1][0]["message"] == "second" diff --git a/tests/test_console_transport.py b/tests/test_transports/test_console_transport.py similarity index 93% rename from tests/test_console_transport.py rename to tests/test_transports/test_console_transport.py index 32e8127..dcb0479 100644 --- a/tests/test_console_transport.py +++ b/tests/test_transports/test_console_transport.py @@ -1,7 +1,7 @@ import io -from logquill.console_transport import ConsoleTransport from logquill.logger import Logger +from logquill.transports.console_transport import ConsoleTransport def test_info_writes_to_stdout_uncolored_by_default_settings() -> None: diff --git a/tests/test_file_transport.py b/tests/test_transports/test_file_transport.py similarity index 95% rename from tests/test_file_transport.py rename to tests/test_transports/test_file_transport.py index 265adb2..01d516d 100644 --- a/tests/test_file_transport.py +++ b/tests/test_transports/test_file_transport.py @@ -1,9 +1,9 @@ from pathlib import Path -from logquill.file_transport import FileTransport from logquill.levels import Level from logquill.logger import Logger from logquill.records import create_record +from logquill.transports.file_transport import FileTransport def test_writes_are_appended_to_the_file(tmp_path: Path) -> None: diff --git a/tests/test_http_transport.py b/tests/test_transports/test_http_transport.py similarity index 95% rename from tests/test_http_transport.py rename to tests/test_transports/test_http_transport.py index ea32548..1fedf85 100644 --- a/tests/test_http_transport.py +++ b/tests/test_transports/test_http_transport.py @@ -1,7 +1,7 @@ from typing import List, Sequence, Tuple -from logquill.http_transport import HTTPTransport from logquill.logger import Logger +from logquill.transports.http_transport import HTTPTransport class FakeSender: diff --git a/tests/test_transport.py b/tests/test_transports/test_transport.py similarity index 93% rename from tests/test_transport.py rename to tests/test_transports/test_transport.py index 7bc1615..c0d9c11 100644 --- a/tests/test_transport.py +++ b/tests/test_transports/test_transport.py @@ -1,5 +1,5 @@ from logquill.logger import Logger -from logquill.transport import CollectingTransport +from logquill.transports.transport import CollectingTransport def test_logger_dispatches_to_attached_transports() -> None: