Conversation
…lates for job materialization with simple unified template + variables approach
- Implemented Unified Payload Model with creation-time materialization.
- Updated DB schema (action_type, payload_template, field_map, and dropped old fields).
- Developed Template Engine for {{variable}} substitution.
- Updated API handlers and Cron Materializer to handle materialization logic.
- Simplified Executors to use static materialized payloads.
- Updated and verified all tests.
…r jobs scheduled to run now - Enhanced smoke test scripts with unified `payload_template` and `field_map` structure. - Modified `Upsert` in ScheduleBucketStore to reset FIRED buckets to PENDING.
…cution
1. Unified Payload Model:
* Transitioned to a "Snapshot" model where jobs are fully materialized at creation.
* Updated database schemas (Postgres and MongoDB) to use action_type, payload_template, and field_map.
* Dropped obsolete delivery_config, delivery_type, context, and payload_override columns.
2. Direct Dispatch Fast-path:
* Implemented logic to bypass the bucket system for jobs scheduled in the past or within a 5-second window.
* Updated the API ScheduleHandler and the CronMaterializer to produce Kafka events directly for these immediate jobs.
* Ensured future-dated jobs continue to use the bucket-based scheduler for efficiency.
3. Code Cleanup & Verification:
* Created a shared internal/executor/types package for consistent dispatch event structures.
* Updated all smoke test scripts and integration tests to align with the new API.
* Verified all functionality with a full run of smoke tests (Phase 5 through 9B).
… storage - Add PostgreSQL range partitioning for `bulk_items` table (migration 0024) - Transition Bulk Ingestor from eager to lazy materialization (storing raw variables) - Optimize ingestion with 1,000-row batch inserts to PostgreSQL - Implement logical batch fan-out (row_start/row_end) to Kafka bulk-records topic - Add BulkProcessor with in-process retries (3x) for 5xx and network errors - Implement deterministic batch idempotency keys for resilient dispatch - Integrate MongoDB for high-volume execution logs with 7-day TTL indexes - Add field mapping contract to template engine for execution-time resolution This architecture isolates high-volume bulk data from the main jobs table, reduces database disk usage by ~90%, and provides instantaneous cleanup via partition dropping.
…sing - Add `error_field` to job definitions for custom error extraction from response. - Implement `preparePayload` for batched payload preparation (single object for 1 item, array for multiple). - Update BulkProcessor to handle batch retries and finalize item results collectively. - Refactor PostgreSQL schema and methods to include `error_field`. - Enhance MongoDB audit logs with response details and latency tracking.
… flexible payload format
- Added ScaledObject for `bulk-executor` with Kafka-based scaling, supporting dynamic replica counts (0-20). - Added ScaledJob for `bulk-ingestor` with Redis-based scaling, supporting up to 50 parallel jobs. - Updated Tiltfile to deploy KEDA infrastructure and integrate scaled resources after CRDs are established. - Refactored bulk-ingestor and bulk-executor deployments to enable KEDA scaling (replica management and resource labels).
…dynamic job handling - Added Scenario A (full success) and Scenario B (partial failure) to thoroughly test bulk job execution. - Refactored dynamic job definition creation and completion detection logic with retry support. - Expanded CSV input variations to validate full and partial processing outcomes. - Improved result verification by checking error file generation for failed cases.
…n and JSON input scenario - Added Scenario C for testing bulk jobs with inline JSON array input (no file upload). - Enhanced Scenario B with WebSocket progress message validation and error CSV spot-checks. - Introduced dynamic CSV/JSON creation using `mktemp` for better test isolation. - Refactored WebSocket message handling to validate status, counts, and structure. - Updated error result CSV verification steps for better coverage and reliability. - Enabled detached context for result file generation to ensure MinIO uploads during pod shutdown.
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.