Tailing Live Logs with Subprocess
The subprocess processor allows you to execute external commands and pipe data through them. One powerful use case is tailing live log files, enabling you to stream log data directly into your Expanso Edge pipeline for real-time processing, filtering, and routing.
Why Use Subprocess for Tailing?
- Process log files in real-time without deploying dedicated log collectors
- Combine with Expanso processors to filter, enrich, parse, and route logs before they leave the edge
- Reduce bandwidth by processing and filtering at the source
- No external agents needed — just the
tailcommand built into Linux/Unix systems
Quick Start
Stream a live log file through a subprocess processor:
input:
generate:
count: 1 # Fire once to start the pipeline
mapping: 'root = ""'
pipeline:
processors:
- subprocess:
name: tail
args:
- "-F"
- "${DATA_DIR:/app/data/sensor.log}"
- mapping: 'root = this.parse_json().catch(deleted())'
output:
resource: my_output
What's happening:
generatewithcount: 1creates a single message to start the pipeline (an interval or a repeatingcount: 0would instead spawn a newtail -Fprocess on every tick)subprocesswithtail -Fruns the tail command, which outputs new log lines as they're written- Each line from tail becomes a new message in the pipeline
- Subsequent processors (like the
mappingprocessor here) parse and transform each line as it arrives - The output resource receives the processed logs
Configuration Details
Basic Tail (No Logs on Startup)
Tail a file starting from the end (skip existing logs):
subprocess:
name: tail
args:
- "-f"
- "/var/log/app.log"
Follow File Descriptor (Handles Log Rotation)
Use -F to follow the file descriptor, automatically picking up rotated logs:
subprocess:
name: tail
args:
- "-F"
- "/var/log/app.log"
This is essential for production log files that get rotated. The -f flag follows filenames (which breaks on rotation), while -F follows file descriptors (continues after rotation).
Start From Beginning
Tail all existing lines plus new ones:
subprocess:
name: tail
args:
- "-F"
- "+1" # Start from line 1
- "/var/log/app.log"
Start From Specific Line Count
Tail the last 100 existing lines plus new ones:
subprocess:
name: tail
args:
- "-F"
- "-n"
- "100"
- "/var/log/app.log"
Dynamic Path with Environment Variables
Use ${ENV_VAR:default} syntax to reference environment variables:
subprocess:
name: tail
args:
- "-F"
- "${LOG_PATH:/var/log/app.log}"
Then set the environment variable when running Expanso Edge:
export LOG_PATH=/custom/path/logs/app.log
expanso-edge run
Or in your configuration:
environment:
LOG_PATH: /custom/path/logs/app.log
Complete Examples
Example 1: Parse and Route JSON Logs
Stream JSON logs, parse them, and route by severity:
input:
generate:
count: 1 # Fire once to start the pipeline
mapping: 'root = ""'
pipeline:
processors:
- subprocess:
name: tail
args:
- "-F"
- "/var/log/app.json"
- mapping: 'root = this.parse_json().catch(deleted())'
- switch:
- check: this.severity == "ERROR"
processors:
- mapping: root.routed_to = "error_channel"
- check: this.severity == "WARN"
processors:
- mapping: root.routed_to = "warning_channel"
- processors:
- mapping: root.routed_to = "info_channel"
output:
switch:
cases:
# errors
- check: this.routed_to == "error_channel"
output:
aws_s3:
bucket: my-bucket
path: 'errors/${! timestamp_unix() }.json'
# warnings
- check: this.routed_to == "warning_channel"
output:
aws_s3:
bucket: my-bucket
path: 'warnings/${! timestamp_unix() }.json'
# info (default)
- output:
aws_s3:
bucket: my-bucket
path: 'info/${! timestamp_unix() }.json'
Example 2: Extract Fields and Batch Upload
Tail logs, parse structured data, enrich, and batch for efficient upload:
input:
generate:
count: 1 # Fire once to start the pipeline
mapping: 'root = ""'
pipeline:
threads: 2
processors:
- subprocess:
name: tail
args:
- "-F"
- "/opt/data/sensor.log"
- mapping: 'root = this.parse_json().catch(deleted())'
- mapping: |
root.timestamp = now().ts_format("2006-01-02T15:04:05Z07:00")
root.node_id = env("NODE_ID").or("unknown")
root.region = env("REGION").or("us-east-1")
output:
aws_s3:
bucket: my-bucket
path: 'logs/year=${! now().ts_format("2006") }/month=${! now().ts_format("01") }/day=${! now().ts_format("02") }/${! timestamp_unix() }.jsonl'
batching:
count: 100
Example 3: Real-time Alert on Error Patterns
Tail logs and trigger alerts on specific error patterns:
input:
generate:
count: 1 # Fire once to start the pipeline
mapping: 'root = ""'
pipeline:
processors:
- subprocess:
name: tail
args:
- "-F"
- "/var/log/application.log"
- mapping: |
root = if content().string().re_match("ERROR.*OutOfMemory|ERROR.*StackOverflow") { this } else { deleted() }
- mapping: |
root.alert = true
root.severity = "critical"
root.timestamp = now()
root.host = hostname()
output:
switch:
cases:
# Alert and back up when alert == true
- check: this.alert == true
output:
broker:
pattern: fan_out
outputs:
- http_client:
url: "https://alerts.example.com/webhook"
verb: POST
- file:
path: /var/log/expanso-critical-errors.jsonl
Common Tail Flags Reference
| Flag | Description | Example |
|---|---|---|
-f | Follow file by name (breaks on rotation) | tail -f app.log |
-F | Follow file descriptor (survives rotation) | tail -F app.log |
-n N | Show last N lines | tail -n 50 app.log |
-n +N | Show from line N onward | tail -n +1 app.log |
--pid=PID | Terminate when PID dies | tail -F --pid=$$ app.log |
-q | Suppress headers | tail -q -F app.log |
--retry | Retry if file is inaccessible | tail -F --retry app.log |
Performance Considerations
1. Output Rate vs Processing Speed
If logs arrive faster than your pipeline can process them, consider:
pipeline:
threads: 4 # Use multiple threads to increase throughput
processors:
- subprocess:
name: tail
args: ["-F", "/var/log/app.log"]
2. Buffer Size
For very large log lines, increase the subprocess buffer:
subprocess:
name: tail
args:
- "-F"
- "/var/log/app.log"
max_buffer: 262144 # Increase from default 65536 (64KB) to 256KB
3. Batching for Efficiency
Batch logs before writing to external systems:
output:
http_client:
url: "https://ingest.example.com/logs"
batching:
period: 5s
Production Readiness
Subprocess Restart Behavior
The subprocess processor is designed to keep long-running processes alive:
- If tail exits early, Expanso Edge automatically restarts it
- Restart is transparent — the pipeline continues processing new logs
- No data loss — thanks to
-Fflag that persists tail state across restarts - Exponential backoff — prevents rapid restart loops if there's a persistent issue
Important: Ensure tail never crashes by using appropriate flags:
subprocess:
name: tail
args:
- "-F" # Follow file descriptor (survives rotation)
- "--retry" # Retry if file temporarily unavailable
- "/var/log/app.log"
Error Handling & Recovery Patterns
Pattern 1: Detect and Mark Tail Failures
If tail fails (file deleted, permissions lost), tag the message so a downstream output can alert or route on it:
pipeline:
processors:
- try:
- subprocess:
name: tail
args:
- "-F"
- "${LOG_FILE:/var/log/app.log}"
- mapping: root.source = "tail"
- catch:
- mapping: |
root.source = "fallback"
root.recovered = true
root.error = error()
This ensures that if tail fails to start or crashes, the affected messages are tagged with the error instead of the pipeline going silent, so a downstream output can alert or route them.
Pattern 2: Health Check with Heartbeat
Add a heartbeat log to monitor tail health:
input:
generate:
count: 1 # Fire once to start the pipeline
mapping: 'root = ""'
pipeline:
processors:
- subprocess:
name: bash
args:
- "-c"
- "tail -F /var/log/app.log & echo '{\"heartbeat\":true,\"timestamp\":\"'$(date -u +%Y-%m-%dT%H:%M:%SZ)'\"}' >> /var/log/expanso-heartbeat.log; wait"
- mapping: |
root.collected_at = now()
output:
switch:
cases:
- check: this.heartbeat == true
output:
http_client:
url: "https://monitoring.example.com/health"
verb: POST
- output:
file:
path: /var/log/expanso-live.jsonl
Pattern 3: Wrap Tail with Health Wrapper Script
For maximum reliability, wrap tail in a script that handles cleanup:
/usr/local/bin/tail-wrapper.sh:
#!/bin/bash
set -e
LOG_FILE="${1:?LOG_FILE required}"
HEALTH_FILE="${2:-/tmp/tail-health.txt}"
# Signal handler for clean shutdown
cleanup() {
rm -f "$HEALTH_FILE"
exit 0
}
trap cleanup EXIT INT TERM
# Monitor the file and track health
tail -F "$LOG_FILE" &
TAIL_PID=$!
# Write health marker every 10 seconds
while kill -0 $TAIL_PID 2>/dev/null; do
echo "$(date -u +%s)" > "$HEALTH_FILE"
sleep 10
done
wait $TAIL_PID
Then use it in your pipeline:
subprocess:
name: /usr/local/bin/tail-wrapper.sh
args:
- "/var/log/app.log"
- "/tmp/tail-${NODE_ID}.health"
Preventing Data Loss
1. Atomic Output Writes
Always use transactional or batch outputs:
output:
aws_s3:
bucket: logs
path: 'batch-${! timestamp_unix() }.jsonl'
batching:
count: 100 # Don't send partial batches
2. Retry on Failure
Add retry logic to your outputs:
output:
retry:
max_retries: 3
backoff:
initial_interval: 1s
max_interval: 30s
output:
http_client:
url: "https://ingest.example.com"
timeout: 10s
3. Local Backup
Always maintain a local backup alongside cloud/remote storage:
output:
broker:
pattern: fan_out
outputs:
- retry:
max_retries: 3
output:
http_client:
url: "https://ingest.example.com"
- file:
path: '/var/log/expanso-backup-${! timestamp_unix() }.jsonl'
Resource Limits
Prevent tail from consuming excessive memory or CPU:
pipeline:
threads: 1 # Limit to 1 thread for subprocess (tail is single-threaded anyway)
processors:
- subprocess:
name: tail
args:
- "-F"
- "/var/log/app.log"
max_buffer: 65536 # Default is good; increase only if needed
Then enforce OS-level limits:
# Limit Expanso Edge process to 512MB RAM, 50% CPU
systemctl set-property expanso.service MemoryLimit=512M
systemctl set-property expanso.service CPUQuota=50%
Monitoring Tail Health
Check Logs for Restarts
Watch Expanso logs for subprocess restart messages:
journalctl -u expanso.service -f | grep -i "subprocess\|restart\|error"
Verify File Permissions
Regular checks prevent permission-related failures:
# Add to crontab
0 * * * * /usr/local/bin/check-log-perms.sh
/usr/local/bin/check-log-perms.sh:
#!/bin/bash
LOG_FILE="/var/log/app.log"
EXPANSO_USER="expanso"
if [ ! -r "$LOG_FILE" ]; then
echo "ERROR: $EXPANSO_USER cannot read $LOG_FILE" | logger
exit 1
fi
Monitor File Handle Limits
If tail can't open more files:
# Check current limits
ulimit -n
# Set higher limit (in systemd service file)
[Service]
LimitNOFILE=65536
Alert on Tail Failure
Create a secondary pipeline to detect when tail stops working:
input:
file:
paths:
- /tmp/tail-health.txt
pipeline:
processors:
- mapping: |
let age_seconds = now().ts_unix() - content().string().number()
root = if $age_seconds > 30 {
"Tail process appears stuck or crashed"
} else {
deleted()
}
output:
http_client:
url: "https://alerts.example.com/critical"
verb: POST
Troubleshooting
"Permission denied" errors
Ensure the Expanso Edge process has read permissions on the log file:
# Check permissions
ls -la /var/log/app.log
# Run Expanso with appropriate user/group
sudo -u loguser expanso-edge run
Subprocess hangs or no output
Check that:
- The log file path is correct
- The file actually exists and is being written to
- Use
-Finstead of-ffor production (handles rotation) - Consider adding
--retryfor reliability:
subprocess:
name: tail
args:
- "-F"
- "--retry"
- "/var/log/app.log"
High CPU usage
If subprocess is consuming too much CPU:
- Reduce pipeline threads
- Check if the log file is being hammered with writes
- Consider filtering earlier in the pipeline to reduce downstream processing
Tail with Other Commands
The subprocess processor works with any command, not just tail. You can combine commands:
# Use tail + grep to filter in-process
subprocess:
name: bash
args:
- "-c"
- "tail -F /var/log/app.log | grep ERROR"
Next Steps
- Subprocess Processor — Full configuration reference
- Log Processing Examples — Production pipeline patterns
- Parse Log Processor — Parse common log formats
- Mapping Processor — Transform parsed data
- Error Handling — Handle failures gracefully