blob: 880cd2c7364f82c936e25418ecab81b80f1727b7 [file]
/**
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing,
* software distributed under the License is distributed on an
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
* KIND, either express or implied. See the License for the
* specific language governing permissions and limitations
* under the License.
*/
#pragma once
#include <pulsar/defines.h>
#include <pulsar/st/Checkpoint.h>
#include <pulsar/st/Future.h>
#include <pulsar/st/detail/MessageCore.h>
#include <chrono>
#include <memory>
#include <string>
#include <string_view>
#include <vector>
namespace pulsar::st {
class CheckpointConsumerImpl;
using CheckpointConsumerImplPtr = std::shared_ptr<CheckpointConsumerImpl>;
namespace detail {
class ClientCore;
/**
* INTERNAL — not part of the public API. Non-templated checkpoint-consumer
* operations over the hidden impl (lib/st). `CheckpointConsumer<T>` wraps it.
*/
class PULSAR_PUBLIC CheckpointConsumerCore {
public:
CheckpointConsumerCore() = default;
Future<MessageCore> receiveAsync() const;
Future<MessageCore> receiveAsync(std::chrono::milliseconds timeout) const;
Future<std::vector<MessageCore>> receiveMultiAsync(int maxMessages,
std::chrono::milliseconds timeout) const;
Checkpoint checkpoint() const;
Future<void> closeAsync() const;
std::string_view topic() const;
std::string_view consumerName() const;
explicit operator bool() const { return static_cast<bool>(impl_); }
private:
friend class ClientCore;
explicit CheckpointConsumerCore(CheckpointConsumerImplPtr impl) : impl_(std::move(impl)) {}
CheckpointConsumerImplPtr impl_;
};
} // namespace detail
} // namespace pulsar::st