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
10 changes: 10 additions & 0 deletions crates/fakecloud-organizations/src/introspection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,16 @@ fn transfer_to_row(org_id: &str, t: &ResponsibilityTransfer) -> ResponsibilityTr
pub fn list_all_responsibility_transfers(
state: &SharedOrganizationsState,
) -> Vec<ResponsibilityTransferRow> {
// Expire past-due handshakes first, exactly as the Organizations API
// does on every call. Without this the route answered `REQUESTED`
// with a live `activeHandshakeId` for a transfer that
// `DescribeResponsibilityTransfer` already reported as `EXPIRED` --
// one object, two contradictory answers, for as long as no API call
// happened to arrive.
let now = Utc::now();
if state.read().has_stale_handshakes(now) {
state.write().expire_stale_handshakes(now);
}
let guard = state.read();
let mut rows: Vec<ResponsibilityTransferRow> = guard
.iter()
Expand Down
19 changes: 19 additions & 0 deletions crates/fakecloud-organizations/src/service/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -475,6 +475,19 @@ impl AwsService for OrganizationsService {

async fn handle(&self, req: AwsRequest) -> Result<AwsResponse, AwsServiceError> {
let mutates = is_mutating_action(&req.action);
// Expire past-due handshakes before answering anything. AWS
// expires a handshake 15 days after it is extended whether or not
// anyone is looking, so the sweep cannot hang off the mutating
// paths alone: a `DescribeHandshake` has to report `EXPIRED` too,
// and an `AcceptHandshake` must be refused rather than reviving
// an offer that lapsed. The read-locked check keeps the write
// lock for the requests that actually have something to expire.
let now = Utc::now();
let expired = if self.state.read().has_stale_handshakes(now) {
self.state.write().expire_stale_handshakes(now)
} else {
0
};
let result = match req.action.as_str() {
"CreateOrganization" => self.create_organization(&req),
"DescribeOrganization" => self.describe_organization(&req),
Expand Down Expand Up @@ -556,6 +569,12 @@ impl AwsService for OrganizationsService {
&req.action,
)),
};
// A sweep is a mutation like any other, even when it happened on
// the way into a read: without this the expiry is lost on
// restart and the handshake comes back OPEN.
if expired > 0 {
self.save_snapshot().await;
}
if mutates && matches!(result.as_ref(), Ok(resp) if resp.status.is_success()) {
self.save_snapshot().await;
// Any successful mutation can have moved an account between OUs,
Expand Down
260 changes: 260 additions & 0 deletions crates/fakecloud-organizations/src/service/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3913,3 +3913,263 @@ async fn deregistering_an_unregistered_administrator_names_the_account() {
"the error must not name the service principal, got: {err}"
);
}

/// AWS gives a handshake 15 days and expires it on its own. Nothing
/// here ever did, so an overdue invitation stayed OPEN and acceptable
/// forever and the `ExpirationTimestamp` fakecloud reported was
/// decoration.
#[tokio::test]
async fn an_overdue_handshake_expires_and_cannot_be_accepted() {
let (svc, state) = OrganizationsService::shared();
create_org_with_root(&svc).await;
let resp = svc
.handle(req_with(
"111111111111",
"InviteAccountToOrganization",
json!({ "Target": {"Id": "222222222222", "Type": "ACCOUNT"} }),
))
.await
.unwrap();
let id = body_json(&resp)["Handshake"]["Id"]
.as_str()
.unwrap()
.to_string();

// Backdate it past its deadline, the way 15 days of wall clock would.
state
.write()
.sole_mut()
.unwrap()
.handshakes
.get_mut(&id)
.unwrap()
.expiration_timestamp = Utc::now() - chrono::Duration::days(1);

let described = svc
.handle(req_with(
"222222222222",
"DescribeHandshake",
json!({ "HandshakeId": id }),
))
.await
.unwrap();
assert_eq!(body_json(&described)["Handshake"]["State"], "EXPIRED");

// And the offer is gone: accepting it is a terminal-transition error,
// not a join.
let err = expect_err(
svc.handle(req_with(
"222222222222",
"AcceptHandshake",
json!({ "HandshakeId": id }),
))
.await,
);
assert_eq!(err.code(), "InvalidHandshakeTransitionException");
assert!(
!state
.read()
.sole()
.unwrap()
.accounts
.contains_key("222222222222"),
"an expired invitation must not enroll its target"
);
}

/// A responsibility transfer rides a handshake, so it ends when that
/// handshake lapses -- leaving it REQUESTED with a live
/// `ActiveHandshakeId` pointing at an EXPIRED handshake made the two
/// records disagree, the same way an unhandled accept once did.
#[tokio::test]
async fn an_expired_handshake_ends_the_transfer_riding_it() {
let (svc, state) = OrganizationsService::shared();
create_org_with_root(&svc).await;
let resp = svc
.handle(req_with(
"111111111111",
"InviteOrganizationToTransferResponsibility",
json!({
"Type": "BILLING",
"SourceName": "handover",
"StartTimestamp": 1893456000.0,
"Target": {"Id": "222222222222", "Type": "ACCOUNT"},
}),
))
.await
.unwrap();
let handshake_id = body_json(&resp)["Handshake"]["Id"]
.as_str()
.unwrap()
.to_string();
let transfer_id = body_json(&resp)["Handshake"]["Resources"]
.as_array()
.unwrap()
.iter()
.find(|r| r["Type"] == "RESPONSIBILITY_TRANSFER")
.unwrap()["Value"]
.as_str()
.unwrap()
.to_string();

let deadline = Utc::now() - chrono::Duration::days(5);
state
.write()
.sole_mut()
.unwrap()
.handshakes
.get_mut(&handshake_id)
.unwrap()
.expiration_timestamp = deadline;

let described = svc
.handle(req_with(
"111111111111",
"DescribeResponsibilityTransfer",
json!({ "Id": transfer_id }),
))
.await
.unwrap();
let transfer = body_json(&described)["ResponsibilityTransfer"].clone();
assert_eq!(transfer["Status"], "EXPIRED");
assert!(
transfer.get("ActiveHandshakeId").is_none(),
"an expired transfer holds no live handshake, got: {transfer}"
);
// The transfer ended when its handshake lapsed, not when the sweep
// happened to notice -- an idle process would otherwise report an
// EndTimestamp days after the ExpirationTimestamp on the same record.
assert_eq!(
transfer["EndTimestamp"].as_f64().unwrap() as i64,
deadline.timestamp()
);

// The introspection route reads the same object and must not
// contradict the API: it answered REQUESTED with a live
// activeHandshakeId until it swept too.
let rows = crate::introspection::list_all_responsibility_transfers(&state);
let row = rows.iter().find(|r| r.id == transfer_id).unwrap();
assert_eq!(row.status, "EXPIRED");
assert!(row.active_handshake_id.is_none());

// The target is free again: the expired offer no longer collides.
svc.handle(req_with(
"111111111111",
"InviteOrganizationToTransferResponsibility",
json!({
"Type": "BILLING",
"SourceName": "handover",
"StartTimestamp": 1893456000.0,
"Target": {"Id": "222222222222", "Type": "ACCOUNT"},
}),
))
.await
.expect("an expired offer must not block a fresh one");
}

/// A handshake still inside its 15 days is untouched by the sweep.
///
/// This guards against OVER-expiry only -- it passes without the sweep
/// too. Its job is to fail if the deadline comparison is ever inverted
/// or the `OPEN | REQUESTED` filter dropped.
#[tokio::test]
async fn a_live_handshake_is_left_alone() {
let (svc, _state) = OrganizationsService::shared();
create_org_with_root(&svc).await;
let resp = svc
.handle(req_with(
"111111111111",
"InviteAccountToOrganization",
json!({ "Target": {"Id": "222222222222", "Type": "ACCOUNT"} }),
))
.await
.unwrap();
let id = body_json(&resp)["Handshake"]["Id"]
.as_str()
.unwrap()
.to_string();
let described = svc
.handle(req_with(
"222222222222",
"DescribeHandshake",
json!({ "HandshakeId": id }),
))
.await
.unwrap();
assert_eq!(body_json(&described)["Handshake"]["State"], "OPEN");
}

/// Recording store: keeps the last bytes written so a test can assert
/// what was actually persisted.
#[derive(Default)]
struct RecordingStore(parking_lot::Mutex<Option<Vec<u8>>>);

impl fakecloud_persistence::SnapshotStore for RecordingStore {
fn load(&self) -> std::io::Result<Option<Vec<u8>>> {
Ok(self.0.lock().clone())
}

fn save(&self, bytes: &[u8]) -> std::io::Result<()> {
*self.0.lock() = Some(bytes.to_vec());
Ok(())
}
}

/// A sweep is a mutation even when it happens on the way into a READ,
/// so it has to be persisted there too. Without that, an expiry noticed
/// by a `DescribeHandshake` is lost on restart and the handshake comes
/// back OPEN -- past its deadline, and acceptable again.
#[tokio::test]
async fn an_expiry_noticed_by_a_read_is_persisted() {
let state: SharedOrganizationsState =
Arc::new(parking_lot::RwLock::new(OrganizationsRegistry::default()));
let store = Arc::new(RecordingStore::default());
let svc = Arc::new(OrganizationsService::new(state.clone()).with_snapshot_store(store.clone()));
svc.handle(req_with("111111111111", "CreateOrganization", json!({})))
.await
.unwrap();
let resp = svc
.handle(req_with(
"111111111111",
"InviteAccountToOrganization",
json!({ "Target": {"Id": "222222222222", "Type": "ACCOUNT"} }),
))
.await
.unwrap();
let id = body_json(&resp)["Handshake"]["Id"]
.as_str()
.unwrap()
.to_string();

state
.write()
.sole_mut()
.unwrap()
.handshakes
.get_mut(&id)
.unwrap()
.expiration_timestamp = Utc::now() - chrono::Duration::days(1);

// A pure read triggers the sweep.
svc.handle(req_with(
"222222222222",
"DescribeHandshake",
json!({ "HandshakeId": id }),
))
.await
.unwrap();

// What landed on disk must carry the expiry, not the OPEN state the
// last mutating call wrote.
let bytes = store
.0
.lock()
.clone()
.expect("the read persisted a snapshot");
let snapshot: OrganizationsSnapshot = serde_json::from_slice(&bytes).unwrap();
let restored = snapshot.into_registry();
assert_eq!(
restored.sole().unwrap().handshakes.get(&id).unwrap().state,
"EXPIRED"
);
}
65 changes: 65 additions & 0 deletions crates/fakecloud-organizations/src/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,28 @@ impl OrganizationsRegistry {
self.orgs.is_empty()
}

/// Does any organization hold a handshake that is past due?
///
/// Cheap enough to run under the read lock on every request, so the
/// write lock the sweep needs is only ever taken when there is
/// something to expire.
pub fn has_stale_handshakes(&self, now: DateTime<Utc>) -> bool {
self.orgs.values().any(|org| {
org.handshakes.values().any(|h| {
matches!(h.state.as_str(), "OPEN" | "REQUESTED") && h.expiration_timestamp <= now
})
})
}

/// Expire every past-due handshake in every organization, and
/// return how many moved.
pub fn expire_stale_handshakes(&mut self, now: DateTime<Utc>) -> usize {
self.orgs
.values_mut()
.map(|org| org.expire_stale_handshakes(now))
.sum()
}

pub fn len(&self) -> usize {
self.orgs.len()
}
Expand Down Expand Up @@ -1084,6 +1106,49 @@ impl OrganizationState {
Ok(snapshot)
}

/// Flip every handshake past its `ExpirationTimestamp` to
/// `EXPIRED`, and return how many moved.
///
/// AWS gives a handshake 15 days and expires it on its own; nothing
/// here ever did, so an overdue invitation stayed `OPEN` and
/// acceptable forever, and the `ExpirationTimestamp` fakecloud
/// reported was decoration. Expiry runs through `resolve_handshake`
/// like every other terminal transition, so a responsibility
/// transfer riding an expired handshake ends with it.
pub fn expire_stale_handshakes(&mut self, now: DateTime<Utc>) -> usize {
let stale: Vec<(String, DateTime<Utc>)> = self
.handshakes
.values()
.filter(|h| {
matches!(h.state.as_str(), "OPEN" | "REQUESTED") && h.expiration_timestamp <= now
})
.map(|h| (h.id.clone(), h.expiration_timestamp))
.collect();
for (id, deadline) in &stale {
// Bind the transfer BEFORE resolving: `resolve_handshake`
// clears the `active_handshake_id` that identifies it.
let riding = self
.responsibility_transfers
.values()
.find(|t| t.active_handshake_id.as_deref() == Some(id.as_str()))
.map(|t| t.id.clone());
// The state check above is exactly the one `resolve_handshake`
// re-applies, so this cannot fail.
let _ = self.resolve_handshake(id, "EXPIRED", None, None);
// A transfer ends when its handshake lapsed, not when the
// sweep noticed. `resolve_handshake` stamps "now", which is
// right for a decline or a cancel -- somebody acted at that
// instant -- but an idle process would otherwise report an
// `EndTimestamp` days after the `ExpirationTimestamp` on the
// same record.
if let Some(transfer) = riding.and_then(|id| self.responsibility_transfers.get_mut(&id))
{
transfer.end_timestamp = Some(*deadline);
}
}
stale.len()
}

/// Every handshake this organization holds.
///
/// Filtering by target is deliberately NOT offered here: deciding
Expand Down
Loading
Loading