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
4 changes: 3 additions & 1 deletion adapters/python/CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
## 0.13.0

No changes.
Changes:

- The manifest and submission payloads name durations with their unit: `delay` and `timeout` are `delay_ms` and `timeout_ms`, matching `max_age_ms` and `backoff_min_ms`/`backoff_max_ms`, which already did. Decorator arguments are unchanged — `delay`, `timeout` and `cache` still take seconds or a `timedelta`.

## 0.12.0

Expand Down
7 changes: 5 additions & 2 deletions adapters/python/coflux/__main__.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,8 +26,11 @@ def main() -> int:
)
discover_parser.add_argument(
"modules",
nargs="+",
help="Python modules to scan for targets",
nargs="*",
help=(
"Python modules to scan for targets; packages are scanned recursively. "
"With none, every module in the working directory is scanned."
),
)

# execute command
Expand Down
8 changes: 4 additions & 4 deletions adapters/python/coflux/context.py
Original file line number Diff line number Diff line change
Expand Up @@ -241,11 +241,11 @@ def submit_execution(
cache: dict[str, Any] | None = None,
defer: dict[str, Any] | None = None,
memo: bool | list[int] | None = None,
delay: float | None = None,
delay_ms: int | None = None,
retries: dict[str, Any] | None = None,
recurrent: bool = False,
requires: dict[str, list[str]] | None = None,
timeout: int = 0,
timeout_ms: int = 0,
streams: dict[str, Any] | None = None,
concurrency: dict[str, Any] | None = None,
) -> dict[str, Any]:
Expand All @@ -267,11 +267,11 @@ def submit_execution(
cache=cache,
defer=defer,
memo=memo,
delay=delay,
delay_ms=delay_ms,
retries=retries,
recurrent=recurrent,
requires=requires,
timeout=timeout,
timeout_ms=timeout_ms,
streams=streams,
concurrency=concurrency,
)
Expand Down
40 changes: 35 additions & 5 deletions adapters/python/coflux/discovery.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
import dataclasses
import importlib
import json
import os
import pkgutil
import sys
import traceback
Expand All @@ -20,14 +21,41 @@
serialize_streams,
)

# Top-level names that are never where targets live, and that tend to have
# import-time side effects or dependencies the worker doesn't have.
_SKIPPED_TOP_LEVEL = frozenset({"setup", "conftest", "tests", "test"})


def _working_directory_modules() -> list[str]:
"""The top-level modules and packages in the working directory.

This is what a worker hosts when it's given no modules: the directory
it was started in, which ``python -m`` puts on ``sys.path``. Private
(``_``-prefixed) names and a few conventional non-targets are left out,
and packages are only found by their ``__init__.py``, as elsewhere.
"""
cwd = os.getcwd()
if cwd not in sys.path and "" not in sys.path:
sys.path.insert(0, cwd)

return sorted(
info.name
for info in pkgutil.iter_modules([cwd])
if not info.name.startswith("_") and info.name not in _SKIPPED_TOP_LEVEL
)


def _expand_modules(module_names: list[str]) -> list[str]:
"""Expand package names into their constituent submodules.

If a name refers to a Python package (has __path__), it is recursively
walked using pkgutil.walk_packages. Plain module names are passed through
unchanged. Submodules whose final component starts with '_' are skipped.
No names at all means every module in the working directory.
"""
if not module_names:
module_names = _working_directory_modules()

result: list[str] = []
seen: set[str] = set()

Expand Down Expand Up @@ -75,11 +103,13 @@ def discover_targets(
"""Discover all targets in the specified modules.

If a module name refers to a package, all submodules are scanned
recursively (private submodules starting with '_' are skipped).
recursively (private submodules starting with '_' are skipped). An
empty list scans every module in the working directory.

Args:
modules: List of module or package names to scan
(e.g., ["myapp.workflows", "myapp.tasks"] or just ["myapp"])
(e.g., ["myapp.workflows", "myapp.tasks"] or just ["myapp"]),
or [] for the working directory

Returns:
A ``(targets, errors)`` pair — target definitions suitable for JSON
Expand Down Expand Up @@ -150,7 +180,7 @@ def _build_target_definition(target: Any, module_name: str) -> dict[str, Any]:
result["defer"] = serialize_defer(definition.defer, definition.parameters)

if definition.delay:
result["delay"] = _to_ms(definition.delay)
result["delay_ms"] = _to_ms(definition.delay)

if definition.memo:
result["memo"] = definition.memo
Expand All @@ -159,7 +189,7 @@ def _build_target_definition(target: Any, module_name: str) -> dict[str, Any]:
result["requires"] = definition.requires

if definition.timeout:
result["timeout"] = _to_ms(definition.timeout)
result["timeout_ms"] = _to_ms(definition.timeout)

if definition.recurrent:
result["recurrent"] = True
Expand Down Expand Up @@ -192,7 +222,7 @@ def run_discovery(modules: list[str]) -> int:
file would mean restarting by hand — has it.

Args:
modules: List of module names to scan.
modules: List of module names to scan, or [] for the working directory.

Returns:
Exit code (0 for success, 1 for error).
Expand Down
12 changes: 6 additions & 6 deletions adapters/python/coflux/protocol.py
Original file line number Diff line number Diff line change
Expand Up @@ -193,11 +193,11 @@ def request_submit_execution(
cache: dict[str, Any] | None = None,
defer: dict[str, Any] | None = None,
memo: bool | list[int] | None = None,
delay: float | None = None,
delay_ms: int | None = None,
retries: dict[str, Any] | None = None,
recurrent: bool = False,
requires: dict[str, list[str]] | None = None,
timeout: int = 0,
timeout_ms: int = 0,
streams: dict[str, Any] | None = None,
concurrency: dict[str, Any] | None = None,
) -> int:
Expand All @@ -220,16 +220,16 @@ def request_submit_execution(
params["defer"] = defer
if memo is not None:
params["memo"] = memo
if delay is not None:
params["delay"] = delay
if delay_ms is not None:
params["delay_ms"] = delay_ms
if retries is not None:
params["retries"] = retries
if recurrent:
params["recurrent"] = recurrent
if requires:
params["requires"] = requires
if timeout:
params["timeout"] = timeout
if timeout_ms:
params["timeout_ms"] = timeout_ms
if streams is not None:
params["streams"] = streams
if concurrency is not None:
Expand Down
6 changes: 4 additions & 2 deletions adapters/python/coflux/target.py
Original file line number Diff line number Diff line change
Expand Up @@ -714,11 +714,13 @@ def submit(self, *args: P.args, **kwargs: P.kwargs) -> Execution[T]:
cache=cache_dict,
defer=defer_dict,
memo=memo_val,
delay=_to_ms(self._definition.delay) if self._definition.delay else None,
delay_ms=_to_ms(self._definition.delay) if self._definition.delay else None,
retries=retries_dict,
recurrent=self._definition.recurrent,
requires=self._definition.requires,
timeout=_to_ms(self._definition.timeout) if self._definition.timeout else 0,
timeout_ms=_to_ms(self._definition.timeout)
if self._definition.timeout
else 0,
streams=streams_dict,
concurrency=concurrency_dict,
)
Expand Down
98 changes: 98 additions & 0 deletions adapters/python/tests/test_discovery.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,98 @@
"""What discovery scans when it's given a package, or nothing at all.

A package name covers everything under it. No names means the working
directory: its top-level modules and packages, less the private ones and
the conventional non-targets that tend to blow up on import.
"""

from __future__ import annotations

import importlib
import sys
import textwrap

import pytest

from coflux.discovery import discover_targets

_TASK = textwrap.dedent(
"""
import coflux as cf

@cf.task()
def {name}():
return 1
"""
)


@pytest.fixture
def project(tmp_path, monkeypatch):
"""A throwaway project as the working directory, cleaned out of
``sys.modules`` afterwards so its names don't leak between tests."""
monkeypatch.chdir(tmp_path)
monkeypatch.syspath_prepend(str(tmp_path))
importlib.invalidate_caches()
before = set(sys.modules)

def write(path, source):
file = tmp_path / path
file.parent.mkdir(parents=True, exist_ok=True)
file.write_text(source)

yield write

for name in set(sys.modules) - before:
sys.modules.pop(name, None)


def _names(targets):
return sorted((t["module"], t["name"]) for t in targets)


def test_a_package_covers_its_submodules(project):
project("dpkg/__init__.py", "")
project("dpkg/flows.py", _TASK.format(name="flow"))
project("dpkg/deep/__init__.py", "")
project("dpkg/deep/jobs.py", _TASK.format(name="job"))
project("dpkg/_private.py", _TASK.format(name="hidden"))

targets, errors = discover_targets(["dpkg"])

assert errors == []
assert _names(targets) == [("dpkg.deep.jobs", "job"), ("dpkg.flows", "flow")]


def test_no_modules_scans_the_working_directory(project):
project("dwd_app/__init__.py", "")
project("dwd_app/flows.py", _TASK.format(name="flow"))
project("dwd_scratch.py", _TASK.format(name="scratch"))
project("_dwd_hidden.py", _TASK.format(name="hidden"))
# Skipped by name, so importing them is never attempted
project("conftest.py", "raise RuntimeError('conftest imported')")
project("setup.py", "raise RuntimeError('setup imported')")
project("tests/__init__.py", "raise RuntimeError('tests imported')")
# Not a package without __init__.py, and not a module at all
project("dwd_data/flows.py", _TASK.format(name="orphan"))
project("notes.txt", "")

targets, errors = discover_targets([])

assert errors == []
assert _names(targets) == [("dwd_app.flows", "flow"), ("dwd_scratch", "scratch")]


def test_an_empty_working_directory_finds_nothing(project):
project("notes.txt", "")

assert discover_targets([]) == ([], [])


def test_a_broken_module_in_the_working_directory_is_reported(project):
project("dwd_ok.py", _TASK.format(name="ok"))
project("dwd_broken.py", "import does_not_exist")

targets, errors = discover_targets([])

assert _names(targets) == [("dwd_ok", "ok")]
assert [e.module for e in errors] == ["dwd_broken"]
12 changes: 11 additions & 1 deletion cli/CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,6 +1,16 @@
## 0.13.0

No changes.
Enhancements:

- Adds `--type ecs` support for `pools create` and `pools update`, including the `roleArn` and `roleExternalId` fields for a role to assume before calling ECS.
- Adds the `idleTimeout` pool field, for how long a pool keeps an idle worker before stopping it (`--set idleTimeout=5m`, exported as `idle_timeout = "5m"`).
- Adds `secrets set`, `secrets list` and `secrets delete`. Pools refer to secrets by name (`tokenSecret`, `credentialsSecret`, `envSecrets`) instead of holding credentials, so `pools export` no longer needs `--include-secrets`. Each secret is set for one or more workspace patterns, given as a required `--workspaces`, in the same language `tokens create --workspaces` uses.
- The worker exits when the server asks it to (a `stop` command over its connection, sent to a pool worker that has been idle past its pool's timeout), draining in-flight executions first as it does on SIGTERM.

Changes:

- Durations are written as durations. `submit --delay` takes `30s` or `5m` rather than a bare number of seconds, and no longer accepts one; `logs.flush_interval` and `metrics.flush_interval` in `coflux.toml` are duration strings (`"500ms"`) rather than floats, and a bare number is now an error rather than being read as nanoseconds. Both need editing by hand.
- Follows the server's rename of duration fields: the manifest and submission payloads carry `delayMs`, `timeoutMs`, `maxAgeMs`, `backoffMinMs` and `backoffMaxMs`, and the worker protocol carries `max_age_ms`, `backoff_min_ms` and `backoff_max_ms`. A CLI of this version needs a server of this version, as before.

## 0.12.0

Expand Down
36 changes: 33 additions & 3 deletions cli/cmd/coflux/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,8 +4,10 @@ import (
"fmt"
"log/slog"
"os"
"reflect"
"runtime"
"strings"
"time"

"github.com/bitroot/coflux/cli/internal/config"
"github.com/bitroot/coflux/cli/internal/version"
Expand Down Expand Up @@ -49,9 +51,9 @@ func init() {
viper.SetDefault("worker.concurrency", min(runtime.NumCPU()+4, 32))
viper.SetDefault("blobs.threshold", 100)
viper.SetDefault("logs.batch_size", 100)
viper.SetDefault("logs.flush_interval", 0.5)
viper.SetDefault("logs.flush_interval", "500ms")
viper.SetDefault("metrics.batch_size", 100)
viper.SetDefault("metrics.flush_interval", 0.2)
viper.SetDefault("metrics.flush_interval", "200ms")
viper.SetDefault("log_level", "info")

// Global flags
Expand Down Expand Up @@ -105,7 +107,8 @@ func init() {
queueCmd.GroupID = "management"
inputsCmd.GroupID = "management"
catalogCmd.GroupID = "management"
rootCmd.AddCommand(workspacesCmd, manifestsCmd, poolsCmd, tokensCmd, assetsCmd, blobsCmd, logsCmd, sessionsCmd, queueCmd, inputsCmd, catalogCmd)
secretsCmd.GroupID = "management"
rootCmd.AddCommand(workspacesCmd, manifestsCmd, poolsCmd, tokensCmd, secretsCmd, assetsCmd, blobsCmd, logsCmd, sessionsCmd, queueCmd, inputsCmd, catalogCmd)
}

func initConfig(cmd *cobra.Command, args []string) error {
Expand Down Expand Up @@ -135,11 +138,38 @@ func initConfig(cmd *cobra.Command, args []string) error {
return nil
}

var durationType = reflect.TypeOf(time.Duration(0))

// stringToDurationHook decodes duration settings from duration strings
// ("500ms", "2s"). A bare number is rejected rather than truncated to
// nanoseconds: before 0.13 these were written as a float of seconds, and
// silently turning `flush_interval = 0.5` into zero would be worse than
// failing.
func stringToDurationHook(from reflect.Type, to reflect.Type, data any) (any, error) {
if to != durationType || from == durationType {
return data, nil
}
if from.Kind() != reflect.String {
return nil, fmt.Errorf("expected a duration string (e.g. \"500ms\"), got %v", data)
}
value, err := time.ParseDuration(data.(string))
if err != nil {
return nil, fmt.Errorf("invalid duration %q (expected e.g. \"500ms\")", data)
}
return value, nil
}

// loadConfig unmarshals viper config into a Config struct
func loadConfig() (*config.Config, error) {
cfg := &config.Config{}
if err := viper.Unmarshal(cfg, func(dc *mapstructure.DecoderConfig) {
dc.ErrorUnused = true
// Replaces viper's default StringToTimeDurationHookFunc with a
// stricter one; the slice hook is viper's default, kept as-is.
dc.DecodeHook = mapstructure.ComposeDecodeHookFunc(
stringToDurationHook,
mapstructure.StringToSliceHookFunc(","),
)
}); err != nil {
return nil, fmt.Errorf("failed to unmarshal config: %w", err)
}
Expand Down
Loading
Loading