/
githubmirror
/
nbs
Обзор
Документация
Войти
/
githubmirror
/
nbs
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
main
yt/cpp/mapreduce/client/yt_poller.h
86 строк
2 KB
Anton Myagkov
Update ydb version to 24.4 (#4579)
07 фев 2026, 11:59
Не верифицирован
07 фев 2026, 11:59
2c21cae
Код
Авторство
О чём код?
#pragma once #include <yt/cpp/mapreduce/common/fwd.h> #include <yt/cpp/mapreduce/http/context.h> #include <yt/cpp/mapreduce/http/requests.h> #include <yt/cpp/mapreduce/interface/client.h> #include <util/generic/list.h> #include <util/system/mutex.h> #include <util/system/thread.h> #include <util/system/condvar.h> namespace NYT { namespace NDetail { namespace NRawClient { class TRawBatchRequest; } //////////////////////////////////////////////////////////////////////////////// class IYtPollerItem : public TThrRefBase { public: enum EStatus { PollContinue, PollBreak, }; public: virtual ~IYtPollerItem() = default; virtual void PrepareRequest(NRawClient::TRawBatchRequest* batchRequest) = 0; // Should return PollContinue if poller should continue polling this item. // Should return PollBreak if poller should stop polling this item. virtual EStatus OnRequestExecuted() = 0; virtual void OnItemDiscarded() = 0; }; using IYtPollerItemPtr = ::TIntrusivePtr<IYtPollerItem>; //////////////////////////////////////////////////////////////////////////////// class TYtPoller : public TThrRefBase { public: TYtPoller(TClientContext context, const IClientRetryPolicyPtr& retryPolicy); ~TYtPoller(); void Watch(IYtPollerItemPtr item); void Stop(); private: void DiscardQueuedItems(); void WatchLoop(); static void* WatchLoopProc(void*); private: struct TItem; const TClientContext Context_; const IClientRetryPolicyPtr ClientRetryPolicy_; TList<IYtPollerItemPtr> InProgress_; TList<IYtPollerItemPtr> Pending_; TThread WaiterThread_; TMutex Lock_; TCondVar HasData_; bool IsRunning_ = true; }; //////////////////////////////////////////////////////////////////////////////// } // namespace NDetail } // namespace NYT