Skip to main content

fakecloud_cloudcontrol/
service.rs

1//! Cloud Control API awsJson1.0 dispatch + provisioner bridge.
2
3use std::sync::Arc;
4
5use async_trait::async_trait;
6use chrono::Utc;
7use http::StatusCode;
8use serde_json::{json, Value};
9use tokio::sync::Mutex as AsyncMutex;
10
11use fakecloud_cloudformation::CloudFormationService;
12use fakecloud_core::service::{AwsRequest, AwsResponse, AwsService, AwsServiceError};
13use fakecloud_persistence::SnapshotStore;
14
15use crate::patch::apply_json_patch;
16use crate::persistence::save_snapshot;
17use crate::state::{CloudControlState, ManagedResource, ResourceRequest, SharedCloudControlState};
18
19/// Every operation name in the Cloud Control Smithy model.
20pub const CLOUDCONTROL_ACTIONS: &[&str] = &[
21    "CreateResource",
22    "GetResource",
23    "UpdateResource",
24    "DeleteResource",
25    "ListResources",
26    "GetResourceRequestStatus",
27    "ListResourceRequests",
28    "CancelResourceRequest",
29];
30
31pub struct CloudControlService {
32    cfn: Arc<CloudFormationService>,
33    state: SharedCloudControlState,
34    snapshot_store: Option<Arc<dyn SnapshotStore>>,
35    snapshot_lock: Arc<AsyncMutex<()>>,
36}
37
38impl CloudControlService {
39    pub fn new(cfn: Arc<CloudFormationService>, state: SharedCloudControlState) -> Self {
40        Self {
41            cfn,
42            state,
43            snapshot_store: None,
44            snapshot_lock: Arc::new(AsyncMutex::new(())),
45        }
46    }
47
48    pub fn with_snapshot_store(mut self, store: Arc<dyn SnapshotStore>) -> Self {
49        self.snapshot_store = Some(store);
50        self
51    }
52
53    /// Persist hook for callers that change this service's state from outside
54    /// (the reset endpoints): writes the current snapshot. `None` in memory
55    /// mode.
56    pub fn snapshot_hook(&self) -> Option<fakecloud_persistence::SnapshotHook> {
57        let store = self.snapshot_store.clone()?;
58        Some(fakecloud_persistence::snapshot_hook(
59            self.state.clone(),
60            store,
61            self.snapshot_lock.clone(),
62            |state, store, lock| async move {
63                save_snapshot(&state, Some(store), &lock).await;
64            },
65        ))
66    }
67
68    async fn persist(&self) {
69        save_snapshot(
70            &self.state,
71            self.snapshot_store.clone(),
72            &self.snapshot_lock,
73        )
74        .await;
75    }
76
77    // --- CreateResource ---------------------------------------------------
78
79    async fn create_resource(&self, req: &AwsRequest) -> Result<AwsResponse, AwsServiceError> {
80        let body = parse_json(&req.body)?;
81        // Enforce the input members' Smithy @length/@pattern constraints before
82        // touching the provisioner.
83        validate_type_name(&body)?;
84        validate_len(&body, "DesiredState", 1, 262144)?;
85        validate_len(&body, "ClientToken", 1, 128)?;
86        validate_len(&body, "RoleArn", 20, 2048)?;
87        let type_name = require_str(&body, "TypeName")?;
88        let desired = require_str(&body, "DesiredState")?;
89        let desired_state: Value = serde_json::from_str(desired)
90            .map_err(|e| invalid_request(&format!("DesiredState is not valid JSON: {e}")))?;
91        let client_token = opt_str(&body, "ClientToken");
92        let fingerprint =
93            fingerprint_of("CREATE", type_name, None, Some(&desired_state.to_string()));
94
95        // ClientToken idempotency: replay the original terminal event when the
96        // parameters match, reject the reuse when they differ.
97        if let Some(token) = &client_token {
98            if let Some(resp) =
99                self.client_token_replay_or_conflict(&req.account_id, token, &fingerprint)
100            {
101                return resp;
102            }
103        }
104
105        let request_token = new_token();
106        let outcome = self.cfn.cloudcontrol_create(
107            type_name,
108            desired_state.clone(),
109            &req.account_id,
110            &req.region,
111        );
112
113        let record = match outcome {
114            Ok(res) => {
115                let managed = ManagedResource {
116                    type_name: type_name.to_string(),
117                    identifier: res.physical_id.clone(),
118                    properties: desired_state.clone(),
119                    attributes: res.attributes.clone(),
120                    created_at: Utc::now(),
121                };
122                let mut accounts = self.state.write();
123                let st = accounts.get_or_create(&req.account_id);
124                st.resources.insert(
125                    CloudControlState::resource_key(type_name, &res.physical_id),
126                    managed,
127                );
128                success_request(
129                    &request_token,
130                    type_name,
131                    Some(res.physical_id),
132                    "CREATE",
133                    Some(desired_state),
134                    client_token,
135                    Some(fingerprint),
136                )
137            }
138            Err(msg) => failed_request(
139                &request_token,
140                type_name,
141                None,
142                "CREATE",
143                &msg,
144                client_token,
145                Some(fingerprint),
146            ),
147        };
148
149        self.store_request(&req.account_id, record.clone());
150        if record.operation_status == "SUCCESS" {
151            // Persist the backing service state the provisioner just mutated,
152            // so a restart doesn't leave this resource with no owning state.
153            self.cfn.cloudcontrol_persist_type(type_name).await;
154        }
155        self.persist().await;
156        Ok(progress_event_response(&record))
157    }
158
159    // --- GetResource ------------------------------------------------------
160
161    fn get_resource(&self, req: &AwsRequest) -> Result<AwsResponse, AwsServiceError> {
162        let body = parse_json(&req.body)?;
163        validate_type_name(&body)?;
164        validate_len(&body, "Identifier", 1, 1024)?;
165        let type_name = require_str(&body, "TypeName")?;
166        let identifier = require_str(&body, "Identifier")?;
167        let accounts = self.state.read();
168        let managed = accounts
169            .get(&req.account_id)
170            .and_then(|s| {
171                s.resources
172                    .get(&CloudControlState::resource_key(type_name, identifier))
173            })
174            .ok_or_else(|| resource_not_found(type_name, identifier))?;
175        Ok(AwsResponse::json_value(
176            StatusCode::OK,
177            json!({
178                "TypeName": type_name,
179                "ResourceDescription": resource_description(managed),
180            }),
181        ))
182    }
183
184    // --- UpdateResource ---------------------------------------------------
185
186    async fn update_resource(&self, req: &AwsRequest) -> Result<AwsResponse, AwsServiceError> {
187        let body = parse_json(&req.body)?;
188        validate_type_name(&body)?;
189        validate_len(&body, "Identifier", 1, 1024)?;
190        validate_len(&body, "PatchDocument", 1, 262144)?;
191        validate_len(&body, "ClientToken", 1, 128)?;
192        validate_len(&body, "RoleArn", 20, 2048)?;
193        let type_name = require_str(&body, "TypeName")?;
194        let identifier = require_str(&body, "Identifier")?;
195        let patch_str = require_str(&body, "PatchDocument")?;
196        let patch: Value = serde_json::from_str(patch_str)
197            .map_err(|e| invalid_request(&format!("PatchDocument is not valid JSON: {e}")))?;
198        let client_token = opt_str(&body, "ClientToken");
199        let fingerprint = fingerprint_of("UPDATE", type_name, Some(identifier), Some(patch_str));
200
201        if let Some(token) = &client_token {
202            if let Some(resp) =
203                self.client_token_replay_or_conflict(&req.account_id, token, &fingerprint)
204            {
205                return resp;
206            }
207        }
208
209        // Load current desired state + attributes.
210        let key = CloudControlState::resource_key(type_name, identifier);
211        let (mut properties, attributes) = {
212            let accounts = self.state.read();
213            let managed = accounts
214                .get(&req.account_id)
215                .and_then(|s| s.resources.get(&key))
216                .ok_or_else(|| resource_not_found(type_name, identifier))?;
217            (managed.properties.clone(), managed.attributes.clone())
218        };
219
220        apply_json_patch(&mut properties, &patch)
221            .map_err(|e| invalid_request(&format!("invalid PatchDocument: {e}")))?;
222
223        let request_token = new_token();
224        let outcome = self.cfn.cloudcontrol_update(
225            type_name,
226            identifier,
227            &attributes,
228            properties.clone(),
229            &req.account_id,
230            &req.region,
231        );
232
233        let record = match outcome {
234            Ok(res) => {
235                let new_id = res.physical_id.clone();
236                let mut accounts = self.state.write();
237                let st = accounts.get_or_create(&req.account_id);
238                if new_id != identifier {
239                    // Replacement update: the provisioner assigned a new
240                    // physical id. Re-key the managed resource so subsequent
241                    // Get/Delete track the new identity instead of the stale one.
242                    let mut managed = st.resources.remove(&key).unwrap_or(ManagedResource {
243                        type_name: type_name.to_string(),
244                        identifier: new_id.clone(),
245                        properties: properties.clone(),
246                        attributes: res.attributes.clone(),
247                        created_at: Utc::now(),
248                    });
249                    managed.identifier = new_id.clone();
250                    managed.properties = properties.clone();
251                    managed.attributes = res.attributes.clone();
252                    st.resources
253                        .insert(CloudControlState::resource_key(type_name, &new_id), managed);
254                } else if let Some(managed) = st.resources.get_mut(&key) {
255                    managed.properties = properties.clone();
256                    managed.attributes = res.attributes.clone();
257                }
258                success_request(
259                    &request_token,
260                    type_name,
261                    Some(new_id),
262                    "UPDATE",
263                    Some(properties),
264                    client_token,
265                    Some(fingerprint),
266                )
267            }
268            Err(msg) => failed_request(
269                &request_token,
270                type_name,
271                Some(identifier.to_string()),
272                "UPDATE",
273                &msg,
274                client_token,
275                Some(fingerprint),
276            ),
277        };
278
279        self.store_request(&req.account_id, record.clone());
280        if record.operation_status == "SUCCESS" {
281            self.cfn.cloudcontrol_persist_type(type_name).await;
282        }
283        self.persist().await;
284        Ok(progress_event_response(&record))
285    }
286
287    // --- DeleteResource ---------------------------------------------------
288
289    async fn delete_resource(&self, req: &AwsRequest) -> Result<AwsResponse, AwsServiceError> {
290        let body = parse_json(&req.body)?;
291        validate_type_name(&body)?;
292        validate_len(&body, "Identifier", 1, 1024)?;
293        validate_len(&body, "ClientToken", 1, 128)?;
294        validate_len(&body, "RoleArn", 20, 2048)?;
295        let type_name = require_str(&body, "TypeName")?;
296        let identifier = require_str(&body, "Identifier")?;
297        let client_token = opt_str(&body, "ClientToken");
298        let fingerprint = fingerprint_of("DELETE", type_name, Some(identifier), None);
299
300        if let Some(token) = &client_token {
301            if let Some(resp) =
302                self.client_token_replay_or_conflict(&req.account_id, token, &fingerprint)
303            {
304                return resp;
305            }
306        }
307
308        let key = CloudControlState::resource_key(type_name, identifier);
309        let attributes = {
310            let accounts = self.state.read();
311            let managed = accounts
312                .get(&req.account_id)
313                .and_then(|s| s.resources.get(&key))
314                .ok_or_else(|| resource_not_found(type_name, identifier))?;
315            managed.attributes.clone()
316        };
317
318        let request_token = new_token();
319        let outcome = self.cfn.cloudcontrol_delete(
320            type_name,
321            identifier,
322            &attributes,
323            &req.account_id,
324            &req.region,
325        );
326
327        let record = match outcome {
328            Ok(()) => {
329                let mut accounts = self.state.write();
330                let st = accounts.get_or_create(&req.account_id);
331                st.resources.remove(&key);
332                success_request(
333                    &request_token,
334                    type_name,
335                    Some(identifier.to_string()),
336                    "DELETE",
337                    None,
338                    client_token,
339                    Some(fingerprint),
340                )
341            }
342            Err(msg) => failed_request(
343                &request_token,
344                type_name,
345                Some(identifier.to_string()),
346                "DELETE",
347                &msg,
348                client_token,
349                Some(fingerprint),
350            ),
351        };
352
353        self.store_request(&req.account_id, record.clone());
354        if record.operation_status == "SUCCESS" {
355            self.cfn.cloudcontrol_persist_type(type_name).await;
356        }
357        self.persist().await;
358        Ok(progress_event_response(&record))
359    }
360
361    // --- ListResources ----------------------------------------------------
362
363    fn list_resources(&self, req: &AwsRequest) -> Result<AwsResponse, AwsServiceError> {
364        let body = parse_json(&req.body)?;
365        // Enforce the input members' Smithy @length/@pattern/@range constraints.
366        validate_type_name(&body)?;
367        validate_len(&body, "TypeVersionId", 1, 128)?;
368        validate_len(&body, "RoleArn", 20, 2048)?;
369        validate_len(&body, "NextToken", 1, 4096)?;
370        validate_len(&body, "ResourceModel", 1, 262144)?;
371        validate_max_results(&body)?;
372        let type_name = require_str(&body, "TypeName")?;
373        let all: Vec<Value> = {
374            let accounts = self.state.read();
375            accounts
376                .get(&req.account_id)
377                .map(|s| {
378                    s.resources
379                        .values()
380                        .filter(|m| m.type_name == type_name)
381                        .map(resource_description)
382                        .collect()
383                })
384                .unwrap_or_default()
385        };
386        let (page, next) = paginate(all, &body);
387        let mut resp = json!({ "TypeName": type_name, "ResourceDescriptions": page });
388        if let Some(nt) = next {
389            resp["NextToken"] = json!(nt);
390        }
391        Ok(AwsResponse::json_value(StatusCode::OK, resp))
392    }
393
394    // --- Request ledger ops -----------------------------------------------
395
396    fn get_request_status(&self, req: &AwsRequest) -> Result<AwsResponse, AwsServiceError> {
397        let body = parse_json(&req.body)?;
398        let token = require_str(&body, "RequestToken")?;
399        let accounts = self.state.read();
400        let record = accounts
401            .get(&req.account_id)
402            .and_then(|s| s.requests.get(token))
403            .ok_or_else(|| request_token_not_found(token))?;
404        Ok(progress_event_response(record))
405    }
406
407    fn list_requests(&self, req: &AwsRequest) -> Result<AwsResponse, AwsServiceError> {
408        let body = parse_json(&req.body)?;
409        validate_len(&body, "NextToken", 1, 2048)?;
410        validate_max_results(&body)?;
411        let filter_statuses: Vec<String> = body
412            .get("ResourceRequestStatusFilter")
413            .and_then(|f| f.get("OperationStatuses"))
414            .and_then(|v| v.as_array())
415            .map(|a| {
416                a.iter()
417                    .filter_map(|s| s.as_str().map(String::from))
418                    .collect()
419            })
420            .unwrap_or_default();
421        let filter_ops: Vec<String> = body
422            .get("ResourceRequestStatusFilter")
423            .and_then(|f| f.get("Operations"))
424            .and_then(|v| v.as_array())
425            .map(|a| {
426                a.iter()
427                    .filter_map(|s| s.as_str().map(String::from))
428                    .collect()
429            })
430            .unwrap_or_default();
431        let all: Vec<Value> = {
432            let accounts = self.state.read();
433            accounts
434                .get(&req.account_id)
435                .map(|s| {
436                    s.requests
437                        .values()
438                        .filter(|r| {
439                            filter_statuses.is_empty()
440                                || filter_statuses.contains(&r.operation_status)
441                        })
442                        .filter(|r| filter_ops.is_empty() || filter_ops.contains(&r.operation))
443                        .map(progress_event_json)
444                        .collect()
445                })
446                .unwrap_or_default()
447        };
448        let (page, next) = paginate(all, &body);
449        let mut resp = json!({ "ResourceRequestStatusSummaries": page });
450        if let Some(nt) = next {
451            resp["NextToken"] = json!(nt);
452        }
453        Ok(AwsResponse::json_value(StatusCode::OK, resp))
454    }
455
456    async fn cancel_request(&self, req: &AwsRequest) -> Result<AwsResponse, AwsServiceError> {
457        let body = parse_json(&req.body)?;
458        let token = require_str(&body, "RequestToken")?;
459        // Provisioning is synchronous, so a request is already terminal by the
460        // time a client could cancel it. Cancelling a terminal request is a
461        // no-op that echoes its final ProgressEvent (AWS rejects cancelling an
462        // already-completed request, but never invents a new state for it).
463        // Mutate (if non-terminal) inside a scoped block so the write guard is
464        // dropped before any await -- and capture the resulting record.
465        let (record, mutated) = {
466            let mut accounts = self.state.write();
467            let st = accounts.get_or_create(&req.account_id);
468            let r = st
469                .requests
470                .get_mut(token)
471                .ok_or_else(|| request_token_not_found(token))?;
472            let mutated = if !is_terminal(&r.operation_status) {
473                r.operation_status = "CANCEL_COMPLETE".to_string();
474                true
475            } else {
476                false
477            };
478            (r.clone(), mutated)
479        };
480        if mutated {
481            self.persist().await;
482        }
483        Ok(progress_event_response(&record))
484    }
485
486    // --- helpers ----------------------------------------------------------
487
488    fn store_request(&self, account_id: &str, record: ResourceRequest) {
489        let mut accounts = self.state.write();
490        let st = accounts.get_or_create(account_id);
491        st.requests.insert(record.request_token.clone(), record);
492    }
493
494    fn find_by_client_token(&self, account_id: &str, token: &str) -> Option<ResourceRequest> {
495        let accounts = self.state.read();
496        accounts.get(account_id).and_then(|s| {
497            s.requests
498                .values()
499                .find(|r| r.client_token.as_deref() == Some(token))
500                .cloned()
501        })
502    }
503
504    /// Resolve a `ClientToken` against the request ledger: `Some(Ok(..))` to
505    /// replay the original terminal event (same parameters), `Some(Err(..))` to
506    /// reject conflicting reuse, or `None` when the token is unseen.
507    fn client_token_replay_or_conflict(
508        &self,
509        account_id: &str,
510        token: &str,
511        fingerprint: &str,
512    ) -> Option<Result<AwsResponse, AwsServiceError>> {
513        let existing = self.find_by_client_token(account_id, token)?;
514        match existing.fingerprint.as_deref() {
515            // Legacy record from a snapshot written before fingerprints existed:
516            // replay rather than invent a conflict we can't actually prove.
517            None => Some(Ok(progress_event_response(&existing))),
518            Some(fp) if fp == fingerprint => Some(Ok(progress_event_response(&existing))),
519            Some(_) => Some(Err(client_token_conflict(token))),
520        }
521    }
522}
523
524#[async_trait]
525impl AwsService for CloudControlService {
526    fn service_name(&self) -> &str {
527        "cloudcontrolapi"
528    }
529
530    fn supported_actions(&self) -> &[&str] {
531        CLOUDCONTROL_ACTIONS
532    }
533
534    async fn handle(&self, req: AwsRequest) -> Result<AwsResponse, AwsServiceError> {
535        match req.action.as_str() {
536            "CreateResource" => self.create_resource(&req).await,
537            "GetResource" => self.get_resource(&req),
538            "UpdateResource" => self.update_resource(&req).await,
539            "DeleteResource" => self.delete_resource(&req).await,
540            "ListResources" => self.list_resources(&req),
541            "GetResourceRequestStatus" => self.get_request_status(&req),
542            "ListResourceRequests" => self.list_requests(&req),
543            "CancelResourceRequest" => self.cancel_request(&req).await,
544            other => Err(AwsServiceError::aws_error(
545                StatusCode::BAD_REQUEST,
546                "InvalidRequestException",
547                format!("Unknown operation: {other}"),
548            )),
549        }
550    }
551}
552
553// ---------------------------------------------------------------------------
554// Request-record + response builders
555// ---------------------------------------------------------------------------
556
557#[allow(clippy::too_many_arguments)]
558fn success_request(
559    token: &str,
560    type_name: &str,
561    identifier: Option<String>,
562    operation: &str,
563    resource_model: Option<Value>,
564    client_token: Option<String>,
565    fingerprint: Option<String>,
566) -> ResourceRequest {
567    ResourceRequest {
568        request_token: token.to_string(),
569        type_name: type_name.to_string(),
570        identifier,
571        operation: operation.to_string(),
572        operation_status: "SUCCESS".to_string(),
573        event_time: Utc::now(),
574        resource_model,
575        status_message: None,
576        error_code: None,
577        client_token,
578        fingerprint,
579    }
580}
581
582fn failed_request(
583    token: &str,
584    type_name: &str,
585    identifier: Option<String>,
586    operation: &str,
587    message: &str,
588    client_token: Option<String>,
589    fingerprint: Option<String>,
590) -> ResourceRequest {
591    ResourceRequest {
592        request_token: token.to_string(),
593        type_name: type_name.to_string(),
594        identifier,
595        operation: operation.to_string(),
596        operation_status: "FAILED".to_string(),
597        event_time: Utc::now(),
598        resource_model: None,
599        status_message: Some(message.to_string()),
600        error_code: Some("GeneralServiceException".to_string()),
601        client_token,
602        fingerprint,
603    }
604}
605
606/// Stable fingerprint of a mutating request's parameters, used to distinguish
607/// an idempotent `ClientToken` replay from a conflicting reuse.
608fn fingerprint_of(
609    operation: &str,
610    type_name: &str,
611    identifier: Option<&str>,
612    payload: Option<&str>,
613) -> String {
614    format!(
615        "{operation}\u{1f}{type_name}\u{1f}{}\u{1f}{}",
616        identifier.unwrap_or(""),
617        payload.unwrap_or(""),
618    )
619}
620
621fn progress_event_json(record: &ResourceRequest) -> Value {
622    let mut ev = json!({
623        "TypeName": record.type_name,
624        "RequestToken": record.request_token,
625        "Operation": record.operation,
626        "OperationStatus": record.operation_status,
627        "EventTime": record.event_time.timestamp() as f64
628            + record.event_time.timestamp_subsec_millis() as f64 / 1000.0,
629    });
630    if let Some(id) = &record.identifier {
631        ev["Identifier"] = json!(id);
632    }
633    if let Some(model) = &record.resource_model {
634        // ResourceModel is the Properties type: a JSON *string*.
635        ev["ResourceModel"] = json!(model.to_string());
636    }
637    if let Some(msg) = &record.status_message {
638        ev["StatusMessage"] = json!(msg);
639    }
640    if let Some(code) = &record.error_code {
641        ev["ErrorCode"] = json!(code);
642    }
643    ev
644}
645
646fn progress_event_response(record: &ResourceRequest) -> AwsResponse {
647    AwsResponse::json_value(
648        StatusCode::OK,
649        json!({ "ProgressEvent": progress_event_json(record) }),
650    )
651}
652
653fn resource_description(managed: &ManagedResource) -> Value {
654    // The stored `properties` is only the caller's DesiredState. The
655    // provisioner-captured `attributes` (Arn, auto-generated names, the primary
656    // identifier) are read-only properties AWS surfaces on read too, so overlay
657    // them onto the returned model. Attributes are GetAtt-resolvable strings;
658    // parse ones that are JSON-encoded (lists/objects) back into structured
659    // values, keep the rest as plain strings.
660    let mut props = managed.properties.clone();
661    if let Value::Object(map) = &mut props {
662        for (k, v) in &managed.attributes {
663            let value = serde_json::from_str::<Value>(v)
664                .ok()
665                .filter(|parsed| parsed.is_array() || parsed.is_object())
666                .unwrap_or_else(|| json!(v));
667            map.insert(k.clone(), value);
668        }
669    }
670    json!({
671        "Identifier": managed.identifier,
672        // Properties is the Properties type: a JSON *string*.
673        "Properties": props.to_string(),
674    })
675}
676
677fn is_terminal(status: &str) -> bool {
678    matches!(status, "SUCCESS" | "FAILED" | "CANCEL_COMPLETE")
679}
680
681// ---------------------------------------------------------------------------
682// Parsing + errors
683// ---------------------------------------------------------------------------
684
685fn parse_json(body: &[u8]) -> Result<Value, AwsServiceError> {
686    if body.is_empty() {
687        return Ok(json!({}));
688    }
689    serde_json::from_slice(body).map_err(|e| invalid_request(&format!("invalid request body: {e}")))
690}
691
692fn require_str<'a>(body: &'a Value, field: &str) -> Result<&'a str, AwsServiceError> {
693    body.get(field)
694        .and_then(|v| v.as_str())
695        .filter(|s| !s.is_empty())
696        .ok_or_else(|| invalid_request(&format!("{field} is required.")))
697}
698
699fn opt_str(body: &Value, field: &str) -> Option<String> {
700    body.get(field).and_then(|v| v.as_str()).map(String::from)
701}
702
703fn new_token() -> String {
704    uuid::Uuid::new_v4().to_string()
705}
706
707/// Enforce a string member's `@length` when present.
708fn validate_len(body: &Value, field: &str, min: usize, max: usize) -> Result<(), AwsServiceError> {
709    if let Some(s) = body.get(field).and_then(|v| v.as_str()) {
710        if s.len() < min || s.len() > max {
711            return Err(invalid_request(&format!(
712                "Value at '{field}' failed to satisfy constraint: Member must have length between {min} and {max}."
713            )));
714        }
715    }
716    Ok(())
717}
718
719/// `TypeName` is required and constrained to the `Service::Type::Name` shape
720/// (`^[A-Za-z0-9]{2,64}::...`, length 10-196).
721fn validate_type_name(body: &Value) -> Result<(), AwsServiceError> {
722    let name = require_str(body, "TypeName")?;
723    if name.len() < 10 || name.len() > 196 || !is_valid_type_name(name) {
724        return Err(invalid_request(
725            "Value at 'TypeName' failed to satisfy constraint: Member must match the resource type name pattern (e.g. AWS::S3::Bucket).",
726        ));
727    }
728    Ok(())
729}
730
731/// `^[A-Za-z0-9]{2,64}::[A-Za-z0-9]{2,64}::[A-Za-z0-9]{2,64}$`.
732fn is_valid_type_name(name: &str) -> bool {
733    let parts: Vec<&str> = name.split("::").collect();
734    parts.len() == 3
735        && parts
736            .iter()
737            .all(|p| (2..=64).contains(&p.len()) && p.chars().all(|c| c.is_ascii_alphanumeric()))
738}
739
740/// `MaxResults` is range 1..=100 when present.
741fn validate_max_results(body: &Value) -> Result<(), AwsServiceError> {
742    if let Some(v) = body.get("MaxResults") {
743        let n = v
744            .as_i64()
745            .ok_or_else(|| invalid_request("MaxResults must be an integer."))?;
746        if !(1..=100).contains(&n) {
747            return Err(invalid_request("MaxResults must be between 1 and 100."));
748        }
749    }
750    Ok(())
751}
752
753/// Slice a full result list by `MaxResults`/`NextToken`. The `NextToken` is an
754/// opaque zero-based offset into the (stably ordered) full list; a continuation
755/// token is returned only when more results remain. A `NextToken` this API
756/// never issued yields an empty terminal page rather than silently restarting
757/// at page one, so a bad token can't loop the caller or replay pages. (The
758/// Smithy model defines no `InvalidRequest` error on a length-valid `NextToken`,
759/// so an empty page keeps this conformant where a 4xx would not.)
760fn paginate(items: Vec<Value>, body: &Value) -> (Vec<Value>, Option<String>) {
761    let max = body
762        .get("MaxResults")
763        .and_then(|v| v.as_i64())
764        .map(|n| n as usize)
765        .unwrap_or(100);
766    let start = match body.get("NextToken").and_then(|v| v.as_str()) {
767        // Unparseable / never-issued token: treat as past the end so a bad
768        // token yields an empty terminal page instead of restarting page one
769        // (which could loop the caller or replay pages).
770        Some(t) => match t.parse::<usize>() {
771            Ok(n) => n,
772            Err(_) => return (Vec::new(), None),
773        },
774        None => 0,
775    };
776    let total = items.len();
777    if start >= total {
778        return (Vec::new(), None);
779    }
780    let end = start.saturating_add(max).min(total);
781    let next = if end < total {
782        Some(end.to_string())
783    } else {
784        None
785    };
786    (items[start..end].to_vec(), next)
787}
788
789fn invalid_request(msg: &str) -> AwsServiceError {
790    AwsServiceError::aws_error(StatusCode::BAD_REQUEST, "InvalidRequestException", msg)
791}
792
793fn resource_not_found(type_name: &str, identifier: &str) -> AwsServiceError {
794    AwsServiceError::aws_error(
795        StatusCode::NOT_FOUND,
796        "ResourceNotFoundException",
797        format!("Resource of type '{type_name}' with identifier '{identifier}' was not found."),
798    )
799}
800
801fn client_token_conflict(token: &str) -> AwsServiceError {
802    AwsServiceError::aws_error(
803        StatusCode::BAD_REQUEST,
804        "ClientTokenConflictException",
805        format!("The client token '{token}' is already in use with different request parameters."),
806    )
807}
808
809fn request_token_not_found(token: &str) -> AwsServiceError {
810    AwsServiceError::aws_error(
811        StatusCode::NOT_FOUND,
812        "RequestTokenNotFoundException",
813        format!("Request token '{token}' was not found."),
814    )
815}
816
817#[cfg(test)]
818mod tests {
819    use super::*;
820    use std::collections::BTreeMap;
821
822    #[test]
823    fn resource_description_surfaces_read_only_attributes() {
824        // GetResource/ListResources previously returned only the caller's
825        // DesiredState; server-generated read-only props (Arn, generated names,
826        // the primary id) captured as attributes were dropped (bug-hunt).
827        let mut attributes = BTreeMap::new();
828        attributes.insert(
829            "Arn".to_string(),
830            "arn:aws:s3:::generated-bucket".to_string(),
831        );
832        attributes.insert("BucketName".to_string(), "generated-bucket".to_string());
833        // A JSON-encoded list attribute round-trips as structured JSON.
834        attributes.insert("Tags".to_string(), "[{\"Key\":\"a\"}]".to_string());
835        let managed = ManagedResource {
836            type_name: "AWS::S3::Bucket".to_string(),
837            identifier: "generated-bucket".to_string(),
838            properties: json!({ "VersioningConfiguration": { "Status": "Enabled" } }),
839            attributes,
840            created_at: Utc::now(),
841        };
842
843        let desc = resource_description(&managed);
844        let props: Value = serde_json::from_str(desc["Properties"].as_str().unwrap()).unwrap();
845        // Caller's desired state preserved.
846        assert_eq!(props["VersioningConfiguration"]["Status"], "Enabled");
847        // Read-only attributes merged in.
848        assert_eq!(props["Arn"], "arn:aws:s3:::generated-bucket");
849        assert_eq!(props["BucketName"], "generated-bucket");
850        // JSON-encoded attribute parsed back to structure.
851        assert_eq!(props["Tags"][0]["Key"], "a");
852    }
853}