ValkaValka

gRPC API

gRPC service definitions and streaming protocols.

Valka's gRPC API runs on port 50051 by default. Proto files are at proto/valka/v1/.

ApiService

The main client-facing API for task management.

service ApiService {
    rpc CreateTask(CreateTaskRequest) returns (CreateTaskResponse);
    rpc GetTask(GetTaskRequest) returns (GetTaskResponse);
    rpc ListTasks(ListTasksRequest) returns (ListTasksResponse);
    rpc CancelTask(CancelTaskRequest) returns (CancelTaskResponse);
    rpc SendSignal(SendSignalRequest) returns (SendSignalResponse);
    rpc SubscribeEvents(SubscribeEventsRequest) returns (stream TaskEvent);
    rpc SubscribeLogs(SubscribeLogsRequest) returns (stream LogEntry);
}

CreateTask

Create a new task in a queue.

FieldTypeDescription
queue_namestringTarget queue
task_namestringTask identifier
inputstring (JSON)Task payload
priorityint32Priority (higher = first)
max_retriesint32Max retry attempts
timeout_secondsint32Lease timeout
idempotency_keystringDedup key
metadatastring (JSON)Arbitrary metadata
scheduled_atTimestampDelayed execution

SubscribeEvents

Server-streaming RPC. Returns a stream of TaskEvent messages for real-time monitoring.

message TaskEvent {
    string event_id = 1;
    string task_id = 2;
    string queue_name = 3;
    TaskStatus new_status = 4;
    int64 timestamp_ms = 5;
}

SubscribeLogs

Server-streaming RPC. Returns a stream of LogEntry messages for a specific task run.

WorkerService

The bidirectional streaming service used by workers.

service WorkerService {
    rpc Session(stream WorkerRequest) returns (stream WorkerResponse);
}

A single Session stream carries all communication between a worker and the server.

Worker → Server Messages

graph LR
  subgraph WorkerRequest
      WH[WorkerHello<br/>Initial handshake]
      TR[TaskResult<br/>Success or failure]
      HB[Heartbeat<br/>Lease extension]
      LB[LogBatch<br/>Log entries]
      GS[GracefulShutdown<br/>Drain signal]
      SA[SignalAck<br/>Signal confirmed]
  end
MessageWhen SentDescription
WorkerHelloOn connectWorker name, queues, concurrency
TaskResultTask doneSuccess/failure with output/error
HeartbeatEvery 30sActive task IDs for lease extension
LogBatchDuring taskStructured log entries
GracefulShutdownShutting downSignals drain mode
SignalAckSignal receivedConfirms signal delivery

Server → Worker Messages

MessageWhen SentDescription
TaskAssignmentTask matchedNew task to execute
TaskCancellationCancel requestCancel a running task
HeartbeatAckAfter heartbeatConfirms heartbeat received
ServerShutdownServer stoppingTells worker to drain
TaskSignalSignal sentReal-time signal for a task

InternalService

Used for inter-node communication in clustered deployments.

service InternalService {
    rpc ForwardTask(ForwardTaskRequest) returns (ForwardTaskResponse);
    rpc RelayLogs(RelayLogsRequest) returns (stream LogEntry);
    rpc ForwardEvent(ForwardEventRequest) returns (ForwardEventResponse);
    rpc Ping(PingRequest) returns (PingResponse);
}

Common Types

TaskStatus

enum TaskStatus {
    TASK_STATUS_UNSPECIFIED = 0;
    TASK_STATUS_PENDING = 1;
    TASK_STATUS_DISPATCHING = 2;
    TASK_STATUS_RUNNING = 3;
    TASK_STATUS_COMPLETED = 4;
    TASK_STATUS_FAILED = 5;
    TASK_STATUS_RETRY = 6;
    TASK_STATUS_DEAD_LETTER = 7;
    TASK_STATUS_CANCELLED = 8;
}

LogLevel

enum LogLevel {
    LOG_LEVEL_UNSPECIFIED = 0;
    LOG_LEVEL_DEBUG = 1;
    LOG_LEVEL_INFO = 2;
    LOG_LEVEL_WARN = 3;
    LOG_LEVEL_ERROR = 4;
}

Using with grpcurl

# Create a task
grpcurl -plaintext -d '{
  "queue_name": "emails",
  "task_name": "send-email",
  "input": "{\"to\": \"user@example.com\"}"
}' localhost:50051 valka.v1.ApiService/CreateTask

# Subscribe to events
grpcurl -plaintext localhost:50051 valka.v1.ApiService/SubscribeEvents

On this page