<?php
declare(strict_types=1);

require_once __DIR__ . '/StayProcessingSelector.php';
require_once __DIR__ . '/StayProcessingTransitionService.php';
require_once __DIR__ . '/StayGeographyWriter.php';
require_once __DIR__ . '/AnchorageGeocoder.php';
require_once __DIR__ . '/SystemSettingsService.php';
require_once __DIR__ . '/../lib/DerivedDb.php';

/**
 * YachtMemory StayProcessingGeographyWorker - Phase 2A.6a
 *
 * Orchestrates due Geography jobs for persisted Stays.
 *
 * Geography belongs to the stable Stay identity (stay_id), not to a Stay
 * revision. A revision change of an OPEN Stay therefore does not invalidate
 * durable Geography and is not used as a stale-worker guard.
 *
 * Two acquisition modes are deliberately exposed:
 *
 * 1. runDue()
 *    The original Phase 2A worker API. Every selected job is processed with
 *    one explicitly supplied acquisition policy. This path remains available
 *    for isolated CACHE_ONLY and NETWORK_ALLOWED tests.
 *
 * 2. runDueCacheFirst()
 *    The bounded production-oriented acquisition policy introduced in
 *    Phase 2A.6a. Every selected job is claimed exactly once, evaluated from
 *    proven cache first, and only an INCOMPLETE result may consume one of the
 *    explicitly bounded network slots.
 *
 * A network slot is consumed before NETWORK_ALLOWED analysis starts. Thus a
 * timeout or other external failure still consumes the slot and cannot cause
 * unbounded pressure on the public geographic service.
 *
 * Processing consistency boundary:
 *
 *   claim PROCESSING
 *       -> acquire/evaluate geographic evidence
 *       -> build persistence plan
 *       -> BEGIN
 *            persist stay_geography
 *            mark geography COMPLETE
 *          COMMIT
 *
 * Potentially slow evidence acquisition happens outside the transaction.
 */
final class StayProcessingGeographyWorker
{
    /**
     * Process currently due Geography jobs using one acquisition policy.
     *
     * This is the existing Phase 2A API and deliberately remains compatible.
     * CACHE_ONLY stays the safe default.
     *
     * @return array{
     *   selected:int,
     *   complete:int,
     *   retry:int,
     *   claim_lost:int,
     *   results:array<int,array<string,mixed>>
     * }
     */
    public static function runDue(
        int $limit = 20,
        ?string $now = null,
        ?int $vesselId = null,
        ?int $retryMinutes = null,
        string $acquisitionPolicy = AnchorageGeocoder::ACQUISITION_CACHE_ONLY
    ): array {
        $now = self::normalizeDateTime($now);
        $retryMinutes = self::resolveRetryMinutes($retryMinutes);

        self::assertLimit($limit);
        self::assertAcquisitionPolicy($acquisitionPolicy);

        $jobs = StayProcessingSelector::selectDueGeography(
            $limit,
            $now,
            $vesselId
        );

        $summary = [
            'selected' => count($jobs),
            'complete' => 0,
            'retry' => 0,
            'claim_lost' => 0,
            'results' => [],
        ];

        /*
         * One geocoder instance is shared by this worker run so every selected
         * job sees the same effective settings and acquisition policy.
         */
        $geocoder = self::createGeocoder($acquisitionPolicy);

        foreach ($jobs as $job) {
            $result = self::processJobSinglePolicy(
                $job,
                $now,
                $retryMinutes,
                $geocoder
            );

            $summary['results'][] = $result;
            $summary[$result['result']]++;
        }

        return $summary;
    }

    /**
     * Process currently due Geography jobs cache-first with a bounded network
     * allowance.
     *
     * Each selected Stay is claimed exactly once.
     *
     * After the claim:
     * - CACHE_ONLY COMPLETE is persisted immediately.
     * - CACHE_ONLY INCOMPLETE may consume one network slot.
     * - NETWORK_ALLOWED COMPLETE is persisted immediately.
     * - Remaining incomplete/error cases transition to RETRY.
     *
     * $networkLimit counts attempted NETWORK_ALLOWED analyses, not successful
     * completions. A failed external request therefore still consumes a slot.
     *
     * Phase 2A.6a only introduces this worker capability. The regular
     * StayProcessingService is intentionally not switched to this method yet.
     *
     * @return array{
     *   selected:int,
     *   complete:int,
     *   retry:int,
     *   claim_lost:int,
     *   network_limit:int,
     *   network_attempted:int,
     *   network_completed:int,
     *   results:array<int,array<string,mixed>>
     * }
     */
    public static function runDueCacheFirst(
        int $limit = 20,
        ?string $now = null,
        ?int $vesselId = null,
        ?int $retryMinutes = null,
        int $networkLimit = 0
    ): array {
        $now = self::normalizeDateTime($now);
        $retryMinutes = self::resolveRetryMinutes($retryMinutes);

        self::assertLimit($limit);
        self::assertNetworkLimit($networkLimit);

        $jobs = StayProcessingSelector::selectDueGeography(
            $limit,
            $now,
            $vesselId
        );

        $summary = [
            'selected' => count($jobs),
            'complete' => 0,
            'retry' => 0,
            'claim_lost' => 0,
            'network_limit' => $networkLimit,
            'network_attempted' => 0,
            'network_completed' => 0,
            'results' => [],
        ];

        /*
         * Keep cache-only and network-enabled geocoders separate. Their
         * acquisition policies are immutable for the lifetime of the object,
         * which makes the safety boundary explicit.
         */
        $cacheGeocoder = self::createGeocoder(
            AnchorageGeocoder::ACQUISITION_CACHE_ONLY
        );

        $networkGeocoder = null;

        foreach ($jobs as $job) {
            $networkAvailable =
                $summary['network_attempted'] < $networkLimit;

            $result = self::processJobCacheFirst(
                $job,
                $now,
                $retryMinutes,
                $cacheGeocoder,
                $networkGeocoder,
                $networkAvailable
            );

            /*
             * A network slot is charged by processJobCacheFirst before the
             * NETWORK_ALLOWED analysis is executed. The returned diagnostic
             * flag lets the outer summary account for that attempt without
             * coupling the processing result to a successful completion.
             */
            if (($result['network_attempted'] ?? false) === true) {
                $summary['network_attempted']++;

                if (($result['result'] ?? '') === 'complete') {
                    $summary['network_completed']++;
                }
            }

            $summary['results'][] = $result;
            $summary[$result['result']]++;
        }

        return $summary;
    }

    /**
     * Process one selected Geography job with one fixed acquisition policy.
     *
     * This preserves the original runDue() semantics.
     *
     * @return array<string,mixed>
     */
    private static function processJobSinglePolicy(
        array $job,
        string $now,
        int $retryMinutes,
        AnchorageGeocoder $geocoder
    ): array {
        $identity = self::jobIdentity($job);

        $claimLost = self::claimJob(
            $identity['stay_id'],
            $identity['stay_number'],
            $now
        );

        if ($claimLost !== null) {
            return $claimLost;
        }

        try {
            $stay = self::loadStay($identity['stay_id']);
            $trace = $geocoder->analyze(self::buildGeocoderInput($stay));

            if (!self::isCompleteTrace($trace)) {
                return self::retryAfterClaim(
                    $identity['stay_id'],
                    $identity['stay_number'],
                    $now,
                    $retryMinutes,
                    self::incompleteTraceMessage($trace)
                );
            }

            return self::completeAfterClaim(
                $identity['stay_id'],
                $identity['stay_number'],
                $stay,
                $trace
            );
        } catch (Throwable $e) {
            return self::retryAfterClaim(
                $identity['stay_id'],
                $identity['stay_number'],
                $now,
                $retryMinutes,
                self::errorText($e)
            );
        }
    }

    /**
     * Process one selected Geography job using the bounded cache-first policy.
     *
     * The job remains PROCESSING between CACHE_ONLY and NETWORK_ALLOWED. There
     * is deliberately no intermediate RETRY transition, because that would
     * make the just-processed job ineligible for the second phase until the
     * retry interval had elapsed.
     *
     * @param ?AnchorageGeocoder $networkGeocoder
     *        Created lazily only when the first network slot is actually used.
     *
     * @return array<string,mixed>
     */
    private static function processJobCacheFirst(
        array $job,
        string $now,
        int $retryMinutes,
        AnchorageGeocoder $cacheGeocoder,
        ?AnchorageGeocoder &$networkGeocoder,
        bool $networkAvailable
    ): array {
        $identity = self::jobIdentity($job);

        $claimLost = self::claimJob(
            $identity['stay_id'],
            $identity['stay_number'],
            $now
        );

        if ($claimLost !== null) {
            $claimLost['network_attempted'] = false;
            return $claimLost;
        }

        try {
            $stay = self::loadStay($identity['stay_id']);
            $input = self::buildGeocoderInput($stay);

            /*
             * Phase A: proven cache only.
             */
            $cacheTrace = $cacheGeocoder->analyze($input);

            if (self::isCompleteTrace($cacheTrace)) {
                $result = self::completeAfterClaim(
                    $identity['stay_id'],
                    $identity['stay_number'],
                    $stay,
                    $cacheTrace
                );

                $result['acquisition_path'] = 'cache_only';
                $result['network_attempted'] = false;

                return $result;
            }

            /*
             * No network allowance remains. Only now do we publish RETRY.
             */
            if (!$networkAvailable) {
                $result = self::retryAfterClaim(
                    $identity['stay_id'],
                    $identity['stay_number'],
                    $now,
                    $retryMinutes,
                    self::incompleteTraceMessage($cacheTrace)
                );

                $result['acquisition_path'] = 'cache_only';
                $result['network_attempted'] = false;

                return $result;
            }

            /*
             * Phase B: consume exactly one network slot before the external
             * acquisition attempt starts.
             *
             * The outer summary observes network_attempted=true even when the
             * analysis subsequently fails or remains incomplete.
             */
            if ($networkGeocoder === null) {
                $networkGeocoder = self::createGeocoder(
                    AnchorageGeocoder::ACQUISITION_NETWORK_ALLOWED
                );
            }

            try {
                $networkTrace = $networkGeocoder->analyze($input);
            } catch (Throwable $e) {
                $result = self::retryAfterClaim(
                    $identity['stay_id'],
                    $identity['stay_number'],
                    $now,
                    $retryMinutes,
                    self::errorText($e)
                );

                $result['acquisition_path'] = 'cache_then_network';
                $result['network_attempted'] = true;

                return $result;
            }

            if (!self::isCompleteTrace($networkTrace)) {
                $result = self::retryAfterClaim(
                    $identity['stay_id'],
                    $identity['stay_number'],
                    $now,
                    $retryMinutes,
                    self::incompleteTraceMessage($networkTrace)
                );

                $result['acquisition_path'] = 'cache_then_network';
                $result['network_attempted'] = true;

                return $result;
            }

            $result = self::completeAfterClaim(
                $identity['stay_id'],
                $identity['stay_number'],
                $stay,
                $networkTrace
            );

            $result['acquisition_path'] = 'cache_then_network';
            $result['network_attempted'] = true;

            return $result;
        } catch (Throwable $e) {
            $result = self::retryAfterClaim(
                $identity['stay_id'],
                $identity['stay_number'],
                $now,
                $retryMinutes,
                self::errorText($e)
            );

            $result['acquisition_path'] = 'cache_first';
            $result['network_attempted'] = false;

            return $result;
        }
    }

    /**
     * Extract and validate the stable identifiers supplied by the selector.
     *
     * @return array{stay_id:int,vessel_id:int,stay_number:int}
     */
    private static function jobIdentity(array $job): array
    {
        $stayId = (int)($job['stay_id'] ?? 0);
        $vesselId = (int)($job['vessel_id'] ?? 0);
        $stayNumber = (int)($job['stay_number'] ?? 0);

        if ($stayId <= 0 || $vesselId <= 0 || $stayNumber <= 0) {
            throw new RuntimeException(
                'Selector returned an invalid Geography job.'
            );
        }

        return [
            'stay_id' => $stayId,
            'vessel_id' => $vesselId,
            'stay_number' => $stayNumber,
        ];
    }

    /**
     * Claim one selected job.
     *
     * Claim loss is an expected concurrency result, not a processing failure.
     *
     * @return ?array<string,mixed>
     */
    private static function claimJob(
        int $stayId,
        int $stayNumber,
        string $now
    ): ?array {
        try {
            StayProcessingTransitionService::beginGeographyAttempt(
                $stayId,
                $now
            );

            return null;
        } catch (RuntimeException $e) {
            return [
                'stay_id' => $stayId,
                'stay_number' => $stayNumber,
                'result' => 'claim_lost',
                'message' => $e->getMessage(),
            ];
        }
    }

    /**
     * Adapt a persisted Stay to the historical V4 geocoder input vocabulary.
     *
     * The word "anchorage" remains confined to this adapter. Durable Geography
     * itself is keyed solely by stay_id.
     *
     * @return array<string,mixed>
     */
    private static function buildGeocoderInput(array $stay): array
    {
        return [
            'vessel_id' => (int)$stay['vessel_id'],
            'anchorage_number' => (int)$stay['stay_number'],
            'anchorage_name' => null,
            'country' => null,
            'region' => null,
            'geocode_status' => null,
            'start_time' => (string)$stay['start_time'],
            'end_time' => (string)$stay['end_time'],
            'lat' => (float)$stay['center_lat'],
            'lon' => (float)$stay['center_lon'],
        ];
    }

    /**
     * A Geography result is durable only when the decision engine reports a
     * complete evidence set. COMPLETE with selected_name=NULL remains valid.
     */
    private static function isCompleteTrace(array $trace): bool
    {
        return strtoupper(trim(
            (string)($trace['decision_completeness'] ?? '')
        )) === 'COMPLETE';
    }

    /**
     * Persist one final COMPLETE Geography result and return its diagnostics.
     *
     * @return array<string,mixed>
     */
    private static function completeAfterClaim(
        int $stayId,
        int $stayNumber,
        array $stay,
        array $trace
    ): array {
        $plan = StayGeographyWriter::plan(
            $stayId,
            (float)$stay['center_lat'],
            (float)$stay['center_lon'],
            $trace,
            AnchorageGeocoder::VERSION
        );

        self::commitComplete($stayId, $plan);

        return [
            'stay_id' => $stayId,
            'stay_number' => $stayNumber,
            'result' => 'complete',
            'writer_action' => (string)($plan['action'] ?? ''),
            'selected_name' => $trace['selected_name'] ?? null,
            'selected_source' => $trace['selected_source'] ?? null,
            'decision_rule' => $trace['decision_rule'] ?? null,
            'confidence' => $trace['confidence'] ?? null,
            'geocoder_version' => AnchorageGeocoder::VERSION,
        ];
    }

    /**
     * Create a geocoder whose acquisition policy cannot change during its
     * lifetime.
     */
    private static function createGeocoder(
        string $acquisitionPolicy
    ): AnchorageGeocoder {
        self::assertAcquisitionPolicy($acquisitionPolicy);

        return new AnchorageGeocoder(
            SystemSettingsService::getGeocodingConfig(),
            null,
            $acquisitionPolicy
        );
    }

    /**
     * Atomically publish final Geography knowledge.
     *
     * There must never be a committed stay_geography row while the processing
     * state still says PROCESSING, nor COMPLETE without the durable Geography
     * row that justifies that state.
     */
    private static function commitComplete(int $stayId, array $plan): void
    {
        $pdo = DerivedDb::getPdo();
        $pdo->beginTransaction();

        try {
            StayGeographyWriter::apply($plan, $pdo);

            StayProcessingTransitionService::markGeographyCompleteUsingPdo(
                $pdo,
                $stayId
            );

            $pdo->commit();
        } catch (Throwable $e) {
            if ($pdo->inTransaction()) {
                $pdo->rollBack();
            }

            throw $e;
        }
    }

    /**
     * Move a claimed job to RETRY.
     *
     * RETRY remains the temporary-failure state. No permanent-failure policy
     * is invented here.
     *
     * @return array<string,mixed>
     */
    private static function retryAfterClaim(
        int $stayId,
        int $stayNumber,
        string $now,
        int $retryMinutes,
        string $message
    ): array {
        $next = (new DateTimeImmutable($now))
            ->modify('+' . $retryMinutes . ' minutes')
            ->format('Y-m-d H:i:s');

        try {
            StayProcessingTransitionService::markGeographyRetry(
                $stayId,
                $next,
                self::limitErrorText($message)
            );
        } catch (Throwable $transitionError) {
            throw new RuntimeException(
                'Geography failed and RETRY transition also failed for Stay #' .
                $stayNumber .
                ': original=' . $message .
                '; transition=' . $transitionError->getMessage(),
                0,
                $transitionError
            );
        }

        return [
            'stay_id' => $stayId,
            'stay_number' => $stayNumber,
            'result' => 'retry',
            'next_attempt_at' => $next,
            'message' => $message,
        ];
    }

    /**
     * Load the current Stay identity and reference position.
     *
     * revision is returned only for diagnostics. It is deliberately not used
     * as a Geography input version or stale-worker guard.
     *
     * @return array<string,mixed>
     */
    private static function loadStay(int $stayId): array
    {
        $pdo = DerivedDb::getPdo();

        $stmt = $pdo->prepare(
            'SELECT
                id,
                vessel_id,
                stay_number,
                start_time,
                end_time,
                center_lat,
                center_lon,
                revision,
                event_state
             FROM stays
             WHERE id = ?'
        );

        $stmt->execute([$stayId]);
        $row = $stmt->fetch(PDO::FETCH_ASSOC);

        if ($row === false) {
            throw new RuntimeException('Stay not found: ' . $stayId);
        }

        return $row;
    }

    /**
     * Explain why a V4 trace is not yet durable.
     *
     * The complete Decision Trace deliberately remains outside
     * stay_processing.
     */
    private static function incompleteTraceMessage(array $trace): string
    {
        $parts = ['Geographic decision incomplete'];

        $sources = $trace['evidence_sources'] ?? null;

        if (is_array($sources)) {
            foreach (['settlement', 'maritime'] as $sourceName) {
                $source = $sources[$sourceName] ?? null;

                if (!is_array($source)) {
                    continue;
                }

                $status = trim((string)($source['status'] ?? ''));
                $error = trim((string)($source['error'] ?? ''));

                if ($status !== '') {
                    $part = $sourceName . '=' . $status;

                    if ($error !== '') {
                        $part .= '(' . $error . ')';
                    }

                    $parts[] = $part;
                }
            }
        }

        return self::limitErrorText(implode('; ', $parts));
    }

    private static function resolveRetryMinutes(?int $retryMinutes): int
    {
        if ($retryMinutes === null) {
            $config = SystemSettingsService::getStayProcessingConfig();
            $retryMinutes = (int)($config['retry_interval_min'] ?? 0);
        }

        if ($retryMinutes < 1) {
            throw new InvalidArgumentException(
                'retryMinutes must be greater than zero.'
            );
        }

        return $retryMinutes;
    }

    private static function assertLimit(int $limit): void
    {
        if ($limit < 1) {
            throw new InvalidArgumentException(
                'limit must be greater than zero.'
            );
        }
    }

    private static function assertNetworkLimit(int $networkLimit): void
    {
        if ($networkLimit < 0) {
            throw new InvalidArgumentException(
                'networkLimit must not be negative.'
            );
        }
    }

    private static function assertAcquisitionPolicy(string $policy): void
    {
        if (!in_array(
            $policy,
            [
                AnchorageGeocoder::ACQUISITION_CACHE_ONLY,
                AnchorageGeocoder::ACQUISITION_NETWORK_ALLOWED,
            ],
            true
        )) {
            throw new InvalidArgumentException(
                'Unknown Geography acquisition policy: ' . $policy
            );
        }
    }

    private static function errorText(Throwable $e): string
    {
        return self::limitErrorText(
            get_class($e) . ': ' . $e->getMessage()
        );
    }

    private static function limitErrorText(string $text): string
    {
        return mb_substr($text, 0, 4000);
    }

    private static function normalizeDateTime(?string $value): string
    {
        if ($value === null) {
            return date('Y-m-d H:i:s');
        }

        $dt = DateTimeImmutable::createFromFormat(
            '!Y-m-d H:i:s',
            $value
        );

        $errors = DateTimeImmutable::getLastErrors();

        if (
            $dt === false
            || (
                $errors !== false
                && (
                    $errors['warning_count'] > 0
                    || $errors['error_count'] > 0
                )
            )
            || $dt->format('Y-m-d H:i:s') !== $value
        ) {
            throw new InvalidArgumentException(
                'Invalid datetime: ' . $value
            );
        }

        return $value;
    }
}
