| 845 | |
| 846 | |
| 847 | void 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 |
nothing calls this directly
no test coverage detected