refactor(proto): isolate gateway storage messages (#3169)

* refactor(proto): isolate gateway storage messages

Move persistence-only protobufs into a server-private versioned package, remove them from generated public SDKs, and gate durable/public schema compatibility with legacy database fixtures.

Closes #3053

Signed-off-by: Varsha Prasad Narsing <varshaprasad96@gmail.com>

* docs(gateway): sync protobuf schema inventory

Signed-off-by: Varsha Prasad Narsing <varshaprasad96@gmail.com>

---------

Signed-off-by: Varsha Prasad Narsing <varshaprasad96@gmail.com>
This commit is contained in:
Varsha
2026-09-10 16:36:07 +00:00
committed by GitHub
parent 67374efdf8
commit 0357daee31
20 changed files with 1821 additions and 2211 deletions
Generated
+2
View File
@@ -4500,6 +4500,7 @@ dependencies = [
"prost",
"prost-reflect",
"prost-types",
"protoc-bin-vendored",
"rand 0.9.4",
"rcgen",
"reqwest 0.12.28",
@@ -4523,6 +4524,7 @@ dependencies = [
"tokio-tungstenite 0.26.2",
"toml",
"tonic",
"tonic-prost-build",
"tower",
"tower-http",
"tracing",
+75
View File
@@ -308,6 +308,81 @@ The storage schema is intentionally narrow:
| `created_at_ms` and `updated_at_ms` | Gateway timestamps used for ordering and list output. |
| `labels` | JSON object carrying Kubernetes-style object labels for filtering and organization. |
### Protobuf API and storage boundaries
Public RPC contracts and durable protobuf formats have separate ownership. The
`openshell.v1.OpenShell` service currently has 74 RPCs. Their request and
response roots, streaming flags, and transitive message closure come from the
public descriptor set generated by `openshell-core`; a fingerprint test in
`openshell-server` requires this inventory to be reviewed whenever it changes.
Compute-driver, credential-driver, gateway-interceptor, and
supervisor-middleware services are compiled contracts for internal extension
boundaries, not public gateway RPCs. The current public inventory has 74
methods, 276 messages, and 12 enums
(`042034fe4d0000279ee4ed27e587ab8e530934b8d5c3aa36dc9769f81dfa6e51`).
Storage-only messages live in the private, versioned
`openshell.storage.v1` package under `crates/openshell-server/proto`. The server
generates these types separately, so the public descriptor set and the Rust,
Go, Python, and TypeScript client generation inputs do not advertise them.
| Storage classification | Protobuf messages | Durable use |
|---|---|---|
| Encoded storage roots | `StoredProviderCredentialRefreshState`, `StoredProviderProfile`, `PolicyRevisionPayload`, `DraftChunkPayload` | Complete protobuf payload stored in an object row or a scoped policy row. |
| Nested storage-only type | `StoredRefreshMaterialDeletion` | Repeated child records inside provider refresh state. |
| SQL materializations | `StoredPolicyRevision`, `StoredDraftChunk` | Server-only typed results assembled from indexed columns and decoded payloads; not public RPC messages. |
| Public messages used directly as encoded storage roots | `Sandbox`, `SandboxWorkloadTemplate`, `Provider`, `Workspace`, `WorkspaceMember`, `SshSession`, `ServiceEndpoint` | The generated public type is also the persisted payload. `SshSession` is not in the current public RPC message closure. |
| Embedded encoded root | `SandboxPolicy` | Stored in policy rows and inside the JSON settings envelope. |
The 12 encoded durable roots above have a closure of 81 messages and eight
enums (`920a5243dfb37ce709f0f562a47d17791a5ede90fd7f662ed01542abd60a0dfb`).
Its intersection with the public RPC closure contains 71 messages and eight
enums (`05add438ba041defc98d791038ae593d3f09352677cae43f2276d494205ce415`).
The descriptor-derived test owns these full inventories; the tables here record
the reviewed roots and classifications.
| Dual-purpose encoded root | Current decision |
|---|---|
| `Sandbox` | Defer a storage twin; govern its complete dependency closure as durable. |
| `SandboxWorkloadTemplate` | Defer a storage twin; govern its complete dependency closure as durable. |
| `Provider` | Defer a storage twin; govern its complete dependency closure as durable. |
| `Workspace` | Defer a storage twin; govern its complete dependency closure as durable. |
| `WorkspaceMember` | Defer a storage twin; govern its complete dependency closure as durable. |
| `SshSession` | Defer a storage twin; govern its complete dependency closure as durable. |
| `ServiceEndpoint` | Defer a storage twin; govern its complete dependency closure as durable. |
| `SandboxPolicy` | Defer a storage twin; govern its complete dependency closure as durable. |
The public/storage overlap is deliberate for the current format. Storage twins
for the public roots are deferred: introducing them would require a broad
conversion boundary, and Prost does not retain unknown fields through a
decode-and-reencode conversion. Each root therefore carries a reviewed decision
to remain dual-purpose, and its complete transitive dependency closure is also
a durable format. Important embedded dependencies include `ObjectMeta`,
`ProviderProfile`, `CredentialHandle`, `SandboxPolicy`, and
`NetworkPolicyRule`. Global and sandbox settings additionally store an encoded
`SandboxPolicy` inside their JSON envelope.
Public API compatibility and storage compatibility are reviewed independently:
- Public compatibility is evaluated from public service descriptors and SDK
generation inputs. Storage-only packages must never enter that closure.
- `openshell.storage.v1` is frozen. Its test fingerprint covers message names,
field numbers, cardinality, scalar wire types, referenced types, map-entry
shapes, and optional presence. Keep its decoder available and introduce a
new versioned package plus an explicit migration or fallback decoder for a
format change; never reuse removed tags or names.
- Checked-in synthetic byte fixtures were encoded with the former
`openshell.v1` declarations. Current storage types must continue to decode
them semantically, which proves the package move does not require a database
rewrite. Protobuf payload bytes do not encode a message's package name.
- A change to a dual-purpose public message or any transitive durable
dependency requires both public-wire review and storage-migration review.
Wire-incompatible changes require a migration or fallback decoder and a
fixture for the earlier format.
- Mixed-version writers are unsupported. An older Prost writer can discard
fields it does not know when it reads and rewrites a record, even when the
newer field is wire-compatible.
Common resources use generic helpers that derive `object_type`, `id`, `name`,
and labels from protobuf metadata traits before encoding the full message into
`payload`. Policy revisions and draft policy chunks use the same table but also
+1
View File
@@ -9,6 +9,7 @@
version: v2
modules:
- path: proto
- path: crates/openshell-server/proto
lint:
use:
- STANDARD
+1 -86
View File
@@ -7,8 +7,7 @@
use crate::proto::{
ObjectForTest, Provider, Sandbox, SandboxStatus, SandboxWorkloadTemplate, ServiceEndpoint,
SshSession, StoredProviderCredentialRefreshState, StoredProviderProfile, Workspace,
WorkspaceMember,
SshSession, Workspace, WorkspaceMember,
};
use std::collections::HashMap;
@@ -233,90 +232,6 @@ impl ObjectWorkspace for Provider {
}
}
// Implementations for StoredProviderProfile
impl ObjectId for StoredProviderProfile {
fn object_id(&self) -> &str {
self.metadata.as_ref().map_or("", |m| m.id.as_str())
}
}
impl ObjectName for StoredProviderProfile {
fn object_name(&self) -> &str {
self.metadata.as_ref().map_or("", |m| m.name.as_str())
}
}
impl ObjectLabels for StoredProviderProfile {
fn object_labels(&self) -> Option<HashMap<String, String>> {
self.metadata.as_ref().map(|m| m.labels.clone())
}
}
impl SetResourceVersion for StoredProviderProfile {
fn set_resource_version(&mut self, version: u64) {
if let Some(meta) = self.metadata.as_mut() {
meta.resource_version = version;
}
}
}
impl GetResourceVersion for StoredProviderProfile {
fn get_resource_version(&self) -> u64 {
self.metadata.as_ref().map_or(0, |m| m.resource_version)
}
}
impl ObjectWorkspace for StoredProviderProfile {
fn object_workspace(&self) -> &str {
self.metadata.as_ref().map_or("", |m| m.workspace.as_str())
}
fn requires_workspace() -> bool {
false
}
}
// Implementations for StoredProviderCredentialRefreshState
impl ObjectId for StoredProviderCredentialRefreshState {
fn object_id(&self) -> &str {
self.metadata.as_ref().map_or("", |m| m.id.as_str())
}
}
impl ObjectName for StoredProviderCredentialRefreshState {
fn object_name(&self) -> &str {
self.metadata.as_ref().map_or("", |m| m.name.as_str())
}
}
impl ObjectLabels for StoredProviderCredentialRefreshState {
fn object_labels(&self) -> Option<HashMap<String, String>> {
self.metadata.as_ref().map(|m| m.labels.clone())
}
}
impl SetResourceVersion for StoredProviderCredentialRefreshState {
fn set_resource_version(&mut self, version: u64) {
if let Some(meta) = self.metadata.as_mut() {
meta.resource_version = version;
}
}
}
impl GetResourceVersion for StoredProviderCredentialRefreshState {
fn get_resource_version(&self) -> u64 {
self.metadata.as_ref().map_or(0, |m| m.resource_version)
}
}
impl ObjectWorkspace for StoredProviderCredentialRefreshState {
fn object_workspace(&self) -> &str {
self.metadata.as_ref().map_or("", |m| m.workspace.as_str())
}
fn requires_workspace() -> bool {
true
}
}
// Implementations for SshSession
impl ObjectId for SshSession {
fn object_id(&self) -> &str {
@@ -377,10 +377,6 @@ mod tests {
("openshell.v1.RevokeSshSessionRequest", "token"),
("openshell.v1.TcpForwardInit", "authorization_token"),
("openshell.v1.SshSession", "token"),
(
"openshell.v1.StoredProviderCredentialRefreshState",
"material",
),
("openshell.v1.ConfigureProviderRefreshRequest", "material"),
(
"openshell.v1.GetSandboxProviderEnvironmentResponse",
+4
View File
@@ -120,6 +120,10 @@ telemetry = ["openshell-core/telemetry"]
bundled-z3 = ["openshell-prover/bundled-z3"]
test-support = []
[build-dependencies]
tonic-prost-build = { workspace = true }
protoc-bin-vendored = { workspace = true }
[dev-dependencies]
base64 = { workspace = true }
hyper-rustls = { version = "0.27", default-features = false, features = ["native-tokio", "http1", "tls12", "logging", "aws-lc-rs"] }
+50
View File
@@ -0,0 +1,50 @@
// SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
// SPDX-License-Identifier: Apache-2.0
use std::env;
use std::path::PathBuf;
fn main() -> Result<(), Box<dyn std::error::Error>> {
let manifest_dir = PathBuf::from(env::var("CARGO_MANIFEST_DIR")?);
let storage_proto_dir = manifest_dir.join("proto");
let public_proto_dir = manifest_dir.join("../../proto");
let storage_proto = storage_proto_dir.join("storage.proto");
println!("cargo:rerun-if-changed={}", storage_proto.display());
for imported_proto in [
"datamodel.proto",
"openshell.proto",
"options.proto",
"sandbox.proto",
] {
println!(
"cargo:rerun-if-changed={}",
public_proto_dir.join(imported_proto).display()
);
}
// SAFETY: Build scripts run in their own single-threaded process.
#[allow(unsafe_code)]
unsafe {
env::set_var("PROTOC", protoc_bin_vendored::protoc_bin_path()?);
env::set_var("PROTOC_INCLUDE", protoc_bin_vendored::include_path()?);
}
let descriptor_path = PathBuf::from(env::var("OUT_DIR")?).join("storage_descriptor.bin");
tonic_prost_build::configure()
.build_server(false)
.build_client(false)
.extern_path(".openshell.v1", "::openshell_core::proto")
.extern_path(
".openshell.datamodel.v1",
"::openshell_core::proto::datamodel::v1",
)
.extern_path(
".openshell.sandbox.v1",
"::openshell_core::proto::sandbox::v1",
)
.file_descriptor_set_path(&descriptor_path)
.compile_protos(&[storage_proto], &[storage_proto_dir, public_proto_dir])?;
Ok(())
}
+170
View File
@@ -0,0 +1,170 @@
// SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
// SPDX-License-Identifier: Apache-2.0
syntax = "proto3";
// Internal durable formats owned by the gateway. This package is not part of
// the public OpenShell API or any generated client SDK.
package openshell.storage.v1;
import "datamodel.proto";
import "openshell.proto";
import "options.proto";
import "sandbox.proto";
message StoredProviderCredentialRefreshState {
openshell.datamodel.v1.ObjectMeta metadata = 1;
string provider_id = 2;
string provider_name = 3;
string credential_key = 4;
openshell.v1.ProviderCredentialRefreshStrategy strategy = 5;
map<string, string> material = 6 [(openshell.options.v1.secret) = true];
// Material names classified as secret. Newly configured values live in the
// active credential driver and are absent from material. Legacy inline values
// are not automatically migrated before OpenShell 0.1.0.
repeated string secret_material_keys = 7;
int64 expires_at_ms = 8;
// int64 max parks the refresh until an explicit rotation or reconfiguration.
int64 next_refresh_at_ms = 9;
int64 last_refresh_at_ms = 10;
string status = 11;
string last_error = 12;
string token_url = 13;
repeated string scopes = 14;
int64 refresh_before_seconds = 15;
int64 max_lifetime_seconds = 16;
// Resolved mapping of strategy-defined output id -> concrete env key, pinned
// at configure time from the profile's additional_outputs. Read by minting,
// collision reservation, and env-key surfacing so later profile edits cannot
// silently redirect writes.
map<string, string> additional_output_keys = 17;
// Opaque gateway-owned authorization epoch for the configured refresh
// grant. Explicit refresh configuration creates a new epoch; automatic and
// manual token rotation preserve it. It is never derived from or exposed
// with refresh material.
string authorization_epoch = 18;
// Secret refresh material is stored through the gateway's active credential
// driver. The persisted refresh state keeps only opaque handles; resolved
// values exist in gateway memory for the duration of one mint operation.
map<string, openshell.datamodel.v1.CredentialHandle> secret_material_handles = 19;
// Handles replaced by reconfiguration or issuer-driven refresh-token
// rotation. This is a repeated entry rather than a material-keyed map so
// multiple superseded generations of the same material remain recoverable.
// Cleanup is retried by the refresh worker so a gateway crash or temporary
// credential-backend outage does not lose the deletion reference.
repeated StoredRefreshMaterialDeletion pending_secret_deletions = 20;
// Structured recovery details for the most recent refresh failure. These
// fields contain only gateway-owned codes and recognized bounded values.
openshell.v1.ProviderCredentialRefreshRecoveryAction recovery_action = 21;
string failure_code = 22;
string provider_error_subtype = 23;
int64 last_error_at_ms = 24;
}
message StoredRefreshMaterialDeletion {
// Original material name used to derive the credential driver's storage key.
string material_key = 1;
// Opaque handle for the superseded secret object.
openshell.datamodel.v1.CredentialHandle handle = 2;
}
// Stored custom provider profile object.
message StoredProviderProfile {
openshell.datamodel.v1.ObjectMeta metadata = 1;
openshell.v1.ProviderProfile profile = 2;
}
// Stored payload for a policy revision row in the generic objects table.
message PolicyRevisionPayload {
// Serialized policy contents.
openshell.sandbox.v1.SandboxPolicy policy = 1;
// Deterministic hash of the policy payload.
string hash = 2;
// Load error reported by the sandbox, if any.
string load_error = 3;
// When the policy version was reported as loaded (ms since epoch). 0 if unset.
int64 loaded_at_ms = 4;
// Immutable provenance supplied when this revision was created.
map<string, string> provenance = 5;
}
// Stored payload for a draft policy chunk row in the generic objects table.
message DraftChunkPayload {
// Proposed network_policies map key.
string rule_name = 1;
// Proposed network policy rule.
openshell.sandbox.v1.NetworkPolicyRule proposed_rule = 2;
// Human-readable explanation of why this rule is proposed.
string rationale = 3;
// Security concerns flagged by analysis (empty if none).
string security_notes = 4;
// Analysis confidence (0.0-1.0). 0 for mechanistic mode.
float confidence = 5;
// When the user approved/rejected (ms since epoch). 0 if undecided.
int64 decided_at_ms = 6;
// Denormalized endpoint host for dedup and display.
string host = 7;
// Denormalized endpoint port for dedup and display.
int32 port = 8;
// Binary path that triggered the denial.
string binary = 9;
// Current draft version for the owning sandbox.
int64 draft_version = 10;
// Gateway prover verdict for this chunk; empty until prover runs.
// Mirrors PolicyChunk.validation_result.
string validation_result = 11;
// Operator-supplied free-form rejection text; empty for non-rejected
// chunks. Mirrors PolicyChunk.rejection_reason.
string rejection_reason = 12;
string application_error = 13;
string review_token = 14;
string current_effective_policy_hash = 15;
string candidate_effective_policy_hash = 16;
openshell.sandbox.v1.SandboxPolicy current_effective_policy = 17;
openshell.sandbox.v1.SandboxPolicy candidate_effective_policy = 18;
}
// Internal stored policy revision row materialized from the generic objects table.
message StoredPolicyRevision {
string id = 1;
string sandbox_id = 2;
int64 version = 3;
bytes policy_payload = 4;
string policy_hash = 5;
string status = 6;
optional string load_error = 7;
int64 created_at_ms = 8;
optional int64 loaded_at_ms = 9;
map<string, string> provenance = 10;
}
// Internal stored draft chunk row materialized from the generic objects table.
message StoredDraftChunk {
string id = 1;
string sandbox_id = 2;
int64 draft_version = 3;
string status = 4;
string rule_name = 5;
bytes proposed_rule = 6;
string rationale = 7;
string security_notes = 8;
double confidence = 9;
int64 created_at_ms = 10;
optional int64 decided_at_ms = 11;
string host = 12;
int32 port = 13;
string binary = 14;
int32 hit_count = 15;
int64 first_seen_ms = 16;
int64 last_seen_ms = 17;
// Gateway prover verdict; empty until the prover runs. See PolicyChunk.
string validation_result = 18;
// Operator-supplied free-form rejection text. See PolicyChunk.
string rejection_reason = 19;
string application_error = 20;
string review_token = 21;
string current_effective_policy_hash = 22;
string candidate_effective_policy_hash = 23;
openshell.sandbox.v1.SandboxPolicy current_effective_policy = 24;
openshell.sandbox.v1.SandboxPolicy candidate_effective_policy = 25;
}
+2 -1
View File
@@ -30,7 +30,7 @@ use openshell_core::proto::credentials::v1::{
GetCredentialDriverCapabilitiesResponse, ResolveCredentialRequest, ResolveCredentialsRequest,
ResolvedCredential, StoreCredentialRequest, credential_driver_client::CredentialDriverClient,
};
use openshell_core::proto::{CredentialHandle, Provider, StoredRefreshMaterialDeletion};
use openshell_core::proto::{CredentialHandle, Provider};
use openshell_core::{Config, Error, Result as CoreResult};
use openshell_driver_db_credstore::{
CredentialObjectWrite, DbCredstoreCredentialDriver, DbCredstoreObjectStore,
@@ -51,6 +51,7 @@ use tower::service_fn;
use tracing::warn;
use crate::persistence::{PersistenceError, Store, WriteCondition};
use crate::storage_proto::StoredRefreshMaterialDeletion;
const DEFAULT_CREDENTIAL_DRIVER_STARTUP_TIMEOUT_SECS: u64 = 10;
const DEFAULT_CREDENTIAL_DRIVER_RPC_TIMEOUT_SECS: u64 = 30;
+39 -51
View File
@@ -22,6 +22,9 @@ use crate::policy_store::{AtomicPolicyRevisionWrite, PolicyStoreExt};
use crate::provider_profile_sources::EffectiveProviderProfileCatalog;
#[cfg(test)]
use crate::provider_profile_sources::ProviderProfileSources;
use crate::storage_proto::StoredProviderCredentialRefreshState;
#[cfg(test)]
use crate::storage_proto::StoredProviderProfile;
use openshell_core::net::{is_always_blocked_ip, is_internal_ip};
use openshell_core::proto::policy_merge_operation;
use openshell_core::proto::setting_value;
@@ -2747,7 +2750,7 @@ fn compute_provider_env_revision_from_records_and_policy_bindings(
}
fn hash_provider_refresh_states(
states: &[openshell_core::proto::StoredProviderCredentialRefreshState],
states: &[StoredProviderCredentialRefreshState],
hasher: &mut Sha256,
) -> Result<(), Status> {
let mut states = states.iter().collect::<Vec<_>>();
@@ -9425,7 +9428,7 @@ mod tests {
.await
.unwrap();
store
.put_message(&openshell_core::proto::StoredProviderProfile {
.put_message(&StoredProviderProfile {
metadata: Some(openshell_core::proto::datamodel::v1::ObjectMeta {
id: "profile-generic".to_string(),
name: "generic".to_string(),
@@ -9476,7 +9479,7 @@ mod tests {
.await
.unwrap();
store
.put_message(&openshell_core::proto::StoredProviderProfile {
.put_message(&StoredProviderProfile {
metadata: Some(openshell_core::proto::datamodel::v1::ObjectMeta {
id: "profile-gh".to_string(),
name: "gh".to_string(),
@@ -9519,7 +9522,7 @@ mod tests {
.await
.unwrap();
store
.put_message(&openshell_core::proto::StoredProviderProfile {
.put_message(&StoredProviderProfile {
metadata: Some(openshell_core::proto::datamodel::v1::ObjectMeta {
id: "profile-custom-api".to_string(),
name: "custom-api".to_string(),
@@ -9591,7 +9594,7 @@ mod tests {
.await
.unwrap();
store
.put_message(&openshell_core::proto::StoredProviderProfile {
.put_message(&StoredProviderProfile {
metadata: Some(openshell_core::proto::datamodel::v1::ObjectMeta {
id: "profile-custom-api".to_string(),
name: "custom-api".to_string(),
@@ -9760,29 +9763,28 @@ mod tests {
async fn provider_policy_layers_respect_profile_workspace_scope() {
let store = test_store().await;
let make_stored_profile =
|id: &str, workspace: &str, host: &str| openshell_core::proto::StoredProviderProfile {
metadata: Some(openshell_core::proto::datamodel::v1::ObjectMeta {
id: format!("profile-{id}-{workspace}"),
name: id.to_string(),
created_at_ms: 1_000_000,
labels: HashMap::new(),
resource_version: 0,
annotations: HashMap::new(),
workspace: workspace.to_string(),
deletion_timestamp_ms: 0,
}),
profile: Some(openshell_core::proto::ProviderProfile {
id: id.to_string(),
display_name: format!("{host} profile"),
endpoints: vec![NetworkEndpoint {
host: host.to_string(),
port: 443,
..Default::default()
}],
let make_stored_profile = |id: &str, workspace: &str, host: &str| StoredProviderProfile {
metadata: Some(openshell_core::proto::datamodel::v1::ObjectMeta {
id: format!("profile-{id}-{workspace}"),
name: id.to_string(),
created_at_ms: 1_000_000,
labels: HashMap::new(),
resource_version: 0,
annotations: HashMap::new(),
workspace: workspace.to_string(),
deletion_timestamp_ms: 0,
}),
profile: Some(openshell_core::proto::ProviderProfile {
id: id.to_string(),
display_name: format!("{host} profile"),
endpoints: vec![NetworkEndpoint {
host: host.to_string(),
port: 443,
..Default::default()
}),
};
}],
..Default::default()
}),
};
store
.put_message(&make_stored_profile(
@@ -9873,9 +9875,7 @@ mod tests {
#[tokio::test]
async fn sandbox_config_materializes_default_mcp_version_after_provider_composition() {
use openshell_core::proto::{
ProviderProfile, ProviderProfileCategory, StoredProviderProfile,
};
use openshell_core::proto::{ProviderProfile, ProviderProfileCategory};
let state = test_server_state().await;
state
@@ -10096,9 +10096,7 @@ mod tests {
#[tokio::test]
async fn update_config_gates_uninspected_endpointless_credential_binding() {
use openshell_core::proto::{
ProviderProfile, ProviderProfileCategory, StoredProviderProfile,
};
use openshell_core::proto::{ProviderProfile, ProviderProfileCategory};
let state = test_server_state().await;
state
@@ -10478,9 +10476,7 @@ mod tests {
#[tokio::test]
async fn provider_attachment_preflight_rejects_composed_ambiguity() {
use openshell_core::proto::{
ProviderProfile, ProviderProfileCategory, StoredProviderProfile,
};
use openshell_core::proto::{ProviderProfile, ProviderProfileCategory};
let state = test_server_state().await;
state
@@ -10551,7 +10547,7 @@ mod tests {
async fn sandbox_config_rejects_invalid_provider_composed_policy() {
use openshell_core::proto::{
MiddlewareEndpointSelector, NetworkMiddlewareConfig, ProviderProfile,
ProviderProfileCategory, StoredProviderProfile,
ProviderProfileCategory,
};
let state = test_server_state().await;
@@ -10712,7 +10708,7 @@ mod tests {
use crate::grpc::provider::handle_update_provider_profiles;
use openshell_core::proto::{
ProviderProfile, ProviderProfileCategory, ProviderProfileImportItem,
StoredProviderProfile, UpdateProviderProfilesRequest,
UpdateProviderProfilesRequest,
};
fn stored_profile(host: &str) -> StoredProviderProfile {
@@ -11043,7 +11039,6 @@ mod tests {
use openshell_core::proto::{
GetSandboxConfigRequest, GetSandboxProviderEnvironmentRequest,
NetworkCredentialBinding, ProviderProfile, ProviderProfileCategory,
StoredProviderProfile,
};
let state = test_server_state().await;
@@ -11265,7 +11260,7 @@ mod tests {
async fn invalid_static_binding_does_not_suppress_valid_dynamic_credentials() {
use openshell_core::proto::{
GetSandboxProviderEnvironmentRequest, ProviderCredentialTokenGrant, ProviderProfile,
ProviderProfileCategory, ProviderProfileCredential, StoredProviderProfile,
ProviderProfileCategory, ProviderProfileCredential,
};
let state = test_server_state().await;
@@ -11358,7 +11353,6 @@ mod tests {
GetSandboxProviderEnvironmentRequest, ProviderCredentialTokenGrant,
ProviderCredentialTokenGrantSubjectToken, ProviderCredentialTokenGrantType,
ProviderProfile, ProviderProfileCategory, ProviderProfileCredential,
StoredProviderProfile,
};
let state = test_server_state().await;
@@ -11563,7 +11557,7 @@ mod tests {
async fn provider_environment_revision_and_payload_share_immutable_record_snapshot() {
use openshell_core::proto::{
ProviderCredentialTokenGrant, ProviderProfile, ProviderProfileCategory,
ProviderProfileCredential, StoredProviderProfile,
ProviderProfileCredential,
};
fn dynamic_profile(
@@ -11751,8 +11745,7 @@ mod tests {
use crate::grpc::provider::handle_update_provider_profiles;
use openshell_core::proto::{
ProviderCredentialTokenGrant, ProviderProfile, ProviderProfileCategory,
ProviderProfileCredential, ProviderProfileImportItem, StoredProviderProfile,
UpdateProviderProfilesRequest,
ProviderProfileCredential, ProviderProfileImportItem, UpdateProviderProfilesRequest,
};
use std::time::Duration;
@@ -11868,9 +11861,7 @@ mod tests {
#[tokio::test]
async fn platform_profile_narrowing_changes_platform_provider_revision_when_shadowed() {
use crate::persistence::WriteCondition;
use openshell_core::proto::{
ProviderProfile, ProviderProfileCategory, StoredProviderProfile,
};
use openshell_core::proto::{ProviderProfile, ProviderProfileCategory};
fn stored_profile(workspace: &str, path: &str) -> StoredProviderProfile {
StoredProviderProfile {
@@ -15625,7 +15616,6 @@ mod tests {
use openshell_core::proto::{
FilesystemPolicy, L7Allow, L7DenyRule, L7Rule, NetworkBinary, NetworkEndpoint,
ProviderProfile, ProviderProfileCategory, SandboxPhase, SandboxPolicy, SandboxSpec,
StoredProviderProfile,
};
let state = test_server_state().await;
@@ -17812,9 +17802,7 @@ mod tests {
}
async fn install_ambiguous_provider_binding(state: &Arc<ServerState>, suffix: &str) {
use openshell_core::proto::{
ProviderProfile, ProviderProfileCategory, StoredProviderProfile,
};
use openshell_core::proto::{ProviderProfile, ProviderProfileCategory};
let profile_name = format!("ambiguous-{suffix}");
let provider_name = format!("provider-{suffix}");
+4 -5
View File
@@ -14,12 +14,13 @@ use crate::provider_profile_sources::{
EffectiveProviderProfileCatalog, ProviderProfileSources, profile_response_payload,
profile_storage_payload, stored_profile_resource_version,
};
use crate::storage_proto::{StoredProviderCredentialRefreshState, StoredProviderProfile};
use openshell_core::metadata::ObjectWorkspace;
use openshell_core::proto::{
CredentialHandle, Provider, ProviderCredentialRefreshStrategy,
ProviderCredentialTokenGrantAudienceOverride, ProviderCredentialTokenGrantType,
ProviderProfile, ProviderProfileCredential, Sandbox, StaticCredentialBinding,
StaticCredentialEndpointBinding, StoredProviderCredentialRefreshState,
StaticCredentialEndpointBinding,
};
use openshell_core::telemetry::{
LifecycleOperation, ProviderProfile as TelemetryProviderProfile, TelemetryOutcome,
@@ -2334,8 +2335,7 @@ use openshell_core::proto::{
ListProviderProfilesResponse, ListProvidersRequest, ListProvidersResponse,
ProviderProfileDiagnostic, ProviderProfileImportItem, ProviderProfileResponse,
ProviderResponse, RotateProviderCredentialRequest, RotateProviderCredentialResponse,
StoredProviderProfile, UpdateProviderProfilesRequest, UpdateProviderProfilesResponse,
UpdateProviderRequest,
UpdateProviderProfilesRequest, UpdateProviderProfilesResponse, UpdateProviderRequest,
};
use openshell_core::spiffe::{
JwtSvidParseError, SpiffeJwtClaims, parse_unverified_jwt_svid_claims,
@@ -4937,8 +4937,7 @@ mod tests {
ProviderCredentialTokenGrantAudienceOverride, ProviderCredentialTokenGrantSubjectToken,
ProviderCredentialTokenGrantType, ProviderProfile, ProviderProfileCategory,
ProviderProfileCredential, ProviderProfileImportItem, RotateProviderCredentialRequest,
Sandbox, SandboxPolicy, SandboxSpec, StoredProviderProfile, UpdateProviderProfilesRequest,
UpdateProviderRequest,
Sandbox, SandboxPolicy, SandboxSpec, UpdateProviderProfilesRequest, UpdateProviderRequest,
};
use openshell_core::{ObjectId, ObjectName};
use tonic::{Code, Request};
@@ -15,8 +15,7 @@ use openshell_core::proto::{
GetWorkspaceResponse, ListWorkspaceMembersRequest, ListWorkspaceMembersResponse,
ListWorkspacesRequest, ListWorkspacesResponse, Provider, RemoveWorkspaceMemberRequest,
RemoveWorkspaceMemberResponse, Sandbox, SandboxWorkloadTemplate, ServiceEndpoint, SshSession,
StoredProviderCredentialRefreshState, StoredProviderProfile, Workspace, WorkspaceMember,
WorkspaceRole,
Workspace, WorkspaceMember, WorkspaceRole,
};
use prost::Message;
use tonic::{Request, Response, Status};
@@ -28,6 +27,7 @@ use crate::persistence::{
DRAFT_CHUNK_OBJECT_TYPE, ObjectLabels, ObjectType, POLICY_OBJECT_TYPE, WriteCondition,
current_time_ms,
};
use crate::storage_proto::{StoredProviderCredentialRefreshState, StoredProviderProfile};
use std::collections::HashMap;
use super::{MAX_PAGE_SIZE, clamp_limit};
+1
View File
@@ -35,6 +35,7 @@ mod sandbox_index;
mod sandbox_watch;
mod service_routing;
mod ssh_sessions;
mod storage_proto;
pub mod supervisor_session;
mod telemetry;
#[cfg(any(test, feature = "test-support"))]
@@ -6,7 +6,7 @@
mod postgres;
mod sqlite;
pub use openshell_core::proto::{
pub use crate::storage_proto::{
StoredDraftChunk as DraftChunkRecord, StoredPolicyRevision as PolicyRecord,
};
@@ -162,8 +162,9 @@ pub trait ObjectType {
fn object_type() -> &'static str;
}
// Import object metadata accessor traits from openshell-core
// (implementations for all proto types are in openshell-core::metadata)
// Import object metadata accessor traits from openshell-core. Implementations
// for public resource types live there; private storage types implement them
// in crate::storage_proto.
pub use openshell_core::{
GetResourceVersion, ObjectId, ObjectLabels, ObjectName, ObjectWorkspace, SetResourceVersion,
};
+2 -4
View File
@@ -4,10 +4,8 @@
use crate::persistence::{
DraftChunkRecord, PersistenceError, PersistenceResult, PolicyRecord, SetResourceVersion, Store,
};
use openshell_core::proto::{
DraftChunkPayload, NetworkPolicyRule, PolicyRevisionPayload, Sandbox,
SandboxPolicy as ProtoSandboxPolicy,
};
use crate::storage_proto::{DraftChunkPayload, PolicyRevisionPayload};
use openshell_core::proto::{NetworkPolicyRule, Sandbox, SandboxPolicy as ProtoSandboxPolicy};
use prost::Message;
use std::collections::HashMap;
@@ -9,7 +9,7 @@ use std::sync::Arc;
use async_trait::async_trait;
use openshell_core::GatewayProviderProfileSourceConfig;
use openshell_core::mcp::normalize_provider_profile_mcp_fields;
use openshell_core::proto::{ProviderProfile, StoredProviderProfile};
use openshell_core::proto::ProviderProfile;
use openshell_gateway_interceptors::{
GatewayInterceptorProfileSource, GatewayInterceptorRuntime,
ProviderProfileSourceSnapshot as InterceptorProfileSnapshot,
@@ -24,6 +24,7 @@ use tonic::Status;
use tracing::debug;
use crate::persistence::{ObjectType, Store};
use crate::storage_proto::StoredProviderProfile;
const BUILTIN_SOURCE_ID: &str = "builtin";
const USER_SOURCE_ID: &str = "user";
@@ -11,7 +11,6 @@ use openshell_core::ObjectWorkspace;
use openshell_core::proto::{
CredentialHandle, Provider, ProviderCredentialRefreshRecoveryAction,
ProviderCredentialRefreshStatus, ProviderCredentialRefreshStrategy,
StoredProviderCredentialRefreshState, StoredRefreshMaterialDeletion,
};
use openshell_core::{ObjectId, ObjectName, SetResourceVersion};
use prost::Message;
@@ -21,6 +20,8 @@ use std::time::Duration;
use tonic::{Code, Status};
use tracing::{info, warn};
use crate::storage_proto::{StoredProviderCredentialRefreshState, StoredRefreshMaterialDeletion};
const DEFAULT_REFRESH_BEFORE_SECONDS: i64 = 300;
const DEFAULT_MAX_LIFETIME_SECONDS: i64 = 3600;
const REFRESH_ERROR_RETRY_SECONDS: i64 = 60;
@@ -2017,12 +2018,12 @@ mod tests {
};
use crate::credentials::CredentialRuntime;
use crate::persistence::{current_time_ms, test_store};
use crate::storage_proto::StoredProviderCredentialRefreshState;
use openshell_core::Config;
use openshell_core::proto::datamodel::v1::ObjectMeta;
use openshell_core::proto::{
CredentialHandle, Provider, ProviderCredentialRefreshRecoveryAction,
ProviderCredentialRefreshStrategy, Sandbox, SandboxSpec,
StoredProviderCredentialRefreshState,
};
use openshell_core::{ObjectId, ObjectName, ObjectWorkspace};
use std::collections::HashMap;
@@ -0,0 +1,721 @@
// SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
// SPDX-License-Identifier: Apache-2.0
//! Gateway-private, versioned protobuf formats for durable storage.
#![allow(
clippy::all,
clippy::pedantic,
clippy::nursery,
dead_code,
unused_imports,
unused_qualifications,
rust_2018_idioms
)]
include!(concat!(env!("OUT_DIR"), "/openshell.storage.v1.rs"));
#[cfg(test)]
const STORAGE_FILE_DESCRIPTOR_SET: &[u8] =
include_bytes!(concat!(env!("OUT_DIR"), "/storage_descriptor.bin"));
use openshell_core::{
GetResourceVersion, ObjectId, ObjectLabels, ObjectName, ObjectWorkspace, SetResourceVersion,
};
use std::collections::HashMap;
impl ObjectId for StoredProviderProfile {
fn object_id(&self) -> &str {
self.metadata.as_ref().map_or("", |m| m.id.as_str())
}
}
impl ObjectName for StoredProviderProfile {
fn object_name(&self) -> &str {
self.metadata.as_ref().map_or("", |m| m.name.as_str())
}
}
impl ObjectLabels for StoredProviderProfile {
fn object_labels(&self) -> Option<HashMap<String, String>> {
self.metadata.as_ref().map(|m| m.labels.clone())
}
}
impl SetResourceVersion for StoredProviderProfile {
fn set_resource_version(&mut self, version: u64) {
if let Some(meta) = self.metadata.as_mut() {
meta.resource_version = version;
}
}
}
impl GetResourceVersion for StoredProviderProfile {
fn get_resource_version(&self) -> u64 {
self.metadata.as_ref().map_or(0, |m| m.resource_version)
}
}
impl ObjectWorkspace for StoredProviderProfile {
fn object_workspace(&self) -> &str {
self.metadata.as_ref().map_or("", |m| m.workspace.as_str())
}
fn requires_workspace() -> bool {
false
}
}
impl ObjectId for StoredProviderCredentialRefreshState {
fn object_id(&self) -> &str {
self.metadata.as_ref().map_or("", |m| m.id.as_str())
}
}
impl ObjectName for StoredProviderCredentialRefreshState {
fn object_name(&self) -> &str {
self.metadata.as_ref().map_or("", |m| m.name.as_str())
}
}
impl ObjectLabels for StoredProviderCredentialRefreshState {
fn object_labels(&self) -> Option<HashMap<String, String>> {
self.metadata.as_ref().map(|m| m.labels.clone())
}
}
impl SetResourceVersion for StoredProviderCredentialRefreshState {
fn set_resource_version(&mut self, version: u64) {
if let Some(meta) = self.metadata.as_mut() {
meta.resource_version = version;
}
}
}
impl GetResourceVersion for StoredProviderCredentialRefreshState {
fn get_resource_version(&self) -> u64 {
self.metadata.as_ref().map_or(0, |m| m.resource_version)
}
}
impl ObjectWorkspace for StoredProviderCredentialRefreshState {
fn object_workspace(&self) -> &str {
self.metadata.as_ref().map_or("", |m| m.workspace.as_str())
}
fn requires_workspace() -> bool {
true
}
}
#[cfg(test)]
mod tests {
use super::*;
use prost::Message;
use prost_types::{DescriptorProto, EnumDescriptorProto, FileDescriptorSet};
use sha2::{Digest, Sha256};
use std::collections::{BTreeMap, BTreeSet, VecDeque};
const STORAGE_V1_SCHEMA_SHA256: &str =
"79c72615d957fc0653c672f61998bf7d8d21b757bc05d07b3fff92bd70fc8f52";
const PUBLIC_RPC_SCHEMA_SHA256: &str =
"042034fe4d0000279ee4ed27e587ab8e530934b8d5c3aa36dc9769f81dfa6e51";
const DURABLE_SCHEMA_SHA256: &str =
"920a5243dfb37ce709f0f562a47d17791a5ede90fd7f662ed01542abd60a0dfb";
const PUBLIC_DURABLE_OVERLAP_SHA256: &str =
"05add438ba041defc98d791038ae593d3f09352677cae43f2276d494205ce415";
// Synthetic payloads generated with the public declarations at v0.0.116,
// before their relocation into openshell.storage.v1. Values are deliberately
// non-secret and the ordinary protobuf bytes contain no package names.
const V0_0_116_REFRESH_STATE: &str = "0a230a096c65676163792d6964120b6c65676163792d6e616d6528073a0764656661756c74120b70726f76696465722d69641a0870726f76696465722205544f4b454e32160a09636c69656e745f6964120973796e7468657469633a0d726566726573685f746f6b656e406448c8015a06616374697665720773636f70652d619a011f0a0d726566726573685f746f6b656e120e0a047465737412066f7061717565a201150a036f6c64120e0a047465737412066f7061717565";
const V0_0_116_DELETION: &str = "0a036f6c64120e0a047465737412066f7061717565";
const V0_0_116_PROFILE: &str = "0a230a096c65676163792d6964120b6c65676163792d6e616d6528073a0764656661756c7412110a0770726f66696c6512064c6567616379";
const V0_0_116_POLICY_PAYLOAD: &str =
"0a0012067368613235361a046e6f6e6520ac022a110a06736f75726365120766697874757265";
const V0_0_116_DRAFT_PAYLOAD: &str =
"0a0472756c651a07666978747572652d0000403f3a0b6578616d706c652e636f6d40bb035002";
const V0_0_116_POLICY_RECORD: &str = "0a09706f6c6963792d6964120a73616e64626f782d6964180222030102032a0673686132353632066c6f616465643a046e6f6e6540fa0148ac0252110a06736f75726365120766697874757265";
const V0_0_116_DRAFT_RECORD: &str = "0a086368756e6b2d6964120a73616e64626f782d69641802220770656e64696e672a0472756c65320204053a076669787475726549000000000000e83f50de02589003620b6578616d706c652e636f6d68bb037801";
const STORAGE_MESSAGE_NAMES: [&str; 7] = [
"DraftChunkPayload",
"PolicyRevisionPayload",
"StoredDraftChunk",
"StoredPolicyRevision",
"StoredProviderCredentialRefreshState",
"StoredProviderProfile",
"StoredRefreshMaterialDeletion",
];
const DURABLE_ROOTS: [&str; 12] = [
".openshell.datamodel.v1.Provider",
".openshell.datamodel.v1.Workspace",
".openshell.sandbox.v1.SandboxPolicy",
".openshell.storage.v1.DraftChunkPayload",
".openshell.storage.v1.PolicyRevisionPayload",
".openshell.storage.v1.StoredProviderCredentialRefreshState",
".openshell.storage.v1.StoredProviderProfile",
".openshell.v1.Sandbox",
".openshell.v1.SandboxWorkloadTemplate",
".openshell.v1.ServiceEndpoint",
".openshell.v1.SshSession",
".openshell.v1.WorkspaceMember",
];
#[derive(Default)]
struct SchemaIndex<'a> {
messages: BTreeMap<String, &'a DescriptorProto>,
enums: BTreeMap<String, &'a EnumDescriptorProto>,
}
#[derive(Debug)]
struct SchemaClosure {
messages: BTreeSet<String>,
enums: BTreeSet<String>,
}
fn qualified_name(prefix: &str, name: &str) -> String {
if prefix.is_empty() {
format!(".{name}")
} else {
format!("{prefix}.{name}")
}
}
fn index_message<'a>(index: &mut SchemaIndex<'a>, prefix: &str, message: &'a DescriptorProto) {
let name = qualified_name(prefix, message.name.as_deref().expect("message name"));
index.messages.insert(name.clone(), message);
for nested in &message.nested_type {
index_message(index, &name, nested);
}
for nested_enum in &message.enum_type {
index.enums.insert(
qualified_name(&name, nested_enum.name.as_deref().expect("enum name")),
nested_enum,
);
}
}
fn index_descriptor<'a>(index: &mut SchemaIndex<'a>, descriptor: &'a FileDescriptorSet) {
for file in &descriptor.file {
let package = file.package.as_deref().unwrap_or_default();
let prefix = if package.is_empty() {
String::new()
} else {
format!(".{package}")
};
for message in &file.message_type {
index_message(index, &prefix, message);
}
for r#enum in &file.enum_type {
index.enums.insert(
qualified_name(&prefix, r#enum.name.as_deref().expect("enum name")),
r#enum,
);
}
}
}
fn schema_closure(
index: &SchemaIndex<'_>,
roots: impl IntoIterator<Item = String>,
) -> SchemaClosure {
let mut messages = BTreeSet::new();
let mut enums = BTreeSet::new();
let mut pending = roots.into_iter().collect::<VecDeque<_>>();
while let Some(name) = pending.pop_front() {
if let Some(message) = index.messages.get(&name) {
if !messages.insert(name) {
continue;
}
for field in &message.field {
let Some(type_name) = field.type_name.as_ref() else {
continue;
};
if index.messages.contains_key(type_name) {
pending.push_back(type_name.clone());
} else if index.enums.contains_key(type_name) {
enums.insert(type_name.clone());
}
}
} else if index.enums.contains_key(&name) {
enums.insert(name);
} else {
panic!("schema root or dependency {name} is missing from descriptors");
}
}
SchemaClosure { messages, enums }
}
fn append_one_message_schema(output: &mut String, full_name: &str, message: &DescriptorProto) {
output.push_str(&format!(
"message|{}|{}|{:?}|{:?}|{:?}\n",
full_name,
message
.options
.as_ref()
.map(|options| hex::encode(options.encode_to_vec()))
.unwrap_or_default(),
message.reserved_range,
message.reserved_name,
message.oneof_decl,
));
for field in &message.field {
output.push_str(&format!(
"field|{}|{}|{}|{}|{}|{}|{}|{}|{}\n",
full_name,
field.number.unwrap_or_default(),
field.name.as_deref().unwrap_or_default(),
field.label.unwrap_or_default(),
field.r#type.unwrap_or_default(),
field.type_name.as_deref().unwrap_or_default(),
field.oneof_index.unwrap_or(-1),
field.proto3_optional.unwrap_or(false),
field
.options
.as_ref()
.map(|options| hex::encode(options.encode_to_vec()))
.unwrap_or_default(),
));
}
}
fn schema_fingerprint(index: &SchemaIndex<'_>, closure: &SchemaClosure) -> String {
let mut schema = String::new();
for name in &closure.messages {
append_one_message_schema(&mut schema, name, index.messages[name]);
}
for name in &closure.enums {
let r#enum = index.enums[name];
schema.push_str(&format!(
"enum|{}|{}|{:?}|{:?}\n",
name,
r#enum
.options
.as_ref()
.map(|options| hex::encode(options.encode_to_vec()))
.unwrap_or_default(),
r#enum.reserved_range,
r#enum.reserved_name,
));
for value in &r#enum.value {
schema.push_str(&format!(
"enum-value|{}|{}|{}|{}\n",
name,
value.number.unwrap_or_default(),
value.name.as_deref().unwrap_or_default(),
value
.options
.as_ref()
.map(|options| hex::encode(options.encode_to_vec()))
.unwrap_or_default(),
));
}
}
format!("{:x}", Sha256::digest(schema.as_bytes()))
}
fn append_message_schema(output: &mut String, prefix: &str, message: &DescriptorProto) {
let name = message.name.as_deref().expect("message name");
let full_name = if prefix.is_empty() {
name.to_string()
} else {
format!("{prefix}.{name}")
};
output.push_str(&format!(
"message|{}|{}|{:?}|{:?}\n",
full_name,
message
.options
.as_ref()
.map(|options| hex::encode(options.encode_to_vec()))
.unwrap_or_default(),
message.reserved_range,
message.reserved_name,
));
for field in &message.field {
output.push_str(&format!(
"field|{}|{}|{}|{}|{}|{}|{}|{}|{}\n",
full_name,
field.number.unwrap_or_default(),
field.name.as_deref().unwrap_or_default(),
field.label.unwrap_or_default(),
field.r#type.unwrap_or_default(),
field.type_name.as_deref().unwrap_or_default(),
field.oneof_index.unwrap_or(-1),
field.proto3_optional.unwrap_or(false),
field
.options
.as_ref()
.map(|options| hex::encode(options.encode_to_vec()))
.unwrap_or_default(),
));
}
for nested in &message.nested_type {
append_message_schema(output, &full_name, nested);
}
}
fn storage_schema() -> (Vec<String>, String) {
let descriptor = FileDescriptorSet::decode(STORAGE_FILE_DESCRIPTOR_SET)
.expect("storage descriptor set must decode");
let file = descriptor
.file
.iter()
.find(|file| file.package.as_deref() == Some("openshell.storage.v1"))
.expect("storage descriptor must contain openshell.storage.v1");
let mut names = file
.message_type
.iter()
.map(|message| message.name.clone().expect("message name"))
.collect::<Vec<_>>();
names.sort();
let mut messages = file.message_type.iter().collect::<Vec<_>>();
messages.sort_by_key(|message| message.name.as_deref().unwrap_or_default());
let mut schema = String::new();
for message in messages {
append_message_schema(&mut schema, "", message);
}
(names, schema)
}
#[test]
fn storage_v1_schema_is_frozen() {
let (names, schema) = storage_schema();
assert_eq!(names, STORAGE_MESSAGE_NAMES);
let actual = format!("{:x}", Sha256::digest(schema.as_bytes()));
assert_eq!(
actual, STORAGE_V1_SCHEMA_SHA256,
"openshell.storage.v1 changed; keep v1 decoders intact and introduce a versioned migration path before updating this reviewed fingerprint"
);
}
#[test]
fn storage_types_are_absent_from_public_descriptor() {
let public = FileDescriptorSet::decode(openshell_core::FILE_DESCRIPTOR_SET)
.expect("public descriptor set must decode");
for file in public.file {
assert_ne!(file.package.as_deref(), Some("openshell.storage.v1"));
for message in file.message_type {
let name = message.name.expect("message name");
assert!(
!STORAGE_MESSAGE_NAMES.contains(&name.as_str()),
"storage-only message {name} leaked into the public descriptor"
);
}
}
}
#[test]
fn public_and_durable_schema_inventories_are_complete() {
let public = FileDescriptorSet::decode(openshell_core::FILE_DESCRIPTOR_SET)
.expect("public descriptor set must decode");
let storage = FileDescriptorSet::decode(STORAGE_FILE_DESCRIPTOR_SET)
.expect("storage descriptor set must decode");
let mut index = SchemaIndex::default();
index_descriptor(&mut index, &storage);
index_descriptor(&mut index, &public);
let mut methods = Vec::new();
let mut public_roots = BTreeSet::new();
let mut compiled_method_count = 0;
for file in &public.file {
let package = file.package.as_deref().unwrap_or_default();
for service in &file.service {
let service_name = service.name.as_deref().expect("service name");
compiled_method_count += service.method.len();
if !matches!(
(package, service_name),
("openshell.v1", "OpenShell") | ("openshell.inference.v1", "Inference")
) {
continue;
}
for method in &service.method {
let input = method.input_type.as_deref().expect("method input type");
let output = method.output_type.as_deref().expect("method output type");
public_roots.insert(input.to_string());
public_roots.insert(output.to_string());
methods.push(format!(
"{package}.{service_name}/{}|{}|{}|{}|{}",
method.name.as_deref().expect("method name"),
input,
output,
method.client_streaming.unwrap_or(false),
method.server_streaming.unwrap_or(false),
));
}
}
}
methods.sort();
assert_eq!(compiled_method_count, 100, "classify every compiled RPC");
assert_eq!(methods.len(), 74, "inventory every public gateway RPC");
assert_eq!(
methods
.iter()
.filter(|method| method.starts_with("openshell.v1.OpenShell/"))
.count(),
74
);
assert!(methods.iter().all(|method| !method.contains(".storage.")));
let public_closure = schema_closure(&index, public_roots);
let durable_closure = schema_closure(&index, DURABLE_ROOTS.into_iter().map(str::to_string));
let overlap_messages = public_closure
.messages
.intersection(&durable_closure.messages)
.cloned()
.collect::<BTreeSet<_>>();
let overlap_enums = public_closure
.enums
.intersection(&durable_closure.enums)
.cloned()
.collect::<BTreeSet<_>>();
let overlap_inventory = overlap_messages
.iter()
.map(|name| format!("message|{name}"))
.chain(overlap_enums.iter().map(|name| format!("enum|{name}")))
.collect::<Vec<_>>()
.join("\n");
let public_schema_hash = schema_fingerprint(&index, &public_closure);
let public_inventory_hash = format!(
"{:x}",
Sha256::digest(format!("{}\n{public_schema_hash}", methods.join("\n")).as_bytes())
);
let durable_inventory_hash = schema_fingerprint(&index, &durable_closure);
let overlap_hash = format!("{:x}", Sha256::digest(overlap_inventory.as_bytes()));
assert_eq!(
(public_closure.messages.len(), public_closure.enums.len()),
(276, 12)
);
assert_eq!(
(durable_closure.messages.len(), durable_closure.enums.len()),
(81, 8)
);
assert_eq!((overlap_messages.len(), overlap_enums.len()), (71, 8));
assert_eq!(
public_inventory_hash, PUBLIC_RPC_SCHEMA_SHA256,
"the public RPC schema closure changed; review API compatibility and update the inventory and architecture/gateway.md"
);
assert_eq!(
durable_inventory_hash, DURABLE_SCHEMA_SHA256,
"a durable protobuf root or transitive dependency changed; record migration handling and a prior-version fixture before updating this fingerprint"
);
assert_eq!(
overlap_hash, PUBLIC_DURABLE_OVERLAP_SHA256,
"the public/durable protobuf overlap changed; review both API and storage compatibility before updating this inventory"
);
}
fn legacy_bytes(encoded: &str) -> Vec<u8> {
hex::decode(encoded).expect("checked-in legacy fixture must be valid hex")
}
#[test]
fn pre_move_storage_payloads_decode_after_package_relocation() {
let refresh = StoredProviderCredentialRefreshState::decode(
legacy_bytes(V0_0_116_REFRESH_STATE).as_slice(),
)
.expect("legacy refresh state must decode");
assert_eq!(refresh.provider_id, "provider-id");
assert_eq!(refresh.provider_name, "provider");
assert_eq!(refresh.credential_key, "TOKEN");
assert_eq!(refresh.material["client_id"], "synthetic");
assert_eq!(refresh.secret_material_keys, ["refresh_token"]);
assert_eq!(refresh.expires_at_ms, 100);
assert_eq!(refresh.next_refresh_at_ms, 200);
assert_eq!(refresh.status, "active");
assert_eq!(refresh.scopes, ["scope-a"]);
assert_eq!(
refresh.secret_material_handles["refresh_token"].driver,
"test"
);
assert_eq!(refresh.pending_secret_deletions[0].material_key, "old");
let deletion =
StoredRefreshMaterialDeletion::decode(legacy_bytes(V0_0_116_DELETION).as_slice())
.expect("legacy deletion must decode");
assert_eq!(deletion.material_key, "old");
assert_eq!(deletion.handle.expect("handle").driver, "test");
let profile = StoredProviderProfile::decode(legacy_bytes(V0_0_116_PROFILE).as_slice())
.expect("legacy provider profile must decode");
let profile = profile.profile.expect("profile");
assert_eq!(profile.id, "profile");
assert_eq!(profile.display_name, "Legacy");
let policy_payload =
PolicyRevisionPayload::decode(legacy_bytes(V0_0_116_POLICY_PAYLOAD).as_slice())
.expect("legacy policy payload must decode");
assert!(policy_payload.policy.is_some());
assert_eq!(policy_payload.hash, "sha256");
assert_eq!(policy_payload.load_error, "none");
assert_eq!(policy_payload.loaded_at_ms, 300);
assert_eq!(policy_payload.provenance["source"], "fixture");
let draft_payload =
DraftChunkPayload::decode(legacy_bytes(V0_0_116_DRAFT_PAYLOAD).as_slice())
.expect("legacy draft payload must decode");
assert_eq!(draft_payload.rule_name, "rule");
assert_eq!(draft_payload.rationale, "fixture");
assert_eq!(draft_payload.confidence, 0.75);
assert_eq!(draft_payload.host, "example.com");
assert_eq!(draft_payload.port, 443);
assert_eq!(draft_payload.draft_version, 2);
let policy = StoredPolicyRevision::decode(legacy_bytes(V0_0_116_POLICY_RECORD).as_slice())
.expect("legacy policy record must decode");
assert_eq!(policy.id, "policy-id");
assert_eq!(policy.sandbox_id, "sandbox-id");
assert_eq!(policy.version, 2);
assert_eq!(policy.policy_payload, [1, 2, 3]);
assert_eq!(policy.policy_hash, "sha256");
assert_eq!(policy.status, "loaded");
assert_eq!(policy.load_error.as_deref(), Some("none"));
assert_eq!(policy.created_at_ms, 250);
assert_eq!(policy.loaded_at_ms, Some(300));
assert_eq!(policy.provenance["source"], "fixture");
let draft = StoredDraftChunk::decode(legacy_bytes(V0_0_116_DRAFT_RECORD).as_slice())
.expect("legacy draft record must decode");
assert_eq!(draft.id, "chunk-id");
assert_eq!(draft.sandbox_id, "sandbox-id");
assert_eq!(draft.draft_version, 2);
assert_eq!(draft.status, "pending");
assert_eq!(draft.rule_name, "rule");
assert_eq!(draft.proposed_rule, [4, 5]);
assert_eq!(draft.rationale, "fixture");
assert_eq!(draft.confidence, 0.75);
assert_eq!(draft.created_at_ms, 350);
assert_eq!(draft.decided_at_ms, Some(400));
assert_eq!(draft.host, "example.com");
assert_eq!(draft.port, 443);
assert_eq!(draft.hit_count, 1);
}
#[tokio::test]
async fn v0_0_116_database_payloads_survive_current_migrations() {
use crate::persistence::{DRAFT_CHUNK_OBJECT_TYPE, POLICY_OBJECT_TYPE, Store};
use crate::policy_store::{draft_chunk_record_from_parts, policy_record_from_parts};
let tempdir = tempfile::tempdir().expect("temporary database directory");
let database_path = tempdir.path().join("v0.0.116.db");
let database_url = format!("sqlite://{}", database_path.display());
// v0.0.116 and this change share migrations 001-006. Populate that
// historical on-disk shape with bytes emitted by the v0.0.116 schema,
// then reopen it through the current migration and loader path.
let old_store = Store::connect(&database_url)
.await
.expect("create legacy database");
let fixtures = [
(
"provider_credential_refresh_state",
"legacy-id",
"legacy-name",
"default",
V0_0_116_REFRESH_STATE,
),
(
"provider_profile",
"legacy-profile-id",
"legacy-profile",
"",
V0_0_116_PROFILE,
),
(
POLICY_OBJECT_TYPE,
"policy-id",
"",
"default",
V0_0_116_POLICY_PAYLOAD,
),
(
DRAFT_CHUNK_OBJECT_TYPE,
"chunk-id",
"",
"default",
V0_0_116_DRAFT_PAYLOAD,
),
];
for (object_type, id, name, workspace, payload) in fixtures {
old_store
.put(
object_type,
id,
name,
workspace,
&legacy_bytes(payload),
None,
)
.await
.expect("insert v0.0.116 fixture");
}
old_store.close().await;
let current_store = Store::connect(&database_url)
.await
.expect("current migrations must accept legacy database");
for (object_type, id, _, _, payload) in fixtures {
let record = current_store
.get(object_type, id)
.await
.expect("load fixture row")
.expect("fixture row must remain present");
assert_eq!(record.payload, legacy_bytes(payload));
}
let refresh = current_store
.get_message::<StoredProviderCredentialRefreshState>("legacy-id")
.await
.expect("decode refresh fixture")
.expect("refresh fixture must remain present");
assert_eq!(refresh.provider_id, "provider-id");
assert_eq!(refresh.get_resource_version(), 1);
let profile = current_store
.get_message::<StoredProviderProfile>("legacy-profile-id")
.await
.expect("decode profile fixture")
.expect("profile fixture must remain present");
assert_eq!(profile.profile.expect("profile").display_name, "Legacy");
let policy_record = current_store
.get(POLICY_OBJECT_TYPE, "policy-id")
.await
.expect("load policy fixture")
.expect("policy fixture must remain present");
let policy = policy_record_from_parts(
policy_record.id,
"sandbox-id".to_string(),
2,
"loaded".to_string(),
&policy_record.payload,
policy_record.created_at_ms,
)
.expect("current policy loader must decode v0.0.116 payload");
assert_eq!(policy.policy_hash, "sha256");
assert_eq!(policy.provenance["source"], "fixture");
let draft_record = current_store
.get(DRAFT_CHUNK_OBJECT_TYPE, "chunk-id")
.await
.expect("load draft fixture")
.expect("draft fixture must remain present");
let draft = draft_chunk_record_from_parts(
draft_record.id,
"sandbox-id".to_string(),
"pending".to_string(),
1,
&draft_record.payload,
draft_record.created_at_ms,
draft_record.updated_at_ms,
)
.expect("current draft loader must decode v0.0.116 payload");
assert_eq!(draft.rule_name, "rule");
assert_eq!(draft.host, "example.com");
assert_eq!(draft.port, 443);
}
}
-157
View File
@@ -1885,62 +1885,6 @@ message ProviderProfileDiscovery {
repeated string credentials = 1;
}
message StoredProviderCredentialRefreshState {
openshell.datamodel.v1.ObjectMeta metadata = 1;
string provider_id = 2;
string provider_name = 3;
string credential_key = 4;
ProviderCredentialRefreshStrategy strategy = 5;
map<string, string> material = 6 [(openshell.options.v1.secret) = true];
// Material names classified as secret. Newly configured values live in the
// active credential driver and are absent from material. Legacy inline values
// are not automatically migrated before OpenShell 0.1.0.
repeated string secret_material_keys = 7;
int64 expires_at_ms = 8;
// int64 max parks the refresh until an explicit rotation or reconfiguration.
int64 next_refresh_at_ms = 9;
int64 last_refresh_at_ms = 10;
string status = 11;
string last_error = 12;
string token_url = 13;
repeated string scopes = 14;
int64 refresh_before_seconds = 15;
int64 max_lifetime_seconds = 16;
// Resolved mapping of strategy-defined output id -> concrete env key, pinned
// at configure time from the profile's additional_outputs. Read by minting,
// collision reservation, and env-key surfacing so later profile edits cannot
// silently redirect writes.
map<string, string> additional_output_keys = 17;
// Opaque gateway-owned authorization epoch for the configured refresh
// grant. Explicit refresh configuration creates a new epoch; automatic and
// manual token rotation preserve it. It is never derived from or exposed
// with refresh material.
string authorization_epoch = 18;
// Secret refresh material is stored through the gateway's active credential
// driver. The persisted refresh state keeps only opaque handles; resolved
// values exist in gateway memory for the duration of one mint operation.
map<string, openshell.datamodel.v1.CredentialHandle> secret_material_handles = 19;
// Handles replaced by reconfiguration or issuer-driven refresh-token
// rotation. This is a repeated entry rather than a material-keyed map so
// multiple superseded generations of the same material remain recoverable.
// Cleanup is retried by the refresh worker so a gateway crash or temporary
// credential-backend outage does not lose the deletion reference.
repeated StoredRefreshMaterialDeletion pending_secret_deletions = 20;
// Structured recovery details for the most recent refresh failure. These
// fields contain only gateway-owned codes and recognized bounded values.
ProviderCredentialRefreshRecoveryAction recovery_action = 21;
string failure_code = 22;
string provider_error_subtype = 23;
int64 last_error_at_ms = 24;
}
message StoredRefreshMaterialDeletion {
// Original material name used to derive the credential driver's storage key.
string material_key = 1;
// Opaque handle for the superseded secret object.
openshell.datamodel.v1.CredentialHandle handle = 2;
}
message GetProviderRefreshStatusRequest {
string provider = 1;
string credential_key = 2;
@@ -2029,12 +1973,6 @@ message ProviderProfile {
string scope = 13;
}
// Stored custom provider profile object.
message StoredProviderProfile {
openshell.datamodel.v1.ObjectMeta metadata = 1;
ProviderProfile profile = 2;
}
// Provider profile response.
message ProviderProfileResponse {
ProviderProfile profile = 1;
@@ -2908,101 +2846,6 @@ message GetDraftHistoryResponse {
repeated DraftHistoryEntry entries = 1;
}
// Stored payload for a policy revision row in the generic objects table.
message PolicyRevisionPayload {
// Serialized policy contents.
openshell.sandbox.v1.SandboxPolicy policy = 1;
// Deterministic hash of the policy payload.
string hash = 2;
// Load error reported by the sandbox, if any.
string load_error = 3;
// When the policy version was reported as loaded (ms since epoch). 0 if unset.
int64 loaded_at_ms = 4;
// Immutable provenance supplied when this revision was created.
map<string, string> provenance = 5;
}
// Stored payload for a draft policy chunk row in the generic objects table.
message DraftChunkPayload {
// Proposed network_policies map key.
string rule_name = 1;
// Proposed network policy rule.
openshell.sandbox.v1.NetworkPolicyRule proposed_rule = 2;
// Human-readable explanation of why this rule is proposed.
string rationale = 3;
// Security concerns flagged by analysis (empty if none).
string security_notes = 4;
// Analysis confidence (0.0-1.0). 0 for mechanistic mode.
float confidence = 5;
// When the user approved/rejected (ms since epoch). 0 if undecided.
int64 decided_at_ms = 6;
// Denormalized endpoint host for dedup and display.
string host = 7;
// Denormalized endpoint port for dedup and display.
int32 port = 8;
// Binary path that triggered the denial.
string binary = 9;
// Current draft version for the owning sandbox.
int64 draft_version = 10;
// Gateway prover verdict for this chunk; empty until prover runs.
// Mirrors PolicyChunk.validation_result.
string validation_result = 11;
// Operator-supplied free-form rejection text; empty for non-rejected
// chunks. Mirrors PolicyChunk.rejection_reason.
string rejection_reason = 12;
string application_error = 13;
string review_token = 14;
string current_effective_policy_hash = 15;
string candidate_effective_policy_hash = 16;
openshell.sandbox.v1.SandboxPolicy current_effective_policy = 17;
openshell.sandbox.v1.SandboxPolicy candidate_effective_policy = 18;
}
// Internal stored policy revision row materialized from the generic objects table.
message StoredPolicyRevision {
string id = 1;
string sandbox_id = 2;
int64 version = 3;
bytes policy_payload = 4;
string policy_hash = 5;
string status = 6;
optional string load_error = 7;
int64 created_at_ms = 8;
optional int64 loaded_at_ms = 9;
map<string, string> provenance = 10;
}
// Internal stored draft chunk row materialized from the generic objects table.
message StoredDraftChunk {
string id = 1;
string sandbox_id = 2;
int64 draft_version = 3;
string status = 4;
string rule_name = 5;
bytes proposed_rule = 6;
string rationale = 7;
string security_notes = 8;
double confidence = 9;
int64 created_at_ms = 10;
optional int64 decided_at_ms = 11;
string host = 12;
int32 port = 13;
string binary = 14;
int32 hit_count = 15;
int64 first_seen_ms = 16;
int64 last_seen_ms = 17;
// Gateway prover verdict; empty until the prover runs. See PolicyChunk.
string validation_result = 18;
// Operator-supplied free-form rejection text. See PolicyChunk.
string rejection_reason = 19;
string application_error = 20;
string review_token = 21;
string current_effective_policy_hash = 22;
string candidate_effective_policy_hash = 23;
openshell.sandbox.v1.SandboxPolicy current_effective_policy = 24;
openshell.sandbox.v1.SandboxPolicy candidate_effective_policy = 25;
}
// ---------------------------------------------------------------------------
// Workspace messages
// ---------------------------------------------------------------------------
File diff suppressed because it is too large Load Diff