Execute raw PostgreSQL frontend protocol bytes and parse backend protocol messages.
(
&mut self,
message: &[u8],
options: ExecProtocolOptions,
)
| 1053 | /// Execute raw PostgreSQL frontend protocol bytes and parse backend |
| 1054 | /// protocol messages. |
| 1055 | pub fn exec_protocol( |
| 1056 | &mut self, |
| 1057 | message: &[u8], |
| 1058 | options: ExecProtocolOptions, |
| 1059 | ) -> Result<ExecProtocolResult> { |
| 1060 | let ExecProtocolOptions { |
| 1061 | sync_to_fs, |
| 1062 | throw_on_error, |
| 1063 | on_notice, |
| 1064 | data_transfer_container, |
| 1065 | } = options; |
| 1066 | |
| 1067 | let data = { |
| 1068 | let _phase = timing::phase("client.protocol_roundtrip"); |
| 1069 | self.exec_protocol_raw_inner(message, sync_to_fs, data_transfer_container)? |
| 1070 | }; |
| 1071 | |
| 1072 | let mut messages = Vec::new(); |
| 1073 | let on_notice_cb = on_notice.clone(); |
| 1074 | let parse_result = { |
| 1075 | let _phase = timing::phase("client.protocol_parse"); |
| 1076 | self.parser.parse(&data, |msg| { |
| 1077 | if let BackendMessage::Error(db_err) = &msg |
| 1078 | && throw_on_error |
| 1079 | { |
| 1080 | return Err(anyhow!(db_err.clone())); |
| 1081 | } |
| 1082 | if let Some(callback) = on_notice_cb.as_ref() |
| 1083 | && let BackendMessage::Notice(notice) = &msg |
| 1084 | { |
| 1085 | callback(notice); |
| 1086 | } |
| 1087 | messages.push(msg); |
| 1088 | Ok(()) |
| 1089 | }) |
| 1090 | }; |
| 1091 | if let Err(err) = parse_result { |
| 1092 | match err.downcast::<DatabaseError>() { |
| 1093 | Ok(db_err) => { |
| 1094 | self.parser = ProtocolParser::new(); |
| 1095 | return Err(anyhow!(db_err)); |
| 1096 | } |
| 1097 | Err(err) => return Err(err), |
| 1098 | } |
| 1099 | } |
| 1100 | |
| 1101 | for message in &messages { |
| 1102 | if let BackendMessage::Notification(note) = message { |
| 1103 | if let Some(listeners) = self.notify_listeners.get(¬e.channel) { |
| 1104 | for listener in listeners { |
| 1105 | (listener.callback)(¬e.payload); |
| 1106 | } |
| 1107 | } |
| 1108 | for listener in &self.global_notify_listeners { |
| 1109 | (listener.callback)(¬e.channel, ¬e.payload); |
| 1110 | } |
| 1111 | } |
| 1112 | } |