Skip to content

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:

  • import
  • export
  • backup
  • reindex
  • remote_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:

  • tasks
  • events rows with entity_type = 'task'
  • import_task_results
  • export_task_outputs
  • remote_call_results

The first two are generic framework storage.

The result/output tables are typed per task kind.

Relevant code:

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:

  • queued
  • validating
  • running
  • succeeded
  • failed
  • partially_succeeded
  • cleanup

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:

  • tasks and task-scoped events rows 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:

  • queued
  • validating
  • running
  • succeeded
  • failed
  • partially_succeeded
  • cancelled

Terminal statuses:

  • succeeded
  • failed
  • partially_succeeded
  • cancelled

For imports:

  • queued means the task row exists and is waiting to be claimed
  • validating means a worker has claimed it and is planning/validating the import
  • running means 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:

  1. The handler serializes the request payload.
  2. It computes a SHA-256 request hash.
  3. It reads and validates Idempotency-Key if present. Keys must contain between 1 and 255 bytes; invalid keys return 400 Bad Request before task persistence.
  4. It either reuses an existing task for that submitter/idempotency key or inserts a new one.
  5. It returns 202 Accepted with Location: /api/v1/tasks/{id}.
  6. It kicks the worker so the queue starts draining immediately.

Task creation itself is generic and implemented in:

When a task is created:

  • the tasks row is inserted with status = queued
  • the full request_payload is stored
  • counters are initialized to zero
  • a queued event is appended

Worker model

The worker implementation lives in:

There are two entry points:

  • ensure_task_worker_running
  • kick_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_WORKERS
  • HUBUUM_TASK_POLL_INTERVAL_MS
  • HUBUUM_TASK_LEASE_SECONDS
  • HUBUUM_TASK_HEARTBEAT_SECONDS
  • HUBUUM_TASK_RECOVERY_INTERVAL_SECONDS
  • HUBUUM_EXPORT_OUTPUT_CLEANUP_INTERVAL_SECONDS

The canonical env-var reference lives in:

Defaults:

  • HUBUUM_ACTIX_WORKERS: detected CPU count
  • HUBUUM_TASK_WORKERS: about half the detected CPU count, minimum 1
  • HUBUUM_TASK_POLL_INTERVAL_MS: 5000
  • HUBUUM_TASK_LEASE_SECONDS: 60
  • HUBUUM_TASK_HEARTBEAT_SECONDS: 20
  • HUBUUM_TASK_RECOVERY_INTERVAL_SECONDS: 30
  • HUBUUM_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 UPDATE
  • SKIP 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 = validating
  • started_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 executor
  • export -> export executor
  • backup -> full-system backup executor
  • reindex -> internal computed-field class rebuild executor
  • remote_call -> remote HTTP invocation executor
  • schema_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:

  1. load payload
  2. plan and validate
  3. execute
  4. 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:

  1. collections
  2. classes
  3. objects
  4. class relations
  5. object relations
  6. 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 items
  • failures: per-item planning failures
  • aborted: 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 = NULL
  • request_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.atomicity
  • mode.collision_policy
  • mode.permission_policy

Current behavior:

  • strict always aborts on the first planning failure
  • mode.permission_policy=abort stops on permission failures
  • mode.permission_policy=continue records permission failures and continues a best-effort import
  • collision_policy=abort records or aborts on collisions depending on atomicity
  • collision_policy=overwrite turns matching collection/class/object collisions into update operations

For relations:

  • overwrite-like behavior is effectively modeled as noop when 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 failed event 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:

The simplest accurate mental model is:

  • tasks is the queue and current status table
  • task-scoped events rows 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.