Skip to content
Merged
Show file tree
Hide file tree
Changes from 14 commits
Commits
Show all changes
33 commits
Select commit Hold shift + click to select a range
cc018ea
feat(framework): Add Automation protos
danielnugraha Jul 8, 2026
27fc868
Add stub
danielnugraha Jul 8, 2026
405b3e6
Trim
danielnugraha Jul 8, 2026
00e23cc
Trim
danielnugraha Jul 8, 2026
93722d3
Merge branch 'main' into add-automation-proto
danielnugraha Jul 8, 2026
e292c36
Merge remote-tracking branch 'origin' into add-automation-proto
danielnugraha Jul 8, 2026
363d7b0
Change to string
danielnugraha Jul 9, 2026
2883025
Merge remote-tracking branch 'refs/remotes/origin/add-automation-prot…
danielnugraha Jul 9, 2026
5e89e8b
Merge remote-tracking branch 'origin' into add-automation-corestate
danielnugraha Jul 9, 2026
4492533
Update
danielnugraha Jul 9, 2026
217ea7d
Trim
danielnugraha Jul 9, 2026
9e7ea72
Trim
danielnugraha Jul 9, 2026
5360be2
Trim
danielnugraha Jul 9, 2026
5573704
feat(framework): Add Automation Control implementation
danielnugraha Jul 9, 2026
dfdebf9
Merge
danielnugraha Jul 16, 2026
5110895
Trim
danielnugraha Jul 16, 2026
25d11a9
Trim
danielnugraha Jul 16, 2026
5ac683b
Merge remote-tracking branch 'origin' into add-automation-control
danielnugraha Jul 27, 2026
ee123a7
Implement feedback
danielnugraha Jul 27, 2026
a64771b
Trim
danielnugraha Jul 27, 2026
9f13d30
Format
danielnugraha Jul 27, 2026
ca60a3b
Format
danielnugraha Jul 27, 2026
16b41d0
Update framework/py/flwr/superlink/servicer/control/control_handlers.py
danielnugraha Jul 27, 2026
5f63fe0
Implement feedback
danielnugraha Jul 27, 2026
1b1a9d0
Trim
danielnugraha Jul 27, 2026
51e4e4e
Merge branch 'main' into add-automation-control
danielnugraha Jul 27, 2026
9c3a0e9
Implement feedback
danielnugraha Jul 27, 2026
ab30e6a
Merge remote-tracking branch 'refs/remotes/origin/add-automation-cont…
danielnugraha Jul 27, 2026
98afdc4
trim
danielnugraha Jul 27, 2026
679ae3f
Merge branch 'main' into add-automation-control
danielnugraha Jul 27, 2026
0685bcf
Implement feedback
danielnugraha Jul 27, 2026
8e0bb0f
Merge remote-tracking branch 'refs/remotes/origin/add-automation-cont…
danielnugraha Jul 27, 2026
4abdca1
Merge branch 'main' into add-automation-control
danieljanes Jul 27, 2026
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
4 changes: 4 additions & 0 deletions framework/py/flwr/supercore/constant.py
Original file line number Diff line number Diff line change
Expand Up @@ -124,6 +124,10 @@
# Default federation names for every Flower account
DEFAULT_FEDERATION_SIMULATION = "workspace"

# Constants for automations
AUTOMATION_STATUS_ACTIVE = "active"
AUTOMATION_STATUS_STOPPED = "stopped"


# Constants for exit handling
FORCE_EXIT_TIMEOUT_SECONDS = 5 # Used in `flwr_exit` function
Expand Down
101 changes: 101 additions & 0 deletions framework/py/flwr/supercore/corestate/corestate.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,9 +17,13 @@

from abc import ABC, abstractmethod
from collections.abc import Sequence
from datetime import datetime
from typing import Literal

from flwr.app import Context, Message
from flwr.app.user_config import UserConfig
from flwr.proto.control_pb2 import Automation # pylint: disable=E0611
from flwr.proto.federation_config_pb2 import SimulationConfig # pylint: disable=E0611
from flwr.proto.runseries_pb2 import RunSeries # pylint: disable=E0611
from flwr.proto.task_pb2 import Task, TaskEvent, TaskUsage # pylint: disable=E0611
from flwr.supercore.fab import Fab
Expand Down Expand Up @@ -132,6 +136,103 @@ def store_run_in_series(
run series.
"""

@abstractmethod
def store_automation( # pylint: disable=too-many-arguments
self,
*,
federation_id: str,
flwr_aid: str,
fab_id: str | None,
fab_version: str | None,
fab_hash: str | None,
override_config: UserConfig,
federation_config: SimulationConfig | None,
primary_task_type: str,
next_run_at: datetime,
fixed_interval: int | None = None,
remaining_runs: int | None = None,
series_id: int | None = None,
) -> Automation:
"""Store an automation and return its metadata.

Parameters
----------
federation_id : str
Federation ID the automation belongs to.
flwr_aid : str
FLWR account ID used to dispatch the automation.
fab_id : str | None
FAB ID used by future runs.
fab_version : str | None
FAB version used by future runs.
fab_hash : str | None
FAB hash used by future runs.
override_config : UserConfig
Run override config used by future runs.
federation_config : SimulationConfig | None
Federation config override used by future runs.
primary_task_type : str
Primary task type used by future runs.
next_run_at : datetime
Next due time.
fixed_interval : int | None (default: None)
Recurring interval in seconds.
remaining_runs : int | None (default: None)
Remaining number of runs, if finite.
series_id : int | None (default: None)
Existing run series to reuse, if any.

Returns
-------
Automation
Stored automation metadata.
"""

@abstractmethod
def list_automations(
self,
*,
federation: str | None = None,
statuses: Sequence[str] | None = None,
due_before: datetime | None = None,
limit: int | None = None,
) -> Sequence[Automation]:
"""Return automations matching the given filters.

Parameters
----------
federation : str | None (default: None)
Federation ID to filter by.
statuses : Sequence[str] | None (default: None)
Automation statuses to filter by.
due_before : datetime | None (default: None)
If set, return only automations with `next_run_at` at or before this
timestamp.
limit : int | None (default: None)
Maximum number of automation records to return.

Returns
-------
Sequence[Automation]
Automation metadata. Records are ordered by `next_run_at` ascending
when `due_before` is set, otherwise by `updated_at` descending.
"""

@abstractmethod
def stop_automation(self, automation_id: int) -> bool:
"""Stop an active automation.

Parameters
----------
automation_id : int
Automation ID to stop.

Returns
-------
bool
True if an active automation was stopped, otherwise False.
"""

@abstractmethod
def add_task_log(self, task_id: int, log_message: str) -> None:
"""Add a log entry to the task logs for the specified `task_id`.
Expand Down
94 changes: 94 additions & 0 deletions framework/py/flwr/supercore/corestate/corestate_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@
Status,
SubStatus,
)
from flwr.proto.control_pb2 import Automation # pylint: disable=E0611
from flwr.proto.task_pb2 import ( # pylint: disable=E0611
TaskEvent,
TaskStatus,
Expand Down Expand Up @@ -78,6 +79,33 @@ def _patch_task_log_datetime_now(self, *timestamps: datetime) -> ExitStack:
mock_datetime.now.side_effect = timestamps
return stack

def store_automation(
self,
state: CoreState,
*,
federation_id: str = "@me/fed-a",
flwr_aid: str = "aid-a",
next_run_at: datetime | None = None,
fixed_interval: int | None = None,
remaining_runs: int | None = 1,
series_id: int | None = None,
) -> Automation:
"""Store a minimal automation."""
return state.store_automation(
federation_id=federation_id,
flwr_aid=flwr_aid,
fab_id=None,
fab_version=None,
fab_hash=None,
override_config={},
federation_config=None,
primary_task_type=TaskType.SERVER_APP,
next_run_at=next_run_at or now(),
fixed_interval=fixed_interval,
remaining_runs=remaining_runs,
series_id=series_id,
)

def test_store_run_in_series_creates_id(self) -> None:
"""Storing a run in a run series should create a nonzero ID."""
state = self.state_factory()
Expand Down Expand Up @@ -157,6 +185,72 @@ def test_get_run_series_filters_by_series_ids_and_federation_ids(self) -> None:
self.assertEqual(state.get_run_series(series_ids=[]), [])
self.assertEqual(state.get_run_series(federation_ids=[]), [])

def test_store_list_and_stop_automation(self) -> None:
"""Automation storage should support list, due filtering, and stop."""
state = self.state_factory()
current = now()
series_id = state.store_run_in_series(
run_id=123, federation_id="@me/fed-a", series_id=None
)
assert series_id is not None

due = self.store_automation(
state,
series_id=series_id,
next_run_at=current - timedelta(seconds=60),
fixed_interval=60,
)
future = self.store_automation(
state,
next_run_at=current + timedelta(seconds=60),
)
_ = self.store_automation(
state,
federation_id="@me/fed-b",
next_run_at=current - timedelta(seconds=30),
)

self.assertEqual(due.series_id, series_id)
self.assertEqual(due.status, "active")
self.assertEqual(due.federation, "@me/fed-a")
self.assertEqual(due.flwr_aid, "aid-a")
self.assertEqual(due.fixed_interval, 60)
self.assertEqual(due.remaining_runs, 1)

listed = state.list_automations(federation="@me/fed-a")
self.assertSetEqual(
{automation.automation_id for automation in listed},
{due.automation_id, future.automation_id},
)

due_list = state.list_automations(
federation="@me/fed-a", statuses=["active"], due_before=current, limit=10
)
self.assertEqual(
[automation.automation_id for automation in due_list], [due.automation_id]
)

self.assertTrue(state.stop_automation(due.automation_id))
self.assertFalse(state.stop_automation(due.automation_id))

stopped = state.list_automations(federation="@me/fed-a", statuses=["stopped"])
self.assertEqual(
[automation.automation_id for automation in stopped], [due.automation_id]
)
self.assertTrue(stopped[0].HasField("stopped_at"))
self.assertFalse(stopped[0].HasField("next_run_at"))

def test_store_automation_rejects_series_in_other_federation(self) -> None:
"""Automation storage should reject a series from another federation."""
state = self.state_factory()
series_id = state.store_run_in_series(
run_id=123, federation_id="@me/fed-b", series_id=None
)
assert series_id is not None

with self.assertRaises(ValueError):
self.store_automation(state, series_id=series_id)

def test_create_and_get_task(self) -> None:
"""Test creating and retrieving a task."""
state = self.state_factory()
Expand Down
Loading
Loading