package pipeline

import (
	"context"
	"fmt"
	"log"
	"os"
	"path/filepath"
	"sort"
	"strings"
	"sync"
	"sync/atomic"

	"pinscrape-allgo/internal/clients/ai"
	"pinscrape-allgo/internal/clients/pinterest"
	"pinscrape-allgo/internal/clients/suggest"
	"pinscrape-allgo/internal/config"
	"pinscrape-allgo/internal/db"
	"pinscrape-allgo/internal/metrics"
	"pinscrape-allgo/internal/util"
)

type Env struct {
	Cfg     *config.Config
	Metrics *metrics.Metrics
	Pin     *pinterest.Client
	Suggest *suggest.Client
	AI      *ai.Client
}

// Convert reads keywords/*.txt into database/<name>.sqlite (parity with
// scripts/convert.py).
func Convert(cfg *config.Config) error {
	kwDir := cfg.Path("keywords")
	entries, err := os.ReadDir(kwDir)
	if err != nil {
		return fmt.Errorf("keywords folder not found: %w", err)
	}
	var files []string
	for _, e := range entries {
		if !e.IsDir() && strings.EqualFold(filepath.Ext(e.Name()), ".txt") {
			files = append(files, filepath.Join(kwDir, e.Name()))
		}
	}
	sort.Strings(files)
	if len(files) == 0 {
		log.Println("No .txt files found in keywords/ folder")
		return nil
	}
	dbDir := cfg.Path(cfg.DataFolder)
	if err := os.MkdirAll(dbDir, 0755); err != nil {
		return err
	}

	for i, f := range files {
		base := strings.TrimSuffix(filepath.Base(f), filepath.Ext(f))
		dbPath := filepath.Join(dbDir, base+".sqlite")

		data, err := os.ReadFile(f)
		if err != nil {
			return err
		}
		var keywords []string
		for _, line := range strings.Split(string(data), "\n") {
			line = strings.TrimSpace(line)
			if line != "" {
				keywords = append(keywords, line)
			}
		}
		if len(keywords) == 0 {
			log.Printf("[%d/%d] No keywords in %s, skipping", i+1, len(files), filepath.Base(f))
			continue
		}

		conn, err := db.Open(dbPath)
		if err != nil {
			return err
		}
		if err := db.EnsureSchema(conn); err != nil {
			db.Close(conn)
			return err
		}

		added, err := appendKeywords(conn, keywords)
		if err != nil {
			db.Close(conn)
			return err
		}
		db.Close(conn)
		log.Printf("[%d/%d] %s -> %s (added %d/%d keywords)", i+1, len(files),
			filepath.Base(f), filepath.Base(dbPath), added, len(keywords))
	}
	return nil
}

// UpdateStatus recomputes status 0/1/2 for all databases (parity with
// pinscrape/database.py update_all_status).
func UpdateStatus(cfg *config.Config) error {
	files, err := db.ListDBFiles(cfg.Path(cfg.DataFolder))
	if err != nil {
		return err
	}
	for _, f := range files {
		conn, err := db.Open(f)
		if err != nil {
			return err
		}
		res, err := db.UpdateAllStatus(conn)
		db.Close(conn)
		if err != nil {
			return err
		}
		log.Printf("%s: status updated -> 0:%d 1:%d 2:%d",
			filepath.Base(f), res["status_0"], res["status_1"], res["status_2"])
	}
	return nil
}

// ScrapeImagesDesc is phase 1 of the scrape: one Pinterest request per
// keyword fills images and descriptions from the same response (merged
// search), so no keyword is ever fetched twice. Nothing else is called in
// this phase; related keywords are phase 2 (ScrapeRelatedKw), so at most one
// request per worker is ever in flight.
func ScrapeImagesDesc(ctx context.Context, env *Env) error {
	files, err := db.ListDBFiles(env.Cfg.Path(env.Cfg.DataFolder))
	if err != nil {
		return err
	}
	for _, f := range files {
		if err := scrapeImagesDescDB(ctx, env, f); err != nil {
			return err
		}
	}
	return nil
}

func scrapeImagesDescDB(ctx context.Context, env *Env, dbPath string) error {
	conn, err := db.Open(dbPath)
	if err != nil {
		return err
	}
	pending, err := db.GetPendingImagesDesc(conn)
	if err != nil {
		db.Close(conn)
		return err
	}
	if len(pending) == 0 {
		log.Printf("No keywords needing images/descriptions in %s", filepath.Base(dbPath))
		db.Close(conn)
		return nil
	}
	log.Printf("Found %d keywords needing images/descriptions in %s", len(pending), filepath.Base(dbPath))

	writer := db.NewWriter(conn, db.WriterOptions{
		BatchSize: env.Cfg.Go.WriterBatchSize,
	})

	sem := make(chan struct{}, env.Cfg.Concurrency)
	var wg sync.WaitGroup
	var done atomic.Int64

	for _, item := range pending {
		if ctx.Err() != nil {
			break
		}
		wg.Add(1)
		sem <- struct{}{}
		go func(item db.WorkRow) {
			defer wg.Done()
			defer func() { <-sem }()
			if ctx.Err() != nil {
				return
			}

			// Descriptions come free in the same response: only ask for
			// them when the description half is still empty (the request
			// is sized up front, see Client.pageSizeFor).
			wantDesc := 0
			if item.NeedDesc {
				wantDesc = env.Cfg.MinSnippet
			}
			results, serr := env.Pin.Search(ctx, item.Keyword, env.Cfg.ImageResult, wantDesc)
			if serr != nil {
				log.Printf("  search warn %s: %v", item.Keyword, serr)
			}

			cleaned := make([]string, 0, len(results.Descriptions))
			for _, d := range results.Descriptions {
				if fd := util.FormatDescription(d); fd != "" {
					cleaned = append(cleaned, fd)
				}
			}

			if item.NeedImages && len(results.Images) > 0 {
				trimmed := results.Images
				if len(trimmed) > env.Cfg.ImageResult {
					trimmed = trimmed[:env.Cfg.ImageResult]
				}
				imagesJSON := marshalImages(trimmed)
				if item.NeedDesc {
					snippet := mergeSnippet(item.Existing, nil, cleaned)
					db.UpdateImagesAndSnippet(ctx, writer, item.ID, imagesJSON, snippet)
				} else {
					db.UpdateImages(ctx, writer, item.ID, imagesJSON)
				}
			} else {
				if item.NeedImages && serr == nil {
					log.Printf("  [done] %s -> 0 images (status stays 0)", item.Keyword)
				}
				switch {
				case !item.NeedDesc:
					// images-only row with no results: nothing to write
				case len(cleaned) == 0:
					if serr == nil {
						log.Printf("  [skip] %s -> no data found", item.Keyword)
					}
				default:
					snippet := mergeSnippet(item.Existing, nil, cleaned)
					db.UpdateSnippet(ctx, writer, item.ID, snippet)
				}
			}

			total := done.Add(1)
			if total%50 == 0 {
				log.Printf("  [%d/%d] progress | %s", total, len(pending), env.Metrics.Snapshot())
			}
		}(item)
	}
	wg.Wait()
	writer.Close()
	db.Close(conn)
	log.Printf("Completed %s: %d keywords processed", filepath.Base(dbPath), done.Load())
	return nil
}

// ScrapeRelatedKw is phase 2 of the scrape: related keywords come from the
// suggest chain (DuckDuckGo first, Google fallback) and never from Pinterest.
// It runs after ScrapeImagesDesc instead of alongside it, so at most one
// request per worker is ever in flight.
func ScrapeRelatedKw(ctx context.Context, env *Env) error {
	files, err := db.ListDBFiles(env.Cfg.Path(env.Cfg.DataFolder))
	if err != nil {
		return err
	}
	for _, f := range files {
		if err := scrapeRelatedKwDB(ctx, env, f); err != nil {
			return err
		}
	}
	return nil
}

func scrapeRelatedKwDB(ctx context.Context, env *Env, dbPath string) error {
	conn, err := db.Open(dbPath)
	if err != nil {
		return err
	}
	pending, err := db.GetPendingRelated(conn)
	if err != nil {
		db.Close(conn)
		return err
	}
	if len(pending) == 0 {
		log.Printf("No keywords needing related keywords in %s", filepath.Base(dbPath))
		db.Close(conn)
		return nil
	}
	log.Printf("Found %d keywords needing related keywords in %s", len(pending), filepath.Base(dbPath))

	writer := db.NewWriter(conn, db.WriterOptions{
		BatchSize: env.Cfg.Go.WriterBatchSize,
	})

	sem := make(chan struct{}, env.Cfg.Concurrency)
	var wg sync.WaitGroup
	var done atomic.Int64

	for _, item := range pending {
		if ctx.Err() != nil {
			break
		}
		wg.Add(1)
		sem <- struct{}{}
		go func(item db.SnippetRow) {
			defer wg.Done()
			defer func() { <-sem }()
			if ctx.Err() != nil {
				return
			}

			related, rerr := env.Suggest.FetchRelatedKeywords(ctx, item.Keyword, env.Cfg.MaxRelatedKw)
			if rerr != nil {
				log.Printf("  suggest warn %s: %v", item.Keyword, rerr)
			}

			if len(related) == 0 {
				if rerr == nil {
					log.Printf("  [skip] %s -> no related keywords found", item.Keyword)
				}
			} else {
				snippet := mergeSnippet(item.Existing, related, nil)
				db.UpdateSnippet(ctx, writer, item.ID, snippet)
			}

			total := done.Add(1)
			if total%50 == 0 {
				log.Printf("  [%d/%d] progress | %s", total, len(pending), env.Metrics.Snapshot())
			}
		}(item)
	}
	wg.Wait()
	writer.Close()
	db.Close(conn)
	log.Printf("Completed %s: %d keywords processed", filepath.Base(dbPath), done.Load())
	return nil
}

// RunAI generates ai_title or ai_content for rows with images but without the
// column filled, using the batched single-writer.
func RunAI(ctx context.Context, env *Env, mode string) error {
	column := "ai_content"
	if mode == "title" {
		column = "ai_title"
	}
	promptTemplate, err := ai.LoadPrompt(env.Cfg.Root, mode)
	if err != nil {
		return fmt.Errorf("prompt not found: %w", err)
	}

	files, err := db.ListDBFiles(env.Cfg.Path(env.Cfg.DataFolder))
	if err != nil {
		return err
	}

	total := 0
	for _, f := range files {
		conn, err := db.Open(f)
		if err != nil {
			return err
		}
		rows, err := db.GetPendingAI(conn, column)
		db.Close(conn)
		if err != nil {
			return err
		}
		total += len(rows)
	}
	if total == 0 {
		log.Printf("No items to process for %s.", mode)
		return nil
	}
	log.Printf("Total items to process: %d", total)

	var counter int64
	for _, f := range files {
		if err := runAIDB(ctx, env, f, mode, column, promptTemplate, &counter, total); err != nil {
			return err
		}
	}
	return nil
}

func runAIDB(ctx context.Context, env *Env, dbPath, mode, column, promptTemplate string, counter *int64, total int) error {
	conn, err := db.Open(dbPath)
	if err != nil {
		return err
	}
	rows, err := db.GetPendingAI(conn, column)
	if err != nil {
		db.Close(conn)
		return err
	}
	if len(rows) == 0 {
		db.Close(conn)
		return nil
	}
	log.Printf("Processing %d items for %s in %s...", len(rows), mode, filepath.Base(dbPath))

	writer := db.NewWriter(conn, db.WriterOptions{
		BatchSize: env.Cfg.Go.AIWriterBatchSize,
	})

	sem := make(chan struct{}, env.Cfg.AI.Concurrency)
	var wg sync.WaitGroup

	for _, row := range rows {
		if ctx.Err() != nil {
			break
		}
		wg.Add(1)
		sem <- struct{}{}
		go func(row db.AIJobRow) {
			defer wg.Done()
			defer func() { <-sem }()
			if ctx.Err() != nil {
				return
			}
			prompt := strings.ReplaceAll(promptTemplate, "{keyword}", row.Keyword)
			content, err := env.AI.Generate(ctx, prompt)
			if err != nil {
				log.Printf("  FAIL %s <%s>: %v", row.Keyword, mode, err)
				return
			}
			if mode == "title" {
				content = ai.CleanTitle(content)
				db.UpdateAITitle(writer, row.ID, content)
			} else {
				db.UpdateAIContent(writer, row.ID, content)
			}
			cur := atomic.AddInt64(counter, 1)
			log.Printf("[%d/%d] Success: %s <%s>", cur, total, row.Keyword, mode)
		}(row)
	}
	wg.Wait()
	writer.Close()
	db.Close(conn)
	return nil
}
