v0.10.0 flow and plugin structure, networking and commands removed, dirty

This commit is contained in:
2025-11-11 16:42:09 -05:00
parent d38908e0f1
commit 46a436baa0
57 changed files with 2637 additions and 7301 deletions
+159
View File
@@ -0,0 +1,159 @@
// FILE: internal/flow/flow.go
package flow
import (
"context"
"fmt"
"sync/atomic"
"logwisp/src/internal/config"
"logwisp/src/internal/core"
"logwisp/src/internal/filter"
"logwisp/src/internal/format"
"github.com/lixenwraith/log"
)
// Flow manages the complete processing pipeline for log entries:
// LogEntry -> Rate Limiter -> Filters -> Formatter -> TransportEvent
type Flow struct {
rateLimiter *RateLimiter
filterChain *filter.Chain
formatter format.Formatter
heartbeat *HeartbeatGenerator
logger *log.Logger
// Statistics
totalProcessed atomic.Uint64
totalDropped atomic.Uint64
totalFormatted atomic.Uint64
}
// NewFlow creates a flow processor from configuration
func NewFlow(cfg *config.FlowConfig, logger *log.Logger) (*Flow, error) {
if cfg == nil {
cfg = &config.FlowConfig{}
}
f := &Flow{
logger: logger,
}
// Create rate limiter if configured
if cfg.RateLimit != nil {
limiter, err := NewRateLimiter(*cfg.RateLimit, logger)
if err != nil {
return nil, fmt.Errorf("failed to create rate limiter: %w", err)
}
f.rateLimiter = limiter
}
// Create filter chain if configured
if len(cfg.Filters) > 0 {
chain, err := filter.NewChain(cfg.Filters, logger)
if err != nil {
return nil, fmt.Errorf("failed to create filter chain: %w", err)
}
f.filterChain = chain
}
// Create formatter, if not configured falls back to raw '\n' delimited
formatter, err := format.NewFormatter(cfg.Format, logger)
if err != nil {
return nil, fmt.Errorf("failed to create formatter: %w", err)
}
f.formatter = formatter
// Create heartbeat generator if configured
if cfg.Heartbeat != nil && cfg.Heartbeat.Enabled {
f.heartbeat = NewHeartbeatGenerator(cfg.Heartbeat, logger)
}
logger.Info("msg", "Flow processor created",
"component", "flow",
"rate_limiter", f.rateLimiter != nil,
"filter_chain", f.filterChain != nil,
"formatter", formatter.Name(),
"heartbeat", f.heartbeat != nil)
return f, nil
}
// Process applies all flow stages to a log entry
// Returns TransportEvent and whether entry passed all stages
func (f *Flow) Process(entry core.LogEntry) (core.TransportEvent, bool) {
f.totalProcessed.Add(1)
// Stage 1: Rate limiting
if f.rateLimiter != nil {
if !f.rateLimiter.Allow(entry) {
f.totalDropped.Add(1)
return core.TransportEvent{}, false
}
}
// Stage 2: Filtering
if f.filterChain != nil {
if !f.filterChain.Apply(entry) {
f.totalDropped.Add(1)
return core.TransportEvent{}, false
}
}
// Stage 3: Formatting
formatted, err := f.formatter.Format(entry)
if err != nil {
f.logger.Error("msg", "Failed to format log entry",
"component", "flow",
"error", err)
f.totalDropped.Add(1)
return core.TransportEvent{}, false
}
f.totalFormatted.Add(1)
// Create transport event
event := core.TransportEvent{
Time: entry.Time,
Payload: formatted,
}
return event, true
}
// StartHeartbeat starts the heartbeat generator if configured
// Returns channel that emits heartbeat events
func (f *Flow) StartHeartbeat(ctx context.Context) <-chan core.TransportEvent {
if f.heartbeat == nil {
return nil
}
return f.heartbeat.Start(ctx)
}
// GetStats returns flow statistics
func (f *Flow) GetStats() map[string]any {
stats := map[string]any{
"total_processed": f.totalProcessed.Load(),
"total_dropped": f.totalDropped.Load(),
"total_formatted": f.totalFormatted.Load(),
}
if f.rateLimiter != nil {
stats["rate_limiter"] = f.rateLimiter.GetStats()
}
if f.filterChain != nil {
stats["filters"] = f.filterChain.GetStats()
}
if f.formatter != nil {
stats["formatter"] = f.formatter.Name()
}
if f.heartbeat != nil {
stats["heartbeat_enabled"] = true
stats["heartbeat_interval_ms"] = f.heartbeat.IntervalMS()
}
return stats
}
+110
View File
@@ -0,0 +1,110 @@
// FILE: src/internal/flow/heartbeat.go
package flow
import (
"context"
"encoding/json"
"fmt"
"sync/atomic"
"time"
"logwisp/src/internal/config"
"logwisp/src/internal/core"
"github.com/lixenwraith/log"
)
// HeartbeatGenerator produces periodic heartbeat events
type HeartbeatGenerator struct {
config *config.HeartbeatConfig
logger *log.Logger
beatCount atomic.Uint64
lastBeat atomic.Value // time.Time
}
// NewHeartbeatGenerator creates a new heartbeat generator
func NewHeartbeatGenerator(cfg *config.HeartbeatConfig, logger *log.Logger) *HeartbeatGenerator {
hg := &HeartbeatGenerator{
config: cfg,
logger: logger,
}
hg.lastBeat.Store(time.Time{})
return hg
}
// Start begins generating heartbeat events
func (hg *HeartbeatGenerator) Start(ctx context.Context) <-chan core.TransportEvent {
ch := make(chan core.TransportEvent)
go func() {
defer close(ch)
ticker := time.NewTicker(time.Duration(hg.config.IntervalMS) * time.Millisecond)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case t := <-ticker.C:
event := hg.generateHeartbeat(t)
select {
case ch <- event:
hg.beatCount.Add(1)
hg.lastBeat.Store(t)
case <-ctx.Done():
return
}
}
}
}()
return ch
}
// generateHeartbeat creates a heartbeat transport event
func (hg *HeartbeatGenerator) generateHeartbeat(t time.Time) core.TransportEvent {
var payload []byte
switch hg.config.Format {
case "json":
data := map[string]any{
"type": "heartbeat",
"timestamp": t.Format(time.RFC3339Nano),
}
if hg.config.IncludeStats {
data["beat_count"] = hg.beatCount.Load()
if last, ok := hg.lastBeat.Load().(time.Time); ok && !last.IsZero() {
data["interval_ms"] = t.Sub(last).Milliseconds()
}
}
payload, _ = json.Marshal(data)
payload = append(payload, '\n')
case "comment":
// SSE-style comment for web streaming
msg := fmt.Sprintf(": heartbeat %s", t.Format(time.RFC3339))
if hg.config.IncludeStats {
msg = fmt.Sprintf("%s [#%d]", msg, hg.beatCount.Load())
}
payload = []byte(msg + "\n")
default:
// Plain text
msg := fmt.Sprintf("heartbeat: %s", t.Format(time.RFC3339))
if hg.config.IncludeStats {
msg = fmt.Sprintf("%s (#%d)", msg, hg.beatCount.Load())
}
payload = []byte(msg + "\n")
}
return core.TransportEvent{
Time: t,
Payload: payload,
}
}
// IntervalMS returns the heartbeat interval in milliseconds
func (hg *HeartbeatGenerator) IntervalMS() int64 {
return hg.config.IntervalMS
}
+5 -5
View File
@@ -12,7 +12,7 @@ import (
"github.com/lixenwraith/log"
)
// RateLimiter enforces rate limits on log entries flowing through a pipeline.
// RateLimiter enforces rate limits on log entries flowing through a pipeline
type RateLimiter struct {
bucket *tokenbucket.TokenBucket
policy config.RateLimitPolicy
@@ -24,7 +24,7 @@ type RateLimiter struct {
droppedCount atomic.Uint64
}
// NewRateLimiter creates a new pipeline-level rate limiter from configuration.
// NewRateLimiter creates a new pipeline-level rate limiter from configuration
func NewRateLimiter(cfg config.RateLimitConfig, logger *log.Logger) (*RateLimiter, error) {
if cfg.Rate <= 0 {
return nil, nil // No rate limit
@@ -53,7 +53,7 @@ func NewRateLimiter(cfg config.RateLimitConfig, logger *log.Logger) (*RateLimite
return l, nil
}
// Allow checks if a log entry is permitted to pass based on the rate limit.
// Allow checks if a log entry is permitted to pass based on the rate limit
func (l *RateLimiter) Allow(entry core.LogEntry) bool {
if l == nil || l.policy == config.PolicyPass {
return true
@@ -79,7 +79,7 @@ func (l *RateLimiter) Allow(entry core.LogEntry) bool {
return true
}
// GetStats returns statistics for the rate limiter.
// GetStats returns statistics for the rate limiter
func (l *RateLimiter) GetStats() map[string]any {
if l == nil {
return map[string]any{
@@ -102,7 +102,7 @@ func (l *RateLimiter) GetStats() map[string]any {
return stats
}
// policyString returns the string representation of a rate limit policy.
// policyString returns the string representation of a rate limit policy
func policyString(p config.RateLimitPolicy) string {
switch p {
case config.PolicyDrop: