diff --git a/.repository-projection.json b/.repository-projection.json index 0a0bef0ee..578f8f253 100644 --- a/.repository-projection.json +++ b/.repository-projection.json @@ -3,11 +3,11 @@ "projection": "deixic-code", "projectionSchemaVersion": 1, "sourceRepository": "dx-corp/mono", - "sourceSha": "36bd59e4bac6e289aaf493f38e26f42fc382c4b3", + "sourceSha": "393b500389b42a02222404dfb7bbe3a3586871f7", "destinationRepository": "dx-corp/code", - "priorProjectedBase": "dccebe04146633d0841585f1d511aceda5ea746d", + "priorProjectedBase": "a775955b5d82cdc7da4d5c41d1c2083b8c208e4d", "definitionDigest": "82936441c776e3e8edb5d215a75007ec9714a233f489d460075d79d5ef5ba32f", - "toolDigest": "45502ff0478e541d02f34cc39ac332935f8d3c0030fa221afe5bcc3b5e51f88e", - "contentDigest": "78b16e2e7b3b23f64ba5345cbae3d36206cdf4d1df28f1756e5ecf53cfb2d4ea", + "toolDigest": "7f2d7ef1750a1809dd90f3bd7e50c4220784a8a9e0b3dfd499355e509132c32f", + "contentDigest": "29e7e643aca0e259d6bf6c423110fa84fe86a8dc49f33b85e0d4770228debcf1", "publicationEligible": true } diff --git a/README.md b/README.md index a263e7e41..945d7d24f 100644 --- a/README.md +++ b/README.md @@ -77,7 +77,7 @@ deixic-code thread attach thread:example --message "Check the latest result" --j ``` The terminal reads the Platform-owned thread, submits messages through the -same tenant-scoped owner API as the Deixic UI, and shows safe execution events. +tenant-scoped public application API, and shows safe execution events. When a turn needs approval or input, use `/respond [text]` after reviewing the request. The CLI never grants an approval automatically. `remote attach` is a separate RunnerSession transport and does not attach to a @@ -85,6 +85,12 @@ Dex thread. A disconnected client can reattach without restarting the resident turn. The one-shot form stops following after two minutes and leaves accepted work running on Platform. +Native thread, readiness, managed setup, and issue-report clients use generated +`deixicpublic.v1.DeixicPublicService` bindings from +[the public contract](proto/deixicpublic/v1/sdk.proto). The release build and staged +release verifier audit the native bytes for private protocol namespaces before +publication. The contract copy is checked against the canonical public source in Mono. + For a multi-feature coding task, prepare a mission contract and run it through the native workflow scheduler. See [Run a local coding mission](docs/MISSION_RUN.md). diff --git a/packages/local-host-rs/build.rs b/packages/local-host-rs/build.rs index 57a871bdf..741aa7de4 100644 --- a/packages/local-host-rs/build.rs +++ b/packages/local-host-rs/build.rs @@ -4,8 +4,8 @@ use std::path::PathBuf; fn main() -> Result<(), Box> { let manifest_dir = PathBuf::from(env::var("CARGO_MANIFEST_DIR")?); let proto_dir = manifest_dir.join("../../proto"); - let readiness_proto = proto_dir.join("console/v1/managed_inference.proto"); - println!("cargo:rerun-if-changed={}", readiness_proto.display()); + let public_proto = proto_dir.join("deixicpublic/v1/sdk.proto"); + println!("cargo:rerun-if-changed={}", public_proto.display()); let proto_file = proto_dir.join("maestro/v1/headless.proto"); println!("cargo:rerun-if-changed=build.rs"); @@ -19,7 +19,7 @@ fn main() -> Result<(), Box> { config.extern_path(".google.protobuf.Struct", "::prost_types::Struct"); config.extern_path(".google.protobuf.ListValue", "::prost_types::ListValue"); config.extern_path(".google.protobuf.NullValue", "::prost_types::NullValue"); - config.compile_protos(&[proto_file, readiness_proto], &[proto_dir])?; + config.compile_protos(&[proto_file, public_proto], &[proto_dir])?; Ok(()) } diff --git a/packages/local-host-rs/src/bug_report.rs b/packages/local-host-rs/src/bug_report.rs index adbfb239c..98dd23df2 100644 --- a/packages/local-host-rs/src/bug_report.rs +++ b/packages/local-host-rs/src/bug_report.rs @@ -15,7 +15,7 @@ use crate::session::CustomEntry; use crate::session::{SessionEntry, SessionManager}; const ENTRY_TYPE: &str = "product_issue_draft_v1"; -const SUBMIT_PATH: &str = "/deixic.v1.DeixicService/SubmitNativeProductIssueReport"; +const SUBMIT_PATH: &str = "/deixicpublic.v1.DeixicPublicService/SubmitIssueReport"; #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] pub struct Destination { @@ -263,7 +263,7 @@ impl FeedbackClient { "The report must be reviewed and saved before sending." ); let request = SubmitRequest { - query: Some(ReportQuery { + scope: Some(ReportQuery { organization_id: self.destination.organization_id.clone(), workspace_id: self.destination.workspace_id.clone(), }), @@ -276,7 +276,7 @@ impl FeedbackClient { }, include_diagnostics: report.include_diagnostics, idempotency_key: report.id.clone(), - context: report.outgoing_context(), + context: report.outgoing_context().map(Into::into), }; let response = self.http.post(&self.destination.endpoint).bearer_auth(&self.token) .header("connect-protocol-version", "1") @@ -345,53 +345,39 @@ fn endpoint(base: &str) -> Result { Ok(url.into()) } -// Bounded native projection of proto/console/v1/console.proto. Wire tags are -// verified against a shared fixture decoded by the production service tests. -#[derive(Clone, PartialEq, Message)] -struct SubmitRequest { - #[prost(message, optional, tag = "1")] - query: Option, - #[prost(string, tag = "2")] - description: String, - #[prost(string, tag = "3")] - expected_behavior: String, - #[prost(string, tag = "5")] - app_version: String, - #[prost(bool, tag = "11")] - include_diagnostics: bool, - #[prost(string, tag = "12")] - idempotency_key: String, - #[prost(message, optional, tag = "13")] - context: Option, -} -#[derive(Clone, PartialEq, Message)] -struct ReportQuery { - #[prost(string, tag = "13")] - organization_id: String, - #[prost(string, tag = "1")] - workspace_id: String, -} -#[derive(Clone, PartialEq, Message)] -struct SubmitResponse { - #[prost(message, optional, tag = "1")] - report: Option, -} -#[derive(Clone, PartialEq, Message)] -struct ReportReceipt { - #[prost(string, tag = "1")] - id: String, - #[prost(string, tag = "2")] - reference: String, +#[cfg(test)] +use crate::public_protocol::IssueReport as ReportReceipt; +use crate::public_protocol::{ + IssueContext, IssueEvidence, Scope as ReportQuery, SubmitIssueReportRequest as SubmitRequest, + SubmitIssueReportResponse as SubmitResponse, +}; + +impl From for IssueContext { + fn from(context: ReportContext) -> Self { + Self { + reproduction_steps: context.reproduction_steps, + model: context.model, + evidence: context + .evidence + .into_iter() + .map(|e| IssueEvidence { + kind: e.kind, + source_id: e.source_id, + text: e.text, + }) + .collect(), + } + } } #[cfg(test)] mod tests { use super::*; #[test] - fn native_wire_matches_the_fixture_read_by_the_product_issue_service() { + fn public_issue_report_wire_matches_reviewed_fixtures() { let mut request = SubmitRequest { context: None, - query: Some(ReportQuery { + scope: Some(ReportQuery { organization_id: "org-1".into(), workspace_id: "workspace-1".into(), }), @@ -408,12 +394,12 @@ mod tests { } assert_eq!( encoded, - include_str!("../../../test/fixtures/product-issue-report-native-v1.hex").trim() + include_str!("../../../test/fixtures/product-issue-report-public-v1.hex").trim() ); - request.context = Some(ReportContext { + request.context = Some(IssueContext { reproduction_steps: "Repeat failing tool".into(), model: "test-model".into(), - evidence: vec![ReportEvidence { + evidence: vec![IssueEvidence { kind: "tool_result".into(), source_id: "call-1".into(), text: "wrong action".into(), @@ -425,7 +411,7 @@ mod tests { } assert_eq!( encoded, - include_str!("../../../test/fixtures/product-issue-report-native-v2.hex").trim() + include_str!("../../../test/fixtures/product-issue-report-public-v2.hex").trim() ); } @@ -546,26 +532,20 @@ fn now_seconds() -> i64 { chrono::Utc::now().timestamp() } -#[derive(Clone, PartialEq, Message, Serialize, Deserialize)] +#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)] pub struct ReportContext { - #[prost(string, tag = "1")] #[serde(default)] pub reproduction_steps: String, - #[prost(string, tag = "2")] #[serde(default)] pub model: String, - #[prost(message, repeated, tag = "3")] #[serde(default)] pub evidence: Vec, } -#[derive(Clone, PartialEq, Message, Serialize, Deserialize)] +#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)] pub struct ReportEvidence { - #[prost(string, tag = "1")] pub kind: String, - #[prost(string, tag = "2")] pub source_id: String, - #[prost(string, tag = "3")] pub text: String, } diff --git a/packages/local-host-rs/src/doctor.rs b/packages/local-host-rs/src/doctor.rs index aad897ee6..3745cc45e 100644 --- a/packages/local-host-rs/src/doctor.rs +++ b/packages/local-host-rs/src/doctor.rs @@ -1694,16 +1694,17 @@ mod tests { .build() .expect("runtime") .block_on(async { - // Independently encode the canonical response tags (version=1, - // organization_id=7, workspace_id=8, MCP policy=5). + // Independently encode the public response: scope=1, version=2, MCP=6. for (organization, expected) in [ ("org-test", CheckStatus::Pass), ("other-org", CheckStatus::Warning), ] { - let mut body = vec![0x08, 1, 0x3a, organization.len() as u8]; - body.extend_from_slice(organization.as_bytes()); - body.extend_from_slice(b"\x42\x0eworkspace-test"); - body.extend_from_slice(&[0x2a, 2, 0x08, 2]); + let mut scope = vec![0x0a, organization.len() as u8]; + scope.extend_from_slice(organization.as_bytes()); + scope.extend_from_slice(b"\x12\x0eworkspace-test"); + let mut body = vec![0x0a, scope.len() as u8]; + body.extend_from_slice(&scope); + body.extend_from_slice(&[0x10, 1, 0x32, 2, 0x08, 2]); let (base, server) = test_server_bytes(200, body, "application/proto", Duration::ZERO).await; std::env::set_var("MAESTRO_MANAGED_SETUP_URL", base); @@ -1717,7 +1718,9 @@ mod tests { assert_eq!(report.status, expected, "{report:?}"); assert!(report.live); let request = server.await.expect("server"); - assert!(request.contains("/console.v1.ManagedSetupService/GetManagedSetup")); + assert!( + request.contains("/deixicpublic.v1.DeixicPublicService/GetClientSetup") + ); assert!( !serde_json::to_string(&report) .expect("report") diff --git a/packages/local-host-rs/src/headless/messages.rs b/packages/local-host-rs/src/headless/messages.rs index 83f299542..97d021d77 100644 --- a/packages/local-host-rs/src/headless/messages.rs +++ b/packages/local-host-rs/src/headless/messages.rs @@ -629,6 +629,10 @@ pub struct GovernedToolGrant { /// Platform-owned process definition instructions installed in the system role. #[serde(default, skip_serializing_if = "Option::is_none")] pub process_system_prompt: Option, + /// Exact Platform-admitted Agent Registry publication and immutable + /// PromptService snapshots for this turn. It grants no tools by itself. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub agent_profile: Option, pub envelope_version: u32, pub grant_id: String, pub grant_version: u64, @@ -657,6 +661,26 @@ pub struct GovernedToolGrant { pub connection_bindings: Vec, } +#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)] +#[serde(default, deny_unknown_fields, rename_all = "camelCase")] +pub struct AgentProfileTurnContext { + pub agent_id: String, + pub config_version: i32, + pub digest_sha256: String, + #[serde(skip_serializing_if = "Vec::is_empty")] + pub instructions: Vec, +} + +#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)] +#[serde(default, deny_unknown_fields, rename_all = "camelCase")] +pub struct AgentProfileInstruction { + pub name: String, + pub prompt_id: String, + pub version_id: String, + pub content_digest_sha256: String, + pub content: String, +} + impl GovernedToolGrant { /// Stable identity used for reconnect replay and duplicate-turn matching. #[must_use] diff --git a/packages/local-host-rs/src/headless/proto.rs b/packages/local-host-rs/src/headless/proto.rs index d65667e33..0237024c4 100644 --- a/packages/local-host-rs/src/headless/proto.rs +++ b/packages/local-host-rs/src/headless/proto.rs @@ -44,6 +44,7 @@ mod tests { RuntimeGovernedToolGrant { process_budget: None, process_system_prompt: None, + agent_profile: None, envelope_version: 2, grant_id: "grant-1".into(), grant_version: 1, diff --git a/packages/local-host-rs/src/headless_server.rs b/packages/local-host-rs/src/headless_server.rs index bffb448b2..f98d3008e 100644 --- a/packages/local-host-rs/src/headless_server.rs +++ b/packages/local-host-rs/src/headless_server.rs @@ -730,6 +730,9 @@ fn governed_authority_material_digest(grant: &GovernedToolGrant) -> Result serde_json::Valu value["process_budget"] = serde_json::json!(budget); value["process_system_prompt"] = serde_json::json!(grant.process_system_prompt); } + if let Some(profile) = &grant.agent_profile { + value["agent_profile"] = serde_json::json!(profile); + } // Preserve the exact v2 canonical form for grants minted before // connection bindings existed. New authority is included whenever used. if !grant.connection_bindings.is_empty() { @@ -1166,6 +1172,29 @@ fn verify_governed_tool_grant_with_keys( if grant.process_budget.is_some() != grant.process_system_prompt.is_some() { anyhow::bail!("process budget requires its signed system instructions"); } + if let Some(profile) = &grant.agent_profile { + if grant.process_budget.is_some() + || profile.agent_id.trim().is_empty() + || profile.config_version <= 0 + || !is_plain_sha256_digest(&profile.digest_sha256) + { + anyhow::bail!("invalid signed agent profile"); + } + let mut names = std::collections::HashSet::new(); + let mut total_bytes = 0usize; + for instruction in &profile.instructions { + total_bytes = total_bytes.saturating_add(instruction.content.len()); + if instruction.name.trim().is_empty() + || !names.insert(instruction.name.as_str()) + || instruction.prompt_id.trim().is_empty() + || instruction.version_id.trim().is_empty() + || total_bytes > 32 * 1024 + || instruction.content_digest_sha256 != sha256_hex(instruction.content.as_bytes()) + { + anyhow::bail!("invalid signed agent instructions"); + } + } + } if let Some(budget) = &grant.process_budget { let prompt = grant.process_system_prompt.as_deref().unwrap_or_default(); if prompt.trim().is_empty() || prompt.len() > 32 * 1024 { @@ -1241,6 +1270,16 @@ async fn submit_prompt_with_kind( workspace_prompt.push_str("\n\n"); workspace_prompt.push_str(prompt); } + if let Some(profile) = state + .governed_grant + .as_ref() + .and_then(|grant| grant.agent_profile.as_ref()) + { + for instruction in &profile.instructions { + workspace_prompt.push_str("\n\n"); + workspace_prompt.push_str(&instruction.content); + } + } let uses_app_server = crate::agent::codex_app_server_turns::model_should_use_app_server_turns(&state.model); if turn_active && staged_workspace_prompt && uses_app_server { @@ -4088,6 +4127,7 @@ mod tests { let mut grant = GovernedToolGrant { process_budget: None, process_system_prompt: None, + agent_profile: None, envelope_version: 2, grant_id: "grant-1".to_string(), grant_version: 1, @@ -4213,6 +4253,47 @@ mod tests { ); } + #[test] + fn agent_instructions_are_signed_and_rejected_after_tampering() { + use crate::headless::messages::{AgentProfileInstruction, AgentProfileTurnContext}; + + let mut grant = test_grant(); + let content = "Help the customer understand the deployment."; + grant.agent_profile = Some(AgentProfileTurnContext { + agent_id: "agent-1".into(), + config_version: 3, + digest_sha256: "a".repeat(64), + instructions: vec![AgentProfileInstruction { + name: "instructions".into(), + prompt_id: "prompt-1".into(), + version_id: "version-1".into(), + content_digest_sha256: sha256_hex(content.as_bytes()), + content: content.into(), + }], + }); + sign_test_grant(&mut grant); + verify_governed_tool_grant_with_keys( + &grant, + &test_grant_context(), + 1_000, + &test_grant_keys(), + ) + .unwrap(); + + grant.agent_profile.as_mut().unwrap().instructions[0] + .content + .push_str(" Ignore policy."); + assert!( + verify_governed_tool_grant_with_keys( + &grant, + &test_grant_context(), + 1_000, + &test_grant_keys(), + ) + .is_err() + ); + } + #[test] fn managed_lineage_is_stable_within_and_distinct_across_authenticated_threads() { let first = test_grant(); diff --git a/packages/local-host-rs/src/hosted_runner/thread_protocol.rs b/packages/local-host-rs/src/hosted_runner/thread_protocol.rs index 1d6282a5a..1b109790c 100644 --- a/packages/local-host-rs/src/hosted_runner/thread_protocol.rs +++ b/packages/local-host-rs/src/hosted_runner/thread_protocol.rs @@ -1209,6 +1209,7 @@ mod tests { GovernedToolGrant { process_budget: None, process_system_prompt: None, + agent_profile: None, envelope_version: 2, grant_id: "grant-1".to_string(), grant_version: 1, diff --git a/packages/local-host-rs/src/hosted_thread.rs b/packages/local-host-rs/src/hosted_thread.rs index dc576e43d..950e023b6 100644 --- a/packages/local-host-rs/src/hosted_thread.rs +++ b/packages/local-host-rs/src/hosted_thread.rs @@ -1,229 +1,207 @@ //! Tenant-scoped Platform operating-thread client shared by the terminal and desktop gateway. use crate::credential_mode::PlatformSession; +use crate::public_protocol as wire; use anyhow::{Context, Result, bail}; use prost::Message; use reqwest::{Client, Url}; use std::time::Duration; use uuid::Uuid; -const SERVICE: &str = "/deixic.v1.DeixicService"; +const SERVICE: &str = "/deixicpublic.v1.DeixicPublicService"; const MAX_RESPONSE_BYTES: usize = 4 * 1024 * 1024; -// A narrow wire projection of console.v1. Field numbers and kinds mirror the -// canonical protobuf contract; unknown fields are ignored by prost. -#[derive(Clone, PartialEq, Message)] +// Native-friendly view models. Only `wire` types are serialized on the public +// connection; these views do not mirror any internal protobuf field numbers. +#[derive(Clone, Debug, PartialEq, Default)] pub struct Query { - #[prost(string, tag = "1")] pub workspace_id: String, - #[prost(string, tag = "13")] pub organization_id: String, } -#[derive(Clone, PartialEq, Message)] +#[derive(Clone, Debug, PartialEq, Default)] pub struct GetRequest { - #[prost(message, optional, tag = "1")] pub query: Option, - #[prost(string, tag = "2")] pub channel_id: String, - #[prost(int32, tag = "3")] pub limit: i32, } -#[derive(Clone, PartialEq, Message)] +#[derive(Clone, Debug, PartialEq, Default)] pub struct Channel { - #[prost(string, tag = "1")] pub id: String, - #[prost(string, tag = "2")] pub label: String, - #[prost(int32, tag = "5")] pub unread_count: i32, - #[prost(int32, tag = "6")] pub open_count: i32, - #[prost(bool, tag = "11")] pub archived: bool, } -#[derive(Clone, PartialEq, Message)] +#[derive(Clone, Debug, PartialEq, Default)] pub struct ListChannelsRequest { - #[prost(message, optional, tag = "1")] pub query: Option, - #[prost(int32, tag = "2")] pub archive_filter: i32, } -#[derive(Clone, PartialEq, Message)] +#[derive(Clone, Debug, PartialEq, Default)] pub struct ListChannelsResponse { - #[prost(message, repeated, tag = "1")] pub channels: Vec, } -#[derive(Clone, PartialEq, Message)] +#[derive(Clone, Debug, PartialEq, Default)] pub struct RenameThreadRequest { - #[prost(message, optional, tag = "1")] pub query: Option, - #[prost(string, tag = "2")] pub channel_id: String, - #[prost(string, tag = "3")] pub title: String, } -#[derive(Clone, PartialEq, Message)] +#[derive(Clone, Debug, PartialEq, Default)] pub struct RenameThreadResponse { - #[prost(message, optional, tag = "1")] pub channel: Option, - #[prost(bool, tag = "2")] pub changed: bool, } -#[derive(Clone, PartialEq, Message)] +#[derive(Clone, Debug, PartialEq, Default)] pub struct ArchiveThreadRequest { - #[prost(message, optional, tag = "1")] pub query: Option, - #[prost(string, tag = "2")] pub channel_id: String, - #[prost(bool, tag = "3")] pub archived: bool, - #[prost(string, tag = "4")] pub idempotency_key: String, } -#[derive(Clone, PartialEq, Message)] +#[derive(Clone, Debug, PartialEq, Default)] pub struct ArchiveThreadResponse { - #[prost(message, optional, tag = "1")] pub channel: Option, - #[prost(bool, tag = "2")] pub changed: bool, - #[prost(bool, tag = "3")] pub replayed: bool, } -#[derive(Clone, PartialEq, Message)] +#[derive(Clone, Debug, PartialEq, Default)] pub struct OperatingMessage { - #[prost(string, tag = "1")] pub id: String, - #[prost(string, tag = "3")] pub role: String, - #[prost(string, tag = "5")] pub body: String, } -#[derive(Clone, PartialEq, Message)] +#[derive(Clone, Debug, PartialEq, Default)] pub struct GetResponse { - #[prost(message, optional, tag = "1")] pub channel: Option, - #[prost(message, repeated, tag = "2")] pub messages: Vec, - #[prost(int64, tag = "7")] pub replay_cursor: i64, } -#[derive(Clone, PartialEq, Message)] +#[derive(Clone, Debug, PartialEq, Default)] pub struct ListRequest { - #[prost(message, optional, tag = "1")] pub query: Option, - #[prost(string, tag = "2")] pub channel_id: String, - #[prost(int64, tag = "3")] pub after_cursor: i64, - #[prost(int32, tag = "4")] pub limit: i32, } -#[derive(Clone, PartialEq, Message)] +#[derive(Clone, Debug, PartialEq, Default)] pub struct Event { - #[prost(int64, tag = "1")] pub cursor: i64, - #[prost(string, tag = "2")] pub event_id: String, - #[prost(string, tag = "3")] pub turn_id: String, - #[prost(int32, tag = "4")] pub kind: i32, - #[prost(string, tag = "5")] pub safe_text: String, - #[prost(string, tag = "9")] pub request_id: String, - #[prost(int32, tag = "10")] pub request_type: i32, - #[prost(string, tag = "12")] pub request_call_id: String, } -#[derive(Clone, PartialEq, Message)] +#[derive(Clone, Debug, PartialEq, Default)] pub struct ListResponse { - #[prost(message, repeated, tag = "1")] pub events: Vec, - #[prost(int64, tag = "2")] pub next_cursor: i64, - #[prost(bool, tag = "3")] pub has_more: bool, - #[prost(bool, tag = "4")] pub reset_required: bool, } -#[derive(Clone, PartialEq, Message)] +#[derive(Clone, Debug, PartialEq, Default)] pub struct SubmitRequest { - #[prost(message, optional, tag = "1")] pub query: Option, - #[prost(string, tag = "2")] pub channel_id: String, - #[prost(string, tag = "3")] pub body: String, - #[prost(string, tag = "4")] pub idempotency_key: String, } -#[derive(Clone, PartialEq, Message)] +#[derive(Clone, Debug, PartialEq, Default)] pub struct Turn { - #[prost(string, tag = "1")] pub turn_id: String, } -#[derive(Clone, PartialEq, Message)] +#[derive(Clone, Debug, PartialEq, Default)] pub struct SubmitResponse { - #[prost(message, optional, tag = "6")] pub accepted_turn: Option, - #[prost(int64, tag = "7")] pub replay_cursor: i64, } -#[derive(Clone, PartialEq, Message)] +#[derive(Clone, Debug, PartialEq, Default)] pub struct ThreadResponse { - #[prost(string, tag = "1")] pub request_id: String, - #[prost(string, tag = "2")] pub call_id: String, - #[prost(int32, tag = "3")] pub request_type: i32, - #[prost(int32, tag = "4")] pub action: i32, - #[prost(string, tag = "5")] pub text: String, - #[prost(string, tag = "7")] pub idempotency_key: String, } -#[derive(Clone, PartialEq, Message)] +#[derive(Clone, Debug, PartialEq, Default)] pub struct RespondRequest { - #[prost(message, optional, tag = "1")] pub query: Option, - #[prost(string, tag = "2")] pub channel_id: String, - #[prost(string, tag = "3")] pub turn_id: String, - #[prost(message, optional, tag = "4")] pub response: Option, - #[prost(string, tag = "5")] pub idempotency_key: String, } -#[derive(Clone, PartialEq, Message)] +#[derive(Clone, Debug, PartialEq, Default)] pub struct RespondResponse { - #[prost(int64, tag = "5")] pub replay_cursor: i64, } +impl From for Channel { + fn from(value: wire::Thread) -> Self { + Self { + id: value.id, + label: value.title, + unread_count: value.unread_count, + open_count: value.open_count, + archived: value.archived, + } + } +} + +impl From for OperatingMessage { + fn from(value: wire::TaskMessage) -> Self { + let role = match value.role { + 1 => "user", + 2 => "assistant", + 3 => "system", + _ => "unknown", + }; + Self { + id: value.id, + role: role.into(), + body: value.body, + } + } +} + +impl From for Event { + fn from(value: wire::TaskEvent) -> Self { + Self { + cursor: value.cursor, + event_id: value.id, + turn_id: value.turn_id, + kind: value.kind, + safe_text: value.text, + request_id: value.request_id, + request_type: value.request_kind, + request_call_id: value.call_id, + } + } +} + pub struct ThreadClient { http: Client, base: Url, @@ -267,6 +245,14 @@ impl ThreadClient { }) } + fn scope(&self) -> Result { + let query = self.query()?; + Ok(wire::Scope { + organization_id: query.organization_id, + workspace_id: query.workspace_id, + }) + } + async fn call( &self, method: &str, @@ -308,16 +294,22 @@ impl ThreadClient { } pub async fn get(&self) -> Result { - let result: GetResponse = self + let result: wire::GetThreadResponse = self .call( - "GetOperatingThread", - GetRequest { - query: Some(self.query()?), - channel_id: self.channel_id.clone(), + "GetThread", + wire::GetThreadRequest { + scope: Some(self.scope()?), + thread_id: self.channel_id.clone(), limit: 50, + ..Default::default() }, ) .await?; + let result = GetResponse { + channel: result.thread.map(Into::into), + messages: result.messages.into_iter().map(Into::into).collect(), + replay_cursor: result.replay_cursor, + }; if result .channel .as_ref() @@ -329,27 +321,35 @@ impl ThreadClient { } pub async fn list_channels(&self, archived: bool) -> Result { - self.call( - "ListOperatingChannels", - ListChannelsRequest { - query: Some(self.query()?), - archive_filter: if archived { 2 } else { 1 }, - }, - ) - .await + let result: wire::ListThreadsResponse = self + .call( + "ListThreads", + wire::ListThreadsRequest { + scope: Some(self.scope()?), + archived, + }, + ) + .await?; + Ok(ListChannelsResponse { + channels: result.threads.into_iter().map(Into::into).collect(), + }) } pub async fn rename(&self, title: String) -> Result { - let result: RenameThreadResponse = self + let result: wire::RenameThreadResponse = self .call( - "RenameOperatingThread", - RenameThreadRequest { - query: Some(self.query()?), - channel_id: self.channel_id.clone(), + "RenameThread", + wire::RenameThreadRequest { + scope: Some(self.scope()?), + thread_id: self.channel_id.clone(), title, }, ) .await?; + let result = RenameThreadResponse { + channel: result.thread.map(Into::into), + changed: result.changed, + }; if result .channel .as_ref() @@ -361,17 +361,22 @@ impl ThreadClient { } pub async fn archive(&self, archived: bool) -> Result { - let result: ArchiveThreadResponse = self + let result: wire::ArchiveThreadResponse = self .call( - "ArchiveOperatingThread", - ArchiveThreadRequest { - query: Some(self.query()?), - channel_id: self.channel_id.clone(), + "ArchiveThread", + wire::ArchiveThreadRequest { + scope: Some(self.scope()?), + thread_id: self.channel_id.clone(), archived, idempotency_key: Uuid::new_v4().to_string(), }, ) .await?; + let result = ArchiveThreadResponse { + channel: result.thread.map(Into::into), + changed: result.changed, + replayed: result.replayed, + }; if result .channel .as_ref() @@ -383,33 +388,47 @@ impl ThreadClient { } pub async fn events(&self, cursor: i64) -> Result { - self.call( - "ListOperatingThreadEvents", - ListRequest { - query: Some(self.query()?), - channel_id: self.channel_id.clone(), - after_cursor: cursor, - limit: 200, - }, - ) - .await + let result: wire::ListEventsResponse = self + .call( + "ListEvents", + wire::ListEventsRequest { + scope: Some(self.scope()?), + thread_id: self.channel_id.clone(), + after_cursor: cursor, + limit: 200, + }, + ) + .await?; + Ok(ListResponse { + events: result.events.into_iter().map(Into::into).collect(), + next_cursor: result.next_cursor, + has_more: result.has_more, + reset_required: result.reset_required, + }) } pub async fn submit(&self, body: String) -> Result { if body.trim().is_empty() || body.len() > 20_000 { bail!("message must contain 1 to 20000 bytes"); } - let result: SubmitResponse = self + let result: wire::SubmitTaskResponse = self .call( - "SubmitOperatingMessage", - SubmitRequest { - query: Some(self.query()?), - channel_id: self.channel_id.clone(), + "SubmitTask", + wire::SubmitTaskRequest { + scope: Some(self.scope()?), + thread_id: self.channel_id.clone(), body, idempotency_key: Uuid::new_v4().to_string(), + ..Default::default() }, ) .await?; + let result = SubmitResponse { + accepted_turn: result.accepted_turn.map(|turn| Turn { + turn_id: turn.turn_id, + }), + replay_cursor: result.replay_cursor, + }; if result .accepted_turn .as_ref() @@ -427,23 +446,25 @@ impl ThreadClient { text: String, ) -> Result { let key = Uuid::new_v4().to_string(); - self.call( - "RespondOperatingThread", - RespondRequest { - query: Some(self.query()?), - channel_id: self.channel_id.clone(), - turn_id: pending.turn_id.clone(), - response: Some(ThreadResponse { + let result: wire::RespondToRequestResponse = self + .call( + "RespondToRequest", + wire::RespondToRequestRequest { + scope: Some(self.scope()?), + thread_id: self.channel_id.clone(), + turn_id: pending.turn_id.clone(), request_id: pending.request_id.clone(), call_id: pending.request_call_id.clone(), - request_type: pending.request_type, + request_kind: pending.request_type, action, text, - idempotency_key: key.clone(), - }), - idempotency_key: key, - }, - ) - .await + idempotency_key: key, + ..Default::default() + }, + ) + .await?; + Ok(RespondResponse { + replay_cursor: result.replay_cursor, + }) } } diff --git a/packages/local-host-rs/src/lib.rs b/packages/local-host-rs/src/lib.rs index 6eabcccd6..c21c4a3c7 100644 --- a/packages/local-host-rs/src/lib.rs +++ b/packages/local-host-rs/src/lib.rs @@ -79,4 +79,8 @@ pub mod hosted_runner_cli; pub mod hosted_runner_conformance; pub mod subagents; +/// Generated bindings for the independently owned public application contract. +pub mod public_protocol { + include!(concat!(env!("OUT_DIR"), "/deixicpublic.v1.rs")); +} pub mod hosted_thread; diff --git a/packages/local-host-rs/src/managed_inference_readiness.rs b/packages/local-host-rs/src/managed_inference_readiness.rs index f2672fe85..a5c249aa2 100644 --- a/packages/local-host-rs/src/managed_inference_readiness.rs +++ b/packages/local-host-rs/src/managed_inference_readiness.rs @@ -1,16 +1,13 @@ //! Read-only, tenant-bound managed inference diagnostics. Never execution authority. use super::{CheckStatus, DoctorCheck}; use crate::credential_mode::{DetectedMode, ENVIRONMENT_ENV, PROVIDER_ENV, PlatformSession}; +use crate::public_protocol::{ + GetModelReadinessRequest, GetModelReadinessResponse, ModelNextAction as Action, + ModelReadinessState as Status, ModelSelection, Scope, +}; use anyhow::{Result, bail}; use prost::Message; -mod wire { - include!(concat!(env!("OUT_DIR"), "/console.v1.rs")); -} use std::{collections::HashMap, io::Read, time::Duration}; -use wire::{ - GetManagedInferenceReadinessRequest, GetManagedInferenceReadinessResponse, - ManagedInferenceNextAction as Action, ManagedInferenceReadinessStatus as Status, -}; pub(super) async fn check( readiness: &Result, @@ -57,29 +54,33 @@ fn request_for( session: &PlatformSession, requested: &str, env: &HashMap, -) -> Result { +) -> Result { let managed = session.managed_env(requested, env)?; let route = session.managed_model_route(requested); - Ok(GetManagedInferenceReadinessRequest { - organization_id: session.organization_id.clone(), - workspace_id: session - .workspace_id - .clone() - .filter(|id| !id.trim().is_empty()) - .ok_or_else(|| anyhow::anyhow!("workspace required"))?, - provider: managed - .get(PROVIDER_ENV) - .cloned() - .ok_or_else(|| anyhow::anyhow!("provider required"))?, - model: route.strip_prefix("evalops/").unwrap_or(&route).to_owned(), + Ok(GetModelReadinessRequest { + scope: Some(Scope { + organization_id: session.organization_id.clone(), + workspace_id: session + .workspace_id + .clone() + .filter(|id| !id.trim().is_empty()) + .ok_or_else(|| anyhow::anyhow!("workspace required"))?, + }), + selection: Some(ModelSelection { + provider: managed + .get(PROVIDER_ENV) + .cloned() + .ok_or_else(|| anyhow::anyhow!("provider required"))?, + model: route.strip_prefix("evalops/").unwrap_or(&route).to_owned(), + }), environment: managed.get(ENVIRONMENT_ENV).cloned().unwrap_or_default(), }) } fn fetch( session: &PlatformSession, - request: &GetManagedInferenceReadinessRequest, -) -> Result { + request: &GetModelReadinessRequest, +) -> Result { let base = crate::managed_setup::platform_base_url() .ok_or_else(|| anyhow::anyhow!("platform address required"))?; fetch_from(session, request, &base) @@ -87,21 +88,25 @@ fn fetch( fn fetch_from( session: &PlatformSession, - request: &GetManagedInferenceReadinessRequest, + request: &GetModelReadinessRequest, base: &str, -) -> Result { +) -> Result { let client = reqwest::blocking::Client::builder() .timeout(Duration::from_secs(5)) .redirect(reqwest::redirect::Policy::none()) .build()?; + let scope = request + .scope + .as_ref() + .ok_or_else(|| anyhow::anyhow!("scope required"))?; let response = client .post(format!( - "{}/deixic.v1.DeixicService/GetManagedInferenceReadiness", + "{}/deixicpublic.v1.DeixicPublicService/GetModelReadiness", base.trim_end_matches('/') )) .bearer_auth(&session.access_token) - .header("x-organization-id", &request.organization_id) - .header("x-workspace-id", &request.workspace_id) + .header("x-organization-id", &scope.organization_id) + .header("x-workspace-id", &scope.workspace_id) .header("connect-protocol-version", "1") .header("content-type", "application/proto") .header("accept", "application/proto") @@ -115,21 +120,18 @@ fn fetch_from( if bytes.len() > 65_536 { bail!("readiness response exceeds diagnostic limit"); } - let response = GetManagedInferenceReadinessResponse::decode(bytes.as_slice())?; + let response = GetModelReadinessResponse::decode(bytes.as_slice())?; if response.environment != request.environment - || response.organization_id != request.organization_id - || response.workspace_id != request.workspace_id - || response.target.as_ref().is_none_or(|target| { - target.provider != request.provider || target.model != request.model - }) + || response.scope != request.scope + || response.selection != request.selection { bail!("readiness scope mismatch"); } Ok(response) } -fn present(response: &GetManagedInferenceReadinessResponse) -> DoctorCheck { - let (status, label) = match Status::try_from(response.status).unwrap_or(Status::Unspecified) { +fn present(response: &GetModelReadinessResponse) -> DoctorCheck { + let (status, label) = match Status::try_from(response.state).unwrap_or(Status::Unspecified) { Status::Ready => (CheckStatus::Pass, "Ready"), Status::ActivationRequired => (CheckStatus::Warning, "Activation required"), Status::FundingRequired => (CheckStatus::Warning, "Funding required"), @@ -139,20 +141,28 @@ fn present(response: &GetManagedInferenceReadinessResponse) -> DoctorCheck { }; let mut details = vec![format!( "Deixic-managed inference: {} / {}; provider {}; model {}", - response.organization_id, - response.workspace_id, response - .target + .scope + .as_ref() + .map(|s| s.organization_id.as_str()) + .unwrap_or("unknown"), + response + .scope + .as_ref() + .map(|s| s.workspace_id.as_str()) + .unwrap_or("unknown"), + response + .selection .as_ref() .map(|target| target.provider.as_str()) .unwrap_or("unknown"), response - .target + .selection .as_ref() .map(|target| target.model.as_str()) .unwrap_or("unknown") )]; - if Status::try_from(response.status).unwrap_or(Status::Unspecified) == Status::Ready { + if Status::try_from(response.state).unwrap_or(Status::Unspecified) == Status::Ready { details.push("No provider key is required. Requests remain subject to your organization’s policy and available funding.".into()); } for action in &response.next_actions { @@ -191,8 +201,8 @@ mod tests { (Status::TemporarilyUnavailable, "Temporarily unavailable"), (Status::Unspecified, "Temporarily unavailable"), ] { - let response = GetManagedInferenceReadinessResponse { - status: state.into(), + let response = GetModelReadinessResponse { + state: state.into(), ..Default::default() }; let report = present(&response); @@ -227,10 +237,10 @@ mod tests { #[test] fn managed_inference_readiness_preserves_upstream_model_and_tenant() { let request = request_for(&session(), "evalops/openai/gpt-5.6", &HashMap::new()).unwrap(); - assert_eq!(request.provider, "openrouter"); - assert_eq!(request.model, "openai/gpt-5.6"); - assert_eq!(request.organization_id, "org-a"); - assert_eq!(request.workspace_id, "workspace-a"); + assert_eq!(request.selection.as_ref().unwrap().provider, "openrouter"); + assert_eq!(request.selection.as_ref().unwrap().model, "openai/gpt-5.6"); + assert_eq!(request.scope.as_ref().unwrap().organization_id, "org-a"); + assert_eq!(request.scope.as_ref().unwrap().workspace_id, "workspace-a"); } #[test] @@ -265,29 +275,26 @@ mod tests { continue; } assert!(headers.starts_with( - "post /deixic.v1.deixicservice/getmanagedinferencereadiness " + "post /deixicpublic.v1.deixicpublicservice/getmodelreadiness " )); assert!(headers.contains("content-type: application/proto")); assert!(headers.contains("connect-protocol-version: 1")); assert!(headers.contains("x-organization-id: org-a")); assert!(headers.contains("x-workspace-id: workspace-a")); - let request = - GetManagedInferenceReadinessRequest::decode(&bytes[end + 4..]).unwrap(); - assert_eq!(request.model, "openai/gpt-5.6"); - let response = GetManagedInferenceReadinessResponse { - organization_id: if matches { - request.organization_id - } else { - "org-other".into() - }, - workspace_id: request.workspace_id, + let request = GetModelReadinessRequest::decode(&bytes[end + 4..]).unwrap(); + assert_eq!(request.selection.as_ref().unwrap().model, "openai/gpt-5.6"); + let response = GetModelReadinessResponse { + scope: Some(Scope { + organization_id: if matches { + request.scope.as_ref().unwrap().organization_id.clone() + } else { + "org-other".into() + }, + workspace_id: request.scope.as_ref().unwrap().workspace_id.clone(), + }), environment: request.environment, - status: Status::Ready.into(), - target: wire::InferenceProviderTarget { - provider: request.provider, - model: request.model, - } - .into(), + state: Status::Ready.into(), + selection: request.selection, ..Default::default() } .encode_to_vec(); diff --git a/packages/local-host-rs/src/managed_setup/mod.rs b/packages/local-host-rs/src/managed_setup/mod.rs index 99b72a219..4eb2f5a59 100644 --- a/packages/local-host-rs/src/managed_setup/mod.rs +++ b/packages/local-host-rs/src/managed_setup/mod.rs @@ -40,7 +40,7 @@ use crate::path_utils; use crate::sandbox_policy::{SandboxPolicyDocument, TeamPolicyProvider, parse_policy_toml}; /// Connect RPC path for the managed setup read. -const GET_MANAGED_SETUP_PATH: &str = "/console.v1.ManagedSetupService/GetManagedSetup"; +const GET_MANAGED_SETUP_PATH: &str = "/deixicpublic.v1.DeixicPublicService/GetClientSetup"; /// Environment variables that name the Deixic platform base URL, in priority /// order. The managed-setup-specific name wins, followed by the shared @@ -924,8 +924,10 @@ fn fetch_managed_setup_from( .build() .map_err(|error| ManagedSetupError::Request(error.to_string()))?; let body = wire::GetManagedSetupRequest { - organization_id: session.organization_id.clone(), - workspace_id: session.workspace_id.clone().unwrap_or_default(), + scope: Some(crate::public_protocol::Scope { + organization_id: session.organization_id.clone(), + workspace_id: session.workspace_id.clone().unwrap_or_default(), + }), } .encode_to_vec(); let request = client diff --git a/packages/local-host-rs/src/managed_setup/tests.rs b/packages/local-host-rs/src/managed_setup/tests.rs index 43f58c338..e40a24c3d 100644 --- a/packages/local-host-rs/src/managed_setup/tests.rs +++ b/packages/local-host-rs/src/managed_setup/tests.rs @@ -752,8 +752,10 @@ fn protobuf_fixture() -> wire::ManagedSetup { seconds: 1_800_000_000, nanos: 123, }), - organization_id: "org-a".into(), - workspace_id: "workspace-a".into(), + scope: Some(crate::public_protocol::Scope { + organization_id: "org-a".into(), + workspace_id: "workspace-a".into(), + }), rules: vec![wire::ManagedRule { id: "rule".into(), title: "Review".into(), @@ -805,7 +807,7 @@ fn protobuf_server_for( } }; let headers = String::from_utf8_lossy(&request[..header_end]).to_ascii_lowercase(); - assert!(headers.starts_with("post /console.v1.managedsetupservice/getmanagedsetup ")); + assert!(headers.starts_with("post /deixicpublic.v1.deixicpublicservice/getclientsetup ")); assert!(headers.contains("content-type: application/proto\r\n")); assert!(headers.contains("accept: application/proto\r\n")); assert!(headers.contains("authorization: bearer access-token\r\n")); @@ -828,17 +830,17 @@ fn protobuf_server_for( request.extend_from_slice(&buffer[..count]); } let expected: &[u8] = if workspace_bound { - b"\x0a\x05org-a\x12\x0bworkspace-a" + b"\x0a\x14\x0a\x05org-a\x12\x0bworkspace-a" } else { - b"\x0a\x05org-a" + b"\x0a\x07\x0a\x05org-a" }; assert_eq!(&request[header_end..], expected); // Decode the actual bytes crossing the HTTP boundary, not the request builder. let decoded = wire::GetManagedSetupRequest::decode(&request[header_end..]).expect("protobuf request"); - assert_eq!(decoded.organization_id, "org-a"); + assert_eq!(decoded.scope.as_ref().unwrap().organization_id, "org-a"); assert_eq!( - decoded.workspace_id, + decoded.scope.as_ref().unwrap().workspace_id, if workspace_bound { "workspace-a" } else { "" } ); write!(stream, "HTTP/1.1 200 OK\r\nContent-Type: application/proto\r\nContent-Length: {}\r\nConnection: close\r\n\r\n", body.len()).expect("headers"); @@ -901,17 +903,20 @@ fn protobuf_transport_preserves_every_policy_field_and_cache_timestamp() { #[test] fn protobuf_transport_rejects_malformed_foreign_and_unknown_policy() { let mut foreign_org = protobuf_fixture(); - foreign_org.organization_id = "other".into(); + foreign_org.scope.as_mut().unwrap().organization_id = "other".into(); let mut foreign_workspace = protobuf_fixture(); - foreign_workspace.workspace_id = "other".into(); + foreign_workspace.scope.as_mut().unwrap().workspace_id = "other".into(); let mut unknown_mode = protobuf_fixture(); unknown_mode.mcp.as_mut().unwrap().mode = 99; let mut unknown_scope = protobuf_fixture(); unknown_scope.rules[0].scope = 99; + let mut missing_scope = protobuf_fixture(); + missing_scope.scope = None; let mut invalid_timestamp = protobuf_fixture(); invalid_timestamp.issued_at.as_mut().unwrap().nanos = -1; for body in [ vec![0x80], + missing_scope.encode_to_vec(), foreign_org.encode_to_vec(), foreign_workspace.encode_to_vec(), unknown_mode.encode_to_vec(), @@ -936,7 +941,7 @@ fn protobuf_transport_rejects_malformed_foreign_and_unknown_policy() { fn protobuf_transport_accepts_organization_policy_for_both_selectors() { for workspace_bound in [false, true] { let mut policy = protobuf_fixture(); - policy.workspace_id.clear(); + policy.scope.as_mut().unwrap().workspace_id.clear(); let (base, server) = protobuf_server_for(policy.encode_to_vec(), workspace_bound); let session = session_for("org-a", workspace_bound.then_some("workspace-a")); let setup = fetch_managed_setup_from(&session, &base).expect("organization policy"); diff --git a/packages/local-host-rs/src/managed_setup/wire.rs b/packages/local-host-rs/src/managed_setup/wire.rs index 056d3d2c8..fd47a7eef 100644 --- a/packages/local-host-rs/src/managed_setup/wire.rs +++ b/packages/local-host-rs/src/managed_setup/wire.rs @@ -1,81 +1,19 @@ -//! Bounded native projection of proto/console/v1/console.proto, matching the -//! existing bug_report adapter. Field tags follow that authoritative contract; -//! enum conversion below rejects unknown policy values rather than widening access. +//! Public client policy conversion rejects unknown values rather than widening access. use super::ManagedSetupError; -use prost::Message; - -#[derive(Clone, PartialEq, Message)] -pub(super) struct GetManagedSetupRequest { - #[prost(string, tag = "1")] - pub organization_id: String, - #[prost(string, tag = "2")] - pub workspace_id: String, -} - -#[derive(Clone, PartialEq, Message)] -pub(super) struct ManagedSetup { - #[prost(uint64, tag = "1")] - pub version: u64, - #[prost(message, optional, tag = "2")] - pub issued_at: Option, - #[prost(message, repeated, tag = "3")] - pub rules: Vec, - #[prost(message, repeated, tag = "4")] - pub skills: Vec, - #[prost(message, optional, tag = "5")] - pub mcp: Option, - #[prost(string, tag = "6")] - pub sandbox_policy_toml: String, - #[prost(string, tag = "7")] - pub organization_id: String, - #[prost(string, tag = "8")] - pub workspace_id: String, -} - -#[derive(Clone, PartialEq, Message)] -pub(super) struct ManagedRule { - #[prost(string, tag = "1")] - pub id: String, - #[prost(string, tag = "2")] - pub title: String, - #[prost(string, tag = "3")] - pub body_markdown: String, - #[prost(int32, tag = "4")] - pub scope: i32, -} - -#[derive(Clone, PartialEq, Message)] -pub(super) struct ManagedSkillRef { - #[prost(string, tag = "1")] - pub id: String, - #[prost(string, tag = "2")] - pub source: String, - #[prost(string, tag = "3")] - pub version: String, - #[prost(bool, tag = "4")] - pub required: bool, -} - -#[derive(Clone, PartialEq, Message)] -pub(super) struct McpPolicy { - #[prost(int32, tag = "1")] - pub mode: i32, - #[prost(message, repeated, tag = "2")] - pub servers: Vec, -} - -#[derive(Clone, PartialEq, Message)] -pub(super) struct McpServerRef { - #[prost(string, tag = "1")] - pub name: String, - #[prost(string, tag = "2")] - pub url_pattern: String, - #[prost(string, tag = "3")] - pub transport: String, -} +#[cfg(test)] +pub(super) use crate::public_protocol::{ + ClientMcpPolicy as McpPolicy, ClientMcpServer as McpServerRef, ClientRule as ManagedRule, + ClientSkill as ManagedSkillRef, +}; +pub(super) use crate::public_protocol::{ + GetClientSetupRequest as GetManagedSetupRequest, GetClientSetupResponse as ManagedSetup, +}; impl ManagedSetup { pub(super) fn try_into_domain(self) -> Result { + let scope = self + .scope + .ok_or_else(|| ManagedSetupError::Decode("missing client setup scope".into()))?; let issued_at = self .issued_at .map(|value| { @@ -128,8 +66,8 @@ impl ManagedSetup { Ok(super::ManagedSetup { version: self.version, issued_at, - organization_id: self.organization_id, - workspace_id: self.workspace_id, + organization_id: scope.organization_id, + workspace_id: scope.workspace_id, rules, skills: self .skills diff --git a/packages/tui-rs/src/thread_cli.rs b/packages/tui-rs/src/thread_cli.rs index c5e84a149..c7920c344 100644 --- a/packages/tui-rs/src/thread_cli.rs +++ b/packages/tui-rs/src/thread_cli.rs @@ -1,7 +1,7 @@ //! Attach the native terminal to a Platform-owned Dex operating thread. //! -//! Wire projections shared with the native desktop retain only fields these -//! clients read. Unknown fields stay opaque; Platform owns execution and approvals. +//! The shared native client uses the generated public application protocol. +//! Platform owns execution and approvals. use std::collections::{HashMap, HashSet}; use std::io::{self, IsTerminal, Write}; @@ -11,13 +11,10 @@ use anyhow::{Context, Result, bail}; #[cfg(test)] use maestro_local_host::credential_mode::PlatformSession; use maestro_local_host::credential_mode::verified_current_identity_session; -#[cfg(test)] -use maestro_local_host::hosted_thread::{ - Channel, GetRequest, GetResponse, ListRequest, ListResponse, OperatingMessage, RespondRequest, - RespondResponse, SubmitRequest, SubmitResponse, Turn, -}; use maestro_local_host::hosted_thread::{Event, ThreadClient}; #[cfg(test)] +use maestro_local_host::public_protocol as wire; +#[cfg(test)] use prost::Message; use serde::Serialize; use tokio::time::sleep; @@ -190,7 +187,7 @@ async fn follow( continue; } cursor = event.cursor; - if !event.request_id.is_empty() && matches!(event.kind, 4 | 5 | 14 | 15) { + if !event.request_id.is_empty() && matches!(event.kind, 4 | 5 | 10 | 11) { pending.insert(event.request_id.clone(), event.clone()); } emit( @@ -219,7 +216,7 @@ async fn follow( } return Ok(cursor); } - if event.turn_id == turn_id && matches!(event.kind, 4 | 5 | 14 | 15) { + if event.turn_id == turn_id && matches!(event.kind, 4 | 5 | 10 | 11) { return Ok(cursor); } } @@ -322,7 +319,7 @@ pub async fn run_thread(args: &[String]) -> Result { break; } for event in page.events { - if !event.request_id.is_empty() && matches!(event.kind, 4 | 5 | 14 | 15) { + if !event.request_id.is_empty() && matches!(event.kind, 4 | 5 | 10 | 11) { pending.insert(event.request_id.clone(), event.clone()); } if matches!(event.kind, 7..=9) { @@ -447,15 +444,18 @@ mod tests { } #[tokio::test] - async fn typed_owner_calls_keep_tenant_and_request_identity() { + async fn public_calls_keep_tenant_and_request_identity() { let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let base = format!("http://{}", listener.local_addr().unwrap()); let server = tokio::spawn(async move { for method in [ - "GetOperatingThread", - "SubmitOperatingMessage", - "ListOperatingThreadEvents", - "RespondOperatingThread", + "GetThread", + "ListThreads", + "RenameThread", + "ArchiveThread", + "SubmitTask", + "ListEvents", + "RespondToRequest", ] { let (mut socket, _) = listener.accept().await.unwrap(); let mut bytes = Vec::new(); @@ -469,7 +469,7 @@ mod tests { if let Some(end) = bytes.windows(4).position(|part| part == b"\r\n\r\n") { let headers = String::from_utf8_lossy(&bytes[..end]).to_ascii_lowercase(); assert!(headers.contains(&format!( - "/deixic.v1.deixicservice/{}", + "/deixicpublic.v1.deixicpublicservice/{}", method.to_ascii_lowercase() ))); assert!(headers.contains("authorization: bearer scoped-token")); @@ -494,67 +494,120 @@ mod tests { } let body = &bytes[body_offset..body_offset + body_length]; let result = match method { - "GetOperatingThread" => { - let request = GetRequest::decode(body).unwrap(); - assert_eq!(request.query.unwrap().organization_id, "org-a"); - assert_eq!(request.channel_id, "thread:one"); - GetResponse { - channel: Some(Channel { + "ListThreads" => { + let request = wire::ListThreadsRequest::decode(body).unwrap(); + assert_eq!(request.scope.unwrap().workspace_id, "ws-a"); + assert!(!request.archived); + wire::ListThreadsResponse { + threads: vec![wire::Thread { id: "thread:one".into(), - ..Channel::default() + title: "Prior title".into(), + unread_count: 2, + ..Default::default() + }], + } + .encode_to_vec() + } + "RenameThread" => { + let request = wire::RenameThreadRequest::decode(body).unwrap(); + assert_eq!(request.thread_id, "thread:one"); + assert_eq!(request.title, "New title"); + wire::RenameThreadResponse { + thread: Some(wire::Thread { + id: "thread:one".into(), + title: "New title".into(), + ..Default::default() }), - messages: vec![OperatingMessage { + changed: true, + } + .encode_to_vec() + } + "ArchiveThread" => { + let request = wire::ArchiveThreadRequest::decode(body).unwrap(); + assert_eq!(request.thread_id, "thread:one"); + assert!(request.archived); + assert!(!request.idempotency_key.is_empty()); + wire::ArchiveThreadResponse { + thread: Some(wire::Thread { + id: "thread:one".into(), + archived: true, + ..Default::default() + }), + changed: true, + ..Default::default() + } + .encode_to_vec() + } + "GetThread" => { + let request = wire::GetThreadRequest::decode(body).unwrap(); + assert_eq!(request.scope.unwrap().organization_id, "org-a"); + assert_eq!(request.thread_id, "thread:one"); + wire::GetThreadResponse { + thread: Some(wire::Thread { + id: "thread:one".into(), + ..Default::default() + }), + messages: vec![wire::TaskMessage { id: "msg-1".into(), - role: "user".into(), + role: wire::MessageRole::User as i32, body: "prior".into(), + ..Default::default() }], replay_cursor: 9, + ..Default::default() } .encode_to_vec() } - "SubmitOperatingMessage" => { - let request = SubmitRequest::decode(body).unwrap(); - assert_eq!(request.query.unwrap().workspace_id, "ws-a"); - assert_eq!(request.channel_id, "thread:one"); + "SubmitTask" => { + let request = wire::SubmitTaskRequest::decode(body).unwrap(); + assert_eq!(request.scope.unwrap().workspace_id, "ws-a"); + assert_eq!(request.thread_id, "thread:one"); assert_eq!(request.body, "continue"); assert!(!request.idempotency_key.is_empty()); - SubmitResponse { - accepted_turn: Some(Turn { + wire::SubmitTaskResponse { + accepted_turn: Some(wire::TaskTurn { turn_id: "turn-2".into(), + ..Default::default() }), replay_cursor: 10, + ..Default::default() } .encode_to_vec() } - "ListOperatingThreadEvents" => { - let request = ListRequest::decode(body).unwrap(); + "ListEvents" => { + let request = wire::ListEventsRequest::decode(body).unwrap(); assert_eq!(request.after_cursor, 9); - ListResponse { - events: vec![Event { + wire::ListEventsResponse { + events: vec![wire::TaskEvent { cursor: 10, - event_id: "event-10".into(), + id: "event-10".into(), turn_id: "turn-2".into(), kind: 4, - safe_text: "Review the action".into(), + text: "Review the action".into(), request_id: "request-1".into(), - request_type: 1, - request_call_id: "call-1".into(), + request_kind: 1, + call_id: "call-1".into(), + ..Default::default() }], next_cursor: 10, has_more: false, reset_required: false, + ..Default::default() } .encode_to_vec() } - "RespondOperatingThread" => { - let request = RespondRequest::decode(body).unwrap(); - assert_eq!(request.query.unwrap().organization_id, "org-a"); + "RespondToRequest" => { + let request = wire::RespondToRequestRequest::decode(body).unwrap(); + assert_eq!(request.scope.unwrap().organization_id, "org-a"); assert_eq!(request.turn_id, "turn-2"); - let response = request.response.unwrap(); - assert_eq!(response.request_id, "request-1"); - assert_eq!(response.action, 2); - assert_eq!(response.idempotency_key, request.idempotency_key); - RespondResponse { replay_cursor: 11 }.encode_to_vec() + assert_eq!(request.request_id, "request-1"); + assert_eq!(request.action, 2); + assert!(!request.idempotency_key.is_empty()); + wire::RespondToRequestResponse { + replay_cursor: 11, + ..Default::default() + } + .encode_to_vec() } _ => unreachable!(), }; @@ -582,6 +635,12 @@ mod tests { .unwrap(); let snapshot = client.get().await.unwrap(); assert_eq!(snapshot.messages[0].body, "prior"); + let listed = client.list_channels(false).await.unwrap(); + assert_eq!(listed.channels[0].unread_count, 2); + let renamed = client.rename("New title".into()).await.unwrap(); + assert_eq!(renamed.channel.unwrap().label, "New title"); + let archived = client.archive(true).await.unwrap(); + assert!(archived.channel.unwrap().archived); let submitted = client.submit("continue".into()).await.unwrap(); assert_eq!(submitted.accepted_turn.unwrap().turn_id, "turn-2"); let page = client.events(9).await.unwrap(); diff --git a/packages/tui-rs/tests/pty_e2e.rs b/packages/tui-rs/tests/pty_e2e.rs index 31ce0da0c..8eb424833 100644 --- a/packages/tui-rs/tests/pty_e2e.rs +++ b/packages/tui-rs/tests/pty_e2e.rs @@ -215,16 +215,16 @@ fn start_mock_managed_setup_server_with_gate( break; }; let request = read_request_body(&mut stream).expect("managed setup request"); - assert_eq!(request.as_bytes(), b"\x0a\x0bpty-e2e-org\x12\x11pty-e2e-workspace"); + assert_eq!(request.as_bytes(), b"\x0a\x20\x0a\x0bpty-e2e-org\x12\x11pty-e2e-workspace"); if let Some(gate) = gate.take() { // The test releases policy only after observing editable input. // Dropping the sender on assertion failure also releases the stub. let _ = gate.recv(); } - // Canonical console.v1.ManagedSetup tags: version=1, mcp=5, - // organization_id=7, workspace_id=8. The real native client + // Public GetClientSetupResponse: scope=1, version=2, mcp=6. + // The real native client // must decode protobuf here, exactly as it does with Platform. - let body = b"\x08\x01\x2a\x02\x08\x02\x3a\x0bpty-e2e-org\x42\x11pty-e2e-workspace"; + let body = b"\x0a\x20\x0a\x0bpty-e2e-org\x12\x11pty-e2e-workspace\x10\x01\x32\x02\x08\x02"; let response = format!( "HTTP/1.1 200 OK\r\ncontent-type: application/proto\r\ncontent-length: {}\r\nconnection: close\r\n\r\n", body.len() diff --git a/proto/console/v1/managed_inference.proto b/proto/console/v1/managed_inference.proto deleted file mode 100644 index f903660a6..000000000 --- a/proto/console/v1/managed_inference.proto +++ /dev/null @@ -1,43 +0,0 @@ -// @generated by scripts/contracts/generate-managed-inference-native.py. DO NOT EDIT. -// A diagnostic projection; unknown owner fields are ignored by protobuf decoding. -syntax = "proto3"; -package console.v1; - -enum ManagedInferenceReadinessStatus { - MANAGED_INFERENCE_READINESS_STATUS_UNSPECIFIED = 0; - MANAGED_INFERENCE_READINESS_STATUS_READY = 1; - MANAGED_INFERENCE_READINESS_STATUS_ACTIVATION_REQUIRED = 2; - MANAGED_INFERENCE_READINESS_STATUS_FUNDING_REQUIRED = 3; - MANAGED_INFERENCE_READINESS_STATUS_LIMIT_REACHED = 4; - MANAGED_INFERENCE_READINESS_STATUS_ACCESS_SUSPENDED = 5; - MANAGED_INFERENCE_READINESS_STATUS_TEMPORARILY_UNAVAILABLE = 6; -} - -enum ManagedInferenceNextAction { - MANAGED_INFERENCE_NEXT_ACTION_UNSPECIFIED = 0; - MANAGED_INFERENCE_NEXT_ACTION_CONTACT_ADMINISTRATOR = 1; - MANAGED_INFERENCE_NEXT_ACTION_VIEW_FUNDING = 2; - MANAGED_INFERENCE_NEXT_ACTION_RETRY = 3; -} - -message GetManagedInferenceReadinessRequest { - string organization_id = 1; - string workspace_id = 2; - string provider = 3; - string model = 4; - string environment = 5; -} - -message GetManagedInferenceReadinessResponse { - string organization_id = 1; - string workspace_id = 2; - ManagedInferenceReadinessStatus status = 3; - repeated ManagedInferenceNextAction next_actions = 5; - InferenceProviderTarget target = 10; - string environment = 12; -} - -message InferenceProviderTarget { - string provider = 1; - string model = 2; -} diff --git a/proto/deixicpublic/v1/sdk.proto b/proto/deixicpublic/v1/sdk.proto new file mode 100644 index 000000000..21cc80c26 --- /dev/null +++ b/proto/deixicpublic/v1/sdk.proto @@ -0,0 +1,604 @@ +syntax = "proto3"; + +package deixicpublic.v1; + +import "google/protobuf/timestamp.proto"; + +option go_package = "github.com/dx-corp/deixic-go/deixicpublic/v1;deixicpublicv1"; + +// This is the complete application SDK wire contract. Public messages are +// projections owned by this service, not aliases of implementation records. +service DeixicPublicService { + rpc GetThread(GetThreadRequest) returns (GetThreadResponse); + rpc ListThreads(ListThreadsRequest) returns (ListThreadsResponse); + rpc RenameThread(RenameThreadRequest) returns (RenameThreadResponse); + rpc ArchiveThread(ArchiveThreadRequest) returns (ArchiveThreadResponse); + rpc ListEvents(ListEventsRequest) returns (ListEventsResponse); + rpc WatchEvents(WatchEventsRequest) returns (stream WatchEventsResponse); + rpc SubmitTask(SubmitTaskRequest) returns (SubmitTaskResponse); + rpc InterruptTask(InterruptTaskRequest) returns (InterruptTaskResponse); + rpc RespondToRequest(RespondToRequestRequest) returns (RespondToRequestResponse); + rpc GetReceipt(GetReceiptRequest) returns (GetReceiptResponse); + rpc ResolveReceiptAction(ResolveReceiptActionRequest) returns (ResolveReceiptActionResponse); + rpc GetModelReadiness(GetModelReadinessRequest) returns (GetModelReadinessResponse); + rpc GetClientSetup(GetClientSetupRequest) returns (GetClientSetupResponse); + rpc SubmitIssueReport(SubmitIssueReportRequest) returns (SubmitIssueReportResponse); + rpc CreateBusinessObject(CreateBusinessObjectRequest) returns (CreateBusinessObjectResponse); + rpc GetBusinessObject(GetBusinessObjectRequest) returns (GetBusinessObjectResponse); + rpc UpdateBusinessObject(UpdateBusinessObjectRequest) returns (UpdateBusinessObjectResponse); + rpc DeleteBusinessObject(DeleteBusinessObjectRequest) returns (DeleteBusinessObjectResponse); +} + +message Scope { + string organization_id = 1; + string workspace_id = 2; +} + +message ErrorDetail { + string code = 1; + string message = 2; + string request_id = 3; + bool retryable = 4; +} + +enum TurnState { + TURN_STATE_UNSPECIFIED = 0; + TURN_STATE_ACCEPTED = 1; + TURN_STATE_RUNNING = 2; + TURN_STATE_WAITING = 3; + TURN_STATE_RESPONDED = 4; + TURN_STATE_COMPLETED = 5; + TURN_STATE_FAILED = 6; + TURN_STATE_INTERRUPTED = 7; +} + +enum WaitingReason { + WAITING_REASON_UNSPECIFIED = 0; + WAITING_REASON_APPROVAL = 1; + WAITING_REASON_USER_INPUT = 2; + WAITING_REASON_CLIENT_TOOL = 3; + WAITING_REASON_EXTERNAL_RETRY = 4; +} + +message TaskTurn { + string turn_id = 1; + int64 sequence = 2; + string user_message_id = 3; + string assistant_message_id = 4; + TurnState state = 5; + WaitingReason waiting_reason = 6; + int64 first_cursor = 7; + int64 last_cursor = 8; + ErrorDetail terminal_error = 9; + google.protobuf.Timestamp created_at = 10; + google.protobuf.Timestamp updated_at = 11; + google.protobuf.Timestamp completed_at = 12; +} + +message Thread { + string id = 1; + string title = 2; + int64 replay_cursor = 3; + int64 active_turn_sequence = 4; + google.protobuf.Timestamp updated_at = 5; + int32 unread_count = 6; + int32 open_count = 7; + bool archived = 8; +} + +enum MessageRole { + MESSAGE_ROLE_UNSPECIFIED = 0; + MESSAGE_ROLE_USER = 1; + MESSAGE_ROLE_ASSISTANT = 2; + MESSAGE_ROLE_SYSTEM = 3; +} + +message TaskMessage { + string id = 1; + MessageRole role = 2; + string body = 3; + google.protobuf.Timestamp created_at = 4; + repeated string receipt_ids = 5; +} + +enum EventKind { + EVENT_KIND_UNSPECIFIED = 0; + EVENT_KIND_TURN_ACCEPTED = 1; + EVENT_KIND_TURN_STARTED = 2; + EVENT_KIND_PROGRESS = 3; + EVENT_KIND_APPROVAL_REQUIRED = 4; + EVENT_KIND_INPUT_REQUIRED = 5; + EVENT_KIND_ARTIFACT_RECORDED = 6; + EVENT_KIND_TURN_COMPLETED = 7; + EVENT_KIND_TURN_FAILED = 8; + EVENT_KIND_TURN_INTERRUPTED = 9; + EVENT_KIND_CLIENT_TOOL_REQUIRED = 10; + EVENT_KIND_EXTERNAL_RETRY_REQUIRED = 11; +} + +enum RequestKind { + REQUEST_KIND_UNSPECIFIED = 0; + REQUEST_KIND_APPROVAL = 1; + REQUEST_KIND_USER_INPUT = 2; + REQUEST_KIND_CLIENT_TOOL = 3; + REQUEST_KIND_EXTERNAL_RETRY = 4; +} + +message EvidenceReference { + string kind = 1; + string id = 2; + string label = 3; + string uri = 4; +} + +message TaskEvent { + int64 cursor = 1; + string id = 2; + string turn_id = 3; + EventKind kind = 4; + string text = 5; + repeated EvidenceReference evidence = 6; + google.protobuf.Timestamp created_at = 7; + string request_id = 8; + RequestKind request_kind = 9; + string tool_name = 10; + string call_id = 11; + ErrorDetail terminal_error = 12; +} + +enum ReceiptState { + RECEIPT_STATE_UNSPECIFIED = 0; + RECEIPT_STATE_UNKNOWN = 1; + RECEIPT_STATE_PROPOSED = 2; + RECEIPT_STATE_QUEUED = 3; + RECEIPT_STATE_WAITING_APPROVAL = 4; + RECEIPT_STATE_BLOCKED = 5; + RECEIPT_STATE_NEEDS_INPUT = 6; + RECEIPT_STATE_UNAVAILABLE = 7; + RECEIPT_STATE_RUNNING = 8; + RECEIPT_STATE_SUCCEEDED = 9; + RECEIPT_STATE_VERIFIED = 10; + RECEIPT_STATE_DENIED = 11; + RECEIPT_STATE_CANCELLED = 12; + RECEIPT_STATE_FAILED = 13; +} + +message ReceiptAction { + string id = 1; + string label = 2; + string kind = 3; + bool requires_confirmation = 4; + string disabled_reason = 5; +} + +message CodingOutput { + string path = 1; + string object_id = 2; + string version_id = 3; + string sha256 = 4; + uint64 size_bytes = 5; + // Included only in an explicitly requested receipt read by the authorized actor. + bytes content = 6; +} + +message CodingAcceptance { + bool accepted = 1; + string contract_digest = 2; + string submission_digest = 3; + repeated string reasons = 4; + string revision = 5; + repeated CodingOutput outputs = 6; +} + +message Receipt { + string id = 1; + string kind = 2; + string title = 3; + string summary = 4; + ReceiptState state = 5; + google.protobuf.Timestamp updated_at = 6; + repeated ReceiptAction allowed_actions = 7; + repeated EvidenceReference evidence = 8; + CodingAcceptance coding_acceptance = 9; +} + +message ModelSelection { + string provider = 1; + string model = 2; +} + +message AvailableModel { + string provider = 1; + string model = 2; + bool ready = 3; + string display_name = 4; +} + +message SetupReadiness { + bool accessible = 1; + repeated string missing_requirements = 2; + ModelSelection selection = 3; + AvailableModel default_model = 4; + repeated AvailableModel available_models = 5; + string next_action = 6; +} + +// Model readiness is an application-level decision; billing and routing +// implementation details stay behind this projection. +enum ModelReadinessState { + MODEL_READINESS_STATE_UNSPECIFIED = 0; + MODEL_READINESS_STATE_READY = 1; + MODEL_READINESS_STATE_ACTIVATION_REQUIRED = 2; + MODEL_READINESS_STATE_FUNDING_REQUIRED = 3; + MODEL_READINESS_STATE_LIMIT_REACHED = 4; + MODEL_READINESS_STATE_ACCESS_SUSPENDED = 5; + MODEL_READINESS_STATE_TEMPORARILY_UNAVAILABLE = 6; +} + +enum ModelNextAction { + MODEL_NEXT_ACTION_UNSPECIFIED = 0; + MODEL_NEXT_ACTION_CONTACT_ADMINISTRATOR = 1; + MODEL_NEXT_ACTION_VIEW_FUNDING = 2; + MODEL_NEXT_ACTION_RETRY = 3; +} + +message GetModelReadinessRequest { + Scope scope = 1; + ModelSelection selection = 2; + string environment = 3; +} + +message GetModelReadinessResponse { + Scope scope = 1; + ModelSelection selection = 2; + string environment = 3; + ModelReadinessState state = 4; + repeated ModelNextAction next_actions = 5; + google.protobuf.Timestamp evaluated_at = 6; +} + +enum ClientRuleScope { + CLIENT_RULE_SCOPE_UNSPECIFIED = 0; + CLIENT_RULE_SCOPE_ORGANIZATION = 1; + CLIENT_RULE_SCOPE_WORKSPACE = 2; +} + +message ClientRule { + string id = 1; + string title = 2; + string body_markdown = 3; + ClientRuleScope scope = 4; +} + +message ClientSkill { + string id = 1; + string source = 2; + string version = 3; + bool required = 4; +} + +enum ClientMcpMode { + CLIENT_MCP_MODE_UNSPECIFIED = 0; + CLIENT_MCP_MODE_OPEN = 1; + CLIENT_MCP_MODE_ALLOWLIST = 2; + CLIENT_MCP_MODE_DENYLIST = 3; +} + +message ClientMcpServer { + string name = 1; + string url_pattern = 2; + string transport = 3; +} + +message ClientMcpPolicy { + ClientMcpMode mode = 1; + repeated ClientMcpServer servers = 2; +} + +message GetClientSetupRequest { + Scope scope = 1; +} + +message GetClientSetupResponse { + Scope scope = 1; + uint64 version = 2; + google.protobuf.Timestamp issued_at = 3; + repeated ClientRule rules = 4; + repeated ClientSkill skills = 5; + ClientMcpPolicy mcp = 6; + string sandbox_policy_toml = 7; +} + +message IssueEvidence { + string kind = 1; + string source_id = 2; + string text = 3; +} + +message IssueContext { + string reproduction_steps = 1; + string model = 2; + repeated IssueEvidence evidence = 3; +} + +message SubmitIssueReportRequest { + Scope scope = 1; + string description = 2; + string expected_behavior = 3; + string app_version = 4; + bool include_diagnostics = 5; + string idempotency_key = 6; + IssueContext context = 7; +} + +message IssueReport { + string id = 1; + string reference = 2; +} + +message SubmitIssueReportResponse { + IssueReport report = 1; +} + +// This resource is the supported Terraform provider contract. Source authority, +// connector state and recovery coordinates belong to the implementation. +message BusinessMoney { + string amount = 1; + string currency = 2; +} +message BusinessObjectReference { + string object_id = 1; +} +message BusinessArtifactReference { + string artifact_version_id = 1; +} +message BusinessTextList { + repeated string values = 1; +} +message BusinessFieldValue { + string field_id = 1; + oneof value { + string text = 2; + int64 integer = 3; + bool boolean = 4; + string decimal = 5; + BusinessMoney money = 6; + string enum_value = 7; + string date = 8; + string timestamp = 9; + BusinessObjectReference reference = 10; + BusinessArtifactReference artifact = 11; + BusinessTextList text_list = 12; + } +} +message BusinessObject { + string object_id = 1; + string type_id = 2; + int64 schema_revision = 3; + int64 revision = 4; + repeated BusinessFieldValue values = 5; + bool deleted = 6; + reserved 7, 8; + string updated_at = 9; + string created_at = 10; +} +message CreateBusinessObjectRequest { + string organization_id = 1; + string workspace_id = 2; + string idempotency_key = 3; + string type_id = 4; + int64 schema_revision = 5; + repeated BusinessFieldValue values = 6; +} +message CreateBusinessObjectResponse { + BusinessObject object = 1; +} +message GetBusinessObjectRequest { + string organization_id = 1; + string workspace_id = 2; + reserved 3; + string object_id = 4; + reserved 5; +} +message GetBusinessObjectResponse { + BusinessObject object = 1; +} +message UpdateBusinessObjectRequest { + string organization_id = 1; + string workspace_id = 2; + string idempotency_key = 3; + string object_id = 4; + int64 expected_revision = 5; + repeated BusinessFieldValue values = 6; + int64 schema_revision = 7; +} +message UpdateBusinessObjectResponse { + BusinessObject object = 1; +} +message DeleteBusinessObjectRequest { + string organization_id = 1; + string workspace_id = 2; + string idempotency_key = 3; + string object_id = 4; + int64 expected_revision = 5; +} +message DeleteBusinessObjectResponse { + BusinessObject object = 1; +} + +message GetThreadRequest { + Scope scope = 1; + string thread_id = 2; + int32 limit = 3; + string page_token = 4; +} + +message GetThreadResponse { + Thread thread = 1; + repeated TaskMessage messages = 2; + repeated Receipt receipts = 3; + repeated TaskTurn turns = 4; + int64 replay_cursor = 5; + string next_page_token = 6; + SetupReadiness setup = 7; +} + +message ListThreadsRequest { + Scope scope = 1; + bool archived = 2; +} + +message ListThreadsResponse { + repeated Thread threads = 1; +} + +message RenameThreadRequest { + Scope scope = 1; + string thread_id = 2; + string title = 3; +} + +message RenameThreadResponse { + Thread thread = 1; + bool changed = 2; +} + +message ArchiveThreadRequest { + Scope scope = 1; + string thread_id = 2; + bool archived = 3; + string idempotency_key = 4; +} + +message ArchiveThreadResponse { + Thread thread = 1; + bool changed = 2; + bool replayed = 3; +} + +message ListEventsRequest { + Scope scope = 1; + string thread_id = 2; + int64 after_cursor = 3; + int32 limit = 4; +} + +message ListEventsResponse { + repeated TaskEvent events = 1; + int64 next_cursor = 2; + bool has_more = 3; + bool reset_required = 4; + Thread snapshot = 5; + repeated TaskTurn snapshot_turns = 6; +} + +message WatchEventsRequest { + Scope scope = 1; + string thread_id = 2; + int64 after_cursor = 3; +} + +message WatchEventsResponse { + repeated TaskEvent events = 1; + int64 next_cursor = 2; + bool reset_required = 3; + Thread snapshot = 4; + repeated TaskTurn snapshot_turns = 5; +} + +message CodingContract { + string repository_id = 1; + uint64 generation = 2; + repeated string required_assertion_ids = 3; + bool require_review = 4; + bool require_behavior = 5; + repeated string readiness_requirements = 6; + repeated string authorized_skips = 7; + repeated string authorized_dispositions = 8; + repeated string output_paths = 9; +} + +message SubmitTaskRequest { + Scope scope = 1; + string thread_id = 2; + string body = 3; + string idempotency_key = 4; + string continuation_task_id = 5; + repeated string reference_task_ids = 6; + ModelSelection model_selection = 7; + CodingContract coding_contract = 8; + string project_resource_id = 9; +} + +message SubmitTaskResponse { + TaskMessage message = 1; + TaskTurn accepted_turn = 2; + int64 replay_cursor = 3; + repeated Receipt receipts = 4; +} + +message InterruptTaskRequest { + Scope scope = 1; + string thread_id = 2; + string turn_id = 3; + string idempotency_key = 4; + string reason = 5; +} + +message InterruptTaskResponse { + string turn_id = 1; + bool replayed = 2; + repeated TaskTurn turns = 3; + int64 replay_cursor = 4; +} + +enum ResponseAction { + RESPONSE_ACTION_UNSPECIFIED = 0; + RESPONSE_ACTION_APPROVE = 1; + RESPONSE_ACTION_DENY = 2; + RESPONSE_ACTION_ANSWER = 3; + RESPONSE_ACTION_RETRY = 4; + RESPONSE_ACTION_SKIP = 5; + RESPONSE_ACTION_ABORT = 6; +} + +message RespondToRequestRequest { + Scope scope = 1; + string thread_id = 2; + string turn_id = 3; + string request_id = 4; + string call_id = 5; + RequestKind request_kind = 6; + ResponseAction action = 7; + string text = 8; + bool is_error = 9; + string idempotency_key = 10; +} + +message RespondToRequestResponse { + bool replayed = 1; + repeated TaskTurn turns = 2; + int64 replay_cursor = 3; +} + +message GetReceiptRequest { + Scope scope = 1; + string receipt_id = 2; + string thread_id = 3; + bool include_coding_output_content = 4; +} + +message GetReceiptResponse { + Receipt receipt = 1; +} + +message ResolveReceiptActionRequest { + Scope scope = 1; + string receipt_id = 2; + string action_id = 3; + string idempotency_key = 4; +} + +message ResolveReceiptActionResponse { + Receipt receipt = 1; +} diff --git a/scripts/build-release-binary.mjs b/scripts/build-release-binary.mjs index 19c94357c..3eef29f0a 100755 --- a/scripts/build-release-binary.mjs +++ b/scripts/build-release-binary.mjs @@ -1,6 +1,7 @@ #!/usr/bin/env node +import { assertPublicProtocolArtifact } from "./check-public-protocol-artifact.mjs"; import { execFileSync } from "node:child_process"; -import { chmodSync, copyFileSync, mkdirSync, statSync } from "node:fs"; +import { chmodSync, copyFileSync, mkdirSync, statSync, readFileSync } from "node:fs"; import { dirname, resolve } from "node:path"; const TARGETS = new Map([ @@ -52,4 +53,5 @@ const source = resolve(`target/${target}/${options.profile}/maestro`); mkdirSync(dirname(outfile), { recursive: true }); copyFileSync(source, outfile); chmodSync(outfile, 0o755); +assertPublicProtocolArtifact(readFileSync(outfile), outfile); console.log(`Built native release ${outfile} (${statSync(outfile).size} bytes).`); diff --git a/scripts/check-public-protocol-artifact.mjs b/scripts/check-public-protocol-artifact.mjs new file mode 100644 index 000000000..9e3b641ce --- /dev/null +++ b/scripts/check-public-protocol-artifact.mjs @@ -0,0 +1,24 @@ +#!/usr/bin/env node +import { readFileSync } from 'node:fs'; +import { resolve } from 'node:path'; +import { pathToFileURL } from 'node:url'; + +// This supplements canonical source parity: prost does not embed a descriptor +// set. Inspect the exact native bytes for private schema and RPC namespaces. +const privateMarkers = [ + 'console.v1', 'console/v1/', 'deixic.v1.DeixicService', + 'evalops.console', 'platform.v1', 'platform/v1/', +]; +export function assertPublicProtocolArtifact(bytes, name = 'native artifact') { + if (bytes.length === 0) throw new Error(`Empty ${name}`); + for (const marker of privateMarkers) { + if (bytes.includes(Buffer.from(marker))) { + throw new Error(`Private protocol marker ${marker} in ${name}`); + } + } +} +if (process.argv[1] && import.meta.url === pathToFileURL(resolve(process.argv[1])).href) { + if (process.argv.length < 3) throw new Error('Usage: check-public-protocol-artifact.mjs BINARY...'); + for (const path of process.argv.slice(2)) assertPublicProtocolArtifact(readFileSync(path), path); + console.log('Native public protocol artifact audit passed'); +} diff --git a/scripts/check-public-protocol-artifact.test.mjs b/scripts/check-public-protocol-artifact.test.mjs new file mode 100644 index 000000000..014e8e72d --- /dev/null +++ b/scripts/check-public-protocol-artifact.test.mjs @@ -0,0 +1,11 @@ +import test from 'node:test'; +import assert from 'node:assert/strict'; +import { assertPublicProtocolArtifact } from './check-public-protocol-artifact.mjs'; + +test('allows public native RPC bytes and rejects private schema and route bytes', () => { + assertPublicProtocolArtifact(Buffer.from('\0/deixicpublic.v1.DeixicPublicService/GetThread\0maestro.v1\0')); + for (const marker of ['console.v1', 'console/v1/console.proto', 'deixic.v1.DeixicService', 'platform/v1/internal.proto']) { + assert.throws(() => assertPublicProtocolArtifact(Buffer.from(`\0${marker}\0`)), /Private protocol/); + } + assert.throws(() => assertPublicProtocolArtifact(Buffer.alloc(0)), /Empty/); +}); diff --git a/scripts/test-capture-tui.py b/scripts/test-capture-tui.py index d67daf645..82044f245 100644 --- a/scripts/test-capture-tui.py +++ b/scripts/test-capture-tui.py @@ -80,9 +80,9 @@ def post(path, body=None): identity = json.loads(post("/v1/tokens/introspect")) self.assertEqual(identity["organization_id"], "capture-org") self.assertEqual(identity["workspace_id"], "capture-workspace") - policy = post("/console.v1.ManagedSetupService/GetManagedSetup") + policy = post("/deixicpublic.v1.DeixicPublicService/GetClientSetup") self.assertEqual( - policy, b"\x08\x01\x2a\x02\x08\x02\x3a\x0bcapture-org" + policy, b"\x0a\x0d\x0a\x0bcapture-org\x10\x01\x32\x02\x08\x02" ) self.assertIn("tool_calls", post("/v1/chat/completions")) self.assertIn( diff --git a/scripts/tui_capture_fixture.py b/scripts/tui_capture_fixture.py index 162666794..b096f23a0 100644 --- a/scripts/tui_capture_fixture.py +++ b/scripts/tui_capture_fixture.py @@ -181,9 +181,9 @@ def chunk(content, reason=None): }} body = "data: " + json.dumps(event) + "\n\n" content_type = "text/event-stream" - elif self.path == "/console.v1.ManagedSetupService/GetManagedSetup": + elif self.path == "/deixicpublic.v1.DeixicPublicService/GetClientSetup": # ManagedSetup v1, empty MCP allowlist, capture organization. - body = b"\x08\x01\x2a\x02\x08\x02\x3a\x0bcapture-org" + body = b"\x0a\x0d\x0a\x0bcapture-org\x10\x01\x32\x02\x08\x02" content_type = "application/proto" elif self.path == "/v1/tokens/introspect": fixture.identity_requests += 1 diff --git a/scripts/verify-staged-release.mjs b/scripts/verify-staged-release.mjs index aff231eb0..5bbb3b308 100644 --- a/scripts/verify-staged-release.mjs +++ b/scripts/verify-staged-release.mjs @@ -4,6 +4,7 @@ import { execFileSync } from 'node:child_process'; import { readFileSync, lstatSync, existsSync } from 'node:fs'; import { resolve, join } from 'node:path'; import { pathToFileURL } from 'node:url'; +import { assertPublicProtocolArtifact } from './check-public-protocol-artifact.mjs'; import { verifySourceManifest } from './release-source-manifest.mjs'; export const platforms = ['linux-x64', 'linux-arm64', 'darwin-x64', 'darwin-arm64']; @@ -29,6 +30,10 @@ export function verifyStagedFiles(dir, version, sourceRoot = process.cwd()) { if (sums.get(name) !== digest) throw new Error(`Checksum mismatch: ${name}`); }; for (const name of stagedFiles) verifyFile(name); + for (const platform of platforms) { + const name = `maestro-${platform}`; + assertPublicProtocolArtifact(readFileSync(join(dir, name)), name); + } const metadata = JSON.parse(readFileSync(join(dir, 'release-metadata.json'), 'utf8')); if (metadata.version !== version || metadata.releaseTag !== `v${version}` || !/^[a-f0-9]{40}$/.test(metadata.receipt?.sourceSha ?? '')) { throw new Error('Staged release version or source does not match'); diff --git a/scripts/verify-staged-release.test.mjs b/scripts/verify-staged-release.test.mjs index d1c177fc7..dd9269a2a 100644 --- a/scripts/verify-staged-release.test.mjs +++ b/scripts/verify-staged-release.test.mjs @@ -94,3 +94,9 @@ test('Code device capability must be a typed per-platform receipt', t => { const {dir,seal}=fixture(t); writeFileSync(join(dir,'code-device-darwin-arm64.json'),JSON.stringify({schemaVersion:1,platform:'darwin-arm64',enabled:'false'})); seal(); assert.throws(() => verifyStagedFiles(dir,'0.10.72',dir), /Invalid Code device/); }); + +test('rejects a sealed native binary containing a private protocol', t => { + const {dir,seal}=fixture(t); + writeFileSync(join(dir,'maestro-linux-x64'), 'binary\0console.v1\0'); seal(); + assert.throws(() => verifyStagedFiles(dir,'0.10.72',dir), /Private protocol/); +}); diff --git a/test/fixtures/product-issue-report-public-v1.hex b/test/fixtures/product-issue-report-public-v1.hex new file mode 100644 index 000000000..e77c3b03a --- /dev/null +++ b/test/fixtures/product-issue-report-public-v1.hex @@ -0,0 +1 @@ +0a140a056f72672d31120b776f726b73706163652d311220546865207465726d696e616c2073746f7070656420726573706f6e64696e672e1a1b546865206e657874207475726e2073686f756c642073746172742e221044656978696320436f64652074657374280132086e61746976652d31 diff --git a/test/fixtures/product-issue-report-public-v2.hex b/test/fixtures/product-issue-report-public-v2.hex new file mode 100644 index 000000000..34f6af718 --- /dev/null +++ b/test/fixtures/product-issue-report-public-v2.hex @@ -0,0 +1 @@ +0a140a056f72672d31120b776f726b73706163652d311220546865207465726d696e616c2073746f7070656420726573706f6e64696e672e1a1b546865206e657874207475726e2073686f756c642073746172742e221044656978696320436f64652074657374280132086e61746976652d313a460a13526570656174206661696c696e6720746f6f6c120a746573742d6d6f64656c1a230a0b746f6f6c5f726573756c74120663616c6c2d311a0c77726f6e6720616374696f6e