<?php

require_once __DIR__ . '/KafkaConfig.php';

/**
 * Producer (PHP 7.0-compatible variant)
 *
 * Identical in behaviour to ../Producer.php, but written so it parses and
 * runs on PHP 7.0+ (no typed properties, no ?type / mixed syntax). Types are
 * documented in docblocks instead.
 *
 * Thin, safe wrapper around \RdKafka\Producer used to publish messages
 * to a topic. One Producer instance can publish to any number of topics
 * and is safe to reuse for the lifetime of a request/worker.
 *
 * A single message published here can be consumed independently by any
 * number of consumer groups (see Consumer) — that is a property of Kafka
 * itself: every distinct group.id keeps its own offset, so the same
 * message is delivered once to each group.
 */
class Producer
{
    /** @var \RdKafka\Producer */
    private $producer;

    /** @var array<string,\RdKafka\ProducerTopic> topic name => handle cache */
    private $topics = [];

    /** @var int how long (ms) to wait for outstanding deliveries on flush */
    private $flushTimeoutMs;

    /**
     * @param KafkaConfig $config
     * @param int         $flushTimeoutMs
     */
    public function __construct(KafkaConfig $config, $flushTimeoutMs = 10000)
    {
        $conf = $config->baseConf();

        // Fail fast if a message can't be delivered rather than hanging.
        $conf->set('message.timeout.ms', '30000');

        $this->producer       = new \RdKafka\Producer($conf);
        $this->flushTimeoutMs = $flushTimeoutMs;
    }

    /**
     * Publish a single message to a topic.
     *
     * @param string      $topic   destination topic name
     * @param mixed       $payload string, or any JSON-encodable value
     * @param string|null $key     partition/ordering key (optional)
     * @param array       $headers optional string headers ['k' => 'v']
     *
     * @return bool true once the message has been queued and flushed
     */
    public function produce($topic, $payload, $key = null, array $headers = [])
    {
        if ($topic === '') {
            throw new \InvalidArgumentException('Topic name cannot be empty.');
        }

        $value = is_string($payload) ? $payload : json_encode($payload, JSON_UNESCAPED_SLASHES);
        if ($value === false) {
            throw new \RuntimeException('Failed to JSON encode Kafka payload: ' . json_last_error_msg());
        }

        $topicHandle = $this->topicHandle($topic);

        // RD_KAFKA_PARTITION_UA lets librdkafka pick the partition
        // (by key hash if a key is given, else round-robin).
        $topicHandle->producev(
            RD_KAFKA_PARTITION_UA,
            0,
            $value,
            $key,
            $headers
        );

        // Serve delivery report callbacks.
        $this->producer->poll(0);

        return $this->flush();
    }

    /**
     * Block until all queued messages are delivered (or timeout).
     * Call this once after a batch of produce() calls if you disable
     * per-message flushing in the future.
     *
     * @return bool
     */
    public function flush()
    {
        $result = $this->producer->flush($this->flushTimeoutMs);

        if (RD_KAFKA_RESP_ERR_NO_ERROR !== $result) {
            throw new \RuntimeException(
                'Kafka flush failed / messages not delivered (err=' . $result . ').'
            );
        }

        return true;
    }

    /**
     * Get (and cache) a topic handle.
     *
     * @param string $topic
     * @return \RdKafka\ProducerTopic
     */
    private function topicHandle($topic)
    {
        if (!isset($this->topics[$topic])) {
            $this->topics[$topic] = $this->producer->newTopic($topic);
        }

        return $this->topics[$topic];
    }
}
