| 9398 | |
| 9399 | |
| 9400 | Future<Nothing> Slave::publishResources( |
| 9401 | const ContainerID& containerId, const Resources& resources) |
| 9402 | { |
| 9403 | hashset<ResourceProviderID> resourceProviderIds; |
| 9404 | foreach (const Resource& resource, resources) { |
| 9405 | if (resource.has_provider_id()) { |
| 9406 | resourceProviderIds.insert(resource.provider_id()); |
| 9407 | } |
| 9408 | } |
| 9409 | |
| 9410 | vector<Future<Nothing>> futures; |
| 9411 | foreach (const ResourceProviderID& resourceProviderId, resourceProviderIds) { |
| 9412 | auto hasResourceProviderId = [&](const Resource& resource) { |
| 9413 | return resource.has_provider_id() && |
| 9414 | resource.provider_id() == resourceProviderId; |
| 9415 | }; |
| 9416 | |
| 9417 | // NOTE: For resources providers that serve quantity-based resources without |
| 9418 | // identifier (such as cpus and mem), we cannot achieve idempotency with |
| 9419 | // diff-based resource publishing, so we have to implement the "ensure-all" |
| 9420 | // semantics, and always calculate the total resources to publish. |
| 9421 | Option<Resources> containerResources; |
| 9422 | Resources complementaryResources; |
| 9423 | foreachvalue (const Framework* framework, frameworks) { |
| 9424 | foreachvalue (const Executor* executor, framework->executors) { |
| 9425 | if (executor->containerId == containerId) { |
| 9426 | containerResources = resources.filter(hasResourceProviderId); |
| 9427 | } else { |
| 9428 | complementaryResources += |
| 9429 | executor->allocatedResources().filter(hasResourceProviderId); |
| 9430 | } |
| 9431 | } |
| 9432 | } |
| 9433 | |
| 9434 | if (containerResources.isNone()) { |
| 9435 | // NOTE: This actually should not happen, as the callers have already |
| 9436 | // ensured the existence of the executor before calling this function |
| 9437 | // synchronously. However we still treat this as a nonfatal error since |
| 9438 | // this might change in the future. |
| 9439 | LOG(WARNING) << "Ignoring publishing resources for container " |
| 9440 | << containerId << ": Executor cannot be found"; |
| 9441 | |
| 9442 | return Nothing(); |
| 9443 | } |
| 9444 | |
| 9445 | // Since we already have resources from any resource provider in the |
| 9446 | // resource pool, the resource provider manager must have been created. |
| 9447 | futures.push_back( |
| 9448 | CHECK_NOTNULL(resourceProviderManager.get()) |
| 9449 | ->publishResources(containerResources.get() + complementaryResources) |
| 9450 | .repair([=](const Future<Nothing>& future) -> Future<Nothing> { |
| 9451 | // TODO(chhsiao): Consider surfacing the set of published resources |
| 9452 | // and only fail if `published - complementaryResources` does not |
| 9453 | // contain `containerResources`. |
| 9454 | return Failure( |
| 9455 | "Failed to publish resources '" + |
| 9456 | stringify(containerResources.get()) + "' for container " + |
| 9457 | stringify(containerId) + ": " + future.failure()); |