feat: integrate workspace memory settings authority
This commit is contained in:
@@ -748,6 +748,8 @@ impl MemoryConfig {
|
|||||||
query_result_limit: upper.query_result_limit.or(self.query_result_limit),
|
query_result_limit: upper.query_result_limit.or(self.query_result_limit),
|
||||||
query_excerpt_lines: upper.query_excerpt_lines.or(self.query_excerpt_lines),
|
query_excerpt_lines: upper.query_excerpt_lines.or(self.query_excerpt_lines),
|
||||||
inject_summary: upper.inject_summary.or(self.inject_summary),
|
inject_summary: upper.inject_summary.or(self.inject_summary),
|
||||||
|
workspace_id: upper.workspace_id.or(self.workspace_id),
|
||||||
|
settings_revision: upper.settings_revision.or(self.settings_revision),
|
||||||
language: upper.language.or(self.language),
|
language: upper.language.or(self.language),
|
||||||
extract_model: upper.extract_model.or(self.extract_model),
|
extract_model: upper.extract_model.or(self.extract_model),
|
||||||
extract_threshold: upper.extract_threshold.or(self.extract_threshold),
|
extract_threshold: upper.extract_threshold.or(self.extract_threshold),
|
||||||
|
|||||||
@@ -95,7 +95,3 @@ pub const COMPACT_DEFAULT_REFERENCE_COUNT: usize = 5;
|
|||||||
/// Optional maximum extract-worker tool-loop depth. `None` means unlimited.
|
/// Optional maximum extract-worker tool-loop depth. `None` means unlimited.
|
||||||
/// See [`crate::MemoryConfig::extract_worker_max_turns`].
|
/// See [`crate::MemoryConfig::extract_worker_max_turns`].
|
||||||
pub const MEMORY_EXTRACT_WORKER_MAX_TURNS: Option<u32> = Some(8);
|
pub const MEMORY_EXTRACT_WORKER_MAX_TURNS: Option<u32> = Some(8);
|
||||||
|
|
||||||
/// Default language used by memory extraction / consolidation workers for
|
|
||||||
/// durable memory text. See [`crate::MemoryConfig::language`].
|
|
||||||
pub const MEMORY_LANGUAGE: &str = "English";
|
|
||||||
|
|||||||
@@ -449,6 +449,29 @@ pub struct WebFetchConfig {
|
|||||||
pub allow_private_addresses: Option<bool>,
|
pub allow_private_addresses: Option<bool>,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Maximum Unicode scalar values accepted in a normalized Workspace Memory language name.
|
||||||
|
pub const MAX_WORKSPACE_MEMORY_LANGUAGE_CHARS: usize = 64;
|
||||||
|
|
||||||
|
/// Return whether a Workspace Memory language is already normalized and safe to persist.
|
||||||
|
pub fn is_normalized_workspace_memory_language(language: &str) -> bool {
|
||||||
|
!language.is_empty()
|
||||||
|
&& language == language.trim()
|
||||||
|
&& language.chars().count() <= MAX_WORKSPACE_MEMORY_LANGUAGE_CHARS
|
||||||
|
&& !language.chars().any(char::is_control)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Immutable Workspace Memory settings bound into a Worker launch snapshot.
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||||
|
#[serde(deny_unknown_fields)]
|
||||||
|
pub struct WorkspaceMemorySettingsSnapshot {
|
||||||
|
/// Workspace that owns the settings revision.
|
||||||
|
pub workspace_id: String,
|
||||||
|
/// Monotonic Workspace Memory settings revision.
|
||||||
|
pub settings_revision: u64,
|
||||||
|
/// Normalized language used for Memory extraction and consolidation output.
|
||||||
|
pub language: String,
|
||||||
|
}
|
||||||
|
|
||||||
/// Memory subsystem configuration. Presence in the manifest enables
|
/// Memory subsystem configuration. Presence in the manifest enables
|
||||||
/// memory; `workspace_root` pins the memory workspace explicitly. When it
|
/// memory; `workspace_root` pins the memory workspace explicitly. When it
|
||||||
/// is absent, memory resolution searches upward from the Worker's pwd for a
|
/// is absent, memory resolution searches upward from the Worker's pwd for a
|
||||||
@@ -477,11 +500,14 @@ pub struct MemoryConfig {
|
|||||||
/// system-prompt section. `None` ⇒ enabled.
|
/// system-prompt section. `None` ⇒ enabled.
|
||||||
#[serde(default)]
|
#[serde(default)]
|
||||||
pub inject_summary: Option<bool>,
|
pub inject_summary: Option<bool>,
|
||||||
/// Language used by memory extraction / consolidation sub_worker for durable
|
/// Workspace that owns the bound Memory settings revision.
|
||||||
/// memory text. Free-form so workspaces can use names like
|
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||||
/// `English`, `Japanese`, or locale tags. `None` ⇒
|
pub workspace_id: Option<String>,
|
||||||
/// [`defaults::MEMORY_LANGUAGE`].
|
/// Monotonic revision of the bound Workspace Memory settings.
|
||||||
#[serde(default)]
|
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||||
|
pub settings_revision: Option<u64>,
|
||||||
|
/// Language from the bound Workspace Memory settings revision.
|
||||||
|
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||||
pub language: Option<String>,
|
pub language: Option<String>,
|
||||||
/// Optional model for the extract worker. When `None`,
|
/// Optional model for the extract worker. When `None`,
|
||||||
/// the main engine model is cloned via `clone_boxed()`. Lightweight
|
/// the main engine model is cloned via `clone_boxed()`. Lightweight
|
||||||
@@ -520,6 +546,24 @@ pub struct MemoryConfig {
|
|||||||
pub consolidation_threshold_bytes: Option<u64>,
|
pub consolidation_threshold_bytes: Option<u64>,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
impl MemoryConfig {
|
||||||
|
/// Replace any untrusted manifest values with a trusted Workspace snapshot.
|
||||||
|
pub fn bind_workspace_settings(&mut self, snapshot: &WorkspaceMemorySettingsSnapshot) {
|
||||||
|
self.workspace_id = Some(snapshot.workspace_id.clone());
|
||||||
|
self.settings_revision = Some(snapshot.settings_revision);
|
||||||
|
self.language = Some(snapshot.language.clone());
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Return the complete bound Workspace settings snapshot, if every field is present.
|
||||||
|
pub fn workspace_settings(&self) -> Option<WorkspaceMemorySettingsSnapshot> {
|
||||||
|
Some(WorkspaceMemorySettingsSnapshot {
|
||||||
|
workspace_id: self.workspace_id.clone()?,
|
||||||
|
settings_revision: self.settings_revision?,
|
||||||
|
language: self.language.clone()?,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Worker metadata.
|
/// Worker metadata.
|
||||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||||
pub struct WorkerMeta {
|
pub struct WorkerMeta {
|
||||||
@@ -1223,6 +1267,18 @@ model_id = "claude-sonnet-4-20250514"
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn workspace_memory_language_validation_is_bounded_free_form_utf8() {
|
||||||
|
assert!(is_normalized_workspace_memory_language("Français"));
|
||||||
|
assert!(is_normalized_workspace_memory_language("日本語"));
|
||||||
|
assert!(!is_normalized_workspace_memory_language(""));
|
||||||
|
assert!(!is_normalized_workspace_memory_language(" English "));
|
||||||
|
assert!(!is_normalized_workspace_memory_language("English\n"));
|
||||||
|
assert!(!is_normalized_workspace_memory_language(
|
||||||
|
&"x".repeat(MAX_WORKSPACE_MEMORY_LANGUAGE_CHARS + 1)
|
||||||
|
));
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn memory_section_with_language() {
|
fn memory_section_with_language() {
|
||||||
let toml = format!("{MINIMAL_REQUIRED}\n[memory]\nlanguage = \"Japanese\"\n");
|
let toml = format!("{MINIMAL_REQUIRED}\n[memory]\nlanguage = \"Japanese\"\n");
|
||||||
|
|||||||
@@ -562,7 +562,7 @@ fn resolve_profile_value(
|
|||||||
mcp: profile.mcp,
|
mcp: profile.mcp,
|
||||||
compaction,
|
compaction,
|
||||||
web: profile.web,
|
web: profile.web,
|
||||||
memory: profile.memory,
|
memory: profile.memory.map(Into::into),
|
||||||
skills: profile.skills,
|
skills: profile.skills,
|
||||||
};
|
};
|
||||||
let config = WorkerManifestConfig::builtin_defaults().merge(config.resolve_paths(profile_dir));
|
let config = WorkerManifestConfig::builtin_defaults().merge(config.resolve_paths(profile_dir));
|
||||||
@@ -582,6 +582,51 @@ fn resolve_profile_value(
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Default, Deserialize)]
|
||||||
|
#[serde(deny_unknown_fields)]
|
||||||
|
struct ProfileMemoryConfig {
|
||||||
|
#[serde(default)]
|
||||||
|
workspace_root: Option<PathBuf>,
|
||||||
|
#[serde(default)]
|
||||||
|
query_result_limit: Option<usize>,
|
||||||
|
#[serde(default)]
|
||||||
|
query_excerpt_lines: Option<usize>,
|
||||||
|
#[serde(default)]
|
||||||
|
inject_summary: Option<bool>,
|
||||||
|
#[serde(default)]
|
||||||
|
extract_model: Option<ModelManifest>,
|
||||||
|
#[serde(default)]
|
||||||
|
extract_threshold: Option<u64>,
|
||||||
|
#[serde(default)]
|
||||||
|
extract_worker_max_turns: Option<u32>,
|
||||||
|
#[serde(default)]
|
||||||
|
consolidation_model: Option<ModelManifest>,
|
||||||
|
#[serde(default)]
|
||||||
|
consolidation_threshold_files: Option<usize>,
|
||||||
|
#[serde(default)]
|
||||||
|
consolidation_threshold_bytes: Option<u64>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl From<ProfileMemoryConfig> for MemoryConfig {
|
||||||
|
fn from(profile: ProfileMemoryConfig) -> Self {
|
||||||
|
Self {
|
||||||
|
workspace_root: profile.workspace_root,
|
||||||
|
query_result_limit: profile.query_result_limit,
|
||||||
|
query_excerpt_lines: profile.query_excerpt_lines,
|
||||||
|
inject_summary: profile.inject_summary,
|
||||||
|
workspace_id: None,
|
||||||
|
settings_revision: None,
|
||||||
|
language: None,
|
||||||
|
extract_model: profile.extract_model,
|
||||||
|
extract_threshold: profile.extract_threshold,
|
||||||
|
extract_worker_max_turns: profile.extract_worker_max_turns,
|
||||||
|
consolidation_model: profile.consolidation_model,
|
||||||
|
consolidation_threshold_files: profile.consolidation_threshold_files,
|
||||||
|
consolidation_threshold_bytes: profile.consolidation_threshold_bytes,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[derive(Debug, Default, Deserialize)]
|
#[derive(Debug, Default, Deserialize)]
|
||||||
#[serde(deny_unknown_fields)]
|
#[serde(deny_unknown_fields)]
|
||||||
struct ProfileConfig {
|
struct ProfileConfig {
|
||||||
@@ -612,7 +657,7 @@ struct ProfileConfig {
|
|||||||
#[serde(default)]
|
#[serde(default)]
|
||||||
web: Option<WebConfig>,
|
web: Option<WebConfig>,
|
||||||
#[serde(default)]
|
#[serde(default)]
|
||||||
memory: Option<MemoryConfig>,
|
memory: Option<ProfileMemoryConfig>,
|
||||||
#[serde(default)]
|
#[serde(default)]
|
||||||
skills: Option<SkillsConfig>,
|
skills: Option<SkillsConfig>,
|
||||||
}
|
}
|
||||||
@@ -1334,6 +1379,34 @@ mod tests {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn profile_rejects_workspace_memory_snapshot_authority_fields() {
|
||||||
|
let tmp = TempDir::new().unwrap();
|
||||||
|
for (field, value) in [
|
||||||
|
("workspace_id", serde_json::json!("workspace-a")),
|
||||||
|
("settings_revision", serde_json::json!(2)),
|
||||||
|
("language", serde_json::json!("Japanese")),
|
||||||
|
] {
|
||||||
|
let artifact = serde_json::json!({ "memory": { (field): value } });
|
||||||
|
let error = resolve_profile_artifact_value(
|
||||||
|
artifact,
|
||||||
|
ProfileSource::Registry {
|
||||||
|
source: ProfileRegistrySource::Builtin,
|
||||||
|
name: "test".to_string(),
|
||||||
|
path: None,
|
||||||
|
provenance: None,
|
||||||
|
},
|
||||||
|
tmp.path(),
|
||||||
|
"test-worker",
|
||||||
|
)
|
||||||
|
.unwrap_err();
|
||||||
|
assert!(
|
||||||
|
error.to_string().contains("unknown field"),
|
||||||
|
"unexpected error for {field}: {error}"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn builtin_companion_can_manage_workdirs() {
|
fn builtin_companion_can_manage_workdirs() {
|
||||||
let tmp = TempDir::new().unwrap();
|
let tmp = TempDir::new().unwrap();
|
||||||
|
|||||||
@@ -177,6 +177,13 @@ impl OperationCounts {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||||
|
pub struct MemorySettingsAudit {
|
||||||
|
pub workspace_id: String,
|
||||||
|
pub settings_revision: u64,
|
||||||
|
pub language: String,
|
||||||
|
}
|
||||||
|
|
||||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||||
pub struct WorkerLifecycleAudit {
|
pub struct WorkerLifecycleAudit {
|
||||||
pub run_id: Uuid,
|
pub run_id: Uuid,
|
||||||
@@ -185,6 +192,8 @@ pub struct WorkerLifecycleAudit {
|
|||||||
pub trigger: AuditTrigger,
|
pub trigger: AuditTrigger,
|
||||||
pub reason: String,
|
pub reason: String,
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||||
|
pub memory_settings: Option<MemorySettingsAudit>,
|
||||||
|
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||||
pub model: Option<ModelAudit>,
|
pub model: Option<ModelAudit>,
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||||
pub usage: Option<UsageAudit>,
|
pub usage: Option<UsageAudit>,
|
||||||
@@ -405,6 +414,7 @@ mod tests {
|
|||||||
status: WorkerLifecycleStatus::Started,
|
status: WorkerLifecycleStatus::Started,
|
||||||
trigger: AuditTrigger::TokenThreshold,
|
trigger: AuditTrigger::TokenThreshold,
|
||||||
reason: "tokens_threshold_reached".to_string(),
|
reason: "tokens_threshold_reached".to_string(),
|
||||||
|
memory_settings: None,
|
||||||
model: None,
|
model: None,
|
||||||
usage: None,
|
usage: None,
|
||||||
extract: None,
|
extract: None,
|
||||||
|
|||||||
@@ -174,6 +174,9 @@ pub struct CreateWorkerRequest {
|
|||||||
pub worker_observation_grants: Vec<RuntimeWorkerRef>,
|
pub worker_observation_grants: Vec<RuntimeWorkerRef>,
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||||
pub workspace_api: Option<WorkspaceApiRef>,
|
pub workspace_api: Option<WorkspaceApiRef>,
|
||||||
|
/// Backend-authored immutable Workspace Memory settings snapshot.
|
||||||
|
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||||
|
pub memory_settings: Option<manifest::WorkspaceMemorySettingsSnapshot>,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Worker lifecycle status for the in-memory embedded runtime.
|
/// Worker lifecycle status for the in-memory embedded runtime.
|
||||||
|
|||||||
@@ -1882,6 +1882,11 @@ mod tests {
|
|||||||
workspace_id: workspace_id.to_string(),
|
workspace_id: workspace_id.to_string(),
|
||||||
base_url: format!("https://workspace.example/{workspace_id}"),
|
base_url: format!("https://workspace.example/{workspace_id}"),
|
||||||
});
|
});
|
||||||
|
request.memory_settings = Some(manifest::WorkspaceMemorySettingsSnapshot {
|
||||||
|
workspace_id: workspace_id.to_string(),
|
||||||
|
settings_revision: 1,
|
||||||
|
language: "English".to_string(),
|
||||||
|
});
|
||||||
request
|
request
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -2193,6 +2198,11 @@ mod tests {
|
|||||||
worker_observation_enabled: false,
|
worker_observation_enabled: false,
|
||||||
worker_observation_grants: Vec::new(),
|
worker_observation_grants: Vec::new(),
|
||||||
workspace_api: None,
|
workspace_api: None,
|
||||||
|
memory_settings: Some(manifest::WorkspaceMemorySettingsSnapshot {
|
||||||
|
workspace_id: "local".to_string(),
|
||||||
|
settings_revision: 1,
|
||||||
|
language: "English".to_string(),
|
||||||
|
}),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -2828,6 +2838,11 @@ mod ws_tests {
|
|||||||
worker_observation_enabled: false,
|
worker_observation_enabled: false,
|
||||||
worker_observation_grants: Vec::new(),
|
worker_observation_grants: Vec::new(),
|
||||||
workspace_api: None,
|
workspace_api: None,
|
||||||
|
memory_settings: Some(manifest::WorkspaceMemorySettingsSnapshot {
|
||||||
|
workspace_id: "local".to_string(),
|
||||||
|
settings_revision: 1,
|
||||||
|
language: "English".to_string(),
|
||||||
|
}),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -2794,6 +2794,27 @@ fn validate_create_workspace_scope(
|
|||||||
)));
|
)));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
let snapshot = request.memory_settings.as_ref().ok_or_else(|| {
|
||||||
|
RuntimeError::InvalidRequest(
|
||||||
|
"Workspace-scoped Worker create requires a bound Memory settings snapshot".to_string(),
|
||||||
|
)
|
||||||
|
})?;
|
||||||
|
if snapshot.workspace_id != workspace_id {
|
||||||
|
return Err(RuntimeError::InvalidRequest(format!(
|
||||||
|
"Memory settings workspace_id {} does not match Runtime auth workspace_id {workspace_id}",
|
||||||
|
snapshot.workspace_id
|
||||||
|
)));
|
||||||
|
}
|
||||||
|
if snapshot.settings_revision == 0 {
|
||||||
|
return Err(RuntimeError::InvalidRequest(
|
||||||
|
"Memory settings revision must be at least 1".to_string(),
|
||||||
|
));
|
||||||
|
}
|
||||||
|
if !manifest::is_normalized_workspace_memory_language(&snapshot.language) {
|
||||||
|
return Err(RuntimeError::InvalidRequest(
|
||||||
|
"Memory settings language must be a normalized bounded UTF-8 value".to_string(),
|
||||||
|
));
|
||||||
|
}
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -3086,6 +3107,11 @@ mod tests {
|
|||||||
worker_observation_enabled: false,
|
worker_observation_enabled: false,
|
||||||
worker_observation_grants: Vec::new(),
|
worker_observation_grants: Vec::new(),
|
||||||
workspace_api: None,
|
workspace_api: None,
|
||||||
|
memory_settings: Some(manifest::WorkspaceMemorySettingsSnapshot {
|
||||||
|
workspace_id: "local".to_string(),
|
||||||
|
settings_revision: 1,
|
||||||
|
language: "English".to_string(),
|
||||||
|
}),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -3095,9 +3121,39 @@ mod tests {
|
|||||||
workspace_id: workspace_id.to_string(),
|
workspace_id: workspace_id.to_string(),
|
||||||
base_url: format!("https://workspace.example/{workspace_id}"),
|
base_url: format!("https://workspace.example/{workspace_id}"),
|
||||||
});
|
});
|
||||||
|
request.memory_settings = Some(manifest::WorkspaceMemorySettingsSnapshot {
|
||||||
|
workspace_id: workspace_id.to_string(),
|
||||||
|
settings_revision: 1,
|
||||||
|
language: "English".to_string(),
|
||||||
|
});
|
||||||
request
|
request
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn workspace_create_requires_matching_normalized_memory_settings_snapshot() {
|
||||||
|
let mut request = scoped_task_request("memory-snapshot", "workspace-a");
|
||||||
|
assert!(validate_create_workspace_scope(&request, Some("workspace-a")).is_ok());
|
||||||
|
request.memory_settings.as_mut().unwrap().language = "Français".to_string();
|
||||||
|
assert!(validate_create_workspace_scope(&request, Some("workspace-a")).is_ok());
|
||||||
|
|
||||||
|
request.memory_settings = None;
|
||||||
|
assert!(validate_create_workspace_scope(&request, Some("workspace-a")).is_err());
|
||||||
|
|
||||||
|
request.memory_settings = Some(manifest::WorkspaceMemorySettingsSnapshot {
|
||||||
|
workspace_id: "workspace-b".to_string(),
|
||||||
|
settings_revision: 1,
|
||||||
|
language: "English".to_string(),
|
||||||
|
});
|
||||||
|
assert!(validate_create_workspace_scope(&request, Some("workspace-a")).is_err());
|
||||||
|
|
||||||
|
request.memory_settings = Some(manifest::WorkspaceMemorySettingsSnapshot {
|
||||||
|
workspace_id: "workspace-a".to_string(),
|
||||||
|
settings_revision: 2,
|
||||||
|
language: " english ".to_string(),
|
||||||
|
});
|
||||||
|
assert!(validate_create_workspace_scope(&request, Some("workspace-a")).is_err());
|
||||||
|
}
|
||||||
|
|
||||||
fn scope(workspace_id: &str, server_id: &str) -> RuntimeWorkspaceScope {
|
fn scope(workspace_id: &str, server_id: &str) -> RuntimeWorkspaceScope {
|
||||||
RuntimeWorkspaceScope::new(workspace_id, server_id)
|
RuntimeWorkspaceScope::new(workspace_id, server_id)
|
||||||
}
|
}
|
||||||
@@ -4484,7 +4540,7 @@ mod tests {
|
|||||||
let worker = runtime
|
let worker = runtime
|
||||||
.create_worker_scoped(
|
.create_worker_scoped(
|
||||||
&RuntimeWorkspaceScope::new("workspace-a", "server"),
|
&RuntimeWorkspaceScope::new("workspace-a", "server"),
|
||||||
task_request("legacy"),
|
scoped_task_request("legacy", "workspace-a"),
|
||||||
)
|
)
|
||||||
.unwrap();
|
.unwrap();
|
||||||
drop(runtime);
|
drop(runtime);
|
||||||
|
|||||||
@@ -642,6 +642,60 @@ fn runtime_local_workdir_session(
|
|||||||
))
|
))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn bind_workspace_memory_settings(
|
||||||
|
manifest: &mut manifest::WorkerManifest,
|
||||||
|
request: &CreateWorkerRequest,
|
||||||
|
) -> Result<(), String> {
|
||||||
|
let Some(snapshot) = request.memory_settings.as_ref() else {
|
||||||
|
if request.workspace_api.is_some() {
|
||||||
|
return Err(
|
||||||
|
"Workspace Worker request is missing its bound Memory settings snapshot"
|
||||||
|
.to_string(),
|
||||||
|
);
|
||||||
|
}
|
||||||
|
return Ok(());
|
||||||
|
};
|
||||||
|
if let Some(workspace_api) = request.workspace_api.as_ref()
|
||||||
|
&& snapshot.workspace_id != workspace_api.workspace_id
|
||||||
|
{
|
||||||
|
return Err(format!(
|
||||||
|
"Memory settings workspace {} does not match Workspace API scope {}",
|
||||||
|
snapshot.workspace_id, workspace_api.workspace_id
|
||||||
|
));
|
||||||
|
}
|
||||||
|
manifest
|
||||||
|
.memory
|
||||||
|
.get_or_insert_with(manifest::MemoryConfig::default)
|
||||||
|
.bind_workspace_settings(snapshot);
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn validate_worker_memory_settings(
|
||||||
|
manifest: &manifest::WorkerManifest,
|
||||||
|
request: &CreateWorkerRequest,
|
||||||
|
) -> Result<(), String> {
|
||||||
|
let Some(expected) = request.memory_settings.as_ref() else {
|
||||||
|
return Ok(());
|
||||||
|
};
|
||||||
|
let actual = manifest
|
||||||
|
.memory
|
||||||
|
.as_ref()
|
||||||
|
.and_then(manifest::MemoryConfig::workspace_settings)
|
||||||
|
.ok_or_else(|| {
|
||||||
|
"Workspace Worker restored without its bound Memory settings snapshot".to_string()
|
||||||
|
})?;
|
||||||
|
if &actual != expected {
|
||||||
|
return Err(format!(
|
||||||
|
"Workspace Worker Memory settings snapshot mismatch: expected {} revision {}, restored {} revision {}",
|
||||||
|
expected.workspace_id,
|
||||||
|
expected.settings_revision,
|
||||||
|
actual.workspace_id,
|
||||||
|
actual.settings_revision
|
||||||
|
));
|
||||||
|
}
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
#[async_trait]
|
#[async_trait]
|
||||||
impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
|
impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
|
||||||
fn observe_workspace_prompt_projection(
|
fn observe_workspace_prompt_projection(
|
||||||
@@ -696,7 +750,7 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
|
|||||||
let archive = self
|
let archive = self
|
||||||
.resolve_profile_source_archive(&request.request.profile_source)
|
.resolve_profile_source_archive(&request.request.profile_source)
|
||||||
.await?;
|
.await?;
|
||||||
let (manifest, mut loader) = {
|
let (mut manifest, mut loader) = {
|
||||||
let manifest = archive
|
let manifest = archive
|
||||||
.resolve_profile(selector, &worker_root, &worker_name)
|
.resolve_profile(selector, &worker_root, &worker_name)
|
||||||
.map_err(|err| format!("failed to resolve profile source archive: {err}"))?;
|
.map_err(|err| format!("failed to resolve profile source archive: {err}"))?;
|
||||||
@@ -714,6 +768,7 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
|
|||||||
)?
|
)?
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
bind_workspace_memory_settings(&mut manifest, &request.request)?;
|
||||||
if let Some(bundle) = request.config_bundle.as_ref()
|
if let Some(bundle) = request.config_bundle.as_ref()
|
||||||
&& let Some(resolution) =
|
&& let Some(resolution) =
|
||||||
self.observe_bundle_prompt_projection(bundle, observation_workspace_id.as_deref())?
|
self.observe_bundle_prompt_projection(bundle, observation_workspace_id.as_deref())?
|
||||||
@@ -750,6 +805,7 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
|
|||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.map_err(|err| format!("failed to create Worker from profile: {err}"))?;
|
.map_err(|err| format!("failed to create Worker from profile: {err}"))?;
|
||||||
|
validate_worker_memory_settings(worker.manifest(), &request.request)?;
|
||||||
if let Some(binding) = request.working_directory.as_ref() {
|
if let Some(binding) = request.working_directory.as_ref() {
|
||||||
worker.bind_workdir_session(Some(runtime_local_workdir_session(
|
worker.bind_workdir_session(Some(runtime_local_workdir_session(
|
||||||
&binding.working_directory.id,
|
&binding.working_directory.id,
|
||||||
@@ -856,7 +912,8 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
|
|||||||
self.embedded_worker_mutation_dispatcher.as_ref(),
|
self.embedded_worker_mutation_dispatcher.as_ref(),
|
||||||
Some(self.prompt_projection_cache.clone()),
|
Some(self.prompt_projection_cache.clone()),
|
||||||
);
|
);
|
||||||
let (manifest, loader) = Self::restore_fallback_manifest(&worker_name)?;
|
let (mut manifest, loader) = Self::restore_fallback_manifest(&worker_name)?;
|
||||||
|
bind_workspace_memory_settings(&mut manifest, &request.request)?;
|
||||||
|
|
||||||
let worker_aggregate_dir = self.worker_aggregate_dir(&request.worker_ref)?;
|
let worker_aggregate_dir = self.worker_aggregate_dir(&request.worker_ref)?;
|
||||||
let session_dir = worker_aggregate_dir.join("session");
|
let session_dir = worker_aggregate_dir.join("session");
|
||||||
@@ -924,6 +981,7 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory {
|
|||||||
}
|
}
|
||||||
Err(err) => return Err(format!("failed to restore Worker from metadata: {err}")),
|
Err(err) => return Err(format!("failed to restore Worker from metadata: {err}")),
|
||||||
};
|
};
|
||||||
|
validate_worker_memory_settings(worker.manifest(), &request.request)?;
|
||||||
let flow_transition_enabled = worker.manifest().feature.flow.enabled;
|
let flow_transition_enabled = worker.manifest().feature.flow.enabled;
|
||||||
if let Some(binding) = request.working_directory.as_ref() {
|
if let Some(binding) = request.working_directory.as_ref() {
|
||||||
worker.bind_workdir_session(Some(runtime_local_workdir_session(
|
worker.bind_workdir_session(Some(runtime_local_workdir_session(
|
||||||
@@ -2504,6 +2562,7 @@ mod tests {
|
|||||||
worker_observation_enabled: false,
|
worker_observation_enabled: false,
|
||||||
worker_observation_grants: Vec::new(),
|
worker_observation_grants: Vec::new(),
|
||||||
workspace_api: None,
|
workspace_api: None,
|
||||||
|
memory_settings: None,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -2765,7 +2824,7 @@ mod tests {
|
|||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[serial_test::serial(worker_allocation)]
|
#[serial_test::serial(worker_allocation)]
|
||||||
async fn restore_pending_workspace_worker_without_system_prompt_fails_closed() {
|
async fn restore_legacy_workspace_worker_without_manifest_snapshot_requires_replacement() {
|
||||||
let root = tempfile::tempdir().unwrap();
|
let root = tempfile::tempdir().unwrap();
|
||||||
let runtime_store_dir = root.path().join("runtime");
|
let runtime_store_dir = root.path().join("runtime");
|
||||||
let worker_ref = WorkerRef::new(crate::identity::WorkerId::from_legacy_u64(1));
|
let worker_ref = WorkerRef::new(crate::identity::WorkerId::from_legacy_u64(1));
|
||||||
@@ -2774,32 +2833,6 @@ mod tests {
|
|||||||
.join(worker_ref.worker_id.to_string());
|
.join(worker_ref.worker_id.to_string());
|
||||||
let worker_name = ProfileRuntimeWorkerFactory::runtime_worker_name_for_ref(&worker_ref);
|
let worker_name = ProfileRuntimeWorkerFactory::runtime_worker_name_for_ref(&worker_ref);
|
||||||
let session_id = session_store::new_session_id();
|
let session_id = session_store::new_session_id();
|
||||||
let manifest = manifest::WorkerManifest::from_toml(&format!(
|
|
||||||
r#"
|
|
||||||
[worker]
|
|
||||||
name = "{}"
|
|
||||||
pwd = "{}"
|
|
||||||
|
|
||||||
[model]
|
|
||||||
scheme = "anthropic"
|
|
||||||
model_id = "test-model"
|
|
||||||
auth = {{ kind = "none" }}
|
|
||||||
|
|
||||||
[engine]
|
|
||||||
max_tokens = 100
|
|
||||||
|
|
||||||
[feature.flow]
|
|
||||||
enabled = true
|
|
||||||
|
|
||||||
[[scope.allow]]
|
|
||||||
target = "{}"
|
|
||||||
permission = "write"
|
|
||||||
"#,
|
|
||||||
worker_name,
|
|
||||||
root.path().display(),
|
|
||||||
root.path().display(),
|
|
||||||
))
|
|
||||||
.unwrap();
|
|
||||||
WorkerAggregateStore::new(&worker_aggregate_dir, &worker_name)
|
WorkerAggregateStore::new(&worker_aggregate_dir, &worker_name)
|
||||||
.unwrap()
|
.unwrap()
|
||||||
.set_active(
|
.set_active(
|
||||||
@@ -2807,7 +2840,7 @@ mod tests {
|
|||||||
Some(session_store::WorkerActiveSegmentRef::pending_segment(
|
Some(session_store::WorkerActiveSegmentRef::pending_segment(
|
||||||
session_id,
|
session_id,
|
||||||
)),
|
)),
|
||||||
Some(serde_json::to_value(&manifest).unwrap()),
|
None,
|
||||||
)
|
)
|
||||||
.unwrap();
|
.unwrap();
|
||||||
|
|
||||||
@@ -2816,6 +2849,11 @@ mod tests {
|
|||||||
workspace_id: "workspace-restore".to_string(),
|
workspace_id: "workspace-restore".to_string(),
|
||||||
base_url: "http://workspace.invalid".to_string(),
|
base_url: "http://workspace.invalid".to_string(),
|
||||||
});
|
});
|
||||||
|
request.memory_settings = Some(manifest::WorkspaceMemorySettingsSnapshot {
|
||||||
|
workspace_id: "workspace-restore".to_string(),
|
||||||
|
settings_revision: 1,
|
||||||
|
language: "English".to_string(),
|
||||||
|
});
|
||||||
let identity = RuntimeIdentityMaterial::generate("runtime-restore").unwrap();
|
let identity = RuntimeIdentityMaterial::generate("runtime-restore").unwrap();
|
||||||
let error = match ProfileRuntimeWorkerFactory::new(root.path())
|
let error = match ProfileRuntimeWorkerFactory::new(root.path())
|
||||||
.with_runtime_store_dir(&runtime_store_dir)
|
.with_runtime_store_dir(&runtime_store_dir)
|
||||||
@@ -2835,10 +2873,10 @@ mod tests {
|
|||||||
})
|
})
|
||||||
.await
|
.await
|
||||||
{
|
{
|
||||||
Ok(_) => panic!("pending Workspace Worker restore unexpectedly succeeded"),
|
Ok(_) => panic!("legacy Workspace Worker restore unexpectedly succeeded"),
|
||||||
Err(error) => error,
|
Err(error) => error,
|
||||||
};
|
};
|
||||||
assert!(error.contains("requires operation-owned launch material"));
|
assert!(error.contains("replacement Worker is required"), "{error}");
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
|
|||||||
+261
-11
@@ -2163,6 +2163,33 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
|||||||
let Some(template) = self.system_prompt_template.take() else {
|
let Some(template) = self.system_prompt_template.take() else {
|
||||||
return Ok(());
|
return Ok(());
|
||||||
};
|
};
|
||||||
|
let is_memory_consolidation = self.manifest.profile.as_ref().is_some_and(|snapshot| {
|
||||||
|
matches!(
|
||||||
|
&snapshot.source,
|
||||||
|
manifest::ProfileSource::Registry {
|
||||||
|
source: manifest::ProfileRegistrySource::Builtin,
|
||||||
|
name,
|
||||||
|
..
|
||||||
|
} if name == "memory-consolidation"
|
||||||
|
)
|
||||||
|
});
|
||||||
|
if is_memory_consolidation {
|
||||||
|
let memory_config = self.manifest.memory.as_ref().ok_or_else(|| {
|
||||||
|
WorkerError::InvalidState(
|
||||||
|
"Memory consolidation Worker has no Memory configuration".to_string(),
|
||||||
|
)
|
||||||
|
})?;
|
||||||
|
let language = memory_language(memory_config)?;
|
||||||
|
let rendered = self
|
||||||
|
.prompts
|
||||||
|
.load_full()
|
||||||
|
.memory_consolidation_system(&language)?;
|
||||||
|
self.engine
|
||||||
|
.as_mut()
|
||||||
|
.expect("worker present")
|
||||||
|
.set_system_prompt(rendered);
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
let alerter = self.alerter.clone();
|
let alerter = self.alerter.clone();
|
||||||
let tool_names: Vec<String> = {
|
let tool_names: Vec<String> = {
|
||||||
let worker = self.engine.as_mut().expect("worker present");
|
let worker = self.engine.as_mut().expect("worker present");
|
||||||
@@ -3750,6 +3777,7 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
|||||||
memory::audit::AuditTrigger::TokenThreshold,
|
memory::audit::AuditTrigger::TokenThreshold,
|
||||||
Some(model_audit_from_manifest(model)),
|
Some(model_audit_from_manifest(model)),
|
||||||
)
|
)
|
||||||
|
.with_memory_settings(&memory_cfg)
|
||||||
.emit(
|
.emit(
|
||||||
self.workspace_client(),
|
self.workspace_client(),
|
||||||
self.event_tx.as_ref(),
|
self.event_tx.as_ref(),
|
||||||
@@ -3780,6 +3808,7 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
|||||||
memory::audit::AuditTrigger::TokenThreshold,
|
memory::audit::AuditTrigger::TokenThreshold,
|
||||||
Some(model_audit_from_manifest(model)),
|
Some(model_audit_from_manifest(model)),
|
||||||
)
|
)
|
||||||
|
.with_memory_settings(&memory_cfg)
|
||||||
.emit(
|
.emit(
|
||||||
self.workspace_client(),
|
self.workspace_client(),
|
||||||
self.event_tx.as_ref(),
|
self.event_tx.as_ref(),
|
||||||
@@ -3845,7 +3874,8 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
|||||||
memory::audit::AuditWorker::MemoryExtract,
|
memory::audit::AuditWorker::MemoryExtract,
|
||||||
memory::audit::AuditTrigger::TokenThreshold,
|
memory::audit::AuditTrigger::TokenThreshold,
|
||||||
Some(model_audit_from_manifest(model)),
|
Some(model_audit_from_manifest(model)),
|
||||||
);
|
)
|
||||||
|
.with_memory_settings(memory_cfg);
|
||||||
let event_tx = self.event_tx.as_ref();
|
let event_tx = self.event_tx.as_ref();
|
||||||
|
|
||||||
let pointer_snapshot = self
|
let pointer_snapshot = self
|
||||||
@@ -3997,11 +4027,11 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
|||||||
return Err(err);
|
return Err(err);
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
let memory_language = memory_language(memory_cfg);
|
let memory_language = memory_language(memory_cfg)?;
|
||||||
let extract_system_prompt = match self
|
let extract_system_prompt = match self
|
||||||
.prompts
|
.prompts
|
||||||
.load_full()
|
.load_full()
|
||||||
.memory_extract_system(memory_language)
|
.memory_extract_system(&memory_language)
|
||||||
{
|
{
|
||||||
Ok(prompt) => prompt,
|
Ok(prompt) => prompt,
|
||||||
Err(err) => {
|
Err(err) => {
|
||||||
@@ -4191,6 +4221,7 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
|||||||
memory::audit::AuditTrigger::StagingBacklog,
|
memory::audit::AuditTrigger::StagingBacklog,
|
||||||
Some(model_audit_from_manifest(model)),
|
Some(model_audit_from_manifest(model)),
|
||||||
)
|
)
|
||||||
|
.with_memory_settings(&memory_cfg)
|
||||||
.emit(
|
.emit(
|
||||||
self.workspace_client(),
|
self.workspace_client(),
|
||||||
self.event_tx.as_ref(),
|
self.event_tx.as_ref(),
|
||||||
@@ -4232,6 +4263,7 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
|||||||
memory::audit::AuditTrigger::StagingBacklog,
|
memory::audit::AuditTrigger::StagingBacklog,
|
||||||
Some(model_audit_from_manifest(model)),
|
Some(model_audit_from_manifest(model)),
|
||||||
)
|
)
|
||||||
|
.with_memory_settings(&memory_cfg)
|
||||||
.emit(
|
.emit(
|
||||||
self.workspace_client(),
|
self.workspace_client(),
|
||||||
self.event_tx.as_ref(),
|
self.event_tx.as_ref(),
|
||||||
@@ -4312,6 +4344,7 @@ struct WorkerAuditBase {
|
|||||||
worker: memory::audit::AuditWorker,
|
worker: memory::audit::AuditWorker,
|
||||||
trigger: memory::audit::AuditTrigger,
|
trigger: memory::audit::AuditTrigger,
|
||||||
model: Option<memory::audit::ModelAudit>,
|
model: Option<memory::audit::ModelAudit>,
|
||||||
|
memory_settings: Option<memory::audit::MemorySettingsAudit>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl WorkerAuditBase {
|
impl WorkerAuditBase {
|
||||||
@@ -4325,9 +4358,22 @@ impl WorkerAuditBase {
|
|||||||
worker,
|
worker,
|
||||||
trigger,
|
trigger,
|
||||||
model,
|
model,
|
||||||
|
memory_settings: None,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn with_memory_settings(mut self, memory_config: &manifest::MemoryConfig) -> Self {
|
||||||
|
self.memory_settings =
|
||||||
|
memory_config
|
||||||
|
.workspace_settings()
|
||||||
|
.map(|snapshot| memory::audit::MemorySettingsAudit {
|
||||||
|
workspace_id: snapshot.workspace_id,
|
||||||
|
settings_revision: snapshot.settings_revision,
|
||||||
|
language: snapshot.language,
|
||||||
|
});
|
||||||
|
self
|
||||||
|
}
|
||||||
|
|
||||||
async fn emit(
|
async fn emit(
|
||||||
&self,
|
&self,
|
||||||
workspace_client: &dyn WorkspaceClient,
|
workspace_client: &dyn WorkspaceClient,
|
||||||
@@ -4345,6 +4391,7 @@ impl WorkerAuditBase {
|
|||||||
status,
|
status,
|
||||||
trigger: self.trigger,
|
trigger: self.trigger,
|
||||||
reason: reason.clone(),
|
reason: reason.clone(),
|
||||||
|
memory_settings: self.memory_settings.clone(),
|
||||||
model: self.model.clone(),
|
model: self.model.clone(),
|
||||||
usage,
|
usage,
|
||||||
extract,
|
extract,
|
||||||
@@ -4391,12 +4438,14 @@ fn is_idle_consolidation_skip_reason(reason: &str) -> bool {
|
|||||||
|| reason.starts_with("threshold_not_reached")
|
|| reason.starts_with("threshold_not_reached")
|
||||||
}
|
}
|
||||||
|
|
||||||
fn memory_language(cfg: &manifest::MemoryConfig) -> &str {
|
fn memory_language(cfg: &manifest::MemoryConfig) -> Result<String, WorkerError> {
|
||||||
cfg.language
|
cfg.workspace_settings()
|
||||||
.as_deref()
|
.map(|snapshot| snapshot.language)
|
||||||
.map(str::trim)
|
.ok_or_else(|| {
|
||||||
.filter(|language| !language.is_empty())
|
WorkerError::InvalidState(
|
||||||
.unwrap_or(manifest::defaults::MEMORY_LANGUAGE)
|
"Memory operation requires a bound Workspace Memory settings snapshot".to_string(),
|
||||||
|
)
|
||||||
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
fn worker_language(cfg: &manifest::EngineManifest) -> &str {
|
fn worker_language(cfg: &manifest::EngineManifest) -> &str {
|
||||||
@@ -4453,6 +4502,7 @@ where
|
|||||||
workspace_context: WorkerWorkspaceContext,
|
workspace_context: WorkerWorkspaceContext,
|
||||||
filesystem_authority: WorkerFilesystemAuthority,
|
filesystem_authority: WorkerFilesystemAuthority,
|
||||||
) -> Result<Self, WorkerError> {
|
) -> Result<Self, WorkerError> {
|
||||||
|
validate_workspace_memory_snapshot(&manifest.worker.name, &manifest, &workspace_context)?;
|
||||||
let common = prepare_worker_common_with_context(
|
let common = prepare_worker_common_with_context(
|
||||||
&manifest,
|
&manifest,
|
||||||
&loader,
|
&loader,
|
||||||
@@ -4654,6 +4704,7 @@ where
|
|||||||
workspace_context: WorkerWorkspaceContext,
|
workspace_context: WorkerWorkspaceContext,
|
||||||
filesystem_authority: WorkerFilesystemAuthority,
|
filesystem_authority: WorkerFilesystemAuthority,
|
||||||
) -> Result<Self, WorkerError> {
|
) -> Result<Self, WorkerError> {
|
||||||
|
validate_workspace_memory_snapshot(&manifest.worker.name, &manifest, &workspace_context)?;
|
||||||
let common = prepare_worker_common_with_context(
|
let common = prepare_worker_common_with_context(
|
||||||
&manifest,
|
&manifest,
|
||||||
&loader,
|
&loader,
|
||||||
@@ -4768,6 +4819,13 @@ where
|
|||||||
.ok_or_else(|| WorkerError::WorkerMetadataMissing {
|
.ok_or_else(|| WorkerError::WorkerMetadataMissing {
|
||||||
worker_name: worker_name.to_string(),
|
worker_name: worker_name.to_string(),
|
||||||
})?;
|
})?;
|
||||||
|
if workspace_context.workspace_id().is_some()
|
||||||
|
&& metadata.resolved_manifest_snapshot.is_none()
|
||||||
|
{
|
||||||
|
return Err(WorkerError::WorkerMetadataManifestSnapshotMissing {
|
||||||
|
worker_name: worker_name.to_string(),
|
||||||
|
});
|
||||||
|
}
|
||||||
let active = metadata
|
let active = metadata
|
||||||
.active
|
.active
|
||||||
.ok_or_else(|| WorkerError::WorkerMetadataInactive {
|
.ok_or_else(|| WorkerError::WorkerMetadataInactive {
|
||||||
@@ -4817,6 +4875,13 @@ where
|
|||||||
.ok_or_else(|| WorkerError::WorkerMetadataMissing {
|
.ok_or_else(|| WorkerError::WorkerMetadataMissing {
|
||||||
worker_name: worker_name.to_string(),
|
worker_name: worker_name.to_string(),
|
||||||
})?;
|
})?;
|
||||||
|
if workspace_context.workspace_id().is_some()
|
||||||
|
&& metadata.resolved_manifest_snapshot.is_none()
|
||||||
|
{
|
||||||
|
return Err(WorkerError::WorkerMetadataManifestSnapshotMissing {
|
||||||
|
worker_name: worker_name.to_string(),
|
||||||
|
});
|
||||||
|
}
|
||||||
let active = metadata
|
let active = metadata
|
||||||
.active
|
.active
|
||||||
.ok_or_else(|| WorkerError::WorkerMetadataInactive {
|
.ok_or_else(|| WorkerError::WorkerMetadataInactive {
|
||||||
@@ -5156,8 +5221,48 @@ fn worker_metadata_for_manifest(
|
|||||||
metadata
|
metadata
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn validate_workspace_memory_snapshot(
|
||||||
|
worker_name: &str,
|
||||||
|
manifest: &WorkerManifest,
|
||||||
|
workspace_context: &WorkerWorkspaceContext,
|
||||||
|
) -> Result<(), WorkerError> {
|
||||||
|
let Some(workspace_id) = workspace_context.workspace_id() else {
|
||||||
|
return Ok(());
|
||||||
|
};
|
||||||
|
let snapshot = manifest
|
||||||
|
.memory
|
||||||
|
.as_ref()
|
||||||
|
.and_then(manifest::MemoryConfig::workspace_settings)
|
||||||
|
.ok_or_else(|| {
|
||||||
|
WorkerError::InvalidState(format!(
|
||||||
|
"Workspace Worker {worker_name} has no complete persisted Memory settings snapshot"
|
||||||
|
))
|
||||||
|
})?;
|
||||||
|
if snapshot.workspace_id != workspace_id.as_str() {
|
||||||
|
return Err(WorkerError::InvalidState(format!(
|
||||||
|
"Workspace Worker {worker_name} Memory settings belong to {} instead of {}",
|
||||||
|
snapshot.workspace_id,
|
||||||
|
workspace_id.as_str()
|
||||||
|
)));
|
||||||
|
}
|
||||||
|
if snapshot.settings_revision == 0
|
||||||
|
|| !manifest::is_normalized_workspace_memory_language(&snapshot.language)
|
||||||
|
{
|
||||||
|
return Err(WorkerError::InvalidState(format!(
|
||||||
|
"Workspace Worker {worker_name} has corrupt Memory settings snapshot metadata"
|
||||||
|
)));
|
||||||
|
}
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
fn should_persist_resolved_manifest_snapshot(manifest: &WorkerManifest) -> bool {
|
fn should_persist_resolved_manifest_snapshot(manifest: &WorkerManifest) -> bool {
|
||||||
manifest.profile.is_some() || manifest.plugins.has_resolved_plan()
|
manifest.profile.is_some()
|
||||||
|
|| manifest.plugins.has_resolved_plan()
|
||||||
|
|| manifest
|
||||||
|
.memory
|
||||||
|
.as_ref()
|
||||||
|
.and_then(manifest::MemoryConfig::workspace_settings)
|
||||||
|
.is_some()
|
||||||
}
|
}
|
||||||
|
|
||||||
fn restore_manifest_from_worker_metadata_snapshot(
|
fn restore_manifest_from_worker_metadata_snapshot(
|
||||||
@@ -5678,7 +5783,9 @@ pub enum WorkerError {
|
|||||||
session_id: SessionId,
|
session_id: SessionId,
|
||||||
},
|
},
|
||||||
|
|
||||||
#[error("worker metadata for {worker_name} does not include a resolved manifest snapshot")]
|
#[error(
|
||||||
|
"worker metadata for {worker_name} does not include a trusted resolved manifest snapshot; a replacement Worker is required"
|
||||||
|
)]
|
||||||
WorkerMetadataManifestSnapshotMissing { worker_name: String },
|
WorkerMetadataManifestSnapshotMissing { worker_name: String },
|
||||||
|
|
||||||
#[error(
|
#[error(
|
||||||
@@ -6179,6 +6286,92 @@ permission = "write"
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn workspace_memory_settings_snapshot_is_persisted_and_scope_checked() {
|
||||||
|
let mut manifest = WorkerManifest::from_toml(
|
||||||
|
r#"
|
||||||
|
[worker]
|
||||||
|
name = "memory-snapshot"
|
||||||
|
|
||||||
|
[model]
|
||||||
|
scheme = "anthropic"
|
||||||
|
model_id = "claude-sonnet-4-20250514"
|
||||||
|
|
||||||
|
[engine]
|
||||||
|
instruction = "default"
|
||||||
|
|
||||||
|
[[scope.allow]]
|
||||||
|
target = "/workspace"
|
||||||
|
permission = "read"
|
||||||
|
|
||||||
|
[[delegation_scope.allow]]
|
||||||
|
target = "/workspace"
|
||||||
|
permission = "read"
|
||||||
|
"#,
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
manifest.memory = Some(manifest::MemoryConfig::default());
|
||||||
|
manifest.memory.as_mut().unwrap().bind_workspace_settings(
|
||||||
|
&manifest::WorkspaceMemorySettingsSnapshot {
|
||||||
|
workspace_id: "workspace-a".to_string(),
|
||||||
|
settings_revision: 7,
|
||||||
|
language: "Japanese".to_string(),
|
||||||
|
},
|
||||||
|
);
|
||||||
|
|
||||||
|
let metadata = worker_metadata_for_manifest(&manifest, None, None, None);
|
||||||
|
let restored: WorkerManifest = serde_json::from_value(
|
||||||
|
metadata
|
||||||
|
.resolved_manifest_snapshot
|
||||||
|
.expect("Memory settings require a resolved manifest snapshot"),
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(
|
||||||
|
restored.memory.unwrap().workspace_settings(),
|
||||||
|
Some(manifest::WorkspaceMemorySettingsSnapshot {
|
||||||
|
workspace_id: "workspace-a".to_string(),
|
||||||
|
settings_revision: 7,
|
||||||
|
language: "Japanese".to_string(),
|
||||||
|
})
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
validate_workspace_memory_snapshot(
|
||||||
|
"memory-snapshot",
|
||||||
|
&manifest,
|
||||||
|
&WorkerWorkspaceContext::unavailable(
|
||||||
|
Some(WorkspaceId::new("workspace-a").unwrap()),
|
||||||
|
"test",
|
||||||
|
)
|
||||||
|
)
|
||||||
|
.is_ok()
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
validate_workspace_memory_snapshot(
|
||||||
|
"memory-snapshot",
|
||||||
|
&manifest,
|
||||||
|
&WorkerWorkspaceContext::unavailable(
|
||||||
|
Some(WorkspaceId::new("workspace-b").unwrap()),
|
||||||
|
"test",
|
||||||
|
)
|
||||||
|
)
|
||||||
|
.is_err()
|
||||||
|
);
|
||||||
|
|
||||||
|
let mut missing = manifest.clone();
|
||||||
|
missing.memory.as_mut().unwrap().settings_revision = None;
|
||||||
|
assert!(
|
||||||
|
validate_workspace_memory_snapshot(
|
||||||
|
"memory-snapshot",
|
||||||
|
&missing,
|
||||||
|
&WorkerWorkspaceContext::unavailable(
|
||||||
|
Some(WorkspaceId::new("workspace-a").unwrap()),
|
||||||
|
"test",
|
||||||
|
)
|
||||||
|
)
|
||||||
|
.is_err()
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn plugin_resolved_manifest_snapshot_is_persisted_without_profile() {
|
fn plugin_resolved_manifest_snapshot_is_persisted_without_profile() {
|
||||||
let mut manifest = WorkerManifest::from_toml(
|
let mut manifest = WorkerManifest::from_toml(
|
||||||
@@ -7064,6 +7257,52 @@ mod build_summary_prompt_tests {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn memory_consolidation_prompt_uses_bound_workspace_language() {
|
||||||
|
let dir = tempfile::tempdir().unwrap();
|
||||||
|
let cwd = dir.path().join("workspace");
|
||||||
|
std::fs::create_dir_all(&cwd).unwrap();
|
||||||
|
let store = session_store::FsStore::new(dir.path().join("sessions")).unwrap();
|
||||||
|
let mut manifest = minimal_manifest();
|
||||||
|
manifest.profile = Some(manifest::ProfileManifestSnapshot {
|
||||||
|
source: manifest::ProfileSource::Registry {
|
||||||
|
source: manifest::ProfileRegistrySource::Builtin,
|
||||||
|
name: "memory-consolidation".to_string(),
|
||||||
|
path: None,
|
||||||
|
provenance: None,
|
||||||
|
},
|
||||||
|
profile: None,
|
||||||
|
});
|
||||||
|
let mut memory = manifest::MemoryConfig::default();
|
||||||
|
memory.bind_workspace_settings(&manifest::WorkspaceMemorySettingsSnapshot {
|
||||||
|
workspace_id: "workspace-test".to_string(),
|
||||||
|
settings_revision: 3,
|
||||||
|
language: "Japanese".to_string(),
|
||||||
|
});
|
||||||
|
manifest.memory = Some(memory);
|
||||||
|
let mut worker = Worker::new(
|
||||||
|
manifest,
|
||||||
|
Engine::new(NoopClient),
|
||||||
|
store,
|
||||||
|
WorkerWorkspaceContext::no_workspace(),
|
||||||
|
WorkerFilesystemAuthority::local(cwd.clone(), cwd.clone()),
|
||||||
|
Scope::writable(&cwd).unwrap(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
worker.set_system_prompt_template(
|
||||||
|
SystemPromptTemplate::parse(
|
||||||
|
"default",
|
||||||
|
crate::prompt::source::PromptCatalogSource::builtins_only(),
|
||||||
|
)
|
||||||
|
.unwrap(),
|
||||||
|
);
|
||||||
|
worker.ensure_system_prompt_materialized().await.unwrap();
|
||||||
|
let prompt = worker.engine().get_system_prompt().unwrap();
|
||||||
|
assert!(prompt.contains("`language`: `Japanese`"));
|
||||||
|
assert!(!prompt.contains("`language`: `English`"));
|
||||||
|
}
|
||||||
|
|
||||||
async fn render_system_prompt_with_summary(
|
async fn render_system_prompt_with_summary(
|
||||||
summary_doc: Option<&str>,
|
summary_doc: Option<&str>,
|
||||||
memory_config: Option<manifest::MemoryConfig>,
|
memory_config: Option<manifest::MemoryConfig>,
|
||||||
@@ -7358,6 +7597,9 @@ mod build_summary_prompt_tests {
|
|||||||
let mut manifest = minimal_manifest();
|
let mut manifest = minimal_manifest();
|
||||||
manifest.memory = Some(manifest::MemoryConfig {
|
manifest.memory = Some(manifest::MemoryConfig {
|
||||||
extract_threshold: Some(1),
|
extract_threshold: Some(1),
|
||||||
|
workspace_id: Some("workspace-test".to_string()),
|
||||||
|
settings_revision: Some(1),
|
||||||
|
language: Some("English".to_string()),
|
||||||
..Default::default()
|
..Default::default()
|
||||||
});
|
});
|
||||||
let memory_config = manifest.memory.clone().unwrap();
|
let memory_config = manifest.memory.clone().unwrap();
|
||||||
@@ -7454,6 +7696,14 @@ mod build_summary_prompt_tests {
|
|||||||
assert_eq!(audits.len(), 2);
|
assert_eq!(audits.len(), 2);
|
||||||
assert_eq!(audits[0].run_id, audits[1].run_id);
|
assert_eq!(audits[0].run_id, audits[1].run_id);
|
||||||
assert_eq!(audits[0].worker, memory::audit::AuditWorker::MemoryExtract);
|
assert_eq!(audits[0].worker, memory::audit::AuditWorker::MemoryExtract);
|
||||||
|
assert!(audits.iter().all(|audit| {
|
||||||
|
audit.memory_settings
|
||||||
|
== Some(memory::audit::MemorySettingsAudit {
|
||||||
|
workspace_id: "workspace-test".to_string(),
|
||||||
|
settings_revision: 1,
|
||||||
|
language: "English".to_string(),
|
||||||
|
})
|
||||||
|
}));
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
audits.iter().map(|audit| audit.status).collect::<Vec<_>>(),
|
audits.iter().map(|audit| audit.status).collect::<Vec<_>>(),
|
||||||
vec![
|
vec![
|
||||||
|
|||||||
@@ -169,6 +169,23 @@ pub struct WorkerRestoreResponse {
|
|||||||
pub result: WorkerRestoreResult,
|
pub result: WorkerRestoreResult,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Workspace-owned Memory settings returned by the shared Server API.
|
||||||
|
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||||
|
#[serde(deny_unknown_fields)]
|
||||||
|
pub struct WorkspaceMemorySettings {
|
||||||
|
pub workspace_id: String,
|
||||||
|
pub settings_revision: u64,
|
||||||
|
pub language: String,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Compare-and-swap update for Workspace-owned Memory settings.
|
||||||
|
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||||
|
#[serde(deny_unknown_fields)]
|
||||||
|
pub struct UpdateWorkspaceMemorySettingsRequest {
|
||||||
|
pub expected_revision: u64,
|
||||||
|
pub language: String,
|
||||||
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::*;
|
use super::*;
|
||||||
|
|||||||
@@ -518,6 +518,9 @@ pub struct WorkerSpawnRequest {
|
|||||||
pub resolved_config_bundle: Option<ConfigBundle>,
|
pub resolved_config_bundle: Option<ConfigBundle>,
|
||||||
#[serde(skip, default)]
|
#[serde(skip, default)]
|
||||||
pub resolved_workspace_api: Option<WorkspaceApiRef>,
|
pub resolved_workspace_api: Option<WorkspaceApiRef>,
|
||||||
|
/// Backend-authored immutable Workspace Memory settings snapshot.
|
||||||
|
#[serde(skip, default)]
|
||||||
|
pub resolved_memory_settings: Option<manifest::WorkspaceMemorySettingsSnapshot>,
|
||||||
/// Backend-owned feature enablement; client input cannot set it.
|
/// Backend-owned feature enablement; client input cannot set it.
|
||||||
#[serde(skip, default)]
|
#[serde(skip, default)]
|
||||||
pub resolved_worker_observation_enabled: bool,
|
pub resolved_worker_observation_enabled: bool,
|
||||||
@@ -2184,6 +2187,7 @@ impl WorkspaceWorkerRuntime for EmbeddedWorkerRuntime {
|
|||||||
worker_observation_enabled: request.resolved_worker_observation_enabled,
|
worker_observation_enabled: request.resolved_worker_observation_enabled,
|
||||||
worker_observation_grants: request.resolved_worker_observation_grants.clone(),
|
worker_observation_grants: request.resolved_worker_observation_grants.clone(),
|
||||||
workspace_api: Some(workspace_api),
|
workspace_api: Some(workspace_api),
|
||||||
|
memory_settings: request.resolved_memory_settings.clone(),
|
||||||
};
|
};
|
||||||
let workspace_scope = RuntimeWorkspaceScope::new(workspace_id, "embedded-backend");
|
let workspace_scope = RuntimeWorkspaceScope::new(workspace_id, "embedded-backend");
|
||||||
match self
|
match self
|
||||||
@@ -3331,6 +3335,7 @@ impl WorkspaceWorkerRuntime for RemoteWorkerRuntime {
|
|||||||
worker_observation_enabled: request.resolved_worker_observation_enabled,
|
worker_observation_enabled: request.resolved_worker_observation_enabled,
|
||||||
worker_observation_grants: request.resolved_worker_observation_grants.clone(),
|
worker_observation_grants: request.resolved_worker_observation_grants.clone(),
|
||||||
workspace_api: Some(workspace_api),
|
workspace_api: Some(workspace_api),
|
||||||
|
memory_settings: request.resolved_memory_settings.clone(),
|
||||||
};
|
};
|
||||||
match self.post_json::<_, RuntimeHttpWorkerResponse>("/v1/workers", &create) {
|
match self.post_json::<_, RuntimeHttpWorkerResponse>("/v1/workers", &create) {
|
||||||
Ok(response) => WorkerSpawnResult {
|
Ok(response) => WorkerSpawnResult {
|
||||||
@@ -4408,6 +4413,14 @@ mod tests {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn test_memory_settings() -> manifest::WorkspaceMemorySettingsSnapshot {
|
||||||
|
manifest::WorkspaceMemorySettingsSnapshot {
|
||||||
|
workspace_id: "workspace-test".to_string(),
|
||||||
|
settings_revision: 1,
|
||||||
|
language: "English".to_string(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn worker_summary_keeps_flat_wire_identity_while_using_structured_internal_identity() {
|
fn worker_summary_keeps_flat_wire_identity_while_using_structured_internal_identity() {
|
||||||
let summary = placeholder_worker("placeholder");
|
let summary = placeholder_worker("placeholder");
|
||||||
@@ -5007,6 +5020,7 @@ mod tests {
|
|||||||
resolved_worker_observation_grants: Vec::new(),
|
resolved_worker_observation_grants: Vec::new(),
|
||||||
resolved_control_operation: None,
|
resolved_control_operation: None,
|
||||||
resolved_workspace_api: Some(test_workspace_api()),
|
resolved_workspace_api: Some(test_workspace_api()),
|
||||||
|
resolved_memory_settings: Some(test_memory_settings()),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -5246,6 +5260,7 @@ mod tests {
|
|||||||
resolved_worker_observation_grants: Vec::new(),
|
resolved_worker_observation_grants: Vec::new(),
|
||||||
resolved_control_operation: None,
|
resolved_control_operation: None,
|
||||||
resolved_workspace_api: Some(test_workspace_api()),
|
resolved_workspace_api: Some(test_workspace_api()),
|
||||||
|
resolved_memory_settings: Some(test_memory_settings()),
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
.unwrap();
|
.unwrap();
|
||||||
@@ -5345,6 +5360,7 @@ mod tests {
|
|||||||
resolved_worker_observation_grants: Vec::new(),
|
resolved_worker_observation_grants: Vec::new(),
|
||||||
resolved_control_operation: None,
|
resolved_control_operation: None,
|
||||||
resolved_workspace_api: Some(test_workspace_api()),
|
resolved_workspace_api: Some(test_workspace_api()),
|
||||||
|
resolved_memory_settings: Some(test_memory_settings()),
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
.unwrap();
|
.unwrap();
|
||||||
@@ -5383,6 +5399,7 @@ mod tests {
|
|||||||
resolved_worker_observation_grants: Vec::new(),
|
resolved_worker_observation_grants: Vec::new(),
|
||||||
resolved_control_operation: None,
|
resolved_control_operation: None,
|
||||||
resolved_workspace_api: Some(test_workspace_api()),
|
resolved_workspace_api: Some(test_workspace_api()),
|
||||||
|
resolved_memory_settings: Some(test_memory_settings()),
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
.unwrap();
|
.unwrap();
|
||||||
|
|||||||
@@ -86,6 +86,11 @@ fn create_request(name: &str) -> CreateWorkerRequest {
|
|||||||
worker_observation_enabled: false,
|
worker_observation_enabled: false,
|
||||||
worker_observation_grants: Vec::new(),
|
worker_observation_grants: Vec::new(),
|
||||||
workspace_api: None,
|
workspace_api: None,
|
||||||
|
memory_settings: Some(manifest::WorkspaceMemorySettingsSnapshot {
|
||||||
|
workspace_id: "local".to_string(),
|
||||||
|
settings_revision: 1,
|
||||||
|
language: "English".to_string(),
|
||||||
|
}),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -1182,8 +1182,11 @@ impl WorkspaceApi {
|
|||||||
&now_registry_timestamp(),
|
&now_registry_timestamp(),
|
||||||
)?;
|
)?;
|
||||||
}
|
}
|
||||||
let create_fingerprint = worker_spawn_create_fingerprint(&request)
|
let request_fingerprint = worker_spawn_create_fingerprint(&request)
|
||||||
.map_err(|message| Error::Config(message.to_string()))?;
|
.map_err(|message| Error::Config(message.to_string()))?;
|
||||||
|
let current_memory_settings = self
|
||||||
|
.config_store
|
||||||
|
.get_workspace_memory_settings(&self.config.workspace_id)?;
|
||||||
let allocation_key = request
|
let allocation_key = request
|
||||||
.resolved_control_operation
|
.resolved_control_operation
|
||||||
.as_ref()
|
.as_ref()
|
||||||
@@ -1195,22 +1198,25 @@ impl WorkspaceApi {
|
|||||||
.map(|assignment| assignment.operation_id.clone())
|
.map(|assignment| assignment.operation_id.clone())
|
||||||
})
|
})
|
||||||
.unwrap_or_else(|| format!("manual:{}", WorkerId::now_v7()));
|
.unwrap_or_else(|| format!("manual:{}", WorkerId::now_v7()));
|
||||||
let worker_id = self
|
let reservation = self
|
||||||
.config_store
|
.config_store
|
||||||
.reserve_worker_create(
|
.reserve_worker_create(
|
||||||
&self.config.workspace_id,
|
&self.config.workspace_id,
|
||||||
runtime_id,
|
runtime_id,
|
||||||
&allocation_key,
|
&allocation_key,
|
||||||
&create_fingerprint,
|
&request_fingerprint,
|
||||||
|
¤t_memory_settings,
|
||||||
)
|
)
|
||||||
.map_err(|error| Error::RuntimeOperationFailed {
|
.map_err(|error| Error::RuntimeOperationFailed {
|
||||||
runtime_id: runtime_id.to_string(),
|
runtime_id: runtime_id.to_string(),
|
||||||
code: "workspace_worker_allocation_conflict".to_string(),
|
code: "workspace_worker_allocation_conflict".to_string(),
|
||||||
message: error.to_string(),
|
message: error.to_string(),
|
||||||
})?;
|
})?;
|
||||||
|
let worker_id = reservation.worker_id;
|
||||||
|
request.resolved_memory_settings = Some(reservation.memory_settings);
|
||||||
let create_binding = WorkerCreateBinding {
|
let create_binding = WorkerCreateBinding {
|
||||||
worker_id,
|
worker_id,
|
||||||
create_fingerprint,
|
create_fingerprint: reservation.create_fingerprint,
|
||||||
};
|
};
|
||||||
let result = match self
|
let result = match self
|
||||||
.runtime
|
.runtime
|
||||||
@@ -1601,6 +1607,11 @@ pub fn build_router(api: WorkspaceApi) -> Router {
|
|||||||
"/api/w/{workspace_id}/settings/workspace",
|
"/api/w/{workspace_id}/settings/workspace",
|
||||||
get(scoped_get_workspace_settings).put(scoped_update_workspace_settings),
|
get(scoped_get_workspace_settings).put(scoped_update_workspace_settings),
|
||||||
)
|
)
|
||||||
|
.route(
|
||||||
|
"/api/w/{workspace_id}/settings/memory",
|
||||||
|
get(scoped_get_workspace_memory_settings)
|
||||||
|
.put(scoped_update_workspace_memory_settings),
|
||||||
|
)
|
||||||
.route(
|
.route(
|
||||||
"/api/w/{workspace_id}/config/source-tree",
|
"/api/w/{workspace_id}/config/source-tree",
|
||||||
get(scoped_get_workspace_config_tree),
|
get(scoped_get_workspace_config_tree),
|
||||||
@@ -2998,6 +3009,43 @@ struct WorkspaceConfigTreeResponse {
|
|||||||
projection_digest: String,
|
projection_digest: String,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn scoped_get_workspace_memory_settings(
|
||||||
|
State(api): State<WorkspaceApi>,
|
||||||
|
AxumPath(workspace_id): AxumPath<String>,
|
||||||
|
) -> ApiResult<Json<workspace_api::WorkspaceMemorySettings>> {
|
||||||
|
validate_workspace_scope(&api, &workspace_id)?;
|
||||||
|
let settings = api
|
||||||
|
.config_store
|
||||||
|
.get_workspace_memory_settings(&workspace_id)
|
||||||
|
.map_err(ApiError::from)?;
|
||||||
|
Ok(Json(workspace_api::WorkspaceMemorySettings {
|
||||||
|
workspace_id: settings.workspace_id,
|
||||||
|
settings_revision: settings.settings_revision,
|
||||||
|
language: settings.language,
|
||||||
|
}))
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn scoped_update_workspace_memory_settings(
|
||||||
|
State(api): State<WorkspaceApi>,
|
||||||
|
AxumPath(workspace_id): AxumPath<String>,
|
||||||
|
Json(request): Json<workspace_api::UpdateWorkspaceMemorySettingsRequest>,
|
||||||
|
) -> ApiResult<Json<workspace_api::WorkspaceMemorySettings>> {
|
||||||
|
validate_workspace_scope(&api, &workspace_id)?;
|
||||||
|
let settings = api
|
||||||
|
.config_store
|
||||||
|
.update_workspace_memory_settings(
|
||||||
|
&workspace_id,
|
||||||
|
request.expected_revision,
|
||||||
|
&request.language,
|
||||||
|
)
|
||||||
|
.map_err(ApiError::from)?;
|
||||||
|
Ok(Json(workspace_api::WorkspaceMemorySettings {
|
||||||
|
workspace_id: settings.workspace_id,
|
||||||
|
settings_revision: settings.settings_revision,
|
||||||
|
language: settings.language,
|
||||||
|
}))
|
||||||
|
}
|
||||||
|
|
||||||
async fn scoped_get_workspace_config_tree(
|
async fn scoped_get_workspace_config_tree(
|
||||||
State(api): State<WorkspaceApi>,
|
State(api): State<WorkspaceApi>,
|
||||||
AxumPath(path): AxumPath<ScopedWorkspacePath>,
|
AxumPath(path): AxumPath<ScopedWorkspacePath>,
|
||||||
@@ -6126,6 +6174,7 @@ fn start_memory_staging_consolidation(
|
|||||||
resolved_worker_observation_grants: Vec::new(),
|
resolved_worker_observation_grants: Vec::new(),
|
||||||
resolved_control_operation: None,
|
resolved_control_operation: None,
|
||||||
resolved_workspace_api: None,
|
resolved_workspace_api: None,
|
||||||
|
resolved_memory_settings: None,
|
||||||
},
|
},
|
||||||
)?;
|
)?;
|
||||||
if result.state != WorkerOperationState::Accepted {
|
if result.state != WorkerOperationState::Accepted {
|
||||||
@@ -7192,6 +7241,7 @@ async fn scoped_start_workspace_orchestrator(
|
|||||||
resolved_worker_observation_grants: Vec::new(),
|
resolved_worker_observation_grants: Vec::new(),
|
||||||
resolved_control_operation: None,
|
resolved_control_operation: None,
|
||||||
resolved_workspace_api: None,
|
resolved_workspace_api: None,
|
||||||
|
resolved_memory_settings: None,
|
||||||
},
|
},
|
||||||
)?;
|
)?;
|
||||||
if result.state != WorkerOperationState::Accepted || result.worker.is_none() {
|
if result.state != WorkerOperationState::Accepted || result.worker.is_none() {
|
||||||
@@ -9736,6 +9786,7 @@ async fn create_workspace_worker(
|
|||||||
resolved_worker_observation_grants: Vec::new(),
|
resolved_worker_observation_grants: Vec::new(),
|
||||||
resolved_control_operation,
|
resolved_control_operation,
|
||||||
resolved_workspace_api: None,
|
resolved_workspace_api: None,
|
||||||
|
resolved_memory_settings: None,
|
||||||
};
|
};
|
||||||
validate_ticket_assignment_spawn(&api, &runtime_id, &request)?;
|
validate_ticket_assignment_spawn(&api, &runtime_id, &request)?;
|
||||||
let assignment = request.ticket_assignment.clone();
|
let assignment = request.ticket_assignment.clone();
|
||||||
@@ -13164,6 +13215,14 @@ mod tests {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn test_worker_memory_settings() -> manifest::WorkspaceMemorySettingsSnapshot {
|
||||||
|
manifest::WorkspaceMemorySettingsSnapshot {
|
||||||
|
workspace_id: TEST_WORKSPACE_ID.to_string(),
|
||||||
|
settings_revision: 1,
|
||||||
|
language: "English".to_string(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn ticket_api_errors_preserve_http_status() {
|
fn ticket_api_errors_preserve_http_status() {
|
||||||
let not_found = ApiError::from(Error::Ticket(ticket::TicketError::NotFound(
|
let not_found = ApiError::from(Error::Ticket(ticket::TicketError::NotFound(
|
||||||
@@ -13432,6 +13491,7 @@ mod tests {
|
|||||||
resolved_worker_observation_grants: Vec::new(),
|
resolved_worker_observation_grants: Vec::new(),
|
||||||
resolved_control_operation: None,
|
resolved_control_operation: None,
|
||||||
resolved_workspace_api: None,
|
resolved_workspace_api: None,
|
||||||
|
resolved_memory_settings: None,
|
||||||
};
|
};
|
||||||
assert!(
|
assert!(
|
||||||
api.validate_worker_spawn_repository_scope(&workdir_flow_launch)
|
api.validate_worker_spawn_repository_scope(&workdir_flow_launch)
|
||||||
@@ -13604,6 +13664,7 @@ mod tests {
|
|||||||
resolved_worker_observation_grants: Vec::new(),
|
resolved_worker_observation_grants: Vec::new(),
|
||||||
resolved_control_operation: None,
|
resolved_control_operation: None,
|
||||||
resolved_workspace_api: None,
|
resolved_workspace_api: None,
|
||||||
|
resolved_memory_settings: None,
|
||||||
};
|
};
|
||||||
|
|
||||||
assert!(
|
assert!(
|
||||||
@@ -15069,6 +15130,7 @@ mod tests {
|
|||||||
resolved_workspace_api: Some(test_worker_workspace_api(
|
resolved_workspace_api: Some(test_worker_workspace_api(
|
||||||
EMBEDDED_WORKER_RUNTIME_ID,
|
EMBEDDED_WORKER_RUNTIME_ID,
|
||||||
)),
|
)),
|
||||||
|
resolved_memory_settings: Some(test_worker_memory_settings()),
|
||||||
resolved_control_operation: None,
|
resolved_control_operation: None,
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
@@ -15282,6 +15344,7 @@ mod tests {
|
|||||||
resolved_workspace_api: Some(test_worker_workspace_api(
|
resolved_workspace_api: Some(test_worker_workspace_api(
|
||||||
EMBEDDED_WORKER_RUNTIME_ID,
|
EMBEDDED_WORKER_RUNTIME_ID,
|
||||||
)),
|
)),
|
||||||
|
resolved_memory_settings: Some(test_worker_memory_settings()),
|
||||||
resolved_control_operation: None,
|
resolved_control_operation: None,
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
@@ -15492,6 +15555,7 @@ mod tests {
|
|||||||
resolved_worker_observation_enabled: false,
|
resolved_worker_observation_enabled: false,
|
||||||
resolved_worker_observation_grants: Vec::new(),
|
resolved_worker_observation_grants: Vec::new(),
|
||||||
resolved_workspace_api: Some(test_worker_workspace_api(EMBEDDED_WORKER_RUNTIME_ID)),
|
resolved_workspace_api: Some(test_worker_workspace_api(EMBEDDED_WORKER_RUNTIME_ID)),
|
||||||
|
resolved_memory_settings: Some(test_worker_memory_settings()),
|
||||||
resolved_control_operation: None,
|
resolved_control_operation: None,
|
||||||
};
|
};
|
||||||
let source_worker = api
|
let source_worker = api
|
||||||
@@ -15756,6 +15820,7 @@ mod tests {
|
|||||||
resolved_workspace_api: Some(test_worker_workspace_api(
|
resolved_workspace_api: Some(test_worker_workspace_api(
|
||||||
EMBEDDED_WORKER_RUNTIME_ID,
|
EMBEDDED_WORKER_RUNTIME_ID,
|
||||||
)),
|
)),
|
||||||
|
resolved_memory_settings: Some(test_worker_memory_settings()),
|
||||||
resolved_control_operation: None,
|
resolved_control_operation: None,
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
@@ -15881,6 +15946,7 @@ mod tests {
|
|||||||
resolved_worker_observation_grants: Vec::new(),
|
resolved_worker_observation_grants: Vec::new(),
|
||||||
resolved_control_operation: None,
|
resolved_control_operation: None,
|
||||||
resolved_workspace_api: None,
|
resolved_workspace_api: None,
|
||||||
|
resolved_memory_settings: None,
|
||||||
};
|
};
|
||||||
let Json(first) = scoped_create_runtime_worker(
|
let Json(first) = scoped_create_runtime_worker(
|
||||||
State(api.clone()),
|
State(api.clone()),
|
||||||
@@ -16045,22 +16111,28 @@ mod tests {
|
|||||||
TEST_CREATED_AT,
|
TEST_CREATED_AT,
|
||||||
)
|
)
|
||||||
.unwrap();
|
.unwrap();
|
||||||
let reserved_worker_id = api
|
let current_memory_settings = api
|
||||||
|
.config_store
|
||||||
|
.get_workspace_memory_settings(TEST_WORKSPACE_ID)
|
||||||
|
.unwrap();
|
||||||
|
let reservation = api
|
||||||
.config_store
|
.config_store
|
||||||
.reserve_worker_create(
|
.reserve_worker_create(
|
||||||
TEST_WORKSPACE_ID,
|
TEST_WORKSPACE_ID,
|
||||||
EMBEDDED_WORKER_RUNTIME_ID,
|
EMBEDDED_WORKER_RUNTIME_ID,
|
||||||
"pending-spawn-operation",
|
"pending-spawn-operation",
|
||||||
&pending_fingerprint,
|
&pending_fingerprint,
|
||||||
|
¤t_memory_settings,
|
||||||
)
|
)
|
||||||
.unwrap();
|
.unwrap();
|
||||||
|
pending_request.resolved_memory_settings = Some(reservation.memory_settings.clone());
|
||||||
let spawned_before_backend_failure = api
|
let spawned_before_backend_failure = api
|
||||||
.runtime
|
.runtime
|
||||||
.spawn_worker(
|
.spawn_worker(
|
||||||
EMBEDDED_WORKER_RUNTIME_ID,
|
EMBEDDED_WORKER_RUNTIME_ID,
|
||||||
WorkerCreateBinding {
|
WorkerCreateBinding {
|
||||||
worker_id: reserved_worker_id,
|
worker_id: reservation.worker_id,
|
||||||
create_fingerprint: pending_fingerprint.clone(),
|
create_fingerprint: reservation.create_fingerprint,
|
||||||
},
|
},
|
||||||
pending_request.clone(),
|
pending_request.clone(),
|
||||||
)
|
)
|
||||||
@@ -16139,6 +16211,7 @@ mod tests {
|
|||||||
resolved_worker_observation_grants: Vec::new(),
|
resolved_worker_observation_grants: Vec::new(),
|
||||||
resolved_control_operation: None,
|
resolved_control_operation: None,
|
||||||
resolved_workspace_api: None,
|
resolved_workspace_api: None,
|
||||||
|
resolved_memory_settings: None,
|
||||||
};
|
};
|
||||||
let Json(created) = scoped_create_runtime_worker(
|
let Json(created) = scoped_create_runtime_worker(
|
||||||
State(api.clone()),
|
State(api.clone()),
|
||||||
@@ -16504,6 +16577,32 @@ mod tests {
|
|||||||
test_api_with_recording_backend(workspace_root).await.0
|
test_api_with_recording_backend(workspace_root).await.0
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn memory_settings_handlers_reject_foreign_workspace_path_scope() {
|
||||||
|
let temp = tempfile::tempdir().unwrap();
|
||||||
|
let api = test_api(temp.path()).await;
|
||||||
|
assert!(
|
||||||
|
scoped_get_workspace_memory_settings(
|
||||||
|
State(api.clone()),
|
||||||
|
AxumPath("workspace-foreign".to_string()),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.is_err()
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
scoped_update_workspace_memory_settings(
|
||||||
|
State(api),
|
||||||
|
AxumPath("workspace-foreign".to_string()),
|
||||||
|
Json(workspace_api::UpdateWorkspaceMemorySettingsRequest {
|
||||||
|
expected_revision: 1,
|
||||||
|
language: "English".to_string(),
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.is_err()
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn destructive_worker_remove_rejects_browser_and_legacy_source_headers() {
|
async fn destructive_worker_remove_rejects_browser_and_legacy_source_headers() {
|
||||||
let headers = HeaderMap::new();
|
let headers = HeaderMap::new();
|
||||||
@@ -16699,6 +16798,7 @@ mod tests {
|
|||||||
resolved_worker_observation_enabled: false,
|
resolved_worker_observation_enabled: false,
|
||||||
resolved_worker_observation_grants: Vec::new(),
|
resolved_worker_observation_grants: Vec::new(),
|
||||||
resolved_workspace_api: None,
|
resolved_workspace_api: None,
|
||||||
|
resolved_memory_settings: None,
|
||||||
resolved_control_operation: None,
|
resolved_control_operation: None,
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
@@ -16755,6 +16855,7 @@ mod tests {
|
|||||||
resolved_worker_observation_enabled: false,
|
resolved_worker_observation_enabled: false,
|
||||||
resolved_worker_observation_grants: Vec::new(),
|
resolved_worker_observation_grants: Vec::new(),
|
||||||
resolved_workspace_api: None,
|
resolved_workspace_api: None,
|
||||||
|
resolved_memory_settings: None,
|
||||||
resolved_control_operation: None,
|
resolved_control_operation: None,
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
@@ -17680,6 +17781,7 @@ mod tests {
|
|||||||
worker_observation_enabled: false,
|
worker_observation_enabled: false,
|
||||||
worker_observation_grants: Vec::new(),
|
worker_observation_grants: Vec::new(),
|
||||||
workspace_api: None,
|
workspace_api: None,
|
||||||
|
memory_settings: None,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -18974,6 +19076,7 @@ mod tests {
|
|||||||
resolved_workspace_api: Some(test_worker_workspace_api(
|
resolved_workspace_api: Some(test_worker_workspace_api(
|
||||||
"embedded-worker-runtime",
|
"embedded-worker-runtime",
|
||||||
)),
|
)),
|
||||||
|
resolved_memory_settings: Some(test_worker_memory_settings()),
|
||||||
resolved_control_operation: None,
|
resolved_control_operation: None,
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
@@ -19493,6 +19596,7 @@ mod tests {
|
|||||||
resolved_worker_observation_grants: Vec::new(),
|
resolved_worker_observation_grants: Vec::new(),
|
||||||
resolved_control_operation: None,
|
resolved_control_operation: None,
|
||||||
resolved_workspace_api: None,
|
resolved_workspace_api: None,
|
||||||
|
resolved_memory_settings: None,
|
||||||
};
|
};
|
||||||
let spawned = api
|
let spawned = api
|
||||||
.spawn_workspace_worker(EMBEDDED_WORKER_RUNTIME_ID, spawn_request)
|
.spawn_workspace_worker(EMBEDDED_WORKER_RUNTIME_ID, spawn_request)
|
||||||
|
|||||||
@@ -8,6 +8,7 @@ use rusqlite::{
|
|||||||
Connection, OpenFlags, OptionalExtension, TransactionBehavior, backup::Backup, params,
|
Connection, OpenFlags, OptionalExtension, TransactionBehavior, backup::Backup, params,
|
||||||
};
|
};
|
||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
|
use sha2::{Digest, Sha256};
|
||||||
use uuid::Uuid;
|
use uuid::Uuid;
|
||||||
|
|
||||||
use worker_runtime::identity::{
|
use worker_runtime::identity::{
|
||||||
@@ -230,6 +231,11 @@ const MIGRATIONS: &[Migration] = &[
|
|||||||
name: "rename Workspace resource keys",
|
name: "rename Workspace resource keys",
|
||||||
apply: verify_workspace_resource_key_schema,
|
apply: verify_workspace_resource_key_schema,
|
||||||
},
|
},
|
||||||
|
Migration {
|
||||||
|
version: 42,
|
||||||
|
name: "create Workspace Memory settings authority",
|
||||||
|
apply: create_workspace_memory_settings_authority,
|
||||||
|
},
|
||||||
];
|
];
|
||||||
|
|
||||||
struct Migration {
|
struct Migration {
|
||||||
@@ -261,6 +267,22 @@ pub struct WorkspaceRecord {
|
|||||||
pub updated_at: String,
|
pub updated_at: String,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||||
|
pub struct WorkspaceMemorySettingsRecord {
|
||||||
|
pub workspace_id: String,
|
||||||
|
pub settings_revision: u64,
|
||||||
|
pub language: String,
|
||||||
|
pub created_at: String,
|
||||||
|
pub updated_at: String,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||||
|
pub struct WorkerCreateReservation {
|
||||||
|
pub worker_id: WorkerId,
|
||||||
|
pub create_fingerprint: String,
|
||||||
|
pub memory_settings: manifest::WorkspaceMemorySettingsSnapshot,
|
||||||
|
}
|
||||||
|
|
||||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||||
pub struct RepositoryRecord {
|
pub struct RepositoryRecord {
|
||||||
pub workspace_id: String,
|
pub workspace_id: String,
|
||||||
@@ -1140,23 +1162,125 @@ impl SqliteWorkspaceStore {
|
|||||||
f(&mut conn)
|
f(&mut conn)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub(crate) fn get_workspace_memory_settings(
|
||||||
|
&self,
|
||||||
|
workspace_id: &str,
|
||||||
|
) -> Result<WorkspaceMemorySettingsRecord> {
|
||||||
|
let record = self.with_conn(|conn| {
|
||||||
|
conn.query_row(
|
||||||
|
"SELECT workspace_id, settings_revision, language, created_at, updated_at \
|
||||||
|
FROM workspace_memory_settings WHERE workspace_id = ?1",
|
||||||
|
params![workspace_id],
|
||||||
|
|row| {
|
||||||
|
let revision = row.get::<_, i64>(1)?;
|
||||||
|
Ok(WorkspaceMemorySettingsRecord {
|
||||||
|
workspace_id: row.get(0)?,
|
||||||
|
settings_revision: revision
|
||||||
|
.try_into()
|
||||||
|
.map_err(|_| rusqlite::Error::IntegralValueOutOfRange(1, revision))?,
|
||||||
|
language: row.get(2)?,
|
||||||
|
created_at: row.get(3)?,
|
||||||
|
updated_at: row.get(4)?,
|
||||||
|
})
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.optional()?
|
||||||
|
.ok_or_else(|| Error::Store("Workspace Memory settings are missing".to_string()))
|
||||||
|
})?;
|
||||||
|
validate_workspace_memory_settings_record(&record, workspace_id)?;
|
||||||
|
Ok(record)
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) fn update_workspace_memory_settings(
|
||||||
|
&self,
|
||||||
|
workspace_id: &str,
|
||||||
|
expected_revision: u64,
|
||||||
|
language: &str,
|
||||||
|
) -> Result<WorkspaceMemorySettingsRecord> {
|
||||||
|
let language = normalize_workspace_memory_language(language)?;
|
||||||
|
self.with_conn_mut(|conn| {
|
||||||
|
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
|
||||||
|
let current = tx
|
||||||
|
.query_row(
|
||||||
|
"SELECT workspace_id, settings_revision, language, created_at, updated_at \
|
||||||
|
FROM workspace_memory_settings WHERE workspace_id = ?1",
|
||||||
|
params![workspace_id],
|
||||||
|
|row| {
|
||||||
|
let revision = row.get::<_, i64>(1)?;
|
||||||
|
Ok(WorkspaceMemorySettingsRecord {
|
||||||
|
workspace_id: row.get(0)?,
|
||||||
|
settings_revision: revision.try_into().map_err(|_| {
|
||||||
|
rusqlite::Error::IntegralValueOutOfRange(1, revision)
|
||||||
|
})?,
|
||||||
|
language: row.get(2)?,
|
||||||
|
created_at: row.get(3)?,
|
||||||
|
updated_at: row.get(4)?,
|
||||||
|
})
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.optional()?
|
||||||
|
.ok_or_else(|| Error::Store("Workspace Memory settings are missing".to_string()))?;
|
||||||
|
validate_workspace_memory_settings_record(¤t, workspace_id)?;
|
||||||
|
let current_revision = current.settings_revision;
|
||||||
|
if current_revision != expected_revision {
|
||||||
|
return Err(Error::WorkspaceConfigConflict(format!(
|
||||||
|
"Workspace Memory settings revision changed: expected {expected_revision}, current {current_revision}"
|
||||||
|
)));
|
||||||
|
}
|
||||||
|
if current.language == language {
|
||||||
|
tx.commit()?;
|
||||||
|
return Ok(current);
|
||||||
|
}
|
||||||
|
let next_revision = current_revision.checked_add(1).ok_or_else(|| {
|
||||||
|
Error::InvalidInput("Workspace Memory settings revision overflow".to_string())
|
||||||
|
})?;
|
||||||
|
let now = chrono::Utc::now().to_rfc3339();
|
||||||
|
tx.execute(
|
||||||
|
"UPDATE workspace_memory_settings \
|
||||||
|
SET settings_revision = ?2, language = ?3, updated_at = ?4 \
|
||||||
|
WHERE workspace_id = ?1",
|
||||||
|
params![workspace_id, next_revision as i64, language, now],
|
||||||
|
)?;
|
||||||
|
let record = tx.query_row(
|
||||||
|
"SELECT workspace_id, settings_revision, language, created_at, updated_at \
|
||||||
|
FROM workspace_memory_settings WHERE workspace_id = ?1",
|
||||||
|
params![workspace_id],
|
||||||
|
|row| {
|
||||||
|
let revision = row.get::<_, i64>(1)?;
|
||||||
|
Ok(WorkspaceMemorySettingsRecord {
|
||||||
|
workspace_id: row.get(0)?,
|
||||||
|
settings_revision: revision as u64,
|
||||||
|
language: row.get(2)?,
|
||||||
|
created_at: row.get(3)?,
|
||||||
|
updated_at: row.get(4)?,
|
||||||
|
})
|
||||||
|
},
|
||||||
|
)?;
|
||||||
|
tx.commit()?;
|
||||||
|
Ok(record)
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
pub(crate) fn reserve_worker_create(
|
pub(crate) fn reserve_worker_create(
|
||||||
&self,
|
&self,
|
||||||
workspace_id: &str,
|
workspace_id: &str,
|
||||||
runtime_id: &str,
|
runtime_id: &str,
|
||||||
allocation_key: &str,
|
allocation_key: &str,
|
||||||
create_fingerprint: &str,
|
request_fingerprint: &str,
|
||||||
) -> Result<WorkerId> {
|
current_memory_settings: &WorkspaceMemorySettingsRecord,
|
||||||
if allocation_key.trim().is_empty() || create_fingerprint.trim().is_empty() {
|
) -> Result<WorkerCreateReservation> {
|
||||||
|
if allocation_key.trim().is_empty() || request_fingerprint.trim().is_empty() {
|
||||||
return Err(Error::InvalidInput(
|
return Err(Error::InvalidInput(
|
||||||
"Worker create allocation key and fingerprint must be non-empty".to_string(),
|
"Worker create allocation key and fingerprint must be non-empty".to_string(),
|
||||||
));
|
));
|
||||||
}
|
}
|
||||||
|
validate_workspace_memory_settings_record(current_memory_settings, workspace_id)?;
|
||||||
self.with_conn_mut(|conn| {
|
self.with_conn_mut(|conn| {
|
||||||
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
|
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
|
||||||
let existing = tx
|
let existing = tx
|
||||||
.query_row(
|
.query_row(
|
||||||
"SELECT worker_id, runtime_id, create_fingerprint \
|
"SELECT worker_id, runtime_id, request_fingerprint, create_fingerprint, \
|
||||||
|
memory_settings_revision, memory_language \
|
||||||
FROM worker_create_reservations \
|
FROM worker_create_reservations \
|
||||||
WHERE workspace_id = ?1 AND allocation_key = ?2",
|
WHERE workspace_id = ?1 AND allocation_key = ?2",
|
||||||
params![workspace_id, allocation_key],
|
params![workspace_id, allocation_key],
|
||||||
@@ -1164,24 +1288,80 @@ impl SqliteWorkspaceStore {
|
|||||||
Ok((
|
Ok((
|
||||||
row.get::<_, String>(0)?,
|
row.get::<_, String>(0)?,
|
||||||
row.get::<_, String>(1)?,
|
row.get::<_, String>(1)?,
|
||||||
row.get::<_, String>(2)?,
|
row.get::<_, Option<String>>(2)?,
|
||||||
|
row.get::<_, String>(3)?,
|
||||||
|
row.get::<_, Option<i64>>(4)?,
|
||||||
|
row.get::<_, Option<String>>(5)?,
|
||||||
))
|
))
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
.optional()?;
|
.optional()?;
|
||||||
if let Some((worker_id, reserved_runtime_id, reserved_fingerprint)) = existing {
|
if let Some((worker_id, reserved_runtime_id, stored_request_fingerprint, create_fingerprint, revision, language)) = existing {
|
||||||
if reserved_runtime_id != runtime_id || reserved_fingerprint != create_fingerprint {
|
if reserved_runtime_id != runtime_id
|
||||||
|
|| stored_request_fingerprint.as_deref() != Some(request_fingerprint)
|
||||||
|
{
|
||||||
return Err(Error::InvalidInput(format!(
|
return Err(Error::InvalidInput(format!(
|
||||||
"Worker create allocation `{allocation_key}` was already used with different input"
|
"Worker create allocation {allocation_key} was already used with different input"
|
||||||
)));
|
)));
|
||||||
}
|
}
|
||||||
return worker_id.parse::<WorkerId>().map_err(|_| {
|
let revision = revision.ok_or_else(|| {
|
||||||
Error::Store(format!(
|
Error::InvalidInput(format!(
|
||||||
"Worker create allocation `{allocation_key}` has a non-UUIDv7 worker id"
|
"Worker create allocation {allocation_key} has no persisted Memory settings snapshot"
|
||||||
))
|
))
|
||||||
|
})?;
|
||||||
|
let language = language.ok_or_else(|| {
|
||||||
|
Error::InvalidInput(format!(
|
||||||
|
"Worker create allocation {allocation_key} has no persisted Memory language"
|
||||||
|
))
|
||||||
|
})?;
|
||||||
|
let worker_id = worker_id.parse::<WorkerId>().map_err(|_| {
|
||||||
|
Error::Store(format!(
|
||||||
|
"Worker create allocation {allocation_key} has a non-UUIDv7 worker id"
|
||||||
|
))
|
||||||
|
})?;
|
||||||
|
let snapshot = manifest::WorkspaceMemorySettingsSnapshot {
|
||||||
|
workspace_id: workspace_id.to_string(),
|
||||||
|
settings_revision: revision.try_into().map_err(|_| {
|
||||||
|
rusqlite::Error::IntegralValueOutOfRange(4, revision)
|
||||||
|
})?,
|
||||||
|
language,
|
||||||
|
};
|
||||||
|
validate_workspace_memory_settings_snapshot(&snapshot, workspace_id)?;
|
||||||
|
return Ok(WorkerCreateReservation {
|
||||||
|
worker_id,
|
||||||
|
create_fingerprint,
|
||||||
|
memory_settings: snapshot,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
let (authoritative_revision, authoritative_language) = tx
|
||||||
|
.query_row(
|
||||||
|
"SELECT settings_revision, language FROM workspace_memory_settings WHERE workspace_id = ?1",
|
||||||
|
params![workspace_id],
|
||||||
|
|row| Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1)?)),
|
||||||
|
)
|
||||||
|
.optional()?
|
||||||
|
.ok_or_else(|| Error::Store("Workspace Memory settings are missing".to_string()))?;
|
||||||
|
let authoritative_revision: u64 = authoritative_revision.try_into().map_err(|_| {
|
||||||
|
rusqlite::Error::IntegralValueOutOfRange(0, authoritative_revision)
|
||||||
|
})?;
|
||||||
|
if authoritative_revision != current_memory_settings.settings_revision
|
||||||
|
|| authoritative_language != current_memory_settings.language
|
||||||
|
{
|
||||||
|
return Err(Error::WorkspaceConfigConflict(
|
||||||
|
"Workspace Memory settings changed while the Worker create reservation was being accepted"
|
||||||
|
.to_string(),
|
||||||
|
));
|
||||||
|
}
|
||||||
|
|
||||||
|
let snapshot = manifest::WorkspaceMemorySettingsSnapshot {
|
||||||
|
workspace_id: workspace_id.to_string(),
|
||||||
|
settings_revision: authoritative_revision,
|
||||||
|
language: authoritative_language,
|
||||||
|
};
|
||||||
|
validate_workspace_memory_settings_snapshot(&snapshot, workspace_id)?;
|
||||||
|
let create_fingerprint =
|
||||||
|
bound_worker_create_fingerprint(request_fingerprint, &snapshot);
|
||||||
let worker_id = WorkerId::now_v7();
|
let worker_id = WorkerId::now_v7();
|
||||||
let now = chrono::Utc::now().to_rfc3339();
|
let now = chrono::Utc::now().to_rfc3339();
|
||||||
allocate_resource_key(
|
allocate_resource_key(
|
||||||
@@ -1194,19 +1374,27 @@ impl SqliteWorkspaceStore {
|
|||||||
tx.execute(
|
tx.execute(
|
||||||
"INSERT INTO worker_create_reservations(\
|
"INSERT INTO worker_create_reservations(\
|
||||||
workspace_id, allocation_key, worker_id, runtime_id, create_fingerprint,\
|
workspace_id, allocation_key, worker_id, runtime_id, create_fingerprint,\
|
||||||
state, created_at, updated_at\
|
state, created_at, updated_at, request_fingerprint,\
|
||||||
) VALUES (?1, ?2, ?3, ?4, ?5, 'reserved', ?6, ?6)",
|
memory_settings_revision, memory_language\
|
||||||
|
) VALUES (?1, ?2, ?3, ?4, ?5, 'reserved', ?6, ?6, ?7, ?8, ?9)",
|
||||||
params![
|
params![
|
||||||
workspace_id,
|
workspace_id,
|
||||||
allocation_key,
|
allocation_key,
|
||||||
worker_id.to_string(),
|
worker_id.to_string(),
|
||||||
runtime_id,
|
runtime_id,
|
||||||
create_fingerprint,
|
create_fingerprint,
|
||||||
now
|
now,
|
||||||
|
request_fingerprint,
|
||||||
|
snapshot.settings_revision as i64,
|
||||||
|
snapshot.language,
|
||||||
],
|
],
|
||||||
)?;
|
)?;
|
||||||
tx.commit()?;
|
tx.commit()?;
|
||||||
Ok(worker_id)
|
Ok(WorkerCreateReservation {
|
||||||
|
worker_id,
|
||||||
|
create_fingerprint,
|
||||||
|
memory_settings: snapshot,
|
||||||
|
})
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1375,8 +1563,9 @@ impl ControlPlaneStore for SqliteWorkspaceStore {
|
|||||||
}
|
}
|
||||||
|
|
||||||
async fn upsert_workspace(&self, record: &WorkspaceRecord) -> Result<()> {
|
async fn upsert_workspace(&self, record: &WorkspaceRecord) -> Result<()> {
|
||||||
self.with_conn(|conn| {
|
self.with_conn_mut(|conn| {
|
||||||
conn.execute(
|
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
|
||||||
|
tx.execute(
|
||||||
r#"INSERT INTO workspaces (
|
r#"INSERT INTO workspaces (
|
||||||
workspace_id, owner_account_id, display_name, state, created_at, updated_at
|
workspace_id, owner_account_id, display_name, state, created_at, updated_at
|
||||||
) VALUES (?1, ?2, ?3, ?4, ?5, ?6)
|
) VALUES (?1, ?2, ?3, ?4, ?5, ?6)
|
||||||
@@ -1394,6 +1583,13 @@ impl ControlPlaneStore for SqliteWorkspaceStore {
|
|||||||
record.updated_at,
|
record.updated_at,
|
||||||
],
|
],
|
||||||
)?;
|
)?;
|
||||||
|
tx.execute(
|
||||||
|
r#"INSERT OR IGNORE INTO workspace_memory_settings (
|
||||||
|
workspace_id, settings_revision, language, created_at, updated_at
|
||||||
|
) VALUES (?1, 1, 'English', ?2, ?3)"#,
|
||||||
|
params![record.workspace_id, record.created_at, record.updated_at],
|
||||||
|
)?;
|
||||||
|
tx.commit()?;
|
||||||
Ok(())
|
Ok(())
|
||||||
})?;
|
})?;
|
||||||
self.materialize_workspace_config(&record.workspace_id, &record.created_at)
|
self.materialize_workspace_config(&record.workspace_id, &record.created_at)
|
||||||
@@ -1518,6 +1714,16 @@ impl ControlPlaneStore for SqliteWorkspaceStore {
|
|||||||
record.workspace.updated_at,
|
record.workspace.updated_at,
|
||||||
],
|
],
|
||||||
)?;
|
)?;
|
||||||
|
tx.execute(
|
||||||
|
r#"INSERT INTO workspace_memory_settings (
|
||||||
|
workspace_id, settings_revision, language, created_at, updated_at
|
||||||
|
) VALUES (?1, 1, 'English', ?2, ?3)"#,
|
||||||
|
params![
|
||||||
|
record.workspace.workspace_id,
|
||||||
|
record.workspace.created_at,
|
||||||
|
record.workspace.updated_at,
|
||||||
|
],
|
||||||
|
)?;
|
||||||
tx.execute(
|
tx.execute(
|
||||||
r#"INSERT INTO repositories (
|
r#"INSERT INTO repositories (
|
||||||
workspace_id, repository_id, name, kind, provider, uri, default_ref,
|
workspace_id, repository_id, name, kind, provider, uri, default_ref,
|
||||||
@@ -5730,6 +5936,67 @@ fn collect_reference_diagnostics(
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn validate_workspace_memory_settings_snapshot(
|
||||||
|
snapshot: &manifest::WorkspaceMemorySettingsSnapshot,
|
||||||
|
expected_workspace_id: &str,
|
||||||
|
) -> Result<()> {
|
||||||
|
if snapshot.workspace_id != expected_workspace_id
|
||||||
|
|| snapshot.settings_revision == 0
|
||||||
|
|| !manifest::is_normalized_workspace_memory_language(&snapshot.language)
|
||||||
|
{
|
||||||
|
return Err(Error::Store(
|
||||||
|
"Workspace Memory settings are corrupt or belong to another Workspace".to_string(),
|
||||||
|
));
|
||||||
|
}
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn validate_workspace_memory_settings_record(
|
||||||
|
record: &WorkspaceMemorySettingsRecord,
|
||||||
|
expected_workspace_id: &str,
|
||||||
|
) -> Result<()> {
|
||||||
|
validate_workspace_memory_settings_snapshot(
|
||||||
|
&manifest::WorkspaceMemorySettingsSnapshot {
|
||||||
|
workspace_id: record.workspace_id.clone(),
|
||||||
|
settings_revision: record.settings_revision,
|
||||||
|
language: record.language.clone(),
|
||||||
|
},
|
||||||
|
expected_workspace_id,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn normalize_workspace_memory_language(language: &str) -> Result<String> {
|
||||||
|
let language = language.trim();
|
||||||
|
if !manifest::is_normalized_workspace_memory_language(language) {
|
||||||
|
return Err(Error::InvalidInput(format!(
|
||||||
|
"Workspace Memory language must be a non-empty UTF-8 string of at most {} characters without control characters",
|
||||||
|
manifest::MAX_WORKSPACE_MEMORY_LANGUAGE_CHARS
|
||||||
|
)));
|
||||||
|
}
|
||||||
|
Ok(language.to_string())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn bound_worker_create_fingerprint(
|
||||||
|
request_fingerprint: &str,
|
||||||
|
snapshot: &manifest::WorkspaceMemorySettingsSnapshot,
|
||||||
|
) -> String {
|
||||||
|
let mut digest = Sha256::new();
|
||||||
|
digest.update(b"workspace-worker-create-v2\0");
|
||||||
|
digest.update(request_fingerprint.as_bytes());
|
||||||
|
digest.update(b"\0");
|
||||||
|
digest.update(snapshot.workspace_id.as_bytes());
|
||||||
|
digest.update(b"\0");
|
||||||
|
digest.update(snapshot.settings_revision.to_be_bytes());
|
||||||
|
digest.update(b"\0");
|
||||||
|
digest.update(snapshot.language.as_bytes());
|
||||||
|
let encoded = digest
|
||||||
|
.finalize()
|
||||||
|
.iter()
|
||||||
|
.map(|byte| format!("{byte:02x}"))
|
||||||
|
.collect::<String>();
|
||||||
|
format!("sha256:{encoded}")
|
||||||
|
}
|
||||||
|
|
||||||
fn current_schema_version(conn: &Connection) -> Result<i64> {
|
fn current_schema_version(conn: &Connection) -> Result<i64> {
|
||||||
conn.query_row(
|
conn.query_row(
|
||||||
"SELECT COALESCE(MAX(version), 0) FROM __yoi_schema_migrations",
|
"SELECT COALESCE(MAX(version), 0) FROM __yoi_schema_migrations",
|
||||||
@@ -6373,6 +6640,42 @@ fn add_workspace_resource_human_keys(conn: &Connection) -> Result<()> {
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn create_workspace_memory_settings_authority(conn: &Connection) -> Result<()> {
|
||||||
|
conn.execute_batch(
|
||||||
|
r#"
|
||||||
|
CREATE TABLE workspace_memory_settings (
|
||||||
|
workspace_id TEXT PRIMARY KEY NOT NULL,
|
||||||
|
settings_revision INTEGER NOT NULL CHECK(settings_revision >= 1),
|
||||||
|
language TEXT NOT NULL CHECK(length(trim(language)) > 0),
|
||||||
|
created_at TEXT NOT NULL,
|
||||||
|
updated_at TEXT NOT NULL,
|
||||||
|
FOREIGN KEY (workspace_id) REFERENCES workspaces(workspace_id) ON DELETE CASCADE
|
||||||
|
);
|
||||||
|
|
||||||
|
INSERT INTO workspace_memory_settings (
|
||||||
|
workspace_id,
|
||||||
|
settings_revision,
|
||||||
|
language,
|
||||||
|
created_at,
|
||||||
|
updated_at
|
||||||
|
)
|
||||||
|
SELECT workspace_id, 1, 'English', created_at, updated_at
|
||||||
|
FROM workspaces;
|
||||||
|
|
||||||
|
ALTER TABLE worker_create_reservations
|
||||||
|
ADD COLUMN request_fingerprint TEXT;
|
||||||
|
ALTER TABLE worker_create_reservations
|
||||||
|
ADD COLUMN memory_settings_revision INTEGER;
|
||||||
|
ALTER TABLE worker_create_reservations
|
||||||
|
ADD COLUMN memory_language TEXT;
|
||||||
|
|
||||||
|
UPDATE worker_create_reservations
|
||||||
|
SET request_fingerprint = create_fingerprint;
|
||||||
|
"#,
|
||||||
|
)?;
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
fn verify_workspace_resource_key_schema(conn: &Connection) -> Result<()> {
|
fn verify_workspace_resource_key_schema(conn: &Connection) -> Result<()> {
|
||||||
ticket::migrate_sqlite_ticket_resource_key_schema_in_transaction(conn).map_err(|error| {
|
ticket::migrate_sqlite_ticket_resource_key_schema_in_transaction(conn).map_err(|error| {
|
||||||
Error::Store(format!(
|
Error::Store(format!(
|
||||||
@@ -7682,7 +7985,7 @@ mod tests {
|
|||||||
let before = std::fs::read(&path).unwrap();
|
let before = std::fs::read(&path).unwrap();
|
||||||
let plan = SqliteWorkspaceStore::migration_plan(&path).unwrap();
|
let plan = SqliteWorkspaceStore::migration_plan(&path).unwrap();
|
||||||
assert_eq!(plan.current_schema_version, 36);
|
assert_eq!(plan.current_schema_version, 36);
|
||||||
assert_eq!(plan.target_schema_version, 41);
|
assert_eq!(plan.target_schema_version, 42);
|
||||||
assert!(plan.migration_required);
|
assert!(plan.migration_required);
|
||||||
assert_eq!(plan.worker_count, 1);
|
assert_eq!(plan.worker_count, 1);
|
||||||
assert_eq!(plan.mappings[0].legacy_worker_id, 7);
|
assert_eq!(plan.mappings[0].legacy_worker_id, 7);
|
||||||
@@ -7696,7 +7999,7 @@ mod tests {
|
|||||||
store
|
store
|
||||||
.with_conn(|conn| {
|
.with_conn(|conn| {
|
||||||
assert!(table_exists(conn, "worker_diagnostics_archives")?);
|
assert!(table_exists(conn, "worker_diagnostics_archives")?);
|
||||||
assert_eq!(current_schema_version(conn)?, 41);
|
assert_eq!(current_schema_version(conn)?, 42);
|
||||||
Ok(())
|
Ok(())
|
||||||
})
|
})
|
||||||
.unwrap();
|
.unwrap();
|
||||||
@@ -7775,7 +8078,7 @@ mod tests {
|
|||||||
),
|
),
|
||||||
]
|
]
|
||||||
);
|
);
|
||||||
assert_eq!(current_schema_version(&conn).unwrap(), 41);
|
assert_eq!(current_schema_version(&conn).unwrap(), 42);
|
||||||
let foreign_key_error: Option<String> = conn
|
let foreign_key_error: Option<String> = conn
|
||||||
.query_row("PRAGMA foreign_key_check", [], |row| row.get(0))
|
.query_row("PRAGMA foreign_key_check", [], |row| row.get(0))
|
||||||
.optional()
|
.optional()
|
||||||
@@ -7904,7 +8207,7 @@ INSERT INTO worker_orphan_diagnostics (
|
|||||||
|
|
||||||
apply_migrations(&conn).unwrap();
|
apply_migrations(&conn).unwrap();
|
||||||
|
|
||||||
assert_eq!(current_schema_version(&conn).unwrap(), 41);
|
assert_eq!(current_schema_version(&conn).unwrap(), 42);
|
||||||
assert!(!table_exists(&conn, "worker_control_delegation_operations").unwrap());
|
assert!(!table_exists(&conn, "worker_control_delegation_operations").unwrap());
|
||||||
let controller_worker_id: String = conn
|
let controller_worker_id: String = conn
|
||||||
.query_row(
|
.query_row(
|
||||||
@@ -8022,10 +8325,36 @@ INSERT INTO worker_orphan_diagnostics (
|
|||||||
|
|
||||||
apply_migrations(&conn).unwrap();
|
apply_migrations(&conn).unwrap();
|
||||||
|
|
||||||
assert_eq!(current_schema_version(&conn).unwrap(), 41);
|
assert_eq!(current_schema_version(&conn).unwrap(), 42);
|
||||||
assert!(table_exists(&conn, "worker_workdir_attachment_reservations").unwrap());
|
assert!(table_exists(&conn, "worker_workdir_attachment_reservations").unwrap());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn schema_v42_initializes_existing_workspaces_with_explicit_english_memory_settings() {
|
||||||
|
let conn = Connection::open_in_memory().unwrap();
|
||||||
|
configure_sqlite(&conn).unwrap();
|
||||||
|
apply_migrations_through(&conn, 41).unwrap();
|
||||||
|
conn.execute(
|
||||||
|
"INSERT INTO workspaces (workspace_id, display_name, state, created_at, updated_at) \
|
||||||
|
VALUES ('workspace-existing', 'Existing', 'active', '1', '1')",
|
||||||
|
[],
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
apply_migrations(&conn).unwrap();
|
||||||
|
|
||||||
|
assert_eq!(current_schema_version(&conn).unwrap(), 42);
|
||||||
|
let settings = conn
|
||||||
|
.query_row(
|
||||||
|
"SELECT settings_revision, language FROM workspace_memory_settings \
|
||||||
|
WHERE workspace_id = 'workspace-existing'",
|
||||||
|
[],
|
||||||
|
|row| Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1)?)),
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(settings, (1, "English".to_string()));
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn schema_v26_removes_legacy_backend_flow_runtime_tables() {
|
fn schema_v26_removes_legacy_backend_flow_runtime_tables() {
|
||||||
let conn = Connection::open_in_memory().unwrap();
|
let conn = Connection::open_in_memory().unwrap();
|
||||||
@@ -8055,7 +8384,7 @@ CREATE TABLE flow_events (event_id TEXT PRIMARY KEY);
|
|||||||
|
|
||||||
apply_migrations(&conn).unwrap();
|
apply_migrations(&conn).unwrap();
|
||||||
|
|
||||||
assert_eq!(current_schema_version(&conn).unwrap(), 41);
|
assert_eq!(current_schema_version(&conn).unwrap(), 42);
|
||||||
assert!(table_exists(&conn, "flow_sources").unwrap());
|
assert!(table_exists(&conn, "flow_sources").unwrap());
|
||||||
assert!(table_exists(&conn, "flow_source_revisions").unwrap());
|
assert!(table_exists(&conn, "flow_source_revisions").unwrap());
|
||||||
assert!(!table_exists(&conn, "flow_instances").unwrap());
|
assert!(!table_exists(&conn, "flow_instances").unwrap());
|
||||||
@@ -8122,7 +8451,7 @@ INSERT INTO worker_workdir_attachment_reservations (
|
|||||||
|
|
||||||
apply_migrations(&conn).unwrap();
|
apply_migrations(&conn).unwrap();
|
||||||
|
|
||||||
assert_eq!(current_schema_version(&conn).unwrap(), 41);
|
assert_eq!(current_schema_version(&conn).unwrap(), 42);
|
||||||
let repositories_sql: String = conn
|
let repositories_sql: String = conn
|
||||||
.query_row(
|
.query_row(
|
||||||
"SELECT sql FROM sqlite_master WHERE type = 'table' AND name = 'repositories'",
|
"SELECT sql FROM sqlite_master WHERE type = 'table' AND name = 'repositories'",
|
||||||
@@ -8300,7 +8629,7 @@ INSERT INTO workdir_registry (
|
|||||||
let db = dir.path().join("control-plane.sqlite");
|
let db = dir.path().join("control-plane.sqlite");
|
||||||
let store = SqliteWorkspaceStore::open(&db).unwrap();
|
let store = SqliteWorkspaceStore::open(&db).unwrap();
|
||||||
|
|
||||||
assert_eq!(store.schema_version().await.unwrap(), 41);
|
assert_eq!(store.schema_version().await.unwrap(), 42);
|
||||||
assert!(
|
assert!(
|
||||||
!store
|
!store
|
||||||
.with_conn(|conn| table_exists(conn, "worker_workspace_credentials"))
|
.with_conn(|conn| table_exists(conn, "worker_workspace_credentials"))
|
||||||
@@ -8317,7 +8646,7 @@ INSERT INTO workdir_registry (
|
|||||||
store.upsert_workspace(&record).await.unwrap();
|
store.upsert_workspace(&record).await.unwrap();
|
||||||
|
|
||||||
let reopened = SqliteWorkspaceStore::open(&db).unwrap();
|
let reopened = SqliteWorkspaceStore::open(&db).unwrap();
|
||||||
assert_eq!(reopened.schema_version().await.unwrap(), 41);
|
assert_eq!(reopened.schema_version().await.unwrap(), 42);
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
reopened.get_workspace("local-dev").await.unwrap(),
|
reopened.get_workspace("local-dev").await.unwrap(),
|
||||||
Some(record)
|
Some(record)
|
||||||
@@ -8386,22 +8715,54 @@ INSERT INTO workdir_registry (
|
|||||||
.await
|
.await
|
||||||
.unwrap();
|
.unwrap();
|
||||||
|
|
||||||
|
let memory_settings = store.get_workspace_memory_settings("workspace-a").unwrap();
|
||||||
|
assert_eq!(memory_settings.settings_revision, 1);
|
||||||
|
assert_eq!(memory_settings.language, "English");
|
||||||
let reserved = store
|
let reserved = store
|
||||||
.reserve_worker_create("workspace-a", "arcadia", "operation-1", "sha256:one")
|
.reserve_worker_create(
|
||||||
|
"workspace-a",
|
||||||
|
"arcadia",
|
||||||
|
"operation-1",
|
||||||
|
"sha256:one",
|
||||||
|
&memory_settings,
|
||||||
|
)
|
||||||
.unwrap();
|
.unwrap();
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
reserved.as_uuid().get_version(),
|
reserved.worker_id.as_uuid().get_version(),
|
||||||
Some(uuid::Version::SortRand)
|
Some(uuid::Version::SortRand)
|
||||||
);
|
);
|
||||||
|
assert_eq!(reserved.memory_settings.settings_revision, 1);
|
||||||
|
assert_eq!(reserved.memory_settings.language, "English");
|
||||||
|
let unchanged_memory_settings = store
|
||||||
|
.update_workspace_memory_settings("workspace-a", 1, " English ")
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(unchanged_memory_settings, memory_settings);
|
||||||
|
let updated_memory_settings = store
|
||||||
|
.update_workspace_memory_settings("workspace-a", 1, " Français ")
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(updated_memory_settings.settings_revision, 2);
|
||||||
|
assert_eq!(updated_memory_settings.language, "Français");
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
store
|
store
|
||||||
.reserve_worker_create("workspace-a", "arcadia", "operation-1", "sha256:one")
|
.reserve_worker_create(
|
||||||
|
"workspace-a",
|
||||||
|
"arcadia",
|
||||||
|
"operation-1",
|
||||||
|
"sha256:one",
|
||||||
|
&updated_memory_settings,
|
||||||
|
)
|
||||||
.unwrap(),
|
.unwrap(),
|
||||||
reserved
|
reserved
|
||||||
);
|
);
|
||||||
assert!(
|
assert!(
|
||||||
store
|
store
|
||||||
.reserve_worker_create("workspace-a", "arcadia", "operation-1", "sha256:different")
|
.reserve_worker_create(
|
||||||
|
"workspace-a",
|
||||||
|
"arcadia",
|
||||||
|
"operation-1",
|
||||||
|
"sha256:different",
|
||||||
|
&updated_memory_settings,
|
||||||
|
)
|
||||||
.is_err()
|
.is_err()
|
||||||
);
|
);
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
@@ -8409,7 +8770,7 @@ INSERT INTO workdir_registry (
|
|||||||
.resource_key(
|
.resource_key(
|
||||||
"workspace-a",
|
"workspace-a",
|
||||||
WorkspaceResourceKind::Worker,
|
WorkspaceResourceKind::Worker,
|
||||||
&reserved.to_string()
|
&reserved.worker_id.to_string()
|
||||||
)
|
)
|
||||||
.unwrap()
|
.unwrap()
|
||||||
.as_deref(),
|
.as_deref(),
|
||||||
@@ -8419,37 +8780,85 @@ INSERT INTO workdir_registry (
|
|||||||
store
|
store
|
||||||
.resolve_resource_reference("workspace-a", WorkspaceResourceKind::Worker, "W-1")
|
.resolve_resource_reference("workspace-a", WorkspaceResourceKind::Worker, "W-1")
|
||||||
.unwrap(),
|
.unwrap(),
|
||||||
Some(reserved.to_string())
|
Some(reserved.worker_id.to_string())
|
||||||
);
|
);
|
||||||
let second = store
|
let second = store
|
||||||
.reserve_worker_create("workspace-a", "arcadia", "operation-2", "sha256:two")
|
.reserve_worker_create(
|
||||||
|
"workspace-a",
|
||||||
|
"arcadia",
|
||||||
|
"operation-2",
|
||||||
|
"sha256:two",
|
||||||
|
&updated_memory_settings,
|
||||||
|
)
|
||||||
.unwrap();
|
.unwrap();
|
||||||
|
assert_eq!(second.memory_settings.settings_revision, 2);
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
store
|
store
|
||||||
.resource_key(
|
.resource_key(
|
||||||
"workspace-a",
|
"workspace-a",
|
||||||
WorkspaceResourceKind::Worker,
|
WorkspaceResourceKind::Worker,
|
||||||
&second.to_string()
|
&second.worker_id.to_string()
|
||||||
)
|
)
|
||||||
.unwrap()
|
.unwrap()
|
||||||
.as_deref(),
|
.as_deref(),
|
||||||
Some("W-2")
|
Some("W-2")
|
||||||
);
|
);
|
||||||
store
|
store
|
||||||
.complete_worker_create_reservation("workspace-a", reserved)
|
.complete_worker_create_reservation("workspace-a", reserved.worker_id)
|
||||||
.unwrap();
|
.unwrap();
|
||||||
let state: String = store
|
let state: String = store
|
||||||
.with_conn(|conn| {
|
.with_conn(|conn| {
|
||||||
conn.query_row(
|
conn.query_row(
|
||||||
"SELECT state FROM worker_create_reservations \
|
"SELECT state FROM worker_create_reservations \
|
||||||
WHERE workspace_id = 'workspace-a' AND worker_id = ?1",
|
WHERE workspace_id = 'workspace-a' AND worker_id = ?1",
|
||||||
[reserved.to_string()],
|
[reserved.worker_id.to_string()],
|
||||||
|row| row.get(0),
|
|row| row.get(0),
|
||||||
)
|
)
|
||||||
.map_err(Error::from)
|
.map_err(Error::from)
|
||||||
})
|
})
|
||||||
.unwrap();
|
.unwrap();
|
||||||
assert_eq!(state, "created");
|
assert_eq!(state, "created");
|
||||||
|
|
||||||
|
store
|
||||||
|
.with_conn(|conn| {
|
||||||
|
conn.execute(
|
||||||
|
"UPDATE workspace_memory_settings SET language = ' English ' \
|
||||||
|
WHERE workspace_id = 'workspace-a'",
|
||||||
|
[],
|
||||||
|
)?;
|
||||||
|
Ok(())
|
||||||
|
})
|
||||||
|
.unwrap();
|
||||||
|
assert!(store.get_workspace_memory_settings("workspace-a").is_err());
|
||||||
|
assert!(
|
||||||
|
store
|
||||||
|
.update_workspace_memory_settings("workspace-a", 2, "Spanish")
|
||||||
|
.is_err()
|
||||||
|
);
|
||||||
|
let mut corrupt = updated_memory_settings.clone();
|
||||||
|
corrupt.language = " English ".to_string();
|
||||||
|
assert!(
|
||||||
|
store
|
||||||
|
.reserve_worker_create(
|
||||||
|
"workspace-a",
|
||||||
|
"arcadia",
|
||||||
|
"operation-corrupt",
|
||||||
|
"sha256:corrupt",
|
||||||
|
&corrupt,
|
||||||
|
)
|
||||||
|
.is_err()
|
||||||
|
);
|
||||||
|
|
||||||
|
store
|
||||||
|
.with_conn(|conn| {
|
||||||
|
conn.execute(
|
||||||
|
"DELETE FROM workspace_memory_settings WHERE workspace_id = 'workspace-a'",
|
||||||
|
[],
|
||||||
|
)?;
|
||||||
|
Ok(())
|
||||||
|
})
|
||||||
|
.unwrap();
|
||||||
|
assert!(store.get_workspace_memory_settings("workspace-a").is_err());
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
@@ -8816,13 +9225,13 @@ INSERT INTO worker_registry (
|
|||||||
configure_sqlite(&conn).unwrap();
|
configure_sqlite(&conn).unwrap();
|
||||||
apply_migrations(&conn).unwrap();
|
apply_migrations(&conn).unwrap();
|
||||||
conn.execute(
|
conn.execute(
|
||||||
"INSERT INTO __yoi_schema_migrations (version, name) VALUES (42, 'future')",
|
"INSERT INTO __yoi_schema_migrations (version, name) VALUES (43, 'future')",
|
||||||
[],
|
[],
|
||||||
)
|
)
|
||||||
.unwrap();
|
.unwrap();
|
||||||
|
|
||||||
let error = apply_migrations(&conn).unwrap_err().to_string();
|
let error = apply_migrations(&conn).unwrap_err().to_string();
|
||||||
assert!(error.contains("schema version 42 is newer"), "{error}");
|
assert!(error.contains("schema version 43 is newer"), "{error}");
|
||||||
assert!(error.contains("refusing to serve"), "{error}");
|
assert!(error.contains("refusing to serve"), "{error}");
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -9043,7 +9452,7 @@ VALUES ('workspace-b', 'ticket-b', 'related', 'ticket-a', NULL, 'tester', '2026-
|
|||||||
|
|
||||||
apply_migrations(&mut conn).unwrap();
|
apply_migrations(&mut conn).unwrap();
|
||||||
|
|
||||||
assert_eq!(current_schema_version(&conn).unwrap(), 41);
|
assert_eq!(current_schema_version(&conn).unwrap(), 42);
|
||||||
let workspace_id: Option<String> = conn
|
let workspace_id: Option<String> = conn
|
||||||
.query_row(
|
.query_row(
|
||||||
"SELECT workspace_id FROM trusted_runtime_records WHERE runtime_id = 'runtime-a'",
|
"SELECT workspace_id FROM trusted_runtime_records WHERE runtime_id = 'runtime-a'",
|
||||||
@@ -9660,7 +10069,7 @@ WHERE workspace_id = 'workspace-a'
|
|||||||
.unwrap();
|
.unwrap();
|
||||||
|
|
||||||
let store = SqliteWorkspaceStore::from_connection(conn).unwrap();
|
let store = SqliteWorkspaceStore::from_connection(conn).unwrap();
|
||||||
assert_eq!(store.schema_version().await.unwrap(), 41);
|
assert_eq!(store.schema_version().await.unwrap(), 42);
|
||||||
|
|
||||||
store
|
store
|
||||||
.with_conn(|conn| {
|
.with_conn(|conn| {
|
||||||
@@ -9849,7 +10258,7 @@ CREATE TABLE ticket_assignment_operations (
|
|||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn repository_records_round_trip() {
|
async fn repository_records_round_trip() {
|
||||||
let store = SqliteWorkspaceStore::in_memory().unwrap();
|
let store = SqliteWorkspaceStore::in_memory().unwrap();
|
||||||
assert_eq!(store.schema_version().await.unwrap(), 41);
|
assert_eq!(store.schema_version().await.unwrap(), 42);
|
||||||
let workspace = WorkspaceRecord {
|
let workspace = WorkspaceRecord {
|
||||||
workspace_id: "local-dev".to_string(),
|
workspace_id: "local-dev".to_string(),
|
||||||
owner_account_id: None,
|
owner_account_id: None,
|
||||||
@@ -9915,7 +10324,7 @@ CREATE TABLE ticket_assignment_operations (
|
|||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn memory_authority_records_round_trip_and_close_staging() {
|
async fn memory_authority_records_round_trip_and_close_staging() {
|
||||||
let store = SqliteWorkspaceStore::in_memory().unwrap();
|
let store = SqliteWorkspaceStore::in_memory().unwrap();
|
||||||
assert_eq!(store.schema_version().await.unwrap(), 41);
|
assert_eq!(store.schema_version().await.unwrap(), 42);
|
||||||
let workspace = WorkspaceRecord {
|
let workspace = WorkspaceRecord {
|
||||||
workspace_id: "local-dev".to_string(),
|
workspace_id: "local-dev".to_string(),
|
||||||
owner_account_id: None,
|
owner_account_id: None,
|
||||||
@@ -10317,7 +10726,7 @@ CREATE TABLE ticket_assignment_operations (
|
|||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn account_and_login_records_round_trip() {
|
async fn account_and_login_records_round_trip() {
|
||||||
let store = SqliteWorkspaceStore::in_memory().unwrap();
|
let store = SqliteWorkspaceStore::in_memory().unwrap();
|
||||||
assert_eq!(store.schema_version().await.unwrap(), 41);
|
assert_eq!(store.schema_version().await.unwrap(), 42);
|
||||||
let now = "2026-07-22T00:00:00Z".to_string();
|
let now = "2026-07-22T00:00:00Z".to_string();
|
||||||
let account = AccountRecord {
|
let account = AccountRecord {
|
||||||
account_id: "acct-user-alice".to_string(),
|
account_id: "acct-user-alice".to_string(),
|
||||||
|
|||||||
Reference in New Issue
Block a user