Task System Internals¶
This document describes how the task system is implemented in Hubuum today.
It is intentionally implementation-focused. For public API behavior, see:
Purpose¶
The task system provides a generic framework for long-running server-side work.
Current task kinds:
importexportbackupreindexremote_call
The goal is to keep task lifecycle, queueing, polling, and audit history generic, while letting each task kind supply its own execution logic and its own typed result storage.
Core model¶
The implementation has:
taskseventsrows withentity_type = 'task'import_task_resultsexport_task_outputsremote_call_results
The first two are generic framework storage.
The result/output tables are typed per task kind.
Relevant code:
- src/models/task.rs
- PostgreSQL task queue operations
- PostgreSQL task execution operations
- PostgreSQL initial migration
tasks¶
tasks is the canonical queue and lifecycle table.
Each row stores:
- task identity:
id,kind - lifecycle:
status,created_at,started_at,finished_at,updated_at - ownership: nullable
submitted_by - durable root provenance: nullable, non-FK
initiator_user_id - submission metadata:
idempotency_key,request_hash - request storage:
request_payload,request_redacted_at - progress counters:
total_items,processed_items,success_items,failed_items - terminal summary:
summary - optional validated OpenTelemetry trace link captured at admission; this is internal provenance and is not part of the public task response
The worker starts task.execute with the admission link rather than treating a
delayed queue claim as a child in the current process. See
Distributed Tracing.
Task lifecycle events¶
Task lifecycle/progress history is stored in the unified events table with
entity_type = 'task' and entity_id = tasks.id.
New lifecycle rows also store task_id and the root task
initiator_user_id. The immediate actor stays distinct: submission is a user
action, execution is a worker action, and recovery or cleanup is a system
action. Queue events record the submitter as both actor and initiator. The task
initiator remains available after submitted_by is cleared by principal
deletion because it is intentionally not a foreign key.
Typical events:
queuedvalidatingrunningsucceededfailedpartially_succeededcleanup
The task row is the current state. Task-scoped events rows are the history of
how the task got there. API enrichment resolves the union of actor and
initiator names in one query per page. Legacy lifecycle rows derive their
initiator from the task's queued event in one additional batch query rather
than one lookup per row.
import_task_results¶
import_task_results is an import-specific typed result table.
Each row records:
- task id
- item ref
- entity kind
- action
- identifier
- outcome
- error
- details
This table is intentionally separate from tasks, because result shape is task-kind-specific.
That is the current architectural rule:
tasksand task-scopedeventsrows are the generic task framework- result persistence is typed per task kind
New task kinds should introduce typed result/output tables where their result shape does not match an existing table.
Status model¶
Generic statuses are defined in src/models/task.rs:
queuedvalidatingrunningsucceededfailedpartially_succeededcancelled
Terminal statuses:
succeededfailedpartially_succeededcancelled
For imports:
queuedmeans the task row exists and is waiting to be claimedvalidatingmeans a worker has claimed it and is planning/validating the importrunningmeans planning is complete and execution is underway, or a dry-run is materializing results- terminal states are set after results are written and summary counters are finalized
How tasks enter the system¶
Imports are created through:
POST /api/v1/imports
Relevant code:
Submission flow:
- The handler serializes the request payload.
- It computes a SHA-256 request hash.
- It reads and validates
Idempotency-Keyif present. Keys must contain between 1 and 255 bytes; invalid keys return400 Bad Requestbefore task persistence. - It either reuses an existing task for that submitter/idempotency key or inserts a new one.
- It returns
202 AcceptedwithLocation: /api/v1/tasks/{id}. - It kicks the worker so the queue starts draining immediately.
Task creation itself is generic and implemented in:
When a task is created:
- the
tasksrow is inserted withstatus = queued - the full
request_payloadis stored - counters are initialized to zero
- a
queuedevent is appended
Worker model¶
The worker implementation lives in:
There are two entry points:
ensure_task_worker_runningkick_task_worker
Startup workers¶
ensure_task_worker_running is called during server startup from:
It starts a fixed number of background worker loops once per process.
Notification-on-submit¶
Queued task transactions send a PostgreSQL notification on commit. Worker
processes listen for those notifications, while same-process submission paths
also call kick_task_worker for an immediate local wakeup. The configured poll
interval is a safety net for listener reconnects and work inserted by an older
server version.
Worker count and polling¶
Worker behavior is configurable via:
HUBUUM_TASK_WORKERSHUBUUM_TASK_POLL_INTERVAL_MSHUBUUM_TASK_LEASE_SECONDSHUBUUM_TASK_HEARTBEAT_SECONDSHUBUUM_TASK_RECOVERY_INTERVAL_SECONDSHUBUUM_EXPORT_OUTPUT_CLEANUP_INTERVAL_SECONDS
The canonical env-var reference lives in:
Defaults:
HUBUUM_ACTIX_WORKERS: detected CPU countHUBUUM_TASK_WORKERS: about half the detected CPU count, minimum1HUBUUM_TASK_POLL_INTERVAL_MS:5000HUBUUM_TASK_LEASE_SECONDS:60HUBUUM_TASK_HEARTBEAT_SECONDS:20HUBUUM_TASK_RECOVERY_INTERVAL_SECONDS:30HUBUUM_EXPORT_OUTPUT_CLEANUP_INTERVAL_SECONDS:300(shared by stored export and backup artifacts; the variable name is retained for compatibility)
The HTTP worker count and background task worker count are intentionally separate.
Queue claiming¶
Task claiming is DB-backed and implemented in:
Claiming uses:
FOR UPDATESKIP LOCKED- ordering by oldest
created_at
That gives these properties:
- only one worker can claim a queued row
- multiple workers in one process are safe
- multiple app instances are also safe
- workers skip rows currently locked by another claimant instead of blocking
Claims also receive a durable lease token. Workers renew the lease while a task is active, and persistence updates are fenced by the token. If a process dies, another worker fails the expired task and records a recovery event without automatically replaying external side effects.
Claiming immediately transitions the task to:
status = validatingstarted_at = now()
After claim, the worker appends a validating event.
Dispatch¶
After a task is claimed, process_one_task dispatches by task.kind.
Current dispatch:
import-> import executorexport-> export executorbackup-> full-system backup executorreindex-> internal computed-field class rebuild executorremote_call-> remote HTTP invocation executorschema_validation-> bounded schema validation executor
This logic is in:
Computed-field reindex tasks carry a server-owned payload with the class, target evaluation revision, and a fixed object-ID upper bound. Definition mutations enqueue them in the same transaction as the revision change. They do not depend on the initiating principal still existing when a worker claims them.
Task workers carry the same permission-backend context as API handlers. Execution therefore rechecks the submitting principal through the configured local or Treetop backend. Import, export, and remote-call workers preserve and enforce the submitted token's scope snapshot against live grants. Backup workers retain their unscoped runtime-administrator requirement.
Import execution pipeline¶
Import execution happens in four major phases:
- load payload
- plan and validate
- execute
- finalize and redact
1. Load payload¶
The worker deserializes tasks.request_payload into ImportRequest.
If the payload is missing or invalid, the task is marked failed.
2. Plan and validate¶
Planning is implemented in:
Planning walks the import graph in dependency order:
- collections
- classes
- objects
- class relations
- object relations
- collection permissions
Planning resolves:
- refs inside the import document
- natural-key selectors for existing records
- permissions for each intended operation
- collisions with existing data
- schema validation for objects when class schema validation is enabled
Planning output is not just “pass/fail”.
It produces:
planned_items: executable work itemsfailures: per-item planning failuresaborted: whether policy required early stop
This is important for best_effort mode. Best-effort imports may now continue with planned work even if some items failed during planning, as long as policy did not require an early abort.
3. Execute¶
Execution mode depends on mode.atomicity.
Strict mode¶
Implemented in:
Behavior:
- all domain mutations run in one SQL transaction
- if one executed item fails, everything rolls back
- if execution succeeds, all planned items are recorded as succeeded
Important boundary:
- the import domain mutations are one transaction
- task bookkeeping is not part of that same transaction
So strict mode means “domain writes are all-or-nothing”, not “all task metadata and domain state live in one giant transaction”.
Best-effort mode¶
Implemented in:
Behavior:
- each executable item runs in its own transaction
- successful items remain committed
- failed items are recorded individually
- planning-time failures are also recorded individually
- continuation depends on permission/collision policy
This is what allows partial success.
4. Finalize¶
After execution:
- import per-item results are inserted into the import-specific typed result table
import_task_results - summary counters are written to
tasks - terminal event is appended
- the original request payload is redacted
Terminal completion, failure, event append, result persistence, and redaction are implemented together in:
Redaction means:
request_payload = NULLrequest_redacted_at = now()
The system keeps summary metadata and typed per-task-kind results, but not the original request body after completion.
Transaction boundaries¶
This is the most important correctness detail.
What is transactional in strict mode¶
The following are inside one transaction:
- collection/class/object/relation/permission mutations for planned import items
What is not in that same transaction¶
The following happen outside that domain mutation transaction:
- claiming the task
- setting
validating - appending lifecycle events
- planning and validation reads
- inserting import-specific typed result rows
- updating task counters and terminal summary
- payload redaction
That separation is intentional:
- task state must survive worker crashes
- clients need to observe progress independently of domain-transaction scope
- task audit/history should not disappear because domain execution rolled back
Collision and permission policy behavior¶
Planning and execution behavior is shaped by:
mode.atomicitymode.collision_policymode.permission_policy
Current behavior:
strictalways aborts on the first planning failuremode.permission_policy=abortstops on permission failuresmode.permission_policy=continuerecords permission failures and continues a best-effort importcollision_policy=abortrecords or aborts on collisions depending on atomicitycollision_policy=overwriteturns matching collection/class/object collisions into update operations
For relations:
- overwrite-like behavior is effectively modeled as
noopwhen a matching relation already exists
Authorization¶
Task visibility:
- owner can view
- admin can view any task
Generic endpoints:
GET /api/v1/tasks/{id}GET /api/v1/tasks/{id}/events
Import-specific endpoints:
GET /api/v1/imports/{id}GET /api/v1/imports/{id}/results
Import and export submission accepts ordinary human and service-account principals, including scoped tokens. The task stores the token's permission and resource boundary, and the worker applies it to resource-by-resource authorization against the configured backend. Revoking a required live grant or disabling the submitting service account while a task is queued prevents the affected work from executing.
Creating collections in an import, and importing identity, template, or integration records, remains restricted to unscoped runtime administrators.
Queue introspection¶
Admin queue state is exposed through:
GET /api/v0/meta/tasks
Relevant code:
This endpoint exports:
- configured worker counts
- poll interval
- task counts by status
- task counts by kind
- total task events
- total import result rows
- oldest queued task timestamp
- oldest active task timestamp
This is intended as an operational view of queue depth and worker activity.
Failure handling¶
If worker dispatch or execution returns an unexpected ApiError:
- task status is set to
failed - a failure summary is stored
- a
failedevent is appended - payload is redacted
This ensures the queue does not leave orphaned active tasks on normal error paths.
Scaling model¶
Today’s design scales in two dimensions:
Vertical within one process¶
Increase:
HUBUUM_TASK_WORKERS
Horizontal across processes¶
Run more app instances pointing at the same database.
This works because claiming is DB-coordinated with SKIP LOCKED.
The design does not require an external queue for correctness.
Limitations and current assumptions¶
- there is no cancellation flow implemented yet
- API and worker duties use the same binary with
HUBUUM_RUNTIME_ROLE; they can run in separate processes or together - progress counters are currently updated at finalize time, not streamed continuously per item
- executors handle imports, exports, backups, computed-field reindexing, and remote calls; each kind persists its appropriate results or artifacts
- task retention is summary/results after payload redaction, not full payload retention
Tests¶
The task system now has coverage in three areas:
Queue mechanics¶
- concurrent-safe claim behavior for queued rows
Execution semantics¶
- strict execution rolls back on runtime failure
- best-effort execution preserves successful items
API behavior¶
- task creation, events, results, and redaction
- idempotency reuse
- collision policy behavior
- permission policy behavior
- non-owner access rejection for task/import reads
- admin queue meta endpoint
See:
- PostgreSQL task queue operations
- PostgreSQL task execution operations
- src/tasks/tests.rs
- tests/api_jobs_suite/imports.rs
- tests/api_platform_suite/meta.rs
Recommended mental model¶
The simplest accurate mental model is:
tasksis the queue and current status table- task-scoped
eventsrows are the audit/history table - result tables are typed per task kind
- task workers claim from Postgres directly
- each task kind supplies its own planner and executor
- strict imports make domain mutations atomically
- best-effort imports trade atomicity for progress
The generated project inventory lists all current task kinds. The runtime hardening guide describes claim fencing and recovery.
Execution limits¶
Every worker admits execution under its live lease and persists an absolute UTC deadline using the storage backend's clock. The first claim's start time anchors it. Queue wait is excluded, and recovery preserves the deadline and original start time. Renewing a lease only establishes ownership; it never extends the execution budget.
| Environment variable | Default seconds |
|---|---|
HUBUUM_TASK_IMPORT_EXECUTION_TIMEOUT_SECONDS |
3600 |
HUBUUM_TASK_EXPORT_EXECUTION_TIMEOUT_SECONDS |
900 |
HUBUUM_TASK_BACKUP_EXECUTION_TIMEOUT_SECONDS |
3600 |
HUBUUM_TASK_REINDEX_EXECUTION_TIMEOUT_SECONDS |
7200 |
HUBUUM_TASK_REMOTE_CALL_EXECUTION_TIMEOUT_SECONDS |
300 |
HUBUUM_TASK_SCHEMA_VALIDATION_EXECUTION_TIMEOUT_SECONDS |
7200 |
Values must be positive whole seconds, at most thirty days. Zero does not disable deadlines. The administrator configuration exposes the effective policy.
A separately scheduled monitor polls durable control every 250 ms and shares a read-only typed execution context with services and adapters. Local deadline checks use a monotonic clock. Lease renewal continues through executor cleanup and terminal acknowledgement. PostgreSQL cancels an in-flight query through its reserved lease pool, drains the operation, and rolls back before releasing the connection. The memory adapter stages strict imports and validates its snapshot again before publication, allowing cancellation and lease renewal while staging. Concurrent state changes invalidate a memory import snapshot rather than allowing it to overwrite another mutation.
Cancellation, success and recovery serialize their final decision under the task lock and live lease. A strict import's receipt fence also checks cancellation and the deadline at commit. Committed receipts survive worker death and preserve factual results. Recovery acknowledges a persisted stop before considering schema checkpoint resumption. A resumed queued task whose original deadline expires can be finalized by recovery without a new claim. A task that has never been claimed has no execution deadline yet.
Cancellation requests and terminal events use bounded reason categories. The first request retains its actor, timestamp and optional private explanation. Terminal cleanup removes the request payload and incomplete export/backup output. See the task API cancellation contract.
Upgrade coordination¶
Apply migration 20260914000001 before starting the new server. Drain and stop
old task workers before migrating; an old worker does not understand cancellation
intent or execution deadlines. Start all replacement workers with consistent
per-kind limits. Existing terminal tasks remain readable and their absent control
metadata means no historical cancellation/deadline evidence was recorded.
Schedule a quiet period for the migration's constraint validation and partial deadline index build. Lock waits are limited to five seconds and each statement to sixty seconds. A timeout rolls back the entire migration; retry during a quieter period with these limits in place.
Operator monitoring¶
Use the shared operator package for Grafana dashboards, Prometheus recording and alerting rules, SLO definitions and response runbooks. The same assets work with the optional single-host stack, independently managed Prometheus/Grafana installations, and Prometheus Operator. Pin the package to your server release and scrape every process directly with deployment labels.