package services

import (
	"context"
	"fmt"
	"strings"

	"github.com/cloudwego/eino/components/model"
	"github.com/cloudwego/eino/components/tool"
	"github.com/cloudwego/eino/schema"
)

const maxUncheckedAnswerRetries = 2

// Keep the native Eino tool loop, but make completion an application decision.
// A prompt alone cannot guarantee that a model calls check_answer before stopping.
func executeAnswerAgent(ctx context.Context, cm model.ToolCallingChatModel, tools []tool.BaseTool, messages []*schema.Message, state *agentTurnState) (*schema.Message, error) {
	return executeEino(ctx, &checkedAnswerModel{inner: cm, state: state}, tools, messages)
}

type checkedAnswerModel struct {
	inner model.ToolCallingChatModel
	state *agentTurnState
	tools []*schema.ToolInfo
}

func (m *checkedAnswerModel) WithTools(tools []*schema.ToolInfo) (model.ToolCallingChatModel, error) {
	inner, err := m.inner.WithTools(tools)
	if err != nil {
		return nil, err
	}
	return &checkedAnswerModel{inner: inner, state: m.state, tools: tools}, nil
}

func (m *checkedAnswerModel) Generate(ctx context.Context, in []*schema.Message, opts ...model.Option) (*schema.Message, error) {
	for {
		if err := ctx.Err(); err != nil {
			return nil, err
		}
		m.state.mu.Lock()
		handoff := m.state.handoffReason != ""
		approved := m.state.approved
		retries := m.state.answerRecoveries
		remaining := max(0, 6-m.state.calls)
		checks := m.state.answerChecks
		sources := joinUintIDs(knowledgeIDs(m.state.knowledge))
		m.state.mu.Unlock()
		if handoff {
			return schema.AssistantMessage("[[ESCALATE]]", nil), nil
		}
		if approved != nil {
			// The final source/version checks still run in chatAgentic. Do not ask
			// the model to rewrite an already checked answer, possibly changing it.
			RecordSimulationEvent(ctx, "answer_ready", "", "completed", nil)
			return schema.AssistantMessage(approved.input.Answer, nil), nil
		}

		messages := in
		activeModel := m.inner
		if remaining == 0 && len(m.tools) > 0 {
			// Stop offering exhausted data tools, while retaining the native
			// transcript and tools that can safely finish the turn.
			completion := make([]*schema.ToolInfo, 0, 2)
			for _, info := range m.tools {
				if info.Name == "request_handoff" || info.Name == "check_answer" && checks < 6 {
					completion = append(completion, info)
				}
			}
			var err error
			activeModel, err = m.inner.WithTools(completion)
			if err != nil {
				return nil, err
			}
			messages = appendAgentInstruction(messages, "\nSELESAIKAN DARI BUKTI TERSEDIA: Jatah tool pencarian/data habis. Jangan mencari lagi. Susun jawaban hanya untuk pertanyaan pelanggan dari sumber yang sudah dibaca dan panggil check_answer. Jika fakta belum cukup, satu klarifikasi tanpa klaim atau request_handoff tetap tersedia. Jangan mengisi kekurangan dari persona.")
			RecordSimulationEvent(ctx, "answer_synthesis", "", "running", map[string]any{"remaining_data_calls": remaining})
		} else if sources != "" {
			messages = appendAgentInstruction(messages, fmt.Sprintf("\nSTATUS GILIRAN: source_id yang sudah dibaca: %s. Sisa tool data: %d. Jika sumber tersebut menjawab pertanyaan, segera check_answer. Cari lagi hanya untuk fakta yang dibutuhkan pelanggan dan belum ditemukan, bukan untuk membuktikan contoh di persona atau menambah rincian yang tidak diminta.", sources, remaining))
		}
		if retries > 0 {
			messages = remindAnswerChecks(messages)
		}
		// DeepSeek thinking rejects required/named tool_choice (HTTP 400).
		// Let the model choose tools natively; enforce completion here instead.
		// The independent JSON evidence reviewer does not use this wrapper.
		callOpts := append(append([]model.Option(nil), opts...), model.WithToolChoice(schema.ToolChoiceAllowed))
		out, err := activeModel.Generate(ctx, messages, callOpts...)
		if err != nil {
			return nil, err
		}
		if out != nil && len(out.ToolCalls) > 0 {
			return out, nil
		}

		// A model can emit plain text before completing the required checks.
		// Discard that unchecked draft and retry the same tool transcript.
		// Never turn it into a customer-visible answer or fabricate tool results.
		m.state.mu.Lock()
		if m.state.lastAnswerIssue == "" {
			m.state.lastAnswerIssue = "answer_check_missing"
		}
		issue := m.state.lastAnswerIssue
		exhausted := m.state.answerRecoveries >= maxUncheckedAnswerRetries
		if exhausted {
			m.state.handoffReason = "answer_validation_failed"
		} else {
			m.state.answerRecoveries++
		}
		m.state.mu.Unlock()
		if exhausted {
			RecordSimulationEvent(ctx, "answer_recovery_stopped", "", "error", map[string]any{"issue": issue})
			return schema.AssistantMessage("[[ESCALATE]]", nil), nil
		}
		RecordSimulationEvent(ctx, "answer_recovery", "", "running", map[string]any{"issue": issue})
	}
}

func remindAnswerChecks(in []*schema.Message) []*schema.Message {
	const reminder = "\nLANJUTKAN PEMERIKSAAN: Respons teks sebelumnya belum boleh dikirim. Pilih dan panggil tool yang diperlukan. Untuk pertanyaan fakta bisnis, cari pengetahuan atau produk terlebih dahulu; jika hasil belum cocok, perbaiki kata kunci. Setelah bukti cukup, panggil check_answer dengan teks jawaban dan sumbernya. Untuk sapaan gunakan check_answer kind=social. Gunakan request_handoff hanya bila memang perlu CS atau sumber yang telah dicari tetap belum cukup. Jangan berhenti dengan teks tanpa pemeriksaan."
	return appendAgentInstruction(in, reminder)
}

func appendAgentInstruction(in []*schema.Message, instruction string) []*schema.Message {
	messages := append([]*schema.Message(nil), in...)
	if len(messages) > 0 && messages[0].Role == schema.System {
		first := *messages[0]
		first.Content = strings.TrimSpace(first.Content) + instruction
		messages[0] = &first
		return messages
	}
	return append([]*schema.Message{schema.SystemMessage(strings.TrimSpace(instruction))}, messages...)
}

func (m *checkedAnswerModel) Stream(ctx context.Context, in []*schema.Message, opts ...model.Option) (*schema.StreamReader[*schema.Message], error) {
	// Never stream unchecked draft tokens to a customer.
	out, err := m.Generate(ctx, in, opts...)
	if err != nil {
		return nil, err
	}
	return schema.StreamReaderFromArray([]*schema.Message{out}), nil
}
