Files
2026-09-01 17:42:24 +02:00

4.8 KiB

Utilisation de ksp-job-api

Cette page décrit la façade publique durable de ksp-job-api. Les consumers utilisent uniquement les exports du crate-root.

Construire une identité de job

fn job_identity() -> ksp_job_api::Result<(ksp_job_api::JobId, ksp_job_api::JobKindCode)> {
    let id = match ksp_job_api::JobId::new("raw-backfill-mainnet-0001") {
        std::result::Result::Ok(value) => value,
        std::result::Result::Err(error) => return std::result::Result::Err(error),
    };
    let kind = match ksp_job_api::JobKindCode::new("raw_transaction_backfill") {
        std::result::Result::Ok(value) => value,
        std::result::Result::Err(error) => return std::result::Result::Err(error),
    };
    return std::result::Result::Ok((id, kind));
}

JobId identifie un job logique et ses reprises contrôlées. JobKindCode identifie une famille de jobs. Les deux sont bornés et utilisent un alphabet sûr ; leur Debug n'est pas une surface destinée à transporter un payload métier.

Piloter un lifecycle passif

fn lifecycle(id: ksp_job_api::JobId, kind: ksp_job_api::JobKindCode) -> ksp_job_api::Result<ksp_job_api::JobLifecycle> {
    let mut lifecycle = ksp_job_api::JobLifecycle::new(id, kind);
    if let std::result::Result::Err(error) = lifecycle.start() {
        return std::result::Result::Err(error);
    }
    if let std::result::Result::Err(error) = lifecycle.complete(ksp_job_api::JobCompletion::Complete) {
        return std::result::Result::Err(error);
    }
    return std::result::Result::Ok(lifecycle);
}

Le producer ne doit pas forcer un état directement. Il utilise les opérations de JobLifecycle, qui refusent les transitions hors matrice.

Pour une annulation observée pendant l'exécution :

if let std::result::Result::Err(error) = lifecycle.mark_cancelling() {
    return std::result::Result::Err(error);
}
if let std::result::Result::Err(error) = lifecycle.mark_cancelled() {
    return std::result::Result::Err(error);
}

Partager une intention d'annulation

let token = ksp_job_api::JobCancellationToken::new();
let listener = token.clone();

assert!(!listener.is_cancellation_requested());
assert!(token.request_cancellation());
assert!(listener.is_cancellation_requested());
assert!(!token.request_cancellation());

Le premier appel qui change l'état retourne true. Les demandes suivantes sont idempotentes et retournent false.

Le token ne doit pas être interprété comme une primitive de kill : le runtime concret décide où l'annulation peut interrompre l'admission ou une attente et quelles opérations déjà engagées doivent être drainées.

Publier une valeur latest-value

Un producer concret peut construire une notification complète :

#[derive(Clone)]
struct Snapshot {
    completed: usize,
}

let id = ksp_job_api::JobId::new("job-0001")?;
let kind = ksp_job_api::JobKindCode::new("example")?;
let sequence = ksp_job_api::JobNotificationSequence::initial();
let notification = ksp_job_api::JobNotification::new(
    id,
    kind,
    sequence,
    ksp_job_api::JobState::Running,
    Snapshot { completed: 0 },
);

assert_eq!(notification.sequence().value(), 0);
assert_eq!(notification.state(), ksp_job_api::JobState::Running);

Le Debug de JobNotification<S> masque volontairement le snapshot. Le type concret S doit lui-même rester sûr à exposer lorsque le consumer accède explicitement à snapshot().

Implémenter une source de snapshots

Une implémentation concrète possède son runtime et expose seulement le contrat JobSnapshotSource :

async fn observe<S>(source: &S)
where
    S: ksp_job_api::JobSnapshotSource,
{
    let current = source.current();
    let observed = current.sequence();
    let newer = source.wait_for_change(observed).await;
    assert!(newer.sequence().is_after(observed));
}

wait_for_change retourne la dernière valeur complète connue après coalescence éventuelle. Un consumer ne doit donc pas supposer qu'il recevra chaque mise à jour intermédiaire.

Faire avancer une séquence

let first = ksp_job_api::JobNotificationSequence::initial();
let second = match first.next() {
    std::result::Result::Ok(value) => value,
    std::result::Result::Err(error) => return std::result::Result::Err(error),
};
assert!(second.is_after(first));

L'épuisement de u64 est une erreur explicite ; la séquence ne wrappe jamais silencieusement.

Frontières à respecter

Ne pas ajouter à ksp-job-api :

Tokio/Futures runtime concret
Transport ou Store
Config/Logging
DTO métier d'un job précis
scheduler, worker ou control plane
persistence de checkpoint

Ces responsabilités appartiennent aux crates concrètes et à la composition supérieure.