areev 0.2.0

Rust SDK for the Areev knowledge database — gRPC and HTTP transports
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
//! `Connectors` resource — Axtion connector catalog + OAuth + credential
//! storage.
//!
//! Mirrors the Python SDK's `client.connectors.*` surface. Backed by the
//! cell's `/api/axtion/*` proxy routes (which forward to Axtion with a
//! per-principal tenant header). The SDK never sees provider tokens — they
//! live in the Axtion credential vault; the SDK only sees catalog metadata,
//! action descriptors, and OAuth handshake state.
//!
//! ```no_run
//! # #[tokio::main]
//! # async fn main() -> areev::Result<()> {
//! let areev = areev::Areev::from_env();
//!
//! let catalog = areev.connectors().list().await?;
//! // → [{name: "gmail", display_name: "Google Gmail",
//! //     category: "Communication", version: "1.0.0"}, ...]
//!
//! let actions = areev.connectors().actions("gmail").await?;
//! // → [{name: "send-email", display_name: "Send Email",
//! //     param_schema: {...}, required: ["to", "subject", "body"]}, ...]
//!
//! // OAuth: mint an authorize URL, redirect the user, then poll.
//! let flow = areev
//!     .connectors()
//!     .authorize("gmail", "https://yourapp.example/oauth/cb")
//!     .await?;
//! // ... user consents in the browser ...
//! let _status = areev.connectors().poll_oauth("gmail", &flow.state).await?;
//! # Ok(())
//! # }
//! ```

use serde::{Deserialize, Serialize};
use serde_json::{json, Value};

use crate::error::Result;
use crate::http::HttpClient;

/// One entry in the connector catalog returned by [`Connectors::list`].
///
/// Parsed out of the cell's A2UI `data_table` (columns `Name`,
/// `Display Name`, `Category`, `Version`).
#[derive(Debug, Clone, PartialEq, Eq, Deserialize, Serialize)]
pub struct ConnectorSpec {
    /// Stable slug — e.g. `"gmail"`. Pass to [`Connectors::get`] /
    /// [`Connectors::actions`] / [`Connectors::authorize`].
    pub name: String,
    /// Human-readable label — e.g. `"Google Gmail"`.
    pub display_name: String,
    /// Catalog category — e.g. `"Communication"`.
    pub category: String,
    /// Connector version string — e.g. `"1.0.0"`.
    pub version: String,
}

/// One bindable action returned by [`Connectors::actions`].
///
/// `name` is the action *key* (the verb passed to
/// [`crate::resources::Tools::bind_axtion`]), **not** the human display
/// name.
#[derive(Debug, Clone, PartialEq, Eq, Deserialize, Serialize)]
pub struct ConnectorAction {
    /// The action key — pass to `bind_axtion`. E.g. `"send-email"`.
    pub name: String,
    /// Human-readable action label — e.g. `"Send Email"`.
    pub display_name: String,
    /// One-line description of what the action does.
    pub description: String,
    /// JSON Schema for the action's input arguments.
    pub param_schema: Value,
    /// The action's required input parameter names (from
    /// `param_schema.required`).
    pub required: Vec<String>,
}

/// Normalized OAuth authorize response from [`Connectors::authorize`].
///
/// The cell wraps the payload under `data` and names the URL `oauth_url`;
/// this surfaces a stable `authorize_url` + `state` while preserving every
/// other field under `extra`.
#[derive(Debug, Clone, Deserialize, Serialize)]
pub struct AuthorizeResponse {
    /// Provider authorize URL — redirect the user here.
    pub authorize_url: String,
    /// Opaque handshake state — pass to [`Connectors::poll_oauth`].
    pub state: String,
    /// Any additional fields the cell returned, preserved verbatim.
    #[serde(flatten)]
    pub extra: serde_json::Map<String, Value>,
}

/// Per-principal connector catalog + OAuth lifecycle.
///
/// Access via [`crate::Areev::connectors`]. OAuth flows and stored
/// credentials are scoped to the calling identity's derived Axtion tenant.
pub struct Connectors<'a> {
    http: &'a HttpClient,
}

impl<'a> Connectors<'a> {
    /// Internal constructor — use [`crate::Areev::connectors`].
    pub(crate) fn new(http: &'a HttpClient) -> Self {
        Self { http }
    }

    /// List the connector catalog.
    ///
    /// The cell returns an A2UI `data_table` (columns `Name`,
    /// `Display Name`, `Category`, `Version`); this parses it into a clean
    /// `Vec<ConnectorSpec>`. The column → field mapping is
    /// case-insensitive.
    ///
    /// # Errors
    ///
    /// Returns [`crate::AreevError::Other`] (`SDK-E020`) if the response is
    /// not a parseable A2UI table — the raw body is surfaced in the error
    /// message so the caller can inspect the unexpected shape.
    pub async fn list(&self) -> Result<Vec<ConnectorSpec>> {
        let body = self.http._get("/axtion/connectors", None).await?;
        parse_a2ui_table(&body).ok_or_else(|| crate::error::AreevError::Other {
            http_status: 0,
            code: Some("SDK-E020".into()),
            message: format!("connectors.list: response is not an A2UI data_table: {body}"),
            body: serde_json::Value::Null,
            request_id: None,
        })
    }

    /// Fetch one connector's raw detail, including its action descriptors
    /// under `data.actions`. Use [`Connectors::actions`] for a cleaned-up
    /// view of just the bindable actions.
    pub async fn get(&self, name: &str) -> Result<Value> {
        let path = format!("/axtion/connectors/{name}");
        self.http._get(&path, None).await
    }

    /// List a connector's bindable actions, cleanly parsed into
    /// [`ConnectorAction`]s.
    ///
    /// Reads the descriptors from `data.actions` (falling back to
    /// `data.tools`, then top-level `actions`/`tools`). Each descriptor's
    /// `key` becomes [`ConnectorAction::name`] (the bind verb) and `name`
    /// becomes the display label.
    pub async fn actions(&self, name: &str) -> Result<Vec<ConnectorAction>> {
        let detail = self.get(name).await?;
        Ok(parse_action_descriptors(&detail))
    }

    /// Mint an OAuth authorize URL for `name` bound to the calling
    /// principal.
    ///
    /// Returns an [`AuthorizeResponse`] — redirect the user to
    /// `authorize_url`; once they consent, poll [`Connectors::poll_oauth`]
    /// with the returned `state` until `status == "success"`.
    pub async fn authorize(&self, name: &str, redirect_uri: &str) -> Result<AuthorizeResponse> {
        let path = format!("/axtion/oauth/{name}/authorize");
        let query = json!({ "redirect_uri": redirect_uri });
        let body = self.http._get(&path, Some(&query)).await?;
        normalize_authorize(body)
    }

    /// Poll for OAuth completion.
    ///
    /// Returns the raw `{status: "pending"|"success"|"error"|"expired",
    /// ...}` body. On `success` the credential is stored in the Axtion
    /// vault and the connector is ready for
    /// [`crate::resources::Tools::bind_axtion`].
    pub async fn poll_oauth(&self, name: &str, state: &str) -> Result<Value> {
        let path = format!("/axtion/oauth/{name}/state/{state}");
        self.http._get(&path, None).await
    }

    /// Store API-key credentials for a non-OAuth connector.
    ///
    /// The `body` is proxied verbatim to Axtion's credential vault (scoped
    /// to the calling principal's tenant). The shape is connector-specific.
    pub async fn store_credentials(&self, body: Value) -> Result<Value> {
        self.http._post("/axtion/credentials", Some(&body)).await
    }
}

// ── Response parsers (defensive — Axtion/cell shapes evolve) ──────────────

/// Map a raw A2UI column header to a [`ConnectorSpec`] field key
/// (case-insensitive, whitespace-trimmed).
fn column_field(raw: &str) -> Option<&'static str> {
    match raw.trim().to_ascii_lowercase().as_str() {
        "name" => Some("name"),
        "display name" => Some("display_name"),
        "category" => Some("category"),
        "version" => Some("version"),
        _ => None,
    }
}

/// Parse the first A2UI `data_table` component into [`ConnectorSpec`]s.
///
/// Returns `None` if no `{columns, rows}` table component is found.
fn parse_a2ui_table(body: &Value) -> Option<Vec<ConnectorSpec>> {
    let components = body.get("a2ui")?.get("components")?.as_array()?;
    for comp in components {
        let (Some(columns), Some(rows)) = (
            comp.get("columns").and_then(Value::as_array),
            comp.get("rows").and_then(Value::as_array),
        ) else {
            continue;
        };
        let keys: Vec<Option<&'static str>> = columns
            .iter()
            .map(|c| c.as_str().and_then(column_field))
            .collect();
        let mut out = Vec::with_capacity(rows.len());
        for row in rows {
            let Some(cells) = row.as_array() else {
                continue;
            };
            let mut spec = ConnectorSpec {
                name: String::new(),
                display_name: String::new(),
                category: String::new(),
                version: String::new(),
            };
            for (i, key) in keys.iter().enumerate() {
                let Some(field) = key else { continue };
                let val = cells
                    .get(i)
                    .and_then(Value::as_str)
                    .unwrap_or("")
                    .to_string();
                match *field {
                    "name" => spec.name = val,
                    "display_name" => spec.display_name = val,
                    "category" => spec.category = val,
                    "version" => spec.version = val,
                    _ => {}
                }
            }
            out.push(spec);
        }
        return Some(out);
    }
    None
}

/// Extract + clean the action descriptors from a connector detail response.
///
/// The cell merges Axtion's tool descriptors under `data.actions`; some
/// shapes expose them under `data.tools`, or top-level `actions`/`tools`.
fn parse_action_descriptors(detail: &Value) -> Vec<ConnectorAction> {
    let raw = detail
        .get("data")
        .and_then(|d| d.get("actions").or_else(|| d.get("tools")))
        .or_else(|| detail.get("actions").or_else(|| detail.get("tools")))
        .and_then(Value::as_array);
    let Some(raw) = raw else {
        return Vec::new();
    };

    raw.iter()
        .filter_map(|d| {
            let obj = d.as_object()?;
            let schema = obj
                .get("input_schema")
                .or_else(|| obj.get("param_schema"))
                .or_else(|| obj.get("inputSchema"))
                .cloned()
                .unwrap_or_else(|| json!({}));
            let required = schema
                .get("required")
                .and_then(Value::as_array)
                .map(|a| {
                    a.iter()
                        .filter_map(|v| v.as_str().map(str::to_string))
                        .collect()
                })
                .unwrap_or_default();
            let key = obj
                .get("key")
                .or_else(|| obj.get("name"))
                .and_then(Value::as_str)
                .unwrap_or("")
                .to_string();
            Some(ConnectorAction {
                name: key,
                display_name: obj
                    .get("name")
                    .and_then(Value::as_str)
                    .unwrap_or("")
                    .to_string(),
                description: obj
                    .get("description")
                    .and_then(Value::as_str)
                    .unwrap_or("")
                    .to_string(),
                param_schema: schema,
                required,
            })
        })
        .collect()
}

/// Normalize the OAuth authorize response to `{authorize_url, state, …}`.
///
/// Axtion wraps the payload under `data` and names the URL `oauth_url`;
/// surface a stable `authorize_url` while preserving the original fields.
fn normalize_authorize(body: Value) -> Result<AuthorizeResponse> {
    // The cell nests the payload under `data`; the top-level may carry
    // `state` (and `data` the url) — merge both so neither is lost.
    let mut merged = serde_json::Map::new();
    if let Some(top) = body.as_object() {
        for (k, v) in top {
            if k == "data" {
                if let Some(inner) = v.as_object() {
                    for (ik, iv) in inner {
                        merged.insert(ik.clone(), iv.clone());
                    }
                }
            } else {
                merged.insert(k.clone(), v.clone());
            }
        }
    }

    let authorize_url = merged
        .get("authorize_url")
        .or_else(|| merged.get("oauth_url"))
        .or_else(|| merged.get("url"))
        .and_then(Value::as_str)
        .map(str::to_string)
        .ok_or_else(|| crate::error::AreevError::Other {
            http_status: 0,
            code: Some("SDK-E021".into()),
            message: format!(
                "connectors.authorize: response has no oauth_url/authorize_url: {body}"
            ),
            body: serde_json::Value::Null,
            request_id: None,
        })?;
    let state = merged
        .get("state")
        .and_then(Value::as_str)
        .unwrap_or("")
        .to_string();

    // Drop the normalized keys from `extra` to avoid duplicating them.
    merged.remove("authorize_url");
    merged.remove("oauth_url");
    merged.remove("url");
    merged.remove("state");

    Ok(AuthorizeResponse {
        authorize_url,
        state,
        extra: merged,
    })
}

#[cfg(test)]
mod tests {
    use super::*;

    fn gmail_detail_fixture() -> Value {
        // Mirror of the cell's `GET /axtion/connectors/gmail` shape: action
        // descriptors under `data.actions`, each with `key`/`name`/
        // `description`/`input_schema`.
        json!({
            "data": {
                "connector": "gmail",
                "actions": [
                    {
                        "connector": "gmail",
                        "key": "send-email",
                        "name": "Send Email",
                        "description": "Send an email on the user's behalf.",
                        "input_schema": {
                            "type": "object",
                            "properties": {
                                "to": {"type": "string"},
                                "subject": {"type": "string"},
                                "body": {"type": "string"}
                            },
                            "required": ["to", "subject", "body"]
                        },
                        "annotations": {}
                    }
                ]
            }
        })
    }

    #[test]
    fn parses_a2ui_catalog_table_case_insensitive() {
        let body = json!({
            "a2ui": {
                "components": [
                    {
                        "columns": ["Name", "Display Name", "Category", "Version"],
                        "rows": [
                            ["gmail", "Google Gmail", "Communication", "1.0.0"],
                            ["google-calendar", "Google Calendar", "Productivity", "2.1.0"]
                        ]
                    }
                ]
            }
        });
        let specs = parse_a2ui_table(&body).expect("parses table");
        assert_eq!(specs.len(), 2);
        assert_eq!(
            specs[0],
            ConnectorSpec {
                name: "gmail".into(),
                display_name: "Google Gmail".into(),
                category: "Communication".into(),
                version: "1.0.0".into(),
            }
        );
        assert_eq!(specs[1].name, "google-calendar");
        assert_eq!(specs[1].version, "2.1.0");
    }

    #[test]
    fn a2ui_parser_returns_none_for_non_table_body() {
        assert!(parse_a2ui_table(&json!({"unexpected": "shape"})).is_none());
        assert!(parse_a2ui_table(&json!({"a2ui": {"components": []}})).is_none());
    }

    #[test]
    fn actions_parses_gmail_send_email_with_required() {
        let actions = parse_action_descriptors(&gmail_detail_fixture());
        assert_eq!(actions.len(), 1);
        let send = &actions[0];
        assert_eq!(send.name, "send-email", "name is the action key");
        assert_eq!(send.display_name, "Send Email");
        assert_eq!(send.description, "Send an email on the user's behalf.");
        assert_eq!(
            send.required,
            vec!["to".to_string(), "subject".to_string(), "body".to_string()]
        );
        assert_eq!(send.param_schema["type"], "object");
        assert!(send.param_schema["properties"]["subject"].is_object());
    }

    #[test]
    fn actions_falls_back_to_data_tools() {
        let detail = json!({
            "data": {
                "tools": [
                    {"key": "lookup", "name": "Lookup", "input_schema": {"type": "object"}}
                ]
            }
        });
        let actions = parse_action_descriptors(&detail);
        assert_eq!(actions.len(), 1);
        assert_eq!(actions[0].name, "lookup");
        assert!(actions[0].required.is_empty());
    }

    #[test]
    fn actions_empty_for_unknown_shape() {
        assert!(parse_action_descriptors(&json!({"nope": true})).is_empty());
    }

    #[test]
    fn normalize_authorize_lifts_oauth_url_and_state() {
        let body = json!({
            "data": {"oauth_url": "https://accounts.google.com/o/oauth2/auth?x=1"},
            "state": "st_abc123"
        });
        let resp = normalize_authorize(body).expect("normalizes");
        assert_eq!(
            resp.authorize_url,
            "https://accounts.google.com/o/oauth2/auth?x=1"
        );
        assert_eq!(resp.state, "st_abc123");
        assert!(!resp.extra.contains_key("oauth_url"));
        assert!(!resp.extra.contains_key("state"));
    }

    #[test]
    fn normalize_authorize_errors_without_url() {
        let err = normalize_authorize(json!({"state": "x"})).unwrap_err();
        assert_eq!(err.code(), Some("SDK-E021"));
    }
}