Skip to content

Implement validated flow graphs with bounded runs and node execution #98

Description

@MasterOfBinary

Outcome

Deliver the first usable flow/ runtime and its concise API design in one cohesive change. The API must support ordinary Go functions, sequence, bounded parallel branches, dependency fan-in and conditional execution. Exported signatures are finalized with compilable examples, not copied from conversation pseudocode.

Required contract

  • Compile an immutable finite graph with stable node IDs and validate duplicate IDs, unknown dependencies, cycles, empty graphs and invalid limits before run admission.
  • Prefer a typed run input and application-defined typed node-result envelope; no heterogeneous runtime registry or named-port schema language.
  • Nodes read immutable input and declared dependency outputs and return owned results. A coordinator publishes immutable views and an explicit join combines them. Document deep ownership of maps/slices/pointers and concurrent use of node functions.
  • Execute a node at most once after its dependencies succeed. Failed or skipped prerequisites skip descendants with reasons; independent branches finish unless the run is canceled. Conditions evaluate once after dependencies, with false distinct from predicate failure.
  • Return every node's terminal success/failure/skipped/canceled outcome and the run summary; never infer success or absence from a zero value; explicit status governs. Recover node/predicate panics to explicit failures while releasing scheduler capacity.
  • Require finite admitted-run and global executing-node bounds per runner, plus maximum graph size. Reject saturated admission before allocating per-run state; no unbounded waiting-run queue or goroutine per waiting node.
  • Bound readiness bookkeeping by admitted runs × graph nodes; a bounded worker pool supplies node execution. Arbitrary user payload sizes remain caller-owned limits.
  • Cancellation prevents new dispatch and requests cancellation of active nodes. A run is complete only after started functions return. Resolve cancel/completion races without duplicate publication.
  • Shutdown closes admission atomically, grants a finite grace budget, then requests cancellation; report incomplete termination and retain accounting if functions ignore context. Caller-owned batchers/adapters outlive dependent runs and are not closed by the runner.
  • Expose minimal dependency-free lifecycle/node events with bounded callback behavior; stable report ordering/tie-breaking does not promise deterministic network timing or effects.

Verification

Race-enabled concurrent-run isolation; dependency and join order; false/error predicates and skipped descendants; malformed graph rejection; global capacity and readiness bounds under load; cancellation at each state; panic recovery; admission/close races and noncooperative-function termination reporting. Include ShitQuant-shaped fixture enrichment/proposal and a nontrading metadata/classification example. Follow normal repo docs/vet/lint checks.

Boundaries and dependencies

Full contract. Parent #97. Does not depend on #71, Redis, OnyxCore, a profitable strategy, or the stream release. Settle relevant lifecycle/observation concepts before implementation, but do not inherit stream collection/error-channel semantics. No automatic retries, optional-edge engine, dynamic graphs, scheduler or persistence in v1.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions