Files
2026-09-05 17:39:50 +02:00

5.9 KiB

Utilisation de ksp-worker-api

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

Construire une identité Worker

fn worker_identity() -> ksp_worker_api::Result<(ksp_worker_api::WorkerId, ksp_worker_api::WorkerKindCode)> {
    let id = match ksp_worker_api::WorkerId::new("raw-ingest-mainnet-0001") {
        std::result::Result::Ok(value) => value,
        std::result::Result::Err(error) => return std::result::Result::Err(error),
    };
    let kind = match ksp_worker_api::WorkerKindCode::new("raw_transaction_ingest") {
        std::result::Result::Ok(value) => value,
        std::result::Result::Err(error) => return std::result::Result::Err(error),
    };
    return std::result::Result::Ok((id, kind));
}

WorkerId identifie une instance logique observée par un lifecycle/source donné. WorkerKindCode identifie une famille de Workers. Les deux sont bornés et utilisent un alphabet sûr. Le Debug de WorkerId masque sa valeur.

Piloter un lifecycle passif

Le lifecycle ne démarre aucun runtime. L'exemple suivant représente uniquement les transitions publiées par un producer concret lorsqu'il entre en exécution :

fn running_lifecycle(id: ksp_worker_api::WorkerId, kind: ksp_worker_api::WorkerKindCode) -> ksp_worker_api::Result<ksp_worker_api::WorkerLifecycle> {
    let mut lifecycle = ksp_worker_api::WorkerLifecycle::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.mark_running() {
        return std::result::Result::Err(error);
    }
    return std::result::Result::Ok(lifecycle);
}

Le producer possède l'autorité de transition. Il ne force jamais un état directement. Une transition invalide retourne une erreur stable et conserve l'état courant.

Pour un shutdown coopératif après observation du token :

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

Pour un fault terminal :

const IO_FAULT: ksp_worker_api::ErrorCode = ksp_worker_api::ErrorCode::new("example_worker", "io_fault");
if let std::result::Result::Err(error) = lifecycle.fault(IO_FAULT) {
    return std::result::Result::Err(error);
}

Stopped et Faulted(ErrorCode) sont terminaux. Un ancien lifecycle terminal ne représente jamais une nouvelle exécution.

Distinguer lifecycle, health et activity

WorkerState décrit la phase du service. WorkerHealth décrit sa qualité opérationnelle. WorkerActivity indique seulement Unknown, Idle ou Active.

let snapshot = ksp_worker_api::WorkerSnapshot::new(
    lifecycle.id().clone(),
    lifecycle.kind().clone(),
    ksp_worker_api::WorkerSnapshotSequence::initial(),
    lifecycle.state(),
    ksp_worker_api::WorkerHealth::Healthy,
    ksp_worker_api::WorkerActivity::Idle,
);

Le snapshot commun ne porte ni pourcentage, ni total, ni backlog, ni slot, ni transaction, ni métrique métier. Une API Worker concrète peut exposer séparément ses propres métriques.

Partager une intention de stop

let token = ksp_worker_api::WorkerStopToken::new();
let listener = token.clone();

assert!(!listener.is_stop_requested());
assert!(token.request_stop());
assert!(listener.is_stop_requested());
assert!(!token.request_stop());

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

Le token n'est pas une primitive de kill et ne garantit aucun délai de shutdown. Timeout, drain, join et retry appartiennent au runtime/caller.

Démarrer et arrêter un Worker concret

ksp-worker-api n'expose volontairement aucune commande runtime universelle. Une crate concrète peut fournir une surface start/stop adaptée à son domaine, mais elle utilise les contrats communs pour publier son identité, ses transitions, son état courant et l'intention de stop.

Un Worker concret ne doit donc pas transformer WorkerLifecycle en handle d'exécution ni ajouter des paramètres métier au contrat générique. Les paramètres/configurations propres à une famille de Workers restent dans cette famille ou dans sa couche de composition.

Observer un snapshot latest-value

Un consumer portable peut travailler directement avec le trait object-safe :

async fn observe(source: &dyn ksp_worker_api::WorkerSnapshotSource) {
    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 disponible après coalescence éventuelle. Un consumer ne doit pas supposer qu'il recevra chaque mise à jour intermédiaire.

Un listener tardif commence par current(). Tant que la source existe, son snapshot terminal courant reste lisible.

Faire avancer une séquence

let first = ksp_worker_api::WorkerSnapshotSequence::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-worker-api :

runtime Tokio/Futures concret
Job lifecycle ou checkpoint/backfill
Transport, Store, Config ou Logging
DTO Solana/provider
restart/retry/scheduler/process manager
payload métier dans WorkerSnapshot

Ces responsabilités appartiennent aux Workers concrets et aux couches de composition/contrôle supérieures.