(&self, query: String)
| 153 | T: ControlStateReadAccess + ControlStateWriteAccess + NodeDelegate + Authorization + Clone, |
| 154 | { |
| 155 | async fn exe_sql(&self, query: String) -> PgWireResult<Vec<Response>> { |
| 156 | let params = self.cached.lock().await.clone().unwrap(); |
| 157 | let db = SqlParams { |
| 158 | name_or_identity: database::NameOrIdentity::Name(DatabaseName(params.database.clone())), |
| 159 | }; |
| 160 | |
| 161 | let sql = match response( |
| 162 | database::sql_direct( |
| 163 | self.ctx.clone(), |
| 164 | db, |
| 165 | SqlQueryParams { confirmed: Some(true) }, |
| 166 | params.caller_identity, |
| 167 | params.caller_auth.clone(), |
| 168 | query.to_string(), |
| 169 | ) |
| 170 | .await, |
| 171 | ¶ms.database, |
| 172 | ) |
| 173 | .await |
| 174 | { |
| 175 | Ok(sql) => sql, |
| 176 | Err(PgError::Pg(PgWireError::UserError(err))) => { |
| 177 | return Ok(vec![Response::Error(err)]); |
| 178 | } |
| 179 | Err(err) => { |
| 180 | return Err(err.into()); |
| 181 | } |
| 182 | }; |
| 183 | |
| 184 | let mut result = Vec::with_capacity(sql.len()); |
| 185 | for sql_result in sql { |
| 186 | let header = row_desc(&sql_result.schema, &Format::UnifiedText); |
| 187 | if !sql_result.schema.is_empty() { |
| 188 | let rows = to_rows(sql_result, header.clone())?; |
| 189 | let q = QueryResponse::new(header, rows); |
| 190 | result.push(Response::Query(q)); |
| 191 | } else { |
| 192 | let tag = Tag::new(&stats(&sql_result)); |
| 193 | result.push(Response::Execution(tag)); |
| 194 | } |
| 195 | } |
| 196 | Ok(result) |
| 197 | } |
| 198 | } |
| 199 | |
| 200 | async fn close_client<C, E>(client: &mut C, err: E) -> PgWireResult<()> |
no test coverage detected