calternal-db
SQLite storage and durable job queue for calternal
SQLite storage for calternal’s Index and durable jobs.
The Index is rebuildable except for Security state. The server is its only
writer. FULL commits preserve acknowledged Security state (#824; DESIGN §2).
Independent PASSIVE maintenance keeps checkpoints off commits (#823).
Db keeps separate ordinary and authority write connections and a read pool.
SQLite permits one write transaction at a time; WAL lets readers continue
during writes. Plugin repositories own their domain traits and keep SQL in their
backend adapters; they must not expose SQLx rows or query types from those
traits. That boundary lets a later Postgres adapter implement the same
plugin-owned trait without adding Postgres here.
The snapshot API takes a trusted destination directory from the server.
This crate has no calternal-fs dependency in this base. The server must
resolve that directory through its directory-handle API before passing it.
Source: crates/calternal-db/src/lib.rs
Modules
Section titled “Modules”| Module | Summary |
|---|---|
cron |
Cron expressions that enqueue idempotent jobs. |
Re-exports
Section titled “Re-exports”pub use cron::CronJobSchedulepub use cron::CronSchedulerpub use db::WriteChangeCounts
Structs
Section titled “Structs”Backoff
Section titled “Backoff”pub struct BackoffExponential retry delay, capped to keep failed work from drifting forever.
Fields
pub initial: Duration: Delay after the first failed attempt.pub maximum: Duration: Maximum delay between retries.
Implements: Clone, Debug, Default
Backoff::delay
Section titled “Backoff::delay”pub fn delay(&self, attempts: u32) -> DurationReturn the delay for a one-based failed attempt count.
Source: crates/calternal-db/src/jobs.rs:55
CheckpointStatus
Section titled “CheckpointStatus”pub struct CheckpointStatusContent-free counters. Negative frame counts mean SQLite could not report a WAL.
Fields
pub runs: u64: Total attempted passes, including the final shutdown pass.pub errors: u64: SQL errors; a later pass retries without changing the writer policy.pub log_frames: i64: Frames present at the last successful pass.pub checkpointed_frames: i64: Frames copied back at the last successful pass.pub elapsed: Duration: Duration of the last attempt.pub busy: bool: The last attempt met a checkpoint lock, pinned reader, or active writer, or it copied back fewer frames than the WAL holds.
Implements: Clone, Copy, Debug, Default
Source: crates/calternal-db/src/checkpoint.rs:21
pub struct DbOne server-owned Index with ordinary and authority write pools and read-only readers.
Implements: Clone, Storage
Db::connect
Section titled “Db::connect”pub async fn connect(path: impl AsRef<Path>, options: DbOptions) -> Result<Self>Open a database file with NORMAL ordinary writes, FULL authority writes, and a busy timeout on both application write connections. Page-cache, statement-cache and reader-count limits stay bounded (#503).
Db::writer_pool
Section titled “Db::writer_pool”pub fn writer_pool(&self) -> &Pool<Sqlite>Return the NORMAL single-connection pool for migrations and ordinary writes.
Plugin repositories should keep this use inside their SQLite adapter. Keep the owning Db alive while these pools are in use so its background checkpoint worker remains active (#823).
Db::authority_pool
Section titled “Db::authority_pool”pub fn authority_pool(&self) -> &Pool<Sqlite>Commit Security state with FULL before acknowledging it (DESIGN §2, #824). Never use this pool for bulk Index or content work. Its immutable policy also applies to implicit transactions and replacement connections.
Db::reader_pool
Section titled “Db::reader_pool”pub fn reader_pool(&self) -> &Pool<Sqlite>Return the read-only pool used by repository reads.
Db::subscribe_job_changes
Section titled “Db::subscribe_job_changes”pub fn subscribe_job_changes(&self) -> broadcast::Receiver<()>Subscribe to queue changes. This is a hint channel; readers reload the durable queue state after every message and after a lagged receiver.
Db::checkpoint_status
Section titled “Db::checkpoint_status”pub fn checkpoint_status(&self) -> CheckpointStatusRead content-free checkpoint progress and failure counters (#823).
Db::close
Section titled “Db::close”pub async fn close(&self)Close application pools, then finish a best-effort maintenance pass. FULL commits do not depend on this pass for durability (#824).
Db::apply_migration_sets
Section titled “Db::apply_migration_sets”pub async fn apply_migration_sets(&self, sets: &[MigrationSet]) -> Result<()>Apply each set in caller-supplied order and each migration in ascending
version order. All owners share _migrations, keyed by namespace and
version, so a Plugin cannot overwrite another Plugin’s migration state.
Db::apply_migration_sets_with_snapshot
Section titled “Db::apply_migration_sets_with_snapshot”pub async fn apply_migration_sets_with_snapshot( &self, sets: &[MigrationSet], options: &crate::SnapshotOptions, ) -> Result<Option<crate::Snapshot>>Snapshot the current Index before applying pending migrations. Unchanged restarts skip the full database copy, and checksum drift fails before a snapshot is written (#549, #748, DESIGN §27).
Db::begin_mutation
Section titled “Db::begin_mutation”pub async fn begin_mutation( &self, owner: &str, id: &str, kind: &str, canonical_input: &[u8], now_ms: i64, ) -> Result<MutationStart<'_>, MutationError>Reserve an intent under the single writer; canonical input includes target, expected revision and every argument. A dropped New transaction rolls back.
Db::file_mutation_replay
Section titled “Db::file_mutation_replay”pub async fn file_mutation_replay( &self, owner: &str, id: &str, kind: &str, canonical_input: &[u8], ) -> Result<Option<FileMutation>, MutationError>Check a file replay before parsing inputs or reading attachment files. A completed create remains valid after later edits or file deletion. The adapter must check access and hold its domain lock (#844).
Db::reserve_file_mutation
Section titled “Db::reserve_file_mutation”pub async fn reserve_file_mutation( &self, owner: &str, id: &str, kind: &str, canonical_input: &[u8], plan: Value, now_ms: i64, ) -> Result<FileMutation, MutationError>Reserve a stable recovery plan using the shared canonical-input receipt. No filesystem I/O holds the Index writer transaction. After a crash the adapter repeats only the stored plan, using stable file item IDs (#844).
Db::complete_file_mutation
Section titled “Db::complete_file_mutation”pub async fn complete_file_mutation( &self, owner: &str, id: &str, kind: &str, result: Value, ) -> Result<(), MutationError>Seal a prepared file intent after every source file and required read projection is durable. A completed replay never recreates deleted data. Access checks and the domain lock belong to the adapter; this receipt carries no inverse and cannot authorize a file Undo (#844; DESIGN §2).
Db::mutation_receipt
Section titled “Db::mutation_receipt”pub async fn mutation_receipt( &self, owner: &str, id: &str, _now_ms: i64, ) -> Result<MutationReceipt, MutationError>Lookup is bounded and User-scoped. Adapters must recheck domain access.
Db::get_instance_secret
Section titled “Db::get_instance_secret”pub async fn get_instance_secret(&self, name: &str) -> Result<Option<Vec<u8>>>Read a named secret for trusted server code. Callers must never return the value to an API response or include it in a log message.
Db::has_instance_secret
Section titled “Db::has_instance_secret”pub async fn has_instance_secret(&self, name: &str) -> Result<bool>Report whether a named secret exists without reading its value.
Db::set_instance_secret
Section titled “Db::set_instance_secret”pub async fn set_instance_secret(&self, name: &str, value: &[u8]) -> Result<()>Replace a named secret in one SQLite transaction.
Db::delete_instance_secret
Section titled “Db::delete_instance_secret”pub async fn delete_instance_secret(&self, name: &str) -> Result<()>Remove a named secret. Missing secrets are already cleared.
Db::get_or_insert_instance_secret
Section titled “Db::get_or_insert_instance_secret”pub async fn get_or_insert_instance_secret( &self, name: &str, candidate: &[u8], ) -> Result<Vec<u8>>Return an existing named secret or persist and return candidate.
Secret names and candidate bytes come from trusted Rust code. Values stay in the private Index and this helper never logs or serializes them.
Db::snapshot
Section titled “Db::snapshot”pub async fn snapshot(&self, options: &SnapshotOptions) -> Result<Snapshot>Write a consistent snapshot with VACUUM INTO, then atomically install
it in the target directory and prune old snapshots.
Source: crates/calternal-db/src/db.rs:249
DbOptions
Section titled “DbOptions”pub struct DbOptionsConnection settings for one SQLite database.
Fields
pub max_readers: u32: Requested reader connections, clamped to the per-database maximum. Zero uses one reader; standalone WAL queries see the latest commit (#503/#512).pub busy_timeout: Duration: Maximum time to wait for SQLite’s file lock.
Implements: Clone, Debug, Default
Source: crates/calternal-db/src/db.rs:52
EncryptedIntegrationCredential
Section titled “EncryptedIntegrationCredential”pub struct EncryptedIntegrationCredentialEncrypted provider credentials stored on the shared account row (#407, DESIGN §49).
Fields
pub nonce: Vec<u8>pub ciphertext: Vec<u8>
Implements: Clone, fmt::Debug
Source: crates/calternal-db/src/integrations.rs:196
EnqueueOptions
Section titled “EnqueueOptions”pub struct EnqueueOptionsValues used to insert a job.
Fields
pub kind: String: Registered handler kind.pub payload: Value: JSON payload for the handler.pub dedup_key: Option<String>: Optional key that allows only one active job with this key.pub priority: i64: Higher values are leased first.pub run_at: DateTime<Utc>: Do not lease the job before this time.pub max_attempts: u32: Maximum number of handler attempts. Must be positive.pub owner_user_id: Option<String>: User who owns this job, when it is private work.
Implements: Clone, Debug
EnqueueOptions::new
Section titled “EnqueueOptions::new”pub fn new(kind: impl Into<String>, payload: Value) -> SelfBuild an immediately runnable job with one allowed attempt.
Source: crates/calternal-db/src/jobs.rs:195
EnqueueOutcome
Section titled “EnqueueOutcome”pub struct EnqueueOutcomeResult of enqueueing a job.
Fields
pub job: Job: Inserted job, or the existing active job with the same dedup key.pub inserted: bool: True if this call inserted the job.
Implements: Clone, Debug
Source: crates/calternal-db/src/jobs.rs:229
FailureOutcome
Section titled “FailureOutcome”pub struct FailureOutcomeResult of recording a handler failure.
Fields
pub state: JobState: State after recording the failure.pub run_at: Option<DateTime<Utc>>: Next eligible time, orNoneafter dead-lettering.
Implements: Clone, Debug
Source: crates/calternal-db/src/jobs.rs:303
IntegrationAccount
Section titled “IntegrationAccount”pub struct IntegrationAccountShared provider account metadata and its encrypted credential (#407, DESIGN §49).
Fields
pub id: Stringpub owner_id: Stringpub provider: IntegrationProviderpub email: Stringpub credential: EncryptedIntegrationCredentialpub mail_enabled: boolpub calendar_enabled: boolpub calendar_supported: boolpub carddav_endpoint: Option<String>pub carddav_supported: boolpub created_ms: i64
Implements: Clone, Debug
Source: crates/calternal-db/src/integrations.rs:307
IntegrationCredential
Section titled “IntegrationCredential”pub struct IntegrationCredentialCleartext fields in one shared envelope; Debug never reveals them (#407, DESIGN §49 I6).
Implements: Clone, Deserialize, Serialize, fmt::Debug
IntegrationCredential::new
Section titled “IntegrationCredential::new”pub fn new(username: String, app_password: String) -> SelfBuild one provider sign-in envelope for a shared account (#407, DESIGN §49).
IntegrationCredential::with_smtp
Section titled “IntegrationCredential::with_smtp”pub fn with_smtp( username: String, app_password: String, smtp_username: String, smtp_app_password: String, ) -> SelfPreserve Mail’s separate SMTP identity inside the shared envelope (#407, DESIGN §49 I6).
IntegrationCredential::with_calendar_credentials
Section titled “IntegrationCredential::with_calendar_credentials”pub fn with_calendar_credentials(mut self, username: String, app_password: String) -> SelfPreserve Calendar’s service-specific login in the shared envelope (#407, DESIGN §49 I6).
IntegrationCredential::with_calendar_oauth
Section titled “IntegrationCredential::with_calendar_oauth”pub fn with_calendar_oauth(mut self, access_token: String, refresh_token: String) -> SelfPreserve Calendar OAuth tokens inside the encrypted envelope (#407, DESIGN §49 I6).
IntegrationCredential::username
Section titled “IntegrationCredential::username”pub fn username(&self) -> &strBorrow the provider login after trusted service code decrypts it (#407, DESIGN §49 I6).
IntegrationCredential::app_password
Section titled “IntegrationCredential::app_password”pub fn app_password(&self) -> &strBorrow the provider password after trusted service code decrypts it (#407, DESIGN §49 I6).
IntegrationCredential::smtp_username
Section titled “IntegrationCredential::smtp_username”pub fn smtp_username(&self) -> &strReturn Mail’s separate SMTP login or the shared login (#407, DESIGN §49 I6).
IntegrationCredential::smtp_app_password
Section titled “IntegrationCredential::smtp_app_password”pub fn smtp_app_password(&self) -> &strReturn Mail’s separate SMTP password or the shared password (#407, DESIGN §49 I6).
IntegrationCredential::calendar_username
Section titled “IntegrationCredential::calendar_username”pub fn calendar_username(&self) -> &strReturn Calendar’s separate login or the shared login (#407, DESIGN §49 I6).
IntegrationCredential::calendar_app_password
Section titled “IntegrationCredential::calendar_app_password”pub fn calendar_app_password(&self) -> &strReturn Calendar’s separate password or the shared password (#407, DESIGN §49 I6).
IntegrationCredential::calendar_oauth_tokens
Section titled “IntegrationCredential::calendar_oauth_tokens”pub fn calendar_oauth_tokens(&self) -> Option<(&str, &str)>Return Calendar OAuth tokens only to trusted service code (#407, DESIGN §49 I6).
Source: crates/calternal-db/src/integrations.rs:64
IntegrationKey
Section titled “IntegrationKey”pub struct IntegrationKey(XChaCha20Poly1305);The Instance key used to encrypt Connected Account credentials (#407, DESIGN §49).
Implements: Clone, fmt::Debug
IntegrationKey::new
Section titled “IntegrationKey::new”pub fn new(key: [u8; 32]) -> Result<Self, IntegrationCredentialError>Reject an uninitialized zero key before it can encrypt account data (#407, DESIGN §49).
IntegrationKey::encrypt
Section titled “IntegrationKey::encrypt”pub fn encrypt( &self, owner_id: &str, account_id: &str, credential: &IntegrationCredential, ) -> Result<EncryptedIntegrationCredential, IntegrationCredentialError>Encrypt credentials with fresh randomness and AAD bound to their User and account ID (#407, DESIGN §49).
IntegrationKey::decrypt
Section titled “IntegrationKey::decrypt”pub fn decrypt( &self, owner_id: &str, account_id: &str, encrypted: &EncryptedIntegrationCredential, ) -> Result<IntegrationCredential, IntegrationCredentialError>Decrypt only the credential envelope bound to this User and account ID (#407, DESIGN §49).
Source: crates/calternal-db/src/integrations.rs:209
IntegrationServiceWrite
Section titled “IntegrationServiceWrite”pub struct IntegrationServiceWriteService switches stored on the central account row (#407, DESIGN §49).
Fields
pub mail_enabled: boolpub calendar_enabled: bool
Implements: Clone, Copy, Debug, Eq, PartialEq
Source: crates/calternal-db/src/integrations.rs:338
IntegrationWrite
Section titled “IntegrationWrite”pub struct IntegrationWrite<'a>Fields stored after provider discovery succeeds (#407, DESIGN §49).
Fields
pub owner_id: &'a strpub id: &'a strpub provider: IntegrationProviderpub email: &'a strpub credential: &'a EncryptedIntegrationCredentialpub mail_enabled: boolpub calendar_enabled: boolpub calendar_supported: boolpub carddav_endpoint: Option<&'a str>pub carddav_supported: boolpub created_ms: i64
Source: crates/calternal-db/src/integrations.rs:322
pub struct JobA row in the durable job queue.
Fields
pub id: String: Stable job identifier.pub kind: String: Registered handler kind.pub payload: Value: JSON value passed to the handler.pub dedup_key: Option<String>: Optional idempotency key shared by active jobs.pub priority: i64: Higher values are leased first.pub run_at: DateTime<Utc>: Earliest eligible time.pub attempts: u32: Number of leases, including expired leases.pub max_attempts: u32: Maximum number of leases before dead-lettering.pub lease_until: Option<DateTime<Utc>>: Current lease expiry, if a worker holds the job.pub leased_by: Option<String>: Current worker id, if a worker holds the job.pub last_error: Option<String>: Last handler failure text.pub state: JobState: Current queue state.pub created_at: DateTime<Utc>: Creation time.pub updated_at: DateTime<Utc>: Last state change time.pub owner_user_id: Option<String>: User who owns this job, if it is private work.pub progress_current: Option<i64>: Completed work units, when the handler reports progress.pub progress_total: Option<i64>: Total work units, when the handler knows the total.pub progress_message: Option<String>: Safe status text supplied by the handler.pub cancellation_requested: bool: A stop request is waiting for the handler’s next checkpoint.
Implements: Clone, Debug, Deserialize, Serialize
Source: crates/calternal-db/src/jobs.rs:124
JobControl
Section titled “JobControl”pub struct JobControlSafe worker operations that a long-running handler calls at checkpoints.
Implements: Clone
JobControl::checkpoint
Section titled “JobControl::checkpoint”pub async fn checkpoint(&self) -> Result<bool>Return true when an admin asked this handler to stop. Call before and after each bounded unit of work, before starting another write.
JobControl::report_progress
Section titled “JobControl::report_progress”pub async fn report_progress( &self, current: Option<i64>, total: Option<i64>, message: Option<&str>, ) -> Result<()>Save progress while this handler still owns its lease.
Source: crates/calternal-db/src/jobs.rs:337
JobKindControl
Section titled “JobKindControl”pub struct JobKindControlPersisted pause state for one registered task kind.
Fields
pub kind: String: Stable queue kind.pub paused: bool: Whether the Worker must keep pending jobs waiting.pub reason: Option<String>: Admin-supplied reason shown while the queue is paused.pub paused_at: Option<DateTime<Utc>>: Time when this pause began.
Implements: Clone, Debug, Deserialize, Eq, PartialEq, Serialize
Source: crates/calternal-db/src/jobs.rs:279
JobKindRegistration
Section titled “JobKindRegistration”pub struct JobKindRegistrationA background task registered with the shared SQLite Worker.
Fields
pub kind: String: Stable queue kind stored on every job row.pub label: String: Human-readable label for the instance Jobs view.pub description: String: One-line explanation shown in queue details.pub concurrency: usize: Maximum number of active handlers for this kind.
Implements: Clone, Debug, Deserialize, Eq, PartialEq, Serialize
JobKindRegistration::new
Section titled “JobKindRegistration::new”pub fn new(kind: impl Into<String>, label: impl Into<String>, concurrency: usize) -> SelfDescribe one task kind and its Worker concurrency limit.
JobKindRegistration::new_described
Section titled “JobKindRegistration::new_described”pub fn new_described( kind: impl Into<String>, label: impl Into<String>, description: impl Into<String>, concurrency: usize, ) -> SelfDescribe one task kind with the copy shown in the Admin Jobs view.
Source: crates/calternal-db/src/jobs.rs:238
JobQueue
Section titled “JobQueue”pub struct JobQueueQueue access with an injectable clock and retry policy.
Implements: Clone
JobQueue::new
Section titled “JobQueue::new”pub fn new(db: Db) -> SelfUse the system clock and default backoff.
JobQueue::with_clock
Section titled “JobQueue::with_clock”pub fn with_clock(db: Db, clock: Arc<dyn Clock>, backoff: Backoff) -> SelfSet a clock and retry policy. A fake Clock makes retry behavior
deterministic in tests and lets a server share one clock source.
JobQueue::register_kind
Section titled “JobQueue::register_kind”pub async fn register_kind(&self, registration: JobKindRegistration) -> Result<()>Register one task kind for the admin Jobs view. Worker startup calls this for every handler, so the queue catalogue and handler limits stay in step without a second registry.
JobQueue::registered_kinds
Section titled “JobQueue::registered_kinds”pub async fn registered_kinds(&self) -> Result<Vec<JobKindRegistration>>Return task kinds registered by the running Worker, sorted by kind.
JobQueue::pause_kind
Section titled “JobQueue::pause_kind”pub async fn pause_kind(&self, kind: &str, reason: &str) -> Result<JobKindControl>Pause a registered task kind. Repeated pauses keep the first reason and timestamp until an admin resumes the queue.
JobQueue::resume_kind
Section titled “JobQueue::resume_kind”pub async fn resume_kind(&self, kind: &str) -> Result<JobKindControl>Resume a registered task kind. Repeated resumes are safe.
JobQueue::job_kind_control
Section titled “JobQueue::job_kind_control”pub async fn job_kind_control(&self, kind: &str) -> Result<Option<JobKindControl>>Read the persisted pause state of a registered task kind.
JobQueue::list_admin_jobs
Section titled “JobQueue::list_admin_jobs”pub async fn list_admin_jobs(&self) -> Result<Vec<JobSummary>>Return active jobs and completed jobs from the last 24 hours for admins.
JobQueue::list_user_jobs
Section titled “JobQueue::list_user_jobs”pub async fn list_user_jobs(&self, user_id: &str) -> Result<Vec<JobSummary>>Return only a User’s active work and terminal jobs from the last day.
JobQueue::run_now
Section titled “JobQueue::run_now”pub async fn run_now(&self, kind: &str) -> Result<u64>Make waiting jobs of one kind eligible now. This does not add another queue entry, so repeated Run now requests remain idempotent.
JobQueue::retry_failed
Section titled “JobQueue::retry_failed”pub async fn retry_failed(&self, kind: &str) -> Result<u64>Retry every failed job of one kind from its first attempt.
JobQueue::clear_failed
Section titled “JobQueue::clear_failed”pub async fn clear_failed(&self, kind: &str) -> Result<u64>Remove failed jobs of one kind after the admin has reviewed them.
JobQueue::enqueue
Section titled “JobQueue::enqueue”pub async fn enqueue(&self, options: EnqueueOptions) -> Result<EnqueueOutcome>Add a job. Active jobs with the same non-empty dedup key return the existing row, including when that row is leased.
JobQueue::enqueue_in_transaction
Section titled “JobQueue::enqueue_in_transaction”pub async fn enqueue_in_transaction( &self, transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>, options: EnqueueOptions, ) -> Result<EnqueueOutcome>Add a job in the caller’s writer transaction (DESIGN §3, #486).
The caller commits its mutation and this wake together, then calls
wake_workers; workers can never observe an uncommitted intent.
JobQueue::wake_workers
Section titled “JobQueue::wake_workers”pub fn wake_workers(&self)Wake workers after committing a caller-owned enqueue transaction (#486). A periodic worker poll also recovers a crash between commit and wake.
JobQueue::get
Section titled “JobQueue::get”pub async fn get(&self, id: &str) -> Result<Option<Job>>Read a job by id.
JobQueue::active_by_dedup_keys
Section titled “JobQueue::active_by_dedup_keys”pub async fn active_by_dedup_keys(&self, keys: &[String]) -> Result<Vec<Job>>Read the active (pending or leased) jobs that hold any of these dedup keys. Handlers that alternate keys so a running job can queue its own successor use this to keep one logical chain per subject: a dedup key alone cannot see the other key of the pair (Mail #753). At most 64 keys are read in one call; an empty list returns no job.
JobQueue::get_summary
Section titled “JobQueue::get_summary”pub async fn get_summary(&self, id: &str) -> Result<Option<JobSummary>>Read one job without deserializing its handler payload.
JobQueue::lease_next
Section titled “JobQueue::lease_next”pub async fn lease_next(&self, worker_id: &str, lease_for: Duration) -> Result<Option<Job>>Lease the highest-priority ready job. worker_id must identify this
attempt uniquely; pass it unchanged to renewal and settlement (#1042).
JobQueue::lease_next_for_kinds
Section titled “JobQueue::lease_next_for_kinds”pub async fn lease_next_for_kinds( &self, worker_id: &str, kinds: &[String], lease_for: Duration, ) -> Result<Option<Job>>Lease the highest-priority ready job whose kind is in kinds.
worker_id must be a fresh attempt identity, as in lease_next (#1042).
An empty list returns no job. The update and selection are one SQLite
statement, so two workers cannot lease the same row.
JobQueue::heartbeat
Section titled “JobQueue::heartbeat”pub async fn heartbeat(&self, id: &str, worker_id: &str, lease_for: Duration) -> Result<bool>Extend an owned lease. Expiry alone does not revoke its token; return false only after ownership is removed or replaced (#1042).
JobQueue::complete
Section titled “JobQueue::complete”pub async fn complete(&self, id: &str, worker_id: &str) -> Result<()>Complete an owned attempt, even after expiry. Recovery must clear the token before ownership is lost (#1042).
JobQueue::fail
Section titled “JobQueue::fail”pub async fn fail(&self, id: &str, worker_id: &str, error: &str) -> Result<FailureOutcome>Record a failure. The job returns to pending with exponential backoff,
or moves to dead when its attempt limit is reached. Expiry alone does
not revoke the caller’s token; recovery must remove it first (#1042).
JobQueue::cancel
Section titled “JobQueue::cancel”pub async fn cancel(&self, id: &str) -> Result<CancelOutcome>Cancel a pending or leased job and release its dedup key.
JobQueue::cancel_for_user
Section titled “JobQueue::cancel_for_user”pub async fn cancel_for_user(&self, id: &str, user_id: &str) -> Result<CancelOutcome>Cancel a job only when it belongs to this User. Unknown and other
Users’ job IDs both return NotActive to avoid an ID oracle.
JobQueue::cancellation_requested
Section titled “JobQueue::cancellation_requested”pub async fn cancellation_requested(&self, id: &str, worker_id: &str) -> Result<bool>Return whether this Worker should stop at its next safe checkpoint. A removed owner token also tells the handler to stop; expiry alone does not (#1042).
JobQueue::update_progress
Section titled “JobQueue::update_progress”pub async fn update_progress( &self, id: &str, worker_id: &str, current: Option<i64>, total: Option<i64>, message: Option<&str>, ) -> Result<bool>Save progress only while the caller still owns the attempt token (#1042).
JobQueue::recover_expired_leases
Section titled “JobQueue::recover_expired_leases”pub async fn recover_expired_leases(&self) -> Result<u64>Return orphaned expired leases to pending, or dead-letter exhausted attempts. Live Worker tokens are excluded even when renewal waits for the single writer. Take the token snapshot after acquiring that writer; new claims register before they wait for it. Stale exclusions last at most until the next recovery pass, and task drop removes tokens (#1042).
Source: crates/calternal-db/src/jobs.rs:312
JobSummary
Section titled “JobSummary”pub struct JobSummaryQueue state safe for an admin or owner view. It excludes the handler payload, which may be large and is never needed to operate a queue.
Fields
pub id: String: Stable job identifier.pub kind: String: Registered handler kind.pub last_error: Option<String>: Last handler failure text.pub state: JobState: Current queue state.pub created_at: DateTime<Utc>: Creation time.pub updated_at: DateTime<Utc>: Last state change time.pub owner_user_id: Option<String>: User who owns this job, if it is private work.pub progress_current: Option<i64>: Completed work units, when the handler reports progress.pub progress_total: Option<i64>: Total work units, when the handler knows the total.pub progress_message: Option<String>: Safe status text supplied by the handler.pub cancellation_requested: bool: A stop request is waiting for the handler’s next checkpoint.
Implements: Clone, Debug
Source: crates/calternal-db/src/jobs.rs:168
Migration
Section titled “Migration”pub struct MigrationOne ordered SQL migration owned by a crate or Plugin.
Fields
pub version: i64: Positive version within its namespace.pub description: String: Short description stored for operators.pub sql: String: SQL statements for this version. Keep each statement separate with a semicolon; the migration runner executes this text as a script.
Implements: Clone, Debug
Migration::new
Section titled “Migration::new”pub fn new(version: i64, description: impl Into<String>, sql: impl Into<String>) -> SelfBuild one migration record.
Source: crates/calternal-db/src/migrations.rs:14
MigrationSet
Section titled “MigrationSet”pub struct MigrationSetA Plugin’s migration namespace and its migrations in ascending order.
Fields
pub namespace: String: Stable crate or Plugin name. Do not reuse a namespace for another owner.pub migrations: Vec<Migration>: Migrations for this namespace.
Implements: Clone, Debug
MigrationSet::new
Section titled “MigrationSet::new”pub fn new(namespace: impl Into<String>, migrations: Vec<Migration>) -> SelfBuild a migration set. The namespace is stored as supplied.
Source: crates/calternal-db/src/migrations.rs:37
MutationReceipt
Section titled “MutationReceipt”pub struct MutationReceiptA committed result. The inverse stays server-side and is never a client command.
Fields
pub operation_id: Stringpub kind: Stringpub result: Valuepub expires_ms: i64pub undone_by: Option<String>
Implements: Clone, Debug, Serialize, Deserialize, PartialEq
Source: crates/calternal-db/src/mutations.rs:46
MutationTransaction
Section titled “MutationTransaction”pub struct MutationTransaction<'a>A transaction cannot acknowledge success without storing the inverse (#667).
MutationTransaction::connection
Section titled “MutationTransaction::connection”pub fn connection(&mut self) -> &mut SqliteConnectionDomain adapters execute writes and queue inserts on this same connection.
MutationTransaction::take_inverse
Section titled “MutationTransaction::take_inverse”pub async fn take_inverse( &mut self, original_id: &str, kind: &str, ) -> Result<Value, MutationError>Load the inverse and mark its receipt in the Undo transaction. Adapter rejection rolls this mark back. A second distinct Undo cannot reuse it.
MutationTransaction::commit
Section titled “MutationTransaction::commit”pub async fn commit( mut self, result: Value, inverse: Value, ) -> Result<MutationReceipt, MutationError>Commit result and inverse atomically with the adapter write. A receipt exists only after this commit. Callers measure durable time after return.
Source: crates/calternal-db/src/mutations.rs:71
Snapshot
Section titled “Snapshot”pub struct SnapshotOne completed database snapshot.
Fields
pub path: PathBuf: Final snapshot path.pub created_at: DateTime<Utc>: Time used in the snapshot filename.
Implements: Clone, Debug
Source: crates/calternal-db/src/snapshot.rs:24
SnapshotOptions
Section titled “SnapshotOptions”pub struct SnapshotOptionsSnapshot output and retention settings.
Fields
pub directory: PathBuf: Existing target directory, resolved and trusted by the server.pub keep_last: usize: Number of newest snapshots to keep. Must be at least one.
Implements: Clone, Debug
Source: crates/calternal-db/src/snapshot.rs:15
SystemClock
Section titled “SystemClock”pub struct SystemClock;Wall clock used by production queue operations.
Implements: Debug, Default, Clock
Source: crates/calternal-db/src/jobs.rs:45
Worker
Section titled “Worker”pub struct WorkerA Tokio runner with per-kind limits and an optional shared handler limit.
Worker::new
Section titled “Worker::new”pub fn new( queue: JobQueue, worker_id: impl Into<String>, lease_for: Duration, poll_interval: Duration, ) -> Result<Self, WorkerError>Create a worker with a diagnostic name. Claims add a unique token so a restarted process can reuse the name without adopting old attempts (#1042).
Worker::with_max_concurrency
Section titled “Worker::with_max_concurrency”pub fn with_max_concurrency(mut self, maximum: usize) -> Result<Self, WorkerError>Bound all active handlers together, in addition to each kind’s limit. The default stays unbounded for existing Workers. The server sets a small shared budget so many due Plugin kinds cannot occupy the Index pools together (#549, DESIGN §3).
Worker::register
Section titled “Worker::register”pub fn register<F, Fut>( &mut self, kind: impl Into<String>, concurrency: usize, handler: F, ) -> Result<(), WorkerError> where F: Fn(Job) -> Fut + Send + Sync + 'static, Fut: Future<Output = std::result::Result<(), JobHandlerError>> + Send + 'static,Register a job kind and its maximum number of in-flight handlers.
Worker::register_named
Section titled “Worker::register_named”pub fn register_named<F, Fut>( &mut self, kind: impl Into<String>, label: impl Into<String>, concurrency: usize, handler: F, ) -> Result<(), WorkerError> where F: Fn(Job) -> Fut + Send + Sync + 'static, Fut: Future<Output = std::result::Result<(), JobHandlerError>> + Send + 'static,Register a named job kind and its maximum number of in-flight handlers.
Worker::register_when
Section titled “Worker::register_when”pub fn register_when<F, G, Fut>( &mut self, kind: impl Into<String>, concurrency: usize, enabled: G, handler: F, ) -> Result<(), WorkerError> where F: Fn(Job) -> Fut + Send + Sync + 'static, G: Fn() -> bool + Send + Sync + 'static, Fut: Future<Output = std::result::Result<(), JobHandlerError>> + Send + 'static,Register a job kind whose claim eligibility can change at runtime.
The worker checks enabled before offering the kind to SQLite, so a
disabled Plugin does not lease new work. Active jobs finish normally.
Worker::register_when_named
Section titled “Worker::register_when_named”pub fn register_when_named<F, G, Fut>( &mut self, kind: impl Into<String>, label: impl Into<String>, concurrency: usize, enabled: G, handler: F, ) -> Result<(), WorkerError> where F: Fn(Job) -> Fut + Send + Sync + 'static, G: Fn() -> bool + Send + Sync + 'static, Fut: Future<Output = std::result::Result<(), JobHandlerError>> + Send + 'static,Register a named job kind whose claim eligibility can change at runtime.
Worker::register_controlled
Section titled “Worker::register_controlled”pub fn register_controlled<F, Fut>( &mut self, kind: impl Into<String>, concurrency: usize, handler: F, ) -> Result<(), WorkerError> where F: Fn(Job, JobControl) -> Fut + Send + Sync + 'static, Fut: Future<Output = std::result::Result<(), JobHandlerError>> + Send + 'static,Register a handler that can observe stop requests and report progress.
Worker::register_when_controlled
Section titled “Worker::register_when_controlled”pub fn register_when_controlled<F, G, Fut>( &mut self, kind: impl Into<String>, concurrency: usize, enabled: G, handler: F, ) -> Result<(), WorkerError> where F: Fn(Job, JobControl) -> Fut + Send + Sync + 'static, G: Fn() -> bool + Send + Sync + 'static, Fut: Future<Output = std::result::Result<(), JobHandlerError>> + Send + 'static,Register a controlled handler whose claim eligibility can change at runtime.
Worker::register_when_controlled_named
Section titled “Worker::register_when_controlled_named”pub fn register_when_controlled_named<F, G, Fut>( &mut self, kind: impl Into<String>, label: impl Into<String>, concurrency: usize, enabled: G, handler: F, ) -> Result<(), WorkerError> where F: Fn(Job, JobControl) -> Fut + Send + Sync + 'static, G: Fn() -> bool + Send + Sync + 'static, Fut: Future<Output = std::result::Result<(), JobHandlerError>> + Send + 'static,Register a named controlled handler with a runtime claim eligibility check.
Worker::register_when_controlled_described
Section titled “Worker::register_when_controlled_described”pub fn register_when_controlled_described<F, G, Fut>( &mut self, kind: impl Into<String>, label: impl Into<String>, description: impl Into<String>, concurrency: usize, enabled: G, handler: F, ) -> Result<(), WorkerError> where F: Fn(Job, JobControl) -> Fut + Send + Sync + 'static, G: Fn() -> bool + Send + Sync + 'static, Fut: Future<Output = std::result::Result<(), JobHandlerError>> + Send + 'static,Register a named handler with copy for its one-line queue description. The Worker persists both strings beside the stable kind so Admin views do not need to duplicate a registry in the frontend.
Worker::run
Section titled “Worker::run”pub async fn run(&self, mut shutdown: watch::Receiver<bool>) -> Result<(), WorkerError>Recover expired leases at startup and during service (#1042). Shutdown stops new leases and waits for current handlers to finish.
Source: crates/calternal-db/src/worker.rs:68
CancelOutcome
Section titled “CancelOutcome”pub enum CancelOutcomeResult of cancelling an active job.
Variants
Cancelled: The job moved to cancelled.CancellationRequested: An active handler will stop at its next checkpoint.NotActive: The job was missing or already terminal.
Implements: Clone, Copy, Debug, Eq, PartialEq
Source: crates/calternal-db/src/jobs.rs:292
DbError
Section titled “DbError”pub enum DbErrorNo doc comment.
Variants
Sqlx(#[from] sqlx::Error)Io(#[from] std::io::Error)Json(#[from] serde_json::Error)InvalidConfig(String)InvalidMigration(String)MigrationChanged { namespace: String, version: i64 }LeaseLostInvalidCron(String)InvalidSnapshotPath
Implements: Debug, Error
Source: crates/calternal-db/src/error.rs:6
FileMutation
Section titled “FileMutation”pub enum FileMutationA file intent records its recovery plan before I/O and its outcome after I/O. Pending is never an acknowledgement. The adapter must hold its domain lock across reserve, idempotent file apply, and complete (#844; DESIGN §§2, 58).
Variants
Pending(Value)Complete(Value)
Implements: Clone, Debug, Serialize, Deserialize, PartialEq
Source: crates/calternal-db/src/mutations.rs:59
IntegrationCredentialError
Section titled “IntegrationCredentialError”pub enum IntegrationCredentialErrorSafe failures returned by shared account credential encryption (#407, DESIGN §49).
Variants
InvalidKeyRandomSourceSerializationAuthentication
Implements: Clone, Copy, Debug, Eq, PartialEq, fmt::Display, std::error::Error
Source: crates/calternal-db/src/integrations.rs:219
IntegrationProvider
Section titled “IntegrationProvider”pub enum IntegrationProviderA provider preset supported by Connected Accounts (#407, DESIGN §49).
Variants
IcloudFastmailGmailYahooOther
Implements: Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize
IntegrationProvider::as_str
Section titled “IntegrationProvider::as_str”pub fn as_str(self) -> &'static strReturn the stable database value used by the shared account row (#407, DESIGN §49).
Source: crates/calternal-db/src/integrations.rs:27
JobState
Section titled “JobState”pub enum JobStateState of a durable job.
Variants
Pending: Waiting for its run time or a retry delay.Leased: Owned by one attempt until settlement or recovery clears its token (#1042).Completed: Finished successfully.Dead: Reached its attempt limit and entered the dead-letter state.Cancelled: Cancelled before completion.
Implements: Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize
Source: crates/calternal-db/src/jobs.rs:86
MutationError
Section titled “MutationError”pub enum MutationErrorSafe outcome distinctions; transport errors do not imply a rejection.
Variants
InvalidConflictExpiredNotFoundDatabase(#[from] sqlx::Error)Json(#[from] serde_json::Error)
Implements: Debug, Error
Source: crates/calternal-db/src/mutations.rs:29
MutationStart
Section titled “MutationStart”pub enum MutationStart<'a>Replay releases the writer; New keeps it until domain write and receipt commit.
Variants
Replay(MutationReceipt)New(MutationTransaction<'a>)
Source: crates/calternal-db/src/mutations.rs:65
WorkerError
Section titled “WorkerError”pub enum WorkerErrorErrors returned by the Worker runner.
Variants
Database(#[from] DbError): Queue operation failed.InvalidConfig(String): Worker configuration is invalid.Task(String): A handler task panicked or was cancelled unexpectedly.
Implements: Debug, Error
Source: crates/calternal-db/src/worker.rs:55
Traits
Section titled “Traits”pub trait Clock: Send + SyncTime source used by queue state transitions.
Clock::now
Section titled “Clock::now”fn now(&self) -> DateTime<Utc>;Return the current UTC time.
Source: crates/calternal-db/src/jobs.rs:38
Storage
Section titled “Storage”pub trait Storage: Send + SyncPools available to a repository adapter.
A plugin should define a domain repository trait, for example
NotesRepository, with methods that accept and return note domain types.
Its SqliteNotesRepository may use this pool, but its trait must not expose
SQL strings, SQLx rows, or SQLite types. A future Postgres implementation
can implement that same plugin-owned trait with its own SQL and pool.
Storage::Database
Section titled “Storage::Database”type Database: Database;The SQLx database used by this adapter.
Storage::readers
Section titled “Storage::readers”fn readers(&self) -> &Pool<Self::Database>;Return the read pool for repository queries.
Storage::writer
Section titled “Storage::writer”fn writer(&self) -> &Pool<Self::Database>;Return the single-writer pool for mutations.
Source: crates/calternal-db/src/storage.rs:14
Type aliases
Section titled “Type aliases”JobHandlerError
Section titled “JobHandlerError”pub type JobHandlerError = String;Text returned by a handler when it cannot complete a job.
Source: crates/calternal-db/src/worker.rs:39
Result
Section titled “Result”pub type Result<T, E = DbError> = std::result::Result<T, E>;No doc comment.
Source: crates/calternal-db/src/error.rs:27
Functions
Section titled “Functions”bounded_sqlite_options
Section titled “bounded_sqlite_options”pub fn bounded_sqlite_options(options: SqliteConnectOptions) -> SqliteConnectOptionsApply shared memory limits to a SQLite connection before it enters a pool.
The page cache and prepared statement cache are per connection. At the reader limit, page caches use at most 64 MiB and statements stay bounded per connection (#503/#512).
Source: crates/calternal-db/src/db.rs:44
built_in_migrations
Section titled “built_in_migrations”pub fn built_in_migrations() -> MigrationSetThe database crate’s migration set. The server should include this set beside each enabled Plugin’s set when it starts the Index.
Source: crates/calternal-db/src/migrations.rs:56
create_integration_account
Section titled “create_integration_account”pub async fn create_integration_account( db: &Db, account: IntegrationWrite<'_>,) -> Result<(), sqlx::Error>Save one provider account and its credential on FULL before acknowledgement. The authority transaction reserves the write lock before reads (#407, #824; DESIGN §§2, 49).
Source: crates/calternal-db/src/integrations.rs:345
delete_integration_account
Section titled “delete_integration_account”pub async fn delete_integration_account( db: &Db, owner_id: &str, id: &str,) -> Result<bool, sqlx::Error>Delete only the owner’s shared row and cascade its Mail links. The server removes Plugin projections in the same FULL transaction (§§2, 49; #407, #824).
Source: crates/calternal-db/src/integrations.rs:463
get_integration_account
Section titled “get_integration_account”pub async fn get_integration_account( db: &Db, owner_id: &str, id: &str,) -> Result<Option<IntegrationAccount>, sqlx::Error>Find one shared account only within the supplied owner’s rows (#407, DESIGN §49).
Source: crates/calternal-db/src/integrations.rs:401
insert_integration_account
Section titled “insert_integration_account”pub async fn insert_integration_account( transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>, account: IntegrationWrite<'_>,) -> Result<(), sqlx::Error>Insert one shared Security row with its enabled service projections (#407, DESIGN §49).
Source: crates/calternal-db/src/integrations.rs:355
integration_account_from_row
Section titled “integration_account_from_row”pub fn integration_account_from_row(row: SqliteRow) -> Result<IntegrationAccount, sqlx::Error>Decode the shared account SELECT shape, rejecting unknown providers and invalid SQL types. The boot reconciler uses this same decoder inside its writer transaction instead of reading from another snapshot (#407, DESIGN §49 I6).
Source: crates/calternal-db/src/integrations.rs:422
is_sqlite_busy_locked
Section titled “is_sqlite_busy_locked”pub fn is_sqlite_busy_locked(error: &Error) -> boolReturn true for SQLite BUSY and LOCKED, including extended result codes.
Source: crates/calternal-db/src/sqlite.rs:46
is_sqlite_transient
Section titled “is_sqlite_transient”pub fn is_sqlite_transient(error: &Error) -> boolReturn true for SQLite lock errors and SQLite pool checkout failures.
These errors describe temporary storage pressure and should become a retryable service response. Pool checkout failures are not retried here: the caller must release its work and let the client retry.
Source: crates/calternal-db/src/sqlite.rs:59
list_integration_accounts
Section titled “list_integration_accounts”pub async fn list_integration_accounts( db: &Db, owner_id: &str,) -> Result<Vec<IntegrationAccount>, sqlx::Error>List one User’s shared accounts in stable creation order, with credentials still encrypted (#407).
Source: crates/calternal-db/src/integrations.rs:384
retry_db_busy
Section titled “retry_db_busy”pub async fn retry_db_busy<T, F, Fut>(operation: F) -> Result<T, crate::DbError>where F: FnMut() -> Fut, Fut: Future<Output = Result<T, crate::DbError>>,Retry a database operation when DbError wraps SQLite BUSY or LOCKED.
Source: crates/calternal-db/src/sqlite.rs:101
retry_sqlite_busy
Section titled “retry_sqlite_busy”pub async fn retry_sqlite_busy<T, F, Fut>(operation: F) -> Result<T, Error>where F: FnMut() -> Fut, Fut: Future<Output = Result<T, Error>>,Retry one complete SQLite write operation after transient lock errors.
The closure must start a fresh transaction on each call and perform every query in that transaction before it returns. This prevents a retry from continuing with a transaction that SQLite has already aborted. A single autocommit statement is also a complete operation and can use this helper.
Source: crates/calternal-db/src/sqlite.rs:92
retry_when
Section titled “retry_when”pub async fn retry_when<T, E, F, Fut, IsRetryable>( mut operation: F, is_retryable: IsRetryable,) -> Result<T, E>where F: FnMut() -> Fut, Fut: Future<Output = Result<T, E>>, IsRetryable: Fn(&E) -> bool,Retry a complete operation with the shared bounded backoff policy (#1021, #1022; DESIGN §2).
Each attempt must begin a fresh transaction. The predicate must select only transient failures, and the operation must be safe to repeat after rollback.
Source: crates/calternal-db/src/sqlite.rs:117
set_integration_services
Section titled “set_integration_services”pub async fn set_integration_services( db: &Db, owner_id: &str, id: &str, services: IntegrationServiceWrite,) -> Result<bool, sqlx::Error>Change service access on FULL; callers update projections atomically (#407, #824).
Source: crates/calternal-db/src/integrations.rs:442
sqlite_transient_kind
Section titled “sqlite_transient_kind”pub fn sqlite_transient_kind(error: &Error) -> Option<&'static str>Name the transient class of an SQLite error for a diagnostic log line.
The value is a fixed label, never the driver message, so a log can say
why a request answered 503 without exposing a path or a bound value
(#960). Returns None when the error is not transient.
Source: crates/calternal-db/src/sqlite.rs:68
validate_operation_id
Section titled “validate_operation_id”pub fn validate_operation_id(id: &str) -> Result<(), MutationError>Restrict IDs before they enter a URL or the Index. UUIDs are accepted.
Source: crates/calternal-db/src/mutations.rs:81
with_debug_slow_statement_logging
Section titled “with_debug_slow_statement_logging”pub fn with_debug_slow_statement_logging(options: SqliteConnectOptions) -> SqliteConnectOptionsEnable SQLx warning logs for slow statements when the #549 debug threshold is set.
The opt-in CALTERNAL_SQLITE_SLOW_STATEMENT_MS value changes only SQLx’s
logging threshold; bound values are not logged. The default connection
logging remains unchanged when the variable is absent or invalid.
Source: crates/calternal-db/src/sqlite.rs:19
Constants
Section titled “Constants”MAX_SQLITE_READER_CONNECTIONS
Section titled “MAX_SQLITE_READER_CONNECTIONS”pub const MAX_SQLITE_READER_CONNECTIONS: u32Maximum number of reader connections in one Instance database pool (#503/#512).
Source: crates/calternal-db/src/db.rs:33
SQLITE_CACHE_SIZE_KIB
Section titled “SQLITE_CACHE_SIZE_KIB”pub const SQLITE_CACHE_SIZE_KIB: usizeLimit each SQLite connection’s page cache to 1 MiB during large Home work (#503).
Source: crates/calternal-db/src/db.rs:35
SQLITE_STATEMENT_CACHE_CAPACITY
Section titled “SQLITE_STATEMENT_CACHE_CAPACITY”pub const SQLITE_STATEMENT_CACHE_CAPACITY: usizeBound prepared statements retained by each SQLite connection (#503).
Source: crates/calternal-db/src/db.rs:37
UNDO_RETENTION_MS
Section titled “UNDO_RETENTION_MS”pub const UNDO_RETENTION_MS: i64Undo may apply an inverse for one day; stored replay history is retained. The UI’s eight-second toast is independent of this server-side deadline.