(&self, event: EventData)
| 690 | } |
| 691 | |
| 692 | pub async fn post_events(&self, event: EventData) -> Result<EventsPostResponse, ErrorResponse> { |
| 693 | let event_data = match decode_multibase_data(&event.data) { |
| 694 | Ok(v) => v, |
| 695 | Err(e) => return Ok(EventsPostResponse::BadRequest(e)), |
| 696 | }; |
| 697 | |
| 698 | let event_id = |
| 699 | match event_id_from_car(self.network.clone(), event_data.as_slice(), &self.model).await |
| 700 | { |
| 701 | Ok(id) => id, |
| 702 | Err(err) => { |
| 703 | return Ok(EventsPostResponse::BadRequest(BadRequestResponse::new( |
| 704 | format!("Failed to parse EventID from event CAR file data: {err}"), |
| 705 | ))) |
| 706 | } |
| 707 | }; |
| 708 | |
| 709 | let (tx, rx) = tokio::sync::oneshot::channel(); |
| 710 | tokio::time::timeout( |
| 711 | INSERT_ENQUEUE_TIMEOUT, |
| 712 | self.insert_task.tx.send(EventInsert { |
| 713 | id: event_id, |
| 714 | data: event_data, |
| 715 | tx, |
| 716 | }), |
| 717 | ) |
| 718 | .map_err(|e| { |
| 719 | ErrorResponse::new(format!( |
| 720 | "Database service queue is too full to accept requests, error: {e}" |
| 721 | )) |
| 722 | }) |
| 723 | .await? |
| 724 | .map_err(|e| ErrorResponse::new(format!("Database service not available, error: {e}")))?; |
| 725 | |
| 726 | let new = tokio::time::timeout(INSERT_REQUEST_TIMEOUT, rx) |
| 727 | .await |
| 728 | .map_err(|e| { |
| 729 | ErrorResponse::new(format!( |
| 730 | "Timeout waiting for database service response, error: {e}" |
| 731 | )) |
| 732 | })? |
| 733 | .map_err(|e| { |
| 734 | ErrorResponse::new(format!( |
| 735 | "No response. Database service crashed with error: {e}" |
| 736 | )) |
| 737 | })? |
| 738 | .map_err(|e| ErrorResponse::new(format!("Failed to insert event: {e}")))?; |
| 739 | |
| 740 | match new { |
| 741 | EventInsertResult::Success(_) => Ok(EventsPostResponse::Success), |
| 742 | EventInsertResult::Failed(_, reason) => Ok(EventsPostResponse::BadRequest( |
| 743 | BadRequestResponse::new(reason), |
| 744 | )), |
| 745 | } |
| 746 | } |
| 747 | |
| 748 | pub async fn post_interests( |
| 749 | &self, |
no test coverage detected