Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
93 changes: 58 additions & 35 deletions SmallPackage/OSlist.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@
- choose the next runnable task by priority

This module keeps those responsibilities together so the scheduler can stay
small and focused. PID lookup uses a sorted list, ready tasks live in one FIFO
small and focused. PID lookup uses a dictionary, ready tasks live in one FIFO
queue per priority, and sleeping tasks live in a wake-time heap.
"""

Expand All @@ -25,68 +25,91 @@

from .SmallTask import SmallTask

from .list_util.binSearchList import insert, search
from .SmallPID import SmallPID


class OSList(SmallPID):
class OSList:
"""
Combined PID registry and queue manager for the cooperative scheduler.
"""

def __init__(self, priors: int = 5, length: int = 2**12) -> None:
"""Create the PID registry plus ready/sleep queue structures."""
SmallPID.__init__(self, length)
self.num_priorities = priors
self.tasks = []
self.ready = [deque() for _ in range(priors)]
self.sleeping = []
self.maxPID = length
self._next_pid = 0
self._tasks_by_pid: dict[int, SmallTask] = {}
# MicroPython requires both an iterable and maxlen. The total task
# capacity is also a safe bound for each queue because a task can be
# present in at most one ready queue once.
self.ready: list[deque[SmallTask]] = [deque((), length) for _ in range(priors)]
self.sleeping: list[tuple[int, int, SmallTask]] = []
self._sleep_seq = 0
self.numWatchers = 0
self.func = lambda data, index: data[index].getID()

def resetCatSel(self):
def resetCatSel(self) -> None:
"""Compatibility no-op kept for older callers."""
return

def _new_pid(self) -> int:
"""Return the next free PID from the bounded PID namespace."""
if len(self._tasks_by_pid) >= self.maxPID:
return -1

pid = self._next_pid
while pid in self._tasks_by_pid:
pid = (pid + 1) % self.maxPID
self._next_pid = (pid + 1) % self.maxPID
return pid

def _is_registered(self, task: SmallTask) -> bool:
"""Check task identity as well as PID to reject stale reused-PID entries."""
return self._tasks_by_pid.get(task.getID()) is task

def _remove_ready_entry(self, task: SmallTask) -> None:
"""Eagerly remove a deleted task so bounded queues cannot retain garbage."""
if not task._queued:
return

priority = task.priority
queue = self.ready[priority]
retained: deque[SmallTask] = deque((), self.maxPID)
while queue:
queued = queue.popleft()
if queued is not task:
retained.append(queued)
self.ready[priority] = retained
task._queued = False

def insert(self, task: SmallTask) -> int:
"""Assign a PID and register a task in the PID-sorted backing list."""
"""Assign a PID and register a task in the PID mapping."""
priority = task.priority
if not 0 < priority < self.num_priorities:
return -1

pid = self.newPID()
pid = self._new_pid()
if pid == -1:
return -1

task.setID(pid)
if task.isWatcher:
self.numWatchers += 1

index = insert(self.tasks, pid, 0, len(self.tasks), func=self.func)
self.tasks.insert(index, task)
self._tasks_by_pid[pid] = task
return pid

def search(self, pid: int) -> SmallTask | Literal[-1]:
"""Look up a task by PID."""
length = len(self.tasks)
index = search(self.tasks, pid, 0, length, self.func)
if index == -1:
return -1
return self.tasks[index]
return self._tasks_by_pid.get(pid, -1)

def delete(self, pid: int) -> int:
"""Remove a task from PID storage and watcher accounting."""
length = len(self.tasks)
index = search(self.tasks, pid, 0, length, self.func)
if index == -1:
task = self._tasks_by_pid.get(pid)
if task is None:
return -1

task = self.tasks[index]
self._remove_ready_entry(task)
if task.isWatcher:
self.numWatchers -= 1
del self.tasks[index]
self.freePID(pid)
del self._tasks_by_pid[pid]
return 0

def enqueue(self, task: SmallTask, front: bool = False) -> int:
Expand All @@ -98,7 +121,7 @@ def enqueue(self, task: SmallTask, front: bool = False) -> int:
"""
if task == -1 or task is None or task.done:
return -1
if self.search(task.getID()) == -1:
if not self._is_registered(task):
return -1
if task._queued:
return 0
Expand All @@ -123,7 +146,7 @@ def pop(self) -> SmallTask | None:
while queue:
task = queue.popleft()
task._queued = False
if self.search(task.getID()) == -1:
if not self._is_registered(task):
continue
if not task.getExeStatus():
continue
Expand All @@ -136,7 +159,7 @@ def has_ready(self) -> bool:
queue = self.ready[priority]
while queue:
task = queue[0]
if self.search(task.getID()) != -1 and task.getExeStatus():
if self._is_registered(task) and task.getExeStatus():
return True
queue.popleft()
task._queued = False
Expand All @@ -157,7 +180,7 @@ def wake_sleeping(self, now: int) -> list[SmallTask]:
ready = []
while self.sleeping and self.sleeping[0][0] <= now:
_, _, task = heapq.heappop(self.sleeping)
if self.search(task.getID()) == -1:
if not self._is_registered(task):
continue
if task.done or task._blocked_reason != "sleep":
continue
Expand All @@ -168,24 +191,24 @@ def next_wake_time(self) -> int | None:
"""Peek at the next valid wake time, discarding stale heap entries."""
while self.sleeping:
wake_time, _, task = self.sleeping[0]
if self.search(task.getID()) == -1 or task.done or task._blocked_reason != "sleep":
if not self._is_registered(task) or task.done or task._blocked_reason != "sleep":
heapq.heappop(self.sleeping)
continue
return wake_time
return None

def list(self) -> list[SmallTask]:
"""Return a snapshot list of currently registered tasks."""
return [task for task in self.tasks]
return [self._tasks_by_pid[pid] for pid in sorted(self._tasks_by_pid)]

def isOnlyWatchers(self) -> bool:
"""Report whether every remaining task is marked as a watcher."""
return len(self.tasks) == self.numWatchers
return len(self._tasks_by_pid) == self.numWatchers

def __len__(self) -> int:
"""Return the number of registered tasks."""
return len(self.tasks)
return len(self._tasks_by_pid)

def __str__(self) -> str:
"""Return a newline-separated dump of all known tasks."""
return "\n".join([str(x) for x in self.tasks])
return "\n".join(str(task) for task in self.list())
2 changes: 1 addition & 1 deletion SmallPackage/SmallOS.py
Original file line number Diff line number Diff line change
Expand Up @@ -1292,7 +1292,7 @@ def cancel_task(self, task: int | SmallTask, recursive: bool = False) -> int:

def __str__(self) -> str:
"""Return a human-readable dump of the currently registered tasks."""
all_tasks = list(self.tasks.tasks)
all_tasks = self.tasks.list()
string = "SmallOS\n"
for count, routine in enumerate(all_tasks):
string += str(count + 1) + ". " + str(routine) + "\n"
Expand Down
27 changes: 5 additions & 22 deletions SmallPackage/SmallSignals.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,6 @@

from .SmallOS import SmallOS
from .SmallTask import SmallTask
from .TaskState import TaskState
from .awaitables import InstructionAwaitable

from .awaitables import (
Expand Down Expand Up @@ -57,7 +56,6 @@ class SmallSignals:

if TYPE_CHECKING:
OS: SmallOS | None
state: TaskState

def __init__(self, OS: SmallOS | None, kwargs: dict[str, Any]) -> None:
"""
Expand Down Expand Up @@ -141,32 +139,17 @@ def acceptSignal(self, sig: int) -> int:
self.handlers(self)
return 0

def sleep(
self, secs: float, state_blob: dict[Any, Any] | None = None
) -> InstructionAwaitable[None]:
"""
Return the awaitable used for cooperative sleeping.

``state_blob`` is preserved for compatibility with the older API style,
where suspension helpers could stash task-local state before yielding.
"""
if state_blob is not None:
self.state.update(state_blob)
def sleep(self, secs: float) -> InstructionAwaitable[None]:
"""Return the awaitable used for cooperative sleeping."""
return sleep_instruction(secs)

def wait_signal(
self, sig: int, state_blob: dict[Any, Any] | None = None
) -> InstructionAwaitable[int]:
def wait_signal(self, sig: int) -> InstructionAwaitable[int]:
"""Return the awaitable used to wait until ``sig`` is delivered."""
if state_blob is not None:
self.state.update(state_blob)
return wait_signal_instruction(sig)

def sigSuspendV2(
self, sig: int, state_blob: dict[Any, Any] | None = None
) -> InstructionAwaitable[int]:
def sigSuspendV2(self, sig: int) -> InstructionAwaitable[int]:
"""Compatibility alias for the older generator-era suspension name."""
return self.wait_signal(sig, state_blob)
return self.wait_signal(sig)

def yield_now(self) -> InstructionAwaitable[None]:
"""Return the awaitable used for an explicit cooperative yield."""
Expand Down
7 changes: 0 additions & 7 deletions SmallPackage/SmallTask.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,6 @@
from .SmallErrors import PIDError, TaskCancelledError
from .SmallSignals import SmallSignals
from .list_util.linkedList import Node
from .TaskState import TaskState


_MISSING = object()
Expand Down Expand Up @@ -60,7 +59,6 @@ def __init__(self, priority: int, routine: TaskRoutine[T], **kwargs: Any) -> Non
self.isWatcher = False
self.parent = None
self.OS: SmallOS | None = None
self.state = TaskState()
self.children = []
self.name = ""
self.args = ()
Expand Down Expand Up @@ -88,8 +86,6 @@ def __init__(self, priority: int, routine: TaskRoutine[T], **kwargs: Any) -> Non
self._adapter_resume_name: str | None = None
self._adapter_resume_job_id: int | None = None

self.state.update({"return_status": 0}, "system")

SmallSignals.__init__(self, self.OS, kwargs)

if kwargs:
Expand Down Expand Up @@ -220,7 +216,6 @@ def complete(self, result: T | None) -> T | None:
self.isReady = 0
self.isWaiting = 0
self.isSleep = 0
self.state.update({"return_status": 0, "result": result}, "system")
return result

def fail(self, exc: BaseException) -> BaseException:
Expand All @@ -230,7 +225,6 @@ def fail(self, exc: BaseException) -> BaseException:
self.isReady = 0
self.isWaiting = 0
self.isSleep = 0
self.state.update({"return_status": -1, "exception": exc}, "system")
return exc

def cancel(self, message: str = "Task cancelled") -> None:
Expand Down Expand Up @@ -270,7 +264,6 @@ def block(self, reason):
self.isReady = 0
self.isWaiting = 1 if reason in ("signal", "join", "join_all", "adapter") else 0
self.isSleep = 1 if reason == "sleep" else 0
self.state.update({"return_status": 1, "blocked_reason": reason}, "system")

def setID(self, pid: int) -> None:
"""Assign the PID chosen by ``SmallOS`` exactly once."""
Expand Down
95 changes: 0 additions & 95 deletions SmallPackage/TaskState.py

This file was deleted.

Loading
Loading