1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
//! Logical sandbox and authoritative workspace-checkpoint contracts.
use async_trait::async_trait;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use uuid::Uuid;
use crate::typed_id::SessionId;
/// How a checkpoint's bytes are stored.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SandboxCheckpointKind {
/// A provider-managed snapshot. A fast resume path only: it dies with the
/// provider resource, so it cannot be the recovery path for physical loss.
ProviderNative,
/// An Everruns-owned archive of the working filesystem. The recovery path
/// for physical sandbox loss and provider migration.
PortableWorkspace,
}
impl SandboxCheckpointKind {
pub const fn as_str(self) -> &'static str {
match self {
Self::ProviderNative => "provider_native",
Self::PortableWorkspace => "portable_workspace",
}
}
}
/// Identity and fencing state of a logical sandbox.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct SandboxRef {
pub id: Uuid,
/// Current incarnation. Writes carry the generation they were issued
/// against so a late response from a replaced sandbox cannot advance state.
pub generation: i64,
}
/// A recorded workspace revision.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct SandboxCheckpoint {
pub id: Uuid,
pub sandbox_id: Uuid,
pub generation: i64,
pub source_turn_id: Option<String>,
pub source_tool_call_id: Option<String>,
pub kind: SandboxCheckpointKind,
pub provider_ref: Option<String>,
pub workspace_revision: String,
/// `None` while the checkpoint is uploaded but not yet authoritative.
pub attached_at: Option<DateTime<Utc>>,
pub created_at: DateTime<Utc>,
}
/// A checkpoint upload that has completed but is not yet authoritative.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct NewSandboxCheckpoint {
pub sandbox_id: Uuid,
pub generation: i64,
pub source_turn_id: Option<String>,
pub source_tool_call_id: Option<String>,
pub kind: SandboxCheckpointKind,
pub provider_ref: Option<String>,
pub workspace_revision: String,
}
#[derive(Debug, thiserror::Error)]
pub enum SandboxCheckpointError {
#[error("sandbox checkpoint storage error: {0}")]
Storage(String),
/// The write was issued against a superseded incarnation. The caller must
/// not retry it against the current generation: the bytes it refers to
/// belong to a sandbox that no longer exists.
#[error("sandbox {sandbox_id} is at generation {current}, write carried {carried}")]
StaleGeneration {
sandbox_id: Uuid,
current: i64,
carried: i64,
},
}
/// Hard upper bound on one garbage-collection batch (TM-DOS). Caps the
/// destructive `collect_unattached_checkpoints` regardless of caller input, so a
/// misconfigured limit cannot request an unbounded delete; large backlogs drain
/// over successive sweeps. Mirrors `MAX_RETENTION_PRUNE_LIMIT` in the server
/// storage backend.
pub const MAX_CHECKPOINT_COLLECT_LIMIT: i64 = 1000;
/// Persistence for logical sandboxes and their checkpoints.
#[async_trait]
pub trait SandboxCheckpointStore: Send + Sync {
/// Resolve the logical sandbox for a session/provider pair, creating it on
/// first use. Idempotent.
async fn ensure_sandbox(
&self,
session_id: SessionId,
provider: &str,
) -> Result<SandboxRef, SandboxCheckpointError>;
/// Record a completed upload. The checkpoint is not authoritative yet, so a
/// crash after this point leaves a collectable orphan rather than a pointer
/// to a revision no committed turn produced.
///
/// Recording the same `workspace_revision` twice returns the existing row.
async fn record_checkpoint(
&self,
checkpoint: NewSandboxCheckpoint,
) -> Result<SandboxCheckpoint, SandboxCheckpointError>;
/// Promote a recorded checkpoint to authoritative, fenced on `generation`.
/// Returns `StaleGeneration` if the sandbox has been replaced since the
/// upload, leaving the previous committed checkpoint in place.
///
/// THREAT[TM-TENANT-012]: takes a bare `sandbox_id` with no org filter.
/// Callers MUST obtain it from [`Self::ensure_sandbox`], which scopes to the
/// session the tool context is executing for. Never accept a `sandbox_id`
/// from request input.
async fn attach_checkpoint(
&self,
sandbox_id: Uuid,
checkpoint_id: Uuid,
generation: i64,
) -> Result<(), SandboxCheckpointError>;
/// The last checkpoint committed as authoritative, if any.
///
/// THREAT[TM-TENANT-012]: bare `sandbox_id`, same caller obligation as
/// [`Self::attach_checkpoint`].
async fn current_checkpoint(
&self,
sandbox_id: Uuid,
) -> Result<Option<SandboxCheckpoint>, SandboxCheckpointError>;
/// Reject `checkpoint_id` as authoritative and fall back to the previous
/// attached checkpoint, returning whatever the sandbox now points at.
///
/// This is the reconciliation half of the crash window: a checkpoint is
/// attached before the tool result it belongs to is settled, so a crash in
/// between leaves the workspace ahead of the conversation. Rolling the
/// pointer back detaches the rejected revision, which returns it to the
/// collectable pool rather than deleting it inline.
///
/// Fenced on `generation`, and a no-op when the sandbox no longer points at
/// `checkpoint_id`: in both cases something newer already decided, and this
/// call must not undo it. `Ok(None)` means the sandbox has no earlier
/// committed revision to fall back to.
///
/// THREAT[TM-TENANT-012]: bare `sandbox_id`, same caller obligation as
/// [`Self::attach_checkpoint`].
async fn rollback_current_checkpoint(
&self,
sandbox_id: Uuid,
checkpoint_id: Uuid,
generation: i64,
) -> Result<Option<SandboxCheckpoint>, SandboxCheckpointError>;
/// Delete up to `limit` uploads that were never attached and are older than
/// `before`. Attached checkpoints are never returned or removed. Returns the
/// `workspace_revision` of each collected row so the caller can drop the
/// matching provider-side artifact.
///
/// `limit` bounds one destructive batch so a large backlog drains over
/// successive sweeps instead of issuing one unbounded delete; implementations
/// clamp it to [`MAX_CHECKPOINT_COLLECT_LIMIT`].
///
/// THREAT[TM-TENANT-012]: bare `sandbox_id`, same caller obligation as
/// [`Self::attach_checkpoint`].
async fn collect_unattached_checkpoints(
&self,
sandbox_id: Uuid,
before: DateTime<Utc>,
limit: i64,
) -> Result<Vec<String>, SandboxCheckpointError>;
}
/// Type-keyed wrapper installed on the tool context by hosted presets.
///
/// Runtime hosts install this service in the generic tool-context extension bag.
#[derive(Clone)]
pub struct SandboxCheckpointStoreExt(pub std::sync::Arc<dyn SandboxCheckpointStore>);