use crate::adapter::package::zip::{ZipMember, scan};
use crate::container::Descriptor;
use crate::error::{Error, Result};
#[cfg(feature = "docx")]
use crate::field::index::SEL_DOCX_MODEL;
#[cfg(feature = "epub")]
use crate::field::index::SEL_EPUB_MODEL;
#[cfg(feature = "odt")]
use crate::field::index::SEL_ODT_MODEL;
#[cfg(feature = "opc")]
use crate::field::index::SEL_OPC_MODEL;
use crate::field::index::{
FsIndexStore, IndexEntry, SEL_PACKAGE_MEMBER_DECODED, SEL_PACKAGE_MEMBER_RAW, SelectorKey,
build, validate,
};
use crate::field::ingest::{MAX_INGEST_INDEX_ENTRIES, MAX_INGEST_NODES, with_observation_index};
use crate::field::manifest::{ABSENT_ROOT, FieldRoot};
use crate::field::node::{NodeKind, SeedNode, object_params, span_params};
use crate::field::resource::{is_shareable_resource, resource_blob_node};
use crate::field::{FieldId, FieldStore, PACKAGE_UNIVERSE};
use crate::limits::Limits;
use crate::parallel::WorkerPool;
use crate::store::NodeId;
#[cfg(feature = "parallel")]
use rayon::prelude::*;
const FLAG_ENCRYPTED: u16 = 0x0001;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PackageIngestReport {
pub format: crate::field::document_format::DocumentFormat,
pub field: FieldId,
pub root_node: NodeId,
pub index_root: Option<NodeId>,
pub node_count: u64,
pub index_node_count: u64,
pub source_len: u64,
pub member_count: u64,
pub raw_nodes: u64,
pub decoded_nodes: u64,
pub declined_decodes: u64,
pub opc_model_nodes: u64,
pub docx_model_nodes: u64,
pub epub_model_nodes: u64,
pub odt_model_nodes: u64,
pub resource_blob_nodes: u64,
pub shared_resource_ids: u64,
pub shared_resource_bytes: u64,
pub nodes_id_shared: u64,
pub seed_bytes_written: u64,
pub index_bytes_written: u64,
}
fn charge_node(node_count: &mut u64) -> Result<()> {
if *node_count >= MAX_INGEST_NODES {
return Err(Error::resource_limit(format!(
"package ingest would exceed {MAX_INGEST_NODES} seed nodes"
)));
}
*node_count += 1;
Ok(())
}
#[derive(Debug, Default)]
struct ShareCounters {
resource_blob_nodes: u64,
shared_resource_ids: u64,
shared_resource_bytes: u64,
nodes_id_shared: u64,
seed_bytes_written: u64,
}
fn put_counted(
store: &mut FieldStore,
node: &SeedNode,
counters: &mut ShareCounters,
) -> Result<(NodeId, bool)> {
let enc = encode_seed(node);
let preexisting = publish_counted(store, &enc, counters)?;
Ok((enc.id, preexisting))
}
struct Encoded {
id: NodeId,
canonical: Vec<u8>,
}
fn encode_seed(node: &SeedNode) -> Encoded {
let canonical = node.encode_canonical();
let id = NodeId::of_node(&canonical);
Encoded { id, canonical }
}
fn publish_counted(
store: &mut FieldStore,
enc: &Encoded,
counters: &mut ShareCounters,
) -> Result<bool> {
let preexisting = store.seeds().contains_node(&enc.id)?;
store.seeds_mut().put_node(&enc.canonical)?;
if preexisting {
counters.nodes_id_shared += 1;
} else {
counters.seed_bytes_written = counters
.seed_bytes_written
.saturating_add(enc.canonical.len() as u64);
}
Ok(preexisting)
}
fn push_entry(entries: &mut Vec<IndexEntry>, entry: IndexEntry) -> Result<()> {
if entries.len() >= MAX_INGEST_INDEX_ENTRIES {
return Err(Error::resource_limit(format!(
"package ingest would exceed {MAX_INGEST_INDEX_ENTRIES} index entries"
)));
}
entries.push(entry);
Ok(())
}
struct MemberEncoded {
raw: Encoded,
blob: Option<(Encoded, u64)>,
decoded: Option<Encoded>,
}
fn encode_member(source: &[u8], member: &ZipMember) -> MemberEncoded {
let ordinal = member.id.ordinal;
let (data_off, data_len) = member.data;
let mut raw = SeedNode::new(
NodeKind::PackageMemberRaw,
data_len,
span_params(data_off, data_len),
Vec::new(),
"pkg:member-raw",
);
raw.limits.max_output_bytes = raw.limits.max_output_bytes.max(data_len);
let raw = encode_seed(&raw);
let blob = if member.method == 0 {
match source.get(data_off as usize..(data_off + data_len) as usize) {
Some(bytes) if is_shareable_resource(bytes) => {
let node = resource_blob_node(bytes);
Some((encode_seed(&node), bytes.len() as u64))
}
_ => None,
}
} else {
None
};
let decodable = member.method == 0 || member.method == 8;
let encrypted = member.flags & FLAG_ENCRYPTED != 0;
let decoded = if !decodable || encrypted || data_len == 0 {
None
} else {
let (decoded_params, decoded_deps) = match &blob {
Some((enc, _)) => (object_params(0, member.method, 0), vec![enc.id]),
None => (object_params(ordinal, member.method, 0), vec![raw.id]),
};
let mut node = SeedNode::new(
NodeKind::PackageMemberDecoded,
member.uncompressed_size,
decoded_params,
decoded_deps,
"pkg:member-decoded",
);
node.limits.max_output_bytes = node.limits.max_output_bytes.max(member.uncompressed_size);
Some(encode_seed(&node))
};
MemberEncoded { raw, blob, decoded }
}
fn encode_members(
pool: Option<&WorkerPool>,
source: &[u8],
members: &[ZipMember],
) -> Vec<MemberEncoded> {
let compute = |m: &ZipMember| encode_member(source, m);
#[cfg(feature = "parallel")]
if let Some(p) = pool.filter(|p| p.workers() > 1) {
return p.install(|| members.par_iter().map(compute).collect());
}
#[cfg(not(feature = "parallel"))]
let _ = pool;
members.iter().map(compute).collect()
}
pub fn ingest_package(
store: &mut FieldStore,
descriptor_bytes: &[u8],
limits: Limits,
) -> Result<PackageIngestReport> {
ingest_package_with(store, descriptor_bytes, limits, None)
}
pub fn ingest_package_with(
store: &mut FieldStore,
descriptor_bytes: &[u8],
limits: Limits,
pool: Option<&WorkerPool>,
) -> Result<PackageIngestReport> {
let observable = with_observation_index(descriptor_bytes, limits)?;
let parsed = Descriptor::parse(&observable, limits)?;
let source = crate::materialize::materialize(&parsed, limits)?;
let source_len = source.len() as u64;
let detected_format = crate::field::document_format::detect_document_format(&source, limits);
let descriptor_id = store.put_descriptor(&observable)?;
let physical = scan(&source, limits)?;
physical.validate(source_len)?;
physical.reemits(&source)?;
let mut root = SeedNode::new(
NodeKind::PackageRoot,
source_len,
Vec::new(),
Vec::new(),
"pkg:root",
);
root.limits.max_output_bytes = root.limits.max_output_bytes.max(source_len);
let mut share = ShareCounters::default();
let (root_id, _) = put_counted(store, &root, &mut share)?;
let mut entries: Vec<IndexEntry> = Vec::new();
let mut node_count: u64 = 1;
let mut raw_nodes: u64 = 0;
let mut decoded_nodes: u64 = 0;
let mut declined_decodes: u64 = 0;
let opc_model_nodes: u64;
let docx_model_nodes: u64;
let epub_model_nodes: u64;
let odt_model_nodes: u64;
let opc_model_id: Option<NodeId>;
let encoded = encode_members(pool, &source, &physical.members);
for (member, enc) in physical.members.iter().zip(encoded.iter()) {
let ordinal = member.id.ordinal;
let (data_off, data_len) = member.data;
charge_node(&mut node_count)?;
publish_counted(store, &enc.raw, &mut share)?;
raw_nodes += 1;
push_entry(
&mut entries,
IndexEntry {
key: SelectorKey::new(SEL_PACKAGE_MEMBER_RAW, ordinal),
out_off: data_off,
out_len: data_len,
node_id: enc.raw.id,
},
)?;
if let Some((blob, blob_len)) = &enc.blob {
let preexisting = publish_counted(store, blob, &mut share)?;
share.resource_blob_nodes += 1;
if preexisting {
share.shared_resource_ids += 1;
share.shared_resource_bytes = share.shared_resource_bytes.saturating_add(*blob_len);
}
}
let Some(decoded) = &enc.decoded else {
declined_decodes += 1;
continue;
};
charge_node(&mut node_count)?;
publish_counted(store, decoded, &mut share)?;
decoded_nodes += 1;
push_entry(
&mut entries,
IndexEntry {
key: SelectorKey::new(SEL_PACKAGE_MEMBER_DECODED, ordinal),
out_off: data_off,
out_len: data_len,
node_id: decoded.id,
},
)?;
}
#[cfg(feature = "opc")]
{
let mut model = SeedNode::new(
NodeKind::PackageOpcModel,
limits.max_output_bytes,
Vec::new(),
vec![root_id],
"pkg:opc-model",
);
model.limits.max_output_bytes = limits.max_output_bytes;
charge_node(&mut node_count)?;
let (model_id, _) = put_counted(store, &model, &mut share)?;
push_entry(
&mut entries,
IndexEntry {
key: SelectorKey::new(SEL_OPC_MODEL, 0),
out_off: 0,
out_len: 0,
node_id: model_id,
},
)?;
opc_model_id = Some(model_id);
opc_model_nodes = 1;
}
#[cfg(not(feature = "opc"))]
{
opc_model_id = None;
opc_model_nodes = 0;
}
#[cfg(feature = "docx")]
{
let model_id = opc_model_id
.ok_or_else(|| Error::internal_invariant("docx requires the OPC model node"))?;
let mut docx_model = SeedNode::new(
NodeKind::DocxModel,
limits.max_output_bytes,
Vec::new(),
vec![model_id],
"pkg:docx-model",
);
docx_model.limits.max_output_bytes = limits.max_output_bytes;
charge_node(&mut node_count)?;
let (docx_id, _) = put_counted(store, &docx_model, &mut share)?;
push_entry(
&mut entries,
IndexEntry {
key: SelectorKey::new(SEL_DOCX_MODEL, 0),
out_off: 0,
out_len: 0,
node_id: docx_id,
},
)?;
docx_model_nodes = 1;
}
#[cfg(not(feature = "docx"))]
{
let _ = opc_model_id;
docx_model_nodes = 0;
}
#[cfg(feature = "epub")]
{
let mut epub_model = SeedNode::new(
NodeKind::EpubModel,
limits.max_output_bytes,
Vec::new(),
vec![root_id],
"pkg:epub-model",
);
epub_model.limits.max_output_bytes = limits.max_output_bytes;
charge_node(&mut node_count)?;
let (epub_id, _) = put_counted(store, &epub_model, &mut share)?;
push_entry(
&mut entries,
IndexEntry {
key: SelectorKey::new(SEL_EPUB_MODEL, 0),
out_off: 0,
out_len: 0,
node_id: epub_id,
},
)?;
epub_model_nodes = 1;
}
#[cfg(not(feature = "epub"))]
{
epub_model_nodes = 0;
}
#[cfg(feature = "odt")]
{
let mut odt_model = SeedNode::new(
NodeKind::OdtModel,
limits.max_output_bytes,
Vec::new(),
vec![root_id],
"pkg:odt-model",
);
odt_model.limits.max_output_bytes = limits.max_output_bytes;
charge_node(&mut node_count)?;
let (odt_id, _) = put_counted(store, &odt_model, &mut share)?;
push_entry(
&mut entries,
IndexEntry {
key: SelectorKey::new(SEL_ODT_MODEL, 0),
out_off: 0,
out_len: 0,
node_id: odt_id,
},
)?;
odt_model_nodes = 1;
}
#[cfg(not(feature = "odt"))]
{
odt_model_nodes = 0;
}
let index_before = dir_bytes(&store.root().join("index"));
let (index_root, index_node_count) = if entries.is_empty() {
(None, 0)
} else {
let mut istore = FsIndexStore::open(store.root())?;
let root = build(&mut istore, &entries)?;
let (count, _depth) = validate(&istore, &root)?;
(Some(root), count)
};
let index_bytes_written = dir_bytes(&store.root().join("index")).saturating_sub(index_before);
let manifest = FieldRoot {
universe_id: crate::container::universe_id_from_str(PACKAGE_UNIVERSE),
source_sha256: parsed.descriptor.source_sha256,
source_len,
descriptor_id: *descriptor_id.as_bytes(),
root_node: root_id,
index_root: index_root.map_or(ABSENT_ROOT, |r| *r.as_bytes()),
node_count,
index_node_count,
provenance: format!(
"{}id_shared={};res_shared={};field:package;members={};raw={};decoded={};declined={};opc={};docx={};epub={};odt={};resource_blobs={}",
detected_format.provenance_prefix(),
share.nodes_id_shared,
share.shared_resource_ids,
physical.members.len(),
raw_nodes,
decoded_nodes,
declined_decodes,
opc_model_nodes,
docx_model_nodes,
epub_model_nodes,
odt_model_nodes,
share.resource_blob_nodes,
),
};
let field = store.put_field(&manifest)?;
Ok(PackageIngestReport {
format: detected_format,
field,
root_node: root_id,
index_root,
node_count,
index_node_count,
source_len,
member_count: physical.members.len() as u64,
raw_nodes,
decoded_nodes,
declined_decodes,
opc_model_nodes,
docx_model_nodes,
epub_model_nodes,
odt_model_nodes,
resource_blob_nodes: share.resource_blob_nodes,
shared_resource_ids: share.shared_resource_ids,
shared_resource_bytes: share.shared_resource_bytes,
nodes_id_shared: share.nodes_id_shared,
seed_bytes_written: share.seed_bytes_written,
index_bytes_written,
})
}
fn dir_bytes(root: &std::path::Path) -> u64 {
let Ok(entries) = std::fs::read_dir(root) else {
return 0;
};
let mut total = 0u64;
for entry in entries.flatten() {
match entry.metadata() {
Ok(meta) if meta.is_file() => total = total.saturating_add(meta.len()),
Ok(meta) if meta.is_dir() => total = total.saturating_add(dir_bytes(&entry.path())),
_ => {}
}
}
total
}