// Package simulation evaluates immutable, bounded alternatives through the
// ordinary Function and Action planning contracts. It owns no executor.
package simulation

import (
	"context"
	"encoding/json"
	"errors"
	"fmt"
	"sort"
	"time"

	"github.com/rafflesia-ai/bijection/internal/blob"

	"github.com/jackc/pgx/v5"

	"github.com/rafflesia-ai/bijection/internal/accesslog"
	"github.com/rafflesia-ai/bijection/internal/actions"
	"github.com/rafflesia-ai/bijection/internal/assurance"
	"github.com/rafflesia-ai/bijection/internal/computecall"
	"github.com/rafflesia-ai/bijection/internal/datasetversion"
	"github.com/rafflesia-ai/bijection/internal/db"
	"github.com/rafflesia-ai/bijection/internal/functionasync"
	"github.com/rafflesia-ai/bijection/internal/jsonvalue"
	"github.com/rafflesia-ai/bijection/internal/ledger"
	"github.com/rafflesia-ai/bijection/internal/obligation"
	"github.com/rafflesia-ai/bijection/internal/ontology"
	"github.com/rafflesia-ai/bijection/internal/security"
	"github.com/rafflesia-ai/bijection/internal/source"
	"github.com/rafflesia-ai/bijection/pkg/api"
)

const MaxEvaluations = 128
const MaxDocumentBytes = 4 << 20
const MaxPolicyFirings = 512

var ErrNotFound = errors.New("simulation not found")
var ErrForbidden = errors.New("function permission is required")
var ErrStale = errors.New("simulation authority, product, or retained evidence changed; create a new baseline")

type Service struct {
	Store    *blob.Store
	Compute  *computecall.Executor
	World    *db.World
	Registry *ontology.Registry
	// Runs executes trajectories larger than the request-local envelope as
	// ordinary durable Function Runs. Without it an oversized simulation is
	// refused rather than silently truncated.
	Runs *functionasync.Service
}

type baseline struct {
	View           db.SimulationView `json:"view"`
	ErasureOffset  int64             `json:"erasure_offset"`
	AuthorityHash  string            `json:"authority_hash"`
	DeploymentHash string            `json:"deployment_hash"`
	// AlternativeViews are the exact admitted private worlds of a durable
	// simulation, indexed by alternative with the baseline first. A partition
	// executes in the world its admission fixed, never in a re-derived one, and
	// they stay in the owner-private experiment rather than in the ledger.
	AlternativeViews []db.SimulationView `json:"alternative_views,omitempty"`
}

// CreateRetained evaluates and persists a simulation while retaining the
// governed physical row sources for its trajectory outputs. The caller owns
// the result and must close it after streaming or materializing the rows.
func (s Service) CreateRetained(ctx context.Context, authority security.Authority, request api.SimulationDefinition) (_ *RetainedSimulation, returnErr error) {
	ctx, cancel := context.WithTimeout(ctx, 2*time.Minute)
	defer cancel()
	if err := validateRequest(request); err != nil {
		return nil, err
	}
	if !authority.CanRunFunctions() {
		return nil, ErrForbidden
	}
	function, err := s.Registry.Function(request.Function)
	if err != nil {
		return nil, err
	}
	if function.Determinism != "deterministic" || len(function.ModelBindings) != 0 {
		return nil, ontology.BadRequest("simulations require deterministic Functions")
	}
	if request.Horizon != nil && function.Trajectory == nil {
		return nil, ontology.BadRequest("simulation horizon requires a Function with an authored trajectory contract")
	}
	if err := validateControlledPolicies(s.Registry, function, request); err != nil {
		return nil, err
	}
	ctx = s.World.FixBasis(ctx)
	checks, datasets, err := s.simulationQualityScope(ctx, authority, request, function)
	if err != nil {
		return nil, err
	}
	var basis baseline
	if request.ParentID != "" {
		var parent *api.Simulation
		parent, basis, err = s.load(ctx, authority, request.ParentID)
		if err == nil && parent.Function != request.Function {
			return nil, ontology.BadRequest("a fork must use the parent simulation's Function")
		}
	} else {
		basis, err = s.capture(ctx, authority, datasets)
	}
	if err != nil {
		return nil, err
	}
	retained, err := s.evaluateRetained(ctx, authority, request, basis, checks)
	// The computed size of this experiment, not a caller-selected mode, decides
	// the execution path. Interactive evaluation is the measurement.
	var oversized *durableRequired
	if errors.As(err, &oversized) {
		return s.createDurable(ctx, authority, request, basis, function, oversized)
	}
	if err != nil {
		return nil, err
	}
	defer func() {
		if returnErr != nil {
			_ = retained.Close()
		}
	}()
	result := retained.Result
	result.ID = api.NewID()
	result.CreatedAt = time.Now().UTC()
	result.Execution = "interactive"
	if err := s.save(ctx, authority, result, basis); err != nil {
		return nil, err
	}
	return retained, nil
}

func (s Service) simulationQualityScope(ctx context.Context, authority security.Authority, request api.SimulationDefinition, function *ontology.FunctionSpec) ([]assurance.DecisionCheck, []string, error) {
	compileParams := cloneParams(request.Params)
	if request.Horizon != nil {
		seed := uint64(0)
		if request.Sampling != nil {
			seed = request.Sampling.Seed
		}
		compileParams[function.Trajectory.StepsParam] = request.Horizon.Steps
		compileParams[function.Trajectory.StepHoursParam] = request.Horizon.StepHours
		compileParams[function.Trajectory.SeedParam] = int64(seed & 0xffff_ffff)
	}
	compiled, err := s.Registry.CompileFunction(ctx, request.Function, compileParams, authority.ModelContext())
	if err != nil {
		return nil, nil, err
	}
	datasets, err := s.sourceDatasets(ctx, authority, request, compiled)
	if err != nil {
		return nil, nil, err
	}
	return assurance.DecisionChecks(ctx, s.Registry, authority, datasets)
}

func (s Service) sourceDatasets(ctx context.Context, authority security.Authority, request api.SimulationDefinition, compiled *ontology.CompiledFunction) ([]string, error) {
	set := map[string]bool{}
	for _, dataset := range s.Registry.FunctionPlanSourceDatasets(compiled) {
		set[dataset] = true
	}
	for _, alternative := range request.Alternatives {
		actionsToCheck := append([]api.SimulationAction(nil), alternative.Actions...)
		for _, policy := range alternative.Policies {
			actionsToCheck = append(actionsToCheck, policy.Action)
		}
		for _, proposed := range actionsToCheck {
			spec, err := s.Registry.Action(proposed.Action)
			if err != nil {
				return nil, err
			}
			if !authority.CanPlanAction(spec) {
				return nil, invalid("Action %s is unavailable to this principal", proposed.Action)
			}
			datasets, err := actions.EvaluationSourceDatasets(ctx, s.World, s.Registry, spec,
				proposed.Request.TargetKey, proposed.Request.Input, "", authority)
			if err != nil {
				return nil, err
			}
			for _, dataset := range datasets {
				set[dataset] = true
			}
		}
	}
	datasets := make([]string, 0, len(set))
	for dataset := range set {
		datasets = append(datasets, dataset)
	}
	sort.Strings(datasets)
	return datasets, nil
}

func (s Service) capture(ctx context.Context, authority security.Authority, datasets []string) (baseline, error) {
	var result baseline
	err := s.World.ReadSources(ctx, source.Datasets(), func(tx pgx.Tx) error {
		basis, err := ledger.Capture(ctx, tx, s.Registry.ModelHash)
		if err != nil {
			return err
		}
		result.View.LedgerOffset = basis.LedgerOffset
		result.View.EvaluatedAt, err = time.Parse(time.RFC3339Nano, basis.AsOf)
		if err != nil {
			return err
		}
		// PostgreSQL timestamptz preserves microseconds. Canonicalize the retained
		// basis before it reaches either execution terminal so step-zero evidence
		// carries the same exact public instant.
		result.View.EvaluatedAt = result.View.EvaluatedAt.UTC().Truncate(time.Microsecond)
		result.AuthorityHash = authority.ContextFingerprint()
		result.DeploymentHash = basis.DeploymentHash
		if err := tx.QueryRow(ctx, `SELECT pg_current_snapshot()::text,
		 coalesce((SELECT max(id) FROM ops.operational_event WHERE kind IN ('SubjectRedacted','FileVersionErased')),0)`).Scan(&result.View.Snapshot, &result.ErasureOffset); err != nil {
			return err
		}
		if len(datasets) > 0 {
			transaction, err := datasetversion.ResolveLatestIn(ctx, tx, s.Registry.ModelHash, datasets)
			if err != nil {
				return err
			}
			result.View.DatasetTransaction = transaction.ID
		}
		return nil
	})
	if err != nil {
		return baseline{}, err
	}
	result.View.Fingerprint, err = fingerprint(result)
	return result, err
}

func (s Service) checkErasure(ctx context.Context, q ledger.Reader, basis baseline) error {
	var offset int64
	if err := q.QueryRow(ctx, `SELECT coalesce(max(id),0) FROM ops.operational_event WHERE kind IN ('SubjectRedacted','FileVersionErased')`).Scan(&offset); err != nil {
		return err
	}
	if offset != basis.ErasureOffset {
		return ErrStale
	}
	return nil
}

func (s Service) loadDocument(ctx context.Context, authority security.Authority, id string) (*api.Simulation, baseline, error) {
	var result api.Simulation
	var basis baseline
	var encoded, encodedBasis []byte
	err := s.World.Pool.QueryRow(ctx, `SELECT document, baseline FROM platform.simulation WHERE id=$1 AND owner_subject=$2`, id, authority.Subject()).Scan(&encoded, &encodedBasis)
	if errors.Is(err, pgx.ErrNoRows) {
		return nil, basis, ErrNotFound
	}
	if err != nil {
		return nil, basis, err
	}
	if err := json.Unmarshal(encoded, &result); err != nil {
		return nil, basis, err
	}
	if err := json.Unmarshal(encodedBasis, &basis); err != nil {
		return nil, basis, err
	}
	if !authority.CanRunFunctions() || !authority.MatchesContextFingerprint(basis.AuthorityHash) || result.Basis.ModelHash != s.Registry.ModelHash || basis.DeploymentHash != s.World.DeploymentHash() {
		return nil, basis, ErrStale
	}
	if err := s.checkErasure(ctx, s.World.Pool, basis); err != nil {
		return nil, basis, err
	}
	if err := s.Registry.CheckMandatoryProtection(authority.ModelContext(), result.MandatoryProtection); err != nil {
		return nil, basis, err
	}
	return &result, basis, nil
}

// load additionally proves the reader may still plan every proposed Action.
// Reading a saved experiment offers those plans; executing one partition does
// not, so the durable partition doorway uses loadDocument above.
func (s Service) load(ctx context.Context, authority security.Authority, id string) (*api.Simulation, baseline, error) {
	result, basis, err := s.loadDocument(ctx, authority, id)
	if err != nil {
		return nil, basis, err
	}
	for _, alternative := range result.Request.Alternatives {
		for _, proposed := range alternative.Actions {
			spec, err := s.Registry.Action(proposed.Action)
			if err != nil || !authority.CanPlanAction(spec) {
				return nil, basis, ErrStale
			}
		}
	}
	return result, basis, nil
}

// GetRetained replays a saved simulation and keeps its governed trajectory
// sources alive through disclosure. The caller must close the result.
func (s Service) GetRetained(ctx context.Context, authority security.Authority, id string) (_ *RetainedSimulation, returnErr error) {
	ctx, cancel := context.WithTimeout(ctx, 2*time.Minute)
	defer cancel()
	result, basis, err := s.load(ctx, authority, id)
	if err != nil {
		return nil, err
	}
	if result.Execution == "durable" {
		return s.getDurable(ctx, authority, result, basis)
	}
	function, err := s.Registry.Function(result.Request.Function)
	if err != nil {
		return nil, err
	}
	checks, _, err := s.simulationQualityScope(ctx, authority, result.Request, function)
	if err != nil {
		return nil, err
	}
	current, err := s.evaluateRetained(ctx, authority, result.Request, basis, checks)
	if err != nil {
		return nil, err
	}
	defer func() {
		if returnErr != nil {
			_ = current.Close()
		}
	}()
	if current.Result.ResultFingerprint != result.ResultFingerprint {
		return nil, ErrStale
	}
	current.Result.SimulationSummary = result.SimulationSummary
	result = current.Result
	err = accesslog.RecordAssetRead(ctx, s.World, authority, accesslog.AssetRead{Kind: accesslog.AssetFunction,
		Asset: result.Function, Mode: "simulation", Outcome: "served", Count: len(result.Runs), MandatoryProtection: result.MandatoryProtection})
	return current, err
}

func (s Service) List(ctx context.Context, authority security.Authority) (api.List[api.SimulationSummary], error) {
	out := api.List[api.SimulationSummary]{Data: []api.SimulationSummary{}}
	if !authority.CanRunFunctions() {
		return out, ErrForbidden
	}
	rows, err := s.World.Pool.Query(ctx, `SELECT id, document->>'name', function_name, created_at FROM platform.simulation
	 WHERE owner_subject=$1 AND model_hash=$2 AND authority_hash=$3 ORDER BY created_at DESC, id LIMIT 100`,
		authority.Subject(), s.Registry.ModelHash, authority.ContextFingerprint())
	if err != nil {
		return out, err
	}
	defer rows.Close()
	for rows.Next() {
		var item api.SimulationSummary
		if err := rows.Scan(&item.ID, &item.Name, &item.Function, &item.CreatedAt); err != nil {
			return out, err
		}
		if s.Registry.FunctionVisibility(item.Function, authority.ModelContext()).IsVisible {
			out.Data = append(out.Data, item)
		}
	}
	return out, rows.Err()
}

func (s Service) Delete(ctx context.Context, authority security.Authority, id string) error {
	// Owner deletion remains possible after a product/authority change; no
	// retained values are returned through this lifecycle operation.
	return db.Mutate(ctx, s.World.Pool, pgx.TxOptions{}, func(tx pgx.Tx) error {
		var result api.Simulation
		var body []byte
		if err := tx.QueryRow(ctx, `SELECT document FROM platform.simulation WHERE id=$1 AND owner_subject=$2 FOR UPDATE`, id, authority.Subject()).Scan(&body); errors.Is(err, pgx.ErrNoRows) {
			return ErrNotFound
		} else if err != nil {
			return err
		}
		if err := json.Unmarshal(body, &result); err != nil {
			return err
		}
		if err := obligation.CheckErasure(ctx, tx, s.Registry, "Simulation", result.ID, result.MandatoryProtection, result.CreatedAt); err != nil {
			return ontology.BadRequest("simulation cannot be deleted: " + err.Error())
		}
		if _, err := tx.Exec(ctx, `DELETE FROM platform.simulation WHERE id=$1 AND owner_subject=$2`, id, authority.Subject()); err != nil {
			return err
		}
		_, err := ledger.Emit(ctx, tx, simulationEvent(authority, ledger.KindSimulationDeleted, &result))
		return err
	})
}

func (s Service) PlansFor(ctx context.Context, authority security.Authority, result *api.Simulation, run int) (api.SimulationPlans, error) {
	if run < 0 || run >= len(result.Runs) {
		return api.SimulationPlans{}, ontology.BadRequest("unknown simulation run")
	}
	out := api.SimulationPlans{Object: "atelier.simulation_plans", Plans: []api.Plan{}}
	for _, proposed := range result.Runs[run].ProposedActions {
		plan, err := actions.Plan(ctx, s.World, s.Registry, proposed.Action, proposed.Request, authority)
		if err != nil {
			return out, err
		}
		out.Plans = append(out.Plans, plan)
	}
	return out, nil
}

func simulationEvent(authority security.Authority, kind ledger.Kind, result *api.Simulation) ledger.Event {
	return ledger.Event{Evidence: authority.Evidence(), Kind: kind, SubjectType: "Simulation", SubjectKey: result.ID,
		Payload: map[string]any{"function": result.Function, "result_fingerprint": result.ResultFingerprint}}
}

func fingerprint(value any) (string, error) {
	valueHash := jsonvalue.Fingerprint(value)
	if valueHash == "" {
		return "", fmt.Errorf("simulation evidence is not finite JSON")
	}
	return valueHash, nil
}

func invalid(format string, args ...any) error {
	return ontology.BadRequest(fmt.Sprintf(format, args...))
}
