1use super::*;
2use ahri_tre_protocol::{ingest, request::ProtocolRequest, upload::UploadRequest};
3use ahri_tre_types::NcName;
4use std::time::Instant;
5
6struct UploadLifecycle<'a> {
7 source_cleanup_safe: &'a mut bool,
8 output_intent: Option<DatasetOutputIntent>,
9 source: ahri_tre_types::VersionId,
10 deadline: Instant,
11 classification: ahri_tre_types::IngestClassification,
12 cancelled: std::sync::Arc<std::sync::atomic::AtomicBool>,
13 monitor: Option<UploadMonitor>,
14 metadata_deadline: Option<ahri_tre_pgmeta::UploadMetadataDeadline>,
15}
16impl datafile_materialization::DatafileLifecycle for UploadLifecycle<'_> {
17 fn accepted(&mut self) {
18 *self.source_cleanup_safe = false;
19 }
20 fn cleanup_finished(&mut self, succeeded: bool) {
21 *self.source_cleanup_safe = succeeded;
22 }
23 fn output_intent(&self) -> Option<DatasetOutputIntent> {
24 self.output_intent
25 }
26 fn ordinary_upload_source(&self) -> Option<ahri_tre_types::VersionId> {
27 Some(self.source)
28 }
29 fn cleanup_started(&mut self) {
30 self.monitor.take();
31 self.metadata_deadline.take();
32 }
33 fn classification(&self) -> ahri_tre_types::IngestClassification {
34 self.classification.clone()
35 }
36 fn record_completion(
37 &mut self,
38 _: &ahri_tre_pgmeta::PgMetadataRepository<'_>,
39 _: &DatasetMaterialization,
40 ) -> Result<(), AppError> {
41 self.check_deadline()
42 }
43 fn checkpoint(
44 &mut self,
45 _: datafile_materialization::DatafileStage,
46 _: &datafile_materialization::ResolvedDatafileMaterialization,
47 ) -> Result<(), AppError> {
48 self.check_deadline()
49 }
50}
51impl UploadLifecycle<'_> {
52 fn check_deadline(&self) -> Result<(), AppError> {
53 if Instant::now() >= self.deadline
54 || self.cancelled.load(std::sync::atomic::Ordering::Acquire)
55 {
56 return Err(AppError::Validation("Upload deadline elapsed".into()));
57 }
58 Ok(())
59 }
60}
61impl AppService {
62 pub async fn ingest_upload(
65 &self,
66 session: &mut DataStoreSession,
67 status: ahri_tre_protocol::session::SessionStatusPayload,
68 upload: UploadRequest,
69 input: ahri_tre_runtime::upload::UploadInput,
70 ) -> Result<serde_json::Value, AppError> {
71 if matches!(upload.request.request, ProtocolRequest::IngestRedcap(_)) {
72 return self
73 .ingest_redcap_upload(session, status, upload, input, None)
74 .await;
75 }
76 self.ingest_upload_with_acquisition(session, status, upload, input, None)
77 .await
78 }
79 pub(super) async fn ingest_upload_with_acquisition(
80 &self,
81 session: &mut DataStoreSession,
82 status: ahri_tre_protocol::session::SessionStatusPayload,
83 mut upload: UploadRequest,
84 input: ahri_tre_runtime::upload::UploadInput,
85 acquisition: Option<ahri_tre_types::AcquisitionPath>,
86 ) -> Result<serde_json::Value, AppError> {
87 let output_intent = match &upload.request.request {
88 ProtocolRequest::IngestSql(request) => Some(DatasetOutputIntent::parse(
89 request.materialization.replace,
90 request.materialization.new_version.as_deref(),
91 )?),
92 _ => None,
93 };
94 let sql_source = if let ProtocolRequest::IngestSql(request) = &upload.request.request {
95 if acquisition.is_none()
96 || upload.metadata.filename != "source.parquet"
97 || upload.metadata.media_type != "application/vnd.apache.parquet"
98 {
99 return Err(AppError::Validation(
100 "SQL results require admitted acquisition".into(),
101 ));
102 }
103 use sha2::Digest;
104 let identity = match &request.source {
105 ingest::SqlSource::Stream { engine } => format!("{engine:?}"),
106 ingest::SqlSource::Remote { endpoint } => format!(
107 "{} endpoint sha256:{}",
108 endpoint.as_str().split(':').next().unwrap_or("sql"),
109 sha2::Sha256::digest(endpoint.as_str().as_bytes())
110 .iter()
111 .map(|byte| format!("{byte:02x}"))
112 .collect::<String>()
113 ),
114 _ => {
115 return Err(AppError::Validation(
116 "SQL source must not contain a local path".into(),
117 ));
118 }
119 };
120 let query_digest = sha2::Sha256::digest(request.sql.as_bytes())
121 .iter()
122 .map(|byte| format!("{byte:02x}"))
123 .collect::<String>();
124 let table = ingest::IngestDatasetFileRequest {
125 classification: request.classification.clone(),
126 session: request.session.clone(),
127 study: request.study.clone(),
128 domain: request.domain.clone(),
129 dataset: request.dataset.clone(),
130 source: ingest::IngestDatasetFileSource::Stream,
131 materialization: ingest::DatasetMaterializationOptions {
132 description: request.materialization.description.clone(),
133 replace: request.materialization.replace,
134 new_version: request.materialization.new_version.clone(),
135 },
136 parse: ingest::DatasetFileParseOptions {
137 format: Some("parquet".into()),
138 sheet: None,
139 header: None,
140 delimiter: None,
141 null_strings: Vec::new(),
142 json_format: None,
143 sample_rows: None,
144 vocabulary_threshold: Some(0),
146 force_metadata_updates: false,
147 },
148 };
149 upload.request.request = ProtocolRequest::IngestDatasetFile(table);
150 Some(format!("SQL {identity} query sha256:{query_digest}"))
151 } else {
152 None
153 };
154 let ahri_tre_runtime::upload::UploadInput {
155 source,
156 budgets,
157 deadline,
158 cancelled,
159 } = input;
160 let preflight_deadline = session
161 .dataset_admission_repository()
162 .map(|repository| repository.upload_deadline(deadline, cancelled.clone()))
163 .transpose()
164 .map_err(|error| core_error("bound upload preflight", error))?;
165 let explicit_study = match &upload.request.request {
166 ProtocolRequest::IngestFile(value) => value.study.clone(),
167 ProtocolRequest::IngestDatasetFile(value) => Some(value.study.clone()),
168 _ => None,
169 };
170 let selected = self
171 .resolve_session_study(
172 self.catalogue_datastore_id()?.as_uuid(),
173 explicit_study,
174 session.current_study().map(|study| study.0),
175 )
176 .await
177 .map_err(|error| match error.code {
178 ahri_tre_protocol::ProtocolErrorCode::Conflict => {
179 AppError::Conflict("Upload Study selectors disagree".into())
180 }
181 ahri_tre_protocol::ProtocolErrorCode::NotFound => {
182 AppError::NotFound("Upload Study is unavailable".into())
183 }
184 _ => AppError::Validation("Upload Study is invalid".into()),
185 })?;
186 let study = StudySelector::Id {
187 study_id: ahri_tre_types::StudyId(selected.id.as_uuid()),
188 };
189 let invalid = || AppError::Validation("Invalid upload request".into());
190 let (file, table) = match upload.request.request {
191 ProtocolRequest::IngestFile(request)
192 if matches!(request.source, ingest::IngestFileSource::Stream) =>
193 {
194 (request, None)
195 }
196 ProtocolRequest::IngestDatasetFile(request)
197 if matches!(request.source, ingest::IngestDatasetFileSource::Stream) =>
198 {
199 if sql_source.is_none()
200 && (request.materialization.replace
201 || request.materialization.new_version.is_some())
202 {
203 return Err(AppError::Validation("Dataset output versions are allocated automatically; replacement and explicit versions are not supported".into()));
204 }
205 let domain = self.catalogue_domain_selector(request.domain.clone())?;
207 self.get_domain(GetDomainRequest { domain })
208 .await?
209 .ok_or_else(invalid)?;
210 let name = derived_source_file_asset_name(
211 &NcName::parse(request.dataset.as_str().to_string()).map_err(|_| invalid())?,
212 );
213 (
214 ingest::IngestFileRequest {
215 classification: request.classification.clone(),
216 session: request.session.clone(),
217 study: Some(request.study.clone()),
218 asset: ahri_tre_protocol::PublicName::new(name.as_str())
219 .map_err(|_| invalid())?,
220 source: ingest::IngestFileSource::Stream,
221 format: request.parse.format.clone().ok_or_else(invalid)?,
222 description: Some(request.materialization.description.clone()),
223 new_version: false,
224 bump_major: false,
225 bump_minor: false,
226 compress: true,
227 encrypt: None,
228 hash_source_local: false,
229 verify_copy: true,
230 },
231 Some(request),
232 )
233 }
234 _ => return Err(invalid()),
235 };
236 let format = ahri_tre_protocol::upload::upload_format(
237 &upload.metadata.filename,
238 &file.format,
239 &upload.metadata.media_type,
240 )
241 .ok_or_else(invalid)?;
242 let source_format = resolve_dataset_file_format(None, Some(format))?;
243 let options = match &table {
244 Some(table) => dataset_file_read_options(
245 table.parse.header,
246 table.parse.delimiter,
247 table.parse.null_strings.clone(),
248 table.parse.json_format.clone(),
249 table.parse.sheet.clone(),
250 ),
251 None => DatasetFileReadOptions::default(),
252 };
253 let metadata = upload.metadata;
254 let mut content = IngestFileContent::new(
255 metadata.filename.clone(),
256 Some(metadata.media_type.clone()),
257 Some(metadata.content_length),
258 Some(metadata.digest.clone()),
259 CancelledReader {
260 source,
261 cancelled: cancelled.clone(),
262 deadline,
263 },
264 );
265 content.acquisition = acquisition;
266 content.upload_validation = Some(ahri_tre_lake::UploadValidation {
267 format: source_format,
268 options,
269 budgets,
270 deadline,
271 cancelled: cancelled.clone(),
272 });
273 drop(preflight_deadline);
274 let ingested = self
275 .ingest_datafile(
276 session,
277 IngestDataFileRequest {
278 classification: file.classification,
279 study: study.clone(),
280 asset_name: NcName::parse(file.asset.as_str().to_string())
281 .map_err(|_| invalid())?,
282 content,
283 format: format.to_owned(),
284 description: file.description,
285 new_version: file.new_version,
286 bump_major: file.bump_major,
287 bump_minor: file.bump_minor,
288 compress: file.compress,
289 encrypt: file.encrypt,
290 },
291 )
292 .await?;
293 let Some(table) = table else {
294 let response = crate::projections::datafile_ingest_response(
295 status,
296 ingested,
297 ingest::IngestSourceSummary {
298 kind: if acquisition.is_some() {
299 ingest::IngestSourceKind::HttpsUri
300 } else {
301 ingest::IngestSourceKind::UploadedContent
302 },
303 request_only: true,
304 upload: Some(ingest::UploadSummary {
305 acquisition,
306 budgets: Some(budgets),
307 filename: metadata.filename,
308 media_type: metadata.media_type,
309 content_length: metadata.content_length,
310 digest: metadata.digest,
311 outcome: ingest::UploadOutcome::Ingested,
312 }),
313 },
314 );
315 return serde_json::to_value(response).map_err(|_| invalid());
316 };
317 let mut source_cleanup_safe = true;
318 let materialized = async {
319 let source_version = ingested.datafile.datafile_id;
320 let metadata_domain = self.catalogue_domain_selector(table.domain)?;
321 let datastore_id = self.catalogue_datastore_id()?;
322 let request = DeriveDatasetFromManagedDataFileRequest {
323 study,
324 metadata_domain,
325 dataset_name: NcName::parse(table.dataset.as_str().to_string())
326 .map_err(|_| invalid())?,
327 source: ingest::ManagedDataFileSource {
328 asset: ahri_tre_protocol::asset::AssetSelector::Id {
329 asset: ingest::version_ref(datastore_id, source_version),
330 study: None,
331 asset_type: Some(ahri_tre_types::AssetType::File),
332 },
333 version: None,
334 },
335 input_format: table.parse.format,
336 header: table.parse.header,
337 delimiter: table.parse.delimiter,
338 null_strings: table.parse.null_strings,
339 json_format: table.parse.json_format,
340 sheet: table.parse.sheet,
341 description: Some(table.materialization.description.clone()),
342 version_note: Some(table.materialization.description),
343 replace: table.materialization.replace,
344 new_version: table.materialization.new_version,
345 transformation: ahri_tre_types::NewTransformationRecord {
346 transformation_type: ahri_tre_types::TransformationType::Transform,
347 description: format!(
348 "materialize ordinary-ingest table from managed Datafile; {}; {}",
349 sql_source.as_deref().unwrap_or("uploaded table"),
350 if sql_source.is_some() {
351 match acquisition {
352 Some(ahri_tre_types::AcquisitionPath::Trusted) => {
353 "Trusted-observed acquisition"
354 }
355 _ => "client-asserted acquisition",
356 }
357 } else {
358 acquisition_attribution(acquisition)
359 }
360 ),
361 repository_url: None,
362 commit_hash: None,
363 file_path: None,
364 date_created: None,
365 created_by: session.authenticated_actor().map(str::to_owned),
366 },
367 sample_rows: table.parse.sample_rows,
368 vocabulary_threshold: table.parse.vocabulary_threshold,
369 force_metadata_updates: table.parse.force_metadata_updates,
370 };
371 let metadata_deadline = session
372 .dataset_admission_repository()
373 .map(|repository| repository.upload_deadline(deadline, cancelled.clone()))
374 .transpose()
375 .map_err(|error| core_error("bound upload metadata", error))?;
376 let monitor = UploadMonitor::new(session, deadline, cancelled.clone());
377 let result = self
378 .derive_dataset_with_lifecycle(
379 session,
380 request,
381 &mut UploadLifecycle {
382 source_cleanup_safe: &mut source_cleanup_safe,
383 output_intent,
384 source: source_version,
385 deadline,
386 classification: table.classification,
387 cancelled,
388 monitor: Some(monitor),
389 metadata_deadline,
390 },
391 )
392 .await?;
393 Ok::<_, AppError>(result)
394 }
395 .await;
396 let result = match materialized {
397 Ok(result) => result,
398 Err(error) => {
399 if sql_source.is_some()
400 && source_cleanup_safe
401 && !matches!(error, AppError::DatasetAdmissionOutcomeUnknown(_))
402 {
403 self.rollback_uploaded_table_source(session, ingested)
404 .await?;
405 }
406 return Err(error);
407 }
408 };
409 let mut response = crate::projections::dataset_materialization_response(
410 status,
411 result.study.name.as_str().to_string(),
412 result.materialization,
413 );
414 response.source = ingest::DatasetMaterializationSourceSummary {
415 upload: Some(ingest::UploadSummary {
416 acquisition,
417 budgets: Some(budgets),
418 filename: metadata.filename,
419 media_type: metadata.media_type,
420 content_length: metadata.content_length,
421 digest: metadata.digest,
422 outcome: ingest::UploadOutcome::Ingested,
423 }),
424 kind: if sql_source.is_some() {
425 ingest::DatasetMaterializationSourceKind::SqlSource
426 } else if acquisition.is_some() {
427 ingest::DatasetMaterializationSourceKind::Https
428 } else {
429 ingest::DatasetMaterializationSourceKind::LocalDatasetFile
430 },
431 request_only: true,
432 };
433 serde_json::to_value(response).map_err(|_| invalid())
434 }
435
436 async fn rollback_uploaded_table_source(
439 &self,
440 session: &DataStoreSession,
441 ingested: crate::DataFileIngestResult,
442 ) -> Result<(), AppError> {
443 let _guard = Arc::clone(&FILE_INGEST_WORKFLOW_LOCK).lock_owned().await;
444 let asset = ingested.catalog.asset.asset_id;
445 let version = ingested.datafile.datafile_id;
446 let datafile = self
447 .assets
448 .get_datafile(version)
449 .await
450 .map_err(|error| core_error("read failed upload source", error))?
451 .ok_or_else(|| AppError::Infrastructure("Failed upload source is missing".into()))?;
452 let versions = self
453 .assets
454 .list_asset_versions(asset)
455 .await
456 .map_err(|error| core_error("read failed upload versions", error))?;
457 let latest = versions.iter().find(|item| item.is_latest == Some(true));
458 let promotion = if latest.is_some_and(|item| item.version_id == version) {
459 Some((
460 asset,
461 versions
462 .iter()
463 .filter(|item| item.version_id != version)
464 .max_by_key(|item| (item.major, item.minor, item.patch))
465 .map(|item| item.version_id),
466 ))
467 } else {
468 None
469 };
470 let lake = DuckLakeAdapter::new(&session.runtime().lake.data_path);
471 run_blocking_lake_work(move || lake.delete_datafile(&DeleteDataFileRequest { datafile }))
472 .await
473 .map_err(|error| infrastructure_error("remove failed upload source payload", error))?;
474 for transformation in self
475 .transformations
476 .list_transformations_producing(version)
477 .await
478 .map_err(|error| core_error("read failed upload provenance", error))?
479 {
480 self.transformations
481 .delete_transformation(transformation.transformation_id)
482 .await
483 .map_err(|error| core_error("remove failed upload provenance", error))?;
484 }
485 self.assets
486 .delete_datafile_version_metadata(version, promotion)
487 .await
488 .map_err(|error| core_error("restore failed upload source catalogue", error))?;
489 if versions.iter().all(|item| item.version_id == version) {
490 self.assets
491 .delete_asset_record(asset)
492 .await
493 .map_err(|error| core_error("remove failed upload source Asset", error))?;
494 }
495 Ok(())
496 }
497}
498
499struct CancelledReader {
500 source: Box<dyn std::io::Read + Send>,
501 cancelled: std::sync::Arc<std::sync::atomic::AtomicBool>,
502 deadline: Instant,
503}
504impl std::io::Read for CancelledReader {
505 fn read(&mut self, out: &mut [u8]) -> std::io::Result<usize> {
506 if Instant::now() >= self.deadline
507 || self.cancelled.load(std::sync::atomic::Ordering::Acquire)
508 {
509 return Err(std::io::Error::other("Upload stopped"));
510 }
511 self.source.read(out)
512 }
513}
514pub(super) struct UploadMonitor {
515 stop: std::sync::mpsc::Sender<()>,
516 worker: Option<std::thread::JoinHandle<()>>,
517}
518impl UploadMonitor {
519 pub(super) fn new(
520 session: &DataStoreSession,
521 deadline: Instant,
522 cancelled: std::sync::Arc<std::sync::atomic::AtomicBool>,
523 ) -> Self {
524 let interrupt = session.lake_connection().interrupt_handle();
525 let (stop, stopped) = std::sync::mpsc::channel();
526 let worker = std::thread::spawn(move || {
527 while stopped
528 .recv_timeout(std::time::Duration::from_millis(25))
529 .is_err()
530 {
531 if Instant::now() >= deadline
532 || cancelled.load(std::sync::atomic::Ordering::Acquire)
533 {
534 interrupt.interrupt();
535 }
536 }
537 });
538 Self {
539 stop,
540 worker: Some(worker),
541 }
542 }
543}
544impl Drop for UploadMonitor {
545 fn drop(&mut self) {
546 let _ = self.stop.send(());
547 if let Some(worker) = self.worker.take() {
548 let _ = worker.join();
549 }
550 }
551}
552
553fn acquisition_attribution(path: Option<ahri_tre_types::AcquisitionPath>) -> &'static str {
554 match path {
555 Some(ahri_tre_types::AcquisitionPath::Trusted) => {
556 "HTTPS source, Trusted-observed acquisition"
557 }
558 Some(ahri_tre_types::AcquisitionPath::Client) => {
559 "HTTPS source, client-asserted acquisition"
560 }
561 None => "client-uploaded source",
562 }
563}