remove semver actions - #3
Merged
Merged
Conversation
GuanyiLi-Craig
approved these changes
Mar 28, 2025
GuanyiLi-Craig
added a commit
that referenced
this pull request
Jun 20, 2026
The per-invocation clone (Defect #3 isolation) guarded against two concurrent invoke() calls on the SAME workflow instance corrupting each other's topic queues/tracker. In current usage each request creates a fresh assistant instance, so instances are never reused concurrently and that bug cannot occur -- the machinery solved a future (instance-reuse) problem. Defer it to that refactor and remove the complexity now (YAGNI). Removed: - NodeBase.rebind_topics (only used by the clone). - EventDrivenWorkflow._spawn_run, _active_runs, the invoke/_invoke_impl split, stop() run-forwarding, and the init_workflow stop-guard. invoke() runs on self again with reset_stop_flag(), as before. - tests/workflow/test_invocation_isolation.py (tested the deferred behavior). Kept (all verified working on self): recovery-offset fix, fan-out per-consumer delivery accounting, the no-progress stuck-detection (AsyncOutputQueue + _progress_possible), span-error fix, EventCodec, provider DRY, and cleanups. 667 tests pass, mypy clean, pre-commit green. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
GuanyiLi-Craig
added a commit
that referenced
this pull request
Jun 20, 2026
…EventCodec, mypy→0) (#97) * Fix event-core correctness defects and make mypy clean Implements Phase 1 (correctness) of the codebase design refactoring plan and brings `mypy grafi` to zero errors. Full suite: 660 passed. Correctness fixes: - Recovery no longer skips an event: add restore_consumer() to the topic-queue port and use it in TopicBase.restore_topic instead of misusing fetch(offset+1), which over-advanced the consumed cursor. - Fan-out quiescence: the workflow counts one pending delivery per consumer (new EventDrivenWorkflow._topic_consumers), registered before an event becomes consumable, and rejects accounting underflow instead of silently clamping. Prevents premature termination when one topic feeds multiple subscribers. - Span error attributes are recorded while the span is still active (inside start_as_current_span), not after it has ended. Typing (mypy 56 -> 0 across 18 files): - Real fixes for ~51 (annotations, Optional, narrowing, SDK omit/Omit sentinel for absent tools=, AsyncStream cast, centralized postgres row->event helper, ParameterSchema.model_validate, builder Optional params). - 5 targeted, commented `# type: ignore` for unverifiable dynamic dispatch / SDK-stub / known-deferred-LSP cases (warn_unused_ignores keeps them honest). Behavior note: OpenAI-compatible and Claude tools now pass the SDK omit sentinel instead of None/NOT_GIVEN when no tools are present. Tests: +7 characterization tests (queue restore, span-before-end, fan-out); updated 3 LLM tool tests for the omit sentinel. Also adds the design audit & refactoring plan doc. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> * Phase 2/5/6/7/8: invocation isolation, provider DRY, EventCodec, cleanups Continues the codebase design refactoring plan. Phase 2 (Defect #3 — invocation isolation): - invoke() runs on a per-invocation clone (_spawn_run) with its own tracker, invoke queue, stop flag, and topic queues; nodes are rebound to the per-run topics (NodeBase.rebind_topics). Concurrent invocations of one workflow no longer share/reset each other's runtime state. stop() forwards to in-flight runs. +tests/workflow/test_invocation_isolation.py (failed pre-fix). Phase 5 (provider DRY): - New OpenAICompatibleTool base owns the shared prepare/invoke/convert logic; OpenAITool/DeepseekTool/OpenRouterTool shrink to ~55-77 lines of config (env var, model, base_url, extra headers, error label). Tests adjusted to patch the shared base and mock the client as an async context manager. Phase 6 (Open/Closed — EventCodec): - New grafi/common/events/event_codec.py owns event-type registration and decoding; EventStore delegates and no longer imports concrete event classes. New event types register on the codec without editing the store. +tests. Phase 7 (cleanups): - generate_manifest returns the written path (+test); removed cloudpickle dep, lockfile entry, and stale jsonpickle/cloudpickle mypy overrides ("no pickle dependency"); corrected the jsonpickle doc; removed dead Command.from_dict. Phase 8a/8b (interfaces/typing): - Mutable Field(default=[])/PrivateAttr(default={}) -> default_factory. - SubscriptionBuilder.build() raises on empty instead of returning None under its -> SubExpr annotation. All green: 668 tests pass, `mypy grafi` clean, pre-commit (black/isort/flake8). Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> * Phase 3/4/9: DI port, provider-neutral function spec, doc accuracy Continues the codebase design refactoring plan. Phase 3 (Dependency Inversion): - Add grafi/runtime/execution_services.py: ExecutionServices(event_store, tracer) port + default_execution_services() that reads the global container (kept as the default composition root) so injection can be adopted incrementally. +tests. - record_base.py no longer imports AssistantBase/Workflow/NodeBase/Tool (typed the decorated component as Any). grafi/common now has no runtime imports of the assistants/nodes/tools/workflows layers. Phase 4 (provider-neutral, safe subset): - Move FunctionSpec.to_openai_tool() to grafi/tools/llms/impl/openai_adapter.py (to_openai_tool(spec)). grafi/common/models/function_spec.py no longer imports the OpenAI SDK. Callers (openai_compatible, ollama, a test) use the adapter. NOTE: the full provider-neutral Message rewrite is intentionally deferred -- Message field types round-trip onto the provider wire (tool_calls etc.), so it needs live-provider integration testing and a deprecation window per the plan. Phase 9 (doc accuracy): - features.md: corrected the idempotency section -- event persistence/replay is idempotent, but node/tool execution is at-least-once (no exactly-once external side effects unless the tool is idempotent). - event-driven-workflow.md: replaced the non-existent MergeIdleQueue with the real AsyncOutputQueue. (A comprehensive doc refresh remains.) All green: 670 tests pass, `mypy grafi` clean, pre-commit (black/isort/flake8). Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> * Fix code-review findings: AND-subscription hang, OpenRouter, stop race + cleanups Addresses the recall code-review on the refactor branch. Correctness: - #1 (HIGH) Parallel workflow no longer hangs when a topic fans out to a firing node AND a node whose AND-subscription can never be satisfied. Per-consumer delivery counting is kept (fixes fan-out), but AsyncOutputQueue now ends iteration when the workflow reports no progress is possible (no active node, nothing consumable, no node can invoke) for two consecutive idle polls, so a parked delivery cannot hang the run. New EventDrivenWorkflow._progress_possible; regression test added (was a confirmed hang). - #2 OpenRouter (and any non-OpenAI-parse provider) no longer routes structured_output to client.beta.chat.completions.parse; new _supports_beta_parse=False keeps it on the standard create call. - #3 OpenRouterTool.from_dict restores `model` (was dropped, reverting to openrouter/auto on round-trip). - #4 stop() forwarded to an in-flight run is no longer undone: the per-run clone no longer calls reset_stop_flag(), and init_workflow skips tracker.reset() when a stop is already requested. - #14 OpenAI-compatible invoke merges request kwargs with chat_params into one dict, so a caller-supplied extra_headers/response_format overrides instead of raising a duplicate-keyword TypeError. Cleanup / robustness: - #10 base_url serialization moved to OpenAICompatibleTool.to_dict (deepseek no longer overrides; openrouter only adds extra_headers). - #8 _spawn_run rebuilds each per-run queue as type(topic.event_queue)() instead of hardcoding InMemTopicEventQueue, preserving a topic's queue kind. - #12 _topic_consumers is memoized (static topology; read per published event). - #9 removed the dead ExecutionServices module + test (unwired DI port). - #5/#6/#7 doc accuracy: tracker underflow (logs+clamps, backstopped by #1), EventStore._codec (shared process-wide default), and tool-statelessness note. 669 tests pass, mypy clean, pre-commit (black/isort/flake8) green. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> * Harden stuck-detection and per-run queue allocation (review follow-up) Second-pass review of the fix commit surfaced two low-severity fragilities; this hardens both. - AsyncOutputQueue stuck-detection no longer depends on the code-layout invariant that a node is never observable between consuming its inputs and calling tracker.enter(). It now also resets the idle counter whenever the tracker's activity count changes between polls (a node processed = real progress), so a busy run is never mistaken for stuck regardless of scheduling. Termination still requires consecutive idle polls with no progress possible. - _spawn_run allocates a fresh InMemTopicEventQueue per topic again, with a comment explaining why: a topic's event queue is a transient per-run buffer (durability/recovery come from the EventStore), so an in-memory queue is correct per run. This removes the prior type(queue)() form that silently required an arg-less constructor and dropped queue configuration. 669 tests pass, mypy clean, pre-commit green. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> * Simplify: drop redundant state in quiescence/kwargs paths Cleanup pass (no behavior change): - _progress_possible: drop the per-node can_invoke() sweep. A node can only invoke when a subscribed topic has an unconsumed event, so a zero pending-consumable count already implies no node can invoke -- the sweep was redundant (and a duplicate lock-acquiring pass over the same topic state). - Remove the _topic_consumers memoization cache (field + per-clone reset + lookup branch): the computation is a tiny pure function over static topology (dict.fromkeys + a type check), read once per published event; the cache added state for negligible savings. - on_messages_committed: collapse the if/else that duplicated the decrement into a single max(0, ...) with the underflow log kept as a tripwire. - openai_compatible.invoke: build call_kwargs directly instead of constructing a throwaway request_kwargs dict and immediately merging it. 669 tests pass, mypy clean, pre-commit green. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> * Revert premature per-invocation isolation The per-invocation clone (Defect #3 isolation) guarded against two concurrent invoke() calls on the SAME workflow instance corrupting each other's topic queues/tracker. In current usage each request creates a fresh assistant instance, so instances are never reused concurrently and that bug cannot occur -- the machinery solved a future (instance-reuse) problem. Defer it to that refactor and remove the complexity now (YAGNI). Removed: - NodeBase.rebind_topics (only used by the clone). - EventDrivenWorkflow._spawn_run, _active_runs, the invoke/_invoke_impl split, stop() run-forwarding, and the init_workflow stop-guard. invoke() runs on self again with reset_stop_flag(), as before. - tests/workflow/test_invocation_isolation.py (tested the deferred behavior). Kept (all verified working on self): recovery-offset fix, fan-out per-consumer delivery accounting, the no-progress stuck-detection (AsyncOutputQueue + _progress_possible), span-error fix, EventCodec, provider DRY, and cleanups. 667 tests pass, mypy clean, pre-commit green. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
GuanyiLi-Craig
added a commit
that referenced
this pull request
Jun 21, 2026
…lowRun) (#98) One EventDrivenWorkflow instance can now handle many invoke() calls concurrently. Previously all runtime state (quiescence tracker, ready-queue, stop flag, topic queues) lived on the shared definition and was reset on every invoke, so two concurrent invocations corrupted each other -- the second reset the first mid-flight and a single stop() halted both (Defect #3). Introduce WorkflowRun, a per-invocation object that owns all mutable runtime state. The definition is read-only during a run; each invoke() allocates a WorkflowRun with its own tracker, ready-queue, stop flag, and per-run topic queues (shallow topic copies with fresh queues, preserving queue kind). Tools, LLM clients, and topology are shared by reference -- no graph cloning. - event_driven_workflow.py: slim immutable definition/facade; invoke() runs on a WorkflowRun tracked in _active_runs; stop() forwards to in-flight runs. - workflow_run.py: WorkflowRun state + shared topology/commit/progress helpers. - run_recovery.py / sequential_engine.py / parallel_engine.py: execution split into focused modules (free functions taking the run; WorkflowRun referenced only under TYPE_CHECKING, so no import cycle). - utils.py: publish_events/get_node_input route I/O through the run's topics. Tests: - tests/workflow/test_invocation_isolation.py: concurrent invokes are isolated (parallel + sequential); fails on the pre-change runtime. - tests/workflow/test_workflow_run.py: 22 unit tests for WorkflowRun internals. - tests_integration: 3 concurrent invoke() examples (parallel, sequential, and a 3-node function-call workflow), plus .env.example. Verified against real providers. - Retargeted 4 coupled unit tests from definition-private fields to WorkflowRun. 692 unit tests pass; mypy grafi clean; pre-commit (black/isort/flake8) clean; full integration suite green (54 passed / 4 skipped). Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
No description provided.