Skip to content
Settings

Client compute nodes

A calculation member is an external gRPC client that participates in entity workflow processing on the Cyoda platform. The platform delegates work to your client over a persistent bidirectional gRPC stream, and your client returns results on the same stream. For the rationale behind preferring gRPC over HTTP for compute nodes, see APIs and surfaces.

β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” gRPC (bidirectional stream) β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
β”‚ Cyoda Platform β”‚ ◄──────────────────────────────────────────►│ Your Calculation β”‚
β”‚ β”‚ CloudEvent (Protobuf, JSON payload) β”‚ Member (Client) β”‚
β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”‚ β”‚
β”‚ β”‚ Workflow Engineβ”‚ β”‚ 1. Client opens stream, sends Join β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚
β”‚ β”‚ β”‚ β”‚ 2. Server responds with Greet β”‚ β”‚ Business Logic β”‚ β”‚
β”‚ β”‚ - Processors │──┼──3. Server pushes Processing/Criteria reqs──┼──│ β”‚ β”‚
β”‚ β”‚ - Criteria β”‚ β”‚ 4. Client returns responses β”‚ β”‚ - Data transforms β”‚ β”‚
β”‚ β”‚ β”‚ β”‚ 5. Keep-alive heartbeats (bidirectional) β”‚ β”‚ - Criteria checks β”‚ β”‚
β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜

Three types of work can be delegated:

Use Case Description Request Type Response Type
Processing Perform actions, such as transforming entity data during a workflow transition, performing CRUD ops on other entities, running reports, interacting with other systems, etc. EntityProcessorCalculationRequest EntityProcessorCalculationResponse
Criteria Evaluation Evaluate a boolean condition (e.g., β€œshould this transition fire?”) EntityCriteriaCalculationRequest EntityCriteriaCalculationResponse
Function Compute and return a declared typed value without mutating anything β€” currently the firing time of a scheduled transition. EntityFunctionCalculationRequest EntityFunctionCalculationResponse
  • Transport: gRPC bidirectional streaming via CloudEventsService.startStreaming
  • Message format: CNCF CloudEvents Protobuf envelope with JSON text_data payload
  • Authentication: Bearer JWT token in gRPC Authorization metadata header
  • Auth context propagation: The platform attaches CloudEvents Auth Context extension attributes to processor and criteria requests, identifying the principal whose action triggered the workflow (see Section 9)
  • Serialization: All payloads are JSON-serialized inside CloudEvent text_data (not binary protobuf)

Your client needs the following proto files to generate gRPC stubs:

  • cloudevents.proto β€” The standard CloudEvents Protobuf message definition (package io.cloudevents.v1)
  • cyoda-cloud-api.proto β€” The Cyoda service definition (package org.cyoda.cloud.api.grpc)

The service definition:

service CloudEventsService {
rpc startStreaming(stream io.cloudevents.v1.CloudEvent) returns (stream io.cloudevents.v1.CloudEvent);
}

The CloudEvent message:

message CloudEvent {
string id = 1; // Unique event ID (UUID recommended)
string source = 2; // URI-reference identifying the event source
string spec_version = 3; // Must be "1.0"
string type = 4; // Event type string (see Section 4)
map<string, CloudEventAttributeValue> attributes = 5;
oneof data {
bytes binary_data = 6;
string text_data = 7; // ← Used by Cyoda (JSON payload)
google.protobuf.Any proto_data = 8;
}
}

Obtain a valid JWT Bearer token from the Cyoda IAM system (OAuth 2.0 client credentials flow). The token must contain:

  • A valid caas_org_id claim (your legal entity ID)
  • Valid user roles

The token is validated on every stream establishment. If the token expires during an active stream, the stream remains valid β€” re-authentication occurs only when a new stream is opened.

For JVM-based clients, the recommended dependencies are:

  • io.grpc:grpc-stub, io.grpc:grpc-protobuf, io.grpc:grpc-netty-shaded β€” gRPC runtime
  • io.cloudevents:cloudevents-protobuf β€” CloudEvents SDK Protobuf format support
  • io.cloudevents:cloudevents-core β€” CloudEvents SDK core
  • com.fasterxml.jackson.core:jackson-databind β€” JSON serialization

ManagedChannel channel = ManagedChannelBuilder
.forAddress("cyoda-host.example.com", 50051)
.usePlaintext() // Use .useTransportSecurity() for TLS in production
.keepAliveTime(30, TimeUnit.SECONDS)
.keepAliveTimeout(10, TimeUnit.SECONDS)
.build();

Production TLS: In production, always use TLS. Replace .usePlaintext() with:

.useTransportSecurity()
.sslContext(/* your SSL context */)

Create a CallCredentials implementation that injects the Authorization header:

CallCredentials callCredentials = new CallCredentials() {
@Override
public void applyRequestMetadata(RequestInfo requestInfo, Executor executor, MetadataApplier applier) {
executor.execute(() -> {
Metadata headers = new Metadata();
headers.put(
Metadata.Key.of("Authorization", Metadata.ASCII_STRING_MARSHALLER),
"Bearer " + jwtTokenSupplier.get() // Always fetch a fresh token
);
applier.apply(headers);
});
}
};
CloudEventsServiceGrpc.CloudEventsServiceStub asyncStub = CloudEventsServiceGrpc
.newStub(channel)
.withCallCredentials(callCredentials)
.withWaitForReady(); // Wait for the channel to become ready before sending

Every message on the stream is a CloudEvent with a type field that determines how to deserialize the JSON text_data. Your client must handle the following types:

CloudEvent type Direction Description
CalculationMemberJoinEvent Client β†’ Server Register as a calculation member
CalculationMemberGreetEvent Server β†’ Client Server confirms registration
CalculationMemberKeepAliveEvent Bidirectional Heartbeat probe and response
EventAckResponse Server β†’ Client Acknowledgment of keep-alive
EntityProcessorCalculationRequest Server β†’ Client Process entity data
EntityProcessorCalculationResponse Client β†’ Server Return processed entity data
EntityCriteriaCalculationRequest Server β†’ Client Evaluate a boolean criterion
EntityCriteriaCalculationResponse Client β†’ Server Return criterion result
EntityFunctionCalculationRequest Server β†’ Client Compute a typed value (no mutation)
EntityFunctionCalculationResponse Client β†’ Server Return the typed result

To send a CloudEvent on the stream (Java/Kotlin with CloudEvents SDK):

// 1. Build the CloudEvents SDK event
io.cloudevents.CloudEvent sdkEvent = CloudEventBuilder.v1()
.withType("CalculationMemberJoinEvent") // Must match the type table above
.withSource(URI.create("my-calculation-member"))
.withId(UUID.randomUUID().toString())
.withData(PojoCloudEventData.wrap(event, e -> objectMapper.writeValueAsBytes(e)))
.build();
// 2. Serialize to Protobuf
EventFormat protobufFormat = EventFormatProvider.getInstance()
.resolveFormat("application/cloudevents+protobuf");
byte[] protoBytes = protobufFormat.serialize(sdkEvent);
// 3. Parse to the gRPC CloudEvent message
io.cloudevents.v1.proto.CloudEvent grpcEvent =
io.cloudevents.v1.proto.CloudEvent.parseFrom(protoBytes);
// From the gRPC StreamObserver<CloudEvent>.onNext(value):
String eventType = value.getType();
String jsonPayload = value.getTextData();
// Deserialize based on type
switch (eventType) {
case "CalculationMemberGreetEvent":
GreetEvent greet = objectMapper.readValue(jsonPayload, GreetEvent.class);
break;
case "EntityProcessorCalculationRequest":
ProcessorRequest req = objectMapper.readValue(jsonPayload, ProcessorRequest.class);
break;
// ... etc
}

StreamObserver<CloudEvent> requestObserver = asyncStub.startStreaming(
new StreamObserver<CloudEvent>() {
@Override
public void onNext(CloudEvent value) {
// Dispatch based on value.getType() β€” see Sections 6–8
}
@Override
public void onError(Throwable t) {
// Connection lost β€” trigger reconnect (see Section 12)
}
@Override
public void onCompleted() {
// Server closed the stream β€” trigger reconnect
}
}
);

Immediately after opening the stream, send a CalculationMemberJoinEvent:

{
"id": "<uuid>",
"tags": ["my-processor-tag", "production"]
}

Tags are critical for routing. The platform routes processing/criteria requests to members whose tags are a superset of the tags configured on the workflow processor/criterion. Tags are case-insensitive (lowercased server-side).

The server responds with a CalculationMemberGreetEvent:

{
"id": "<uuid>",
"success": true,
"memberId": "<server-assigned-member-uuid>",
"joinedLegalEntityId": "<your-legal-entity-id>"
}

Store the memberId β€” you will need it for keep-alive messages.

If success is false, inspect the error object for the failure reason (e.g., subscription limit exceeded, invalid token).

The platform periodically probes your member with CalculationMemberKeepAliveEvent messages to verify liveness. You must respond to each probe with an EventAckResponse.

Server-initiated keep-alive probe (Server β†’ Client):

{
"id": "<probe-uuid>",
"memberId": "<your-member-id>"
}

Required response (Client β†’ Server):

{
"id": "<new-uuid>",
"sourceEventId": "<probe-uuid>",
"success": true
}

You may also send client-initiated keep-alive messages to confirm your own liveness.

⚠️ The server does not reply to an inbound member keep-alive. It is liveness-only: it refreshes your member’s last-seen timestamp and produces no response event. Each side pings on its own ticker.

Do not write a handler that echoes an inbound keep-alive. Two peers that each echo what they receive form a zero-delay feedback loop, which pins both processes at ~100% CPU with nothing above Debug in the logs. If your handler waits for an ack after it sends a keep-alive, remove that wait. No ack arrives.

Anything you send counts as activity: a keep-alive, a processor response, a criteria response, or an EventAckResponse all refresh liveness. The server closes the stream if it sees nothing from you within CYODA_KEEPALIVE_TIMEOUT seconds (default 30), and sends its own probes every CYODA_KEEPALIVE_INTERVAL seconds (default 10).

Timing parameters (server-side defaults):

Variable Default Description
CYODA_KEEPALIVE_INTERVAL 10 s Interval between server-sent keep-alive probes.
CYODA_KEEPALIVE_TIMEOUT 30 s Inactivity after which the server terminates the stream.

⚠️ Going idle costs you the stream, not just routing. There is no intermediate β€œmarked not alive but still registered” state that a later probe recovers from β€” exceeding CYODA_KEEPALIVE_TIMEOUT closes the stream with DeadlineExceeded. Keep your keep-alive handler fast and non-blocking, and implement the reconnection logic in Section 12.1; it is the only way back.

The server evicts a member for two conditions. The first is inbound silence for CYODA_KEEPALIVE_TIMEOUT. The second is a stalled write: one server-to-client write blocked for CYODA_KEEPALIVE_TIMEOUT. The second condition evicts a node that continues to send keep-alives while its application is blocked.

A slow consumer applies back pressure to its own stream only, not to the streams of other members.


When an entity reaches a workflow transition with an externalized processor configured to match your member’s tags, the platform sends an EntityProcessorCalculationRequest.

{
"id": "<event-uuid>",
"requestId": "<correlation-id>",
"entityId": "<entity-uuid>",
"processorId": "<processor-uuid>",
"processorName": "<configured-processor-name>",
"transactionId": "<transaction-uuid>",
"workflow": {
"id": "<workflow-uuid>",
"name": "<workflow-name>"
},
"transition": {
"id": "<transition-uuid>",
"name": "<transition-name>",
"stateFrom": "<source-state>",
"stateTo": "<target-state>"
},
"parameters": { /* arbitrary JSON configured on the processor */ },
"payload": {
"type": "TREE",
"data": { /* entity data as JSON β€” present only if attachEntity=true */ },
"meta": { /* entity metadata */ }
}
}

Key fields:

  • requestId β€” You must echo this back in the response for correlation.
  • entityId β€” The entity being processed. Echo this back.
  • processorName β€” Use this to dispatch to different business logic handlers.
  • parameters β€” Arbitrary JSON configured on the processor in the workflow definition (the context field). Use for passing configuration to your handler.
  • payload.data β€” The entity data. Only present when attachEntity is true in the workflow configuration.

πŸ’‘ Auth context: The CloudEvent envelope for this request also carries auth context extension attributes (authtype, authid, authclaims) identifying the principal whose action triggered the workflow. See Section 9 for details on how to extract them.

{
"id": "<new-uuid>",
"requestId": "<echo-request-id>",
"entityId": "<echo-entity-id>",
"success": true,
"payload": {
"type": "TREE",
"data": { /* modified entity data to write back */ }
}
}

Rules:

  1. requestId must exactly match the value from the request.
  2. entityId must exactly match the value from the request.
  3. If you set success: true, the platform applies your payload.data to the entity.
  4. If you set success: false, the platform treats this as a processing failure. Include an error object.
  5. The payload field is optional. If omitted (or payload.data is null), no data modification occurs.
{
"id": "<new-uuid>",
"requestId": "<echo-request-id>",
"entityId": "<echo-entity-id>",
"success": false,
"error": {
"code": "BUSINESS_ERROR",
"message": "Detailed error description",
"retryable": true
}
}

The error.retryable flag tells the platform whether it should retry the request on a different member (if a retry policy is configured). Set to true for transient failures and false for permanent failures.

Separately from errors you return, the platform raises its own codes when a callout cannot be delivered at all. These surface uniformly as a retryable 503 across all three callout kinds β€” processor, criteria, and function:

Code Meaning
NO_COMPUTE_MEMBER_FOR_TAG No connected member carries the required calculationNodesTags.
COMPUTE_MEMBER_DISCONNECTED The selected member dropped mid-dispatch.
DISPATCH_TIMEOUT No response within the configured response timeout.
DISPATCH_FORWARD_FAILED Cross-node forwarding to the member’s owner failed.

Previously some of these β€” a missing compute member most visibly β€” fell through to a misleading 400 WORKFLOW_FAILED, which looks like a caller mistake rather than a deployment problem. If you have error handling that treats WORKFLOW_FAILED as non-retryable, re-check it: these now arrive as 503 and are worth retrying.


When a workflow transition has an externalized criterion configured as a function, the platform sends an EntityCriteriaCalculationRequest.

{
"id": "<event-uuid>",
"requestId": "<correlation-id>",
"entityId": "<entity-uuid>",
"criteriaId": "<criteria-uuid>",
"criteriaName": "<configured-function-name>",
"target": "TRANSITION",
"transactionId": "<transaction-uuid>",
"workflow": {
"id": "<workflow-uuid>",
"name": "<workflow-name>"
},
"transition": {
"id": "<transition-uuid>",
"name": "<transition-name>",
"stateFrom": "<source-state>",
"stateTo": "<target-state>"
},
"processor": {
"id": "<processor-uuid>",
"name": "<processor-name>"
},
"parameters": { /* arbitrary JSON */ },
"payload": {
"type": "TREE",
"data": { /* entity data */ }
}
}

The target field indicates what the criterion is attached to:

Target Meaning Available Context
WORKFLOW Workflow-level criterion (selects which workflow applies) workflow
TRANSITION Transition-level criterion (should this transition fire?) workflow, transition
PROCESSOR Processor-level criterion (should this processor run?) workflow, transition, processor
NA Reserved for future use β€”

πŸ’‘ Auth context: Like processor requests, criteria requests also carry auth context extension attributes on the CloudEvent envelope. See Section 9.

{
"id": "<new-uuid>",
"requestId": "<echo-request-id>",
"entityId": "<echo-entity-id>",
"success": true,
"matches": true,
"reason": "Entity meets all validation criteria"
}

Key fields:

  • requestId β€” Must exactly match the request.
  • entityId β€” Must exactly match the request.
  • matches β€” The boolean result: true means the criterion is satisfied (transition fires / processor runs), false means it is not.
  • reason β€” Optional explanation for a false result; see below.

If success: false, the platform treats it as a criteria evaluation failure (the criterion evaluates to false by default).

reason reaches the caller and the audit trail, so populate it on every false.

Where it surfaces depends on how the criterion was reached:

  • A manual, explicitly-requested transition rejected by its criterion appends the reason to the 400 WORKFLOW_FAILED detail: transition "<name>" criterion not matched: <reason>. This is the guaranteed, backend-independent surface.
  • Automated cascade and workflow-selection paths additionally record it durably on the state-machine audit trail β€” TRANSITION_NOT_MATCH_CRITERION carries {workflowName, transition, criterion, reason} in its data, and WORKFLOW_SKIP carries {workflowName, reason}.

Reasons are capped at 2 KiB. If you omit one, the audit trail falls back to "criterion did not match" while the 400 detail is left bare β€” so an omitted reason really does leave the caller with nothing. Write the specific cause: "order total 4200 exceeds the 1000 auto-approval limit" beats "validation failed".


A Function callout is the third request shape your compute node may receive, alongside Processor and Criteria. The three differ in what they are allowed to do:

Callout Returns Mutates the entity?
Processor an updated entity payload Yes
Criteria a boolean matches No
Function a declared typed value No

A Function computes and returns a value without side effects. The response carries a resultKind discriminator naming the shape of result, so the caller can validate what it receives.

The request is an EntityFunctionCalculationRequest. It mirrors a processor request, but the callout target is named by functionId / functionName β€” there is no processorName on this message, so route on functionName.

Field Required Description
requestId YES Echo back unchanged.
entityId YES Echo back unchanged.
functionId YES Identifier of the invoked function.
functionName YES The routing key β€” the registered schedule.function.name. Dispatch on this.
workflow YES The workflow this callout belongs to.
transition NO The transition carrying the schedule.
transactionId NO The originating transaction.
payload NO The entity data β€” present only when attachEntity is true.
parameters NO The configured context string, forwarded verbatim.
{
"id": "<new-uuid>",
"requestId": "<echo-request-id>",
"entityId": "<echo-entity-id>",
"success": true,
"resultKind": "Schedule",
"result": { "fireAfterMs": 3600000, "expireAfterMs": 600000 }
}
  • resultKind β€” names the shape of result. "Schedule" is currently the only defined kind.
  • result β€” the typed value itself.
  • success: false, or an error, fails the dispatch exactly as a processor or criteria failure does.

Schedule drives a scheduled transition’s schedule.function, computing when a transition should fire for one specific entity:

  • Fire time (required) β€” exactly one of fireAt (absolute epoch-ms) or fireAfterMs (relative to arm time).
  • Expiry (optional) β€” at most one of expireAt or expireAfterMs, the latter relative to the resolved fire time rather than to arm time.

⚠️ This callout runs inside the caller’s write transaction. It is invoked synchronously while the entity write that arms the timer is still open, and it is re-invoked on every re-arm β€” not once per entity. A slow handler slows every write to that entity; a failing or unreachable handler fails the write with a retryable 503. Keep it fast, side-effect free, and dependent only on data you already have.

Returning a malformed result, or one that does not match the declared resultKind, fails the write with 500 SCHEDULE_FUNCTION_INVALID_RESULT.


The platform attaches CloudEvents Auth Context extension attributes to every EntityProcessorCalculationRequest and EntityCriteriaCalculationRequest. These attributes identify the authenticated principal whose action triggered the workflow execution (e.g., the user who created or updated the entity).

The auth context is carried as CloudEvent extension attributes in the Protobuf attributes map β€” not inside the JSON text_data payload.

Attribute Type Required Description
authtype String YES Principal kind. Exactly one of: user, service, system
authid String NO Unique identifier of the principal (UUID). Absent for system.
authclaims String NO Comma-separated role names on cyoda-go (e.g. ROLE_ADMIN,ROLE_USER). Absent when the principal has no roles. Never contains credentials.
authtype Value Meaning
user A regular authenticated user (JWT-based login)
service A machine-to-machine (M2M) technical account
system An internal platform trigger with no user context

The value is driven by the principal’s explicit kind, recorded when the principal is established. The platform does not infer it from the roles the principal carries.

The attribute is always present and always faithful. If a principal’s kind is unset or unrecognized, the platform fails the callout dispatch rather than emitting a normalized placeholder β€” a bogus authtype never reaches your compute node.

⚠️ The only values are user, service and system. There is no service_account, unauthenticated or unknown. A handler that switches on one of those strings never matches it. Do not rely on a default or else branch to catch an unroutable principal either: the platform fails the dispatch rather than sending one.

authtype and authid describe the principal the action is attributed to, which is not necessarily the identity it executes with. When a user’s action sets off a follow-on β€” a cascade write from a processor, or a scheduled transition firing later β€” the platform executes it with system or service authority, borrowing no user permissions, while still attributing it to the principal that caused it. Treat the auth context as provenance for audit and business logic, not as proof of the executing identity’s privileges.

The attributes are available in the Protobuf CloudEvent’s attributes map. The keys are the attribute names listed above (no prefix):

// From the gRPC StreamObserver<CloudEvent>.onNext(value):
String authType = value.getAttributesMap().get("authtype").getCeString();
String authId = value.getAttributesMap().containsKey("authid")
? value.getAttributesMap().get("authid").getCeString()
: null;
String authClaims = value.getAttributesMap().containsKey("authclaims")
? value.getAttributesMap().get("authclaims").getCeString()
: null;
// cyoda-go sends roles as a comma-separated string, NOT JSON.
List<String> roles = (authClaims == null || authClaims.isEmpty())
? List.of()
: Arrays.asList(authClaims.split(","));

The exact accessor depends on your gRPC tooling β€” in Go, use the generated message’s GetAttributes() method; in Python, dict-like indexing on .attributes. See your language’s generated proto bindings.

Go clients can skip the manual extraction. cyoda-go ships an api/grpc/authctx helper exposing Type, ID and Roles readers plus Require(ce, role). Require is a fail-closed role gate β€” it returns false for a nil event, for absent or empty claims, and for authtype == system. Prefer it over hand-rolled claim parsing, which tends to fail open on exactly those cases.

On cyoda-go, authclaims is a flat comma-separated list of role names:

ROLE_USER,ROLE_SUPER_USER

For a service (M2M) principal:

ROLE_M2M

⚠️ Do not parse this as JSON. objectMapper.readValue(authClaims, Map.class) throws on ROLE_USER,ROLE_SUPER_USER. Split on ,. Go clients should use authctx.Roles, which does exactly this.

Cyoda Cloud has historically sent a richer JSON claims object including legalEntityId. If you target both, branch on whether the value starts with { rather than assuming either shape.

  • Audit logging: Record which user triggered the processing for compliance.
  • Authorization decisions: Apply different business logic based on the caller’s roles or legal entity.
  • Multi-tenant isolation: Verify the triggering principal belongs to the expected tenant.
  • Debugging: Trace processing failures back to the originating user action.

⚠️ Note: The authclaims field never contains credentials (passwords, tokens, secrets). It contains only identity and authorization metadata.


Your calculation member does not exist in isolation β€” it is invoked by workflow configurations on the platform side. This section describes how workflows reference externalized processors and criteria, so you understand the relationship between your member’s tags/handlers and the platform configuration.

{
"workflows": [{
"version": "1",
"name": "my-workflow",
"initialState": "start",
"states": {
"start": {
"transitions": [{
"name": "process-data",
"next": "processed",
"manual": false,
"processors": [{
"name": "my-processor-function",
"executionMode": "SYNC",
"config": {
"attachEntity": true,
"calculationNodesTags": "my-processor-tag",
"responseTimeoutMs": 60000,
"retryPolicy": "FIXED",
"context": "{\"key\": \"value\"}"
}
}]
}]
},
"processed": {}
}
}]
}
Field Type Default Description
name string β€” Required. The processor name. Sent as processorName in the request.
executionMode string β€” Required. One of SYNC, ASYNC_SAME_TX, ASYNC_NEW_TX, COMMIT_BEFORE_DISPATCH.
config.attachEntity boolean true Whether to send entity data in the request payload. A processor that omits this field is imported with attachEntity: true. Set it to false explicitly to opt out.
config.calculationNodesTags string "" Comma/semicolon-separated tags. Only members whose tags are a superset are eligible.
config.responseTimeoutMs long 30000 How long the platform waits for your response before timing out.
config.retryPolicy string FIXED NONE β€” no retry. FIXED β€” retry with fixed delay (default: 3 retries, 500ms delay).
config.context string null Arbitrary string passed as parameters in the request. Use for handler-specific configuration.
config.asyncResult boolean false Enable async response processing (advanced).
config.crossoverToAsyncMs long 5000 Time before switching from sync to async response handling (advanced).
startNewTxOnDispatch boolean false Whether a fresh transaction is started when this processor is dispatched. Valid only when executionMode is COMMIT_BEFORE_DISPATCH; the validator rejects true for any other mode. With COMMIT_BEFORE_DISPATCH and startNewTxOnDispatch=false (the default), callbacks run standalone instead of joining the originating transaction (see Β§10.3.1).
Mode Behavior
SYNC The workflow engine waits for your response within the same transaction. The transition completes only after your response is applied.
ASYNC_SAME_TX The engine sends the request and can process other work. Your response is applied within the same entity transaction.
ASYNC_NEW_TX Like ASYNC_SAME_TX, but your response is applied in a new transaction. Useful for long-running computations.
COMMIT_BEFORE_DISPATCH Commits the originating transaction before dispatching the request, releasing the storage connection for the duration of the external compute; the processor runs outside any transaction unless startNewTxOnDispatch=true opens one for it (see Β§10.3.1). On return, the result is reapplied via CompareAndSave in a new transaction. Relieves connection-pool pressure for slow processors; supersedes ASYNC_NEW_TX as the recommended mode for slow external work.

For most use cases, SYNC is the simplest and recommended starting point.

Processor and criteria-evaluation callbacks from a compute node join the originating workflow transaction T rather than running standalone. Before each dispatch the engine mints a signed HMAC transaction token and attaches it to the outbound CloudEvent as the cyodatxtoken extension attribute. Your compute node echoes it back on callbacks β€” as the X-Tx-Token HTTP header or the tx-token gRPC metadata β€” and the receiving node verifies the HMAC and routes the callback to the transaction owner (a local join when the owner is this node, or an HTTP reverse-proxy / gRPC forward otherwise).

The practical consequences:

  • Callbacks see the cascade’s uncommitted writes, and their acks stay provisional until T commits.
  • ASYNC_NEW_TX callbacks join T via a savepoint, so a processor failure discards its own writes without aborting the whole cascade.
  • If no token is present the callback falls back to standalone execution β€” the normal behaviour for COMMIT_BEFORE_DISPATCH with startNewTxOnDispatch=false.
  • The token covers the search RPCs as well as the write RPCs. A callback that presents a valid tx-token on EntitySearch or EntitySearchCollection joins T, so a processor’s searches see the writes the same processor has already made in that transaction.
  • pointInTime selects committed state, even inside a joined transaction. A point-in-time read is committed-only on every backend, so it does not see the writes T has not yet committed. Pass it when you want the committed state as at an instant, which is the usual reason to read history from a processor. Omit it when you want to read back what this transaction has written: a current-state read inside T is read-your-own-writes correct. An entity that this transaction created, read back with pointInTime, answers 404 ENTITY_NOT_FOUND.
  • A request that joins an open transaction cannot carry its own deadline. The server rejects transactionTimeoutMillis, transactionSize and search’s timeoutMillis with 400 on a joined callback. The callback does not own the transaction.
  • The model governs the data you return. See the model governs the data a processor returns.

Three environment variables tune this (full list in the configuration reference):

  • CYODA_TX_TOKEN_TTL β€” validity of the signed token (default 1m30s).
  • CYODA_GRPC_NODE_ADDR β€” this node’s gRPC endpoint advertised to peers (host:port, no scheme), used for Bβ†’A forwarding.
  • CYODA_COMPUTE_HTTP_BASE β€” base URL of the cyoda instance a compute node calls back into (compute-client side).

10.4 Externalized Criteria (type: function) in Workflow JSON

Section titled β€œ10.4 Externalized Criteria (type: function) in Workflow JSON”
{
"transitions": [{
"name": "conditional-transition",
"next": "target-state",
"manual": false,
"criterion": {
"type": "function",
"function": {
"name": "my-criteria-function",
"config": {
"attachEntity": true,
"calculationNodesTags": "my-processor-tag",
"responseTimeoutMs": 5000,
"retryPolicy": "NONE"
}
}
}
}]
}

Criteria functions use the same config fields as processors (except asyncResult and crossoverToAsyncMs, which are not applicable to criteria).

Policy Behavior
NONE No retry. If the member fails or times out, the processing fails.
FIXED Retries up to N times (default: 3) with a fixed delay (default: 500ms) between retries. Each retry attempts a different member if available (the failed member is excluded from selection).

All events on the stream extend the BaseEvent schema:

{
"id": "<string, required>",
"success": true,
"error": {
"code": "<string, required if error>",
"message": "<string, required if error>",
"retryable": false
},
"warnings": ["<optional array of warning strings>"]
}
  • id β€” Every event must have a unique ID (UUID recommended).
  • success β€” Defaults to true. Set to false to indicate an error.
  • error β€” Only relevant when success is false. The code and message fields are required within the error object.
  • warnings β€” Optional array of warning strings.

gRPC streams can be terminated by network issues, server restarts, or load balancer timeouts. Implement automatic reconnection:

  1. Detect disconnection via onError or onCompleted on the response observer.
  2. Back off exponentially β€” start at 1 second, cap at 60 seconds.
  3. Re-join after reconnect β€” every new stream requires a fresh CalculationMemberJoinEvent.
  4. Refresh the JWT token before reconnecting if it is near expiry.
β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β” onError/onCompleted β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” delay β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” success β”Œβ”€β”€β”€β”€β”€β”€β”
β”‚ Connectedβ”‚ ──────────────────────► β”‚ Backoff β”‚ ────────► β”‚ Reconnecting β”‚ ────────────► β”‚ Join β”‚
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”˜
β–² β”‚ failure β”‚
β”‚ β–Ό β”‚
β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚
β”‚ β”‚ Backoff β”‚ (increase delay) β”‚
β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
Greet received

The gRPC StreamObserver is not thread-safe. If your business logic runs on multiple threads, synchronize all calls to observer.onNext():

synchronized (requestObserver) {
requestObserver.onNext(cloudEvent);
}

Your client must respond within the configured responseTimeoutMs (default 30 seconds β€” defaultResponseTimeoutMs in the engine). If you exceed this:

  • The platform considers the request failed.
  • If retry policy is FIXED, the platform retries with a different member.
  • Late responses are silently discarded.

Design your business logic to complete well within the timeout, accounting for network latency.

In edge cases (e.g., network partitions, retries), you may receive the same request more than once. Use the requestId as an idempotency key to avoid processing the same request twice.

When shutting down your client:

  1. Stop accepting new requests (drain in-flight work).
  2. Complete any pending responses and send them.
  3. Close the gRPC stream via requestObserver.onCompleted().
  4. Shut down the ManagedChannel with a grace period:
    channel.shutdown().awaitTermination(10, TimeUnit.SECONDS);

The platform will detect the stream closure and broadcast a member-offline event to the cluster. Pending requests that were in-flight will time out and may be retried on other members.

You can run multiple calculation member instances (same or different processes) with the same tags for horizontal scaling and high availability. The platform selects one eligible member per request, preferring members connected to the local cluster node. Running at least two members ensures continued processing if one goes down.

Track these metrics in your client:

  • Request count by type (processor vs. criteria) and result (success vs. failure)
  • Response latency (time from receiving request to sending response)
  • Keep-alive response time
  • Reconnection count and frequency
  • Stream errors (by gRPC status code)

Client Server
β”‚ β”‚
│──── startStreaming() ─────────────────────────►│ (open bidirectional stream)
β”‚ β”‚
│──── CalculationMemberJoinEvent ───────────────►│ (register with tags)
│◄─── CalculationMemberGreetEvent ───────────────│ (server confirms, assigns memberId)
β”‚ β”‚
│◄─── CalculationMemberKeepAliveEvent ───────────│ (periodic heartbeat probe)
│──── EventAckResponse ─────────────────────────►│ (ack the probe)
β”‚ β”‚
│◄─── EntityProcessorCalculationRequest ─────────│ (process this entity)
│──── EntityProcessorCalculationResponse ───────►│ (here's the result)
β”‚ β”‚
│◄─── EntityCriteriaCalculationRequest ──────────│ (evaluate this criterion)
│──── EntityCriteriaCalculationResponse ────────►│ (matches: true/false)
β”‚ β”‚
│──── CalculationMemberKeepAliveEvent ──────────►│ (client-initiated heartbeat)
│◄─── EventAckResponse ─────────────────────────│ (server acks)
β”‚ β”‚

Symptom Likely Cause Fix
UNAUTHENTICATED on stream open Missing/invalid/expired JWT token Refresh JWT before connecting. Ensure Authorization: Bearer <token> header.
NOT_FOUND after JWT validation User not found in Cyoda for the given JWT Verify user enrollment and legal entity configuration.
Greet event has success: false Subscription limit exceeded (max client nodes) Check your subscription plan limits.
Stream closed with DeadlineExceeded Nothing sent within CYODA_KEEPALIVE_TIMEOUT (default 30 s) Ensure a non-blocking, fast keep-alive handler; check network latency. Reconnect per Section 12.1.
Requests not arriving Tags mismatch Verify your member’s tags are a superset of the workflow processor’s calculationNodesTags. Tags are case-insensitive.
Requests not arriving Member on wrong legal entity Requests only route to members in the same legal entity as the entity owner.
Request timeout Business logic too slow Optimize processing time or increase responseTimeoutMs in workflow config.
Duplicate requests Retry policy triggered Implement idempotency using requestId.
Stream drops unexpectedly Server restart, network issue, idle timeout Implement reconnection with exponential backoff (Section 12.1).
authtype is system unexpectedly Workflow triggered by an internal platform action (e.g., scheduled transition) or no user context was available This is expected for system-initiated workflows. If you expect a user context, verify the originating API call is authenticated.
authclaims is missing The triggering principal is a plain IUser without extended claims, or the auth type is system Only user and service auth types include claims. Check authtype before parsing claims.
authtype is service, not service_account service_account is not a value the platform emits Switch on service instead. See Section 9.2.
Callout never arrives; dispatch fails The originating principal’s kind is unset or unrecognized Dispatch fails closed rather than sending a bogus authtype. Check the principal’s kind on the platform side.