@@ -41,6 +41,13 @@ use crate::{
4141 } ,
4242} ;
4343
44+ fn get_version ( metadata : & serde_json:: Value ) -> Option < & str > {
45+ metadata
46+ . as_object ( )
47+ . and_then ( |meta| meta. get ( "version" ) )
48+ . and_then ( |version| version. as_str ( ) )
49+ }
50+
4451/// Migrate the metdata from v1 or v2 to v3
4552/// This is a one time migration
4653pub async fn run_metadata_migration (
@@ -55,13 +62,6 @@ pub async fn run_metadata_migration(
5562 }
5663 let staging_metadata = get_staging_metadata ( config) ?;
5764
58- fn get_version ( metadata : & serde_json:: Value ) -> Option < & str > {
59- metadata
60- . as_object ( )
61- . and_then ( |meta| meta. get ( "version" ) )
62- . and_then ( |version| version. as_str ( ) )
63- }
64-
6565 // if storage metadata is none do nothing
6666 if let Some ( storage_metadata) = storage_metadata {
6767 match get_version ( & storage_metadata) {
@@ -117,37 +117,42 @@ pub async fn run_metadata_migration(
117117
118118 // if staging metadata is none do nothing
119119 if let Some ( staging_metadata) = staging_metadata {
120- match get_version ( & staging_metadata) {
121- Some ( "v1" ) => {
122- let mut metadata = metadata_migration:: v1_v3 ( staging_metadata) ;
123- metadata = metadata_migration:: v3_v4 ( metadata) ;
124- put_staging_metadata ( config, & metadata) ?;
125- }
126- Some ( "v2" ) => {
127- let mut metadata = metadata_migration:: v2_v3 ( staging_metadata) ;
128- metadata = metadata_migration:: v3_v4 ( metadata) ;
129- put_staging_metadata ( config, & metadata) ?;
130- }
131- Some ( "v3" ) => {
132- let metadata = metadata_migration:: v3_v4 ( staging_metadata) ;
133- put_staging_metadata ( config, & metadata) ?;
134- }
135- Some ( "v4" ) => {
136- let metadata = metadata_migration:: v4_v5 ( staging_metadata) ;
137- let metadata = metadata_migration:: v5_v6 ( metadata) ;
138- put_staging_metadata ( config, & metadata) ?;
139- }
140- Some ( "v5" ) => {
141- let metadata = metadata_migration:: v5_v6 ( staging_metadata) ;
142- put_staging_metadata ( config, & metadata) ?;
143- }
144- _ => ( ) ,
145- }
120+ migrate_staging ( config, staging_metadata) ?;
146121 }
147122
148123 Ok ( ( ) )
149124}
150125
126+ fn migrate_staging ( config : & Parseable , staging_metadata : Value ) -> anyhow:: Result < ( ) > {
127+ match get_version ( & staging_metadata) {
128+ Some ( "v1" ) => {
129+ let mut metadata = metadata_migration:: v1_v3 ( staging_metadata) ;
130+ metadata = metadata_migration:: v3_v4 ( metadata) ;
131+ put_staging_metadata ( config, & metadata) ?;
132+ }
133+ Some ( "v2" ) => {
134+ let mut metadata = metadata_migration:: v2_v3 ( staging_metadata) ;
135+ metadata = metadata_migration:: v3_v4 ( metadata) ;
136+ put_staging_metadata ( config, & metadata) ?;
137+ }
138+ Some ( "v3" ) => {
139+ let metadata = metadata_migration:: v3_v4 ( staging_metadata) ;
140+ put_staging_metadata ( config, & metadata) ?;
141+ }
142+ Some ( "v4" ) => {
143+ let metadata = metadata_migration:: v4_v5 ( staging_metadata) ;
144+ let metadata = metadata_migration:: v5_v6 ( metadata) ;
145+ put_staging_metadata ( config, & metadata) ?;
146+ }
147+ Some ( "v5" ) => {
148+ let metadata = metadata_migration:: v5_v6 ( staging_metadata) ;
149+ put_staging_metadata ( config, & metadata) ?;
150+ }
151+ _ => ( ) ,
152+ }
153+ Ok ( ( ) )
154+ }
155+
151156/// run the migration for all streams concurrently
152157pub async fn run_migration ( config : & Parseable ) -> anyhow:: Result < ( ) > {
153158 let storage = config. storage . get_object_store ( ) ;
0 commit comments