Skip to content
9 changes: 5 additions & 4 deletions .pipelines/azure_pipeline_mergedbranches.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -1072,7 +1072,7 @@ extends:
azureClientId: $(AksGenevaIntegrationMultiTenancyClientId)
azureTenantId: $(CI_BUILD_AZURE_TENANT_ID)
teamsWebhookUri: $(TeamsWebhookUri)
additionalTestParams: 'GenevaIntegration=true'
additionalTestParams: 'GenevaIntegration=true AgentTelemetryResourceId=$(AGENT_TELEMETRY_RESOURCE_ID) AgentTelemetryVersion=$(linuxImageTagUnderTest)'

# ============================================================
# Cluster: ci-logs-prod-aks-work-load-identity — Deploy via Helm
Expand Down Expand Up @@ -1101,7 +1101,7 @@ extends:
azureClientId: $(AksWorkLoadIdentityClientId)
azureTenantId: $(CI_BUILD_AZURE_TENANT_ID)
teamsWebhookUri: $(TeamsWebhookUri)
additionalTestParams: 'LinuxTestsOnly=true'
additionalTestParams: 'LinuxTestsOnly=true AgentTelemetryResourceId=$(AGENT_TELEMETRY_RESOURCE_ID) AgentTelemetryVersion=$(linuxImageTagUnderTest)'

# ============================================================
# Cluster: ci-logs-prod-wcus-fips — Deploy via Helm
Expand Down Expand Up @@ -1130,6 +1130,7 @@ extends:
azureClientId: $(WcusFipsClientId)
azureTenantId: $(CI_BUILD_AZURE_TENANT_ID)
teamsWebhookUri: $(TeamsWebhookUri)
additionalTestParams: 'AgentTelemetryResourceId=$(AGENT_TELEMETRY_RESOURCE_ID) AgentTelemetryVersion=$(linuxImageTagUnderTest)'

# ============================================================
# Cluster: ci-logs-prod-aks-networkflowlogs — Deploy via Helm
Expand Down Expand Up @@ -1159,7 +1160,7 @@ extends:
azureClientId: $(NetworkFlowLogsClientId)
azureTenantId: $(CI_BUILD_AZURE_TENANT_ID)
teamsWebhookUri: $(TeamsWebhookUri)
additionalTestParams: 'LinuxTestsOnly=true'
additionalTestParams: 'LinuxTestsOnly=true AgentTelemetryResourceId=$(AGENT_TELEMETRY_RESOURCE_ID) AgentTelemetryVersion=$(linuxImageTagUnderTest)'

# ============================================================
# Cluster: ci-logs-dev-aks-std-prof-config-test1 — Deploy via Helm
Expand Down Expand Up @@ -1211,4 +1212,4 @@ extends:
azureClientId: $(AllNodesClientId)
azureTenantId: $(CI_BUILD_AZURE_TENANT_ID)
teamsWebhookUri: $(TeamsWebhookUri)
additionalTestParams: 'PerNodeLogCoverage=true'
additionalTestParams: 'PerNodeLogCoverage=true AgentTelemetryResourceId=$(AGENT_TELEMETRY_RESOURCE_ID) AgentTelemetryVersion=$(linuxImageTagUnderTest)'
4 changes: 4 additions & 0 deletions test/ginkgo-e2e/querylogs/querylogs_suite_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,8 @@ var AKSResourceId string
var RetinaNetworkFlowLogsEnabled string
var GenevaIntegrationEnabled string
var PerNodeLogCoverageEnabled string
var AgentTelemetryResourceId string
var AgentTelemetryVersion string
var Cfg *rest.Config

func TestQuerylogs(t *testing.T) {
Expand All @@ -36,6 +38,8 @@ var _ = BeforeSuite(func() {
Expect(err).NotTo(HaveOccurred())
GenevaIntegrationEnabled = os.Getenv("GENEVA_INTEGRATION")
PerNodeLogCoverageEnabled = os.Getenv("PER_NODE_LOG_COVERAGE")
AgentTelemetryResourceId = os.Getenv("AGENT_TELEMETRY_RESOURCE_ID")
AgentTelemetryVersion = os.Getenv("AGENT_TELEMETRY_VERSION")
LogsClient, err = utils.SetupLogsClient()
Expect(err).NotTo(HaveOccurred())
})
Expand Down
29 changes: 29 additions & 0 deletions test/ginkgo-e2e/querylogs/querylogs_test.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package querylogs_test

import (
"fmt"
"strings"

. "github.com/onsi/ginkgo/v2"
Expand Down Expand Up @@ -86,3 +87,31 @@ var _ = Describe("When querying the number of resources of the cluster", func()
Entry("Nodes", "KubeNodeInventory"),
)
})

var _ = Describe("When querying the agent telemetry", func() {
DescribeTable("Every node should report telemetry from the new agent version",
func(telemetrySource string) {
if AgentTelemetryResourceId == "" {
Skip("Agent telemetry checks skipped because AGENT_TELEMETRY_RESOURCE_ID is not set")
}

Expect(AgentTelemetryVersion).NotTo(BeEmpty(), "AGENT_TELEMETRY_VERSION must be set to the newly deployed image tag")

expectedNodes, err := utils.GetExpectedAmaLogsNodes(K8sClient)
Expect(err).NotTo(HaveOccurred())

query := fmt.Sprintf(`%s
| where timestamp > ago(30m)
| extend ClusterId = iff(isnotempty(tostring(customDimensions.ID)), tostring(customDimensions.ID), tostring(customDimensions.AKS_RESOURCE_ID))
| where ClusterId =~ %q
| where tostring(customDimensions.Version) in (%q, %q)
| distinct Computer = tolower(tostring(customDimensions.Computer))`, telemetrySource, AKSResourceId, AgentTelemetryVersion, "win-"+AgentTelemetryVersion)

err = utils.CompareResourcesHelper(LogsClient, AgentTelemetryResourceId, query, expectedNodes)
Expect(err).NotTo(HaveOccurred())
},
Entry("customMetrics", "customMetrics"),
Entry("traces", "traces"),
Entry("heartbeat", `customEvents | where name == "ContainerLogDaemonSetHeartbeatEvent"`),
)
})
2 changes: 2 additions & 0 deletions test/ginkgo-e2e/utils/constants.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,8 @@ var (
"GetAgentConfigurations",
"RefreshConfigurations",
"canceled by user",
"(deleted)",
"errno=2] No such file or directory",
}
)

Expand Down
2 changes: 1 addition & 1 deletion test/ginkgo-e2e/utils/kubernetes_api_utils.go
Original file line number Diff line number Diff line change
Expand Up @@ -141,7 +141,7 @@ func GetContainerEnvVars(clientset *kubernetes.Clientset, namespace string, labe
}
}

return nil, fmt.Errorf("container %s not found in pod %s", containerName, &pods[0].Name)
return nil, fmt.Errorf("container %s not found in pod %s", containerName, pods[0].Name)
}

func GetAKSResourceID(clientset *kubernetes.Clientset, namespace string, labelKey string, labelValue string, containerName string) (string, error) {
Expand Down
19 changes: 12 additions & 7 deletions test/ginkgo-e2e/utils/query_logs_api_utils.go
Original file line number Diff line number Diff line change
Expand Up @@ -151,7 +151,6 @@ func CompareResourcesInLogsAndKubeAPI(K8sClient *kubernetes.Clientset, logsClien
return CompareResourcesHelper(logsClient, resourceID, query, resources)
}


func GetComputerFromContainerLog(logsClient *azquery.LogsClient, resourceID string, window string) (map[string]int64, error) {
counts, v2Err := queryCountsByComputer(logsClient, resourceID, "ContainerLogV2", window)
if v2Err == nil {
Expand Down Expand Up @@ -189,12 +188,11 @@ func queryCountsByComputer(logsClient *azquery.LogsClient, resourceID string, ta
return counts, nil
}

// AssertContainerLogNodeCoverage returns nil if every expected node appears
// in the per-Computer count map with a positive row count (compared
// case-insensitively), or an error listing the missing nodes otherwise.
func AssertContainerLogNodeCoverage(expectedNodes []string, observedCountsByComputer map[string]int64) error {
// AssertNodeCoverage returns nil if every expected node appears in the per-Computer count map
// with a positive count (compared case-insensitively), or an error listing the missing nodes.
func AssertNodeCoverage(signal string, expectedNodes []string, observedCountsByComputer map[string]int64) error {
if len(expectedNodes) == 0 {
return fmt.Errorf("no expected nodes provided; cannot verify ContainerLogV2 coverage")
return fmt.Errorf("no expected nodes provided; cannot verify %s coverage", signal)
}

var missing []string
Expand All @@ -204,7 +202,14 @@ func AssertContainerLogNodeCoverage(expectedNodes []string, observedCountsByComp
}
}
if len(missing) > 0 {
return fmt.Errorf("ContainerLogV2 ingestion is missing for %d/%d expected node(s): %s", len(missing), len(expectedNodes), strings.Join(missing, ", "))
return fmt.Errorf("%s is missing for %d/%d expected node(s): %s", signal, len(missing), len(expectedNodes), strings.Join(missing, ", "))
}
return nil
}

// AssertContainerLogNodeCoverage returns nil if every expected node appears
// in the per-Computer count map with a positive row count (compared
// case-insensitively), or an error listing the missing nodes otherwise.
func AssertContainerLogNodeCoverage(expectedNodes []string, observedCountsByComputer map[string]int64) error {
return AssertNodeCoverage("ContainerLogV2", expectedNodes, observedCountsByComputer)
}
51 changes: 37 additions & 14 deletions test/testkube/install-and-execute-testkube-tests.sh
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,8 @@ do
LinuxTestsOnly) LinuxTestsOnly=$VALUE ;;
GenevaIntegration) GenevaIntegration=$VALUE ;;
PerNodeLogCoverage) PerNodeLogCoverage=$VALUE ;;
AgentTelemetryResourceId) AgentTelemetryResourceId=$VALUE ;;
AgentTelemetryVersion) AgentTelemetryVersion=$VALUE ;;
*)
esac
done
Expand Down Expand Up @@ -68,6 +70,8 @@ export AZURE_TENANT_ID=$AzureTenantId
export WEBHOOK_URI=$TeamsWebhookUri
export GENEVA_INTEGRATION=$GenevaIntegration
export PER_NODE_LOG_COVERAGE=$PerNodeLogCoverage
export AGENT_TELEMETRY_RESOURCE_ID=$AgentTelemetryResourceId
export AGENT_TELEMETRY_VERSION=$AgentTelemetryVersion
kubectl apply -f ./api-server-permissions.yaml
kubectl apply -f ./testkube-test-crs.yaml

Expand All @@ -91,6 +95,8 @@ for wf in "${workflows[@]}"; do
kubectl testkube run testworkflow "$wf" \
--config GENEVA_INTEGRATION="$GENEVA_INTEGRATION" \
--config PER_NODE_LOG_COVERAGE="$PER_NODE_LOG_COVERAGE" \
--config AGENT_TELEMETRY_RESOURCE_ID="$AGENT_TELEMETRY_RESOURCE_ID" \
--config AGENT_TELEMETRY_VERSION="$AGENT_TELEMETRY_VERSION" \
--config AZURE_TENANT_ID="$AZURE_TENANT_ID" \
--config AZURE_CLIENT_ID="$AZURE_CLIENT_ID" \
--config GOTOOLCHAIN="auto" \
Expand All @@ -111,22 +117,39 @@ for wf in "${workflows[@]}"; do
exit 1
fi

# Watch until the testworkflow finishes
# Watch until the testworkflow finishes. The exit code is the authoritative result:
# the CLI returns non-zero when the execution fails.
kubectl testkube watch testworkflowexecution $execution_id
watch_rc=$?

# Get the results as a formatted json file.
# The execution status is not necessarily final the moment `watch` returns, so poll briefly
# for a terminal one instead of reading a status that is still "running" and mistaking it
# for a result. An empty file satisfies `jq empty`, so the document is also confirmed to be
# an object before any field is read out of it. The poll is kept short because the status is
# only ever corroboration: it can add a failure, never clear one.
wf_status=""
for attempt in $(seq 1 10); do
kubectl testkube get testworkflowexecution $execution_id --output json > "testkube-results-${wf}.json"
if [[ -s "testkube-results-${wf}.json" ]] && jq -e 'type == "object"' "testkube-results-${wf}.json" >/dev/null 2>&1; then
wf_status=$(jq -r '.result.status // empty' "testkube-results-${wf}.json")
fi
case "$wf_status" in
passed|failed|aborted|canceled) break ;;
esac
sleep 1
done
echo "TestWorkflow $wf finished with exit code $watch_rc and status '${wf_status:-unknown}'"

# The status only decides the outcome once it is terminal. When it never became terminal,
# or no usable JSON was returned at all, the exit code of `watch` is the only signal left,
# and it is what stops a failing workflow from being reported as a successful one.
status_failed=0
case "$wf_status" in
failed|aborted|canceled) status_failed=1 ;;
esac

# Get the results as a formatted json file
kubectl testkube get testworkflowexecution $execution_id --output json > "testkube-results-${wf}.json"

# Verify the JSON is valid
if ! jq empty "testkube-results-${wf}.json" 2>/dev/null; then
echo "Error: Failed to get valid JSON results from testkube for $wf"
echo "Contents of testkube-results-${wf}.json:"
cat "testkube-results-${wf}.json"
exit 1
fi

# For any test that has failed, print out the logs
if [[ $(jq -r '.result.status' "testkube-results-${wf}.json") == "failed" ]]; then
if [[ $watch_rc -ne 0 || $status_failed -eq 1 ]]; then

echo "TestWorkflow failed. Execution ID: $execution_id"

Expand Down
10 changes: 10 additions & 0 deletions test/testkube/testkube-test-crs.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -133,6 +133,12 @@ spec:
PER_NODE_LOG_COVERAGE:
type: string
default: "false"
AGENT_TELEMETRY_RESOURCE_ID:
type: string
default: ""
AGENT_TELEMETRY_VERSION:
type: string
default: ""
GOTOOLCHAIN:
type: string
default: ""
Expand All @@ -156,6 +162,10 @@ spec:
value: "{{config.GENEVA_INTEGRATION}}"
- name: PER_NODE_LOG_COVERAGE
value: "{{config.PER_NODE_LOG_COVERAGE}}"
- name: AGENT_TELEMETRY_RESOURCE_ID
value: "{{config.AGENT_TELEMETRY_RESOURCE_ID}}"
- name: AGENT_TELEMETRY_VERSION
value: "{{config.AGENT_TELEMETRY_VERSION}}"
- name: GOTOOLCHAIN
value: "{{config.GOTOOLCHAIN}}"
shell: ginkgo ./querylogs
Expand Down
Loading