Skip to content
Merged
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
11 changes: 7 additions & 4 deletions client/tests/integration/tasks.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,15 +9,18 @@ async fn task_list(app: Arc<DivviupApi>, account: Account, client: DivviupClient
fixtures::task(&app, &account).await,
];
let response_tasks = client.tasks(account.id).await?;
assert_same_json_representation(&tasks, &response_tasks);
assert_eq!(tasks.len(), response_tasks.len());
for (task, response_task) in tasks.iter().zip(response_tasks.iter()) {
assert_same_json_representation_ignoring_query_type(task, response_task);
}
Ok(())
}

#[test(harness = with_configured_client)]
async fn get_task(app: Arc<DivviupApi>, account: Account, client: DivviupClient) -> TestResult {
let task = fixtures::task(&app, &account).await;
let response_task = client.task(&task.id).await?;
assert_same_json_representation(&task, &response_task);
assert_same_json_representation_ignoring_query_type(&task, &response_task);
Ok(())
}

Expand Down Expand Up @@ -45,7 +48,7 @@ async fn create_task(app: Arc<DivviupApi>, account: Account, client: DivviupClie
.one(app.db())
.await?
.unwrap();
assert_same_json_representation(&task_from_db, &response_task);
assert_same_json_representation_ignoring_query_type(&task_from_db, &response_task);
Ok(())
}

Expand Down Expand Up @@ -85,7 +88,7 @@ async fn create_task_time_bucketed_fixed_size(
.one(app.db())
.await?
.unwrap();
assert_same_json_representation(&task_from_db, &response_task);
assert_same_json_representation_ignoring_query_type(&task_from_db, &response_task);
Ok(())
}

Expand Down
7 changes: 7 additions & 0 deletions compose.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -104,6 +104,10 @@ services:
--name=helper --api-url=http://janus_2_aggregator:8080/aggregator-api \
--bearer-token=0000 && \
touch /tmp/done)
volumes:
- type: volume
source: pair_aggregator_state
target: /tmp
network_mode: service:divviup_api
depends_on:
divviup_api:
Expand Down Expand Up @@ -277,6 +281,9 @@ services:
CONFIG_FILE: /janus_2_garbage_collector.yaml
<<: *janus_environment

volumes:
pair_aggregator_state:

configs:
postgres_init:
content: |
Expand Down
12 changes: 10 additions & 2 deletions documentation/openapi.yml
Original file line number Diff line number Diff line change
Expand Up @@ -276,6 +276,9 @@ paths:
format: uuid
vdaf:
$ref: "#/components/schemas/Vdaf"
query_type:
type: string
enum: [TimeInterval, FixedSize]
min_batch_size:
type: number
max_batch_size:
Expand Down Expand Up @@ -694,6 +697,9 @@ components:
format: uuid
vdaf:
$ref: "#/components/schemas/Vdaf"
query_type:
type: string
enum: [TimeInterval, FixedSize]
min_batch_size:
type: number
max_batch_size:
Expand Down Expand Up @@ -929,8 +935,10 @@ components:
is_first_party:
type: boolean
query_types:
type: string
enum: [TimeInterval, FixedSize]
type: array
items:
type: string
enum: [TimeInterval, FixedSize]
Comment thread
divergentdave marked this conversation as resolved.
vdafs:
type: string
examples:
Expand Down
2 changes: 1 addition & 1 deletion migration/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ This is the standard migrator CLI that comes with SeaORM.

- Generate a new migration file
```sh
cargo run -- migrate generate MIGRATION_NAME
cargo run -- generate MIGRATION_NAME
```
- Apply all pending migrations
```sh
Expand Down
2 changes: 2 additions & 0 deletions migration/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ mod m20240214_215101_upload_metrics;
mod m20240411_195358_time_bucketed_fixed_size;
mod m20240416_172920_task_deleted_at;
mod m20250801_164739_aggregation_job_metrics;
mod m20260921_223229_add_query_type_column;

pub struct Migrator;

Expand Down Expand Up @@ -59,6 +60,7 @@ impl MigratorTrait for Migrator {
Box::new(m20240411_195358_time_bucketed_fixed_size::Migration),
Box::new(m20240416_172920_task_deleted_at::Migration),
Box::new(m20250801_164739_aggregation_job_metrics::Migration),
Box::new(m20260921_223229_add_query_type_column::Migration),
]
}
}
105 changes: 105 additions & 0 deletions migration/src/m20260921_223229_add_query_type_column.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,105 @@
use sea_orm::{sea_query::extension::postgres::Type, DbBackend};
use sea_orm_migration::prelude::*;

#[derive(DeriveMigrationName)]
pub struct Migration;

#[async_trait::async_trait]
impl MigrationTrait for Migration {
async fn up(&self, manager: &SchemaManager) -> Result<(), DbErr> {
let db = manager.get_connection();
// Use an enum in Postgres, and a string in SQLite.
let mut column_def = if db.get_database_backend() == DbBackend::Postgres {
manager
.create_type(
Type::create()
.as_enum(QueryType::Enum)
.values([QueryType::TimeInterval, QueryType::FixedSize])
.to_owned(),
)
.await?;
ColumnDef::new(Task::QueryType)
.custom(QueryType::Enum)
.to_owned()
} else {
ColumnDef::new(Task::QueryType).string().to_owned()
};
// Add the column as a nullable column.
manager
.alter_table(
Table::alter()
.table(Task::Table)
.add_column(column_def.null().default(Expr::null()))
.to_owned(),
)
.await?;
// Backfill the column.
manager
.execute(
Query::update()
.table(Task::Table)
.value(
Task::QueryType,
Expr::case(
Expr::column(Task::MaxBatchSize).is_not_null(),
Expr::cast_as(
Expr::value(QueryType::FixedSize.unquoted()),
QueryType::Enum,
),
)
.finally(Expr::cast_as(
Expr::value(QueryType::TimeInterval.unquoted()),
QueryType::Enum,
)),
)
.to_owned(),
)
.await?;
// Change the column to be not nullable.
manager
.alter_table(
Table::alter()
.table(Task::Table)
.modify_column(column_def.not_null())
.to_owned(),
)
.await
}

async fn down(&self, manager: &SchemaManager) -> Result<(), DbErr> {
manager
.alter_table(
Table::alter()
.table(Task::Table)
.drop_column(Task::QueryType)
.to_owned(),
)
.await?;
let db = manager.get_connection();
if db.get_database_backend() == DbBackend::Postgres {
manager
.drop_type(Type::drop().name(QueryType::Enum).to_owned())
.await?;
}
Ok(())
}
}

#[derive(Iden)]
enum Task {
Table,

MaxBatchSize,
QueryType,
}

#[derive(Iden)]
pub enum QueryType {
#[iden = "query_type"]
Enum,

#[iden = "TIME_INTERVAL"]
TimeInterval,
#[iden = "FIXED_SIZE"]
FixedSize,
}
3 changes: 2 additions & 1 deletion src/clients/aggregator_client/api_types.rs
Original file line number Diff line number Diff line change
Expand Up @@ -154,7 +154,8 @@ impl From<AggregatorVdaf> for Vdaf {
pub enum QueryType {
TimeInterval,
FixedSize {
max_batch_size: u64,
#[serde(skip_serializing_if = "Option::is_none")]
max_batch_size: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
batch_time_window_size: Option<u64>,
},
Expand Down
11 changes: 11 additions & 0 deletions src/entity/aggregator/query_type_name.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,8 @@ use std::{
str::FromStr,
};

use crate::entity::task;

/// https://www.ietf.org/archive/id/draft-ietf-ppm-dap-05.html#name-queries
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Hash, PartialOrd, Ord)]
pub enum QueryTypeName {
Expand Down Expand Up @@ -61,6 +63,15 @@ impl From<&str> for QueryTypeName {
}
}

impl From<task::QueryType> for QueryTypeName {
fn from(value: task::QueryType) -> Self {
match value {
task::QueryType::TimeInterval => Self::TimeInterval,
task::QueryType::FixedSize => Self::FixedSize,
}
}
}

#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, Eq)]
pub struct QueryTypeNameSet(Set<QueryTypeName>);

Expand Down
2 changes: 2 additions & 0 deletions src/entity/task.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,8 @@ mod provisionable_task;
pub use provisionable_task::ProvisionableTask;
pub mod model;
pub use model::*;
pub mod query_type;
pub use query_type::*;

pub const DEFAULT_EXPIRATION_DURATION: Duration = Duration::days(365);

Expand Down
5 changes: 3 additions & 2 deletions src/entity/task/model.rs
Original file line number Diff line number Diff line change
@@ -1,8 +1,8 @@
use crate::{
clients::aggregator_client::{api_types::TaskAggregationJobMetrics, TaskUploadMetrics},
entity::{
account, json::Json, membership, AccountColumn, Accounts, Aggregator, AggregatorColumn,
Aggregators, CollectorCredentialColumn, CollectorCredentials,
account, json::Json, membership, task::QueryType, AccountColumn, Accounts, Aggregator,
AggregatorColumn, Aggregators, CollectorCredentialColumn, CollectorCredentials,
},
};
use sea_orm::{
Expand All @@ -26,6 +26,7 @@ pub struct Model {
pub account_id: Uuid,
pub name: String,
pub vdaf: Json<Vdaf>,
pub query_type: QueryType,
pub min_batch_size: i64,
pub max_batch_size: Option<i64>,
pub batch_time_window_size_seconds: Option<i64>,
Expand Down
Loading
Loading