Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -86,7 +86,7 @@ This crate is a thin async HTTP SDK around a DataHub-style REST API. Entry point
`update` echoes a typed `Node` rather than a flat `Resource`.
- `datasets` (`src/datasets/`)
- `files` (`src/files/`) — raw-`PUT` upload via `execute_file_upload_request` (content is the body, metadata rides in `X-Datahub-*` headers), plus directory listing, get/search, `FileUpdate` (rename/move/re-dataset), trash + restore, delete, and download (`download` in memory, `download_to_path` streamed)
- `subscriptions` (`src/subscriptions/`) — subscription CRUD, plus `listen.rs`: WebSocket listening against the api's subscription-listen endpoint (`tokio-tungstenite`). Reads follow the same split as every other collection: `list(limit)` over `GET /subscriptions?limit=` and `filter(form)` over `POST /subscriptions/filter`. Both are recent — `POST /subscriptions/list` was subscriptions-only (a `limit` defaulting to 100 where the api defaulted to 1000, an unvalidated sort property, and no cursor, so a tenant past one page could not reach the rest) and the api removed it rather than aliasing it, so a client that has not moved gets a 404. `SubscriptionFilterForm` is a strict subset of what `/filter` now accepts: it carries no `cursor`, and `SubscriptionFilter` has only `timeseries`, not the `id`, `externalId`, `name`, `createdTime` and `lastUpdatedTime` criteria the api's filter grew beside it.
- `subscriptions` (`src/subscriptions/`) — subscription CRUD, plus `listen.rs`: WebSocket listening against the api's subscription-listen endpoint (`tokio-tungstenite`). Reads follow the same split as every other collection: `list(limit)` over `GET /subscriptions?limit=` and `filter(form)` over `POST /subscriptions/filter`. Both are recent — `POST /subscriptions/list` was subscriptions-only (a `limit` defaulting to 100 where the api defaulted to 1000, an unvalidated sort property, and no cursor, so a tenant past one page could not reach the rest) and the api removed it rather than aliasing it, so a client that has not moved gets a 404. `SubscriptionFilterForm` has the family shape (`filter`, `limit: Option<u64>`, flattened `PageRequest`), and `SubscriptionFilter` carries everything the api's does: `id`, `externalId`, `name`, `timeseries`, `createdTime`, `lastUpdatedTime`. It is deliberately not a `NodeFilter` — a subscription has no `source`, `labels` or `metadata` — and note the criteria are `createdTime`/`lastUpdatedTime` while the entity spells its timestamps `dateCreated`/`lastUpdated`; that asymmetry is the api's. In Python, `filter()` takes the criteria as keywords or `filter=` and returns a `Page`, like `datasets.filter`.
- `functions` (`src/functions/`) — `create`, `list`, `get_by_id`, `update`, `delete`, plus a
client-side `by_ids`/`by_external_id`. The api **does** serve `/functions/byids`, `/filter` and
`/search` (platform #131) — **the SDK has not wired them yet**, so `by_ids` still lists and
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -100,8 +100,12 @@ class DataHubException(Exception):

All three are `None` when the API answered with something that is not a
problem document — an empty 401, a stack trace, or plain text.

A subscription listener raises it too, for a failure on its WebSocket. There
is no HTTP response behind that, so `status_code` is `None` along with the
problem attributes. Using a listener after `close()` raises `ValueError`.
"""
status_code: int
status_code: int | None
message: str
problem: dict | None
problem_type: str | None
Expand Down Expand Up @@ -1947,6 +1951,9 @@ class INode:


class FileUpload:
"""A local file to upload. Constructing one reads the file's metadata, so a
missing path raises `FileNotFoundError` and a directory `IsADirectoryError`."""

def __init__(
self,
path: str,
Expand Down Expand Up @@ -2075,36 +2082,29 @@ SubscriptionTimeseriesId = Union[TimeSeries, IdCollection, int, str]


class SubscriptionFilter:
def __init__(self, timeseries: list[SubscriptionTimeseriesId] | None = None) -> None: ...
@property
def timeseries(self) -> list[IdCollection]: ...

"""AND-combined criteria for ``subscriptions.filter``.

class DataSort:
``external_id`` and ``name`` are pattern lists — see ``PatternList``. ``timeseries`` matches
subscriptions bound to at least one of the given timeseries. A subscription is not a node, so
there is no ``source``, ``labels`` or ``metadata``.
"""
def __init__(
self,
property: list[str] | None = None,
order: str | None = None,
id: Sequence[int] | None = None,
external_id: PatternList | None = None,
name: PatternList | None = None,
timeseries: list[SubscriptionTimeseriesId] | None = None,
created_time: TimeFilter | None = None,
last_updated_time: TimeFilter | None = None,
) -> None: ...
@property
def property(self) -> list[str]: ...
def id(self) -> list[int] | None: ...
@property
def order(self) -> str | None: ...


class SubscriptionFilterForm:
def __init__(
self,
filter: SubscriptionFilter | None = None,
limit: int | None = None,
sort: DataSort | None = None,
) -> None: ...
@property
def filter(self) -> SubscriptionFilter: ...
def external_id(self) -> list[str] | None: ...
@property
def limit(self) -> int: ...
def name(self) -> list[str] | None: ...
@property
def sort(self) -> DataSort | None: ...
def timeseries(self) -> list[IdCollection]: ...


class EventAction:
Expand Down Expand Up @@ -2196,11 +2196,23 @@ class SubscriptionsServiceSync:
def list(self, limit: int | None = None) -> list[Subscription]: ...
def filter(
self,
form: SubscriptionFilterForm | None = None,
*,
filter: SubscriptionFilter | None = None,
id: Sequence[int] | None = None,
external_id: PatternList | None = None,
name: PatternList | None = None,
timeseries: list[SubscriptionTimeseriesId] | None = None,
created_time: TimeFilter | None = None,
last_updated_time: TimeFilter | None = None,
limit: int | None = None,
sort: DataSort | None = None,
) -> list[Subscription]: ...
sort_by: SortBy | None = None,
sort_order: str | None = None,
cursor: str | None = None,
) -> Page:
"""Pass either ``filter=`` or the individual criteria keywords; passing both is a
``TypeError``. Sortable by ``id``, ``externalId``, ``name``, ``createdTime`` and
``lastUpdatedTime``; the default is newest created first.
"""
def delete(self, input: list[SubscriptionIdentifiable]) -> None: ...
def listen(self, subscription_external_ids: list[str]) -> SubscriptionListener: ...

Expand All @@ -2210,11 +2222,23 @@ class SubscriptionsServiceAsync:
async def list(self, limit: int | None = None) -> list[Subscription]: ...
async def filter(
self,
form: SubscriptionFilterForm | None = None,
*,
filter: SubscriptionFilter | None = None,
id: Sequence[int] | None = None,
external_id: PatternList | None = None,
name: PatternList | None = None,
timeseries: list[SubscriptionTimeseriesId] | None = None,
created_time: TimeFilter | None = None,
last_updated_time: TimeFilter | None = None,
limit: int | None = None,
sort: DataSort | None = None,
) -> list[Subscription]: ...
sort_by: SortBy | None = None,
sort_order: str | None = None,
cursor: str | None = None,
) -> Page:
"""Pass either ``filter=`` or the individual criteria keywords; passing both is a
``TypeError``. Sortable by ``id``, ``externalId``, ``name``, ``createdTime`` and
``lastUpdatedTime``; the default is newest created first.
"""
async def delete(self, input: list[SubscriptionIdentifiable]) -> None: ...
async def listen(self, subscription_external_ids: list[str]) -> SubscriptionListenerAsync: ...

Expand Down
2 changes: 1 addition & 1 deletion datahub_python_bindings/src/datasets/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -484,7 +484,7 @@ impl PyDatasetFilter {
///
/// Shared by the sync and async services so the accepted keywords cannot drift apart between them.
///
/// `limit` defaults to the server's 100 and may not exceed 10000 — above that the request is
/// `limit` defaults to 1000 and may not exceed 10000 — above that the request is
/// rejected.
#[allow(clippy::too_many_arguments)]
pub fn dataset_filter_form(
Expand Down
4 changes: 3 additions & 1 deletion datahub_python_bindings/src/events/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -120,7 +120,9 @@ pub fn event_filter_form(
from_keywords,
any_keyword,
)?);
form.set_limit(limit.unwrap_or(100));
if let Some(limit) = limit {
form.set_limit(limit);
}
// A bare string is a one-element list here as it is on every other filter field; only one
// property is used either way.
if let Some(property) = sort_by {
Expand Down
6 changes: 3 additions & 3 deletions datahub_python_bindings/src/files/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -304,7 +304,7 @@ impl PyFileUpload {
data_set_id: Option<u64>,
related_resources: Option<Vec<u64>>,
) -> PyResult<Self> {
let mut file_upload = FileUpload::new(path);
let mut file_upload = FileUpload::new(path)?;
if let Some(external_id) = external_id {
file_upload.external_id = external_id.to_string();
}
Expand Down Expand Up @@ -334,7 +334,7 @@ impl PyFileUpload {
#[classmethod]
pub fn from_path(_py: Py<PyType>, path: &str) -> PyResult<Self> {
Ok(Self {
inner: FileUpload::new(path),
inner: FileUpload::new(path)?,
})
}
#[classmethod]
Expand All @@ -344,7 +344,7 @@ impl PyFileUpload {
destination_path: &str,
) -> PyResult<Self> {
Ok(Self {
inner: FileUpload::new_with_destination_path(path, destination_path),
inner: FileUpload::new_with_destination_path(path, destination_path)?,
})
}
#[getter]
Expand Down
16 changes: 16 additions & 0 deletions datahub_python_bindings/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,22 @@ create_exception!(
/// Convert an SDK `ResponseError` into a `DataHubException` that exposes the HTTP
/// `status_code` and `message` as attributes, so Python callers can branch on the
/// status code (e.g. `except DataHubException as e: if e.status_code == 409: ...`).
/// Convert a WebSocket listener failure into a `DataHubException`. There is no HTTP response behind
/// one, so `status_code` and the problem attributes are `None`, set anyway so an `except` block can
/// read them without `hasattr`.
pub(crate) fn listen_err(e: intellistream_datahub_sdk::ListenError) -> PyErr {
Python::attach(|py| {
let err = DataHubException::new_err(e.to_string());
let value = err.value(py);
let _ = value.setattr("status_code", py.None());
let _ = value.setattr("message", e.to_string());
let _ = value.setattr("problem", py.None());
let _ = value.setattr("problem_type", py.None());
let _ = value.setattr("problem_slug", py.None());
err
})
}

/// Map a "not found" into Python's `None`.
///
/// Every single-resource `GET` in the API answers an unknown id with **404** (batch `/byids`
Expand Down
45 changes: 27 additions & 18 deletions datahub_python_bindings/src/subscriptions/async_service.rs
Original file line number Diff line number Diff line change
@@ -1,13 +1,11 @@
use crate::subscriptions::listener::{PySubscriptionListenerAsync, shared_listener};
use crate::subscriptions::sync_service::build_filter_form;
use crate::subscriptions::{
PyDataSort, PySubscription, PySubscriptionFilterForm, SubscriptionIdentifyable,
SubscriptionTimeseriesId,
PySubscription, PySubscriptionFilter, SubscriptionIdentifyable, SubscriptionTimeseriesId,
subscription_filter_form,
};
use intellistream_datahub_sdk::ApiService;
use intellistream_datahub_sdk::generic::IdAndExtId;
use intellistream_datahub_sdk::subscriptions::Subscription;
use pyo3::exceptions::PyException;
use pyo3::prelude::*;
use pyo3_async_runtimes::tokio::future_into_py;
use std::sync::Arc;
Expand Down Expand Up @@ -62,30 +60,41 @@ impl PySubscriptionsServiceAsync {
})
}

/// Subscriptions matching every criterion on the filter.
#[pyo3(signature=(form=None, *, timeseries=None, limit=None, sort=None))]
/// Subscriptions matching every criterion on the filter, newest first.
#[pyo3(signature = (filter=None, id=None, external_id=None, name=None, timeseries=None,
created_time=None, last_updated_time=None, limit=None, sort_by=None,
sort_order=None, cursor=None))]
#[allow(clippy::too_many_arguments)]
fn filter<'py>(
&self,
py: Python<'py>,
form: Option<PySubscriptionFilterForm>,
filter: Option<PySubscriptionFilter>,
id: Option<Vec<u64>>,
external_id: Option<crate::StringOrList>,
name: Option<crate::StringOrList>,
timeseries: Option<Vec<SubscriptionTimeseriesId>>,
limit: Option<u32>,
sort: Option<PyDataSort>,
created_time: Option<crate::events::PyTimeFilter>,
last_updated_time: Option<crate::events::PyTimeFilter>,
limit: Option<u64>,
sort_by: Option<crate::StringOrList>,
sort_order: Option<String>,
cursor: Option<String>,
) -> PyResult<Bound<'py, PyAny>> {
let form = build_filter_form(form, timeseries, limit, sort)?;
let form = subscription_filter_form(
filter, id, external_id, name, timeseries, created_time, last_updated_time, limit,
sort_by, sort_order, cursor,
)?;
let service = self.api_service.clone();
future_into_py(py, async move {
let result = service
.subscriptions
.filter(&form)
.await
.map_err(|e| crate::datahub_err(e))?;
Ok(result
.get_items()
.iter()
.cloned()
.map(PySubscription::from)
.collect::<Vec<_>>())
.map_err(crate::datahub_err)?;
let next_cursor = result.next_cursor().map(str::to_string);
let items: Vec<PySubscription> =
result.get_items().iter().cloned().map(PySubscription::from).collect();
Python::attach(|py| crate::PyPage::new(py, items, next_cursor))
})
}

Expand Down Expand Up @@ -117,7 +126,7 @@ impl PySubscriptionsServiceAsync {
.subscriptions
.listen(&subscription_external_ids)
.await
.map_err(|e| PyException::new_err(e.to_string()))?;
.map_err(crate::listen_err)?;
Ok(PySubscriptionListenerAsync {
listener: shared_listener(listener),
})
Expand Down
Loading
Loading