Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions api/v1alpha1/crds.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Comment thread
sourcery-ai[bot] marked this conversation as resolved.

// TrustedExecutionClusterSpec defines the desired state of TrustedExecutionCluster
// +kubebuilder:validation:XValidation:rule="!has(oldSelf.publicAttestationKeyRegisterAddr) || has(self.publicAttestationKeyRegisterAddr)", message="Value is required once set"
Expand Down
53 changes: 48 additions & 5 deletions attestation-key-register/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,16 +8,19 @@ 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;
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)]
Expand All @@ -32,6 +35,12 @@ struct Args {
key_path: Option<String>,
}

#[derive(Clone)]
struct AppState {
client: Client,
recorder: Recorder,
}

#[derive(Debug, Deserialize, Serialize)]
struct AttestationKeyRegistration {
/// Public attestation key
Expand All @@ -44,10 +53,12 @@ struct AttestationKeyRegistration {
}

async fn handle_registration(
State(client): State<Client>,
State(state): State<AppState>,
Json(registration): Json<AttestationKeyRegistration>,
) -> 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;
Expand Down Expand Up @@ -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!({
Expand All @@ -93,7 +117,7 @@ async fn handle_registration(
Err(e) => {
return internal_error(
anyhow::Error::from(e).context("Failed to check for existing keys"),
)
);
}
}

Expand All @@ -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",
}));
Expand All @@ -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(),

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sorry for question after merge. I know this example suggests the same https://docs.rs/kube/latest/kube/runtime/events/struct.Recorder.html but is this expected to ever be set?

};
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();

Expand Down
1 change: 1 addition & 0 deletions lib/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
23 changes: 23 additions & 0 deletions lib/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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]
Expand Down Expand Up @@ -112,6 +114,27 @@ pub fn committed_condition(
}
}

pub async fn record_event(
recorder: &Recorder,
reference: &ObjectReference,
type_: EventType,
reason: &str,
note: String,
action: &str,
secondary: Option<ObjectReference>,
) {
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<T: Resource<DynamicType = ()>>(
object: &T,
Expand Down
31 changes: 27 additions & 4 deletions operator/src/attestation_key_register.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
};
Expand All @@ -15,15 +16,15 @@ use k8s_openapi::apimachinery::pkg::{
};
use kube::api::{Patch, PatchParams};
use kube::runtime::{Controller, controller::Action, reflector::ObjectRef, watcher};
use kube::runtime::{finalizer, finalizer::Event};
use kube::runtime::{events::EventType, finalizer, finalizer::Event};
use kube::{Api, Client, Resource};
use log::{info, warn};
use serde_json::json;
use std::{collections::BTreeMap, sync::Arc};

use trusted_cluster_operator_lib::conditions::ATTESTATION_KEY_MACHINE_APPROVE;
use trusted_cluster_operator_lib::endpoints::*;
use trusted_cluster_operator_lib::{AttestationKey, Machine};
use trusted_cluster_operator_lib::{AttestationKey, Machine, record_event};

use crate::conditions::{attestation_key_approved_condition, machine_ak_approved_condition};
use crate::trustee;
Expand Down Expand Up @@ -201,6 +202,7 @@ async fn approve_ak(ak: &AttestationKey, machine: &Machine, ctx: &OperatorContex
let client = &ctx.client;
let aks: Api<AttestationKey> = Api::default_namespaced(client.clone());

let machine_name = machine.metadata.name.clone().unwrap_or_default();
let condition = attestation_key_approved_condition(
ATTESTATION_KEY_MACHINE_APPROVE,
ak.metadata.generation,
Expand All @@ -217,9 +219,30 @@ async fn approve_ak(ak: &AttestationKey, machine: &Machine, ctx: &OperatorContex
.await?
{
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
Expand Down
12 changes: 12 additions & 0 deletions operator/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ use k8s_openapi::api::core::v1::{ConfigMap, Secret, SecretVolumeSource, Volume,
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};
Expand All @@ -34,6 +35,7 @@ use trusted_cluster_operator_lib::{
/// Stores give local cache access to avoid repeated API-server reads.
pub struct OperatorContext {
pub client: Client,
pub recorder: Recorder,
pub tec_store: Store<TrustedExecutionCluster>,
pub cm_store: Store<ConfigMap>,
pub machine_store: Store<Machine>,
Expand All @@ -45,8 +47,10 @@ pub struct OperatorContext {

impl OperatorContext {
pub fn new(client: Client) -> Self {
let recorder = new_recorder(client.clone(), "operator");
Self {
client,
recorder,
tec_store: reflector::store().0,
cm_store: reflector::store().0,
machine_store: reflector::store().0,
Expand Down Expand Up @@ -88,6 +92,14 @@ pub async fn controller_info<T: Debug, E: Debug>(res: Result<T, E>) {
}
}

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)
}

#[macro_export]
macro_rules! create_or_info_if_exists {
($client:expr, $type:ident, $resource:ident) => {
Expand Down
66 changes: 61 additions & 5 deletions operator/src/reference_values.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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},
Expand All @@ -15,9 +16,14 @@ use k8s_openapi::{
jiff::Timestamp,
};
use kube::api::{DeleteParams, ListParams, ObjectMeta, Patch, PatchParams};
use kube::runtime::controller::{Action, Controller};
use kube::runtime::{finalizer, finalizer::Event};
use kube::runtime::{reflector::ObjectRef, watcher};
use kube::runtime::{
controller::{Action, Controller},
events::EventType,
finalizer,
finalizer::Event,
reflector::ObjectRef,
watcher,
};
use kube::{Api, Client, Resource};
use log::{info, warn};
use oci_client::secrets::RegistryAuth;
Expand All @@ -31,7 +37,7 @@ use crate::COMPONENT_VERSION;
use crate::trustee::{self, get_image_pcrs};
use operator::{ControllerError, KIND_LABEL_KEY, LONG_REQUEUE, OperatorContext, upsert_condition};
use operator::{controller_error_policy, controller_info, create_or_info_if_exists};
use trusted_cluster_operator_lib::{conditions::*, reference_values::*, *};
use trusted_cluster_operator_lib::{conditions::*, record_event, reference_values::*, *};

const APPROVED_IMAGE_ANNOTATION: &str = "approved-image";
const PCR_COMMAND_NAME: &str = "compute-pcrs";
Expand Down Expand Up @@ -136,6 +142,31 @@ async fn job_reconcile(
delete.map_err(Into::<anyhow::Error>::into)?;
let image_pcrs = cached_image_pcrs(&ctx)?;
trustee::update_reference_values(&ctx, image_pcrs).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(LONG_REQUEUE)
}

Expand Down Expand Up @@ -297,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 image_ref: ObjectReference = image.object_ref(&());
let (action, reason) = match handle_new_image(ctx, image).await {
Ok(reason) => (LONG_REQUEUE, reason),
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)
}
Expand Down
Loading