5.1 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
fn start_worker(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.
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.