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")));
+}