diff --git a/Cargo.lock b/Cargo.lock index 729123cdcc..c1e7b14c3c 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1089,7 +1089,7 @@ dependencies = [ [[package]] name = "bifrost-benchpress" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "anyhow", "bytes", @@ -1711,7 +1711,7 @@ checksum = "3f88a43d011fc4a6876cb7344703e297c71dda42494fee094d5f7c76bf13f746" [[package]] name = "codederror" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "codederror-derive", "restate-workspace-hack", @@ -5253,7 +5253,7 @@ checksum = "5e5032e24019045c762d3c0f28f5b6b8bbf38563a65908389bf7978758920897" [[package]] name = "logserver-bench" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "anyhow", "axum", @@ -5536,7 +5536,7 @@ checksum = "f2b8f3a258db515d5e91a904ce4ae3f73e091149b90cadbdb93d210bee07f63b" [[package]] name = "mock-service-endpoint" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "assert2", "async-stream", @@ -6392,7 +6392,7 @@ checksum = "439ee305def115ba05938db6eb1644ff94165c5ab5e9420d1c1bcedbba909391" [[package]] name = "pp-bench" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "anyhow", "axum", @@ -7374,7 +7374,7 @@ checksum = "1e061d1b48cb8d38042de4ae0a7a6401009d6143dc80d2e2d6f31f0bdd6470c7" [[package]] name = "restate-admin" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "ahash", "anyhow", @@ -7435,7 +7435,7 @@ dependencies = [ [[package]] name = "restate-admin-rest-model" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "bytes", "bytestring", @@ -7458,7 +7458,7 @@ dependencies = [ [[package]] name = "restate-base64-util" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "base64 0.22.1", "restate-workspace-hack", @@ -7466,7 +7466,7 @@ dependencies = [ [[package]] name = "restate-benchmarks" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "anyhow", "criterion", @@ -7493,7 +7493,7 @@ dependencies = [ [[package]] name = "restate-bifrost" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "adaptive-timeout", "ahash", @@ -7550,7 +7550,7 @@ dependencies = [ [[package]] name = "restate-cli" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "anyhow", "arc-swap", @@ -7621,7 +7621,7 @@ dependencies = [ [[package]] name = "restate-cli-util" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "anyhow", "arc-swap", @@ -7650,7 +7650,7 @@ dependencies = [ [[package]] name = "restate-clock" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "bilrost", "criterion", @@ -7669,7 +7669,7 @@ dependencies = [ [[package]] name = "restate-cloud-tunnel-client" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "anyhow", "bs58", @@ -7689,7 +7689,7 @@ dependencies = [ [[package]] name = "restate-core" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "ahash", "anyhow", @@ -7758,7 +7758,7 @@ dependencies = [ [[package]] name = "restate-doctor" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "anyhow", "arrow", @@ -7811,7 +7811,7 @@ dependencies = [ [[package]] name = "restate-encoding" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "bilrost", "bytes", @@ -7832,7 +7832,7 @@ dependencies = [ [[package]] name = "restate-errors" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "codederror", "paste", @@ -7844,7 +7844,7 @@ dependencies = [ [[package]] name = "restate-fs-util" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "restate-workspace-hack", "tokio", @@ -7854,7 +7854,7 @@ dependencies = [ [[package]] name = "restate-futures-util" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "anyhow", "async-channel", @@ -7874,7 +7874,7 @@ dependencies = [ [[package]] name = "restate-hyper-uds" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "http 1.4.2", "hyper-util", @@ -7886,7 +7886,7 @@ dependencies = [ [[package]] name = "restate-ingestion-client" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "bytes", "dashmap", @@ -7907,7 +7907,7 @@ dependencies = [ [[package]] name = "restate-ingress-http" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "anyhow", "bytes", @@ -7954,7 +7954,7 @@ dependencies = [ [[package]] name = "restate-ingress-kafka" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "anyhow", "base64 0.22.1", @@ -7984,7 +7984,7 @@ dependencies = [ [[package]] name = "restate-invoker-impl" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "anyhow", "bytes", @@ -8032,7 +8032,7 @@ dependencies = [ [[package]] name = "restate-limiter" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "bilrost", "bytes", @@ -8050,7 +8050,7 @@ dependencies = [ [[package]] name = "restate-lite" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "anyhow", "http 1.4.2", @@ -8079,7 +8079,7 @@ dependencies = [ [[package]] name = "restate-local-cluster-runner" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "anyhow", "arc-swap", @@ -8116,7 +8116,7 @@ dependencies = [ [[package]] name = "restate-log-server" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "ahash", "anyhow", @@ -8158,7 +8158,7 @@ dependencies = [ [[package]] name = "restate-log-server-grpc" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "prost", "restate-types", @@ -8170,7 +8170,7 @@ dependencies = [ [[package]] name = "restate-memory" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "bytes", "futures", @@ -8186,7 +8186,7 @@ dependencies = [ [[package]] name = "restate-metadata-providers" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "anyhow", "async-trait", @@ -8222,7 +8222,7 @@ dependencies = [ [[package]] name = "restate-metadata-server" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "anyhow", "arc-swap", @@ -8267,7 +8267,7 @@ dependencies = [ [[package]] name = "restate-metadata-server-grpc" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "bytes", "bytestring", @@ -8288,7 +8288,7 @@ dependencies = [ [[package]] name = "restate-metadata-store" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "async-trait", "bytes", @@ -8310,7 +8310,7 @@ dependencies = [ [[package]] name = "restate-node" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "ahash", "anyhow", @@ -8373,7 +8373,7 @@ dependencies = [ [[package]] name = "restate-object-store-util" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "anyhow", "aws-config", @@ -8391,7 +8391,7 @@ dependencies = [ [[package]] name = "restate-partition-store" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "ahash", "anyhow", @@ -8447,7 +8447,7 @@ dependencies = [ [[package]] name = "restate-platform" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "bilrost", "bytes", @@ -8463,7 +8463,7 @@ dependencies = [ [[package]] name = "restate-queue" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "bincode", "criterion", @@ -8478,7 +8478,7 @@ dependencies = [ [[package]] name = "restate-rocksdb" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "anyhow", "bytes", @@ -8509,7 +8509,7 @@ dependencies = [ [[package]] name = "restate-serde-util" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "bytes", "http 1.4.2", @@ -8525,7 +8525,7 @@ dependencies = [ [[package]] name = "restate-server" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "anyhow", "bytestring", @@ -8575,7 +8575,7 @@ dependencies = [ [[package]] name = "restate-service-client" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "ahash", "arc-swap", @@ -8629,7 +8629,7 @@ dependencies = [ [[package]] name = "restate-service-protocol" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "bytes", "bytes-utils", @@ -8650,7 +8650,7 @@ dependencies = [ [[package]] name = "restate-service-protocol-v4" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "assert2", "bytes", @@ -8685,7 +8685,7 @@ dependencies = [ [[package]] name = "restate-sharding" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "bilrost", "derive_more", @@ -8699,7 +8699,7 @@ dependencies = [ [[package]] name = "restate-storage-api" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "ahash", "anyhow", @@ -8727,7 +8727,7 @@ dependencies = [ [[package]] name = "restate-storage-query-datafusion" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "ahash", "anyhow", @@ -8772,7 +8772,7 @@ dependencies = [ [[package]] name = "restate-test-util" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "assert2", "bytes", @@ -8787,7 +8787,7 @@ dependencies = [ [[package]] name = "restate-timer" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "ahash", "futures-util", @@ -8804,7 +8804,7 @@ dependencies = [ [[package]] name = "restate-timer-queue" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "futures", "restate-workspace-hack", @@ -8813,7 +8813,7 @@ dependencies = [ [[package]] name = "restate-tracing-instrumentation" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "console-subscriber", "criterion", @@ -8843,7 +8843,7 @@ dependencies = [ [[package]] name = "restate-types" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "adaptive-timeout", "ahash", @@ -8942,7 +8942,7 @@ dependencies = [ [[package]] name = "restate-util-bytecount" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "bilrost", "bytesize", @@ -8958,14 +8958,14 @@ dependencies = [ [[package]] name = "restate-util-random" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "restate-workspace-hack", ] [[package]] name = "restate-util-string" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "ahash", "bilrost", @@ -8983,7 +8983,7 @@ dependencies = [ [[package]] name = "restate-util-time" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "jiff", "rand 0.10.1", @@ -8997,7 +8997,7 @@ dependencies = [ [[package]] name = "restate-utoipa" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "indexmap 2.14.0", "restate-utoipa", @@ -9008,7 +9008,7 @@ dependencies = [ [[package]] name = "restate-vqueues" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "arrayvec", "bilrost", @@ -9047,7 +9047,7 @@ dependencies = [ [[package]] name = "restate-wal-protocol" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "anyhow", "bilrost", @@ -9078,7 +9078,7 @@ dependencies = [ [[package]] name = "restate-worker" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "ahash", "anyhow", @@ -9142,7 +9142,7 @@ dependencies = [ [[package]] name = "restate-worker-api" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "bytes", "codederror", @@ -9326,7 +9326,7 @@ dependencies = [ [[package]] name = "restatectl" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "anyhow", "arrow", @@ -9911,7 +9911,7 @@ dependencies = [ [[package]] name = "service-protocol-wireshark-dissector" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "bytes", "mlua", @@ -12144,7 +12144,7 @@ checksum = "66fee0b777b0f5ac1c69bb06d361268faafa61cd4682ae064a171c16c433e9e4" [[package]] name = "xtask" -version = "1.7.2-dev" +version = "1.8.0-dev" dependencies = [ "anyhow", "reqwest 0.12.28", diff --git a/Cargo.toml b/Cargo.toml index 41c81749d7..f8fdadb8b2 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -32,7 +32,7 @@ default-members = [ resolver = "2" [workspace.package] -version = "1.7.2-dev" +version = "1.8.0-dev" authors = ["restate.dev"] edition = "2024" rust-version = "1.95.0" diff --git a/charts/restate-helm/Chart.yaml b/charts/restate-helm/Chart.yaml index 966f8a9833..c8e4dc2083 100644 --- a/charts/restate-helm/Chart.yaml +++ b/charts/restate-helm/Chart.yaml @@ -2,7 +2,7 @@ apiVersion: v2 name: restate-helm description: Deploy a Restate cluster on Kubernetes type: application -version: "1.7.2-dev" +version: "1.8.0-dev" home: https://restate.dev sources: - https://github.com/restatedev/restate diff --git a/crates/admin/src/rest_api/invocations.rs b/crates/admin/src/rest_api/invocations.rs index e50f23e9ef..107ced6cec 100644 --- a/crates/admin/src/rest_api/invocations.rs +++ b/crates/admin/src/rest_api/invocations.rs @@ -12,6 +12,7 @@ use axum::Json; use axum::extract::{Path, Query, State}; use axum::http::StatusCode; use futures::future; +use serde::Deserialize; use tracing::warn; use restate_admin_rest_model::invocations::{ @@ -28,12 +29,12 @@ use restate_types::invocation::client::{ }; use restate_types::invocation::{InvocationTermination, PurgeInvocationRequest, TerminationFlavor}; use restate_types::journal_v2::EntryIndex; -use restate_wal_protocol::{Command, Envelope}; -use serde::Deserialize; +use restate_types::logs::{BodyWithKeys, Keys}; +use restate_wal_protocol::v2; +use restate_wal_protocol::v2::commands; use super::error::*; use crate::generate_meta_api_error; -use crate::rest_api::create_envelope_header; use crate::state::AdminServiceState; #[derive(Debug, Default, Deserialize, utoipa::ToSchema)] @@ -88,30 +89,43 @@ where .parse::() .map_err(|e| MetaApiError::InvalidField("invocation_id", e.to_string()))?; - let cmd = match mode.unwrap_or_default() { - DeletionMode::Cancel => Command::TerminateInvocation(InvocationTermination { - invocation_id, - flavor: TerminationFlavor::Cancel, - response_sink: None, - }), - DeletionMode::Kill => Command::TerminateInvocation(InvocationTermination { - invocation_id, - flavor: TerminationFlavor::Kill, - response_sink: None, - }), - DeletionMode::Purge => Command::PurgeInvocation(PurgeInvocationRequest { - invocation_id, - response_sink: None, - }), + let envelope = match mode.unwrap_or_default() { + DeletionMode::Cancel => v2::Envelope::new( + v2::Dedup::None, + commands::TerminateInvocationCommand::from(InvocationTermination { + invocation_id, + flavor: TerminationFlavor::Cancel, + response_sink: None, + }), + ) + .into_raw(), + DeletionMode::Kill => v2::Envelope::new( + v2::Dedup::None, + commands::TerminateInvocationCommand::from(InvocationTermination { + invocation_id, + flavor: TerminationFlavor::Kill, + response_sink: None, + }), + ) + .into_raw(), + DeletionMode::Purge => v2::Envelope::new( + v2::Dedup::None, + commands::PurgeInvocationCommand::from(PurgeInvocationRequest { + invocation_id, + response_sink: None, + }), + ) + .into_raw(), }; let partition_key = invocation_id.partition_key(); - let envelope = Envelope::new(create_envelope_header(partition_key), cmd); - let result = state .ingestion_client - .ingest(partition_key, envelope) + .ingest( + partition_key, + BodyWithKeys::new(envelope, Keys::Single(partition_key)), + ) .await .map_err(|err| { warn!("Could not ingest invocation termination command: {err}"); diff --git a/crates/admin/src/rest_api/mod.rs b/crates/admin/src/rest_api/mod.rs index 83c9841425..70d55310c9 100644 --- a/crates/admin/src/rest_api/mod.rs +++ b/crates/admin/src/rest_api/mod.rs @@ -29,10 +29,8 @@ use utoipa::OpenApi; use utoipa_axum::{router::OpenApiRouter, routes}; use restate_core::network::TransportConnect; -use restate_types::identifiers::PartitionKey; use restate_types::invocation::client::InvocationClient; use restate_types::schema::registry::{DiscoveryClient, MetadataService, TelemetryClient}; -use restate_wal_protocol::{Destination, Header, Source}; use crate::state::AdminServiceState; @@ -187,16 +185,6 @@ where .with_state(state) } -fn create_envelope_header(partition_key: PartitionKey) -> Header { - Header { - source: Source::ControlPlane {}, - dest: Destination::Processor { - partition_key, - dedup: None, - }, - } -} - /// # Error description response /// /// Error details of the response diff --git a/crates/admin/src/rest_api/services.rs b/crates/admin/src/rest_api/services.rs index 62277ff099..d59fe4ef48 100644 --- a/crates/admin/src/rest_api/services.rs +++ b/crates/admin/src/rest_api/services.rs @@ -8,12 +8,11 @@ // the Business Source License, use of this software will be governed // by the Apache License, Version 2.0. -use tracing::{debug, warn}; - use axum::Json; use axum::extract::{Path, State}; use bytes::Bytes; use http::StatusCode; +use tracing::{debug, warn}; use restate_admin_rest_model::services::ListServicesResponse; use restate_admin_rest_model::services::*; @@ -22,13 +21,14 @@ use restate_core::network::TransportConnect; use restate_errors::warn_it; use restate_types::config::Configuration; use restate_types::identifiers::{ServiceId, WithPartitionKey}; +use restate_types::logs::{BodyWithKeys, Keys}; use restate_types::schema::registry::MetadataService; use restate_types::schema::service::ServiceMetadata; use restate_types::state_mut::ExternalStateMutation; use restate_types::{Scope, schema}; -use restate_wal_protocol::{Command, Envelope}; +use restate_wal_protocol::v2; +use restate_wal_protocol::v2::commands; -use super::create_envelope_header; use super::error::*; use crate::state::AdminServiceState; @@ -250,14 +250,17 @@ where state: new_state, }; - let envelope = Envelope::new( - create_envelope_header(partition_key), - Command::PatchState(patch_state), + let envelop = v2::Envelope::new( + v2::Dedup::None, + commands::PatchStateCommand::from(patch_state), ); let result = state .ingestion_client - .ingest(partition_key, envelope) + .ingest( + partition_key, + BodyWithKeys::new(envelop.into_raw(), Keys::Single(partition_key)), + ) .await .map_err(|err| { warn!("Could not ingest state patching command: {err}"); diff --git a/crates/admin/src/service.rs b/crates/admin/src/service.rs index 930c41dcee..a0d1714796 100644 --- a/crates/admin/src/service.rs +++ b/crates/admin/src/service.rs @@ -13,8 +13,6 @@ use std::time::Duration; use axum::error_handling::HandleErrorLayer; use http::{Request, Response, StatusCode}; -use restate_ingestion_client::IngestionClient; -use restate_wal_protocol::Envelope; use tower::ServiceBuilder; use tower_http::classify::ServerErrorsFailureClass; use tower_http::compression::CompressionLayer; @@ -24,6 +22,7 @@ use tracing::{Span, debug, info, info_span}; use restate_admin_rest_model::version::AdminApiVersion; use restate_core::network::{TransportConnect, net_util}; use restate_core::{MetadataWriter, TaskCenter}; +use restate_ingestion_client::IngestionClient; use restate_limiter::rule_book::RuleBookObserver; use restate_metadata_store::MetadataStoreClient; use restate_service_client::HttpClient; @@ -36,6 +35,7 @@ use restate_types::net::address::AdminPort; use restate_types::net::listener::Listeners; use restate_types::schema::registry::SchemaRegistry; use restate_util_time::DurationExt; +use restate_wal_protocol::v2::{Envelope, Raw}; use crate::rest_api::{MAX_ADMIN_API_VERSION, MIN_ADMIN_API_VERSION}; use crate::schema_registry_integration::{MetadataService, TelemetryClient}; @@ -47,7 +47,7 @@ pub struct BuildError(#[from] restate_service_client::BuildError); pub struct AdminService { listeners: Listeners, - ingestion_client: IngestionClient, + ingestion_client: IngestionClient>, schema_registry: SchemaRegistry, serdes_client: SerdesClient, invocation_client: Invocations, @@ -65,7 +65,7 @@ where pub fn new( listeners: Listeners, metadata_writer: MetadataWriter, - ingestion_client: IngestionClient, + ingestion_client: IngestionClient>, invocation_client: Invocations, serdes_client: SerdesClient, service_discovery: ServiceDiscovery, diff --git a/crates/admin/src/state.rs b/crates/admin/src/state.rs index 568ae5590d..f9930e4896 100644 --- a/crates/admin/src/state.rs +++ b/crates/admin/src/state.rs @@ -8,6 +8,8 @@ // the Business Source License, use of this software will be governed // by the Apache License, Version 2.0. +use std::sync::Arc; + use restate_core::network::TransportConnect; use restate_ingestion_client::IngestionClient; use restate_limiter::rule_book::RuleBookObserver; @@ -15,15 +17,14 @@ use restate_metadata_store::MetadataStoreClient; use restate_service_protocol_v4::serdes::SerdesClient; use restate_storage_query_datafusion::context::QueryContext; use restate_types::schema::registry::SchemaRegistry; -use restate_wal_protocol::Envelope; -use std::sync::Arc; +use restate_wal_protocol::v2::{Envelope, Raw}; #[derive(Clone, derive_builder::Builder)] pub struct AdminServiceState { pub schema_registry: SchemaRegistry, pub serdes_client: SerdesClient, pub invocation_client: Invocations, - pub ingestion_client: IngestionClient, + pub ingestion_client: IngestionClient>, /// Used by handlers that mutate cluster-global metadata-store keys /// directly (e.g. the rule book) via `read_modify_write`. pub metadata_store_client: MetadataStoreClient, @@ -41,7 +42,7 @@ where schema_registry: SchemaRegistry, serdes_client: SerdesClient, invocation_client: Invocations, - ingestion_client: IngestionClient, + ingestion_client: IngestionClient>, metadata_store_client: MetadataStoreClient, query_context: Option, rule_book_observer: Option>, diff --git a/crates/ingestion-client/src/client.rs b/crates/ingestion-client/src/client.rs index e92a784d1e..ca7b850662 100644 --- a/crates/ingestion-client/src/client.rs +++ b/crates/ingestion-client/src/client.rs @@ -21,7 +21,7 @@ use restate_core::{ use restate_types::{ identifiers::PartitionKey, live::Live, - logs::{HasRecordKeys, Keys}, + logs::{BodyWithKeys, HasRecordKeys, Keys}, net::ingest::IngestRecord, partitions::{FindPartition, PartitionTable, PartitionTableError}, storage::{StorageCodec, StorageEncode}, @@ -271,6 +271,16 @@ impl InputRecord { } } +impl From> for InputRecord +where + T: StorageEncode, +{ + fn from(value: BodyWithKeys) -> Self { + let (keys, record) = value.split(); + Self { keys, record } + } +} + #[cfg(test)] mod test { use std::{num::NonZeroUsize, time::Duration}; diff --git a/crates/ingress-kafka/src/builder.rs b/crates/ingress-kafka/src/builder.rs index fa4a65561d..44efb67428 100644 --- a/crates/ingress-kafka/src/builder.rs +++ b/crates/ingress-kafka/src/builder.rs @@ -18,21 +18,20 @@ use opentelemetry::trace::{Span, SpanContext, TraceContextExt}; use opentelemetry_sdk::propagation::TraceContextPropagator; use rdkafka::Message; use rdkafka::message::BorrowedMessage; -use tracing::{info_span, trace}; - use rdkafka::message::Headers; +use tracing::{info_span, trace}; -use restate_storage_api::deduplication_table::DedupInformation; use restate_types::Scope; -use restate_types::identifiers::{InvocationId, WithPartitionKey, partitioner}; +use restate_types::identifiers::{InvocationId, partitioner}; use restate_types::invocation::{Header, InvocationTarget, ServiceInvocation, SpanRelation}; use restate_types::limit_key::LimitKey; use restate_types::live::Live; use restate_types::schema::Schema; use restate_types::schema::invocation_target::{DeploymentStatus, InvocationTargetResolver}; use restate_types::schema::subscriptions::{EventInvocationTargetTemplate, Sink, Subscription}; +use restate_types::sharding::{PartitionKey, WithPartitionKey}; use restate_util_string::{ReString, RestateString, RestrictedValueError}; -use restate_wal_protocol::{Command, Destination, Envelope, Source}; +use restate_wal_protocol::v2::{Dedup, Envelope, commands}; use crate::Error; @@ -62,7 +61,7 @@ impl EnvelopeBuilder { producer_id: u128, consumer_group_id: &str, msg: BorrowedMessage<'_>, - ) -> Result { + ) -> Result<(PartitionKey, Envelope), Error> { // Prepare ingress span let ingress_span = info_span!( "kafka_ingress_consume", @@ -106,7 +105,11 @@ impl EnvelopeBuilder { (None, LimitKey::None) }; - let dedup = DedupInformation::producer(producer_id, msg.offset() as u64); + let dedup = Dedup::Arbitrary { + prefix: None, + producer_id: producer_id.into(), + seq: msg.offset() as u64, + }; let invocation = InvocationBuilder::create( &self.subscription, @@ -130,23 +133,11 @@ impl EnvelopeBuilder { cause, })?; - Ok(self.wrap_service_invocation_in_envelope(invocation, dedup)) - } - - fn wrap_service_invocation_in_envelope( - &self, - service_invocation: Box, - dedup_information: DedupInformation, - ) -> Envelope { - let header = restate_wal_protocol::Header { - source: Source::Ingress {}, - dest: Destination::Processor { - partition_key: service_invocation.partition_key(), - dedup: Some(dedup_information), - }, - }; - - Envelope::new(header, Command::Invoke(service_invocation)) + let partition_key = invocation.partition_key(); + Ok(( + partition_key, + Envelope::new(dedup, commands::InvokeCommand::from(invocation)), + )) } fn generate_events_attributes(msg: &impl Message, subscription_id: &str) -> Vec
{ @@ -225,7 +216,7 @@ impl InvocationBuilder { topic: &str, partition: i32, offset: i64, - ) -> Result, anyhow::Error> { + ) -> Result { let Sink::Invocation { event_invocation_target_template, } = subscription.sink(); @@ -325,11 +316,11 @@ impl InvocationBuilder { ); // Finally generate service invocation - let mut service_invocation = Box::new(ServiceInvocation::initialize( + let mut service_invocation = ServiceInvocation::initialize( invocation_id, invocation_target, restate_types::invocation::Source::Subscription(subscription.id()), - )); + ); service_invocation.with_related_span(SpanRelation::parent(ingress_span_context)); service_invocation.argument = payload; service_invocation.headers = headers; diff --git a/crates/ingress-kafka/src/consumer_task.rs b/crates/ingress-kafka/src/consumer_task.rs index 7d51cce86b..8ef7a60a98 100644 --- a/crates/ingress-kafka/src/consumer_task.rs +++ b/crates/ingress-kafka/src/consumer_task.rs @@ -26,19 +26,20 @@ use rdkafka::error::KafkaError; use rdkafka::topic_partition_list::TopicPartitionListElem; use rdkafka::types::RDKafkaErrorCode; use rdkafka::{ClientConfig, ClientContext, Message, Statistics}; +use tokio::sync::{mpsc, oneshot}; +use tracing::{debug, instrument, trace, warn}; use restate_core::network::{NetworkSender, Swimlane, TransportConnect}; use restate_core::{Metadata, TaskCenter, TaskHandle, TaskKind, task_center}; use restate_ingestion_client::{IngestionClient, IngestionError, RecordCommit}; +use restate_types::identifiers::SubscriptionId; use restate_types::identifiers::partitioner::HashPartitioner; -use restate_types::identifiers::{SubscriptionId, WithPartitionKey}; +use restate_types::logs::{BodyWithKeys, Keys}; use restate_types::net::ingest::{DedupSequenceNrQueryRequest, ProducerId, ResponseStatus}; use restate_types::partitions::FindPartition; use restate_types::retries::RetryPolicy; use restate_types::schema::subscriptions::{EventInvocationTargetTemplate, Sink}; -use restate_wal_protocol::Envelope; -use tokio::sync::{mpsc, oneshot}; -use tracing::{debug, instrument, trace, warn}; +use restate_wal_protocol::v2::{Envelope, Raw}; use crate::Error; use crate::builder::EnvelopeBuilder; @@ -50,7 +51,7 @@ type MessageConsumer = StreamConsumer>; pub struct ConsumerTask { client_config: ClientConfig, topics: Vec, - ingestion: IngestionClient, + ingestion: IngestionClient>, builder: EnvelopeBuilder, } @@ -61,7 +62,7 @@ where pub fn new( client_config: ClientConfig, topics: Vec, - ingestion: IngestionClient, + ingestion: IngestionClient>, builder: EnvelopeBuilder, ) -> Self { Self { @@ -167,7 +168,7 @@ struct RebalanceContext { consumer: OnceLock>>, topic_partition_tasks: parking_lot::Mutex>, failures_tx: mpsc::UnboundedSender, - ingestion: IngestionClient, + ingestion: IngestionClient>, builder: EnvelopeBuilder, consumer_group_id: String, } @@ -324,7 +325,7 @@ where T: TransportConnect, C: ConsumerContext, { - ingestion: IngestionClient, + ingestion: IngestionClient>, builder: EnvelopeBuilder, topic_partition: TopicPartition, topic_partition_consumer: StreamPartitionQueue, @@ -339,7 +340,7 @@ where C: ConsumerContext, { fn new( - ingestion: IngestionClient, + ingestion: IngestionClient>, builder: EnvelopeBuilder, topic_partition: TopicPartition, topic_partition_consumer: StreamPartitionQueue, @@ -530,11 +531,11 @@ where "Ingesting kafka message" ); - let envelope = self.builder.build(producer_id, &self.consumer_group_id, msg)?; + let (partition_key, envelope) = self.builder.build(producer_id, &self.consumer_group_id, msg)?; let commit_token = self .ingestion - .ingest(envelope.partition_key(), envelope) + .ingest(partition_key, BodyWithKeys::new(envelope.into_raw(), Keys::Single(partition_key))) .await? .map(|_| offset); diff --git a/crates/ingress-kafka/src/subscription_controller.rs b/crates/ingress-kafka/src/subscription_controller.rs index 0cd20771d9..6de0175c68 100644 --- a/crates/ingress-kafka/src/subscription_controller.rs +++ b/crates/ingress-kafka/src/subscription_controller.rs @@ -11,7 +11,6 @@ use std::collections::{HashMap, HashSet}; use std::time::Duration; -use restate_wal_protocol::Envelope; use tokio::sync::mpsc; use tracing::{error, warn}; @@ -24,6 +23,7 @@ use restate_types::retries::RetryPolicy; use restate_types::schema::Schema; use restate_types::schema::kafka::KafkaCluster; use restate_types::schema::subscriptions::{Source, Subscription}; +use restate_wal_protocol::v2::{Envelope, Raw}; use super::*; use crate::builder::EnvelopeBuilder; @@ -32,7 +32,7 @@ use crate::subscription_controller::task_orchestrator::TaskOrchestrator; // For simplicity of the current implementation, this currently lives in this module // In future versions, we should either pull this out in a separate process, or generify it and move it to the worker, or an ad-hoc module pub struct Service { - ingestion: IngestionClient, + ingestion: IngestionClient>, schema: Live, commands_tx: SubscriptionCommandSender, @@ -43,7 +43,7 @@ impl Service where T: TransportConnect, { - pub fn new(ingestion: IngestionClient, schema: Live) -> Self { + pub fn new(ingestion: IngestionClient>, schema: Live) -> Self { metric_definitions::describe_metrics(); let (commands_tx, commands_rx) = mpsc::channel(10); diff --git a/crates/node/src/roles/admin.rs b/crates/node/src/roles/admin.rs index f90a8ee216..4ea6ad4940 100644 --- a/crates/node/src/roles/admin.rs +++ b/crates/node/src/roles/admin.rs @@ -45,7 +45,8 @@ use restate_types::partition_table::PartitionTable; use restate_types::partitions::state::PartitionReplicaSetStates; use restate_types::protobuf::common::AdminStatus; use restate_types::retries::RetryPolicy; -use restate_wal_protocol::Envelope; +use restate_wal_protocol::v2::Envelope; +use restate_wal_protocol::v2::Raw; use restate_worker_api::PartitionProcessorInvocationClient; #[derive(Debug, thiserror::Error, CodedError)] @@ -82,7 +83,7 @@ impl AdminRole { pub async fn create( health_status: HealthStatus, bifrost: Bifrost, - ingestion_client: IngestionClient, + ingestion_client: IngestionClient>, updateable_config: Live, partition_routing: PartitionRouting, partition_table: Live, diff --git a/crates/node/src/roles/worker.rs b/crates/node/src/roles/worker.rs index 48e174605a..b46474d0d6 100644 --- a/crates/node/src/roles/worker.rs +++ b/crates/node/src/roles/worker.rs @@ -24,7 +24,8 @@ use restate_storage_query_datafusion::remote_query_scanner_manager::RemoteScanne use restate_types::health::HealthStatus; use restate_types::partitions::state::PartitionReplicaSetStates; use restate_types::protobuf::common::WorkerStatus; -use restate_wal_protocol::Envelope; +use restate_wal_protocol::v2::Envelope; +use restate_wal_protocol::v2::Raw; use restate_worker::{RuleBookCacheHandle, Worker}; use restate_worker_api::ProcessorsManagerHandle; @@ -54,7 +55,7 @@ where partition_store_manager: Arc, networking: Networking, bifrost: Bifrost, - ingestion_client: IngestionClient, + ingestion_client: IngestionClient>, metadata_writer: MetadataWriter, remote_scanner_manager: RemoteScannerManager, ) -> Result { diff --git a/crates/types/src/cluster_marker.rs b/crates/types/src/cluster_marker.rs index d7b1b680ff..54d02e9ea6 100644 --- a/crates/types/src/cluster_marker.rs +++ b/crates/types/src/cluster_marker.rs @@ -33,6 +33,7 @@ const TMP_CLUSTER_MARKER_FILE_NAME: &str = ".tmp-cluster-marker"; /// /// To reduce the risk of unexpected incompatibility issues, the minimum tracks the /// version we allow restate to downgrade to according to the compatibility policy. +// todo(azmy): Bump to 1.7 before Restate v1.8.0 because of envelope v2 const COMPATIBILITY_INFORMATION: CompatibilityInformation = CompatibilityInformation::new( SemanticRestateVersion::new(1, 6, 0), SemanticRestateVersion::new(1, 6, 0), diff --git a/crates/wal-protocol/src/lib.rs b/crates/wal-protocol/src/lib.rs index 7b27dd2e6a..eb658e86c8 100644 --- a/crates/wal-protocol/src/lib.rs +++ b/crates/wal-protocol/src/lib.rs @@ -15,4 +15,5 @@ pub mod v1; pub mod v2; pub mod vqueues; +// Drop v1 in v1.9 pub use v1::{Command, Destination, Envelope, Header, Source}; diff --git a/crates/wal-protocol/src/v2/commands.rs b/crates/wal-protocol/src/v2/commands.rs index 86e53adbba..725a8dfb8e 100644 --- a/crates/wal-protocol/src/v2/commands.rs +++ b/crates/wal-protocol/src/v2/commands.rs @@ -135,15 +135,6 @@ impl HasRecordKeys for InvokeCommand { pub struct TruncateOutboxCommand { #[bilrost(1)] pub index: MessageIndex, - - #[bilrost(2)] - pub partition_key_range: Keys, -} - -impl HasRecordKeys for TruncateOutboxCommand { - fn record_keys(&self) -> Keys { - self.partition_key_range.clone() - } } bilrost_storage_encode_decode!(TruncateOutboxCommand); diff --git a/crates/wal-protocol/src/v2/compatibility.rs b/crates/wal-protocol/src/v2/compatibility.rs index d01bc87ccf..03e7fc1ab8 100644 --- a/crates/wal-protocol/src/v2/compatibility.rs +++ b/crates/wal-protocol/src/v2/compatibility.rs @@ -21,6 +21,11 @@ use crate::{ v2::{self, Envelope, commands::TruncateOutboxCommand}, }; +// TODO: Keep for backward compatibility only. Do not extend +// v1 Commands with new commands. New commands should be +// added to v2 only. +// +// Drop in v1.9 impl TryFrom for v2::Envelope { type Error = anyhow::Error; @@ -121,17 +126,9 @@ impl TryFrom for v2::Envelope { v1::Command::Timer(payload) => { Envelope::new(dedup, commands::TimerCommand::from(payload)).into_raw() } - v1::Command::TruncateOutbox(payload) => Envelope::new( - dedup, - TruncateOutboxCommand { - index: payload, - // this actually should be a key-range but v1 unfortunately - // only hold the "start" of the range. - // will be fixed in v2 - partition_key_range: Keys::Single(partition_key), - }, - ) - .into_raw(), + v1::Command::TruncateOutbox(payload) => { + Envelope::new(dedup, TruncateOutboxCommand { index: payload }).into_raw() + } v1::Command::UpdatePartitionDurability(payload) => { Envelope::new(dedup, payload).into_raw() } diff --git a/crates/worker/src/lib.rs b/crates/worker/src/lib.rs index 7c477a630b..47c2689ef4 100644 --- a/crates/worker/src/lib.rs +++ b/crates/worker/src/lib.rs @@ -26,10 +26,6 @@ mod subscription_integration; use std::sync::Arc; use codederror::CodedError; -use restate_core::network::Swimlane; -use restate_ingestion_client::SessionOptions; -use restate_types::net::connect_opts::GrpcConnectionOptions; -use restate_wal_protocol::Envelope; use tracing::info; use restate_bifrost::Bifrost; @@ -37,11 +33,13 @@ use restate_core::MetadataKind; use restate_core::cancellation_watcher; use restate_core::network::MessageRouterBuilder; use restate_core::network::Networking; +use restate_core::network::Swimlane; use restate_core::network::TransportConnect; use restate_core::partitions::PartitionRouting; use restate_core::{Metadata, TaskKind}; use restate_core::{MetadataWriter, TaskCenter}; use restate_ingestion_client::IngestionClient; +use restate_ingestion_client::SessionOptions; use restate_ingress_kafka::Service as IngressKafkaService; use restate_partition_store::PartitionStoreManager; use restate_partition_store::snapshots::SnapshotRepository; @@ -51,11 +49,13 @@ use restate_types::Version; use restate_types::Versioned; use restate_types::config::Configuration; use restate_types::health::HealthStatus; +use restate_types::net::connect_opts::GrpcConnectionOptions; use restate_types::partitions::state::PartitionReplicaSetStates; use restate_types::protobuf::common::WorkerStatus; use restate_types::schema::Redaction; use restate_types::schema::kafka::KafkaClusterResolver; use restate_types::schema::subscriptions::SubscriptionResolver; +use restate_wal_protocol::v2::{Envelope, Raw}; use restate_worker_api::ProcessorsManagerHandle; use crate::partition_processor_manager::PartitionProcessorManager; @@ -113,7 +113,7 @@ where partition_store_manager: Arc, networking: Networking, bifrost: Bifrost, - ingestion_client: IngestionClient, + ingestion_client: IngestionClient>, router_builder: &mut MessageRouterBuilder, metadata_writer: MetadataWriter, remote_scanner_manager: RemoteScannerManager, @@ -131,7 +131,7 @@ where let schema = metadata.updateable_schema(); // ingress_kafka - let ingress_kafka = IngressKafkaService::new(ingestion_client.clone(), schema.clone()); + let ingress_kafka = IngressKafkaService::new(ingestion_client, schema.clone()); let subscription_controller_handle = SubscriptionControllerHandle::new(ingress_kafka.create_command_sender()); diff --git a/crates/worker/src/partition/leadership/leader_state.rs b/crates/worker/src/partition/leadership/leader_state.rs index 66b5259bab..00993462a3 100644 --- a/crates/worker/src/partition/leadership/leader_state.rs +++ b/crates/worker/src/partition/leadership/leader_state.rs @@ -16,7 +16,6 @@ use std::pin::Pin; use std::sync::Arc; use std::task::{Context, Poll, ready}; -use bytes::BytesMut; use futures::future::OptionFuture; use futures::stream::FuturesUnordered; use futures::{FutureExt, StreamExt, stream}; @@ -33,12 +32,11 @@ use restate_limiter::RuleBook; use restate_partition_store::PartitionDb; use restate_storage_api::vqueue_table::scheduler::SchedulerDecisionsCommand; use restate_types::identifiers::{ - InvocationId, LeaderEpoch, PartitionId, PartitionKey, PartitionProcessorRpcRequestId, - WithPartitionKey, + InvocationId, LeaderEpoch, PartitionId, PartitionProcessorRpcRequestId, WithPartitionKey, }; use restate_types::invocation::client::{InvocationOutput, SubmittedInvocationNotification}; use restate_types::invocation::{FencingToken, PurgeInvocationRequest}; -use restate_types::logs::Keys; +use restate_types::logs::{BodyWithKeys, Keys}; use restate_types::net::ingest::IngestRecord; use restate_types::net::partition_processor::{ PartitionProcessorRpcError, PartitionProcessorRpcResponse, @@ -48,9 +46,9 @@ use restate_types::{RESTATE_VERSION_1_7_0, SemanticRestateVersion, Version, Vers use restate_vqueues::VQueueEvent; use restate_vqueues::scheduler::Decisions; use restate_vqueues::{SchedulerService, VQueuesMeta}; -use restate_wal_protocol::Command; use restate_wal_protocol::control::UpsertSchemaCommand; -use restate_wal_protocol::v1::UpsertRuleBookCommandWrapper; +use restate_wal_protocol::invocation::PauseInvocationCommand; +use restate_wal_protocol::v2::{Command, CommandWithKeys, commands}; use restate_worker_api::invoker::InvokerHandle; use restate_worker_api::resources::ReservedResources; use restate_worker_api::{SchedulerStatusEntry, UserLimitCounterEntry}; @@ -359,7 +357,6 @@ impl LeaderState { &mut self, action_effects: impl IntoIterator, ) -> Result<(), Error> { - let mut arena = BytesMut::with_capacity(128 * 1024); for effect in action_effects { match effect { ActionEffect::Scheduler(decisions) => { @@ -379,17 +376,11 @@ impl LeaderState { .into_iter() // one command per partition key .map(|(partition_key, group)| { - let decisions = SchedulerDecisionsCommand { - qids: group.collect(), - }; - - arena.reserve(decisions.encoded_len()); - // safe to unwrap because we reserved enough space - decisions.bilrost_encode(&mut arena).unwrap(); - - ( - partition_key, - Command::VQSchedulerDecisions(arena.split().freeze()), + BodyWithKeys::new( + SchedulerDecisionsCommand { + qids: group.collect(), + }, + Keys::Single(partition_key), ) // an action goes into command, and resources are popped }) @@ -404,10 +395,10 @@ impl LeaderState { // based on configuration, whether to consider partition-local durability in // the replica-set as a sufficient source of durability, or only snapshots. self.self_proposer - .self_propose( - self.partition_key_range.start(), - Command::UpdatePartitionDurability(partition_durability), - ) + .self_propose(BodyWithKeys::new( + partition_durability, + Keys::RangeInclusive(self.partition_key_range.into()), + )) .await?; } ActionEffect::Invoker(fenced) => { @@ -424,10 +415,7 @@ impl LeaderState { self.fencing_tokens.remove(&invocation_id); } self.self_proposer - .self_propose( - invocation_id.partition_key(), - Command::InvokerEffect(fenced.effect), - ) + .self_propose(commands::InvokerEffectCommand::from(*fenced.effect)) .await?; } else { debug!( @@ -440,70 +428,58 @@ impl LeaderState { // todo: Until we support partition splits we need to get rid of outboxes or introduce partition // specific destination messages that are identified by a partition_id self.self_proposer - .self_propose( - self.partition_key_range.start(), - Command::TruncateOutbox(outbox_truncation.index()), - ) + .self_propose(BodyWithKeys::new( + commands::TruncateOutboxCommand { + index: outbox_truncation.index(), + }, + Keys::RangeInclusive(self.partition_key_range.into()), + )) .await?; } ActionEffect::Timer(timer) => { self.self_proposer - .self_propose(timer.invocation_id().partition_key(), Command::Timer(timer)) + .self_propose(commands::TimerCommand::from(timer)) .await?; } - ActionEffect::Cleaner(effect) => { - let (invocation_id, cmd) = match effect { - CleanerEffect::PurgeJournal(invocation_id) => ( + ActionEffect::Cleaner(effect) => match effect { + CleanerEffect::PurgeJournal(invocation_id) => { + let record = commands::PurgeJournalCommand::from(PurgeInvocationRequest { invocation_id, - Command::PurgeJournal(PurgeInvocationRequest { - invocation_id, - response_sink: None, - }), - ), - CleanerEffect::PurgeInvocation(invocation_id) => ( - invocation_id, - Command::PurgeInvocation(PurgeInvocationRequest { + response_sink: None, + }); + + self.self_proposer.self_propose(record).await?; + } + CleanerEffect::PurgeInvocation(invocation_id) => { + let record = + commands::PurgeInvocationCommand::from(PurgeInvocationRequest { invocation_id, response_sink: None, - }), - ), - }; + }); - self.self_proposer - .self_propose(invocation_id.partition_key(), cmd) - .await?; - } + self.self_proposer.self_propose(record).await?; + } + }, ActionEffect::UpsertSchema(schema) => { if SemanticRestateVersion::current() .is_equal_or_newer_than(&RESTATE_VERSION_1_7_0) { self.self_proposer - .self_propose( - self.partition_key_range.start(), - Command::UpsertSchema(UpsertSchemaCommand { - partition_key_range: Keys::RangeInclusive( - self.partition_key_range.into(), - ), - schema, - }), - ) + .self_propose(UpsertSchemaCommand { + partition_key_range: Keys::RangeInclusive( + self.partition_key_range.into(), + ), + schema, + }) .await?; } } ActionEffect::UpsertRuleBook(rule_book) => { - let cmd = restate_wal_protocol::control::UpsertRuleBookCommand { rule_book }; - arena.reserve(cmd.encoded_len()); - // safe to unwrap because we reserved enough space - cmd.bilrost_encode(&mut arena).unwrap(); - self.self_proposer - .self_propose( - self.partition_key_range.start(), - Command::UpsertRuleBook(UpsertRuleBookCommandWrapper { - partition_key_range: self.partition_key_range, - command: arena.split().freeze(), - }), - ) + .self_propose(BodyWithKeys::new( + restate_wal_protocol::control::UpsertRuleBookCommand { rule_book }, + Keys::RangeInclusive(self.partition_key_range.into()), + )) .await?; } ActionEffect::AwaitingRpcSelfProposeDone => { @@ -515,14 +491,13 @@ impl LeaderState { Ok(()) } - pub async fn handle_rpc_proposal_command( + pub async fn handle_rpc_proposal_command( &mut self, request_id: PartitionProcessorRpcRequestId, reciprocal: Reciprocal< Oneshot>, >, - partition_key: PartitionKey, - cmd: Command, + cmd: impl CommandWithKeys, ) { match self.awaiting_rpc_actions.entry(request_id) { Entry::Occupied(mut o) => { @@ -536,7 +511,7 @@ impl LeaderState { } Entry::Vacant(v) => { // In this case, no one proposed this command yet, let's try to propose it - if let Err(e) = self.self_proposer.self_propose(partition_key, cmd).await { + if let Err(e) = self.self_proposer.self_propose(cmd).await { reciprocal.send(Err(PartitionProcessorRpcError::Internal(e.to_string()))); } else { v.insert(reciprocal); @@ -562,8 +537,7 @@ impl LeaderState { reciprocal: Reciprocal< Oneshot>, >, - invocation_id: InvocationId, - cmd: Command, + cmd: PauseInvocationCommand, ) { match self.awaiting_rpc_actions.entry(request_id) { Entry::Occupied(mut o) => { @@ -577,9 +551,13 @@ impl LeaderState { ))); } Entry::Vacant(v) => { + let invocation_id = cmd.invocation_id; if let Err(e) = self .self_proposer - .self_propose(invocation_id.partition_key(), cmd) + .self_propose(BodyWithKeys::new( + cmd, + Keys::Single(invocation_id.partition_key()), + )) .await { reciprocal.send(Err(PartitionProcessorRpcError::Internal(e.to_string()))); @@ -597,18 +575,13 @@ impl LeaderState { /// Records appended this way are never filtered by the dedup mechanism during leadership /// transitions, making this safe for fire-and-forget ingress commands (signals, invocation /// responses). - pub async fn append_and_respond_asynchronously( + pub async fn append_and_respond_asynchronously( &mut self, - partition_key: PartitionKey, - cmd: Command, + cmd: impl CommandWithKeys, reciprocal: RpcReciprocal, success_response: PartitionProcessorRpcResponse, ) { - match self - .self_proposer - .append_with_notification(partition_key, cmd) - .await - { + match self.self_proposer.append_with_notification(cmd).await { Ok(commit_token) => { self.awaiting_rpc_self_propose.push(SelfAppendFuture::new( commit_token, diff --git a/crates/worker/src/partition/leadership/mod.rs b/crates/worker/src/partition/leadership/mod.rs index b60ddc2972..7dd9b61910 100644 --- a/crates/worker/src/partition/leadership/mod.rs +++ b/crates/worker/src/partition/leadership/mod.rs @@ -49,11 +49,11 @@ use restate_timer::TokioClock; use restate_types::cluster::cluster_state::RunMode; use restate_types::config::Configuration; use restate_types::errors::GenericError; +use restate_types::identifiers::PartitionProcessorRpcRequestId; use restate_types::identifiers::{InvocationId, LeaderEpoch, PartitionId}; -use restate_types::identifiers::{PartitionKey, PartitionProcessorRpcRequestId}; use restate_types::invocation::FencingToken; use restate_types::live::LiveLoadExt; -use restate_types::logs::Keys; +use restate_types::logs::{HasRecordKeys, Keys}; use restate_types::message::MessageIndex; use restate_types::net::ingest::IngestRecord; use restate_types::net::partition_processor::{ @@ -70,8 +70,9 @@ use restate_vqueues::{ResourceManager, SchedulerService, VQueuesMeta, VQueuesMet use restate_wal_protocol::control::{ AnnounceLeaderCommand, UpdatePartitionDurabilityCommand, VersionBarrierCommand, }; +use restate_wal_protocol::invocation::PauseInvocationCommand; use restate_wal_protocol::timer::TimerKeyValue; -use restate_wal_protocol::{Command, Envelope}; +use restate_wal_protocol::v2::{Command, CommandWithKeys, Envelope, Raw}; use restate_worker_api::invoker::InvokerHandle; use restate_worker_api::invoker::capacity::InvokerCapacity; use restate_worker_api::{ @@ -182,7 +183,7 @@ pub(crate) struct LeadershipState { last_seen_leader_epoch: Option, partition: Arc, - ingestion_client: IngestionClient, + ingestion_client: IngestionClient>, invoker_capacity: InvokerCapacity, bifrost: Bifrost, trim_queue: TrimQueue, @@ -199,7 +200,7 @@ where pub(crate) fn new( partition: Arc, invoker_capacity: InvokerCapacity, - ingestion_client: IngestionClient, + ingestion_client: IngestionClient>, bifrost: Bifrost, last_seen_leader_epoch: Option, trim_queue: TrimQueue, @@ -268,14 +269,14 @@ where ) -> Result<(), Error> { let leader_epoch = leadership_info.leader_epoch; - let announce_leader = Command::AnnounceLeader(Box::new(AnnounceLeaderCommand { + let announce_leader = AnnounceLeaderCommand { node_id: my_node_id(), leader_epoch, epoch_version: Some(leadership_info.version), partition_key_range: self.partition.key_range, current_config: Some(leadership_info.current_config), next_config: leadership_info.next_config, - })); + }; let mut self_proposer = SelfProposer::new( self.partition.log_id(), @@ -283,9 +284,7 @@ where &self.bifrost, )?; - self_proposer - .self_propose(self.partition.key_range.start(), announce_leader) - .await?; + self_proposer.self_propose(announce_leader).await?; self.state = State::Candidate { leader_epoch, @@ -508,7 +507,7 @@ where let (shuffle_tx, shuffle_rx) = mpsc::channel(config.worker.internal_queue_length()); let shuffle = Shuffle::new( - ShuffleMetadata::new(self.partition.partition_id, *leader_epoch), + ShuffleMetadata::new(self.partition.partition_id), OutboxReader::from(partition_store.clone()), shuffle_tx, config.worker.internal_queue_length(), @@ -597,17 +596,12 @@ where .clone(); self_proposer - .self_propose( - self.partition.key_range.start(), - Command::VersionBarrier(VersionBarrierCommand { - version: barrier_version, - partition_key_range: Keys::RangeInclusive( - self.partition.key_range.into(), - ), - human_reason: Some("Apply state-machine feature changes".to_owned()), - feature_changes: feature_changes.into_iter().map(|c| c.id()).collect(), - }), - ) + .self_propose(VersionBarrierCommand { + version: barrier_version, + partition_key_range: Keys::RangeInclusive(self.partition.key_range.into()), + human_reason: Some("Apply state-machine feature changes".to_owned()), + feature_changes: feature_changes.into_iter().map(|c| c.id()).collect(), + }) .await?; } let last_reported_durable_lsn = partition_store @@ -803,14 +797,13 @@ impl LeadershipState { } } - pub async fn handle_rpc_proposal_command( + pub async fn handle_rpc_proposal_command( &mut self, request_id: PartitionProcessorRpcRequestId, reciprocal: Reciprocal< Oneshot>, >, - partition_key: PartitionKey, - cmd: Command, + cmd: impl CommandWithKeys, ) { match &mut self.state { State::Follower | State::Candidate { .. } => { @@ -821,7 +814,7 @@ impl LeadershipState { } State::Leader(leader_state) => { leader_state - .handle_rpc_proposal_command(request_id, reciprocal, partition_key, cmd) + .handle_rpc_proposal_command(request_id, reciprocal, cmd) .await; } } @@ -833,8 +826,7 @@ impl LeadershipState { reciprocal: Reciprocal< Oneshot>, >, - invocation_id: InvocationId, - cmd: Command, + cmd: PauseInvocationCommand, ) { match &mut self.state { State::Follower | State::Candidate { .. } => { @@ -845,17 +837,16 @@ impl LeadershipState { } State::Leader(leader_state) => { leader_state - .propose_pause_and_fence(request_id, reciprocal, invocation_id, cmd) + .propose_pause_and_fence(request_id, reciprocal, cmd) .await; } } } /// Append a command to Bifrost without dedup information, responding on Bifrost commit. - pub async fn append_and_respond_asynchronously( + pub async fn append_and_respond_asynchronously( &mut self, - partition_key: PartitionKey, - cmd: Command, + cmd: C, reciprocal: Reciprocal< Oneshot>, >, @@ -867,12 +858,7 @@ impl LeadershipState { )), State::Leader(leader_state) => { leader_state - .append_and_respond_asynchronously( - partition_key, - cmd, - reciprocal, - success_response, - ) + .append_and_respond_asynchronously(cmd, reciprocal, success_response) .await; } } @@ -949,7 +935,6 @@ mod tests { use crate::partition::state_machine::StateMachineFeatures; use crate::partition_processor_manager::PartitionLeaderHandlesRegistry; use crate::rule_book_cache::RuleBookCacheHandle; - use assert2::let_assert; use restate_bifrost::Bifrost; use restate_core::partitions::PartitionRouting; use restate_core::{TaskCenter, TestCoreEnv}; @@ -965,8 +950,7 @@ mod tests { use restate_types::sharding::KeyRange; use restate_types::{GenerationalNodeId, SemanticRestateVersion, Version}; use restate_vqueues::VQueuesMetaCache; - use restate_wal_protocol::Command; - use restate_wal_protocol::Envelope; + use restate_wal_protocol::v2::{CommandKind, Envelope, Raw, commands}; use restate_worker_api::invoker::capacity::InvokerCapacity; use std::num::NonZeroUsize; use std::sync::Arc; @@ -1046,9 +1030,11 @@ mod tests { .await .unwrap()?; - let envelope = record.try_decode::().unwrap()?; + let envelope = record.try_decode::>().unwrap()?; - let_assert!(Command::AnnounceLeader(announce_leader) = envelope.command); + assert_eq!(CommandKind::AnnounceLeader, envelope.kind()); + let envelope = envelope.into_typed::(); + let announce_leader = envelope.into_inner().unwrap(); assert_eq!(announce_leader.node_id, NODE_ID); assert_eq!(announce_leader.leader_epoch, leader_epoch); assert_eq!(announce_leader.partition_key_range, PARTITION_KEY_RANGE); diff --git a/crates/worker/src/partition/leadership/self_proposer.rs b/crates/worker/src/partition/leadership/self_proposer.rs index 15d1a0de8e..7724a2db43 100644 --- a/crates/worker/src/partition/leadership/self_proposer.rs +++ b/crates/worker/src/partition/leadership/self_proposer.rs @@ -8,16 +8,16 @@ // the Business Source License, use of this software will be governed // by the Apache License, Version 2.0. -use std::sync::Arc; - use futures::never::Never; use restate_bifrost::{Bifrost, CommitToken, EnqueueError, ErrorRecoveryStrategy, InputRecord}; -use restate_storage_api::deduplication_table::{DedupInformation, EpochSequenceNumber}; +use restate_storage_api::deduplication_table::EpochSequenceNumber; use restate_types::{ - identifiers::PartitionKey, logs::LogId, net::ingest::IngestRecord, time::NanosSinceEpoch, + logs::{BodyWithKeys, LogId}, + net::ingest::IngestRecord, + time::NanosSinceEpoch, }; -use restate_wal_protocol::{Command, Destination, Envelope, Header, Source}; +use restate_wal_protocol::v2::{Command, CommandWithKeys, Dedup, Envelope, Raw}; use crate::partition::leadership::Error; @@ -34,7 +34,7 @@ static BIFROST_APPENDER_TASK: &str = "bifrost-appender"; pub struct SelfProposer { epoch_sequence_number: EpochSequenceNumber, - bifrost_appender: restate_bifrost::AppenderHandle, + bifrost_appender: restate_bifrost::AppenderHandle>, } impl SelfProposer { @@ -72,32 +72,29 @@ impl SelfProposer { /// /// Note that self_propose_many will return an error if the number of commands is greater than /// the internal channel's max capacity. - pub async fn self_propose_many( - &mut self, - cmds: impl ExactSizeIterator, - ) -> Result<(), Error> { + pub async fn self_propose_many(&mut self, records: I) -> Result<(), Error> + where + I: ExactSizeIterator, + T: CommandWithKeys, + C: Command, + { // allocate a sequence number range for the batch let leader_epoch = self.epoch_sequence_number.leader_epoch; let start_seq = self.epoch_sequence_number.sequence_number; - let end_seq = start_seq + cmds.len() as u64; + let end_seq = start_seq + records.len() as u64; - let envelopes = cmds.enumerate().map(|(idx, (partition_key, cmd))| { - let esn = EpochSequenceNumber { - leader_epoch, - sequence_number: start_seq + idx as u64, - }; - let header = Header { - dest: Destination::Processor { - partition_key, - dedup: Some(DedupInformation::self_proposal(esn)), - }, - source: Source::Processor { - partition_key: Some(partition_key), + let envelopes = records.enumerate().map(|(idx, command)| { + let keys = command.keys(); + let envelope = Envelope::new( + Dedup::SelfProposal { leader_epoch, + seq: start_seq + idx as u64, }, - }; - Arc::new(Envelope::new(header, cmd)) + command.inner(), + ) + .into_raw(); + BodyWithKeys::new(envelope, keys) }); // Only blocks if background append is pushing back (queue full) @@ -116,21 +113,24 @@ impl SelfProposer { Ok(()) } - /// Self-propose a single command to Bifrost, attaching ESN-based dedup information. - pub async fn self_propose( + pub async fn self_propose( &mut self, - partition_key: PartitionKey, - cmd: Command, + command: impl CommandWithKeys, ) -> Result<(), Error> { - let envelope = Envelope::new(self.create_self_propose_header(partition_key), cmd); + let esn = self.next_esn(); + let dedup = Dedup::SelfProposal { + leader_epoch: esn.leader_epoch, + seq: esn.sequence_number, + }; + + let keys = command.keys(); + let envelope = Envelope::new(dedup, command.inner()); - // Only blocks if background append is pushing back (queue full) self.bifrost_appender .sender() - .enqueue(Arc::new(envelope)) + .enqueue(BodyWithKeys::new(envelope.into_raw(), keys)) .await .map_err(|e| Error::SelfProposer(e.to_string()))?; - Ok(()) } @@ -139,27 +139,17 @@ impl SelfProposer { /// Unlike [`Self::self_propose`], this does not attach an epoch sequence number. Records /// appended this way are never filtered by the dedup mechanism during leadership transitions, /// which makes them safe for fire-and-forget ingress commands (signals, invocation responses). - pub async fn append_with_notification( + pub async fn append_with_notification( &mut self, - partition_key: PartitionKey, - cmd: Command, + command: impl CommandWithKeys, ) -> Result { - let header = Header { - dest: Destination::Processor { - partition_key, - dedup: None, - }, - source: Source::Processor { - partition_key: Some(partition_key), - leader_epoch: self.epoch_sequence_number.leader_epoch, - }, - }; - let envelope = Envelope::new(header, cmd); + let keys = command.keys(); + let envelope = Envelope::new(Dedup::None, command.inner()).into_raw(); let commit_token = self .bifrost_appender .sender() - .enqueue_with_notification(Arc::new(envelope)) + .enqueue_with_notification(BodyWithKeys::new(envelope, keys)) .await .map_err(|e| Error::SelfProposer(e.to_string()))?; @@ -210,22 +200,6 @@ impl SelfProposer { sender.notify_committed().await } - fn create_self_propose_header(&mut self, partition_key: PartitionKey) -> Header { - let esn = self.epoch_sequence_number; - self.epoch_sequence_number = self.epoch_sequence_number.next(); - - Header { - dest: Destination::Processor { - partition_key, - dedup: Some(DedupInformation::self_proposal(esn)), - }, - source: Source::Processor { - partition_key: Some(partition_key), - leader_epoch: self.epoch_sequence_number.leader_epoch, - }, - } - } - /// Waits for self proposer to fail. This method will only complete with an error if the self /// proposer has failed. There is no guarantee up to which point the self proposer has finished /// processing the proposed commands. @@ -237,4 +211,10 @@ impl SelfProposer { Err(err) => Error::task_failed(BIFROST_APPENDER_TASK, err), }) } + + fn next_esn(&mut self) -> EpochSequenceNumber { + let esn = self.epoch_sequence_number; + self.epoch_sequence_number = self.epoch_sequence_number.next(); + esn + } } diff --git a/crates/worker/src/partition/mod.rs b/crates/worker/src/partition/mod.rs index 7c7b5e9d79..cb922e1e1e 100644 --- a/crates/worker/src/partition/mod.rs +++ b/crates/worker/src/partition/mod.rs @@ -89,8 +89,8 @@ use restate_vqueues::{VQueuesMeta, VQueuesMetaCache}; use restate_wal_protocol::control::{ AnnounceLeaderCommand, CurrentReplicaSetConfiguration, NextReplicaSetConfiguration, }; -use restate_wal_protocol::v2::commands; -use restate_wal_protocol::{Envelope, v2}; +use restate_wal_protocol::v1; +use restate_wal_protocol::v2::{CommandKind, Envelope, Raw, commands}; use restate_worker_api::invoker::capacity::InvokerCapacity; use restate_worker_api::{LeaderQueryCommand, LeaderQueryReceiver}; @@ -175,7 +175,7 @@ impl PartitionProcessorBuilder { pub async fn build( self, bifrost: Bifrost, - ingestion_client: IngestionClient, + ingestion_client: IngestionClient>, mut partition_store: PartitionStore, replica_set_states: PartitionReplicaSetStates, ) -> Result, state_machine::Error> @@ -400,7 +400,7 @@ struct LsnEnvelope { pub lsn: Lsn, pub keys: Keys, pub created_at: NanosSinceEpoch, - pub envelope: Arc>, + pub envelope: Arc>, } /// OrderedOperations are scheduled operations that @@ -463,11 +463,11 @@ where /// Decode record tries to decode the record first as v2 Envelope, if it failed, /// it decodes as v1 Envelope then converts into v2. - fn decode_record(record: Record) -> Result>, StorageDecodeError> { - let envelope = match record.decode_arc::>() { + fn decode_record(record: Record) -> Result>, StorageDecodeError> { + let envelope = match record.decode_arc::>() { Ok(envelope) => envelope, Err(RecordDecodeError::TypedValueMismatch(v1_envelope)) => { - let v1_envelope: Arc = v1_envelope + let v1_envelope: Arc = v1_envelope .downcast_arc() .map_err(|_| StorageDecodeError::DecodeValue("Type mismatch. Record value in PolyBytes::Typed does not match requested type".into()))?; @@ -476,7 +476,7 @@ where Err(arc) => arc.as_ref().clone(), }; - let envelope: v2::Envelope = v1_envelope + let envelope: Envelope = v1_envelope .try_into() .map_err(|err: anyhow::Error| StorageDecodeError::DecodeValue(err.into()))?; @@ -1056,12 +1056,12 @@ where let envelope = Arc::unwrap_or_clone(record.envelope); match envelope.kind() { - v2::CommandKind::AnnounceLeader => { + CommandKind::AnnounceLeader => { let envelope = envelope.into_typed::(); let announce_leader = envelope.into_inner()?; return Ok(Some(Box::new(announce_leader))); } - v2::CommandKind::UpdatePartitionDurability => { + CommandKind::UpdatePartitionDurability => { let envelope = envelope.into_typed::(); let partition_durability = envelope.into_inner()?; if partition_durability.partition_id != self.partition_store.partition_id() { diff --git a/crates/worker/src/partition/rpc/append_invocation.rs b/crates/worker/src/partition/rpc/append_invocation.rs index 50e8fa0360..9c5f6a03ee 100644 --- a/crates/worker/src/partition/rpc/append_invocation.rs +++ b/crates/worker/src/partition/rpc/append_invocation.rs @@ -8,12 +8,13 @@ // the Business Source License, use of this software will be governed // by the Apache License, Version 2.0. -use super::*; -use restate_types::identifiers::WithPartitionKey; use restate_types::invocation; use restate_types::invocation::{ ServiceInvocation, ServiceInvocationResponseSink, SubmitNotificationSink, }; +use restate_wal_protocol::v2::commands; + +use super::*; pub(super) struct Request { pub(super) request_id: PartitionProcessorRpcRequestId, @@ -55,15 +56,13 @@ impl<'a, TActuator: Actuator, TSchemas, TStorage> RpcHandler } }; - let partition_key = service_invocation.partition_key(); - let cmd = Command::Invoke(Box::new(service_invocation)); + let record = commands::InvokeCommand::from(service_invocation); match append_invocation_reply_on { AppendInvocationReplyOn::Appended => { self.proposer .append_and_respond_asynchronously( - partition_key, - cmd, + record, replier, PartitionProcessorRpcResponse::Appended, ) @@ -71,7 +70,7 @@ impl<'a, TActuator: Actuator, TSchemas, TStorage> RpcHandler } AppendInvocationReplyOn::Submitted | AppendInvocationReplyOn::Output => { self.proposer - .handle_rpc_proposal_command(partition_key, cmd, request_id, replier) + .handle_rpc_proposal_command(record, request_id, replier) .await; } } @@ -87,7 +86,6 @@ mod tests { use crate::partition::rpc::MockActuator; use futures::FutureExt; use googletest::prelude::*; - use restate_test_util::let_assert; use std::future::ready; use test_log::test; @@ -95,20 +93,21 @@ mod tests { async fn reply_on_appended() { let mut proposer = MockActuator::new(); proposer - .expect_append_and_respond_asynchronously::() - .return_once_st(|_, cmd, _, _| { - let_assert!(Command::Invoke(service_invocation) = cmd); + .expect_append_and_respond_asynchronously::() + .return_once_st(|cmd, _, _| { + let service_invocation: ServiceInvocation = cmd.into(); + assert_that!( service_invocation, - points_to(all!( + all!( field!(ServiceInvocation.response_sink, none()), field!(ServiceInvocation.submit_notification_sink, none()), - )) + ) ); ready(()).boxed() }); proposer - .expect_handle_rpc_proposal_command::() + .expect_handle_rpc_proposal_command::() .never(); let (tx, _rx) = Reciprocal::mock(); @@ -130,21 +129,22 @@ mod tests { let request_id = PartitionProcessorRpcRequestId::new(); let mut proposer = MockActuator::new(); proposer - .expect_append_and_respond_asynchronously::() + .expect_append_and_respond_asynchronously::() .never(); proposer - .expect_handle_rpc_proposal_command::() - .return_once_st(|_, cmd, req_id, _| { - let_assert!(Command::Invoke(service_invocation) = cmd); + .expect_handle_rpc_proposal_command::() + .return_once_st(|cmd, req_id, _| { + let service_invocation: ServiceInvocation = cmd.into(); + assert_that!( service_invocation, - points_to(all!( + all!( field!(ServiceInvocation.response_sink, none()), field!( ServiceInvocation.submit_notification_sink, some(eq(SubmitNotificationSink::Ingress { request_id: req_id })) ), - )) + ) ); ready(()).boxed() }); @@ -168,15 +168,16 @@ mod tests { let request_id = PartitionProcessorRpcRequestId::new(); let mut proposer = MockActuator::new(); proposer - .expect_append_and_respond_asynchronously::() + .expect_append_and_respond_asynchronously::() .never(); proposer - .expect_handle_rpc_proposal_command::() - .return_once_st(|_, cmd, req_id, _| { - let_assert!(Command::Invoke(service_invocation) = cmd); + .expect_handle_rpc_proposal_command::() + .return_once_st(|cmd, req_id, _| { + let service_invocation: ServiceInvocation= cmd.into(); + assert_that!( service_invocation, - points_to(all!( + all!( field!( ServiceInvocation.response_sink, some(eq(ServiceInvocationResponseSink::Ingress { @@ -184,7 +185,7 @@ mod tests { })) ), field!(ServiceInvocation.submit_notification_sink, none()), - )) + ) ); ready(()).boxed() }); diff --git a/crates/worker/src/partition/rpc/append_invocation_response.rs b/crates/worker/src/partition/rpc/append_invocation_response.rs index 7f6842a644..657000a5fe 100644 --- a/crates/worker/src/partition/rpc/append_invocation_response.rs +++ b/crates/worker/src/partition/rpc/append_invocation_response.rs @@ -8,11 +8,11 @@ // the Business Source License, use of this software will be governed // by the Apache License, Version 2.0. -use super::*; -use restate_types::identifiers::WithPartitionKey; use restate_types::invocation::InvocationResponse; use restate_types::net::partition_processor::PartitionProcessorRpcResponse; -use restate_wal_protocol::Command; +use restate_wal_protocol::v2::commands; + +use super::*; pub(super) struct Request { pub(super) invocation_response: InvocationResponse, @@ -33,8 +33,7 @@ impl<'a, TActuator: Actuator, TSchemas, TStorage> RpcHandler ) -> Result<(), Self::Error> { self.proposer .append_and_respond_asynchronously( - invocation_response.partition_key(), - Command::InvocationResponse(invocation_response), + commands::InvocationResponseCommand::from(invocation_response), replier, PartitionProcessorRpcResponse::Appended, ) diff --git a/crates/worker/src/partition/rpc/append_signal.rs b/crates/worker/src/partition/rpc/append_signal.rs index d8f08d0df6..7ad400673a 100644 --- a/crates/worker/src/partition/rpc/append_signal.rs +++ b/crates/worker/src/partition/rpc/append_signal.rs @@ -8,12 +8,13 @@ // the Business Source License, use of this software will be governed // by the Apache License, Version 2.0. -use super::*; -use restate_types::identifiers::{InvocationId, WithPartitionKey}; +use restate_types::identifiers::InvocationId; use restate_types::invocation::NotifySignalRequest; use restate_types::journal_v2::Signal; use restate_types::net::partition_processor::PartitionProcessorRpcResponse; -use restate_wal_protocol::Command; +use restate_wal_protocol::v2::commands; + +use super::*; pub(super) struct Request { pub(super) invocation_id: InvocationId, @@ -36,8 +37,7 @@ impl<'a, TActuator: Actuator, TSchemas, TStorage> RpcHandler ) -> Result<(), Self::Error> { self.proposer .append_and_respond_asynchronously( - invocation_id.partition_key(), - Command::NotifySignal(NotifySignalRequest { + commands::NotifySignalCommand::from(NotifySignalRequest { invocation_id, signal, }), diff --git a/crates/worker/src/partition/rpc/cancel_invocation.rs b/crates/worker/src/partition/rpc/cancel_invocation.rs index d7829db624..49bd51eae1 100644 --- a/crates/worker/src/partition/rpc/cancel_invocation.rs +++ b/crates/worker/src/partition/rpc/cancel_invocation.rs @@ -8,14 +8,15 @@ // the Business Source License, use of this software will be governed // by the Apache License, Version 2.0. -use super::*; -use restate_types::identifiers::{InvocationId, WithPartitionKey}; +use restate_types::identifiers::InvocationId; use restate_types::invocation::{ IngressInvocationResponseSink, InvocationMutationResponseSink, InvocationTermination, TerminationFlavor, }; use restate_types::net::partition_processor::CancelInvocationRpcResponse; -use restate_wal_protocol::Command; +use restate_wal_protocol::v2::commands; + +use super::*; pub(super) struct Request { pub(super) request_id: PartitionProcessorRpcRequestId, @@ -38,8 +39,7 @@ impl<'a, TActuator: Actuator, TSchemas, TStorage> RpcHandler ) -> Result<(), Self::Error> { self.proposer .handle_rpc_proposal_command( - invocation_id.partition_key(), - Command::TerminateInvocation(InvocationTermination { + commands::TerminateInvocationCommand::from(InvocationTermination { invocation_id, flavor: TerminationFlavor::Cancel, response_sink: Some(InvocationMutationResponseSink::Ingress( diff --git a/crates/worker/src/partition/rpc/get_invocation_output.rs b/crates/worker/src/partition/rpc/get_invocation_output.rs index f864239a6a..8e59b20fb6 100644 --- a/crates/worker/src/partition/rpc/get_invocation_output.rs +++ b/crates/worker/src/partition/rpc/get_invocation_output.rs @@ -8,10 +8,8 @@ // the Business Source License, use of this software will be governed // by the Apache License, Version 2.0. -use super::*; use restate_storage_api::StorageError; use restate_storage_api::invocation_status_table::{InvocationStatus, ReadInvocationStatusTable}; -use restate_types::identifiers::WithPartitionKey; use restate_types::invocation; use restate_types::invocation::client::{InvocationOutput, InvocationOutputResponse}; use restate_types::invocation::{ @@ -20,7 +18,9 @@ use restate_types::invocation::{ use restate_types::net::partition_processor::{ GetInvocationOutputResponseMode, PartitionProcessorRpcError, PartitionProcessorRpcResponse, }; -use restate_wal_protocol::Command; +use restate_wal_protocol::v2::commands; + +use super::*; pub(super) struct Request { pub(super) request_id: PartitionProcessorRpcRequestId, @@ -95,8 +95,7 @@ where self.proposer .handle_rpc_proposal_command( - invocation_query.partition_key(), - Command::AttachInvocation(AttachInvocationRequest { + commands::AttachInvocationCommand::from(AttachInvocationRequest { invocation_query, block_on_inflight: true, response_sink: ServiceInvocationResponseSink::Ingress { request_id }, diff --git a/crates/worker/src/partition/rpc/kill_invocation.rs b/crates/worker/src/partition/rpc/kill_invocation.rs index 4401199918..9c7a3f1a62 100644 --- a/crates/worker/src/partition/rpc/kill_invocation.rs +++ b/crates/worker/src/partition/rpc/kill_invocation.rs @@ -8,14 +8,15 @@ // the Business Source License, use of this software will be governed // by the Apache License, Version 2.0. -use super::*; -use restate_types::identifiers::{InvocationId, WithPartitionKey}; +use restate_types::identifiers::InvocationId; use restate_types::invocation::{ IngressInvocationResponseSink, InvocationMutationResponseSink, InvocationTermination, TerminationFlavor, }; use restate_types::net::partition_processor::KillInvocationRpcResponse; -use restate_wal_protocol::Command; +use restate_wal_protocol::v2::commands; + +use super::*; pub(super) struct Request { pub(super) request_id: PartitionProcessorRpcRequestId, @@ -38,8 +39,7 @@ impl<'a, TActuator: Actuator, TSchemas, TStorage> RpcHandler ) -> Result<(), Self::Error> { self.proposer .handle_rpc_proposal_command( - invocation_id.partition_key(), - Command::TerminateInvocation(InvocationTermination { + commands::TerminateInvocationCommand::from(InvocationTermination { invocation_id, flavor: TerminationFlavor::Kill, response_sink: Some(InvocationMutationResponseSink::Ingress( diff --git a/crates/worker/src/partition/rpc/mod.rs b/crates/worker/src/partition/rpc/mod.rs index 4ca6e5f565..c7663a194f 100644 --- a/crates/worker/src/partition/rpc/mod.rs +++ b/crates/worker/src/partition/rpc/mod.rs @@ -27,37 +27,39 @@ use restate_core::network::{Oneshot, Reciprocal, TransportConnect}; use restate_storage_api::invocation_status_table::ReadInvocationStatusTable; use restate_storage_api::journal_table as journal_table_v1; use restate_storage_api::journal_table_v2::ReadJournalTable; -use restate_types::identifiers::{ - InvocationId, PartitionId, PartitionKey, PartitionProcessorRpcRequestId, -}; +use restate_types::identifiers::{InvocationId, PartitionId, PartitionProcessorRpcRequestId}; use restate_types::invocation::InvocationRequest; +use restate_types::logs::HasRecordKeys; use restate_types::net::partition_processor::{ AppendInvocationReplyOn, PartitionProcessorRpcError, PartitionProcessorRpcRequest, PartitionProcessorRpcRequestInner, PartitionProcessorRpcResponse, }; use restate_types::schema::deployment::DeploymentResolver; -use restate_wal_protocol::Command; +use restate_wal_protocol::invocation::PauseInvocationCommand; +use restate_wal_protocol::v2::{Command, CommandWithKeys}; use restate_worker_api::invoker::InvokerHandle; use crate::partition::leadership::LeadershipState; #[cfg_attr(test, mockall::automock)] pub(super) trait Actuator { - fn handle_rpc_proposal_command( + fn handle_rpc_proposal_command( &mut self, - partition_key: PartitionKey, - cmd: Command, + cmd: K, request_id: PartitionProcessorRpcRequestId, replier: Replier, - ) -> impl Future; + ) -> impl Future + where + K: CommandWithKeys + 'static, + C: Command, + O: 'static; /// Like [`Self::handle_rpc_proposal_command`], but for the PauseInvocation command: once the /// command is appended, it clears `invocation_id`'s in-memory fencing token so a straggler /// effect from the attempt being paused is dropped at write time. fn propose_pause_and_fence( &mut self, - invocation_id: InvocationId, - cmd: Command, + cmd: PauseInvocationCommand, request_id: PartitionProcessorRpcRequestId, replier: Replier, ) -> impl Future; @@ -68,11 +70,11 @@ pub(super) trait Actuator { /// transitions, making this safe for fire-and-forget ingress commands (signals, invocation /// responses). fn append_and_respond_asynchronously< + C: Command + HasRecordKeys, O: 'static + Into + Send + Sync, >( &mut self, - partition_key: PartitionKey, - cmd: Command, + cmd: C, replier: Replier, success_response: O, ) -> impl Future; @@ -90,16 +92,17 @@ impl Actuator for LeadershipState where T: TransportConnect, { - async fn append_and_respond_asynchronously>( + async fn append_and_respond_asynchronously< + C: Command + HasRecordKeys, + O: Into, + >( &mut self, - partition_key: PartitionKey, - cmd: Command, + cmd: C, replier: Replier, on_proposed_response: O, ) { LeadershipState::append_and_respond_asynchronously( self, - partition_key, cmd, replier.0, on_proposed_response.into(), @@ -107,32 +110,25 @@ where .await } - async fn handle_rpc_proposal_command( + async fn handle_rpc_proposal_command( &mut self, - partition_key: PartitionKey, - cmd: Command, + cmd: K, request_id: PartitionProcessorRpcRequestId, replier: Replier, - ) { - LeadershipState::handle_rpc_proposal_command( - self, - request_id, - replier.0, - partition_key, - cmd, - ) - .await + ) where + K: CommandWithKeys, + C: Command, + { + LeadershipState::handle_rpc_proposal_command(self, request_id, replier.0, cmd).await } async fn propose_pause_and_fence( &mut self, - invocation_id: InvocationId, - cmd: Command, + cmd: PauseInvocationCommand, request_id: PartitionProcessorRpcRequestId, replier: Replier, ) { - LeadershipState::propose_pause_and_fence(self, request_id, replier.0, invocation_id, cmd) - .await + LeadershipState::propose_pause_and_fence(self, request_id, replier.0, cmd).await } fn notify_invoker_to_retry_now(&mut self, invocation_id: InvocationId) { diff --git a/crates/worker/src/partition/rpc/pause_invocation.rs b/crates/worker/src/partition/rpc/pause_invocation.rs index b058b45918..7917eb050c 100644 --- a/crates/worker/src/partition/rpc/pause_invocation.rs +++ b/crates/worker/src/partition/rpc/pause_invocation.rs @@ -72,14 +72,10 @@ where // straggler effect from the attempt we are pausing is dropped at write time. self.proposer .propose_pause_and_fence( - invocation_id, - Command::PauseInvocation( - PauseInvocationCommand { - invocation_id, - request_id: Some(request_id), - } - .bilrost_encode_to_bytes(), - ), + PauseInvocationCommand { + invocation_id, + request_id: Some(request_id), + }, request_id, replier, ) @@ -240,12 +236,7 @@ mod tests { proposer.expect_is_leader().return_const(true); proposer .expect_propose_pause_and_fence::() - .return_once_st(move |got_invocation_id, cmd, request_id, replier| { - assert_eq!(got_invocation_id, invocation_id); - let Command::PauseInvocation(bytes) = cmd else { - panic!("expected a PauseInvocation command"); - }; - let pause = PauseInvocationCommand::bilrost_decode(bytes).unwrap(); + .return_once_st(move |pause, request_id, replier| { assert_eq!(pause.invocation_id, invocation_id); assert_eq!(pause.request_id, Some(request_id)); replier.send(PauseInvocationRpcResponse::Accepted); diff --git a/crates/worker/src/partition/rpc/purge_invocation.rs b/crates/worker/src/partition/rpc/purge_invocation.rs index 291c268c3d..e1ff77358d 100644 --- a/crates/worker/src/partition/rpc/purge_invocation.rs +++ b/crates/worker/src/partition/rpc/purge_invocation.rs @@ -8,13 +8,14 @@ // the Business Source License, use of this software will be governed // by the Apache License, Version 2.0. -use super::*; -use restate_types::identifiers::{InvocationId, WithPartitionKey}; +use restate_types::identifiers::InvocationId; use restate_types::invocation::{ IngressInvocationResponseSink, InvocationMutationResponseSink, PurgeInvocationRequest, }; use restate_types::net::partition_processor::PurgeInvocationRpcResponse; -use restate_wal_protocol::Command; +use restate_wal_protocol::v2::commands; + +use super::*; pub(super) struct Request { pub(super) request_id: PartitionProcessorRpcRequestId, @@ -37,8 +38,7 @@ impl<'a, TActuator: Actuator, TSchemas, TStorage> RpcHandler ) -> Result<(), Self::Error> { self.proposer .handle_rpc_proposal_command( - invocation_id.partition_key(), - Command::PurgeInvocation(PurgeInvocationRequest { + commands::PurgeInvocationCommand::from(PurgeInvocationRequest { invocation_id, response_sink: Some(InvocationMutationResponseSink::Ingress( IngressInvocationResponseSink { request_id }, diff --git a/crates/worker/src/partition/rpc/purge_journal.rs b/crates/worker/src/partition/rpc/purge_journal.rs index 7630c28038..25b6337657 100644 --- a/crates/worker/src/partition/rpc/purge_journal.rs +++ b/crates/worker/src/partition/rpc/purge_journal.rs @@ -8,13 +8,14 @@ // the Business Source License, use of this software will be governed // by the Apache License, Version 2.0. -use super::*; -use restate_types::identifiers::{InvocationId, WithPartitionKey}; +use restate_types::identifiers::InvocationId; use restate_types::invocation::{ IngressInvocationResponseSink, InvocationMutationResponseSink, PurgeInvocationRequest, }; use restate_types::net::partition_processor::PurgeInvocationRpcResponse; -use restate_wal_protocol::Command; +use restate_wal_protocol::v2::commands; + +use super::*; pub(super) struct Request { pub(super) request_id: PartitionProcessorRpcRequestId, @@ -37,8 +38,7 @@ impl<'a, TActuator: Actuator, TSchemas, TStorage> RpcHandler ) -> Result<(), Self::Error> { self.proposer .handle_rpc_proposal_command( - invocation_id.partition_key(), - Command::PurgeJournal(PurgeInvocationRequest { + commands::PurgeJournalCommand::from(PurgeInvocationRequest { invocation_id, response_sink: Some(InvocationMutationResponseSink::Ingress( IngressInvocationResponseSink { request_id }, diff --git a/crates/worker/src/partition/rpc/restart_as_new_invocation.rs b/crates/worker/src/partition/rpc/restart_as_new_invocation.rs index 79fb672343..5d0563548a 100644 --- a/crates/worker/src/partition/rpc/restart_as_new_invocation.rs +++ b/crates/worker/src/partition/rpc/restart_as_new_invocation.rs @@ -8,9 +8,9 @@ // the Business Source License, use of this software will be governed // by the Apache License, Version 2.0. -use super::*; use assert2::let_assert; use opentelemetry::trace::Span; + use restate_service_protocol::codec::ProtobufRawEntryCodec as OldProtocolEntryCodec; use restate_service_protocol_v4::entry_codec::ServiceProtocolV4Codec; use restate_storage_api::invocation_status_table::{InvocationStatus, ReadInvocationStatusTable}; @@ -28,6 +28,9 @@ use restate_types::journal_v2::{CommandMetadata, EntryMetadata, EntryType}; use restate_types::net::partition_processor::RestartAsNewInvocationRpcResponse; use restate_types::service_protocol::ServiceProtocolVersion; use restate_types::{invocation, journal_v2}; +use restate_wal_protocol::v2::commands; + +use super::*; pub(super) struct Request { pub(super) request_id: PartitionProcessorRpcRequestId, @@ -240,14 +243,13 @@ where ); // Propose the usual Invoke command - let cmd = Command::Invoke(Box::new(service_invocation)); + let record = commands::InvokeCommand::from(service_invocation); // Propose and done // This path should be no longer needed once we switch to the journal v2 by default. self.proposer .append_and_respond_asynchronously( - invocation_id.partition_key(), - cmd, + record, replier, RestartAsNewInvocationRpcResponse::Ok { new_invocation_id }, ) @@ -372,7 +374,7 @@ where } // Pass the ball to the state machine, the PP will reply to the RPC request. - let cmd = Command::RestartAsNewInvocation(RestartAsNewInvocationRequest { + let record = commands::RestartAsNewInvocationCommand::from(RestartAsNewInvocationRequest { invocation_id, new_invocation_id, copy_prefix_up_to_index_included, @@ -382,7 +384,7 @@ where )), }); self.proposer - .handle_rpc_proposal_command(invocation_id.partition_key(), cmd, request_id, replier) + .handle_rpc_proposal_command(record, request_id, replier) .await; Ok(()) @@ -394,7 +396,6 @@ mod tests { use std::collections::HashMap; use std::future::ready; - use assert2::let_assert; use bytes::Bytes; use futures::{FutureExt, Stream, StreamExt, stream}; use googletest::prelude::*; @@ -735,12 +736,12 @@ mod tests { let headers_clone = vec![Header::new("key", "value")]; let payload_clone = payload.clone(); proposer - .expect_append_and_respond_asynchronously::() - .return_once_st(move |_, cmd, _, response| { - let_assert!(Command::Invoke(service_invocation) = cmd); + .expect_append_and_respond_asynchronously::() + .return_once_st(move |cmd, _, response| { + let service_invocation: ServiceInvocation = cmd.into(); assert_that!( service_invocation, - points_to(all!( + all!( field!(ServiceInvocation.invocation_id, not(eq(old_invocation_id))), field!(ServiceInvocation.argument, eq(payload_clone)), field!(ServiceInvocation.headers, eq(headers_clone)), @@ -750,7 +751,7 @@ mod tests { ), field!(ServiceInvocation.response_sink, none()), field!(ServiceInvocation.submit_notification_sink, none()), - )) + ) ); assert_that!( response, @@ -761,7 +762,7 @@ mod tests { ready(()).boxed() }); proposer - .expect_handle_rpc_proposal_command::() + .expect_handle_rpc_proposal_command::() .never(); let mut storage = MockStorage::new_with_input_v1( @@ -803,10 +804,10 @@ mod tests { let mut proposer = MockActuator::new(); proposer.expect_is_leader().return_const(true); proposer - .expect_append_and_respond_asynchronously::() + .expect_append_and_respond_asynchronously::() .never(); proposer - .expect_handle_rpc_proposal_command::() + .expect_handle_rpc_proposal_command::() .never(); // Completed with no pinned deployment triggers v1 workaround @@ -858,10 +859,10 @@ mod tests { let mut proposer = MockActuator::new(); proposer.expect_is_leader().return_const(true); proposer - .expect_append_and_respond_asynchronously::() + .expect_append_and_respond_asynchronously::() .never(); proposer - .expect_handle_rpc_proposal_command::() + .expect_handle_rpc_proposal_command::() .never(); let status = InvocationStatus::Completed(CompletedInvocation { @@ -905,10 +906,10 @@ mod tests { let mut proposer = MockActuator::new(); proposer.expect_is_leader().return_const(true); proposer - .expect_append_and_respond_asynchronously::() + .expect_append_and_respond_asynchronously::() .never(); proposer - .expect_handle_rpc_proposal_command::() + .expect_handle_rpc_proposal_command::() .never(); let status = InvocationStatus::Completed(CompletedInvocation { @@ -959,10 +960,10 @@ mod tests { let mut proposer = MockActuator::new(); proposer.expect_is_leader().return_const(true); proposer - .expect_append_and_respond_asynchronously::() + .expect_append_and_respond_asynchronously::() .never(); proposer - .expect_handle_rpc_proposal_command::() + .expect_handle_rpc_proposal_command::() .never(); // Completed with journal length>0 but no v1 input present in storage @@ -1015,21 +1016,20 @@ mod tests { let mut proposer = MockActuator::new(); proposer.expect_is_leader().return_const(true); proposer - .expect_handle_rpc_proposal_command::() - .return_once_st(move |_, cmd, _, _| { + .expect_handle_rpc_proposal_command::() + .return_once_st(move |cmd, _, _| { + let request: RestartAsNewInvocationRequest = cmd.into(); assert_that!( - cmd, - pat!(Command::RestartAsNewInvocation(pat!( - RestartAsNewInvocationRequest { - copy_prefix_up_to_index_included: eq(0), - patch_deployment_id: none() - } - ))) + request, + pat!(RestartAsNewInvocationRequest { + copy_prefix_up_to_index_included: eq(0), + patch_deployment_id: none() + }) ); ready(()).boxed() }); proposer - .expect_append_and_respond_asynchronously::() + .expect_append_and_respond_asynchronously::() .never(); let (tx, rx) = Reciprocal::mock(); @@ -1077,21 +1077,20 @@ mod tests { let mut proposer = MockActuator::new(); proposer.expect_is_leader().return_const(true); proposer - .expect_handle_rpc_proposal_command::() - .return_once_st(move |_, cmd, _, _| { + .expect_handle_rpc_proposal_command::() + .return_once_st(move |cmd, _, _| { + let request: RestartAsNewInvocationRequest =cmd.into(); assert_that!( - cmd, - pat!(Command::RestartAsNewInvocation(pat!( - RestartAsNewInvocationRequest { - copy_prefix_up_to_index_included: eq(0), - patch_deployment_id: none() - } - ))) + request, + pat!(RestartAsNewInvocationRequest { + copy_prefix_up_to_index_included: eq(0), + patch_deployment_id: none() + }) ); ready(()).boxed() }); proposer - .expect_append_and_respond_asynchronously::() + .expect_append_and_respond_asynchronously::() .never(); let (tx, rx) = Reciprocal::mock(); @@ -1122,10 +1121,10 @@ mod tests { let mut proposer = MockActuator::new(); proposer.expect_is_leader().return_const(true); proposer - .expect_append_and_respond_asynchronously::() + .expect_append_and_respond_asynchronously::() .never(); proposer - .expect_handle_rpc_proposal_command::() + .expect_handle_rpc_proposal_command::() .never(); let mut storage = MockStorage::new_without_journal(invocation_id, Default::default()); @@ -1159,10 +1158,10 @@ mod tests { let mut proposer = MockActuator::new(); proposer.expect_is_leader().return_const(true); proposer - .expect_append_and_respond_asynchronously::() + .expect_append_and_respond_asynchronously::() .never(); proposer - .expect_handle_rpc_proposal_command::() + .expect_handle_rpc_proposal_command::() .never(); let mut storage = MockStorage::new_with_input_v1( @@ -1218,10 +1217,10 @@ mod tests { let mut proposer = MockActuator::new(); proposer.expect_is_leader().return_const(true); proposer - .expect_append_and_respond_asynchronously::() + .expect_append_and_respond_asynchronously::() .never(); proposer - .expect_handle_rpc_proposal_command::() + .expect_handle_rpc_proposal_command::() .never(); let mut storage = MockStorage::new_without_journal(invocation_id, status); @@ -1267,10 +1266,10 @@ mod tests { let mut proposer = MockActuator::new(); proposer.expect_is_leader().return_const(true); proposer - .expect_append_and_respond_asynchronously::() + .expect_append_and_respond_asynchronously::() .never(); proposer - .expect_handle_rpc_proposal_command::() + .expect_handle_rpc_proposal_command::() .never(); let mut storage = MockStorage::new_without_journal(invocation_id, status); @@ -1351,21 +1350,20 @@ mod tests { let mut proposer = MockActuator::new(); proposer.expect_is_leader().return_const(true); proposer - .expect_handle_rpc_proposal_command::() - .return_once_st(move |_, cmd, _, _| { + .expect_handle_rpc_proposal_command::() + .return_once_st(move |cmd, _, _| { + let request: RestartAsNewInvocationRequest = cmd.into(); assert_that!( - cmd, - pat!(Command::RestartAsNewInvocation(pat!( - RestartAsNewInvocationRequest { - copy_prefix_up_to_index_included: eq(1), - patch_deployment_id: none() - } - ))) + request, + pat!(RestartAsNewInvocationRequest { + copy_prefix_up_to_index_included: eq(1), + patch_deployment_id: none() + }) ); ready(()).boxed() }); proposer - .expect_append_and_respond_asynchronously::() + .expect_append_and_respond_asynchronously::() .never(); let (tx, rx) = Reciprocal::mock(); @@ -1415,21 +1413,20 @@ mod tests { let mut proposer = MockActuator::new(); proposer.expect_is_leader().return_const(true); proposer - .expect_handle_rpc_proposal_command::() - .return_once_st(move |_, cmd, _, _| { + .expect_handle_rpc_proposal_command::() + .return_once_st(move |cmd, _, _| { + let request: RestartAsNewInvocationRequest = cmd.into(); assert_that!( - cmd, - pat!(Command::RestartAsNewInvocation(pat!( - RestartAsNewInvocationRequest { - copy_prefix_up_to_index_included: eq(1), - patch_deployment_id: some(eq(latest_id)) - } - ))) + request, + pat!(RestartAsNewInvocationRequest { + copy_prefix_up_to_index_included: eq(1), + patch_deployment_id: some(eq(latest_id)) + }) ); ready(()).boxed() }); proposer - .expect_append_and_respond_asynchronously::() + .expect_append_and_respond_asynchronously::() .never(); let (tx, rx) = Reciprocal::mock(); @@ -1473,10 +1470,10 @@ mod tests { let mut proposer = MockActuator::new(); proposer.expect_is_leader().return_const(true); proposer - .expect_handle_rpc_proposal_command::() + .expect_handle_rpc_proposal_command::() .never(); proposer - .expect_append_and_respond_asynchronously::() + .expect_append_and_respond_asynchronously::() .never(); let (tx, rx) = Reciprocal::mock(); @@ -1524,7 +1521,7 @@ mod tests { let mut proposer = MockActuator::new(); proposer.expect_is_leader().return_const(true); proposer - .expect_handle_rpc_proposal_command::() + .expect_handle_rpc_proposal_command::() .never(); let (tx, rx) = Reciprocal::mock(); @@ -1565,7 +1562,7 @@ mod tests { let mut proposer = MockActuator::new(); proposer.expect_is_leader().return_const(true); proposer - .expect_handle_rpc_proposal_command::() + .expect_handle_rpc_proposal_command::() .never(); let (tx, rx) = Reciprocal::mock(); @@ -1606,7 +1603,7 @@ mod tests { let mut proposer = MockActuator::new(); proposer.expect_is_leader().return_const(true); proposer - .expect_handle_rpc_proposal_command::() + .expect_handle_rpc_proposal_command::() .never(); let (tx, rx) = Reciprocal::mock(); @@ -1645,10 +1642,10 @@ mod tests { .expect_partition_id() .return_const(PartitionId::from(0)); proposer - .expect_append_and_respond_asynchronously::() + .expect_append_and_respond_asynchronously::() .never(); proposer - .expect_handle_rpc_proposal_command::() + .expect_handle_rpc_proposal_command::() .never(); let mut storage = MockStorage::new_without_journal(invocation_id, Default::default()); diff --git a/crates/worker/src/partition/rpc/resume_invocation.rs b/crates/worker/src/partition/rpc/resume_invocation.rs index 5d04e9a417..00eadf2c43 100644 --- a/crates/worker/src/partition/rpc/resume_invocation.rs +++ b/crates/worker/src/partition/rpc/resume_invocation.rs @@ -8,16 +8,18 @@ // the Business Source License, use of this software will be governed // by the Apache License, Version 2.0. -use super::*; -use crate::partition::state_machine::resolve_pinned_deployment; use restate_storage_api::invocation_status_table::{InvocationStatus, ReadInvocationStatusTable}; -use restate_types::identifiers::{InvocationId, WithPartitionKey}; +use restate_types::identifiers::InvocationId; use restate_types::invocation::client::PatchDeploymentId; use restate_types::invocation::{ IngressInvocationResponseSink, InvocationMutationResponseSink, ResumeInvocationRequest, }; use restate_types::net::partition_processor::ResumeInvocationRpcResponse; use restate_types::schema::deployment::DeploymentResolver; +use restate_wal_protocol::v2::commands; + +use super::*; +use crate::partition::state_machine::resolve_pinned_deployment; pub(super) struct Request { pub(super) request_id: PartitionProcessorRpcRequestId, @@ -68,8 +70,7 @@ where // is safe here: vqueues being enabled implies a cluster min version >= 1.7.0. self.proposer .handle_rpc_proposal_command( - invocation_id.partition_key(), - Command::ResumeInvocation(ResumeInvocationRequest { + commands::ResumeInvocationCommand::from(ResumeInvocationRequest { invocation_id, update_deployment_id: Some(update_deployment_id), update_pinned_deployment_id: None, @@ -101,8 +102,7 @@ where // `update_deployment_id` field -- vqueues imply a cluster min version >= 1.7.0. self.proposer .handle_rpc_proposal_command( - invocation_id.partition_key(), - Command::ResumeInvocation(ResumeInvocationRequest { + commands::ResumeInvocationCommand::from(ResumeInvocationRequest { invocation_id, update_deployment_id: Some(update_deployment_id), update_pinned_deployment_id: None, @@ -137,8 +137,7 @@ where self.proposer .handle_rpc_proposal_command( - invocation_id.partition_key(), - Command::ResumeInvocation(ResumeInvocationRequest { + commands::ResumeInvocationCommand::from(ResumeInvocationRequest { invocation_id, update_deployment_id: None, update_pinned_deployment_id, @@ -177,7 +176,6 @@ mod tests { use super::*; use crate::partition::rpc::MockActuator; - use assert2::let_assert; use futures::FutureExt; use googletest::prelude::*; use restate_storage_api::invocation_status_table::{ @@ -191,6 +189,7 @@ mod tests { use restate_types::schema::deployment::Deployment; use restate_types::schema::deployment::test_util::MockDeploymentMetadataRegistry; use restate_types::service_protocol::ServiceProtocolVersion; + use restate_types::sharding::WithPartitionKey; use rstest::rstest; use std::future::ready; use test_log::test; @@ -264,12 +263,11 @@ mod tests { // The invoker must NOT be poked for a VQueue invocation. proposer.expect_notify_invoker_to_retry_now().never(); proposer - .expect_handle_rpc_proposal_command::() - .return_once_st(move |_, cmd, _request_id, replier| { + .expect_handle_rpc_proposal_command::() + .return_once_st(move |request, _request_id, replier| { // The command shape (response_sink, deployment patching) is covered by // `propose_resume_command_on_paused_and_suspended`; here we only care that the // Invoked+VQueue path proposes ResumeInvocation rather than poking the invoker. - let_assert!(Command::ResumeInvocation(request) = cmd); assert_eq!(request.invocation_id, invocation_id); replier.send(ResumeInvocationRpcResponse::Ok); ready(()).boxed() @@ -318,9 +316,9 @@ mod tests { let mut proposer = MockActuator::new(); proposer.expect_is_leader().return_const(true); proposer - .expect_handle_rpc_proposal_command::() - .return_once_st(move |_, cmd, request_id, replier| { - let_assert!(Command::ResumeInvocation(resume_invocation_request) = cmd); + .expect_handle_rpc_proposal_command::() + .return_once_st(move |cmd, request_id, replier| { + let resume_invocation_request: ResumeInvocationRequest = cmd.into(); assert_that!( resume_invocation_request, all!( @@ -368,7 +366,7 @@ mod tests { let mut proposer = MockActuator::new(); proposer.expect_is_leader().return_const(true); proposer - .expect_handle_rpc_proposal_command::() + .expect_handle_rpc_proposal_command::() .never(); let mut storage = MockStorage { @@ -402,7 +400,7 @@ mod tests { let mut proposer = MockActuator::new(); proposer.expect_is_leader().return_const(true); proposer - .expect_handle_rpc_proposal_command::() + .expect_handle_rpc_proposal_command::() .never(); let mut storage = MockStorage { @@ -448,7 +446,7 @@ mod tests { let mut proposer = MockActuator::new(); proposer.expect_is_leader().return_const(true); proposer - .expect_handle_rpc_proposal_command::() + .expect_handle_rpc_proposal_command::() .never(); let mut storage = MockStorage { @@ -501,7 +499,7 @@ mod tests { let mut proposer = MockActuator::new(); proposer.expect_is_leader().return_const(true); proposer - .expect_handle_rpc_proposal_command::() + .expect_handle_rpc_proposal_command::() .never(); let metadata = InFlightInvocationMetadata { @@ -565,7 +563,7 @@ mod tests { proposer.expect_is_leader().return_const(true); proposer.expect_notify_invoker_to_retry_now().never(); proposer - .expect_handle_rpc_proposal_command::() + .expect_handle_rpc_proposal_command::() .never(); let mut storage = MockStorage { @@ -604,7 +602,7 @@ mod tests { .expect_partition_id() .return_const(PartitionId::from(0)); proposer - .expect_handle_rpc_proposal_command::() + .expect_handle_rpc_proposal_command::() .never(); let mut storage = MockStorage { diff --git a/crates/worker/src/partition/shuffle.rs b/crates/worker/src/partition/shuffle.rs index 50e2baabff..92b0c2ef17 100644 --- a/crates/worker/src/partition/shuffle.rs +++ b/crates/worker/src/partition/shuffle.rs @@ -18,16 +18,18 @@ use tracing::debug; use restate_core::cancellation_token; use restate_core::network::TransportConnect; use restate_ingestion_client::IngestionClient; -use restate_storage_api::deduplication_table::DedupInformation; use restate_storage_api::outbox_table::OutboxMessage; -use restate_types::identifiers::{LeaderEpoch, PartitionId, PartitionKey, WithPartitionKey}; +use restate_types::identifiers::PartitionId; use restate_types::message::MessageIndex; -use restate_wal_protocol::{Destination, Envelope, Header, Source}; +use restate_wal_protocol::v2::commands::{ + AttachInvocationCommand, InvocationResponseCommand, InvokeCommand, NotifySignalCommand, + TerminateInvocationCommand, +}; +use restate_wal_protocol::v2::{Dedup, Envelope, Raw}; use crate::metric_definitions::{ PARTITION_LABEL, PARTITION_SHUFFLE_INFLIGHT_COUNT, PARTITION_SHUFFLE_MESSAGE_COUNT, }; -use crate::partition::types::OutboxMessageExt; #[derive(Debug)] pub(crate) struct NewOutboxMessage { @@ -61,30 +63,32 @@ pub(crate) fn wrap_outbox_message_in_envelope( message: OutboxMessage, seq_number: MessageIndex, shuffle_metadata: &ShuffleMetadata, -) -> Envelope { - Envelope::new( - create_header(message.partition_key(), seq_number, shuffle_metadata), - message.to_command(), - ) +) -> Envelope { + let dedup = create_dedup(seq_number, shuffle_metadata); + + match message { + OutboxMessage::AttachInvocation(payload) => { + Envelope::new(dedup, AttachInvocationCommand::from(payload)).into_raw() + } + OutboxMessage::InvocationTermination(payload) => { + Envelope::new(dedup, TerminateInvocationCommand::from(payload)).into_raw() + } + OutboxMessage::ServiceInvocation(payload) => { + Envelope::new(dedup, InvokeCommand::from(*payload)).into_raw() + } + OutboxMessage::NotifySignal(payload) => { + Envelope::new(dedup, NotifySignalCommand::from(payload)).into_raw() + } + OutboxMessage::ServiceResponse(payload) => { + Envelope::new(dedup, InvocationResponseCommand::from(payload)).into_raw() + } + } } -fn create_header( - dest_partition_key: PartitionKey, - seq_number: MessageIndex, - shuffle_metadata: &ShuffleMetadata, -) -> Header { - Header { - source: Source::Processor { - partition_key: None, - leader_epoch: shuffle_metadata.leader_epoch, - }, - dest: Destination::Processor { - partition_key: dest_partition_key, - dedup: Some(DedupInformation::cross_partition( - shuffle_metadata.partition_id, - seq_number, - )), - }, +fn create_dedup(seq_number: MessageIndex, shuffle_metadata: &ShuffleMetadata) -> Dedup { + Dedup::ForeignPartition { + partition: shuffle_metadata.partition_id, + seq: seq_number, } } @@ -151,22 +155,18 @@ impl HintSender { #[derive(Debug, Copy, Clone)] pub(crate) struct ShuffleMetadata { partition_id: PartitionId, - leader_epoch: LeaderEpoch, } impl ShuffleMetadata { - pub(crate) fn new(partition_id: PartitionId, leader_epoch: LeaderEpoch) -> Self { - ShuffleMetadata { - partition_id, - leader_epoch, - } + pub(crate) fn new(partition_id: PartitionId) -> Self { + ShuffleMetadata { partition_id } } } pub(crate) struct Shuffle { metadata: ShuffleMetadata, outbox_reader: OR, - ingestion_client: IngestionClient, + ingestion_client: IngestionClient>, // used to tell partition processor about outbox truncations truncation_tx: mpsc::Sender, hint_rx: async_channel::Receiver, @@ -184,7 +184,7 @@ where outbox_reader: OR, truncation_tx: mpsc::Sender, channel_size: usize, - ingestion_client: IngestionClient, + ingestion_client: IngestionClient>, ) -> Self { let (hint_tx, hint_rx) = async_channel::bounded(channel_size); @@ -270,8 +270,12 @@ mod state_machine { use restate_core::network::TransportConnect; use restate_ingestion_client::{IngestFuture, IngestionClient, RecordCommit}; use restate_storage_api::outbox_table::OutboxMessage; - use restate_types::{identifiers::WithPartitionKey, message::MessageIndex}; - use restate_wal_protocol::Envelope; + use restate_types::{ + identifiers::WithPartitionKey, + logs::{BodyWithKeys, Keys}, + message::MessageIndex, + }; + use restate_wal_protocol::v2::{Envelope, Raw}; use crate::partition::shuffle::{ NewOutboxMessage, OutboxReaderError, ShuffleMetadata, wrap_outbox_message_in_envelope, @@ -304,7 +308,7 @@ mod state_machine { pub struct StateMachine { metadata: ShuffleMetadata, - ingestion: IngestionClient, + ingestion: IngestionClient>, hint_rx: async_channel::Receiver, reader: Option, read_fut: ReadFuture, @@ -319,7 +323,7 @@ mod state_machine { { pub fn new( metadata: ShuffleMetadata, - ingestion: IngestionClient, + ingestion: IngestionClient>, reader: R, hint_rx: async_channel::Receiver, ) -> Self { @@ -349,12 +353,15 @@ mod state_machine { match sn.cmp(&self.next_sequence_number) { Ordering::Equal => { + let partition_key = message.partition_key(); let envelope = wrap_outbox_message_in_envelope(message, sn, &self.metadata); + self.state = State::Ingesting { - ingest: self - .ingestion - .ingest(envelope.partition_key(), envelope), + ingest: self.ingestion.ingest( + partition_key, + BodyWithKeys::new(envelope, Keys::Single(partition_key)), + ), sn, }; } @@ -398,12 +405,14 @@ mod state_machine { sn >= self.next_sequence_number, "message sequence numbers must not decrease" ); + let partition_key = message.partition_key(); let envelope = wrap_outbox_message_in_envelope(message, sn, &self.metadata); self.state = State::Ingesting { - ingest: self - .ingestion - .ingest(envelope.partition_key(), envelope), + ingest: self.ingestion.ingest( + partition_key, + BodyWithKeys::new(envelope, Keys::Single(partition_key)), + ), sn, }; } @@ -423,7 +432,6 @@ mod ingestion_client_tests { use std::sync::atomic::{AtomicUsize, Ordering}; use anyhow::{Context, anyhow}; - use assert2::let_assert; use futures::StreamExt; use restate_core::partitions::PartitionRouting; use restate_ingestion_client::{IngestionClient, SessionOptions}; @@ -432,6 +440,8 @@ mod ingestion_client_tests { use restate_types::net::partition_processor::PartitionLeaderService; use restate_types::partitions::state::{LeadershipState, PartitionReplicaSetStates}; use restate_types::storage::StorageCodec; + use restate_wal_protocol::v2::commands::InvokeCommand; + use restate_wal_protocol::v2::{CommandKind, Envelope, Raw}; use test_log::test; use tokio::sync::mpsc; @@ -447,7 +457,6 @@ mod ingestion_client_tests { use restate_types::invocation::ServiceInvocation; use restate_types::message::MessageIndex; use restate_types::partition_table::PartitionTable; - use restate_wal_protocol::{Command, Envelope}; use crate::partition::shuffle::{OutboxReader, OutboxReaderError, Shuffle, ShuffleMetadata}; @@ -569,12 +578,15 @@ mod ingestion_client_tests { let (r, body) = incoming.split(); r.send(ResponseStatus::Ack.into()); for mut record in body.records { - let envelope = StorageCodec::decode::(&mut record.record) + let envelope = StorageCodec::decode::, _>(&mut record.record) .context("Failed to decode envelope")?; - let_assert!(Command::Invoke(service_invocation) = envelope.command); + assert!(CommandKind::Invoke == envelope.kind()); + let envelope: Envelope = envelope.into_typed(); + let service_invocation: ServiceInvocation = envelope.into_inner().unwrap().into(); + let invocation_id = service_invocation.invocation_id; - messages.push(*service_invocation); + messages.push(service_invocation); if last_invocation_id == invocation_id { break 'out; @@ -615,7 +627,7 @@ mod ingestion_client_tests { #[allow(dead_code)] env: TestCoreEnv, stream: ServiceStream, - ingestion: IngestionClient, + ingestion: IngestionClient>, shuffle: Shuffle, } @@ -627,7 +639,7 @@ mod ingestion_client_tests { PartitionTable::with_equally_sized_partitions(Version::MIN, 1), ); - let metadata = ShuffleMetadata::new(PartitionId::from(0), LeaderEpoch::from(0)); + let metadata = ShuffleMetadata::new(PartitionId::from(0)); let partition_replica_set_states = PartitionReplicaSetStates::default(); diff --git a/crates/worker/src/partition/state_machine/tests/mod.rs b/crates/worker/src/partition/state_machine/tests/mod.rs index dcff06e8e2..8945205277 100644 --- a/crates/worker/src/partition/state_machine/tests/mod.rs +++ b/crates/worker/src/partition/state_machine/tests/mod.rs @@ -62,7 +62,7 @@ use restate_types::journal::{CompleteAwakeableEntry, EntryResult, InvokeRequest} use restate_types::journal::{Entry, EntryType}; use restate_types::journal_events::Event; use restate_types::journal_v2::raw::TryFromEntry; -use restate_types::logs::{Keys, SequenceNumber}; +use restate_types::logs::SequenceNumber; use restate_types::partitions::Partition; use restate_types::state_mut::ExternalStateMutation; use restate_wal_protocol::v2::Command; @@ -1054,7 +1054,6 @@ async fn truncate_outbox_from_empty() -> Result<(), Error> { .apply(commands::TruncateOutboxCommand::test_envelope( commands::TruncateOutboxCommand { index: outbox_index, - partition_key_range: Keys::None, }, )) .await; @@ -1095,7 +1094,6 @@ async fn truncate_outbox_with_gap() -> Result<(), Error> { .apply(commands::TruncateOutboxCommand::test_envelope( commands::TruncateOutboxCommand { index: outbox_tail_index, - partition_key_range: Keys::None, }, )) .await; diff --git a/crates/worker/src/partition/types.rs b/crates/worker/src/partition/types.rs index 7e6d53d52b..7a52ae4ac1 100644 --- a/crates/worker/src/partition/types.rs +++ b/crates/worker/src/partition/types.rs @@ -11,7 +11,6 @@ use restate_storage_api::outbox_table::OutboxMessage; use restate_types::identifiers::{EntryIndex, InvocationId}; use restate_types::invocation::{InvocationResponse, JournalCompletionTarget, ResponseResult}; -use restate_wal_protocol::Command; pub(crate) type InvokerEffect = restate_worker_api::invoker::FencedEffect; pub(crate) type InvokerEffectKind = restate_worker_api::invoker::EffectKind; @@ -23,8 +22,6 @@ pub(crate) trait OutboxMessageExt { entry_index: EntryIndex, result: ResponseResult, ) -> OutboxMessage; - - fn to_command(self) -> Command; } impl OutboxMessageExt for OutboxMessage { @@ -41,14 +38,4 @@ impl OutboxMessageExt for OutboxMessage { result, }) } - - fn to_command(self) -> Command { - match self { - OutboxMessage::ServiceInvocation(si) => Command::Invoke(si), - OutboxMessage::ServiceResponse(sr) => Command::InvocationResponse(sr), - OutboxMessage::InvocationTermination(it) => Command::TerminateInvocation(it), - OutboxMessage::AttachInvocation(ai) => Command::AttachInvocation(ai), - OutboxMessage::NotifySignal(notify_signal) => Command::NotifySignal(notify_signal), - } - } } diff --git a/crates/worker/src/partition_processor_manager.rs b/crates/worker/src/partition_processor_manager.rs index 7af8c1c1cb..b6f7c66b92 100644 --- a/crates/worker/src/partition_processor_manager.rs +++ b/crates/worker/src/partition_processor_manager.rs @@ -80,7 +80,7 @@ use restate_types::partitions::state::PartitionReplicaSetStates; use restate_types::protobuf::common::WorkerStatus; use restate_util_string::format_restring; use restate_util_time::DurationExt; -use restate_wal_protocol::Envelope; +use restate_wal_protocol::v2::{Envelope, Raw}; use restate_worker_api::invoker::capacity::InvokerCapacity; use restate_worker_api::{ProcessorsManagerCommand, ProcessorsManagerHandle}; @@ -131,7 +131,7 @@ pub struct PartitionProcessorManager { invoker_capacity: InvokerCapacity, - ingestion_client: IngestionClient, + ingestion_client: IngestionClient>, /// Built in `new`; the polling task is spawned at the start of `run`. rule_book_cache_task: Option, @@ -208,7 +208,7 @@ where router_builder: &mut MessageRouterBuilder, bifrost: Bifrost, snapshot_repository: Option, - ingestion_client: IngestionClient, + ingestion_client: IngestionClient>, ) -> Self { let config = updateable_config.pinned(); let ppm_svc_rx = router_builder.register_service(BackPressureMode::Lossy); diff --git a/crates/worker/src/partition_processor_manager/spawn_processor_task.rs b/crates/worker/src/partition_processor_manager/spawn_processor_task.rs index 872f1f3f9f..5f7a160ab1 100644 --- a/crates/worker/src/partition_processor_manager/spawn_processor_task.rs +++ b/crates/worker/src/partition_processor_manager/spawn_processor_task.rs @@ -26,7 +26,7 @@ use restate_types::cluster::cluster_state::PartitionProcessorStatus; use restate_types::logs::Lsn; use restate_types::partitions::Partition; use restate_types::partitions::state::PartitionReplicaSetStates; -use restate_wal_protocol::Envelope; +use restate_wal_protocol::v2::{Envelope, Raw}; use restate_worker_api::invoker::capacity::InvokerCapacity; use crate::PartitionProcessorBuilder; @@ -43,7 +43,7 @@ pub struct SpawnPartitionProcessorTask { partition_store_manager: Arc, fast_forward_lsn: Option, invoker_capacity: InvokerCapacity, - ingestion_client: IngestionClient, + ingestion_client: IngestionClient>, leader_handles_registry: PartitionLeaderHandlesRegistry, rule_book_cache: RuleBookCacheHandle, } @@ -61,7 +61,7 @@ where partition_store_manager: Arc, fast_forward_lsn: Option, invoker_capacity: InvokerCapacity, - ingestion_client: IngestionClient, + ingestion_client: IngestionClient>, leader_handles_registry: PartitionLeaderHandlesRegistry, rule_book_cache: RuleBookCacheHandle, ) -> Self {