Skip to content

Latest commit

 

History

3 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 

Repository files navigation

Coges

A Python library for implementing non-linear behavior and stream processing using a predicate-action pattern.

Beta Software - Proceed with caution

Features

  • 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+

Installation

pip install coges

Or with Poetry:

poetry add coges

Quick Start

import 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())

Core Concepts

Coge

A Coge is the fundamental building block - a combination of a predicate and an action:

  • Predicate: An async function returning bool that 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_result

Machine

The 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:

  1. Get next tick from async generator
  2. Run all predicates concurrently
  3. Filter coges where predicate returned True
  4. Run all active actions concurrently
  5. Update state with results

State

State is a dictionary passed to all predicates and actions. It contains:

KeyDescription
tickCurrent tick value from the generator
resultsDictionary of {coge_name: result} from last tick
active_cogesList 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)

Dependency Injection

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 yield runs on cleanup
  • Automatic injection: Matches parameter names to dependency names

API Reference

coges.coge

create_coge(name: str) -> Coge

Factory function to create a Coge.

coge = create_coge("my_coge")

Coge.predicate(fn: PredicateFn) -> None

Decorator to set the predicate function.

@coge.predicate
async def check(tick, **state) -> bool:
    return True

Coge.action(fn: ActionFn) -> None

Decorator to set the action function.

@coge.action
async def execute(tick, **state) -> Any:
    return result

coges.machine

create_machine(coges, tick_fn, di_resolver=identity, initial_state={}) -> Machine

Creates a machine that runs coges on each tick.

ParameterTypeDescription
cogeslist[Coge]List of coges to run
tick_fnAsyncGeneratorAsync generator yielding ticks
di_resolverCallableDI resolver function (optional)
initial_statedictInitial state values (optional)

Returns an async function that runs the machine.

coges.di

create_dependency_injector() -> tuple[add_dependency, resolve]

Creates a dependency injector, returning two functions:

  • add_dependency: Decorator to register dependencies
  • resolve: Function to resolve dependencies for a function

Examples

Chat Bot

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())

Data Pipeline

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())

Logging

Enable debug logging to see machine execution:

import logging

logging.basicConfig(level=logging.DEBUG)
logging.getLogger("coges.machine").setLevel(logging.DEBUG)

Changelog

0.9.0

  • tick_fn now requires an AsyncGenerator instead of a regular generator

0.8.0

  • Initial release

License

MIT - Jellyfish.tech

About

Library for implementing non-linear behaviour and stream processing

Resources

Stars

1 star

Watchers

1 watching

Forks

Releases

Packages

Used by

Contributors

Languages