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 | |
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
effective_check_ttl(task)
¶
The check TTL of the task: its own declaration when set, else the workflow default (see 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
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.
add_cache_listener(listener)
¶
Attach the listener to the caches this workflow can see: the workflow cache and the global cache.
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
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
get_cache(name)
¶
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
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
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
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
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
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
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
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 | |
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
get_initialization_constraints(workflow_config)
classmethod
¶
Override this to select should_be_initialized constraints by env or by config.
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
is_eligible()
¶
can_run()
¶
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.
check()
¶
run()
¶
dry_run()
¶
The runner calls this instead of run in dry-run mode. Override it to simulate the change. The default makes no change.
TaskFilter¶
winslow.TaskFilter
¶
Bases: Registerable
Source code in src/winslow/filter/base.py
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
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
winslow.ConstraintType
¶
Bases: Enum
Source code in src/winslow/constraints.py
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
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
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
winslow.ConfigOption
dataclass
¶
Source code in src/winslow/descriptors.py
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
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
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
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 | |
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
describe_storage()
¶
The human label of the storage of this cache. It reads no entry and takes no lock (see BaseStorage.describe).
inspect()
¶
Return one CacheEntryInfo per declared entry. The values stay out; a detail view fetches one on demand with peek.
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
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
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
winslow.cache.get_workflow_cache()
¶
The workflow container of the current context.
Source code in src/winslow/cache/runtime.py
winslow.cache.get_global_cache()
¶
The process-level container. It exists after the first workflow initialization.
Source code in src/winslow/cache/runtime.py
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
read(key)
¶
Return the StorageRecord of the key, or the MISSING sentinel. None cannot mark a miss: it is a legal cached value.
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
describe()
¶
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.
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
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
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
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
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
get_handler(orchestrator_config)
¶
Return the TelemetryHandler to register, or None to stay inactive. An exception fails the start of the run loudly.
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
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.
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
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 | |
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
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
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
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.
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
winslow.client.SessionClient
¶
Bases: Port
One session: the reads, the subscriptions, and the actions.
Source code in src/winslow/client/base.py
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 | |
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
unsubscribe(topic, handler)
¶
Disconnect the handler (see subscribe). An unknown pair is a no-op, so a teardown path can run twice.
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.
unsubscribe_task_log(task_key, handler)
¶
submit(action)
¶
Submit one action dataclass (see winslow.actions) and return its ack. A refused action answers an ack that carries the reason.
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
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
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
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
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
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
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
CheckTasks
dataclass
¶
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
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
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
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 | |
submit(action)
¶
The one entry point: dispatch on the action class.
Source code in src/winslow/actions.py
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
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
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 | |
get_event_classes()
classmethod
¶
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
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
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
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
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
TaskStatusEvent
dataclass
¶
ExecutionStatusEvent
dataclass
¶
One status write on the record store of one batch.
Source code in src/winslow/events.py
BatchCreatedEvent
dataclass
¶
BatchCompletedEvent
dataclass
¶
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
ErrorOrigin
¶
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
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
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
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
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
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
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
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 | |
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
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
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
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
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
TaskStatusSummary
dataclass
¶
The (completed, problematic, total) counts of one session (see Session.task_status_summary).
Source code in src/winslow/model.py
SessionInfo
dataclass
¶
One row of the session list (see AppClient.sessions).
Source code in src/winslow/model.py
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
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 | |
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
PhaseInfo
dataclass
¶
One entry of a record's phase timeline (see PhaseSpan).
Source code in src/winslow/model.py
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
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
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
CachesPayload
dataclass
¶
Every cache of one session, in name order.
Source code in src/winslow/model.py
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
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
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
ManifestInfo
dataclass
¶
One restorable session manifest (see SessionManifest).
Source code in src/winslow/model.py
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
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
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
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
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
TaskLogEvent
dataclass
¶
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
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
resolve(session_id)
¶
The live session under the id. Raises KeyError with direction.
Source code in src/winslow/session.py
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
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 | |