Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
38 commits
Select commit Hold shift + click to select a range
3358dcb
Added PythonTask output to message column of logs in maestro.db
DavideCorigliano-Unimib Nov 11, 2025
bd7aebd
PythonTask message added to logs table
DavideCorigliano-Unimib Nov 11, 2025
2a6a83c
Added BashTask messages to logs table
DavideCorigliano-Unimib Nov 11, 2025
2ac3899
Fields in executions table completed
DavideCorigliano-Unimib Nov 12, 2025
d6d38cb
Logs in client implemented successfully
DavideCorigliano-Unimib Nov 12, 2025
9a4977a
Added a new set of brand new examples
DavideCorigliano-Unimib Nov 15, 2025
5ad6e62
Fixed 2 broken examples
DavideCorigliano-Unimib Nov 15, 2025
ceac3cb
Fixed endpoints in api_client.py
DavideCorigliano-Unimib Nov 15, 2025
08a3150
Fixed again api_client.py and now 'maestro attach...' works... but it…
DavideCorigliano-Unimib Nov 15, 2025
8c18dc1
Patches on server components
DavideCorigliano-Unimib Nov 16, 2025
961439f
Fixed Python messages timestamp on live log streaming
DavideCorigliano-Unimib Nov 16, 2025
0f4931b
Fixed also Bash and Print messages on live log streaming
DavideCorigliano-Unimib Nov 16, 2025
27a4801
'maestro attach...' command finally works
DavideCorigliano-Unimib Nov 16, 2025
022f3d4
Attachment inverted successfully
DavideCorigliano-Unimib Nov 16, 2025
a94d43f
Fine tuning - Inversion of attachment of client console - completed
DavideCorigliano-Unimib Nov 16, 2025
fe8efc7
Fixed logs in parallel execution
DavideCorigliano-Unimib Nov 18, 2025
6f78122
Patched retries bugs
DavideCorigliano-Unimib Nov 18, 2025
6d7b850
Fixed retries in Bash tasks
DavideCorigliano-Unimib Nov 25, 2025
b447460
Updated examples
DavideCorigliano-Unimib Nov 25, 2025
36a55a7
Updated examples + testing passage of files
DavideCorigliano-Unimib Dec 2, 2025
f6296ae
Checked task.py and status_manager.py -'skipped' logic already implem…
DavideCorigliano-Unimib Dec 2, 2025
205201a
Added output column to TaskORM for storing task results
DavideCorigliano-Unimib Dec 2, 2025
13e36ce
Added set_task_output and get_task_output methods to status_manager.py
DavideCorigliano-Unimib Dec 2, 2025
91f8d5b
Added structured output support with __output__, result, and output v…
DavideCorigliano-Unimib Dec 2, 2025
41a0851
Added parsing and support for task-level 'condition' in YAML, patchin…
DavideCorigliano-Unimib Dec 2, 2025
1cf7240
Patched orchestrator.py
DavideCorigliano-Unimib Dec 2, 2025
f04029b
Intermediate commit
DavideCorigliano-Unimib Dec 2, 2025
3265b30
Intermediate passage 2
DavideCorigliano-Unimib Dec 2, 2025
eb89c90
Intermediate checkpoint 3 - basic functioning
DavideCorigliano-Unimib Dec 2, 2025
1ea4c85
First tests - on main Tasks
DavideCorigliano-Unimib Dec 5, 2025
4595da4
Added test for dag_loader.py and task_registry.py
DavideCorigliano-Unimib Dec 6, 2025
734923d
Test on status_manager-py - work-in-progress
DavideCorigliano-Unimib Dec 6, 2025
0061b06
Phase 1 completed - patched base.py
DavideCorigliano-Unimib Dec 6, 2025
bcb9d1b
Phase 2 completed - still test on dag_loader.py on its way
DavideCorigliano-Unimib Dec 6, 2025
23add31
Phase 3 completed - updated orchestrator.py
DavideCorigliano-Unimib Dec 6, 2025
dd9d4eb
Phase 4 completed - Patched set_task_status() in status_manager.py an…
DavideCorigliano-Unimib Dec 6, 2025
97fe7a5
Fine tuning of orchestrator.py completed
DavideCorigliano-Unimib Dec 6, 2025
5c0c7c7
Fixed last example of conditional branching
DavideCorigliano-Unimib Dec 6, 2025
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .vscode/settings.json
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@
"git.enableSmartCommit": true,
"git.confirmSync": false,
"git.autofetch": true,
"workbench.colorTheme": "Dracula Theme",
"workbench.colorTheme": "Ubuntu Color VSCode Dark Highlight",
"workbench.iconTheme": "material-icon-theme",
"errorLens.enabled": true,
"breadcrumbs.enabled": true,
Expand Down
1 change: 0 additions & 1 deletion data.json

This file was deleted.

Binary file removed examples/1_Old_examples/terraform/maestro.db
Binary file not shown.
50 changes: 50 additions & 0 deletions examples/2_New_examples/1.2.1.Long_waits.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
# -------------------------------------------------------------------------------------
# Intermediate Long Waits
# -------------------------------------------------------------------------------------
# This DAG implements a linear pipeline composed of three consecutive steps, each characterized by an artificial wait phase (a 20-second sleep). The goal is to simulate a sequence of slow processing, allowing testing of the orchestrator's "attached" mode for a total duration of approximately 60 seconds. The structure is simple, but useful for verifying the behavior of streaming logs in a linear flow.
# -------------------------------------------------------------------------------------

dag:
name: "long_waits"
tasks:
- task_id: "start"
type: "PrintTask"
params:
message: "Starting long wait DAG"
dependencies: []

- task_id: "step1"
type: "PythonTask"
params:
code: |
import time
print("Step 1 running...")
time.sleep(20)
print("Step 1 - Complete successfully!")
dependencies: ["start"]

- task_id: "step2"
type: "PythonTask"
params:
code: |
import time
print("Step 2 running...")
time.sleep(20)
print("Step 2 - Complete successfully!")
dependencies: ["step1"]

- task_id: "step3"
type: "PythonTask"
params:
code: |
import time
print("Step 3 running...")
time.sleep(20)
print("Step 3 - Complete successfully!")
dependencies: ["step2"]

- task_id: "finish"
type: "PrintTask"
params:
message: "All steps completed!"
dependencies: ["step3"]
50 changes: 50 additions & 0 deletions examples/2_New_examples/1.2.2.Parallel_then_join.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
# -------------------------------------------------------------------------------------
# Parallel then Join
# -------------------------------------------------------------------------------------
# This DAG involves three parallel tasks of different durations, all dependent on a single entry point. Once completed, they converge into a single final merge task. The maximum duration of the slowest branch (approximately 40 seconds) defines the total execution time. It is particularly suitable for testing streaming visualization of parallel tasks and synchronization during the join phase.
# -------------------------------------------------------------------------------------

dag:
name: "parallel_then_join"
tasks:
- task_id: "start"
type: "PrintTask"
params:
message: "Starting parallel DAG"
dependencies: []

- task_id: "worker_a"
type: "PythonTask"
params:
code: |
import time
print("Worker A...")
time.sleep(30)
print("Worker A - Task complete.")
dependencies: ["start"]

- task_id: "worker_b"
type: "PythonTask"
params:
code: |
import time
print("Worker B...")
time.sleep(40)
print("Worker B - Task complete.")
dependencies: ["start"]

- task_id: "worker_c"
type: "PythonTask"
params:
code: |
import time
print("Worker C...")
time.sleep(30)
print("Worker C - Task complete.")
dependencies: ["start"]

- task_id: "merge"
type: "PrintTask"
params:
message: "Workers completed!"
dependencies: ["worker_a", "worker_b", "worker_c"]
40 changes: 40 additions & 0 deletions examples/2_New_examples/1.2.3.Bash_python_mix.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
# -------------------------------------------------------------------------------------
# Bash/Python Mix
# -------------------------------------------------------------------------------------
# Pipeline mista che alterna task Bash e Python, entrambe con ritardi significativi
# (sleep di 25 secondi ciascuna). Serve a testare la coesistenza e il comportamento
# dei log provenienti da esecutori differenti, oltre alla gestione di tempi di attesa
# in sequenza. Durata complessiva di circa 50 secondi.
# -------------------------------------------------------------------------------------

dag:
name: "intermediate_bash_python_mix"
tasks:
- task_id: "start"
type: "PrintTask"
params:
message: "Starting mixed DAG"
dependencies: []

- task_id: "bash_wait"
type: "BashTask"
params:
command: |
echo "Bash waiting 25s..."
sleep 25
dependencies: ["start"]

- task_id: "python_wait"
type: "PythonTask"
params:
code: |
import time
print("Python waiting 25s...")
time.sleep(25)
dependencies: ["bash_wait"]

- task_id: "finish"
type: "PrintTask"
params:
message: "Finished mixed DAG!"
dependencies: ["python_wait"]
56 changes: 56 additions & 0 deletions examples/2_New_examples/1.3.1.Branching_waits_and_joins.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,56 @@
# -------------------------------------------------------------------------------------
# Branching, Heavy Waits, and Multi-Join
# -------------------------------------------------------------------------------------
# This complex DAG combines multiple branching, slow processing, and a final join phase that depends on distinct processing paths. Each branch contains slow subtasks (20–30 seconds) that are combined into a complex pipeline. It is designed to test system scalability, the management of richer graphs, and the correct synchronization of dependencies in the log stream.
# -------------------------------------------------------------------------------------


dag:
name: "branching_waits_and_joins"
tasks:
- task_id: "start"
type: "PrintTask"
params:
message: "Complex DAG start"
dependencies: []

- task_id: "A1"
type: "PythonTask"
params:
code: |
import time; print("Running A1..."); time.sleep(20); print("Task A1 completed.")
dependencies: ["start"]

- task_id: "A2"
type: "PythonTask"
params:
code: |
import time; print("Running A2..."); time.sleep(25); print("Task A2 completed.")
dependencies: ["start"]

- task_id: "B1"
type: "PythonTask"
params:
code: |
import time; print("Running B1..."); time.sleep(30); print("Task B1 completed.")
dependencies: ["A1"]

- task_id: "B2"
type: "PythonTask"
params:
code: |
import time; print("Running B2..."); time.sleep(30); print("Task B2 completed.")
dependencies: ["A2"]

- task_id: "combine"
type: "PrintTask"
params:
message: "Combination done"
dependencies: ["B1", "B2"]

- task_id: "final"
type: "PythonTask"
params:
code: |
import time; print("Finalizing..."); time.sleep(15)
dependencies: ["combine"]
46 changes: 46 additions & 0 deletions examples/2_New_examples/1.3.2.Cascade_retry_delays.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,46 @@
# -------------------------------------------------------------------------------------
# Cascade with Retries and Delay Logic
# -------------------------------------------------------------------------------------
# This DAG tests the system by integrating unstable tasks, automatic retries with configured delays, and a slow pipeline that depends on the success of the critical task. It includes simulated failures (with a 50% probability) to stress the orchestrator's retry logic. The total duration varies based on the retries, ranging from 70 to 100+ seconds. Ideal for validating the robustness of StatusManager, TaskStatus, and logging under error conditions.
# -------------------------------------------------------------------------------------

dag:
name: "cascade_retry_delays"
tasks:
- task_id: "start"
type: "PrintTask"
params:
message: "Start cascade"
dependencies: []

- task_id: "unstable_step"
type: "PythonTask"
params:
code: |
import time, random
print("Unstable step (may fail)...")
# time.sleep(15)
generate_random = random.random()
print(f"Generated: {round(generate_random, 2)}")
if generate_random < 0.67:
print("Failure occurred")
raise Exception("Simulated failure!")
print("Unstable succeeded")
retries: 3
retry_delay: 10
dependencies: ["start"]

- task_id: "slow_processing"
type: "PythonTask"
params:
code: |
import time
print("Slow processing...")
# time.sleep(30)
dependencies: ["unstable_step"]

- task_id: "finish"
type: "PrintTask"
params:
message: "Cascade completed"
dependencies: ["slow_processing"]
39 changes: 39 additions & 0 deletions examples/2_New_examples/1.3.3.Parallel_heavy_bash_python.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
# -------------------------------------------------------------------------------------
# Heavy Parallel Bash/Python
# -------------------------------------------------------------------------------------
# A DAG designed to generate a heavy, concurrent load on two separate executors: Bash and Python. Both main tasks start in parallel, each with a significant wait time (45–50 seconds). The final join allows you to check the logger's behavior, the smoothness of the streaming, and the correct handling of very slow parallel tasks.
# -------------------------------------------------------------------------------------

dag:
name: "parallel_heavy_bash_python"
tasks:
- task_id: "start"
type: "PrintTask"
params:
message: "Begin heavy parallel sequence"
dependencies: []

- task_id: "bash_long"
type: "BashTask"
params:
command: |
echo "Bash running long task..."
sleep 45
echo "Bash task completed."
dependencies: ["start"]

- task_id: "python_long"
type: "PythonTask"
params:
code: |
import time
print("Python running long task...")
time.sleep(50)
print("Python task completed.")
dependencies: ["start"]

- task_id: "final_merge"
type: "PrintTask"
params:
message: "All heavy tasks completed"
dependencies: ["bash_long", "python_long"]
51 changes: 51 additions & 0 deletions examples/2_New_examples/2.1.1.Bash_chain.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,51 @@
dag:
name: "bash_chain"
tasks:

- task_id: "start_list"
type: "PrintTask"
params:
message: "Listing files..."
dependencies: []

- task_id: "list_files"
type: "BashTask"
params:
command: "ls -l"
dependencies: ["start_list"]

- task_id: "start_counting_elements"
type: "PrintTask"
params:
message: "Counting elements..."
dependencies: ["list_files"]

- task_id: "counting_elements"
type: "BashTask"
params:
command: "ls -l | wc -l"
dependencies: ["start_counting_elements"]

- task_id: "start_counting_lines"
type: "PrintTask"
params:
message: "Counting lines..."
dependencies: ["start_counting_elements"]

- task_id: "count_lines"
type: "BashTask"
params:
command: "echo 'Counting lines from previous output' | wc -l"
dependencies: ["start_counting_lines"]

- task_id: "end"
type: "PrintTask"
params:
message: "Bash chain complete!"
dependencies: ["count_lines"]






Loading
Loading