package main

import (
	"context"
	"encoding/json"
	"flag"
	"fmt"
	"log"
	"os"
	"time"

	"github.com/ably/ably-go/ably"
	"github.com/joho/godotenv"
	"github.com/rabbitmq/amqp091-go"
)

type AblyMessage struct {
	ID     string                 `json:"id"`
	Token  string                 `json:"token"`
	Params map[string]interface{} `json:"params"`
}

func PublishMessage(client *ably.Realtime, channelName, messageName, messageData string) string {

	// Get the channel
	channel := client.Channels.Get(channelName)

	var ablyMsg AblyMessage
	if err := json.Unmarshal([]byte(messageData), &ablyMsg); err != nil {
		log.Fatalf("Failed to unmarshal message: %v", err)
	}

	// Publish a message to the channel
	err := channel.Publish(context.Background(), messageName, messageData)
	if err != nil {
		log.Fatalf("Failed to publish message: %v", err)
	}

	log.Printf("Message published to channel %s: %s", channelName, messageData)

	return ablyMsg.ID
}

func PublishMessageToQueue(queueName, messageName, messageData string) string {
	// Connect to the Ably queue via AMQP
	apiKey := os.Getenv("ABLY_API_KEY")
	amqpURL := "amqps://" + apiKey + "@us-east-1-a-queue.ably.io/shared"
	conn, err := amqp091.Dial(amqpURL)
	if err != nil {
		log.Fatalf("Failed to connect to Ably queue: %v", err)
	}
	defer conn.Close()

	// Create a channel
	ch, err := conn.Channel()
	if err != nil {
		log.Fatalf("Failed to open a channel: %v", err)
	}
	defer ch.Close()

	// Prepare the message
	var ablyMsg AblyMessage
	if err := json.Unmarshal([]byte(messageData), &ablyMsg); err != nil {
		log.Fatalf("Failed to unmarshal message: %v", err)
	}

	// Marshal the message back to JSON for publishing
	messageBytes, err := json.Marshal(ablyMsg)
	if err != nil {
		log.Fatalf("Failed to marshal message: %v", err)
	}

	fmt.Println(os.Getenv("ABLY_API_KEY"))
	fmt.Println(string(messageBytes))

	// Publish the message to the queue
	err = ch.Publish(
		"",        // exchange (empty for default exchange)
		queueName, // routing key (queue name)
		false,     // mandatory
		false,     // immediate
		amqp091.Publishing{
			ContentType: "application/json",
			Body:        messageBytes,
		})
	if err != nil {
		log.Fatalf("Failed to publish message: %v", err)
	}

	log.Printf("Message published to queue %s: %s", queueName, messageData)

	return ablyMsg.ID
}

func SubscribeToChannel(client *ably.Realtime, channelName string, messageID string) func() {
	channel := client.Channels.Get("reply-channel-" + messageID)
	unsubscribe, err := channel.SubscribeAll(context.Background(), func(msg *ably.Message) {
		log.Printf("Received message from channel %s: %s", channelName, msg.Data)
	})
	if err != nil {
		log.Fatalf("Failed to subscribe to channel %s: %v", channelName, err)
	}
	return unsubscribe
}

func main() {
	// Define command-line flags
	channelName := flag.String("channel", "", "The name of the Ably channel")
	messageName := flag.String("name", "", "The name of the message event")
	messageData := flag.String("data", "", "The data of the message")

	// Parse command-line flags
	flag.Parse()

	// Validate required flags
	if *channelName == "" || *messageName == "" || *messageData == "" {
		log.Fatalf("All flags -channel, -name, and -data are required")
	}

	// Load environment variables
	err := godotenv.Load()
	if err != nil {
		log.Fatalf("Error loading .env file: %v", err)
	}

	// Create an Ably client
	client, err := ably.NewRealtime(ably.WithKey(os.Getenv("ABLY_API_KEY")))
	if err != nil {
		log.Fatalf("Failed to create Ably client: %v", err)
	}

	// Publish the message
	messageID := PublishMessage(client, *channelName, *messageName, *messageData)

	// messageID := PublishMessageToQueue(*channelName, *messageName, *messageData)

	// Subscribe to the channel
	unsubscribe := SubscribeToChannel(client, *channelName, messageID)
	time.Sleep(10 * time.Minute)
	defer unsubscribe()
}
