pub struct SeppClient { /* private fields */ }Expand description
A handle to a Sepp server.
Cloning is cheap — clones share the same underlying connection and retry
policy — so clone freely to use the client across tasks. Every RPC method
takes &self.
Implementations§
Source§impl SeppClient
impl SeppClient
Sourcepub async fn connect(addr: impl Into<String>) -> Result<Self, ClientError>
pub async fn connect(addr: impl Into<String>) -> Result<Self, ClientError>
Connects to a Sepp server over plaintext with no authentication.
addr is a URI such as http://127.0.0.1:50051. For API-key auth, TLS,
or a custom RetryPolicy, use builder instead.
Sourcepub fn builder(addr: impl Into<String>) -> SeppClientBuilder
pub fn builder(addr: impl Into<String>) -> SeppClientBuilder
Starts building a client for addr, allowing authentication, TLS, and
retry configuration before connect.
Sourcepub fn from_channel(channel: Channel) -> Self
pub fn from_channel(channel: Channel) -> Self
Wraps an already-established tonic Channel, with no authentication
and the default RetryPolicy.
Use this to share a channel or apply custom tonic transport configuration the builder does not expose.
Sourcepub async fn enqueue_batch(
&self,
jobs: impl IntoIterator<Item = EnqueueRequest>,
) -> Result<Vec<Result<EnqueueAck, JobRejection>>, ClientError>
pub async fn enqueue_batch( &self, jobs: impl IntoIterator<Item = EnqueueRequest>, ) -> Result<Vec<Result<EnqueueAck, JobRejection>>, ClientError>
Enqueues a batch of jobs on a best-effort basis.
Each job is accepted or rejected independently: the returned vector has
one entry per submitted job, in the same order, where the inner Result
is Ok for an accepted job or Err for a per-job JobRejection. The
outer Err is reserved for whole-call failures (empty batch, transport
error, protocol violation). Transient failures are retried per the
client’s RetryPolicy; note that retried enqueues can duplicate jobs
that carry no idempotency key, so when any job in the batch lacks one,
the ambiguous-commit codes DeadlineExceeded and Aborted are not
retried (see RetryPolicy).
For all-or-nothing semantics, use enqueue_atomic.
Sourcepub async fn enqueue(
&self,
job: EnqueueRequest,
) -> Result<EnqueueAck, EnqueueError>
pub async fn enqueue( &self, job: EnqueueRequest, ) -> Result<EnqueueAck, EnqueueError>
Enqueues a single job.
A convenience wrapper over enqueue_batch that
flattens the result: a per-job rejection becomes
EnqueueError::Rejected. Its retry behavior — including that a
retried enqueue can duplicate a job that carries no idempotency key —
is inherited from enqueue_batch.
Sourcepub async fn enqueue_atomic(
&self,
jobs: impl IntoIterator<Item = EnqueueRequest>,
) -> Result<Vec<EnqueueAck>, AtomicEnqueueError>
pub async fn enqueue_atomic( &self, jobs: impl IntoIterator<Item = EnqueueRequest>, ) -> Result<Vec<EnqueueAck>, AtomicEnqueueError>
Enqueues a batch of jobs atomically: either all are accepted or none are.
On success, returns one EnqueueAck per job, in order. If any job
fails validation, nothing is enqueued and every failure is returned
together as AtomicEnqueueError::Validation.
Use this when the jobs are coordinated steps and a partial enqueue would
leave the system inconsistent.
Transient failures are retried per the client’s RetryPolicy; note
that retried enqueues can duplicate jobs that carry no idempotency key,
so when any job in the batch lacks one, the ambiguous-commit codes
DeadlineExceeded and Aborted are not retried (see RetryPolicy).
Sourcepub async fn reserve(
&self,
opts: &ReserveOptions,
) -> Result<Option<Vec<Job>>, ReserveError>
pub async fn reserve( &self, opts: &ReserveOptions, ) -> Result<Option<Vec<Job>>, ReserveError>
Long-polls for jobs to process.
Blocks up to the options’ wait_timeout
for at least one job. Returns Ok(Some(jobs)) with one or more leased
Jobs, or Ok(None) if the wait elapsed with nothing available (poll
again). Each returned job must be acked,
nacked, or extended before its lease
expires.
Unlike the other RPCs, reserve is not retried by the
RetryPolicy: as a long poll, an empty return is the normal idle
outcome and the caller loops anyway. A malformed job in the response is
logged and skipped rather than failing the whole batch.
Sourcepub async fn ack(&self, ctx: &JobCtx) -> Result<(), LeaseError>
pub async fn ack(&self, ctx: &JobCtx) -> Result<(), LeaseError>
Acknowledges that a job completed successfully, removing it from the queue.
The attempt carried by ctx guards against acking a job whose lease
was already reassigned — that surfaces as
LeaseError::AttemptMismatch or LeaseError::JobNotFound.
Sourcepub async fn nack(
&self,
ctx: &JobCtx,
retry: RetryDirective,
reason: impl Into<String>,
) -> Result<bool, LeaseError>
pub async fn nack( &self, ctx: &JobCtx, retry: RetryDirective, reason: impl Into<String>, ) -> Result<bool, LeaseError>
Negatively acknowledges a job, signalling that processing failed.
retry selects what the server does next (see RetryDirective) and
reason is recorded for debugging and metrics; an empty reason is
omitted from the request rather than sent as an empty string. Returns
true if this nack moved the job to the dead-letter queue (because
DeadLetter was requested or max_attempts was reached), false if
it will be retried.
Sourcepub async fn extend(
&self,
ctx: &JobCtx,
extension: Duration,
) -> Result<SystemTime, LeaseError>
pub async fn extend( &self, ctx: &JobCtx, extension: Duration, ) -> Result<SystemTime, LeaseError>
Extends a job’s lease by extension, measured from now, returning the
new expiry.
Call this when a handler needs longer than the original lease. Equivalent
to JobCtx::extend; a Worker with
with_auto_extend does it
automatically.
Sourcepub async fn get_server_info(&self) -> Result<ServerInfo, ClientError>
pub async fn get_server_info(&self) -> Result<ServerInfo, ClientError>
Fetches the server’s ServerInfo: version, capabilities, and limits.
Useful once at startup so a producer can validate jobs locally against the advertised limits and avoid round-trips that would only be rejected.
Sourcepub async fn drain_dead_letters(
&self,
queue: Option<&str>,
max: u32,
) -> Result<Vec<DeadLetterRecord>, ClientError>
pub async fn drain_dead_letters( &self, queue: Option<&str>, max: u32, ) -> Result<Vec<DeadLetterRecord>, ClientError>
Drains dead-lettered jobs for inspection and manual replay.
Returns up to max DeadLetterRecords (oldest-first, optionally
filtered to one queue) and removes them from the server; a max
of 0 returns an empty vector without making an RPC. This is
destructive: the records are gone once returned, so a dropped response
loses exactly that batch — for that reason it is not retried by the
RetryPolicy. Inspect each record, then replay any you want with
DeadLetterRecord::to_enqueue_request.
An empty result means nothing matched, which is indistinguishable from
dead-letter retention being disabled — check
ServerInfo::dead_letter_retention_enabled.
Trait Implementations§
Source§impl Clone for SeppClient
impl Clone for SeppClient
Source§fn clone(&self) -> SeppClient
fn clone(&self) -> SeppClient
1.0.0 (const: unstable) · Source§fn clone_from(&mut self, source: &Self)
fn clone_from(&mut self, source: &Self)
source. Read moreAuto Trait Implementations§
impl !Freeze for SeppClient
impl !RefUnwindSafe for SeppClient
impl !UnwindSafe for SeppClient
impl Send for SeppClient
impl Sync for SeppClient
impl Unpin for SeppClient
impl UnsafeUnpin for SeppClient
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Source§impl<T> CloneToUninit for Twhere
T: Clone,
impl<T> CloneToUninit for Twhere
T: Clone,
Source§impl<T> FutureExt for T
impl<T> FutureExt for T
Source§fn with_context(self, otel_cx: Context) -> WithContext<Self>
fn with_context(self, otel_cx: Context) -> WithContext<Self>
Source§fn with_current_context(self) -> WithContext<Self>
fn with_current_context(self) -> WithContext<Self>
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
Source§fn in_current_span(self) -> Instrumented<Self>
fn in_current_span(self) -> Instrumented<Self>
Source§impl<T> IntoRequest<T> for T
impl<T> IntoRequest<T> for T
Source§fn into_request(self) -> Request<T>
fn into_request(self) -> Request<T>
T in a tonic::Request