-
Logs
+
+
{poolHeaderText(pools, "Logs", "Logs")}
+ {pools.length > 1 ? (
+
+ navigate({
+ to: ".",
+ search: (old) => ({
+ ...old,
+ pool: name,
+ }),
+ })
+ }
+ />
+ ) : null}
+
@@ -149,7 +190,7 @@ function RouteComponent() {
,
pub create_in_region: Option,
pub create_with_input: Option,
+ /// Overrides the client's configured pool name for this actor if it is created.
+ pub pool_name: Option,
}
#[derive(Default)]
@@ -32,6 +34,8 @@ pub struct CreateOptions {
pub params: Option,
pub region: Option,
pub input: Option,
+ /// Overrides the client's configured pool name for this actor.
+ pub pool_name: Option,
}
pub struct ClientConfig {
@@ -219,6 +223,7 @@ impl Client {
key: key,
input,
region,
+ pool_name: opts.pool_name,
},
};
@@ -236,7 +241,10 @@ impl Client {
let input = opts.input;
let _region = opts.region;
- let actor_id = self.remote_manager.create_actor(name, &key, input).await?;
+ let actor_id = self
+ .remote_manager
+ .create_actor(name, &key, input, opts.pool_name)
+ .await?;
let get_query = ActorQuery::GetForId {
get_for_id: GetForIdRequest {
diff --git a/rivetkit-rust/packages/client/src/protocol/query.rs b/rivetkit-rust/packages/client/src/protocol/query.rs
index 85ec36db90..c837c4bedf 100644
--- a/rivetkit-rust/packages/client/src/protocol/query.rs
+++ b/rivetkit-rust/packages/client/src/protocol/query.rs
@@ -11,6 +11,8 @@ pub struct CreateRequest {
pub input: Option,
#[serde(skip_serializing_if = "Option::is_none")]
pub region: Option,
+ #[serde(rename = "poolName", skip_serializing_if = "Option::is_none")]
+ pub pool_name: Option,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
@@ -34,6 +36,8 @@ pub struct GetOrCreateRequest {
pub input: Option,
#[serde(skip_serializing_if = "Option::is_none")]
pub region: Option,
+ #[serde(rename = "poolName", skip_serializing_if = "Option::is_none")]
+ pub pool_name: Option,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
diff --git a/rivetkit-rust/packages/client/src/remote_manager.rs b/rivetkit-rust/packages/client/src/remote_manager.rs
index c1525f59f7..a15bd23e88 100644
--- a/rivetkit-rust/packages/client/src/remote_manager.rs
+++ b/rivetkit-rust/packages/client/src/remote_manager.rs
@@ -278,6 +278,7 @@ impl RemoteManager {
name: &str,
key: &ActorKey,
input: Option,
+ pool_name: Option,
) -> Result {
let config = self.resolved_config().await?;
// Canonical slash-escaped key format (matches the TS SDK and the
@@ -295,7 +296,7 @@ impl RemoteManager {
name: name.to_string(),
key: key_str,
input: input_encoded,
- runner_name_selector: self.pool_name.clone(),
+ runner_name_selector: pool_name.unwrap_or_else(|| self.pool_name.clone()),
crash_policy: "destroy".to_string(),
};
@@ -327,6 +328,7 @@ impl RemoteManager {
name: &str,
key: &ActorKey,
input: Option,
+ pool_name: Option,
) -> Result {
let config = self.resolved_config().await?;
// Canonical slash-escaped key format (matches the TS SDK and the
@@ -344,7 +346,7 @@ impl RemoteManager {
name: name.to_string(),
key: key_str,
input: input_encoded,
- runner_name_selector: self.pool_name.clone(),
+ runner_name_selector: pool_name.unwrap_or_else(|| self.pool_name.clone()),
crash_policy: "destroy".to_string(),
};
@@ -388,12 +390,18 @@ impl RemoteManager {
&get_or_create_for_key.name,
&get_or_create_for_key.key,
get_or_create_for_key.input.clone(),
+ get_or_create_for_key.pool_name.clone(),
)
.await
}
ActorQuery::Create { create } => {
- self.create_actor(&create.name, &create.key, create.input.clone())
- .await
+ self.create_actor(
+ &create.name,
+ &create.key,
+ create.input.clone(),
+ create.pool_name.clone(),
+ )
+ .await
}
}
}
@@ -438,6 +446,7 @@ impl RemoteManager {
Some(&get_for_key.key),
None,
None,
+ None,
),
ActorQuery::GetOrCreateForKey {
get_or_create_for_key,
@@ -447,6 +456,7 @@ impl RemoteManager {
Some(&get_or_create_for_key.key),
get_or_create_for_key.input.as_ref(),
get_or_create_for_key.region.as_deref(),
+ get_or_create_for_key.pool_name.as_deref(),
),
ActorQuery::Create { .. } => {
Err(anyhow!("gateway URL does not support create actor queries"))
@@ -488,6 +498,7 @@ impl RemoteManager {
key: Option<&ActorKey>,
input: Option<&serde_json::Value>,
region: Option<&str>,
+ pool_name: Option<&str>,
) -> Result {
if self.namespace.is_empty() {
return Err(anyhow!("actor query namespace must not be empty"));
@@ -512,7 +523,7 @@ impl RemoteManager {
push_query_param(&mut params, "rvt-input", &URL_SAFE_NO_PAD.encode(encoded));
}
if method == "getOrCreate" {
- push_query_param(&mut params, "rvt-runner", &self.pool_name);
+ push_query_param(&mut params, "rvt-runner", pool_name.unwrap_or(&self.pool_name));
push_query_param(&mut params, "rvt-crash-policy", "sleep");
}
if let Some(region) = region {
diff --git a/rivetkit-rust/packages/client/tests/bare.rs b/rivetkit-rust/packages/client/tests/bare.rs
index 05de82ef9d..e17834fcae 100644
--- a/rivetkit-rust/packages/client/tests/bare.rs
+++ b/rivetkit-rust/packages/client/tests/bare.rs
@@ -25,8 +25,8 @@ use reqwest::{
Method, Url,
};
use rivetkit_client::{
- Client, ClientConfig, ConnectionStatus, EncodingKind, GetOptions, GetOrCreateOptions,
- QueueSendStatus, SendAndWaitOpts, SendOpts,
+ Client, ClientConfig, ConnectionStatus, CreateOptions, EncodingKind, GetOptions,
+ GetOrCreateOptions, QueueSendStatus, SendAndWaitOpts, SendOpts,
};
use rivetkit_client_protocol as wire;
use serde::{Deserialize, Serialize};
@@ -95,6 +95,43 @@ struct ActorResponse {
created: bool,
}
+#[derive(Deserialize)]
+struct RunnerSelectorRequest {
+ name: String,
+ key: String,
+ runner_name_selector: String,
+}
+
+#[derive(Clone, Default)]
+struct RunnerSelectorState {
+ seen: Arc>>,
+}
+
+impl RunnerSelectorState {
+ fn take(&self) -> Vec {
+ std::mem::take(&mut self.seen.lock().unwrap())
+ }
+}
+
+async fn capture_runner_selector(
+ State(state): State,
+ Json(request): Json,
+) -> impl IntoResponse {
+ state
+ .seen
+ .lock()
+ .unwrap()
+ .push(request.runner_name_selector);
+ Json(ActorResponse {
+ actor: Actor {
+ actor_id: "actor-1",
+ name: request.name,
+ key: request.key,
+ },
+ created: true,
+ })
+}
+
#[tokio::test]
async fn default_bare_action_round_trips_against_test_actor() {
assert_eq!(EncodingKind::default(), EncodingKind::Bare);
@@ -611,6 +648,167 @@ fn gateway_url_uses_query_backed_get_or_create_target() {
.is_some_and(|value| !value.is_empty()));
}
+#[tokio::test]
+async fn create_sends_per_call_pool_name_as_runner_name_selector() {
+ let state = RunnerSelectorState::default();
+ let app = Router::new()
+ .route("/actors", post(capture_runner_selector))
+ .with_state(state.clone());
+
+ let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
+ let addr = listener.local_addr().unwrap();
+ let server = tokio::spawn(async move {
+ axum::serve(listener, app).await.unwrap();
+ });
+
+ let client = Client::new(
+ ClientConfig::new(endpoint(addr))
+ .namespace("ns")
+ .pool_name("config-pool")
+ .disable_metadata_lookup(true),
+ );
+ client
+ .create(
+ "counter",
+ vec!["k".to_owned()],
+ CreateOptions {
+ pool_name: Some("call-pool".to_owned()),
+ ..Default::default()
+ },
+ )
+ .await
+ .unwrap();
+
+ assert_eq!(state.take(), vec!["call-pool".to_owned()]);
+
+ server.abort();
+}
+
+#[tokio::test]
+async fn create_falls_back_to_config_pool_name() {
+ let state = RunnerSelectorState::default();
+ let app = Router::new()
+ .route("/actors", post(capture_runner_selector))
+ .with_state(state.clone());
+
+ let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
+ let addr = listener.local_addr().unwrap();
+ let server = tokio::spawn(async move {
+ axum::serve(listener, app).await.unwrap();
+ });
+
+ let client = Client::new(
+ ClientConfig::new(endpoint(addr))
+ .namespace("ns")
+ .pool_name("config-pool")
+ .disable_metadata_lookup(true),
+ );
+ client
+ .create("counter", vec!["k".to_owned()], CreateOptions::default())
+ .await
+ .unwrap();
+
+ assert_eq!(state.take(), vec!["config-pool".to_owned()]);
+
+ server.abort();
+}
+
+#[tokio::test]
+async fn get_or_create_sends_per_call_pool_name_as_runner_name_selector() {
+ let state = RunnerSelectorState::default();
+ let app = Router::new()
+ .route("/actors", put(capture_runner_selector))
+ .with_state(state.clone());
+
+ let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
+ let addr = listener.local_addr().unwrap();
+ let server = tokio::spawn(async move {
+ axum::serve(listener, app).await.unwrap();
+ });
+
+ let client = Client::new(
+ ClientConfig::new(endpoint(addr))
+ .namespace("ns")
+ .pool_name("config-pool")
+ .disable_metadata_lookup(true),
+ );
+ let actor = client
+ .get_or_create(
+ "counter",
+ vec!["k".to_owned()],
+ GetOrCreateOptions {
+ pool_name: Some("call-pool".to_owned()),
+ ..Default::default()
+ },
+ )
+ .unwrap();
+ // Resolving the handle to an actor ID triggers the PUT /actors request
+ // that carries runner_name_selector.
+ actor.resolve().await.unwrap();
+
+ assert_eq!(state.take(), vec!["call-pool".to_owned()]);
+
+ server.abort();
+}
+
+#[tokio::test]
+async fn get_or_create_falls_back_to_config_pool_name() {
+ let state = RunnerSelectorState::default();
+ let app = Router::new()
+ .route("/actors", put(capture_runner_selector))
+ .with_state(state.clone());
+
+ let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
+ let addr = listener.local_addr().unwrap();
+ let server = tokio::spawn(async move {
+ axum::serve(listener, app).await.unwrap();
+ });
+
+ let client = Client::new(
+ ClientConfig::new(endpoint(addr))
+ .namespace("ns")
+ .pool_name("config-pool")
+ .disable_metadata_lookup(true),
+ );
+ let actor = client
+ .get_or_create("counter", vec!["k".to_owned()], GetOrCreateOptions::default())
+ .unwrap();
+ // Resolving the handle to an actor ID triggers the PUT /actors request
+ // that carries runner_name_selector.
+ actor.resolve().await.unwrap();
+
+ assert_eq!(state.take(), vec!["config-pool".to_owned()]);
+
+ server.abort();
+}
+
+#[test]
+fn gateway_url_get_or_create_uses_per_call_pool_name() {
+ let client = Client::new(
+ ClientConfig::new("http://127.0.0.1:6420")
+ .namespace("ns")
+ .pool_name("config-pool")
+ .disable_metadata_lookup(true),
+ );
+ let actor = client
+ .get_or_create(
+ "chat room",
+ vec!["tenant".to_owned()],
+ GetOrCreateOptions {
+ pool_name: Some("call-pool".to_owned()),
+ ..Default::default()
+ },
+ )
+ .unwrap();
+
+ let url = Url::parse(&actor.gateway_url().unwrap()).unwrap();
+ let params = query_params(&url);
+ assert_eq!(
+ params.get("rvt-runner").map(String::as_str),
+ Some("call-pool")
+ );
+}
+
#[tokio::test]
async fn metadata_lookup_overrides_endpoint_before_requests() {
let target_state = TestState {
diff --git a/rivetkit-typescript/packages/rivetkit/src/client/client.ts b/rivetkit-typescript/packages/rivetkit/src/client/client.ts
index 7130f325b0..33262b0096 100644
--- a/rivetkit-typescript/packages/rivetkit/src/client/client.ts
+++ b/rivetkit-typescript/packages/rivetkit/src/client/client.ts
@@ -128,6 +128,11 @@ export interface GetOrCreateOptions extends QueryOptions {
createInRegion?: string;
/** Input data to pass to the actor. */
createWithInput?: unknown;
+ /**
+ * Name of the envoy pool to select for this actor if it is created.
+ * Overrides the client's configured `poolName` for this call.
+ */
+ poolName?: string;
}
/**
@@ -140,6 +145,11 @@ export interface CreateOptions extends QueryOptions {
region?: string;
/** Input data to pass to the actor. */
input?: unknown;
+ /**
+ * Name of the envoy pool to select for this actor.
+ * Overrides the client's configured `poolName` for this call.
+ */
+ poolName?: string;
}
/**
@@ -300,6 +310,7 @@ export class ClientRaw {
key: keyArray,
input: opts?.createWithInput,
region: opts?.createInRegion,
+ poolName: opts?.poolName,
},
};
diff --git a/rivetkit-typescript/packages/rivetkit/src/client/query.ts b/rivetkit-typescript/packages/rivetkit/src/client/query.ts
index 02e23dced4..dc7a83dae6 100644
--- a/rivetkit-typescript/packages/rivetkit/src/client/query.ts
+++ b/rivetkit-typescript/packages/rivetkit/src/client/query.ts
@@ -31,6 +31,7 @@ export const CreateRequestSchema = z.object({
key: ActorKeySchema,
input: z.unknown().optional(),
region: z.string().optional(),
+ poolName: z.string().optional(),
});
export const GetForKeyRequestSchema = z.object({
@@ -43,6 +44,7 @@ export const GetOrCreateRequestSchema = z.object({
key: ActorKeySchema,
input: z.unknown().optional(),
region: z.string().optional(),
+ poolName: z.string().optional(),
});
export const ActorQuerySchema = z.union([
diff --git a/rivetkit-typescript/packages/rivetkit/src/client/resolve-gateway-target.ts b/rivetkit-typescript/packages/rivetkit/src/client/resolve-gateway-target.ts
index b45005eea8..c089802470 100644
--- a/rivetkit-typescript/packages/rivetkit/src/client/resolve-gateway-target.ts
+++ b/rivetkit-typescript/packages/rivetkit/src/client/resolve-gateway-target.ts
@@ -41,6 +41,7 @@ export async function resolveGatewayTarget(
key: target.getOrCreateForKey.key,
input: target.getOrCreateForKey.input,
region: target.getOrCreateForKey.region,
+ poolName: target.getOrCreateForKey.poolName,
});
return output.actorId;
}
@@ -51,6 +52,7 @@ export async function resolveGatewayTarget(
key: target.create.key,
input: target.create.input,
region: target.create.region,
+ poolName: target.create.poolName,
});
return output.actorId;
}
diff --git a/rivetkit-typescript/packages/rivetkit/src/engine-client/driver.ts b/rivetkit-typescript/packages/rivetkit/src/engine-client/driver.ts
index c447b3dd1c..14f36b31ef 100644
--- a/rivetkit-typescript/packages/rivetkit/src/engine-client/driver.ts
+++ b/rivetkit-typescript/packages/rivetkit/src/engine-client/driver.ts
@@ -86,6 +86,8 @@ export interface GetOrCreateWithKeyInput {
input?: unknown;
region?: string;
crashPolicy?: CrashPolicy;
+ /** Overrides the client's configured pool name for this actor. */
+ poolName?: string;
}
export interface CreateInput {
@@ -95,6 +97,8 @@ export interface CreateInput {
input?: unknown;
region?: string;
crashPolicy?: CrashPolicy;
+ /** Overrides the client's configured pool name for this actor. */
+ poolName?: string;
}
export interface ListActorsInput {
diff --git a/rivetkit-typescript/packages/rivetkit/src/engine-client/mod.ts b/rivetkit-typescript/packages/rivetkit/src/engine-client/mod.ts
index 6e0d7e78f6..8efd0ca552 100644
--- a/rivetkit-typescript/packages/rivetkit/src/engine-client/mod.ts
+++ b/rivetkit-typescript/packages/rivetkit/src/engine-client/mod.ts
@@ -167,7 +167,14 @@ export class RemoteEngineControlClient implements EngineControlClient {
): Promise {
await this.#metadataPromise;
- const { name, key, input: actorInput, region, crashPolicy } = input;
+ const {
+ name,
+ key,
+ input: actorInput,
+ region,
+ crashPolicy,
+ poolName,
+ } = input;
logger().info({
msg: "getOrCreateWithKey: getting or creating actor via engine api",
@@ -179,7 +186,7 @@ export class RemoteEngineControlClient implements EngineControlClient {
datacenter: region,
name,
key: serializeActorKey(key),
- runner_name_selector: this.#config.poolName,
+ runner_name_selector: poolName ?? this.#config.poolName,
input: actorInput
? uint8ArrayToBase64(
encodeCborCompat(actorInput as JsonCompatValue),
@@ -205,6 +212,7 @@ export class RemoteEngineControlClient implements EngineControlClient {
input,
region,
crashPolicy,
+ poolName,
}: CreateInput): Promise {
await this.#metadataPromise;
@@ -214,7 +222,7 @@ export class RemoteEngineControlClient implements EngineControlClient {
const result = await createActor(this.#config, {
datacenter: region,
name,
- runner_name_selector: this.#config.poolName,
+ runner_name_selector: poolName ?? this.#config.poolName,
key: serializeActorKey(key),
input: input
? uint8ArrayToBase64(encodeCborCompat(input as JsonCompatValue))
@@ -422,7 +430,8 @@ export class RemoteEngineControlClient implements EngineControlClient {
this.#config.maxInputSize,
undefined,
"getOrCreateForKey" in target
- ? this.#config.poolName
+ ? (target.getOrCreateForKey.poolName ??
+ this.#config.poolName)
: undefined,
options,
);
diff --git a/rivetkit-typescript/packages/rivetkit/tests/remote-engine-client-pool-name.test.ts b/rivetkit-typescript/packages/rivetkit/tests/remote-engine-client-pool-name.test.ts
new file mode 100644
index 0000000000..7150fd30aa
--- /dev/null
+++ b/rivetkit-typescript/packages/rivetkit/tests/remote-engine-client-pool-name.test.ts
@@ -0,0 +1,118 @@
+import { afterEach, beforeEach, describe, expect, test, vi } from "vitest";
+import { ClientConfigSchema } from "@/client/config";
+import { RemoteEngineControlClient } from "@/engine-client/mod";
+
+describe.sequential("RemoteEngineControlClient poolName selection", () => {
+ beforeEach(() => {
+ vi.restoreAllMocks();
+ });
+
+ afterEach(() => {
+ vi.unstubAllGlobals();
+ });
+
+ function makeDriver(): RemoteEngineControlClient {
+ return new RemoteEngineControlClient(
+ ClientConfigSchema.parse({
+ endpoint: "https://api.rivet.dev",
+ namespace: "default",
+ poolName: "config-pool",
+ disableMetadataLookup: true,
+ }),
+ );
+ }
+
+ function stubActorFetch(): Request[] {
+ const requests: Request[] = [];
+ vi.stubGlobal(
+ "fetch",
+ vi.fn(async (input: Request) => {
+ requests.push(input);
+ return new Response(
+ JSON.stringify({
+ actor: { actor_id: "act-1", name: "counter", key: "a" },
+ created: true,
+ }),
+ { headers: { "content-type": "application/json" } },
+ );
+ }),
+ );
+ return requests;
+ }
+
+ test("createActor uses per-call poolName as runner_name_selector", async () => {
+ const requests = stubActorFetch();
+ const driver = makeDriver();
+
+ await driver.createActor({
+ name: "counter",
+ key: ["a"],
+ poolName: "call-pool",
+ });
+
+ expect(requests).toHaveLength(1);
+ const body = await requests[0]!.json();
+ expect(body.runner_name_selector).toBe("call-pool");
+ });
+
+ test("createActor falls back to configured poolName", async () => {
+ const requests = stubActorFetch();
+ const driver = makeDriver();
+
+ await driver.createActor({ name: "counter", key: ["a"] });
+
+ expect(requests).toHaveLength(1);
+ const body = await requests[0]!.json();
+ expect(body.runner_name_selector).toBe("config-pool");
+ });
+
+ test("getOrCreateWithKey uses per-call poolName as runner_name_selector", async () => {
+ const requests = stubActorFetch();
+ const driver = makeDriver();
+
+ await driver.getOrCreateWithKey({
+ name: "counter",
+ key: ["a"],
+ poolName: "call-pool",
+ });
+
+ expect(requests).toHaveLength(1);
+ const body = await requests[0]!.json();
+ expect(body.runner_name_selector).toBe("call-pool");
+ });
+
+ test("getOrCreateWithKey falls back to configured poolName", async () => {
+ const requests = stubActorFetch();
+ const driver = makeDriver();
+
+ await driver.getOrCreateWithKey({ name: "counter", key: ["a"] });
+
+ expect(requests).toHaveLength(1);
+ const body = await requests[0]!.json();
+ expect(body.runner_name_selector).toBe("config-pool");
+ });
+
+ test("getOrCreate gateway URL uses per-call poolName for rvt-runner", async () => {
+ const driver = makeDriver();
+
+ const url = await driver.buildGatewayUrl({
+ getOrCreateForKey: {
+ name: "room",
+ key: ["a"],
+ poolName: "call-pool",
+ },
+ });
+
+ expect(new URL(url).searchParams.get("rvt-runner")).toBe("call-pool");
+ });
+
+ test("getOrCreate gateway URL falls back to configured poolName for rvt-runner", async () => {
+ const driver = makeDriver();
+
+ const url = await driver.buildGatewayUrl({
+ getOrCreateForKey: { name: "room", key: ["a"] },
+ });
+
+ expect(new URL(url).searchParams.get("rvt-runner")).toBe("config-pool");
+ });
+});
diff --git a/website/src/content/docs/clients/rust.mdx b/website/src/content/docs/clients/rust.mdx
index 9eb2fa2ceb..d5f89014ad 100644
--- a/website/src/content/docs/clients/rust.mdx
+++ b/website/src/content/docs/clients/rust.mdx
@@ -122,7 +122,7 @@ async fn main() -> Result<()> {
}
```
-`get_typed_default` / `get_or_create_typed_default` use default options. The non-default variants (`get_typed` / `get_or_create_typed`) take `GetOptions` / `GetOrCreateOptions` for connection parameters, input, and region.
+`get_typed_default` / `get_or_create_typed_default` use default options. The non-default variants (`get_typed` / `get_or_create_typed`) take `GetOptions` / `GetOrCreateOptions` for connection parameters, input, and region. Set `pool_name` on `GetOrCreateOptions` or `CreateOptions` to override the client's configured pool for a single actor.
## Connection Parameters