TT Lab
Get started
Learn Learning paths Courses

Irreversible Changes

The Alien Festival Cancellation Stopped at 25%

Continue in TT Lab

Goal

You cancel the 40 orders of the alien dessert festival 10 at a time. The approved list is frozen, and even if the process disappears midway, you resume from right after the chunk that was committed.

Why it matters

If you change a full rollback into per-chunk commits, the business contract changes too. You have to get the partial completion explicitly approved and check that the checkpoint matches the actual business change and the audit. First study Python functions and exceptions, SQL transactions, and the earlier approved-version and compensation labs. The expected time is 120 minutes, so extend with +time before it expires. Files disappear when the session ends. Keep the code you need separately.

Environment and common contract

The deliverable is /root/chunks/worker.py. PostgreSQL 16, psycopg 3.2.3, and Python 3 are in the image, and there is no runtime installation. You can write to /root as the postgres user, and no extra capability or user switch is needed.

The grader prepares the tables below and the fictional orders in a separate temporary schema of the local labdb, and cleans up only the schema it created. The student functions use the search_path and the DSN of the connection they are given. Do not change the public tables or hardcode the schema, customers, IDs, or the DSN. Pass SQL values as parameters.

CREATE TABLE orders(id integer PRIMARY KEY,tenant text NOT NULL,
 qty integer NOT NULL CHECK(qty BETWEEN 1 AND 1000),
 state text NOT NULL CHECK(state IN ('pending','paid','cancelled')),
 revision integer NOT NULL CHECK(revision>=0));
CREATE TABLE jobs(job_id text PRIMARY KEY,tenant text NOT NULL,targets jsonb NOT NULL,
 chunk_size integer NOT NULL CHECK(chunk_size BETWEEN 1 AND 10),
 next_index integer NOT NULL CHECK(next_index>=0));
CREATE TABLE job_audit(job_id text REFERENCES jobs(job_id),ordinal integer NOT NULL,
 id integer NOT NULL,previous_revision integer NOT NULL,new_revision integer NOT NULL,
 qty integer NOT NULL,PRIMARY KEY(job_id,ordinal),UNIQUE(job_id,id));

The identifiers job_id and tenant are exact str values of 1–64 ASCII letters, digits, underscores, and hyphens. items/targets is an exact list of 1–64 items, and each item is an exact dict that has only id, revision, and qty. id is an int 1–2147483647, revision is an int 0–2147483646, and qty is an int 1–1000. bool is not allowed as an integer. Duplicate IDs are rejected, and the list is normalized into a deep copy in ID order. chunk_size must be an exact int 1–10, and allow_partial must be exactly True. A wrong direct input is a ValueError before any write and is not corrected automatically.

The contract is that after registration, the tenant, targets, and chunk_size of jobs and job_audit are immutable. next_index is 0 or more and at most the number of targets, and unless it is the completion position, it must be a multiple of chunk_size. The stored targets must equal the normalized input. The audit, in ordinal order, must exactly equal [0:next_index] of the original approval array in ID, previous revision, new revision=previous+1, and qty. A leftover audit is also corruption. A later change to an original order is not corruption of the past audit, so you do not overwrite the past approval with the current values of the earlier chunks.

The section that run_chunk will process is [next_index:min(next_index+chunk_size, number of targets)]. You call fault with each string argument only when it exists. Do not hide hook errors. The lock and statement limits apply per SQL statement and are not a limit on the elapsed time of the whole chunk. After success or failure, restore the settings of the original connection.

The borrowed con is autocommit=True and Read Committed, and there is no open transaction at the start of an external call. Your functions do not close the connection and leave no transaction after success or failure. The internal nested calls of cancel_chunk must preserve the outer transaction. run_chunk calls the hook after the real commit, so do not wrap it in another outer transaction. Only chunk_file owns a connection.

Steps

  1. Freeze the partial completion approval — implement a Conflict that subclasses Exception and manifest(job_id,tenant,items,chunk_size,allow_partial). Validate the input contract below and return a new dict of job_id, tenant, targets, chunk_size, and allow_partial. targets is a new list of new dicts in ID order. Only an explicit True approval is allowed, and a wrong input is a ValueError.
  2. Do not overwrite the approval even when re-registering — register(con,plan) validates and normalizes an exact dict that has only manifest's five keys, then registers it in jobs with next_index=0 and returns True. The same customer, normalized targets, and chunk_size for the same job ID changes nothing and returns False, and different contents are a Conflict. An input error is a ValueError before any write, and it does not change orders or the audit.
  3. Read the approval and the progress separately — inspect_job(con,job_id) validates the ID and returns None if it does not exist, and otherwise a job_id, tenant, targets, chunk_size, and next_index dict. It is read-only, and modifying the returned value does not affect the original record.
  4. A single-row failure inside a chunk rolls everything back — cancel_chunk(con,tenant,items) validates the target input first and performs a conditional UPDATE in ID order on the current customer, ID, revision, qty, and pending. It changes to cancelled, raises the revision by 1, and returns a list of id, previous_revision, new_revision, and qty dicts. If even one row does not match, it rolls back the whole chunk with a Conflict. It does not write jobs or job_audit, and a nested call does not commit the outer transaction early.
  5. Commit the change, the audit, and the position together — run_chunk(con,job_id,fault=None) first sets its own transaction's lock limit of 500ms and statement limit of 2000ms and then locks and reads the jobs row. If it does not exist, or the stored approval, next_index, or audit differs from the contract below, it is a Conflict. After changing the next section with cancel_chunk, it calls fault(after-orders); after storing all of the section's ordinals, IDs, before and after versions, and quantities in job_audit, it calls fault(after-audit); after changing next_index to the end of the section, it calls fault(after-checkpoint); and after the real COMMIT, it calls fault(after-commit). The return value is processed (an ID-ordered list), next_index, and done (a bool). If it is already complete, processed=[] and done=True, and it does not call the hooks. An error before the commit rolls back only this chunk and preserves the earlier chunks. An error after the commit preserves the committed state and propagates the original error.
  6. Report the past completion and the present difference — report(con,job_id) is a Conflict if it does not exist, and otherwise returns a dict with ID-ordered lists for approved, committed, remaining, matching, drifted, and missing. approved is the original approval, committed is the part before next_index, and remaining is the rest. If the current row of a committed ID matches in customer, qty, cancelled, and the original revision+1, it is matching; if the ID is missing, it is missing; everything else is drifted. It looks up the progress and the current rows in a single SELECT and does not fix any data.
  7. Continue the remaining section after a client termination — chunk_file(dsn,job_id,fault=None) opens an owned connection with psycopg.connect(dsn,autocommit=True,connect_timeout=2), calls run_chunk, and returns the same result. It closes the connection on both success and failure and does not hide errors. The check covers the results of resuming after a real client termination at each of the four hook points.
  8. Advance only a limited number of times with two workers — drain(dsn,job_id,max_chunks) validates the ID and an exact int max_chunks of 1–64 before creating a connection. It calls chunk_file at most max_chunks times, but ends immediately if done=True. It returns calls (the number of calls), the list of the processed IDs that this call actually processed, joined together, the last next_index, and the last done. It propagates any error immediately and does no hidden retry. When two processes process the same job to the end, every approved ID must be changed exactly once.

Notes

Freeze the partial completion approval

Implement a Conflict that subclasses Exception and manifest(job_id,tenant,items,chunk_size,allow_partial). Validate the input contract below and return a new dict of job_id, tenant, targets, chunk_size, and allow_partial. targets is a new list of new dicts in ID order. Only an explicit True approval is allowed, and a wrong input is a ValueError.

The promise that all 40 are atomic and the promise of committing 10 at a time are different. Distinguish bool from int too.

Do not overwrite the approval even when re-registering

register(con,plan) validates and normalizes an exact dict that has only manifest's five keys, then registers it in jobs with next_index=0 and returns True. The same customer, normalized targets, and chunk_size for the same job ID changes nothing and returns False, and different contents are a Conflict. An input error is a ValueError before any write, and it does not change orders or the audit.

An upsert that overwrites the existing next_index with 0 is not a re-registration but a loss of progress.

Read the approval and the progress separately

inspect_job(con,job_id) validates the ID and returns None if it does not exist, and otherwise a job_id, tenant, targets, chunk_size, and next_index dict. It is read-only, and modifying the returned value does not affect the original record.

Do not confuse the approval array with the current orders. A registration that has not been executed yet is also a valid record.

A single-row failure inside a chunk rolls everything back

cancel_chunk(con,tenant,items) validates the target input first and performs a conditional UPDATE in ID order on the current customer, ID, revision, qty, and pending. It changes to cancelled, raises the revision by 1, and returns a list of id, previous_revision, new_revision, and qty dicts. If even one row does not match, it rolls back the whole chunk with a Conflict. It does not write jobs or job_audit, and a nested call does not commit the outer transaction early.

Even if the number of rows the UPDATE affected is 0, the SQL is a success. Judge whether having no returned row is a business conflict.

Commit the change, the audit, and the position together

run_chunk(con,job_id,fault=None) first sets its own transaction's lock limit of 500ms and statement limit of 2000ms and then locks and reads the jobs row. If it does not exist, or the stored approval, next_index, or audit differs from the contract below, it is a Conflict. After changing the next section with cancel_chunk, it calls fault(after-orders); after storing all of the section's ordinals, IDs, before and after versions, and quantities in job_audit, it calls fault(after-audit); after changing next_index to the end of the section, it calls fault(after-checkpoint); and after the real COMMIT, it calls fault(after-commit). The return value is processed (an ID-ordered list), next_index, and done (a bool). If it is already complete, processed=[] and done=True, and it does not call the hooks. An error before the commit rolls back only this chunk and preserves the earlier chunks. An error after the commit preserves the committed state and propagates the original error.

Compare the audit sequence numbers with the prefix set of the approval array. The job row lock is needed from before the selection of the next section.

Report the past completion and the present difference

report(con,job_id) is a Conflict if it does not exist, and otherwise returns a dict with ID-ordered lists for approved, committed, remaining, matching, drifted, and missing. approved is the original approval, committed is the part before next_index, and remaining is the rest. If the current row of a committed ID matches in customer, qty, cancelled, and the original revision+1, it is matching; if the ID is missing, it is missing; everything else is drifted. It looks up the progress and the current rows in a single SELECT and does not fix any data.

The same values of a different ID at present cannot stand in for a missing original ID. Advancing to the next chunk does not overwrite a later change to the earlier chunks.

Continue the remaining section after a client termination

chunk_file(dsn,job_id,fault=None) opens an owned connection with psycopg.connect(dsn,autocommit=True,connect_timeout=2), calls run_chunk, and returns the same result. It closes the connection on both success and failure and does not hide errors. The check covers the results of resuming after a real client termination at each of the four hook points.

If you lost the response after the commit, the next advance of the same job is the next chunk, not the chunk that has already finished.

Advance only a limited number of times with two workers

drain(dsn,job_id,max_chunks) validates the ID and an exact int max_chunks of 1–64 before creating a connection. It calls chunk_file at most max_chunks times, but ends immediately if done=True. It returns calls (the number of calls), the list of the processed IDs that this call actually processed, joined together, the last next_index, and the last done. It propagates any error immediately and does no hidden retry. When two processes process the same job to the end, every approved ID must be changed exactly once.

Compare the IDs returned by the different workers with the independent DB audit. An empty call made to receive the completion response is also included in calls.