cs-api-client
v0.1.3
Published
TypeScript client for OGC API Connected Systems Parts 1, 2, and 3
Readme
cs-api-client
A TypeScript client for OGC API — Connected Systems Part 1 (systems, procedures, deployments, sampling features, properties, collections), Part 2 (datastreams, observations, control streams, commands, system events/history), and the Part 3 Publish/Subscribe draft.
- Fully typed models for every resource, in every encoding it supports (GeoJSON, SensorML-JSON, and Part 2 JSON), validated at runtime with Zod.
- A common model per resource (
System,Procedure,Deployment, ...) that is encoding-independent: fetch a System as SensorML or GeoJSON and get the same TypeScript shape back, with fields absent from the source encoding simply leftundefined. - One client, grouped by resource:
client.systems,client.datastreams.observations(id), etc. — plain data in, plain data out, no smart/stateful resource objects.
Install
npm install cs-api-clientRequires Node 18+ (uses the global fetch).
MQTT and RxJS integrations are optional peers. Install them only if you use the async pub/sub adapters:
npm install mqtt rxjsQuick start
import { CSApiClient } from "cs-api-client";
const client = new CSApiClient({ baseUrl: "https://data.example.org/api" });
// Default format ("common") fetches the richest encoding (SensorML for System/Procedure/
// Deployment/Property, GeoJSON for SamplingFeature) and maps it to the shared model.
const system = await client.systems.get("abc123");
console.log(system.uniqueId, system.label, system.inputs);
// Ask for a specific wire encoding instead — return type narrows accordingly.
const geo = await client.systems.get("abc123", { format: "geojson" });
console.log(geo.geometry, geo.properties.assetType);
const sml = await client.systems.get("abc123", { format: "sml" });
console.log(sml.characteristics);Listing & pagination
const page = await client.systems.list({ bbox: [-10, 40, 10, 60], limit: 50 });
console.log(page.items, page.links);
// Or drain every page automatically:
for await (const system of client.systems.listAll({ q: "weather" })) {
console.log(system.label);
}Creating & updating
const id = await client.systems.create({
uniqueId: "urn:x-example:systems:001",
label: "Outdoor Thermometer",
featureType: "http://www.w3.org/ns/sosa/Sensor",
assetType: "Equipment",
}, { format: "geojson" });
await client.systems.update(id, { ...updatedFields }, { format: "geojson" });
await client.systems.delete(id, { cascade: true });Common-model fields with no representation in the target encoding are dropped on write (documented lossy behavior) — write with the same format a resource was read in to avoid surprises.
Sub-resources
const subsystems = await client.systems.subsystems(systemId);
const deployments = await client.systems.deployments(systemId);
const samplingFeatures = await client.systems.samplingFeatures(systemId);
const datastreams = await client.datastreams.list({ system: [systemId] });
const observations = await client.datastreams.observations(dsId, { phenomenonTime: { start: "2024-01-01T00:00:00Z" } });
const controlStreams = await client.controlstreams.list({ system: [systemId] });
const commands = await client.controlstreams.commands(csId);
const status = await client.commands.status(commandId);Part 2: datastreams & observations
Observation result shape is defined by the datastream's schema, fetched separately:
const schema = await client.datastreams.schema(dsId, "application/json"); // obsFormat is required
console.log(schema.resultSchema); // SWE Common AnyComponent describing the result
const obsId = await client.datastreams.createObservation(dsId, {
resultTime: new Date().toISOString(),
result: 21.4,
});Only application/json (OM-JSON) and application/swe+json result/command payloads are decoded today; SWE-Text/CSV/Binary and Protobuf schemas parse structurally but their result bodies are treated as opaque (unknown).
Part 3: publish/subscribe
The transport-neutral Part 3 client separates lifecycle notifications from full resource messages:
resourceEventssubscribes to CloudEvents notifications for one resource operation.batchResourceEventssubscribes to aggregated CloudEvents notifications.resourceDatasubscribes to and publishes complete CS API resource representations.
The current Part 3 working draft does not yet contain a normative MQTT channel layout, discovery mechanism, or flow-control rules. Every operation therefore takes an explicit channel. MQTT topics are one protocol-specific form of those channels; the client does not invent Part 3 defaults.
import { CSApiClient, resourceDataCodecs } from "cs-api-client";
import { createMqttTransport } from "cs-api-client/mqtt";
const client = new CSApiClient({
baseUrl: "https://data.example.org/api",
pubsub: {
transport: createMqttTransport({
url: "wss://broker.example.org/mqtt",
clientOptions: { username: "user", password: "token" },
}),
},
});
const sub = await client.pubsub!.resourceData.subscribe(
"datastreams/ds1/observations",
resourceDataCodecs.observation.omJson,
{
next: (obs) => console.log(obs.resultTime, obs.result),
error: console.error,
},
);
await client.pubsub!.resourceData.publish(
"controlstreams/control1/commands",
{ parameters: { pan: 10, tilt: 5 } },
resourceDataCodecs.command.cmdJson,
);
await sub.close();
await client.pubsub!.close();Resource and Batch Resource Events are server-published, so their high-level APIs are subscribe-only:
client.pubsub!.resourceEvents.subscribe("channels/resource-events", handlers);
client.pubsub!.batchResourceEvents.subscribe("channels/batch-events", handlers);The event schemas validate the portable message contract: CloudEvents version and event vocabulary, non-empty IDs, URLs, RFC 3339 timestamps, JSON data/content-type coupling, and batch count/time-range rules. Server-context requirements cannot be established from one message alone and are therefore documented rather than structurally enforced: the server owns event publication, each source/ID pair must be unique, parentId depends on resource context, and the working draft's collection paths are internally inconsistent and require server/channel context.
Typed codecs cover all eight Part 3 Resource Data types. Systems have GeoJSON and SensorML JSON codecs; the other resource types use their Part 2 models, with OM JSON, CMD JSON, and SensorML JSON media-type variants where applicable:
resourceDataCodecs.system.geoJson;
resourceDataCodecs.system.smlJson;
resourceDataCodecs.dataStream.json;
resourceDataCodecs.controlStream.json;
resourceDataCodecs.observation.json;
resourceDataCodecs.observation.omJson;
resourceDataCodecs.command.json;
resourceDataCodecs.command.cmdJson;
resourceDataCodecs.commandStatus.json;
resourceDataCodecs.commandStatus.cmdJson;
resourceDataCodecs.commandResult.json;
resourceDataCodecs.systemEvent.json;
resourceDataCodecs.systemEvent.smlJson;Custom encodings use the same codec interface. createJsonPubSubCodec, createTextPubSubCodec, and createBinaryPubSubCodec cover JSON validation and opaque SWE Text/Binary data. formatNameFromMediaType("application/swe+binary; version=1") returns the Part 3 format name swe-binary; it does not append that name to a channel.
The selected codec is authoritative. An MQTT v3 message without content-type metadata uses the codec's media type; an MQTT v5 message advertising a different media type is rejected. Publishing always advertises the selected codec's media type.
RxJS stays optional:
import { toObservable } from "cs-api-client/rxjs";
const observations$ = toObservable(() =>
client.pubsub!.resourceData.subscribe(
"datastreams/ds1/observations",
resourceDataCodecs.observation.json,
),
);
const rxSub = observations$.subscribe((obs) => console.log(obs.result));
rxSub.unsubscribe(); // closes the underlying CS subscriptionError handling
import { HttpError, NotFoundError, PubSubError, ValidationError } from "cs-api-client";
try {
await client.systems.get("missing");
} catch (err) {
if (err instanceof NotFoundError) { /* 404 */ }
else if (err instanceof HttpError) { /* other non-2xx, err.problem has the problem+json body if any */ }
else if (err instanceof ValidationError) { /* response didn't match the expected schema, err.zodError has details */ }
else if (err instanceof PubSubError) { /* pub/sub connection, parse, subscribe, or publish failure */ }
}Auth, hooks & custom fetch
const client = new CSApiClient({
baseUrl: "https://data.example.org/api",
auth: { type: "bearer", token: async () => getToken() },
fetch: myFetchImplementation, // for testing, polyfills, proxies, etc.
});Built-in auth helpers:
// Reuse an app-managed token
new CSApiClient({
baseUrl,
auth: { type: "bearer", token: () => authStore.accessToken },
});
// Basic auth
new CSApiClient({
baseUrl,
auth: { type: "basic", username: "user", password: "pass" },
});
// OAuth token refresh. `expiresAt` is epoch milliseconds.
new CSApiClient({
baseUrl,
auth: {
type: "oauth2",
token: () => authStore.token,
refresh: (token) => refreshOAuthToken(token?.refreshToken),
setToken: (token) => authStore.save(token),
refreshBeforeExpiresIn: 60,
},
});Hooks are async-capable and run for every request:
new CSApiClient({
baseUrl,
hooks: {
beforeRequest: ({ init }) => {
(init.headers as Record<string, string>)["x-trace-id"] = traceId();
},
afterResponse: ({ response }) => {
if (response.status === 401) notifyAuthLayer();
},
onError: ({ error }) => {
reportClientError(error);
},
},
});afterResponse receives a cloned Response, so it can inspect the body without consuming the client parser. For OAuth, the client refreshes before expiry and retries once after a 401 unless retryOnUnauthorized: false.
Project layout
src/
├── models/
│ ├── common/ # Link, TimePeriod, Geometry, well-known URIs
│ ├── swe/ # SWE Common data components (the AnyComponent recursion knot)
│ ├── sensorml/ # SensorML wire encodings and shared SensorML building blocks
│ ├── geojson/ # GeoJSON wire encodings
│ └── resources/ # CS API resources: System, DataStream, Command, Observation, ...
├── codec/ # @link/@id ↔ camelCase key mapping, wire ↔ common model mappers
├── http/ # fetch wrapper, errors, query serialization, pagination
├── pubsub/ # transport-neutral Part 3 events, resource data, and codecs
└── api/ # Per-resource endpoint classes + CSApiClientDevelopment
npm test # vitest, includes fixture-parsing tests against real spec examples
npm run typecheck
npm run build