<?php

require_once __DIR__ . '/KafkaConfig.php';

/**
 * Consumer
 *
 * Wrapper around the high-level \RdKafka\KafkaConsumer.
 *
 * The high-level consumer is keyed by a group.id. This is what makes the
 * "one message, many consumers, different actions" requirement work:
 *
 *   - All consumers that share the SAME group.id form one group and the
 *     partitions (and therefore messages) are split between them — this is
 *     how you scale one logical job horizontally.
 *   - Every DISTINCT group.id gets its OWN independent copy of every
 *     message and its own committed offsets.
 *
 * So if you push one message to topic "orders" and run:
 *     Consumer(group "email-service")      -> sends an email
 *     Consumer(group "analytics-service")  -> records a metric
 *     Consumer(group "audit-service")      -> writes an audit log
 * each group receives that same message and performs a different action.
 *
 * Usage:
 *     $c = new Consumer($config, 'email-service', ['orders']);
 *     $c->consume(function (array $msg) {
 *         // $msg = ['topic','partition','offset','key','payload','headers']
 *         sendEmail($msg['payload']);
 *     });
 */
class Consumer
{
    private \RdKafka\KafkaConsumer $consumer;

    /** @var string[] topics this consumer is subscribed to */
    private array $topics;

    private string $groupId;

    /** @var bool flag used to stop the consume() loop gracefully */
    private bool $running = false;

    /**
     * @param KafkaConfig $config
     * @param string      $groupId  consumer group id (defines the "who" —
     *                              use a different id per independent action)
     * @param string[]    $topics   topics to subscribe to
     * @param string      $offsetReset 'earliest' (from start) or 'latest'
     */
    public function __construct(
        KafkaConfig $config,
        string $groupId,
        array $topics,
        string $offsetReset = 'earliest'
    ) {
        if ($groupId === '') {
            throw new \InvalidArgumentException('Consumer group id cannot be empty.');
        }
        if ($topics === []) {
            throw new \InvalidArgumentException('At least one topic is required.');
        }

        $conf = $config->baseConf();
        $conf->set('group.id', $groupId);
        // Where a brand-new group starts reading from.
        $conf->set('auto.offset.reset', $offsetReset);
        // Commit offsets manually only after the handler succeeds, so a
        // crash mid-processing re-delivers the message instead of losing it.
        $conf->set('enable.auto.commit', 'false');

        $this->consumer = new \RdKafka\KafkaConsumer($conf);
        $this->groupId  = $groupId;
        $this->topics   = $topics;

        $this->consumer->subscribe($topics);
    }

    /**
     * Read the next single message, hand it to $handler, then commit.
     * Returns true if a message was processed, false on timeout/no message.
     *
     * @param callable $handler   fn(array $message): void
     * @param int      $timeoutMs how long to wait for a message
     */
    public function consumeOnce(callable $handler, int $timeoutMs = 10000): bool
    {
        $message = $this->consumer->consume($timeoutMs);

        switch ($message->err) {
            case RD_KAFKA_RESP_ERR_NO_ERROR:
                $handler($this->normalise($message));
                // Only commit after the handler returns without throwing.
                $this->consumer->commit($message);
                return true;

            case RD_KAFKA_RESP_ERR__PARTITION_EOF:
            case RD_KAFKA_RESP_ERR__TIMED_OUT:
                // No message available right now — normal, not an error.
                return false;

            default:
                throw new \RuntimeException(
                    'Kafka consume error: ' . $message->errstr() . ' (' . $message->err . ')'
                );
        }
    }

    /**
     * Continuously consume messages, invoking $handler for each one, until
     * stop() is called or (optionally) $maxMessages have been processed.
     *
     * @param callable $handler     fn(array $message): void
     * @param int|null $maxMessages stop after this many (null = forever)
     * @param int      $timeoutMs   per-poll wait
     *
     * @return int number of messages processed
     */
    public function consume(callable $handler, ?int $maxMessages = null, int $timeoutMs = 10000): int
    {
        $this->running = true;
        $processed     = 0;

        while ($this->running) {
            if ($this->consumeOnce($handler, $timeoutMs)) {
                $processed++;
                if ($maxMessages !== null && $processed >= $maxMessages) {
                    break;
                }
            }
        }

        return $processed;
    }

    /** Ask the consume() loop to stop after the current message. */
    public function stop(): void
    {
        $this->running = false;
    }

    public function getGroupId(): string
    {
        return $this->groupId;
    }

    /**
     * Cleanly leave the group so partitions are reassigned quickly.
     */
    public function close(): void
    {
        $this->consumer->unsubscribe();
        $this->consumer->close();
    }

    /**
     * Turn an \RdKafka\Message into a plain array so handlers never
     * have to touch the rdkafka classes directly (easier to test/mock).
     *
     * @return array{topic:string,partition:int,offset:int,key:?string,payload:?string,headers:array}
     */
    private function normalise(\RdKafka\Message $message): array
    {
        return [
            'topic'     => $message->topic_name,
            'partition' => $message->partition,
            'offset'    => $message->offset,
            'key'       => $message->key,
            'payload'   => $message->payload,
            'headers'   => $message->headers ?? [],
        ];
    }
}
