| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103 |
- package main
- import (
- "encoding/json"
- "flag"
- "log"
- "os"
- "time"
- amqp "github.com/rabbitmq/amqp091-go"
- )
- type Transaction map[string]interface{}
- func main() {
- var (
- amqpURL = flag.String("amqp", "amqp://admin:ward-tootled-wick-intrados@localhost:5672/", "AMQP URL")
- queueName = flag.String("queue", "ingestion", "RabbitMQ queue to consume")
- persistF = flag.String("persist-file", os.Getenv("PERSIST_FILE"), "file to append persisted transactions")
- )
- flag.Parse()
- // set PERSIST_FILE from flag if provided
- if *persistF != "" {
- os.Setenv("PERSIST_FILE", *persistF)
- }
- // start AMQP consumer (blocking)
- for {
- err := consumeAMQP(*amqpURL, *queueName)
- if err != nil {
- log.Printf("AMQP consumer error: %v; retrying in 5s", err)
- time.Sleep(5 * time.Second)
- continue
- }
- break
- }
- }
- func consumeAMQP(amqpURL, queue string) error {
- conn, err := amqp.Dial(amqpURL)
- if err != nil {
- return err
- }
- defer conn.Close()
- ch, err := conn.Channel()
- if err != nil {
- return err
- }
- defer ch.Close()
- _, err = ch.QueueDeclare(
- queue,
- true,
- false,
- false,
- false,
- nil,
- )
- if err != nil {
- return err
- }
- msgs, err := ch.Consume(
- queue,
- "",
- false, // autoAck false
- false,
- false,
- false,
- nil,
- )
- if err != nil {
- return err
- }
- log.Printf("AMQP consumer started on queue %s", queue)
- for d := range msgs {
- var t Transaction
- if err := json.Unmarshal(d.Body, &t); err != nil {
- log.Printf("failed to unmarshal message: %v", err)
- d.Nack(false, false)
- continue
- }
- // persist same as HTTP handler: append to file if set
- if path := os.Getenv("PERSIST_FILE"); path != "" {
- f, err := os.OpenFile(path, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0o644)
- if err != nil {
- log.Printf("failed to open persist file: %v", err)
- } else {
- f.Write(d.Body)
- f.Write([]byte("\n"))
- f.Close()
- }
- }
- d.Ack(false)
- }
- return nil
- }
|