add sources/ with some old files
This commit is contained in:
226
sources/git-arch-sources/260911A-work-writup.md
Normal file
226
sources/git-arch-sources/260911A-work-writup.md
Normal file
@@ -0,0 +1,226 @@
|
||||
|
||||
• Implemented durable, at-least-once processing for the voice-cloning worker. The central change is that an SQS message is no longer deleted before training begins.
|
||||
|
||||
## SQS visibility and acknowledgment
|
||||
|
||||
Previously, the worker deleted the message immediately after receiving it. A crash during download, training, MongoDB updates, or S3 upload permanently lost the job.
|
||||
|
||||
The new lifecycle is:
|
||||
|
||||
Receive message
|
||||
→ establish visibility lease
|
||||
→ renew lease during processing
|
||||
→ persist assets and completion state
|
||||
→ stop heartbeat
|
||||
→ delete message
|
||||
|
||||
On failure:
|
||||
|
||||
Processing error
|
||||
→ record error state where possible
|
||||
→ do not delete message
|
||||
→ set retry visibility delay
|
||||
→ SQS delivers it again later
|
||||
|
||||
On a hard crash:
|
||||
|
||||
Worker dies
|
||||
→ heartbeat stops
|
||||
→ latest visibility lease expires
|
||||
→ SQS redelivers the message
|
||||
|
||||
### Visibility heartbeat
|
||||
|
||||
The worker immediately extends a received message’s visibility to 300 seconds by default. It then renews that lease every 60 seconds while training runs.
|
||||
|
||||
Each renewal resets the remaining visibility window to 300 seconds; it does not add 300 seconds cumulatively. Therefore, if the worker crashes, the message becomes available no later
|
||||
than roughly five minutes after the last successful renewal.
|
||||
|
||||
The initial visibility extension must succeed before MongoDB or training work starts. Periodic renewal failures are reported, and the next heartbeat attempts another renewal.
|
||||
|
||||
The heartbeat is stopped before acknowledgment so there is no renewal racing with message deletion.
|
||||
|
||||
### Failure backoff
|
||||
|
||||
The worker requests ApproximateReceiveCount when receiving messages. Caught failures use that count to apply exponential visibility backoff:
|
||||
|
||||
Receive count Retry delay
|
||||
━━━━━━━━━━━━━━━ ━━━━━━━━━━━━━━━━━━━━━
|
||||
1 30 seconds
|
||||
─────────────── ─────────────────────
|
||||
2 60 seconds
|
||||
─────────────── ─────────────────────
|
||||
3 120 seconds
|
||||
─────────────── ─────────────────────
|
||||
4 240 seconds
|
||||
─────────────── ─────────────────────
|
||||
5 480 seconds
|
||||
─────────────── ─────────────────────
|
||||
6+ 900 seconds maximum
|
||||
|
||||
If changing visibility for the retry also fails, the message is still not acknowledged. It naturally reappears when its existing lease expires.
|
||||
|
||||
### Configurable visibility settings
|
||||
|
||||
The following environment variables were added:
|
||||
|
||||
- SQS_VISIBILITY_TIMEOUT_SECONDS — default 300
|
||||
- SQS_VISIBILITY_HEARTBEAT_INTERVAL_MS — default 60000
|
||||
- SQS_RETRY_VISIBILITY_BASE_SECONDS — default 30
|
||||
- SQS_RETRY_VISIBILITY_MAX_SECONDS — default 900
|
||||
|
||||
The worker rejects a configuration where the heartbeat interval is equal to or longer than the visibility timeout.
|
||||
|
||||
This provides at-least-once rather than exactly-once delivery. SQS can still deliver duplicates, so the processing path was also made idempotent.
|
||||
|
||||
## Durable completion and idempotent retries
|
||||
|
||||
Before doing work, the worker reads both the VoiceCloning record and its UserAudioProfile.
|
||||
|
||||
A job is considered fully complete only when:
|
||||
|
||||
- Both records have status: completed.
|
||||
- The profile contains all five local model paths.
|
||||
- The profile contains all five corresponding S3 paths.
|
||||
|
||||
The required assets are:
|
||||
|
||||
- Full voice model
|
||||
- Full model configuration
|
||||
- Speaker embeddings
|
||||
- Lightweight voice model
|
||||
- Lightweight model configuration
|
||||
|
||||
If all completion data already exists, a redelivered message skips training and is simply acknowledged.
|
||||
|
||||
For newly completed work, persistence now occurs in this order:
|
||||
|
||||
1. Verify all local model files exist.
|
||||
2. Upload all assets to S3.
|
||||
3. Update the user profile with local and S3 paths.
|
||||
4. Mark the user profile completed.
|
||||
5. Mark the voice-cloning record completed as the final commit marker.
|
||||
6. Delete the SQS message.
|
||||
|
||||
MongoDB updates are also checked for a returned record. If an update resolves with null, the message is not acknowledged.
|
||||
|
||||
If SQS deletion fails after completion, the completed states are preserved rather than changed to error. On redelivery, the worker recognizes completion, skips training, and retries
|
||||
only the acknowledgment.
|
||||
|
||||
## Recovery from partially completed jobs
|
||||
|
||||
The training pipeline now attempts to reuse durable work left behind by a crashed worker.
|
||||
|
||||
It first checks:
|
||||
|
||||
1. Model paths already stored on the user profile.
|
||||
2. Completed model artifacts under the job’s EFS output directory.
|
||||
|
||||
If all expected files exist, training is skipped. Existing S3 paths are also reused when they correspond to the same local asset map.
|
||||
|
||||
If only partial artifacts exist, the worker removes the job-scoped temporary dataset, archive, and incomplete model output before retrying. This prevents files such as a half-written
|
||||
speakers.pth or checkpoint from poisoning every subsequent delivery.
|
||||
|
||||
Logs remain outside the cleaned model output and are preserved across retries.
|
||||
|
||||
## MongoDB retry handling
|
||||
|
||||
The original recursive connection retry could leave the outer promise unresolved forever after an initial failure.
|
||||
|
||||
It was replaced with a bounded retry loop:
|
||||
|
||||
- Seven attempts by default.
|
||||
- Linear delay between attempts.
|
||||
- Proper rejection after exhaustion.
|
||||
- The final error retains the original connection failure as its cause.
|
||||
|
||||
Configuration:
|
||||
|
||||
- MONGO_CONNECT_MAX_ATTEMPTS — default 7
|
||||
- MONGO_CONNECT_RETRY_DELAY_MS — default 1000
|
||||
|
||||
MongoDB connections are closed only after a successful connection and closure errors are reported without hiding the processing result.
|
||||
|
||||
## Download and process error handling
|
||||
|
||||
The training pipeline was extracted into voice-cloning-job-handler/training_pipeline.js.
|
||||
|
||||
Audio downloads now handle:
|
||||
|
||||
- Non-2xx HTTP responses
|
||||
- Up to three redirects
|
||||
- Network errors
|
||||
- Stream/write failures
|
||||
- A 60-second timeout
|
||||
- Removal of partially downloaded files
|
||||
|
||||
Training commands now use execFile with argument arrays rather than interpolated shell command strings. This gives reliable exit-code handling and avoids shell interpretation of job-
|
||||
derived paths.
|
||||
|
||||
Command output is appended to timestamped stage logs. A non-zero child-process exit now reliably rejects the pipeline after stdout and stderr have been retained.
|
||||
|
||||
The generated model directory and all five expected output files are verified before the job can be completed.
|
||||
|
||||
## Job validation
|
||||
|
||||
Messages are validated before processing:
|
||||
|
||||
- Body must be valid JSON.
|
||||
- _doc, job ID, profile ID, metadata, and input are required.
|
||||
- Environment must be development, staging, or production.
|
||||
- Input cannot be empty.
|
||||
- Recording URLs must be valid HTTPS URLs.
|
||||
- Original transcript text must be present.
|
||||
- directoryName must be safe for filesystem paths.
|
||||
|
||||
Malformed messages are not deleted. They remain eligible for the queue’s retry and dead-letter behavior.
|
||||
|
||||
## Worker lifecycle
|
||||
|
||||
voice-cloning-job-handler/index.js is now primarily responsible for dependency and environment configuration.
|
||||
|
||||
The worker now:
|
||||
|
||||
- Handles SIGTERM and SIGINT.
|
||||
- Stops fetching new work after a shutdown request.
|
||||
- Allows the current job to finish when possible.
|
||||
- Sets a failure exit code if initialization unexpectedly stops.
|
||||
- Does not automatically start when imported by tests.
|
||||
|
||||
Errors are reported through Bugsnag, while failures in error reporting itself are contained.
|
||||
|
||||
## Shared SQS support
|
||||
|
||||
The shared SQS service (app/services/sqs/sqs_service.js) now:
|
||||
|
||||
- Requests ApproximateReceiveCount.
|
||||
- Exposes changeMessageVisibility.
|
||||
- Continues to expose receive, delete, and send operations.
|
||||
|
||||
## Testing and documentation
|
||||
|
||||
Added 18 automated tests covering:
|
||||
|
||||
- Acknowledgment ordering
|
||||
- Processing failure retention
|
||||
- Retry visibility backoff
|
||||
- Visibility heartbeat renewal
|
||||
- Initial visibility-extension failure
|
||||
- Safe completed-job redelivery
|
||||
- SQS deletion failure
|
||||
- Missing MongoDB update targets
|
||||
- Malformed messages
|
||||
- MongoDB retry success and exhaustion
|
||||
- EFS artifact recovery
|
||||
- Partial artifact cleanup
|
||||
- Full pipeline execution
|
||||
- Child-process failure logging
|
||||
- Input/path validation
|
||||
|
||||
npm test passes all 18 tests. Syntax and whitespace checks also pass.
|
||||
|
||||
The tests use mocked SQS/MongoDB/S3 integrations and temporary filesystem fixtures; no live AWS resources or GPU training run were invoked.
|
||||
|
||||
Finally, README.md documents the visibility and retry settings. The SQS queue should have an external redrive policy and dead-letter queue configured so permanently invalid jobs do
|
||||
not retry indefinitely.
|
||||
|
||||
Reference in New Issue
Block a user