diff options
| author | zverevgeny <[email protected]> | 2026-07-22 14:11:15 +0300 |
|---|---|---|
| committer | GitHub <[email protected]> | 2026-07-22 14:11:15 +0300 |
| commit | cef090cff683d97f35becb1cd1099b4860feaf18 (patch) | |
| tree | cd41fea9d8f34eb68297008eacea0da9039be6be | |
| parent | 05e1ec9189d97fcdb7484e419be4951ab91a87b9 (diff) | |
Decouple WorkloadManager and Kqp (#46638)
117 files changed, 846 insertions, 807 deletions
diff --git a/ydb/core/base/events.h b/ydb/core/base/events.h index 750fee9746d..9a6f13caeab 100644 --- a/ydb/core/base/events.h +++ b/ydb/core/base/events.h @@ -201,6 +201,7 @@ struct TKikimrEvents : TEvents { ES_SET_COLUMN_CONSTRAINT = 4278, ES_EXTERNAL_IDP_PROVIDER = 4279, ES_PQ_DEFERRED_PUBLISH = 4280, + ES_WORKLOAD_MANAGER = 4281, }; }; diff --git a/ydb/core/driver_lib/run/kikimr_services_initializers.cpp b/ydb/core/driver_lib/run/kikimr_services_initializers.cpp index 97b2ad601b8..64c129f2d5f 100644 --- a/ydb/core/driver_lib/run/kikimr_services_initializers.cpp +++ b/ydb/core/driver_lib/run/kikimr_services_initializers.cpp @@ -84,6 +84,7 @@ #include <ydb/core/kafka_proxy/kafka_proxy.h> #include <ydb/core/kafka_proxy/kafka_transactions_coordinator.h> +#include <ydb/services/workload_manager/service/service.h> #include <ydb/core/kqp/common/kqp.h> #include <ydb/core/kqp/proxy_service/kqp_proxy_service.h> #include <ydb/core/kqp/rm_service/kqp_rm_service.h> @@ -2428,6 +2429,17 @@ void TQuoterServiceInitializer::InitializeServices(NActors::TActorSystemSetup* s ); } +TWorkloadManagerServiceInitializer::TWorkloadManagerServiceInitializer(const TKikimrRunConfig& runConfig) + : IKikimrServicesInitializer(runConfig) +{} + +void TWorkloadManagerServiceInitializer::InitializeServices(NActors::TActorSystemSetup* setup, const NKikimr::TAppData* appData) { + auto workloadManager = NWorkloadManager::CreateService(NWorkloadManager::GetWorkloadManagerCounters(appData->Counters)); + setup->LocalServices.push_back(std::make_pair( + NWorkloadManager::MakeServiceId(NodeId), + TActorSetupCmd(workloadManager, TMailboxType::HTSwap, appData->UserPoolId))); +} + TKqpServiceInitializer::TKqpServiceInitializer( const TKikimrRunConfig& runConfig, std::shared_ptr<TModuleFactories> factories, diff --git a/ydb/core/driver_lib/run/kikimr_services_initializers.h b/ydb/core/driver_lib/run/kikimr_services_initializers.h index 284a53b8108..b67e8da3bca 100644 --- a/ydb/core/driver_lib/run/kikimr_services_initializers.h +++ b/ydb/core/driver_lib/run/kikimr_services_initializers.h @@ -408,6 +408,13 @@ public: void InitializeServices(NActors::TActorSystemSetup* setup, const NKikimr::TAppData* appData) override; }; +class TWorkloadManagerServiceInitializer : public IKikimrServicesInitializer { +public: + TWorkloadManagerServiceInitializer(const TKikimrRunConfig& runConfig); + + void InitializeServices(NActors::TActorSystemSetup* setup, const NKikimr::TAppData* appData) override; +}; + class TKqpServiceInitializer : public IKikimrServicesInitializer { public: TKqpServiceInitializer(const TKikimrRunConfig& runConfig, std::shared_ptr<TModuleFactories> factories, diff --git a/ydb/core/driver_lib/run/run.cpp b/ydb/core/driver_lib/run/run.cpp index 7718289fd78..e3718a301db 100644 --- a/ydb/core/driver_lib/run/run.cpp +++ b/ydb/core/driver_lib/run/run.cpp @@ -2127,6 +2127,10 @@ TIntrusivePtr<TServiceInitializersList> TKikimrRunner::CreateServiceInitializers sil->AddServiceInitializer(new TMemoryControllerInitializer(runConfig, ProcessMemoryInfoProvider)); + if (serviceMask.EnableWorkloadManagerService) { + sil->AddServiceInitializer(new TWorkloadManagerServiceInitializer(runConfig)); + } + if (serviceMask.EnableKqp) { sil->AddServiceInitializer(new TKqpServiceInitializer(runConfig, ModuleFactories, *this)); } diff --git a/ydb/core/driver_lib/run/service_mask.h b/ydb/core/driver_lib/run/service_mask.h index 0db11ca100a..812834bdbd1 100644 --- a/ydb/core/driver_lib/run/service_mask.h +++ b/ydb/core/driver_lib/run/service_mask.h @@ -86,6 +86,7 @@ union TBasicKikimrServicesMask { bool EnableCountersInfoProvider : 1; bool EnableNBSService : 1; bool EnableUdfStore : 1; + bool EnableWorkloadManagerService : 1; }; struct { diff --git a/ydb/core/driver_lib/run/ya.make b/ydb/core/driver_lib/run/ya.make index ad8c280ef8f..3b751848e0b 100644 --- a/ydb/core/driver_lib/run/ya.make +++ b/ydb/core/driver_lib/run/ya.make @@ -182,6 +182,7 @@ PEERDIR( ydb/services/tablet ydb/services/test_shard ydb/services/view + ydb/services/workload_manager/service ydb/services/ydb yql/essentials/minikql/comp_nodes/llvm16 yql/essentials/public/udf/service/exception_policy diff --git a/ydb/core/fq/libs/compute/ydb/control_plane/database_monitoring.cpp b/ydb/core/fq/libs/compute/ydb/control_plane/database_monitoring.cpp index f0077feedee..b128d269e44 100644 --- a/ydb/core/fq/libs/compute/ydb/control_plane/database_monitoring.cpp +++ b/ydb/core/fq/libs/compute/ydb/control_plane/database_monitoring.cpp @@ -1,7 +1,7 @@ #include <ydb/core/fq/libs/compute/ydb/events/events.h> #include <ydb/core/fq/libs/control_plane_storage/util.h> -#include <ydb/core/kqp/workload_service/common/cpu_quota_manager.h> +#include <ydb/services/workload_manager/common/cpu_quota_manager.h> #include <ydb/library/services/services.pb.h> @@ -193,7 +193,7 @@ private: const ui32 PendingQueueSize; const bool Strict; - NKikimr::NKqp::NWorkload::TCpuQuotaManager CpuQuotaManager; + NKikimr::NWorkloadManager::TCpuQuotaManager CpuQuotaManager; TQueue<TEvYdbCompute::TEvCpuQuotaRequest::TPtr> PendingQueue; TInstant StartCpuLoad; diff --git a/ydb/core/fq/libs/compute/ydb/control_plane/ya.make b/ydb/core/fq/libs/compute/ydb/control_plane/ya.make index 1c1019d3247..feb602740df 100644 --- a/ydb/core/fq/libs/compute/ydb/control_plane/ya.make +++ b/ydb/core/fq/libs/compute/ydb/control_plane/ya.make @@ -15,7 +15,7 @@ PEERDIR( ydb/core/fq/libs/compute/ydb/synchronization_service ydb/core/fq/libs/control_plane_storage/proto ydb/core/fq/libs/quota_manager/proto - ydb/core/kqp/workload_service/common + ydb/services/workload_manager/common ydb/core/protos ydb/library/actors/core ydb/library/actors/protos diff --git a/ydb/core/kqp/common/events/query.h b/ydb/core/kqp/common/events/query.h index cb9fba143df..8c883c04074 100644 --- a/ydb/core/kqp/common/events/query.h +++ b/ydb/core/kqp/common/events/query.h @@ -18,7 +18,7 @@ #include <memory> -namespace NKikimr::NKqp::NWorkload { +namespace NKikimr::NWorkloadManager { class ISessionUpdater; class IQueryClassifier; } @@ -368,19 +368,19 @@ public: return UserRequestContext; } - void SetWmSessionUpdater(const std::shared_ptr<NWorkload::ISessionUpdater>& wmSessionUpdater) { + void SetWmSessionUpdater(const std::shared_ptr<NWorkloadManager::ISessionUpdater>& wmSessionUpdater) { WmSessionUpdater = wmSessionUpdater; } - std::shared_ptr<NWorkload::ISessionUpdater> GetWmSessionUpdater() const { + std::shared_ptr<NWorkloadManager::ISessionUpdater> GetWmSessionUpdater() const { return WmSessionUpdater; } - void SetWmQueryClassifier(std::shared_ptr<NWorkload::IQueryClassifier> classifier) { + void SetWmQueryClassifier(std::shared_ptr<NWorkloadManager::IQueryClassifier> classifier) { WmQueryClassifier = std::move(classifier); } - std::shared_ptr<NWorkload::IQueryClassifier> GetWmQueryClassifier() const { + std::shared_ptr<NWorkloadManager::IQueryClassifier> GetWmQueryClassifier() const { return WmQueryClassifier; } @@ -521,8 +521,8 @@ private: std::shared_ptr<const NKikimrKqp::TQueryPhysicalGraph> QueryPhysicalGraph; i64 Generation = 0; bool DisableDefaultTimeout = false; - std::shared_ptr<NWorkload::ISessionUpdater> WmSessionUpdater; - std::shared_ptr<NWorkload::IQueryClassifier> WmQueryClassifier; + std::shared_ptr<NWorkloadManager::ISessionUpdater> WmSessionUpdater; + std::shared_ptr<NWorkloadManager::IQueryClassifier> WmQueryClassifier; }; struct TEvDataQueryStreamPart: public TEventPB<TEvDataQueryStreamPart, diff --git a/ydb/core/kqp/common/events/workload_service.cpp b/ydb/core/kqp/common/events/workload_service.cpp deleted file mode 100644 index 690db6a4904..00000000000 --- a/ydb/core/kqp/common/events/workload_service.cpp +++ /dev/null @@ -1 +0,0 @@ -#include "workload_service.h" diff --git a/ydb/core/kqp/common/events/ya.make b/ydb/core/kqp/common/events/ya.make index 260928b73a4..160a1ccefa0 100644 --- a/ydb/core/kqp/common/events/ya.make +++ b/ydb/core/kqp/common/events/ya.make @@ -3,7 +3,6 @@ LIBRARY() SRCS( events.cpp query.cpp - workload_service.cpp ) PEERDIR( diff --git a/ydb/core/kqp/common/simple/kqp_event_ids.h b/ydb/core/kqp/common/simple/kqp_event_ids.h index d269f585847..3dda9315aa6 100644 --- a/ydb/core/kqp/common/simple/kqp_event_ids.h +++ b/ydb/core/kqp/common/simple/kqp_event_ids.h @@ -189,17 +189,7 @@ struct TKqpResourceInfoExchangerEvents { }; }; -struct TKqpWorkloadServiceEvents { - enum EKqpWorkloadServiceEvents { - EvPlaceRequestIntoPool = EventSpaceBegin(TKikimrEvents::ES_KQP) + 700, - EvContinueRequest, - EvCleanupRequest, - EvCleanupResponse, - EvUpdatePoolInfo, - EvSubscribeOnPoolChanges, - EvFetchDatabaseResponse, - }; -}; + struct TKqpBufferWriterEvents { enum EKqpBufferWriterEvents { diff --git a/ydb/core/kqp/common/simple/services.h b/ydb/core/kqp/common/simple/services.h index 2a53d208071..b978b0ea593 100644 --- a/ydb/core/kqp/common/simple/services.h +++ b/ydb/core/kqp/common/simple/services.h @@ -41,10 +41,6 @@ inline NActors::TActorId MakeKqpFinalizeScriptServiceId(ui32 nodeId) { return NActors::TActorId(nodeId, TStringBuf(name, 12)); } -inline NActors::TActorId MakeKqpWorkloadServiceId(ui32 nodeId) { - const char name[12] = "kqp_workld"; - return NActors::TActorId(nodeId, TStringBuf(name, 12)); -} inline NActors::TActorId MakeKqpSchedulerServiceId(ui32 nodeId) { const char name[12] = "kqp_schdlr"; diff --git a/ydb/core/kqp/gateway/behaviour/ya.make b/ydb/core/kqp/gateway/behaviour/ya.make index cc2c79e609f..374a4c9eab2 100644 --- a/ydb/core/kqp/gateway/behaviour/ya.make +++ b/ydb/core/kqp/gateway/behaviour/ya.make @@ -1,7 +1,5 @@ RECURSE( external_data_source - resource_pool - resource_pool_classifier streaming_query table tablestore diff --git a/ydb/core/kqp/gateway/ya.make b/ydb/core/kqp/gateway/ya.make index 737a4db8a46..c536651ceb5 100644 --- a/ydb/core/kqp/gateway/ya.make +++ b/ydb/core/kqp/gateway/ya.make @@ -14,8 +14,9 @@ PEERDIR( ydb/core/kqp/federated_query/actors ydb/core/kqp/gateway/actors ydb/core/kqp/gateway/behaviour/external_data_source - ydb/core/kqp/gateway/behaviour/resource_pool - ydb/core/kqp/gateway/behaviour/resource_pool_classifier + ydb/services/workload_manager/metadata_subscription + ydb/services/workload_manager/metadata_subscription/resource_pool_classifier + ydb/services/workload_manager/service ydb/core/kqp/gateway/behaviour/streaming_query ydb/core/kqp/gateway/behaviour/table ydb/core/kqp/gateway/behaviour/tablestore diff --git a/ydb/core/kqp/proxy_service/kqp_proxy_databases_cache.cpp b/ydb/core/kqp/proxy_service/kqp_proxy_databases_cache.cpp index 4b949031032..3253fb50a87 100644 --- a/ydb/core/kqp/proxy_service/kqp_proxy_databases_cache.cpp +++ b/ydb/core/kqp/proxy_service/kqp_proxy_databases_cache.cpp @@ -1,6 +1,6 @@ #include "kqp_proxy_service_impl.h" -#include <ydb/core/kqp/workload_service/actors/actors.h> +#include <ydb/services/workload_manager/actors/actors.h> #include <ydb/core/tx/scheme_cache/scheme_cache.h> @@ -69,7 +69,7 @@ public: if (databaseStateIt == DatabaseStates.End()) { DatabaseStates.Insert({database, TDatabaseState{.Database = database}}); - Register(NWorkload::CreateDatabaseFetcherActor(SelfId(), database)); + Register(NWorkloadManager::CreateDatabaseFetcherActor(SelfId(), database)); StartIdleCheck(); return; } @@ -87,7 +87,7 @@ public: } } - void Handle(NWorkload::TEvFetchDatabaseResponse::TPtr& ev) { + void Handle(NWorkloadManager::TEvFetchDatabaseResponse::TPtr& ev) { auto databaseStateIt = DatabaseStates.Find(ev->Get()->Database); if (databaseStateIt == DatabaseStates.End()) { return; @@ -144,7 +144,7 @@ public: STRICT_STFUNC(StateFunc, hFunc(TEvPrivate::TEvSubscribeOnDatabase, Handle); hFunc(TEvPrivate::TEvPingDatabaseSubscription, Handle); - hFunc(NWorkload::TEvFetchDatabaseResponse, Handle); + hFunc(NWorkloadManager::TEvFetchDatabaseResponse, Handle); sFunc(TEvents::TEvPoison, HandlePoison); sFunc(TEvents::TEvWakeup, HandleWakeup); diff --git a/ydb/core/kqp/proxy_service/kqp_proxy_service.cpp b/ydb/core/kqp/proxy_service/kqp_proxy_service.cpp index e7b464503f5..ad76fb9cb92 100644 --- a/ydb/core/kqp/proxy_service/kqp_proxy_service.cpp +++ b/ydb/core/kqp/proxy_service/kqp_proxy_service.cpp @@ -15,7 +15,7 @@ #include <ydb/core/fq/libs/row_dispatcher/events/data_plane.h> #include <ydb/core/fq/libs/row_dispatcher/row_dispatcher_service.h> #include <ydb/core/kqp/common/events/script_executions.h> -#include <ydb/core/kqp/common/events/workload_service.h> +#include <ydb/services/workload_manager/events.h> #include <ydb/core/kqp/common/kqp_lwtrace_probes.h> #include <ydb/core/kqp/common/kqp_timeouts.h> #include <ydb/core/kqp/compile_service/kqp_compile_service.h> @@ -27,11 +27,10 @@ #include <ydb/core/kqp/finalize_script_service/kqp_finalize_script_service.h> #include <ydb/core/kqp/gateway/behaviour/streaming_query/behaviour.h> #include <ydb/core/kqp/node_service/kqp_node_service.h> -#include <ydb/core/kqp/workload_service/kqp_query_classifier.h> +#include <ydb/services/workload_manager/query_classifier.h> #include <ydb/core/kqp/proxy_service/kqp_query_text_cache_service.h> #include <ydb/core/kqp/rm_service/kqp_rm_service.h> #include <ydb/core/kqp/session_actor/kqp_worker_common.h> -#include <ydb/core/kqp/workload_service/kqp_workload_service.h> #include <ydb/core/mon/mon.h> #include <ydb/core/node_whiteboard/node_whiteboard.h> #include <ydb/core/protos/workload_manager_config.pb.h> @@ -390,10 +389,6 @@ public: TActivationContext::ActorSystem()->RegisterLocalService( MakeKqpNodeServiceID(SelfId().NodeId()), KqpNodeService); - KqpWorkloadService = TActivationContext::Register(CreateKqpWorkloadService(Counters->GetWorkloadManagerCounters())); - TActivationContext::ActorSystem()->RegisterLocalService( - MakeKqpWorkloadServiceId(SelfId().NodeId()), KqpWorkloadService); - auto updateFairSharePeriod = TDuration::MilliSeconds(TableServiceConfig.GetComputeSchedulerSettings().GetUpdateFairShareMs()); KqpComputeSchedulerService = TActivationContext::Register(CreateKqpComputeSchedulerService(updateFairSharePeriod)); TActivationContext::ActorSystem()->RegisterLocalService( @@ -516,7 +511,6 @@ public: Send(SpillingService, new TEvents::TEvPoison); Send(KqpNodeService, new TEvents::TEvPoison); - Send(KqpWorkloadService, new TEvents::TEvPoison()); Send(KqpComputeSchedulerService, new TEvents::TEvPoison()); Send(KqpQueryTextCacheService, new TEvents::TEvPoison()); if (RowDispatcherService) { @@ -1445,7 +1439,7 @@ public: hFunc(TEvInterconnect::TEvNodeDisconnected, Handle); hFunc(TEvKqp::TEvListSessionsRequest, Handle); hFunc(TEvKqp::TEvListProxyNodesRequest, Handle); - hFunc(NWorkload::TEvUpdatePoolInfo, Handle); + hFunc(NWorkloadManager::TEvUpdatePoolInfo, Handle); hFunc(TEvKqp::TEvUpdateDatabaseInfo, Handle); hFunc(TEvKqp::TEvDelayedRequestError, Handle); hFunc(NMetadata::NProvider::TEvRefreshSubscriberData, Handle); @@ -1663,14 +1657,14 @@ private: auto poolId = ev->Get()->GetPoolId(); ResourcePoolsCache.GetPoolInfo(databaseId, poolId ? poolId : NResourcePool::DEFAULT_POOL_ID, ActorContext()); - auto context = TClassifyContext{ + auto context = NWorkloadManager::TClassifyContext{ // This parameter is set when user creates a request with an explicit PoolId .PoolId = poolId, .AppName = sessionInfo ? sessionInfo->ClientApplicationName : "", .UserToken = ev->Get()->GetUserToken() }; - auto classifier = NWorkload::CreateQueryClassifier( + auto classifier = NWorkloadManager::CreateQueryClassifier( ResourcePoolsCache.GetLastResourcePoolMapSnapshot(), ResourcePoolsCache.GetClassifierViewFor(databaseId), databaseId, @@ -1941,7 +1935,7 @@ private: Send(ev->Sender, result.release(), 0, ev->Cookie); } - void Handle(NWorkload::TEvUpdatePoolInfo::TPtr& ev) { + void Handle(NWorkloadManager::TEvUpdatePoolInfo::TPtr& ev) { ResourcePoolsCache.UpdatePoolInfo(ev->Get()->DatabaseId, ev->Get()->PoolId, ev->Get()->Config, ev->Get()->SecurityObject, ActorContext()); } @@ -1958,7 +1952,7 @@ private: } void Handle(NMetadata::NProvider::TEvRefreshSubscriberData::TPtr& ev) { - ResourcePoolsCache.UpdateResourcePoolClassifiersInfo(ev->Get()->GetValidatedSnapshotAs<TResourcePoolClassifierSnapshot>(), ActorContext()); + ResourcePoolsCache.UpdateResourcePoolClassifiersInfo(ev->Get()->GetValidatedSnapshotAs<NWorkloadManager::TResourcePoolClassifierSnapshot>(), ActorContext()); } void InitSharedReading() { @@ -2077,7 +2071,6 @@ private: TActorId KqpNodeService; TActorId SpillingService; TActorId WhiteBoardService; - TActorId KqpWorkloadService; TActorId KqpComputeSchedulerService; TActorId KqpQueryTextCacheService; TActorId RowDispatcherService; diff --git a/ydb/core/kqp/proxy_service/kqp_proxy_service_impl.h b/ydb/core/kqp/proxy_service/kqp_proxy_service_impl.h index d34a10f4929..d7e469534ef 100644 --- a/ydb/core/kqp/proxy_service/kqp_proxy_service_impl.h +++ b/ydb/core/kqp/proxy_service/kqp_proxy_service_impl.h @@ -3,14 +3,14 @@ #include <ydb/core/base/appdata.h> #include <ydb/core/base/path.h> #include <ydb/core/kqp/common/kqp.h> -#include <ydb/core/kqp/common/events/workload_service.h> +#include <ydb/services/workload_manager/events.h> #include <ydb/core/kqp/counters/kqp_counters.h> -#include <ydb/core/kqp/workload_service/kqp_query_classifier.h> -#include <ydb/core/kqp/proxy_service/kqp_session_state.h> -#include <ydb/core/kqp/gateway/behaviour/resource_pool_classifier/fetcher.h> +#include <ydb/services/workload_manager/query_classifier.h> +#include <ydb/services/workload_manager/session_updater.h> +#include <ydb/services/workload_manager/service/service.h> +#include <ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/fetcher.h> #include <ydb/core/kqp/rm_service/kqp_rm_service.h> #include <ydb/core/kqp/runtime/scheduler/kqp_compute_scheduler_service.h> -#include <ydb/core/kqp/workload_service/kqp_workload_service.h> #include <ydb/core/protos/feature_flags.pb.h> #include <ydb/core/protos/kqp.pb.h> #include <ydb/core/protos/workload_manager_config.pb.h> @@ -79,9 +79,9 @@ public: } }; -class TWmSessionUpdater final : public NWorkload::ISessionUpdater { +class TWmSessionUpdater final : public NWorkloadManager::ISessionUpdater { public: - using EState = NWorkload::ISessionUpdater::EState; + using EState = NWorkloadManager::ISessionUpdater::EState; void SetRequestState(EState state, TInstant timestamp) override { const ui64 ts = timestamp.MicroSeconds(); @@ -511,12 +511,12 @@ public: } std::optional<TPoolInfo> GetPoolInfo(const TString& databaseId, const TString& poolId, TActorContext actorContext) const { - auto it = PoolsCache.find(GetPoolKey(databaseId, poolId)); + auto it = PoolsCache.find(NWorkloadManager::GetPoolKey(databaseId, poolId)); if (it == PoolsCache.end()) { Y_ASSERT(!poolId.empty()); actorContext.Send(MakeKqpSchedulerServiceId(actorContext.SelfID.NodeId()), new NScheduler::TEvAddPool(databaseId, poolId)); - actorContext.Send(MakeKqpWorkloadServiceId(actorContext.SelfID.NodeId()), new NWorkload::TEvSubscribeOnPoolChanges(databaseId, poolId)); + actorContext.Send(NWorkloadManager::MakeServiceId(actorContext.SelfID.NodeId()), new NWorkloadManager::TEvSubscribeOnPoolChanges(databaseId, poolId)); return std::nullopt; } return it->second; @@ -533,7 +533,7 @@ public: } void UpdatePoolInfo(const TString& databaseId, const TString& poolId, const std::optional<NResourcePool::TPoolSettings>& config, const std::optional<NACLib::TSecurityObject>& securityObject, TActorContext actorContext) { - const TString& poolKey = GetPoolKey(databaseId, poolId); + const TString& poolKey = NWorkloadManager::GetPoolKey(databaseId, poolId); if (!config) { auto it = PoolsCache.find(poolKey); if (it == PoolsCache.end()) { @@ -545,7 +545,7 @@ public: } else { // Refresh pool subscription it->second.Expired = true; - actorContext.Send(MakeKqpWorkloadServiceId(actorContext.SelfID.NodeId()), new NWorkload::TEvSubscribeOnPoolChanges(databaseId, poolId)); + actorContext.Send(NWorkloadManager::MakeServiceId(actorContext.SelfID.NodeId()), new NWorkloadManager::TEvSubscribeOnPoolChanges(databaseId, poolId)); } } else { auto& poolInfo = PoolsCache[poolKey]; @@ -557,7 +557,7 @@ public: BuildResourcePoolMapSnapshot(); } - void UpdateResourcePoolClassifiersInfo(std::shared_ptr<TResourcePoolClassifierSnapshot> snapshot, TActorContext actorContext) { + void UpdateResourcePoolClassifiersInfo(std::shared_ptr<NWorkloadManager::TResourcePoolClassifierSnapshot> snapshot, TActorContext actorContext) { LastClassifierSnapshot = snapshot; for (const auto& [databaseId, info] : snapshot->GetResourcePoolClassifierConfigs()) { for (const auto& [_, classifier] : info.ByName) { @@ -569,23 +569,23 @@ public: void UnsubscribeFromResourcePoolClassifiers(TActorContext actorContext) { if (SubscribedOnResourcePoolClassifiers) { SubscribedOnResourcePoolClassifiers = false; - actorContext.Send(NMetadata::NProvider::MakeServiceId(actorContext.SelfID.NodeId()), new NMetadata::NProvider::TEvUnsubscribeExternal(std::make_shared<TResourcePoolClassifierSnapshotsFetcher>())); + actorContext.Send(NMetadata::NProvider::MakeServiceId(actorContext.SelfID.NodeId()), new NMetadata::NProvider::TEvUnsubscribeExternal(std::make_shared<NWorkloadManager::TResourcePoolClassifierSnapshotsFetcher>())); } } - TClassifierConfigsView GetClassifierViewFor(const TString& databaseId) const { - return TClassifierConfigsView(LastClassifierSnapshot, databaseId); + NWorkloadManager::TClassifierConfigsView GetClassifierViewFor(const TString& databaseId) const { + return NWorkloadManager::TClassifierConfigsView(LastClassifierSnapshot, databaseId); } private: void BuildResourcePoolMapSnapshot() { - auto pools = std::make_shared<TResourcePoolMap>(); + auto pools = std::make_shared<NWorkloadManager::TResourcePoolMap>(); pools->reserve(PoolsCache.size()); for (const auto& [key, info] : PoolsCache) { if (!info.Expired) { - pools->emplace(key, TResourcePoolEntry{info.Config, info.SecurityObject}); + pools->emplace(key, NWorkloadManager::TResourcePoolEntry{info.Config, info.SecurityObject}); } } @@ -603,7 +603,7 @@ private: void SubscribeOnResourcePoolClassifiers(TActorContext actorContext) { if (!SubscribedOnResourcePoolClassifiers && NMetadata::NProvider::TServiceOperator::IsEnabled()) { SubscribedOnResourcePoolClassifiers = true; - actorContext.Send(NMetadata::NProvider::MakeServiceId(actorContext.SelfID.NodeId()), new NMetadata::NProvider::TEvSubscribeExternal(std::make_shared<TResourcePoolClassifierSnapshotsFetcher>())); + actorContext.Send(NMetadata::NProvider::MakeServiceId(actorContext.SelfID.NodeId()), new NMetadata::NProvider::TEvSubscribeExternal(std::make_shared<NWorkloadManager::TResourcePoolClassifierSnapshotsFetcher>())); } } @@ -620,17 +620,17 @@ private: } public: - const std::shared_ptr<const TResourcePoolClassifierSnapshot>& GetLastClassifierSnapshot() const { + const std::shared_ptr<const NWorkloadManager::TResourcePoolClassifierSnapshot>& GetLastClassifierSnapshot() const { return LastClassifierSnapshot; } - const std::shared_ptr<const TResourcePoolMap>& GetLastResourcePoolMapSnapshot() const { + const std::shared_ptr<const NWorkloadManager::TResourcePoolMap>& GetLastResourcePoolMapSnapshot() const { return LastResourcePoolMapSnapshot; } private: - std::shared_ptr<const TResourcePoolClassifierSnapshot> LastClassifierSnapshot; - std::shared_ptr<const TResourcePoolMap> LastResourcePoolMapSnapshot; + std::shared_ptr<const NWorkloadManager::TResourcePoolClassifierSnapshot> LastClassifierSnapshot; + std::shared_ptr<const NWorkloadManager::TResourcePoolMap> LastResourcePoolMapSnapshot; std::unordered_map<TString, TPoolInfo> PoolsCache; std::unordered_map<TString, TDatabaseInfo> DatabasesCache; diff --git a/ydb/core/kqp/proxy_service/kqp_proxy_ut.cpp b/ydb/core/kqp/proxy_service/kqp_proxy_ut.cpp index 0d9be16a32f..dc5f1853399 100644 --- a/ydb/core/kqp/proxy_service/kqp_proxy_ut.cpp +++ b/ydb/core/kqp/proxy_service/kqp_proxy_ut.cpp @@ -4,7 +4,7 @@ #include <ydb/core/kqp/proxy_service/kqp_proxy_service.h> #include <ydb/core/kqp/proxy_service/kqp_proxy_service_impl.h> #include <ydb/core/kqp/common/kqp.h> -#include <ydb/core/kqp/workload_service/ut/common/kqp_workload_service_ut_common.h> +#include <ydb/services/workload_manager/ut/common/workload_service_ut_common.h> #include <ydb/core/protos/config.pb.h> #include <ydb/core/protos/kqp.pb.h> #include <ydb/core/testlib/test_client.h> @@ -644,7 +644,7 @@ Y_UNIT_TEST_SUITE(KqpProxy) { } Y_UNIT_TEST(DatabasesCacheForServerless) { - auto ydb = NWorkload::TYdbSetupSettings() + auto ydb = NWorkloadManager::TYdbSetupSettings() .CreateSampleTenants(true) .Create(); diff --git a/ydb/core/kqp/proxy_service/kqp_session_info.cpp b/ydb/core/kqp/proxy_service/kqp_session_info.cpp index bcddfa27619..125b468a940 100644 --- a/ydb/core/kqp/proxy_service/kqp_session_info.cpp +++ b/ydb/core/kqp/proxy_service/kqp_session_info.cpp @@ -84,7 +84,7 @@ void TKqpSessionInfo::SerializeTo(::NKikimrKqp::TSessionInfo* proto, const TFiel if (fieldsMap.NeedField(VSessions::WmState::ColumnId)) { // 18 if (WmState) { - using EWmState = NWorkload::ISessionUpdater::EState; + using EWmState = NWorkloadManager::ISessionUpdater::EState; switch(WmState->GetState()) { case EWmState::NONE: { proto->SetWmState("NONE"); diff --git a/ydb/core/kqp/proxy_service/ut/ya.make b/ydb/core/kqp/proxy_service/ut/ya.make index eeca3a762a1..6bbc1d8fe4c 100644 --- a/ydb/core/kqp/proxy_service/ut/ya.make +++ b/ydb/core/kqp/proxy_service/ut/ya.make @@ -19,7 +19,7 @@ PEERDIR( ydb/core/kqp/run_script_actor ydb/core/kqp/proxy_service ydb/core/kqp/ut/common - ydb/core/kqp/workload_service/ut/common + ydb/services/workload_manager/ut/common ydb/public/lib/ut_helpers ydb/public/sdk/cpp/src/client/driver ydb/public/sdk/cpp/src/client/query diff --git a/ydb/core/kqp/proxy_service/ya.make b/ydb/core/kqp/proxy_service/ya.make index f2f29746d2a..221f57cbfab 100644 --- a/ydb/core/kqp/proxy_service/ya.make +++ b/ydb/core/kqp/proxy_service/ya.make @@ -20,12 +20,12 @@ PEERDIR( ydb/core/kqp/common/events ydb/core/kqp/compile_service ydb/core/kqp/counters - ydb/core/kqp/gateway/behaviour/resource_pool_classifier + ydb/services/workload_manager/metadata_subscription/resource_pool_classifier ydb/core/kqp/gateway/behaviour/streaming_query ydb/core/kqp/proxy_service/proto ydb/core/kqp/proxy_service/script_executions_utils ydb/core/kqp/run_script_actor - ydb/core/kqp/workload_service + ydb/services/workload_manager ydb/core/mind ydb/core/mon ydb/core/protos diff --git a/ydb/core/kqp/runtime/scheduler/kqp_compute_scheduler_service.cpp b/ydb/core/kqp/runtime/scheduler/kqp_compute_scheduler_service.cpp index ddd3f482dda..8288d5f2b5b 100644 --- a/ydb/core/kqp/runtime/scheduler/kqp_compute_scheduler_service.cpp +++ b/ydb/core/kqp/runtime/scheduler/kqp_compute_scheduler_service.cpp @@ -7,7 +7,8 @@ #include <ydb/core/base/feature_flags.h> #include <ydb/core/cms/console/configs_dispatcher.h> #include <ydb/core/cms/console/console.h> -#include <ydb/core/kqp/common/events/workload_service.h> +#include <ydb/services/workload_manager/events.h> +#include <ydb/services/workload_manager/service/service.h> #include <ydb/core/kqp/common/simple/services.h> #include <ydb/core/protos/feature_flags.pb.h> #include <ydb/core/protos/table_service_config.pb.h> @@ -58,7 +59,7 @@ public: hFunc(TEvAddDatabase, Handle); hFunc(TEvRemoveDatabase, Handle); hFunc(TEvAddPool, Handle); - hFunc(NWorkload::TEvUpdatePoolInfo, Handle); + hFunc(NWorkloadManager::TEvUpdatePoolInfo, Handle); hFunc(TEvRemovePool, Handle); hFunc(TEvAddQuery, Handle); hFunc(TEvRemoveQuery, Handle); @@ -126,7 +127,7 @@ public: if (PoolSubscribtions.insert({std::make_pair(databaseId, poolId), {.IsFirstRemoval=false, .ExternalWeight=resourceWeight}}).second) { PoolExternalWeightSum += resourceWeight; Scheduler->AddOrUpdatePool(databaseId, poolId, attrs); - Send(MakeKqpWorkloadServiceId(SelfId().NodeId()), new NWorkload::TEvSubscribeOnPoolChanges(databaseId, poolId)); + Send(NWorkloadManager::MakeServiceId(SelfId().NodeId()), new NWorkloadManager::TEvSubscribeOnPoolChanges(databaseId, poolId)); if (resourceWeight > Epsilon) { UpdatePoolsGuarantee(); } @@ -137,7 +138,7 @@ public: Y_ABORT("Unsupported yet"); } - void Handle(NWorkload::TEvUpdatePoolInfo::TPtr& ev) { + void Handle(NWorkloadManager::TEvUpdatePoolInfo::TPtr& ev) { const auto& databaseId = ev->Get()->DatabaseId; const auto& poolId = ev->Get()->PoolId; auto poolIt = PoolSubscribtions.find(std::make_pair(databaseId, poolId)); @@ -172,7 +173,7 @@ public: if (!poolIt->second.IsFirstRemoval) { // The first removal - try to re-subscribe in case it's just the pool removal from cache. poolIt->second.IsFirstRemoval = true; - Send(MakeKqpWorkloadServiceId(SelfId().NodeId()), new NWorkload::TEvSubscribeOnPoolChanges(databaseId, poolId)); + Send(NWorkloadManager::MakeServiceId(SelfId().NodeId()), new NWorkloadManager::TEvSubscribeOnPoolChanges(databaseId, poolId)); } else { // The second removal - the pool was really removed. PoolSubscribtions.erase(poolIt); diff --git a/ydb/core/kqp/runtime/scheduler/kqp_compute_scheduler_service_ut.cpp b/ydb/core/kqp/runtime/scheduler/kqp_compute_scheduler_service_ut.cpp index a58729c92ff..6841a75a0d8 100644 --- a/ydb/core/kqp/runtime/scheduler/kqp_compute_scheduler_service_ut.cpp +++ b/ydb/core/kqp/runtime/scheduler/kqp_compute_scheduler_service_ut.cpp @@ -1,10 +1,10 @@ -#include <ydb/core/kqp/workload_service/ut/common/kqp_workload_service_ut_common.h> +#include <ydb/services/workload_manager/ut/common/workload_service_ut_common.h> #include <ydb/library/testlib/helpers.h> namespace NKikimr::NKqp { -using namespace NWorkload; +using namespace NWorkloadManager; Y_UNIT_TEST_SUITE(KqpComputeSchedulerService) { diff --git a/ydb/core/kqp/runtime/ut/ya.make b/ydb/core/kqp/runtime/ut/ya.make index dd8197c8a6b..6259bfafde1 100644 --- a/ydb/core/kqp/runtime/ut/ya.make +++ b/ydb/core/kqp/runtime/ut/ya.make @@ -18,7 +18,7 @@ PEERDIR( library/cpp/testing/unittest ydb/core/kqp/common ydb/core/kqp/ut/common - ydb/core/kqp/workload_service/ut/common + ydb/services/workload_manager/ut/common ydb/core/testlib/basics/pg yql/essentials/minikql/comp_nodes/llvm16 yql/essentials/public/udf/service/exception_policy diff --git a/ydb/core/kqp/session_actor/kqp_query_state.h b/ydb/core/kqp/session_actor/kqp_query_state.h index 0c9ebb1c74d..22091a6159b 100644 --- a/ydb/core/kqp/session_actor/kqp_query_state.h +++ b/ydb/core/kqp/session_actor/kqp_query_state.h @@ -27,11 +27,13 @@ #include <map> #include <memory> -namespace NKikimr::NKqp { +namespace NKikimr { -namespace NWorkload { +namespace NWorkloadManager { class IQueryClassifier; -} // namespace NWorkload +} // namespace NWorkloadManager + +namespace NKqp { class TKqpQueryCache; @@ -150,7 +152,7 @@ public: TActorId Sender; ui64 ProxyRequestId = 0; std::unique_ptr<TEvKqp::TEvQueryRequest> RequestEv; - std::shared_ptr<NWorkload::IQueryClassifier> QueryClassifier; + std::shared_ptr<NWorkloadManager::IQueryClassifier> QueryClassifier; ui64 ParametersSize = 0; TPreparedQueryHolder::TConstPtr PreparedQuery; TKqpCompileResult::TConstPtr CompileResult; @@ -715,5 +717,6 @@ public: }; +} //namespace NKqp -} +} //namespace NKikimr diff --git a/ydb/core/kqp/session_actor/kqp_session_actor.cpp b/ydb/core/kqp/session_actor/kqp_session_actor.cpp index 68d9515531e..7f981c0289d 100644 --- a/ydb/core/kqp/session_actor/kqp_session_actor.cpp +++ b/ydb/core/kqp/session_actor/kqp_session_actor.cpp @@ -14,7 +14,7 @@ #include <ydb/core/kqp/common/kqp_timeouts.h> #include <ydb/core/kqp/common/kqp_tx.h> #include <ydb/core/kqp/common/kqp.h> -#include <ydb/core/kqp/common/events/workload_service.h> +#include <ydb/services/workload_manager/events.h> #include <ydb/core/kqp/common/simple/query_ast.h> #include <ydb/core/kqp/compile_service/kqp_compile_service.h> #include <ydb/core/kqp/executer_actor/kqp_executer.h> @@ -26,7 +26,7 @@ #include <ydb/core/kqp/opt/kqp_query_plan.h> #include <ydb/core/kqp/provider/yql_kikimr_provider.h> #include <ydb/core/kqp/provider/yql_kikimr_results.h> -#include <ydb/core/kqp/workload_service/kqp_query_classifier.h> +#include <ydb/services/workload_manager/query_classifier.h> #include <ydb/core/kqp/rm_service/kqp_snapshot_manager.h> #include <ydb/core/ydb_convert/ydb_convert.h> #include <ydb/core/tx/schemeshard/schemeshard.h> @@ -46,6 +46,7 @@ #include <ydb/core/protos/kqp.pb.h> #include <ydb/core/sys_view/service/sysview_service.h> #include <ydb/core/tx/tx_proxy/proxy.h> +#include <ydb/services/workload_manager/service/service.h> #include <ydb/library/actors/async/wait_for_event.h> #include <ydb/library/actors/core/actor_bootstrapped.h> @@ -347,7 +348,7 @@ public: Y_VALIDATE(!QueryState->UserRequestContext->PoolConfig, "Cannot send to workload manager: PoolConfig is already resolved"); - Send(MakeKqpWorkloadServiceId(SelfId().NodeId()), new NWorkload::TEvPlaceRequestIntoPool( + Send(NWorkloadManager::MakeServiceId(SelfId().NodeId()), new NWorkloadManager::TEvPlaceRequestIntoPool( QueryState->UserRequestContext->DatabaseId, SessionId, QueryState->UserRequestContext->PoolId, @@ -356,7 +357,7 @@ public: QueryState->RequestEv->GetWmSessionUpdater() ), IEventHandle::FlagTrackDelivery); - QueryState->PoolHandlerActor = MakeKqpWorkloadServiceId(SelfId().NodeId()); + QueryState->PoolHandlerActor = NWorkloadManager::MakeServiceId(SelfId().NodeId()); Become(&TKqpSessionActor::ExecuteState); } @@ -598,7 +599,7 @@ public: using TError = std::optional<std::pair<Ydb::StatusIds::StatusCode, TString>>; auto error = std::visit(TOverloaded { - [this, &sent](const NWorkload::IQueryClassifier::TResolvedPoolId& s) -> TError { + [this, &sent](const NWorkloadManager::IQueryClassifier::TResolvedPoolId& s) -> TError { STLOG_D("PreCompile Classify resolved", (pool_id, s.PoolId), (skip_admission, s.SkipAdmission), @@ -612,18 +613,18 @@ public: PassRequestToResourcePool(); return std::nullopt; }, - [this](const NWorkload::IQueryClassifier::TReject& r) -> TError { + [this](const NWorkloadManager::IQueryClassifier::TReject& r) -> TError { STLOG_N("PreCompile Classify rejected", (trace_id, TraceId())); return std::make_pair(r.Code, r.Message); }, - [this](const NWorkload::IQueryClassifier::TBypass&) -> TError { + [this](const NWorkloadManager::IQueryClassifier::TBypass&) -> TError { STLOG_D("PreCompile Classify bypass, compiling", (trace_id, TraceId())); QueryState->UserRequestContext->PoolId = NResourcePool::DEFAULT_POOL_ID; return std::nullopt; }, - [this](const NWorkload::IQueryClassifier::TPendingCompilation&) -> TError { + [this](const NWorkloadManager::IQueryClassifier::TPendingCompilation&) -> TError { STLOG_D("PreCompile Classify pending, compiling", (trace_id, TraceId())); return std::nullopt; @@ -639,7 +640,7 @@ public: } void Handle(TEvents::TEvUndelivered::TPtr& ev) { - if (ev->Get()->SourceType == TKqpWorkloadServiceEvents::EvPlaceRequestIntoPool) { + if (ev->Get()->SourceType == NWorkloadManager::TWorkloadManagerEvents::EvPlaceRequestIntoPool) { STLOG_W("Failed to deliver request to workload service, bypassing WLM", (trace_id, TraceId())); ContinueAfterWmAdmission(); @@ -659,7 +660,7 @@ public: } } - void Handle(NWorkload::TEvContinueRequest::TPtr& ev) { + void Handle(NWorkloadManager::TEvContinueRequest::TPtr& ev) { YQL_ENSURE(QueryState); QueryState->ContinueTime = TInstant::Now(); @@ -705,10 +706,10 @@ public: auto classifier = QueryState->QueryClassifier; auto state = classifier->GetState(); - if (state == NWorkload::IQueryClassifier::EState::PreCompileDone) { + if (state == NWorkloadManager::IQueryClassifier::EState::PreCompileDone) { STLOG_D("Pre-compile admission completed, compiling", (trace_id, TraceId())); CompileQuery(); - } else if (state == NWorkload::IQueryClassifier::EState::PostCompileDone) { + } else if (state == NWorkloadManager::IQueryClassifier::EState::PostCompileDone) { STLOG_D("Post-compile admission completed, executing", (trace_id, TraceId())); OnSuccessCompileRequest(); } else { @@ -1002,7 +1003,7 @@ public: bool WmPostCompileClassify() { auto classifier = QueryState->QueryClassifier; - if (!classifier || classifier->GetState() != NWorkload::IQueryClassifier::EState::WaitCompile) { + if (!classifier || classifier->GetState() != NWorkloadManager::IQueryClassifier::EState::WaitCompile) { return false; } @@ -1011,7 +1012,7 @@ public: using TError = std::optional<std::pair<Ydb::StatusIds::StatusCode, TString>>; auto error = std::visit(TOverloaded { - [this, &sent](const NWorkload::IQueryClassifier::TResolvedPoolId& r) -> TError { + [this, &sent](const NWorkloadManager::IQueryClassifier::TResolvedPoolId& r) -> TError { STLOG_D("PostCompile Classify resolved", (pool_id, r.PoolId), (skip_admission, r.SkipAdmission), @@ -1025,12 +1026,12 @@ public: PassRequestToResourcePool(); return std::nullopt; }, - [this](const NWorkload::IQueryClassifier::TBypass&) -> TError { + [this](const NWorkloadManager::IQueryClassifier::TBypass&) -> TError { STLOG_D("PostCompile Classify bypass", (trace_id, TraceId())); return std::nullopt; }, - [this](const NWorkload::IQueryClassifier::TReject& r) -> TError { + [this](const NWorkloadManager::IQueryClassifier::TReject& r) -> TError { STLOG_N("PostCompile Classify rejected", (trace_id, TraceId())); return std::make_pair(r.Code, r.Message); @@ -3520,12 +3521,12 @@ public: CleanupCtx->IsWaitingForWorkloadServiceCleanup = true; const auto& stats = QueryState->QueryStats; - auto event = std::make_unique<NWorkload::TEvCleanupRequest>( + auto event = std::make_unique<NWorkloadManager::TEvCleanupRequest>( QueryState->UserRequestContext->DatabaseId, SessionId, QueryState->UserRequestContext->PoolId, TDuration::MicroSeconds(stats.DurationUs), TDuration::MicroSeconds(stats.WorkerCpuTimeUs) ); - auto forwardId = MakeKqpWorkloadServiceId(SelfId().NodeId()); + auto forwardId = NWorkloadManager::MakeServiceId(SelfId().NodeId()); Send(new IEventHandle(*QueryState->PoolHandlerActor, SelfId(), event.release(), IEventHandle::FlagForwardOnNondelivery, 0, &forwardId)); QueryState->PoolHandlerActor = Nothing(); } @@ -3601,7 +3602,7 @@ public: } } - void HandleCleanup(NWorkload::TEvCleanupResponse::TPtr& ev) { + void HandleCleanup(NWorkloadManager::TEvCleanupResponse::TPtr& ev) { YQL_ENSURE(CleanupCtx); CleanupCtx->IsWaitingForWorkloadServiceCleanup = false; @@ -3758,7 +3759,7 @@ public: hFunc(TEvKqpExecuter::TEvExecuterProgress, HandleNoop) hFunc(TEvTxProxySchemeCache::TEvNavigateKeySetResult, HandleNoop); hFunc(TEvents::TEvUndelivered, HandleNoop); - hFunc(NWorkload::TEvContinueRequest, HandleNoop); + hFunc(NWorkloadManager::TEvContinueRequest, HandleNoop); // message from KQP proxy in case of our reply just after kqp proxy timer tick hFunc(NYql::NDq::TEvDq::TEvAbortExecution, HandleNoop); // A finished request's client may be lost after we already replied and @@ -3793,7 +3794,7 @@ public: hFunc(TEvTxUserProxy::TEvAllocateTxIdResult, Handle); hFunc(TEvents::TEvUndelivered, Handle); - hFunc(NWorkload::TEvContinueRequest, Handle); + hFunc(NWorkloadManager::TEvContinueRequest, Handle); hFunc(TEvKqpExecuter::TEvTxResponse, HandleExecute); hFunc(TEvKqpExecuter::TEvExecuterProgress, HandleExecute) @@ -3839,7 +3840,7 @@ public: hFunc(TEvKqp::TEvQueryRequest, Handle); hFunc(TEvKqpExecuter::TEvTxResponse, HandleCleanup); - hFunc(NWorkload::TEvCleanupResponse, HandleCleanup); + hFunc(NWorkloadManager::TEvCleanupResponse, HandleCleanup); hFunc(TEvKqp::TEvCloseSessionRequest, HandleCleanup); hFunc(NGRpcService::TEvClientLost, HandleNoop); @@ -3854,7 +3855,7 @@ public: hFunc(TEvents::TEvUndelivered, HandleNoop); hFunc(TEvTxUserProxy::TEvAllocateTxIdResult, HandleNoop); hFunc(TEvKqpExecuter::TEvStreamData, HandleNoop); - hFunc(NWorkload::TEvContinueRequest, HandleNoop); + hFunc(NWorkloadManager::TEvContinueRequest, HandleNoop); // always come from WorkerActor hFunc(TEvKqp::TEvCloseSessionResponse, HandleCleanup); @@ -3877,7 +3878,7 @@ public: hFunc(TEvents::TEvGone, HandleFinalCleanup); hFunc(TEvents::TEvUndelivered, HandleNoop); hFunc(TEvKqpSnapshot::TEvCreateSnapshotResponse, Handle); - hFunc(NWorkload::TEvContinueRequest, HandleNoop); + hFunc(NWorkloadManager::TEvContinueRequest, HandleNoop); hFunc(TEvKqp::TEvQueryRequest, HandleFinalCleanup); } } catch (const yexception& ex) { diff --git a/ydb/core/kqp/session_actor/ya.make b/ydb/core/kqp/session_actor/ya.make index 0398b0248e1..64f45b6ae0c 100644 --- a/ydb/core/kqp/session_actor/ya.make +++ b/ydb/core/kqp/session_actor/ya.make @@ -18,6 +18,7 @@ PEERDIR( ydb/library/security ydb/public/sdk/cpp/src/library/operation_id ydb/core/tx/schemeshard + ydb/services/workload_manager/service ) YQL_LAST_ABI_VERSION() diff --git a/ydb/core/kqp/ut/federated_query/datastreams/kqp_has_path_ut.cpp b/ydb/core/kqp/ut/federated_query/datastreams/kqp_has_path_ut.cpp index 3d3b992c865..9d0e4f288b0 100644 --- a/ydb/core/kqp/ut/federated_query/datastreams/kqp_has_path_ut.cpp +++ b/ydb/core/kqp/ut/federated_query/datastreams/kqp_has_path_ut.cpp @@ -66,7 +66,7 @@ void WaitClassifierVisible(TStreamingTestFixture& fixture, // HAS_PATH end-to-end coverage for the object kinds whose fixture requirements // (real local PQ, mock connector, mock PQ gateway, http gateway) are only met // by TStreamingTestFixture. Cheaper kinds (regular tables, sysview, secondary -// index, view underlying) live in ydb/core/kqp/workload_service/ut. +// index, view underlying) live in ydb/services/workload_manager/ut. Y_UNIT_TEST_SUITE(HasPathDatastreams) { // KindTopic — direct read of a local topic (no EDS in the chain). diff --git a/ydb/core/kqp/ut/scheme/kqp_scheme_ut.cpp b/ydb/core/kqp/ut/scheme/kqp_scheme_ut.cpp index 70352f4f6d9..31681ea506a 100644 --- a/ydb/core/kqp/ut/scheme/kqp_scheme_ut.cpp +++ b/ydb/core/kqp/ut/scheme/kqp_scheme_ut.cpp @@ -1,13 +1,13 @@ #include <ydb/core/base/tablet_resolver.h> #include <ydb/core/formats/arrow/arrow_helpers.h> #include <ydb/core/kqp/gateway/actors/scheme.h> -#include <ydb/core/kqp/gateway/behaviour/resource_pool_classifier/fetcher.h> +#include <ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/fetcher.h> #include <ydb/core/kqp/gateway/kqp_gateway.h> #include <ydb/core/kqp/ut/common/kqp_ut_common.h> #include <ydb/core/kqp/ut/common/columnshard.h> #include <ydb/core/kqp/ut/common/olap_indexes_enums.h> -#include <ydb/core/kqp/workload_service/actors/actors.h> -#include <ydb/core/kqp/workload_service/ut/common/kqp_workload_service_ut_common.h> +#include <ydb/services/workload_manager/actors/actors.h> +#include <ydb/services/workload_manager/ut/common/workload_service_ut_common.h> #include <ydb/core/protos/schemeshard/operations.pb.h> #include <ydb/core/testlib/cs_helper.h> #include <ydb/core/testlib/common_helper.h> @@ -8552,7 +8552,7 @@ Y_UNIT_TEST_SUITE(KqpScheme) { } Y_UNIT_TEST(DisableExternalDataSourcesOnServerless) { - auto ydb = NWorkload::TYdbSetupSettings() + auto ydb = NWorkloadManager::TYdbSetupSettings() .CreateSampleTenants(true) .EnableExternalDataSourcesOnServerless(false) .Create(); @@ -8589,21 +8589,21 @@ Y_UNIT_TEST_SUITE(KqpScheme) { const auto& dropTableSql = "DROP EXTERNAL TABLE MyExternalTable;"; - auto settings = NWorkload::TQueryRunnerSettings().PoolId(""); + auto settings = NWorkloadManager::TQueryRunnerSettings().PoolId(""); // Dedicated, enabled settings.Database(ydb->GetSettings().GetDedicatedTenantName()).NodeIndex(ydb->GetDedicatedTenantInfo().NodeIdx); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(createSourceSql, settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(createTableSql, settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(dropTableSql, settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(dropSourceSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(createSourceSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(createTableSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(dropTableSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(dropSourceSql, settings)); // Shared, enabled settings.Database(ydb->GetSettings().GetSharedTenantName()).NodeIndex(ydb->GetSharedTenantInfo().NodeIdx); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(createSourceSql, settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(createTableSql, settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(dropTableSql, settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(dropSourceSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(createSourceSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(createTableSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(dropTableSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(dropSourceSql, settings)); // Serverless, disabled settings.Database(ydb->GetSettings().GetServerlessTenantName()).NodeIndex(ydb->GetServerlessTenantInfo().NodeIdx); @@ -12427,7 +12427,7 @@ Y_UNIT_TEST_SUITE(KqpScheme) { } Y_UNIT_TEST(DisableResourcePoolsOnServerless) { - auto ydb = NWorkload::TYdbSetupSettings() + auto ydb = NWorkloadManager::TYdbSetupSettings() .CreateSampleTenants(true) .EnableResourcePoolsOnServerless(false) .Create(); @@ -12457,19 +12457,19 @@ Y_UNIT_TEST_SUITE(KqpScheme) { const auto& dropSql = "DROP RESOURCE POOL MyResourcePool;"; - auto settings = NWorkload::TQueryRunnerSettings().PoolId(""); + auto settings = NWorkloadManager::TQueryRunnerSettings().PoolId(""); // Dedicated, enabled settings.Database(ydb->GetSettings().GetDedicatedTenantName()).NodeIndex(ydb->GetDedicatedTenantInfo().NodeIdx); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(createSql, settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(alterSql, settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(dropSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(createSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(alterSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(dropSql, settings)); // Shared, enabled settings.Database(ydb->GetSettings().GetSharedTenantName()).NodeIndex(ydb->GetSharedTenantInfo().NodeIdx); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(createSql, settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(alterSql, settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(dropSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(createSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(alterSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(dropSql, settings)); // Serverless, disabled settings.Database(ydb->GetSettings().GetServerlessTenantName()).NodeIndex(ydb->GetServerlessTenantInfo().NodeIdx); @@ -12756,7 +12756,7 @@ Y_UNIT_TEST_SUITE(KqpScheme) { } Y_UNIT_TEST(DisableResourcePoolClassifiersOnServerless) { - auto ydb = NWorkload::TYdbSetupSettings() + auto ydb = NWorkloadManager::TYdbSetupSettings() .CreateSampleTenants(true) .EnableResourcePoolsOnServerless(false) .Create(); @@ -12785,23 +12785,23 @@ Y_UNIT_TEST_SUITE(KqpScheme) { const auto& dropSql = "DROP RESOURCE POOL CLASSIFIER MyResourcePoolClassifier;"; - auto settings = NWorkload::TQueryRunnerSettings().PoolId(""); + auto settings = NWorkloadManager::TQueryRunnerSettings().PoolId(""); // Dedicated, enabled settings.Database(ydb->GetSettings().GetDedicatedTenantName()).NodeIndex(ydb->GetDedicatedTenantInfo().NodeIdx); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery("CREATE RESOURCE POOL test_pool WITH (CONCURRENT_QUERY_LIMIT=10);", settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery("CREATE RESOURCE POOL test WITH (CONCURRENT_QUERY_LIMIT=10);", settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(createSql, settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(alterSql, settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(dropSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery("CREATE RESOURCE POOL test_pool WITH (CONCURRENT_QUERY_LIMIT=10);", settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery("CREATE RESOURCE POOL test WITH (CONCURRENT_QUERY_LIMIT=10);", settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(createSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(alterSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(dropSql, settings)); // Shared, enabled settings.Database(ydb->GetSettings().GetSharedTenantName()).NodeIndex(ydb->GetSharedTenantInfo().NodeIdx); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery("CREATE RESOURCE POOL test_pool WITH (CONCURRENT_QUERY_LIMIT=10);", settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery("CREATE RESOURCE POOL test WITH (CONCURRENT_QUERY_LIMIT=10);", settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(createSql, settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(alterSql, settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(dropSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery("CREATE RESOURCE POOL test_pool WITH (CONCURRENT_QUERY_LIMIT=10);", settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery("CREATE RESOURCE POOL test WITH (CONCURRENT_QUERY_LIMIT=10);", settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(createSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(alterSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(dropSql, settings)); // Serverless, disabled settings.Database(ydb->GetSettings().GetServerlessTenantName()).NodeIndex(ydb->GetServerlessTenantInfo().NodeIdx); @@ -12945,11 +12945,11 @@ Y_UNIT_TEST_SUITE(KqpScheme) { TString FetchResourcePoolClassifiers(TTestActorRuntime& runtime, ui32 nodeIndex) { const TActorId edgeActor = runtime.AllocateEdgeActor(nodeIndex); - runtime.Send(NMetadata::NProvider::MakeServiceId(runtime.GetNodeId(nodeIndex)), edgeActor, new NMetadata::NProvider::TEvAskSnapshot(std::make_shared<TResourcePoolClassifierSnapshotsFetcher>()), nodeIndex); + runtime.Send(NMetadata::NProvider::MakeServiceId(runtime.GetNodeId(nodeIndex)), edgeActor, new NMetadata::NProvider::TEvAskSnapshot(std::make_shared<NWorkloadManager::TResourcePoolClassifierSnapshotsFetcher>()), nodeIndex); const auto response = runtime.GrabEdgeEvent<NMetadata::NProvider::TEvRefreshSubscriberData>(edgeActor); UNIT_ASSERT(response); - return response->Get()->GetSnapshotAs<TResourcePoolClassifierSnapshot>()->SerializeToString(); + return response->Get()->GetSnapshotAs<NWorkloadManager::TResourcePoolClassifierSnapshot>()->SerializeToString(); } TString FetchResourcePoolClassifiers(TKikimrRunner& kikimr) { @@ -12994,7 +12994,7 @@ Y_UNIT_TEST_SUITE(KqpScheme) { } Y_UNIT_TEST(CreateResourcePoolClassifierOnServerless) { - auto ydb = NWorkload::TYdbSetupSettings() + auto ydb = NWorkloadManager::TYdbSetupSettings() .CreateSampleTenants(true) .EnableResourcePoolsOnServerless(true) .Create(); @@ -13003,7 +13003,7 @@ Y_UNIT_TEST_SUITE(KqpScheme) { const auto& serverlessTenant = ydb->GetSettings().GetServerlessTenantName(); ydb->ExecuteQueryRetry("Wait EnableResourcePools on Serverless", R"( CREATE RESOURCE POOL test_pool WITH (CONCURRENT_QUERY_LIMIT=10);)", - NWorkload::TQueryRunnerSettings() + NWorkloadManager::TQueryRunnerSettings() .PoolId("") .Database(serverlessTenant) .NodeIndex(nodeIdx) @@ -13013,7 +13013,7 @@ Y_UNIT_TEST_SUITE(KqpScheme) { RANK=20, RESOURCE_POOL="test_pool" );)", - NWorkload::TQueryRunnerSettings() + NWorkloadManager::TQueryRunnerSettings() .PoolId("") .Database(serverlessTenant) .NodeIndex(nodeIdx) @@ -13201,7 +13201,7 @@ Y_UNIT_TEST_SUITE(KqpScheme) { } Y_UNIT_TEST(DisableMetadataObjectsOnServerless) { - auto ydb = NWorkload::TYdbSetupSettings() + auto ydb = NWorkloadManager::TYdbSetupSettings() .CreateSampleTenants(true) .EnableMetadataObjectsOnServerless(false) .Create(); @@ -13216,28 +13216,28 @@ Y_UNIT_TEST_SUITE(KqpScheme) { const auto& upsertSql = "UPSERT OBJECT MySecretObject (TYPE SECRET) WITH value = \"edcba\";"; const auto& dropSql = "DROP OBJECT MySecretObject (TYPE SECRET);"; - auto settings = NWorkload::TQueryRunnerSettings().PoolId(""); + auto settings = NWorkloadManager::TQueryRunnerSettings().PoolId(""); // Dedicated, enabled settings.Database(ydb->GetSettings().GetDedicatedTenantName()).NodeIndex(ydb->GetDedicatedTenantInfo().NodeIdx); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(createSql, settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(alterSql, settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(upsertSql, settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(dropSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(createSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(alterSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(upsertSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(dropSql, settings)); // Shared, enabled settings.Database(ydb->GetSettings().GetSharedTenantName()).NodeIndex(ydb->GetSharedTenantInfo().NodeIdx); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(createSql, settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(alterSql, settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(upsertSql, settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(dropSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(createSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(alterSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(upsertSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(dropSql, settings)); // Serverless, disabled settings.Database(ydb->GetSettings().GetServerlessTenantName()).NodeIndex(ydb->GetServerlessTenantInfo().NodeIdx); checkDisabled(ydb->ExecuteQuery(createSql, settings)); checkDisabled(ydb->ExecuteQuery(alterSql, settings)); checkDisabled(ydb->ExecuteQuery(upsertSql, settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(dropSql, settings)); + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(dropSql, settings)); } Y_UNIT_TEST(CreateBackupCollectionDisabledByDefault) { @@ -14459,17 +14459,17 @@ END DO)", } Y_UNIT_TEST(StreamingQueriesOnServerless) { - auto ydb = NWorkload::TYdbSetupSettings() + auto ydb = NWorkloadManager::TYdbSetupSettings() .CreateSampleTenants(true) .Create(); const auto& tenantName = ydb->GetSettings().GetServerlessTenantName(); - const auto settings = NWorkload::TQueryRunnerSettings() + const auto settings = NWorkloadManager::TQueryRunnerSettings() .PoolId("") .Database(tenantName) .NodeIndex(ydb->GetServerlessTenantInfo().NodeIdx); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(fmt::format(R"( + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(fmt::format(R"( CREATE TOPIC MyTopic; CREATE EXTERNAL DATA SOURCE MySource WITH ( SOURCE_TYPE = "Ydb", @@ -14482,7 +14482,7 @@ END DO)", "database"_a = tenantName ), settings)); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(R"( + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(R"( CREATE STREAMING QUERY MyStreamingQuery WITH ( RUN = TRUE ) AS DO BEGIN INSERT INTO MySource.MyTopic SELECT * FROM MySource.MyTopic END DO @@ -14500,7 +14500,7 @@ END DO)", const auto& result = ydb->ExecuteQuery( TStringBuilder() << "SELECT * FROM `.sys/streaming_queries` " << filter , settings); - NWorkload::TSampleQueries::CheckSuccess(result); + NWorkloadManager::TSampleQueries::CheckSuccess(result); UNIT_ASSERT_VALUES_EQUAL(result.ResultSets.size(), 1); NYdb::TResultSetParser resultParser(result.ResultSets[0]); @@ -14534,7 +14534,7 @@ END DO)", checkSysView(queryText, false, TStringBuilder() << "WHERE Path > '" << queryName << "'"); checkSysView(queryText, false, TStringBuilder() << "WHERE Path < '" << queryName << "'"); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(R"( + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(R"( ALTER STREAMING QUERY MyStreamingQuery SET ( FORCE = TRUE ) AS DO BEGIN INSERT INTO MySource.MyTopic SELECT /* hint */ * FROM MySource.MyTopic END DO @@ -14550,7 +14550,7 @@ END DO)", Sleep(TDuration::Seconds(2)); checkSysView(queryText); - NWorkload::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(R"( + NWorkloadManager::TSampleQueries::CheckSuccess(ydb->ExecuteQuery(R"( DROP STREAMING QUERY MyStreamingQuery )", settings)); diff --git a/ydb/core/kqp/ut/scheme/ya.make b/ydb/core/kqp/ut/scheme/ya.make index 7974f7c2e1a..d9464d45057 100644 --- a/ydb/core/kqp/ut/scheme/ya.make +++ b/ydb/core/kqp/ut/scheme/ya.make @@ -25,7 +25,7 @@ PEERDIR( library/cpp/threading/local_executor ydb/core/kqp ydb/core/kqp/ut/common - ydb/core/kqp/workload_service/ut/common + ydb/services/workload_manager/ut/common ydb/core/tx/columnshard/hooks/testing ydb/public/sdk/cpp/src/client/arrow ydb/public/sdk/cpp/src/client/draft diff --git a/ydb/core/kqp/workload_service/kqp_has_stream_matcher.cpp b/ydb/core/kqp/workload_service/kqp_has_stream_matcher.cpp deleted file mode 100644 index 5b6a0b5ee9d..00000000000 --- a/ydb/core/kqp/workload_service/kqp_has_stream_matcher.cpp +++ /dev/null @@ -1,14 +0,0 @@ -#include "kqp_has_stream_matcher.h" -#include <ydb/core/kqp/common/kqp_user_request_context.h> - -namespace NKikimr::NKqp::NWorkload { - -bool MatchesStream(const std::optional<bool>& hasStream, - const TUserRequestContext& userRequestContext) { - if (!hasStream) { - return true; - } - return *hasStream ? userRequestContext.IsStreamingQuery : !userRequestContext.IsStreamingQuery; -} - -} // namespace NKikimr::NKqp::NWorkload diff --git a/ydb/core/kqp/workload_service/kqp_workload_service.h b/ydb/core/kqp/workload_service/kqp_workload_service.h deleted file mode 100644 index 75a36db5cb1..00000000000 --- a/ydb/core/kqp/workload_service/kqp_workload_service.h +++ /dev/null @@ -1,12 +0,0 @@ -#pragma once - -#include <ydb/core/resource_pools/resource_pool_settings.h> - -#include <ydb/library/actors/core/actor.h> - - -namespace NKikimr::NKqp { - -NActors::IActor* CreateKqpWorkloadService(NMonitoring::TDynamicCounterPtr counters); - -} // namespace NKikimr::NKqp diff --git a/ydb/core/kqp/workload_service/ut/common/ya.make b/ydb/core/kqp/workload_service/ut/common/ya.make deleted file mode 100644 index 087874d1a15..00000000000 --- a/ydb/core/kqp/workload_service/ut/common/ya.make +++ /dev/null @@ -1,16 +0,0 @@ -LIBRARY() - -SRCS( - kqp_query_classifier_ut_common.h - kqp_workload_service_ut_common.cpp -) - -PEERDIR( - ydb/core/kqp/gateway/behaviour/resource_pool_classifier - ydb/core/kqp/ut/common - ydb/services/metadata -) - -YQL_LAST_ABI_VERSION() - -END() diff --git a/ydb/core/kqp/workload_service/ut/kqp_has_stream_matcher_ut.cpp b/ydb/core/kqp/workload_service/ut/kqp_has_stream_matcher_ut.cpp deleted file mode 100644 index 12f77152c6b..00000000000 --- a/ydb/core/kqp/workload_service/ut/kqp_has_stream_matcher_ut.cpp +++ /dev/null @@ -1,52 +0,0 @@ -#include <ydb/core/kqp/workload_service/kqp_has_stream_matcher.h> - -#include <library/cpp/testing/unittest/registar.h> - - -namespace NKikimr::NKqp { - - -Y_UNIT_TEST_SUITE(TStreamMatcherPredicate) { - - - Y_UNIT_TEST(NulloptAcceptsAny) { - // No filter — every query passes regardless of streaming ops. - { - TUserRequestContext userRequestContext; - userRequestContext.IsStreamingQuery = false; - UNIT_ASSERT(NWorkload::MatchesStream(std::nullopt, userRequestContext)); - } - { - TUserRequestContext userRequestContext; - userRequestContext.IsStreamingQuery = true; - UNIT_ASSERT(NWorkload::MatchesStream(std::nullopt, userRequestContext)); - } - } - - Y_UNIT_TEST(TrueMatchesStreamingQuery) { - TUserRequestContext userRequestContext; - userRequestContext.IsStreamingQuery = true; - - UNIT_ASSERT(NWorkload::MatchesStream(true, userRequestContext)); - } - - Y_UNIT_TEST(TrueDoesNotMatchNonStreamingQuery) { - TUserRequestContext userRequestContext; - userRequestContext.IsStreamingQuery = false; - UNIT_ASSERT(!NWorkload::MatchesStream(true, userRequestContext)); - } - - Y_UNIT_TEST(FalseMatchesNonStreamingQuery) { - TUserRequestContext userRequestContext; - userRequestContext.IsStreamingQuery = false; - UNIT_ASSERT(NWorkload::MatchesStream(false, userRequestContext)); - } - - Y_UNIT_TEST(FalseDoesNotMatchStreamingQuery) { - TUserRequestContext userRequestContext; - userRequestContext.IsStreamingQuery = true; - UNIT_ASSERT(!NWorkload::MatchesStream(false, userRequestContext)); - } -} - -} // namespace NKikimr::NKqp diff --git a/ydb/core/kqp/workload_service/ut/ya.make b/ydb/core/kqp/workload_service/ut/ya.make deleted file mode 100644 index 97b51dfe873..00000000000 --- a/ydb/core/kqp/workload_service/ut/ya.make +++ /dev/null @@ -1,40 +0,0 @@ -UNITTEST_FOR(ydb/core/kqp/workload_service) - -FORK_SUBTESTS() - -SIZE(MEDIUM) -IF (SANITIZER_TYPE) - REQUIREMENTS(cpu:4) -ELSE() - REQUIREMENTS(cpu:2) -ENDIF() - -SRCS( - kqp_action_reject_ut.cpp - kqp_has_app_name_ut.cpp - kqp_has_full_scan_matcher_ut.cpp - kqp_has_full_scan_ut.cpp - kqp_has_path_ddl_ut.cpp - kqp_has_path_matcher_ut.cpp - kqp_has_path_ut.cpp - kqp_stream_query_classification_ut.cpp - kqp_has_stream_matcher_ut.cpp - kqp_has_stream_ut.cpp - kqp_member_name_ut.cpp - kqp_query_classifier_match_ut.cpp - kqp_query_classifier_ut.cpp - kqp_workload_service_actors_ut.cpp - kqp_workload_service_query_sessions_ut.cpp - kqp_workload_service_tables_ut.cpp - kqp_workload_service_ut.cpp -) - -PEERDIR( - ydb/core/kqp/workload_service/ut/common - - yql/essentials/sql/pg_dummy -) - -YQL_LAST_ABI_VERSION() - -END() diff --git a/ydb/core/kqp/workload_service/ya.make b/ydb/core/kqp/workload_service/ya.make deleted file mode 100644 index ba5c5bd0056..00000000000 --- a/ydb/core/kqp/workload_service/ya.make +++ /dev/null @@ -1,44 +0,0 @@ -LIBRARY() - -SRCS( - kqp_has_full_scan_matcher.cpp - kqp_has_path_matcher.cpp - kqp_has_stream_matcher.cpp - kqp_query_classifier.cpp - kqp_workload_service.cpp -) - -PEERDIR( - ydb/core/cms/console - - ydb/core/fq/libs/compute/common - - ydb/core/kqp/common - ydb/core/kqp/gateway/behaviour/resource_pool_classifier - ydb/core/kqp/query_data - ydb/core/kqp/workload_service/actors - - ydb/core/mind - - ydb/core/resource_pools - - ydb/library/actors/interconnect - ydb/library/aclib - ydb/library/yql/providers/pq/proto - - ydb/public/api/protos -) - -YQL_LAST_ABI_VERSION() - -END() - -RECURSE( - actors - common - tables -) - -RECURSE_FOR_TESTS( - ut -) diff --git a/ydb/core/kqp/ya.make b/ydb/core/kqp/ya.make index d2cd8434ed2..5b281c3b21f 100644 --- a/ydb/core/kqp/ya.make +++ b/ydb/core/kqp/ya.make @@ -74,7 +74,6 @@ RECURSE( runtime session_actor tests - workload_service ) RECURSE_FOR_TESTS( diff --git a/ydb/core/sys_view/resource_pool_classifiers/resource_pool_classifiers.cpp b/ydb/core/sys_view/resource_pool_classifiers/resource_pool_classifiers.cpp index f9e2c703a68..dc46181dff8 100644 --- a/ydb/core/sys_view/resource_pool_classifiers/resource_pool_classifiers.cpp +++ b/ydb/core/sys_view/resource_pool_classifiers/resource_pool_classifiers.cpp @@ -1,8 +1,8 @@ #include "resource_pool_classifiers.h" -#include <ydb/core/kqp/gateway/behaviour/resource_pool_classifier/fetcher.h> -#include <ydb/core/kqp/gateway/behaviour/resource_pool_classifier/snapshot.h> -#include <ydb/core/kqp/workload_service/actors/actors.h> +#include <ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/fetcher.h> +#include <ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/snapshot.h> +#include <ydb/services/workload_manager/actors/actors.h> #include <ydb/core/node_whiteboard/node_whiteboard.h> #include <ydb/core/sys_view/common/events.h> #include <ydb/core/sys_view/common/registry.h> @@ -46,7 +46,7 @@ public: switch (ev->GetTypeRewrite()) { sFunc(NKqp::TEvKqpCompute::TEvScanDataAck, HandleAck); hFunc(NMetadata::NProvider::TEvRefreshSubscriberData, Handle) - hFunc(NKqp::NWorkload::TEvFetchDatabaseResponse, Handle); + hFunc(NWorkloadManager::TEvFetchDatabaseResponse, Handle); hFunc(NKqp::TEvKqp::TEvAbortExecution, HandleAbortExecution); cFunc(TEvents::TEvWakeup::EventType, HandleTimeout); cFunc(TEvents::TEvPoison::EventType, PassAway); @@ -66,63 +66,63 @@ private: if (!NMetadata::NProvider::TServiceOperator::IsEnabled()) { ReplyEmptyAndDie(); } - Register(NKqp::NWorkload::CreateDatabaseFetcherActor(SelfId(), TenantName, UserToken, NACLib::EAccessRights::GenericUse)); + Register(NWorkloadManager::CreateDatabaseFetcherActor(SelfId(), TenantName, UserToken, NACLib::EAccessRights::GenericUse)); } - void Handle(NKqp::NWorkload::TEvFetchDatabaseResponse::TPtr& ev) { + void Handle(NWorkloadManager::TEvFetchDatabaseResponse::TPtr& ev) { auto& event = *ev->Get(); if (event.Status != Ydb::StatusIds::SUCCESS) { ReplyErrorAndDie(event.Status, event.Issues.ToOneLineString()); return; } DatabaseId = event.DatabaseId; - Send(NMetadata::NProvider::MakeServiceId(SelfId().NodeId()), new NMetadata::NProvider::TEvAskSnapshot(std::make_shared<NKqp::TResourcePoolClassifierSnapshotsFetcher>())); + Send(NMetadata::NProvider::MakeServiceId(SelfId().NodeId()), new NMetadata::NProvider::TEvAskSnapshot(std::make_shared<NWorkloadManager::TResourcePoolClassifierSnapshotsFetcher>())); } void Handle(NMetadata::NProvider::TEvRefreshSubscriberData::TPtr& ev) { - using TExtractor = std::function<TCell(const NKqp::TResourcePoolClassifierConfig&)>; + using TExtractor = std::function<TCell(const NWorkloadManager::TResourcePoolClassifierConfig&)>; using TSchema = Schema::ResourcePoolClassifiers; struct TExtractorsMap : public THashMap<NTable::TTag, TExtractor> { TExtractorsMap() { - insert({TSchema::Name::ColumnId, [] (const NKqp::TResourcePoolClassifierConfig& config) { + insert({TSchema::Name::ColumnId, [] (const NWorkloadManager::TResourcePoolClassifierConfig& config) { return TCell(config.GetName().data(), config.GetName().size()); }}); - insert({TSchema::Rank::ColumnId, [] (const NKqp::TResourcePoolClassifierConfig& config) { + insert({TSchema::Rank::ColumnId, [] (const NWorkloadManager::TResourcePoolClassifierConfig& config) { return TCell::Make<i64>(config.GetRank()); }}); - insert({TSchema::MemberName::ColumnId, [] (const NKqp::TResourcePoolClassifierConfig& config) { + insert({TSchema::MemberName::ColumnId, [] (const NWorkloadManager::TResourcePoolClassifierConfig& config) { const auto& memberName = config.GetConfigJson()["member_name"].GetString(); return TCell(memberName.data(), memberName.size()); }}); - insert({TSchema::ResourcePool::ColumnId, [] (const NKqp::TResourcePoolClassifierConfig& config) { + insert({TSchema::ResourcePool::ColumnId, [] (const NWorkloadManager::TResourcePoolClassifierConfig& config) { const auto& resourcePool = config.GetConfigJson()["resource_pool"].GetString(); return TCell(resourcePool.data(), resourcePool.size()); }}); - insert({TSchema::HasAppName::ColumnId, [] (const NKqp::TResourcePoolClassifierConfig& config) { + insert({TSchema::HasAppName::ColumnId, [] (const NWorkloadManager::TResourcePoolClassifierConfig& config) { const auto& hasAppName = config.GetConfigJson()["has_app_name"].GetString(); return TCell(hasAppName.data(), hasAppName.size()); }}); - insert({TSchema::Action::ColumnId, [] (const NKqp::TResourcePoolClassifierConfig& config) { + insert({TSchema::Action::ColumnId, [] (const NWorkloadManager::TResourcePoolClassifierConfig& config) { const auto& action = config.GetConfigJson()["action"].GetString(); return TCell(action.data(), action.size()); }}); - insert({TSchema::HasFullScan::ColumnId, [] (const NKqp::TResourcePoolClassifierConfig& config) { + insert({TSchema::HasFullScan::ColumnId, [] (const NWorkloadManager::TResourcePoolClassifierConfig& config) { const auto& hasFullScan = config.GetConfigJson()["has_full_scan"].GetString(); return TCell(hasFullScan.data(), hasFullScan.size()); }}); - insert({TSchema::HasPath::ColumnId, [] (const NKqp::TResourcePoolClassifierConfig& config) { + insert({TSchema::HasPath::ColumnId, [] (const NWorkloadManager::TResourcePoolClassifierConfig& config) { const auto& hasPath = config.GetConfigJson()["has_path"].GetString(); return TCell(hasPath.data(), hasPath.size()); }}); - insert({TSchema::HasStream::ColumnId, [] (const NKqp::TResourcePoolClassifierConfig& config) { + insert({TSchema::HasStream::ColumnId, [] (const NWorkloadManager::TResourcePoolClassifierConfig& config) { return TCell::Make<bool>(config.GetConfigJson()["has_stream"].GetBoolean()); }}); } }; static TExtractorsMap extractors; - const auto& snapshot = ev->Get()->GetSnapshotAs<NKqp::TResourcePoolClassifierSnapshot>(); + const auto& snapshot = ev->Get()->GetSnapshotAs<NWorkloadManager::TResourcePoolClassifierSnapshot>(); const auto& config = snapshot->GetResourcePoolClassifierConfigs(); auto resourcePoolsIt = config.find(DatabaseId); if (resourcePoolsIt == config.end()) { diff --git a/ydb/core/sys_view/streaming_queries/streaming_queries.cpp b/ydb/core/sys_view/streaming_queries/streaming_queries.cpp index ef935174e2b..aa8d27ba170 100644 --- a/ydb/core/sys_view/streaming_queries/streaming_queries.cpp +++ b/ydb/core/sys_view/streaming_queries/streaming_queries.cpp @@ -3,7 +3,7 @@ #include <ydb/core/kqp/common/events/script_executions.h> #include <ydb/core/kqp/common/kqp_script_executions.h> #include <ydb/core/kqp/gateway/behaviour/streaming_query/common/utils.h> -#include <ydb/core/kqp/workload_service/actors/actors.h> +#include <ydb/services/workload_manager/actors/actors.h> #include <ydb/core/sys_view/common/registry.h> #include <ydb/core/sys_view/common/scan_actor_base_impl.h> #include <ydb/library/query_actor/query_actor.h> @@ -613,7 +613,7 @@ public: STFUNC(StateScan) { switch (ev->GetTypeRewrite()) { hFunc(NKqp::TEvKqpCompute::TEvScanDataAck, Handle); - hFunc(NKqp::NWorkload::TEvFetchDatabaseResponse, Handle); + hFunc(NWorkloadManager::TEvFetchDatabaseResponse, Handle); hFunc(TEvPrivate::TEvCheckStreamingQueriesTables, Handle); hFunc(TEvPrivate::TEvFetchStreamingQueriesResult, Handle); hFunc(TEvPrivate::TEvDescribeStreamingQueriesResult, Handle); @@ -633,7 +633,7 @@ public: ContinueScan(); } - void Handle(NKqp::NWorkload::TEvFetchDatabaseResponse::TPtr& ev) { + void Handle(NWorkloadManager::TEvFetchDatabaseResponse::TPtr& ev) { HasInflightOperation = false; if (const auto status = ev->Get()->Status; status != Ydb::StatusIds::SUCCESS) { @@ -834,7 +834,7 @@ private: if (!DatabaseId) { HasInflightOperation = true; - const auto& databaseFetcher = Register(NKqp::NWorkload::CreateDatabaseFetcherActor(SelfId(), DatabaseName, MakeIntrusive<NACLib::TUserToken>(BUILTIN_ACL_METADATA, TVector<NACLib::TSID>{}))); + const auto& databaseFetcher = Register(NWorkloadManager::CreateDatabaseFetcherActor(SelfId(), DatabaseName, MakeIntrusive<NACLib::TUserToken>(BUILTIN_ACL_METADATA, TVector<NACLib::TSID>{}))); LOG_D("Start database fetcher " << databaseFetcher); return; } diff --git a/ydb/core/sys_view/streaming_queries/ya.make b/ydb/core/sys_view/streaming_queries/ya.make index 7ce9098465f..d819f52862f 100644 --- a/ydb/core/sys_view/streaming_queries/ya.make +++ b/ydb/core/sys_view/streaming_queries/ya.make @@ -9,7 +9,7 @@ PEERDIR( ydb/core/kqp/common/events ydb/core/kqp/gateway/behaviour/streaming_query/common ydb/core/kqp/runtime - ydb/core/kqp/workload_service/actors + ydb/services/workload_manager/actors ydb/core/sys_view/common ydb/library/actors/core ydb/library/query_actor diff --git a/ydb/core/testlib/test_client.cpp b/ydb/core/testlib/test_client.cpp index ca57c2a7c37..aee2907175d 100644 --- a/ydb/core/testlib/test_client.cpp +++ b/ydb/core/testlib/test_client.cpp @@ -76,6 +76,7 @@ #include <ydb/core/kqp/common/kqp.h> #include <ydb/core/kqp/rm_service/kqp_rm_service.h> #include <ydb/core/kqp/proxy_service/kqp_proxy_service.h> +#include <ydb/services/workload_manager/service/service.h> #include <ydb/core/kqp/finalize_script_service/kqp_finalize_script_service.h> #include <ydb/core/metering/metering.h> #include <ydb/core/protos/replication.pb.h> @@ -1352,6 +1353,14 @@ namespace Tests { TActorId describeSchemaSecretsServiceId = Runtime->Register(describeSchemaSecretsService, nodeIdx, userPoolId); Runtime->RegisterService(NSecret::MakeDescribeSchemaSecretServiceId(Runtime->GetNodeId(nodeIdx)), describeSchemaSecretsServiceId, nodeIdx); } + + { + const auto& appData = Runtime->GetAppData(nodeIdx); + IActor* workloadManager = NWorkloadManager::CreateService(NWorkloadManager::GetWorkloadManagerCounters(appData.Counters)); + TActorId workloadManagerId = Runtime->Register(workloadManager, nodeIdx, userPoolId, TMailboxType::HTSwap, 0); + Runtime->RegisterService(NWorkloadManager::MakeServiceId(Runtime->GetNodeId(nodeIdx)), workloadManagerId, nodeIdx); + } + { auto kqpProxySharedResources = std::make_shared<NKqp::TKqpProxySharedResources>(); diff --git a/ydb/core/testlib/ya.make b/ydb/core/testlib/ya.make index 9a63cfcd979..7b25d62c931 100644 --- a/ydb/core/testlib/ya.make +++ b/ydb/core/testlib/ya.make @@ -57,6 +57,7 @@ PEERDIR( ydb/core/kqp ydb/core/kqp/federated_query ydb/services/scheme_secret + ydb/services/workload_manager/service ydb/core/kqp/finalize_script_service ydb/core/kqp/proxy_service ydb/core/metering diff --git a/ydb/core/kqp/workload_service/actors/actors.h b/ydb/services/workload_manager/actors/actors.h index 1dfcaa85911..82b68c5667f 100644 --- a/ydb/core/kqp/workload_service/actors/actors.h +++ b/ydb/services/workload_manager/actors/actors.h @@ -1,10 +1,10 @@ #pragma once -#include <ydb/core/kqp/common/events/workload_service.h> +#include <ydb/services/workload_manager/events.h> #include <ydb/core/protos/workload_manager_config.pb.h> -namespace NKikimr::NKqp::NWorkload { +namespace NKikimr::NWorkloadManager { // Pool state holder NActors::IActor* CreatePoolHandlerActor(const TString& databaseId, const TString& poolId, const NResourcePool::TPoolSettings& poolConfig, NMonitoring::TDynamicCounterPtr counters); @@ -22,4 +22,4 @@ NActors::IActor* CreateDatabaseFetcherActor(const NActors::TActorId& replyActorI // Cpu load fetcher actor NActors::IActor* CreateCpuLoadFetcherActor(const NActors::TActorId& replyActorId); -} // NKikimr::NKqp::NWorkload +} // NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/actors/cpu_load_actors.cpp b/ydb/services/workload_manager/actors/cpu_load_actors.cpp index 761a4628f21..466e90da0b5 100644 --- a/ydb/core/kqp/workload_service/actors/cpu_load_actors.cpp +++ b/ydb/services/workload_manager/actors/cpu_load_actors.cpp @@ -1,11 +1,11 @@ #include "actors.h" -#include <ydb/core/kqp/workload_service/common/events.h> +#include <ydb/services/workload_manager/common/events.h> #include <ydb/library/query_actor/query_actor.h> -namespace NKikimr::NKqp::NWorkload { +namespace NKikimr::NWorkloadManager { namespace { @@ -74,4 +74,4 @@ IActor* CreateCpuLoadFetcherActor(const TActorId& replyActorId) { return new TQueryRetryActor<TCpuLoadFetcherActor, TEvPrivate::TEvCpuLoadResponse>(replyActorId); } -} // NKikimr::NKqp::NWorkload +} // NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/actors/pool_handlers_actors.cpp b/ydb/services/workload_manager/actors/pool_handlers_actors.cpp index 136c137197c..b3eac3e270e 100644 --- a/ydb/core/kqp/workload_service/actors/pool_handlers_actors.cpp +++ b/ydb/services/workload_manager/actors/pool_handlers_actors.cpp @@ -4,13 +4,15 @@ #include <ydb/core/kqp/common/events/events.h> #include <ydb/core/kqp/common/simple/services.h> -#include <ydb/core/kqp/proxy_service/kqp_session_state.h> -#include <ydb/core/kqp/workload_service/common/events.h> -#include <ydb/core/kqp/workload_service/common/helpers.h> -#include <ydb/core/kqp/workload_service/tables/table_queries.h> +#include <ydb/services/workload_manager/session_updater.h> +#include <ydb/services/workload_manager/common/events.h> +#include <ydb/services/workload_manager/common/helpers.h> +#include <ydb/services/workload_manager/tables/table_queries.h> +#include <ydb/services/workload_manager/service/service.h> +#include <ydb/core/kqp/common/events/events.h> -namespace NKikimr::NKqp::NWorkload { +namespace NKikimr::NWorkloadManager { namespace { @@ -151,7 +153,7 @@ public: // Worker actor events hFunc(TEvCleanupRequest, Handle); - IgnoreFunc(TEvKqp::TEvCancelQueryResponse); + IgnoreFunc(NKqp::TEvKqp::TEvCancelQueryResponse); // Schemeboard events hFunc(TEvTxProxySchemeCache::TEvWatchNotifyUpdated, Handle); @@ -165,7 +167,7 @@ public: } SendPoolInfoUpdate(std::nullopt, std::nullopt, Subscribers); - this->Send(MakeKqpWorkloadServiceId(this->SelfId().NodeId()), new TEvPrivate::TEvStopPoolHandlerResponse(DatabaseId, PoolId)); + this->Send(NWorkloadManager::MakeServiceId(this->SelfId().NodeId()), new TEvPrivate::TEvStopPoolHandlerResponse(DatabaseId, PoolId)); Counters.OnCleanup(ResetCountersOnStrop); @@ -193,7 +195,7 @@ private: const TString& sessionId = event->Get()->SessionId; const TString& requestText = event->Get()->RequestText; auto wmSessionUpdater = event->Get()->WmSessionUpdater; - this->Send(MakeKqpWorkloadServiceId(this->SelfId().NodeId()), new TEvPrivate::TEvPlaceRequestIntoPoolResponse(DatabaseId, PoolId, sessionId)); + this->Send(NWorkloadManager::MakeServiceId(this->SelfId().NodeId()), new TEvPrivate::TEvPlaceRequestIntoPoolResponse(DatabaseId, PoolId, sessionId)); const TActorId& workerActorId = event->Sender; if (!InFlightLimit) { @@ -406,7 +408,7 @@ protected: auto event = std::make_unique<TEvPrivate::TEvFinishRequestInPool>( DatabaseId, PoolId, request->Duration, request->CpuConsumed, request->UsedCpuQuota ); - this->Send(MakeKqpWorkloadServiceId(this->SelfId().NodeId()), event.release()); + this->Send(NWorkloadManager::MakeServiceId(this->SelfId().NodeId()), event.release()); LocalSessions.erase(request->SessionId); if (StopHandler && LocalSessions.empty()) { @@ -456,9 +458,9 @@ private: } void ReplyCancel(const TRequest* request) { - auto ev = std::make_unique<TEvKqp::TEvCancelQueryRequest>(); + auto ev = std::make_unique<NKqp::TEvKqp::TEvCancelQueryRequest>(); ev->Record.MutableRequest()->SetSessionId(request->SessionId); - this->Send(MakeKqpProxyID(this->SelfId().NodeId()), ev.release()); + this->Send(NKqp::MakeKqpProxyID(this->SelfId().NodeId()), ev.release()); Counters.Cancelled->Inc(); Counters.CollectRequestLatency(request->ContinueTime); @@ -497,7 +499,7 @@ private: if (ShouldResign()) { const TActorId& newHandler = this->RegisterWithSameMailbox(CreatePoolHandlerActor(DatabaseId, PoolId, poolConfig, Counters.CountersRoot)); - this->Send(MakeKqpWorkloadServiceId(this->SelfId().NodeId()), new TEvPrivate::TEvResignPoolHandler(DatabaseId, PoolId, newHandler)); + this->Send(NWorkloadManager::MakeServiceId(this->SelfId().NodeId()), new TEvPrivate::TEvResignPoolHandler(DatabaseId, PoolId, newHandler)); } } @@ -664,7 +666,7 @@ protected: } if (!PreparingFinished) { - this->Send(MakeKqpWorkloadServiceId(this->SelfId().NodeId()), new TEvPrivate::TEvPrepareTablesRequest(DatabaseId, PoolId)); + this->Send(NWorkloadManager::MakeServiceId(this->SelfId().NodeId()), new TEvPrivate::TEvPrepareTablesRequest(DatabaseId, PoolId)); } RefreshState(); @@ -682,7 +684,7 @@ protected: void RefreshState(bool refreshRequired = false) override { if (!WaitingNodesInfo && TInstant::Now() - LastNodesInfoRefreshTime > LEASE_DURATION) { WaitingNodesInfo = true; - this->Send(MakeKqpWorkloadServiceId(this->SelfId().NodeId()), new TEvPrivate::TEvNodesInfoRequest()); + this->Send(NWorkloadManager::MakeServiceId(this->SelfId().NodeId()), new TEvPrivate::TEvNodesInfoRequest()); } RefreshRequired |= refreshRequired; @@ -880,7 +882,7 @@ private: auto event = std::make_unique<TEvPrivate::TEvRefreshPoolState>(); event->Record.SetPoolId(PoolId); event->Record.SetDatabase(DatabaseId); - this->Send(MakeKqpWorkloadServiceId(nodeId), std::move(event)); + this->Send(NWorkloadManager::MakeServiceId(nodeId), std::move(event)); RefreshState(); return; } @@ -1004,7 +1006,7 @@ private: } void RequestCpuQuota(double loadCpuThreshold, EStartRequestCase requestCase) const { - this->Send(MakeKqpWorkloadServiceId(this->SelfId().NodeId()), new TEvPrivate::TEvCpuQuotaRequest(loadCpuThreshold / 100.0), 0, static_cast<ui64>(requestCase)); + this->Send(NWorkloadManager::MakeServiceId(this->SelfId().NodeId()), new TEvPrivate::TEvCpuQuotaRequest(loadCpuThreshold / 100.0), 0, static_cast<ui64>(requestCase)); } private: @@ -1094,4 +1096,4 @@ IActor* CreatePoolHandlerActor(const TString& databaseId, const TString& poolId, return new TFifoPoolHandlerActor(databaseId, poolId, poolConfig, counters); } -} // NKikimr::NKqp::NWorkload +} // NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/actors/scheme_actors.cpp b/ydb/services/workload_manager/actors/scheme_actors.cpp index fd00b160ea2..ab94b9a9509 100644 --- a/ydb/core/kqp/workload_service/actors/scheme_actors.cpp +++ b/ydb/services/workload_manager/actors/scheme_actors.cpp @@ -6,9 +6,9 @@ #include <ydb/core/protos/schemeshard/operations.pb.h> #include <ydb/core/protos/workload_manager_config.pb.h> -#include <ydb/core/kqp/common/simple/services.h> -#include <ydb/core/kqp/workload_service/common/events.h> -#include <ydb/core/kqp/workload_service/common/helpers.h> +#include <ydb/services/workload_manager/service/service.h> +#include <ydb/services/workload_manager/common/events.h> +#include <ydb/services/workload_manager/common/helpers.h> #include <ydb/core/tx/schemeshard/schemeshard.h> #include <ydb/core/tx/tx_proxy/proxy.h> @@ -16,7 +16,7 @@ #include <ydb/library/table_creator/table_creator.h> -namespace NKikimr::NKqp::NWorkload { +namespace NKikimr::NWorkloadManager { namespace { @@ -106,14 +106,14 @@ private: void Reply(NResourcePool::TPoolSettings poolConfig, TPathId pathId) { LOG_D("Pool info successfully resolved"); - Send(MakeKqpWorkloadServiceId(SelfId().NodeId()), new TEvPrivate::TEvResolvePoolResponse(Ydb::StatusIds::SUCCESS, poolConfig, pathId, DefaultPoolCreated, std::move(Event))); + Send(NWorkloadManager::MakeServiceId(SelfId().NodeId()), new TEvPrivate::TEvResolvePoolResponse(Ydb::StatusIds::SUCCESS, poolConfig, pathId, DefaultPoolCreated, std::move(Event))); PassAway(); } void Reply(Ydb::StatusIds::StatusCode status, NYql::TIssues issues) { LOG_W("Failed to resolve pool, " << status << ", issues: " << issues.ToOneLineString()); - Send(MakeKqpWorkloadServiceId(SelfId().NodeId()), new TEvPrivate::TEvResolvePoolResponse(status, {}, {}, DefaultPoolCreated, std::move(Event), std::move(issues))); + Send(NWorkloadManager::MakeServiceId(SelfId().NodeId()), new TEvPrivate::TEvResolvePoolResponse(status, {}, {}, DefaultPoolCreated, std::move(Event), std::move(issues))); PassAway(); } @@ -642,4 +642,4 @@ IActor* CreateDatabaseFetcherActor(const TActorId& replyActorId, const TString& return new TDatabaseFetcherActor(replyActorId, database, userToken, checkAccess); } -} // NKikimr::NKqp::NWorkload +} // NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/actors/ya.make b/ydb/services/workload_manager/actors/ya.make index 1cb4c527b5a..b90c9e64e4f 100644 --- a/ydb/core/kqp/workload_service/actors/ya.make +++ b/ydb/services/workload_manager/actors/ya.make @@ -7,8 +7,8 @@ SRCS( ) PEERDIR( - ydb/core/kqp/workload_service/common - ydb/core/kqp/workload_service/tables + ydb/services/workload_manager/common + ydb/services/workload_manager/tables ydb/core/tx/tx_proxy ) diff --git a/ydb/core/kqp/workload_service/common/cpu_quota_manager.cpp b/ydb/services/workload_manager/common/cpu_quota_manager.cpp index dd3a6618342..b8b4b9c8968 100644 --- a/ydb/core/kqp/workload_service/common/cpu_quota_manager.cpp +++ b/ydb/services/workload_manager/common/cpu_quota_manager.cpp @@ -3,7 +3,7 @@ #include <util/string/builder.h> -namespace NKikimr::NKqp::NWorkload { +namespace NKikimr::NWorkloadManager { //// TCpuQuotaManager::TCounters @@ -153,4 +153,4 @@ void TCpuQuotaManager::AdjustCpuQuota(double quota, TDuration duration, double c } } -} // namespace NKikimr::NKqp::NWorkload +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/common/cpu_quota_manager.h b/ydb/services/workload_manager/common/cpu_quota_manager.h index ee9e8c2a246..3fe2dee520d 100644 --- a/ydb/core/kqp/workload_service/common/cpu_quota_manager.h +++ b/ydb/services/workload_manager/common/cpu_quota_manager.h @@ -7,7 +7,7 @@ #include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/types/status_codes.h> -namespace NKikimr::NKqp::NWorkload { +namespace NKikimr::NWorkloadManager { class TCpuQuotaManager { struct TCounters { @@ -73,4 +73,4 @@ private: bool Ready = false; }; -} // namespace NKikimr::NKqp::NWorkload +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/common/events.cpp b/ydb/services/workload_manager/common/events.cpp index ca8d9c554ac..33213f1ef04 100644 --- a/ydb/core/kqp/workload_service/common/events.cpp +++ b/ydb/services/workload_manager/common/events.cpp @@ -1,10 +1,10 @@ #include "events.h" -namespace NKikimr::NKqp::NWorkload { +namespace NKikimr::NWorkloadManager { ui64 TPoolStateDescription::AmountRequests() const { return DelayedRequests + RunningRequests; } -} // NKikimr::NKqp::NWorkload +} // NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/common/events.h b/ydb/services/workload_manager/common/events.h index 911dccd36f0..d91b2fe39cc 100644 --- a/ydb/core/kqp/workload_service/common/events.h +++ b/ydb/services/workload_manager/common/events.h @@ -1,13 +1,13 @@ #pragma once -#include <ydb/core/kqp/common/events/workload_service.h> +#include <ydb/services/workload_manager/events.h> #include <ydb/core/scheme/scheme_pathid.h> #include <ydb/core/protos/kqp.pb.h> -namespace NKikimr::NKqp::NWorkload { +namespace NKikimr::NWorkloadManager { struct TPoolStateDescription { ui64 DelayedRequests = 0; @@ -319,4 +319,4 @@ struct TEvPrivate { }; }; -} // NKikimr::NKqp::NWorkload +} // NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/common/helpers.cpp b/ydb/services/workload_manager/common/helpers.cpp index ec551ef9ebb..147eda4402b 100644 --- a/ydb/core/kqp/workload_service/common/helpers.cpp +++ b/ydb/services/workload_manager/common/helpers.cpp @@ -4,7 +4,7 @@ #include <ydb/core/base/path.h> -namespace NKikimr::NKqp::NWorkload { +namespace NKikimr::NWorkloadManager { TString CreateDatabaseId(const TString& database, bool serverless, TPathId pathId) { TString databasePath = CanonizePath(database); @@ -93,4 +93,4 @@ NResourcePool::TPoolSettings PoolSettingsFromConfig(const NKikimrConfig::TWorklo return poolSettings; } -} // NKikimr::NKqp::NWorkload +} // NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/common/helpers.h b/ydb/services/workload_manager/common/helpers.h index 2ef0e27c320..88dc669872d 100644 --- a/ydb/core/kqp/workload_service/common/helpers.h +++ b/ydb/services/workload_manager/common/helpers.h @@ -2,6 +2,7 @@ #include <library/cpp/retry/retry_policy.h> +#include <ydb/core/base/path.h> #include <ydb/core/protos/workload_manager_config.pb.h> #include <ydb/core/resource_pools/resource_pool_settings.h> #include <ydb/core/tx/scheme_cache/scheme_cache.h> @@ -13,8 +14,15 @@ #include <ydb/public/api/protos/ydb_status_codes.pb.h> +#include <util/string/builder.h> + + +namespace NKikimr::NWorkloadManager { + +inline TString GetPoolKey(const TString& databaseId, const TString& poolId) { + return CanonizePath(TStringBuilder() << databaseId << "/" << poolId); +} -namespace NKikimr::NKqp::NWorkload { #define LOG_T(stream) LOG_TRACE_S(*TlsActivationContext, NKikimrServices::KQP_WORKLOAD_SERVICE, "[WorkloadService] " << LogPrefix() << stream) #define LOG_D(stream) LOG_DEBUG_S(*TlsActivationContext, NKikimrServices::KQP_WORKLOAD_SERVICE, "[WorkloadService] " << LogPrefix() << stream) @@ -111,4 +119,4 @@ ui64 SaturationSub(ui64 x, ui64 y); NResourcePool::TPoolSettings PoolSettingsFromConfig(const NKikimrConfig::TWorkloadManagerConfig& workloadManagerConfig); -} // NKikimr::NKqp::NWorkload +} // NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/common/ya.make b/ydb/services/workload_manager/common/ya.make index a8eb736baaa..1a6a565f35b 100644 --- a/ydb/core/kqp/workload_service/common/ya.make +++ b/ydb/services/workload_manager/common/ya.make @@ -7,11 +7,10 @@ SRCS( ) PEERDIR( - ydb/core/kqp/common/events + ydb/core/base ydb/core/scheme - ydb/core/tx/scheme_cache ydb/library/actors/core diff --git a/ydb/core/kqp/common/events/workload_service.h b/ydb/services/workload_manager/events.h index 5bbdeb5707c..f04a474feea 100644 --- a/ydb/core/kqp/common/events/workload_service.h +++ b/ydb/services/workload_manager/events.h @@ -1,8 +1,8 @@ #pragma once -#include <ydb/core/kqp/common/simple/kqp_event_ids.h> - +#include <ydb/core/base/events.h> #include <ydb/core/resource_pools/resource_pool_settings.h> +#include "session_updater.h" #include <ydb/core/scheme/scheme_pathid.h> #include <ydb/library/aclib/aclib.h> @@ -14,11 +14,22 @@ #include <memory> -namespace NKikimr::NKqp::NWorkload { +namespace NKikimr::NWorkloadManager { + +struct TWorkloadManagerEvents { + enum EEvents { + EvPlaceRequestIntoPool = EventSpaceBegin(TKikimrEvents::ES_WORKLOAD_MANAGER), + EvContinueRequest, + EvCleanupRequest, + EvCleanupResponse, + EvUpdatePoolInfo, + EvSubscribeOnPoolChanges, + EvFetchDatabaseResponse, + }; +}; -class ISessionUpdater; -struct TEvSubscribeOnPoolChanges : public NActors::TEventLocal<TEvSubscribeOnPoolChanges, TKqpWorkloadServiceEvents::EvSubscribeOnPoolChanges> { +struct TEvSubscribeOnPoolChanges : public NActors::TEventLocal<TEvSubscribeOnPoolChanges, TWorkloadManagerEvents::EvSubscribeOnPoolChanges> { TEvSubscribeOnPoolChanges(const TString& databaseId, const TString& poolId) : DatabaseId(databaseId) , PoolId(poolId) @@ -28,7 +39,7 @@ struct TEvSubscribeOnPoolChanges : public NActors::TEventLocal<TEvSubscribeOnPoo const TString PoolId; }; -struct TEvPlaceRequestIntoPool : public NActors::TEventLocal<TEvPlaceRequestIntoPool, TKqpWorkloadServiceEvents::EvPlaceRequestIntoPool> { +struct TEvPlaceRequestIntoPool : public NActors::TEventLocal<TEvPlaceRequestIntoPool, TWorkloadManagerEvents::EvPlaceRequestIntoPool> { TEvPlaceRequestIntoPool(const TString& databaseId, const TString& sessionId, const TString& poolId, TIntrusiveConstPtr<NACLib::TUserToken> userToken, const TString& requestText = "", std::shared_ptr<ISessionUpdater> wmSessionUpdater = nullptr) : DatabaseId(databaseId) , SessionId(sessionId) @@ -46,7 +57,7 @@ struct TEvPlaceRequestIntoPool : public NActors::TEventLocal<TEvPlaceRequestInto std::shared_ptr<ISessionUpdater> WmSessionUpdater; }; -struct TEvContinueRequest : public NActors::TEventLocal<TEvContinueRequest, TKqpWorkloadServiceEvents::EvContinueRequest> { +struct TEvContinueRequest : public NActors::TEventLocal<TEvContinueRequest, TWorkloadManagerEvents::EvContinueRequest> { TEvContinueRequest(Ydb::StatusIds::StatusCode status, const TString& poolId, const NResourcePool::TPoolSettings& poolConfig, NYql::TIssues issues = {}) : Status(status) , PoolId(poolId) @@ -71,7 +82,7 @@ struct TEvContinueRequest : public NActors::TEventLocal<TEvContinueRequest, TKqp const NYql::TIssues Issues; }; -struct TEvCleanupRequest : public NActors::TEventLocal<TEvCleanupRequest, TKqpWorkloadServiceEvents::EvCleanupRequest> { +struct TEvCleanupRequest : public NActors::TEventLocal<TEvCleanupRequest, TWorkloadManagerEvents::EvCleanupRequest> { TEvCleanupRequest(const TString& databaseId, const TString& sessionId, const TString& poolId, TDuration duration, TDuration cpuConsumed) : DatabaseId(databaseId) , SessionId(sessionId) @@ -87,7 +98,7 @@ struct TEvCleanupRequest : public NActors::TEventLocal<TEvCleanupRequest, TKqpWo const TDuration CpuConsumed; }; -struct TEvCleanupResponse : public NActors::TEventLocal<TEvCleanupResponse, TKqpWorkloadServiceEvents::EvCleanupResponse> { +struct TEvCleanupResponse : public NActors::TEventLocal<TEvCleanupResponse, TWorkloadManagerEvents::EvCleanupResponse> { explicit TEvCleanupResponse(Ydb::StatusIds::StatusCode status, NYql::TIssues issues = {}) : Status(status) , Issues(std::move(issues)) @@ -97,7 +108,7 @@ struct TEvCleanupResponse : public NActors::TEventLocal<TEvCleanupResponse, TKqp const NYql::TIssues Issues; }; -struct TEvUpdatePoolInfo : public NActors::TEventLocal<TEvUpdatePoolInfo, TKqpWorkloadServiceEvents::EvUpdatePoolInfo> { +struct TEvUpdatePoolInfo : public NActors::TEventLocal<TEvUpdatePoolInfo, TWorkloadManagerEvents::EvUpdatePoolInfo> { TEvUpdatePoolInfo(const TString& databaseId, const TString& poolId, const std::optional<NResourcePool::TPoolSettings>& config, const std::optional<NACLib::TSecurityObject>& securityObject) : DatabaseId(databaseId) , PoolId(poolId) @@ -111,7 +122,7 @@ struct TEvUpdatePoolInfo : public NActors::TEventLocal<TEvUpdatePoolInfo, TKqpWo const std::optional<NACLib::TSecurityObject> SecurityObject; }; -struct TEvFetchDatabaseResponse : public NActors::TEventLocal<TEvFetchDatabaseResponse, TKqpWorkloadServiceEvents::EvFetchDatabaseResponse> { +struct TEvFetchDatabaseResponse : public NActors::TEventLocal<TEvFetchDatabaseResponse, TWorkloadManagerEvents::EvFetchDatabaseResponse> { TEvFetchDatabaseResponse(Ydb::StatusIds::StatusCode status, const TString& database, const TString& databaseId, bool serverless, TPathId pathId, NYql::TIssues issues) : Status(status) , Database(database) @@ -129,4 +140,4 @@ struct TEvFetchDatabaseResponse : public NActors::TEventLocal<TEvFetchDatabaseRe const NYql::TIssues Issues; }; -} // NKikimr::NKqp::NWorkload +} // NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/kqp_has_full_scan_matcher.cpp b/ydb/services/workload_manager/has_full_scan_matcher.cpp index 85caab8d008..915844085b2 100644 --- a/ydb/core/kqp/workload_service/kqp_has_full_scan_matcher.cpp +++ b/ydb/services/workload_manager/has_full_scan_matcher.cpp @@ -1,9 +1,9 @@ -#include "kqp_has_full_scan_matcher.h" +#include "has_full_scan_matcher.h" #include <ydb/core/base/path.h> -namespace NKikimr::NKqp::NWorkload { +namespace NKikimr::NWorkloadManager { namespace { @@ -78,4 +78,4 @@ bool MatchesFullScan(const std::optional<NResourcePool::TRegexPredicate>& predic return false; } -} // namespace NKikimr::NKqp::NWorkload +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/kqp_has_full_scan_matcher.h b/ydb/services/workload_manager/has_full_scan_matcher.h index 3ae7374b40d..e4976c550e3 100644 --- a/ydb/core/kqp/workload_service/kqp_has_full_scan_matcher.h +++ b/ydb/services/workload_manager/has_full_scan_matcher.h @@ -6,7 +6,7 @@ #include <optional> -namespace NKikimr::NKqp::NWorkload { +namespace NKikimr::NWorkloadManager { /// /// Returns true if the classifier should accept the query based on the HAS_FULL_SCAN filter. @@ -18,4 +18,4 @@ namespace NKikimr::NKqp::NWorkload { bool MatchesFullScan(const std::optional<NResourcePool::TRegexPredicate>& predicate, const NKqpProto::TKqpPhyQuery& phyQuery); -} // namespace NKikimr::NKqp::NWorkload +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/kqp_has_path_matcher.cpp b/ydb/services/workload_manager/has_path_matcher.cpp index fd674228c7c..f6f35735113 100644 --- a/ydb/core/kqp/workload_service/kqp_has_path_matcher.cpp +++ b/ydb/services/workload_manager/has_path_matcher.cpp @@ -1,4 +1,4 @@ -#include "kqp_has_path_matcher.h" +#include "has_path_matcher.h" #include <ydb/core/base/path.h> @@ -10,7 +10,7 @@ #include <util/generic/vector.h> -namespace NKikimr::NKqp::NWorkload { +namespace NKikimr::NWorkloadManager { namespace { @@ -95,4 +95,4 @@ bool MatchesPath(const std::optional<NResourcePool::TRegexPredicate>& predicate, return false; } -} // namespace NKikimr::NKqp::NWorkload +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/kqp_has_path_matcher.h b/ydb/services/workload_manager/has_path_matcher.h index c1a835eba0c..07215e5764d 100644 --- a/ydb/core/kqp/workload_service/kqp_has_path_matcher.h +++ b/ydb/services/workload_manager/has_path_matcher.h @@ -9,7 +9,7 @@ #include <optional> -namespace NKikimr::NKqp::NWorkload { +namespace NKikimr::NWorkloadManager { /// /// Returns true if the classifier should accept the query based on the HAS_PATH filter. @@ -26,4 +26,4 @@ bool MatchesPath(const std::optional<NResourcePool::TRegexPredicate>& predicate, const TVector<TString>& queryTables, const NKqpProto::TKqpPhyQuery& phyQuery); -} // namespace NKikimr::NKqp::NWorkload +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/services/workload_manager/has_stream_matcher.cpp b/ydb/services/workload_manager/has_stream_matcher.cpp new file mode 100644 index 00000000000..dde424bf792 --- /dev/null +++ b/ydb/services/workload_manager/has_stream_matcher.cpp @@ -0,0 +1,14 @@ +#include "has_stream_matcher.h" + + +namespace NKikimr::NWorkloadManager { + +bool MatchesStream(const std::optional<bool>& hasStream, + const NKqp::TUserRequestContext& userRequestContext) { + if (!hasStream) { + return true; + } + return *hasStream ? userRequestContext.IsStreamingQuery : !userRequestContext.IsStreamingQuery; +} + +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/kqp_has_stream_matcher.h b/ydb/services/workload_manager/has_stream_matcher.h index c11e867734c..65e115c19f0 100644 --- a/ydb/core/kqp/workload_service/kqp_has_stream_matcher.h +++ b/ydb/services/workload_manager/has_stream_matcher.h @@ -5,10 +5,10 @@ #include <optional> -namespace NKikimr::NKqp::NWorkload { +namespace NKikimr::NWorkloadManager { /// Returns true if the classifier should accept the query based on the HAS_STREAM filter. bool MatchesStream(const std::optional<bool>& hasStream, - const TUserRequestContext& userRequestContext); + const NKqp::TUserRequestContext& userRequestContext); -} // namespace NKikimr::NKqp::NWorkload +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/gateway/behaviour/resource_pool/behaviour.cpp b/ydb/services/workload_manager/metadata_subscription/behaviour.cpp index a1b172b8741..0854488a194 100644 --- a/ydb/core/kqp/gateway/behaviour/resource_pool/behaviour.cpp +++ b/ydb/services/workload_manager/metadata_subscription/behaviour.cpp @@ -4,7 +4,7 @@ #include <ydb/services/metadata/abstract/initialization.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { TResourcePoolBehaviour::TFactory::TRegistrator<TResourcePoolBehaviour> TResourcePoolBehaviour::Registrator(TResourcePoolConfig::GetTypeId()); @@ -25,4 +25,4 @@ TString TResourcePoolBehaviour::GetTypeId() const { return TResourcePoolConfig::GetTypeId(); } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/gateway/behaviour/resource_pool/behaviour.h b/ydb/services/workload_manager/metadata_subscription/behaviour.h index 01587224e41..24f98a2a0e1 100644 --- a/ydb/core/kqp/gateway/behaviour/resource_pool/behaviour.h +++ b/ydb/services/workload_manager/metadata_subscription/behaviour.h @@ -3,7 +3,7 @@ #include <ydb/services/metadata/abstract/kqp_common.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { class TResourcePoolConfig { public: @@ -24,4 +24,4 @@ public: virtual TString GetTypeId() const override; }; -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/gateway/behaviour/resource_pool/manager.cpp b/ydb/services/workload_manager/metadata_subscription/manager.cpp index bbaa2f8f585..fbd524f446c 100644 --- a/ydb/core/kqp/gateway/behaviour/resource_pool/manager.cpp +++ b/ydb/services/workload_manager/metadata_subscription/manager.cpp @@ -9,9 +9,10 @@ #include <ydb/core/protos/feature_flags.pb.h> #include <ydb/core/protos/schemeshard/operations.pb.h> #include <ydb/core/resource_pools/resource_pool_settings.h> +#include <ydb/core/kqp/gateway/utils/metadata_helpers.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { namespace { @@ -20,7 +21,7 @@ using namespace NResourcePool; using TYqlConclusionStatus = TResourcePoolManager::TYqlConclusionStatus; using TAsyncStatus = TResourcePoolManager::TAsyncStatus; -struct TFeatureFlagExtractor : public IFeatureFlagExtractor { +struct TFeatureFlagExtractor : public NKqp::IFeatureFlagExtractor { bool IsEnabled(const NKikimrConfig::TFeatureFlags& flags) const override { return flags.GetEnableResourcePools(); } @@ -225,13 +226,13 @@ TAsyncStatus TResourcePoolManager::ExecutePrepared(const NKqpProto::TKqpSchemeOp TAsyncStatus TResourcePoolManager::ExecuteSchemeRequest(const NKikimrSchemeOp::TModifyScheme& schemeTx, const TExternalModificationContext& context, ui32 nodeId, NKqpProto::TKqpSchemeOperation::OperationCase operationCase) const { TAsyncStatus validationFuture = NThreading::MakeFuture<TYqlConclusionStatus>(TYqlConclusionStatus::Success()); if (operationCase != NKqpProto::TKqpSchemeOperation::kDropResourcePool) { - validationFuture = ChainFeatures(validationFuture, [context, nodeId] { + validationFuture = NKqp::ChainFeatures(validationFuture, [context, nodeId] { return CheckFeatureFlag(nodeId, MakeIntrusive<TFeatureFlagExtractor>(), context); }); } - return ChainFeatures(validationFuture, [schemeTx, context] { - return SendSchemeRequest(schemeTx, context); + return NKqp::ChainFeatures(validationFuture, [schemeTx, context] { + return NKqp::SendSchemeRequest(schemeTx, context); }); } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/gateway/behaviour/resource_pool/manager.h b/ydb/services/workload_manager/metadata_subscription/manager.h index 0c6ac1bf60e..6fbe30b5038 100644 --- a/ydb/core/kqp/gateway/behaviour/resource_pool/manager.h +++ b/ydb/services/workload_manager/metadata_subscription/manager.h @@ -3,7 +3,7 @@ #include <ydb/services/metadata/manager/abstract.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { class TResourcePoolManager : public NMetadata::NModifications::IOperationsManager { public: @@ -33,4 +33,4 @@ private: TAsyncStatus ExecuteSchemeRequest(const NKikimrSchemeOp::TModifyScheme& schemeTx, const TExternalModificationContext& context, ui32 nodeId, NKqpProto::TKqpSchemeOperation::OperationCase operationCase) const; }; -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/behaviour.cpp b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/behaviour.cpp index aad9f783100..7b3122053d7 100644 --- a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/behaviour.cpp +++ b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/behaviour.cpp @@ -3,7 +3,7 @@ #include "manager.h" -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { TResourcePoolClassifierBehaviour::TFactory::TRegistrator<TResourcePoolClassifierBehaviour> TResourcePoolClassifierBehaviour::Registrator(TResourcePoolClassifierConfig::GetTypeId()); @@ -28,4 +28,4 @@ NMetadata::IClassBehaviour::TPtr TResourcePoolClassifierBehaviour::GetInstance() return result; } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/behaviour.h b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/behaviour.h index 42c8440f80e..580ec31b8a0 100644 --- a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/behaviour.h +++ b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/behaviour.h @@ -6,7 +6,7 @@ #include <ydb/services/metadata/abstract/kqp_common.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { class TResourcePoolClassifierBehaviour : public NMetadata::TClassBehaviour<TResourcePoolClassifierConfig> { static TFactory::TRegistrator<TResourcePoolClassifierBehaviour> Registrator; @@ -22,4 +22,4 @@ public: static IClassBehaviour::TPtr GetInstance(); }; -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/checker.cpp b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/checker.cpp index 3e440fd148b..2e5a932ce38 100644 --- a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/checker.cpp +++ b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/checker.cpp @@ -3,8 +3,8 @@ #include <ydb/core/base/path.h> #include <ydb/core/cms/console/configs_dispatcher.h> -#include <ydb/core/kqp/workload_service/actors/actors.h> -#include <ydb/core/kqp/workload_service/common/events.h> +#include <ydb/services/workload_manager/actors/actors.h> +#include <ydb/services/workload_manager/common/events.h> #include <ydb/core/protos/console_config.pb.h> #include <ydb/core/protos/feature_flags.pb.h> #include <ydb/core/resource_pools/resource_pool_classifier_settings.h> @@ -14,13 +14,12 @@ #include <unordered_set> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { namespace { using namespace NActors; using namespace NResourcePool; -using namespace NWorkload; struct TEvCheckerPrivate { @@ -249,7 +248,7 @@ public: TryFinish(); } - void Handle(NWorkload::TEvPrivate::TEvFetchPoolResponse::TPtr& ev) { + void Handle(TEvPrivate::TEvFetchPoolResponse::TPtr& ev) { if (ev->Get()->Status == Ydb::StatusIds::NOT_FOUND) { FailAndPassAway(TStringBuilder() << "Resource pool '" << ev->Get()->PoolId << "' not found"); return; @@ -273,7 +272,7 @@ public: hFunc(TEvents::TEvUndelivered, Handle); hFunc(NConsole::TEvConfigsDispatcher::TEvGetConfigResponse, Handle); hFunc(NMetadata::NProvider::TEvRefreshSubscriberData, Handle); - hFunc(NWorkload::TEvPrivate::TEvFetchPoolResponse, Handle); + hFunc(TEvPrivate::TEvFetchPoolResponse, Handle); ) private: @@ -386,7 +385,7 @@ private: const TString& databaseId = externalContext.GetDatabaseId(); const auto& workloadManagerConfig = AppData()->WorkloadManagerConfig; for (const auto& poolId : poolIds) { - Register(NWorkload::CreatePoolFetcherActor(SelfId(), databaseId, poolId, userToken, workloadManagerConfig)); + Register(CreatePoolFetcherActor(SelfId(), databaseId, poolId, userToken, workloadManagerConfig)); } } @@ -429,4 +428,4 @@ IActor* CreateResourcePoolClassifierPreparationActor(std::vector<TResourcePoolCl return new TResourcePoolClassifierPreparationActor(std::move(patchedObjects), std::move(controller), context, alterContext); } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/checker.h b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/checker.h index 3018087f755..db7bd0ee391 100644 --- a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/checker.h +++ b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/checker.h @@ -5,8 +5,8 @@ #include <ydb/services/metadata/manager/generic_manager.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { NActors::IActor* CreateResourcePoolClassifierPreparationActor(std::vector<TResourcePoolClassifierConfig>&& patchedObjects, NMetadata::NModifications::IAlterPreparationController<TResourcePoolClassifierConfig>::TPtr controller, const NMetadata::NModifications::IOperationsManager::TInternalModificationContext& context, const NMetadata::NModifications::TAlterOperationContext& alterContext); -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/fetcher.cpp b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/fetcher.cpp index e9f7f3d59bb..23d164e7b85 100644 --- a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/fetcher.cpp +++ b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/fetcher.cpp @@ -1,10 +1,10 @@ #include "fetcher.h" -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { std::vector<NMetadata::IClassBehaviour::TPtr> TResourcePoolClassifierSnapshotsFetcher::DoGetManagers() const { return {TResourcePoolClassifierConfig::GetBehaviour()}; } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/fetcher.h b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/fetcher.h index 29611f0cbeb..9e99bea784b 100644 --- a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/fetcher.h +++ b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/fetcher.h @@ -3,11 +3,11 @@ #include "snapshot.h" -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { class TResourcePoolClassifierSnapshotsFetcher : public NMetadata::NFetcher::TSnapshotsFetcher<TResourcePoolClassifierSnapshot> { protected: virtual std::vector<NMetadata::IClassBehaviour::TPtr> DoGetManagers() const override; }; -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/initializer.cpp b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/initializer.cpp index 39230296e4b..2ba7c4336eb 100644 --- a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/initializer.cpp +++ b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/initializer.cpp @@ -2,7 +2,7 @@ #include "object.h" -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { namespace { @@ -38,4 +38,4 @@ void TResourcePoolClassifierInitializer::DoPrepare(NMetadata::NInitializer::IIni controller->OnPreparationFinished(result); } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/initializer.h b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/initializer.h index c1743cfcbcf..493d15f649f 100644 --- a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/initializer.h +++ b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/initializer.h @@ -3,11 +3,11 @@ #include <ydb/services/metadata/abstract/initialization.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { class TResourcePoolClassifierInitializer : public NMetadata::NInitializer::IInitializationBehaviour { protected: virtual void DoPrepare(NMetadata::NInitializer::IInitializerInput::TPtr controller) const override; }; -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/manager.cpp b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/manager.cpp index 7b0d7c7c882..ee78bd7b199 100644 --- a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/manager.cpp +++ b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/manager.cpp @@ -5,7 +5,7 @@ #include <ydb/core/resource_pools/resource_pool_classifier_settings.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { namespace { @@ -118,4 +118,4 @@ void TResourcePoolClassifierManager::DoPrepareObjectsBeforeModification(std::vec actorSystem->Register(CreateResourcePoolClassifierPreparationActor(std::move(patchedObjects), std::move(controller), context, alterContext)); } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/manager.h b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/manager.h index 947069f8a1a..fba335bb305 100644 --- a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/manager.h +++ b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/manager.h @@ -5,7 +5,7 @@ #include <ydb/services/metadata/manager/generic_manager.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { class TResourcePoolClassifierManager : public NMetadata::NModifications::TGenericOperationsManager<TResourcePoolClassifierConfig> { protected: @@ -18,4 +18,4 @@ private: NMetadata::NModifications::TOperationParsingResult FillDropInfo(const NYql::TObjectSettingsImpl& settings, const TInternalModificationContext& context) const; }; -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/object.cpp b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/object.cpp index 93112f35b3b..6efd775c162 100644 --- a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/object.cpp +++ b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/object.cpp @@ -4,7 +4,7 @@ #include <library/cpp/json/json_reader.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { namespace { @@ -133,4 +133,4 @@ TString TResourcePoolClassifierConfig::GetTypeId() { return "RESOURCE_POOL_CLASSIFIER"; } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/object.h b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/object.h index 3c11d0b9efd..31422f20893 100644 --- a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/object.h +++ b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/object.h @@ -6,7 +6,7 @@ #include <ydb/services/metadata/manager/object.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { class TResourcePoolClassifierConfig : public NMetadata::NModifications::TObject<TResourcePoolClassifierConfig> { using TBase = NMetadata::NModifications::TObject<TResourcePoolClassifierConfig>; @@ -53,4 +53,4 @@ public: static TString GetTypeId(); }; -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/snapshot.cpp b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/snapshot.cpp index 501cf3ec781..40b13d8533e 100644 --- a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/snapshot.cpp +++ b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/snapshot.cpp @@ -1,7 +1,7 @@ #include "snapshot.h" -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { bool TResourcePoolClassifierSnapshot::DoDeserializeFromResultSet(const Ydb::Table::ExecuteQueryResult& rawData) { Y_ABORT_UNLESS(rawData.result_sets().size() == 1); @@ -41,4 +41,4 @@ std::optional<TResourcePoolClassifierConfig> TResourcePoolClassifierSnapshot::Ge return configIt->second; } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/snapshot.h b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/snapshot.h index 8d6eb6406f2..556cc76053d 100644 --- a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/snapshot.h +++ b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/snapshot.h @@ -5,7 +5,7 @@ #include <ydb/services/metadata/abstract/fetcher.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { class TResourcePoolClassifierSnapshot : public NMetadata::NFetcher::ISnapshot { using TBase = NMetadata::NFetcher::ISnapshot; @@ -61,4 +61,4 @@ private: const TConfigsByRankMap* Configs = nullptr; }; -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/ya.make b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/ya.make index 53592996430..056947f4114 100644 --- a/ydb/core/kqp/gateway/behaviour/resource_pool_classifier/ya.make +++ b/ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/ya.make @@ -12,7 +12,7 @@ SRCS( PEERDIR( ydb/core/cms/console - ydb/core/kqp/workload_service/actors + ydb/services/workload_manager/actors ydb/core/protos ydb/core/resource_pools ydb/library/query_actor diff --git a/ydb/core/kqp/gateway/behaviour/resource_pool/ya.make b/ydb/services/workload_manager/metadata_subscription/ya.make index d3333254ed3..d3333254ed3 100644 --- a/ydb/core/kqp/gateway/behaviour/resource_pool/ya.make +++ b/ydb/services/workload_manager/metadata_subscription/ya.make diff --git a/ydb/core/kqp/workload_service/kqp_query_classifier.cpp b/ydb/services/workload_manager/query_classifier.cpp index 88bdf504efa..d5951780686 100644 --- a/ydb/core/kqp/workload_service/kqp_query_classifier.cpp +++ b/ydb/services/workload_manager/query_classifier.cpp @@ -1,18 +1,13 @@ -#include "kqp_has_full_scan_matcher.h" -#include "kqp_has_path_matcher.h" -#include "kqp_has_stream_matcher.h" -#include "kqp_query_classifier.h" -#include "kqp_workload_service.h" +#include "has_full_scan_matcher.h" +#include "has_path_matcher.h" +#include "has_stream_matcher.h" +#include "query_classifier.h" -#include <ydb/core/kqp/gateway/behaviour/resource_pool_classifier/object.h> - -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { inline constexpr char RESOLVER_IS_USER[] = "User request"; inline constexpr char DEFAULT_RESOLVER[] = "Default"; -namespace NWorkload { - class TQueryClassifier : public IQueryClassifier { public: TQueryClassifier(TResourcePoolMapPtr resourcePoolMap, @@ -70,7 +65,7 @@ public: } [[nodiscard]] - TPostCompileClassifyResult PostCompileClassify(const TPreparedQueryHolder& preparedQuery, const TUserRequestContext& userRequestContext) override { + TPostCompileClassifyResult PostCompileClassify(const NKqp::TPreparedQueryHolder& preparedQuery, const NKqp::TUserRequestContext& userRequestContext) override { Y_VALIDATE(ClassifierView, "Post compile classify without configuration"); Y_VALIDATE(PreClassifyResult.has_value() && std::holds_alternative<TPendingCompilation>(*PreClassifyResult), "Post compile classify requires TPendingCompilation from pre-classification"); @@ -189,7 +184,7 @@ private: /// - Involves SQL analysis, plan building, or computations. /// - Depends on actual query structure and execution characteristics. /// - bool MatchesDynamic(const NResourcePool::TClassifierSettings& settings, const TPreparedQueryHolder& preparedQuery, const TUserRequestContext& userRequestContext) const { + bool MatchesDynamic(const NResourcePool::TClassifierSettings& settings, const NKqp::TPreparedQueryHolder& preparedQuery, const NKqp::TUserRequestContext& userRequestContext) const { return MatchesFullScan(settings.HasFullScan, preparedQuery.GetPhysicalQuery()) && MatchesPath(settings.HasPath, preparedQuery.GetQueryTables(), preparedQuery.GetPhysicalQuery()) && MatchesStream(settings.HasStream, userRequestContext); @@ -300,6 +295,4 @@ std::shared_ptr<IQueryClassifier> CreateQueryClassifier(TResourcePoolMapPtr reso return std::make_shared<TQueryClassifier>(std::move(resourcePoolMap), std::move(classifierView), databaseId, std::move(context)); } -} // namespace NWorkload - -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/kqp_query_classifier.h b/ydb/services/workload_manager/query_classifier.h index 808a30f45c2..79d35d5f87e 100644 --- a/ydb/core/kqp/workload_service/kqp_query_classifier.h +++ b/ydb/services/workload_manager/query_classifier.h @@ -1,17 +1,18 @@ #pragma once #include <ydb/core/kqp/common/simple/helpers.h> -#include <ydb/core/kqp/gateway/behaviour/resource_pool_classifier/snapshot.h> #include <ydb/core/kqp/query_data/kqp_prepared_query.h> -#include <ydb/core/kqp/workload_service/kqp_workload_service.h> #include <ydb/core/resource_pools/resource_pool_classifier_settings.h> +#include <ydb/services/workload_manager/common/helpers.h> +#include <ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/snapshot.h> #include <ydb/library/aclib/aclib.h> -#include <util/string/builder.h> namespace NKikimr::NKqp { - struct TUserRequestContext; +} // namespace NKikimr::NKqp + +namespace NKikimr::NWorkloadManager { struct TClassifyContext { const TString PoolId; @@ -27,12 +28,6 @@ struct TResourcePoolEntry { using TResourcePoolMap = std::unordered_map<TString, TResourcePoolEntry>; using TResourcePoolMapPtr = std::shared_ptr<const TResourcePoolMap>; -inline TString GetPoolKey(const TString& databaseId, const TString& poolId) { - return TStringBuilder() << databaseId << "/" << poolId; -} - -namespace NWorkload { - /// /// Manages per-query workload manager policies /// @@ -79,7 +74,7 @@ public: /// Refines classification once the query plan is available [[nodiscard]] - virtual TPostCompileClassifyResult PostCompileClassify(const TPreparedQueryHolder& preparedQuery, const TUserRequestContext& userRequestContext) = 0; + virtual TPostCompileClassifyResult PostCompileClassify(const NKqp::TPreparedQueryHolder& preparedQuery, const NKqp::TUserRequestContext& userRequestContext) = 0; /// Get the current classification state [[nodiscard]] @@ -91,6 +86,4 @@ std::shared_ptr<IQueryClassifier> CreateQueryClassifier(TResourcePoolMapPtr reso const TString& databaseId, TClassifyContext context); -} // namespace NWorkload - -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/kqp_workload_service.cpp b/ydb/services/workload_manager/service/service.cpp index 0af61f06b83..c83334309a7 100644 --- a/ydb/core/kqp/workload_service/kqp_workload_service.cpp +++ b/ydb/services/workload_manager/service/service.cpp @@ -1,16 +1,17 @@ -#include "kqp_workload_service.h" -#include "kqp_workload_service_impl.h" +#include "service.h" +#include "workload_service_impl.h" #include <ydb/core/base/appdata_fwd.h> +#include <ydb/core/base/counters.h> #include <ydb/core/base/feature_flags.h> #include <ydb/core/base/path.h> #include <ydb/core/cms/console/configs_dispatcher.h> #include <ydb/core/cms/console/console.h> -#include <ydb/core/kqp/workload_service/actors/actors.h> -#include <ydb/core/kqp/workload_service/common/helpers.h> -#include <ydb/core/kqp/workload_service/tables/table_queries.h> +#include <ydb/services/workload_manager/actors/actors.h> +#include <ydb/services/workload_manager/common/helpers.h> +#include <ydb/services/workload_manager/tables/table_queries.h> #include <ydb/core/mind/tenant_node_enumeration.h> @@ -21,16 +22,14 @@ #include <ydb/library/actors/interconnect/interconnect.h> -namespace NKikimr::NKqp { - -namespace NWorkload { +namespace NKikimr::NWorkloadManager { namespace { using namespace NActors; -class TKqpWorkloadService : public TActorBootstrapped<TKqpWorkloadService> { +class TWorkloadService : public TActorBootstrapped<TWorkloadService> { struct TCounters { const NMonitoring::TDynamicCounterPtr Counters; @@ -62,12 +61,12 @@ class TKqpWorkloadService : public TActorBootstrapped<TKqpWorkloadService> { }; public: - explicit TKqpWorkloadService(NMonitoring::TDynamicCounterPtr counters) + explicit TWorkloadService(NMonitoring::TDynamicCounterPtr counters) : Counters(counters) {} void Bootstrap() { - Become(&TKqpWorkloadService::MainState); + Become(&TWorkloadService::MainState); // Subscribe for FeatureFlags and WorkloadManagerConfig Send(NConsole::MakeConfigsDispatcherID(SelfId().NodeId()), new NConsole::TEvConfigsDispatcher::TEvSetConfigSubscriptionRequest({ @@ -693,10 +692,6 @@ private: return nullptr; } - static TString GetPoolKey(const TString& databaseId, const TString& poolId) { - return CanonizePath(TStringBuilder() << databaseId << "/" << poolId); - } - TString LogPrefix() const { return "[Service] "; } @@ -722,10 +717,18 @@ private: } // anonymous namespace -} // namespace NWorkload +IActor* CreateService(NMonitoring::TDynamicCounterPtr counters) { + return new NWorkloadManager::TWorkloadService(counters); +} -IActor* CreateKqpWorkloadService(NMonitoring::TDynamicCounterPtr counters) { - return new NWorkload::TKqpWorkloadService(counters); +NMonitoring::TDynamicCounterPtr GetWorkloadManagerCounters(NMonitoring::TDynamicCounterPtr rootCounters) { + return GetServiceCounters(rootCounters, "kqp")->GetSubgroup("subsystem", "workload_manager"); } -} // namespace NKikimr::NKqp +NActors::TActorId MakeServiceId(ui32 nodeId) { + const char name[12] = "kqp_workld"; + return NActors::TActorId(nodeId, TStringBuf(name, 12)); +} + + +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/services/workload_manager/service/service.h b/ydb/services/workload_manager/service/service.h new file mode 100644 index 00000000000..fc8f2a307d4 --- /dev/null +++ b/ydb/services/workload_manager/service/service.h @@ -0,0 +1,16 @@ +#pragma once + +#include <ydb/core/resource_pools/resource_pool_settings.h> + +#include <ydb/library/actors/core/actor.h> +#include <library/cpp/monlib/dynamic_counters/counters.h> + +namespace NKikimr::NWorkloadManager { + +NActors::TActorId MakeServiceId(ui32 nodeId); + +NMonitoring::TDynamicCounterPtr GetWorkloadManagerCounters(NMonitoring::TDynamicCounterPtr rootCounters); + +NActors::IActor* CreateService(NMonitoring::TDynamicCounterPtr counters); + +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/kqp_workload_service_impl.h b/ydb/services/workload_manager/service/workload_service_impl.h index 72856267bde..51b9d2a2835 100644 --- a/ydb/core/kqp/workload_service/kqp_workload_service_impl.h +++ b/ydb/services/workload_manager/service/workload_service_impl.h @@ -4,13 +4,13 @@ #include <ydb/core/kqp/common/events/events.h> #include <ydb/core/kqp/common/simple/services.h> -#include <ydb/core/kqp/workload_service/actors/actors.h> -#include <ydb/core/kqp/workload_service/common/cpu_quota_manager.h> -#include <ydb/core/kqp/workload_service/common/events.h> -#include <ydb/core/kqp/workload_service/common/helpers.h> +#include <ydb/services/workload_manager/actors/actors.h> +#include <ydb/services/workload_manager/common/cpu_quota_manager.h> +#include <ydb/services/workload_manager/common/events.h> +#include <ydb/services/workload_manager/common/helpers.h> -namespace NKikimr::NKqp::NWorkload { +namespace NKikimr::NWorkloadManager { constexpr TDuration IDLE_DURATION = TDuration::Seconds(60); @@ -81,7 +81,7 @@ struct TDatabaseState { } if (Serverless != ev->Get()->Serverless) { - TActivationContext::Send(MakeKqpProxyID(SelfId.NodeId()), std::make_unique<TEvKqp::TEvUpdateDatabaseInfo>(ev->Get()->Database, ev->Get()->DatabaseId, ev->Get()->Serverless)); + TActivationContext::Send(NKqp::MakeKqpProxyID(SelfId.NodeId()), std::make_unique<NKqp::TEvKqp::TEvUpdateDatabaseInfo>(ev->Get()->Database, ev->Get()->DatabaseId, ev->Get()->Serverless)); } LastUpdateTime = TInstant::Now(); @@ -249,4 +249,4 @@ private: std::unordered_map<TActorId, double> HandlersLimits; }; -} // namespace NKikimr::NKqp::NWorkload +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/services/workload_manager/service/ya.make b/ydb/services/workload_manager/service/ya.make new file mode 100644 index 00000000000..959821993e6 --- /dev/null +++ b/ydb/services/workload_manager/service/ya.make @@ -0,0 +1,22 @@ +LIBRARY() + +SRCS( + service.cpp +) + +PEERDIR( + ydb/core/base + ydb/core/cms/console + ydb/core/mind + ydb/core/resource_pools + ydb/library/aclib + ydb/library/actors/interconnect + ydb/services/workload_manager/actors + ydb/services/workload_manager/common + ydb/services/workload_manager/tables +) + +YQL_LAST_ABI_VERSION() + +END() + diff --git a/ydb/core/kqp/proxy_service/kqp_session_state.h b/ydb/services/workload_manager/session_updater.h index bdb165a7b69..5c62b94b15c 100644 --- a/ydb/core/kqp/proxy_service/kqp_session_state.h +++ b/ydb/services/workload_manager/session_updater.h @@ -3,7 +3,7 @@ #include <util/datetime/base.h> #include <util/generic/string.h> -namespace NKikimr::NKqp::NWorkload { +namespace NKikimr::NWorkloadManager { /// /// Interface for updating the execution state of a KQP session within the Workload Manager (WM). @@ -24,4 +24,4 @@ public: virtual void SetPoolId(TString poolId) = 0; }; -} // namespace NKikimr::NKqp::NWorkload +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/tables/table_queries.cpp b/ydb/services/workload_manager/tables/table_queries.cpp index 0c6f9780e3d..4365bf8fab7 100644 --- a/ydb/core/kqp/workload_service/tables/table_queries.cpp +++ b/ydb/services/workload_manager/tables/table_queries.cpp @@ -3,15 +3,16 @@ #include <ydb/core/base/path.h> #include <ydb/core/kqp/common/simple/services.h> -#include <ydb/core/kqp/workload_service/common/events.h> -#include <ydb/core/kqp/workload_service/common/helpers.h> +#include <ydb/services/workload_manager/common/events.h> +#include <ydb/services/workload_manager/common/helpers.h> +#include <ydb/services/workload_manager/service/service.h> #include <ydb/library/query_actor/query_actor.h> #include <ydb/library/table_creator/table_creator.h> -namespace NKikimr::NKqp::NWorkload { +namespace NKikimr::NWorkloadManager { namespace { @@ -289,7 +290,7 @@ private: LOG_E("Failed with issues: " << Issues.ToOneLineString()); } - Send(MakeKqpWorkloadServiceId(SelfId().NodeId()), new TEvPrivate::TEvCleanupTablesFinished(Success, TablesExists, std::move(Issues))); + Send(NWorkloadManager::MakeServiceId(SelfId().NodeId()), new TEvPrivate::TEvCleanupTablesFinished(Success, TablesExists, std::move(Issues))); PassAway(); } @@ -862,4 +863,4 @@ IActor* CreateCleanupRequestsActor(const TActorId& replyActorId, const TString& return TCleanupRequestsQuery::MakeRetry(replyActorId, databaseId, poolId, sessionIds, counters); } -} // NKikimr::NKqp::NWorkload +} // NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/tables/table_queries.h b/ydb/services/workload_manager/tables/table_queries.h index 599b6ba08a9..dc2cbfbc258 100644 --- a/ydb/core/kqp/workload_service/tables/table_queries.h +++ b/ydb/services/workload_manager/tables/table_queries.h @@ -3,7 +3,7 @@ #include <ydb/library/actors/core/actor.h> -namespace NKikimr::NKqp::NWorkload { +namespace NKikimr::NWorkloadManager { // Creates all needed tables for workload service NActors::IActor* CreateTablesCreator(); @@ -19,4 +19,4 @@ NActors::IActor* CreateDelayRequestActor(const NActors::TActorId& replyActorId, NActors::IActor* CreateStartRequestActor(const NActors::TActorId& replyActorId, const TString& databaseId, const TString& poolId, const std::optional<TString>& sessionId, TDuration leaseDuration, NMonitoring::TDynamicCounterPtr counters); NActors::IActor* CreateCleanupRequestsActor(const NActors::TActorId& replyActorId, const TString& databaseId, const TString& poolId, const std::vector<TString>& sessionIds, NMonitoring::TDynamicCounterPtr counters); -} // NKikimr::NKqp::NWorkload +} // NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/tables/ya.make b/ydb/services/workload_manager/tables/ya.make index 1f82c244588..9f7caa93fab 100644 --- a/ydb/core/kqp/workload_service/tables/ya.make +++ b/ydb/services/workload_manager/tables/ya.make @@ -5,7 +5,7 @@ SRCS( ) PEERDIR( - ydb/core/kqp/workload_service/common + ydb/services/workload_manager/common ydb/library/query_actor diff --git a/ydb/core/kqp/workload_service/ut/kqp_action_reject_ut.cpp b/ydb/services/workload_manager/ut/action_reject_ut.cpp index 1a0b8ab264b..818d75b9ff9 100644 --- a/ydb/core/kqp/workload_service/ut/kqp_action_reject_ut.cpp +++ b/ydb/services/workload_manager/ut/action_reject_ut.cpp @@ -1,13 +1,13 @@ #include <ydb/core/base/path.h> -#include <ydb/core/kqp/workload_service/ut/common/kqp_query_classifier_ut_common.h> -#include <ydb/core/kqp/workload_service/ut/common/kqp_workload_service_ut_common.h> +#include <ydb/services/workload_manager/ut/common/query_classifier_ut_common.h> +#include <ydb/services/workload_manager/ut/common/workload_service_ut_common.h> #include <library/cpp/testing/unittest/registar.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { -using namespace NWorkload; +using namespace NWorkloadManager; using namespace NYdb; namespace { @@ -225,4 +225,4 @@ Y_UNIT_TEST_SUITE(ActionRejectDdl) { } } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/ut/common/kqp_query_classifier_ut_common.h b/ydb/services/workload_manager/ut/common/query_classifier_ut_common.h index df4ff58edd0..12058b2ea72 100644 --- a/ydb/core/kqp/workload_service/ut/common/kqp_query_classifier_ut_common.h +++ b/ydb/services/workload_manager/ut/common/query_classifier_ut_common.h @@ -2,8 +2,8 @@ #include <ydb/core/base/path.h> #include <ydb/core/kqp/common/kqp_user_request_context.h> -#include <ydb/core/kqp/gateway/behaviour/resource_pool_classifier/snapshot.h> -#include <ydb/core/kqp/workload_service/kqp_query_classifier.h> +#include <ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/snapshot.h> +#include <ydb/services/workload_manager/query_classifier.h> #include <ydb/core/kqp/query_data/kqp_prepared_query.h> #include <ydb/library/aclib/aclib.h> @@ -14,7 +14,7 @@ #include <vector> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { inline constexpr char TEST_DB[] = "/Root/testdb"; @@ -120,7 +120,7 @@ struct TClassifyTestCase { }; std::vector<TExtraClassifier> ExtraClassifiers; - std::shared_ptr<NWorkload::IQueryClassifier> BuildClassifier() const { + std::shared_ptr<NWorkloadManager::IQueryClassifier> BuildClassifier() const { std::vector<TResourcePoolClassifierConfig> configs; configs.push_back(MakeClassifierConfig( TEST_DB, "c_main", Rank, ResourcePool, @@ -157,10 +157,10 @@ struct TClassifyTestCase { .UserToken = token, }; - return NWorkload::CreateQueryClassifier(poolSnap, TClassifierConfigsView(classifierSnap, TEST_DB), TEST_DB, std::move(ctx)); + return NWorkloadManager::CreateQueryClassifier(poolSnap, TClassifierConfigsView(classifierSnap, TEST_DB), TEST_DB, std::move(ctx)); } - NWorkload::IQueryClassifier::TPreCompileClassifyResult RunPreClassify() const { + NWorkloadManager::IQueryClassifier::TPreCompileClassifyResult RunPreClassify() const { auto classifier = BuildClassifier(); return classifier->PreCompileClassify(); } @@ -170,7 +170,7 @@ struct TClassifyTestCase { /// TKqpPhyQuery holding one table op on `queryTablePath` that either performs /// a full scan (`isFullScan=true`) or a bounded read. /// - NWorkload::IQueryClassifier::TPostCompileClassifyResult RunPostClassify( + NWorkloadManager::IQueryClassifier::TPostCompileClassifyResult RunPostClassify( const TString& queryTablePath, bool isFullScan) const { auto classifier = BuildClassifier(); @@ -189,8 +189,8 @@ struct TClassifyTestCase { op->MutableReadRange()->MutableKeyRange()->MutableFrom()->AddValues(); } - TPreparedQueryHolder holder(proto.release(), nullptr, /*noFillTables=*/true); - TUserRequestContext userRequestContext; + NKqp::TPreparedQueryHolder holder(proto.release(), nullptr, /*noFillTables=*/true); + NKqp::TUserRequestContext userRequestContext; return classifier->PostCompileClassify(holder, userRequestContext); } @@ -200,7 +200,7 @@ struct TClassifyTestCase { /// Exercises HAS_PATH's (B) walk over `tx.GetTables()`. One shape is enough /// for wiring verification; the matcher UT covers the full walk surface. /// - NWorkload::IQueryClassifier::TPostCompileClassifyResult RunPostClassifyForPath( + NWorkloadManager::IQueryClassifier::TPostCompileClassifyResult RunPostClassifyForPath( const TString& queryTablePath) const { auto classifier = BuildClassifier(); @@ -211,30 +211,30 @@ struct TClassifyTestCase { auto* tx = phyQuery->AddTransactions(); tx->AddTables()->MutableId()->SetPath(queryTablePath); - TPreparedQueryHolder holder(proto.release(), nullptr, /*noFillTables=*/true); - TUserRequestContext userRequestContext; + NKqp::TPreparedQueryHolder holder(proto.release(), nullptr, /*noFillTables=*/true); + NKqp::TUserRequestContext userRequestContext; return classifier->PostCompileClassify(holder, userRequestContext); } - NWorkload::IQueryClassifier::TPostCompileClassifyResult RunPostClassifyForStream( + NWorkloadManager::IQueryClassifier::TPostCompileClassifyResult RunPostClassifyForStream( bool isStreamingQuery) const { auto classifier = BuildClassifier(); (void)classifier->PreCompileClassify(); auto proto = std::make_unique<NKikimrKqp::TPreparedQuery>(); - TPreparedQueryHolder holder(proto.release(), nullptr, /*noFillTables=*/true); - TUserRequestContext userRequestContext; + NKqp::TPreparedQueryHolder holder(proto.release(), nullptr, /*noFillTables=*/true); + NKqp::TUserRequestContext userRequestContext; userRequestContext.IsStreamingQuery = isStreamingQuery; return classifier->PostCompileClassify(holder, userRequestContext); } }; -inline TString GetPoolId(const NWorkload::IQueryClassifier::TPreCompileClassifyResult& result) { - UNIT_ASSERT_C(std::holds_alternative<NWorkload::IQueryClassifier::TResolvedPoolId>(result), +inline TString GetPoolId(const NWorkloadManager::IQueryClassifier::TPreCompileClassifyResult& result) { + UNIT_ASSERT_C(std::holds_alternative<NWorkloadManager::IQueryClassifier::TResolvedPoolId>(result), TStringBuilder() << "Expected TResolvedPoolId, got variant with index: " << result.index() ); - return std::get<NWorkload::IQueryClassifier::TResolvedPoolId>(result).PoolId; + return std::get<NWorkloadManager::IQueryClassifier::TResolvedPoolId>(result).PoolId; } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/ut/common/kqp_workload_service_ut_common.cpp b/ydb/services/workload_manager/ut/common/workload_service_ut_common.cpp index db35abb6bb3..c814295509e 100644 --- a/ydb/core/kqp/workload_service/ut/common/kqp_workload_service_ut_common.cpp +++ b/ydb/services/workload_manager/ut/common/workload_service_ut_common.cpp @@ -1,14 +1,14 @@ -#include "kqp_workload_service_ut_common.h" +#include "workload_service_ut_common.h" #include <ydb/core/base/backtrace.h> #include <ydb/core/base/counters.h> #include <ydb/core/kqp/common/events/events.h> -#include <ydb/core/kqp/common/simple/services.h> +#include <ydb/services/workload_manager/service/service.h> #include <ydb/core/kqp/executer_actor/kqp_executer.h> -#include <ydb/core/kqp/gateway/behaviour/resource_pool_classifier/fetcher.h> -#include <ydb/core/kqp/workload_service/actors/actors.h> -#include <ydb/core/kqp/workload_service/tables/table_queries.h> +#include <ydb/services/workload_manager/metadata_subscription/resource_pool_classifier/fetcher.h> +#include <ydb/services/workload_manager/actors/actors.h> +#include <ydb/services/workload_manager/tables/table_queries.h> #include <ydb/core/kqp/ut/common/kqp_ut_common.h> #include <ydb/core/node_whiteboard/node_whiteboard.h> @@ -16,7 +16,7 @@ #include <ydb/library/yql/utils/actor_log/log.h> #include <ydb/services/metadata/service.h> -namespace NKikimr::NKqp::NWorkload { +namespace NKikimr::NWorkloadManager { namespace { @@ -63,7 +63,7 @@ class TQueryRunnerActor : public TActorBootstrapped<TQueryRunnerActor> { using TBase = TActorBootstrapped<TQueryRunnerActor>; public: - TQueryRunnerActor(std::unique_ptr<TEvKqp::TEvQueryRequest> request, TPromise<TQueryRunnerResult> promise, const TQueryRunnerSettings& settings, ui32 targetNodeId) + TQueryRunnerActor(std::unique_ptr<NKqp::TEvKqp::TEvQueryRequest> request, TPromise<TQueryRunnerResult> promise, const TQueryRunnerSettings& settings, ui32 targetNodeId) : Request_(std::move(request)) , Promise_(promise) , Settings_(settings) @@ -77,19 +77,19 @@ public: void Bootstrap() { ActorIdToProto(SelfId(), Request_->Record.MutableRequestActorId()); - Send(MakeKqpProxyID(TargetNodeId_), std::move(Request_)); + Send(NKqp::MakeKqpProxyID(TargetNodeId_), std::move(Request_)); Become(&TQueryRunnerActor::StateFunc); } - void Handle(TEvKqpExecuter::TEvStreamData::TPtr& ev) { + void Handle(NKqp::TEvKqpExecuter::TEvStreamData::TPtr& ev) { UNIT_ASSERT_C(Settings_.ExecutionExpected_ || ExecutionContinued, "Unexpected stream data, execution is not expected and was not continued"); if (!ExecutionStartReported) { ExecutionStartReported = true; SendNotification<TEvQueryRunner::TEvExecutionStarted>(); } - auto response = std::make_unique<TEvKqpExecuter::TEvStreamDataAck>(ev->Get()->Record.GetSeqNo(), ev->Get()->Record.GetChannelId()); + auto response = std::make_unique<NKqp::TEvKqpExecuter::TEvStreamDataAck>(ev->Get()->Record.GetSeqNo(), ev->Get()->Record.GetChannelId()); response->Record.SetFreeSpace(std::numeric_limits<i64>::max()); auto resultSetIndex = ev->Get()->Record.GetQueryResultIndex(); @@ -109,7 +109,7 @@ public: } } - void Handle(TEvKqp::TEvQueryResponse::TPtr& ev) { + void Handle(NKqp::TEvKqp::TEvQueryResponse::TPtr& ev) { SendNotification<TEvQueryRunner::TEvExecutionFinished>(); Result_.Response = ev->Get()->Record; @@ -129,8 +129,8 @@ public: } STRICT_STFUNC(StateFunc, - hFunc(TEvKqpExecuter::TEvStreamData, Handle); - hFunc(TEvKqp::TEvQueryResponse, Handle); + hFunc(NKqp::TEvKqpExecuter::TEvStreamData, Handle); + hFunc(NKqp::TEvKqp::TEvQueryResponse, Handle); hFunc(TEvQueryRunner::TEvContinueExecution, Handle); ) @@ -144,7 +144,7 @@ private: } private: - std::unique_ptr<TEvKqp::TEvQueryRequest> Request_; + std::unique_ptr<NKqp::TEvKqp::TEvQueryRequest> Request_; TPromise<TQueryRunnerResult> Promise_; const TQueryRunnerSettings Settings_; const ui32 TargetNodeId_; @@ -662,14 +662,14 @@ public: UNIT_ASSERT_C(response, "Timed out waiting for resource pool classifier snapshot refresh"); runtime->Send( - MakeKqpProxyID(nodeId), + NKqp::MakeKqpProxyID(nodeId), edgeActor, new NMetadata::NProvider::TEvRefreshSubscriberData(response->Get()->GetSnapshot()), nodeIndex); } void StopWorkloadService(ui64 nodeIndex = 0) const override { - GetRuntime()->Send(MakeKqpWorkloadServiceId(GetRuntime()->GetNodeId(nodeIndex)), GetRuntime()->AllocateEdgeActor(), new TEvents::TEvPoison()); + GetRuntime()->Send(MakeServiceId(GetRuntime()->GetNodeId(nodeIndex)), GetRuntime()->AllocateEdgeActor(), new TEvents::TEvPoison()); Sleep(TDuration::Seconds(1)); } @@ -759,10 +759,10 @@ private: } } - std::unique_ptr<TEvKqp::TEvQueryRequest> GetQueryRequest(const TString& query, const TQueryRunnerSettings& settings) const { + std::unique_ptr<NKqp::TEvKqp::TEvQueryRequest> GetQueryRequest(const TString& query, const TQueryRunnerSettings& settings) const { UNIT_ASSERT_C(settings.PoolId_, "Query pool id is not specified"); - auto event = std::make_unique<TEvKqp::TEvQueryRequest>(); + auto event = std::make_unique<NKqp::TEvKqp::TEvQueryRequest>(); event->Record.SetUserToken(NACLib::TUserToken("", settings.UserSID_, settings.GroupSIDs_).SerializeAsString()); auto request = event->Record.MutableRequest(); @@ -779,8 +779,7 @@ private: } NMonitoring::TDynamicCounterPtr GetWorkloadManagerCounters(ui32 nodeIndex) const { - return GetServiceCounters(GetRuntime()->GetAppData(nodeIndex).Counters, "kqp") - ->GetSubgroup("subsystem", "workload_manager"); + return NWorkloadManager::GetWorkloadManagerCounters(GetRuntime()->GetAppData(nodeIndex).Counters); } static void CheckCommonCounters(NMonitoring::TDynamicCounterPtr subgroup, const TString& description) { @@ -937,10 +936,5 @@ void WaitForClassifierSuccess(TIntrusivePtr<IYdbSetup> ydb, const TQueryRunnerSe WaitForClassifierSuccess(ydb, TSampleQueries::TSelect42::Query, settings); } -//// TSampleQueriess -void TSampleQueries::CompareYson(const TString& expected, const TString& actual) { - NKqp::CompareYson(expected, actual); -} - -} // NKikimr::NKqp::NWorkload +} // NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/ut/common/kqp_workload_service_ut_common.h b/ydb/services/workload_manager/ut/common/workload_service_ut_common.h index 1b6b2d020bd..a5a1a5571b4 100644 --- a/ydb/core/kqp/workload_service/ut/common/kqp_workload_service_ut_common.h +++ b/ydb/services/workload_manager/ut/common/workload_service_ut_common.h @@ -1,7 +1,7 @@ #include <library/cpp/testing/unittest/registar.h> -#include <ydb/core/kqp/workload_service/common/events.h> - +#include <ydb/services/workload_manager/common/events.h> +#include <ydb/core/kqp/ut/common/kqp_ut_common.h> #include <ydb/core/testlib/actors/test_runtime.h> #include <ydb/core/testlib/test_client.h> @@ -16,7 +16,7 @@ #include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/table/table.h> -namespace NKikimr::NKqp::NWorkload { +namespace NKikimr::NWorkloadManager { inline constexpr TDuration FUTURE_WAIT_TIMEOUT = TDuration::Seconds(60); @@ -195,12 +195,10 @@ struct TSampleQueries { static void CheckResult(const TResult& result) { CheckSuccess(result); UNIT_ASSERT_VALUES_EQUAL_C(result.GetResultSets().size(), 1, "Unexpected result set size"); - CompareYson("[[42]]", NYdb::FormatResultSetYson(result.GetResultSet(0))); + NKqp::CompareYson("[[42]]", NYdb::FormatResultSetYson(result.GetResultSet(0))); } }; -private: - static void CompareYson(const TString& expected, const TString& actual); }; -} // namespace NKikimr::NKqp::NWorkload +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/services/workload_manager/ut/common/ya.make b/ydb/services/workload_manager/ut/common/ya.make new file mode 100644 index 00000000000..d0c96fcd48c --- /dev/null +++ b/ydb/services/workload_manager/ut/common/ya.make @@ -0,0 +1,16 @@ +LIBRARY() + +SRCS( + query_classifier_ut_common.h + workload_service_ut_common.cpp +) + +PEERDIR( + ydb/core/kqp/ut/common + ydb/services/metadata + ydb/services/workload_manager/metadata_subscription/resource_pool_classifier +) + +YQL_LAST_ABI_VERSION() + +END() diff --git a/ydb/core/kqp/workload_service/ut/kqp_has_app_name_ut.cpp b/ydb/services/workload_manager/ut/has_app_name_ut.cpp index 58298f71107..6773b3ca83a 100644 --- a/ydb/core/kqp/workload_service/ut/kqp_has_app_name_ut.cpp +++ b/ydb/services/workload_manager/ut/has_app_name_ut.cpp @@ -1,13 +1,13 @@ #include <ydb/core/base/path.h> -#include <ydb/core/kqp/workload_service/ut/common/kqp_query_classifier_ut_common.h> -#include <ydb/core/kqp/workload_service/ut/common/kqp_workload_service_ut_common.h> +#include <ydb/services/workload_manager/ut/common/query_classifier_ut_common.h> +#include <ydb/services/workload_manager/ut/common/workload_service_ut_common.h> #include <library/cpp/testing/unittest/registar.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { -using namespace NWorkload; +using namespace NWorkloadManager; using namespace NYdb; namespace { @@ -129,4 +129,4 @@ Y_UNIT_TEST_SUITE(HasAppNameDdl) { } } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/ut/kqp_has_full_scan_matcher_ut.cpp b/ydb/services/workload_manager/ut/has_full_scan_matcher_ut.cpp index 28c135ff3e3..77e9269f7c6 100644 --- a/ydb/core/kqp/workload_service/ut/kqp_has_full_scan_matcher_ut.cpp +++ b/ydb/services/workload_manager/ut/has_full_scan_matcher_ut.cpp @@ -1,10 +1,10 @@ -#include <ydb/core/kqp/workload_service/kqp_has_full_scan_matcher.h> +#include <ydb/services/workload_manager/has_full_scan_matcher.h> #include <ydb/core/resource_pools/regex_predicate.h> #include <library/cpp/testing/unittest/registar.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { using NResourcePool::TRegexPredicate; using NKqpProto::TKqpPhyQuery; @@ -36,7 +36,7 @@ TKqpReadRangesSource* AddSource(TKqpPhyQuery& phy, const TString& path) { const auto Matches = [](auto builder, auto configure) { TKqpPhyQuery phy; configure(builder(phy, TString("/Root/t"))); - return NWorkload::MatchesFullScan(TRegexPredicate::FromGlob("/Root/t"), phy); + return MatchesFullScan(TRegexPredicate::FromGlob("/Root/t"), phy); }; const auto AssertHasFullScan = [](auto c) { UNIT_ASSERT( Matches(AddOp, c)); }; @@ -51,21 +51,21 @@ Y_UNIT_TEST_SUITE(TFullScanMatcherPredicate) { Y_UNIT_TEST(EmptyPredicateAcceptsAny) { TKqpPhyQuery phy; - UNIT_ASSERT(NWorkload::MatchesFullScan(std::nullopt, phy)); + UNIT_ASSERT(MatchesFullScan(std::nullopt, phy)); } Y_UNIT_TEST(NonMatchingPathIsIgnored) { auto pred = TRegexPredicate::FromGlob("/Root/critical"); TKqpPhyQuery phy; AddOp(phy, "/Root/other")->MutableReadRange()->MutableKeyRange(); - UNIT_ASSERT(!NWorkload::MatchesFullScan(pred, phy)); + UNIT_ASSERT(!MatchesFullScan(pred, phy)); } Y_UNIT_TEST(GlobPatternMatches) { auto pred = TRegexPredicate::FromGlob("/Root/db/orders_archive*"); TKqpPhyQuery phy; AddOp(phy, "/Root/db/orders_archive_2024")->MutableReadRange()->MutableKeyRange(); - UNIT_ASSERT(NWorkload::MatchesFullScan(pred, phy)); + UNIT_ASSERT(MatchesFullScan(pred, phy)); } Y_UNIT_TEST(MixedOpsReturnTrueOnAnyMatch) { @@ -73,7 +73,7 @@ Y_UNIT_TEST_SUITE(TFullScanMatcherPredicate) { TKqpPhyQuery phy; AddOp(phy, "/Root/other")->MutableReadRange()->MutableKeyRange(); AddOp(phy, "/Root/t")->MutableReadRange()->MutableKeyRange(); - UNIT_ASSERT(NWorkload::MatchesFullScan(pred, phy)); + UNIT_ASSERT(MatchesFullScan(pred, phy)); } } @@ -161,4 +161,4 @@ Y_UNIT_TEST_SUITE(TFullScanMatcherLimit) { } } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/ut/kqp_has_full_scan_ut.cpp b/ydb/services/workload_manager/ut/has_full_scan_ut.cpp index ee69f7f5fa4..05c47bfda46 100644 --- a/ydb/core/kqp/workload_service/ut/kqp_has_full_scan_ut.cpp +++ b/ydb/services/workload_manager/ut/has_full_scan_ut.cpp @@ -1,13 +1,13 @@ #include <ydb/core/base/path.h> -#include <ydb/core/kqp/workload_service/ut/common/kqp_query_classifier_ut_common.h> -#include <ydb/core/kqp/workload_service/ut/common/kqp_workload_service_ut_common.h> +#include <ydb/services/workload_manager/ut/common/query_classifier_ut_common.h> +#include <ydb/services/workload_manager/ut/common/workload_service_ut_common.h> #include <library/cpp/testing/unittest/registar.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { -using namespace NWorkload; +using namespace NWorkloadManager; using namespace NYdb; @@ -66,7 +66,7 @@ namespace { // Creates `orders(Id, Status)` with a global secondary index on Status, // grants the given user full access, and inserts a couple of rows so the planner // doesn't take an empty-table shortcut. -void SetupOrdersWithIndex(TIntrusivePtr<NWorkload::IYdbSetup> ydb, const TString& userSID) { +void SetupOrdersWithIndex(TIntrusivePtr<IYdbSetup> ydb, const TString& userSID) { ydb->ExecuteSchemeQuery(TStringBuilder() << R"( GRANT ALL ON `/)" << ydb->GetSettings().DomainName_ << R"(` TO `)" << userSID << R"(`; CREATE TABLE orders ( @@ -77,7 +77,7 @@ void SetupOrdersWithIndex(TIntrusivePtr<NWorkload::IYdbSetup> ydb, const TString ); )"); - auto insertSettings = NWorkload::TQueryRunnerSettings().PoolId(NResourcePool::DEFAULT_POOL_ID); + auto insertSettings = TQueryRunnerSettings().PoolId(NResourcePool::DEFAULT_POOL_ID); auto insertResult = ydb->ExecuteQuery(R"( UPSERT INTO orders (Id, Status) VALUES (1u, "a"), @@ -91,7 +91,7 @@ void SetupOrdersWithIndex(TIntrusivePtr<NWorkload::IYdbSetup> ydb, const TString // classifier targeting it, and inserts a couple of rows. The classifier's target pool // has CONCURRENT_QUERY_LIMIT=0 so any match rejects with PRECONDITION_FAILED. void SetupOlapEventsAndClassifier( - TIntrusivePtr<NWorkload::IYdbSetup> ydb, const TString& userSID, const TString& poolId) + TIntrusivePtr<IYdbSetup> ydb, const TString& userSID, const TString& poolId) { ydb->ExecuteSchemeQuery(TStringBuilder() << R"( GRANT ALL ON `/)" << ydb->GetSettings().DomainName_ << R"(` TO `)" << userSID << R"(`; @@ -111,7 +111,7 @@ void SetupOlapEventsAndClassifier( ); )"); - auto insertSettings = NWorkload::TQueryRunnerSettings().PoolId(NResourcePool::DEFAULT_POOL_ID); + auto insertSettings = TQueryRunnerSettings().PoolId(NResourcePool::DEFAULT_POOL_ID); auto insertResult = ydb->ExecuteQuery(R"( UPSERT INTO olap_events (Id, Payload) VALUES (1u, "a"), (2u, "b"); )", insertSettings); @@ -128,7 +128,7 @@ Y_UNIT_TEST_SUITE(HasFullScanSecondaryIndexOltpDdl) { // main table is not accessed. A classifier targeting only `/Root/orders` (not the impl path) // must NOT fire — the impl full scan happens on a path outside the classifier's regex. Y_UNIT_TEST(TestClassifierOnMainTableIgnoresIndexScan) { - auto ydb = NWorkload::TYdbSetupSettings().Create(); + auto ydb = TYdbSetupSettings().Create(); const TString& poolId = "reject_pool"; const TString& userSID = "test@user"; @@ -145,8 +145,8 @@ Y_UNIT_TEST_SUITE(HasFullScanSecondaryIndexOltpDdl) { )"); const TString query = "SELECT * FROM orders VIEW by_status;"; - auto settings = NWorkload::TQueryRunnerSettings().PoolId("").UserSID(userSID); - NWorkload::WaitForClassifierSuccess(ydb, query, settings); + auto settings = TQueryRunnerSettings().PoolId("").UserSID(userSID); + WaitForClassifierSuccess(ydb, query, settings); } // Same query, but the classifier targets the index impl table path directly. @@ -157,7 +157,7 @@ Y_UNIT_TEST_SUITE(HasFullScanSecondaryIndexOltpDdl) { // `SELECT * FROM orders VIEW by_status` (no predicate) would trigger the planner // warning "Given predicate is not suitable for used index" and fall back to a main-table scan. Y_UNIT_TEST(TestClassifierOnIndexImplTableDetectsFullScan) { - auto ydb = NWorkload::TYdbSetupSettings().Create(); + auto ydb = TYdbSetupSettings().Create(); const TString& poolId = "reject_pool"; const TString& userSID = "test@user"; @@ -174,8 +174,8 @@ Y_UNIT_TEST_SUITE(HasFullScanSecondaryIndexOltpDdl) { )"); const TString query = "SELECT * FROM orders VIEW by_status;"; - auto settings = NWorkload::TQueryRunnerSettings().PoolId("").UserSID(userSID); - NWorkload::WaitForClassifierFail(ydb, query, settings, poolId); + auto settings = TQueryRunnerSettings().PoolId("").UserSID(userSID); + WaitForClassifierFail(ydb, query, settings, poolId); } } @@ -184,7 +184,7 @@ Y_UNIT_TEST_SUITE(HasFullScanIgnoresLimitDdl) { // LIMIT does not bypass full-scan detection. `SELECT * FROM orders LIMIT 1` still has an // unbounded ReadRange on /Root/orders — the classifier fires and rejects into the target pool. Y_UNIT_TEST(TestLimitDoesNotBypassClassifier) { - auto ydb = NWorkload::TYdbSetupSettings().Create(); + auto ydb = TYdbSetupSettings().Create(); const TString& poolId = "reject_pool"; const TString& userSID = "test@user"; @@ -201,22 +201,22 @@ Y_UNIT_TEST_SUITE(HasFullScanIgnoresLimitDdl) { )"); const TString query = "SELECT * FROM orders LIMIT 1;"; - auto settings = NWorkload::TQueryRunnerSettings().PoolId("").UserSID(userSID); - NWorkload::WaitForClassifierFail(ydb, query, settings, poolId); + auto settings = TQueryRunnerSettings().PoolId("").UserSID(userSID); + WaitForClassifierFail(ydb, query, settings, poolId); } // Same LIMIT semantic applied to OLAP: `ReadOlapRange` with an unbounded param name and // HasItemsLimit() set still counts as a full scan. Y_UNIT_TEST(TestLimitDoesNotBypassClassifierOnOlapTable) { - auto ydb = NWorkload::TYdbSetupSettings().Create(); + auto ydb = TYdbSetupSettings().Create(); const TString& poolId = "reject_pool"; const TString& userSID = "test@user"; SetupOlapEventsAndClassifier(ydb, userSID, poolId); const TString query = "SELECT * FROM olap_events LIMIT 1;"; - auto settings = NWorkload::TQueryRunnerSettings().PoolId("").UserSID(userSID); - NWorkload::WaitForClassifierFail(ydb, query, settings, poolId); + auto settings = TQueryRunnerSettings().PoolId("").UserSID(userSID); + WaitForClassifierFail(ydb, query, settings, poolId); } // Zero-selectivity filter behind a LIMIT: at runtime the executor scans the whole table @@ -225,15 +225,15 @@ Y_UNIT_TEST_SUITE(HasFullScanIgnoresLimitDdl) { // exactly the runtime full-scan the LIMIT rule used to miss. OLAP is used because there is // no secondary index to route the filter through, so the plan shape is unambiguous. Y_UNIT_TEST(TestFilterWithLimitDoesNotBypassClassifierOnOlapTable) { - auto ydb = NWorkload::TYdbSetupSettings().Create(); + auto ydb = TYdbSetupSettings().Create(); const TString& poolId = "reject_pool"; const TString& userSID = "test@user"; SetupOlapEventsAndClassifier(ydb, userSID, poolId); const TString query = "SELECT * FROM olap_events WHERE Payload = 'nonexistent' LIMIT 1;"; - auto settings = NWorkload::TQueryRunnerSettings().PoolId("").UserSID(userSID); - NWorkload::WaitForClassifierFail(ydb, query, settings, poolId); + auto settings = TQueryRunnerSettings().PoolId("").UserSID(userSID); + WaitForClassifierFail(ydb, query, settings, poolId); } } @@ -243,30 +243,30 @@ Y_UNIT_TEST_SUITE(HasFullScanOlapDdl) { // row-store's `ReadRange`/`ReadRangesSource`. This test proves the matcher's OLAP branch // is invoked end-to-end by the real planner, complementing the ReadOlapRange matcher UT. Y_UNIT_TEST(TestClassifierDetectsFullScanOnOlapTable) { - auto ydb = NWorkload::TYdbSetupSettings().Create(); + auto ydb = TYdbSetupSettings().Create(); const TString& poolId = "reject_pool"; const TString& userSID = "test@user"; SetupOlapEventsAndClassifier(ydb, userSID, poolId); const TString query = "SELECT * FROM olap_events;"; - auto settings = NWorkload::TQueryRunnerSettings().PoolId("").UserSID(userSID); - NWorkload::WaitForClassifierFail(ydb, query, settings, poolId); + auto settings = TQueryRunnerSettings().PoolId("").UserSID(userSID); + WaitForClassifierFail(ydb, query, settings, poolId); } // PK point lookup on OLAP compiles to a parametrized ReadOlapRange with a non-empty // KeyRanges.ParamName. Matcher must not flag it — proves the OLAP negative branch e2e. Y_UNIT_TEST(TestClassifierIgnoresPkPointLookupOnOlapTable) { - auto ydb = NWorkload::TYdbSetupSettings().Create(); + auto ydb = TYdbSetupSettings().Create(); const TString& poolId = "reject_pool"; const TString& userSID = "test@user"; SetupOlapEventsAndClassifier(ydb, userSID, poolId); const TString query = "SELECT * FROM olap_events WHERE Id = 1u;"; - auto settings = NWorkload::TQueryRunnerSettings().PoolId("").UserSID(userSID); - NWorkload::WaitForClassifierSuccess(ydb, query, settings); + auto settings = TQueryRunnerSettings().PoolId("").UserSID(userSID); + WaitForClassifierSuccess(ydb, query, settings); } } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/ut/kqp_has_path_ddl_ut.cpp b/ydb/services/workload_manager/ut/has_path_ddl_ut.cpp index 114721f8651..260701a254d 100644 --- a/ydb/core/kqp/workload_service/ut/kqp_has_path_ddl_ut.cpp +++ b/ydb/services/workload_manager/ut/has_path_ddl_ut.cpp @@ -1,20 +1,20 @@ #include <ydb/core/base/path.h> -#include <ydb/core/kqp/workload_service/ut/common/kqp_query_classifier_ut_common.h> -#include <ydb/core/kqp/workload_service/ut/common/kqp_workload_service_ut_common.h> +#include <ydb/services/workload_manager/ut/common/query_classifier_ut_common.h> +#include <ydb/services/workload_manager/ut/common/workload_service_ut_common.h> #include <library/cpp/testing/unittest/registar.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { -using namespace NWorkload; +using namespace NWorkloadManager; using namespace NYdb; namespace { // Creates a simple row-store `t(Id, Payload)` and grants the given user full access. -void SetupRowStoreTable(TIntrusivePtr<NWorkload::IYdbSetup> ydb, const TString& userSID) { +void SetupRowStoreTable(TIntrusivePtr<IYdbSetup> ydb, const TString& userSID) { ydb->ExecuteSchemeQuery(TStringBuilder() << R"( GRANT ALL ON `/)" << ydb->GetSettings().DomainName_ << R"(` TO `)" << userSID << R"(`; CREATE TABLE t ( @@ -24,7 +24,7 @@ void SetupRowStoreTable(TIntrusivePtr<NWorkload::IYdbSetup> ydb, const TString& ); )"); - auto insertSettings = NWorkload::TQueryRunnerSettings().PoolId(NResourcePool::DEFAULT_POOL_ID); + auto insertSettings = TQueryRunnerSettings().PoolId(NResourcePool::DEFAULT_POOL_ID); auto insertResult = ydb->ExecuteQuery(R"( UPSERT INTO t (Id, Payload) VALUES (1u, "a"), (2u, "b"); )", insertSettings); @@ -33,7 +33,7 @@ void SetupRowStoreTable(TIntrusivePtr<NWorkload::IYdbSetup> ydb, const TString& } // Creates a column-store `olap_t(Id, Payload)` and grants full access. -void SetupColumnStoreTable(TIntrusivePtr<NWorkload::IYdbSetup> ydb, const TString& userSID) { +void SetupColumnStoreTable(TIntrusivePtr<IYdbSetup> ydb, const TString& userSID) { ydb->ExecuteSchemeQuery(TStringBuilder() << R"( GRANT ALL ON `/)" << ydb->GetSettings().DomainName_ << R"(` TO `)" << userSID << R"(`; CREATE TABLE olap_t ( @@ -45,7 +45,7 @@ void SetupColumnStoreTable(TIntrusivePtr<NWorkload::IYdbSetup> ydb, const TStrin ); )"); - auto insertSettings = NWorkload::TQueryRunnerSettings().PoolId(NResourcePool::DEFAULT_POOL_ID); + auto insertSettings = TQueryRunnerSettings().PoolId(NResourcePool::DEFAULT_POOL_ID); auto insertResult = ydb->ExecuteQuery(R"( UPSERT INTO olap_t (Id, Payload) VALUES (1u, "a"), (2u, "b"); )", insertSettings); @@ -54,7 +54,7 @@ void SetupColumnStoreTable(TIntrusivePtr<NWorkload::IYdbSetup> ydb, const TStrin } // Creates `orders(Id, Status)` with a global secondary index on Status. -void SetupOrdersWithIndex(TIntrusivePtr<NWorkload::IYdbSetup> ydb, const TString& userSID) { +void SetupOrdersWithIndex(TIntrusivePtr<IYdbSetup> ydb, const TString& userSID) { ydb->ExecuteSchemeQuery(TStringBuilder() << R"( GRANT ALL ON `/)" << ydb->GetSettings().DomainName_ << R"(` TO `)" << userSID << R"(`; CREATE TABLE orders ( @@ -65,7 +65,7 @@ void SetupOrdersWithIndex(TIntrusivePtr<NWorkload::IYdbSetup> ydb, const TString ); )"); - auto insertSettings = NWorkload::TQueryRunnerSettings().PoolId(NResourcePool::DEFAULT_POOL_ID); + auto insertSettings = TQueryRunnerSettings().PoolId(NResourcePool::DEFAULT_POOL_ID); auto insertResult = ydb->ExecuteQuery(R"( UPSERT INTO orders (Id, Status) VALUES (1u, "a"), (2u, "b"); )", insertSettings); @@ -76,7 +76,7 @@ void SetupOrdersWithIndex(TIntrusivePtr<NWorkload::IYdbSetup> ydb, const TString // Creates `t(...)` and a view `v` over it, grants full access. // The view body must fully qualify the underlying table — the invoker's // session context doesn't resolve unqualified `t` when the view is expanded. -void SetupTableAndView(TIntrusivePtr<NWorkload::IYdbSetup> ydb, const TString& userSID) { +void SetupTableAndView(TIntrusivePtr<IYdbSetup> ydb, const TString& userSID) { ydb->ExecuteSchemeQuery(TStringBuilder() << R"( GRANT ALL ON `/)" << ydb->GetSettings().DomainName_ << R"(` TO `)" << userSID << R"(`; CREATE TABLE t ( @@ -87,7 +87,7 @@ void SetupTableAndView(TIntrusivePtr<NWorkload::IYdbSetup> ydb, const TString& u CREATE VIEW v WITH (security_invoker = true) AS SELECT * FROM `/Root/t`; )"); - auto insertSettings = NWorkload::TQueryRunnerSettings().PoolId(NResourcePool::DEFAULT_POOL_ID); + auto insertSettings = TQueryRunnerSettings().PoolId(NResourcePool::DEFAULT_POOL_ID); auto insertResult = ydb->ExecuteQuery(R"( UPSERT INTO t (Id, Payload) VALUES (1u, "a"), (2u, "b"); )", insertSettings); @@ -96,7 +96,7 @@ void SetupTableAndView(TIntrusivePtr<NWorkload::IYdbSetup> ydb, const TString& u } // EDS/ET/Topic/CDC DDL ITs live in `ydb/core/kqp/ut/federated_query/datastreams/` -// (see `kqp_has_path_ut.cpp` there). Those kinds need TStreamingTestFixture's +// (see `has_path_ut.cpp` there). Those kinds need TStreamingTestFixture's // mock connector + mock PQ gateway + http gateway wiring, which is not present // in workload_service_ut's setup. This file covers the kinds whose fixture is // cheap: regular tables, sysview, secondary index, view (underlying-table gate). @@ -120,7 +120,7 @@ TString RejectClassifierDdl(const TString& poolId, const TString& classifierName Y_UNIT_TEST_SUITE(HasPathTableOltp) { Y_UNIT_TEST(TestClassifierMatchesTablePath) { - auto ydb = NWorkload::TYdbSetupSettings().Create(); + auto ydb = TYdbSetupSettings().Create(); const TString& poolId = "reject_pool"; const TString& userSID = "test@user"; SetupRowStoreTable(ydb, userSID); @@ -128,12 +128,12 @@ Y_UNIT_TEST_SUITE(HasPathTableOltp) { ydb->ExecuteSchemeQuery(RejectClassifierDdl(poolId, "row_table_classifier", "/Root/t")); const TString query = "SELECT * FROM t;"; - auto settings = NWorkload::TQueryRunnerSettings().PoolId("").UserSID(userSID); - NWorkload::WaitForClassifierFail(ydb, query, settings, poolId); + auto settings = TQueryRunnerSettings().PoolId("").UserSID(userSID); + WaitForClassifierFail(ydb, query, settings, poolId); } Y_UNIT_TEST(TestClassifierIgnoresOtherTable) { - auto ydb = NWorkload::TYdbSetupSettings().Create(); + auto ydb = TYdbSetupSettings().Create(); const TString& poolId = "reject_pool"; const TString& userSID = "test@user"; SetupRowStoreTable(ydb, userSID); @@ -141,8 +141,8 @@ Y_UNIT_TEST_SUITE(HasPathTableOltp) { ydb->ExecuteSchemeQuery(RejectClassifierDdl(poolId, "row_table_classifier", "/Root/other_table")); const TString query = "SELECT * FROM t;"; - auto settings = NWorkload::TQueryRunnerSettings().PoolId("").UserSID(userSID); - NWorkload::WaitForClassifierSuccess(ydb, query, settings); + auto settings = TQueryRunnerSettings().PoolId("").UserSID(userSID); + WaitForClassifierSuccess(ydb, query, settings); } } @@ -152,7 +152,7 @@ Y_UNIT_TEST_SUITE(HasPathTableOltp) { Y_UNIT_TEST_SUITE(HasPathColumnTableOlap) { Y_UNIT_TEST(TestClassifierMatchesOlapTablePath) { - auto ydb = NWorkload::TYdbSetupSettings().Create(); + auto ydb = TYdbSetupSettings().Create(); const TString& poolId = "reject_pool"; const TString& userSID = "test@user"; SetupColumnStoreTable(ydb, userSID); @@ -160,14 +160,14 @@ Y_UNIT_TEST_SUITE(HasPathColumnTableOlap) { ydb->ExecuteSchemeQuery(RejectClassifierDdl(poolId, "olap_table_classifier", "/Root/olap_t")); const TString query = "SELECT * FROM olap_t;"; - auto settings = NWorkload::TQueryRunnerSettings().PoolId("").UserSID(userSID); - NWorkload::WaitForClassifierFail(ydb, query, settings, poolId); + auto settings = TQueryRunnerSettings().PoolId("").UserSID(userSID); + WaitForClassifierFail(ydb, query, settings, poolId); } Y_UNIT_TEST(TestClassifierMatchesOlapPointLookup) { // Unlike HAS_FULL_SCAN, HAS_PATH doesn't care about read shape — a PK lookup // still touches the table and must fire. - auto ydb = NWorkload::TYdbSetupSettings().Create(); + auto ydb = TYdbSetupSettings().Create(); const TString& poolId = "reject_pool"; const TString& userSID = "test@user"; SetupColumnStoreTable(ydb, userSID); @@ -175,8 +175,8 @@ Y_UNIT_TEST_SUITE(HasPathColumnTableOlap) { ydb->ExecuteSchemeQuery(RejectClassifierDdl(poolId, "olap_table_classifier", "/Root/olap_t")); const TString query = "SELECT * FROM olap_t WHERE Id = 1u;"; - auto settings = NWorkload::TQueryRunnerSettings().PoolId("").UserSID(userSID); - NWorkload::WaitForClassifierFail(ydb, query, settings, poolId); + auto settings = TQueryRunnerSettings().PoolId("").UserSID(userSID); + WaitForClassifierFail(ydb, query, settings, poolId); } } @@ -185,7 +185,7 @@ Y_UNIT_TEST_SUITE(HasPathColumnTableOlap) { Y_UNIT_TEST_SUITE(HasPathSecondaryIndex) { Y_UNIT_TEST(TestClassifierMatchesIndexImplPath) { - auto ydb = NWorkload::TYdbSetupSettings().Create(); + auto ydb = TYdbSetupSettings().Create(); const TString& poolId = "reject_pool"; const TString& userSID = "test@user"; SetupOrdersWithIndex(ydb, userSID); @@ -195,14 +195,14 @@ Y_UNIT_TEST_SUITE(HasPathSecondaryIndex) { // Covering scan on the index — planner emits a full read of indexImplTable. const TString query = "SELECT Status FROM orders VIEW by_status;"; - auto settings = NWorkload::TQueryRunnerSettings().PoolId("").UserSID(userSID); - NWorkload::WaitForClassifierFail(ydb, query, settings, poolId); + auto settings = TQueryRunnerSettings().PoolId("").UserSID(userSID); + WaitForClassifierFail(ydb, query, settings, poolId); } Y_UNIT_TEST(TestClassifierMatchesIndexOnWrite) { // UPSERT into an indexed table fans out to write both the main table and // the impl table. HAS_PATH targeting the impl path still fires. - auto ydb = NWorkload::TYdbSetupSettings().Create(); + auto ydb = TYdbSetupSettings().Create(); const TString& poolId = "reject_pool"; const TString& userSID = "test@user"; SetupOrdersWithIndex(ydb, userSID); @@ -211,8 +211,8 @@ Y_UNIT_TEST_SUITE(HasPathSecondaryIndex) { poolId, "index_write_classifier", "/Root/orders/by_status/indexImplTable")); const TString query = "UPSERT INTO orders (Id, Status) VALUES (3u, \"c\");"; - auto settings = NWorkload::TQueryRunnerSettings().PoolId("").UserSID(userSID); - NWorkload::WaitForClassifierFail(ydb, query, settings, poolId); + auto settings = TQueryRunnerSettings().PoolId("").UserSID(userSID); + WaitForClassifierFail(ydb, query, settings, poolId); } } @@ -221,7 +221,7 @@ Y_UNIT_TEST_SUITE(HasPathSecondaryIndex) { Y_UNIT_TEST_SUITE(HasPathSysView) { Y_UNIT_TEST(TestClassifierMatchesSysViewPath) { - auto ydb = NWorkload::TYdbSetupSettings().Create(); + auto ydb = TYdbSetupSettings().Create(); const TString& poolId = "reject_pool"; const TString& userSID = "test@user"; SetupRowStoreTable(ydb, userSID); @@ -230,12 +230,12 @@ Y_UNIT_TEST_SUITE(HasPathSysView) { poolId, "sysview_classifier", "/Root/.sys/partition_stats")); const TString query = "SELECT * FROM `.sys/partition_stats`;"; - auto settings = NWorkload::TQueryRunnerSettings().PoolId("").UserSID(userSID); - NWorkload::WaitForClassifierFail(ydb, query, settings, poolId); + auto settings = TQueryRunnerSettings().PoolId("").UserSID(userSID); + WaitForClassifierFail(ydb, query, settings, poolId); } Y_UNIT_TEST(TestClassifierIgnoresRegularTableWhenSysViewTargeted) { - auto ydb = NWorkload::TYdbSetupSettings().Create(); + auto ydb = TYdbSetupSettings().Create(); const TString& poolId = "reject_pool"; const TString& userSID = "test@user"; SetupRowStoreTable(ydb, userSID); @@ -244,8 +244,8 @@ Y_UNIT_TEST_SUITE(HasPathSysView) { poolId, "sysview_classifier", "/Root/.sys/*")); const TString query = "SELECT * FROM t;"; - auto settings = NWorkload::TQueryRunnerSettings().PoolId("").UserSID(userSID); - NWorkload::WaitForClassifierSuccess(ydb, query, settings); + auto settings = TQueryRunnerSettings().PoolId("").UserSID(userSID); + WaitForClassifierSuccess(ydb, query, settings); } } @@ -260,7 +260,7 @@ Y_UNIT_TEST_SUITE(HasPathView) { // A regex on the view path `/Root/v` never fires. /* Y_UNIT_TEST(TestClassifierMatchesViewPath) { - auto ydb = NWorkload::TYdbSetupSettings().Create(); + auto ydb = TYdbSetupSettings().Create(); const TString& poolId = "reject_pool"; const TString& userSID = "test@user"; SetupTableAndView(ydb, userSID); @@ -268,15 +268,15 @@ Y_UNIT_TEST_SUITE(HasPathView) { ydb->ExecuteSchemeQuery(RejectClassifierDdl(poolId, "view_classifier", "/Root/v")); const TString query = "SELECT * FROM v;"; - auto settings = NWorkload::TQueryRunnerSettings().PoolId("").UserSID(userSID); - NWorkload::WaitForClassifierFail(ydb, query, settings, poolId); + auto settings = TQueryRunnerSettings().PoolId("").UserSID(userSID); + WaitForClassifierFail(ydb, query, settings, poolId); } */ Y_UNIT_TEST(TestClassifierMatchesUnderlyingTableUnderView) { // Views are inlined — the underlying table's path also appears in QueryTables. // A regex targeting the underlying table matches queries going through the view. - auto ydb = NWorkload::TYdbSetupSettings().Create(); + auto ydb = TYdbSetupSettings().Create(); const TString& poolId = "reject_pool"; const TString& userSID = "test@user"; SetupTableAndView(ydb, userSID); @@ -284,9 +284,9 @@ Y_UNIT_TEST_SUITE(HasPathView) { ydb->ExecuteSchemeQuery(RejectClassifierDdl(poolId, "underlying_classifier", "/Root/t")); const TString query = "SELECT * FROM v;"; - auto settings = NWorkload::TQueryRunnerSettings().PoolId("").UserSID(userSID); - NWorkload::WaitForClassifierFail(ydb, query, settings, poolId); + auto settings = TQueryRunnerSettings().PoolId("").UserSID(userSID); + WaitForClassifierFail(ydb, query, settings, poolId); } } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/ut/kqp_has_path_matcher_ut.cpp b/ydb/services/workload_manager/ut/has_path_matcher_ut.cpp index 55cd3c86ef8..1cc1980b870 100644 --- a/ydb/core/kqp/workload_service/ut/kqp_has_path_matcher_ut.cpp +++ b/ydb/services/workload_manager/ut/has_path_matcher_ut.cpp @@ -1,4 +1,4 @@ -#include <ydb/core/kqp/workload_service/kqp_has_path_matcher.h> +#include <ydb/services/workload_manager/has_path_matcher.h> #include <ydb/core/resource_pools/regex_predicate.h> #include <ydb/library/yql/providers/pq/proto/dq_io.pb.h> @@ -8,7 +8,7 @@ #include <library/cpp/testing/unittest/registar.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { using NResourcePool::TRegexPredicate; using NKqpProto::TKqpPhyQuery; @@ -41,7 +41,7 @@ void AddPqTopicSink(TKqpPhyQuery& phy, const TString& topicPath) { } // Fixed-pattern helpers: each suite pins the predicate so tests read as pure -// build-and-check one-liners (mirrors kqp_has_full_scan_matcher_ut.cpp style). +// build-and-check one-liners (mirrors has_full_scan_matcher_ut.cpp style). // Runs the matcher against a phyQuery holding one tx-level table entry. // Predicate + path both vary so predicate-behaviour tests can exercise glob @@ -49,27 +49,27 @@ void AddPqTopicSink(TKqpPhyQuery& phy, const TString& topicPath) { const auto MatchesTxTable = [](const TString& pattern, const TString& path) { TKqpPhyQuery phy; AddTxTable(phy, path); - return NWorkload::MatchesPath(TRegexPredicate::FromGlob(pattern), {}, phy); + return MatchesPath(TRegexPredicate::FromGlob(pattern), {}, phy); }; // Runs the matcher with pattern `/Root/t*` against the given queryTables list. const auto MatchesQueryTables = [](const TVector<TString>& queryTables) { TKqpPhyQuery phy; - return NWorkload::MatchesPath(TRegexPredicate::FromGlob("/Root/t*"), queryTables, phy); + return MatchesPath(TRegexPredicate::FromGlob("/Root/t*"), queryTables, phy); }; // Runs the matcher with pattern `/Root/db/target` against a phyQuery built by `build`. const auto MatchesTx = [](auto build) { TKqpPhyQuery phy; build(phy); - return NWorkload::MatchesPath(TRegexPredicate::FromGlob("/Root/db/target"), {}, phy); + return MatchesPath(TRegexPredicate::FromGlob("/Root/db/target"), {}, phy); }; // Runs the matcher with pattern `/Root/db/topic` against a phyQuery built by `build`. const auto MatchesTopics = [](auto build) { TKqpPhyQuery phy; build(phy); - return NWorkload::MatchesPath(TRegexPredicate::FromGlob("/Root/db/topic"), {}, phy); + return MatchesPath(TRegexPredicate::FromGlob("/Root/db/topic"), {}, phy); }; } // anonymous namespace @@ -80,7 +80,7 @@ Y_UNIT_TEST_SUITE(TPathMatcherPredicate) { Y_UNIT_TEST(EmptyPredicateAcceptsAny) { TKqpPhyQuery phy; AddTxTable(phy, "/Root/anything"); - UNIT_ASSERT(NWorkload::MatchesPath(std::nullopt, {}, phy)); + UNIT_ASSERT(MatchesPath(std::nullopt, {}, phy)); } Y_UNIT_TEST(NonMatchingPathIsNotAccepted) { @@ -166,4 +166,4 @@ Y_UNIT_TEST_SUITE(TPathMatcherTopics) { } } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/ut/kqp_has_path_ut.cpp b/ydb/services/workload_manager/ut/has_path_ut.cpp index 44f22c21661..760313b4c7f 100644 --- a/ydb/core/kqp/workload_service/ut/kqp_has_path_ut.cpp +++ b/ydb/services/workload_manager/ut/has_path_ut.cpp @@ -1,13 +1,13 @@ #include <ydb/core/base/path.h> -#include <ydb/core/kqp/workload_service/ut/common/kqp_query_classifier_ut_common.h> -#include <ydb/core/kqp/workload_service/ut/common/kqp_workload_service_ut_common.h> +#include <ydb/services/workload_manager/ut/common/query_classifier_ut_common.h> +#include <ydb/services/workload_manager/ut/common/workload_service_ut_common.h> #include <library/cpp/testing/unittest/registar.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { -using namespace NWorkload; +using namespace NWorkloadManager; namespace { @@ -61,4 +61,4 @@ Y_UNIT_TEST_SUITE(TQueryClassifierHasPath) { } } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/services/workload_manager/ut/has_stream_matcher_ut.cpp b/ydb/services/workload_manager/ut/has_stream_matcher_ut.cpp new file mode 100644 index 00000000000..89ee1ab6dac --- /dev/null +++ b/ydb/services/workload_manager/ut/has_stream_matcher_ut.cpp @@ -0,0 +1,52 @@ +#include <ydb/services/workload_manager/has_stream_matcher.h> + +#include <library/cpp/testing/unittest/registar.h> + + +namespace NKikimr::NWorkloadManager { + + +Y_UNIT_TEST_SUITE(TStreamMatcherPredicate) { + + + Y_UNIT_TEST(NulloptAcceptsAny) { + // No filter — every query passes regardless of streaming ops. + { + NKqp::TUserRequestContext userRequestContext; + userRequestContext.IsStreamingQuery = false; + UNIT_ASSERT(MatchesStream(std::nullopt, userRequestContext)); + } + { + NKqp::TUserRequestContext userRequestContext; + userRequestContext.IsStreamingQuery = true; + UNIT_ASSERT(MatchesStream(std::nullopt, userRequestContext)); + } + } + + Y_UNIT_TEST(TrueMatchesStreamingQuery) { + NKqp::TUserRequestContext userRequestContext; + userRequestContext.IsStreamingQuery = true; + + UNIT_ASSERT(MatchesStream(true, userRequestContext)); + } + + Y_UNIT_TEST(TrueDoesNotMatchNonStreamingQuery) { + NKqp::TUserRequestContext userRequestContext; + userRequestContext.IsStreamingQuery = false; + UNIT_ASSERT(!MatchesStream(true, userRequestContext)); + } + + Y_UNIT_TEST(FalseMatchesNonStreamingQuery) { + NKqp::TUserRequestContext userRequestContext; + userRequestContext.IsStreamingQuery = false; + UNIT_ASSERT(MatchesStream(false, userRequestContext)); + } + + Y_UNIT_TEST(FalseDoesNotMatchStreamingQuery) { + NKqp::TUserRequestContext userRequestContext; + userRequestContext.IsStreamingQuery = true; + UNIT_ASSERT(!MatchesStream(false, userRequestContext)); + } +} + +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/ut/kqp_has_stream_ut.cpp b/ydb/services/workload_manager/ut/has_stream_ut.cpp index 2674ecbd350..f12fdae8bd8 100644 --- a/ydb/core/kqp/workload_service/ut/kqp_has_stream_ut.cpp +++ b/ydb/services/workload_manager/ut/has_stream_ut.cpp @@ -1,11 +1,11 @@ -#include <ydb/core/kqp/workload_service/ut/common/kqp_query_classifier_ut_common.h> +#include <ydb/services/workload_manager/ut/common/query_classifier_ut_common.h> #include <library/cpp/testing/unittest/registar.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { -using namespace NWorkload; +using namespace NWorkloadManager; namespace { @@ -58,4 +58,4 @@ Y_UNIT_TEST_SUITE(TQueryClassifierHasStream) { } } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/ut/kqp_member_name_ut.cpp b/ydb/services/workload_manager/ut/member_name_ut.cpp index 995713933cf..f51b41afa10 100644 --- a/ydb/core/kqp/workload_service/ut/kqp_member_name_ut.cpp +++ b/ydb/services/workload_manager/ut/member_name_ut.cpp @@ -1,6 +1,6 @@ -#include <ydb/core/kqp/workload_service/ut/common/kqp_query_classifier_ut_common.h> +#include <ydb/services/workload_manager/ut/common/query_classifier_ut_common.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { Y_UNIT_TEST_SUITE(TQueryClassifierMemberName) { @@ -30,7 +30,7 @@ Y_UNIT_TEST_SUITE(TQueryClassifierMemberName) { }; auto view = TClassifierConfigsView(classifierSnap, TEST_DB); - auto classifier = NWorkload::CreateQueryClassifier(poolSnap, view, TEST_DB, ctx); + auto classifier = NWorkloadManager::CreateQueryClassifier(poolSnap, view, TEST_DB, ctx); auto result = classifier->PreCompileClassify(); UNIT_ASSERT_VALUES_EQUAL(GetPoolId(result), "pool_target"); } @@ -59,10 +59,10 @@ Y_UNIT_TEST_SUITE(TQueryClassifierMemberName) { }; auto view = TClassifierConfigsView(classifierSnap, TEST_DB); - auto classifier = NWorkload::CreateQueryClassifier(poolSnap, view, TEST_DB, ctx); + auto classifier = NWorkloadManager::CreateQueryClassifier(poolSnap, view, TEST_DB, ctx); auto result = classifier->PreCompileClassify(); UNIT_ASSERT_VALUES_EQUAL(GetPoolId(result), "pool_target"); } } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/ut/kqp_query_classifier_match_ut.cpp b/ydb/services/workload_manager/ut/query_classifier_match_ut.cpp index 7f4954171ba..378c589636b 100644 --- a/ydb/core/kqp/workload_service/ut/kqp_query_classifier_match_ut.cpp +++ b/ydb/services/workload_manager/ut/query_classifier_match_ut.cpp @@ -1,6 +1,6 @@ -#include <ydb/core/kqp/workload_service/ut/common/kqp_query_classifier_ut_common.h> +#include <ydb/services/workload_manager/ut/common/query_classifier_ut_common.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { namespace { @@ -23,27 +23,27 @@ void CheckDecision(std::function<void(NResourcePool::TPoolSettings&)> setup, EEx }); TClassifyContext ctx{.PoolId = "", .AppName = "", .UserToken = nullptr}; - auto classifier = NWorkload::CreateQueryClassifier( + auto classifier = NWorkloadManager::CreateQueryClassifier( poolSnap, TClassifierConfigsView(classifierSnap, TEST_DB), TEST_DB, std::move(ctx)); auto result = classifier->PreCompileClassify(); switch (expected) { case EExpected::Bypass: - UNIT_ASSERT_C(std::holds_alternative<NWorkload::IQueryClassifier::TBypass>(result), + UNIT_ASSERT_C(std::holds_alternative<NWorkloadManager::IQueryClassifier::TBypass>(result), TStringBuilder() << "Expected TBypass, got variant index " << result.index()); break; case EExpected::ResolvedAdmission: { - UNIT_ASSERT_C(std::holds_alternative<NWorkload::IQueryClassifier::TResolvedPoolId>(result), + UNIT_ASSERT_C(std::holds_alternative<NWorkloadManager::IQueryClassifier::TResolvedPoolId>(result), TStringBuilder() << "Expected TResolvedPoolId, got variant index " << result.index()); - UNIT_ASSERT_C(!std::get<NWorkload::IQueryClassifier::TResolvedPoolId>(result).SkipAdmission, + UNIT_ASSERT_C(!std::get<NWorkloadManager::IQueryClassifier::TResolvedPoolId>(result).SkipAdmission, "Expected SkipAdmission=false"); break; } case EExpected::ResolvedSkipAdmission: { - UNIT_ASSERT_C(std::holds_alternative<NWorkload::IQueryClassifier::TResolvedPoolId>(result), + UNIT_ASSERT_C(std::holds_alternative<NWorkloadManager::IQueryClassifier::TResolvedPoolId>(result), TStringBuilder() << "Expected TResolvedPoolId, got variant index " << result.index()); - UNIT_ASSERT_C(std::get<NWorkload::IQueryClassifier::TResolvedPoolId>(result).SkipAdmission, + UNIT_ASSERT_C(std::get<NWorkloadManager::IQueryClassifier::TResolvedPoolId>(result).SkipAdmission, "Expected SkipAdmission=true"); break; } @@ -134,4 +134,4 @@ Y_UNIT_TEST_SUITE(TQueryClassifierDecision) { } } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/ut/kqp_query_classifier_ut.cpp b/ydb/services/workload_manager/ut/query_classifier_ut.cpp index bccbbd71df6..4e1afd27e27 100644 --- a/ydb/core/kqp/workload_service/ut/kqp_query_classifier_ut.cpp +++ b/ydb/services/workload_manager/ut/query_classifier_ut.cpp @@ -1,16 +1,16 @@ #include <fmt/format.h> #include <library/cpp/testing/unittest/registar.h> -#include <ydb/core/kqp/common/events/workload_service.h> +#include <ydb/services/workload_manager/events.h> #include <ydb/core/kqp/common/simple/services.h> -#include <ydb/core/kqp/workload_service/kqp_query_classifier.h> +#include <ydb/services/workload_manager/query_classifier.h> #include <ydb/core/testlib/test_client.h> #include <ydb/core/kqp/common/events/events.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { namespace { -class TestQueryClassifier : public NWorkload::IQueryClassifier { +class TestQueryClassifier : public NWorkloadManager::IQueryClassifier { public: TestQueryClassifier(TPreCompileClassifyResult preClassifyResult, TPostCompileClassifyResult postClassifyResult) : PostCompileCalled(false) @@ -32,7 +32,7 @@ public: return EState::PreCompileDone; } - TPostCompileClassifyResult PostCompileClassify(const TPreparedQueryHolder&, const TUserRequestContext&) override { + TPostCompileClassifyResult PostCompileClassify(const NKqp::TPreparedQueryHolder&, const NKqp::TUserRequestContext&) override { PostCompileCalled = true; return PostClassifyResult; } @@ -44,7 +44,7 @@ private: }; template<typename TPreResult, typename TPostResult> -TEvKqp::TEvQueryResponse::TPtr RunQueryWith(TPreResult preResult, TPostResult postResult) { +NKqp::TEvKqp::TEvQueryResponse::TPtr RunQueryWith(TPreResult preResult, TPostResult postResult) { TPortManager tp; auto mbusport = tp.GetPort(2134); auto settings = Tests::TServerSettings(mbusport); @@ -65,7 +65,7 @@ TEvKqp::TEvQueryResponse::TPtr RunQueryWith(TPreResult preResult, TPostResult po auto captureEvents = [&](TTestActorRuntimeBase&, TAutoPtr<IEventHandle>& ev) { if (ev->GetTypeRewrite() == NKqp::TEvKqp::TEvQueryRequest::EventType) { // Replace the classifier after proxy set it - auto* request = ev->Get<TEvKqp::TEvQueryRequest>(); + auto* request = ev->Get<NKqp::TEvKqp::TEvQueryRequest>(); request->SetWmQueryClassifier(std::make_shared<TestQueryClassifier>(preResult, postResult)); } @@ -80,20 +80,20 @@ TEvKqp::TEvQueryResponse::TPtr RunQueryWith(TPreResult preResult, TPostResult po ev->Record.MutableRequest()->SetQuery("SELECT 1;"); runtime.Send(new IEventHandle(proxyId, sender, ev.Release())); - return runtime.GrabEdgeEventRethrow<TEvKqp::TEvQueryResponse>(sender); + return runtime.GrabEdgeEventRethrow<NKqp::TEvKqp::TEvQueryResponse>(sender); } template<typename TPreResult> -TEvKqp::TEvQueryResponse::TPtr RunQueryWithPreClassify(TPreResult preResult) { - return RunQueryWith(preResult, NWorkload::IQueryClassifier::TBypass{}); +NKqp::TEvKqp::TEvQueryResponse::TPtr RunQueryWithPreClassify(TPreResult preResult) { + return RunQueryWith(preResult, NWorkloadManager::IQueryClassifier::TBypass{}); } template<typename TPostResult> -TEvKqp::TEvQueryResponse::TPtr RunQueryWithPostClassify(TPostResult postResult) { - return RunQueryWith(NWorkload::IQueryClassifier::TPendingCompilation{}, postResult); +NKqp::TEvKqp::TEvQueryResponse::TPtr RunQueryWithPostClassify(TPostResult postResult) { + return RunQueryWith(NWorkloadManager::IQueryClassifier::TPendingCompilation{}, postResult); } -TString GetErrorMessageFromResponse(TEvKqp::TEvQueryResponse::TPtr r) { +TString GetErrorMessageFromResponse(NKqp::TEvKqp::TEvQueryResponse::TPtr r) { NYql::TIssues issues; NYql::IssuesFromMessage( r->Get()->Record.GetResponse().GetQueryIssues(), issues); UNIT_ASSERT(issues.Size()); @@ -104,14 +104,14 @@ TString GetErrorMessageFromResponse(TEvKqp::TEvQueryResponse::TPtr r) { Y_UNIT_TEST_SUITE(KqpQueryPreClassifier) { Y_UNIT_TEST(ShouldBypassOnPreClassify) { - auto reply = RunQueryWithPreClassify(NWorkload::IQueryClassifier::TBypass()); + auto reply = RunQueryWithPreClassify(NWorkloadManager::IQueryClassifier::TBypass()); auto status = reply->Get()->Record.GetYdbStatus(); UNIT_ASSERT_EQUAL(status, Ydb::StatusIds::SUCCESS); } Y_UNIT_TEST(ShouldResolveDefaultOnPreClassify) { - auto resolve = NWorkload::IQueryClassifier::TResolvedPoolId{ + auto resolve = NWorkloadManager::IQueryClassifier::TResolvedPoolId{ .PoolId = NResourcePool::DEFAULT_POOL_ID }; auto reply = RunQueryWithPreClassify(resolve); @@ -119,7 +119,7 @@ Y_UNIT_TEST_SUITE(KqpQueryPreClassifier) { } Y_UNIT_TEST(ShouldRejectOnPreClassify) { - auto reject = NWorkload::IQueryClassifier::TReject{ + auto reject = NWorkloadManager::IQueryClassifier::TReject{ .Code = Ydb::StatusIds::ABORTED, .Message = "Reject by ShouldRejectOnPreClassify" }; @@ -132,14 +132,14 @@ Y_UNIT_TEST_SUITE(KqpQueryPreClassifier) { Y_UNIT_TEST_SUITE(KqpQueryPostClassifier) { Y_UNIT_TEST(ShouldBypassOnPostClassify) { - auto reply = RunQueryWithPostClassify(NWorkload::IQueryClassifier::TBypass()); + auto reply = RunQueryWithPostClassify(NWorkloadManager::IQueryClassifier::TBypass()); auto status = reply->Get()->Record.GetYdbStatus(); UNIT_ASSERT_EQUAL(status, Ydb::StatusIds::SUCCESS); } Y_UNIT_TEST(ShouldResolveDefaultOnPostClassify) { - auto resolve = NWorkload::IQueryClassifier::TResolvedPoolId{ + auto resolve = NWorkloadManager::IQueryClassifier::TResolvedPoolId{ .PoolId = NResourcePool::DEFAULT_POOL_ID }; auto reply = RunQueryWithPostClassify(resolve); @@ -147,7 +147,7 @@ Y_UNIT_TEST_SUITE(KqpQueryPostClassifier) { } Y_UNIT_TEST(ShouldRejectOnPostClassify) { - auto reject = NWorkload::IQueryClassifier::TReject{ + auto reject = NWorkloadManager::IQueryClassifier::TReject{ .Code = Ydb::StatusIds::ABORTED, .Message = "Rejected by ShouldRejectOnPostClassify" }; @@ -158,4 +158,4 @@ Y_UNIT_TEST_SUITE(KqpQueryPostClassifier) { } } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/ut/kqp_stream_query_classification_ut.cpp b/ydb/services/workload_manager/ut/stream_query_classification_ut.cpp index 45ccc9a364f..f7d4dda2c07 100644 --- a/ydb/core/kqp/workload_service/ut/kqp_stream_query_classification_ut.cpp +++ b/ydb/services/workload_manager/ut/stream_query_classification_ut.cpp @@ -1,25 +1,25 @@ -#include <ydb/core/kqp/workload_service/ut/common/kqp_workload_service_ut_common.h> +#include <ydb/services/workload_manager/ut/common/workload_service_ut_common.h> #include <library/cpp/testing/unittest/registar.h> #include <chrono> #include <thread> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { -using namespace NWorkload; +using namespace NWorkloadManager; using namespace NYdb; namespace { -TIntrusivePtr<NWorkload::IYdbSetup> MakeStreamingYdb() { - return NWorkload::TYdbSetupSettings() +TIntrusivePtr<IYdbSetup> MakeStreamingYdb() { + return TYdbSetupSettings() .EnableHasPredicatesInResourcePoolClassifiers(true) .Create([](auto) {}); } -void CreateTopic(TIntrusivePtr<NWorkload::IYdbSetup> ydb, TString name) { +void CreateTopic(TIntrusivePtr<IYdbSetup> ydb, TString name) { const auto& result = ydb->ExecuteQuery(TStringBuilder() << "CREATE TOPIC " << name); UNIT_ASSERT_VALUES_EQUAL_C(result.GetStatus(), NYdb::EStatus::SUCCESS, result.GetIssues().ToOneLineString()); } @@ -112,4 +112,4 @@ Y_UNIT_TEST_SUITE(StreamingQueryClassification) { } } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/ut/kqp_workload_service_actors_ut.cpp b/ydb/services/workload_manager/ut/workload_service_actors_ut.cpp index acd8278a23c..71a2e363c5f 100644 --- a/ydb/core/kqp/workload_service/ut/kqp_workload_service_actors_ut.cpp +++ b/ydb/services/workload_manager/ut/workload_service_actors_ut.cpp @@ -1,15 +1,15 @@ #include <ydb/core/base/appdata_fwd.h> -#include <ydb/core/kqp/common/simple/services.h> -#include <ydb/core/kqp/workload_service/actors/actors.h> -#include <ydb/core/kqp/workload_service/ut/common/kqp_workload_service_ut_common.h> +#include <ydb/services/workload_manager/service/service.h> +#include <ydb/services/workload_manager/actors/actors.h> +#include <ydb/services/workload_manager/ut/common/workload_service_ut_common.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { namespace { -using namespace NWorkload; +using namespace NWorkloadManager; TEvPrivate::TEvFetchPoolResponse::TPtr FetchPool(TIntrusivePtr<IYdbSetup> ydb, const TString& poolId = "", const TString& userSID = "user@" BUILTIN_SYSTEM_DOMAIN) { @@ -171,7 +171,7 @@ Y_UNIT_TEST_SUITE(KqpWorkloadServiceSubscriptions) { auto& runtime = *ydb->GetRuntime(); const auto& edgeActor = runtime.AllocateEdgeActor(); - runtime.Send(MakeKqpWorkloadServiceId(runtime.GetNodeId()), edgeActor, new TEvSubscribeOnPoolChanges(settings.DomainName_, settings.PoolId_)); + runtime.Send(MakeServiceId(runtime.GetNodeId()), edgeActor, new TEvSubscribeOnPoolChanges(settings.DomainName_, settings.PoolId_)); const auto& response = runtime.GrabEdgeEvent<TEvUpdatePoolInfo>(edgeActor, FUTURE_WAIT_TIMEOUT); UNIT_ASSERT_C(response, "Subscription update not found"); @@ -253,4 +253,4 @@ Y_UNIT_TEST_SUITE(KqpWorkloadServiceSubscriptions) { } } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/ut/kqp_workload_service_query_sessions_ut.cpp b/ydb/services/workload_manager/ut/workload_service_query_sessions_ut.cpp index e89e2d64e63..53adbee7fe0 100644 --- a/ydb/core/kqp/workload_service/ut/kqp_workload_service_query_sessions_ut.cpp +++ b/ydb/services/workload_manager/ut/workload_service_query_sessions_ut.cpp @@ -1,16 +1,16 @@ #include <fmt/format.h> -#include <ydb/core/kqp/common/events/workload_service.h> -#include <ydb/core/kqp/common/simple/services.h> -#include <ydb/core/kqp/proxy_service/kqp_session_state.h> -#include <ydb/core/kqp/workload_service/ut/common/kqp_workload_service_ut_common.h> +#include <ydb/services/workload_manager/events.h> +#include <ydb/services/workload_manager/service/service.h> +#include <ydb/services/workload_manager/session_updater.h> +#include <ydb/services/workload_manager/ut/common/workload_service_ut_common.h> #include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/table/table.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { namespace { -using namespace NWorkload; +using namespace NWorkloadManager; using namespace NYdb; using namespace NActors; @@ -62,13 +62,13 @@ public: STFUNC(StateWork) { switch (ev->GetTypeRewrite()) { - hFunc(NWorkload::TEvContinueRequest, HandleContinueRequest); + hFunc(NWorkloadManager::TEvContinueRequest, HandleContinueRequest); default: Send(ev->Forward(SessionActorId)); } } - void HandleContinueRequest(NWorkload::TEvContinueRequest::TPtr& ev) { + void HandleContinueRequest(NWorkloadManager::TEvContinueRequest::TPtr& ev) { senderId = ev->Sender; Send(new IEventHandle(EdgeActorId, SelfId(), ev->Release().Release())); Become(&TSessionProxyActor::StateWait); @@ -76,13 +76,13 @@ public: STFUNC(StateWait) { switch (ev->GetTypeRewrite()) { - hFunc(NWorkload::TEvContinueRequest, HandleRelease); + hFunc(NWorkloadManager::TEvContinueRequest, HandleRelease); default: Send(ev->Forward(SessionActorId)); } } - void HandleRelease(NWorkload::TEvContinueRequest::TPtr& ev) { + void HandleRelease(NWorkloadManager::TEvContinueRequest::TPtr& ev) { Send(new IEventHandle(SessionActorId, senderId, ev->Release().Release())); } @@ -136,17 +136,17 @@ public: STFUNC(StateWork) { switch (ev->GetTypeRewrite()) { - hFunc(NWorkload::TEvPlaceRequestIntoPool, HandlePlaceRequest); + hFunc(NWorkloadManager::TEvPlaceRequestIntoPool, HandlePlaceRequest); default: Send(ev->Forward(WorkloadServiceId)); } } - void HandlePlaceRequest(NWorkload::TEvPlaceRequestIntoPool::TPtr& ev) { + void HandlePlaceRequest(NWorkloadManager::TEvPlaceRequestIntoPool::TPtr& ev) { auto* msg = ev->Get(); auto wrapper = std::make_shared<TWmSessionUpdaterWrapper>(FinalState, msg->WmSessionUpdater); - auto* proxyMsg = new NWorkload::TEvPlaceRequestIntoPool( + auto* proxyMsg = new NWorkloadManager::TEvPlaceRequestIntoPool( msg->DatabaseId, msg->SessionId, msg->PoolId, @@ -256,7 +256,7 @@ public: .Create(); auto& runtime = *Ydb->GetRuntime(); - auto workloadServiceId = MakeKqpWorkloadServiceId(runtime.GetNodeId(0)); + auto workloadServiceId = MakeServiceId(runtime.GetNodeId(0)); auto realWorkloadServiceId = runtime.GetLocalServiceId(workloadServiceId); auto proxyActor = new TKqpWorkloadProxyActor(State, realWorkloadServiceId, Rules); @@ -309,7 +309,7 @@ Y_UNIT_TEST_SUITE(KqpWorkloadServiceQuerySessions) { // Stop a request before a real execution to prevent read races from a .sys/query_sessions // The request is now 'parked' in the interceptor actor - auto ev = runtime->GrabEdgeEvent<NWorkload::TEvContinueRequest>(edge); + auto ev = runtime->GrabEdgeEvent<NWorkloadManager::TEvContinueRequest>(edge); TQuerySessionReader reader(f.GetYdb()); reader.FetchAll(query); @@ -374,7 +374,7 @@ Y_UNIT_TEST_SUITE(KqpWorkloadServiceQuerySessions) { // Stop a first request before a real execution to fill limits auto hangingRequest = f.GetYdb()->ExecuteQueryAsync(qHanging, myPool); - auto evHanging = runtime.GrabEdgeEvent<NWorkload::TEvContinueRequest>(edgeHanging); + auto evHanging = runtime.GrabEdgeEvent<NWorkloadManager::TEvContinueRequest>(edgeHanging); // Run a second request which has to be placed in a Delayed Queue auto delayedRequest = f.GetYdb()->ExecuteQueryAsync(qDelayed, myPool); @@ -392,7 +392,7 @@ Y_UNIT_TEST_SUITE(KqpWorkloadServiceQuerySessions) { // Continue the first request and wait for the second request to be about to execute runtime.Send(new IEventHandle(evHanging->Sender, edgeHanging, evHanging->Release().Release())); - auto evDelayed = runtime.GrabEdgeEvent<NWorkload::TEvContinueRequest>(edgeDelayed); + auto evDelayed = runtime.GrabEdgeEvent<NWorkloadManager::TEvContinueRequest>(edgeDelayed); TQuerySessionReader reader2(f.GetYdb()); reader2.FetchAll(qDelayed); @@ -432,7 +432,7 @@ Y_UNIT_TEST_SUITE(KqpWorkloadServiceQuerySessions) { const TString qFirst = "SELECT 11;"; TActorId edge1 = f.SetupInterceptor(qFirst); auto future1 = session.ExecuteDataQuery(qFirst, TTxControl::BeginTx().CommitTx()); - auto ev1 = runtime.GrabEdgeEvent<NWorkload::TEvContinueRequest>(edge1); + auto ev1 = runtime.GrabEdgeEvent<NWorkloadManager::TEvContinueRequest>(edge1); TQuerySessionReader reader(f.GetYdb()); reader.FetchAll(qFirst); @@ -450,7 +450,7 @@ Y_UNIT_TEST_SUITE(KqpWorkloadServiceQuerySessions) { const TString qSecond = "SELECT 12;"; TActorId edge2 = f.SetupInterceptor(qSecond); auto future2 = session.ExecuteDataQuery(qSecond, TTxControl::BeginTx().CommitTx()); - auto ev2 = runtime.GrabEdgeEvent<NWorkload::TEvContinueRequest>(edge2); + auto ev2 = runtime.GrabEdgeEvent<NWorkloadManager::TEvContinueRequest>(edge2); TQuerySessionReader reader2(f.GetYdb()); reader2.FetchAll(qSecond); @@ -467,4 +467,4 @@ Y_UNIT_TEST_SUITE(KqpWorkloadServiceQuerySessions) { } } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/ut/kqp_workload_service_tables_ut.cpp b/ydb/services/workload_manager/ut/workload_service_tables_ut.cpp index ec7fc061ad6..1f33ded14ce 100644 --- a/ydb/core/kqp/workload_service/ut/kqp_workload_service_tables_ut.cpp +++ b/ydb/services/workload_manager/ut/workload_service_tables_ut.cpp @@ -2,17 +2,17 @@ #include <ydb/core/base/path.h> #include <ydb/core/kqp/common/simple/services.h> -#include <ydb/core/kqp/workload_service/kqp_workload_service.h> -#include <ydb/core/kqp/workload_service/common/events.h> -#include <ydb/core/kqp/workload_service/tables/table_queries.h> -#include <ydb/core/kqp/workload_service/ut/common/kqp_workload_service_ut_common.h> +#include <ydb/services/workload_manager/service/service.h> +#include <ydb/services/workload_manager/common/events.h> +#include <ydb/services/workload_manager/tables/table_queries.h> +#include <ydb/services/workload_manager/ut/common/workload_service_ut_common.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { namespace { using namespace NActors; -using namespace NWorkload; +using namespace NWorkloadManager; void CheckPoolDescription(TIntrusivePtr<IYdbSetup> ydb, ui64 expectedQueueSize, ui64 expectedInFlight, TDuration leaseDuration = FUTURE_WAIT_TIMEOUT) { const auto& description = ydb->GetPoolDescription(leaseDuration); @@ -121,7 +121,7 @@ Y_UNIT_TEST_SUITE(KqpWorkloadServiceTables) { // Restart workload service ydb->StopWorkloadService(); auto runtime = ydb->GetRuntime(); - runtime->Register(CreateKqpWorkloadService(runtime->GetAppData().Counters)); + runtime->Register(CreateService(NWorkloadManager::GetWorkloadManagerCounters(runtime->GetAppData().Counters))); // Check that tables will be cleanuped ydb->WaitPoolState({.DelayedRequests = 0, .RunningRequests = 0}); @@ -176,4 +176,4 @@ Y_UNIT_TEST_SUITE(KqpWorkloadServiceTables) { } } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/core/kqp/workload_service/ut/kqp_workload_service_ut.cpp b/ydb/services/workload_manager/ut/workload_service_ut.cpp index 85490c90c74..d48bc6a6fad 100644 --- a/ydb/core/kqp/workload_service/ut/kqp_workload_service_ut.cpp +++ b/ydb/services/workload_manager/ut/workload_service_ut.cpp @@ -1,13 +1,13 @@ #include <ydb/core/base/appdata_fwd.h> #include <ydb/core/base/path.h> -#include <ydb/core/kqp/workload_service/ut/common/kqp_workload_service_ut_common.h> +#include <ydb/services/workload_manager/ut/common/workload_service_ut_common.h> -namespace NKikimr::NKqp { +namespace NKikimr::NWorkloadManager { namespace { -using namespace NWorkload; +using namespace NWorkloadManager; using namespace NYdb; @@ -888,11 +888,11 @@ Y_UNIT_TEST_SUITE(ResourcePoolClassifiersDdl) { } void WaitForFail(TIntrusivePtr<IYdbSetup> ydb, const TQueryRunnerSettings& settings, const TString& poolId) { - NWorkload::WaitForClassifierFail(ydb, settings, poolId); + NWorkloadManager::WaitForClassifierFail(ydb, settings, poolId); } void WaitForSuccess(TIntrusivePtr<IYdbSetup> ydb, const TQueryRunnerSettings& settings) { - NWorkload::WaitForClassifierSuccess(ydb, settings); + NWorkloadManager::WaitForClassifierSuccess(ydb, settings); } Y_UNIT_TEST(TestCreateResourcePoolClassifier) { @@ -2227,4 +2227,4 @@ Y_UNIT_TEST_SUITE(RejectActionFeatureFlagEnabled) { } } -} // namespace NKikimr::NKqp +} // namespace NKikimr::NWorkloadManager diff --git a/ydb/services/workload_manager/ut/ya.make b/ydb/services/workload_manager/ut/ya.make new file mode 100644 index 00000000000..525a582d59a --- /dev/null +++ b/ydb/services/workload_manager/ut/ya.make @@ -0,0 +1,40 @@ +UNITTEST_FOR(ydb/services/workload_manager) + +FORK_SUBTESTS() + +SIZE(MEDIUM) +IF (SANITIZER_TYPE) + REQUIREMENTS(cpu:4) +ELSE() + REQUIREMENTS(cpu:2) +ENDIF() + +SRCS( + has_app_name_ut.cpp + action_reject_ut.cpp + has_full_scan_matcher_ut.cpp + has_full_scan_ut.cpp + has_path_ddl_ut.cpp + has_path_matcher_ut.cpp + has_path_ut.cpp + stream_query_classification_ut.cpp + has_stream_matcher_ut.cpp + has_stream_ut.cpp + member_name_ut.cpp + query_classifier_match_ut.cpp + query_classifier_ut.cpp + workload_service_actors_ut.cpp + workload_service_query_sessions_ut.cpp + workload_service_tables_ut.cpp + workload_service_ut.cpp +) + +PEERDIR( + ydb/services/workload_manager/ut/common + + yql/essentials/sql/pg_dummy +) + +YQL_LAST_ABI_VERSION() + +END() diff --git a/ydb/services/workload_manager/ya.make b/ydb/services/workload_manager/ya.make new file mode 100644 index 00000000000..ca99fdb9abe --- /dev/null +++ b/ydb/services/workload_manager/ya.make @@ -0,0 +1,38 @@ +LIBRARY() + +SRCS( + has_full_scan_matcher.cpp + has_path_matcher.cpp + has_stream_matcher.cpp + query_classifier.cpp +) + +PEERDIR( + ydb/core/cms/console + + ydb/core/kqp/common + ydb/core/kqp/query_data + + ydb/core/mind + + ydb/core/resource_pools + + ydb/library/aclib + +) + +YQL_LAST_ABI_VERSION() + +END() + +RECURSE( + actors + common + metadata_subscription + tables + service +) + +RECURSE_FOR_TESTS( + ut +) diff --git a/ydb/tests/tools/kqprun/src/actors.cpp b/ydb/tests/tools/kqprun/src/actors.cpp index 4adafc08c76..9691ca2919d 100644 --- a/ydb/tests/tools/kqprun/src/actors.cpp +++ b/ydb/tests/tools/kqprun/src/actors.cpp @@ -4,7 +4,7 @@ #include <ydb/core/kqp/common/simple/services.h> #include <ydb/core/kqp/rm_service/kqp_rm_service.h> -#include <ydb/core/kqp/workload_service/actors/actors.h> +#include <ydb/services/workload_manager/actors/actors.h> using namespace NKikimrRun; @@ -180,7 +180,7 @@ public: } } - void Handle(NKikimr::NKqp::NWorkload::TEvFetchDatabaseResponse::TPtr& ev) { + void Handle(NKikimr::NWorkloadManager::TEvFetchDatabaseResponse::TPtr& ev) { const auto status = ev->Get()->Status; if (status == Ydb::StatusIds::SUCCESS) { HealthCheckStage_ = EHealthCheck::ScriptRequest; @@ -203,7 +203,7 @@ public: STRICT_STFUNC(StateFunc, sFunc(NActors::TEvents::TEvWakeup, DoHealthCheck); - hFunc(NKikimr::NKqp::NWorkload::TEvFetchDatabaseResponse, Handle); + hFunc(NKikimr::NWorkloadManager::TEvFetchDatabaseResponse, Handle); hFunc(NKikimr::NKqp::TEvKqp::TEvScriptResponse, Handle); ) @@ -229,7 +229,7 @@ private: } void FetchDatabase() { - Register(NKikimr::NKqp::NWorkload::CreateDatabaseFetcherActor(SelfId(), Settings_.Database)); + Register(NKikimr::NWorkloadManager::CreateDatabaseFetcherActor(SelfId(), Settings_.Database)); } void StartScriptQuery() { diff --git a/ydb/tests/tools/kqprun/src/ya.make b/ydb/tests/tools/kqprun/src/ya.make index b9a88f52573..3a9ff2e3bd4 100644 --- a/ydb/tests/tools/kqprun/src/ya.make +++ b/ydb/tests/tools/kqprun/src/ya.make @@ -10,7 +10,7 @@ PEERDIR( library/cpp/protobuf/json ydb/core/client/server ydb/core/grpc_services - ydb/core/kqp/workload_service/actors + ydb/services/workload_manager/actors ydb/core/testlib ydb/core/util ydb/library/aclib |
