package appleconsumeranalytics

import (
	"bufio"
	"context"
	"encoding/json"
	"log"
	"math"
	"time"

	"github.com/aws/aws-sdk-go/aws"
	"github.com/aws/aws-sdk-go/service/s3"
	"github.com/filtr/go-apollo/pkg/daterange"
	pb "github.com/filtr/go-apollo/pkg/leveldbgrpc"
	"google.golang.org/grpc"
	grpcCodes "google.golang.org/grpc/codes"
)

type ingestValue struct {
	Timestamp string `json:"timestamp"`
}

type dbTrackStreamMap map[string]dbTrackStreamValue
type dbTrackDemographicsMap map[string]dbTrackDemographicsValue

type dbTrackStreamValue struct {
	EndReason0         int `json:"end_reason_0,omitempty"`
	EndReason1         int `json:"end_reason_1,omitempty"`
	EndReason2         int `json:"end_reason_2,omitempty"`
	EndReason3         int `json:"end_reason_3,omitempty"`
	EndReason4         int `json:"end_reason_4,omitempty"`
	ContainerSubType0  int `json:"container_subtype_0,omitempty"`
	ContainerSubType1  int `json:"container_subtype_1,omitempty"`
	ContainerSubType2  int `json:"container_subtype_2,omitempty"`
	ContainerSubType3  int `json:"container_subtype_3,omitempty"`
	ContainerSubType4  int `json:"container_subtype_4,omitempty"`
	ContainerSubType5  int `json:"container_subtype_5,omitempty"`
	ContainerSubType6  int `json:"container_subtype_6,omitempty"`
	ContainerSubType7  int `json:"container_subtype_7,omitempty"`
	ContainerSubType8  int `json:"container_subtype_8,omitempty"`
	ContainerSubType9  int `json:"container_subtype_9,omitempty"`
	ContainerSubType10 int `json:"container_subtype_10,omitempty"`
	ContainerSubType11 int `json:"container_subtype_11,omitempty"`
	ContainerSubType12 int `json:"container_subtype_12,omitempty"`
	ContainerSubType13 int `json:"container_subtype_13,omitempty"`
	ContainerSubType14 int `json:"container_subtype_14,omitempty"`
	ContainerSubType15 int `json:"container_subtype_15,omitempty"`
	ContainerType0     int `json:"container_type_0,omitempty"`
	ContainerType1     int `json:"container_type_1,omitempty"`
	ContainerType2     int `json:"container_type_2,omitempty"`
	ContainerType3     int `json:"container_type_3,omitempty"`
	ContainerType4     int `json:"container_type_4,omitempty"`
	Offline0           int `json:"offline_0,omitempty"`
	Offline1           int `json:"offline_1,omitempty"`
	Source0            int `json:"source_0,omitempty"`
	Source1            int `json:"source_1,omitempty"`
	Source2            int `json:"source_2,omitempty"`
	Source3            int `json:"source_3,omitempty"`
	Source4            int `json:"source_4,omitempty"`
	Source5            int `json:"source_5,omitempty"`
	Source6            int `json:"source_6,omitempty"`
	Source7            int `json:"source_7,omitempty"`
	Streams            int `json:"streams"`
}

type dbTrackContainersValue struct {
	EndReason0 int `json:"end_reason_0,omitempty"`
	EndReason1 int `json:"end_reason_1,omitempty"`
	EndReason2 int `json:"end_reason_2,omitempty"`
	EndReason3 int `json:"end_reason_3,omitempty"`
	EndReason4 int `json:"end_reason_4,omitempty"`
	Source0    int `json:"source_0,omitempty"`
	Source1    int `json:"source_1,omitempty"`
	Source2    int `json:"source_2,omitempty"`
	Source3    int `json:"source_3,omitempty"`
	Source4    int `json:"source_4,omitempty"`
	Source5    int `json:"source_5,omitempty"`
	Source6    int `json:"source_6,omitempty"`
	Source7    int `json:"source_7,omitempty"`
	Streams    int `json:"streams"`
}

type dbTrackDemographicsValue struct {
	ListenersGenderFemaleAge0_17     int `json:"listeners_gender_female_age_0_17,omitempty"`
	ListenersGenderMaleAge0_17       int `json:"listeners_gender_male_age_0_17,omitempty"`
	ListenersGenderUnknownAge0_17    int `json:"listeners_gender_unknown_age_0_17,omitempty"`
	ListenersGenderFemaleAge18_24    int `json:"listeners_gender_female_age_18_24,omitempty"`
	ListenersGenderMaleAge18_24      int `json:"listeners_gender_male_age_18_24,omitempty"`
	ListenersGenderUnknownAge18_24   int `json:"listeners_gender_unknown_age_18_24,omitempty"`
	ListenersGenderFemaleAge25_34    int `json:"listeners_gender_female_age_25_34,omitempty"`
	ListenersGenderMaleAge25_34      int `json:"listeners_gender_male_age_25_34,omitempty"`
	ListenersGenderUnknownAge25_34   int `json:"listeners_gender_unknown_age_25_34,omitempty"`
	ListenersGenderFemaleAge35_44    int `json:"listeners_gender_female_age_35_44,omitempty"`
	ListenersGenderMaleAge35_44      int `json:"listeners_gender_male_age_35_44,omitempty"`
	ListenersGenderUnknownAge35_44   int `json:"listeners_gender_unknown_age_35_44,omitempty"`
	ListenersGenderFemaleAge45_54    int `json:"listeners_gender_female_age_45_54,omitempty"`
	ListenersGenderMaleAge45_54      int `json:"listeners_gender_male_age_45_54,omitempty"`
	ListenersGenderUnknownAge45_54   int `json:"listeners_gender_unknown_age_45_54,omitempty"`
	ListenersGenderFemaleAge55_64    int `json:"listeners_gender_female_age_55_64,omitempty"`
	ListenersGenderMaleAge55_64      int `json:"listeners_gender_male_age_55_64,omitempty"`
	ListenersGenderUnknownAge55_64   int `json:"listeners_gender_unknown_age_55_64,omitempty"`
	ListenersGenderFemaleAge65Plus   int `json:"listeners_gender_female_age_65_plus,omitempty"`
	ListenersGenderMaleAge65Plus     int `json:"listeners_gender_male_age_65_plus,omitempty"`
	ListenersGenderUnknownAge65Plus  int `json:"listeners_gender_unknown_age_65_plus,omitempty"`
	ListenersGenderFemaleAgeUnknown  int `json:"listeners_gender_female_age_unknown,omitempty"`
	ListenersGenderMaleAgeUnknown    int `json:"listeners_gender_male_age_unknown,omitempty"`
	ListenersGenderUnknownAgeUnknown int `json:"listeners_gender_unknown_age_unknown,omitempty"`
	StreamsGenderFemaleAge0_17       int `json:"streams_gender_female_age_0_17,omitempty"`
	StreamsGenderMaleAge0_17         int `json:"streams_gender_male_age_0_17,omitempty"`
	StreamsGenderUnknownAge0_17      int `json:"streams_gender_unknown_age_0_17,omitempty"`
	StreamsGenderFemaleAge18_24      int `json:"streams_gender_female_age_18_24,omitempty"`
	StreamsGenderMaleAge18_24        int `json:"streams_gender_male_age_18_24,omitempty"`
	StreamsGenderUnknownAge18_24     int `json:"streams_gender_unknown_age_18_24,omitempty"`
	StreamsGenderFemaleAge25_34      int `json:"streams_gender_female_age_25_34,omitempty"`
	StreamsGenderMaleAge25_34        int `json:"streams_gender_male_age_25_34,omitempty"`
	StreamsGenderUnknownAge25_34     int `json:"streams_gender_unknown_age_25_34,omitempty"`
	StreamsGenderFemaleAge35_44      int `json:"streams_gender_female_age_35_44,omitempty"`
	StreamsGenderMaleAge35_44        int `json:"streams_gender_male_age_35_44,omitempty"`
	StreamsGenderUnknownAge35_44     int `json:"streams_gender_unknown_age_35_44,omitempty"`
	StreamsGenderFemaleAge45_54      int `json:"streams_gender_female_age_45_54,omitempty"`
	StreamsGenderMaleAge45_54        int `json:"streams_gender_male_age_45_54,omitempty"`
	StreamsGenderUnknownAge45_54     int `json:"streams_gender_unknown_age_45_54,omitempty"`
	StreamsGenderFemaleAge55_64      int `json:"streams_gender_female_age_55_64,omitempty"`
	StreamsGenderMaleAge55_64        int `json:"streams_gender_male_age_55_64,omitempty"`
	StreamsGenderUnknownAge55_64     int `json:"streams_gender_unknown_age_55_64,omitempty"`
	StreamsGenderFemaleAge65Plus     int `json:"streams_gender_female_age_65_plus,omitempty"`
	StreamsGenderMaleAge65Plus       int `json:"streams_gender_male_age_65_plus,omitempty"`
	StreamsGenderUnknownAge65Plus    int `json:"streams_gender_unknown_age_65_plus,omitempty"`
	StreamsGenderFemaleAgeUnknown    int `json:"streams_gender_female_age_unknown,omitempty"`
	StreamsGenderMaleAgeUnknown      int `json:"streams_gender_male_age_unknown,omitempty"`
	StreamsGenderUnknownAgeUnknown   int `json:"streams_gender_unknown_age_unknown,omitempty"`
	Listeners                        int `json:"listeners"`
	Streams                          int `json:"streams"`
}

func getRows(s3Service *s3.S3, s3TargetBucket, sourceKey string) ([]string, error) {
	input := &s3.GetObjectInput{
		Bucket: aws.String(s3TargetBucket),
		Key:    aws.String(sourceKey),
	}

	result, err := s3Service.GetObject(input)
	if err != nil {
		log.Fatal(err)
	}

	var rows []string

	scanner := bufio.NewScanner(result.Body)
	for scanner.Scan() {
		var text = scanner.Text()
		rows = append(rows, text)
	}
	if err := scanner.Err(); err != nil {
		return nil, err
	}

	return rows, nil
}

func putIngestItem(client pb.LevelDBClient, key []byte) error {
	iv := ingestValue{
		Timestamp: time.Now().UTC().Format(time.RFC3339Nano),
	}
	jsonIv, err1 := json.Marshal(iv)
	if err1 != nil {
		return err1
	}
	_, err2 := client.Put(context.Background(), &pb.PutRequest{Key: key, Value: jsonIv})
	if err2 != nil {
		return err2
	}
	return nil
}

func isIngested(client pb.LevelDBClient, key []byte) bool {
	_, err := client.Get(context.Background(), &pb.GetRequest{Key: key})
	return err == nil
}

func makeTrackKey(r trackRecord) string {
	return "Aÿ" + r.ISRC + "ÿ" + r.Date + "ÿ" + r.CountryCode + "ÿ" + r.Licensor + "ÿ" + r.ReportDate + "ÿ" + r.ReportVendorID
}

func makeTrackContainersKey(r containerTrackRecord) string {
	return "Bÿ" + r.ISRC + "ÿ" + r.Date + "ÿ" + r.CountryCode + "ÿ" + r.Licensor + "ÿ" + r.ReportDate + "ÿ" + r.ReportVendorID
}

func makeTrackDemographicsKey(r contentDemographicsRecord) string {
	return "Cÿ" + r.ISRC + "ÿ" + r.Date + "ÿ" + r.CountryCode + "ÿ" + r.Licensor + "ÿ" + r.ReportDate + "ÿ" + r.ReportVendorID
}

func makeTrackRecordKey(countryCode, licensor, reportDate, vendorID string) string {
	return countryCode + "ÿ" + licensor + "ÿ" + reportDate + "ÿ" + vendorID
}

func makeTrackDemographicsRecordKey(countryCode, licensor, reportDate, vendorID string) string {
	return countryCode + "ÿ" + licensor + "ÿ" + reportDate + "ÿ" + vendorID
}

func makeIngestTracksDBKey(date, licensor, vendorID string) []byte {
	return []byte("Zÿtracksÿ" + date + "ÿ" + licensor + "ÿ" + vendorID)
}

func makeIngestTrackContainersDBKey(date, licensor, vendorID string) []byte {
	return []byte("Zÿtrack_containersÿ" + date + "ÿ" + licensor + "ÿ" + vendorID)
}

func makeIngestTrackDemographicsDBKey(date, licensor, vendorID string) []byte {
	return []byte("Zÿtrack_demographicsÿ" + date + "ÿ" + licensor + "ÿ" + vendorID)
}

func batchIngest(client pb.LevelDBClient, list []*pb.BatchPutItem) {
	chunks := int(math.Ceil(float64(len(list)) / 100))
	count := 0

	for i := 0; i < chunks; i++ {
		newCount := count + 100
		start := int(math.Min(float64(count), float64(len(list))))
		end := int(math.Min(float64(newCount), float64(len(list))))
		items := list[start:end]

		r, err := client.BatchPut(context.Background(), &pb.BatchPutRequest{Items: items})
		if err != nil {
			log.Fatal(err)
		}
		_ = r
		count = newCount
	}
}

func reduceTrackStreamItem(client pb.LevelDBClient, key string, value dbTrackStreamMap) *pb.BatchPutItem {
	newValue := make(map[string]dbTrackStreamValue)
	var prevValue map[string]dbTrackStreamValue
	itemKey := []byte(key)

	prevItem, err := client.Get(context.Background(), &pb.GetRequest{Key: itemKey})
	if err != nil && grpc.Code(err) != grpcCodes.NotFound {
		log.Fatal(err)
	}

	if prevItem != nil {
		json.Unmarshal(prevItem.Value, &prevValue)
		for pk, pv := range prevValue {
			newValue[pk] = pv
		}
	}

	for nk, nv := range value {
		newValue[nk] = nv
	}

	newItemValue, err := json.Marshal(newValue)
	if err != nil {
		log.Fatal(err)
	}
	return &pb.BatchPutItem{Key: itemKey, Value: newItemValue}
}

func reduceTrackDemographicsItem(client pb.LevelDBClient, key string, value dbTrackDemographicsMap) *pb.BatchPutItem {
	newValue := make(map[string]dbTrackDemographicsValue)
	var prevValue map[string]dbTrackDemographicsValue
	itemKey := []byte(key)

	prevItem, err := client.Get(context.Background(), &pb.GetRequest{Key: itemKey})
	if err != nil && grpc.Code(err) != grpcCodes.NotFound {
		log.Fatal(err)
	}

	if prevItem != nil {
		json.Unmarshal(prevItem.Value, &prevValue)
		for pk, pv := range prevValue {
			newValue[pk] = pv
		}
	}

	for nk, nv := range value {
		newValue[nk] = nv
	}

	newItemValue, err := json.Marshal(newValue)
	if err != nil {
		log.Fatal(err)
	}
	return &pb.BatchPutItem{Key: itemKey, Value: newItemValue}
}

func processTrackDemographics(client pb.LevelDBClient, s3Service *s3.S3, s3TargetBucket, reportDate, licensor, vendorID string) {
	ingestKey := makeIngestTrackDemographicsDBKey(reportDate, licensor, vendorID)
	if isIngested(client, ingestKey) == true {
		log.Println("track_demographics - already ingested", reportDate, licensor, vendorID)
		return
	}

	sourceKey := makeContentDemographicsTargetS3Key(reportDate, licensor, vendorID)
	sourceExists := targetFileExists(s3Service, s3TargetBucket, sourceKey)

	if sourceExists != true {
		log.Println("track_demographics - source doesn't exist", reportDate, licensor, vendorID)
		return
	}

	log.Println("track_demographics - starting", reportDate, licensor, vendorID)

	var list []*pb.BatchPutItem

	rows, err := getRows(s3Service, s3TargetBucket, sourceKey)
	if err != nil {
		log.Fatal(err)
	}

	for _, row := range rows {
		var r contentDemographicsRecord
		var v dbTrackDemographicsValue
		json.Unmarshal([]byte(row), &r)
		json.Unmarshal([]byte(row), &v)
		key := makeTrackDemographicsKey(r)
		value, err := json.Marshal(v)
		if err != nil {
			log.Fatal(err)
		}
		item := &pb.BatchPutItem{Key: []byte(key), Value: value}
		list = append(list, item)
	}

	if len(list) > 0 {
		batchIngest(client, list)
	}

	putIngestItem(client, ingestKey)
}

func processTrackStreams(client pb.LevelDBClient, s3Service *s3.S3, s3TargetBucket, reportDate, licensor, vendorID string) {
	ingestKey := makeIngestTracksDBKey(reportDate, licensor, vendorID)
	if isIngested(client, ingestKey) == true {
		log.Println("track_streams - already ingested", reportDate, licensor, vendorID)
		return
	}

	sourceKey := makeStreamsTracksTargetS3Key(reportDate, licensor, vendorID)
	sourceExists := targetFileExists(s3Service, s3TargetBucket, sourceKey)

	if sourceExists != true {
		log.Println("track_streams - source doesn't exist", reportDate, licensor, vendorID)
		return
	}

	log.Println("track_streams - starting", reportDate, licensor, vendorID)

	var list []*pb.BatchPutItem

	rows, err := getRows(s3Service, s3TargetBucket, sourceKey)
	if err != nil {
		log.Fatal(err)
	}

	for _, row := range rows {
		var r trackRecord
		var v dbTrackStreamValue
		json.Unmarshal([]byte(row), &r)
		json.Unmarshal([]byte(row), &v)
		key := makeTrackKey(r)
		value, err := json.Marshal(v)
		if err != nil {
			log.Fatal(err)
		}
		item := &pb.BatchPutItem{Key: []byte(key), Value: value}
		list = append(list, item)
	}

	if len(list) > 0 {
		batchIngest(client, list)
	}

	putIngestItem(client, ingestKey)
}

func processTrackContainerStreams(client pb.LevelDBClient, s3Service *s3.S3, s3TargetBucket, reportDate, licensor, vendorID string) {
	ingestKey := makeIngestTrackContainersDBKey(reportDate, licensor, vendorID)
	if isIngested(client, ingestKey) == true {
		log.Println("track_container_streams - already ingested", reportDate, licensor, vendorID)
		return
	}

	sourceKey := makeStreamsContainerTracksTargetS3Key(reportDate, licensor, vendorID)
	sourceExists := targetFileExists(s3Service, s3TargetBucket, sourceKey)

	if sourceExists != true {
		log.Println("track_container_streams - source doesn't exist", reportDate, licensor, vendorID)
		return
	}

	log.Println("track_container_streams - starting", reportDate, licensor, vendorID)

	var list []*pb.BatchPutItem

	rows, err := getRows(s3Service, s3TargetBucket, sourceKey)
	if err != nil {
		log.Fatal(err)
	}

	items := make(map[string]map[string]dbTrackContainersValue)

	for _, row := range rows {
		var r containerTrackRecord
		json.Unmarshal([]byte(row), &r)
		key := makeTrackContainersKey(r)

		item, ok := items[key]
		if ok == false {
			item = make(map[string]dbTrackContainersValue)
		}

		trackContainer, ok := item[r.ContainerID]
		if ok == false {
			trackContainer = dbTrackContainersValue{
				EndReason0: r.EndReason0,
				EndReason1: r.EndReason1,
				EndReason2: r.EndReason2,
				EndReason3: r.EndReason3,
				EndReason4: r.EndReason4,
				Source0:    r.Source0,
				Source1:    r.Source1,
				Source2:    r.Source2,
				Source3:    r.Source3,
				Source4:    r.Source4,
				Source5:    r.Source5,
				Source6:    r.Source6,
				Source7:    r.Source7,
				Streams:    r.Streams,
			}
		} else {
			trackContainer = dbTrackContainersValue{
				EndReason0: r.EndReason0 + trackContainer.EndReason0,
				EndReason1: r.EndReason1 + trackContainer.EndReason1,
				EndReason2: r.EndReason2 + trackContainer.EndReason2,
				EndReason3: r.EndReason3 + trackContainer.EndReason3,
				EndReason4: r.EndReason4 + trackContainer.EndReason4,
				Source0:    r.Source0 + trackContainer.Source0,
				Source1:    r.Source1 + trackContainer.Source1,
				Source2:    r.Source2 + trackContainer.Source2,
				Source3:    r.Source3 + trackContainer.Source3,
				Source4:    r.Source4 + trackContainer.Source4,
				Source5:    r.Source5 + trackContainer.Source5,
				Source6:    r.Source6 + trackContainer.Source6,
				Source7:    r.Source7 + trackContainer.Source7,
				Streams:    r.Streams + trackContainer.Streams,
			}
		}

		item[r.ContainerID] = trackContainer
		items[key] = item
	}

	for key, v := range items {
		newItemValue, err := json.Marshal(v)
		if err != nil {
			log.Fatal(err)
		}
		item := &pb.BatchPutItem{Key: []byte(key), Value: newItemValue}
		list = append(list, item)
	}

	if len(list) > 0 {
		batchIngest(client, list)
	}

	putIngestItem(client, ingestKey)
}

// DBIngest ingests Apple data
func DBIngest(s3Service *s3.S3, startDate, endDate time.Time, licensors []string, dbAddress, report, s3Bucket string) {
	conn, err := grpc.Dial(dbAddress, grpc.WithInsecure())
	if err != nil {
		log.Fatalf("did not connect: %v", err)
	}
	defer conn.Close()
	client := pb.NewLevelDBClient(conn)
	dates := daterange.DateRange(startDate, endDate)

	for _, date := range dates {
		for _, licensor := range licensors {
			for _, vendorID := range VendorMap[licensor] {
				if report == "" {
					processTrackStreams(client, s3Service, s3Bucket, date, licensor, vendorID)
					processTrackContainerStreams(client, s3Service, s3Bucket, date, licensor, vendorID)
					if date >= "2017-08-02" {
						processTrackDemographics(client, s3Service, s3Bucket, date, licensor, vendorID)
					}
				}
				if report == "tracks_streams" {
					processTrackStreams(client, s3Service, s3Bucket, date, licensor, vendorID)
				}
				if report == "tracks_container_streams" {
					processTrackContainerStreams(client, s3Service, s3Bucket, date, licensor, vendorID)
				}
				if report == "tracks_demographics" && date >= "2017-08-02" {
					processTrackDemographics(client, s3Service, s3Bucket, date, licensor, vendorID)
				}
			}
		}
	}
}
