diff --git a/source/plugins/ruby/in_kube_podinventory.rb b/source/plugins/ruby/in_kube_podinventory.rb index 584e6bfbc9..4bc378da35 100644 --- a/source/plugins/ruby/in_kube_podinventory.rb +++ b/source/plugins/ruby/in_kube_podinventory.rb @@ -56,6 +56,7 @@ def initialize @watchWinNodesThread = nil @windowsNodeNameListCache = [] + @podCacheRefreshRequested = false @windowsContainerRecordsCacheSizeBytes = 0 @kubeservicesTag = "oneagent.containerInsights.KUBE_SERVICES_BLOB" @@ -742,25 +743,22 @@ def getServiceNameFromLabels(namespace, labels, serviceRecords) def watch_pods $log.info("in_kube_podinventory::watch_pods:Start @ #{Time.now.utc.iso8601}") + return if @finished + podsResourceVersion = nil - # invoke getWindowsNodes to handle scenario where windowsNodeNameCache not populated yet on containerstart - winNodes = KubernetesApiClient.getWindowsNodesArray() - if winNodes.length > 0 - @windowsNodeNameCacheMutex.synchronize { - @windowsNodeNameListCache = winNodes.dup - } - end loop do begin + @windowsNodeNameCacheMutex.synchronize { + if @podCacheRefreshRequested + podsResourceVersion = nil + @podCacheRefreshRequested = false + end + } if podsResourceVersion.nil? # clear cache before filling the cache with list @podCacheMutex.synchronize { @podItemsCache.clear() } - currentWindowsNodeNameList = [] - @windowsNodeNameCacheMutex.synchronize { - currentWindowsNodeNameList = @windowsNodeNameListCache.dup - } continuationToken = nil resourceUri = "pods?limit=#{@PODS_CHUNK_SIZE}" $log.info("in_kube_podinventory::watch_pods:Getting pods from Kube API: #{resourceUri} @ #{Time.now.utc.iso8601}") @@ -777,13 +775,9 @@ def watch_pods key = item["metadata"]["uid"] if !key.nil? && !key.empty? nodeName = (!item["spec"].nil? && !item["spec"]["nodeName"].nil?) ? item["spec"]["nodeName"] : "" - isWindowsPodItem = false - if !nodeName.empty? && - !currentWindowsNodeNameList.nil? && - !currentWindowsNodeNameList.empty? && - currentWindowsNodeNameList.include?(nodeName) - isWindowsPodItem = true - end + isWindowsPodItem = @windowsNodeNameCacheMutex.synchronize { + @windowsNodeNameListCache.include?(nodeName) + } podItem = KubernetesApiClient.getOptimizedItem("pods", item, isWindowsPodItem) if !podItem.nil? && !podItem.empty? @podCacheMutex.synchronize { @@ -816,13 +810,9 @@ def watch_pods key = item["metadata"]["uid"] if !key.nil? && !key.empty? nodeName = (!item["spec"].nil? && !item["spec"]["nodeName"].nil?) ? item["spec"]["nodeName"] : "" - isWindowsPodItem = false - if !nodeName.empty? && - !currentWindowsNodeNameList.nil? && - !currentWindowsNodeNameList.empty? && - currentWindowsNodeNameList.include?(nodeName) - isWindowsPodItem = true - end + isWindowsPodItem = @windowsNodeNameCacheMutex.synchronize { + @windowsNodeNameListCache.include?(nodeName) + } podItem = KubernetesApiClient.getOptimizedItem("pods", item, isWindowsPodItem) if !podItem.nil? && !podItem.empty? @podCacheMutex.synchronize { @@ -843,6 +833,8 @@ def watch_pods end end end + next if @windowsNodeNameCacheMutex.synchronize { @podCacheRefreshRequested } + if podsResourceVersion.nil? || podsResourceVersion.empty? || podsResourceVersion == "0" # https://github.com/kubernetes/kubernetes/issues/74022 $log.warn("in_kube_podinventory::watch_pods:received podsResourceVersion either nil or empty or 0 @ #{Time.now.utc.iso8601}") @@ -856,6 +848,8 @@ def watch_pods $log.warn("in_kube_podinventory::watch_pods:watch API returned nil watcher for watch connection with resource version: #{podsResourceVersion} @ #{Time.now.utc.iso8601}") else watcher.each do |notice| + break if @windowsNodeNameCacheMutex.synchronize { @podCacheRefreshRequested } + case notice["type"] when "ADDED", "MODIFIED", "DELETED", "BOOKMARK" item = notice["object"] @@ -874,18 +868,10 @@ def watch_pods if ((notice["type"] == "ADDED") || (notice["type"] == "MODIFIED")) key = item["metadata"]["uid"] if !key.nil? && !key.empty? - currentWindowsNodeNameList = [] - @windowsNodeNameCacheMutex.synchronize { - currentWindowsNodeNameList = @windowsNodeNameListCache.dup - } - isWindowsPodItem = false nodeName = (!item["spec"].nil? && !item["spec"]["nodeName"].nil?) ? item["spec"]["nodeName"] : "" - if !nodeName.empty? && - !currentWindowsNodeNameList.nil? && - !currentWindowsNodeNameList.empty? && - currentWindowsNodeNameList.include?(nodeName) - isWindowsPodItem = true - end + isWindowsPodItem = @windowsNodeNameCacheMutex.synchronize { + @windowsNodeNameListCache.include?(nodeName) + } podItem = KubernetesApiClient.getOptimizedItem("pods", item, isWindowsPodItem) if !podItem.nil? && !podItem.empty? @podCacheMutex.synchronize { @@ -1070,67 +1056,39 @@ def watch_windows_nodes loop do begin if nodesResourceVersion.nil? - @windowsNodeNameCacheMutex.synchronize { - @windowsNodeNameListCache.clear() - } + windowsNodeNameList = [] + listResourceVersion = nil continuationToken = nil resourceUri = KubernetesApiClient.getNodesResourceUri("nodes?labelSelector=kubernetes.io%2Fos%3Dwindows&limit=#{@NODES_CHUNK_SIZE}") - $log.info("in_kube_podinventory::watch_windows_nodes:Getting windows nodes from Kube API: #{resourceUri} @ #{Time.now.utc.iso8601}") - continuationToken, nodeInventory, responseCode = KubernetesApiClient.getResourcesAndContinuationTokenV2(resourceUri) - if responseCode.nil? || responseCode != "200" - $log.info("in_kube_podinventory::watch_windows_nodes:Getting windows nodes from Kube API: #{resourceUri} failed with statuscode: #{responseCode} @ #{Time.now.utc.iso8601}") - else - $log.info("in_kube_podinventory::watch_windows_nodes:Done getting windows nodes from Kube API @ #{Time.now.utc.iso8601}") - if (!nodeInventory.nil? && !nodeInventory.empty?) - nodesResourceVersion = nodeInventory["metadata"]["resourceVersion"] - if (nodeInventory.key?("items") && !nodeInventory["items"].nil? && !nodeInventory["items"].empty?) - $log.info("in_kube_podinventory::watch_windows_nodes: number of windows node items :#{nodeInventory["items"].length} from Kube API @ #{Time.now.utc.iso8601}") - nodeInventory["items"].each do |item| - key = item["metadata"]["name"] - if !key.nil? && !key.empty? - @windowsNodeNameCacheMutex.synchronize { - if !@windowsNodeNameListCache.include?(key) - @windowsNodeNameListCache.push(key) - end - } - else - $log.warn "in_kube_podinventory::watch_windows_nodes:Received node name either nil or empty @ #{Time.now.utc.iso8601}" - end - end - end - else - $log.warn "in_kube_podinventory::watch_windows_nodes:Received empty nodeInventory @ #{Time.now.utc.iso8601}" + begin + pageResourceUri = resourceUri + if !continuationToken.nil? && !continuationToken.empty? + pageResourceUri += "&continue=#{continuationToken}" end - while (!continuationToken.nil? && !continuationToken.empty?) - continuationToken, nodeInventory, responseCode = KubernetesApiClient.getResourcesAndContinuationTokenV2(resourceUri + "&continue=#{continuationToken}") - if responseCode.nil? || responseCode != "200" - $log.info("in_kube_podinventory::watch_windows_nodes:Getting windows nodes from Kube API: #{resourceUri}&continue=#{continuationToken} failed with statuscode: #{responseCode} @ #{Time.now.utc.iso8601}") - nodesResourceVersion = nil - break # break, if any of the pagination call failed so that full cache can be rebuild with LIST again - else - if (!nodeInventory.nil? && !nodeInventory.empty?) - nodesResourceVersion = nodeInventory["metadata"]["resourceVersion"] - if (nodeInventory.key?("items") && !nodeInventory["items"].nil? && !nodeInventory["items"].empty?) - $log.info("in_kube_podinventory::watch_windows_nodes : number of windows node items :#{nodeInventory["items"].length} from Kube API @ #{Time.now.utc.iso8601}") - nodeInventory["items"].each do |item| - key = item["metadata"]["name"] - if !key.nil? && !key.empty? - @windowsNodeNameCacheMutex.synchronize { - if !@windowsNodeNameListCache.include?(key) - @windowsNodeNameListCache.push(key) - end - } - else - $log.warn "in_kube_podinventory::watch_windows_nodes:Received node name either nil or empty @ #{Time.now.utc.iso8601}" - end - end - end - else - $log.warn "in_kube_podinventory::watch_windows_nodes:Received empty nodeInventory @ #{Time.now.utc.iso8601}" - end - end + $log.info("in_kube_podinventory::watch_windows_nodes:Getting windows nodes from Kube API: #{pageResourceUri} @ #{Time.now.utc.iso8601}") + continuationToken, nodeInventory, responseCode = KubernetesApiClient.getResourcesAndContinuationTokenV2(pageResourceUri) + unless responseCode == "200" && nodeInventory.is_a?(Hash) && + nodeInventory["metadata"].is_a?(Hash) && nodeInventory["items"].is_a?(Array) + raise "Invalid windows node LIST response for #{pageResourceUri}, statuscode: #{responseCode}" end - end + pageResourceVersion = nodeInventory["metadata"]["resourceVersion"] + unless pageResourceVersion.is_a?(String) && !pageResourceVersion.empty? && pageResourceVersion != "0" && + (listResourceVersion.nil? || listResourceVersion == pageResourceVersion) + raise "Invalid or inconsistent windows node LIST resourceVersion for #{pageResourceUri}: #{pageResourceVersion}" + end + listResourceVersion = pageResourceVersion + nodeInventory["items"].each do |item| + nodeName = item["metadata"]["name"] + raise "Invalid windows node name in LIST response for #{pageResourceUri}" unless nodeName.is_a?(String) && !nodeName.empty? + windowsNodeNameList.push(nodeName) + end + end while !continuationToken.nil? && !continuationToken.empty? + windowsNodeNameList.uniq! + @windowsNodeNameCacheMutex.synchronize { + @podCacheRefreshRequested = true unless (windowsNodeNameList - @windowsNodeNameListCache).empty? + @windowsNodeNameListCache = windowsNodeNameList + } + nodesResourceVersion = listResourceVersion end if nodesResourceVersion.nil? || nodesResourceVersion.empty? || nodesResourceVersion == "0" # https://github.com/kubernetes/kubernetes/issues/74022 @@ -1165,6 +1123,7 @@ def watch_windows_nodes @windowsNodeNameCacheMutex.synchronize { if !@windowsNodeNameListCache.include?(key) @windowsNodeNameListCache.push(key) + @podCacheRefreshRequested = true end } elsif notice["type"] == "DELETED" @@ -1196,6 +1155,7 @@ def watch_windows_nodes rescue => errorStr $log.warn("in_kube_podinventory::watch_windows_nodes:failed with an error: #{errorStr} @ #{Time.now.utc.iso8601}") nodesResourceVersion = nil + sleep(30) end end $log.info("in_kube_podinventory::watch_windows_nodes:End @ #{Time.now.utc.iso8601}") diff --git a/source/plugins/ruby/in_kube_podinventory_test.rb b/source/plugins/ruby/in_kube_podinventory_test.rb new file mode 100644 index 0000000000..3bffe99b3b --- /dev/null +++ b/source/plugins/ruby/in_kube_podinventory_test.rb @@ -0,0 +1,301 @@ +require "minitest/autorun" +require "fluent/test" +require "logger" +require "net/http" +require "time" +require "timeout" +require_relative "in_kube_podinventory" + +class InKubePodInventoryTests < Minitest::Test + Watch = Struct.new(:events) do + def each + events.each do |event| + yield(event.respond_to?(:call) ? event.call : event) + end + end + + def finish + end + end + + class ApiClient + attr_accessor :responses, :node_events, :pod_events + attr_reader :requests, :snapshots, :classifications, :watches + + def initialize(&cache_reader) + @cache_reader = cache_reader + @responses = [] + @node_events = [] + @pod_events = [] + @requests = [] + @snapshots = [] + @classifications = {} + @watches = [] + end + + def getNodesResourceUri(uri) + uri + end + + def getResourcesAndContinuationTokenV2(uri) + @requests << uri + @snapshots << @cache_reader.call + response = @responses.shift + raise Minitest::Assertion, "Unexpected API request: #{uri}" unless response + response.respond_to?(:call) ? response.call(uri) : response + end + + def getWindowsNodesArray + raise Minitest::Assertion, "Pod watcher must not fetch a separate node snapshot" + end + + def getOptimizedItem(_resource, item, is_windows) + @classifications[item["metadata"]["uid"]] = is_windows + item + end + + def watch(resource, **options) + @watches << [resource, options] + Watch.new(resource == "nodes" ? @node_events : @pod_events) + end + end + + def setup + Fluent::Test.setup + @previous_log = $log + $log = Logger.new(File::NULL) + @plugin_class = Fluent::Plugin::Kube_PodInventory_Input + @plugin = @plugin_class.allocate + @plugin.instance_variable_set(:@windowsNodeNameListCache, ["win-cached"]) + @plugin.instance_variable_set(:@podCacheRefreshRequested, false) + @plugin.instance_variable_set(:@windowsNodeNameCacheMutex, Mutex.new) + @plugin.instance_variable_set(:@finished, false) + @plugin.instance_variable_set(:@podItemsCache, {}) + @plugin.instance_variable_set(:@podCacheMutex, Mutex.new) + @plugin.instance_variable_set(:@NODES_CHUNK_SIZE, 2) + @plugin.instance_variable_set(:@PODS_CHUNK_SIZE, 2) + @plugin.define_singleton_method(:loop) { |&iteration| iteration.call } + retry_delays = @retry_delays = [] + @plugin.define_singleton_method(:sleep) { |duration| retry_delays << duration } + @api = ApiClient.new { cached_nodes } + @plugin_class.const_set(:KubernetesApiClient, @api) + end + + def teardown + @plugin_class.send(:remove_const, :KubernetesApiClient) + $log = @previous_log + end + + def cached_nodes + @plugin.instance_variable_get(:@windowsNodeNameListCache).dup + end + + def node_inventory(names, version = "10") + { + "metadata" => { "resourceVersion" => version }, + "items" => names.map { |name| { "metadata" => { "name" => name } } }, + } + end + + def pod_inventory(node_name = "win-new") + { + "metadata" => { "resourceVersion" => "11" }, + "items" => [{ "metadata" => { "uid" => "pod-1" }, "spec" => { "nodeName" => node_name } }], + } + end + + def node_event(type, name) + { "type" => type, "object" => { "metadata" => { "name" => name, "resourceVersion" => "11" } } } + end + + def test_paginated_node_list_is_published_atomically + @api.responses = [["next-page", node_inventory(["win-a"]), "200"], [nil, node_inventory(["win-b"]), "200"]] + + @plugin.watch_windows_nodes + + assert_equal [["win-cached"], ["win-cached"]], @api.snapshots + assert_equal ["win-a", "win-b"], cached_nodes + assert_equal @api.requests.first + "&continue=next-page", @api.requests.last + assert_equal "10", @api.watches.first.last[:resource_version] + assert @plugin.instance_variable_get(:@podCacheRefreshRequested) + end + + def test_valid_empty_node_list_clears_cache_without_requesting_a_relist + @api.responses = [[nil, node_inventory([]), "200"]] + + @plugin.watch_windows_nodes + + assert_empty cached_nodes + refute @plugin.instance_variable_get(:@podCacheRefreshRequested) + assert_equal "nodes", @api.watches.first.first + end + + def test_invalid_first_page_preserves_cache_without_requesting_a_relist + @api.responses = [[nil, nil, "503"]] + + @plugin.watch_windows_nodes + + assert_equal ["win-cached"], cached_nodes + refute @plugin.instance_variable_get(:@podCacheRefreshRequested) + assert_empty @api.watches + assert_equal [30], @retry_delays + end + + def test_invalid_later_pages_preserve_complete_cache + invalid_responses = [ + [nil, nil, "503"], + [nil, nil, "200"], + [nil, {}, "200"], + [nil, { "metadata" => { "resourceVersion" => "10" } }, "200"], + [nil, { "items" => [] }, "200"], + [nil, node_inventory([], nil), "200"], + [nil, node_inventory([], ""), "200"], + [nil, node_inventory([], "0"), "200"], + [nil, node_inventory(["win-b"], "11"), "200"], + [nil, node_inventory([nil]), "200"], + ] + invalid_responses.each do |response| + @api.responses = [["next-page", node_inventory(["win-a"]), "200"], response] + + @plugin.watch_windows_nodes + + assert_equal ["win-cached"], cached_nodes, "Invalid page: #{response.inspect}" + assert_empty @api.watches + refute @plugin.instance_variable_get(:@podCacheRefreshRequested) + end + assert_equal [30] * invalid_responses.length, @retry_delays + end + + def test_node_watch_addition_during_pod_list_is_preserved + @api.node_events = [node_event("ADDED", "win-new")] + @api.responses = [lambda { |_uri| @plugin.watch_windows_nodes; [nil, pod_inventory, "200"] }, [nil, node_inventory(["win-cached"]), "200"]] + + @plugin.watch_pods + + assert_equal ["win-cached", "win-new"], cached_nodes + assert_equal true, @api.classifications.fetch("pod-1") + end + + def test_node_watch_deletion_during_pod_list_is_preserved + @api.node_events = [node_event("DELETED", "win-cached")] + @api.responses = [lambda { |_uri| @plugin.watch_windows_nodes; [nil, pod_inventory("win-cached"), "200"] }, [nil, node_inventory(["win-cached"]), "200"]] + + @plugin.watch_pods + + assert_empty cached_nodes + assert_equal false, @api.classifications.fetch("pod-1") + end + + def test_pod_relists_do_not_fetch_or_replace_the_node_cache + original_cache = @plugin.instance_variable_get(:@windowsNodeNameListCache) + @plugin.define_singleton_method(:loop) { |&iteration| 2.times(&iteration) } + @api.pod_events = [{ "type" => "ERROR", "object" => {} }] + @api.responses = [[nil, pod_inventory("win-cached"), "200"], [nil, pod_inventory("win-cached"), "200"]] + + @plugin.watch_pods + + assert_equal ["pods?limit=2", "pods?limit=2"], @api.requests + assert_same original_cache, @plugin.instance_variable_get(:@windowsNodeNameListCache) + assert_equal true, @api.classifications.fetch("pod-1") + end + + def test_linux_pods_start_before_windows_node_discovery + @plugin.instance_variable_set(:@windowsNodeNameListCache, []) + linux_pods = pod_inventory("linux-node") + @api.responses = [[nil, linux_pods, "200"]] + + Timeout.timeout(5) { @plugin.watch_pods } + + assert_equal ["pods?limit=2"], @api.requests + assert_equal false, @api.classifications.fetch("pod-1") + assert_equal linux_pods["items"].first, @plugin.instance_variable_get(:@podItemsCache).fetch("pod-1") + end + + def test_linux_pods_continue_after_windows_node_discovery_fails + @plugin.instance_variable_set(:@windowsNodeNameListCache, []) + 2.times do + @api.responses = [[nil, nil, "503"]] + @plugin.watch_windows_nodes + end + @api.responses = [[nil, pod_inventory("linux-node"), "200"]] + + Timeout.timeout(5) { @plugin.watch_pods } + + assert_equal false, @api.classifications.fetch("pod-1") + assert_equal [30, 30], @retry_delays + assert_equal "pods", @api.watches.last.first + assert_equal ["pod-1"], @plugin.instance_variable_get(:@podItemsCache).keys + end + + def test_linux_only_cluster_does_not_relist_pods_after_empty_node_discovery + @plugin.instance_variable_set(:@windowsNodeNameListCache, []) + @plugin.define_singleton_method(:loop) { |&iteration| 2.times(&iteration) } + @api.responses = [[nil, node_inventory([]), "200"]] + @plugin.watch_windows_nodes + @api.responses = [[nil, pod_inventory("linux-node"), "200"]] + + @plugin.watch_pods + + assert_equal 1, @api.requests.count { |uri| uri.start_with?("pods?") } + assert_equal 2, @api.watches.count { |resource, _options| resource == "pods" } + assert_equal false, @api.classifications.fetch("pod-1") + end + + def test_windows_node_discovery_during_pod_watch_triggers_recovery + @plugin.instance_variable_set(:@windowsNodeNameListCache, []) + @plugin.define_singleton_method(:loop) { |&iteration| 2.times(&iteration) } + linux_pod = pod_inventory("linux-node")["items"].first + linux_pod["metadata"]["uid"] = "pod-linux" + mixed_pods = pod_inventory + mixed_pods["items"] << linux_pod + @api.responses = [[nil, mixed_pods, "200"], [nil, node_inventory(["win-new"]), "200"], [nil, mixed_pods, "200"]] + @api.pod_events = [lambda { + assert_equal false, @api.classifications.fetch("pod-1") + @plugin.watch_windows_nodes + @api.pod_events = [] + { "type" => "BOOKMARK", "object" => { "metadata" => { "resourceVersion" => "12" } } } + }] + + @plugin.watch_pods + + assert_equal true, @api.classifications.fetch("pod-1") + assert_equal false, @api.classifications.fetch("pod-linux") + assert_equal 2, @api.requests.count { |uri| uri.start_with?("pods?") } + assert_equal linux_pod, @plugin.instance_variable_get(:@podItemsCache).fetch("pod-linux") + refute @plugin.instance_variable_get(:@podCacheRefreshRequested) + end + + def test_windows_node_discovery_recovers_pods_after_idle_watch_timeout + @plugin.instance_variable_set(:@windowsNodeNameListCache, []) + @plugin.define_singleton_method(:loop) { |&iteration| 2.times(&iteration) } + @api.responses = [[nil, pod_inventory, "200"], [nil, node_inventory(["win-new"]), "200"], [nil, pod_inventory, "200"]] + @api.pod_events = [lambda { + @plugin.watch_windows_nodes + @api.pod_events = [] + raise Net::ReadTimeout + }] + + @plugin.watch_pods + + assert_equal true, @api.classifications.fetch("pod-1") + assert_equal 2, @api.requests.count { |uri| uri.start_with?("pods?") } + end + + def test_unchanged_windows_nodes_do_not_request_another_pod_relist + @api.responses = [[nil, node_inventory(["win-cached"]), "200"]] + @api.node_events = [node_event("ADDED", "win-cached")] + + @plugin.watch_windows_nodes + + refute @plugin.instance_variable_get(:@podCacheRefreshRequested) + end + + def test_pod_watcher_exits_if_shutdown_was_requested + @plugin.instance_variable_set(:@finished, true) + + @plugin.watch_pods + + assert_empty @api.requests + assert_empty @api.watches + end +end \ No newline at end of file