package events import ( "context" "encoding/json" "fmt" "time" "github.com/segmentio/kafka-go" ) // KafkaClient wraps kafka-go for BOC event streaming type KafkaClient struct { writer *kafka.Writer reader *kafka.Reader brokers []string } // Event represents a domain event type Event struct { ID string `json:"id"` Type string `json:"type"` TenantID string `json:"tenant_id"` EntityID string `json:"entity_id"` EntityType string `json:"entity_type"` Action string `json:"action"` Data map[string]interface{} `json:"data"` Timestamp time.Time `json:"timestamp"` UserID string `json:"user_id,omitempty"` } // Event types const ( EventCustomerCreated = "customer.created" EventCustomerUpdated = "customer.updated" EventCustomerDeleted = "customer.deleted" EventDealCreated = "deal.created" EventDealUpdated = "deal.updated" EventDealClosed = "deal.closed" EventInvoiceCreated = "invoice.created" EventInvoicePaid = "invoice.paid" EventTicketCreated = "ticket.created" EventTicketResolved = "ticket.resolved" EventEmployeeCreated = "employee.created" EventContractRenewal = "contract.renewal_due" EventWorkflowTriggered = "workflow.triggered" EventReportGenerated = "report.generated" ) // Topic names const ( TopicBOCEvents = "boc.events" TopicAuditLog = "boc.audit" TopicAnalytics = "boc.analytics" TopicNotifications = "boc.notifications" ) // NewKafkaClient creates a new Kafka client func NewKafkaClient(brokers []string) (*KafkaClient, error) { writer := &kafka.Writer{ Addr: kafka.TCP(brokers...), Balancer: &kafka.LeastBytes{}, RequiredAcks: kafka.RequireAll, } // Test connection conn, err := kafka.Dial("tcp", brokers[0]) if err != nil { return nil, fmt.Errorf("kafka connection failed: %w", err) } conn.Close() return &KafkaClient{ writer: writer, brokers: brokers, }, nil } // Close closes the Kafka client func (k *KafkaClient) Close() error { return k.writer.Close() } // Publish sends an event to Kafka func (k *KafkaClient) Publish(ctx context.Context, topic string, event Event) error { data, err := json.Marshal(event) if err != nil { return fmt.Errorf("marshal event: %w", err) } return k.writer.WriteMessages(ctx, kafka.Message{ Topic: topic, Key: []byte(event.EntityID), Value: data, Time: event.Timestamp, }) } // PublishAsync sends an event asynchronously func (k *KafkaClient) PublishAsync(topic string, event Event) { go func() { ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() if err := k.Publish(ctx, topic, event); err != nil { fmt.Printf("Failed to publish event: %v\n", err) } }() } // CreateReader creates a new Kafka reader for a topic func (k *KafkaClient) CreateReader(topic, groupID string) *kafka.Reader { return kafka.NewReader(kafka.ReaderConfig{ Brokers: k.brokers, Topic: topic, GroupID: groupID, MinBytes: 10e3, // 10KB MaxBytes: 10e6, // 10MB }) } // CreateEvent creates a new event with defaults func CreateEvent(eventType, tenantID, entityType, entityID, action string, data map[string]interface{}) Event { return Event{ ID: fmt.Sprintf("%d-%s", time.Now().UnixNano(), entityID), Type: eventType, TenantID: tenantID, EntityID: entityID, EntityType: entityType, Action: action, Data: data, Timestamp: time.Now().UTC(), } } // EnsureTopics creates topics if they don't exist func (k *KafkaClient) EnsureTopics() error { conn, err := kafka.Dial("tcp", k.brokers[0]) if err != nil { return err } defer conn.Close() topics := []string{TopicBOCEvents, TopicAuditLog, TopicAnalytics, TopicNotifications} for _, topic := range topics { topicConfigs := []kafka.TopicConfig{ { Topic: topic, NumPartitions: 3, ReplicationFactor: 1, }, } if err := conn.CreateTopics(topicConfigs...); err != nil { // Topic might already exist, continue continue } } return nil } // EventConsumer handles consuming events from Kafka type EventConsumer struct { reader *kafka.Reader handlers map[string]func(Event) error } // NewEventConsumer creates a new event consumer func NewEventConsumer(brokers []string, topic, groupID string) *EventConsumer { reader := kafka.NewReader(kafka.ReaderConfig{ Brokers: brokers, Topic: topic, GroupID: groupID, MinBytes: 10e3, MaxBytes: 10e6, }) return &EventConsumer{ reader: reader, handlers: make(map[string]func(Event) error), } } // RegisterHandler registers a handler for an event type func (c *EventConsumer) RegisterHandler(eventType string, handler func(Event) error) { c.handlers[eventType] = handler } // Start begins consuming events func (c *EventConsumer) Start(ctx context.Context) { go func() { for { msg, err := c.reader.ReadMessage(ctx) if err != nil { if ctx.Err() != nil { return // Context cancelled } fmt.Printf("Error reading message: %v\n", err) continue } var event Event if err := json.Unmarshal(msg.Value, &event); err != nil { fmt.Printf("Error unmarshaling event: %v\n", err) continue } if handler, ok := c.handlers[event.Type]; ok { if err := handler(event); err != nil { fmt.Printf("Error handling event %s: %v\n", event.Type, err) } } } }() } // Close closes the consumer func (c *EventConsumer) Close() error { return c.reader.Close() }