diff --git a/api/v1alpha1/crds.go b/api/v1alpha1/crds.go index f5a488d0..e1d43c22 100644 --- a/api/v1alpha1/crds.go +++ b/api/v1alpha1/crds.go @@ -33,6 +33,7 @@ var ( // +kubebuilder:rbac:groups=trusted-execution-clusters.io,resources=trustedexecutionclusters;machines;approvedimages;attestationkeys,verbs=create;delete;get;list;patch;update;watch // +kubebuilder:rbac:groups=trusted-execution-clusters.io,resources=trustedexecutionclusters/finalizers;machines/finalizers;attestationkeys/finalizers;approvedimages/finalizers,verbs=update // +kubebuilder:rbac:groups=trusted-execution-clusters.io,resources=trustedexecutionclusters/status;machines/status;approvedimages/status;attestationkeys/status,verbs=get;patch;update +// +kubebuilder:rbac:groups=events.k8s.io,resources=events,verbs=create;patch // TrustedExecutionClusterSpec defines the desired state of TrustedExecutionCluster // +kubebuilder:validation:XValidation:rule="!has(oldSelf.publicAttestationKeyRegisterAddr) || has(self.publicAttestationKeyRegisterAddr)", message="Value is required once set" diff --git a/attestation-key-register/src/main.rs b/attestation-key-register/src/main.rs index 3effee0a..f0c5c4bd 100644 --- a/attestation-key-register/src/main.rs +++ b/attestation-key-register/src/main.rs @@ -8,8 +8,10 @@ use axum::{http::StatusCode, routing::put, Router}; use axum_server::tls_openssl::OpenSSLConfig; use clap::Parser; use env_logger::Env; +use k8s_openapi::api::core::v1::ObjectReference; use k8s_openapi::apimachinery::pkg::apis::meta::v1::ObjectMeta; -use kube::{Api, Client}; +use kube::runtime::events::{EventType, Recorder, Reporter}; +use kube::{Api, Client, Resource}; use log::{error, info}; use serde::{Deserialize, Serialize}; use std::net::SocketAddr; @@ -17,7 +19,8 @@ use uuid::Uuid; use trusted_cluster_operator_lib::endpoints::ATTESTATION_KEY_REGISTER_RESOURCE; use trusted_cluster_operator_lib::{ - generate_owner_reference, get_trusted_execution_cluster, AttestationKey, AttestationKeySpec, + generate_owner_reference, get_trusted_execution_cluster, record_event, AttestationKey, + AttestationKeySpec, }; #[derive(Parser)] @@ -32,6 +35,12 @@ struct Args { key_path: Option, } +#[derive(Clone)] +struct AppState { + client: Client, + recorder: Recorder, +} + #[derive(Debug, Deserialize, Serialize)] struct AttestationKeyRegistration { /// Public attestation key @@ -44,10 +53,12 @@ struct AttestationKeyRegistration { } async fn handle_registration( - State(client): State, + State(state): State, Json(registration): Json, ) -> impl IntoResponse { info!("Received registration request: {registration:?}"); + let client = state.client; + let recorder = state.recorder; let internal_error = |e: anyhow::Error| { let code = StatusCode::INTERNAL_SERVER_ERROR; @@ -76,10 +87,23 @@ async fn handle_registration( Ok(existing_keys) => { for key in existing_keys.items { if key.spec.public_key == registration.public_key { + let key_ref: ObjectReference = key.object_ref(&()); let existing_name = key.metadata.name.unwrap_or_default(); error!( "Duplicate public key detected: already exists in AttestationKey '{existing_name}'" ); + record_event( + &recorder, + &key_ref, + EventType::Warning, + "DuplicateKeyRejected", + format!( + "Duplicate registration attempt for AttestationKey '{existing_name}'" + ), + "Registering", + None, + ) + .await; return ( StatusCode::CONFLICT, Json(serde_json::json!({ @@ -93,7 +117,7 @@ async fn handle_registration( Err(e) => { return internal_error( anyhow::Error::from(e).context("Failed to check for existing keys"), - ) + ); } } @@ -113,8 +137,19 @@ async fn handle_registration( match api.create(&Default::default(), &attestation_key).await { Ok(created) => { + let created_ref: ObjectReference = created.object_ref(&()); let name = created.metadata.name.unwrap_or_default(); info!("Successfully created AttestationKey: {name}",); + record_event( + &recorder, + &created_ref, + EventType::Normal, + "AttestationKeyRegistered", + format!("AttestationKey '{name}' registered"), + "Registering", + None, + ) + .await; let json = Json(serde_json::json!({ "status": "success", })); @@ -132,9 +167,17 @@ async fn main() { let endpoint = format!("/{ATTESTATION_KEY_REGISTER_RESOURCE}"); let err = "failed to create Kubernetes client"; let client = Client::try_default().await.expect(err); + let reporter = Reporter { + controller: "attestation-key-register".into(), + instance: std::env::var("CONTROLLER_POD_NAME").ok(), + }; + let state = AppState { + recorder: Recorder::new(client.clone(), reporter), + client, + }; let app = Router::new() .route(&endpoint, put(handle_registration)) - .with_state(client); + .with_state(state); let addr = SocketAddr::from(([0, 0, 0, 0], args.port)); let service = app.into_make_service(); diff --git a/lib/Cargo.toml b/lib/Cargo.toml index a7620cbc..36a05333 100644 --- a/lib/Cargo.toml +++ b/lib/Cargo.toml @@ -15,6 +15,7 @@ anyhow.workspace = true compute-pcrs-lib.workspace = true k8s-openapi.workspace = true kube.workspace = true +log.workspace = true serde.workspace = true serde_json.workspace = true diff --git a/lib/src/lib.rs b/lib/src/lib.rs index 03adb539..f4c469fb 100644 --- a/lib/src/lib.rs +++ b/lib/src/lib.rs @@ -25,7 +25,9 @@ pub use vendor_kopium::virtualmachines; use anyhow::{Context, Result, anyhow}; use conditions::*; +use k8s_openapi::api::core::v1::ObjectReference; use k8s_openapi::apimachinery::pkg::apis::meta::v1::{Condition, OwnerReference, Time}; +use kube::runtime::events::{Event as K8sEvent, EventType, Recorder}; use kube::{Api, Client, Resource}; #[macro_export] @@ -106,6 +108,27 @@ pub fn committed_condition( } } +pub async fn record_event( + recorder: &Recorder, + reference: &ObjectReference, + type_: EventType, + reason: &str, + note: String, + action: &str, + secondary: Option, +) { + let ev = K8sEvent { + type_, + reason: reason.into(), + note: Some(note), + action: action.into(), + secondary, + }; + if let Err(e) = recorder.publish(&ev, reference).await { + log::warn!("Failed to publish event: {e}"); + } +} + /// Generate an OwnerReference for any Kubernetes resource pub fn generate_owner_reference>( object: &T, diff --git a/operator/src/attestation_key_register.rs b/operator/src/attestation_key_register.rs index a8c9efc3..71fba6b2 100644 --- a/operator/src/attestation_key_register.rs +++ b/operator/src/attestation_key_register.rs @@ -6,6 +6,7 @@ use anyhow::{Result, anyhow}; use futures_util::StreamExt; use k8s_openapi::ByteString; use k8s_openapi::api::apps::v1::{Deployment, DeploymentSpec}; +use k8s_openapi::api::core::v1::ObjectReference; use k8s_openapi::api::core::v1::{ Container, ContainerPort, PodSpec, PodTemplateSpec, Secret, Service, ServicePort, ServiceSpec, }; @@ -19,6 +20,7 @@ use kube::{ runtime::{ Controller, controller::Action, + events::{EventType, Recorder}, finalizer, finalizer::Event, reflector::{self, ObjectRef, Store}, @@ -37,11 +39,13 @@ use crate::conditions::attestation_key_approved_condition; use crate::trustee; use operator::{ControllerError, LONG_REQUEUE, TLS_DIR, controller_error_policy}; use operator::{create_or_info_if_exists, read_certificate, upsert_condition}; +use trusted_cluster_operator_lib::record_event; /// Shared context for the three attestation-key controllers. /// Stores give local cache access to avoid repeated API-server reads. pub struct AkContextData { pub client: Client, + pub recorder: Recorder, pub machine_store: Store, pub ak_store: Store, pub secret_store: Store, @@ -60,8 +64,10 @@ impl AkContextData { crate::spawn_reflector::(secret_writer, client.clone(), "Secret"); crate::spawn_reflector::(deployment_writer, client.clone(), "Deployment"); + let recorder = operator::new_recorder(client.clone(), "ak-controller"); Self { client, + recorder, machine_store, ak_store, secret_store, @@ -238,14 +244,36 @@ async fn approve_ak(ak: &AttestationKey, machine: &Machine, ctx: &AkContextData) let condition = attestation_key_approved_condition(approve_reason, generation, &ak.status); let mut conditions = ak.status.as_ref().and_then(|s| s.conditions.clone()); let changed = upsert_condition(&mut conditions, condition); + let machine_name = machine.metadata.name.clone().unwrap_or_default(); if changed { let status = AttestationKeyStatus { conditions }; update_status!(aks, &name, status)?; info!("Approved attestation key {name}"); - } - let machine_name = machine.metadata.name.clone().unwrap_or_default(); + let ak_ref: ObjectReference = ak.object_ref(&()); + let machine_ref: ObjectReference = machine.object_ref(&()); + record_event( + &ctx.recorder, + &ak_ref, + EventType::Normal, + "AttestationKeyApproved", + format!("Attestation key {name} approved for machine {machine_name}"), + "Approving", + Some(machine_ref.clone()), + ) + .await; + record_event( + &ctx.recorder, + &machine_ref, + EventType::Normal, + "AttestationKeyApproved", + format!("Machine {machine_name} matched attestation key {name}"), + "Approving", + Some(ak_ref), + ) + .await; + } let has_machine_owner = ak .metadata .owner_references diff --git a/operator/src/lib.rs b/operator/src/lib.rs index d77bc31e..bb3e5119 100644 --- a/operator/src/lib.rs +++ b/operator/src/lib.rs @@ -14,6 +14,7 @@ use k8s_openapi::api::core::v1::{Secret, SecretVolumeSource, Volume, VolumeMount use k8s_openapi::apimachinery::pkg::apis::meta::v1::{Condition, Time}; use k8s_openapi::jiff::Timestamp; use kube::Resource; +use kube::runtime::events::{Recorder, Reporter}; use kube::runtime::reflector::{self, Store}; use kube::runtime::watcher::watcher; use kube::{Api, Client, runtime::controller::Action}; @@ -43,6 +44,19 @@ pub async fn controller_info(res: Result) { } } +pub fn new_recorder(client: Client, controller_name: &str) -> Recorder { + let reporter = Reporter { + controller: controller_name.into(), + instance: std::env::var("CONTROLLER_POD_NAME").ok(), + }; + Recorder::new(client, reporter) +} + +pub struct ControllerContext { + pub client: Client, + pub recorder: Recorder, +} + #[macro_export] macro_rules! create_or_info_if_exists { ($client:expr, $type:ident, $resource:ident) => { diff --git a/operator/src/reference_values.rs b/operator/src/reference_values.rs index 3350de0a..154c4892 100644 --- a/operator/src/reference_values.rs +++ b/operator/src/reference_values.rs @@ -6,6 +6,7 @@ use anyhow::{Context, Result, anyhow}; use compute_pcrs_lib::Pcr; use futures_util::StreamExt; +use k8s_openapi::api::core::v1::ObjectReference; use k8s_openapi::{ api::{ batch::v1::{Job, JobSpec}, @@ -17,6 +18,7 @@ use k8s_openapi::{ use kube::api::{DeleteParams, ListParams, ObjectMeta, Patch}; use kube::runtime::{ controller::{Action, Controller}, + events::EventType, finalizer, finalizer::Event, watcher, @@ -32,9 +34,9 @@ use std::{collections::BTreeMap, sync::Arc, time::Duration}; use crate::COMPONENT_VERSION; use crate::trustee::{self, get_image_pcrs}; -use operator::{ControllerError, LONG_REQUEUE, upsert_condition}; -use operator::{controller_error_policy, controller_info, create_or_info_if_exists}; -use trusted_cluster_operator_lib::{conditions::*, reference_values::*, *}; +use operator::{ControllerContext, ControllerError, LONG_REQUEUE, upsert_condition}; +use operator::{controller_error_policy, controller_info, create_or_info_if_exists, new_recorder}; +use trusted_cluster_operator_lib::{conditions::*, record_event, reference_values::*, *}; const JOB_LABEL_KEY: &str = "kind"; const APPROVED_IMAGE_ANNOTATION: &str = "approved-image"; @@ -115,12 +117,15 @@ fn build_compute_pcrs_pod_spec( } } -async fn job_reconcile(job: Arc, client: Arc) -> Result { +async fn job_reconcile( + job: Arc, + ctx: Arc, +) -> Result { let err = "Job changed, but had no name"; let name = &job.metadata.name.clone().context(err)?; let err = format!("Job {name} changed, but had no status"); let status = &job.status.clone().context(err)?; - let kube_client = Arc::unwrap_or_clone(client); + let kube_client = ctx.client.clone(); if status.completion_time.is_none() { info!("Job {name} changed, but had not completed"); return Ok(Action::requeue(Duration::from_secs(300))); @@ -130,6 +135,31 @@ async fn job_reconcile(job: Arc, client: Arc) -> Result::into)?; trustee::update_reference_values(kube_client).await?; + if let Some(owner) = job + .metadata + .owner_references + .as_ref() + .and_then(|refs| refs.iter().find(|r| r.kind == "ApprovedImage")) + { + let image_ref = ObjectReference { + api_version: Some(owner.api_version.clone()), + kind: Some(owner.kind.clone()), + name: Some(owner.name.clone()), + namespace: job.metadata.namespace.clone(), + uid: Some(owner.uid.clone()), + ..Default::default() + }; + record_event( + &ctx.recorder, + &image_ref, + EventType::Normal, + "ComputationCompleted", + format!("Reference values computed for {}", owner.name), + "Computing", + None, + ) + .await; + } Ok(Action::await_change()) } @@ -139,9 +169,13 @@ pub async fn launch_rv_job_controller(client: Client) { label_selector: Some(format!("{JOB_LABEL_KEY}={PCR_COMMAND_NAME}")), ..Default::default() }; + let ctx = Arc::new(ControllerContext { + recorder: new_recorder(client.clone(), "rv-controller"), + client, + }); tokio::spawn( Controller::new(jobs, watcher) - .run(job_reconcile, controller_error_policy, Arc::new(client)) + .run(job_reconcile, controller_error_policy, ctx) .for_each(controller_info), ); } @@ -238,9 +272,9 @@ pub async fn adopt_approved_images( async fn image_reconcile( image: Arc, - client: Arc, + ctx: Arc, ) -> Result { - let kube_client = Arc::::unwrap_or_clone(client); + let kube_client = ctx.client.clone(); let err = "ApprovedImage had no name"; let name = image.metadata.name.clone().context(err)?; let cluster = get_opt_trusted_execution_cluster(kube_client.clone()) @@ -267,7 +301,7 @@ async fn image_reconcile( let images: Api = Api::default_namespaced(kube_client.clone()); finalizer(&images, APPROVED_IMAGE_FINALIZER, image, |ev| async { match ev { - Event::Apply(image) => image_add_reconcile(kube_client, &image, cluster) + Event::Apply(image) => image_add_reconcile(&ctx, &image, cluster) .await .map_err(|e| finalizer::Error::::ApplyFailed(e.into())), Event::Cleanup(image) => image_remove_reconcile(kube_client, image, cluster) @@ -280,7 +314,7 @@ async fn image_reconcile( } async fn image_add_reconcile( - client: Client, + ctx: &ControllerContext, image: &ApprovedImage, cluster: Option, ) -> Result { @@ -294,10 +328,35 @@ async fn image_add_reconcile( info!("TrustedExecutionCluster is being deleted, deferring image processing for {name}"); return Ok(Action::requeue(Duration::from_secs(5))); } - let (action, reason) = match handle_new_image(client.clone(), image).await { - Ok(reason) => (LONG_REQUEUE, reason), + let image_ref: ObjectReference = image.object_ref(&()); + let (action, reason) = match handle_new_image(ctx.client.clone(), image).await { + Ok(reason) => { + if reason == NOT_COMMITTED_REASON_COMPUTING { + record_event( + &ctx.recorder, + &image_ref, + EventType::Normal, + "ComputationStarted", + format!("PCR computation started for {name}"), + "Computing", + None, + ) + .await; + } + (LONG_REQUEUE, reason) + } Err(e) => { warn!("PCR computation for {name} failed: {e}"); + record_event( + &ctx.recorder, + &image_ref, + EventType::Warning, + "ComputationFailed", + format!("PCR computation for {name} failed: {e}"), + "Computing", + None, + ) + .await; let action = Action::requeue(Duration::from_secs(60)); (action, NOT_COMMITTED_REASON_FAILED) } @@ -308,7 +367,7 @@ async fn image_add_reconcile( let mut conditions = image.status.as_ref().and_then(|s| s.conditions.clone()); let changed = upsert_condition(&mut conditions, committed); if changed { - let images: Api = Api::default_namespaced(client); + let images: Api = Api::default_namespaced(ctx.client.clone()); update_status!(images, &name, ApprovedImageStatus { conditions }) .map_err(|e| finalizer::Error::::ApplyFailed(e.into()))?; } @@ -341,9 +400,13 @@ async fn image_remove_reconcile( pub async fn launch_rv_image_controller(client: Client) { let images: Api = Api::default_namespaced(client.clone()); + let ctx = Arc::new(ControllerContext { + recorder: new_recorder(client.clone(), "rv-controller"), + client, + }); tokio::spawn( Controller::new(images, Default::default()) - .run(image_reconcile, controller_error_policy, Arc::new(client)) + .run(image_reconcile, controller_error_policy, ctx) .for_each(controller_info), ); } @@ -503,7 +566,11 @@ mod tests { }; count_check!(5, clos, |client| { let job = Arc::new(dummy_job()); - let result = job_reconcile(job, Arc::new(client)).await.unwrap(); + let ctx = Arc::new(ControllerContext { + recorder: new_recorder(client.clone(), "test"), + client, + }); + let result = job_reconcile(job, ctx).await.unwrap(); assert_eq!(result, Action::await_change()); }); } @@ -515,7 +582,11 @@ mod tests { let mut job = dummy_job(); let status = job.status.as_mut().unwrap(); status.completion_time = None; - let result = job_reconcile(Arc::new(job), Arc::new(client)).await; + let ctx = Arc::new(ControllerContext { + recorder: new_recorder(client.clone(), "test"), + client, + }); + let result = job_reconcile(Arc::new(job), ctx).await; assert_eq!(result.unwrap(), Action::requeue(Duration::from_secs(300))); }); } diff --git a/operator/src/register_server.rs b/operator/src/register_server.rs index 3ee1351c..f8e9beb5 100644 --- a/operator/src/register_server.rs +++ b/operator/src/register_server.rs @@ -6,6 +6,7 @@ use anyhow::{Result, anyhow}; use futures_util::StreamExt; use k8s_openapi::api::apps::v1::{Deployment, DeploymentSpec}; +use k8s_openapi::api::core::v1::ObjectReference; use k8s_openapi::api::core::v1::{ Container, ContainerPort, PodSpec, PodTemplateSpec, Service, ServicePort, ServiceSpec, }; @@ -15,6 +16,7 @@ use k8s_openapi::apimachinery::pkg::{ }; use kube::runtime::{ controller::{Action, Controller}, + events::EventType, finalizer, finalizer::Event, }; @@ -24,7 +26,7 @@ use std::{collections::BTreeMap, sync::Arc}; use crate::trustee; use operator::*; -use trusted_cluster_operator_lib::{Machine, TrustedExecutionCluster, endpoints::*}; +use trusted_cluster_operator_lib::{Machine, TrustedExecutionCluster, endpoints::*, record_event}; /// Finalizer name to discard decryption keys when a machine is deleted const MACHINE_FINALIZER: &str = "finalizer.machine.trusted-execution-clusters.io"; @@ -127,26 +129,52 @@ pub async fn create_register_server_service( async fn keygen_reconcile( machine: Arc, - client: Arc, + ctx: Arc, ) -> Result { - let kube_client_clone = Arc::unwrap_or_clone(client.clone()); - let machines: Api = Api::default_namespaced(kube_client_clone.clone()); + let machines: Api = Api::default_namespaced(ctx.client.clone()); + let recorder = ctx.recorder.clone(); finalizer(&machines, MACHINE_FINALIZER, machine, |ev| async move { match ev { Event::Apply(machine) => { - let kube_client = Arc::unwrap_or_clone(client); + let machine_ref: ObjectReference = machine.object_ref(&()); + let kube_client = ctx.client.clone(); let id = &machine.spec.id.clone(); - async { + let result = async { let owner_reference = generate_owner_reference(&Arc::unwrap_or_clone(machine))?; trustee::generate_secret(kube_client.clone(), id, owner_reference).await?; trustee::send_secret(kube_client, id).await } - .await - .map(|_| LONG_REQUEUE) - .map_err(|e| finalizer::Error::::ApplyFailed(e.into())) + .await; + if result.is_ok() { + record_event( + &recorder, + &machine_ref, + EventType::Normal, + "KeyProvisioned", + format!("Decryption key provisioned for machine {id}"), + "Provisioning", + None, + ) + .await; + } else { + record_event( + &recorder, + &machine_ref, + EventType::Warning, + "KeyProvisioningFailed", + format!("Failed to provision decryption key for machine {id}"), + "Provisioning", + None, + ) + .await; + } + result + .map(|_| LONG_REQUEUE) + .map_err(|e| finalizer::Error::::ApplyFailed(e.into())) } Event::Cleanup(machine) => { - let kube_client = Arc::unwrap_or_clone(client); + let machine_ref: ObjectReference = machine.object_ref(&()); + let kube_client = ctx.client.clone(); let id = &machine.spec.id; // Check if the TrustedExecutionCluster is being deleted @@ -184,8 +212,20 @@ async fn keygen_reconcile( } } } - trustee::delete_secret(kube_client, id) - .await + let result = trustee::delete_secret(kube_client, id).await; + if result.is_ok() { + record_event( + &recorder, + &machine_ref, + EventType::Normal, + "KeyRevoked", + format!("Decryption key revoked for machine {id}"), + "Revoking", + None, + ) + .await; + } + result .map(|_| LONG_REQUEUE) .map_err(|e| finalizer::Error::::CleanupFailed(e.into())) } @@ -197,9 +237,13 @@ async fn keygen_reconcile( pub async fn launch_keygen_controller(client: Client) { let machines: Api = Api::default_namespaced(client.clone()); + let ctx = Arc::new(ControllerContext { + recorder: new_recorder(client.clone(), "keygen-controller"), + client, + }); tokio::spawn( Controller::new(machines, Default::default()) - .run(keygen_reconcile, controller_error_policy, Arc::new(client)) + .run(keygen_reconcile, controller_error_policy, ctx) .for_each(controller_info), ); } diff --git a/register-server/src/main.rs b/register-server/src/main.rs index 6db836b2..5826783c 100644 --- a/register-server/src/main.rs +++ b/register-server/src/main.rs @@ -16,16 +16,18 @@ use env_logger::Env; use ignition_config::v3_6::{ Clevis, ClevisCustom, Config as IgnitionConfig, Filesystem, Luks, Storage, }; +use k8s_openapi::api::core::v1::ObjectReference; use k8s_openapi::api::core::v1::Secret; use k8s_openapi::apimachinery::pkg::apis::meta::v1::{ObjectMeta, OwnerReference}; -use kube::{Api, Client}; +use kube::runtime::events::{EventType, Recorder, Reporter}; +use kube::{Api, Client, Resource}; use log::{error, info}; use std::net::SocketAddr; use uuid::Uuid; use trusted_cluster_operator_lib::endpoints::*; use trusted_cluster_operator_lib::{ - generate_owner_reference, get_trusted_execution_cluster, Machine, MachineSpec, + generate_owner_reference, get_trusted_execution_cluster, record_event, Machine, MachineSpec, }; #[derive(Parser)] @@ -42,6 +44,12 @@ struct Args { key_path: Option, } +#[derive(Clone)] +struct AppState { + client: Client, + recorder: Recorder, +} + /// Information about endpoints for clevis configuration struct EndpointInfo { /// The public address of the Trustee server @@ -163,8 +171,10 @@ fn generate_ignition(id: &str, endpoint_info: &EndpointInfo) -> IgnitionConfig { } } -async fn register_handler(State(kube_client): State) -> impl IntoResponse { +async fn register_handler(State(state): State) -> impl IntoResponse { let id = Uuid::new_v4().to_string(); + let kube_client = state.client; + let recorder = state.recorder; let internal_error = |e: anyhow::Error| { let code = StatusCode::INTERNAL_SERVER_ERROR; error!("{e:?}"); @@ -187,7 +197,20 @@ async fn register_handler(State(kube_client): State) -> impl IntoRespons }; match create_machine(kube_client.clone(), &id, owner_reference).await { - Ok(_) => info!("Machine created successfully: machine-{id}"), + Ok(machine) => { + let machine_ref: ObjectReference = machine.object_ref(&()); + info!("Machine created successfully: machine-{id}"); + record_event( + &recorder, + &machine_ref, + EventType::Normal, + "MachineRegistered", + format!("Machine machine-{id} registered"), + "Registering", + None, + ) + .await; + } Err(e) => return internal_error(e.context("Failed to create machine")), } let endpoint_info = match EndpointInfo::create(kube_client).await { @@ -208,7 +231,7 @@ async fn create_machine( client: Client, uuid: &str, owner_reference: OwnerReference, -) -> anyhow::Result<()> { +) -> anyhow::Result { let machine_name = format!("machine-{uuid}"); let machine = Machine { metadata: ObjectMeta { @@ -223,9 +246,9 @@ async fn create_machine( }; let machines: Api = Api::default_namespaced(client); - machines.create(&Default::default(), &machine).await?; + let created = machines.create(&Default::default(), &machine).await?; info!("Created Machine: {machine_name} with UUID: {uuid}"); - Ok(()) + Ok(created) } #[tokio::main] @@ -235,9 +258,18 @@ async fn main() { let args = Args::parse(); let endpoint = format!("/{REGISTER_SERVER_RESOURCE}"); let err = "failed to create Kubernetes client"; + let client = Client::try_default().await.expect(err); + let reporter = Reporter { + controller: "register-server".into(), + instance: std::env::var("CONTROLLER_POD_NAME").ok(), + }; + let state = AppState { + recorder: Recorder::new(client.clone(), reporter), + client, + }; let app = Router::new() .route(&endpoint, get(register_handler)) - .with_state(Client::try_default().await.expect(err)); + .with_state(state); let addr = SocketAddr::from(([0, 0, 0, 0], args.port)); let service = app.into_make_service(); @@ -357,9 +389,10 @@ mod tests { async fn test_create_machine() { let clos = async |_, _| Ok(serde_json::to_string(&dummy_machine()).unwrap()); count_check!(1, clos, |client| { - assert!(create_machine(client, "test", dummy_owner_reference()) + let machine = create_machine(client, "test", dummy_owner_reference()) .await - .is_ok()); + .unwrap(); + assert_eq!(machine.metadata.name, Some("test".to_string())); }); } diff --git a/test_utils/src/lib.rs b/test_utils/src/lib.rs index 39de2d9e..118177cf 100644 --- a/test_utils/src/lib.rs +++ b/test_utils/src/lib.rs @@ -12,8 +12,9 @@ use k8s_openapi::api::core::v1::{ ConfigMap, LoadBalancerStatus, Namespace, Secret, Service, ServicePort, ServiceSpec, ServiceStatus, }; +use k8s_openapi::api::events::v1::Event; use k8s_openapi::apimachinery::pkg::apis::meta::v1::Condition; -use kube::api::{DeleteParams, ObjectMeta, Patch}; +use kube::api::{DeleteParams, ListParams, ObjectMeta, Patch}; use kube::runtime::wait::await_condition; use kube::{Api, Client}; use percent_encoding::{NON_ALPHANUMERIC, utf8_percent_encode}; @@ -1183,3 +1184,39 @@ where timeout(duration, done).await.context(ctx)??; Ok(()) } + +pub async fn wait_for_event( + client: &Client, + namespace: &str, + regarding_name: &str, + reason: &str, + timeout_secs: u64, +) -> Result { + let events_api: Api = Api::namespaced(client.clone(), namespace); + Poller::new() + .with_timeout(Duration::from_secs(timeout_secs)) + .with_error_message(format!( + "waiting for event reason='{reason}' regarding='{regarding_name}'" + )) + .poll_async(|| { + let api = events_api.clone(); + let name = regarding_name.to_string(); + let rsn = reason.to_string(); + async move { + let list = api.list(&ListParams::default()).await?; + list.items + .into_iter() + .find(|e| { + let name_match = e + .regarding + .as_ref() + .and_then(|r| r.name.as_ref()) + .is_some_and(|n| n == &name); + let reason_match = e.reason.as_deref() == Some(&rsn); + name_match && reason_match + }) + .ok_or_else(|| anyhow!("not found yet")) + } + }) + .await +} diff --git a/tests/attestation.rs b/tests/attestation.rs index 6b04b6c1..d5f64c47 100644 --- a/tests/attestation.rs +++ b/tests/attestation.rs @@ -10,8 +10,11 @@ if #[cfg(feature = "virtualization")] { use anyhow::Result; use k8s_openapi::api::apps::v1::Deployment; use kube::Api; -use trusted_cluster_operator_lib::Machine; +use std::time::Duration; +use trusted_cluster_operator_lib::{AttestationKey, Machine, TrustedExecutionCluster}; +use trusted_cluster_operator_test_utils::constants::APPROVED_IMAGE_NAME; use trusted_cluster_operator_test_utils::virt::{self, VmBackend}; +use trusted_cluster_operator_test_utils::{Poller, wait_for_event}; const ENCRYPTED_ROOT_ASSERT: &str = "should have an encrypted root device (attestation failed)"; const ENCRYPTED_ROOT_WARN: &str = "Backend reports that Machine IDs cannot be correlated to IP \ @@ -258,3 +261,89 @@ async fn test_vm_restart_operator_new() -> anyhow::Result<()> { Ok(()) } } + +virt_test! { +async fn test_attestation_events() -> anyhow::Result<()> { + let test_ctx = setup!().await?; + let client = test_ctx.client(); + let namespace = test_ctx.namespace(); + + let vm_name = "test-coreos-events"; + let att_ctx = SingleAttestationContext::new(vm_name, &test_ctx).await?; + + test_ctx.info("Verifying encrypted root device"); + let has_encrypted_root = att_ctx.verify_encrypted_root().await?; + assert!(has_encrypted_root, "VM should have an encrypted root device"); + test_ctx.info("Attestation successful, verifying Kubernetes events"); + + let tecs: Api = Api::namespaced(client.clone(), namespace); + let tec = tecs.get("trusted-execution-cluster").await?; + let has_ak_register = tec.spec.public_attestation_key_register_addr.is_some(); + + let machines: Api = Api::namespaced(client.clone(), namespace); + let machine_list = machines.list(&Default::default()).await?; + let machine_name = machine_list.items.first() + .expect("No Machine found in namespace") + .metadata + .name + .as_ref() + .expect("Machine should have a name"); + + wait_for_event(client, namespace, APPROVED_IMAGE_NAME, "ComputationStarted", scaled_timeout(30)).await?; + test_ctx.info("Event ComputationStarted verified"); + + // 3x the controller error policy requeue (60s) to survive a retry + wait_for_event(client, namespace, APPROVED_IMAGE_NAME, "ComputationCompleted", scaled_timeout(180)).await?; + test_ctx.info("Event ComputationCompleted verified"); + + wait_for_event(client, namespace, machine_name, "MachineRegistered", scaled_timeout(30)).await?; + test_ctx.info("Event MachineRegistered verified"); + + wait_for_event(client, namespace, machine_name, "KeyProvisioned", scaled_timeout(30)).await?; + test_ctx.info("Event KeyProvisioned verified"); + + if has_ak_register { + let aks: Api = Api::namespaced(client.clone(), namespace); + let ak = Poller::new() + .with_timeout(Duration::from_secs(scaled_timeout(60))) + .with_error_message("waiting for AttestationKey to be created") + .poll_async(|| { + let api = aks.clone(); + async move { + let list = api.list(&Default::default()).await?; + list.items + .into_iter() + .next() + .ok_or_else(|| anyhow::anyhow!("No AttestationKey found yet")) + } + }) + .await?; + let ak_name = ak.metadata.name.as_ref().expect("AttestationKey should have a name"); + + test_ctx.info(format!("Found AttestationKey: {ak_name}")); + + wait_for_event(client, namespace, ak_name, "AttestationKeyRegistered", scaled_timeout(30)).await?; + test_ctx.info("Event AttestationKeyRegistered verified"); + + wait_for_event(client, namespace, ak_name, "AttestationKeyApproved", scaled_timeout(30)).await?; + test_ctx.info("Event AttestationKeyApproved on AttestationKey verified"); + + wait_for_event(client, namespace, machine_name, "AttestationKeyApproved", scaled_timeout(30)).await?; + test_ctx.info("Event AttestationKeyApproved on Machine verified"); + } else { + test_ctx.info("Skipping AttestationKey events (publicAttestationKeyRegisterAddr not set)"); + } + + test_ctx.info(format!("Deleting Machine {machine_name} to trigger KeyRevoked")); + machines.delete(machine_name, &Default::default()).await?; + wait_for_resource_deleted(&machines, machine_name, scaled_timeout(120)).await?; + + wait_for_event(client, namespace, machine_name, "KeyRevoked", scaled_timeout(60)).await?; + test_ctx.info("Event KeyRevoked verified"); + + test_ctx.info("All expected Kubernetes events verified"); + att_ctx.cleanup().await?; + test_ctx.cleanup().await?; + Ok(()) +} +}