Rust + Axum + lapin service
Task Gateway is a small HTTP API service that works as a task bus between clients and downstream processing services.
The service accepts task requests through HTTP, assigns a task id, publishes the
request to RabbitMQ, and registers the initial task state through the configured
state manager. Downstream services consume messages from their own queues and
process the task asynchronously.
A successful publish response means that the task was published to RabbitMQ and registered in the state manager. It does not mean that the downstream service has completed the task.
Task Gateway is connected to RabbitMQ for task delivery and to Webhook Manager as the state manager for task creation and cancellation (see Webhook Manager). It routes tasks to these downstream service domains:
| Service domain | Service name in task key | Exchange | Queue | Task types |
|---|---|---|---|---|
| Image generation | image-generation |
images.tasks |
images.queue |
images.generate, images.edit |
| Video generation | video-generation |
videos.tasks |
videos.queue |
videos.generate, videos.animate |
The public HTTP endpoints are:
POST /api/v1/broker/publish
POST /api/v1/tasks/cancel?task_id={task_key}POST /api/v1/broker/publish
Content-Type: application/jsonExample request:
{
"user_id": "12345",
"task_type": "images.generate",
"payload": {
"model": "openrouter::google/gemini-3.1-flash-image-preview",
"prompt": "post-apocalyptic warrior standing in a ruined city",
"image_name": "warrior"
}
}Example response:
{
"task_key": "12345:image-generation:550e8400-e29b-41d4-a716-446655440000"
}The task_key format is:
user_id:service_name:task_uuid
Clients should store this key if they need to track the task in downstream APIs.
Pass the task_key returned by the publish endpoint as the task_id query
parameter:
POST /api/v1/tasks/cancel?task_id=12345:image-generation:550e8400-e29b-41d4-a716-446655440000Successful response:
{
"code": 200,
"message": "ok"
}The endpoint passes task_id to the configured state manager. The webhook
implementation updates the task progress to CANCELLED with progress 0.0.
After RabbitMQ accepts a message, Task Gateway extracts service_name from the
returned task_key, serializes the original request payload into
response_data, builds a pending TaskState, and submits it to the state
manager.
The publish operation is performed before state registration. If publishing succeeds but the state manager call fails, the API returns the state manager error even though the broker message has already been accepted.
The state manager implementation used by Task Gateway is
Webhook Manager — a Redis-backed
service that stores the state of long-running user tasks and exposes it over an
HTTP API, a WebSocket stream, and Swagger UI. It listens on port 10001 by
default and keeps task records with a TTL (3600 seconds by default).
Task Gateway talks to it through the StateManager trait
(src/modules/state_manager/), so the service can
be replaced by any other backend implementing the same two operations.
Both services use the same key format, so the task_key returned by
POST /api/v1/broker/publish can be used directly against Webhook Manager:
user_id:service_name:task_uuid
Task Gateway builds it from user_id, the route's service_name, and the
generated task UUID. service_name therefore must not contain :.
Task Gateway is a write-only client — it never reads task state back.
| Task Gateway action | Method | Endpoint (configurable) | Body |
|---|---|---|---|
| Register a published task | POST |
TASK_GATEWAY__STATE_MANAGER__CREATE_TASK_ENDPOINT (/api/v2/storage/task) |
{ "task": { ... } } |
| Cancel a task | PATCH |
TASK_GATEWAY__STATE_MANAGER__UPDATE_PROGRESS_ENDPOINT (/api/v1/storage/update_progress) |
{ "key": "...", "progress": { ... } } |
Create request body (TaskState wrapped in task):
{
"task": {
"task_id": "550e8400-e29b-41d4-a716-446655440000",
"user_id": "12345",
"service": "image-generation",
"progress": {
"status": "PENDING",
"progress": 0.0
},
"response_data": "{\"prompt\":\"post-apocalyptic warrior\"}"
}
}response_data is the original request payload serialized into a JSON string.
Cancel request body:
{
"key": "12345:image-generation:550e8400-e29b-41d4-a716-446655440000",
"progress": {
"status": "CANCELLED",
"progress": 0.0
}
}Any non-2xx response from Webhook Manager is propagated to the client as a Task Gateway error.
PENDING, AWAITING, PROCESSING, READY, ERROR, CANCELLED.
Task Gateway writes only PENDING (on publish) and CANCELLED (on cancel).
The remaining transitions are made by the downstream services that consume the
RabbitMQ queues.
These endpoints are served by Webhook Manager itself, not by Task Gateway:
| Method | Endpoint | Purpose |
|---|---|---|
GET |
/api/v1/storage/task?key={task_key} |
read a single task |
GET |
/api/v1/storage/tasks?user_id={user_id} |
list user tasks (also accepts X-User-ID) |
POST |
/api/v1/storage/task |
create a task |
DELETE |
/api/v1/storage/task?key={task_key} |
delete a task |
PATCH |
/api/v1/storage/update_progress |
update status and progress |
PATCH |
/api/v1/storage/update_response_data |
write the task result |
GET |
/api/v1/storage/ws |
WebSocket stream of the user's task list |
Typical flow: the client publishes a task through Task Gateway, stores the
returned task_key, and then polls /api/v1/storage/task or subscribes to
/api/v1/storage/ws on Webhook Manager while the downstream service updates
progress and response data.
Webhook Manager is deployed separately — it is not part of
docker-compose/docker-compose.yaml. The Docker environment expects it at
webhook-manager:10001; point TASK_GATEWAY__STATE_MANAGER__ADDRESS at the
actual host and adjust the endpoint variables if the deployed instance exposes
a different API version.
RabbitMQ topology is configured on the broker side, not by Task Gateway.
Task Gateway does not create queues, exchanges, or bindings. The service only checks that the selected exchange already exists and publishes a message to it. In code this is done with passive exchange declaration before publishing.
For Docker Compose, the broker topology is loaded from:
docker-compose/rabbitmq/definitions.json
RabbitMQ is configured to load that file through:
docker-compose/rabbitmq/rabbitmq.conf
Current broker configuration includes:
- exchanges:
images.tasks,videos.tasks - queues:
images.queue,videos.queue - bindings:
images.tasks->images.queuewithimages.generateimages.tasks->images.queuewithimages.editvideos.tasks->videos.queuewithvideos.generatevideos.tasks->videos.queuewithvideos.animate
If a task type points to an exchange or routing key that is not configured in RabbitMQ, publishing will fail. With mandatory: true, unroutable messages are returned by RabbitMQ and Task Gateway converts that into an error.
Task Gateway configuration contains the HTTP server address, RabbitMQ connection address, state manager endpoints, and the task routing table.
Default local configuration is stored in:
config/development.toml
Docker Compose environment values are stored in:
docker-compose/.env.task-gateway
Important variables:
TASK_GATEWAY__RUN_MODE=development
TASK_GATEWAY__SERVER__ADDRESS=0.0.0.0:10010
TASK_GATEWAY__BROKER__ADDRESS=amqp://rabbitmq:5672
TASK_GATEWAY__STATE_MANAGER__ADDRESS=http://webhook-manager:10001
TASK_GATEWAY__STATE_MANAGER__CREATE_TASK_ENDPOINT=/api/v2/storage/task
TASK_GATEWAY__STATE_MANAGER__UPDATE_PROGRESS_ENDPOINT=/api/v1/storage/update_progressTASK_GATEWAY__STATE_MANAGER__CREATE_TASK_ENDPOINT is used to register newly
published tasks. TASK_GATEWAY__STATE_MANAGER__UPDATE_PROGRESS_ENDPOINT is
used by the cancel endpoint to set the task state to CANCELLED.
The Docker environment expects the state manager to be reachable as
webhook-manager:10001. Override TASK_GATEWAY__STATE_MANAGER__ADDRESS when
the service uses a different host or port.
Routes are defined in TOML because arrays of route tables are easier to maintain there than through environment variables:
[[broker.routes]]
task_type = "images.generate"
exchange = "images.tasks"
service_name = "image-generation"task_type is also used as the RabbitMQ routing key. service_name becomes the
middle segment of the returned task_key and therefore must not contain :.
Duplicate task types and empty route fields cause startup to fail.
Queues, exchanges, and bindings are still owned by RabbitMQ and configured in
docker-compose/rabbitmq/definitions.json. They are not created by Task Gateway.
To add a new task type for an existing service domain, update the Task Gateway configuration and RabbitMQ definitions. Recompiling Task Gateway is not required.
Example: add images.upscale to the existing image service.
- Add the route to
config/development.tomlor the active run-mode file:
[[broker.routes]]
task_type = "images.upscale"
exchange = "images.tasks"
service_name = "image-generation"- Add a RabbitMQ binding in
docker-compose/rabbitmq/definitions.json:
{
"source": "images.tasks",
"vhost": "/",
"destination": "images.queue",
"destination_type": "queue",
"routing_key": "images.upscale",
"arguments": {}
}- Update API documentation and tests so the new task type is visible to clients.
To connect a new downstream service domain, add a Task Gateway route and create the RabbitMQ topology on the broker side.
Example: add an audio service with audio.generate.
- Add the route to the active Task Gateway configuration:
[[broker.routes]]
task_type = "audio.generate"
exchange = "audio.tasks"
service_name = "audio-generation"- Configure RabbitMQ in
docker-compose/rabbitmq/definitions.json:
{
"name": "audio.queue",
"vhost": "/",
"durable": true,
"auto_delete": false,
"arguments": {
"x-queue-type": "classic"
}
}{
"name": "audio.tasks",
"vhost": "/",
"type": "direct",
"durable": true,
"auto_delete": false,
"internal": false,
"arguments": {}
}{
"source": "audio.tasks",
"vhost": "/",
"destination": "audio.queue",
"destination_type": "queue",
"routing_key": "audio.generate",
"arguments": {}
}-
Make sure the new downstream service consumes from
audio.queue. -
Update Swagger descriptions and tests.
Task Gateway publishes messages with:
- direct exchange routing
- routing key equal to
task_type - persistent delivery mode
- publisher confirms enabled
mandatory: true
The message body is a serialized PublishMessage:
{
"task_id": "550e8400-e29b-41d4-a716-446655440000",
"user_id": "12345",
"task_type": "images.generate",
"payload": {
"prompt": "Generate an image"
}
}The payload object is service-specific. Task Gateway forwards it unchanged to
the selected downstream service. The same object is serialized as a JSON string
when it is stored in TaskState.response_data.
Run tests:
cargo testRun the service locally:
cargo run --bin run_serverRun through Docker Compose:
cd docker-compose
docker compose upSwagger UI is available at:
/docs