Skip to main content

KevyCommands

Struct KevyCommands 

Source
pub struct KevyCommands;
Expand description

kevy’s command set, plugged into the kevy-rt runtime. Stateless — the keyspace lives in each shard’s Store, so this is a zero-sized clone target.

Trait Implementations§

Source§

impl Clone for KevyCommands

Source§

fn clone(&self) -> KevyCommands

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more
Source§

impl Commands for KevyCommands

Source§

fn resolve_block_argv<A: ArgvView + ?Sized>( &self, store: &mut Store, args: &A, kind: BlockKind, ) -> Argv

Freeze $ IDs in an XREAD BLOCK argv at park time. Default would leave $ literal in the parked argv; the wake retry would then re-resolve $ to the post-XADD last_id, miss the new entry, and time out instead of returning it. Other block kinds (BLPOP / BRPOP / XREADGROUP >) have no state-dependent argv — they fall through to the trait default.

Source§

fn resolve<A: ArgvView + ?Sized>(&self, args: &A) -> ResolvedCmd

One-pass verb resolution — the reactor calls this once per cmd and reads back txn_kind / route / is_quit / is_write without re-scanning the verb. This is kevy-rt’s primary hot-path optimization: every match arm uses the same upper buffer. Body in cmd_resolve.

Source§

fn route<A: ArgvView + ?Sized>(&self, args: &A) -> Route

Classify how a command is routed across shards.
Source§

fn dispatch<A: ArgvView + ?Sized>(&self, store: &mut Store, args: &A) -> Vec<u8>

Execute a full command against one shard’s store, returning RESP bytes.
Source§

fn dispatch_into<A: ArgvView + ?Sized>( &self, store: &mut Store, args: &A, out: &mut Vec<u8>, )

Execute a command, appending the RESP reply to out. The in-order local fast path uses this to write straight into the connection’s output buffer (no per-command reply Vec). Default: delegate to dispatch.
Source§

fn dispatch_resp3<A: ArgvView + ?Sized>( &self, store: &mut Store, args: &A, ) -> Vec<u8>

RESP3 variant of Self::dispatch — called when the connection has negotiated HELLO 3. Default: delegate to the RESP2 path (the cross-shard forward carries a per-cmd RespVersion so a V2 client and a V3 client can share the owning shard).
Source§

fn dispatch_into_resp3<A: ArgvView + ?Sized>( &self, store: &mut Store, args: &A, out: &mut Vec<u8>, )

RESP3 variant of Self::dispatch_into — called when the connection has negotiated HELLO 3. Default: delegate to the RESP2 path (so a server that hasn’t migrated any replies still works correctly with a RESP3 client, per spec). Override per command to emit RESP3 shapes (Map / Set / Double / …).
Source§

fn is_quit<A: ArgvView + ?Sized>(&self, args: &A) -> bool

Whether this command should close the connection (QUIT).
Source§

fn on_shard_init(&self, store: &mut Store)

Called once per shard, immediately after Store::new, before the reactor enters its event loop. Implementations install per-shard configuration that the runtime doesn’t know about — currently the maxmemory + eviction-policy pair, which kevy ships via its own process-wide config snapshot. Default: no-op so non-kevy embedders aren’t forced to override.
Source§

fn on_shard_start(&self, shard: usize)

Called once on the shard’s own thread, first thing in the reactor entry (both reactors), before restore/replay. Implementations that need per-shard identity at dispatch time (e.g. kevy’s CLUSTER MYID / CLUSTER NODES myself flag) stash shard in a thread-local here — in a thread-per-core runtime the current thread is the shard. Default: no-op.
Source§

fn on_persist_stats(&self, in_flight: bool, aof_rewrites_total: u64)

Per-tick persistence-stats publication: whether this shard has a background save/rewrite in flight and how many AOF rewrites have completed since open. Command layers that serve INFO persistence stash these in a thread-local (thread-per-core: the answering thread is the shard, same pattern as Self::on_shard_start). Default: no-op.
Source§

fn on_replication_view( &self, master_repl_offset: u64, replicas: Vec<(Ipv4Addr, u16, u64)>, )

Per-tick replication-view publication: the answering shard’s current master_repl_offset (== ReplicationSource::next_offset()) plus the per-replica (ipv4, port, sent_offset) triple for every handshake-complete replica (in AckSent, Streaming, or SnapshotShipping). connected_slaves for INFO / ROLE is derived as replicas.len(). Only called when this shard has a ReplicationSource installed (i.e. Runtime::with_replication(true, ...) was requested); standalone setups pay nothing. Command layers that serve ROLE / INFO replication stash the values in a thread-local (thread-per-core: the answering thread is the shard, same pattern as Self::on_persist_stats). Default no-op.
Source§

fn on_command(&self)

Called once per client command at dispatch entry (before routing / fan-out, so a multi-key command counts once). kevy uses it for INFO stats: total_commands_processed. Hot path — keep it to a single thread-local bump. Default no-op so non-kevy embedders pay nothing.
Source§

fn on_connection(&self)

Called once per accepted client connection. kevy uses it for INFO stats: total_connections_received. Default no-op.
Source§

fn shard_tick_interval_ms(&self) -> u64

Interval between Self::on_shard_tick calls. Default 100 ms (matching Redis’s hz = 10). 0 disables ticking entirely.
Source§

fn on_shard_tick(&self, store: &mut Store)

Periodic shard housekeeping (the equivalent of Redis’s serverCron). kevy uses this to run Store::tick_expire at the configured [expiry].hz. Default no-op so non-kevy embedders / runtimes can ignore it.
Source§

fn live_runtime_config(&self) -> LiveRuntimeConfig

Snapshot of the runtime-owned knobs that can be hot-modified (the kevy server wires this to CONFIG SET). Called once per shard tick — each Some value is applied to the shard’s live state; each None keeps the existing setting untouched. Read more
Source§

fn hello_reply<A: ArgvView + ?Sized>( &self, args: &A, current_proto: RespVersion, ) -> (RespVersion, Vec<u8>)

Handle HELLO — return the new connection protocol version + the reply bytes. The runtime applies the new version to the conn before scheduling the reply, so a HELLO 3 ack itself comes out shaped as a RESP3 Map (the new protocol is in effect for its own reply). Read more
Source§

fn is_write<A: ArgvView + ?Sized>(&self, args: &A) -> bool

Whether this command mutates the keyspace (so it must be logged to the AOF).
Source§

fn notify_class<A: ArgvView + ?Sized>(&self, args: &A) -> Option<NotifyClass>

Classify a command for keyspace notifications. Returns Some for write commands that should fire a notification when the corresponding flag is enabled; None for read-only / no-op / not-yet-classified commands (those never publish). Default None so non-kevy embedders pay nothing.
Source§

fn txn_kind<A: ArgvView + ?Sized>(&self, args: &A) -> TxnKind

Transaction-control classification (MULTI/EXEC/DISCARD vs anything else).
Source§

fn block_serve_argv<A: ArgvView + ?Sized>( &self, args: &A, kind: BlockKind, key: &[u8], ) -> Argv

Build the single-key command the dispatcher will replay to satisfy one watched key of a (possibly multi-key) blocking command. args is the original command; key is one of its watched keys. Returns an Argv that, when dispatched, pops / reads only key — e.g. BLPOP k1 k2 0 watching k2 yields BLPOP k2 0; XREAD … STREAMS s1 s2 id1 id2 watching s2 yields XREAD … STREAMS s2 id2. Read more
Source§

fn block_ready<A: ArgvView + ?Sized>( &self, store: &mut Store, serve_argv: &A, kind: BlockKind, ) -> bool

Non-destructive readiness peek for a parked waiter: would replaying serve_argv (built by Self::block_serve_argv, $ already frozen) produce a reply right now? Runs on the key’s owning shard when arming and is the gate for emitting a cross-shard wake. Must NOT mutate the store (no pop / no group-cursor advance). Default false so non-blocking embedders never spuriously wake.
Source§

fn wake_idx<A: ArgvView + ?Sized>(&self, args: &A) -> Option<u8>

Index into args of the key whose write may wake a blocked waiter (LPUSH / RPUSH feed BLPOP / BRPOP; XADD feeds the stream blocks). Some(1) for those verbs, None for everything else. The in-shard fast path reads this off ResolvedCmd::wake_idx; the cross-shard write path (exec_op, where a forwarded write lands on the key’s owning shard) re-derives it via this method since the forwarded envelope doesn’t carry the resolved hint. Default None so non-blocking embedders pay nothing.
Source§

fn block_hint<A>(&self, _args: &A) -> BlockHint
where A: ArgvView + ?Sized,

Classify a command for blocking semantics. BlockHint::None (default) is the zero-cost answer for every non-blocking verb; the dispatcher only registers a waiter when this returns BlockHint::Block and the command’s dispatch_into produced no reply (i.e. it could not satisfy itself immediately — e.g. BLPOP on an empty list). Concrete impls should fold this into their override of Self::resolve so the verb-table lookup happens once per command.
Source§

impl Copy for KevyCommands

Source§

impl Default for KevyCommands

Source§

fn default() -> KevyCommands

Returns the “default value” for a type. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.