Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions crates/xmtp_mls/src/worker/tasks.rs
Original file line number Diff line number Diff line change
Expand Up @@ -392,6 +392,16 @@ where
.db()
.pull_in_task_deadline(&p.target_data_hash, p.not_later_than_ns)?;
}
Some(xmtp_proto::xmtp::mls::database::task::Task::KpRotation(_))
| Some(xmtp_proto::xmtp::mls::database::task::Task::KpDeletion(_)) => {
// Minimal arms: real handlers land with the KP-consumer impl.
// Nothing seeds these singletons yet, so dropping is safe.
tracing::warn!(
"KP task {} received before handler landed; dropping",
task.id
);
context.db().delete_task(task.id)?;
}
None => {
tracing::error!("Task {} has no data. Deleting.", task.id);
context.db().delete_task(task.id)?;
Expand Down
Binary file modified crates/xmtp_proto/src/gen/proto_descriptor.bin
Binary file not shown.
34 changes: 33 additions & 1 deletion crates/xmtp_proto/src/gen/xmtp.mls.database.rs
Original file line number Diff line number Diff line change
Expand Up @@ -783,7 +783,7 @@ impl PermissionPolicyOption {
}
#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
pub struct Task {
#[prost(oneof = "task::Task", tags = "1, 2, 3, 4")]
#[prost(oneof = "task::Task", tags = "1, 2, 3, 4, 5, 6")]
pub task: ::core::option::Option<task::Task>,
}
/// Nested message and enum types in `Task`.
Expand All @@ -798,6 +798,10 @@ pub mod task {
ProcessPendingSelfRemove(super::ProcessPendingSelfRemove),
#[prost(message, tag = "4")]
PullInDeadline(super::PullInDeadline),
#[prost(message, tag = "5")]
KpRotation(super::KpRotation),
#[prost(message, tag = "6")]
KpDeletion(super::KpDeletion),
}
}
impl ::prost::Name for Task {
Expand Down Expand Up @@ -831,6 +835,34 @@ impl ::prost::Name for PullInDeadline {
"/xmtp.mls.database.PullInDeadline".into()
}
}
/// Recurring singleton: rotate + upload a fresh key package when the identity's
/// rotation deadline is due. Empty payload => stable data_hash for pull-ins.
#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
pub struct KpRotation {}
impl ::prost::Name for KpRotation {
const NAME: &'static str = "KpRotation";
const PACKAGE: &'static str = "xmtp.mls.database";
fn full_name() -> ::prost::alloc::string::String {
"xmtp.mls.database.KpRotation".into()
}
fn type_url() -> ::prost::alloc::string::String {
"/xmtp.mls.database.KpRotation".into()
}
}
/// Recurring singleton: delete superseded local key-package material whose
/// delete_at_ns has passed. Empty payload => stable data_hash for pull-ins.
#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
pub struct KpDeletion {}
impl ::prost::Name for KpDeletion {
const NAME: &'static str = "KpDeletion";
const PACKAGE: &'static str = "xmtp.mls.database";
fn full_name() -> ::prost::alloc::string::String {
"xmtp.mls.database.KpDeletion".into()
}
fn type_url() -> ::prost::alloc::string::String {
"/xmtp.mls.database.KpDeletion".into()
}
}
#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
pub struct SendSyncArchive {
#[prost(message, optional, tag = "1")]
Expand Down
172 changes: 172 additions & 0 deletions crates/xmtp_proto/src/gen/xmtp.mls.database.serde.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1002,6 +1002,150 @@ impl<'de> serde::Deserialize<'de> for InstallationIds {
deserializer.deserialize_struct("xmtp.mls.database.InstallationIds", FIELDS, GeneratedVisitor)
}
}
impl serde::Serialize for KpDeletion {
#[allow(deprecated)]
fn serialize<S>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
use serde::ser::SerializeStruct;
let len = 0;
let struct_ser = serializer.serialize_struct("xmtp.mls.database.KpDeletion", len)?;
struct_ser.end()
}
}
impl<'de> serde::Deserialize<'de> for KpDeletion {
#[allow(deprecated)]
fn deserialize<D>(deserializer: D) -> std::result::Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
const FIELDS: &[&str] = &[
];

#[allow(clippy::enum_variant_names)]
enum GeneratedField {
__SkipField__,
}
impl<'de> serde::Deserialize<'de> for GeneratedField {
fn deserialize<D>(deserializer: D) -> std::result::Result<GeneratedField, D::Error>
where
D: serde::Deserializer<'de>,
{
struct GeneratedVisitor;

impl serde::de::Visitor<'_> for GeneratedVisitor {
type Value = GeneratedField;

fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(formatter, "expected one of: {:?}", &FIELDS)
}

#[allow(unused_variables)]
fn visit_str<E>(self, value: &str) -> std::result::Result<GeneratedField, E>
where
E: serde::de::Error,
{
Ok(GeneratedField::__SkipField__)
}
}
deserializer.deserialize_identifier(GeneratedVisitor)
}
}
struct GeneratedVisitor;
impl<'de> serde::de::Visitor<'de> for GeneratedVisitor {
type Value = KpDeletion;

fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str("struct xmtp.mls.database.KpDeletion")
}

fn visit_map<V>(self, mut map_: V) -> std::result::Result<KpDeletion, V::Error>
where
V: serde::de::MapAccess<'de>,
{
while map_.next_key::<GeneratedField>()?.is_some() {
let _ = map_.next_value::<serde::de::IgnoredAny>()?;
}
Ok(KpDeletion {
})
}
}
deserializer.deserialize_struct("xmtp.mls.database.KpDeletion", FIELDS, GeneratedVisitor)
}
}
impl serde::Serialize for KpRotation {
#[allow(deprecated)]
fn serialize<S>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
use serde::ser::SerializeStruct;
let len = 0;
let struct_ser = serializer.serialize_struct("xmtp.mls.database.KpRotation", len)?;
struct_ser.end()
}
}
impl<'de> serde::Deserialize<'de> for KpRotation {
#[allow(deprecated)]
fn deserialize<D>(deserializer: D) -> std::result::Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
const FIELDS: &[&str] = &[
];

#[allow(clippy::enum_variant_names)]
enum GeneratedField {
__SkipField__,
}
impl<'de> serde::Deserialize<'de> for GeneratedField {
fn deserialize<D>(deserializer: D) -> std::result::Result<GeneratedField, D::Error>
where
D: serde::Deserializer<'de>,
{
struct GeneratedVisitor;

impl serde::de::Visitor<'_> for GeneratedVisitor {
type Value = GeneratedField;

fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(formatter, "expected one of: {:?}", &FIELDS)
}

#[allow(unused_variables)]
fn visit_str<E>(self, value: &str) -> std::result::Result<GeneratedField, E>
where
E: serde::de::Error,
{
Ok(GeneratedField::__SkipField__)
}
}
deserializer.deserialize_identifier(GeneratedVisitor)
}
}
struct GeneratedVisitor;
impl<'de> serde::de::Visitor<'de> for GeneratedVisitor {
type Value = KpRotation;

fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str("struct xmtp.mls.database.KpRotation")
}

fn visit_map<V>(self, mut map_: V) -> std::result::Result<KpRotation, V::Error>
where
V: serde::de::MapAccess<'de>,
{
while map_.next_key::<GeneratedField>()?.is_some() {
let _ = map_.next_value::<serde::de::IgnoredAny>()?;
}
Ok(KpRotation {
})
}
}
deserializer.deserialize_struct("xmtp.mls.database.KpRotation", FIELDS, GeneratedVisitor)
}
}
impl serde::Serialize for PermissionPolicyOption {
#[allow(deprecated)]
fn serialize<S>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error>
Expand Down Expand Up @@ -2958,6 +3102,12 @@ impl serde::Serialize for Task {
task::Task::PullInDeadline(v) => {
struct_ser.serialize_field("pull_in_deadline", v)?;
}
task::Task::KpRotation(v) => {
struct_ser.serialize_field("kp_rotation", v)?;
}
task::Task::KpDeletion(v) => {
struct_ser.serialize_field("kp_deletion", v)?;
}
}
}
struct_ser.end()
Expand All @@ -2978,6 +3128,10 @@ impl<'de> serde::Deserialize<'de> for Task {
"processPendingSelfRemove",
"pull_in_deadline",
"pullInDeadline",
"kp_rotation",
"kpRotation",
"kp_deletion",
"kpDeletion",
];

#[allow(clippy::enum_variant_names)]
Expand All @@ -2986,6 +3140,8 @@ impl<'de> serde::Deserialize<'de> for Task {
SendSyncArchive,
ProcessPendingSelfRemove,
PullInDeadline,
KpRotation,
KpDeletion,
__SkipField__,
}
impl<'de> serde::Deserialize<'de> for GeneratedField {
Expand All @@ -3012,6 +3168,8 @@ impl<'de> serde::Deserialize<'de> for Task {
"sendSyncArchive" | "send_sync_archive" => Ok(GeneratedField::SendSyncArchive),
"processPendingSelfRemove" | "process_pending_self_remove" => Ok(GeneratedField::ProcessPendingSelfRemove),
"pullInDeadline" | "pull_in_deadline" => Ok(GeneratedField::PullInDeadline),
"kpRotation" | "kp_rotation" => Ok(GeneratedField::KpRotation),
"kpDeletion" | "kp_deletion" => Ok(GeneratedField::KpDeletion),
_ => Ok(GeneratedField::__SkipField__),
}
}
Expand Down Expand Up @@ -3060,6 +3218,20 @@ impl<'de> serde::Deserialize<'de> for Task {
return Err(serde::de::Error::duplicate_field("pullInDeadline"));
}
task__ = map_.next_value::<::std::option::Option<_>>()?.map(task::Task::PullInDeadline)
;
}
GeneratedField::KpRotation => {
if task__.is_some() {
return Err(serde::de::Error::duplicate_field("kpRotation"));
}
task__ = map_.next_value::<::std::option::Option<_>>()?.map(task::Task::KpRotation)
;
}
GeneratedField::KpDeletion => {
if task__.is_some() {
return Err(serde::de::Error::duplicate_field("kpDeletion"));
}
task__ = map_.next_value::<::std::option::Option<_>>()?.map(task::Task::KpDeletion)
;
}
GeneratedField::__SkipField__ => {
Expand Down
Loading