name: FileProcessing
description: |
Async file processing workflow including:
- Image thumbnail generation
- Video thumbnail generation
- Document preview generation
- File compression
- Virus scanning
- Content deduplication checking
- CDN URL generation
version: 2
trigger:
event: ProcessingJobCreatedEvent
extract:
job_id: "event.job_id"
file_id: "event.file_id"
job_type: "event.job_type"
config:
timeout: 30m
persistence: true
retry:
max_attempts: 3
backoff: exponential
context:
job: null
file: null
bucket: null
result: null
steps:
- name: load_job
type: action
action: query
entity: ProcessingJob
params:
id: "{{ context.job_id }}"
on_success:
set:
job: "{{ result }}"
next: mark_running
on_failure:
next: job_not_found
- name: mark_running
type: action
action: update
entity: ProcessingJob
params:
id: "{{ context.job.id }}"
status: running
started_at: "{{ now() }}"
on_success:
next: load_file
on_failure:
next: mark_failed
- name: load_file
type: action
action: query
entity: StoredFile
params:
id: "{{ context.job.file_id }}"
on_success:
set:
file: "{{ result }}"
next: load_bucket
on_failure:
next: file_not_found
- name: load_bucket
type: action
action: query
entity: Bucket
params:
id: "{{ context.file.bucket_id }}"
on_success:
set:
bucket: "{{ result }}"
next: route_by_job_type
- name: route_by_job_type
type: condition
conditions:
- if: "context.job.job_type == 'thumbnail_generation'"
next: process_image_thumbnail
- if: "context.job.job_type == 'video_thumbnail'"
next: process_video_thumbnail
- if: "context.job.job_type == 'document_preview'"
next: process_document_preview
- if: "context.job.job_type == 'compression'"
next: process_compression
- if: "context.job.job_type == 'virus_scan'"
next: process_virus_scan
- if: "context.job.job_type == 'deduplication_check'"
next: process_deduplication
- else: true
next: unknown_job_type
- name: process_image_thumbnail
type: action
action: generate_image_thumbnail
params:
storage_key: "{{ context.file.storage_key }}"
backend: "{{ context.bucket.storage_backend }}"
mime_type: "{{ context.file.mime_type }}"
sizes: [256, 512]
timeout: 30s
on_success:
set:
result: "{{ result }}"
next: store_image_thumbnail
on_failure:
next: retry_or_fail
- name: store_image_thumbnail
type: action
action: store_file
params:
key: "{{ context.file.storage_key }}_thumb"
content: "{{ context.result.thumbnail }}"
backend: "{{ context.bucket.storage_backend }}"
on_success:
next: update_file_thumbnail
- name: update_file_thumbnail
type: action
action: update
entity: StoredFile
params:
id: "{{ context.file.id }}"
has_thumbnail: true
thumbnail_path: "{{ context.file.storage_key }}_thumb"
processing_status: thumbnails_ready
on_success:
next: mark_complete
- name: process_video_thumbnail
type: action
action: extract_video_frame
params:
storage_key: "{{ context.file.storage_key }}"
backend: "{{ context.bucket.storage_backend }}"
timestamp: "00:00:05"
count: 3
timeout: 120s
on_success:
set:
result: "{{ result }}"
next: store_video_thumbnail
on_failure:
next: retry_or_fail
- name: store_video_thumbnail
type: action
action: store_file
params:
key: "{{ context.file.storage_key }}_video_thumb"
content: "{{ context.result.thumbnail }}"
backend: "{{ context.bucket.storage_backend }}"
on_success:
next: update_file_video_thumbnail
- name: update_file_video_thumbnail
type: action
action: update
entity: StoredFile
params:
id: "{{ context.file.id }}"
has_video_thumbnail: true
thumbnail_path: "{{ context.file.storage_key }}_video_thumb"
processing_status: thumbnails_ready
on_success:
next: mark_complete
- name: process_document_preview
type: action
action: generate_document_preview
params:
storage_key: "{{ context.file.storage_key }}"
backend: "{{ context.bucket.storage_backend }}"
mime_type: "{{ context.file.mime_type }}"
max_pages: 10
dpi: 150
timeout: 180s
on_success:
set:
result: "{{ result }}"
next: store_document_preview
on_failure:
next: retry_or_fail
- name: store_document_preview
type: action
action: store_file
params:
key: "{{ context.file.storage_key }}_preview"
content: "{{ context.result.preview }}"
backend: "{{ context.bucket.storage_backend }}"
on_success:
next: update_file_document_preview
- name: update_file_document_preview
type: action
action: update
entity: StoredFile
params:
id: "{{ context.file.id }}"
has_document_preview: true
preview_path: "{{ context.file.storage_key }}_preview"
on_success:
next: mark_complete
- name: process_compression
type: action
action: compress_file
params:
storage_key: "{{ context.file.storage_key }}"
backend: "{{ context.bucket.storage_backend }}"
mime_type: "{{ context.file.mime_type }}"
algorithm: gzip
timeout: 300s
on_success:
set:
result: "{{ result }}"
next: store_compressed_file
- name: store_compressed_file
type: action
action: store_file
params:
key: "{{ context.file.storage_key }}_compressed"
content: "{{ context.result.compressed }}"
backend: "{{ context.bucket.storage_backend }}"
on_success:
next: update_file_compression
- name: update_file_compression
type: action
action: update
entity: StoredFile
params:
id: "{{ context.file.id }}"
is_compressed: true
original_size: "{{ context.file.size_bytes }}"
size_bytes: "{{ context.result.compressed_size }}"
compression_algorithm: "{{ context.result.algorithm }}"
storage_key: "{{ context.file.storage_key }}_compressed"
on_success:
next: update_bucket_compression_stats
- name: update_bucket_compression_stats
type: action
action: update
entity: Bucket
params:
id: "{{ context.bucket.id }}"
total_size_bytes: "{{ context.bucket.total_size_bytes - (context.file.original_size - context.result.compressed_size) }}"
on_success:
next: mark_complete
- name: process_virus_scan
type: action
action: scan_file_for_threats
params:
storage_key: "{{ context.file.storage_key }}"
backend: "{{ context.bucket.storage_backend }}"
filename: "{{ context.file.original_name }}"
timeout: 300s
on_success:
set:
result: "{{ result }}"
next: check_threat_level
- name: check_threat_level
type: condition
conditions:
- if: "context.result.threat_level == 'safe' || context.result.threat_level == null"
next: update_scan_safe
- if: "context.result.threat_level IN ['low', 'medium']"
next: update_scan_warning
- else: true
next: update_scan_threat
- name: update_scan_safe
type: action
action: update
entity: StoredFile
params:
id: "{{ context.file.id }}"
is_scanned: true
scan_result: "{{ context.result.details }}"
threat_level: safe
processing_status: scan_complete
on_success:
next: mark_complete
- name: update_scan_warning
type: action
action: update
entity: StoredFile
params:
id: "{{ context.file.id }}"
is_scanned: true
scan_result: "{{ context.result.details }}"
threat_level: "{{ context.result.threat_level }}"
on_success:
next: notify_scan_warning
- name: notify_scan_warning
type: action
action: send_notification
params:
template: scan_warning
recipients: "[admin, '{{ context.file.owner_id }}']"
data:
file_id: "{{ context.file.id }}"
threat_level: "{{ context.result.threat_level }}"
on_success:
next: mark_complete
- name: update_scan_threat
type: action
action: update
entity: StoredFile
params:
id: "{{ context.file.id }}"
is_scanned: true
scan_result: "{{ context.result.details }}"
threat_level: "{{ context.result.threat_level }}"
status: quarantined
processing_status: failed
on_success:
next: notify_threat_detected
- name: notify_threat_detected
type: action
action: send_notification
params:
template: threat_detected
recipients: "[admin]"
data:
file_id: "{{ context.file.id }}"
threats: "{{ context.result.threats }}"
on_success:
next: mark_complete
- name: process_deduplication
type: action
action: calculate_file_hash
params:
storage_key: "{{ context.file.storage_key }}"
backend: "{{ context.bucket.storage_backend }}"
algorithm: sha256
timeout: 600s
on_success:
set:
result: "{{ result }}"
next: check_existing_hash
- name: check_existing_hash
type: action
action: query
entity: ContentHash
params:
hash: "{{ context.result.hash }}"
on_success:
set:
content_hash: "{{ result }}"
next: update_existing_hash
on_failure:
next: create_new_hash
- name: update_existing_hash
type: action
action: update
entity: ContentHash
params:
id: "{{ context.content_hash.id }}"
reference_count: "{{ context.content_hash.reference_count + 1 }}"
last_referenced_at: "{{ now() }}"
on_success:
next: link_file_hash
- name: link_file_hash
type: action
action: update
entity: StoredFile
params:
id: "{{ context.file.id }}"
content_hash_id: "{{ context.content_hash.id }}"
checksum: "{{ context.result.hash }}"
on_success:
next: mark_complete
- name: create_new_hash
type: action
action: create
entity: ContentHash
params:
hash: "{{ context.result.hash }}"
size_bytes: "{{ context.file.size_bytes }}"
storage_key: "{{ context.file.storage_key }}"
storage_backend: "{{ context.bucket.storage_backend }}"
reference_count: 1
on_success:
set:
content_hash: "{{ result }}"
next: link_file_hash
- name: mark_complete
type: action
action: update
entity: ProcessingJob
params:
id: "{{ context.job.id }}"
status: completed
completed_at: "{{ now() }}"
progress: 100
result: "{{ context.result }}"
on_success:
next: emit_complete_event
- name: emit_complete_event
type: action
action: emit_event
params:
event: ProcessingJobCompletedEvent
data:
job_id: "{{ context.job.id }}"
file_id: "{{ context.file.id }}"
job_type: "{{ context.job.job_type }}"
result: "{{ context.result }}"
on_success:
next: complete
- name: complete
type: terminal
status: completed
result:
job: "{{ context.job }}"
file: "{{ context.file }}"
result: "{{ context.result }}"
- name: retry_or_fail
type: condition
conditions:
- if: "context.job.retry_count < 3"
next: schedule_retry
- else: true
next: mark_failed
- name: schedule_retry
type: action
action: update
entity: ProcessingJob
params:
id: "{{ context.job.id }}"
status: pending
retry_count: "{{ context.job.retry_count + 1 }}"
retry_after: "{{ now() + exponential_backoff(context.job.retry_count) }}"
error_message: "{{ context.error }}"
on_success:
next: retry_terminal
- name: retry_terminal
type: terminal
status: pending
result:
retry_after: "{{ context.job.retry_after }}"
- name: mark_failed
type: action
action: update
entity: ProcessingJob
params:
id: "{{ context.job.id }}"
status: failed
completed_at: "{{ now() }}"
error_message: "{{ context.error }}"
on_success:
next: emit_failed_event
- name: emit_failed_event
type: action
action: emit_event
params:
event: ProcessingJobFailedEvent
data:
job_id: "{{ context.job.id }}"
file_id: "{{ context.job.file_id }}"
job_type: "{{ context.job.job_type }}"
error: "{{ context.error }}"
on_success:
next: failed_terminal
- name: failed_terminal
type: terminal
status: failed
reason: "{{ context.error }}"
- name: job_not_found
type: terminal
status: failed
reason: "Processing job not found"
- name: file_not_found
type: terminal
status: failed
reason: "File not found"
- name: unknown_job_type
type: terminal
status: failed
reason: "Unknown job type: {{ context.job.job_type }}"
compensation:
- name: cleanup_partial_storage
condition: "steps.store_compressed_file.completed == true && steps.mark_complete.completed == false"
action: delete_file
params:
key: "{{ context.file.storage_key }}_compressed"
backend: "{{ context.bucket.storage_backend }}"
- name: cleanup_thumbnails
condition: "steps.store_image_thumbnail.completed == true && steps.mark_complete.completed == false"
action: delete_file
params:
key: "{{ context.file.storage_key }}_thumb"
backend: "{{ context.bucket.storage_backend }}"