Files
Bernt (LandveX AI) 67a69ab073 feat(boc): v1.0 - Complete Business Operations Center
- Go backend API with full CRUD for all modules
- Rust analytics service with parallel processing
- C runtime with POSIX shared memory IPC
- PostgreSQL schema with 30+ tables
- Redis cache, Kafka event streaming
- WebSocket hub, automation engine
- PDF generation, Resend email integration
- JWT auth, multi-tenant
- Docker Compose deployment
- Nginx reverse proxy

Refs: BOC-001
2026-07-12 13:21:10 +00:00

176 lines
4.5 KiB
Go

package cache
import (
"context"
"encoding/json"
"fmt"
"time"
"github.com/redis/go-redis/v9"
)
// RedisClient wraps go-redis with BOC-specific operations
type RedisClient struct {
client *redis.Client
ctx context.Context
}
// NewRedisClient creates a new Redis client
func NewRedisClient(addr string) (*RedisClient, error) {
client := redis.NewClient(&redis.Options{
Addr: addr,
Password: "", // no password
DB: 0, // default DB
PoolSize: 10,
})
ctx := context.Background()
if err := client.Ping(ctx).Err(); err != nil {
return nil, fmt.Errorf("redis ping failed: %w", err)
}
return &RedisClient{
client: client,
ctx: ctx,
}, nil
}
// Close closes the Redis connection
func (r *RedisClient) Close() error {
return r.client.Close()
}
// Get retrieves a value from cache
func (r *RedisClient) Get(key string, dest interface{}) error {
data, err := r.client.Get(r.ctx, key).Bytes()
if err == redis.Nil {
return fmt.Errorf("cache miss")
}
if err != nil {
return err
}
return json.Unmarshal(data, dest)
}
// Set stores a value in cache with TTL
func (r *RedisClient) Set(key string, value interface{}, ttl time.Duration) error {
data, err := json.Marshal(value)
if err != nil {
return err
}
return r.client.Set(r.ctx, key, data, ttl).Err()
}
// Delete removes a key from cache
func (r *RedisClient) Delete(key string) error {
return r.client.Del(r.ctx, key).Err()
}
// DeletePattern removes keys matching a pattern
func (r *RedisClient) DeletePattern(pattern string) error {
keys, err := r.client.Keys(r.ctx, pattern).Result()
if err != nil {
return err
}
if len(keys) > 0 {
return r.client.Del(r.ctx, keys...).Err()
}
return nil
}
// Exists checks if a key exists
func (r *RedisClient) Exists(key string) bool {
n, err := r.client.Exists(r.ctx, key).Result()
return err == nil && n > 0
}
// Increment atomically increments a counter
func (r *RedisClient) Increment(key string) (int64, error) {
return r.client.Incr(r.ctx, key).Result()
}
// Expire sets a TTL on a key
func (r *RedisClient) Expire(key string, ttl time.Duration) error {
return r.client.Expire(r.ctx, key, ttl).Err()
}
// Cache analytics result
func (r *RedisClient) CacheAnalytics(tenantID, metric, period string, data interface{}) error {
key := fmt.Sprintf("analytics:%s:%s:%s", tenantID, metric, period)
return r.Set(key, data, 5*time.Minute)
}
// GetCachedAnalytics retrieves cached analytics
func (r *RedisClient) GetCachedAnalytics(tenantID, metric, period string, dest interface{}) error {
key := fmt.Sprintf("analytics:%s:%s:%s", tenantID, metric, period)
return r.Get(key, dest)
}
// Cache dashboard data
func (r *RedisClient) CacheDashboard(tenantID string, data interface{}) error {
key := fmt.Sprintf("dashboard:%s", tenantID)
return r.Set(key, data, 1*time.Minute)
}
// GetCachedDashboard retrieves cached dashboard
func (r *RedisClient) GetCachedDashboard(tenantID string, dest interface{}) error {
key := fmt.Sprintf("dashboard:%s", tenantID)
return r.Get(key, dest)
}
// Rate limiting
func (r *RedisClient) RateLimit(key string, maxRequests int, window time.Duration) (bool, error) {
pipe := r.client.Pipeline()
now := time.Now().Unix()
windowStart := now - int64(window.Seconds())
// Remove old entries
pipe.ZRemRangeByScore(r.ctx, key, "0", fmt.Sprintf("%d", windowStart))
// Count current entries
pipe.ZCard(r.ctx, key)
// Add current request
pipe.ZAdd(r.ctx, key, redis.Z{Score: float64(now), Member: now})
// Set expiry on the key
pipe.Expire(r.ctx, key, window)
cmders, err := pipe.Exec(r.ctx)
if err != nil {
return false, err
}
// cmders[1] is ZCard result
count := cmders[1].(*redis.IntCmd).Val()
return count <= int64(maxRequests), nil
}
// Session management
func (r *RedisClient) SetSession(sessionID string, data map[string]interface{}, ttl time.Duration) error {
key := fmt.Sprintf("session:%s", sessionID)
return r.Set(key, data, ttl)
}
func (r *RedisClient) GetSession(sessionID string) (map[string]interface{}, error) {
key := fmt.Sprintf("session:%s", sessionID)
var data map[string]interface{}
err := r.Get(key, &data)
return data, err
}
func (r *RedisClient) DeleteSession(sessionID string) error {
key := fmt.Sprintf("session:%s", sessionID)
return r.Delete(key)
}
// Pub/Sub for real-time events
func (r *RedisClient) Publish(channel string, message interface{}) error {
data, err := json.Marshal(message)
if err != nil {
return err
}
return r.client.Publish(r.ctx, channel, data).Err()
}
func (r *RedisClient) Subscribe(channel string) *redis.PubSub {
return r.client.Subscribe(r.ctx, channel)
}