// Run after completing the Apache Fluss Flink Quickstart.
//
//	go mod init fluss-go-quickstart
//	go get github.com/pletorco/fluss-go/pkg/fgo@v0.1.0-beta.10
//	go run log-table.go
//
// Set FLUSS_BOOTSTRAP when the Coordinator is not at localhost:9123.
package main

import (
	"context"
	"encoding/json"
	"log"
	"math/big"
	"os"
	"time"

	"github.com/pletorco/fluss-go/pkg/fgo"
)

// orderEvent is the JSON-friendly representation printed by this example.
// total_price stays a string in JSON so decimal precision is never lost.
type orderEvent struct {
	OrderID       int64  `json:"order_id"`
	CustomerID    int32  `json:"customer_id"`
	TotalPrice    string `json:"total_price"`
	OrderedOn     string `json:"ordered_on"`
	OrderPriority string `json:"order_priority"`
	Clerk         string `json:"clerk"`
}

func main() {
	ctx := context.Background()
	log.SetFlags(0)

	// One client owns the Coordinator and TabletServer connections. Reuse it for
	// both the append and the scan below.
	client, err := fgo.Open(
		ctx,
		fgo.WithBootstrapServers(bootstrapAddress()),
		fgo.WithClientSoftware("quickstart-go", "0.1.0"),
	)
	if err != nil {
		log.Fatal(err)
	}
	defer func() {
		if err := client.Close(); err != nil {
			log.Printf("close Fluss client: %v", err)
		}
	}()

	// Opening the table reads its schema and bucket-routing metadata. The
	// Quickstart's order_events table has no primary key, so it is a Log Table.
	events, err := client.GetTable(ctx, fgo.TablePath{
		Database: "demo",
		Table:    "order_events",
	})
	if err != nil {
		log.Fatal(err)
	}

	// A Log Table only appends. A repeat of the same business values is another
	// event, not an update to an existing row.
	writer, err := client.NewAppendWriter(
		ctx,
		events,
		fgo.WithAppendBatchLimits(1<<20, 500),
		fgo.WithAppendBatchTimeout(5*time.Millisecond),
	)
	if err != nil {
		log.Fatal(err)
	}
	defer func() {
		if err := writer.Close(ctx); err != nil {
			log.Printf("close append writer: %v", err)
		}
	}()

	// Queue four rows before waiting. With the Quickstart's one bucket, the writer
	// preserves their order and can send them as one collected batch.
	orderIDBase := time.Now().UnixNano()
	eventsToAppend := []orderEvent{
		{orderIDBase, 999999, "42.15", "2026-08-01", "low", "Go Client Clerk"},
		{orderIDBase + 1, 999999, "85.60", "2026-08-02", "medium", "Go Client Clerk"},
		{orderIDBase + 2, 999999, "129.90", "2026-08-03", "high", "Go Client Clerk"},
		{orderIDBase + 3, 999999, "24.50", "2026-08-04", "low", "Go Client Clerk"},
	}
	printJSON("[1/3] Queue these four new events:", eventsToAppend)

	futures := make([]*fgo.WriteFuture, len(eventsToAppend))
	for i, event := range eventsToAppend {
		futures[i] = writer.Append(ctx, event.row())
	}
	if err := writer.Flush(ctx); err != nil {
		log.Fatal(err)
	}

	results := make([]fgo.WriteResult, len(futures))
	for i, future := range futures {
		results[i] = future.Await(ctx)
		if results[i].Err != nil {
			log.Fatal(results[i].Err)
		}
		if !results[i].OffsetKnown {
			log.Fatal("append succeeded but did not return an offset")
		}
	}
	first := results[0]
	for i, result := range results {
		expectedOffset := first.BaseOffset + int64(i)
		if result.Bucket != first.Bucket || result.BaseOffset != expectedOffset {
			log.Fatalf("events do not occupy one contiguous range: result %d is bucket=%d offset=%d, want bucket=%d offset=%d", i, result.Bucket, result.BaseOffset, first.Bucket, expectedOffset)
		}
	}
	log.Printf("[2/3] Fluss stored four events in bucket=%d at offsets %d through %d", first.Bucket, first.BaseOffset, results[len(results)-1].BaseOffset)

	// The Flink Quickstart creates order_events with one bucket. Starting at the
	// first returned offset and limiting the scan to four rows reads this group.
	scanner, err := client.NewLogScanner(
		ctx,
		events,
		fgo.AtOffset(first.BaseOffset),
		fgo.WithScanRowLimit(int64(len(eventsToAppend))),
		fgo.WithScanStoppingOffsets(map[int32]int64{
			first.Bucket: results[len(results)-1].BaseOffset + 1,
		}),
	)
	if err != nil {
		log.Fatal(err)
	}
	defer func() {
		if err := scanner.Close(); err != nil {
			log.Printf("close log scanner: %v", err)
		}
	}()

	// Poll can wait for records. The row limit makes this loop finish after the
	// four events just appended have been returned.
	readEvents := make([]orderEvent, 0, len(eventsToAppend))
	for !scanner.Done() {
		batch, err := scanner.Poll(ctx)
		if err != nil {
			log.Fatal(err)
		}
		for _, record := range batch.Records {
			readEvents = append(readEvents, orderEventFromRow(record.Record.Value))
		}
		// Release returns the batch buffers once every record has been handled.
		batch.Release()
	}
	if len(readEvents) != len(eventsToAppend) {
		log.Fatalf("read %d events, want %d", len(readEvents), len(eventsToAppend))
	}
	printJSON("[3/3] Read the appended events:", readEvents)
}

func (event orderEvent) row() fgo.Row {
	totalPrice, ok := new(big.Rat).SetString(event.TotalPrice)
	if !ok {
		log.Fatalf("invalid total_price: %q", event.TotalPrice)
	}
	orderedOn, err := time.Parse(time.DateOnly, event.OrderedOn)
	if err != nil {
		log.Fatalf("invalid ordered_on: %q", event.OrderedOn)
	}
	return fgo.Row{event.OrderID, event.CustomerID, totalPrice, orderedOn, event.OrderPriority, event.Clerk}
}

func orderEventFromRow(row fgo.Row) orderEvent {
	totalPrice, ok := row[2].(*big.Rat)
	if !ok {
		log.Fatalf("unexpected total_price type: %T", row[2])
	}
	return orderEvent{
		OrderID:       row[0].(int64),
		CustomerID:    row[1].(int32),
		TotalPrice:    totalPrice.FloatString(2),
		OrderedOn:     row[3].(time.Time).Format(time.DateOnly),
		OrderPriority: row[4].(string),
		Clerk:         row[5].(string),
	}
}

func printJSON(step string, value any) {
	encoded, err := json.MarshalIndent(value, "", "  ")
	if err != nil {
		log.Fatal(err)
	}
	log.Printf("%s\n%s", step, encoded)
}

func bootstrapAddress() string {
	if address := os.Getenv("FLUSS_BOOTSTRAP"); address != "" {
		return address
	}
	return "localhost:9123"
}
