// Sensor enrichment with fluss-go v0.1.0-beta.10.
//
// The flow follows the IoT scenario in Apache Fluss's Java Client article.
// In Java, Connection creates Admin and Table objects; here fgo.Client is the
// shared entry point, while fadm.New provides the administration API.
//
// Install and run:
//
//	go mod init fluss-go-sensor-enrichment
//	go get github.com/pletorco/fluss-go@v0.1.0-beta.10
//	go run sensor-enrichment.go
//
// Set FLUSS_BOOTSTRAP when the Coordinator is not localhost:9123.
package main

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

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

const database = "go_sensor_demo"

var (
	sensorInfoPath = fgo.TablePath{Database: database, Table: "sensor_info"}
	readingPath    = fgo.TablePath{Database: database, Table: "sensor_readings"}
)

type sensorInfo struct {
	SensorID int32  `json:"sensor_id"`
	Name     string `json:"name"`
	Location string `json:"location"`
	State    string `json:"state"`
}

type sensorReading struct {
	SensorID     int32     `json:"sensor_id"`
	MeasuredAt   time.Time `json:"measured_at"`
	TemperatureC float64   `json:"temperature_c"`
	HumidityPct  float64   `json:"humidity_pct"`
}

type enrichedReading struct {
	SensorID     int32     `json:"sensor_id"`
	MeasuredAt   time.Time `json:"measured_at"`
	TemperatureC float64   `json:"temperature_c"`
	HumidityPct  float64   `json:"humidity_pct"`
	SensorName   string    `json:"sensor_name"`
	Location     string    `json:"location"`
	State        string    `json:"state"`
}

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

	// Java equivalent: ConnectionFactory.createConnection(...).
	// fgo.Client is the shared connection used by both the data and admin APIs.
	client, err := fgo.Open(
		ctx,
		fgo.WithBootstrapServers(bootstrapAddress()),
		fgo.WithClientSoftware("sensor-enrichment", "0.1.0"),
	)
	if err != nil {
		log.Fatal(err)
	}
	defer closeClient(client)

	// Java equivalent: connection.getAdmin().
	admin, err := fadm.New(client)
	if err != nil {
		log.Fatal(err)
	}
	// Repeated tutorial runs reuse the same database and tables.
	prepareTables(ctx, admin)

	// Java equivalent: connection.getTable(...). The Go client loads the table
	// metadata before it creates data-operation objects for that table.
	infos, err := client.GetTable(ctx, sensorInfoPath)
	if err != nil {
		log.Fatal(err)
	}
	readings, err := client.GetTable(ctx, readingPath)
	if err != nil {
		log.Fatal(err)
	}

	// First write current sensor state, then append immutable measurements.
	writeSensorInfo(ctx, client, infos)
	first, last, readingCount := appendReadings(ctx, client, readings)
	scanAndEnrich(ctx, client, infos, readings, first, last, readingCount)
}

func prepareTables(ctx context.Context, admin *fadm.Client) {
	// Java equivalent: admin.createTable(...). The final true is ignoreIfExists.
	if err := admin.CreateDatabase(ctx, database, fadm.DatabaseDescriptor{
		Comment: "sensor-enrichment tutorial",
	}, true); err != nil {
		log.Fatal(err)
	}
	if err := admin.CreateTable(ctx, sensorInfoPath, fadm.TableDescriptor{
		Comment:     "current sensor metadata",
		BucketCount: 1,
		Schema: fgo.Schema{
			Columns: []fgo.Column{
				{Name: "sensor_id", Type: fgo.IntType},
				{Name: "name", Type: fgo.StringType},
				{Name: "location", Type: fgo.StringType},
				{Name: "state", Type: fgo.StringType},
			},
			PrimaryKey: []string{"sensor_id"},
			BucketKey:  []string{"sensor_id"},
		},
	}, true); err != nil {
		log.Fatal(err)
	}
	if err := admin.CreateTable(ctx, readingPath, fadm.TableDescriptor{
		Comment:     "append-only sensor readings",
		BucketCount: 1,
		Schema: fgo.Schema{
			Columns: []fgo.Column{
				{Name: "sensor_id", Type: fgo.IntType},
				{
					Name: "measured_at",
					Type: fgo.TimestampType,
					// Fluss 0.9.1 needs the timestamp precision in schema JSON.
					LogicalType: &fgo.LogicalType{
						Root:      "TIMESTAMP_WITHOUT_TIME_ZONE",
						Precision: 3,
					},
				},
				{Name: "temperature_c", Type: fgo.DoubleType},
				{Name: "humidity_pct", Type: fgo.DoubleType},
			},
			BucketKey: []string{"sensor_id"},
		},
	}, true); err != nil {
		log.Fatal(err)
	}
	log.Printf("[1/4] Prepared %s.sensor_info and %s.sensor_readings", database, database)
}

func writeSensorInfo(ctx context.Context, client *fgo.Client, table fgo.Table) {
	infos := []sensorInfo{
		// Ten initial sensor states.
		{1, "Roof temperature sensor", "roof", "OK"},
		{2, "Lobby humidity sensor", "lobby", "ERROR"},
		{3, "Server room sensor", "server-room", "MAINTENANCE"},
		{4, "Warehouse pressure sensor", "warehouse", "OK"},
		{5, "Conference room humidity sensor", "conference-room", "OK"},
		{6, "Office 1 temperature sensor", "office-1", "LOW_BATTERY"},
		{7, "Office 2 humidity sensor", "office-2", "OK"},
		{8, "Lab pressure sensor", "lab", "ERROR"},
		{9, "Parking pressure sensor", "parking", "OK"},
		{10, "Backyard temperature sensor", "backyard", "OK"},
		// These later rows overwrite the current state for existing keys.
		{2, "Lobby humidity sensor", "lobby", "OK"},
		{8, "Lab pressure sensor", "lab", "CALIBRATING"},
	}
	printJSON("[2/4] Upsert these current sensor states:", infos)

	// Java equivalent: sensorInfoTable.newUpsert().createWriter().
	writer, err := client.NewUpsertWriter(ctx, table, fgo.WithUpsertBatchLimits(1<<20, len(infos)))
	if err != nil {
		log.Fatal(err)
	}
	defer closeUpsertWriter(ctx, writer)
	for _, info := range infos {
		// Await each result so the example reports a failed write immediately.
		result := writer.Upsert(ctx, fgo.Row{info.SensorID, info.Name, info.Location, info.State}).Await(ctx)
		if result.Err != nil {
			log.Fatal(result.Err)
		}
	}
	// Flush waits for writes already accepted by the writer to finish.
	if err := writer.Flush(ctx); err != nil {
		log.Fatal(err)
	}
}

func appendReadings(ctx context.Context, client *fgo.Client, table fgo.Table) (fgo.WriteResult, fgo.WriteResult, int) {
	// The timestamps are deterministic so the printed JSON is easy to compare.
	base := time.Date(2026, time.August, 4, 9, 0, 0, 0, time.UTC)
	readings := []sensorReading{
		{1, base, 24.1, 43.2},
		{2, base.Add(15 * time.Minute), 22.7, 52.8},
		{3, base.Add(30 * time.Minute), 20.4, 38.5},
		{4, base.Add(45 * time.Minute), 18.9, 48.1},
		{5, base.Add(60 * time.Minute), 23.5, 46.3},
		{6, base.Add(75 * time.Minute), 21.8, 44.9},
		{7, base.Add(90 * time.Minute), 22.1, 47.5},
		{8, base.Add(105 * time.Minute), 20.7, 49.2},
		{9, base.Add(120 * time.Minute), 19.6, 51.7},
		{10, base.Add(135 * time.Minute), 25.0, 41.8},
	}
	printJSON("[3/4] Append these immutable readings:", readings)

	// Java equivalent: readingsTable.newAppend().createWriter().
	writer, err := client.NewAppendWriter(ctx, table,
		fgo.WithAppendBatchLimits(1<<20, len(readings)),
		fgo.WithAppendBatchTimeout(5*time.Millisecond),
	)
	if err != nil {
		log.Fatal(err)
	}
	defer closeAppendWriter(ctx, writer)

	futures := make([]*fgo.WriteFuture, len(readings))
	for i, reading := range readings {
		futures[i] = writer.Append(ctx, fgo.Row{
			reading.SensorID, reading.MeasuredAt, reading.TemperatureC, reading.HumidityPct,
		})
	}
	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 || !results[i].OffsetKnown {
			log.Fatalf("append %d: %v", i, results[i].Err)
		}
	}
	// This tutorial uses one bucket. Confirm that these writes form one
	// contiguous range before using the first and last offsets for the scan.
	for i, result := range results {
		if result.Bucket != results[0].Bucket || result.BaseOffset != results[0].BaseOffset+int64(i) {
			log.Fatal("this single-bucket tutorial expected one contiguous offset range")
		}
	}
	return results[0], results[len(results)-1], len(readings)
}

func scanAndEnrich(ctx context.Context, client *fgo.Client, infos, readings fgo.Table, first, last fgo.WriteResult, readingCount int) {
	// Java equivalent: sensorInfoTable.newLookup().createLookuper().
	lookuper, err := client.NewLookuper(ctx, infos, fgo.WithLookupBatchLimits(100, 4))
	if err != nil {
		log.Fatal(err)
	}
	defer closeLookuper(lookuper)

	// Java equivalent: readingsTable.newScan().createLogScanner(). Start at
	// this run's first offset and stop just after its final offset, so
	// rows appended by another run are not included in the demonstration.
	scanner, err := client.NewLogScanner(ctx, readings, fgo.AtOffset(first.BaseOffset),
		fgo.WithScanRowLimit(int64(readingCount)),
		fgo.WithScanStoppingOffsets(map[int32]int64{first.Bucket: last.BaseOffset + 1}),
	)
	if err != nil {
		log.Fatal(err)
	}
	defer closeScanner(scanner)

	enriched := make([]enrichedReading, 0, readingCount)
	for !scanner.Done() {
		batch, err := scanner.Poll(ctx)
		if err != nil {
			log.Fatal(err)
		}
		for _, record := range batch.Records {
			reading := readingFromRow(record.Record.Value)
			// A lookup returns the state that is current when this code runs,
			// not necessarily the state that existed when the event was written.
			lookups := lookuper.Lookup(ctx, fgo.PrimaryKey{reading.SensorID})
			if len(lookups) != 1 {
				log.Fatalf("sensor %d: expected one lookup result, got %d", reading.SensorID, len(lookups))
			}
			lookup := lookups[0]
			switch {
			case errors.Is(lookup.Err, fgo.ErrNotFound):
				log.Fatalf("sensor %d has no current metadata", reading.SensorID)
			case lookup.Err != nil:
				log.Fatal(lookup.Err)
			}
			info := infoFromRow(lookup.Row)
			enriched = append(enriched, enrichedReading{
				SensorID: reading.SensorID, MeasuredAt: reading.MeasuredAt,
				TemperatureC: reading.TemperatureC, HumidityPct: reading.HumidityPct,
				SensorName: info.Name, Location: info.Location, State: info.State,
			})
		}
		// Release Arrow-backed buffers once every record in the batch is copied.
		batch.Release()
	}
	printJSON("[4/4] Scan readings and enrich each with current sensor metadata:", enriched)
}

func readingFromRow(row fgo.Row) sensorReading {
	return sensorReading{row[0].(int32), row[1].(time.Time), row[2].(float64), row[3].(float64)}
}

func infoFromRow(row fgo.Row) sensorInfo {
	return sensorInfo{row[0].(int32), row[1].(string), row[2].(string), row[3].(string)}
}

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

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

func closeClient(client *fgo.Client) {
	if err := client.Close(); err != nil {
		log.Printf("close Fluss client: %v", err)
	}
}

func closeUpsertWriter(ctx context.Context, writer *fgo.UpsertWriter) {
	if err := writer.Close(ctx); err != nil {
		log.Printf("close upsert writer: %v", err)
	}
}

func closeAppendWriter(ctx context.Context, writer *fgo.AppendWriter) {
	if err := writer.Close(ctx); err != nil {
		log.Printf("close append writer: %v", err)
	}
}

func closeLookuper(lookuper *fgo.Lookuper) {
	if err := lookuper.Close(); err != nil {
		log.Printf("close lookuper: %v", err)
	}
}

func closeScanner(scanner *fgo.LogScanner) {
	if err := scanner.Close(); err != nil {
		log.Printf("close log scanner: %v", err)
	}
}
