Skip to content

API reference

The public API of Winslow: the names that import winslow exposes, then the client API a UI pane or an agent uses to read and drive sessions, in process or over the wire (see Serve and connect and UI plugins).

Workflow

winslow.Workflow

Bases: _ConfigBase

Source code in src/winslow/workflow/workflow.py
 38
 39
 40
 41
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
class Workflow(_ConfigBase):
    name = None

    # If True, the interactive app initializes this workflow at launch with the
    # default args. It does not show the selector, the form or the confirmation.
    # This is useful for a test. More than one workflow can set it, and the app
    # initializes each of them.
    auto_init = False

    store_classes = {
        Mode.HEADLESS: TaskStore,
        Mode.TUI: TaskStore,
    }

    runner_classes = {
        Mode.HEADLESS: HeadlessRunner,
        Mode.TUI: InteractiveRunner,
    }

    # The default check TTL of the tasks, in seconds: a passing check younger
    # than this counts as verified without a probe. None means always probe.
    # A task overrides it with its own check_ttl (see BaseRunner).
    check_ttl = None

    graph_class = Graph
    registry_class = TaskRegistry

    filter_registry_class = FilterRegistry

    cache_registry_class = WorkflowCacheRegistry

    def __init__(
        self,
        orchestrator_config,
        workflow_config=None,
        store=None,
        logger=LOGGER,
        root_dir=None,
    ):

        super().__init__(orchestrator_config)
        # The project root of the process (see Orchestrator.directory). The
        # task detail renders source paths relative to it.
        self.root_dir = root_dir

        self.workflow_config = (
            workflow_config if workflow_config is not None else Namespace()
        )
        # The caches see only the workflow config, so the identity travels on
        # it (see WorkflowCache._storage_namespace). The name is safe: a config
        # option cannot bind it, because it clashes with the property.
        self.workflow_config.cache_namespace = self.cache_namespace
        self.registry = self.registry_class(
            orchestrator_config=orchestrator_config,
            workflow_config=self.workflow_config,
        )
        self.graph = self.graph_class(
            orchestrator_config=orchestrator_config,
            workflow_config=self.workflow_config,
            logger=logger,
        )
        self.filter_registry = self.filter_registry_class(
            orchestrator_config=orchestrator_config,
            workflow_config=self.workflow_config,
        )
        self.logger = logger
        # The Session does not exist at construction time. It is created after
        # the workflow init and attaches itself here.
        self._session = None
        # The workflow owns its persistence adapter and stale sweeper. None
        # until init_state attaches them; archive_state detaches them.
        self.persistence_listener = None
        self.stale_sweeper = None

        # initialize_tasks builds the containers, before it builds the graph.
        self._workflow_cache = None
        self._global_cache = None
        self._task_index = None
        # The workflow owns the prepared tasks; the task index holds weak
        # references. initialize fills the list, release_tasks clears it.
        self.tasks = None

        # The session baseline of the batch options, from the CLI. A submit can
        # carry its own values; this never changes (see BatchOptions).
        self.batch_options = BatchOptions(
            dry_run=orchestrator_config.dry_run,
            force_run=orchestrator_config.force_run,
            force_success=orchestrator_config.force_success,
            disable_concurrency=orchestrator_config.disable_concurrency,
        )

        if store is None:
            self.logger.debug(f"Auto-initializing store for {self}")
            self.bus = SessionBus()
            self.store = self.generate_store(
                bus=self.bus,
                orchestrator_config=orchestrator_config,
                workflow_config=self.workflow_config,
            )
        else:
            # A given store carries its bus, so the workflow adopts it: one
            # bus per session, however the store was built.
            self.store = store
            self.bus = store.bus
        if orchestrator_config.mode is Mode.HEADLESS:
            # The store publishes each transition; the log line is a
            # subscriber, and the mode decides who listens (see log_task_status).
            self.bus.subscribe(TaskStatusEvent, log_task_status)
        self.runner = self.runner_classes[orchestrator_config.mode](
            orchestrator_config=orchestrator_config,
            workflow=self,
            store=self.store,
            logger=self.logger,
            workflow_config=self.workflow_config,
            batch_options=self.batch_options,
        )

    @classmethod
    def generate_store(cls, bus, orchestrator_config, workflow_config):
        return cls.store_classes[orchestrator_config.mode](bus)

    @property
    def check(self):
        return self.orchestrator_config.check

    @property
    def dry_run(self):
        return self.batch_options.dry_run

    @property
    def force_run(self):
        return self.batch_options.force_run

    @property
    def force_success(self):
        return self.batch_options.force_success

    @property
    def disable_concurrency(self):
        return self.batch_options.disable_concurrency

    @property
    def session(self):
        return self._session

    @property
    def session_id(self):
        # A log record can be emitted before a session attaches, so "no session"
        # must be a legal state. The stamp filter accepts None.
        return self._session.session_id if self._session else None

    @property
    def identifiers_dict_safe(self):
        """{option name: display value} for the config options with
        identifier=True. identifier implies required, so each option always
        has a value and the code reads it directly."""
        return {
            name: option.format_value(getattr(self.workflow_config, name))
            for name, option in self.config_meta.items()
            if option.identifier
        }

    @property
    def identifier_suffix(self):
        """A key=value list of the identifier options. This part makes two
        runs of the same workflow different; empty without such options."""
        return " | ".join(f"{k}={v}" for k, v in self.identifiers_dict_safe.items())

    @property
    def identity_prefix(self):
        """A readable prefix: the instance name plus the scalar identifier
        values. Not unique - a structured value (a multiselect list, a tuple)
        is dropped, so pair it with identity_hash, which covers the full dict."""
        scalars = [
            formatted
            for name, formatted in self.identifiers_dict_safe.items()
            if isinstance(getattr(self.workflow_config, name), (str, int, float, bool))
        ]
        return slugify("-".join([self.instance_name, *scalars]))

    @property
    def identity_hash(self):
        """A digest of the full identity: two runs collide only when the name
        and every identifier match, however identity_prefix flattened them."""
        return identity_digest(self.instance_name, self.identifiers_dict_safe)

    @property
    def cache_namespace(self):
        """The directory of the persistent cache tiers of this run (see
        JsonFileStorage): readable prefix plus digest, stable across sessions."""
        return f"{self.identity_prefix}-{self.identity_hash}"

    @cached_property
    def run_nonce(self):
        """The nonce separates two concurrent runs of one workflow in the log
        routing (see Task.log_key). The property owns the name, so a config
        option cannot bind it."""
        return new_uuid()

    def __str__(self):
        """The display form of the run: the name plus the identifier options,
        for example "etl (client=acme)" - the same shape as str(task)."""
        # _Base.__init__ logs through __str__ before workflow_config is set.
        if getattr(self, "workflow_config", None) is None:
            return self.instance_name
        if not self.identifier_suffix:
            return self.instance_name
        return f"{self.instance_name} ({self.identifier_suffix})"

    @classmethod
    def should_be_initialized(cls, orchestrator_config, parameters=None):
        return True

    @property
    def global_cache(self):
        if self._global_cache is None:
            raise InitializationError(
                f"{self} caches read before initialize_tasks built them."
            )
        return self._global_cache

    @property
    def workflow_cache(self):
        if self._workflow_cache is None:
            raise InitializationError(
                f"{self} caches read before initialize_tasks built them."
            )
        return self._workflow_cache

    @property
    def task_index(self):
        """Resolve an identity key to a live task (see TaskIndex)."""
        if self._task_index is None:
            raise InitializationError(
                f"{self} task index read before initialize_tasks built it."
            )
        return self._task_index

    def _initialize_caches(self):
        """Build and populate the containers of both scopes: the graph and the
        pre-graph hooks read them (see docs/caching.md)."""
        self._global_cache = initialize_global_cache(
            self.orchestrator_config,
            self.disable_concurrency,
            clear=self.orchestrator_config.clear_cache,
        )
        registry = self.cache_registry_class(self.orchestrator_config)
        registry.collect_classes(self.module_directory)
        instances = {
            kls.get_name(): kls(self.workflow_config) for kls in registry.classes
        }
        self._workflow_cache = CacheContainer(instances, scope=WORKFLOW_SCOPE)
        if self.orchestrator_config.clear_cache:
            # Before the population, so the eager loaders run fresh and a
            # persistent tier rewrites. A memory tier is cold anyway.
            self._workflow_cache.clear_all()
        self._workflow_cache.populate_eager_entries(self.disable_concurrency)

    def initialize_tasks(self, logger=LOGGER):
        if self.graph is None:
            raise InitializationError(
                f"{self} tasks already initialized - initialize_tasks is one-shot."
            )

        logger.debug(f"{self} initializing tasks.")

        self._initialize_caches()
        self.registry.collect_classes(self.module_directory)

        # The context serves the classmethod hooks that run before a task
        # instance exists (see winslow.cache.get_workflow_cache).
        with workflow_cache_context(self._workflow_cache):
            tasks = self.graph.generate_pipeline(self.registry)

        for task in tasks:
            # A batch thread reads the stamp, because the thread pool does not
            # propagate a context variable.
            task._workflow_cache_container = self._workflow_cache
            task._global_cache_container = self._global_cache
            # The nonce goes first: the buffer registers under log_key, and
            # the nonce prefixes that key.
            task._run_nonce = self.run_nonce
            if self.orchestrator_config.is_interactive:
                task._enable_log_buffer()
            self.store[task] = TaskStatus.INITIALIZED

        self.tasks = sorted(tasks, key=operator.attrgetter("_index"))
        self._task_index = TaskIndex(tasks)

        logger.debug(f"{self} initialized {len(self.store)} tasks.")

        # The graph did its one job. Drop it, so initialize_tasks stays one-shot.
        self.graph = None

    def release_tasks(self):
        # The single release point, at session end: history holds values and
        # uuids, so the workflow task list is the last owner of each task.
        # Batch errors go first, because their traceback frames reference tasks.
        self.runner.release_batch_errors()
        self.store.clear()
        self.tasks = None
        # The container dies with the session. Only a new session builds fresh
        # WorkflowCache instances.
        self._workflow_cache = None

    def effective_check_ttl(self, task):
        """The check TTL of the task: its own declaration when set, else the
        workflow default (see check_ttl)."""
        return task.check_ttl if task.check_ttl is not None else self.check_ttl

    def task_info(self, task, **kwargs):
        """The TaskInfo of the task, with the trust fields of the check_ttl
        rule filled from the stored snapshot (see TaskInfo.from_task)."""
        entry = self.load_snapshot(task.identity_key)
        return TaskInfo.from_task(
            task,
            checked_at=entry.checked_at if entry is not None else None,
            effective_ttl=self.effective_check_ttl(task),
            **kwargs,
        )

    def subscribe(self, event_type, handler):
        """Subscribe the handler to the session events of this workflow."""
        self.bus.subscribe(event_type, handler)

    def unsubscribe(self, event_type, handler):
        """Disconnect the handler (see subscribe). An unknown handler is a
        no-op, so a teardown path can run twice."""
        self.bus.unsubscribe(event_type, handler)

    def add_cache_listener(self, listener):
        """Attach the listener to the caches this workflow can see: the
        workflow cache and the global cache."""
        self.workflow_cache.add_listener(listener)
        self.global_cache.add_listener(listener)

    def remove_cache_listener(self, listener):
        """Detach the listener from both caches (see add_cache_listener). A
        no-op on the workflow cache once the session end has released it
        (see release_tasks): a teardown that races the session end must not
        raise."""
        if self._workflow_cache is not None:
            self._workflow_cache.remove_listener(listener)
        self.global_cache.remove_listener(listener)

    def caches(self):
        """The caches this workflow can see, workflow scope first. After the
        session end the read refuses with direction: the caches are released.
        One read of the container, so the message matches the state."""
        container = self._workflow_cache
        if container is None:
            if self.session is not None and self.session.has_ended:
                raise ValueError(
                    f"{self} has ended and released its caches - the recorded "
                    f"cache reads live in the execution history (see "
                    f"record_detail)."
                )
            raise InitializationError(
                f"{self} caches read before initialize_tasks built them."
            )
        return (*container.caches(), *self.global_cache.caches())

    def get_cache(self, name):
        """The live cache named `name`, or None."""
        for cache in self.caches():
            if cache.get_name() == name:
                return cache
        return None

    def init_state(
        self,
        state_store,
        origin=None,
        orchestrator_overrides=None,
        workflow_values=None,
    ):
        """Start persistence for the attached session: the manifest, and the
        persistence listener on the store. Call this once the pipeline is
        runnable, after the eligibility pass: the manifest marks the session
        as a restore candidate. A failure degrades to a run without state."""
        if self.persistence_listener is not None:
            # A second registration doubles the writes.
            return
        session = self._session
        adapter = sweeper = None
        try:
            adapter = SessionPersistenceAdapter(state_store, session.session_id)
            adapter.attach(self)
            self.persistence_listener = adapter
            sweeper = StaleSweeper(self)
            self.bus.subscribe(TaskStatusEvent, sweeper.on_task_status)
            self.stale_sweeper = sweeper
            # The manifest lands last: a failure before this point leaves no
            # durable state.
            state_store.save_manifest(
                SessionManifest(
                    session_id=session.session_id,
                    workflow_class=type(self).get_name(),
                    workflow_namespace=self.cache_namespace,
                    orchestrator_overrides=orchestrator_overrides,
                    workflow_values=workflow_values,
                    origin=origin,
                    started_at=session.start,
                )
            )
        except Exception:
            # A persistence failure must not break the session start: the
            # session degrades to a run without state. The locals include a
            # subscriber whose registration failed.
            if adapter is not None:
                adapter.detach(self)
                adapter.close()
            if sweeper is not None:
                self.bus.unsubscribe(TaskStatusEvent, sweeper.on_task_status)
                sweeper.close()
            self.persistence_listener = None
            self.stale_sweeper = None
            self.logger.error(
                f"Could not persist the manifest of {session.session_id} - "
                f"the session runs without state",
                exc_info=True,
            )

    def load_snapshot(self, key):
        """The latest snapshot of the key, or None. The listener overlays the
        writes of this session on the initial state, so the read stays in
        memory (see SessionPersistenceAdapter.get)."""
        listener = self.persistence_listener
        return listener.get(key) if listener is not None else None

    def seed_from_state(self):
        """Feed the persisted state of the session onto the store: the
        snapshots replay as statuses, and the open batch records register as
        INTERRUPTED. Call this after the eligibility pass: that pass
        overwrites earlier status writes."""
        listener = self.persistence_listener
        if listener is None:
            return
        self._seed_task_statuses(listener.initial_state)
        self.runner.seed_interrupted_batches(listener.load_open_batches())

    def _seed_task_statuses(self, snapshots):
        """Replay the last terminal status of each task that the eligibility
        pass left READY_TO_PROCESS."""
        for task in self.tasks:
            if self.store[task] is not TaskStatus.READY_TO_PROCESS:
                continue
            entry = snapshots.get(task.identity_key)
            status = TaskStatus.__members__.get(entry.status) if entry else None
            if status is None:
                continue
            if status in PASSING_STATUSES and not is_trusted(
                entry.checked_at,
                self.effective_check_ttl(task),
                self._session.start,
                time.time(),
            ):
                # An untrusted success seeds as STALE: the next touch
                # re-verifies it (see TaskStatus.STALE).
                status = TaskStatus.STALE
            # An ordinary store write, so the subscribers see a normal
            # event. The SEED origin keeps checked_at where the probe
            # left it (see SessionPersistenceAdapter).
            self.runner.set_status(task, status, None, origin=Origin.SEED)

    def archive_state(self):
        """End persistence: stop the sweeper and the writer, unsubscribe them,
        then stamp and archive the manifest. After the durable writes the bus
        publishes SessionEndedEvent and closes, which disconnects every
        remaining subscriber. The session end calls this, and a persistence
        failure must not break the end. A session in ERROR archives as failed
        (see StateStore.mark_errored)."""
        if (sweeper := self.stale_sweeper) is not None:
            sweeper.close()
            self.bus.unsubscribe(TaskStatusEvent, sweeper.on_task_status)
            self.stale_sweeper = None
        if (listener := self.persistence_listener) is not None:
            listener.close()
            self.bus.unsubscribe(TaskStatusEvent, listener.on_task_status)
            self.persistence_listener = None
            try:
                if self.session.status is SessionStatus.ERROR:
                    listener.mark_errored()
                else:
                    listener.mark_ended()
            except Exception:
                self.logger.error(
                    f"Could not archive the manifest of {self.session_id}",
                    exc_info=True,
                )
        self.bus.publish(SessionEndedEvent(session_id=self.session_id))
        self.bus.close()

    def check_pipeline_eligibility(self, logger=LOGGER):
        tasks = self.tasks
        logger.debug(f"Checking eligibility for {len(tasks)} tasks.")
        self.runner.check_eligibility(tasks)

    @property
    def settings_snapshot(self):
        return {
            "dry_run": self.dry_run,
            "force_run": self.force_run,
            "force_success": self.force_success,
            "check": self.check,
            "disable_concurrency": self.disable_concurrency,
            "env": self.env,
        }

    @property
    def filter(self):
        return getattr(self.orchestrator_config, "filter", None)

    def get_filtered_tasks(self, query=None):
        tasks = self.tasks
        query = query or self.filter
        if not query:
            return tasks
        try:
            return self.filter_registry.parse(query).apply(tasks)
        except ValueError as e:
            # Report the bad filter and do not run everything silently. The parse
            # error message names the exact part that is wrong.
            raise MisconfigurationError(f"Invalid filter: {e}") from e

    def roster_tasks(self):
        """The tasks a roster read serves, in launch-filter order. A bad
        launch filter logs and answers every task, so an interactive client
        still renders the list (see get_filtered_tasks)."""
        try:
            return self.get_filtered_tasks()
        except MisconfigurationError:
            self.logger.error(
                "The launch filter does not parse - the roster lists every task.",
                exc_info=True,
            )
            return self.tasks

    def record_infos(self):
        """One TaskInfo per task with an execution record, across every
        batch. The record stores survive the session end, so a history
        search works after the task release (see ExecutionRecordStore)."""
        return tuple(
            {
                record.info.key: record.info
                for store in self.runner.record_stores()
                for record in store.records
            }.values()
        )

    def filter_keys(self, query, scope="tasks"):
        """The identity keys the query matches over the named corpus: 'tasks'
        applies the full registry over the live tasks, 'history' the builtin
        filters over the record infos. Raises ValueError with direction."""
        if scope not in ("tasks", "history"):
            raise ValueError(
                f"{scope!r} names no filter scope - the scopes are "
                f"'tasks' and 'history'."
            )
        parsed = self.filter_registry.parse(query)
        if scope == "history":
            enforce_builtin_only(parsed)
            return tuple(info.key for info in parsed.apply(self.record_infos()))
        if self.tasks is None:
            raise ValueError(
                f"{self} has ended and released its tasks - search the "
                f"execution records with scope='history'."
            )
        return tuple(task.identity_key for task in parsed.apply(self.tasks))

    def headless_run(self):
        # This looks unused, but the construction of the Session attaches it as
        # _session, and the runner needs it as its logging identity. The
        # interactive path gets its session from the UI, so a headless run must
        # make its own.
        Session(self)

        # Validate and resolve the filter first, so a bad --filter fails
        # immediately, before the eligibility pass runs and logs for each task.
        filtered_tasks = self.get_filtered_tasks()

        if self.filter and not filtered_tasks:
            raise MisconfigurationError(
                f"Filter {self.filter!r} matched no tasks in {self} - nothing to run."
            )

        self.runner.check_eligibility(self.tasks)

        if self.check:
            self.runner.check(filtered_tasks)
        else:
            self.runner.run(filtered_tasks)

        statuses = list(self.store.values())
        completed = sum(1 for s in statuses if s in PASSING_STATUSES)
        unsuccessful = sum(1 for s in statuses if s in UNSUCCESSFUL_STATUSES)

        self.logger.info(
            f"{self} finished - {completed} completed, "
            f"{unsuccessful} unsuccessful, {len(statuses)} total."
        )

        flagged = [
            key for key, s in self.store.items() if s is TaskStatus.COMPLETED_WITH_ERROR
        ]
        if flagged:
            self.logger.warning(
                f"{len(flagged)} task(s) completed despite errors during "
                f"processing - check task logs: {', '.join(map(str, flagged))}"
            )
        return unsuccessful == 0

    @classmethod
    def get_parser(cls, lenient=False):
        if not cls.config_meta:
            return None

        parser = ArgumentParser(
            prog=f"Workflow - {cls.get_name()}",
            description="Sub parser for a workflow",
        )

        for arg_ctx in cls.get_argparse_context(lenient=lenient):
            name = arg_ctx.pop("arg_name")
            parser.add_argument(name, **arg_ctx)

        return parser

identifiers_dict_safe property

{option name: display value} for the config options with identifier=True. identifier implies required, so each option always has a value and the code reads it directly.

identifier_suffix property

A key=value list of the identifier options. This part makes two runs of the same workflow different; empty without such options.

identity_prefix property

A readable prefix: the instance name plus the scalar identifier values. Not unique - a structured value (a multiselect list, a tuple) is dropped, so pair it with identity_hash, which covers the full dict.

identity_hash property

A digest of the full identity: two runs collide only when the name and every identifier match, however identity_prefix flattened them.

cache_namespace property

The directory of the persistent cache tiers of this run (see JsonFileStorage): readable prefix plus digest, stable across sessions.

run_nonce cached property

The nonce separates two concurrent runs of one workflow in the log routing (see Task.log_key). The property owns the name, so a config option cannot bind it.

task_index property

Resolve an identity key to a live task (see TaskIndex).

__str__()

The display form of the run: the name plus the identifier options, for example "etl (client=acme)" - the same shape as str(task).

Source code in src/winslow/workflow/workflow.py
def __str__(self):
    """The display form of the run: the name plus the identifier options,
    for example "etl (client=acme)" - the same shape as str(task)."""
    # _Base.__init__ logs through __str__ before workflow_config is set.
    if getattr(self, "workflow_config", None) is None:
        return self.instance_name
    if not self.identifier_suffix:
        return self.instance_name
    return f"{self.instance_name} ({self.identifier_suffix})"

effective_check_ttl(task)

The check TTL of the task: its own declaration when set, else the workflow default (see check_ttl).

Source code in src/winslow/workflow/workflow.py
def effective_check_ttl(self, task):
    """The check TTL of the task: its own declaration when set, else the
    workflow default (see check_ttl)."""
    return task.check_ttl if task.check_ttl is not None else self.check_ttl

task_info(task, **kwargs)

The TaskInfo of the task, with the trust fields of the check_ttl rule filled from the stored snapshot (see TaskInfo.from_task).

Source code in src/winslow/workflow/workflow.py
def task_info(self, task, **kwargs):
    """The TaskInfo of the task, with the trust fields of the check_ttl
    rule filled from the stored snapshot (see TaskInfo.from_task)."""
    entry = self.load_snapshot(task.identity_key)
    return TaskInfo.from_task(
        task,
        checked_at=entry.checked_at if entry is not None else None,
        effective_ttl=self.effective_check_ttl(task),
        **kwargs,
    )

subscribe(event_type, handler)

Subscribe the handler to the session events of this workflow.

Source code in src/winslow/workflow/workflow.py
def subscribe(self, event_type, handler):
    """Subscribe the handler to the session events of this workflow."""
    self.bus.subscribe(event_type, handler)

unsubscribe(event_type, handler)

Disconnect the handler (see subscribe). An unknown handler is a no-op, so a teardown path can run twice.

Source code in src/winslow/workflow/workflow.py
def unsubscribe(self, event_type, handler):
    """Disconnect the handler (see subscribe). An unknown handler is a
    no-op, so a teardown path can run twice."""
    self.bus.unsubscribe(event_type, handler)

add_cache_listener(listener)

Attach the listener to the caches this workflow can see: the workflow cache and the global cache.

Source code in src/winslow/workflow/workflow.py
def add_cache_listener(self, listener):
    """Attach the listener to the caches this workflow can see: the
    workflow cache and the global cache."""
    self.workflow_cache.add_listener(listener)
    self.global_cache.add_listener(listener)

remove_cache_listener(listener)

Detach the listener from both caches (see add_cache_listener). A no-op on the workflow cache once the session end has released it (see release_tasks): a teardown that races the session end must not raise.

Source code in src/winslow/workflow/workflow.py
def remove_cache_listener(self, listener):
    """Detach the listener from both caches (see add_cache_listener). A
    no-op on the workflow cache once the session end has released it
    (see release_tasks): a teardown that races the session end must not
    raise."""
    if self._workflow_cache is not None:
        self._workflow_cache.remove_listener(listener)
    self.global_cache.remove_listener(listener)

caches()

The caches this workflow can see, workflow scope first. After the session end the read refuses with direction: the caches are released. One read of the container, so the message matches the state.

Source code in src/winslow/workflow/workflow.py
def caches(self):
    """The caches this workflow can see, workflow scope first. After the
    session end the read refuses with direction: the caches are released.
    One read of the container, so the message matches the state."""
    container = self._workflow_cache
    if container is None:
        if self.session is not None and self.session.has_ended:
            raise ValueError(
                f"{self} has ended and released its caches - the recorded "
                f"cache reads live in the execution history (see "
                f"record_detail)."
            )
        raise InitializationError(
            f"{self} caches read before initialize_tasks built them."
        )
    return (*container.caches(), *self.global_cache.caches())

get_cache(name)

The live cache named name, or None.

Source code in src/winslow/workflow/workflow.py
def get_cache(self, name):
    """The live cache named `name`, or None."""
    for cache in self.caches():
        if cache.get_name() == name:
            return cache
    return None

init_state(state_store, origin=None, orchestrator_overrides=None, workflow_values=None)

Start persistence for the attached session: the manifest, and the persistence listener on the store. Call this once the pipeline is runnable, after the eligibility pass: the manifest marks the session as a restore candidate. A failure degrades to a run without state.

Source code in src/winslow/workflow/workflow.py
def init_state(
    self,
    state_store,
    origin=None,
    orchestrator_overrides=None,
    workflow_values=None,
):
    """Start persistence for the attached session: the manifest, and the
    persistence listener on the store. Call this once the pipeline is
    runnable, after the eligibility pass: the manifest marks the session
    as a restore candidate. A failure degrades to a run without state."""
    if self.persistence_listener is not None:
        # A second registration doubles the writes.
        return
    session = self._session
    adapter = sweeper = None
    try:
        adapter = SessionPersistenceAdapter(state_store, session.session_id)
        adapter.attach(self)
        self.persistence_listener = adapter
        sweeper = StaleSweeper(self)
        self.bus.subscribe(TaskStatusEvent, sweeper.on_task_status)
        self.stale_sweeper = sweeper
        # The manifest lands last: a failure before this point leaves no
        # durable state.
        state_store.save_manifest(
            SessionManifest(
                session_id=session.session_id,
                workflow_class=type(self).get_name(),
                workflow_namespace=self.cache_namespace,
                orchestrator_overrides=orchestrator_overrides,
                workflow_values=workflow_values,
                origin=origin,
                started_at=session.start,
            )
        )
    except Exception:
        # A persistence failure must not break the session start: the
        # session degrades to a run without state. The locals include a
        # subscriber whose registration failed.
        if adapter is not None:
            adapter.detach(self)
            adapter.close()
        if sweeper is not None:
            self.bus.unsubscribe(TaskStatusEvent, sweeper.on_task_status)
            sweeper.close()
        self.persistence_listener = None
        self.stale_sweeper = None
        self.logger.error(
            f"Could not persist the manifest of {session.session_id} - "
            f"the session runs without state",
            exc_info=True,
        )

load_snapshot(key)

The latest snapshot of the key, or None. The listener overlays the writes of this session on the initial state, so the read stays in memory (see SessionPersistenceAdapter.get).

Source code in src/winslow/workflow/workflow.py
def load_snapshot(self, key):
    """The latest snapshot of the key, or None. The listener overlays the
    writes of this session on the initial state, so the read stays in
    memory (see SessionPersistenceAdapter.get)."""
    listener = self.persistence_listener
    return listener.get(key) if listener is not None else None

seed_from_state()

Feed the persisted state of the session onto the store: the snapshots replay as statuses, and the open batch records register as INTERRUPTED. Call this after the eligibility pass: that pass overwrites earlier status writes.

Source code in src/winslow/workflow/workflow.py
def seed_from_state(self):
    """Feed the persisted state of the session onto the store: the
    snapshots replay as statuses, and the open batch records register as
    INTERRUPTED. Call this after the eligibility pass: that pass
    overwrites earlier status writes."""
    listener = self.persistence_listener
    if listener is None:
        return
    self._seed_task_statuses(listener.initial_state)
    self.runner.seed_interrupted_batches(listener.load_open_batches())

archive_state()

End persistence: stop the sweeper and the writer, unsubscribe them, then stamp and archive the manifest. After the durable writes the bus publishes SessionEndedEvent and closes, which disconnects every remaining subscriber. The session end calls this, and a persistence failure must not break the end. A session in ERROR archives as failed (see StateStore.mark_errored).

Source code in src/winslow/workflow/workflow.py
def archive_state(self):
    """End persistence: stop the sweeper and the writer, unsubscribe them,
    then stamp and archive the manifest. After the durable writes the bus
    publishes SessionEndedEvent and closes, which disconnects every
    remaining subscriber. The session end calls this, and a persistence
    failure must not break the end. A session in ERROR archives as failed
    (see StateStore.mark_errored)."""
    if (sweeper := self.stale_sweeper) is not None:
        sweeper.close()
        self.bus.unsubscribe(TaskStatusEvent, sweeper.on_task_status)
        self.stale_sweeper = None
    if (listener := self.persistence_listener) is not None:
        listener.close()
        self.bus.unsubscribe(TaskStatusEvent, listener.on_task_status)
        self.persistence_listener = None
        try:
            if self.session.status is SessionStatus.ERROR:
                listener.mark_errored()
            else:
                listener.mark_ended()
        except Exception:
            self.logger.error(
                f"Could not archive the manifest of {self.session_id}",
                exc_info=True,
            )
    self.bus.publish(SessionEndedEvent(session_id=self.session_id))
    self.bus.close()

roster_tasks()

The tasks a roster read serves, in launch-filter order. A bad launch filter logs and answers every task, so an interactive client still renders the list (see get_filtered_tasks).

Source code in src/winslow/workflow/workflow.py
def roster_tasks(self):
    """The tasks a roster read serves, in launch-filter order. A bad
    launch filter logs and answers every task, so an interactive client
    still renders the list (see get_filtered_tasks)."""
    try:
        return self.get_filtered_tasks()
    except MisconfigurationError:
        self.logger.error(
            "The launch filter does not parse - the roster lists every task.",
            exc_info=True,
        )
        return self.tasks

record_infos()

One TaskInfo per task with an execution record, across every batch. The record stores survive the session end, so a history search works after the task release (see ExecutionRecordStore).

Source code in src/winslow/workflow/workflow.py
def record_infos(self):
    """One TaskInfo per task with an execution record, across every
    batch. The record stores survive the session end, so a history
    search works after the task release (see ExecutionRecordStore)."""
    return tuple(
        {
            record.info.key: record.info
            for store in self.runner.record_stores()
            for record in store.records
        }.values()
    )

filter_keys(query, scope='tasks')

The identity keys the query matches over the named corpus: 'tasks' applies the full registry over the live tasks, 'history' the builtin filters over the record infos. Raises ValueError with direction.

Source code in src/winslow/workflow/workflow.py
def filter_keys(self, query, scope="tasks"):
    """The identity keys the query matches over the named corpus: 'tasks'
    applies the full registry over the live tasks, 'history' the builtin
    filters over the record infos. Raises ValueError with direction."""
    if scope not in ("tasks", "history"):
        raise ValueError(
            f"{scope!r} names no filter scope - the scopes are "
            f"'tasks' and 'history'."
        )
    parsed = self.filter_registry.parse(query)
    if scope == "history":
        enforce_builtin_only(parsed)
        return tuple(info.key for info in parsed.apply(self.record_infos()))
    if self.tasks is None:
        raise ValueError(
            f"{self} has ended and released its tasks - search the "
            f"execution records with scope='history'."
        )
    return tuple(task.identity_key for task in parsed.apply(self.tasks))

Task

winslow.Task

Bases: _ParameterizationBase

Source code in src/winslow/task/task.py
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
class Task(_ParameterizationBase):
    class Meta:
        abstract = True

    # The name attribute comes from _Base.
    groups = None

    # Examples of dependency declarations:
    # dependencies = MyTask
    # dependencies = (MyTask1, "MyTask2")
    # dependencies = "MyTask"

    dependencies = None

    # A premier task runs before all tasks that are not premier. A premier task
    # is a foundation step with no explicit dependencies downstream.
    is_premier = False

    # A terminal task runs last, independent of its priority.
    is_terminal = False

    # Tasks in the same priority group are processed concurrently. Set these to
    # False to prevent the run or the completion check from overlapping other
    # work in the batch.
    can_run_parallel = True
    can_check_parallel = True

    # The check TTL of this task, in seconds. It overrides the check_ttl of
    # the workflow; None defers to that default (see Workflow.check_ttl).
    check_ttl = None

    # Composable constraints: the Constraint classes that gate the matching
    # hook (see winslow.constraints). A single class or a collection. None
    # means that no constraint is declared. Declare these constraints or
    # override the hook directly. Both run when the gate is evaluated.
    eligibility_constraints = None
    runnability_constraints = None
    checkability_constraints = None
    success_constraints = None
    initialization_constraints = None

    def __init_subclass__(cls, **kwargs):
        super().__init_subclass__(**kwargs)
        # True when the task defines success. _evaluate_check reads this flag
        # and never calls the default check, which raises NotImplementedError.
        cls._check_overridden = cls.check is not Task.check

    @classmethod
    def _validate_constraint(cls, c, ctype):
        # One validation flow for every constraint that the get_*_constraints
        # methods return, applied at resolve time. The base class follows the
        # lifecycle: an instance hook takes Constraint, and graph init
        # (should_be_initialized) takes ClassConstraint.
        expected = ClassConstraint if ctype in _CLASS_CONSTRAINTS else Constraint
        if isinstance(c, _ConstraintBase):
            raise MisconfigurationError(
                f"{cls.__name__}: {ctype.name} constraint declared as an instance "
                f"({c!r}) - pass the class; the framework initializes it."
            )
        if not (isinstance(c, type) and issubclass(c, expected)):
            raise MisconfigurationError(
                f"{cls.__name__}: {ctype.name} constraint {c!r} must be a "
                f"{expected.__name__} subclass."
            )
        valid_types = c.get_valid_types()
        if valid_types and ctype not in valid_types:
            valid_names = ", ".join(t.name for t in valid_types)
            raise MisconfigurationError(
                f"{cls.__name__}: {c.__name__!r} is not valid as a {ctype.name} "
                f"constraint (valid_types=[{valid_names}])."
            )

    def __init__(self, workflow_config, parameters=None):

        super().__init__(parameters=parameters)

        self.workflow_config = workflow_config

        if self.is_premier and self.is_terminal:
            raise MisconfigurationError(
                f"{self} cannot be declared as premier and terminal at the same time."
            )

        # The graph sets these after it initializes the tasks.
        self._index = None
        self._dependent_tasks = None
        self._priority = None

        # The workflow stamps these at graph build (see
        # Workflow.initialize_tasks).
        self._workflow_cache_container = None
        self._global_cache_container = None
        self._run_nonce = None

        self._log_buffer = None

        # Local flag that marks the task as eligible or ineligible. The graph can
        # set it while it assigns the dependencies. The runner then does not
        # check the eligibility again.
        self._is_eligible_result = None

    def _enable_log_buffer(self):
        """Buffer the log records of this task for the Logs tab of the info
        modal. The graph calls this at setup time, in interactive mode only. The
        finalizer drops the buffer when the task is collected. A released task
        therefore frees its buffer, and a task that history retains keeps it."""
        self._log_buffer = collections.deque(maxlen=TASK_LOG_BUFFER_SIZE)
        dispatcher = get_task_dispatcher()
        dispatcher.register_buffer(self.log_key, self._log_buffer)
        weakref.finalize(self, dispatcher.unregister, self.log_key)

    @cached_property
    def logger(self):
        """All tasks share one logger. The adapter stamps log_key onto each
        record, so the dispatcher and a log view route by that key. The
        property is lazy: only a task that logs pays the key derivation."""
        return logging.LoggerAdapter(
            logging.getLogger(TASK_LOGGER_NAME), {"task_id": self.log_key}
        )

    @cached_property
    def log_key(self):
        """The process-local log routing key. The run nonce of the workflow
        separates two concurrent runs, so the log dispatcher does not mix
        their records. Session-durable structures use identity_key alone."""
        nonce = self._run_nonce
        return f"{nonce}:{self.identity_key}" if nonce else self.identity_key

    @property
    def buffered_logs(self):
        """Buffered log records for the Logs tab of the info modal. None when the
        task is not interactive."""
        return self._log_buffer

    @property
    def _active_execution_context(self):
        ctx = get_execution_context()
        if ctx is None:
            raise RuntimeError(
                f"Execution flags for {self} are only available during batch "
                f"execution (run/check) - not at import, eligibility or "
                f"graph-build time."
            )
        return ctx

    is_dry_run = _ExecutionFlag("dry_run")
    is_force_run = _ExecutionFlag("force_run")
    is_force_success = _ExecutionFlag("force_success")

    # The stamp attributes are set in __init__; the fallback serves the
    # graph-build hooks.
    workflow_cache = CacheContainerRef("_workflow_cache_container", get_workflow_cache)
    global_cache = CacheContainerRef("_global_cache_container", get_global_cache)

    @property
    def is_noop(self):
        """True for a task with no run action, which only checks completion. The
        value is inferred from run, because the base run does nothing. It thus
        inherits without a declaration: if any class in the MRO overrides run,
        the task is no longer noop."""
        return type(self).run is Task.run

    @property
    def can_process_parallel(self):
        # Processing includes the checks and the run, so it needs both gates.
        # A noop task only checks.
        if self.is_noop:
            return self.can_check_parallel
        return self.can_run_parallel and self.can_check_parallel

    @classmethod
    def get_parameters(cls, workflow_config):

        if not cls._is_parameterized:
            raise ValueError(
                f"{cls.__name__} declares no parameters - get_parameters "
                f"applies only to parameterized tasks."
            )

        raise _GetParametersNotImplemented(
            "get_parameters should be explicitly overridden for parametrized tasks to return a"
            " list of dictionaries that match the parameterization fields."
        )

    @classmethod
    def get_groups(cls):
        groups = cls.groups if cls.groups else []
        return frozenset(to_tuple(groups))

    @property
    def groups_readable(self):
        return ", ".join(sorted(self.get_groups()))

    @cached_property
    def _str_cached(self):

        label = str(self.instance_name)

        if self._is_parameterized:
            param_values_readable = ", ".join(
                safe_repr(v) for v in self._parameters_dict.values()
            )

            label = f"{label} ({param_values_readable})"

        return label

    def __str__(self):
        return self._str_cached

    @property
    def dependent_tasks(self):
        return self._dependent_tasks or tuple()

    @property
    def info(self):
        return TaskInfo.from_task(self)

    @classmethod
    def _get_dependency_context(cls):
        """Return the full dependency context of the task in flat form."""
        dependencies = cls.dependencies

        if not dependencies:
            return tuple()

        return flatten(to_tuple(dependencies))

    @property
    def priority(self):
        """A lower priority is processed first. The graph stamps this value with
        the topological generation of the task. This property does not
        calculate it."""
        if self._priority is None:
            raise MisconfigurationError(
                f"{self} priority read before the graph assigned it - read it "
                f"after initialize_tasks."
            )
        return self._priority

    @property
    def execution_order(self):
        return (not self.is_premier, self.is_terminal, self.priority)

    @classmethod
    def should_be_initialized(cls, workflow_config, parameters=None):
        """
        Override this to prevent the creation of the task instance.

        This is stronger than is_eligible, which filters the tasks after their
        creation. Use it with task parameterization when most of the
        parameterized tasks are skipped.

        A task that is not initialized does not show on the UI. This can confuse
        a user who declared the task.
        """
        return True

    @classmethod
    def get_initialization_constraints(cls, workflow_config):
        """Override this to select should_be_initialized constraints by env or by
        config."""
        return cls.initialization_constraints

    @classmethod
    def _evaluate_should_be_initialized(cls, workflow_config, parameters=None):
        # Class-level constraints take (task_class, parameters). The get_* method
        # supplies them. They go through the same validate and resolve flow, and
        # they run with the override.
        for c in to_tuple(cls.get_initialization_constraints(workflow_config) or ()):
            if not cls._prepare_constraint(
                c, ConstraintType.INITIALIZATION, workflow_config
            )(cls, parameters):
                return False
        return cls.should_be_initialized(workflow_config, parameters=parameters)

    def depends_on(self, task):
        """
        Override this to control if a task instance depends on another task
        instance.

        Use it to prevent a cyclic dependency between two task classes that
        depend on each other. Parameterization makes this more frequent.
        """
        return True

    def is_eligible(self):
        """Override this to control the eligibility of the task for the workflow."""
        return True

    def can_run(self):
        """Override this to block a task that does not satisfy your constraints."""
        return True

    def can_check(self):
        """Override this to gate the completion check. Return False, or call
        self.block, to mark the task BLOCKED and not run check. Use it when the
        check depends on data that is not always available."""
        return True

    def check(self):
        """Override this to implement the success check."""
        raise NotImplementedError

    # -- constraint evaluation -------------------------------------------------
    # The runner calls the sealed _evaluate_* methods and never the hooks. The
    # declared constraints thus always run with an override. The constraints run
    # first, so a cheap declarative guard can stop the evaluation before the
    # custom logic.

    # The get_*_constraints methods return the constraints for each gate. The
    # default is the declared class attribute. Override them to select the
    # constraints by env or by self.workflow_config.
    def get_eligibility_constraints(self):
        return self.eligibility_constraints

    def get_runnability_constraints(self):
        return self.runnability_constraints

    def get_checkability_constraints(self):
        return self.checkability_constraints

    def get_success_constraints(self):
        return self.success_constraints

    @cached_property
    def _constraint_instances(self):
        """Resolve the instance-level constraints one time, from the
        get_*_constraints methods. Validate each constraint, then instantiate
        it."""
        return {
            ctype: tuple(
                self._prepare_constraint(c, ctype, self.workflow_config)
                for c in to_tuple(getattr(self, f"get_{attr}")() or ())
            )
            for ctype, attr in _INSTANCE_CONSTRAINTS.items()
        }

    @classmethod
    def _prepare_constraint(cls, c, ctype, workflow_config):
        """Validate the constraint, then instantiate it. Every constraint from
        get_*_constraints goes through this one step."""
        cls._validate_constraint(c, ctype)
        return c(workflow_config)

    def _check_constraints(self, ctype):
        return all(c(self) for c in self._constraint_instances[ctype])

    def _evaluate_is_eligible(self):
        return (
            self._check_constraints(ConstraintType.ELIGIBILITY) and self.is_eligible()
        )

    def _check_eligibility(self, logger=LOGGER):
        """Resolve the eligibility one time and keep the result. The graph and
        the runner both call this. A crash in is_eligible aborts the run."""
        if self._is_eligible_result is not None:
            logger.debug(
                f"Eligibility already resolved for {self} - is_eligible will not be called."
            )
            return self._is_eligible_result

        logger.debug(f"Checking eligibility: {self}")
        try:
            result = self._evaluate_is_eligible()
        except exceptions.TaskSkip as e:
            logger.info(e)
            result = False
        except Exception as e:
            # A crash is not an answer. A False here drops the task silently, and
            # its dependents run without it (see Graph).
            logger.error(f"is_eligible crashed for {self}", exc_info=True)
            raise exceptions.EligibilityError(
                f"is_eligible crashed for {self}: {e}"
            ) from e
        self._is_eligible_result = result
        return result

    def _evaluate_can_run(self):
        return self._check_constraints(ConstraintType.RUNNABILITY) and self.can_run()

    def _evaluate_can_check(self):
        return self._check_constraints(ConstraintType.CHECKABILITY) and self.can_check()

    def _evaluate_check(self):
        constraints = self._constraint_instances[ConstraintType.SUCCESS]
        # The task must define success. If it does not, all([]) reports completed
        # with no check. This runs here and not at import, because
        # get_success_constraints supplies the constraints and can depend on the
        # env.
        if not self._check_overridden and not constraints:
            raise MisconfigurationError(
                f"{type(self).__name__} must define success: override check() "
                f"or provide success_constraints."
            )
        result = all(c(self) for c in constraints)
        if self._check_overridden:
            result = result and self.check()
        return result

    def run(self):
        """Make the system change. The check method then confirms the change."""
        pass

    def dry_run(self):
        """The runner calls this instead of run in dry-run mode. Override it to
        simulate the change. The default makes no change."""
        pass

    def block(self, msg):
        raise exceptions.TaskBlock(msg)

    def skip(self, msg):
        raise exceptions.TaskSkip(msg)

    def fail(self, msg):
        raise exceptions.TaskFailure(msg)

    def require_action(self, msg):
        raise TaskActionRequired(msg)

    def _check_and_raise(self, condition, msg, method):
        result = condition() if callable(condition) else condition

        if result:
            method(msg)

    def block_if(self, condition, msg):
        self._check_and_raise(condition, msg, self.block)

    def skip_if(self, condition, msg):
        self._check_and_raise(condition, msg, self.skip)

    def fail_if(self, condition, msg):
        self._check_and_raise(condition, msg, self.fail)

logger cached property

All tasks share one logger. The adapter stamps log_key onto each record, so the dispatcher and a log view route by that key. The property is lazy: only a task that logs pays the key derivation.

log_key cached property

The process-local log routing key. The run nonce of the workflow separates two concurrent runs, so the log dispatcher does not mix their records. Session-durable structures use identity_key alone.

buffered_logs property

Buffered log records for the Logs tab of the info modal. None when the task is not interactive.

is_noop property

True for a task with no run action, which only checks completion. The value is inferred from run, because the base run does nothing. It thus inherits without a declaration: if any class in the MRO overrides run, the task is no longer noop.

priority property

A lower priority is processed first. The graph stamps this value with the topological generation of the task. This property does not calculate it.

should_be_initialized(workflow_config, parameters=None) classmethod

Override this to prevent the creation of the task instance.

This is stronger than is_eligible, which filters the tasks after their creation. Use it with task parameterization when most of the parameterized tasks are skipped.

A task that is not initialized does not show on the UI. This can confuse a user who declared the task.

Source code in src/winslow/task/task.py
@classmethod
def should_be_initialized(cls, workflow_config, parameters=None):
    """
    Override this to prevent the creation of the task instance.

    This is stronger than is_eligible, which filters the tasks after their
    creation. Use it with task parameterization when most of the
    parameterized tasks are skipped.

    A task that is not initialized does not show on the UI. This can confuse
    a user who declared the task.
    """
    return True

get_initialization_constraints(workflow_config) classmethod

Override this to select should_be_initialized constraints by env or by config.

Source code in src/winslow/task/task.py
@classmethod
def get_initialization_constraints(cls, workflow_config):
    """Override this to select should_be_initialized constraints by env or by
    config."""
    return cls.initialization_constraints

depends_on(task)

Override this to control if a task instance depends on another task instance.

Use it to prevent a cyclic dependency between two task classes that depend on each other. Parameterization makes this more frequent.

Source code in src/winslow/task/task.py
def depends_on(self, task):
    """
    Override this to control if a task instance depends on another task
    instance.

    Use it to prevent a cyclic dependency between two task classes that
    depend on each other. Parameterization makes this more frequent.
    """
    return True

is_eligible()

Override this to control the eligibility of the task for the workflow.

Source code in src/winslow/task/task.py
def is_eligible(self):
    """Override this to control the eligibility of the task for the workflow."""
    return True

can_run()

Override this to block a task that does not satisfy your constraints.

Source code in src/winslow/task/task.py
def can_run(self):
    """Override this to block a task that does not satisfy your constraints."""
    return True

can_check()

Override this to gate the completion check. Return False, or call self.block, to mark the task BLOCKED and not run check. Use it when the check depends on data that is not always available.

Source code in src/winslow/task/task.py
def can_check(self):
    """Override this to gate the completion check. Return False, or call
    self.block, to mark the task BLOCKED and not run check. Use it when the
    check depends on data that is not always available."""
    return True

check()

Override this to implement the success check.

Source code in src/winslow/task/task.py
def check(self):
    """Override this to implement the success check."""
    raise NotImplementedError

run()

Make the system change. The check method then confirms the change.

Source code in src/winslow/task/task.py
def run(self):
    """Make the system change. The check method then confirms the change."""
    pass

dry_run()

The runner calls this instead of run in dry-run mode. Override it to simulate the change. The default makes no change.

Source code in src/winslow/task/task.py
def dry_run(self):
    """The runner calls this instead of run in dry-run mode. Override it to
    simulate the change. The default makes no change."""
    pass

TaskFilter

winslow.TaskFilter

Bases: Registerable

Source code in src/winslow/filter/base.py
class TaskFilter(Registerable):
    short_command = None
    long_command = None

    def __init__(self, value):
        self.value = value

    @classmethod
    @lru_cache
    def get_name(cls):
        # Use long_command first, then short_command, then the kebab class name.
        return cls.long_command or cls.short_command or super().get_name()

    def matches(self, task) -> bool:
        raise NotImplementedError

    def explain(self) -> str:
        raise NotImplementedError

Constraints

winslow.Constraint

Bases: _ConstraintBase

Constraint for the instance hooks: is_eligible, can_run, can_check and check. It is called with the live task.

Source code in src/winslow/constraints.py
class Constraint(_ConstraintBase):
    """Constraint for the instance hooks: is_eligible, can_run,
    can_check and check. It is called with the live task."""

    def __call__(self, task):
        return self.apply(task)

    def apply(self, task):
        raise NotImplementedError

winslow.ClassConstraint

Bases: _ConstraintBase

Constraint for should_be_initialized. It is evaluated at graph init, before an instance exists, and is called with the task class and the parameters.

Source code in src/winslow/constraints.py
class ClassConstraint(_ConstraintBase):
    """Constraint for should_be_initialized. It is evaluated at graph init, before
    an instance exists, and is called with the task class and the parameters."""

    def __call__(self, task_class, parameters=None):
        return self.apply(task_class, parameters)

    def apply(self, task_class, parameters=None):
        raise NotImplementedError

winslow.ConstraintType

Bases: Enum

Source code in src/winslow/constraints.py
class ConstraintType(Enum):
    # The task gate for each constraint list. INITIALIZATION is evaluated at
    # graph init, on the class (should_be_initialized), so its constraints take
    # (task_class, parameters). Each other type gates a live instance and takes
    # (task). CHECKABILITY gates the completion check and blocks the task if it
    # fails. SUCCESS adds to the success predicate (check).
    INITIALIZATION = auto()
    ELIGIBILITY = auto()
    RUNNABILITY = auto()
    CHECKABILITY = auto()
    SUCCESS = auto()

Descriptors

winslow.Parameter dataclass

Declare a parameter on a Task class. An instance is a Python descriptor. Access on the class returns the Parameter. Access on an instance returns the resolved value from task._params.

One Parameter usually binds one task attribute. from_tuple and from_dict declare a compound parameter, which binds more than one attribute. See those classmethods. name binds the value to an attribute with a different name than the declared one.

Source code in src/winslow/descriptors.py
@dataclass
class Parameter:
    """
    Declare a parameter on a Task class. An instance is a Python descriptor.
    Access on the class returns the Parameter. Access on an instance returns the
    resolved value from task._params.

    One Parameter usually binds one task attribute. `from_tuple` and `from_dict`
    declare a compound parameter, which binds more than one attribute. See those
    classmethods. `name` binds the value to an attribute with a different name
    than the declared one.
    """

    values: Union[Callable, Any, None] = None
    allowed_types: Optional[Tuple[type, ...]] = None
    param_style: ParameterStyle = ParameterStyle.SEQUENTIAL
    name: Optional[str] = None

    # An InitVar, to process the raw initialization value.
    _raw_values: InitVar[Union[Callable, Any, None]] = None

    # Markers for a compound parameter (from_tuple or from_dict). These are plain
    # class attributes and not fields.
    # _compound: this Parameter expands into more than one attribute.
    # _compound_kind: "tuple" or "dict".
    # _compound_names: the explicit attribute names from names= in from_tuple, or
    # None to derive them from the declared name.
    # _compound_name_map: the {dict-key: attribute-name} translation from
    # from_dict, or None.
    _compound = False
    _compound_kind = None
    _compound_names = None
    _compound_name_map = None

    def __post_init__(self, _raw_values):
        if _raw_values is None:
            _raw_values = self.values

        if (
            _raw_values is not None
            and not callable(_raw_values)
            and self.allowed_types is not None
        ):
            for value in to_tuple(_raw_values):
                if not isinstance(value, self.allowed_types):
                    raise ParameterizationError(
                        f"{value} is not of type: {self.allowed_types}."
                    )

    @classmethod
    def from_tuple(cls, values, names=None):
        """Declare more than one attribute from rows of aligned values.

        `foo_bar = Parameter.from_tuple([(1, 2), (3, 4)])` gives two tasks:
        `(foo=1, bar=2)` and `(foo=3, bar=4)`. The attribute names come from the
        declared name, which is split on `_`, or from `names=("foo", "bar")`.
        `values` is a list of rows, or a callable(workflow_config) that returns
        one. The rows are aligned, so each row is one combination. They combine
        with each other parameter as an independent cartesian axis.
        """
        param = cls(values=values)
        param._compound = True
        param._compound_kind = "tuple"
        param._compound_names = tuple(names) if names is not None else None
        return param

    @classmethod
    def from_dict(cls, values, name_map=None):
        """Declare more than one attribute from rows of dicts.

        `foo_bar = Parameter.from_dict([{"foo": 1, "bar": 2}, {"foo": 3, "bar": 4}])`
        gives two tasks: `(foo=1, bar=2)` and `(foo=3, bar=4)`. The attribute
        names come from the declared name, as in `from_tuple`. Each row must have
        the same keys. If the dict keys are different from the attribute names,
        pass `name_map={dict_key: attribute_name}` to translate them. `values` is
        a list of dicts, or a callable(workflow_config) that returns one. The rows
        combine with each other parameter as an independent cartesian axis.
        """
        param = cls(values=values)
        param._compound = True
        param._compound_kind = "dict"
        param._compound_name_map = dict(name_map) if name_map is not None else None
        return param

    def __set_name__(self, owner, name):
        self._attr_name = name

    def __get__(self, obj, objtype=None):
        if obj is None:
            return self
        return getattr(obj._params, self._attr_name)

    def __set__(self, obj, value):
        raise ParameterizationError(
            f"'{self._attr_name}' is a parameter and cannot be assigned."
        )

from_tuple(values, names=None) classmethod

Declare more than one attribute from rows of aligned values.

foo_bar = Parameter.from_tuple([(1, 2), (3, 4)]) gives two tasks: (foo=1, bar=2) and (foo=3, bar=4). The attribute names come from the declared name, which is split on _, or from names=("foo", "bar"). values is a list of rows, or a callable(workflow_config) that returns one. The rows are aligned, so each row is one combination. They combine with each other parameter as an independent cartesian axis.

Source code in src/winslow/descriptors.py
@classmethod
def from_tuple(cls, values, names=None):
    """Declare more than one attribute from rows of aligned values.

    `foo_bar = Parameter.from_tuple([(1, 2), (3, 4)])` gives two tasks:
    `(foo=1, bar=2)` and `(foo=3, bar=4)`. The attribute names come from the
    declared name, which is split on `_`, or from `names=("foo", "bar")`.
    `values` is a list of rows, or a callable(workflow_config) that returns
    one. The rows are aligned, so each row is one combination. They combine
    with each other parameter as an independent cartesian axis.
    """
    param = cls(values=values)
    param._compound = True
    param._compound_kind = "tuple"
    param._compound_names = tuple(names) if names is not None else None
    return param

from_dict(values, name_map=None) classmethod

Declare more than one attribute from rows of dicts.

foo_bar = Parameter.from_dict([{"foo": 1, "bar": 2}, {"foo": 3, "bar": 4}]) gives two tasks: (foo=1, bar=2) and (foo=3, bar=4). The attribute names come from the declared name, as in from_tuple. Each row must have the same keys. If the dict keys are different from the attribute names, pass name_map={dict_key: attribute_name} to translate them. values is a list of dicts, or a callable(workflow_config) that returns one. The rows combine with each other parameter as an independent cartesian axis.

Source code in src/winslow/descriptors.py
@classmethod
def from_dict(cls, values, name_map=None):
    """Declare more than one attribute from rows of dicts.

    `foo_bar = Parameter.from_dict([{"foo": 1, "bar": 2}, {"foo": 3, "bar": 4}])`
    gives two tasks: `(foo=1, bar=2)` and `(foo=3, bar=4)`. The attribute
    names come from the declared name, as in `from_tuple`. Each row must have
    the same keys. If the dict keys are different from the attribute names,
    pass `name_map={dict_key: attribute_name}` to translate them. `values` is
    a list of dicts, or a callable(workflow_config) that returns one. The rows
    combine with each other parameter as an independent cartesian axis.
    """
    param = cls(values=values)
    param._compound = True
    param._compound_kind = "dict"
    param._compound_name_map = dict(name_map) if name_map is not None else None
    return param

winslow.ConfigOption dataclass

Source code in src/winslow/descriptors.py
@dataclass
class ConfigOption:
    type: Optional[Type] = None
    default: Optional[Any] = None
    # The inverse of `type`. str() cannot invert it for a value that is not a
    # scalar: a list default becomes "[1]", which does not parse again. An option
    # can thus supply its own serializer from a value to a string.
    serializer: Optional[Callable] = None
    help_text: str = ""
    required: bool = False
    choices: Optional[List[Any]] = None
    multiselect: bool = False
    action: Optional[str] = None
    const: Optional[Any] = None  # used only if action == "store_const"
    show_on_ui: bool = True
    # Part of the display label of the instance. It implies required.
    identifier: bool = False
    subcommands: Tuple[Any, ...] = field(default_factory=tuple)
    depends_on: Tuple[str, ...] = field(default_factory=tuple)
    name: Optional[str] = None

    def __post_init__(self):
        self.subcommands = to_tuple(self.subcommands) if self.subcommands else tuple()
        self.depends_on = to_tuple(self.depends_on) if self.depends_on else tuple()
        if self.identifier:
            # An identifier with no value distinguishes nothing, so the user must
            # supply it.
            self.required = True

    def format_value(self, value):
        if value is None:
            return None
        if self.serializer:
            return self.serializer(value)
        if isinstance(value, (list, tuple)):
            # A multiselect value is a list. str() of a list is noisy in a label.
            return ", ".join(str(v) for v in value)
        return str(value)

    def to_arg(self, name=None, lenient=False):
        """Make the argparse argument definition. `lenient` removes `required`, so
        a parser that is shared between workflows does not abort on a mandatory
        option of one workflow. The strict parse of a single workflow enforces the
        value."""
        name = (name or self.name).replace("_", "-")

        common = dict(
            arg_name=f"--{name}",
            help=self.help_text,
            # A default supplies the value, so the command line does not demand
            # the option again.
            required=self.required and self.default is None and not lenient,
            default=self.default,
        )

        if self.action:
            if self.action == "store_const":
                return dict(common, action=self.action, const=self.const)
            else:
                return dict(common, action=self.action)

        result = dict(
            common,
            type=self.type,
            choices=self.choices,
        )
        if self.multiselect:
            # The form takes more than one value, so the command line must
            # take them too.
            result["nargs"] = "+"
        return result

to_arg(name=None, lenient=False)

Make the argparse argument definition. lenient removes required, so a parser that is shared between workflows does not abort on a mandatory option of one workflow. The strict parse of a single workflow enforces the value.

Source code in src/winslow/descriptors.py
def to_arg(self, name=None, lenient=False):
    """Make the argparse argument definition. `lenient` removes `required`, so
    a parser that is shared between workflows does not abort on a mandatory
    option of one workflow. The strict parse of a single workflow enforces the
    value."""
    name = (name or self.name).replace("_", "-")

    common = dict(
        arg_name=f"--{name}",
        help=self.help_text,
        # A default supplies the value, so the command line does not demand
        # the option again.
        required=self.required and self.default is None and not lenient,
        default=self.default,
    )

    if self.action:
        if self.action == "store_const":
            return dict(common, action=self.action, const=self.const)
        else:
            return dict(common, action=self.action)

    result = dict(
        common,
        type=self.type,
        choices=self.choices,
    )
    if self.multiselect:
        # The form takes more than one value, so the command line must
        # take them too.
        result["nargs"] = "+"
    return result

Decorators

winslow.decorators.transient_property = TransientProperty module-attribute

Caching

winslow.cache.entry(func=None, *, eager=False, depends_on=None, ttl=None, display_style=DisplayStyle.RAW)

Declare a cached field: cached_property plus a lock (see docs/caching.md). Treat every value as immutable: a fresh value comes from a recomputation.

Source code in src/winslow/cache/base.py
def entry(
    func=None, *, eager=False, depends_on=None, ttl=None, display_style=DisplayStyle.RAW
):
    """Declare a cached field: cached_property plus a lock (see docs/caching.md).
    Treat every value as immutable: a fresh value comes from a recomputation."""
    if func is not None:
        if (
            eager
            or depends_on
            or ttl is not None
            or display_style is not DisplayStyle.RAW
        ):
            raise MisconfigurationError(
                "entry(func, ...) drops the options - decorate with "
                "@entry(eager=..., depends_on=..., ttl=...) so they apply."
            )
        return Entry(func)

    def wrap(f):
        return Entry(
            f, eager=eager, depends_on=depends_on, ttl=ttl, display_style=display_style
        )

    return wrap

winslow.cache.base.BaseCache

The base of the declarative caches: the registries discover the concrete subclasses and validate them at collection (see docs/caching.md).

Source code in src/winslow/cache/base.py
class BaseCache:
    """The base of the declarative caches: the registries discover the
    concrete subclasses and validate them at collection (see docs/caching.md)."""

    class Meta:
        abstract = True

    # The attribute name on the container (see get_name).
    name = None

    # The scope label of the cache. The scope bases declare it.
    scope = None

    # The backend for the values of one instance, in line with graph_class and
    # registry_class on Workflow.
    storage_class = MemoryStorage

    # The byte cap of one rendered read in history. None resolves to the
    # WINSLOW_CACHE_SNAPSHOT_SIZE_BYTES setting; 0 records a summary only.
    snapshot_size_bytes = None

    def __init__(self):
        self._entries = declared_entries(type(self))
        self._storage = self.storage_class(self.get_name(), self._storage_namespace())
        # One lock per field, so two threads that hit a cold field compute it
        # once. Vanilla cached_property has no lock since Python 3.12.
        self._locks = {name: threading.Lock() for name in self._entries}
        # The error context per entry (see CacheEntryError). Mutated under the
        # entry's field lock; a successful write of the entry clears it.
        self._errors = {}
        # The container attaches itself after the construction. A cache built
        # outside a container keeps None and emits into the void.
        self._container = None

    def _attach_container(self, container):
        self._container = container

    def _emit(self, event, *args):
        """Send one listener event through the container (see CacheListener)."""
        if self._container is not None:
            self._container._emit(event, *args)

    def __str__(self):
        return self.get_name()

    @classmethod
    def get_name(cls):
        name = cls.name if cls.name is not None else camel_to_snake(cls.__name__)
        if not isinstance(name, str) or not name:
            raise MisconfigurationError(
                f"Invalid name definition for {cls} ({name!r}) - needs to be a non-empty string."
            )
        if not _CACHE_NAME_PATTERN.match(name):
            raise MisconfigurationError(
                f"Invalid cache name '{name}' for {cls}: must match [a-z][a-z0-9_]*, "
                f"because the name is an attribute on the cache container."
            )
        return name

    @property
    def logger(self):
        """The logger for a loader emission, resolved per access from the
        ambient context - never stored, because a GlobalCache outlives every
        session (see cache_logger)."""
        return cache_logger()

    def peek(self, name):
        """Return the storage record of the entry, MISSING, or
        EntryState.COMPUTING. No computation, no tier promotion, and a bounded
        wait: a running loader cannot stall an observer."""
        if name not in self._entries:
            raise AttributeError(
                f"Cache '{self}' has no entry {name!r} - "
                f"known entries: {sorted(self._entries)}"
            )
        if self._locks[name].acquire(timeout=_PEEK_GRACE_SECONDS):
            try:
                return self._storage.peek(name)
            finally:
                self._locks[name].release()
        # The grace absorbs every brief hold: only a loader holds a lock
        # longer (see _PEEK_GRACE_SECONDS).
        return EntryState.COMPUTING

    def _set_error(self, name, origin, tier, message):
        """Record the error context of the entry and report it. The caller
        runs inside an except block, so the traceback formats here. The next
        successful write of the entry clears the context."""
        error = CacheEntryError(
            origin=origin,
            tier=tier,
            message=message,
            at=time.time(),
            traceback=traceback.format_exc(),
        )
        self._errors[name] = error
        self._emit(
            CacheListener.on_entry_error, self.scope, self.get_name(), name, error
        )

    def describe_storage(self):
        """The human label of the storage of this cache. It reads no entry
        and takes no lock (see BaseStorage.describe)."""
        return self._storage.describe()

    def _failing_tier(self, exc):
        """The storage layer reference for a failed drop's error context."""
        if isinstance(exc, StorageError) and exc.tiers:
            return ", ".join(exc.tiers)
        return self.describe_storage()

    def inspect(self):
        """Return one CacheEntryInfo per declared entry. The values stay out;
        a detail view fetches one on demand with peek."""
        return tuple(self._entry_info(name, self.peek(name)) for name in self._entries)

    def _entry_info(self, name, record):
        """Build the projection from a record the caller peeked. No lock: an
        emission caller holds the field lock, and the inspect path accepts a
        record and an error context from different instants."""
        entry = self._entries[name]
        error = self._errors.get(name)
        state = self._entry_state(entry, record, error)
        return CacheEntryInfo(
            scope=self.scope,
            cache_name=self.get_name(),
            entry_name=name,
            state=state,
            # An ERRORED entry keeps its raw record observable: the UI shows
            # the quarantined value next to the error context.
            written_at=record.written_at if isinstance(record, StorageRecord) else None,
            ttl=entry.ttl,
            eager=entry.eager,
            depends_on=entry.depends_on,
            storage=self.describe_storage(),
            # Only an ERRORED entry carries its context: a healthy state next
            # to an error would contradict itself.
            error=error if state is EntryState.ERRORED else None,
        )

    @classmethod
    def _entry_state(cls, entry, record, error):
        """The freshness of one entry. One classifier serves the projection
        and the serve decision, so the two cannot diverge (see _entry_value)."""
        if record is EntryState.COMPUTING:
            return EntryState.COMPUTING
        # A record left behind by a failed drop is untrusted: it must not
        # serve, the same way an expired record must not.
        distrusted = error is not None and error.origin is ErrorOrigin.DELETE
        if record is not MISSING and not distrusted and not entry.is_expired(record):
            return EntryState.WARM
        if error is not None:
            return EntryState.ERRORED
        if record is MISSING:
            return EntryState.COLD
        return EntryState.STALE

    def _storage_namespace(self):
        """The key prefix that separates the scopes and the workflows in a
        persistent backend (see JsonFileStorage). The scope bases define it."""
        raise NotImplementedError

    def _own_chain(self):
        """The entries of this cache that the current logical thread is
        computing (see _read_chain)."""
        return tuple(name for cache, name in _read_chain.get() if cache is self)

    def _entry_value(self, entry):
        """Return the value of one entry. A cold or expired entry computes
        under its field lock and writes through the storage layer."""
        name = entry._attr_name
        chain = self._own_chain()
        if name in chain:
            # The thread already holds the lock of this field, so a second
            # acquisition would deadlock it silently and forever.
            cycle = " -> ".join((*chain[chain.index(name) :], name))
            raise CacheReentrancyError(
                f"Cache '{self}': the loader of '{chain[-1]}' reads '{name}' "
                f"while '{name}' is computing on the same thread ({cycle}) - "
                f"an undeclared read cycle."
            )
        with self._locks[name]:
            record = self._storage.read(name)
            error = self._errors.get(name)
            state = self._entry_state(entry, record, error)
            if state is EntryState.WARM:
                # The serve proves the value: a load mark cannot survive it.
                self._errors.pop(name, None)
                return record.value
            if state is EntryState.STALE:
                self.logger.info(
                    f"Cache '{self}': entry '{name}' went stale "
                    f"(age {time.time() - record.written_at:.1f}s, ttl {entry.ttl}s) - recomputing."
                )
            token = _read_chain.set((*_read_chain.get(), (self, name)))
            try:
                value = entry.func(self)
            except Exception as exc:
                # Only the outermost frame emits: a nested read raises again
                # through its caller, and each error reaches the backends once.
                # The quarantine comes first: a raising log handler or
                # telemetry backend must not leave the entry unmarked.
                if not chain:
                    self._set_error(name, ErrorOrigin.LOAD, None, str(exc))
                    self.logger.error(
                        f"Cache '{self}': the loader of '{name}' failed.",
                        exc_info=True,
                    )
                    emit_lazy_error(self, name, exc)
                raise
            finally:
                _read_chain.reset(token)
            # Serve what the storage stored: a serializing backend returns the
            # normalized round trip, so the shape never changes after a restart.
            stored = self._storage.write(
                name, StorageRecord(value=value, written_at=time.time())
            )
            # The write proves every writable tier holds the fresh value, so
            # any error context of the entry is obsolete.
            self._errors.pop(name, None)
            # The container guard keeps a bare cache free of the projection
            # cost; _emit would drop the event anyway.
            if self._container is not None:
                info = self._entry_info(name, stored)
                self._emit(CacheListener.on_entry_computed, info, state)
            return stored.value

    def invalidate(self, *names):
        """Drop the entries and their declared dependents, transitively. The
        next access recomputes each dropped entry, exactly like a ttl expiry."""
        if not names:
            raise TypeError(
                f"Cache '{self}': invalidate() takes at least one entry name - "
                f"to drop every entry, call invalidate_all()."
            )
        if unknown := [n for n in names if n not in self._entries]:
            raise AttributeError(
                f"Cache '{self}' has no entry {', '.join(repr(n) for n in unknown)} - "
                f"known entries: {sorted(self._entries)}"
            )
        graph = _dependency_graph(self._entries)
        affected = set(names).union(*(nx.descendants(graph, n) for n in names))
        order = [n for n in nx.topological_sort(graph) if n in affected]
        self._emit_invalidated(self._drop_entries(order, trigger=names), names)

    def invalidate_all(self):
        """Drop every entry of the instance."""
        self._emit_invalidated(self._drop_all(), None)

    def _drop_all(self):
        """Drop every entry and return the live drops, without an emission:
        clear_all on the container aggregates the caches into one event."""
        order = list(nx.topological_sort(_dependency_graph(self._entries)))
        return self._drop_entries(order, trigger=None)

    def _emit_invalidated(self, dropped, trigger):
        # Live drops only, matching the log line. The event fires once per
        # cascade, outside the field locks (see CacheListener).
        if dropped:
            self._emit(
                CacheListener.on_entries_invalidated,
                self.scope,
                {self.get_name(): dropped},
                _trigger_label(trigger),
            )

    def _drop_entries(self, names, trigger):
        """Drop upstream first: the other order lets a reader recompute a
        dependent from the stale upstream and keep it. Returns the names that
        held a live value."""
        chain = self._own_chain()
        if blocked := [name for name in names if name in chain]:
            # The thread holds the locks of its own chain, so the drop would
            # deadlock silently and forever (compare _entry_value).
            raise CacheReentrancyError(
                f"Cache '{self}': an invalidation from the loader of "
                f"'{chain[-1]}' reaches {', '.join(repr(n) for n in blocked)}, "
                f"which is computing on the same thread."
            )
        dropped = tuple(name for name in names if self._drop_entry(name))
        if dropped:
            self.logger.info(
                f"Cache '{self}': {_trigger_label(trigger)} dropped "
                f"{', '.join(repr(n) for n in dropped)}."
            )
        return dropped

    def _drop_entry(self, name):
        """Drop one entry and report whether a live value was present. One lock
        at a time: two held locks could deadlock against a nested computation.
        A storage failure quarantines the entry instead of aborting the cascade."""
        with self._locks[name]:
            try:
                record = self._storage.read(name)
                self._storage.delete(name)
            except Exception as exc:
                # The quarantine comes first: a raising log handler must not
                # leave the kept record servable.
                self._set_error(
                    name, ErrorOrigin.DELETE, self._failing_tier(exc), str(exc)
                )
                self.logger.error(
                    f"Cache '{self}': the drop of '{name}' failed - the entry "
                    f"is quarantined until a recompute overwrites it.",
                    exc_info=True,
                )
                # Counted as a live drop: the quarantine guarantees that the
                # kept record is never served.
                return True
            return record is not MISSING and not self._entries[name].is_expired(record)

logger property

The logger for a loader emission, resolved per access from the ambient context - never stored, because a GlobalCache outlives every session (see cache_logger).

peek(name)

Return the storage record of the entry, MISSING, or EntryState.COMPUTING. No computation, no tier promotion, and a bounded wait: a running loader cannot stall an observer.

Source code in src/winslow/cache/base.py
def peek(self, name):
    """Return the storage record of the entry, MISSING, or
    EntryState.COMPUTING. No computation, no tier promotion, and a bounded
    wait: a running loader cannot stall an observer."""
    if name not in self._entries:
        raise AttributeError(
            f"Cache '{self}' has no entry {name!r} - "
            f"known entries: {sorted(self._entries)}"
        )
    if self._locks[name].acquire(timeout=_PEEK_GRACE_SECONDS):
        try:
            return self._storage.peek(name)
        finally:
            self._locks[name].release()
    # The grace absorbs every brief hold: only a loader holds a lock
    # longer (see _PEEK_GRACE_SECONDS).
    return EntryState.COMPUTING

describe_storage()

The human label of the storage of this cache. It reads no entry and takes no lock (see BaseStorage.describe).

Source code in src/winslow/cache/base.py
def describe_storage(self):
    """The human label of the storage of this cache. It reads no entry
    and takes no lock (see BaseStorage.describe)."""
    return self._storage.describe()

inspect()

Return one CacheEntryInfo per declared entry. The values stay out; a detail view fetches one on demand with peek.

Source code in src/winslow/cache/base.py
def inspect(self):
    """Return one CacheEntryInfo per declared entry. The values stay out;
    a detail view fetches one on demand with peek."""
    return tuple(self._entry_info(name, self.peek(name)) for name in self._entries)

invalidate(*names)

Drop the entries and their declared dependents, transitively. The next access recomputes each dropped entry, exactly like a ttl expiry.

Source code in src/winslow/cache/base.py
def invalidate(self, *names):
    """Drop the entries and their declared dependents, transitively. The
    next access recomputes each dropped entry, exactly like a ttl expiry."""
    if not names:
        raise TypeError(
            f"Cache '{self}': invalidate() takes at least one entry name - "
            f"to drop every entry, call invalidate_all()."
        )
    if unknown := [n for n in names if n not in self._entries]:
        raise AttributeError(
            f"Cache '{self}' has no entry {', '.join(repr(n) for n in unknown)} - "
            f"known entries: {sorted(self._entries)}"
        )
    graph = _dependency_graph(self._entries)
    affected = set(names).union(*(nx.descendants(graph, n) for n in names))
    order = [n for n in nx.topological_sort(graph) if n in affected]
    self._emit_invalidated(self._drop_entries(order, trigger=names), names)

invalidate_all()

Drop every entry of the instance.

Source code in src/winslow/cache/base.py
def invalidate_all(self):
    """Drop every entry of the instance."""
    self._emit_invalidated(self._drop_all(), None)

winslow.cache.GlobalCache

Bases: BaseCache

A cache with process scope, shared by the workflows. The instances live for the process (see winslow.cache.get_global_cache).

Source code in src/winslow/cache/base.py
class GlobalCache(BaseCache):
    """A cache with process scope, shared by the workflows. The instances live
    for the process (see winslow.cache.get_global_cache)."""

    class Meta:
        abstract = True

    def __init__(self, orchestrator_config):
        super().__init__()
        self.orchestrator_config = orchestrator_config

    scope = GLOBAL_SCOPE

    def _storage_namespace(self):
        return "global"

winslow.cache.WorkflowCache

Bases: BaseCache

A cache with session scope. A new session builds fresh instances, and the container is dropped when the session ends.

Source code in src/winslow/cache/base.py
class WorkflowCache(BaseCache):
    """A cache with session scope. A new session builds fresh instances, and
    the container is dropped when the session ends."""

    class Meta:
        abstract = True

    def __init__(self, workflow_config):
        # Before super(), which builds the storage from the namespace.
        self.workflow_config = workflow_config
        super().__init__()

    scope = WORKFLOW_SCOPE

    def _storage_namespace(self):
        # The workflow stamps its identity (see Workflow.cache_namespace). The
        # workflows/ segment keeps any stamp out of the global scope.
        namespace = getattr(self.workflow_config, "cache_namespace", None)
        if namespace is None:
            LOGGER.warning(
                f"{type(self).__name__} was built outside a workflow - its "
                f"storage shares the 'workflows/_unscoped' namespace."
            )
            namespace = "_unscoped"
        return f"workflows/{namespace}"

winslow.cache.get_workflow_cache()

The workflow container of the current context.

Source code in src/winslow/cache/runtime.py
def get_workflow_cache():
    """The workflow container of the current context."""
    container = _workflow_container.get()
    if container is None:
        raise RuntimeError(
            "No active workflow cache container in this context - it is "
            "available during the workflow initialization and on the tasks."
        )
    return container

winslow.cache.get_global_cache()

The process-level container. It exists after the first workflow initialization.

Source code in src/winslow/cache/runtime.py
def get_global_cache():
    """The process-level container. It exists after the first workflow
    initialization."""
    if _global_container is None:
        raise RuntimeError(
            "The global cache container does not exist yet - it is built at "
            "the first workflow initialization."
        )
    return _global_container

winslow.cache.BaseStorage

The storage contract of a cache. The built-in backends subclass it; a custom backend overrides the three methods (see docs/caching.md).

Source code in src/winslow/cache/storage.py
class BaseStorage:
    """The storage contract of a cache. The built-in backends subclass it; a
    custom backend overrides the three methods (see docs/caching.md)."""

    # A read-only backend is a pure source: ComposedStorage skips it on write
    # and delete.
    read_only = False

    def __init__(self, cache_name, namespace):
        self.cache_name = cache_name
        self.namespace = namespace

    def read(self, key):
        """Return the StorageRecord of the key, or the MISSING sentinel. None
        cannot mark a miss: it is a legal cached value."""
        raise NotImplementedError

    def peek(self, key):
        """Return the record of the key, or MISSING, with no side effect. The
        default delegates to read; a backend whose read has a side effect
        overrides this, so an observation never changes the stored state."""
        return self.read(key)

    def describe(self):
        """The human label of the backend for the UI projections."""
        return type(self).__name__

    def write(self, key, record):
        """Store the record and return the stored record. The caller serves
        the returned value, so a normalizing backend returns its round trip."""
        raise NotImplementedError

    def delete(self, key):
        """Remove the record of the key, also when the key holds no record."""
        raise NotImplementedError

read(key)

Return the StorageRecord of the key, or the MISSING sentinel. None cannot mark a miss: it is a legal cached value.

Source code in src/winslow/cache/storage.py
def read(self, key):
    """Return the StorageRecord of the key, or the MISSING sentinel. None
    cannot mark a miss: it is a legal cached value."""
    raise NotImplementedError

peek(key)

Return the record of the key, or MISSING, with no side effect. The default delegates to read; a backend whose read has a side effect overrides this, so an observation never changes the stored state.

Source code in src/winslow/cache/storage.py
def peek(self, key):
    """Return the record of the key, or MISSING, with no side effect. The
    default delegates to read; a backend whose read has a side effect
    overrides this, so an observation never changes the stored state."""
    return self.read(key)

describe()

The human label of the backend for the UI projections.

Source code in src/winslow/cache/storage.py
def describe(self):
    """The human label of the backend for the UI projections."""
    return type(self).__name__

write(key, record)

Store the record and return the stored record. The caller serves the returned value, so a normalizing backend returns its round trip.

Source code in src/winslow/cache/storage.py
def write(self, key, record):
    """Store the record and return the stored record. The caller serves
    the returned value, so a normalizing backend returns its round trip."""
    raise NotImplementedError

delete(key)

Remove the record of the key, also when the key holds no record.

Source code in src/winslow/cache/storage.py
def delete(self, key):
    """Remove the record of the key, also when the key holds no record."""
    raise NotImplementedError

winslow.cache.MemoryStorage

Bases: BaseStorage

The default backend: a dict per cache instance. It puts no constraint on the values.

Source code in src/winslow/cache/storage.py
class MemoryStorage(BaseStorage):
    """The default backend: a dict per cache instance. It puts no constraint
    on the values."""

    def __init__(self, cache_name, namespace):
        super().__init__(cache_name, namespace)
        self._records = {}

    def read(self, key):
        """Return the record of the key, or MISSING."""
        return self._records.get(key, MISSING)

    def write(self, key, record):
        """Store the record verbatim and return it."""
        self._records[key] = record
        return record

    def delete(self, key):
        """Remove the record of the key. An absent key is a no-op."""
        self._records.pop(key, None)

read(key)

Return the record of the key, or MISSING.

Source code in src/winslow/cache/storage.py
def read(self, key):
    """Return the record of the key, or MISSING."""
    return self._records.get(key, MISSING)

write(key, record)

Store the record verbatim and return it.

Source code in src/winslow/cache/storage.py
def write(self, key, record):
    """Store the record verbatim and return it."""
    self._records[key] = record
    return record

delete(key)

Remove the record of the key. An absent key is a no-op.

Source code in src/winslow/cache/storage.py
def delete(self, key):
    """Remove the record of the key. An absent key is a no-op."""
    self._records.pop(key, None)

winslow.cache.JsonFileStorage

Bases: BaseStorage

One JSON file per entry, under a directory per scope and cache. The write is strict and atomic; a corrupt or unreadable file reads as MISSING.

Source code in src/winslow/cache/storage.py
class JsonFileStorage(BaseStorage):
    """One JSON file per entry, under a directory per scope and cache. The
    write is strict and atomic; a corrupt or unreadable file reads as MISSING."""

    # None resolves to the WINSLOW_CACHE_DIR setting. A test or a project can
    # override it on a subclass. The default is relative to the CWD.
    base_directory = None

    def __init__(self, cache_name, namespace):
        _validate_path_components(cache_name, namespace)
        super().__init__(cache_name, namespace)
        base = self.base_directory or settings.CACHE_DIR
        # The namespace keeps same-named caches of different scopes apart.
        self.directory = Path(base) / namespace / cache_name

    def _path(self, key):
        return self.directory / f"{key}.json"

    def read(self, key):
        try:
            return self._decode(self._path(key).read_text(encoding="utf-8"))
        except FileNotFoundError:
            return MISSING
        except (OSError, DeserializationError):
            LOGGER.error(
                f"Cache file {self._path(key)} is unreadable - "
                f"treating the entry as cold.",
                exc_info=True,
            )
            return MISSING

    @classmethod
    def _decode(cls, text):
        """Decode one stored record. A corrupt payload raises, and the caller
        owns the policy: read() serves a cold miss instead."""
        try:
            payload = json.loads(text)
            return StorageRecord(
                value=payload["value"], written_at=payload["written_at"]
            )
        except (ValueError, KeyError, TypeError) as exc:
            raise DeserializationError(
                f"the record is not a valid cache payload ({exc})"
            ) from exc

    def write(self, key, record):
        """Store the record and return its JSON round trip: the caller serves
        the normalized value, so a value never changes shape after a restart."""
        try:
            text = json.dumps({"written_at": record.written_at, "value": record.value})
        except TypeError as exc:
            # Strict on purpose: a default=str fallback would coerce silently,
            # and a lossy cache file is not debuggable.
            raise SerializationError(
                f"Cache '{self.cache_name}', entry '{key}': "
                f"the value is not JSON-serializable ({exc})."
            ) from exc
        self.directory.mkdir(parents=True, exist_ok=True)
        # A private temp name per writer: two processes that write the same
        # entry cannot publish each other's bytes through a shared temp file.
        with tempfile.NamedTemporaryFile(
            "w", dir=self.directory, suffix=".json.tmp", delete=False, encoding="utf-8"
        ) as temp:
            temp.write(text)
        os.replace(temp.name, self._path(key))
        return StorageRecord(
            value=json.loads(text)["value"], written_at=record.written_at
        )

    def delete(self, key):
        self._path(key).unlink(missing_ok=True)

write(key, record)

Store the record and return its JSON round trip: the caller serves the normalized value, so a value never changes shape after a restart.

Source code in src/winslow/cache/storage.py
def write(self, key, record):
    """Store the record and return its JSON round trip: the caller serves
    the normalized value, so a value never changes shape after a restart."""
    try:
        text = json.dumps({"written_at": record.written_at, "value": record.value})
    except TypeError as exc:
        # Strict on purpose: a default=str fallback would coerce silently,
        # and a lossy cache file is not debuggable.
        raise SerializationError(
            f"Cache '{self.cache_name}', entry '{key}': "
            f"the value is not JSON-serializable ({exc})."
        ) from exc
    self.directory.mkdir(parents=True, exist_ok=True)
    # A private temp name per writer: two processes that write the same
    # entry cannot publish each other's bytes through a shared temp file.
    with tempfile.NamedTemporaryFile(
        "w", dir=self.directory, suffix=".json.tmp", delete=False, encoding="utf-8"
    ) as temp:
        temp.write(text)
    os.replace(temp.name, self._path(key))
    return StorageRecord(
        value=json.loads(text)["value"], written_at=record.written_at
    )

winslow.cache.compose(*storage_classes)

A storage class with the given tiers, read order first to last (see docs/caching.md).

Source code in src/winslow/cache/storage.py
def compose(*storage_classes):
    """A storage class with the given tiers, read order first to last (see
    docs/caching.md)."""
    if not storage_classes:
        raise MisconfigurationError("compose() needs at least one storage class.")
    if len(storage_classes) == 1:
        return storage_classes[0]
    # A real class, never a partial: functools.partial binds the instance as
    # a descriptor on Python 3.14, which breaks the class-attribute call.
    return type(
        "ComposedStorage", (ComposedStorage,), {"storage_classes": storage_classes}
    )

Telemetry

winslow.telemetry.TelemetryConfiguration

The configuration-first activation seam: declared in the telemetry.py files of a repo, driven by the orchestrator (see docs/telemetry.md).

Source code in src/winslow/telemetry.py
class TelemetryConfiguration:
    """The configuration-first activation seam: declared in the telemetry.py
    files of a repo, driven by the orchestrator (see docs/telemetry.md)."""

    def get_handler(self, orchestrator_config):
        """Return the TelemetryHandler to register, or None to stay
        inactive. An exception fails the start of the run loudly."""
        return None

    def shutdown(self):
        """Flush and release the backend, once, before the process exits."""

    @classmethod
    def get_name(cls):
        return derive_name(cls)

get_handler(orchestrator_config)

Return the TelemetryHandler to register, or None to stay inactive. An exception fails the start of the run loudly.

Source code in src/winslow/telemetry.py
def get_handler(self, orchestrator_config):
    """Return the TelemetryHandler to register, or None to stay
    inactive. An exception fails the start of the run loudly."""
    return None

shutdown()

Flush and release the backend, once, before the process exits.

Source code in src/winslow/telemetry.py
def shutdown(self):
    """Flush and release the backend, once, before the process exits."""

winslow.telemetry.TelemetryHandler

The interface of a telemetry backend. Each callback default is a no-op, so a subclass overrides only the callbacks that it consumes.

Source code in src/winslow/telemetry.py
class TelemetryHandler:
    """The interface of a telemetry backend. Each callback default is a
    no-op, so a subclass overrides only the callbacks that it consumes."""

    def on_task_error(self, workflow, task, exc, batch_uuid, phase):
        """An errored task step (see BaseRunner.task_scope). Each errored
        step of a task emits one call. Do not retain the live objects."""

    def on_unscoped_error(
        self,
        exc,
        workflow_name=None,
        session_id=None,
        workflow_instance=None,
        workflow_class=None,
    ):
        """An error outside every task scope. workflow_instance carries
        str(workflow); each argument is None when it is not known yet."""

on_task_error(workflow, task, exc, batch_uuid, phase)

An errored task step (see BaseRunner.task_scope). Each errored step of a task emits one call. Do not retain the live objects.

Source code in src/winslow/telemetry.py
def on_task_error(self, workflow, task, exc, batch_uuid, phase):
    """An errored task step (see BaseRunner.task_scope). Each errored
    step of a task emits one call. Do not retain the live objects."""

on_unscoped_error(exc, workflow_name=None, session_id=None, workflow_instance=None, workflow_class=None)

An error outside every task scope. workflow_instance carries str(workflow); each argument is None when it is not known yet.

Source code in src/winslow/telemetry.py
def on_unscoped_error(
    self,
    exc,
    workflow_name=None,
    session_id=None,
    workflow_instance=None,
    workflow_class=None,
):
    """An error outside every task scope. workflow_instance carries
    str(workflow); each argument is None when it is not known yet."""

Orchestrator

winslow.Orchestrator

Bases: _ConfigBase

Source code in src/winslow/orchestrator.py
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
class Orchestrator(_ConfigBase):
    workflow_registry_class = WorkflowRegistry
    telemetry_registry_class = TelemetryRegistry
    global_cache_registry_class = GlobalCacheRegistry

    workflow = ConfigOption(
        help_text="Name of the workflow to view / run.",
        subcommands=(
            Command.RUN.value,
            Command.SHOW.value,
        ),
        # The UI has a workflow selection widget, so it needs no text input.
        show_on_ui=False,
    )

    host = ConfigOption(
        help_text="The bind address of the serve process. Loopback needs no credential.",
        default="127.0.0.1",
        subcommands=Command.SERVE.value,
        show_on_ui=False,
    )

    port = ConfigOption(
        help_text="The port of the serve process.",
        type=int,
        default=8866,
        subcommands=Command.SERVE.value,
        show_on_ui=False,
    )

    endpoints = ConfigOption(
        help_text="The endpoints to serve: ws, mcp, or both. mcp requires the mcp extra.",
        choices=list(ENDPOINTS),
        multiselect=True,
        default=["ws"],
        subcommands=Command.SERVE.value,
        show_on_ui=False,
    )

    filter = ConfigOption(
        help_text=(
            "Filter expression for tasks. Supports name matching (bare text), "
            "filter commands (!g <group>, !group <group>), "
            "boolean operators (& |), negation (~), grouping (()), "
            "and comma-separated OR shorthand (foo,bar)."
        ),
        subcommands=(
            Command.RUN.value,
            Command.SHOW.value,
        ),
        depends_on="initialize",
        show_on_ui=False,
    )

    initialize = ConfigOption(
        action="store_true",
        default=False,
        help_text=(
            "Initialize a single workflow (chosen with --workflow) and list its "
            "tasks. Required in order to use --filter with show."
        ),
        subcommands=Command.SHOW.value,
        show_on_ui=False,
    )

    with_deps = ConfigOption(
        action="store_true",
        default=False,
        help_text=(
            "With --initialize, also list each task's dependencies, in the order "
            "they run."
        ),
        subcommands=Command.SHOW.value,
        depends_on="initialize",
        show_on_ui=False,
    )

    mode = ConfigOption(
        type=_parse_mode,
        choices=tuple(Mode),
        help_text=(
            "How to run the workflow(s): 'tui' launches the interactive terminal"
            " UI; 'headless' runs single-thread/single-process with no UI, useful"
            " for CI and debugging."
        ),
        default=Mode.TUI,
        subcommands=Command.RUN.value,
        show_on_ui=False,
    )

    check = ConfigOption(
        action="store_true",
        help_text=(
            "Checks the tasks in the workflow instead of running them"
            " - available in headless run mode."
        ),
        default=False,
        subcommands=Command.RUN.value,
        show_on_ui=False,
    )

    disable_concurrency = ConfigOption(
        action="store_true",
        help_text=(
            "Disables all concurrency during task runs (e.g. task dependencies and task "
            "eligibility will be checked sequentially)."
        ),
        default=False,
        subcommands=Command.RUN.value,
    )

    clear_cache = ConfigOption(
        action="store_true",
        help_text=(
            "Invalidate every cache entry at workflow initialization, before "
            "the eager population. Meaningful for persistent cache storage; "
            "a memory cache starts cold anyway. A read-only storage tier "
            "keeps its records, and a later read can promote them back."
        ),
        default=False,
        subcommands=Command.RUN.value,
    )

    dry_run = ConfigOption(
        action="store_true",
        help_text=(
            "Run tasks in dry-run mode (dry_run method will be "
            "called on tasks, instead of run)"
        ),
        default=False,
        subcommands=Command.RUN.value,
    )

    force_run = ConfigOption(
        action="store_true",
        help_text=(
            "Skips implicit success checks before a task run. "
            "Runnability and dependency checks will still be in effect."
        ),
        default=False,
        subcommands=Command.RUN.value,
    )

    force_success = ConfigOption(
        action="store_true",
        help_text=(
            "Mark tasks as successful (FORCE_SUCCESS) without running any checks "
            "or the task itself. Ineligible (skipped) tasks are left untouched."
        ),
        default=False,
        subcommands=Command.RUN.value,
    )

    reraise_errors = ConfigOption(
        action="store_true",
        help_text=(
            "Re-raise unexpected task errors after marking the task ERROR, "
            "aborting the run instead of continuing. For CI and debugging."
        ),
        default=False,
        subcommands=Command.RUN.value,
        show_on_ui=False,
    )

    debug = ConfigOption(
        action="store_true",
        help_text="Enable debug mode and logging.",
        default=False,
        subcommands=(
            Command.RUN.value,
            Command.SHOW.value,
        ),
        # The command line enables the debug mode.
        show_on_ui=False,
    )

    def __init__(self, orchestrator_config, directory=None, unknown_args=None):
        self.directory = directory or os.getcwd()
        # argparse reads sys.argv if this is None. Normalize it at the boundary.
        self.unknown_args = list(unknown_args or ())

        super().__init__(orchestrator_config)

        self.workflow_registry = self.workflow_registry_class(self.orchestrator_config)
        self.telemetry_registry = self.telemetry_registry_class(
            self.orchestrator_config
        )
        self.global_cache_registry = self.global_cache_registry_class(
            self.orchestrator_config
        )

        # The TUI application. run and connect create it.

        self.logger = LOGGER

    @classmethod
    def _add_arguments(cls, parser, subcommand=None):
        for arg_ctx in cls.get_argparse_context(subcommand):
            try:
                name = arg_ctx.pop("arg_name")
                parser.add_argument(name, **arg_ctx)
            except (TypeError, KeyError):
                cmd = subcommand if subcommand else "base"
                LOGGER.error(
                    f"Could not add argument for {cmd} parser {parser}, context: '{name}' - {arg_ctx}"
                )
                raise

    @classmethod
    def _generate_subcommand(
        cls,
        subparsers,
        action,
        help_text,
    ):
        sub_parser = subparsers.add_parser(action.value, help=help_text)

        cls._add_arguments(sub_parser, subcommand=action.value)
        sub_parser.set_defaults(action=action)

        return sub_parser

    @classmethod
    def get_base_parser(cls):

        parser = argparse.ArgumentParser(
            prog="Winslow",
            description="A state and workflow management framework with a terminal UI",
        )

        cls._add_arguments(parser)

        subparsers = parser.add_subparsers(help="Available subcommands")

        all_defaults = {
            name: conf.default
            for name, conf in cls.config_meta.items()
            if conf.subcommands
        }
        parser.set_defaults(**all_defaults)

        cls._generate_subcommand(
            subparsers,
            action=Command.SHOW,
            help_text="Show workflow / task information.",
        )

        cls._generate_subcommand(
            subparsers, action=Command.RUN, help_text="Run workflow(s)."
        )

        cls._generate_subcommand(
            subparsers,
            action=Command.SERVE,
            help_text="Serve the live sessions over one websocket endpoint.",
        )

        connect_parser = cls._generate_subcommand(
            subparsers,
            action=Command.CONNECT,
            help_text="Run the TUI against a serve process.",
        )
        # The ConfigOption machinery declares only --flag options, so the
        # URL is added here as the one positional argument of the CLI.
        connect_parser.add_argument(
            "url",
            help="The websocket URL of the serve process, e.g. ws://host:8866. "
            "A non-loopback server reads the bearer token from WINSLOW_TOKEN.",
        )

        return parser

    @classmethod
    def parse_args(cls):
        base_parser = cls.get_base_parser()
        base_args, unknown_args = base_parser.parse_known_args(
            namespace=OrchestratorConfig()
        )
        return base_args, unknown_args

    @property
    def sorted_workflow_classes(self):
        return sorted(
            self.workflow_registry.classes,
            key=lambda kls: (kls.get_module_path(), kls.__name__),
        )

    @classmethod
    def get_from_cwd(cls):
        klasses = []
        for scoped_name in iter_dir_module_names(os.getcwd(), recursive=False):
            try:
                module = importlib.import_module(scoped_name)
            # Catch SystemExit too. A sys.exit() in an unguarded script must not
            # stop the CLI.
            except (Exception, SystemExit) as e:
                # The orchestrator is optional. A broken file that has no
                # relation to it must not stop the CLI.
                LOGGER.warning(
                    f"Skipping {scoped_name.split('.', 1)[1]}.py - "
                    f"{type(e).__name__}: {e}",
                    exc_info=True,
                )
                continue
            klasses.extend(classes_in_module(module, cls))

        if len(klasses) > 1:
            names = ", ".join(kls.__name__ for kls in klasses)

            raise MisconfigurationError(
                f"Multiple concrete orchestrator classes found in {os.getcwd()}: {names}. "
                f"Please define one orchestrator at most per directory."
            )

        return klasses[0] if klasses else None

    def _parse_known_workflow_args(self, workflow_kls, unknown):
        """Returns the parsed args, or None, and the tokens that this workflow did
        not claim."""
        parser = workflow_kls.get_parser(lenient=True)
        if parser is None:
            return None, tuple(unknown)
        try:
            # Lenient: the args of another workflow are not for this parser.
            args, leftover = parser.parse_known_args(unknown)
        except SystemExit:
            # argparse already printed the message. Send the exit to the clean
            # error path.
            raise MisconfigurationError(
                f"Invalid arguments for workflow '{workflow_kls.get_name()}' (see above)."
            )
        return args, tuple(leftover)

    @classmethod
    def _leftover_indices(cls, unknown, leftover):
        """The positions in `unknown` that argparse did not claim for this
        workflow. The match is by position, because two workflows can claim
        equal tokens: `--my-integer 7 --batch-size 7` leaves a `7` in each
        leftover list, and a match by value would report that `7` as unclaimed."""
        positions = set()
        cursor = 0
        for i, token in enumerate(unknown):
            if cursor < len(leftover) and token == leftover[cursor]:
                positions.add(i)
                cursor += 1
        return positions

    def collect_workflow_args(self):
        """{workflow_kls: parsed args or None}, validated together.

        `show`, the interactive start, and a serve process's descriptors
        request call this. They parse the args of each discovered workflow
        at one time, so a start form (local or remote) prefills from the
        same CLI-supplied values. A headless run has one named workflow and
        parses its args strictly (see `_handle_headless_run`), so the
        collision between workflows that this method handles cannot occur
        there."""
        parsed = {
            kls: self._parse_known_workflow_args(kls, self.unknown_args)
            for kls in self.sorted_workflow_classes
        }

        # A token is unrecognized if no workflow claimed its position. A position
        # that one workflow claimed leaves the intersection.
        unclaimed = set(range(len(self.unknown_args)))
        for _, leftover in parsed.values():
            unclaimed &= self._leftover_indices(self.unknown_args, leftover)

        if unclaimed:
            tokens = [self.unknown_args[i] for i in sorted(unclaimed)]
            raise MisconfigurationError(f"Unrecognized arguments: {' '.join(tokens)}")

        return {kls: args for kls, (args, _) in parsed.items()}

    def _display_workflow(self, idx, workflow, show_tasks):
        self.logger.info(
            f"{idx + 1}. {workflow.__class__.__name__} ({workflow.instance_name})"
            f" - {workflow.module_directory}"
        )

        if show_tasks:
            self._display_tasks(workflow)
        else:
            self._display_task_classes(workflow)

    def _display_task_classes(self, workflow):
        for i, kls in enumerate(
            sorted(workflow.registry.classes, key=lambda k: k.get_name())
        ):
            parameterized = getattr(kls, "_is_parameterized", False)
            marker = "  [parameterized]" if parameterized else ""
            self.logger.info(f"{INDENT}{i + 1}. {kls.__name__}{marker}")

            if parameterized:
                for declared_name, param in kls._parameterization_meta.items():
                    self.logger.info(
                        f"{INDENT * 2}- {self._describe_parameter(kls, declared_name, param)}"
                    )

    @classmethod
    def _describe_parameter(cls, kls, declared_name, param):
        style = param.param_style.name.lower()
        if param._compound:
            names = ", ".join(kls._parameter_attrs[declared_name])
            return f"{declared_name} -> {names}  ({style})"
        return f"{declared_name}  ({style})"

    def _display_tasks(self, workflow):
        for task in workflow.get_filtered_tasks():
            self.logger.info(f"{INDENT}{task._index + 1}. {str(task)}")

            if self.orchestrator_config.with_deps and task.dependent_tasks:
                self.logger.info(f"{INDENT * 2}Dependencies:")
                for dep in sorted(task.dependent_tasks, key=lambda d: d._index):
                    self.logger.info(f"{INDENT * 3}{dep._index + 1}. {str(dep)}")

    def check_filters(self, workflow):
        wanted = getattr(self.orchestrator_config, "workflow", None)
        return not wanted or workflow.get_name() == wanted

    def _handle_show(self):
        if self.orchestrator_config.initialize:
            return self._handle_initialize()
        self._list_workflow_classes()

    def _list_workflow_classes(self):
        self.logger.debug("Listing workflow classes")
        workflow_args_map = self.collect_workflow_args()
        workflows = []

        for workflow_kls in self.sorted_workflow_classes:
            if not workflow_kls.should_be_initialized(self.orchestrator_config):
                continue

            workflow = workflow_kls(
                self.orchestrator_config,
                workflow_args_map[workflow_kls],
                root_dir=self.directory,
            )

            if not self.check_filters(workflow):
                continue

            workflow.registry.collect_classes(workflow.module_directory)
            workflows.append(workflow)

        for idx, workflow in enumerate(workflows):
            self._display_workflow(idx, workflow, show_tasks=False)

    def _handle_initialize(self):
        self.logger.debug("Initializing a single workflow")
        workflow_name = self.orchestrator_config.workflow

        if not workflow_name:
            raise MisconfigurationError(
                "--initialize requires a single workflow - pass --workflow <name>."
            )

        if workflow_name not in self.workflow_registry:
            raise MisconfigurationError(
                f"{workflow_name} not found in the workflow registry"
                f" - the available workflows are {self.workflow_registry.names}."
            )

        workflow_kls = self.workflow_registry[workflow_name]
        workflow_args = self._parse_workflow_args(workflow_kls)

        if not workflow_kls.should_be_initialized(self.orchestrator_config):
            raise InitializationError(
                f"Cannot initialize workflow {workflow_kls} - its should_be_initialized answered False for this configuration."
            )

        workflow = workflow_kls(
            self.orchestrator_config, workflow_args, root_dir=self.directory
        )
        workflow.initialize_tasks()
        self._display_workflow(0, workflow, show_tasks=True)

    def _parse_workflow_args(self, workflow_kls):
        """Parse strictly for the one workflow that is selected to run. `required`
        and depends_on are thus enforced here. The lenient collective parse, which
        only lists the workflows, does not enforce them."""
        parser = workflow_kls.get_parser()
        if parser is None:
            if self.unknown_args:
                raise MisconfigurationError(
                    f"Unrecognized arguments: {' '.join(self.unknown_args)}"
                )
            return None
        args = parser.parse_args(self.unknown_args)
        workflow_kls.validate_option_dependencies(args)
        return args

    def _handle_run(self):
        if self.orchestrator_config.mode is Mode.HEADLESS:
            return self._handle_headless_run()
        self._handle_interactive_run()

    def _handle_headless_run(self):
        self.logger.debug("Headless run")

        workflow_name = self.orchestrator_config.workflow

        if not workflow_name:
            raise MisconfigurationError(
                f"a headless run initializes one workflow - pass --workflow "
                f"<name>; the collected workflows are "
                f"{self.workflow_registry.names}."
            )

        if workflow_name not in self.workflow_registry:
            raise MisconfigurationError(
                f"{workflow_name} not found in the workflow registry"
                f" - the available workflows are {self.workflow_registry.names}."
            )

        workflow_kls = self.workflow_registry[workflow_name]
        workflow_args = self._parse_workflow_args(workflow_kls)

        if not workflow_kls.should_be_initialized(self.orchestrator_config):
            error = InitializationError(
                f"Cannot initialize workflow {workflow_kls} - its should_be_initialized answered False for this configuration."
            )
            emit_unscoped_error(
                error, workflow_name=workflow_name, workflow_class=workflow_kls.__name__
            )
            raise error

        workflow = workflow_kls(
            self.orchestrator_config, workflow_args, root_dir=self.directory
        )

        # A run that does not start is invisible under cron, so the telemetry
        # hook must report it. The catch is narrow on purpose:
        # MisconfigurationError is bad input, and a reraise_errors escape is
        # already reported at the task boundary.
        try:
            workflow.initialize_tasks()
            return workflow.headless_run()
        except (EligibilityError, InitializationError) as e:
            emit_unscoped_error(
                e,
                workflow_name=workflow.instance_name,
                workflow_instance=str(workflow),
                workflow_class=type(workflow).__name__,
                session_id=workflow.session_id,
            )
            raise

    def _handle_serve(self):
        try:
            import uvicorn
            from winslow.serve.app import create_app
        except ImportError as e:
            raise MisconfigurationError(
                "Serve mode requires the serve extra - install with: "
                "pip install 'winslow[serve]'"
            ) from e
        from winslow.serve.auth import Credentials
        from winslow.session import SessionRegistry
        from winslow.state import create_state_store

        config = self.orchestrator_config
        self.logger.info(f"Serving on {config.host}:{config.port}")
        # The boundary must exist before the first session logs (see
        # setup_run_logging). WINSLOW_LOG_JSON sends the run lane to stdout
        # for the log store of a pod; the default keeps the session files.
        setup_run_logging(sinks=[stdout_json_sink()] if settings.LOG_JSON else None)
        registry = SessionRegistry()
        state_store = create_state_store(config)
        self._restore_sessions(registry, state_store)
        self._auto_init_sessions(registry, state_store)
        app = create_app(
            registry,
            Credentials.from_env(config.host),
            orchestrator=self,
            state_store=state_store,
            endpoints=config.endpoints,
            base_url=f"http://{config.host}:{config.port}",
        )
        try:
            # log_config=None: uvicorn's loggers propagate to the root
            # handler, so its lines share the winslow format - JSON under
            # WINSLOW_LOG_JSON (see _console_handler).
            uvicorn.run(app, host=config.host, port=config.port, log_config=None)
        finally:
            shutdown_run_logging()

    def _restore_sessions(self, registry, state_store):
        """Rebuild every open manifest at serve startup: the sessions of a
        dead process come back without a client, the way a local user's
        restore brings them back. A failed rebuild logs and skips."""
        from winslow.client import LocalAppClient

        port = LocalAppClient(registry, orchestrator=self, state_store=state_store)
        for manifest in port.manifests():
            self.logger.info(f"restore: rebuilding {manifest.session_id}")
            try:
                port.restore_session(manifest.session_id)
            except Exception:
                self.logger.error(
                    f"restore: the rebuild of '{manifest.session_id}' failed.",
                    exc_info=True,
                )

    def _auto_init_sessions(self, registry, state_store):
        """One session per auto_init workflow: the process that owns the
        sessions runs auto_init, so a connecting client starts none. A
        failed initialization logs and skips; the process serves on."""
        from winslow.session import create_session

        # A restored session satisfies auto_init: the workflow already runs.
        live = {session.workflow.instance_name for session in registry.sessions()}
        for name in self.workflow_registry.names:
            workflow_kls = self.workflow_registry[name]
            if not workflow_kls.auto_init or name in live:
                continue
            if not workflow_kls.should_be_initialized(self.orchestrator_config):
                continue
            self.logger.info(f"auto_init: initializing {name}")
            try:
                create_session(self, state_store, registry, name)
            except Exception:
                self.logger.error(
                    f"auto_init: the initialization of '{name}' failed.",
                    exc_info=True,
                )

    def _handle_connect(self):
        """The remote TUI: the same app over the wire transport of the
        session port. The serve process owns the workflows and the state."""
        try:
            from winslow.client.websocket import RemoteAppClient
            from winslow.ui import Winslow
        except ImportError as e:
            raise MisconfigurationError(
                "Connect mode requires the connect extra - install with: "
                "pip install 'winslow[connect]'"
            ) from e

        config = self.orchestrator_config
        client = RemoteAppClient(config.url, token=settings.SERVE_TOKEN)
        client.connect()
        self.logger.info(f"Connected to {config.url}")

        # The tasks run on the serve process, so no run record reaches this
        # one: run-log wiring here would only create an empty log directory.
        app = Winslow(client=client, logger=self.logger, owns_sessions=False)
        try:
            app.run()
        finally:
            client.close()

    def _handle_interactive_run(self):
        self.logger.debug("Interactive run")

        try:
            from winslow.ui import Winslow
        except ImportError as e:
            raise MisconfigurationError(
                "Interactive mode requires the UI extra - install with: pip install 'winslow[tui]'"
            ) from e

        # Validate the CLI args before any effect, so a typo fails immediately.
        # The descriptors read of the app parses them again per request.
        self.collect_workflow_args()

        # Set up the winslow.runs sink and the propagate=False boundary BEFORE a
        # workflow logger or a task logger starts to propagate. The run logs then
        # go to the sink and not to the console. This is interactive only. A
        # headless run keeps the console output.
        setup_run_logging()

        from winslow.client import LocalAppClient
        from winslow.session import SessionRegistry
        from winslow.state import create_state_store

        # The composition root of the local TUI: this process owns the registry
        # and the durable store (see winslow.state); the app consumes the port.
        local_client = LocalAppClient(
            SessionRegistry(),
            orchestrator=self,
            state_store=create_state_store(self.orchestrator_config),
        )
        app = Winslow(client=local_client, logger=self.logger, owns_sessions=True)

        # app.run() blocks until the TUI stops. Then flush and stop the
        # run-logging listener. This is in a finally clause, so it also occurs
        # after an error. If it does not occur, the listener thread and the open
        # file handles stay, and the buffered lines are lost.
        try:
            app.run()
        finally:
            shutdown_run_logging()

    def initialize_workflow(
        self,
        workflow_kls,
        orchestrator_overrides,
        workflow_values,
        workflow_base=None,
        task_store=None,
        logger=LOGGER,
    ):
        # Put the UI values on the parsed base and do not build a new config from
        # them. The base holds the default of each declared option, and also of
        # an option with show_on_ui=False. The form path thus cannot drop an
        # option that it did not show.
        orchestrator_config = _merged_config(
            self.orchestrator_config, orchestrator_overrides
        )
        workflow_params = _merged_config(workflow_base, workflow_values)

        return workflow_kls(
            orchestrator_config,
            workflow_params,
            store=task_store,
            logger=logger,
            root_dir=self.directory,
        )

    def _collect_caches(self):
        """Collect the GlobalCache classes for the workflow initializations. A
        WorkflowCache outside every workflow directory only gets a warning."""
        self.global_cache_registry.collect_classes(self.directory)
        set_global_cache_registry(self.global_cache_registry)

        workflow_directories = [
            os.path.dirname(os.path.abspath(sys.modules[kls.__module__].__file__))
            for kls in self.workflow_registry.classes
        ]
        for kls in stray_workflow_caches(self.directory, workflow_directories):
            LOGGER.warning(
                f"WorkflowCache {kls.__name__} is outside every workflow "
                f"directory - no workflow collects it."
            )

    def start(self):
        """
        Do the action that self.orchestrator_config (base_args) and unknown_args
        select:

        - Start the UI.
        - Start a simple run.
        - Show a summary of the workflows and the tasks.
        """
        self.validate_option_dependencies(
            self.orchestrator_config, subcommand=self.orchestrator_config.action.value
        )

        if self.orchestrator_config.action is Command.CONNECT:
            # A remote TUI reads everything over the wire, so the local
            # workflow and cache collection is skipped.
            return self._handle_connect()

        self.workflow_registry.collect_classes(self.directory)
        self._collect_caches()

        if self.orchestrator_config.action is Command.SHOW:
            self._handle_show()
        elif self.orchestrator_config.action is Command.SERVE:
            self._handle_serve()
        elif self.orchestrator_config.action is Command.RUN:
            # Runs only: a show produces no errors worth a backend. The finally
            # flushes and unregisters, so an embedding process can start again.
            self.telemetry_registry.collect_classes(self.directory)
            active_telemetry = activate_telemetry_configurations(
                self.telemetry_registry.classes, self.orchestrator_config
            )
            try:
                return self._handle_run()
            finally:
                shutdown_telemetry_configurations(active_telemetry)

collect_workflow_args()

{workflow_kls: parsed args or None}, validated together.

show, the interactive start, and a serve process's descriptors request call this. They parse the args of each discovered workflow at one time, so a start form (local or remote) prefills from the same CLI-supplied values. A headless run has one named workflow and parses its args strictly (see _handle_headless_run), so the collision between workflows that this method handles cannot occur there.

Source code in src/winslow/orchestrator.py
def collect_workflow_args(self):
    """{workflow_kls: parsed args or None}, validated together.

    `show`, the interactive start, and a serve process's descriptors
    request call this. They parse the args of each discovered workflow
    at one time, so a start form (local or remote) prefills from the
    same CLI-supplied values. A headless run has one named workflow and
    parses its args strictly (see `_handle_headless_run`), so the
    collision between workflows that this method handles cannot occur
    there."""
    parsed = {
        kls: self._parse_known_workflow_args(kls, self.unknown_args)
        for kls in self.sorted_workflow_classes
    }

    # A token is unrecognized if no workflow claimed its position. A position
    # that one workflow claimed leaves the intersection.
    unclaimed = set(range(len(self.unknown_args)))
    for _, leftover in parsed.values():
        unclaimed &= self._leftover_indices(self.unknown_args, leftover)

    if unclaimed:
        tokens = [self.unknown_args[i] for i in sorted(unclaimed)]
        raise MisconfigurationError(f"Unrecognized arguments: {' '.join(tokens)}")

    return {kls: args for kls, (args, _) in parsed.items()}

start()

Do the action that self.orchestrator_config (base_args) and unknown_args select:

  • Start the UI.
  • Start a simple run.
  • Show a summary of the workflows and the tasks.
Source code in src/winslow/orchestrator.py
def start(self):
    """
    Do the action that self.orchestrator_config (base_args) and unknown_args
    select:

    - Start the UI.
    - Start a simple run.
    - Show a summary of the workflows and the tasks.
    """
    self.validate_option_dependencies(
        self.orchestrator_config, subcommand=self.orchestrator_config.action.value
    )

    if self.orchestrator_config.action is Command.CONNECT:
        # A remote TUI reads everything over the wire, so the local
        # workflow and cache collection is skipped.
        return self._handle_connect()

    self.workflow_registry.collect_classes(self.directory)
    self._collect_caches()

    if self.orchestrator_config.action is Command.SHOW:
        self._handle_show()
    elif self.orchestrator_config.action is Command.SERVE:
        self._handle_serve()
    elif self.orchestrator_config.action is Command.RUN:
        # Runs only: a show produces no errors worth a backend. The finally
        # flushes and unregisters, so an embedding process can start again.
        self.telemetry_registry.collect_classes(self.directory)
        active_telemetry = activate_telemetry_configurations(
            self.telemetry_registry.classes, self.orchestrator_config
        )
        try:
            return self._handle_run()
        finally:
            shutdown_telemetry_configurations(active_telemetry)

Client API

winslow.client.AppClient

Bases: Port

The dashboard scope: the live sessions, descriptors, open manifests, create and restore. It hands out SessionClients (see session).

Source code in src/winslow/client/base.py
class AppClient(Port):
    """The dashboard scope: the live sessions, descriptors, open manifests,
    create and restore. It hands out SessionClients (see session)."""

    scope = "app"

    sessions = Read(
        result=SessionInfo,
        many=True,
        doc="Return a SessionInfo per live session.",
    )
    descriptors = Read(
        result=Descriptors,
        doc="Return the Descriptors of the process: the collected workflows "
        "with their form options, plus the orchestrator overrides.",
    )
    manifests = Read(
        result=ManifestInfo,
        many=True,
        doc="Return a ManifestInfo per restorable session: the open manifests "
        "that name no live session.",
    )
    create_session = Read(
        workflow=str,
        overrides=(dict | None, None),
        values=(dict | None, None),
        result=SessionInfo,
        traceback=True,
        doc="Build, initialize, persist and register one session. Return its "
        "SessionInfo. overrides and values default to {} at the handler, so "
        "None and an absent field behave the same.",
    )
    restore_session = Read(
        session_id=str,
        result=SessionInfo,
        traceback=True,
        doc="Re-create the session of an open manifest under its stored id, "
        "seeded from state. Return its SessionInfo.",
    )

    def session(self, session_id):
        """Return the SessionClient of one live session. The local transport
        builds a fresh client per call. The wire transport returns the one
        shared lane of the session."""
        raise NotImplementedError

    def subscribe_connection(self, handler):
        """Connect the handler to the ConnectionEvent lane (see
        winslow.model). The wire transport emits on a drop and on the
        reconnect. The in-process transport keeps the handler idle. A
        handler can run on any thread."""
        raise NotImplementedError

session(session_id)

Return the SessionClient of one live session. The local transport builds a fresh client per call. The wire transport returns the one shared lane of the session.

Source code in src/winslow/client/base.py
def session(self, session_id):
    """Return the SessionClient of one live session. The local transport
    builds a fresh client per call. The wire transport returns the one
    shared lane of the session."""
    raise NotImplementedError

subscribe_connection(handler)

Connect the handler to the ConnectionEvent lane (see winslow.model). The wire transport emits on a drop and on the reconnect. The in-process transport keeps the handler idle. A handler can run on any thread.

Source code in src/winslow/client/base.py
def subscribe_connection(self, handler):
    """Connect the handler to the ConnectionEvent lane (see
    winslow.model). The wire transport emits on a drop and on the
    reconnect. The in-process transport keeps the handler idle. A
    handler can run on any thread."""
    raise NotImplementedError

winslow.client.SessionClient

Bases: Port

One session: the reads, the subscriptions, and the actions.

Source code in src/winslow/client/base.py
class SessionClient(Port):
    """One session: the reads, the subscriptions, and the actions."""

    scope = "session"

    # --- reads ---------------------------------------------------------------

    snapshot = Read(
        result=SessionSnapshot,
        doc="Return the SessionSnapshot: statuses, batch rows, session log "
        "backlog, session meta.",
    )
    roster = Read(
        result=TaskInfo,
        many=True,
        doc="Return a stub TaskInfo per task, in the order of the launch filter.",
    )
    task_detail = Read(
        task_key=str,
        result=TaskInfo,
        doc="Return the full TaskInfo of one task, evaluated and with the "
        "trust fields filled (see Workflow.task_info).",
    )
    record_detail = Read(
        batch_uuid=str,
        task_key=str,
        result=RecordDetail,
        doc="Return the RecordDetail of one execution record.",
    )
    history = Read(
        result=BatchOutcome,
        many=True,
        doc="Return a BatchOutcome per batch, with the per-task outcomes.",
    )
    log_tail = Read(
        batch_uuid=str,
        task_key=str,
        limit=(int | None, None),
        result=list,
        doc="Return the last `limit` stored log lines of one record. None "
        "reads the default window of the record store.",
    )
    caches = Read(
        result=CacheInfo,
        many=True,
        doc="Return a CacheInfo per cache of the session.",
    )
    cache_value = Read(
        cache_name=str,
        entry_name=str,
        result=CacheValueView,
        doc="Return the CacheValueView of one entry: the rendered value, its "
        "encoding, and the error context.",
    )
    apply_filter = Read(
        query=str,
        scope=(str, "tasks"),
        result=tuple,
        doc="Return the identity keys the query matches over the named corpus, "
        "'tasks' or 'history' (see Workflow.filter_keys). A bad query raises "
        "RequestError with the parse error.",
    )
    batch_options = Read(
        result=dict,
        doc="Return the baseline batch option values of the session as a dict. "
        "A fresh client prefills its toggles from it. Each submit then carries "
        "the values of the client (see RunTasks.options).",
    )
    session_params = Read(
        result=SessionParams,
        doc="Return the SessionParams: settings_snapshot plus the resolved "
        "workflow_config values.",
    )

    # --- subscriptions ---------------------------------------------------------

    def subscribe(self, topic, handler):
        """Connect the handler to one event topic of the session. The topics
        are the session bus event classes (see winslow.events) plus
        CacheUpdatedEvent and SessionLogEvent (see winslow.model). A handler
        can run on any thread. The dispatch order between handlers is
        undefined."""
        raise NotImplementedError

    def unsubscribe(self, topic, handler):
        """Disconnect the handler (see subscribe). An unknown pair is a
        no-op, so a teardown path can run twice."""
        raise NotImplementedError

    def subscribe_task_log(self, task_key, handler):
        """Connect the handler to the live log stream of one task, outside
        any batch, and return the buffered backlog lines. The handler
        receives TaskLogEvent values."""
        raise NotImplementedError

    def unsubscribe_task_log(self, task_key, handler):
        """Disconnect the task log handler (see subscribe_task_log)."""
        raise NotImplementedError

    # --- actions ----------------------------------------------------------------

    def submit(self, action):
        """Submit one action dataclass (see winslow.actions) and return its
        ack. A refused action answers an ack that carries the reason."""
        raise NotImplementedError

subscribe(topic, handler)

Connect the handler to one event topic of the session. The topics are the session bus event classes (see winslow.events) plus CacheUpdatedEvent and SessionLogEvent (see winslow.model). A handler can run on any thread. The dispatch order between handlers is undefined.

Source code in src/winslow/client/base.py
def subscribe(self, topic, handler):
    """Connect the handler to one event topic of the session. The topics
    are the session bus event classes (see winslow.events) plus
    CacheUpdatedEvent and SessionLogEvent (see winslow.model). A handler
    can run on any thread. The dispatch order between handlers is
    undefined."""
    raise NotImplementedError

unsubscribe(topic, handler)

Disconnect the handler (see subscribe). An unknown pair is a no-op, so a teardown path can run twice.

Source code in src/winslow/client/base.py
def unsubscribe(self, topic, handler):
    """Disconnect the handler (see subscribe). An unknown pair is a
    no-op, so a teardown path can run twice."""
    raise NotImplementedError

subscribe_task_log(task_key, handler)

Connect the handler to the live log stream of one task, outside any batch, and return the buffered backlog lines. The handler receives TaskLogEvent values.

Source code in src/winslow/client/base.py
def subscribe_task_log(self, task_key, handler):
    """Connect the handler to the live log stream of one task, outside
    any batch, and return the buffered backlog lines. The handler
    receives TaskLogEvent values."""
    raise NotImplementedError

unsubscribe_task_log(task_key, handler)

Disconnect the task log handler (see subscribe_task_log).

Source code in src/winslow/client/base.py
def unsubscribe_task_log(self, task_key, handler):
    """Disconnect the task log handler (see subscribe_task_log)."""
    raise NotImplementedError

submit(action)

Submit one action dataclass (see winslow.actions) and return its ack. A refused action answers an ack that carries the reason.

Source code in src/winslow/client/base.py
def submit(self, action):
    """Submit one action dataclass (see winslow.actions) and return its
    ack. A refused action answers an ack that carries the reason."""
    raise NotImplementedError

winslow.client.LocalAppClient

Bases: AppClient

The dashboard scope over a live SessionRegistry. orchestrator and state_store serve descriptors, manifests, create and restore. A client without them serves the registry reads only.

Source code in src/winslow/client/local.py
class LocalAppClient(AppClient):
    """The dashboard scope over a live SessionRegistry. orchestrator and
    state_store serve descriptors, manifests, create and restore. A client
    without them serves the registry reads only."""

    def __init__(self, registry, orchestrator, state_store):
        self.registry = registry
        self.orchestrator = orchestrator
        self.state_store = state_store

    def sessions(self):
        return tuple(
            SessionInfo.from_session(session) for session in self.registry.sessions()
        )

    def descriptors(self):
        return Descriptors.from_orchestrator(self.orchestrator)

    def manifests(self):
        return tuple(
            ManifestInfo.from_manifest(manifest)
            for manifest in self.state_store.list_open_manifests()
            if manifest.session_id not in self.registry
        )

    def create_session(self, workflow, overrides=None, values=None):
        try:
            session = create_session(
                self.orchestrator,
                self.state_store,
                self.registry,
                workflow,
                overrides,
                values,
                origin="local",
            )
        except (KeyError, ValueError) as exc:
            # A refusal: an unknown workflow, a bad option value. An init
            # bug propagates unchanged, so an in-process caller keeps its traceback.
            raise RequestError(exc.args[0] if exc.args else str(exc)) from None
        return SessionInfo.from_session(session)

    def restore_session(self, session_id):
        if session_id in self.registry:
            raise RequestError(f"{session_id!r} is already a live session.")
        manifest = next(
            (
                m
                for m in self.state_store.list_open_manifests()
                if m.session_id == session_id
            ),
            None,
        )
        if manifest is None:
            raise RequestError(f"{session_id!r} names no open manifest to restore.")
        if manifest.workflow_class not in self.orchestrator.workflow_registry.names:
            raise RequestError(
                f"the manifest names workflow {manifest.workflow_class!r}, "
                f"which this process does not collect."
            )
        session = create_session(
            self.orchestrator,
            self.state_store,
            self.registry,
            manifest.workflow_class,
            manifest.orchestrator_overrides or {},
            manifest.workflow_values or {},
            session_id=manifest.session_id,
            seed=True,
            origin="local",
        )
        return SessionInfo.from_session(session)

    def session(self, session_id):
        return LocalSessionClient(self.registry.resolve(session_id))

    def subscribe_connection(self, handler):
        # The in-process transport has no connection to lose.
        pass

winslow.exceptions.RequestError

Bases: WinslowError

A session port read the server or the session refused: an unknown key, an ended session, a bad query. Both transports raise it, so one catch site serves both modes (see winslow.client.base). detail can carry a server traceback for the error modal.

Source code in src/winslow/exceptions.py
class RequestError(WinslowError):
    """A session port read the server or the session refused: an unknown key,
    an ended session, a bad query. Both transports raise it, so one catch site
    serves both modes (see winslow.client.base). detail can carry a server
    traceback for the error modal."""

    def __init__(self, reason, detail=None):
        super().__init__(reason)
        self.detail = detail

Actions

winslow.actions

The inbound action path of a session. A presentation layer translates user input into one action dataclass and submits it to the ActionHandler of the session. Action fields are values only: identity keys, scalars (the payload rule, see winslow.events).

Ack dataclass

The synchronous answer to an action: accepted or refused. The reason names what refused the action.

Source code in src/winslow/actions.py
@dataclass(frozen=True)
class Ack:
    """The synchronous answer to an action: accepted or refused. The reason
    names what refused the action."""

    accepted: bool
    reason: str | None = None

BatchAck dataclass

Bases: Ack

The answer to a batch submit. An accepted submit carries the uuid of the created batch; the batch worker threads do the work.

Source code in src/winslow/actions.py
@dataclass(frozen=True)
class BatchAck(Ack):
    """The answer to a batch submit. An accepted submit carries the uuid of
    the created batch; the batch worker threads do the work."""

    batch_uuid: str | None = None

Action

The base of every action. name is the wire name (see ActionFrame). The actions register by name at class creation (see build).

Source code in src/winslow/actions.py
class Action:
    """The base of every action. name is the wire name (see ActionFrame). The
    actions register by name at class creation (see build)."""

    name = None
    by_name = {}

    def __init_subclass__(cls, **kwargs):
        super().__init_subclass__(**kwargs)
        # A base with no name is an action kind and stays unregistered (see Lane).
        if cls.name is not None:
            Action.by_name[cls.name] = cls

    @classmethod
    def build(cls, name, payload):
        """The action instance of one wire name and its field dict (see
        ActionFrame). The codec validates the payload against the dataclass."""
        # The codec needs pydantic, which only the serve and connect extras install.
        from winslow.protocol.codec import CODEC, ValidationError, report

        action_class = cls.by_name.get(name)
        if action_class is None:
            raise ValueError(
                f"{name!r} names no action. The actions are {sorted(cls.by_name)}."
            )
        try:
            return CODEC.decode(action_class, payload or {})
        except ValidationError as exc:
            raise ValueError(
                f"bad fields for {name} - {report(exc)}. The fields of {name} are "
                f"{[f.name for f in fields(action_class)]}."
            ) from None

build(name, payload) classmethod

The action instance of one wire name and its field dict (see ActionFrame). The codec validates the payload against the dataclass.

Source code in src/winslow/actions.py
@classmethod
def build(cls, name, payload):
    """The action instance of one wire name and its field dict (see
    ActionFrame). The codec validates the payload against the dataclass."""
    # The codec needs pydantic, which only the serve and connect extras install.
    from winslow.protocol.codec import CODEC, ValidationError, report

    action_class = cls.by_name.get(name)
    if action_class is None:
        raise ValueError(
            f"{name!r} names no action. The actions are {sorted(cls.by_name)}."
        )
    try:
        return CODEC.decode(action_class, payload or {})
    except ValidationError as exc:
        raise ValueError(
            f"bad fields for {name} - {report(exc)}. The fields of {name} are "
            f"{[f.name for f in fields(action_class)]}."
        ) from None

RunTasks dataclass

Bases: Action

options carries the batch options of this submit: {name: bool}, over the session baseline. They are client view state, so each client sends its own with every batch (see BatchOptions).

Source code in src/winslow/actions.py
@dataclass(frozen=True)
class RunTasks(Action):
    """options carries the batch options of this submit: {name: bool},
    over the session baseline. They are client view state, so each client
    sends its own with every batch (see BatchOptions)."""

    name = "run_tasks"
    keys: tuple
    options: dict | None = None

CheckTasks dataclass

Bases: Action

The check form of RunTasks; options works the same.

Source code in src/winslow/actions.py
@dataclass(frozen=True)
class CheckTasks(Action):
    """The check form of RunTasks; options works the same."""

    name = "check_tasks"
    keys: tuple
    options: dict | None = None

LoadCacheEntries dataclass

Bases: Action

Bulk-only, like RunTasks: a single selection sends a one-pair list. Each pair is (cache_name, entry_name).

Source code in src/winslow/actions.py
@dataclass(frozen=True)
class LoadCacheEntries(Action):
    """Bulk-only, like RunTasks: a single selection sends a one-pair list.
    Each pair is (cache_name, entry_name)."""

    name = "load_cache_entries"
    entries: tuple

ClearCacheEntries dataclass

Bases: Action

Bulk-only, like LoadCacheEntries. "Clear" is the action verb on the wire; the handler calls cache.invalidate internally.

Source code in src/winslow/actions.py
@dataclass(frozen=True)
class ClearCacheEntries(Action):
    """Bulk-only, like LoadCacheEntries. "Clear" is the action verb on the
    wire; the handler calls cache.invalidate internally."""

    name = "clear_cache_entries"
    entries: tuple

ActionHandler

One per session: the inbound half of the session boundary (the bus is the outbound half, see SessionBus). The handler accepts one action, resolves its values to live objects, gates it, delegates it to the runner or the session, and answers with an ack. It refuses with an ack, never with an exception: a transport forwards the reason as-is.

Source code in src/winslow/actions.py
class ActionHandler:
    """One per session: the inbound half of the session boundary (the bus is
    the outbound half, see SessionBus). The handler accepts one action,
    resolves its values to live objects, gates it, delegates it to the runner
    or the session, and answers with an ack. It refuses with an ack, never
    with an exception: a transport forwards the reason as-is."""

    def __init__(self, session):
        self.session = session

    @property
    def _workflow(self):
        return self.session.workflow

    @property
    def _runner(self):
        return self.session.workflow.runner

    def submit(self, action):
        """The one entry point: dispatch on the action class."""
        method = self._methods.get(type(action))
        if method is None:
            return self._refuse(
                action,
                f"{type(action).__name__} names no action of this session. "
                f"The actions are {sorted(k.__name__ for k in self._methods)}.",
            )
        if self.session.has_ended:
            return self._refuse(
                action, f"{self.session.session_id} has ended and accepts no action."
            )
        return method(self, action)

    def submit_guarded(self, action):
        """submit for a wire transport: an unexpected raise becomes a refused
        ack with the traceback in the session log, so no exception crosses the
        wire boundary. The TUI calls submit and keeps the real traceback."""
        try:
            return self.submit(action)
        except Exception:
            self._workflow.logger.error(
                f"{type(action).__name__} failed inside the session.", exc_info=True
            )
            return self._refuse(
                action,
                f"{type(action).__name__} failed inside the session - "
                f"the session log has the traceback.",
            )

    @classmethod
    def _refuse(cls, action, reason):
        ack_class = BatchAck if isinstance(action, (RunTasks, CheckTasks)) else Ack
        return ack_class(accepted=False, reason=reason)

    @handles(RunTasks)
    def run_tasks(self, action):
        return self._submit_batch(
            action, self._runner.submit_run_single, self._runner.submit_run
        )

    @handles(CheckTasks)
    def check_tasks(self, action):
        return self._submit_batch(
            action, self._runner.submit_check_single, self._runner.submit_check
        )

    def _batch_options_for(self, action):
        """The BatchOptions of one submit: the session baseline with the
        action's values on top, or (None, reason) on an unknown option."""
        overrides = action.options or {}
        known = {field.name for field in fields(BatchOptions)}
        unknown = sorted(set(overrides) - known)
        if unknown:
            return None, (
                f"{', '.join(repr(name) for name in unknown)} names no batch "
                f"option - the options are {sorted(known)}."
            )
        baseline = asdict(self._workflow.batch_options)
        return BatchOptions(**{**baseline, **overrides}), None

    def _submit_batch(self, action, submit_single, submit_bulk):
        keys = tuple(dict.fromkeys(action.keys))
        if len(keys) != len(action.keys):
            # A wire client can repeat a key; a batch runs each task once.
            dupes = sorted({key for key in keys if action.keys.count(key) > 1})
            self._workflow.logger.warning(
                f"{type(action).__name__} repeats {dupes}; each task enters the batch once."
            )
        options, reason = self._batch_options_for(action)
        if reason is not None:
            return self._refuse(action, reason)
        try:
            tasks = [self._workflow.task_index.resolve(key) for key in keys]
        except KeyError as exc:
            return self._refuse(action, exc.args[0])
        try:
            if len(tasks) == 1:
                batch = submit_single(tasks[0], options=options)
            else:
                batch = submit_bulk(tasks, options=options)
        except SessionEndingError as exc:
            return self._refuse(action, str(exc))
        if batch is None:
            # The admission filtered every task out (see _open_batch).
            return self._refuse(
                action,
                (
                    "no requested task is eligible for this action right now - "
                    "the current statuses exclude all of them; the task list "
                    "shows each status."
                ),
            )
        return BatchAck(accepted=True, batch_uuid=batch.uuid)

    @handles(StopBatch)
    def stop_batch(self, action):
        batch = self._runner.get_batch(action.batch_uuid)
        if batch is None:
            return self._refuse(
                action, f"{action.batch_uuid} names no batch of this session."
            )
        # The acceptance means "stop requested". The batch drains on its own
        # worker threads (see ExecutionBatch.request_stop).
        batch.request_stop()
        return Ack(accepted=True)

    @handles(EndSession)
    def end_session(self, action):
        if action.force:
            self.session.force_end()
        else:
            self.session.end()
        return Ack(accepted=True)

    def _resolve_cache_entries(self, action):
        """(cache, entry_name) pairs for the wire pairs of the action, or a
        refusal reason naming the first unknown cache or entry."""
        caches_by_name = {cache.get_name(): cache for cache in self._workflow.caches()}
        resolved = []
        for cache_name, entry_name in action.entries:
            cache = caches_by_name.get(cache_name)
            if cache is None:
                return None, f"{cache_name!r} names no cache of this session."
            if entry_name not in declared_entries(type(cache)):
                return None, f"{cache} has no entry {entry_name!r}."
            resolved.append((cache, entry_name))
        return resolved, None

    def _load_cache_entry(self, cache, entry_name):
        # A loader failure is data, not an action failure: the entry reports
        # ERRORED and the session log carries the traceback (see
        # BaseCache._entry_value). The ack still accepts.
        try:
            getattr(cache, entry_name)
        except Exception:
            self._workflow.logger.error(
                f"Cache '{cache.get_name()}': the load of '{entry_name}' failed.",
                exc_info=True,
            )

    def _clear_cache_entry(self, cache, entry_name):
        cache.invalidate(entry_name)

    def _run_cache_entries(self, work, resolved):
        # A bare thread carries no LogContext. The scope routes the loader
        # and invalidation lines to the session log (see Session.log_scope).
        with self.session.log_scope():
            execute_in_threads(work, resolved)

    def _cache_entries_action(self, action, work):
        if not action.entries:
            return self._refuse(action, "the entries list is empty - nothing to do.")
        resolved, reason = self._resolve_cache_entries(action)
        if reason is not None:
            return self._refuse(action, reason)
        # The ack means "started", like RunTasks: execute_in_threads blocks
        # until every entry finishes, so a background thread runs it and
        # the caller does not wait. cache_updated events report progress.
        threading.Thread(
            target=self._run_cache_entries, args=(work, resolved), daemon=True
        ).start()
        return Ack(accepted=True)

    @handles(LoadCacheEntries)
    def load_cache_entries(self, action):
        return self._cache_entries_action(action, self._load_cache_entry)

    @handles(ClearCacheEntries)
    def clear_cache_entries(self, action):
        return self._cache_entries_action(action, self._clear_cache_entry)

submit(action)

The one entry point: dispatch on the action class.

Source code in src/winslow/actions.py
def submit(self, action):
    """The one entry point: dispatch on the action class."""
    method = self._methods.get(type(action))
    if method is None:
        return self._refuse(
            action,
            f"{type(action).__name__} names no action of this session. "
            f"The actions are {sorted(k.__name__ for k in self._methods)}.",
        )
    if self.session.has_ended:
        return self._refuse(
            action, f"{self.session.session_id} has ended and accepts no action."
        )
    return method(self, action)

submit_guarded(action)

submit for a wire transport: an unexpected raise becomes a refused ack with the traceback in the session log, so no exception crosses the wire boundary. The TUI calls submit and keeps the real traceback.

Source code in src/winslow/actions.py
def submit_guarded(self, action):
    """submit for a wire transport: an unexpected raise becomes a refused
    ack with the traceback in the session log, so no exception crosses the
    wire boundary. The TUI calls submit and keeps the real traceback."""
    try:
        return self.submit(action)
    except Exception:
        self._workflow.logger.error(
            f"{type(action).__name__} failed inside the session.", exc_info=True
        )
        return self._refuse(
            action,
            f"{type(action).__name__} failed inside the session - "
            f"the session log has the traceback.",
        )

Events

winslow.bus.SessionBus

One bus per session. Every component that observes the session subscribes here, by event class, and the session-end sweep disconnects every remaining subscriber (see close).

publish dispatches synchronously on the calling thread, outside the store lock (see ReactiveDict.set). Dispatch order between events is undefined: a callback that renders state reads the store for the latest view. A callback must return immediately, because it runs on the thread that produced the event.

Example scenario, the workflow screen of the TUI (see WorkflowScreen):

1. A worker thread completes a task and writes
   store[task] = COMPLETED, then publishes with the lock released.
2. The TaskStatusEvent callback runs on that worker thread.
3. The screen handler posts the event to the UI thread and returns
   immediately. The worker continues.

A slow body in step 3, for example a synchronous render of the task table, stalls that worker at each write.

Source code in src/winslow/bus.py
class SessionBus:
    """One bus per session. Every component that observes the session
    subscribes here, by event class, and the session-end sweep disconnects
    every remaining subscriber (see close).

    publish dispatches synchronously on the calling thread, outside the store
    lock (see ReactiveDict.set). Dispatch order between events is undefined:
    a callback that renders state reads the store for the latest view.
    A callback must return immediately, because it runs on the thread that
    produced the event.

    Example scenario, the workflow screen of the TUI (see WorkflowScreen):

        1. A worker thread completes a task and writes
           store[task] = COMPLETED, then publishes with the lock released.
        2. The TaskStatusEvent callback runs on that worker thread.
        3. The screen handler posts the event to the UI thread and returns
           immediately. The worker continues.

    A slow body in step 3, for example a synchronous render of the task
    table, stalls that worker at each write."""

    # The event vocabulary of the bus. A subclass overrides this to extend it.
    event_classes = (
        TaskStatusEvent,
        ExecutionStatusEvent,
        BatchCreatedEvent,
        BatchCompletedEvent,
        LogLineEvent,
        SessionEndedEvent,
    )

    @classmethod
    def get_event_classes(cls):
        """The declared events of this bus (see event_classes)."""
        return cls.event_classes

    def __init__(self):
        # The lock guards the subscription table. publish runs outside the
        # lock: blinker iterates a snapshot, so a subscriber can unsubscribe
        # from another thread during a dispatch.
        self._lock = threading.Lock()
        self._signals = {
            event_class: Signal() for event_class in self.get_event_classes()
        }
        self._receivers = {}
        self._closed = False

    def _refuse_undeclared(self, event_class):
        raise RegistrationError(
            f"{event_class.__name__} is not a declared event of this bus. "
            f"The declared events are "
            f"{sorted(k.__name__ for k in self.get_event_classes())}. "
            f"A subclass extends event_classes."
        )

    def subscribe(self, event_class, callback):
        """Connect the callback to the event class. The bus holds the callback
        strongly until unsubscribe or close, so a bound method stays alive."""
        receiver = _isolate(callback)
        with self._lock:
            if self._closed:
                raise RegistrationError(
                    f"The session bus is closed - it accepts no subscription "
                    f"to {event_class.__name__}. Subscribe before the session "
                    f"ends."
                )
            if event_class not in self._signals:
                self._refuse_undeclared(event_class)
            if (event_class, callback) in self._receivers:
                raise RegistrationError(
                    f"{callback!r} is already subscribed to "
                    f"{event_class.__name__}. Unsubscribe it before a second "
                    f"subscription."
                )
            self._receivers[(event_class, callback)] = receiver
            self._signals[event_class].connect(receiver, weak=False)

    def unsubscribe(self, event_class, callback):
        """Disconnect the callback. An unknown callback is a no-op, so a
        teardown path can run twice."""
        with self._lock:
            receiver = self._receivers.pop((event_class, callback), None)
            if receiver is not None:
                self._signals[event_class].disconnect(receiver)

    def publish(self, event):
        """Dispatch the event to its subscribers, on the calling thread. The
        dispatch order is undefined: a subscriber must not depend on another
        subscriber. A publish on a closed bus is a no-op: the session end can
        race a draining worker."""
        with self._lock:
            if self._closed:
                return
            signal = self._signals.get(type(event))
        if signal is None:
            self._refuse_undeclared(type(event))
        signal.send(self, event=event)

    def close(self):
        """Disconnect every remaining subscriber. The session end calls this
        after the SessionEndedEvent dispatch. This method is idempotent."""
        with self._lock:
            for (event_class, _), receiver in self._receivers.items():
                self._signals[event_class].disconnect(receiver)
            self._receivers.clear()
            # An empty table frees the signals with the session.
            self._signals.clear()
            self._closed = True

get_event_classes() classmethod

The declared events of this bus (see event_classes).

Source code in src/winslow/bus.py
@classmethod
def get_event_classes(cls):
    """The declared events of this bus (see event_classes)."""
    return cls.event_classes

subscribe(event_class, callback)

Connect the callback to the event class. The bus holds the callback strongly until unsubscribe or close, so a bound method stays alive.

Source code in src/winslow/bus.py
def subscribe(self, event_class, callback):
    """Connect the callback to the event class. The bus holds the callback
    strongly until unsubscribe or close, so a bound method stays alive."""
    receiver = _isolate(callback)
    with self._lock:
        if self._closed:
            raise RegistrationError(
                f"The session bus is closed - it accepts no subscription "
                f"to {event_class.__name__}. Subscribe before the session "
                f"ends."
            )
        if event_class not in self._signals:
            self._refuse_undeclared(event_class)
        if (event_class, callback) in self._receivers:
            raise RegistrationError(
                f"{callback!r} is already subscribed to "
                f"{event_class.__name__}. Unsubscribe it before a second "
                f"subscription."
            )
        self._receivers[(event_class, callback)] = receiver
        self._signals[event_class].connect(receiver, weak=False)

unsubscribe(event_class, callback)

Disconnect the callback. An unknown callback is a no-op, so a teardown path can run twice.

Source code in src/winslow/bus.py
def unsubscribe(self, event_class, callback):
    """Disconnect the callback. An unknown callback is a no-op, so a
    teardown path can run twice."""
    with self._lock:
        receiver = self._receivers.pop((event_class, callback), None)
        if receiver is not None:
            self._signals[event_class].disconnect(receiver)

publish(event)

Dispatch the event to its subscribers, on the calling thread. The dispatch order is undefined: a subscriber must not depend on another subscriber. A publish on a closed bus is a no-op: the session end can race a draining worker.

Source code in src/winslow/bus.py
def publish(self, event):
    """Dispatch the event to its subscribers, on the calling thread. The
    dispatch order is undefined: a subscriber must not depend on another
    subscriber. A publish on a closed bus is a no-op: the session end can
    race a draining worker."""
    with self._lock:
        if self._closed:
            return
        signal = self._signals.get(type(event))
    if signal is None:
        self._refuse_undeclared(type(event))
    signal.send(self, event=event)

close()

Disconnect every remaining subscriber. The session end calls this after the SessionEndedEvent dispatch. This method is idempotent.

Source code in src/winslow/bus.py
def close(self):
    """Disconnect every remaining subscriber. The session end calls this
    after the SessionEndedEvent dispatch. This method is idempotent."""
    with self._lock:
        for (event_class, _), receiver in self._receivers.items():
            self._signals[event_class].disconnect(receiver)
        self._receivers.clear()
        # An empty table frees the signals with the session.
        self._signals.clear()
        self._closed = True

winslow.events

The event payloads of the session bus (see SessionBus). An event carries values, never a live object (the payload rule, see docs/ui-plugins.md). The event class is the topic: a subscriber passes it to SessionBus.subscribe.

Origin

Bases: Enum

Why a store write happened. RUN is a live transition. SEED is a restore write (see Workflow.seed_from_state).

Source code in src/winslow/events.py
class Origin(Enum):
    """Why a store write happened. RUN is a live transition. SEED is a
    restore write (see Workflow.seed_from_state)."""

    RUN = "run"
    SEED = "seed"

TaskStatusEvent dataclass

One status write on the task store of the session.

Source code in src/winslow/events.py
@dataclass(frozen=True)
class TaskStatusEvent:
    """One status write on the task store of the session."""

    key: str
    status: "TaskStatus"  # noqa: F821
    origin: Origin = Origin.RUN

ExecutionStatusEvent dataclass

One status write on the record store of one batch.

Source code in src/winslow/events.py
@dataclass(frozen=True)
class ExecutionStatusEvent:
    """One status write on the record store of one batch."""

    task_key: str
    status: "TaskStatus"  # noqa: F821
    batch_uuid: str
    origin: Origin = Origin.RUN

BatchCreatedEvent dataclass

One admitted batch, published before its first task work.

Source code in src/winslow/events.py
@dataclass(frozen=True)
class BatchCreatedEvent:
    """One admitted batch, published before its first task work."""

    info: "BatchInfo"  # noqa: F821

BatchCompletedEvent dataclass

One completed batch, published after its final status is set.

Source code in src/winslow/events.py
@dataclass(frozen=True)
class BatchCompletedEvent:
    """One completed batch, published after its final status is set."""

    info: "BatchInfo"  # noqa: F821

LogLineEvent dataclass

One captured log line of one task in one batch.

Source code in src/winslow/events.py
@dataclass(frozen=True)
class LogLineEvent:
    """One captured log line of one task in one batch."""

    task_key: str
    batch_uuid: str
    line: str

Model

winslow.model

The data model of the session port: every value shape that crosses the port or the serve boundary (the payload rule, see winslow.events). Every field is JSON-safe. The local adapter hands these instances through in-process (see winslow.client.local). winslow.protocol.codec writes each to the wire and decodes the payload back, so a wire shape has exactly one declaration.

The producing side keeps the from_x classmethods, for example from_task and from_batch. A constructor that needs core machinery imports it inside the method, so this module imports nothing from winslow at module level.

EntryState

Bases: StrEnum

The freshness of one entry, derived at peek time (see BaseCache.inspect).

Source code in src/winslow/model.py
class EntryState(StrEnum):
    """The freshness of one entry, derived at peek time (see BaseCache.inspect)."""

    COLD = "cold"
    WARM = "warm"
    STALE = "stale"
    # A loader produces the value right now. A live observation only: the
    # state never appears in a history snapshot.
    COMPUTING = "computing"
    # A delete or a loader failed on the entry. The error context of the
    # projection names the operation and the layer (see CacheEntryError).
    ERRORED = "errored"

ErrorOrigin

Bases: StrEnum

The operation that left an entry in the ERRORED state.

Source code in src/winslow/model.py
class ErrorOrigin(StrEnum):
    """The operation that left an entry in the ERRORED state."""

    DELETE = "delete"
    LOAD = "load"

CacheEntryError dataclass

The error context of one entry: the failed operation and its layer. Plain strings only, so the context stays wire-ready. A successful write of the entry clears it (see BaseCache._entry_value).

Source code in src/winslow/model.py
@dataclass(frozen=True)
class CacheEntryError:
    """The error context of one entry: the failed operation and its layer.
    Plain strings only, so the context stays wire-ready. A successful write
    of the entry clears it (see BaseCache._entry_value)."""

    origin: ErrorOrigin
    tier: Optional[str]  # the failing storage layer; None for a loader error
    message: str
    at: float
    # The traceback, formatted at failure time: a string retains no frame,
    # so the context stays GC-safe and a value view can show the full cause.
    traceback: Optional[str] = None

SnapshotEncoding

Bases: StrEnum

What the rendered field of a snapshot holds. TEXT displays as it is; JSON deserializes back into a value for a tree view.

Source code in src/winslow/model.py
class SnapshotEncoding(StrEnum):
    """What the rendered field of a snapshot holds. TEXT displays as it is;
    JSON deserializes back into a value for a tree view."""

    TEXT = "text"
    JSON = "json"

CacheReadSnapshot dataclass

One recorded cache read of a task phase, rendered and bounded. Plain strings only, so a history record outlives the session and its caches. A summary marks a bounded rendering; a full rendering carries none.

Source code in src/winslow/model.py
@dataclass(frozen=True)
class CacheReadSnapshot:
    """One recorded cache read of a task phase, rendered and bounded. Plain
    strings only, so a history record outlives the session and its caches.
    A summary marks a bounded rendering; a full rendering carries none."""

    scope: str
    cache_name: str
    entry_name: str
    written_at: float
    rendered: str
    summary: Optional[str]
    encoding: SnapshotEncoding

CacheEntryInfo dataclass

The projection of one cache entry for a UI layer. Plain scalar fields only, so the projection is wire-ready. It never carries the value.

Source code in src/winslow/model.py
@dataclass(frozen=True)
class CacheEntryInfo:
    """The projection of one cache entry for a UI layer. Plain scalar fields
    only, so the projection is wire-ready. It never carries the value."""

    scope: str
    cache_name: str
    entry_name: str
    state: EntryState
    written_at: Optional[float]
    ttl: Optional[float]
    eager: bool
    depends_on: tuple[str, ...]
    storage: str
    error: Optional[CacheEntryError]

SourceNode dataclass

A node in the inheritance source tree of a task. label and location are display-ready, so a consumer renders the tree without the origin rules of winslow.task.info.

Source code in src/winslow/model.py
@dataclass(frozen=True)
class SourceNode:
    """A node in the inheritance source tree of a task. label and location
    are display-ready, so a consumer renders the tree without the origin
    rules of winslow.task.info."""

    name: str
    module: str
    source: str
    path: str | None  # the absolute source file
    children: tuple["SourceNode", ...]
    # The name, with the location in parentheses only for an ambiguous name.
    label: str = ""
    # The project-relative path, or the dotted module outside the project.
    location: str = ""

TaskRef dataclass

A renderable pointer to another task: what a dependency row needs. No nested dependencies, so a TaskInfo stays bounded on a deep graph.

Source code in src/winslow/model.py
@dataclass(frozen=True)
class TaskRef:
    """A renderable pointer to another task: what a dependency row needs. No
    nested dependencies, so a TaskInfo stays bounded on a deep graph."""

    key: str
    label: str
    is_premier: bool
    is_terminal: bool
    is_noop: bool

    def __str__(self):
        return self.label

    @classmethod
    def from_task(cls, task):
        return cls(
            key=task.identity_key,
            label=str(task),
            is_premier=task.is_premier,
            is_terminal=task.is_terminal,
            is_noop=task.is_noop,
        )

TaskInfo dataclass

The value view model of a task: plain values only, so asdict is JSON-serializable and history can hold it without a retention of the task.

from_task has two depths. The stub, the default, carries the identity and the dependency refs. The full capture adds attributes, docs, source and transients, and it evaluates a getter only with evaluate=True, which only the on-demand detail view passes. Equality and hash use the key.

Source code in src/winslow/model.py
@dataclass(frozen=True, eq=False)
class TaskInfo:
    """The value view model of a task: plain values only, so asdict is
    JSON-serializable and history can hold it without a retention of the task.

    from_task has two depths. The stub, the default, carries the identity and
    the dependency refs. The full capture adds attributes, docs, source and
    transients, and it evaluates a getter only with evaluate=True, which only
    the on-demand detail view passes. Equality and hash use the key."""

    key: str
    label: str
    name: str
    is_premier: bool
    is_terminal: bool
    is_noop: bool
    task_class: str
    index: int
    groups: tuple[str, ...] = ()
    parameters: dict | None = None
    dependencies: tuple[TaskRef, ...] = ()
    premier_dependencies: tuple[TaskRef, ...] = ()
    terminal_dependencies: tuple[TaskRef, ...] = ()
    # The trust fields of the check_ttl rule, from the session snapshots. None
    # means no verification on record, or no TTL (see Workflow.task_info).
    checked_at: float | None = None
    effective_ttl: float | None = None
    # Full-capture fields. None marks a stub, an empty tuple marks a capture
    # that found nothing. A doc is a (title, markdown) pair.
    attributes: tuple[AttributeSection, ...] | None = None
    docs: tuple[tuple[str, str], ...] | None = None
    source: SourceNode | None = None
    transients: tuple[str, ...] | None = None

    def __str__(self):
        return self.label

    def __hash__(self):
        return hash(self.key)

    def __eq__(self, other):
        if not isinstance(other, TaskInfo):
            return NotImplemented
        return self.key == other.key

    def get_name(self):
        return self.name

    def get_groups(self):
        return frozenset(self.groups)

    @property
    def groups_readable(self):
        return ", ".join(self.groups) if self.groups else None

    @classmethod
    def from_task(
        cls,
        task,
        full=False,
        evaluate=False,
        root_dir=None,
        checked_at=None,
        effective_ttl=None,
    ):
        # The capture machinery walks live classes and files, so it stays in
        # winslow.task.info; only the value shape lives here.
        from winslow.decorators import declared_transient_properties
        from winslow.task.info import (
            _attribute_sections,
            _display_parameters,
            _safe_sourcefile,
            _task_docs,
            labeled_source_tree,
        )

        full = full or evaluate
        deps = task.dependent_tasks
        task_cls = task.__class__

        return cls(
            key=task.identity_key,
            checked_at=checked_at,
            effective_ttl=effective_ttl,
            label=str(task),
            name=task.instance_name,
            is_terminal=task.is_terminal,
            is_premier=task.is_premier,
            is_noop=task.is_noop,
            index=task._index,
            # The real source file through inspect, and not the synthetic scoped
            # __module__.
            task_class=f"{task_cls.__qualname__} ({_safe_sourcefile(task_cls) or task_cls.__module__})",
            groups=tuple(sorted(task.get_groups())),
            parameters=_display_parameters(task) or None,
            dependencies=tuple(
                TaskRef.from_task(d)
                for d in deps
                if not (d.is_premier or d.is_terminal)
            ),
            premier_dependencies=tuple(
                TaskRef.from_task(d) for d in deps if d.is_premier
            ),
            terminal_dependencies=tuple(
                TaskRef.from_task(d) for d in deps if d.is_terminal
            ),
            attributes=_attribute_sections(task, root_dir, evaluate) if full else None,
            docs=_task_docs(task) if full else None,
            source=labeled_source_tree(task_cls, root_dir) if full else None,
            transients=(
                tuple(sorted(declared_transient_properties(task_cls))) if full else None
            ),
        )

BatchInfo dataclass

The value snapshot of one batch, for events and the wire (the payload rule, see winslow.events). Enum values travel by name, timestamps as epoch seconds.

Source code in src/winslow/model.py
@dataclass(frozen=True)
class BatchInfo:
    """The value snapshot of one batch, for events and the wire (the payload
    rule, see winslow.events). Enum values travel by name, timestamps as epoch
    seconds."""

    uuid: str
    action: str
    status: str
    task_count: int
    # The roster, {identity key: label} (see BatchRecord.tasks).
    tasks: dict[str, str]
    # The batch option snapshot of the execution context, without the uuid.
    options: dict | None
    created_at: float
    started_at: float | None
    completed_at: float | None
    # The message of the framework error that aborted the batch, or None.
    error: str | None = None

    @classmethod
    def from_batch(cls, batch, tasks):
        return cls._build(batch, {task.identity_key: str(task) for task in tasks})

    @classmethod
    def from_stored(cls, batch, store):
        """The info of a stored batch: the labels come from its records, so
        a snapshot carries the value the created event carried. The
        reconnect heal re-emits it (see RemoteSessionClient._on_snapshot)."""
        labels = (
            {key: store.get_record(key).info.label for key, _ in store.items()}
            if store is not None
            else {}
        )
        return cls._build(batch, labels)

    @classmethod
    def _build(cls, batch, tasks):
        return cls(
            uuid=batch.uuid,
            action=batch.action.name,
            status=batch.status.name,
            task_count=batch.task_count,
            tasks=tasks,
            options=_batch_options(batch),
            created_at=batch.created_at.timestamp(),
            started_at=(batch.started_at.timestamp() if batch.started_at else None),
            completed_at=(
                batch.completed_at.timestamp() if batch.completed_at else None
            ),
            error=batch.error,
        )

from_stored(batch, store) classmethod

The info of a stored batch: the labels come from its records, so a snapshot carries the value the created event carried. The reconnect heal re-emits it (see RemoteSessionClient._on_snapshot).

Source code in src/winslow/model.py
@classmethod
def from_stored(cls, batch, store):
    """The info of a stored batch: the labels come from its records, so
    a snapshot carries the value the created event carried. The
    reconnect heal re-emits it (see RemoteSessionClient._on_snapshot)."""
    labels = (
        {key: store.get_record(key).info.label for key, _ in store.items()}
        if store is not None
        else {}
    )
    return cls._build(batch, labels)

TaskOutcome dataclass

The outcome of one task in one batch: the status of the record store plus the record fields a history row shows.

Source code in src/winslow/model.py
@dataclass(frozen=True)
class TaskOutcome:
    """The outcome of one task in one batch: the status of the record store
    plus the record fields a history row shows."""

    status: str
    started_at: float | None
    duration: float | None
    last_log: str

    @classmethod
    def from_record(cls, status, record):
        return cls(
            status=status.name,
            started_at=(record.started_at.timestamp() if record.started_at else None),
            duration=record.duration,
            last_log=record.last_log,
        )

BatchOutcome dataclass

One batch of the history, with the per-task outcomes of its record store. A client that subscribes after the batch renders these rows without one log_tail call per task.

Source code in src/winslow/model.py
@dataclass(frozen=True)
class BatchOutcome:
    """One batch of the history, with the per-task outcomes of its record
    store. A client that subscribes after the batch renders these rows
    without one log_tail call per task."""

    uuid: str
    action: str
    status: str
    task_count: int
    created_at: float
    completed_at: float | None
    # The batch option snapshot, without the uuid (see BatchInfo.options).
    options: dict | None
    tasks: dict[str, TaskOutcome]  # keyed by identity key

    @classmethod
    def from_batch(cls, batch, store):
        return cls(
            uuid=batch.uuid,
            action=batch.action.name,
            status=batch.status.name,
            task_count=batch.task_count,
            created_at=batch.created_at.timestamp(),
            completed_at=(
                batch.completed_at.timestamp() if batch.completed_at else None
            ),
            options=_batch_options(batch),
            tasks=(
                {
                    key: TaskOutcome.from_record(status, store.get_record(key))
                    for key, status in store.items()
                }
                if store is not None
                else {}
            ),
        )

StatusSnapshot dataclass

The latest snapshot of one task in one session. The key is the task identity key; the status is a TaskStatus name; checked_at is a wall-clock epoch.

Source code in src/winslow/model.py
@dataclass(frozen=True)
class StatusSnapshot:
    """The latest snapshot of one task in one session. The key is the task
    identity key; the status is a TaskStatus name; checked_at is a wall-clock
    epoch."""

    key: str
    status: str
    checked_at: float

TaskStatusSummary dataclass

The (completed, problematic, total) counts of one session (see Session.task_status_summary).

Source code in src/winslow/model.py
@dataclass(frozen=True)
class TaskStatusSummary:
    """The (completed, problematic, total) counts of one session (see
    Session.task_status_summary)."""

    completed: int
    problematic: int
    total: int

SessionInfo dataclass

One row of the session list (see AppClient.sessions).

Source code in src/winslow/model.py
@dataclass(frozen=True)
class SessionInfo:
    """One row of the session list (see AppClient.sessions)."""

    session_id: str
    workflow: str
    status: str
    display_name: str
    instance_name: str
    identifier_suffix: str
    started_at: float
    elapsed: float
    task_status_summary: TaskStatusSummary
    # The project root of the serving process. A client shortens the server
    # source paths with it (see TaskDetailRenderContext.root_dir).
    root_dir: str | None = None

    @classmethod
    def from_session(cls, session):
        workflow = session.workflow
        completed, problematic, total = session.task_status_summary
        return cls(
            session_id=session.session_id,
            workflow=str(workflow),
            status=session.status.name,
            display_name=workflow.get_display_name(),
            instance_name=workflow.instance_name,
            identifier_suffix=workflow.identifier_suffix,
            started_at=session.start,
            elapsed=session.elapsed,
            task_status_summary=TaskStatusSummary(
                completed=completed, problematic=problematic, total=total
            ),
            root_dir=workflow.root_dir,
        )

SessionSnapshot dataclass

The current state of one session: the task statuses, the batch rows, the session log backlog, and the session meta. The subscribe reply and the local snapshot read serve the same shape (see EventBridge.snapshot).

Source code in src/winslow/model.py
@dataclass(frozen=True)
class SessionSnapshot:
    """The current state of one session: the task statuses, the batch rows,
    the session log backlog, and the session meta. The subscribe reply and
    the local snapshot read serve the same shape (see EventBridge.snapshot)."""

    session_id: str
    workflow: str
    status: str
    tasks: dict[str, str]  # {identity key: TaskStatus name}
    session_log_backlog: tuple[str, ...]
    batches: tuple[BatchInfo, ...]
    # The cache names of the session, so a client can decide whether to show
    # a caches pane before the first caches read. None once the session has
    # ended and released its caches; an empty tuple means no registered caches.
    cache_names: tuple[str, ...] | None = None

    @classmethod
    def _cache_names(cls, workflow):
        try:
            return tuple(cache.get_name() for cache in workflow.caches())
        except ValueError:
            # The session ended between the status read and this read.
            return None

    @classmethod
    def from_session(cls, session):
        workflow = session.workflow
        return cls(
            session_id=session.session_id,
            workflow=str(workflow),
            status=session.status.name,
            tasks={key: status.name for key, status in workflow.store.items()},
            session_log_backlog=(
                tuple(session.log_buffer.lines)
                if session.log_buffer is not None
                else ()
            ),
            batches=tuple(
                BatchInfo.from_stored(batch, workflow.runner.record_store(batch.uuid))
                for batch in workflow.runner.batches
            ),
            cache_names=cls._cache_names(workflow),
        )

    def as_events(self):
        """The events a live subscriber saw to reach this state: each batch
        created, then completed if it is, then every task status, then the
        end if the session ended. A subscriber that heals a sequence gap
        replays them (see RemoteSessionClient._on_snapshot)."""
        from winslow.events import (
            BatchCompletedEvent,
            BatchCreatedEvent,
            Origin,
            SessionEndedEvent,
            TaskStatusEvent,
        )
        from winslow.task.status import TaskStatus

        # A created event precedes the statuses of its tasks, as on the live bus.
        batches = [
            event
            for info in self.batches
            for event in (
                BatchCreatedEvent(info=info),
                *(
                    [BatchCompletedEvent(info=info)]
                    if info.completed_at is not None
                    else []
                ),
            )
        ]
        statuses = [
            TaskStatusEvent(key=key, status=TaskStatus[name], origin=Origin.RUN)
            for key, name in self.tasks.items()
        ]
        end = (
            [SessionEndedEvent(session_id=self.session_id)]
            if self.status == "ENDED"
            else []
        )
        return tuple(batches + statuses + end)

as_events()

The events a live subscriber saw to reach this state: each batch created, then completed if it is, then every task status, then the end if the session ended. A subscriber that heals a sequence gap replays them (see RemoteSessionClient._on_snapshot).

Source code in src/winslow/model.py
def as_events(self):
    """The events a live subscriber saw to reach this state: each batch
    created, then completed if it is, then every task status, then the
    end if the session ended. A subscriber that heals a sequence gap
    replays them (see RemoteSessionClient._on_snapshot)."""
    from winslow.events import (
        BatchCompletedEvent,
        BatchCreatedEvent,
        Origin,
        SessionEndedEvent,
        TaskStatusEvent,
    )
    from winslow.task.status import TaskStatus

    # A created event precedes the statuses of its tasks, as on the live bus.
    batches = [
        event
        for info in self.batches
        for event in (
            BatchCreatedEvent(info=info),
            *(
                [BatchCompletedEvent(info=info)]
                if info.completed_at is not None
                else []
            ),
        )
    ]
    statuses = [
        TaskStatusEvent(key=key, status=TaskStatus[name], origin=Origin.RUN)
        for key, name in self.tasks.items()
    ]
    end = (
        [SessionEndedEvent(session_id=self.session_id)]
        if self.status == "ENDED"
        else []
    )
    return tuple(batches + statuses + end)

PhaseInfo dataclass

One entry of a record's phase timeline (see PhaseSpan).

Source code in src/winslow/model.py
@dataclass(frozen=True)
class PhaseInfo:
    """One entry of a record's phase timeline (see PhaseSpan)."""

    phase: str
    started_at: float
    completed_at: float | None
    duration: float | None

    @classmethod
    def from_span(cls, span):
        return cls(
            phase=span.phase.value,
            started_at=span.started_at.timestamp(),
            completed_at=(span.completed_at.timestamp() if span.completed_at else None),
            duration=span.duration,
        )

RecordDetail dataclass

The full capture of one execution record: its TaskInfo, its phase timeline, and its transient and cache snapshots (see ExecutionRecord). The snapshot dicts key by the phase name.

Source code in src/winslow/model.py
@dataclass(frozen=True)
class RecordDetail:
    """The full capture of one execution record: its TaskInfo, its phase
    timeline, and its transient and cache snapshots (see ExecutionRecord).
    The snapshot dicts key by the phase name."""

    info: TaskInfo
    phases: tuple[PhaseInfo, ...]
    transient_snapshots: dict
    cache_snapshots: dict[str, tuple[CacheReadSnapshot, ...]]  # by phase name

    @classmethod
    def from_record(cls, record):
        return cls(
            info=record.info,
            phases=tuple(PhaseInfo.from_span(span) for span in record.phases),
            transient_snapshots={
                phase.value: snapshot
                for phase, snapshot in record.transient_snapshots.items()
            },
            cache_snapshots={
                phase.value: tuple(snapshots)
                for phase, snapshots in record.cache_snapshots.items()
            },
        )

CacheInfo dataclass

One cache: identity, storage, the declared entry names, and the current value preview of each written entry (see unobservable for the error form).

Source code in src/winslow/model.py
@dataclass(frozen=True)
class CacheInfo:
    """One cache: identity, storage, the declared entry names, and the current
    value preview of each written entry (see unobservable for the error form)."""

    name: str
    scope: str
    docstring: str | None
    storage: str
    entries: tuple[str, ...]
    info: tuple[CacheEntryInfo, ...]
    values: dict[str, str | None]  # {entry name: safe_repr preview}
    error: str | None = None

    @classmethod
    def _entry_names(cls, cache):
        from winslow.cache import declared_entries

        return tuple(declared_entries(type(cache)))

    @classmethod
    def from_cache(cls, cache):
        infos = tuple(cache.inspect())
        return cls(
            name=cache.get_name(),
            scope=cache.scope,
            docstring=type(cache).__doc__,
            storage=cache.describe_storage(),
            entries=cls._entry_names(cache),
            info=infos,
            values={
                info.entry_name: _entry_value_preview(cache, info.entry_name)
                for info in infos
                if info.state in PREVIEWABLE_STATES
            },
        )

    @classmethod
    def unobservable(cls, cache, error):
        """The info of a cache whose inspect or peek raised. The declarations
        stand, and the states and the values stay empty."""
        return cls(
            name=cache.get_name(),
            scope=cache.scope,
            docstring=type(cache).__doc__,
            storage=cache.describe_storage(),
            entries=cls._entry_names(cache),
            info=(),
            values={},
            error=error,
        )

unobservable(cache, error) classmethod

The info of a cache whose inspect or peek raised. The declarations stand, and the states and the values stay empty.

Source code in src/winslow/model.py
@classmethod
def unobservable(cls, cache, error):
    """The info of a cache whose inspect or peek raised. The declarations
    stand, and the states and the values stay empty."""
    return cls(
        name=cache.get_name(),
        scope=cache.scope,
        docstring=type(cache).__doc__,
        storage=cache.describe_storage(),
        entries=cls._entry_names(cache),
        info=(),
        values={},
        error=error,
    )

CachesPayload dataclass

Every cache of one session, in name order.

Source code in src/winslow/model.py
@dataclass(frozen=True)
class CachesPayload:
    """Every cache of one session, in name order."""

    caches: tuple[CacheInfo, ...]

    @classmethod
    def from_workflow(cls, workflow):
        """One card per cache. A cache whose storage raises answers an
        unobservable card, so one broken storage cannot hide the others."""

        def card(cache):
            try:
                return CacheInfo.from_cache(cache)
            except Exception as exc:
                workflow.logger.debug(
                    f"The caches read cannot observe '{cache.get_name()}'.",
                    exc_info=True,
                )
                return CacheInfo.unobservable(cache, str(exc))

        return cls(caches=tuple(card(cache) for cache in workflow.caches()))

from_workflow(workflow) classmethod

One card per cache. A cache whose storage raises answers an unobservable card, so one broken storage cannot hide the others.

Source code in src/winslow/model.py
@classmethod
def from_workflow(cls, workflow):
    """One card per cache. A cache whose storage raises answers an
    unobservable card, so one broken storage cannot hide the others."""

    def card(cache):
        try:
            return CacheInfo.from_cache(cache)
        except Exception as exc:
            workflow.logger.debug(
                f"The caches read cannot observe '{cache.get_name()}'.",
                exc_info=True,
            )
            return CacheInfo.unobservable(cache, str(exc))

    return cls(caches=tuple(card(cache) for cache in workflow.caches()))

CacheValueView dataclass

The rendered form of one cache entry value, built server-side. The live modal and a wire client thus render the same text the history path already serves (see CacheReadSnapshot). encoding and rendered stay None for a cold or a computing entry.

Source code in src/winslow/model.py
@dataclass(frozen=True)
class CacheValueView:
    """The rendered form of one cache entry value, built server-side. The
    live modal and a wire client thus render the same text the history path
    already serves (see CacheReadSnapshot). encoding and rendered stay None
    for a cold or a computing entry."""

    cache_name: str
    entry_name: str
    state: str
    encoding: str | None
    rendered: str | None
    summary: str | None
    written_at: float | None
    error: CacheEntryError | None

    @classmethod
    def from_entry(cls, cache, entry_name):
        from winslow.cache import (
            MISSING,
            declared_entries,
            render_value,
            resolve_snapshot_cap,
        )

        info = next(i for i in cache.inspect() if i.entry_name == entry_name)
        record = cache.peek(entry_name)
        if record is MISSING:
            return cls._unrendered(cache, entry_name, EntryState.COLD, info.error)
        if record is EntryState.COMPUTING:
            return cls._unrendered(cache, entry_name, EntryState.COMPUTING, info.error)
        display_style = declared_entries(type(cache))[entry_name].display_style
        rendered, summary, encoding = render_value(
            record.value, resolve_snapshot_cap(type(cache)), display_style
        )
        return cls(
            cache_name=cache.get_name(),
            entry_name=entry_name,
            state=info.state.value,
            encoding=encoding.value,
            rendered=rendered,
            summary=summary,
            written_at=record.written_at,
            error=info.error,
        )

    @classmethod
    def _unrendered(cls, cache, entry_name, state, error):
        """The view of an entry with no value to render: cold, or a loader
        runs right now."""
        return cls(
            cache_name=cache.get_name(),
            entry_name=entry_name,
            state=state.value,
            encoding=None,
            rendered=None,
            summary=None,
            written_at=None,
            error=error,
        )

SessionParams dataclass

settings_snapshot plus the resolved workflow_config values of one session (see WorkflowParams, the local modal with the same content).

Source code in src/winslow/model.py
@dataclass(frozen=True)
class SessionParams:
    """settings_snapshot plus the resolved workflow_config values of one
    session (see WorkflowParams, the local modal with the same content)."""

    settings: dict
    workflow_config: dict

    @classmethod
    def from_workflow(cls, workflow):
        return cls(
            settings=workflow.settings_snapshot,
            workflow_config={
                name: workflow.config_meta[name].format_value(
                    getattr(workflow.workflow_config, name, None)
                )
                for name in workflow.config_option_names
            },
        )

ManifestInfo dataclass

One restorable session manifest (see SessionManifest).

Source code in src/winslow/model.py
@dataclass(frozen=True)
class ManifestInfo:
    """One restorable session manifest (see SessionManifest)."""

    session_id: str
    workflow_class: str
    orchestrator_overrides: dict | None
    workflow_values: dict | None

    @classmethod
    def from_manifest(cls, manifest):
        return cls(
            session_id=manifest.session_id,
            workflow_class=manifest.workflow_class,
            orchestrator_overrides=manifest.orchestrator_overrides,
            workflow_values=manifest.workflow_values,
        )

OptionInfo dataclass

One ConfigOption as form metadata. Defaults travel as formatted strings: right for a form, accepted for an agent (see serve-spec 6.1). initial names the value a form should prefill. It is the live parsed value when the caller passes one, an orchestrator override already parsed from the CLI, and the declared default otherwise.

Source code in src/winslow/model.py
@dataclass(frozen=True)
class OptionInfo:
    """One ConfigOption as form metadata. Defaults travel as formatted
    strings: right for a form, accepted for an agent (see serve-spec 6.1).
    initial names the value a form should prefill. It is the live parsed
    value when the caller passes one, an orchestrator override already
    parsed from the CLI, and the declared default otherwise."""

    name: str
    help: str | None
    default: str | None
    initial: str | None
    required: bool
    choices: tuple[str, ...] | None
    multiselect: bool
    type: str | None
    identifier: bool
    depends_on: tuple[str, ...]
    action: str | None
    const: object
    # The initial value of a multiselect option, one string per selected
    # choice. The joined `initial` string cannot split a choice that holds a
    # comma, so a form preselects from this field.
    initial_selection: tuple[str, ...] | None = None

    @classmethod
    def from_option(cls, name, option, current=None):
        initial = option.default if current is None else current
        return cls(
            name=name,
            help=option.help_text,
            default=option.format_value(option.default),
            initial=option.format_value(initial),
            initial_selection=(
                tuple(str(value) for value in initial)
                if option.multiselect and initial is not None
                else None
            ),
            required=option.required,
            choices=(
                tuple(str(choice) for choice in option.choices)
                if option.choices
                else None
            ),
            multiselect=option.multiselect,
            type=option.type.__name__ if option.type else None,
            identifier=option.identifier,
            depends_on=tuple(option.depends_on),
            action=option.action,
            const=option.const,
        )

WorkflowDescriptor dataclass

One collected workflow and the options its start form shows. auto_init marks a workflow the dashboard starts without a form (see Workflow.auto_init).

Source code in src/winslow/model.py
@dataclass(frozen=True)
class WorkflowDescriptor:
    """One collected workflow and the options its start form shows.
    auto_init marks a workflow the dashboard starts without a form (see
    Workflow.auto_init)."""

    workflow: str
    options: tuple[OptionInfo, ...]
    auto_init: bool = False

Descriptors dataclass

The parameter context of one process: a descriptor per collected workflow (the values of create_session), plus the orchestrator options the start form shows (the overrides). A workflow's initial values come from the same CLI-supplied args that prefill the local start form (see Orchestrator.collect_workflow_args). A workflow the orchestrator config excludes from the local selector is excluded here too.

Source code in src/winslow/model.py
@dataclass(frozen=True)
class Descriptors:
    """The parameter context of one process: a descriptor per collected
    workflow (the `values` of create_session), plus the orchestrator options
    the start form shows (the `overrides`). A workflow's initial values come
    from the same CLI-supplied args that prefill the local start form (see
    Orchestrator.collect_workflow_args). A workflow the orchestrator config
    excludes from the local selector is excluded here too."""

    workflows: tuple[WorkflowDescriptor, ...]
    overrides: tuple[OptionInfo, ...]

    @classmethod
    def from_orchestrator(cls, orchestrator):
        workflow_args = orchestrator.collect_workflow_args()
        workflows = []
        for name in orchestrator.workflow_registry.names:
            workflow_kls = orchestrator.workflow_registry[name]
            if not workflow_kls.should_be_initialized(orchestrator.orchestrator_config):
                continue
            parsed = workflow_args.get(workflow_kls)
            workflows.append(
                WorkflowDescriptor(
                    workflow=name,
                    auto_init=workflow_kls.auto_init,
                    options=tuple(
                        OptionInfo.from_option(
                            option_name,
                            option,
                            current=getattr(parsed, option_name, None),
                        )
                        for option_name, option in workflow_kls.config_meta.items()
                        if option.show_on_ui
                    ),
                )
            )
        overrides = tuple(
            OptionInfo.from_option(
                option_name,
                option,
                current=getattr(orchestrator.orchestrator_config, option_name, None),
            )
            for option_name, option in orchestrator.config_meta.items()
            if option.show_on_ui
        )
        return cls(workflows=tuple(workflows), overrides=overrides)

CacheUpdatedEvent dataclass

The repaint trigger of one cache. The port synthesizes it from the CacheListener callbacks (see winslow.client.local) or from cache_updated frames.

Source code in src/winslow/model.py
@dataclass(frozen=True)
class CacheUpdatedEvent:
    """The repaint trigger of one cache. The port synthesizes it from the
    CacheListener callbacks (see winslow.client.local) or from cache_updated
    frames."""

    cache_name: str

SessionLogEvent dataclass

One line of the session log lane: the workflow logger stream (init, eligibility, cache lines). LogLineEvent stays task-scoped.

Source code in src/winslow/model.py
@dataclass(frozen=True)
class SessionLogEvent:
    """One line of the session log lane: the workflow logger stream (init,
    eligibility, cache lines). LogLineEvent stays task-scoped."""

    line: str

TaskLogEvent dataclass

One live line of a task log subscription (see SessionClient.subscribe_task_log).

Source code in src/winslow/model.py
@dataclass(frozen=True)
class TaskLogEvent:
    """One live line of a task log subscription (see
    SessionClient.subscribe_task_log)."""

    task_key: str
    line: str

ConnectionEvent dataclass

The connection state of the wire transport: connected is False on a drop and True once the reconnect lands. Only the wire transport emits it (see AppClient.subscribe_connection).

Source code in src/winslow/model.py
@dataclass(frozen=True)
class ConnectionEvent:
    """The connection state of the wire transport: connected is False on a
    drop and True once the reconnect lands. Only the wire transport emits
    it (see AppClient.subscribe_connection)."""

    connected: bool

Sessions

winslow.session.SessionRegistry

The live sessions of one process, by session id. One registry serves every consumer of the process: the TUI app, and the serve transports (websocket, MCP), so each resolves the same map.

Source code in src/winslow/session.py
class SessionRegistry:
    """The live sessions of one process, by session id. One registry serves
    every consumer of the process: the TUI app, and the serve transports
    (websocket, MCP), so each resolves the same map."""

    def __init__(self):
        self._sessions = {}
        self._lock = threading.Lock()

    def register(self, session):
        with self._lock:
            self._sessions[session.session_id] = session

    def resolve(self, session_id):
        """The live session under the id. Raises KeyError with direction."""
        session = self.get(session_id)
        if session is None:
            raise KeyError(
                f"session id {session_id!r} does not resolve to a live session - "
                f"it ended, or it belongs to another process."
            )
        return session

    def get(self, session_id):
        return self._sessions.get(session_id)

    def remove(self, session_id):
        """Drop and return the session, or None: a teardown can run twice."""
        with self._lock:
            return self._sessions.pop(session_id, None)

    def sessions(self):
        # A tuple, so iteration survives a registration from another thread.
        return tuple(self._sessions.values())

    def __contains__(self, session_id):
        return session_id in self._sessions

    def __len__(self):
        return len(self._sessions)

resolve(session_id)

The live session under the id. Raises KeyError with direction.

Source code in src/winslow/session.py
def resolve(self, session_id):
    """The live session under the id. Raises KeyError with direction."""
    session = self.get(session_id)
    if session is None:
        raise KeyError(
            f"session id {session_id!r} does not resolve to a live session - "
            f"it ended, or it belongs to another process."
        )
    return session

remove(session_id)

Drop and return the session, or None: a teardown can run twice.

Source code in src/winslow/session.py
def remove(self, session_id):
    """Drop and return the session, or None: a teardown can run twice."""
    with self._lock:
        return self._sessions.pop(session_id, None)

winslow.session.create_session(orchestrator, state_store, registry, workflow_name, orchestrator_overrides=None, workflow_values=None, session_id=None, seed=False, origin='serve')

Build, initialize, persist, and register one session: the shared flow behind the serve create_session request and the local AppClient. origin stamps the manifest with the endpoint that created the session. Raises with a directional message on an unknown workflow; a failure after registration marks the session errored and unregisters it.

session_id and seed serve a restore: the caller passes the id of the stored manifest, so the session rebuilds under it, and seed=True replays the stored snapshots onto the store after the eligibility pass (see Workflow.seed_from_state).

Source code in src/winslow/session.py
def create_session(
    orchestrator,
    state_store,
    registry,
    workflow_name,
    orchestrator_overrides=None,
    workflow_values=None,
    session_id=None,
    seed=False,
    origin="serve",
):
    """Build, initialize, persist, and register one session: the shared flow
    behind the serve create_session request and the local AppClient. origin
    stamps the manifest with the endpoint that created the session. Raises with
    a directional message on an unknown workflow; a failure after
    registration marks the session errored and unregisters it.

    session_id and seed serve a restore: the caller passes the id of the
    stored manifest, so the session rebuilds under it, and seed=True replays
    the stored snapshots onto the store after the eligibility pass (see
    Workflow.seed_from_state)."""
    try:
        workflow_kls = orchestrator.workflow_registry[workflow_name]
    except KeyError:
        raise KeyError(
            f"workflow {workflow_name!r} names no collected workflow. "
            f"The workflows are {orchestrator.workflow_registry.names}."
        ) from None

    orchestrator_overrides = orchestrator_overrides or {}
    workflow_values = workflow_values or {}
    # The parsed CLI base fills every option the caller did not send, so an
    # option outside the form keeps its command-line value (see
    # Orchestrator.collect_workflow_args).
    workflow_base = orchestrator.collect_workflow_args().get(workflow_kls)
    workflow_values, orchestrator_overrides = validate_values(
        workflow_name,
        workflow_kls,
        orchestrator,
        workflow_values,
        orchestrator_overrides,
        workflow_base=workflow_base,
    )
    session_id = session_id or generate_id(workflow_name)
    workflow_logger = logging.getLogger(run_logger_name(session_id))
    workflow_logger.propagate = True
    # Attached before any initialization work runs, so init and eligibility
    # lines survive until a client subscribes (see SessionLogBuffer).
    log_buffer = SessionLogBuffer()
    workflow_logger.addHandler(log_buffer)

    init_log_ctx = LogContext(
        session_id=session_id,
        workflow_name=workflow_name,
        workflow_instance=workflow_name,
        task_name=None,
        task_instance=None,
        batch_uuid=None,
    )
    with scoped_log_context(init_log_ctx):
        workflow = orchestrator.initialize_workflow(
            workflow_kls=workflow_kls,
            orchestrator_overrides=orchestrator_overrides,
            workflow_values=workflow_values,
            workflow_base=workflow_base,
            logger=workflow_logger,
        )
        session = Session(workflow, session_id=session_id, log_buffer=log_buffer)
        registry.register(session)
        try:
            workflow.initialize_tasks(logger=workflow.logger)
            workflow.check_pipeline_eligibility(logger=workflow.logger)
            # Persistence starts only once the pipeline is runnable: a kill
            # during the initialization above leaves no restore candidate.
            # The manifest stores the effective workflow values, so a restore
            # does not depend on the argv of its process (spec decision 9).
            workflow.init_state(
                state_store,
                origin=origin,
                orchestrator_overrides=orchestrator_overrides,
                workflow_values=effective_workflow_values(
                    workflow_kls, workflow_base, workflow_values
                ),
            )
            if seed:
                # After the eligibility pass: that pass overwrites earlier
                # status writes (see Workflow.seed_from_state).
                workflow.seed_from_state()
        except Exception as exc:
            registry.remove(session_id)
            session.mark_error(exc)
            raise
    return session