1use std::{
2 future::Future,
3 ops::Deref,
4 pin::Pin,
5 sync::atomic::{AtomicUsize, Ordering::Relaxed},
6 task::{Context, Poll},
7};
8
9use pin_project::{pin_project, pinned_drop};
10use saluki_common::resource_tracking::{ResourceGroupRegistry, ResourceGroupToken, Track as _, Tracked};
11use stringtheory::MetaString;
12use tracing::{debug_span, instrument::Instrumented, Instrument as _};
13
14use super::state::{DataspaceRegistry, CURRENT_DATASPACE};
15
16static GLOBAL_PROCESS_ID_COUNTER: AtomicUsize = AtomicUsize::new(1);
17
18#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
24pub struct Id(usize);
25
26tokio::task_local! {
27 pub(crate) static CURRENT_PROCESS_ID: Id;
28}
29
30impl Id {
31 pub const ROOT: Self = Self(0);
33
34 pub fn new() -> Self {
36 let id = GLOBAL_PROCESS_ID_COUNTER.fetch_add(1, Relaxed);
37 Self(id)
38 }
39
40 pub fn current() -> Self {
44 CURRENT_PROCESS_ID.try_with(|id| *id).unwrap_or(Self::ROOT)
45 }
46
47 pub fn as_usize(&self) -> usize {
49 self.0
50 }
51}
52
53#[derive(Clone, Debug, PartialEq, Eq, Hash)]
70pub struct Name(MetaString);
71
72impl Name {
73 pub(crate) fn root<N: AsRef<str>>(name: N) -> Option<Self> {
74 sanitize_scoped_name(name.as_ref()).map(Self)
75 }
76
77 pub(crate) fn scoped<N: AsRef<str>>(parent: &Name, name: N) -> Option<Self> {
78 let child = sanitize_scoped_name(name.as_ref())?;
79 Some(Self(format!("{}.{}", parent.0, child).into()))
80 }
81}
82
83impl Deref for Name {
84 type Target = str;
85
86 fn deref(&self) -> &Self::Target {
87 &self.0
88 }
89}
90
91#[derive(Clone)]
93pub struct Process {
94 id: Id,
95 name: Name,
96 resource_group_token: ResourceGroupToken,
97 dataspace: DataspaceRegistry,
98}
99
100impl Process {
101 pub(crate) fn supervisor<N: AsRef<str>>(name: N, parent: Option<&Process>) -> Option<Self> {
102 let name = parent
103 .and_then(|p| Name::scoped(&p.name, &name))
104 .or_else(|| Name::root(name))?;
105 let resource_group_token = ResourceGroupRegistry::global().register_resource_group(&*name);
106 let dataspace = parent.map(|p| p.dataspace.clone()).unwrap_or_default();
107 Some(Self::from_parts(Id::new(), name, resource_group_token, dataspace))
108 }
109
110 pub(crate) fn supervisor_with_dataspace<N: AsRef<str>>(
111 name: N, parent: Option<&Process>, dataspace: Option<DataspaceRegistry>,
112 ) -> Option<Self> {
113 let name = parent
114 .and_then(|p| Name::scoped(&p.name, &name))
115 .or_else(|| Name::root(name))?;
116 let resource_group_token = ResourceGroupRegistry::global().register_resource_group(&*name);
117 let dataspace = dataspace
118 .or_else(|| parent.map(|p| p.dataspace.clone()))
119 .unwrap_or_default();
120 Some(Self::from_parts(Id::new(), name, resource_group_token, dataspace))
121 }
122
123 pub(crate) fn worker<N: AsRef<str>>(name: N, parent: &Process) -> Option<Self> {
124 let name = Name::scoped(&parent.name, name)?;
125 Some(Self::from_parts(
126 Id::new(),
127 name,
128 parent.resource_group_token,
129 parent.dataspace.clone(),
130 ))
131 }
132
133 fn from_parts(id: Id, name: Name, resource_group_token: ResourceGroupToken, dataspace: DataspaceRegistry) -> Self {
134 Self {
135 id,
136 name,
137 resource_group_token,
138 dataspace,
139 }
140 }
141
142 pub fn id(&self) -> &Id {
144 &self.id
145 }
146
147 pub(crate) fn name(&self) -> &str {
150 &self.name
151 }
152
153 pub(crate) fn dataspace(&self) -> &DataspaceRegistry {
155 &self.dataspace
156 }
157
158 pub fn into_process_future<F>(self, inner: F) -> ProcessFuture<F>
159 where
160 F: Future,
161 {
162 ProcessFuture::new(self, inner)
163 }
164}
165
166#[pin_project(PinnedDrop)]
171pub struct ProcessFuture<F> {
172 process_id: Id,
173 dataspace: DataspaceRegistry,
174 #[pin]
175 inner: Instrumented<Tracked<F>>,
176}
177
178impl<F> ProcessFuture<F>
179where
180 F: Future,
181{
182 pub(crate) fn new(process: Process, inner: F) -> Self {
183 let span = debug_span!(
184 "process",
185 process_id = process.id().as_usize(),
186 process_name = &*process.name,
187 );
188
189 let process_id = process.id;
190 let dataspace = process.dataspace.clone();
191 let inner = inner.track_resources(process.resource_group_token).instrument(span);
192
193 Self {
194 process_id,
195 dataspace,
196 inner,
197 }
198 }
199}
200
201impl<F> Future for ProcessFuture<F>
202where
203 F: Future,
204{
205 type Output = F::Output;
206
207 fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
208 let this = self.project();
209 CURRENT_PROCESS_ID.sync_scope(*this.process_id, || {
210 CURRENT_DATASPACE.sync_scope(this.dataspace.clone(), || this.inner.poll(cx))
211 })
212 }
213}
214
215#[pinned_drop]
216impl<F> PinnedDrop for ProcessFuture<F> {
217 fn drop(self: Pin<&mut Self>) {
218 let this = self.project();
219 this.dataspace.retract_all_for_process(*this.process_id);
220 }
221}
222
223pub trait ProcessExt {
225 fn into_process_future(self, process: Process) -> ProcessFuture<Self>
227 where
228 Self: Future + Sized;
229}
230
231impl<F> ProcessExt for F
232where
233 F: Future,
234{
235 fn into_process_future(self, process: Process) -> ProcessFuture<Self>
236 where
237 Self: Future + Sized,
238 {
239 process.into_process_future(self)
240 }
241}
242
243fn is_process_name_segment_valid(name: &str) -> bool {
244 if name.is_empty() {
246 return false;
247 }
248
249 if !name.starts_with(|c: char| c.is_alphanumeric()) || !name.ends_with(|c: char| c.is_alphanumeric()) {
251 return false;
252 }
253
254 for c in name.chars() {
258 if !c.is_alphanumeric() && c != '_' {
259 return false;
260 }
261 }
262
263 true
264}
265
266pub fn get_sanitized_name(name: &str) -> MetaString {
284 if is_process_name_segment_valid(name) {
285 name.into()
286 } else {
287 let raw_sanitized = name
289 .chars()
290 .map(|c| if c.is_alphanumeric() || c == '_' { c } else { '_' });
291 let mut sanitized = String::with_capacity(name.len());
292
293 let mut last_was_underscore = true;
294 for c in raw_sanitized {
295 if c == '_' {
296 if !last_was_underscore {
297 sanitized.push(c);
298 last_was_underscore = true;
299 }
300 } else {
301 sanitized.push(c);
302 last_was_underscore = false;
303 }
304 }
305
306 let trimmed = sanitized.trim_matches(|c: char| !c.is_alphanumeric());
308 trimmed.into()
309 }
310}
311
312fn sanitize_scoped_name(name: &str) -> Option<MetaString> {
318 let mut rendered = String::new();
319 for segment in name.split('.') {
320 let sanitized = get_sanitized_name(segment);
321 if sanitized.is_empty() {
322 continue;
323 }
324
325 if !rendered.is_empty() {
326 rendered.push('.');
327 }
328 rendered.push_str(&sanitized);
329 }
330
331 if rendered.is_empty() {
332 None
333 } else {
334 Some(rendered.into())
335 }
336}
337
338#[cfg(test)]
339mod tests {
340 use super::*;
341
342 #[test]
343 fn test_process_name_root() {
344 let cases = [
345 ("topology_sup", Some("topology_sup")),
346 ("worker", Some("worker")),
347 ("worker.", Some("worker")),
348 ("_worker_", Some("worker")),
349 ("worker-123", Some("worker_123")),
350 ("--worker_123", Some("worker_123")),
351 ("worker 123", Some("worker_123")),
352 ("worker===123", Some("worker_123")),
353 ("topology.worker", Some("topology.worker")),
354 ("", None),
355 ];
356
357 for (input, expected) in cases {
358 let name = Name::root(input);
359 assert_eq!(name.as_deref(), expected);
360 }
361 }
362
363 #[test]
364 fn test_process_name_scoped() {
365 let parent = Name::root("topology_sup").unwrap();
366 let cases = [
367 ("worker", Some("topology_sup.worker")),
368 ("worker.", Some("topology_sup.worker")),
369 ("_worker_", Some("topology_sup.worker")),
370 ("worker-123", Some("topology_sup.worker_123")),
371 ("--worker_123", Some("topology_sup.worker_123")),
372 ("worker 123", Some("topology_sup.worker_123")),
373 ("worker===123", Some("topology_sup.worker_123")),
374 ("nested.worker", Some("topology_sup.nested.worker")),
375 ("", None),
376 ];
377
378 for (input, expected) in cases {
379 let name = Name::scoped(&parent, input);
380 assert_eq!(name.as_deref(), expected);
381 }
382 }
383}