From d5a7714696a878ad8ca9453df8e96085513fda8a Mon Sep 17 00:00:00 2001 From: Kacy Fortner Date: Tue, 10 Feb 2026 15:29:38 -0500 Subject: [PATCH] feat: PROTO.GETFIELD for field-level protobuf reads adds server-side field extraction from stored protobuf values, returning native RESP3 types. supports dot-separated nested paths (e.g., `address.city`). scalars map to native frames; complex types (message/list/map) return an error directing clients to use PROTO.GET. - FieldNotFound error variant in SchemaError - get_field(), resolve_field_path(), value_to_frame() in schema registry - ProtoGetField command variant + parser - handlers in both sharded and concurrent modes - unit tests for all field types and edge cases - integration tests for sharded + concurrent modes --- crates/ember-core/src/schema.rs | 351 +++++++++++++++++- crates/ember-protocol/src/command.rs | 52 +++ crates/ember-server/src/concurrent_handler.rs | 30 +- crates/ember-server/src/connection.rs | 31 +- tests/integration/src/proto.rs | 90 +++++ 5 files changed, 550 insertions(+), 4 deletions(-) diff --git a/crates/ember-core/src/schema.rs b/crates/ember-core/src/schema.rs index 271ecad4..f1891ef6 100644 --- a/crates/ember-core/src/schema.rs +++ b/crates/ember-core/src/schema.rs @@ -11,7 +11,10 @@ use std::collections::HashMap; use std::sync::{Arc, RwLock}; use bytes::Bytes; -use prost_reflect::{DescriptorPool, DynamicMessage, MessageDescriptor}; +use ember_protocol::Frame; +use prost_reflect::{ + DescriptorPool, DynamicMessage, FieldDescriptor, Kind, MessageDescriptor, ReflectMessage, +}; use thiserror::Error; /// Errors that can occur during schema operations. @@ -28,6 +31,9 @@ pub enum SchemaError { #[error("schema already registered: {0}")] AlreadyExists(String), + + #[error("field not found: {0}")] + FieldNotFound(String), } /// A registered schema: the raw descriptor bytes and the parsed pool. @@ -169,6 +175,25 @@ impl SchemaRegistry { ); } + /// Reads a single field from an encoded protobuf message. + /// + /// Decodes the message using the schema registry, walks the dot-separated + /// `field_path` to the target field, and converts the value to a RESP3 + /// frame. Returns an error for complex types (message, list, map) — those + /// require `PROTO.GET` for full deserialization. + pub fn get_field( + &self, + type_name: &str, + data: &[u8], + field_path: &str, + ) -> Result { + let descriptor = self.find_message(type_name)?; + let msg = DynamicMessage::decode(descriptor, data) + .map_err(|e| SchemaError::ValidationFailed(e.to_string()))?; + let (value, field_desc) = resolve_field_path(&msg, field_path)?; + value_to_frame(&value, &field_desc) + } + /// Looks up a message descriptor by full name across all schemas. fn find_message(&self, message_type: &str) -> Result { for schema in self.schemas.values() { @@ -180,6 +205,105 @@ impl SchemaRegistry { } } +/// Walks a dot-separated field path through a `DynamicMessage`, returning +/// the leaf value (owned) and its field descriptor. +/// +/// Intermediate path segments must be message-typed fields. The leaf +/// segment is the target field whose value is returned. +fn resolve_field_path( + msg: &DynamicMessage, + path: &str, +) -> Result<(prost_reflect::Value, FieldDescriptor), SchemaError> { + if path.is_empty() { + return Err(SchemaError::FieldNotFound("empty field path".into())); + } + + let segments: Vec<&str> = path.split('.').collect(); + for seg in &segments { + if seg.is_empty() { + return Err(SchemaError::FieldNotFound(format!( + "invalid field path '{path}': empty segment" + ))); + } + } + + let mut current_msg = msg.clone(); + + for (i, segment) in segments.iter().enumerate() { + let field_desc = current_msg + .descriptor() + .get_field_by_name(segment) + .ok_or_else(|| SchemaError::FieldNotFound(segment.to_string()))?; + + let value = current_msg.get_field(&field_desc).into_owned(); + + if i == segments.len() - 1 { + return Ok((value, field_desc)); + } + + // intermediate segment — must be a message type + match value { + prost_reflect::Value::Message(nested) => { + current_msg = nested; + } + _ => { + return Err(SchemaError::FieldNotFound(format!( + "'{segment}' is not a message field, cannot traverse further" + ))); + } + } + } + + unreachable!("loop always returns at the leaf segment") +} + +/// Converts a `prost_reflect::Value` + its field descriptor into a RESP3 frame. +/// +/// Scalar types are mapped to native RESP3 types. Complex types (message, +/// repeated, map) return an error directing clients to use `PROTO.GET`. +fn value_to_frame( + value: &prost_reflect::Value, + field_desc: &FieldDescriptor, +) -> Result { + // reject repeated and map fields up front + if field_desc.is_list() || field_desc.is_map() { + return Err(SchemaError::ValidationFailed( + "use PROTO.GET for repeated/map fields".into(), + )); + } + + match value { + prost_reflect::Value::String(s) => Ok(Frame::Bulk(Bytes::from(s.clone()))), + prost_reflect::Value::Bytes(b) => Ok(Frame::Bulk(b.clone())), + prost_reflect::Value::I32(n) => Ok(Frame::Integer(i64::from(*n))), + prost_reflect::Value::I64(n) => Ok(Frame::Integer(*n)), + prost_reflect::Value::U32(n) => Ok(Frame::Integer(i64::from(*n))), + prost_reflect::Value::U64(n) => Ok(Frame::Integer(*n as i64)), + prost_reflect::Value::F32(n) => Ok(Frame::Bulk(Bytes::from(format!("{n}")))), + prost_reflect::Value::F64(n) => Ok(Frame::Bulk(Bytes::from(format!("{n}")))), + prost_reflect::Value::Bool(b) => Ok(Frame::Integer(if *b { 1 } else { 0 })), + prost_reflect::Value::EnumNumber(n) => { + // look up the enum value name from the descriptor + if let Kind::Enum(enum_desc) = field_desc.kind() { + if let Some(val) = enum_desc.get_value(*n) { + return Ok(Frame::Bulk(Bytes::from(val.name().to_owned()))); + } + } + // fallback: return the numeric value + Ok(Frame::Integer(i64::from(*n))) + } + prost_reflect::Value::Message(_) => Err(SchemaError::ValidationFailed( + "use PROTO.GET for nested message fields".into(), + )), + prost_reflect::Value::List(_) => Err(SchemaError::ValidationFailed( + "use PROTO.GET for repeated fields".into(), + )), + prost_reflect::Value::Map(_) => Err(SchemaError::ValidationFailed( + "use PROTO.GET for map fields".into(), + )), + } +} + impl std::fmt::Debug for SchemaRegistry { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.debug_struct("SchemaRegistry") @@ -334,4 +458,229 @@ mod tests { pairs.sort(); assert_eq!(pairs, vec!["alpha", "beta"]); } + + // --- get_field tests --- + + /// Helper: encode a test.User message with the given name. + fn encode_user(registry: &SchemaRegistry, name: &str) -> Vec { + let pool = ®istry.schemas["users"].pool; + let msg_desc = pool.get_message_by_name("test.User").unwrap(); + let mut msg = DynamicMessage::new(msg_desc); + msg.set_field_by_name("name", prost_reflect::Value::String(name.into())); + let mut buf = Vec::new(); + use prost_reflect::prost::Message; + msg.encode(&mut buf).unwrap(); + buf + } + + #[test] + fn get_field_string() { + let mut registry = SchemaRegistry::new(); + let desc = make_descriptor("test", "User", "name"); + registry.register("users".into(), desc).unwrap(); + + let data = encode_user(®istry, "alice"); + let frame = registry.get_field("test.User", &data, "name").unwrap(); + assert_eq!(frame, Frame::Bulk(Bytes::from("alice"))); + } + + #[test] + fn get_field_default_value() { + let mut registry = SchemaRegistry::new(); + let desc = make_descriptor("test", "User", "name"); + registry.register("users".into(), desc).unwrap(); + + // encode an empty message (no fields set) + let pool = ®istry.schemas["users"].pool; + let msg_desc = pool.get_message_by_name("test.User").unwrap(); + let msg = DynamicMessage::new(msg_desc); + let mut buf = Vec::new(); + use prost_reflect::prost::Message; + msg.encode(&mut buf).unwrap(); + + // default string should be empty + let frame = registry.get_field("test.User", &buf, "name").unwrap(); + assert_eq!(frame, Frame::Bulk(Bytes::from(""))); + } + + #[test] + fn get_field_int() { + use prost_reflect::prost_types::{ + DescriptorProto, FieldDescriptorProto, FileDescriptorProto, FileDescriptorSet, + }; + + let fds = FileDescriptorSet { + file: vec![FileDescriptorProto { + name: Some("test.proto".into()), + package: Some("test".into()), + message_type: vec![DescriptorProto { + name: Some("Counter".into()), + field: vec![FieldDescriptorProto { + name: Some("count".into()), + number: Some(1), + r#type: Some(5), // TYPE_INT32 + label: Some(1), + ..Default::default() + }], + ..Default::default() + }], + ..Default::default() + }], + }; + let mut desc_buf = Vec::new(); + use prost_reflect::prost::Message; + fds.encode(&mut desc_buf).unwrap(); + let desc = Bytes::from(desc_buf); + + let mut registry = SchemaRegistry::new(); + registry.register("counters".into(), desc.clone()).unwrap(); + + let pool = ®istry.schemas["counters"].pool; + let msg_desc = pool.get_message_by_name("test.Counter").unwrap(); + let mut msg = DynamicMessage::new(msg_desc); + msg.set_field_by_name("count", prost_reflect::Value::I32(42)); + let mut buf = Vec::new(); + msg.encode(&mut buf).unwrap(); + + let frame = registry.get_field("test.Counter", &buf, "count").unwrap(); + assert_eq!(frame, Frame::Integer(42)); + } + + #[test] + fn get_field_bool() { + use prost_reflect::prost_types::{ + DescriptorProto, FieldDescriptorProto, FileDescriptorProto, FileDescriptorSet, + }; + + let fds = FileDescriptorSet { + file: vec![FileDescriptorProto { + name: Some("test.proto".into()), + package: Some("test".into()), + message_type: vec![DescriptorProto { + name: Some("Flag".into()), + field: vec![FieldDescriptorProto { + name: Some("active".into()), + number: Some(1), + r#type: Some(8), // TYPE_BOOL + label: Some(1), + ..Default::default() + }], + ..Default::default() + }], + ..Default::default() + }], + }; + let mut desc_buf = Vec::new(); + use prost_reflect::prost::Message; + fds.encode(&mut desc_buf).unwrap(); + let desc = Bytes::from(desc_buf); + + let mut registry = SchemaRegistry::new(); + registry.register("flags".into(), desc).unwrap(); + + let pool = ®istry.schemas["flags"].pool; + let msg_desc = pool.get_message_by_name("test.Flag").unwrap(); + let mut msg = DynamicMessage::new(msg_desc); + msg.set_field_by_name("active", prost_reflect::Value::Bool(true)); + let mut buf = Vec::new(); + msg.encode(&mut buf).unwrap(); + + let frame = registry.get_field("test.Flag", &buf, "active").unwrap(); + assert_eq!(frame, Frame::Integer(1)); + } + + /// Builds a descriptor with a nested message: Outer { Inner inner = 1; } + /// where Inner { string value = 1; } + fn make_nested_descriptor() -> Bytes { + use prost_reflect::prost_types::{ + DescriptorProto, FieldDescriptorProto, FileDescriptorProto, FileDescriptorSet, + }; + + let fds = FileDescriptorSet { + file: vec![FileDescriptorProto { + name: Some("test.proto".into()), + package: Some("test".into()), + message_type: vec![ + DescriptorProto { + name: Some("Inner".into()), + field: vec![FieldDescriptorProto { + name: Some("value".into()), + number: Some(1), + r#type: Some(9), // TYPE_STRING + label: Some(1), + ..Default::default() + }], + ..Default::default() + }, + DescriptorProto { + name: Some("Outer".into()), + field: vec![FieldDescriptorProto { + name: Some("inner".into()), + number: Some(1), + r#type: Some(11), // TYPE_MESSAGE + label: Some(1), + type_name: Some(".test.Inner".into()), + ..Default::default() + }], + ..Default::default() + }, + ], + ..Default::default() + }], + }; + let mut buf = Vec::new(); + use prost_reflect::prost::Message; + fds.encode(&mut buf).unwrap(); + Bytes::from(buf) + } + + #[test] + fn get_field_nested_path() { + let desc = make_nested_descriptor(); + let mut registry = SchemaRegistry::new(); + registry.register("nested".into(), desc).unwrap(); + + let pool = ®istry.schemas["nested"].pool; + let outer_desc = pool.get_message_by_name("test.Outer").unwrap(); + let inner_desc = pool.get_message_by_name("test.Inner").unwrap(); + + let mut inner = DynamicMessage::new(inner_desc); + inner.set_field_by_name("value", prost_reflect::Value::String("hello".into())); + + let mut outer = DynamicMessage::new(outer_desc); + outer.set_field_by_name("inner", prost_reflect::Value::Message(inner)); + + let mut buf = Vec::new(); + use prost_reflect::prost::Message; + outer.encode(&mut buf).unwrap(); + + let frame = registry + .get_field("test.Outer", &buf, "inner.value") + .unwrap(); + assert_eq!(frame, Frame::Bulk(Bytes::from("hello"))); + } + + #[test] + fn get_field_nonexistent() { + let mut registry = SchemaRegistry::new(); + let desc = make_descriptor("test", "User", "name"); + registry.register("users".into(), desc).unwrap(); + + let data = encode_user(®istry, "alice"); + let err = registry + .get_field("test.User", &data, "nonexistent") + .unwrap_err(); + assert!(matches!(err, SchemaError::FieldNotFound(_))); + } + + #[test] + fn get_field_empty_path() { + let mut registry = SchemaRegistry::new(); + let desc = make_descriptor("test", "User", "name"); + registry.register("users".into(), desc).unwrap(); + + let data = encode_user(®istry, "alice"); + let err = registry.get_field("test.User", &data, "").unwrap_err(); + assert!(matches!(err, SchemaError::FieldNotFound(_))); + } } diff --git a/crates/ember-protocol/src/command.rs b/crates/ember-protocol/src/command.rs index b6eaac8c..1b873842 100644 --- a/crates/ember-protocol/src/command.rs +++ b/crates/ember-protocol/src/command.rs @@ -351,6 +351,10 @@ pub enum Command { /// PROTO.DESCRIBE `name`. Lists message types in a registered schema. ProtoDescribe { name: String }, + /// PROTO.GETFIELD `key` `field_path`. Reads a single field from a + /// protobuf value, returning it as a native RESP3 type. + ProtoGetField { key: String, field_path: String }, + /// AUTH \[username\] password. Authenticate the connection. Auth { /// Username for ACL-style auth. None for legacy AUTH. @@ -483,6 +487,7 @@ impl Command { Command::ProtoType { .. } => "proto.type", Command::ProtoSchemas => "proto.schemas", Command::ProtoDescribe { .. } => "proto.describe", + Command::ProtoGetField { .. } => "proto.getfield", Command::Auth { .. } => "auth", Command::Quit => "quit", Command::Unknown(_) => "unknown", @@ -586,6 +591,7 @@ impl Command { "PROTO.TYPE" => parse_proto_type(&frames[1..]), "PROTO.SCHEMAS" => parse_proto_schemas(&frames[1..]), "PROTO.DESCRIBE" => parse_proto_describe(&frames[1..]), + "PROTO.GETFIELD" => parse_proto_getfield(&frames[1..]), "AUTH" => parse_auth(&frames[1..]), "QUIT" => parse_quit(&frames[1..]), _ => Ok(Command::Unknown(name)), @@ -1855,6 +1861,15 @@ fn parse_proto_describe(args: &[Frame]) -> Result { Ok(Command::ProtoDescribe { name }) } +fn parse_proto_getfield(args: &[Frame]) -> Result { + if args.len() != 2 { + return Err(ProtocolError::WrongArity("PROTO.GETFIELD".into())); + } + let key = extract_string(&args[0])?; + let field_path = extract_string(&args[1])?; + Ok(Command::ProtoGetField { key, field_path }) +} + fn parse_auth(args: &[Frame]) -> Result { match args.len() { 1 => { @@ -4238,4 +4253,41 @@ mod tests { let err = Command::from_frame(cmd(&["PROTO.DESCRIBE"])).unwrap_err(); assert!(matches!(err, ProtocolError::WrongArity(_))); } + + // --- proto.getfield --- + + #[test] + fn proto_getfield_basic() { + assert_eq!( + Command::from_frame(cmd(&["PROTO.GETFIELD", "user:1", "name"])).unwrap(), + Command::ProtoGetField { + key: "user:1".into(), + field_path: "name".into(), + }, + ); + } + + #[test] + fn proto_getfield_nested_path() { + assert_eq!( + Command::from_frame(cmd(&["PROTO.GETFIELD", "key", "address.city"])).unwrap(), + Command::ProtoGetField { + key: "key".into(), + field_path: "address.city".into(), + }, + ); + } + + #[test] + fn proto_getfield_wrong_arity() { + let err = Command::from_frame(cmd(&["PROTO.GETFIELD"])).unwrap_err(); + assert!(matches!(err, ProtocolError::WrongArity(_))); + + let err = Command::from_frame(cmd(&["PROTO.GETFIELD", "key"])).unwrap_err(); + assert!(matches!(err, ProtocolError::WrongArity(_))); + + let err = + Command::from_frame(cmd(&["PROTO.GETFIELD", "key", "field", "extra"])).unwrap_err(); + assert!(matches!(err, ProtocolError::WrongArity(_))); + } } diff --git a/crates/ember-server/src/concurrent_handler.rs b/crates/ember-server/src/concurrent_handler.rs index 060db6c6..84f67971 100644 --- a/crates/ember-server/src/concurrent_handler.rs +++ b/crates/ember-server/src/concurrent_handler.rs @@ -448,13 +448,41 @@ async fn execute_concurrent( } } + #[cfg(feature = "protobuf")] + Command::ProtoGetField { key, field_path } => { + let registry = match _engine.schema_registry() { + Some(r) => r, + None => return Frame::Error("ERR protobuf support is not enabled".into()), + }; + let req = ember_core::ShardRequest::ProtoGet { key: key.clone() }; + match _engine.route(&key, req).await { + Ok(ember_core::ShardResponse::ProtoValue(Some((type_name, data)))) => { + let reg = match registry.read() { + Ok(r) => r, + Err(_) => return Frame::Error("ERR schema registry lock poisoned".into()), + }; + match reg.get_field(&type_name, &data, &field_path) { + Ok(frame) => frame, + Err(e) => Frame::Error(format!("ERR {e}")), + } + } + Ok(ember_core::ShardResponse::ProtoValue(None)) => Frame::Null, + Ok(ember_core::ShardResponse::WrongType) => Frame::Error( + "WRONGTYPE Operation against a key holding the wrong kind of value".into(), + ), + Ok(other) => Frame::Error(format!("ERR unexpected shard response: {other:?}")), + Err(e) => Frame::Error(format!("ERR {e}")), + } + } + #[cfg(not(feature = "protobuf"))] Command::ProtoRegister { .. } | Command::ProtoSet { .. } | Command::ProtoGet { .. } | Command::ProtoType { .. } | Command::ProtoSchemas - | Command::ProtoDescribe { .. } => { + | Command::ProtoDescribe { .. } + | Command::ProtoGetField { .. } => { Frame::Error("ERR unknown command (protobuf support not compiled)".into()) } diff --git a/crates/ember-server/src/connection.rs b/crates/ember-server/src/connection.rs index c28664aa..95859822 100644 --- a/crates/ember-server/src/connection.rs +++ b/crates/ember-server/src/connection.rs @@ -557,7 +557,8 @@ async fn cluster_slot_check(ctx: &ServerContext, cmd: &Command) -> Option | Command::SCard { ref key } | Command::ProtoSet { ref key, .. } | Command::ProtoGet { ref key } - | Command::ProtoType { ref key } => cluster.check_slot(key.as_bytes()).await, + | Command::ProtoType { ref key } + | Command::ProtoGetField { ref key, .. } => cluster.check_slot(key.as_bytes()).await, // multi-key commands — crossslot validation + slot ownership Command::Del { ref keys } @@ -1774,6 +1775,31 @@ async fn execute( } } + #[cfg(feature = "protobuf")] + Command::ProtoGetField { key, field_path } => { + let registry = match engine.schema_registry() { + Some(r) => r, + None => return Frame::Error("ERR protobuf support is not enabled".into()), + }; + let req = ShardRequest::ProtoGet { key: key.clone() }; + match engine.route(&key, req).await { + Ok(ShardResponse::ProtoValue(Some((type_name, data)))) => { + let reg = match registry.read() { + Ok(r) => r, + Err(_) => return Frame::Error("ERR schema registry lock poisoned".into()), + }; + match reg.get_field(&type_name, &data, &field_path) { + Ok(frame) => frame, + Err(e) => Frame::Error(format!("ERR {e}")), + } + } + Ok(ShardResponse::ProtoValue(None)) => Frame::Null, + Ok(ShardResponse::WrongType) => wrongtype_error(), + Ok(other) => Frame::Error(format!("ERR unexpected shard response: {other:?}")), + Err(e) => Frame::Error(format!("ERR {e}")), + } + } + // when protobuf feature is disabled, proto commands are unknown #[cfg(not(feature = "protobuf"))] Command::ProtoRegister { .. } @@ -1781,7 +1807,8 @@ async fn execute( | Command::ProtoGet { .. } | Command::ProtoType { .. } | Command::ProtoSchemas - | Command::ProtoDescribe { .. } => { + | Command::ProtoDescribe { .. } + | Command::ProtoGetField { .. } => { Frame::Error("ERR unknown command (protobuf support not compiled)".into()) } diff --git a/tests/integration/src/proto.rs b/tests/integration/src/proto.rs index 60b02bbd..9fb2b274 100644 --- a/tests/integration/src/proto.rs +++ b/tests/integration/src/proto.rs @@ -305,6 +305,80 @@ async fn persistence_recovery() { drop(data_dir); } +// ---- PROTO.GETFIELD sharded tests ---- + +#[tokio::test] +async fn getfield_string() { + let server = start_proto_server(false); + let mut c = server.connect().await; + + let desc = make_descriptor("test", "User", "name"); + c.cmd_raw(&[b"PROTO.REGISTER", b"users", &desc]).await; + + let data = encode_message(&desc, "test.User", "name", "alice"); + c.cmd_raw(&[b"PROTO.SET", b"user:1", b"test.User", &data]) + .await; + + let resp = c.cmd(&["PROTO.GETFIELD", "user:1", "name"]).await; + assert_eq!(resp, Frame::Bulk(Bytes::from("alice"))); +} + +#[tokio::test] +async fn getfield_missing_key() { + let server = start_proto_server(false); + let mut c = server.connect().await; + + let desc = make_descriptor("test", "User", "name"); + c.cmd_raw(&[b"PROTO.REGISTER", b"users", &desc]).await; + + let resp = c.cmd(&["PROTO.GETFIELD", "nonexistent", "name"]).await; + assert!(matches!(resp, Frame::Null)); +} + +#[tokio::test] +async fn getfield_wrong_type() { + let server = start_proto_server(false); + let mut c = server.connect().await; + + c.ok(&["SET", "str:key", "hello"]).await; + + let resp = c.cmd(&["PROTO.GETFIELD", "str:key", "name"]).await; + assert!(matches!(resp, Frame::Error(ref s) if s.starts_with("WRONGTYPE"))); +} + +#[tokio::test] +async fn getfield_nonexistent_field() { + let server = start_proto_server(false); + let mut c = server.connect().await; + + let desc = make_descriptor("test", "User", "name"); + c.cmd_raw(&[b"PROTO.REGISTER", b"users", &desc]).await; + + let data = encode_message(&desc, "test.User", "name", "alice"); + c.cmd_raw(&[b"PROTO.SET", b"user:1", b"test.User", &data]) + .await; + + let resp = c.cmd(&["PROTO.GETFIELD", "user:1", "nonexistent"]).await; + assert!(matches!(resp, Frame::Error(_))); +} + +#[tokio::test] +async fn getfield_default_value() { + let server = start_proto_server(false); + let mut c = server.connect().await; + + let desc = make_descriptor("test", "User", "name"); + c.cmd_raw(&[b"PROTO.REGISTER", b"users", &desc]).await; + + // encode a message with the field at its default (empty string) + let data = encode_message(&desc, "test.User", "name", ""); + c.cmd_raw(&[b"PROTO.SET", b"user:1", b"test.User", &data]) + .await; + + let resp = c.cmd(&["PROTO.GETFIELD", "user:1", "name"]).await; + assert_eq!(resp, Frame::Bulk(Bytes::from(""))); +} + // ---- concurrent mode tests ---- // These mirror the core sharded tests to verify the concurrent handler's // proto command routing through the engine fallback path. @@ -437,3 +511,19 @@ async fn concurrent_get_missing_key_returns_null() { let resp = c.cmd(&["PROTO.TYPE", "nonexistent"]).await; assert!(matches!(resp, Frame::Null)); } + +#[tokio::test] +async fn concurrent_getfield_string() { + let server = start_proto_server(true); + let mut c = server.connect().await; + + let desc = make_descriptor("test", "User", "name"); + c.cmd_raw(&[b"PROTO.REGISTER", b"users", &desc]).await; + + let data = encode_message(&desc, "test.User", "name", "alice"); + c.cmd_raw(&[b"PROTO.SET", b"user:1", b"test.User", &data]) + .await; + + let resp = c.cmd(&["PROTO.GETFIELD", "user:1", "name"]).await; + assert_eq!(resp, Frame::Bulk(Bytes::from("alice"))); +}