/* Definition of the pqxx::internal::stream_query class.
 *
 * Enables optimized batch reads from a database table.
 *
 * DO NOT INCLUDE THIS FILE DIRECTLY; include pqxx/transaction_base instead.
 *
 * Copyright (c) 2000-2026, Jeroen T. Vermeulen.
 *
 * See COPYING for copyright license.  If you did not receive a file called
 * COPYING with this source code, please notify the distributor of this
 * mistake, or contact the author.
 */
#ifndef PQXX_INTERNAL_STREAM_QUERY_HXX
#define PQXX_INTERNAL_STREAM_QUERY_HXX

#if !defined(PQXX_HEADER_PRE)
#  error "Include libpqxx headers as <pqxx/header>, not <pqxx/header.hxx>."
#endif

#include "pqxx/internal/encodings.hxx"
#include "pqxx/internal/gates/connection-stream_from.hxx"
#include "pqxx/transaction_focus.hxx"


namespace pqxx
{
class transaction_base;
} // namespace pqxx


namespace pqxx::internal
{
/// The `end()` iterator for a `stream_query`.
class stream_query_end_iterator
{};


// TODO: Can we use generators, and maybe get speedup from HALO?
/// Stream query results from the database.  Used by transaction_base::stream.
/** For larger data sets, retrieving data this way is likely to be faster than
 * executing a query and then iterating and converting the rows' fields.  You
 * will also be able to start processing before all of the data has come in.
 * (For smaller result sets though, a stream is likely to be a bit slower.)
 *
 * A `stream_query` stream is strongly typed.  You specify the columns' types
 * while instantiating the stream_query template.
 *
 * Not all kinds of query will work in a stream.  But straightforward `SELECT`
 * and `UPDATE ... RETURNING` queries should work.  The class uses PostgreSQL's
 * `COPY` command, so see the documentation for that command to get the full
 * details.
 *
 * There are other downsides.  If the stream encounters an error, it may leave
 * the entire connection in an unusable state, so you'll have to give the
 * whole thing up.  Finally, opening a stream puts the connection in a special
 * state, so you won't be able to do many other things with the connection or
 * the transaction while the stream is open.
 *
 * Usually you'll want the `stream` convenience wrapper in
 * @ref transaction_base, so you don't need to deal with this class directly.
 *
 * @warning While a stream is active, you cannot execute queries, open a
 * pipeline, etc. on the same transaction.  A transaction can have at most one
 * object of a type derived from @ref pqxx::transaction_focus active on it at a
 * time.
 */
template<typename... TYPE> class stream_query final : transaction_focus
{
public:
  using line_handle = std::unique_ptr<char[], void (*)(void const *)>;

  /// Execute `query` on `tx`, stream results.
  inline stream_query(
    transaction_base &tx, std::string_view query, conversion_context c);

  stream_query(stream_query const &) = delete;
  stream_query(stream_query &&) = delete;
  stream_query &operator=(stream_query const &) = delete;
  stream_query &operator=(stream_query &&) = delete;

  ~stream_query() noexcept
  {
    try
    {
      close();
    }
    catch (std::exception const &e)
    {
      reg_pending_error(e.what(), sl::current());
    }
  }

  /// Has this stream reached the end of its data?
  [[nodiscard]] bool done() const & noexcept
  {
    return m_char_finder == nullptr;
  }

  /// Begin iterator.  Only for use by "range for."
  [[nodiscard]] inline auto begin() &;

  /// End iterator.  Only for use by "range for."
  /** The end iterator is a different type than the regular iterator.  It
   * simplifies the comparisons: we know at compile time that we're comparing
   * to the end pointer.
   */
  [[nodiscard]] auto end() const & { return stream_query_end_iterator{}; }

  /// Parse and convert the latest line of data we received.
  std::tuple<TYPE...> parse_line(std::string_view line) &
  {
    assert(not done());

    auto const line_size{std::size(line)};

    // This function uses m_row as a buffer, across calls.  The only reason for
    // it to carry over across calls is to avoid reallocation.

    // Make room for unescaping the line.  It's a pessimistic size.
    // Unusually, we're storing terminating zeroes *inside* the string.
    // This is the only place where we modify m_row.  MAKE SURE THE BUFFER DOES
    // NOT GET RESIZED while we're working, because we're working with views
    // into its buffer.
    m_row.resize(line_size + 1);

    std::size_t offset{0u};
    char *write{m_row.data()};

    // DO NOT shrink m_row to fit.  We're carrying views pointing into the
    // buffer.  (Also, how useful would shrinking really be?)

    // Folding expression: scan and unescape each field, and convert it to its
    // requested type.
    std::tuple<TYPE...> data{parse_field<TYPE>(line, offset, write, m_ctx)...};

    assert(offset == line_size + 1u);
    return data;
  }

  /// Read a COPY line from the server.
  std::pair<line_handle, std::size_t> read_line(sl) &;

private:
  /// Look up a char_finder_func.
  /** This is the only encoding-dependent code in the class.  All we need to
   * store after that is this function pointer.
   */
  PQXX_PURE PQXX_RETURNS_NONNULL static inline char_finder_func *
  get_finder(transaction_base const &tx, sl);

  /// Scan and unescape a field into the row buffer.
  /** The row buffer is `m_row`.
   *
   * @param line The line of COPY output.
   * @param offset The current scanning position inside `line`.
   * @param write The current writing position in the row buffer.
   *
   * @return new `offset`; new `write`; and a `string_view` on the unescaped
   * field text in the row buffer.
   *
   * The `string_view`'s data pointer will be nullptr for a null field.
   *
   * After reading the final field in a row, if all goes well, offset should be
   * one greater than the size of the line, pointing at the terminating zero.
   */
  std::tuple<std::size_t, char *, std::string_view>
  read_field(std::string_view line, std::size_t offset, char *write, ctx c)
  {
#if !defined(NDEBUG)
    auto const line_size{std::size(line)};
#endif

    assert(offset <= line_size);

    char const *lp{std::data(line)};

    // The COPY line now ends in a tab.  (We replace the trailing newline with
    // that to simplify the loop here.)
    assert(lp[line_size] == '\t');
    assert(lp[line_size + 1] == '\0');

    if ((lp[offset] == '\\') and (lp[offset + 1] == 'N'))
    {
      // Null field.  Consume the "\N" and the field separator.
      offset += 3;
      assert(offset <= (line_size + 1));
      assert(lp[offset - 1] == '\t');
      // Return a null value.  There's nothing to write into m_row.
      return {offset, write, {}};
    }

    // Beginning of the field text in the row buffer.
    char const *const field_begin{write};

    // We're relying on several assumptions just for making the main loop
    // condition work:
    // * The COPY line ends in a newline.
    // * Multibyte characters never start with an ASCII-range byte.
    // * We can index a view beyond its bounds (but within its address space).
    //
    // Effectively, the newline acts as a final field separator.
    while (lp[offset] != '\t')
    {
      assert(lp[offset] != '\0');

      // Beginning of the next character of interest (or the end of the line).
      // It may be right where we start searching, and this won't loop forever
      // since the previous iteration (if any) put us right _after_ the
      // previous character of interest.
      auto const stop_char{m_char_finder(line, offset, c.loc)};
      PQXX_ASSUME(stop_char >= offset);
      assert(stop_char < (line_size + 1));

      // Copy the text we have so far.  It's got no special characters in it.
      std::memcpy(write, &lp[offset], stop_char - offset);
      write += (stop_char - offset);
      offset = stop_char;

      // We're still within the line.
      char const special{lp[offset]};
      if (special == '\\')
      {
        // Escape sequence.
        // Consume the backslash.
        ++offset;
        assert(offset < line_size);

        // The database will only escape ASCII characters, so we assume that
        // we're dealing with a single-byte character.
        char const escaped{lp[offset]};

        // I think this is a valid way to check for the high bit: the shift
        // may be signed or unsigned (implementation-defined for char), but
        // either way we get a zero if the bit is clear or nonzero if it's set.
        assert((escaped >> 7) == 0);
        ++offset;
        *write++ = unescape_char(escaped);
      }
      else
      {
        // Field separator.  Fall out of the loop.
        assert(special == '\t');
      }
    }

    // Hit the end of the field.
    assert(lp[offset] == '\t');
    *write = '\0';
    ++write;
    ++offset;
    return {
      offset,
      write,
      {field_begin, static_cast<std::size_t>(write - field_begin - 1)}};
  }

  /// Parse the next field.
  /** Unescapes the field into the row buffer (m_row), and converts it to its
   * TARGET type.
   *
   * Using non-const reference parameters here, so we can propagate side
   * effects across a fold expression.
   *
   * @param line The latest COPY line.
   * @param offset The current parsing offset in `line`.  The function will
   *   update this value.
   * @param write The current writing position in the row buffer.  The
   *   function will update this value.
   * @return Field value converted to TARGET type.
   */
  template<typename TARGET>
  TARGET
  parse_field(std::string_view line, std::size_t &offset, char *&write, ctx c)
  {
    using field_type = std::remove_cvref_t<TARGET>;

    assert(offset <= std::size(line));

    auto [new_offset, new_write, text]{read_field(line, offset, write, c)};
    PQXX_ASSUME(new_offset > offset);
    PQXX_ASSUME(new_write >= write);
    offset = new_offset;
    write = new_write;
    if constexpr (pqxx::always_null<TARGET>())
    {
      if (std::data(text) != nullptr)
        throw conversion_error{std::format(
          "Streaming a non-null value into a {}, which must always be null.",
          name_type<field_type>())};
    }
    else if (std::data(text) == nullptr)
    {
      if constexpr (has_null<TARGET>())
        return make_null<TARGET>();
      else
        internal::throw_null_conversion(name_type<field_type>(), c.loc);
    }
    else [[likely]]
    {
      // Don't ever try to convert a non-null value to nullptr_t!
      return from_string<field_type>(text, c);
    }
  }

  /// If this stream isn't already closed, close it now.
  void close() noexcept
  {
    if (not done())
    {
      m_char_finder = nullptr;
      unregister_me();
    }
  }

  /// Callback for finding next special character (or end of line).
  /** This pointer doubles as an indication that we're done.  We set it to
   * `nullptr` when the iteration is finished, and that's how we can know that
   * there are no more rows to be iterated.
   */
  char_finder_func *m_char_finder;

  /// Current row's fields' text, combined into one reusable string.
  /** We carry this buffer over from one invocation to the next, not because we
   * need the data, but just so we can re-use the space.  It saves us having to
   * re-allocate it every time.
   */
  std::string m_row;

  /// Caller source location, encoding group, possibly more.
  conversion_context const m_ctx;
};
} // namespace pqxx::internal
#endif
