Skip to content
Closed
46 changes: 44 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -35,8 +35,50 @@ logger:
## Features

- JSON event streaming to message brokers
- RabbitMQ stream queue support via STOMP queue headers
- SSL/TLS support
- Configurable formatters (default + JLab SWF schema)
- Configurable formatters (default + JLab SWF schema + ComprehensiveEventFormatter)
- Event filtering (include/exclude)
- Heartbeat management
- Environment variable support for secrets
- Optional application-level consumer heartbeat events
- Environment variable support for secrets

## Consumer Heartbeat Events

In addition to STOMP transport heartbeats, you can emit periodic heartbeat messages
that are delivered to consumers on the configured destination.

Set `consumer_heartbeat_interval` in seconds:

```yaml
logger:
stomp:
consumer_heartbeat_interval: 30
```

Use `0` (default) to disable consumer heartbeat events.

## RabbitMQ Streams

RabbitMQ streams are supported as a queue type when the broker is accessed via the
RabbitMQ STOMP adapter.

Example profile configuration:

```yaml
logger:
stomp:
host: "rabbitmq.example.com"
port: 61613
user: "${STOMP_USER}"
password: "${STOMP_PASSWORD}"
queue: "/queue/snakemake.stream"
use_stream: true
stream_filter_by_workflow: true
```

Notes:

- RabbitMQ stream support in this plugin is publish-only.
- RabbitMQ queue type is immutable. If a destination already exists as a classic queue,
it must be deleted and recreated as a stream before `use_stream: true` will work.
81 changes: 81 additions & 0 deletions docs/further.md
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,40 @@ snakemake --logger stomp \

The `SWF_RUN_NUMBER` environment variable can be set to populate the `correlation.run_number` field for tracking related workflow runs.

### ComprehensiveEventFormatter

Human-readable event view with rich passthrough of Snakemake event attributes.
This formatter keeps top-level context fields and adds an `event_details` section
containing nearly all event-specific data available on the log record.

```json
{
"timestamp": "2024-01-15T14:30:00Z",
"hostname": "compute-node-01",
"workflow_id": "550e8400-e29b-41d4-a716-446655440000",
"event_type": "job_finished",
"level": "INFO",
"message": "Rule completed",
"summary": "Job Finished | rule=align_reads | job=5 | wildcards={'sample': 'A'} | Rule completed",
"event_details": {
"event": "job_finished",
"jobid": 5,
"name": "align_reads",
"wildcards": {"sample": "A"},
"threads": 4,
"resources": {"mem_mb": 16000},
"benchmark": "bench/align_reads.txt"
}
}
```

**To use ComprehensiveEventFormatter:**

```bash
snakemake --logger stomp \
--logger-stomp-formatter-class "snakemake_logger_plugin_stomp.formatters.ComprehensiveEventFormatter"
```

## Custom Formatters

Create your own formatter by implementing the `BaseFormatter` interface:
Expand Down Expand Up @@ -198,6 +232,53 @@ queue: "/topic/snakemake.broadcast"

Use **queues** when you want load balancing across multiple consumers. Use **topics** when multiple systems need to receive all events (e.g., monitoring + archival + alerting).

## RabbitMQ Streams

RabbitMQ streams are a queue type exposed by the RabbitMQ STOMP adapter. This plugin
supports stream publishing and stream declaration for queue destinations.

### Creating a Stream

Use a `/queue/<name>` destination together with `use_stream: true` to have RabbitMQ
declare the destination as a stream on first publish:

```yaml
logger:
stomp:
queue: "/queue/snakemake.stream"
use_stream: true
```

### Publishing to an Existing Stream

Use `/amq/queue/<name>` when the stream already exists and should not be redeclared:

```yaml
logger:
stomp:
queue: "/amq/queue/snakemake.stream"
use_stream: true
```

### Stream Filter Headers

RabbitMQ streams support outbound filter values on published messages. Use
`stream_filter_by_workflow: true` to set `x-stream-filter-value` to the current
Snakemake `workflow_id` on each published event once the workflow ID is available:

```yaml
logger:
stomp:
queue: "/queue/snakemake.stream"
use_stream: true
stream_filter_by_workflow: true
```

### Constraints

- Stream support is publish-only in this plugin.
- RabbitMQ queue type is immutable, so an existing classic queue cannot be converted in place.

## Complete Production Example

A full production configuration combining SSL, event filtering, custom formatter, and heartbeat tuning:
Expand Down
22 changes: 21 additions & 1 deletion docs/intro.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,8 +5,9 @@ Stream Snakemake workflow events to STOMP message brokers (ActiveMQ, RabbitMQ, A
## Key Features

- **Real-time event streaming** to STOMP-compatible message brokers
- **RabbitMQ stream queue support** for append-only log style destinations
- **SSL/TLS encryption** for secure production deployments
- **Flexible formatters** - default flat JSON or JLab Scientific Workflow schema
- **Flexible formatters** - default flat JSON, JLab Scientific Workflow schema, or comprehensive Snakemake logs
- **Event filtering** - include/exclude specific event types to reduce noise
- **Custom formatters** - plugin your own formatter classes
- **Environment variable support** for credentials and sensitive configuration
Expand Down Expand Up @@ -48,6 +49,25 @@ Then run:
snakemake --profile profiles/stomp
```

### RabbitMQ Stream Destinations

When using RabbitMQ's STOMP adapter, the plugin can publish to stream queues.

```yaml
logger:
stomp:
host: "rabbitmq.example.com"
port: 61613
user: "${STOMP_USER}"
password: "${STOMP_PASSWORD}"
queue: "/queue/snakemake.stream"
use_stream: true
stream_filter_by_workflow: true
```

Use `/queue/<name>` to create a new stream on first publish. Use `/amq/queue/<name>`
to publish to an existing stream without redeclaring it.

### Environment Variables

Secure credentials using environment variables:
Expand Down
1 change: 1 addition & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ dependencies = [
"snakemake-interface-common>=1.22.0,<2",
"snakemake-interface-logger-plugins>=2.0.0,<3",
"stomp-py>=8.2.0",
"paramiko>=3.4.0",
]

[[project.authors]]
Expand Down
Loading
Loading