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
Original file line number Diff line number Diff line change
Expand Up @@ -13,8 +13,8 @@
import ai.intellistream.datahub.models.analysis.AnalysisResult;
import ai.intellistream.datahub.models.forms.AnalysisForm;
import ai.intellistream.datahub.sdk.client.DatahubClient;
import ai.intellistream.datahub.sdk.services.ResourceService;
import ai.intellistream.datahub.sdk.services.TimeseriesService;
import ai.intellistream.datahub.sdk.client.ResourceService;
import ai.intellistream.datahub.sdk.client.TimeseriesService;
import ai.intellistream.datahub.timeseries.Timeseries;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
Expand Down
57 changes: 33 additions & 24 deletions datahub-java-sdk/AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ Thin, synchronous Java client for the DataHub Platform REST API, published as
- **Binary ingest is its own method.** `ingestBinary(...)` and `binaryBuffer()` on
`TimeseriesService` go to `POST /timeseries/data/binary`; the JSON `ingest(...)` is untouched
and the durable spool applies to it only. The binary path resolves series through
`/timeseries/byids` (`ingest/SeriesResolver`), so it needs read access to the dataset too.
`/timeseries/byids` (`client/SeriesResolver`), so it needs read access to the dataset too.
- **Branch on the problem `type`, never on a substring of the body.** The api answers every
failure with one RFC 9457 shape whose `type` URI is the contract; `detail` and `title` are prose
for a human and may be reworded. `Problem.of(status, body)` never throws and never returns null,
Expand All @@ -36,7 +36,7 @@ Thin, synchronous Java client for the DataHub Platform REST API, published as
first send, so a retry carries the same id and collapses in ClickHouse
(`ReplacingMergeTree ORDER BY id`). Never switch to random v4 ids for events — they scatter
the sort key and degrade insert/merge/query performance.
- **The durable spool must stay memory-safe** (`ingest/DurableSpool`): append to a plain NDJSON
- **The durable spool must stay memory-safe** (`client/DurableSpool`): append to a plain NDJSON
active segment, gzip-seal at ~50 MiB rollover, stream sealed segments in fixed-size chunks on
flush — a multi-gigabyte spool never loads into memory. Buffer only retryable failures:
unreachable (network error, 429, 5xx) and auth (401/403). Terminal errors such as 400 are
Expand All @@ -47,28 +47,37 @@ Thin, synchronous Java client for the DataHub Platform REST API, published as

## Layout (`ai.intellistream.datahub.sdk`)

- `client/` — `DatahubClient` (entry point, one accessor per service), `DatahubConfig`
(builder; `fromEnv()` on `BASE_URL` + `TOKEN` or `CLIENT_ID`/`CLIENT_SECRET`/`TOKEN_URI`,
optionally `SCOPE`/`AUDIENCE` and the `ASSERTION*` keys that select the `jwt-bearer` grant;
Vault variants via `VaultSecretLoader`, a JDK-HttpClient KV v2 read supporting token and
AppRole auth).
- `auth/` — `TokenProvider`: static token pass-through, or a cached single-flight exchange
refreshed ~30 s before expiry — client-credentials, or the RFC 7523 `jwt-bearer` grant when an
assertion source is configured. The assertion is re-requested per exchange, never cached,
because providers commonly reject a replayed one.
- `http/` — shared plumbing: `ApiHttp` request helpers, `DatahubApiException` error mapping.
Every non-2xx is read as the api's RFC 9457 problem document through
`DatahubApiException.problem()` (`ai.intellistream.datahub.api.errors.Problem`, in api-model).
- `services/` — one class per API area: resources, assets, functions, timeseries, datasets,
events, labels, policies, governance, tenant, units, files, subscriptions. `assets` and
`functions` are the typed views of the `ASSET`/`FUNCTION` corners of the same graph `resources`
serves polymorphically; `labels` reads and writes through `LabelForm`, which is
`@Schema(name = "Label")` and is the label wire shape on both sides.
- `ingest/` — batched ingestion plus the durable disk spool (`DatapointIngestor`,
`EventIngestor`, `DurableSpool`, `BatchExecutor`), and the binary path
(`BinaryDatapointIngestor`, `BinaryIngestOptions`, `BinaryIngestBuffer`, `SeriesResolver`).
- `subscriptions/` — `SubscriptionListener`: durable subscription listening over the api's
WebSocket endpoint with per-subscription ack/nack.
**`DatahubClient` is the only way in.** Everything a caller does is reached through it. The
plumbing is package-private, and so are the service constructors, which is why the services and
the plumbing share one package: Java can only hide a constructor from other packages. Keep new
plumbing package-private in `client/`, and keep new service constructors package-private.
Narrowing a published public type later is a breaking change. The other packages hold only
value types a caller names.

- `client/` — the entry point and everything behind it:
- public: `DatahubClient` (one accessor per service); `DatahubConfig` (builder; `fromEnv()`
on `BASE_URL` + `TOKEN` or `CLIENT_ID`/`CLIENT_SECRET`/`TOKEN_URI`, optionally
`SCOPE`/`AUDIENCE` and the `ASSERTION*` keys that select the `jwt-bearer` grant; Vault
variants via `VaultSecretLoader`, a JDK-HttpClient KV v2 read supporting token and AppRole
auth); one `*Service` per API area (resources, assets, functions, timeseries, datasets,
events, labels, policies, governance, tenant, units, files, subscriptions); and the handles
a service returns, `BinaryIngestBuffer` and `SubscriptionListener`. `assets` and `functions`
are the typed views of the `ASSET`/`FUNCTION` corners of the same graph `resources` serves
polymorphically; `labels` reads and writes through `LabelForm`, which is
`@Schema(name = "Label")` and is the label wire shape on both sides.
- package-private: `ApiHttp` (request helpers; every non-2xx becomes a
`DatahubApiException`); `TokenProvider` (static token pass-through, or a cached single-flight
exchange refreshed ~30 s before expiry, either client-credentials or the RFC 7523
`jwt-bearer` grant when an assertion source is configured; the assertion is re-requested per
exchange, never cached, because providers commonly reject a replayed one); the ingest
machinery (`DatapointIngestor`, `EventIngestor`, `BatchExecutor`, `DurableSpool`,
`DatapointSpool`, and for the binary path `BinaryDatapointIngestor` and `SeriesResolver`).
- `http/` — `DatahubApiException`. Every refusal is read as the api's RFC 9457 problem
document through `problem()` (`ai.intellistream.datahub.api.errors.Problem`, in api-model).
- `ingest/` — `IngestOptions`, `BinaryIngestOptions`, `IngestResult`.
- `subscriptions/` — `SubscriptionMessage`, `SubscriptionError`, delivered by
`client/SubscriptionListener` (durable subscription listening over the api's WebSocket endpoint
with per-subscription ack/nack).
- `timeseries/`, `util/` — `Datapoint` model, UUID v7 generator.

## Tests
Expand Down
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
// SPDX-License-Identifier: Apache-2.0
package ai.intellistream.datahub.sdk.http;
package ai.intellistream.datahub.sdk.client;

import ai.intellistream.datahub.sdk.auth.TokenProvider;
import ai.intellistream.datahub.sdk.http.DatahubApiException;
import tools.jackson.databind.JavaType;
import tools.jackson.databind.json.JsonMapper;
import tools.jackson.databind.type.TypeFactory;
Expand All @@ -19,7 +19,7 @@
* (de)serializes JSON via Jackson, and maps non-2xx responses to {@link DatahubApiException}.
* Thread-safe and meant to be shared.
*/
public final class ApiHttp {
final class ApiHttp {

/**
* The API answers a failure with {@code application/problem+json} (RFC 9457), a different media
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
// SPDX-License-Identifier: Apache-2.0
package ai.intellistream.datahub.sdk.services;
package ai.intellistream.datahub.sdk.client;

import ai.intellistream.datahub.api.responses.DataWrapper;
import ai.intellistream.datahub.api.responses.GraphDataWrapper;
Expand All @@ -12,7 +12,6 @@
import ai.intellistream.datahub.models.UpdateRelForm;
import ai.intellistream.datahub.models.UpdateAssetForm;
import ai.intellistream.datahub.models.datafilters.ResourceFilter;
import ai.intellistream.datahub.sdk.http.ApiHttp;
import tools.jackson.databind.JavaType;
import tools.jackson.databind.type.TypeFactory;

Expand All @@ -32,7 +31,7 @@ public final class AssetService {
private final JavaType assets; // DataWrapper<Asset>
private final JavaType nodeGraph; // GraphDataWrapper<NodeModel, EdgeProxy> — the update echo

public AssetService(ApiHttp http) {
AssetService(ApiHttp http) {
this.http = http;
TypeFactory tf = http.typeFactory();
this.assets = tf.constructParametricType(DataWrapper.class, Asset.class);
Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
// SPDX-License-Identifier: Apache-2.0
package ai.intellistream.datahub.sdk.ingest;
package ai.intellistream.datahub.sdk.client;

import ai.intellistream.datahub.sdk.ingest.IngestOptions;
import ai.intellistream.datahub.sdk.ingest.IngestResult;
import ai.intellistream.datahub.api.errors.Problem;
import ai.intellistream.datahub.sdk.http.DatahubApiException;

Expand Down
Original file line number Diff line number Diff line change
@@ -1,17 +1,19 @@
// SPDX-License-Identifier: Apache-2.0
package ai.intellistream.datahub.sdk.ingest;
package ai.intellistream.datahub.sdk.client;

import ai.intellistream.datahub.sdk.ingest.IngestOptions;
import ai.intellistream.datahub.sdk.ingest.IngestResult;
import ai.intellistream.datahub.sdk.ingest.BinaryIngestOptions;
import ai.intellistream.datahub.api.binary.DatapointFrameWriter;
import ai.intellistream.datahub.api.binary.DatapointValueType;
import ai.intellistream.datahub.api.binary.FrameLimits;
import ai.intellistream.datahub.api.binary.ZstdPayloadCodec;
import ai.intellistream.datahub.api.responses.DatapointString;
import ai.intellistream.datahub.api.responses.DatapointsCollection;
import ai.intellistream.datahub.helpers.datetime.DateTimeHandler;
import ai.intellistream.datahub.sdk.http.ApiHttp;
import ai.intellistream.datahub.api.errors.Problem;
import ai.intellistream.datahub.sdk.http.DatahubApiException;
import ai.intellistream.datahub.sdk.ingest.SeriesResolver.Resolved;
import ai.intellistream.datahub.sdk.client.SeriesResolver.Resolved;

import java.io.ByteArrayOutputStream;
import java.nio.ByteBuffer;
Expand All @@ -36,7 +38,7 @@
* {@link BatchExecutor}. A request the server refuses because a series is unknown or renamed
* evicts those series from the resolver, and the points it carried are rebuilt and sent once more.
*/
public final class BinaryDatapointIngestor {
final class BinaryDatapointIngestor {

static final String PATH = "/timeseries/data/binary";
/** Leave headroom under the 4 MiB raw cap so the estimate never lands on the wrong side of it. */
Expand Down Expand Up @@ -166,7 +168,13 @@ private IngestResult send(List<Request> requests, BinaryIngestOptions options,
}
}));
}
return BatchExecutor.execute(tasks, options.executorOptions());
// Same retry, parallelism and fail-fast rules as the JSON path's executor.
IngestOptions executorOptions = IngestOptions.builder()
.parallelism(options.parallelism())
.maxRetries(options.maxRetries())
.failFast(options.failFast())
.build();
return BatchExecutor.execute(tasks, executorOptions);
}

/**
Expand Down
Original file line number Diff line number Diff line change
@@ -1,7 +1,8 @@
// SPDX-License-Identifier: Apache-2.0
package ai.intellistream.datahub.sdk.ingest;
package ai.intellistream.datahub.sdk.client;

import ai.intellistream.datahub.sdk.services.TimeseriesService;
import ai.intellistream.datahub.sdk.ingest.IngestResult;
import ai.intellistream.datahub.sdk.ingest.BinaryIngestOptions;
import ai.intellistream.datahub.sdk.timeseries.Datapoint;

import java.time.Duration;
Expand Down Expand Up @@ -51,7 +52,7 @@ public final class BinaryIngestBuffer implements AutoCloseable {
private volatile IngestResult lastResult;
private volatile boolean closed;

public BinaryIngestBuffer(TimeseriesService service, BinaryIngestOptions options,
BinaryIngestBuffer(TimeseriesService service, BinaryIngestOptions options,
int maxPoints, Duration maxAge, Consumer<IngestResult> onFlush) {
if (maxPoints <= 0) {
throw new IllegalArgumentException("maxPoints must be > 0");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,24 +3,6 @@

import ai.intellistream.datahub.models.NodeModel;
import ai.intellistream.datahub.models.EventModel;
import ai.intellistream.datahub.sdk.auth.TokenProvider;
import ai.intellistream.datahub.sdk.http.ApiHttp;
import ai.intellistream.datahub.sdk.ingest.DatapointSpool;
import ai.intellistream.datahub.sdk.ingest.DurableSpool;
import ai.intellistream.datahub.sdk.services.DatasetService;
import ai.intellistream.datahub.sdk.services.AssetService;
import ai.intellistream.datahub.sdk.services.EdgeService;
import ai.intellistream.datahub.sdk.services.FunctionService;
import ai.intellistream.datahub.sdk.services.GovernanceService;
import ai.intellistream.datahub.sdk.services.LabelService;
import ai.intellistream.datahub.sdk.services.PolicyService;
import ai.intellistream.datahub.sdk.services.TenantService;
import ai.intellistream.datahub.sdk.services.EventService;
import ai.intellistream.datahub.sdk.services.FileService;
import ai.intellistream.datahub.sdk.services.ResourceService;
import ai.intellistream.datahub.sdk.services.SubscriptionService;
import ai.intellistream.datahub.sdk.services.TimeseriesService;
import ai.intellistream.datahub.sdk.services.UnitService;
import ai.intellistream.datahub.models.NodeModelSubtypes;
import tools.jackson.databind.json.JsonMapper;

Expand Down
Original file line number Diff line number Diff line change
@@ -1,10 +1,11 @@
// SPDX-License-Identifier: Apache-2.0
package ai.intellistream.datahub.sdk.ingest;
package ai.intellistream.datahub.sdk.client;

import ai.intellistream.datahub.sdk.ingest.IngestOptions;
import ai.intellistream.datahub.sdk.ingest.IngestResult;
import ai.intellistream.datahub.api.responses.DataWrapper;
import ai.intellistream.datahub.api.responses.DatapointString;
import ai.intellistream.datahub.api.responses.DatapointsCollection;
import ai.intellistream.datahub.sdk.http.ApiHttp;
import tools.jackson.databind.JavaType;

import java.util.ArrayList;
Expand All @@ -15,7 +16,7 @@
* large collections) and runs the batches through {@link BatchExecutor}. Batches are independent,
* so there is no cross-batch ordering guarantee — which is fine for timestamped data.
*/
public final class DatapointIngestor {
final class DatapointIngestor {

private final ApiHttp http;
private final String path;
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
// SPDX-License-Identifier: Apache-2.0
package ai.intellistream.datahub.sdk.ingest;
package ai.intellistream.datahub.sdk.client;

import ai.intellistream.datahub.api.responses.DatapointString;
import ai.intellistream.datahub.api.responses.DatapointsCollection;
Expand All @@ -17,7 +17,7 @@
* time/size retention bounds apply per datapoint), with the flattening to and grouping from
* {@link DatapointsCollection} that the timeseries ingest path needs.
*/
public final class DatapointSpool {
final class DatapointSpool {

/** One spooled datapoint: its series external id and the wire timestamp/value (both strings). */
public record SpoolLine(String externalId, String timestamp, String value) {
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
// SPDX-License-Identifier: Apache-2.0
package ai.intellistream.datahub.sdk.services;
package ai.intellistream.datahub.sdk.client;

import ai.intellistream.datahub.api.responses.DataWrapper;
import ai.intellistream.datahub.models.DataSetModel;
Expand All @@ -8,7 +8,6 @@
import ai.intellistream.datahub.models.Resource;
import ai.intellistream.datahub.models.datafilters.DataSetFilter;
import ai.intellistream.datahub.models.forms.DataSetForm;
import ai.intellistream.datahub.sdk.http.ApiHttp;
import ai.intellistream.datahub.models.SearchBody;
import tools.jackson.databind.JavaType;

Expand All @@ -21,7 +20,7 @@ public final class DatasetService {
private final JavaType datasets; // DataWrapper<DataSetModel>
private final JavaType policyNodes; // DataWrapper<Resource>

public DatasetService(ApiHttp http) {
DatasetService(ApiHttp http) {
this.http = http;
this.datasets = http.typeFactory().constructParametricType(DataWrapper.class, DataSetModel.class);
this.policyNodes = http.typeFactory().constructParametricType(DataWrapper.class, Resource.class);
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
// SPDX-License-Identifier: Apache-2.0
package ai.intellistream.datahub.sdk.ingest;
package ai.intellistream.datahub.sdk.client;

import tools.jackson.databind.JavaType;
import tools.jackson.databind.json.JsonMapper;
Expand Down Expand Up @@ -53,7 +53,7 @@
*
* @param <T> the spooled item type (e.g. a datapoint line or an event)
*/
public final class DurableSpool<T> {
final class DurableSpool<T> {

private static final Logger log = LoggerFactory.getLogger(DurableSpool.class);

Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
// SPDX-License-Identifier: Apache-2.0
package ai.intellistream.datahub.sdk.services;
package ai.intellistream.datahub.sdk.client;

import ai.intellistream.datahub.api.responses.DataWrapper;
import ai.intellistream.datahub.api.responses.GraphDataWrapper;
Expand All @@ -9,7 +9,6 @@
import ai.intellistream.datahub.models.RelationshipType;
import ai.intellistream.datahub.models.Resource;
import ai.intellistream.datahub.resource.RelTypeForm;
import ai.intellistream.datahub.sdk.http.ApiHttp;
import tools.jackson.databind.JavaType;
import tools.jackson.databind.type.TypeFactory;

Expand All @@ -30,7 +29,7 @@ public final class EdgeService {
private final JavaType edgeGraph; // GraphDataWrapper<Resource, EdgeProxy>
private final JavaType relationshipTypes; // DataWrapper<RelationshipType>

public EdgeService(ApiHttp http) {
EdgeService(ApiHttp http) {
this.http = http;
TypeFactory tf = http.typeFactory();
this.edges = tf.constructParametricType(DataWrapper.class, EdgeProxy.class);
Expand Down
Original file line number Diff line number Diff line change
@@ -1,9 +1,10 @@
// SPDX-License-Identifier: Apache-2.0
package ai.intellistream.datahub.sdk.ingest;
package ai.intellistream.datahub.sdk.client;

import ai.intellistream.datahub.sdk.ingest.IngestOptions;
import ai.intellistream.datahub.sdk.ingest.IngestResult;
import ai.intellistream.datahub.api.responses.DataWrapper;
import ai.intellistream.datahub.models.EventModel;
import ai.intellistream.datahub.sdk.http.ApiHttp;
import tools.jackson.databind.JavaType;

import java.util.ArrayList;
Expand All @@ -13,7 +14,7 @@
* Sends events concurrently: chunks them into batches of at most {@code batchSize} and runs the
* batches through {@link BatchExecutor}.
*/
public final class EventIngestor {
final class EventIngestor {

private final ApiHttp http;
private final String path;
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
// SPDX-License-Identifier: Apache-2.0
package ai.intellistream.datahub.sdk.services;
package ai.intellistream.datahub.sdk.client;

import ai.intellistream.datahub.api.responses.DataWrapper;
import ai.intellistream.datahub.models.EventModel;
Expand All @@ -8,9 +8,6 @@
import ai.intellistream.datahub.models.UpdateEventForm;
import ai.intellistream.datahub.models.events.EventFilter;
import ai.intellistream.datahub.models.events.EventRetreiver;
import ai.intellistream.datahub.sdk.http.ApiHttp;
import ai.intellistream.datahub.sdk.ingest.DurableSpool;
import ai.intellistream.datahub.sdk.ingest.EventIngestor;
import ai.intellistream.datahub.sdk.ingest.IngestOptions;
import ai.intellistream.datahub.sdk.ingest.IngestResult;
import ai.intellistream.datahub.sdk.util.UuidV7;
Expand Down Expand Up @@ -38,11 +35,11 @@ public final class EventService {
private final JavaType stringWrapper; // DataWrapper<String>
private final JavaType countType; // Map<String, Object> — /events/count returns {"count": N}

public EventService(ApiHttp http) {
EventService(ApiHttp http) {
this(http, null);
}

public EventService(ApiHttp http, DurableSpool<EventModel> spool) {
EventService(ApiHttp http, DurableSpool<EventModel> spool) {
this.http = http;
this.spool = spool;
this.ingestor = new EventIngestor(http, CREATE_PATH);
Expand Down
Loading
Loading