A Python library for implementing non-linear behavior and stream processing using a predicate-action pattern.
Beta Software - Proceed with caution
- Predicate → Action based control flow
- Multi-state and single-state finite machine capabilities
- Built-in dependency injector with lazy resolving
- Concurrent predicate and action execution
- Functional approach using decorators
- Zero external dependencies
- Python 3.10+
pip install cogesOr with Poetry:
poetry add cogesimport asyncio
from coges.coge import create_coge
from coges.machine import create_machine
# Create a coge (predicate + action pair)
greeter = create_coge("greeter")
@greeter.predicate
async def should_greet(tick, **state):
return tick.get("type") == "greeting"
@greeter.action
async def greet(tick, **state):
return f"Hello, {tick.get('name', 'World')}!"
# Create tick generator
async def ticks():
yield {"type": "greeting", "name": "Alice"}
yield {"type": "other"}
yield {"type": "greeting", "name": "Bob"}
# Create and run machine
machine = create_machine(
coges=[greeter],
tick_fn=ticks()
)
asyncio.run(machine())A Coge is the fundamental building block - a combination of a predicate and an action:
- Predicate: An async function returning
boolthat determines if the action should run - Action: An async function that executes when the predicate returns
True
from coges.coge import create_coge
my_coge = create_coge("my_coge")
@my_coge.predicate
async def check_condition(tick, **state):
return some_condition(tick)
@my_coge.action
async def do_something(tick, results, **state):
return some_resultThe Machine orchestrates Coges, running them on each tick from an async generator:
from coges.machine import create_machine
machine = create_machine(
coges=[coge1, coge2, coge3],
tick_fn=my_async_generator(),
di_resolver=resolve, # optional
initial_state={"key": "value"} # optional
)
await machine()Execution flow per tick:
- Get next tick from async generator
- Run all predicates concurrently
- Filter coges where predicate returned
True - Run all active actions concurrently
- Update state with results
State is a dictionary passed to all predicates and actions. It contains:
| Key | Description |
|---|---|
tick | Current tick value from the generator |
results | Dictionary of {coge_name: result} from last tick |
active_coges | List of coge names that were active last tick |
... | Any keys from initial_state |
@my_coge.action
async def handler(tick, results, active_coges, **state):
previous_result = results.get("other_coge")
return process(tick, previous_result)Coges includes a built-in DI container with lazy resolution:
from coges.di import create_dependency_injector
add_dependency, resolve = create_dependency_injector()
# Define dependencies as generator functions (single yield)
@add_dependency
def database():
db = Database.connect()
yield db
db.close() # cleanup after yield
@add_dependency("custom_name")
def my_service():
yield Service()
# Dependencies are injected by parameter name
@my_coge.action
async def handler(tick, database, custom_name, **state):
return database.query(tick["id"])
# Pass resolver to machine
machine = create_machine(
coges=[my_coge],
tick_fn=ticks(),
di_resolver=resolve
)Key features:
- Lazy resolution: Dependencies instantiated only when needed
- Generator-based lifecycle: Code after
yieldruns on cleanup - Automatic injection: Matches parameter names to dependency names
Factory function to create a Coge.
coge = create_coge("my_coge")Decorator to set the predicate function.
@coge.predicate
async def check(tick, **state) -> bool:
return TrueDecorator to set the action function.
@coge.action
async def execute(tick, **state) -> Any:
return resultCreates a machine that runs coges on each tick.
| Parameter | Type | Description |
|---|---|---|
coges | list[Coge] | List of coges to run |
tick_fn | AsyncGenerator | Async generator yielding ticks |
di_resolver | Callable | DI resolver function (optional) |
initial_state | dict | Initial state values (optional) |
Returns an async function that runs the machine.
Creates a dependency injector, returning two functions:
add_dependency: Decorator to register dependenciesresolve: Function to resolve dependencies for a function
import asyncio
from coges.coge import create_coge
from coges.machine import create_machine
# Define coges for different intents
help_handler = create_coge("help")
order_handler = create_coge("order")
fallback_handler = create_coge("fallback")
@help_handler.predicate
async def is_help(tick, **s):
return "help" in tick["message"].lower()
@help_handler.action
async def send_help(tick, **s):
return "Here's how I can help..."
@order_handler.predicate
async def is_order(tick, **s):
return "order" in tick["message"].lower()
@order_handler.action
async def process_order(tick, **s):
return "Processing your order..."
@fallback_handler.predicate
async def always(tick, **s):
return True
@fallback_handler.action
async def default_response(tick, active_coges, **s):
if len(active_coges) == 1: # only fallback active
return "I don't understand"
return None
# Message stream
async def messages():
yield {"message": "I need help"}
yield {"message": "Place an order"}
yield {"message": "Random text"}
machine = create_machine(
coges=[help_handler, order_handler, fallback_handler],
tick_fn=messages()
)
asyncio.run(machine())import asyncio
from coges.coge import create_coge
from coges.machine import create_machine
from coges.di import create_dependency_injector
add_dependency, resolve = create_dependency_injector()
@add_dependency
def logger():
import logging
yield logging.getLogger("pipeline")
validator = create_coge("validator")
transformer = create_coge("transformer")
loader = create_coge("loader")
@validator.predicate
async def needs_validation(tick, **s):
return "data" in tick
@validator.action
async def validate(tick, logger, **s):
logger.info(f"Validating: {tick['data']}")
return {"valid": True, "data": tick["data"]}
@transformer.predicate
async def can_transform(tick, results, **s):
return results.get("validator", {}).get("valid", False)
@transformer.action
async def transform(tick, results, **s):
data = results["validator"]["data"]
return {"transformed": data.upper()}
@loader.predicate
async def can_load(tick, results, **s):
return "transformer" in results
@loader.action
async def load(tick, results, logger, **s):
logger.info(f"Loading: {results['transformer']}")
return {"loaded": True}
async def data_stream():
yield {"data": "record1"}
yield {"data": "record2"}
yield {"skip": True}
machine = create_machine(
coges=[validator, transformer, loader],
tick_fn=data_stream(),
di_resolver=resolve
)
asyncio.run(machine())Enable debug logging to see machine execution:
import logging
logging.basicConfig(level=logging.DEBUG)
logging.getLogger("coges.machine").setLevel(logging.DEBUG)tick_fnnow requires anAsyncGeneratorinstead of a regular generator
- Initial release
MIT - Jellyfish.tech