Skip to main content

micromegas_analytics/lakehouse/
jit_partitions.rs

1use super::{
2    block_partition_spec::{BlockPartitionSpec, BlockProcessorMap},
3    blocks_view::BlocksView,
4    lakehouse_context::LakehouseContext,
5    partition_cache::{LivePartitionProvider, QueryPartitionProvider},
6    partition_source_data::{PartitionSourceBlock, SourceDataBlocksInMemory},
7    view::{View, ViewMetadata},
8};
9use crate::{
10    dfext::{
11        string_column_accessor::string_column_by_name,
12        typed_column::{get_single_row_primitive_value, typed_column_by_name},
13    },
14    lakehouse::{
15        partition_cache::PartitionCache, partition_source_data::hash_to_object_count,
16        query::query_partitions, view::PartitionSpec,
17    },
18    metadata::{ProcessMetadata, StreamMetadata, block_from_batch_row},
19    properties::properties_column_accessor::properties_column_by_name,
20    response_writer::ResponseWriter,
21    time::TimeRange,
22};
23use anyhow::{Context, Result};
24use chrono::DurationRound;
25use chrono::{DateTime, TimeDelta, Utc};
26use datafusion::arrow::array::{BinaryArray, GenericListArray, StringArray};
27use datafusion::arrow::datatypes::{Schema, TimestampNanosecondType};
28use micromegas_ingestion::data_lake_connection::DataLakeConnection;
29use micromegas_tracing::prelude::*;
30use sqlx::Row;
31use std::sync::Arc;
32use uuid::Uuid;
33
34/// Configuration for Just-In-Time (JIT) partition generation.
35pub struct JitPartitionConfig {
36    pub max_nb_objects: i64,
37    pub max_insert_time_slice: TimeDelta,
38}
39
40impl Default for JitPartitionConfig {
41    fn default() -> Self {
42        JitPartitionConfig {
43            max_nb_objects: 20 * 1024 * 1024,
44            max_insert_time_slice: TimeDelta::hours(1),
45        }
46    }
47}
48
49async fn get_insert_time_range(
50    lakehouse: Arc<LakehouseContext>,
51    blocks_view: &BlocksView,
52    query_time_range: &TimeRange,
53    stream: Arc<StreamMetadata>,
54) -> Result<Option<TimeRange>> {
55    // we would need a PartitionCache built from event time range and then filtered for insert time range
56    let part_provider = LivePartitionProvider::new(lakehouse.lake().db_pool.clone());
57    let partitions = part_provider
58        .fetch(
59            &blocks_view.get_view_set_name(),
60            &blocks_view.get_view_instance_id(),
61            Some(*query_time_range),
62            blocks_view.get_file_schema_hash(),
63        )
64        .await?;
65    let stream_id = &stream.stream_id;
66    let begin_range_iso = query_time_range.begin.to_rfc3339();
67    let end_range_iso = query_time_range.end.to_rfc3339();
68    let sql = format!(
69        r#"SELECT MIN(insert_time) as min_insert_time, MAX(insert_time) as max_insert_time
70        FROM source
71        WHERE stream_id = '{stream_id}'
72        AND begin_time <= '{end_range_iso}'
73        AND end_time >= '{begin_range_iso}';"#
74    );
75    let reader_factory = lakehouse.reader_factory().clone();
76    let rbs = query_partitions(
77        lakehouse.runtime().clone(),
78        reader_factory,
79        lakehouse.lake().blob_storage.inner(),
80        blocks_view.get_file_schema(),
81        Arc::new(partitions),
82        &sql,
83    )
84    .await?
85    .collect()
86    .await?;
87    if rbs.is_empty() {
88        return Ok(None);
89    }
90    if rbs[0].num_rows() == 0 {
91        return Ok(None);
92    }
93    let min_insert_time = get_single_row_primitive_value::<TimestampNanosecondType>(&rbs, 0)?;
94    let max_insert_time = get_single_row_primitive_value::<TimestampNanosecondType>(&rbs, 1)?;
95    Ok(Some(TimeRange::new(
96        DateTime::from_timestamp_nanos(min_insert_time),
97        DateTime::from_timestamp_nanos(max_insert_time),
98    )))
99}
100
101/// Generates a segment of JIT partitions.
102pub async fn generate_stream_jit_partitions_segment(
103    config: &JitPartitionConfig,
104    lakehouse: Arc<LakehouseContext>,
105    blocks_view: &BlocksView,
106    partitions: &PartitionCache,
107    insert_time_range: &TimeRange,
108    stream: Arc<StreamMetadata>,
109    process: Arc<ProcessMetadata>,
110) -> Result<Vec<SourceDataBlocksInMemory>> {
111    let partitions = partitions
112        .filter_insert_range(*insert_time_range)
113        .partitions;
114
115    let stream_id = &stream.stream_id;
116    let begin_range_iso = insert_time_range.begin.to_rfc3339();
117    let end_range_iso = insert_time_range.end.to_rfc3339();
118    let sql = format!(
119        r#"SELECT block_id, stream_id, process_id, begin_time, end_time, begin_ticks, end_ticks, nb_objects, object_offset, payload_size, insert_time, "streams.format"
120             FROM source
121             WHERE stream_id = '{stream_id}'
122             AND insert_time >= '{begin_range_iso}'
123             AND insert_time < '{end_range_iso}'
124             ORDER BY insert_time, block_id;"#
125    );
126
127    let reader_factory = lakehouse.reader_factory().clone();
128    let rbs = query_partitions(
129        lakehouse.runtime().clone(),
130        reader_factory,
131        lakehouse.lake().blob_storage.inner(),
132        blocks_view.get_file_schema(),
133        Arc::new(partitions),
134        &sql,
135    )
136    .await?
137    .collect()
138    .await?;
139
140    let mut partitions = vec![];
141    let mut partition_blocks = vec![];
142    let mut partition_nb_objects: i64 = 0;
143    for rb in rbs {
144        let format_column = string_column_by_name(&rb, "streams.format")?;
145        for ir in 0..rb.num_rows() {
146            let block = block_from_batch_row(&rb, ir).with_context(|| "block_from_batch_row")?;
147            let block_nb_objects = block.nb_objects as i64;
148            let format = format_column.value(ir)?.to_string();
149
150            // Check if adding this block would exceed the limit
151            if partition_nb_objects + block_nb_objects > config.max_nb_objects
152                && !partition_blocks.is_empty()
153            {
154                // Push current partition without this block
155                partitions.push(SourceDataBlocksInMemory {
156                    blocks: partition_blocks,
157                    block_ids_hash: partition_nb_objects.to_le_bytes().to_vec(),
158                });
159                // Start new partition with this block
160                partition_blocks = vec![Arc::new(PartitionSourceBlock {
161                    block,
162                    stream: stream.clone(),
163                    process: process.clone(),
164                    format,
165                })];
166                partition_nb_objects = block_nb_objects;
167            } else {
168                // Add block to current partition
169                partition_nb_objects += block_nb_objects;
170                partition_blocks.push(Arc::new(PartitionSourceBlock {
171                    block,
172                    stream: stream.clone(),
173                    process: process.clone(),
174                    format,
175                }));
176            }
177        }
178    }
179    if partition_nb_objects != 0 {
180        partitions.push(SourceDataBlocksInMemory {
181            blocks: partition_blocks,
182            block_ids_hash: partition_nb_objects.to_le_bytes().to_vec(),
183        });
184    }
185
186    Ok(partitions)
187}
188
189/// generate_stream_jit_partitions lists the partitiions that are needed to cover a time span
190/// these partitions may not exist or they could be out of date
191/// Generates JIT partitions for a given time range.
192pub async fn generate_stream_jit_partitions(
193    config: &JitPartitionConfig,
194    lakehouse: Arc<LakehouseContext>,
195    blocks_view: &BlocksView,
196    query_time_range: &TimeRange,
197    stream: Arc<StreamMetadata>,
198    process: Arc<ProcessMetadata>,
199) -> Result<Vec<SourceDataBlocksInMemory>> {
200    let insert_time_range = get_insert_time_range(
201        lakehouse.clone(),
202        blocks_view,
203        query_time_range,
204        stream.clone(),
205    )
206    .await?;
207    if insert_time_range.is_none() {
208        return Ok(vec![]);
209    }
210    let insert_time_range = insert_time_range.with_context(|| "missing insert_time_range")?;
211    let insert_time_range = TimeRange::new(
212        insert_time_range
213            .begin
214            .duration_trunc(config.max_insert_time_slice)?,
215        insert_time_range
216            .end
217            .duration_trunc(config.max_insert_time_slice)?
218            + config.max_insert_time_slice,
219    );
220    let segment_source_partitions = instrument_named!(
221        PartitionCache::fetch_overlapping_insert_range_for_view(
222            &lakehouse.lake().db_pool,
223            blocks_view.get_view_set_name(),
224            blocks_view.get_view_instance_id(),
225            insert_time_range,
226        ),
227        "fetch_overlapping_insert_range_for_view"
228    )
229    .await?;
230
231    let mut begin_segment = insert_time_range.begin;
232    let mut end_segment = begin_segment + config.max_insert_time_slice;
233    let mut partitions = vec![];
234    while end_segment <= insert_time_range.end {
235        let insert_time_range = TimeRange::new(begin_segment, end_segment);
236        let mut segment_partitions = generate_stream_jit_partitions_segment(
237            config,
238            lakehouse.clone(),
239            blocks_view,
240            &segment_source_partitions,
241            &insert_time_range,
242            stream.clone(),
243            process.clone(),
244        )
245        .await?;
246        partitions.append(&mut segment_partitions);
247        begin_segment = end_segment;
248        end_segment = begin_segment + config.max_insert_time_slice;
249    }
250    Ok(partitions)
251}
252
253/// Generates a segment of JIT partitions filtered by process.
254#[span_fn]
255pub async fn generate_process_jit_partitions_segment(
256    config: &JitPartitionConfig,
257    lakehouse: Arc<LakehouseContext>,
258    blocks_view: &BlocksView,
259    partitions: &PartitionCache,
260    insert_time_range: &TimeRange,
261    process: Arc<ProcessMetadata>,
262    stream_tag: &str,
263) -> Result<Vec<SourceDataBlocksInMemory>> {
264    let partitions = partitions
265        .filter_insert_range(*insert_time_range)
266        .partitions;
267
268    let process_id = &process.process_id;
269    let begin_range_iso = insert_time_range.begin.to_rfc3339();
270    let end_range_iso = insert_time_range.end.to_rfc3339();
271    let sql = format!(
272        r#"SELECT block_id, stream_id, process_id, begin_time, end_time, begin_ticks, end_ticks, nb_objects, object_offset, payload_size, insert_time,
273             "streams.dependencies_metadata", "streams.objects_metadata", "streams.tags", "streams.properties", "streams.format"
274             FROM source
275             WHERE process_id = '{process_id}'
276             AND array_has( "streams.tags", '{stream_tag}' )
277             AND insert_time >= '{begin_range_iso}'
278             AND insert_time < '{end_range_iso}'
279             ORDER BY insert_time, block_id;"#
280    );
281
282    let reader_factory = lakehouse.reader_factory().clone();
283    let df = instrument_named!(
284        query_partitions(
285            lakehouse.runtime().clone(),
286            reader_factory,
287            lakehouse.lake().blob_storage.inner(),
288            blocks_view.get_file_schema(),
289            Arc::new(partitions),
290            &sql,
291        ),
292        "query_partitions"
293    )
294    .await?;
295    let rbs = instrument_named!(df.collect(), "collect_partition_blocks").await?;
296
297    let mut partitions = vec![];
298    let mut partition_blocks = vec![];
299    let mut partition_nb_objects: i64 = 0;
300
301    for rb in rbs {
302        for ir in 0..rb.num_rows() {
303            let block = block_from_batch_row(&rb, ir).with_context(|| "block_from_batch_row")?;
304            let block_nb_objects = block.nb_objects as i64;
305
306            // Build StreamMetadata from the query results
307            let stream_id_column = string_column_by_name(&rb, "stream_id")?;
308            let stream_process_id_column = string_column_by_name(&rb, "process_id")?;
309            let dependencies_metadata_column: &BinaryArray =
310                typed_column_by_name(&rb, "streams.dependencies_metadata")?;
311            let objects_metadata_column: &BinaryArray =
312                typed_column_by_name(&rb, "streams.objects_metadata")?;
313            let stream_tags_column: &GenericListArray<i32> =
314                typed_column_by_name(&rb, "streams.tags")?;
315            let stream_properties_accessor = properties_column_by_name(&rb, "streams.properties")?;
316            let stream_format_column = string_column_by_name(&rb, "streams.format")?;
317
318            let stream_id = Uuid::parse_str(stream_id_column.value(ir)?)
319                .with_context(|| "parsing stream_id")?;
320            let stream_process_id = Uuid::parse_str(stream_process_id_column.value(ir)?)
321                .with_context(|| "parsing stream process_id")?;
322
323            let dependencies_metadata = dependencies_metadata_column.value(ir);
324            let objects_metadata = objects_metadata_column.value(ir);
325            let stream_tags = stream_tags_column
326                .value(ir)
327                .as_any()
328                .downcast_ref::<StringArray>()
329                .with_context(|| "casting stream_tags")?
330                .iter()
331                .map(|item| String::from(item.unwrap_or_default()))
332                .collect();
333
334            // Get pre-serialized JSONB properties directly from accessor
335            let stream_properties_jsonb = stream_properties_accessor.jsonb_value(ir)?;
336
337            let stream = Arc::new(StreamMetadata {
338                stream_id,
339                process_id: stream_process_id,
340                dependencies_metadata: ciborium::from_reader(dependencies_metadata)
341                    .with_context(|| "decoding dependencies_metadata")?,
342                objects_metadata: ciborium::from_reader(objects_metadata)
343                    .with_context(|| "decoding objects_metadata")?,
344                tags: stream_tags,
345                properties: Arc::new(stream_properties_jsonb),
346            });
347
348            let format = stream_format_column.value(ir)?.to_string();
349
350            // Check if adding this block would exceed the limit
351            if partition_nb_objects + block_nb_objects > config.max_nb_objects
352                && !partition_blocks.is_empty()
353            {
354                // Push current partition without this block
355                partitions.push(SourceDataBlocksInMemory {
356                    blocks: partition_blocks,
357                    block_ids_hash: partition_nb_objects.to_le_bytes().to_vec(),
358                });
359                // Start new partition with this block
360                partition_blocks = vec![Arc::new(PartitionSourceBlock {
361                    block,
362                    stream: stream.clone(),
363                    process: process.clone(),
364                    format,
365                })];
366                partition_nb_objects = block_nb_objects;
367            } else {
368                // Add block to current partition
369                partition_nb_objects += block_nb_objects;
370                partition_blocks.push(Arc::new(PartitionSourceBlock {
371                    block,
372                    stream: stream.clone(),
373                    process: process.clone(),
374                    format,
375                }));
376            }
377        }
378    }
379    if partition_nb_objects != 0 {
380        partitions.push(SourceDataBlocksInMemory {
381            blocks: partition_blocks,
382            block_ids_hash: partition_nb_objects.to_le_bytes().to_vec(),
383        });
384    }
385    Ok(partitions)
386}
387
388/// generate_process_jit_partitions lists the partitions that are needed to cover a time span for a specific process
389/// these partitions may not exist or they could be out of date
390/// Generates JIT partitions for a given time range filtered by process.
391#[span_fn]
392pub async fn generate_process_jit_partitions(
393    config: &JitPartitionConfig,
394    lakehouse: Arc<LakehouseContext>,
395    blocks_view: &BlocksView,
396    query_time_range: &TimeRange,
397    process: Arc<ProcessMetadata>,
398    stream_tag: &str,
399) -> Result<Vec<SourceDataBlocksInMemory>> {
400    // Get insert time range for all blocks in this process
401    let part_provider = LivePartitionProvider::new(lakehouse.lake().db_pool.clone());
402    let view_set_name = blocks_view.get_view_set_name();
403    let view_instance_id = blocks_view.get_view_instance_id();
404    let partitions = instrument_named!(
405        part_provider.fetch(
406            &view_set_name,
407            &view_instance_id,
408            Some(*query_time_range),
409            blocks_view.get_file_schema_hash(),
410        ),
411        "live_partition_provider_fetch"
412    )
413    .await?;
414
415    let process_id = &process.process_id;
416    let begin_range_iso = query_time_range.begin.to_rfc3339();
417    let end_range_iso = query_time_range.end.to_rfc3339();
418    let sql = format!(
419        r#"SELECT MIN(insert_time) as min_insert_time, MAX(insert_time) as max_insert_time
420        FROM source
421        WHERE process_id = '{process_id}'
422        AND array_has( "streams.tags", '{stream_tag}' )
423        AND begin_time <= '{end_range_iso}'
424        AND end_time >= '{begin_range_iso}';"#
425    );
426
427    let reader_factory = lakehouse.reader_factory().clone();
428    let df = instrument_named!(
429        query_partitions(
430            lakehouse.runtime().clone(),
431            reader_factory,
432            lakehouse.lake().blob_storage.inner(),
433            blocks_view.get_file_schema(),
434            Arc::new(partitions),
435            &sql,
436        ),
437        "query_partitions"
438    )
439    .await?;
440    let rbs = instrument_named!(df.collect(), "collect_insert_time_range").await?;
441
442    if rbs.is_empty() || rbs[0].num_rows() == 0 {
443        return Ok(vec![]);
444    }
445
446    let min_insert_time = get_single_row_primitive_value::<TimestampNanosecondType>(&rbs, 0)?;
447    let max_insert_time = get_single_row_primitive_value::<TimestampNanosecondType>(&rbs, 1)?;
448
449    if min_insert_time == 0 || max_insert_time == 0 {
450        return Ok(vec![]);
451    }
452
453    let insert_time_range = TimeRange::new(
454        DateTime::from_timestamp_nanos(min_insert_time)
455            .duration_trunc(config.max_insert_time_slice)?,
456        DateTime::from_timestamp_nanos(max_insert_time)
457            .duration_trunc(config.max_insert_time_slice)?
458            + config.max_insert_time_slice,
459    );
460
461    let segment_source_partitions = instrument_named!(
462        PartitionCache::fetch_overlapping_insert_range_for_view(
463            &lakehouse.lake().db_pool,
464            blocks_view.get_view_set_name(),
465            blocks_view.get_view_instance_id(),
466            insert_time_range,
467        ),
468        "fetch_overlapping_insert_range_for_view"
469    )
470    .await?;
471
472    let mut begin_segment = insert_time_range.begin;
473    let mut end_segment = begin_segment + config.max_insert_time_slice;
474    let mut partitions = vec![];
475
476    while end_segment <= insert_time_range.end {
477        let insert_time_range = TimeRange::new(begin_segment, end_segment);
478        let mut segment_partitions = generate_process_jit_partitions_segment(
479            config,
480            lakehouse.clone(),
481            blocks_view,
482            &segment_source_partitions,
483            &insert_time_range,
484            process.clone(),
485            stream_tag,
486        )
487        .await?;
488        partitions.append(&mut segment_partitions);
489        begin_segment = end_segment;
490        end_segment = begin_segment + config.max_insert_time_slice;
491    }
492    Ok(partitions)
493}
494
495/// is_jit_partition_up_to_date compares a partition spec with the partitions that exist to know if it should be recreated
496/// Checks if a JIT partition is up to date.
497#[span_fn]
498pub async fn is_jit_partition_up_to_date(
499    pool: &sqlx::PgPool,
500    view_meta: ViewMetadata,
501    spec: &SourceDataBlocksInMemory,
502) -> Result<bool> {
503    let (min_insert_time, max_insert_time) =
504        get_part_insert_time_range(spec).with_context(|| "get_event_time_range")?;
505    let desc = format!(
506        "[{}, {}] {} {}",
507        min_insert_time.to_rfc3339(),
508        max_insert_time.to_rfc3339(),
509        *view_meta.view_set_name,
510        *view_meta.view_instance_id,
511    );
512
513    // CRITICAL: Use inclusive inequalities (<=, >=) to prevent race conditions.
514    // With exclusive inequalities (<, >), identical time ranges never match, causing
515    // partitions to be unnecessarily recreated on every query, leading to non-deterministic
516    // results. See: https://github.com/madesroches/micromegas/issues/488
517    //
518    // ADDITIONAL FIX: For identical timestamps (min_insert_time == max_insert_time),
519    // we need exact equality matching to handle single-timestamp partitions correctly.
520    let rows = instrument_named!(
521        if min_insert_time == max_insert_time {
522            // For identical timestamps, look for exact matches
523            sqlx::query(
524                "SELECT file_schema_hash, source_data_hash
525             FROM lakehouse_partitions
526             WHERE view_set_name = $1
527             AND view_instance_id = $2
528             AND begin_insert_time = $3
529             AND end_insert_time = $3
530             ;",
531            )
532            .bind(&*view_meta.view_set_name)
533            .bind(&*view_meta.view_instance_id)
534            .bind(min_insert_time)
535        } else {
536            // For time ranges, use inclusive inequalities
537            sqlx::query(
538                "SELECT file_schema_hash, source_data_hash
539             FROM lakehouse_partitions
540             WHERE view_set_name = $1
541             AND view_instance_id = $2
542             AND begin_insert_time <= $3
543             AND end_insert_time >= $4
544             ;",
545            )
546            .bind(&*view_meta.view_set_name)
547            .bind(&*view_meta.view_instance_id)
548            .bind(max_insert_time)
549            .bind(min_insert_time)
550        }
551        .fetch_all(pool),
552        "sql_select_matching_partitions"
553    )
554    .await
555    .with_context(|| "fetching matching partitions")?;
556    if rows.len() != 1 {
557        debug!("{desc}: found {} partitions (expected 1)", rows.len());
558        for (i, row) in rows.iter().enumerate() {
559            let part_file_schema: Vec<u8> = row.try_get("file_schema_hash")?;
560            let part_source_data: Vec<u8> = row.try_get("source_data_hash")?;
561            let source_row_count = hash_to_object_count(&part_source_data)?;
562            debug!(
563                "{desc}: partition {}: file_schema_hash={:?}, source_rows={}",
564                i, part_file_schema, source_row_count
565            );
566        }
567        info!("{desc}: found {} partitions", rows.len());
568        return Ok(false);
569    }
570    let r = &rows[0];
571    let part_file_schema: Vec<u8> = r.try_get("file_schema_hash")?;
572    if part_file_schema != view_meta.file_schema_hash {
573        // this is dangerous because we could be creating a new partition smaller than the old one, which is not supported.
574        // let's make sure there is no old data loitering
575        warn!("{desc}: found matching partition with different file schema");
576        return Ok(false);
577    }
578    let part_source_data: Vec<u8> = r.try_get("source_data_hash")?;
579    let existing_count = hash_to_object_count(&part_source_data)?;
580    let required_count = hash_to_object_count(&spec.block_ids_hash)?;
581    if existing_count < required_count {
582        info!("{desc}: existing partition lacks source data: creating a new partition");
583        return Ok(false);
584    }
585    info!("{desc}: partition up to date");
586    Ok(true)
587}
588
589/// get_event_time_range returns the time range covered by a partition spec
590/// Returns the event time range covered by a partition spec.
591fn get_part_insert_time_range(
592    spec: &SourceDataBlocksInMemory,
593) -> Result<(DateTime<Utc>, DateTime<Utc>)> {
594    if spec.blocks.is_empty() {
595        anyhow::bail!("empty partition should not exist");
596    }
597    // blocks need to be sorted by (event & insert) time
598    let min_insert_time = spec.blocks[0].block.insert_time;
599    let max_insert_time = spec.blocks[spec.blocks.len() - 1].block.insert_time;
600    Ok((min_insert_time, max_insert_time))
601}
602
603/// Writes a partition from a set of blocks.
604///
605/// `block_processors` keys per-block dispatch on the source block's `format`;
606/// see `BlockPartitionSpec::write` for the unknown-format behavior.
607#[span_fn]
608pub async fn write_partition_from_blocks(
609    lake: Arc<DataLakeConnection>,
610    view_metadata: ViewMetadata,
611    schema: Arc<Schema>,
612    source_data: SourceDataBlocksInMemory,
613    block_processors: Arc<BlockProcessorMap>,
614) -> Result<()> {
615    if source_data.blocks.is_empty() {
616        anyhow::bail!("empty partition spec");
617    }
618    // blocks need to be sorted by (event & insert) time
619    let min_insert_time = source_data.blocks[0].block.insert_time;
620    let max_insert_time = source_data.blocks[source_data.blocks.len() - 1]
621        .block
622        .insert_time;
623    let block_spec = BlockPartitionSpec {
624        view_metadata,
625        schema,
626        insert_range: TimeRange::new(min_insert_time, max_insert_time),
627        source_data: Arc::new(source_data),
628        block_processors,
629    };
630    let null_response_writer = Arc::new(ResponseWriter::new(None));
631    block_spec
632        .write(lake, null_response_writer)
633        .await
634        .with_context(|| "block_spec.write")?;
635    Ok(())
636}