Skip to content
1 change: 1 addition & 0 deletions crates/ironclaw_conversations/src/inbound.rs
Original file line number Diff line number Diff line change
Expand Up @@ -243,6 +243,7 @@ where
let turn_submission_result = self
.turn_coordinator
.submit_turn(SubmitTurnRequest {
requested_model: None,
scope: resolution.turn_scope.clone(),
actor: accepted_message.actor.clone(),
accepted_message_ref: accepted_message.message_ref.clone(),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2264,6 +2264,7 @@ pub(crate) fn http_without_body_then_operation_failed_wat() -> String {
#[cfg(feature = "libsql")]
pub(crate) fn submit_turn_request(thread: &str, idempotency_key: &str) -> SubmitTurnRequest {
SubmitTurnRequest {
requested_model: None,
scope: TurnScope::new(
TenantId::new("tenant1").unwrap(),
Some(AgentId::new("agent1").unwrap()),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ fn actor() -> TurnActor {

fn submit_request(thread: &str, idempotency_key: &str) -> SubmitTurnRequest {
SubmitTurnRequest {
requested_model: None,
scope: scope(thread),
actor: actor(),
accepted_message_ref: AcceptedMessageRef::new(format!("message-{idempotency_key}"))
Expand Down
100 changes: 98 additions & 2 deletions crates/ironclaw_product_adapters/src/inbound.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ use crate::outbound::ProjectionCursor;
use crate::redaction::RedactedString;

const USER_MESSAGE_TEXT_MAX_BYTES: usize = 64 * 1024;
const REQUESTED_MODEL_MAX_BYTES: usize = 256;
const COMMAND_MAX_BYTES: usize = 256;
const COMMAND_ARGUMENTS_MAX_BYTES: usize = 64 * 1024;
const THREAD_HINT_MAX_BYTES: usize = 512;
Expand Down Expand Up @@ -99,6 +100,13 @@ pub struct UserMessagePayload {
pub text: String,
pub attachments: Vec<ProductAttachmentDescriptor>,
pub trigger: ProductTriggerReason,
/// Caller-requested model for this turn (e.g. an OpenAI-compatible client's
/// `model` field). A model *hint*, not authority: the coordinator routes to
/// it only when the operator has it configured, otherwise it falls back to
/// the deployment's active model. `None` for surfaces that don't select a
/// model (chat UI, channels).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub requested_model: Option<String>,
}

impl UserMessagePayload {
Expand All @@ -111,13 +119,25 @@ impl UserMessagePayload {
text: text.into(),
attachments,
trigger,
requested_model: None,
};
payload.validate()?;
Ok(payload)
}

/// Attach a caller-requested model to this payload. See
/// [`UserMessagePayload::requested_model`].
pub fn with_requested_model(mut self, requested_model: Option<String>) -> Self {
self.requested_model = requested_model.filter(|model| !model.is_empty());
self
}

pub fn validate(&self) -> Result<(), ProductAdapterError> {
validate_payload_string("user message text", &self.text, USER_MESSAGE_TEXT_MAX_BYTES)
validate_payload_string("user message text", &self.text, USER_MESSAGE_TEXT_MAX_BYTES)?;
if let Some(model) = &self.requested_model {
validate_payload_string("requested model", model, REQUESTED_MODEL_MAX_BYTES)?;
}
Ok(())
}
}

Expand All @@ -126,6 +146,8 @@ struct UserMessagePayloadWire {
text: String,
attachments: Vec<ProductAttachmentDescriptor>,
trigger: ProductTriggerReason,
#[serde(default)]
requested_model: Option<String>,
}

impl<'de> Deserialize<'de> for UserMessagePayload {
Expand All @@ -134,7 +156,14 @@ impl<'de> Deserialize<'de> for UserMessagePayload {
D: Deserializer<'de>,
{
let wire = UserMessagePayloadWire::deserialize(deserializer)?;
Self::new(wire.text, wire.attachments, wire.trigger).map_err(serde::de::Error::custom)
let payload = Self::new(wire.text, wire.attachments, wire.trigger)
.map(|payload| payload.with_requested_model(wire.requested_model))
.map_err(serde::de::Error::custom)?;
// `new` validated the payload while `requested_model` was still `None`;
// re-validate the assembled value so the wire-supplied model hint is
// bounded like every other ingress field (bypass flagged in PR review).
payload.validate().map_err(serde::de::Error::custom)?;
Ok(payload)
}
}

Expand Down Expand Up @@ -882,6 +911,73 @@ mod tests {
use crate::auth::AuthRequirement;
use crate::external::{ExternalActorRef, ExternalConversationRef, ExternalEventId};

#[test]
fn user_message_payload_round_trips_and_filters_requested_model() {
let with_model = UserMessagePayload::new("hi", vec![], ProductTriggerReason::DirectChat)
.unwrap()
.with_requested_model(Some("gpt-4o".to_string()));
assert_eq!(with_model.requested_model.as_deref(), Some("gpt-4o"));
// Round-trips over the wire (custom Deserialize via the wire struct).
let decoded: UserMessagePayload =
serde_json::from_str(&serde_json::to_string(&with_model).unwrap()).unwrap();
assert_eq!(decoded.requested_model.as_deref(), Some("gpt-4o"));

// Omitted → None, and not serialized when absent.
let without =
UserMessagePayload::new("hi", vec![], ProductTriggerReason::DirectChat).unwrap();
assert!(without.requested_model.is_none());
assert!(
!serde_json::to_string(&without)
.unwrap()
.contains("requested_model")
);

// An empty requested model is filtered to None.
assert!(
UserMessagePayload::new("hi", vec![], ProductTriggerReason::DirectChat)
.unwrap()
.with_requested_model(Some(String::new()))
.requested_model
.is_none()
);
}

#[test]
fn user_message_payload_bounds_requested_model_on_every_path() {
let over_limit = "m".repeat(REQUESTED_MODEL_MAX_BYTES + 1);

// Explicit validation after the builder rejects an over-long hint.
let built = UserMessagePayload::new("hi", vec![], ProductTriggerReason::DirectChat)
.unwrap()
.with_requested_model(Some(over_limit.clone()));
assert!(built.validate().is_err());

// Deserialization must not smuggle an unbounded hint past validation:
// the wire path attaches `requested_model` after `new`, so it re-validates.
let wire = serde_json::json!({
"text": "hi",
"attachments": [],
"trigger": "direct_chat",
"requested_model": over_limit,
})
.to_string();
let decoded: Result<UserMessagePayload, _> = serde_json::from_str(&wire);
assert!(
decoded.is_err(),
"an over-long requested_model must be rejected during deserialization"
);

// A hint at the cap is accepted on both paths.
let at_cap = "m".repeat(REQUESTED_MODEL_MAX_BYTES);
assert!(
UserMessagePayload::new("hi", vec![], ProductTriggerReason::DirectChat)
.unwrap()
.with_requested_model(Some(at_cap))
.validate()
.is_ok()
);
}

fn sample_context() -> TrustedInboundContext {
let evidence = ProtocolAuthEvidence::test_verified(
AuthRequirement::SharedSecretHeader {
Expand Down
1 change: 1 addition & 0 deletions crates/ironclaw_product_workflow/src/auth_continuation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -795,6 +795,7 @@ mod tests {
let actor = TurnActor::new(UserId::new("alice").unwrap());
let submit = coordinator
.submit_turn(SubmitTurnRequest {
requested_model: None,
scope: scope.clone(),
actor: actor.clone(),
accepted_message_ref: AcceptedMessageRef::new("message-auth-real").unwrap(),
Expand Down
8 changes: 8 additions & 0 deletions crates/ironclaw_product_workflow/src/inbound_turn.rs
Original file line number Diff line number Diff line change
Expand Up @@ -491,6 +491,7 @@ where
received_at: envelope.received_at(),
adapter_id: prepared.adapter_id,
surface_type: prepared.surface_type,
requested_model: payload.requested_model.clone(),
}))
.submit_or_replay(&self.thread_service, &self.turn_coordinator)
.await
Expand Down Expand Up @@ -654,6 +655,10 @@ impl ProductInboundTurnHandoff {
received_at,
adapter_id,
surface_type,
// The requested model is not persisted in the message store, so an
// idempotent resubmission of an accepted message falls back to the
// deployment's active model rather than recovering the original hint.
requested_model: None,
},
)))
}
Expand Down Expand Up @@ -703,6 +708,7 @@ struct AcceptedProductInboundTurn {
received_at: DateTime<Utc>,
adapter_id: ProductAdapterId,
surface_type: TurnSurfaceType,
requested_model: Option<String>,
}

impl AcceptedProductInboundTurn {
Expand All @@ -725,6 +731,7 @@ impl AcceptedProductInboundTurn {
received_at,
adapter_id,
surface_type,
requested_model,
} = self;
let turn_scope = TurnScope::new_with_owner(
binding.tenant_id.clone(),
Expand Down Expand Up @@ -779,6 +786,7 @@ impl AcceptedProductInboundTurn {
source_binding_ref,
reply_target_binding_ref,
requested_run_profile: None,
requested_model,
idempotency_key,
received_at,
requested_run_id: None,
Expand Down
1 change: 1 addition & 0 deletions crates/ironclaw_product_workflow/src/reborn_services.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3755,6 +3755,7 @@ impl RebornServicesApi for RebornServices {
)?;
let product_context = ironclaw_product_context::resolve_web_ui(scope.product_owner(&actor));
let submit = SubmitTurnRequest {
requested_model: None,
scope: scope.clone(),
actor,
accepted_message_ref: accepted_message_ref.clone(),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2995,6 +2995,7 @@ async fn before_inbound_policy_path_probes_replay_once() {
async fn before_inbound_policy_rewrite_revalidates_payload_before_turn_path() {
let (workflow, inbound, ledger, policy) = build_workflow_with_policy();
policy.rewrite_user_message(UserMessagePayload {
requested_model: None,
text: "a".repeat(64 * 1024 + 1),
attachments: vec![],
trigger: ProductTriggerReason::DirectChat,
Expand Down
2 changes: 2 additions & 0 deletions crates/ironclaw_reborn_composition/src/factory.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6709,6 +6709,7 @@ mod tests {
Some(owner.clone()),
);
let submit = ironclaw_turns::SubmitTurnRequest {
requested_model: None,
scope,
actor: ironclaw_turns::TurnActor::new(owner),
accepted_message_ref: ironclaw_turns::AcceptedMessageRef::new("configured-message-ref")
Expand Down Expand Up @@ -6832,6 +6833,7 @@ mod tests {
Some(owner.clone()),
);
let submit = ironclaw_turns::SubmitTurnRequest {
requested_model: None,
scope,
actor: ironclaw_turns::TurnActor::new(owner),
accepted_message_ref: ironclaw_turns::AcceptedMessageRef::new("default-message-ref")
Expand Down
3 changes: 3 additions & 0 deletions crates/ironclaw_reborn_composition/src/factory/auth_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -193,6 +193,7 @@ async fn local_dev_oauth_turn_gate_callback_resumes_default_turn_coordinator() {
let actor = TurnActor::new(UserId::new("alice").unwrap());
let submit = turn_coordinator
.submit_turn(SubmitTurnRequest {
requested_model: None,
scope: scope.clone(),
actor: actor.clone(),
accepted_message_ref: AcceptedMessageRef::new("message-auth-callback").unwrap(),
Expand Down Expand Up @@ -956,6 +957,7 @@ async fn submit_and_block_provider_auth_run(
) -> TurnRunId {
let submit = turn_coordinator
.submit_turn(SubmitTurnRequest {
requested_model: None,
scope: scope.clone(),
actor,
accepted_message_ref: AcceptedMessageRef::new(format!("message-fanout-{suffix}"))
Expand Down Expand Up @@ -1075,6 +1077,7 @@ async fn submit_and_block_auth_run(
) -> ironclaw_turns::TurnRunId {
let submit = turn_coordinator
.submit_turn(SubmitTurnRequest {
requested_model: None,
scope: scope.clone(),
actor,
accepted_message_ref: AcceptedMessageRef::new("message-auth-callback-2").unwrap(),
Expand Down
4 changes: 4 additions & 0 deletions crates/ironclaw_reborn_composition/src/runtime.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2346,6 +2346,7 @@ impl RebornRuntime {
let response = match self
.turn_coordinator
.submit_turn(SubmitTurnRequest {
requested_model: None,
scope: scope.clone(),
actor: TurnActor::new(self.actor_user_id.clone()),
accepted_message_ref: accepted_message_ref.clone(),
Expand Down Expand Up @@ -8177,6 +8178,7 @@ output_schema_ref = "schemas/write.output.json"
let parent = runtime
.turn_coordinator
.submit_turn(SubmitTurnRequest {
requested_model: None,
scope: parent_scope.clone(),
actor: actor.clone(),
accepted_message_ref: AcceptedMessageRef::new("msg:cancel-parent").unwrap(),
Expand Down Expand Up @@ -10521,6 +10523,7 @@ output_schema_ref = "schemas/write.output.json"
let submitted = runtime
.turn_coordinator
.submit_turn(SubmitTurnRequest {
requested_model: None,
scope: scope.clone(),
actor: actor.clone(),
accepted_message_ref: AcceptedMessageRef::new("msg:audit").unwrap(),
Expand Down Expand Up @@ -11102,6 +11105,7 @@ output_schema_ref = "schemas/write.output.json"
let submitted_a = runtime
.turn_coordinator
.submit_turn(SubmitTurnRequest {
requested_model: None,
scope: scope.clone(),
actor: actor.clone(),
accepted_message_ref: AcceptedMessageRef::new("msg:rejected-busy-a").unwrap(),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -199,6 +199,7 @@ async fn submit_and_block_auth_run(
.turn_state
.submit_turn(
SubmitTurnRequest {
requested_model: None,
scope: scope.clone(),
actor,
accepted_message_ref: AcceptedMessageRef::new("message-runtime-auth-read-model")
Expand Down
8 changes: 7 additions & 1 deletion crates/ironclaw_reborn_openai_compat/src/chat_workflow.rs
Original file line number Diff line number Diff line change
Expand Up @@ -795,7 +795,13 @@ fn chat_user_message_and_attachments(
bytes: image.bytes,
})
.collect();
let payload = UserMessagePayload::new(text, vec![], ProductTriggerReason::DirectChat)?;
let payload = UserMessagePayload::new(text, vec![], ProductTriggerReason::DirectChat)?
.with_requested_model(crate::model_validation::requested_model_hint(
&request.model,
));
Comment thread
coderabbitai[bot] marked this conversation as resolved.
// The builder attaches the model hint after `new`'s validation, so bound the
// assembled payload before it is submitted.
payload.validate()?;
Ok((payload, attachments))
}

Expand Down
39 changes: 39 additions & 0 deletions crates/ironclaw_reborn_openai_compat/src/model_validation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,27 @@ use crate::OpenAiCompatHttpError;
/// Maximum accepted `model` string length, in bytes.
pub(crate) const MAX_MODEL_NAME_BYTES: usize = 256;

/// The OpenAI-compatible alias every client may send to mean "use the server's
/// active/default model" rather than naming a concrete one. The models listing
/// advertises it, so it is not a routable model id.
const DEFAULT_MODEL_ALIAS: &str = "default";

/// Map a validated client `model` string to an optional caller-requested model
/// *hint* for turn routing.
///
/// Returns `None` for the [`DEFAULT_MODEL_ALIAS`] sentinel (and defensively for
/// empty), so a client asking for the server default does not pin an advisory
/// route to the non-routable `"default"` id — which the model gateway rejects as
/// non-concrete and which would fail route resolution on routed hosts. A
/// concrete model name is forwarded as `Some`.
pub(crate) fn requested_model_hint(model: &str) -> Option<String> {
let trimmed = model.trim();
if trimmed.is_empty() || trimmed.eq_ignore_ascii_case(DEFAULT_MODEL_ALIAS) {
return None;
}
Some(trimmed.to_string())
}

/// Validate the client-supplied `model` string before it is carried as a
/// projection/policy hint.
///
Expand Down Expand Up @@ -82,4 +103,22 @@ mod tests {
let at_cap = "m".repeat(MAX_MODEL_NAME_BYTES);
assert!(validate_model_name(&at_cap).is_ok());
}

#[test]
fn requested_model_hint_drops_default_sentinel() {
assert_eq!(requested_model_hint("default"), None);
assert_eq!(requested_model_hint("DEFAULT"), None);
assert_eq!(requested_model_hint("Default"), None);
assert_eq!(requested_model_hint(""), None);
assert_eq!(requested_model_hint(" "), None);
}

#[test]
fn requested_model_hint_forwards_concrete_model() {
assert_eq!(requested_model_hint("gpt-4o"), Some("gpt-4o".to_string()));
assert_eq!(
requested_model_hint("anthropic/claude-opus-4"),
Some("anthropic/claude-opus-4".to_string())
);
}
}
Loading
Loading