uarp-sdk 0.5.7

Async Rust client for the UARP (Snaga) Universal Agent Runtime Platform API
Documentation
// Code generated by @uarp/codegen from spec/openapi.json. DO NOT EDIT.
//!
//! Activity feed and real-time event streaming

#![allow(unused_imports, clippy::too_many_arguments)]

use reqwest::Method;
use serde::{Deserialize, Serialize};
use futures_core::Stream;

use crate::client::{Client, Request, NO_BODY, NO_QUERY};
use crate::error::Result;
use crate::generated::models;
use crate::multipart::{field_text, FilePart};
use crate::pagination::CursorGuard;
use crate::sse::EventStream;
use crate::util::encode_path;

/// Query and header parameters for `getActivityFeed`.
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
pub struct GetActivityFeedParams {
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub limit: Option<i64>,
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub cursor: Option<String>,
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub agent_id: Option<String>,
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub company_id: Option<String>,
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub team_id: Option<String>,
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub types: Option<String>,
}

/// Query and header parameters for `streamActivityFeed`.
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
pub struct StreamActivityFeedParams {
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub limit: Option<i64>,
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub agent_id: Option<String>,
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub company_id: Option<String>,
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub team_id: Option<String>,
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub types: Option<String>,
}

/// Activity feed and real-time event streaming
#[derive(Debug, Clone)]
pub struct FeedApi {
    pub(crate) client: Client,
}

impl Client {
    /// Activity feed and real-time event streaming
    pub fn feed(&self) -> FeedApi {
        FeedApi { client: self.clone() }
    }
}

impl FeedApi {
    /// Activity feed
    ///
    /// `GET /api/v1/feed`
    ///
    /// Required scopes: `runs:read`.
    pub async fn get_activity_feed(&self, params: &GetActivityFeedParams) -> Result<models::GetActivityFeedResponse> {
        self.client
            .request_json(Request {
                method: Method::GET,
                path: "/api/v1/feed".to_string(),
                query: Some(params),
                body: NO_BODY,
                headers: Vec::new(),
                idempotent: false,
            })
            .await
    }

    /// Stream every item returned by `getActivityFeed`, following the `cursor` cursor until the
    /// server reports no further pages.
    pub fn get_activity_feed_all<'a>(&'a self, params: &'a GetActivityFeedParams) -> impl Stream<Item = Result<serde_json::Map<String, serde_json::Value>>> + 'a {
        async_stream::try_stream! {
            let mut guard = CursorGuard::new();
            let mut cursor = params.cursor.clone();
            loop {
                let mut page_params = params.clone();
                page_params.cursor = cursor.clone();
                let page = self.get_activity_feed(&page_params).await?;
                let items = page.entries.unwrap_or_default();
                let was_empty = items.is_empty();
                for item in items {
                    yield item;
                }
                match guard.advance(page.cursor, None, was_empty) {
                    Some(next) => cursor = Some(next),
                    None => break,
                }
            }
        }
    }

    /// SSE activity feed stream
    ///
    /// `GET /api/v1/feed/stream`
    ///
    /// Required scopes: `runs:read`.
    ///
    /// Returns a server-sent event stream.
    pub fn stream_activity_feed(&self, params: &StreamActivityFeedParams) -> EventStream {
        self.client.request_stream(
            "/api/v1/feed/stream",
            Some(params),
            Vec::new(),
        )
    }
}