Skip to content

Commit 904c94b

Browse files
committed
fix(compute): preserve docker tracing after decoupling
Signed-off-by: Drew Newberry <385+drew@users.noreply.github.com>
1 parent 8d0eed7 commit 904c94b

2 files changed

Lines changed: 160 additions & 42 deletions

File tree

  • crates

crates/openshell-driver-docker/src/lib.rs

Lines changed: 156 additions & 36 deletions
Original file line numberDiff line numberDiff line change
@@ -55,6 +55,7 @@ use openshell_core::proto_struct::{
5555
deserialize_optional_non_empty_string_list, struct_to_json_value,
5656
};
5757
use openshell_core::{Error, Result as CoreResult};
58+
use opentelemetry::trace::TraceContextExt as _;
5859
use std::collections::{HashMap, HashSet};
5960
use std::future::Future;
6061
use std::net::{IpAddr, Ipv4Addr, SocketAddr};
@@ -67,7 +68,8 @@ use tokio::sync::{Mutex, broadcast, mpsc};
6768
use tokio::task::JoinHandle;
6869
use tokio_stream::wrappers::ReceiverStream;
6970
use tonic::{Request, Response, Status};
70-
use tracing::{debug, info, warn};
71+
use tracing::{Instrument as _, debug, info, warn};
72+
use tracing_opentelemetry::OpenTelemetrySpanExt as _;
7173
use url::Url;
7274

7375
const WATCH_BUFFER: usize = 128;
@@ -84,6 +86,28 @@ const HOST_OPENSHELL_INTERNAL: &str = "host.openshell.internal";
8486
const HOST_DOCKER_INTERNAL: &str = "host.docker.internal";
8587
const DOCKER_NETWORK_DRIVER: &str = "bridge";
8688

89+
fn provisioning_span(
90+
parent: &opentelemetry::Context,
91+
sandbox: &DriverSandbox,
92+
image_ref: &str,
93+
) -> tracing::Span {
94+
let span = tracing::info_span!(
95+
parent: None,
96+
"docker.provision",
97+
otel.name = "docker.provision",
98+
otel.status_code = tracing::field::Empty,
99+
sandbox.id = %sandbox.id,
100+
sandbox.name = %sandbox.name,
101+
image.ref = %image_ref,
102+
);
103+
let parent_span_context = parent.span().span_context().clone();
104+
if parent_span_context.is_valid() {
105+
let parent = opentelemetry::Context::new().with_remote_span_context(parent_span_context);
106+
let _ = span.set_parent(parent);
107+
}
108+
span
109+
}
110+
87111
/// Gateway-local configuration for the Docker compute driver.
88112
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
89113
#[serde(default, deny_unknown_fields)]
@@ -877,7 +901,7 @@ impl DockerComputeDriver {
877901
&sandbox.id,
878902
"Scheduled",
879903
format!("Docker sandbox accepted for image \"{image}\""),
880-
HashMap::from([("image_ref".to_string(), image)]),
904+
HashMap::from([("image_ref".to_string(), image.clone())]),
881905
);
882906
self.publish_sandbox_snapshot(pending_sandbox_snapshot(
883907
sandbox,
@@ -889,9 +913,14 @@ impl DockerComputeDriver {
889913
let driver = self.clone();
890914
let sandbox_for_task = sandbox.clone();
891915
let sandbox_id = sandbox.id.clone();
892-
let task = tokio::spawn(async move {
893-
driver.provision_sandbox(sandbox_for_task).await;
894-
});
916+
let parent = tracing::Span::current().context();
917+
let provisioning_span = provisioning_span(&parent, sandbox, &image);
918+
let task = tokio::spawn(
919+
async move {
920+
driver.provision_sandbox(sandbox_for_task).await;
921+
}
922+
.instrument(provisioning_span),
923+
);
895924

896925
let mut pending = self.pending.lock().await;
897926
if let Some(record) = pending.get_mut(&sandbox_id) {
@@ -914,20 +943,42 @@ impl DockerComputeDriver {
914943
}
915944
}
916945

946+
#[tracing::instrument(
947+
name = "docker.provision_sandbox",
948+
skip(self, sandbox),
949+
fields(
950+
otel.name = "docker.provision_sandbox",
951+
otel.status_code = tracing::field::Empty,
952+
sandbox.id = %sandbox.id,
953+
sandbox.name = %sandbox.name,
954+
)
955+
)]
917956
async fn provision_sandbox_inner(
918957
&self,
919958
sandbox: &DriverSandbox,
920959
) -> Result<(), DockerProvisioningFailure> {
960+
let span_status = openshell_otel::ErrorStatusGuard::current();
921961
let validated = Self::validated_sandbox(sandbox, &self.config).map_err(|status| {
922962
DockerProvisioningFailure::new("ContainerCreateFailed", status.message())
923963
})?;
924964
let template = validated.template;
925-
let image = self
926-
.ensure_image_available(&sandbox.id, &template.image)
927-
.await
928-
.map_err(|status| {
929-
DockerProvisioningFailure::new("ImagePullFailed", status.message())
930-
})?;
965+
let image = async {
966+
openshell_otel::record_error_result(
967+
self.ensure_image_available(&sandbox.id, &template.image)
968+
.await
969+
.map_err(|status| {
970+
DockerProvisioningFailure::new("ImagePullFailed", status.message())
971+
}),
972+
)
973+
}
974+
.instrument(tracing::info_span!(
975+
"docker.prepare_image",
976+
otel.name = "docker.prepare_image",
977+
otel.status_code = tracing::field::Empty,
978+
sandbox.id = %sandbox.id,
979+
image.ref = %template.image,
980+
))
981+
.await?;
931982
let token_file_created = write_sandbox_token_file(sandbox, &self.config)
932983
.await
933984
.map_err(|status| {
@@ -961,33 +1012,58 @@ impl DockerComputeDriver {
9611012
}
9621013
DockerProvisioningFailure::new("ContainerCreateFailed", status.message())
9631014
})?;
964-
self.docker
965-
.create_container(
966-
Some(
967-
CreateContainerOptionsBuilder::default()
968-
.name(container_name.as_str())
969-
.build(),
970-
),
971-
create_body,
1015+
async {
1016+
openshell_otel::record_error_result(
1017+
self.docker
1018+
.create_container(
1019+
Some(
1020+
CreateContainerOptionsBuilder::default()
1021+
.name(container_name.as_str())
1022+
.build(),
1023+
),
1024+
create_body,
1025+
)
1026+
.await
1027+
.map_err(|err| {
1028+
if token_file_created {
1029+
cleanup_sandbox_token_file(sandbox, &self.config);
1030+
}
1031+
DockerProvisioningFailure::from_status(
1032+
"ContainerCreateFailed",
1033+
create_status_from_docker_error("create docker sandbox container", err),
1034+
)
1035+
}),
9721036
)
973-
.await
974-
.map_err(|err| {
975-
if token_file_created {
976-
cleanup_sandbox_token_file(sandbox, &self.config);
977-
}
978-
DockerProvisioningFailure::from_status(
979-
"ContainerCreateFailed",
980-
create_status_from_docker_error("create docker sandbox container", err),
981-
)
982-
})?;
1037+
}
1038+
.instrument(tracing::info_span!(
1039+
"docker.create_container",
1040+
otel.name = "docker.create_container",
1041+
otel.status_code = tracing::field::Empty,
1042+
sandbox.id = %sandbox.id,
1043+
container.name = %container_name,
1044+
))
1045+
.await?;
9831046
self.publish_docker_progress(
9841047
&sandbox.id,
9851048
"Created",
9861049
format!("Created Docker container \"{container_name}\""),
9871050
HashMap::from([("container_name".to_string(), container_name.clone())]),
9881051
);
9891052

990-
if let Err(err) = self.docker.start_container(&container_name, None).await {
1053+
let start_result = async {
1054+
openshell_otel::record_error_result(
1055+
self.docker.start_container(&container_name, None).await,
1056+
)
1057+
}
1058+
.instrument(tracing::info_span!(
1059+
"docker.start_container",
1060+
otel.name = "docker.start_container",
1061+
otel.status_code = tracing::field::Empty,
1062+
sandbox.id = %sandbox.id,
1063+
container.name = %container_name,
1064+
))
1065+
.await;
1066+
if let Err(err) = start_result {
9911067
let cleanup = self
9921068
.docker
9931069
.remove_container(
@@ -1028,7 +1104,7 @@ impl DockerComputeDriver {
10281104
);
10291105
}
10301106

1031-
Ok(())
1107+
span_status.finish(Ok(()))
10321108
}
10331109

10341110
async fn delete_sandbox_inner(
@@ -1142,17 +1218,29 @@ impl DockerComputeDriver {
11421218
/// Returns `Ok(true)` when a container existed and was started (or was
11431219
/// already running), `Ok(false)` when no managed container is found for
11441220
/// the sandbox, and `Err(...)` for any Docker failure.
1221+
#[tracing::instrument(
1222+
name = "docker.start_sandbox",
1223+
skip(self),
1224+
fields(
1225+
otel.name = "docker.start_sandbox",
1226+
otel.status_code = tracing::field::Empty,
1227+
sandbox.id = %sandbox_id,
1228+
sandbox.name = %sandbox_name,
1229+
)
1230+
)]
11451231
pub async fn start_sandbox(
11461232
&self,
11471233
sandbox_id: &str,
11481234
sandbox_name: &str,
11491235
) -> Result<bool, Status> {
1236+
let span_status = openshell_otel::ErrorStatusGuard::current();
1237+
require_sandbox_identifier(sandbox_id, sandbox_name)?;
11501238
self.lifecycle_event_fences.begin_start(sandbox_id);
11511239
let result = self
11521240
.start_sandbox_with_lifecycle_fence(sandbox_id, sandbox_name)
11531241
.await;
11541242
self.lifecycle_event_fences.finish_start(sandbox_id);
1155-
result
1243+
span_status.finish(result)
11561244
}
11571245

11581246
async fn start_sandbox_with_lifecycle_fence(
@@ -1924,38 +2012,59 @@ impl ComputeDriver for DockerComputeDriver {
19242012
}))
19252013
}
19262014

2015+
#[tracing::instrument(
2016+
name = "docker.schedule_sandbox",
2017+
skip(self, request),
2018+
fields(
2019+
otel.name = "docker.schedule_sandbox",
2020+
otel.status_code = tracing::field::Empty,
2021+
sandbox.id = %request.get_ref().sandbox.as_ref().map_or("", |sandbox| sandbox.id.as_str()),
2022+
sandbox.name = %request.get_ref().sandbox.as_ref().map_or("", |sandbox| sandbox.name.as_str()),
2023+
)
2024+
)]
19272025
async fn create_sandbox(
19282026
&self,
19292027
request: Request<CreateSandboxRequest>,
19302028
) -> Result<Response<CreateSandboxResponse>, Status> {
2029+
let span_status = openshell_otel::ErrorStatusGuard::current();
19312030
let sandbox = request
19322031
.into_inner()
19332032
.sandbox
19342033
.ok_or_else(|| Status::invalid_argument("sandbox is required"))?;
19352034
self.create_sandbox_inner(&sandbox).await?;
1936-
Ok(Response::new(CreateSandboxResponse {}))
2035+
span_status.finish(Ok(Response::new(CreateSandboxResponse {})))
19372036
}
19382037

2038+
#[tracing::instrument(
2039+
name = "docker.stop_sandbox",
2040+
skip(self, request),
2041+
fields(
2042+
otel.name = "docker.stop_sandbox",
2043+
otel.status_code = tracing::field::Empty,
2044+
sandbox.id = %request.get_ref().sandbox_id,
2045+
sandbox.name = %request.get_ref().sandbox_name,
2046+
)
2047+
)]
19392048
async fn stop_sandbox(
19402049
&self,
19412050
request: Request<StopSandboxRequest>,
19422051
) -> Result<Response<StopSandboxResponse>, Status> {
2052+
let span_status = openshell_otel::ErrorStatusGuard::current();
19432053
let request = request.into_inner();
19442054
require_sandbox_identifier(&request.sandbox_id, &request.sandbox_name)?;
19452055

19462056
self.stop_sandbox_inner(&request.sandbox_id, &request.sandbox_name)
19472057
.await?;
19482058
self.publish_container_snapshot(&request.sandbox_id, &request.sandbox_name)
19492059
.await?;
1950-
Ok(Response::new(StopSandboxResponse {}))
2060+
span_status.finish(Ok(Response::new(StopSandboxResponse {})))
19512061
}
19522062

19532063
async fn start_sandbox(
19542064
&self,
19552065
request: Request<StartSandboxRequest>,
19562066
) -> Result<Response<StartSandboxResponse>, Status> {
19572067
let request = request.into_inner();
1958-
require_sandbox_identifier(&request.sandbox_id, &request.sandbox_name)?;
19592068
if !Self::start_sandbox(self, &request.sandbox_id, &request.sandbox_name).await? {
19602069
return Err(Status::not_found("sandbox not found"));
19612070
}
@@ -1964,10 +2073,21 @@ impl ComputeDriver for DockerComputeDriver {
19642073
Ok(Response::new(StartSandboxResponse {}))
19652074
}
19662075

2076+
#[tracing::instrument(
2077+
name = "docker.delete_sandbox",
2078+
skip(self, request),
2079+
fields(
2080+
otel.name = "docker.delete_sandbox",
2081+
otel.status_code = tracing::field::Empty,
2082+
sandbox.id = %request.get_ref().sandbox_id,
2083+
sandbox.name = %request.get_ref().sandbox_name,
2084+
)
2085+
)]
19672086
async fn delete_sandbox(
19682087
&self,
19692088
request: Request<DeleteSandboxRequest>,
19702089
) -> Result<Response<DeleteSandboxResponse>, Status> {
2090+
let span_status = openshell_otel::ErrorStatusGuard::current();
19712091
let request = request.into_inner();
19722092
require_sandbox_identifier(&request.sandbox_id, &request.sandbox_name)?;
19732093

@@ -1986,7 +2106,7 @@ impl ComputeDriver for DockerComputeDriver {
19862106
});
19872107
}
19882108

1989-
Ok(Response::new(DeleteSandboxResponse { deleted }))
2109+
span_status.finish(Ok(Response::new(DeleteSandboxResponse { deleted })))
19902110
}
19912111

19922112
async fn watch_sandboxes(

crates/openshell-server/src/lib.rs

Lines changed: 4 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1409,7 +1409,7 @@ async fn build_compute_runtime(
14091409
inherited_config_keys: registration.inherited_config_keys,
14101410
};
14111411
let instance = registration.factory.build(build_context).await?;
1412-
let runtime = match instance {
1412+
match instance {
14131413
ComputeDriverInstance::InProcess(driver) => ComputeRuntime::from_driver(
14141414
registration.name,
14151415
driver,
@@ -1439,8 +1439,7 @@ async fn build_compute_runtime(
14391439
Error::execution(format!("failed to create compute runtime: {error}"))
14401440
})?
14411441
}
1442-
};
1443-
runtime
1442+
}
14441443
}
14451444
ConfiguredComputeDriver::Remote { name } => {
14461445
let remote_config =
@@ -1453,7 +1452,7 @@ async fn build_compute_runtime(
14531452
let endpoint = compute::connect_remote_compute_driver(name, &remote_config.socket_path)
14541453
.await
14551454
.map_err(|e| Error::execution(format!("failed to create compute runtime: {e}")))?;
1456-
let runtime = ComputeRuntime::new_remote_driver(
1455+
ComputeRuntime::new_remote_driver(
14571456
endpoint,
14581457
store,
14591458
sandbox_index,
@@ -1462,8 +1461,7 @@ async fn build_compute_runtime(
14621461
supervisor_sessions,
14631462
)
14641463
.await
1465-
.map_err(|e| Error::execution(format!("failed to create compute runtime: {e}")))?;
1466-
runtime
1464+
.map_err(|e| Error::execution(format!("failed to create compute runtime: {e}")))?
14671465
}
14681466
};
14691467

0 commit comments

Comments
 (0)