MCPcopy Create free account
hub / github.com/apache/mesos / _subscribe

Method _subscribe

src/resource_provider/manager.cpp:847–930  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

845
846
847void ResourceProviderManagerProcess::_subscribe(
848 const Future<bool>& admitResourceProvider,
849 Owned<ResourceProvider> resourceProvider)
850{
851 if (!admitResourceProvider.isReady()) {
852 LOG(INFO)
853 << "Not subscribing resource provider " << resourceProvider->info.id()
854 << " as registry update did not succeed: " << admitResourceProvider;
855
856 return;
857 }
858
859 CHECK(admitResourceProvider.get())
860 << "Could not admit resource provider " << resourceProvider->info.id()
861 << " as registry update was rejected";
862
863 const ResourceProviderID& resourceProviderId = resourceProvider->info.id();
864
865 resourceProvider->http.closed()
866 .onAny(defer(self(), [=](const Future<Nothing>& future) {
867 // Iff the remote side closes the HTTP connection, the future will be
868 // ready. We will remove the resource provider in that case.
869 // This side closes the HTTP connection only when removing a resource
870 // provider, therefore we shouldn't try to remove it again here.
871 if (future.isReady()) {
872 CHECK(resourceProviders.subscribed.contains(resourceProviderId));
873
874 // NOTE: All pending futures of publish requests for the resource
875 // provider will become failed.
876 resourceProviders.subscribed.erase(resourceProviderId);
877 }
878
879 ResourceProviderMessage::Disconnect disconnect{resourceProviderId};
880
881 ResourceProviderMessage message;
882 message.type = ResourceProviderMessage::Type::DISCONNECT;
883 message.disconnect = std::move(disconnect);
884
885 messages.put(std::move(message));
886
887 ++metrics.disconnections;
888 }));
889
890 if (!resourceProviders.known.contains(resourceProviderId)) {
891 mesos::resource_provider::registry::ResourceProvider resourceProvider_ =
892 createRegistryResourceProvider(resourceProvider->info);
893
894 resourceProviders.known.put(
895 resourceProviderId,
896 std::move(resourceProvider_));
897 }
898
899 ResourceProviderMessage::Subscribe subscribe{resourceProvider->info};
900
901 ResourceProviderMessage message;
902 message.type = ResourceProviderMessage::Type::SUBSCRIBE;
903 message.subscribe = std::move(subscribe);
904

Callers

nothing calls this directly

Calls 11

deferFunction · 0.85
isReadyMethod · 0.80
CopyFromMethod · 0.80
sendMethod · 0.65
idMethod · 0.45
getMethod · 0.45
closedMethod · 0.45
containsMethod · 0.45
eraseMethod · 0.45
putMethod · 0.45

Tested by

no test coverage detected