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/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/mod.rs b/crates/worker/src/partition/leadership/mod.rs index b60ddc2972..be5939a720 100644 --- a/crates/worker/src/partition/leadership/mod.rs +++ b/crates/worker/src/partition/leadership/mod.rs @@ -67,11 +67,12 @@ use restate_types::{GenerationalNodeId, SemanticRestateVersion}; use restate_util_time::DurationExt; use restate_vqueues::scheduler::{self}; use restate_vqueues::{ResourceManager, SchedulerService, VQueuesMeta, VQueuesMetaCache}; +use restate_wal_protocol::Command; use restate_wal_protocol::control::{ AnnounceLeaderCommand, UpdatePartitionDurabilityCommand, VersionBarrierCommand, }; use restate_wal_protocol::timer::TimerKeyValue; -use restate_wal_protocol::{Command, Envelope}; +use restate_wal_protocol::v2::{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, @@ -508,7 +509,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(), diff --git a/crates/worker/src/partition/mod.rs b/crates/worker/src/partition/mod.rs index 7c7b5e9d79..749f5fdb64 100644 --- a/crates/worker/src/partition/mod.rs +++ b/crates/worker/src/partition/mod.rs @@ -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> 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/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 {