From d35d6e53b7827205cfc7745d7380efbd39afc2a2 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sat, 29 Aug 2026 23:43:02 -0700 Subject: [PATCH 1/6] Add tests for @corbits/migration-runner advisory-lock guard Covers the shared runner that will replace six identical migration runners (bench, insights, preferences, inference-catalog, evals, access-policy): idempotent apply, two callers racing the same migration set, and recovery after a lock holder's connection dies without releasing it. --- bun.lock | 16 +- packages/migration-runner/LICENSE | 176 ++++++++++++++++++ packages/migration-runner/README.md | 49 +++++ packages/migration-runner/package.json | 23 +++ .../test/migration-runner.test.ts | 171 +++++++++++++++++ packages/migration-runner/tsconfig.json | 7 + 6 files changed, 440 insertions(+), 2 deletions(-) create mode 100644 packages/migration-runner/LICENSE create mode 100644 packages/migration-runner/README.md create mode 100644 packages/migration-runner/package.json create mode 100644 packages/migration-runner/test/migration-runner.test.ts create mode 100644 packages/migration-runner/tsconfig.json diff --git a/bun.lock b/bun.lock index cc2eb879b..2b596e798 100644 --- a/bun.lock +++ b/bun.lock @@ -1022,6 +1022,18 @@ "typescript": "catalog:", }, }, + "packages/migration-runner": { + "name": "@corbits/migration-runner", + "version": "0.0.1", + "dependencies": { + "@intx/log": "0.3.0", + "postgres": "catalog:", + }, + "devDependencies": { + "@types/bun": "catalog:", + "typescript": "catalog:", + }, + }, "packages/mocks": { "name": "@corbits/mocks", "version": "0.0.1", @@ -2227,6 +2239,8 @@ "@corbits/memory-tools": ["@corbits/memory-tools@workspace:packages/memory-tools"], + "@corbits/migration-runner": ["@corbits/migration-runner@workspace:packages/migration-runner"], + "@corbits/mocks": ["@corbits/mocks@workspace:packages/mocks"], "@corbits/morning-brief-workflow": ["@corbits/morning-brief-workflow@workspace:workflows/morning-brief"], @@ -3589,8 +3603,6 @@ "@typescript-eslint/eslint-plugin/ignore": ["ignore@7.0.6", "", {}, "sha512-BAg6QkE8W+TuQLrrw0Ugr7HegXduRuuj8/ti2kSOc+jz1dmx8/WNcjr6XGnq5YpDWxFwwaavqD0+jIUOKelTsw=="], - "@workbench/hub/@corbits/mailbox": ["@corbits/mailbox@github:corbitsdev/corbits-mailbox#caa5214", { "dependencies": { "@hono/standard-validator": "0.2.3", "@standard-community/standard-json": "0.3.5", "@standard-community/standard-openapi": "0.2.9", "arktype": "2.1.29", "hono-openapi": "1.3.1" }, "peerDependencies": { "@intx/log": "^0.2.2", "@intx/mime": "^0.2.2", "@intx/types": "^0.2.2", "drizzle-orm": "^0.45.2", "hono": "^4.12.0", "postgres": "^3.4.0" } }, "corbitsdev-corbits-mailbox-caa5214", "sha512-z8DRBFgA4ukM8p29COeaMjfKZYe5jAUF4OBMiaIQFuW592+DGD/y6Ws6SjGlXmR9azkHNWh8oTzjlWlRP24vsQ=="], - "ajv-formats/ajv": ["ajv@8.20.0", "", { "dependencies": { "fast-deep-equal": "^3.1.3", "fast-uri": "^3.0.1", "json-schema-traverse": "^1.0.0", "require-from-string": "^2.0.2" } }, "sha512-Thbli+OlOj+iMPYFBVBfJ3OmCAnaSyNn4M1vz9T6Gka5Jt9ba/HIR56joy65tY6kx/FCF5VXNB819Y7/GUrBGA=="], "better-call/@better-auth/utils": ["@better-auth/utils@0.5.0", "", { "dependencies": { "@noble/hashes": "^2.0.1" } }, "sha512-BL8W4EfIZFwlu0r54m3v1ztjDhu6dDe/amLTm0xybmbZaNgYUqhD3SjpAsnq0q8YD6/ki4iwIgxJNLP/N3TxiA=="], diff --git a/packages/migration-runner/LICENSE b/packages/migration-runner/LICENSE new file mode 100644 index 000000000..c6487f4fd --- /dev/null +++ b/packages/migration-runner/LICENSE @@ -0,0 +1,176 @@ +GNU LESSER GENERAL PUBLIC LICENSE + +Version 2.1, February 1999 + +Copyright (C) 1991, 1999 Free Software Foundation, Inc. +51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA + +Everyone is permitted to copy and distribute verbatim copies of this license document, but changing it is not allowed. + +[This is the first released version of the Lesser GPL. It also counts as the successor of the GNU Library Public License, version 2, hence the version number 2.1.] + +Preamble + +The licenses for most software are designed to take away your freedom to share and change it. By contrast, the GNU General Public Licenses are intended to guarantee your freedom to share and change free software--to make sure the software is free for all its users. + +This license, the Lesser General Public License, applies to some specially designated software packages--typically libraries--of the Free Software Foundation and other authors who decide to use it. You can use it too, but we suggest you first think carefully about whether this license or the ordinary General Public License is the better strategy to use in any particular case, based on the explanations below. + +When we speak of free software, we are referring to freedom of use, not price. Our General Public Licenses are designed to make sure that you have the freedom to distribute copies of free software (and charge for this service if you wish); that you receive source code or can get it if you want it; that you can change the software and use pieces of it in new free programs; and that you are informed that you can do these things. + +To protect your rights, we need to make restrictions that forbid distributors to deny you these rights or to ask you to surrender these rights. These restrictions translate to certain responsibilities for you if you distribute copies of the library or if you modify it. + +For example, if you distribute copies of the library, whether gratis or for a fee, you must give the recipients all the rights that we gave you. You must make sure that they, too, receive or can get the source code. If you link other code with the library, you must provide complete object files to the recipients, so that they can relink them with the library after making changes to the library and recompiling it. And you must show them these terms so they know their rights. + +We protect your rights with a two-step method: (1) we copyright the library, and (2) we offer you this license, which gives you legal permission to copy, distribute and/or modify the library. + +To protect each distributor, we want to make it very clear that there is no warranty for the free library. Also, if the library is modified by someone else and passed on, the recipients should know that what they have is not the original version, so that the original author's reputation will not be affected by problems that might be introduced by others. + +Finally, software patents pose a constant threat to the existence of any free program. We wish to make sure that a company cannot effectively restrict the users of a free program by obtaining a restrictive license from a patent holder. Therefore, we insist that any patent license obtained for a version of the library must be consistent with the full freedom of use specified in this license. + +Most GNU software, including some libraries, is covered by the ordinary GNU General Public License. This license, the GNU Lesser General Public License, applies to certain designated libraries, and is quite different from the ordinary General Public License. We use this license for certain libraries in order to permit linking those libraries into non-free programs. + +When a program is linked with a library, whether statically or using a shared library, the combination of the two is legally speaking a combined work, a derivative of the original library. The ordinary General Public License therefore permits such linking only if the entire combination fits its criteria of freedom. The Lesser General Public License permits more lax criteria for linking other code with the library. + +We call this license the "Lesser" General Public License because it does Less to protect the user's freedom than the ordinary General Public License. It also provides other free software developers Less of an advantage over competing non-free programs. These disadvantages are the reason we use the ordinary General Public License for many libraries. However, the Lesser license provides advantages in certain special circumstances. + +For example, on rare occasions, there may be a special need to encourage the widest possible use of a certain library, so that it becomes a de-facto standard. To achieve this, non-free programs must be allowed to use the library. A more frequent case is that a free library does the same job as widely used non-free libraries. In this case, there is little to gain by limiting the free library to free software only, so we use the Lesser General Public License. + +In other cases, permission to use a particular library in non-free programs enables a greater number of people to use a large body of free software. For example, permission to use the GNU C Library in non-free programs enables many more people to use the whole GNU operating system, as well as its variant, the GNU/Linux operating system. + +Although the Lesser General Public License is Less protective of the users' freedom, it does ensure that the user of a program that is linked with the Library has the freedom and the wherewithal to run that program using a modified version of the Library. + +The precise terms and conditions for copying, distribution and modification follow. Pay close attention to the difference between a "work based on the library" and a "work that uses the library". The former contains code derived from the library, whereas the latter must be combined with the library in order to run. + +GNU LESSER GENERAL PUBLIC LICENSE +TERMS AND CONDITIONS FOR COPYING, DISTRIBUTION AND MODIFICATION + +0. This License Agreement applies to any software library or other program which contains a notice placed by the copyright holder or other authorized party saying it may be distributed under the terms of this Lesser General Public License (also called "this License"). Each licensee is addressed as "you". + +A "library" means a collection of software functions and/or data prepared so as to be conveniently linked with application programs (which use some of those functions and data) to form executables. + +The "Library", below, refers to any such software library or work which has been distributed under these terms. A "work based on the Library" means either the Library or any derivative work under copyright law: that is to say, a work containing the Library or a portion of it, either verbatim or with modifications and/or translated straightforwardly into another language. (Hereinafter, translation is included without limitation in the term "modification".) + +"Source code" for a work means the preferred form of the work for making modifications to it. For a library, complete source code means all the source code for all modules it contains, plus any associated interface definition files, plus the scripts used to control compilation and installation of the library. + +Activities other than copying, distribution and modification are not covered by this License; they are outside its scope. The act of running a program using the Library is not restricted, and output from such a program is covered only if its contents constitute a work based on the Library (independent of the use of the Library in a tool for writing it). Whether that is true depends on what the Library does and what the program that uses the Library does. + +1. You may copy and distribute verbatim copies of the Library's complete source code as you receive it, in any medium, provided that you conspicuously and appropriately publish on each copy an appropriate copyright notice and disclaimer of warranty; keep intact all the notices that refer to this License and to the absence of any warranty; and distribute a copy of this License along with the Library. + +You may charge a fee for the physical act of transferring a copy, and you may at your option offer warranty protection in exchange for a fee. + +2. You may modify your copy or copies of the Library or any portion of it, thus forming a work based on the Library, and copy and distribute such modifications or work under the terms of Section 1 above, provided that you also meet all of these conditions: + + a) The modified work must itself be a software library. + + b) You must cause the files modified to carry prominent notices stating that you changed the files and the date of any change. + + c) You must cause the whole of the work to be licensed at no charge to all third parties under the terms of this License. + + d) If a facility in the modified Library refers to a function or a table of data to be supplied by an application program that uses the facility, other than as an argument passed when the facility is invoked, then you must make a good faith effort to ensure that, in the event an application does not supply such function or table, the facility still operates, and performs whatever part of its purpose remains meaningful. + +(For example, a function in a library to compute square roots has a purpose that is entirely well-defined independent of the application. Therefore, Subsection 2d requires that any application-supplied function or table used by this function must be optional: if the application does not supply it, the square root function must still compute square roots.) + +These requirements apply to the modified work as a whole. If identifiable sections of that work are not derived from the Library, and can be reasonably considered independent and separate works in themselves, then this License, and its terms, do not apply to those sections when you distribute them as separate works. But when you distribute the same sections as part of a whole which is a work based on the Library, the distribution of the whole must be on the terms of this License, whose permissions for other licensees extend to the entire whole, and thus to each and every part regardless of who wrote it. + +Thus, it is not the intent of this section to claim rights or contest your rights to work written entirely by you; rather, the intent is to exercise the right to control the distribution of derivative or collective works based on the Library. + +In addition, mere aggregation of another work not based on the Library with the Library (or with a work based on the Library) on a volume of a storage or distribution medium does not bring the other work under the scope of this License. + +3. You may opt to apply the terms of the ordinary GNU General Public License instead of this License to a given copy of the Library. To do this, you must alter all the notices that refer to this License, so that they refer to the ordinary GNU General Public License, version 2, instead of to this License. (If a newer version than version 2 of the ordinary GNU General Public License has appeared, then you can specify that version instead if you wish.) Do not make any other change in these notices. + +Once this change is made in a given copy, it is irreversible for that copy, so the ordinary GNU General Public License applies to all subsequent copies and derivative works made from that copy. + +This option is useful when you wish to copy part of the code of the Library into a program that is not a library. + +4. You may copy and distribute the Library (or a portion or derivative of it, under Section 2) in object code or executable form under the terms of Sections 1 and 2 above provided that you accompany it with the complete corresponding machine-readable source code, which must be distributed under the terms of Sections 1 and 2 above on a medium customarily used for software interchange. + +If distribution of object code is made by offering access to copy from a designated place, then offering equivalent access to copy the source code from the same place satisfies the requirement to distribute the source code, even though third parties are not compelled to copy the source along with the object code. + +5. A program that contains no derivative of any portion of the Library, but is designed to work with the Library by being compiled or linked with it, is called a "work that uses the Library". Such a work, in isolation, is not a derivative work of the Library, and therefore falls outside the scope of this License. + +However, linking a "work that uses the Library" with the Library creates an executable that is a derivative of the Library (because it contains portions of the Library), rather than a "work that uses the library". The executable is therefore covered by this License. Section 6 states terms for distribution of such executables. + +When a "work that uses the Library" uses material from a header file that is part of the Library, the object code for the work may be a derivative work of the Library even though the source code is not. Whether this is true is especially significant if the work can be linked without the Library, or if the work is itself a library. The threshold for this to be true is not precisely defined by law. + +If such an object file uses only numerical parameters, data structure layouts and accessors, and small macros and small inline functions (ten lines or less in length), then the use of the object file is unrestricted, regardless of whether it is legally a derivative work. (Executables containing this object code plus portions of the Library will still fall under Section 6.) + +Otherwise, if the work is a derivative of the Library, you may distribute the object code for the work under the terms of Section 6. Any executables containing that work also fall under Section 6, whether or not they are linked directly with the Library itself. + +6. As an exception to the Sections above, you may also combine or link a "work that uses the Library" with the Library to produce a work containing portions of the Library, and distribute that work under terms of your choice, provided that the terms permit modification of the work for the customer's own use and reverse engineering for debugging such modifications. + +You must give prominent notice with each copy of the work that the Library is used in it and that the Library and its use are covered by this License. You must supply a copy of this License. If the work during execution displays copyright notices, you must include the copyright notice for the Library among them, as well as a reference directing the user to the copy of this License. Also, you must do one of these things: + + a) Accompany the work with the complete corresponding machine-readable source code for the Library including whatever changes were used in the work (which must be distributed under Sections 1 and 2 above); and, if the work is an executable linked with the Library, with the complete machine-readable "work that uses the Library", as object code and/or source code, so that the user can modify the Library and then relink to produce a modified executable containing the modified Library. (It is understood that the user who changes the contents of definitions files in the Library will not necessarily be able to recompile the application to use the modified definitions.) + + b) Use a suitable shared library mechanism for linking with the Library. A suitable mechanism is one that (1) uses at run time a copy of the library already present on the user's computer system, rather than copying library functions into the executable, and (2) will operate properly with a modified version of the library, if the user installs one, as long as the modified version is interface-compatible with the version that the work was made with. + + c) Accompany the work with a written offer, valid for at least three years, to give the same user the materials specified in Subsection 6a, above, for a charge no more than the cost of performing this distribution. + + d) If distribution of the work is made by offering access to copy from a designated place, offer equivalent access to copy the above specified materials from the same place. + + e) Verify that the user has already received a copy of these materials or that you have already sent this user a copy. + +For an executable, the required form of the "work that uses the Library" must include any data and utility programs needed for reproducing the executable from it. However, as a special exception, the materials to be distributed need not include anything that is normally distributed (in either source or binary form) with the major components (compiler, kernel, and so on) of the operating system on which the executable runs, unless that component itself accompanies the executable. + +It may happen that this requirement contradicts the license restrictions of other proprietary libraries that do not normally accompany the operating system. Such a contradiction means you cannot use both them and the Library together in an executable that you distribute. + +7. You may place library facilities that are a work based on the Library side-by-side in a single library together with other library facilities not covered by this License, and distribute such a combined library, provided that the separate distribution of the work based on the Library and of the other library facilities is otherwise permitted, and provided that you do these two things: + + a) Accompany the combined library with a copy of the same work based on the Library, uncombined with any other library facilities. This must be distributed under the terms of the Sections above. + + b) Give prominent notice with the combined library of the fact that part of it is a work based on the Library, and explaining where to find the accompanying uncombined form of the same work. + +8. You may not copy, modify, sublicense, link with, or distribute the Library except as expressly provided under this License. Any attempt otherwise to copy, modify, sublicense, link with, or distribute the Library is void, and will automatically terminate your rights under this License. However, parties who have received copies, or rights, from you under this License will not have their licenses terminated so long as such parties remain in full compliance. + +9. You are not required to accept this License, since you have not signed it. However, nothing else grants you permission to modify or distribute the Library or its derivative works. These actions are prohibited by law if you do not accept this License. Therefore, by modifying or distributing the Library (or any work based on the Library), you indicate your acceptance of this License to do so, and all its terms and conditions for copying, distributing or modifying the Library or works based on it. + +10. Each time you redistribute the Library (or any work based on the Library), the recipient automatically receives a license from the original licensor to copy, distribute, link with or modify the Library subject to these terms and conditions. You may not impose any further restrictions on the recipients' exercise of the rights granted herein. You are not responsible for enforcing compliance by third parties with this License. + +11. If, as a consequence of a court judgment or allegation of patent infringement or for any other reason (not limited to patent issues), conditions are imposed on you (whether by court order, agreement or otherwise) that contradict the conditions of this License, they do not excuse you from the conditions of this License. If you cannot distribute so as to satisfy simultaneously your obligations under this License and any other pertinent obligations, then as a consequence you may not distribute the Library at all. For example, if a patent license would not permit royalty-free redistribution of the Library by all those who receive copies directly or indirectly through you, then the only way you could satisfy both it and this License would be to refrain entirely from distribution of the Library. + +If any portion of this section is held invalid or unenforceable under any particular circumstance, the balance of the section is intended to apply, and the section as a whole is intended to apply in other circumstances. + +It is not the purpose of this section to induce you to infringe any patents or other property right claims or to contest validity of any such claims; this section has the sole purpose of protecting the integrity of the free software distribution system which is implemented by public license practices. Many people have made generous contributions to the wide range of software distributed through that system in reliance on consistent application of that system; it is up to the author/donor to decide if he or she is willing to distribute software through any other system and a licensee cannot impose that choice. + +This section is intended to make thoroughly clear what is believed to be a consequence of the rest of this License. + +12. If the distribution and/or use of the Library is restricted in certain countries either by patents or by copyrighted interfaces, the original copyright holder who places the Library under this License may add an explicit geographical distribution limitation excluding those countries, so that distribution is permitted only in or among countries not thus excluded. In such case, this License incorporates the limitation as if written in the body of this License. + +13. The Free Software Foundation may publish revised and/or new versions of the Lesser General Public License from time to time. Such new versions will be similar in spirit to the present version, but may differ in detail to address new problems or concerns. + +Each version is given a distinguishing version number. If the Library specifies a version number of this License which applies to it and "any later version", you have the option of following the terms and conditions either of that version or of any later version published by the Free Software Foundation. If the Library does not specify a license version number, you may choose any version ever published by the Free Software Foundation. + +14. If you wish to incorporate parts of the Library into other free programs whose distribution conditions are incompatible with these, write to the author to ask for permission. For software which is copyrighted by the Free Software Foundation, write to the Free Software Foundation; we sometimes make exceptions for this. Our decision will be guided by the two goals of preserving the free status of all derivatives of our free software and of promoting the sharing and reuse of software generally. + +NO WARRANTY + +15. BECAUSE THE LIBRARY IS LICENSED FREE OF CHARGE, THERE IS NO WARRANTY FOR THE LIBRARY, TO THE EXTENT PERMITTED BY APPLICABLE LAW. EXCEPT WHEN OTHERWISE STATED IN WRITING THE COPYRIGHT HOLDERS AND/OR OTHER PARTIES PROVIDE THE LIBRARY "AS IS" WITHOUT WARRANTY OF ANY KIND, EITHER EXPRESSED OR IMPLIED, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE. THE ENTIRE RISK AS TO THE QUALITY AND PERFORMANCE OF THE LIBRARY IS WITH YOU. SHOULD THE LIBRARY PROVE DEFECTIVE, YOU ASSUME THE COST OF ALL NECESSARY SERVICING, REPAIR OR CORRECTION. + +16. IN NO EVENT UNLESS REQUIRED BY APPLICABLE LAW OR AGREED TO IN WRITING WILL ANY COPYRIGHT HOLDER, OR ANY OTHER PARTY WHO MAY MODIFY AND/OR REDISTRIBUTE THE LIBRARY AS PERMITTED ABOVE, BE LIABLE TO YOU FOR DAMAGES, INCLUDING ANY GENERAL, SPECIAL, INCIDENTAL OR CONSEQUENTIAL DAMAGES ARISING OUT OF THE USE OR INABILITY TO USE THE LIBRARY (INCLUDING BUT NOT LIMITED TO LOSS OF DATA OR DATA BEING RENDERED INACCURATE OR LOSSES SUSTAINED BY YOU OR THIRD PARTIES OR A FAILURE OF THE LIBRARY TO OPERATE WITH ANY OTHER SOFTWARE), EVEN IF SUCH HOLDER OR OTHER PARTY HAS BEEN ADVISED OF THE POSSIBILITY OF SUCH DAMAGES. + +END OF TERMS AND CONDITIONS + +How to Apply These Terms to Your New Libraries + +If you develop a new library, and you want it to be of the greatest possible use to the public, we recommend making it free software that everyone can redistribute and change. You can do so by permitting redistribution under these terms (or, alternatively, under the terms of the ordinary General Public License). + +To apply these terms, attach the following notices to the library. It is safest to attach them to the start of each source file to most effectively convey the exclusion of warranty; and each file should have at least the "copyright" line and a pointer to where the full notice is found. + + one line to give the library's name and an idea of what it does. + Copyright (C) year name of author + + This library is free software; you can redistribute it and/or modify it under the terms of the GNU Lesser General Public License as published by the Free Software Foundation; either version 2.1 of the License, or (at your option) any later version. + + This library is distributed in the hope that it will be useful, but WITHOUT ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU Lesser General Public License for more details. + + You should have received a copy of the GNU Lesser General Public License along with this library; if not, write to the Free Software Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA Also add information on how to contact you by electronic and paper mail. + +You should also get your employer (if you work as a programmer) or your school, if any, to sign a "copyright disclaimer" for the library, if necessary. Here is a sample; alter the names: + +Yoyodyne, Inc., hereby disclaims all copyright interest in +the library `Frob' (a library for tweaking knobs) written +by James Random Hacker. + +signature of Ty Coon, 1 April 1990 +Ty Coon, President of Vice +That's all there is to it! diff --git a/packages/migration-runner/README.md b/packages/migration-runner/README.md new file mode 100644 index 000000000..e28073cc0 --- /dev/null +++ b/packages/migration-runner/README.md @@ -0,0 +1,49 @@ +# @corbits/migration-runner + +Shared mechanics for the "transactional, self-contained" package migration +shape documented in [docs/package-migrations.md](../../docs/package-migrations.md): +schema and ledger bootstrap, a per-migration `sql.begin` transaction, and a +session-level Postgres advisory lock held around the whole run. + +## Why the lock + +Every caller opens its own single-connection (`max: 1`) client and calls +`applyPackageMigrations` at hub boot. With two hub replicas booting at once, +both would otherwise `SELECT` the ledger before either had inserted a row, +both would see "not yet applied," and the loser's ledger `INSERT` would +violate the ledger table's primary key and crash that replica's boot. +`applyPackageMigrations` takes `pg_advisory_lock(hashtext(ledgerTable))` +immediately after connecting and holds it — across every migration's own +transaction — until the run finishes, releasing it in a `finally` before the +connection closes. The loser blocks until the winner finishes, then re-checks +the ledger and finds every migration already applied. Postgres also releases +a session-level advisory lock automatically if the holding connection dies, +so a crashed replica never wedges the lock for the next boot. + +The lock key is derived from the ledger table name, which is already unique +per package (`bench_migrations`, `insights_migrations`, ...), so distinct +packages never contend on the same lock. + +## What it does not change + +Each migration still applies and records itself in its own transaction, and +a failure still wraps and rethrows as `` `${packageLabel} migration ${name} +failed: ...` `` with the original error as `cause` — this package only closes +the race around the check-and-insert, not the per-migration failure +semantics packages already relied on. + +## Usage + +```ts +import { applyPackageMigrations } from "@corbits/migration-runner"; + +export async function applyBenchMigrations(databaseUrl: string) { + return applyPackageMigrations({ + databaseUrl, + schema: "bench", + ledgerTable: "bench_migrations", + migrations: benchMigrations, + packageLabel: "bench", + }); +} +``` diff --git a/packages/migration-runner/package.json b/packages/migration-runner/package.json new file mode 100644 index 000000000..3aff768d7 --- /dev/null +++ b/packages/migration-runner/package.json @@ -0,0 +1,23 @@ +{ + "name": "@corbits/migration-runner", + "private": true, + "description": "Shared runner for package-owned Postgres migrations: schema/ledger bootstrap, per-migration transactional apply, and a session-level advisory lock around the whole run so concurrent hub replicas race-check the ledger safely instead of crashing on a duplicate insert.", + "version": "0.0.1", + "license": "LGPL-2.1-or-later", + "type": "module", + "exports": { + ".": "./src/index.ts" + }, + "scripts": { + "typecheck": "tsc --noEmit", + "test": "bun test" + }, + "dependencies": { + "@intx/log": "0.3.0", + "postgres": "catalog:" + }, + "devDependencies": { + "@types/bun": "catalog:", + "typescript": "catalog:" + } +} diff --git a/packages/migration-runner/test/migration-runner.test.ts b/packages/migration-runner/test/migration-runner.test.ts new file mode 100644 index 000000000..8a9188183 --- /dev/null +++ b/packages/migration-runner/test/migration-runner.test.ts @@ -0,0 +1,171 @@ +// DB-gated: skipped when no DATABASE_URL is reachable (a fresh checkout +// still runs the unit gates), mirroring the six package migration test +// suites this runner replaces the mechanics of. Runs against its own +// scratch database, never the developer's or the walking-skeleton suite's. +import { afterAll, beforeAll, describe, expect, test } from "bun:test"; +import postgres from "postgres"; + +import { e2eDatabaseUrl } from "../../../scripts/e2e/harness"; +import { + applyPackageMigrations, + type PackageMigration, +} from "../src/index"; + +function scratchUrlFor(e2eUrl: string): string { + const url = new URL(e2eUrl); + const database = url.pathname.replace(/^\//, ""); + url.pathname = `/${database}_migration_runner_test`; + return url.toString(); +} + +const databaseUrl = e2eDatabaseUrl(); +const describeIfDb = databaseUrl === undefined ? describe.skip : describe; + +const SCHEMA = "migration_runner_test"; +const LEDGER_TABLE = "migration_runner_test_migrations"; + +const migrations: readonly PackageMigration[] = [ + { + name: "0001_widget", + sql: ` + CREATE TABLE IF NOT EXISTS "${SCHEMA}"."widget" ( + "id" text PRIMARY KEY + ); + `, + }, + { + name: "0002_widget_label", + sql: ` + ALTER TABLE "${SCHEMA}"."widget" + ADD COLUMN IF NOT EXISTS "label" text; + `, + }, +]; + +describeIfDb("applyPackageMigrations", () => { + const scratchUrl = scratchUrlFor( + databaseUrl ?? "postgres://localhost:5432/unused", + ); + const scratchDatabase = new URL(scratchUrl).pathname.replace(/^\//, ""); + + async function withMaintenanceConnection( + run: (sql: postgres.Sql) => Promise, + ): Promise { + const maintenanceUrl = new URL(scratchUrl); + maintenanceUrl.pathname = "/postgres"; + const maintenance = postgres(maintenanceUrl.toString(), { + max: 1, + onnotice: () => undefined, + }); + try { + return await run(maintenance); + } finally { + await maintenance.end(); + } + } + + beforeAll(async () => { + await withMaintenanceConnection(async (maintenance) => { + await maintenance.unsafe( + `DROP DATABASE IF EXISTS "${scratchDatabase}"`, + ); + await maintenance.unsafe(`CREATE DATABASE "${scratchDatabase}"`); + }); + }, 20000); + + afterAll(async () => { + await withMaintenanceConnection(async (maintenance) => { + await maintenance.unsafe( + `DROP DATABASE IF EXISTS "${scratchDatabase}"`, + ); + }); + }, 20000); + + test("applies migrations into the package schema and is idempotent on a second run", async () => { + const first = await applyPackageMigrations({ + databaseUrl: scratchUrl, + schema: SCHEMA, + ledgerTable: LEDGER_TABLE, + migrations, + packageLabel: "migration-runner-test", + }); + expect(first.applied).toEqual(["0001_widget", "0002_widget_label"]); + expect(first.alreadyApplied).toEqual([]); + + const second = await applyPackageMigrations({ + databaseUrl: scratchUrl, + schema: SCHEMA, + ledgerTable: LEDGER_TABLE, + migrations, + packageLabel: "migration-runner-test", + }); + expect(second.applied).toEqual([]); + expect(second.alreadyApplied.sort()).toEqual( + ["0001_widget", "0002_widget_label"].sort(), + ); + }); + + test("two replicas booting concurrently both complete without either crashing on a duplicate ledger insert", async () => { + const concurrentSchema = `${SCHEMA}_concurrent`; + const concurrentLedger = `${LEDGER_TABLE}_concurrent`; + const options = { + databaseUrl: scratchUrl, + schema: concurrentSchema, + ledgerTable: concurrentLedger, + migrations, + packageLabel: "migration-runner-concurrent-test", + }; + + const [first, second] = await Promise.all([ + applyPackageMigrations(options), + applyPackageMigrations(options), + ]); + + const appliedNames = [...first.applied, ...second.applied].sort(); + const alreadyAppliedNames = [ + ...first.alreadyApplied, + ...second.alreadyApplied, + ].sort(); + + expect(new Set(appliedNames).size).toBe(appliedNames.length); + expect([...appliedNames, ...alreadyAppliedNames].sort()).toEqual( + [...migrations, ...migrations] + .map((migration) => migration.name) + .sort(), + ); + + const sql = postgres(scratchUrl, { max: 1, onnotice: () => undefined }); + try { + const ledgerRows = await sql.unsafe( + `SELECT name FROM "${concurrentSchema}"."${concurrentLedger}" ORDER BY name`, + ); + expect(ledgerRows.map((row) => String(row["name"]))).toEqual([ + "0001_widget", + "0002_widget_label", + ]); + } finally { + await sql.end(); + } + }); + + test("recovers the lock after the holding connection dies without releasing it", async () => { + const crashSchema = `${SCHEMA}_crash`; + const crashLedger = `${LEDGER_TABLE}_crash`; + + const holder = postgres(scratchUrl, { max: 1, onnotice: () => undefined }); + await holder.unsafe(`SELECT pg_advisory_lock(hashtext($1)::bigint)`, [ + crashLedger, + ]); + await holder.end({ timeout: 0 }); + + const result = await applyPackageMigrations({ + databaseUrl: scratchUrl, + schema: crashSchema, + ledgerTable: crashLedger, + migrations, + packageLabel: "migration-runner-crash-test", + }); + + expect(result.applied).toEqual(["0001_widget", "0002_widget_label"]); + }, 10000); +}); diff --git a/packages/migration-runner/tsconfig.json b/packages/migration-runner/tsconfig.json new file mode 100644 index 000000000..e956ddd88 --- /dev/null +++ b/packages/migration-runner/tsconfig.json @@ -0,0 +1,7 @@ +{ + "extends": "../../tsconfig.base.json", + "compilerOptions": { + "types": ["bun"] + }, + "include": ["src", "test"] +} From 3a2c7c6ab3bba71a5a8f9d96e8f83438a2ab0b7b Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sat, 29 Aug 2026 23:55:05 -0700 Subject: [PATCH 2/6] migration-runner: shared advisory-lock guard for package migrations A session-level pg_advisory_lock held across the check-and-insert closes the race where two hub replicas booting together both see a migration as unapplied and one crashes on the ledger's primary-key violation. Six package migration runners delegate to this next. --- packages/migration-runner/src/index.ts | 100 +++++++++++++++++++++++++ 1 file changed, 100 insertions(+) create mode 100644 packages/migration-runner/src/index.ts diff --git a/packages/migration-runner/src/index.ts b/packages/migration-runner/src/index.ts new file mode 100644 index 000000000..62a8c4145 --- /dev/null +++ b/packages/migration-runner/src/index.ts @@ -0,0 +1,100 @@ +// Shared mechanics for the "transactional, self-contained" package +// migration shape in docs/package-migrations.md: schema/ledger bootstrap, +// per-migration transactional apply, and a session-level advisory lock +// held around the whole run so concurrent hub replicas can't both see a +// migration as unapplied and race its ledger insert. See README.md for +// why the lock is session-level rather than per-transaction. +import { getLogger } from "@intx/log"; +import postgres from "postgres"; + +const log = getLogger(["migration-runner"]); + +export interface PackageMigration { + name: string; + sql: string; +} + +export interface ApplyPackageMigrationsReport { + applied: string[]; + alreadyApplied: string[]; +} + +export interface ApplyPackageMigrationsOptions { + databaseUrl: string; + schema: string; + ledgerTable: string; + migrations: readonly PackageMigration[]; + packageLabel: string; +} + +function quoteIdentifier(name: string): string { + return `"${name.replace(/"/g, '""')}"`; +} + +function quoteQualified(schema: string, name: string): string { + return `${quoteIdentifier(schema)}.${quoteIdentifier(name)}`; +} + +export async function applyPackageMigrations( + options: ApplyPackageMigrationsOptions, +): Promise { + const { databaseUrl, schema, ledgerTable, migrations, packageLabel } = + options; + const sql = postgres(databaseUrl, { max: 1, onnotice: () => undefined }); + let holdsLock = false; + try { + log.info`waiting for the ${packageLabel} migration lock (${ledgerTable})`; + await sql.unsafe(`SELECT pg_advisory_lock(hashtext($1)::bigint)`, [ + ledgerTable, + ]); + holdsLock = true; + + await sql.unsafe(`CREATE SCHEMA IF NOT EXISTS ${quoteIdentifier(schema)}`); + await sql.unsafe( + `CREATE TABLE IF NOT EXISTS ${quoteQualified(schema, ledgerTable)} (` + + `name text PRIMARY KEY, applied_at timestamptz NOT NULL DEFAULT now())`, + ); + + const applied: string[] = []; + const alreadyApplied: string[] = []; + + for (const migration of migrations) { + const existing = await sql.unsafe( + `SELECT 1 FROM ${quoteQualified(schema, ledgerTable)} WHERE name = $1`, + [migration.name], + ); + if (existing.length > 0) { + alreadyApplied.push(migration.name); + continue; + } + try { + await sql.begin(async (tx) => { + await tx.unsafe(migration.sql); + await tx.unsafe( + `INSERT INTO ${quoteQualified(schema, ledgerTable)} (name) VALUES ($1)`, + [migration.name], + ); + }); + applied.push(migration.name); + } catch (err) { + const message = err instanceof Error ? err.message : String(err); + throw new Error( + `${packageLabel} migration ${migration.name} failed: ${message}`, + { cause: err }, + ); + } + } + + return { applied, alreadyApplied }; + } finally { + try { + if (holdsLock) { + await sql.unsafe(`SELECT pg_advisory_unlock(hashtext($1)::bigint)`, [ + ledgerTable, + ]); + } + } finally { + await sql.end({ timeout: 5 }); + } + } +} From eceae36cb02f508f038a55e911592c58fcc0e7a8 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 04:42:32 -0700 Subject: [PATCH 3/6] migration-runner: adopt shared advisory-lock guard in six packages bench, insights, preferences, inference-catalog, evals, and access-policy each carried their own copy of the migration-apply loop without a lock, so two hub replicas booting together could both see a migration as unapplied and race the ledger insert. Point all six at @corbits/migration-runner's applyPackageMigrations instead of their own copy; schema and ledger table names are unchanged so existing ledgers keep working. --- bun.lock | 22 +++++ packages/access-policy/package.json | 1 + packages/access-policy/src/migrations.ts | 79 ++++-------------- packages/bench/package.json | 1 + packages/bench/src/migrations.ts | 81 +++++------------- packages/evals/package.json | 1 + packages/evals/src/store/migrations.ts | 83 +++++-------------- packages/inference-catalog/package.json | 1 + packages/inference-catalog/src/migrations.ts | 79 ++++-------------- packages/insights/package.json | 1 + packages/insights/src/migrations.ts | 79 ++++-------------- .../test/migration-runner.test.ts | 17 +--- packages/preferences/package.json | 1 + packages/preferences/src/migrations.ts | 81 +++++------------- 14 files changed, 143 insertions(+), 384 deletions(-) diff --git a/bun.lock b/bun.lock index 2b596e798..e94fd8011 100644 --- a/bun.lock +++ b/bun.lock @@ -192,6 +192,7 @@ "name": "@workbench/access-policy", "version": "0.0.1", "dependencies": { + "@corbits/migration-runner": "workspace:*", "@intx/hub-api": "workspace:*", "@workbench/hub-client": "workspace:*", "arktype": "catalog:", @@ -378,6 +379,7 @@ "version": "0.0.1", "dependencies": { "@corbits/api-query": "workspace:*", + "@corbits/migration-runner": "workspace:*", "@intx/hub-api": "workspace:*", "arktype": "catalog:", "drizzle-orm": "catalog:", @@ -687,6 +689,7 @@ "dependencies": { "@corbits/agent-directory": "workspace:*", "@corbits/assistant-workflow": "workspace:*", + "@corbits/migration-runner": "workspace:*", "@corbits/routines": "workspace:*", "@corbits/webhook-triggers": "workspace:*", "@corbits/workflow-catalog": "workspace:*", @@ -863,6 +866,7 @@ "name": "@corbits/inference-catalog", "version": "0.0.1", "dependencies": { + "@corbits/migration-runner": "workspace:*", "@corbits/ollama-adapter": "workspace:*", "@intx/db": "workspace:*", "@intx/hub-api": "workspace:*", @@ -895,6 +899,7 @@ "name": "@corbits/insights", "version": "0.0.1", "dependencies": { + "@corbits/migration-runner": "workspace:*", "@corbits/provider-pricing": "workspace:*", "@intx/db": "workspace:*", "@intx/hub-api": "workspace:*", @@ -1123,6 +1128,7 @@ "version": "0.0.1", "dependencies": { "@corbits/api-query": "workspace:*", + "@corbits/migration-runner": "workspace:*", "@intx/hub-api": "workspace:*", "arktype": "catalog:", "drizzle-orm": "catalog:", @@ -3581,6 +3587,18 @@ "@babel/helper-compilation-targets/semver": ["semver@6.3.1", "", { "bin": { "semver": "bin/semver.js" } }, "sha512-BR7VvDCVHO+q2xBEWskxS6DJE1qRnb7DxzUrogb71CWoSficBxYsiAGd+Kl0mmq/MprG9yArRkyrQxTO6XjMzA=="], + "@corbits/artifact-ui/@corbits/react-ui": ["@corbits/react-ui@github:corbitsdev/react-ui#3b12281", { "dependencies": { "class-variance-authority": "^0.7.1", "clsx": "^2.1.1", "tailwind-merge": "^3.3.1" }, "peerDependencies": { "@radix-ui/react-dialog": "^1.1.15", "@radix-ui/react-dropdown-menu": "^2.1.16", "@radix-ui/react-slot": "^1.2.3", "@radix-ui/react-tooltip": "^1.2.8", "@tanstack/react-query": "^5.90.2", "lucide-react": "^0.545.0 || ^1.0.0", "react": "^18.2.0 || ^19.0.0", "react-dom": "^18.2.0 || ^19.0.0", "sonner": "^2.0.7" }, "optionalPeers": ["@tanstack/react-query"] }, "corbitsdev-react-ui-3b12281", "sha512-Abvm/DO0Gqg0ITHGT9355ZxyKRPMVJLSSQSjpd3a8qt4JPrSMOLIOS4sX8ZMNNaArIbnY9F+VKrOWkUJUyO4Nw=="], + + "@corbits/chat-ui/@corbits/react-ui": ["@corbits/react-ui@github:corbitsdev/react-ui#3b12281", { "dependencies": { "class-variance-authority": "^0.7.1", "clsx": "^2.1.1", "tailwind-merge": "^3.3.1" }, "peerDependencies": { "@radix-ui/react-dialog": "^1.1.15", "@radix-ui/react-dropdown-menu": "^2.1.16", "@radix-ui/react-slot": "^1.2.3", "@radix-ui/react-tooltip": "^1.2.8", "@tanstack/react-query": "^5.90.2", "lucide-react": "^0.545.0 || ^1.0.0", "react": "^18.2.0 || ^19.0.0", "react-dom": "^18.2.0 || ^19.0.0", "sonner": "^2.0.7" }, "optionalPeers": ["@tanstack/react-query"] }, "corbitsdev-react-ui-3b12281", "sha512-Abvm/DO0Gqg0ITHGT9355ZxyKRPMVJLSSQSjpd3a8qt4JPrSMOLIOS4sX8ZMNNaArIbnY9F+VKrOWkUJUyO4Nw=="], + + "@corbits/context-menu/@corbits/react-ui": ["@corbits/react-ui@github:corbitsdev/react-ui#3b12281", { "dependencies": { "class-variance-authority": "^0.7.1", "clsx": "^2.1.1", "tailwind-merge": "^3.3.1" }, "peerDependencies": { "@radix-ui/react-dialog": "^1.1.15", "@radix-ui/react-dropdown-menu": "^2.1.16", "@radix-ui/react-slot": "^1.2.3", "@radix-ui/react-tooltip": "^1.2.8", "@tanstack/react-query": "^5.90.2", "lucide-react": "^0.545.0 || ^1.0.0", "react": "^18.2.0 || ^19.0.0", "react-dom": "^18.2.0 || ^19.0.0", "sonner": "^2.0.7" }, "optionalPeers": ["@tanstack/react-query"] }, "corbitsdev-react-ui-3b12281", "sha512-Abvm/DO0Gqg0ITHGT9355ZxyKRPMVJLSSQSjpd3a8qt4JPrSMOLIOS4sX8ZMNNaArIbnY9F+VKrOWkUJUyO4Nw=="], + + "@corbits/memory-hub/@corbits/memory": ["@corbits/memory@github:corbitsdev/corbits-memory#9e6f213", { "dependencies": { "@intx/agent": "0.2.2", "@intx/authz": "0.2.2", "@intx/hub-api": "0.2.2", "@intx/log": "0.2.2", "@intx/workflow": "0.2.2", "arktype": "^2.1.29", "drizzle-orm": "^0.45.1", "hono": "^4.9.0", "hono-openapi": "^1.3.1", "postgres": "^3.4.7" } }, "corbitsdev-corbits-memory-9e6f213", "sha512-utnM4ZT2zmslcPXYWAAqxlDNLcpGsXFiTOtj8h7+OXnhCP0Eaw8yl25+yCTyHpvt3jcdeG4h5uFsSj7ou0BZCA=="], + + "@corbits/plugins-ui/@corbits/react-ui": ["@corbits/react-ui@github:corbitsdev/react-ui#3b12281", { "dependencies": { "class-variance-authority": "^0.7.1", "clsx": "^2.1.1", "tailwind-merge": "^3.3.1" }, "peerDependencies": { "@radix-ui/react-dialog": "^1.1.15", "@radix-ui/react-dropdown-menu": "^2.1.16", "@radix-ui/react-slot": "^1.2.3", "@radix-ui/react-tooltip": "^1.2.8", "@tanstack/react-query": "^5.90.2", "lucide-react": "^0.545.0 || ^1.0.0", "react": "^18.2.0 || ^19.0.0", "react-dom": "^18.2.0 || ^19.0.0", "sonner": "^2.0.7" }, "optionalPeers": ["@tanstack/react-query"] }, "corbitsdev-react-ui-3b12281", "sha512-Abvm/DO0Gqg0ITHGT9355ZxyKRPMVJLSSQSjpd3a8qt4JPrSMOLIOS4sX8ZMNNaArIbnY9F+VKrOWkUJUyO4Nw=="], + + "@corbits/settings-ui/@corbits/react-ui": ["@corbits/react-ui@github:corbitsdev/react-ui#3b12281", { "dependencies": { "class-variance-authority": "^0.7.1", "clsx": "^2.1.1", "tailwind-merge": "^3.3.1" }, "peerDependencies": { "@radix-ui/react-dialog": "^1.1.15", "@radix-ui/react-dropdown-menu": "^2.1.16", "@radix-ui/react-slot": "^1.2.3", "@radix-ui/react-tooltip": "^1.2.8", "@tanstack/react-query": "^5.90.2", "lucide-react": "^0.545.0 || ^1.0.0", "react": "^18.2.0 || ^19.0.0", "react-dom": "^18.2.0 || ^19.0.0", "sonner": "^2.0.7" }, "optionalPeers": ["@tanstack/react-query"] }, "corbitsdev-react-ui-3b12281", "sha512-Abvm/DO0Gqg0ITHGT9355ZxyKRPMVJLSSQSjpd3a8qt4JPrSMOLIOS4sX8ZMNNaArIbnY9F+VKrOWkUJUyO4Nw=="], + "@esbuild-kit/core-utils/esbuild": ["esbuild@0.18.20", "", { "optionalDependencies": { "@esbuild/android-arm": "0.18.20", "@esbuild/android-arm64": "0.18.20", "@esbuild/android-x64": "0.18.20", "@esbuild/darwin-arm64": "0.18.20", "@esbuild/darwin-x64": "0.18.20", "@esbuild/freebsd-arm64": "0.18.20", "@esbuild/freebsd-x64": "0.18.20", "@esbuild/linux-arm": "0.18.20", "@esbuild/linux-arm64": "0.18.20", "@esbuild/linux-ia32": "0.18.20", "@esbuild/linux-loong64": "0.18.20", "@esbuild/linux-mips64el": "0.18.20", "@esbuild/linux-ppc64": "0.18.20", "@esbuild/linux-riscv64": "0.18.20", "@esbuild/linux-s390x": "0.18.20", "@esbuild/linux-x64": "0.18.20", "@esbuild/netbsd-x64": "0.18.20", "@esbuild/openbsd-x64": "0.18.20", "@esbuild/sunos-x64": "0.18.20", "@esbuild/win32-arm64": "0.18.20", "@esbuild/win32-ia32": "0.18.20", "@esbuild/win32-x64": "0.18.20" }, "bin": { "esbuild": "bin/esbuild" } }, "sha512-ceqxoedUrcayh7Y7ZX6NdbbDzGROiyVBgC4PriJThBKSVPWnnFHZAkfI1lJT8QFkOwH4qOS2SJkS4wvpGl8BpA=="], "@eslint-community/eslint-utils/eslint-visitor-keys": ["eslint-visitor-keys@3.4.3", "", {}, "sha512-wpc+LXeiyiisxPlEkUzU6svyS1frIO3Mgxj1fdy7Pm8Ygzguax2N3Fa/D/ag1WqbOprdI+uY6wMUl8/a2G+iag=="], @@ -3603,6 +3621,10 @@ "@typescript-eslint/eslint-plugin/ignore": ["ignore@7.0.6", "", {}, "sha512-BAg6QkE8W+TuQLrrw0Ugr7HegXduRuuj8/ti2kSOc+jz1dmx8/WNcjr6XGnq5YpDWxFwwaavqD0+jIUOKelTsw=="], + "@workbench/hub/@corbits/memory": ["@corbits/memory@github:corbitsdev/corbits-memory#9e6f213", { "dependencies": { "@intx/agent": "0.2.2", "@intx/authz": "0.2.2", "@intx/hub-api": "0.2.2", "@intx/log": "0.2.2", "@intx/workflow": "0.2.2", "arktype": "^2.1.29", "drizzle-orm": "^0.45.1", "hono": "^4.9.0", "hono-openapi": "^1.3.1", "postgres": "^3.4.7" } }, "corbitsdev-corbits-memory-9e6f213", "sha512-utnM4ZT2zmslcPXYWAAqxlDNLcpGsXFiTOtj8h7+OXnhCP0Eaw8yl25+yCTyHpvt3jcdeG4h5uFsSj7ou0BZCA=="], + + "@workbench/web/@corbits/react-ui": ["@corbits/react-ui@github:corbitsdev/react-ui#3b12281", { "dependencies": { "class-variance-authority": "^0.7.1", "clsx": "^2.1.1", "tailwind-merge": "^3.3.1" }, "peerDependencies": { "@radix-ui/react-dialog": "^1.1.15", "@radix-ui/react-dropdown-menu": "^2.1.16", "@radix-ui/react-slot": "^1.2.3", "@radix-ui/react-tooltip": "^1.2.8", "@tanstack/react-query": "^5.90.2", "lucide-react": "^0.545.0 || ^1.0.0", "react": "^18.2.0 || ^19.0.0", "react-dom": "^18.2.0 || ^19.0.0", "sonner": "^2.0.7" }, "optionalPeers": ["@tanstack/react-query"] }, "corbitsdev-react-ui-3b12281", "sha512-Abvm/DO0Gqg0ITHGT9355ZxyKRPMVJLSSQSjpd3a8qt4JPrSMOLIOS4sX8ZMNNaArIbnY9F+VKrOWkUJUyO4Nw=="], + "ajv-formats/ajv": ["ajv@8.20.0", "", { "dependencies": { "fast-deep-equal": "^3.1.3", "fast-uri": "^3.0.1", "json-schema-traverse": "^1.0.0", "require-from-string": "^2.0.2" } }, "sha512-Thbli+OlOj+iMPYFBVBfJ3OmCAnaSyNn4M1vz9T6Gka5Jt9ba/HIR56joy65tY6kx/FCF5VXNB819Y7/GUrBGA=="], "better-call/@better-auth/utils": ["@better-auth/utils@0.5.0", "", { "dependencies": { "@noble/hashes": "^2.0.1" } }, "sha512-BL8W4EfIZFwlu0r54m3v1ztjDhu6dDe/amLTm0xybmbZaNgYUqhD3SjpAsnq0q8YD6/ki4iwIgxJNLP/N3TxiA=="], diff --git a/packages/access-policy/package.json b/packages/access-policy/package.json index cf976fb0f..ff4e9b5f1 100644 --- a/packages/access-policy/package.json +++ b/packages/access-policy/package.json @@ -14,6 +14,7 @@ "test": "bun test" }, "dependencies": { + "@corbits/migration-runner": "workspace:*", "@intx/hub-api": "workspace:*", "@workbench/hub-client": "workspace:*", "arktype": "catalog:", diff --git a/packages/access-policy/src/migrations.ts b/packages/access-policy/src/migrations.ts index 60114d2f1..9d54af3f7 100644 --- a/packages/access-policy/src/migrations.ts +++ b/packages/access-policy/src/migrations.ts @@ -3,13 +3,17 @@ // disentangling history from the platform drizzle journal. Every table // this package owns — including its ledger — lives in its own // `access_policy` Postgres schema, never `public`; see -// docs/package-migrations.md. -import postgres from "postgres"; +// docs/package-migrations.md. Mechanics (schema/ledger bootstrap, +// transactional apply, the advisory lock across concurrent hub +// replicas) live in @corbits/migration-runner — this file owns only +// the domain SQL. +import { + applyPackageMigrations, + type ApplyPackageMigrationsReport, + type PackageMigration, +} from "@corbits/migration-runner"; -export interface AccessPolicyMigration { - name: string; - sql: string; -} +export type AccessPolicyMigration = PackageMigration; const SCHEMA = "access_policy"; @@ -56,63 +60,16 @@ export const accessPolicyMigrations: readonly AccessPolicyMigration[] = [ const LEDGER_TABLE = "access_policy_migrations"; -function quoteIdentifier(name: string): string { - return `"${name.replace(/"/g, '""')}"`; -} - -function quoteQualified(schema: string, name: string): string { - return `${quoteIdentifier(schema)}.${quoteIdentifier(name)}`; -} - -export interface ApplyAccessPolicyMigrationsReport { - applied: string[]; - alreadyApplied: string[]; -} +export type ApplyAccessPolicyMigrationsReport = ApplyPackageMigrationsReport; export async function applyAccessPolicyMigrations( databaseUrl: string, ): Promise { - const sql = postgres(databaseUrl, { max: 1, onnotice: () => undefined }); - try { - await sql.unsafe(`CREATE SCHEMA IF NOT EXISTS ${quoteIdentifier(SCHEMA)}`); - - await sql.unsafe( - `CREATE TABLE IF NOT EXISTS ${quoteQualified(SCHEMA, LEDGER_TABLE)} (` + - `name text PRIMARY KEY, applied_at timestamptz NOT NULL DEFAULT now())`, - ); - - const applied: string[] = []; - const alreadyApplied: string[] = []; - - for (const migration of accessPolicyMigrations) { - const existing = await sql.unsafe( - `SELECT 1 FROM ${quoteQualified(SCHEMA, LEDGER_TABLE)} WHERE name = $1`, - [migration.name], - ); - if (existing.length > 0) { - alreadyApplied.push(migration.name); - continue; - } - try { - await sql.begin(async (tx) => { - await tx.unsafe(migration.sql); - await tx.unsafe( - `INSERT INTO ${quoteQualified(SCHEMA, LEDGER_TABLE)} (name) VALUES ($1)`, - [migration.name], - ); - }); - applied.push(migration.name); - } catch (err) { - const message = err instanceof Error ? err.message : String(err); - throw new Error( - `access_policy migration ${migration.name} failed: ${message}`, - { cause: err }, - ); - } - } - - return { applied, alreadyApplied }; - } finally { - await sql.end({ timeout: 5 }); - } + return applyPackageMigrations({ + databaseUrl, + schema: SCHEMA, + ledgerTable: LEDGER_TABLE, + migrations: accessPolicyMigrations, + packageLabel: "access_policy", + }); } diff --git a/packages/bench/package.json b/packages/bench/package.json index c66e4992d..133a7ded4 100644 --- a/packages/bench/package.json +++ b/packages/bench/package.json @@ -16,6 +16,7 @@ }, "dependencies": { "@corbits/api-query": "workspace:*", + "@corbits/migration-runner": "workspace:*", "@intx/hub-api": "workspace:*", "arktype": "catalog:", "drizzle-orm": "catalog:", diff --git a/packages/bench/src/migrations.ts b/packages/bench/src/migrations.ts index e03770d10..8289a976f 100644 --- a/packages/bench/src/migrations.ts +++ b/packages/bench/src/migrations.ts @@ -2,13 +2,17 @@ // ledger table so the package can be extracted without disentangling // history from the platform drizzle journal. Every table this package // owns — including its ledger — lives in its own `bench` Postgres -// schema, never `public`; see docs/package-migrations.md. -import postgres from "postgres"; - -export interface BenchMigration { - name: string; - sql: string; -} +// schema, never `public`; see docs/package-migrations.md. Mechanics +// (schema/ledger bootstrap, transactional apply, the advisory lock +// across concurrent hub replicas) live in @corbits/migration-runner — +// this file owns only the domain SQL. +import { + applyPackageMigrations, + type ApplyPackageMigrationsReport, + type PackageMigration, +} from "@corbits/migration-runner"; + +export type BenchMigration = PackageMigration; const SCHEMA = "bench"; @@ -28,63 +32,16 @@ export const benchMigrations: readonly BenchMigration[] = [ const LEDGER_TABLE = "bench_migrations"; -function quoteIdentifier(name: string): string { - return `"${name.replace(/"/g, '""')}"`; -} - -function quoteQualified(schema: string, name: string): string { - return `${quoteIdentifier(schema)}.${quoteIdentifier(name)}`; -} - -export interface ApplyBenchMigrationsReport { - applied: string[]; - alreadyApplied: string[]; -} +export type ApplyBenchMigrationsReport = ApplyPackageMigrationsReport; export async function applyBenchMigrations( databaseUrl: string, ): Promise { - const sql = postgres(databaseUrl, { max: 1, onnotice: () => undefined }); - try { - await sql.unsafe(`CREATE SCHEMA IF NOT EXISTS ${quoteIdentifier(SCHEMA)}`); - - await sql.unsafe( - `CREATE TABLE IF NOT EXISTS ${quoteQualified(SCHEMA, LEDGER_TABLE)} (` + - `name text PRIMARY KEY, applied_at timestamptz NOT NULL DEFAULT now())`, - ); - - const applied: string[] = []; - const alreadyApplied: string[] = []; - - for (const migration of benchMigrations) { - const existing = await sql.unsafe( - `SELECT 1 FROM ${quoteQualified(SCHEMA, LEDGER_TABLE)} WHERE name = $1`, - [migration.name], - ); - if (existing.length > 0) { - alreadyApplied.push(migration.name); - continue; - } - try { - await sql.begin(async (tx) => { - await tx.unsafe(migration.sql); - await tx.unsafe( - `INSERT INTO ${quoteQualified(SCHEMA, LEDGER_TABLE)} (name) VALUES ($1)`, - [migration.name], - ); - }); - applied.push(migration.name); - } catch (err) { - const message = err instanceof Error ? err.message : String(err); - throw new Error( - `bench migration ${migration.name} failed: ${message}`, - { cause: err }, - ); - } - } - - return { applied, alreadyApplied }; - } finally { - await sql.end({ timeout: 5 }); - } + return applyPackageMigrations({ + databaseUrl, + schema: SCHEMA, + ledgerTable: LEDGER_TABLE, + migrations: benchMigrations, + packageLabel: "bench", + }); } diff --git a/packages/evals/package.json b/packages/evals/package.json index 3e62e5c51..85a46e4c8 100644 --- a/packages/evals/package.json +++ b/packages/evals/package.json @@ -15,6 +15,7 @@ "dependencies": { "@corbits/agent-directory": "workspace:*", "@corbits/assistant-workflow": "workspace:*", + "@corbits/migration-runner": "workspace:*", "@corbits/routines": "workspace:*", "@corbits/webhook-triggers": "workspace:*", "@corbits/workflow-catalog": "workspace:*", diff --git a/packages/evals/src/store/migrations.ts b/packages/evals/src/store/migrations.ts index a46513d6d..89abed7d8 100644 --- a/packages/evals/src/store/migrations.ts +++ b/packages/evals/src/store/migrations.ts @@ -3,13 +3,17 @@ // disentangling history from the platform drizzle journal. Every // table this package owns — including its ledger — lives in its own // `evals` Postgres schema, never `public`; see -// docs/package-migrations.md. -import postgres from "postgres"; - -export interface EvalsMigration { - name: string; - sql: string; -} +// docs/package-migrations.md. Mechanics (schema/ledger bootstrap, +// transactional apply, the advisory lock across concurrent hub +// replicas) live in @corbits/migration-runner — this file owns only +// the domain SQL. +import { + applyPackageMigrations, + type ApplyPackageMigrationsReport, + type PackageMigration, +} from "@corbits/migration-runner"; + +export type EvalsMigration = PackageMigration; const SCHEMA = "evals"; @@ -34,65 +38,16 @@ export const evalsMigrations: readonly EvalsMigration[] = [ const LEDGER_TABLE = "evals_migrations"; -function quoteIdentifier(name: string): string { - return `"${name.replace(/"/g, '""')}"`; -} - -function quoteQualified(schema: string, name: string): string { - return `${quoteIdentifier(schema)}.${quoteIdentifier(name)}`; -} - -export interface ApplyEvalsMigrationsReport { - applied: string[]; - alreadyApplied: string[]; -} +export type ApplyEvalsMigrationsReport = ApplyPackageMigrationsReport; export async function applyEvalsMigrations( databaseUrl: string, ): Promise { - const sql = postgres(databaseUrl, { max: 1, onnotice: () => undefined }); - try { - await sql.unsafe(`CREATE SCHEMA IF NOT EXISTS ${quoteIdentifier(SCHEMA)}`); - - await sql.unsafe( - `CREATE TABLE IF NOT EXISTS ${quoteQualified(SCHEMA, LEDGER_TABLE)} (` + - `name text PRIMARY KEY, applied_at timestamptz NOT NULL DEFAULT now())`, - ); - - const applied: string[] = []; - const alreadyApplied: string[] = []; - - for (const migration of evalsMigrations) { - const existing = await sql.unsafe( - `SELECT 1 FROM ${quoteQualified(SCHEMA, LEDGER_TABLE)} WHERE name = $1`, - [migration.name], - ); - if (existing.length > 0) { - alreadyApplied.push(migration.name); - continue; - } - try { - await sql.begin(async (tx) => { - await tx.unsafe(migration.sql); - await tx.unsafe( - `INSERT INTO ${quoteQualified(SCHEMA, LEDGER_TABLE)} (name) VALUES ($1)`, - [migration.name], - ); - }); - applied.push(migration.name); - } catch (err) { - const message = err instanceof Error ? err.message : String(err); - throw new Error( - `evals migration ${migration.name} failed: ${message}`, - { - cause: err, - }, - ); - } - } - - return { applied, alreadyApplied }; - } finally { - await sql.end({ timeout: 5 }); - } + return applyPackageMigrations({ + databaseUrl, + schema: SCHEMA, + ledgerTable: LEDGER_TABLE, + migrations: evalsMigrations, + packageLabel: "evals", + }); } diff --git a/packages/inference-catalog/package.json b/packages/inference-catalog/package.json index e0840a392..2320979c2 100644 --- a/packages/inference-catalog/package.json +++ b/packages/inference-catalog/package.json @@ -18,6 +18,7 @@ "test": "bun test" }, "dependencies": { + "@corbits/migration-runner": "workspace:*", "@corbits/ollama-adapter": "workspace:*", "@intx/db": "workspace:*", "@intx/hub-api": "workspace:*", diff --git a/packages/inference-catalog/src/migrations.ts b/packages/inference-catalog/src/migrations.ts index a454dbcc1..b9db427b6 100644 --- a/packages/inference-catalog/src/migrations.ts +++ b/packages/inference-catalog/src/migrations.ts @@ -1,12 +1,15 @@ // Package-owned migrations for @corbits/inference-catalog. Own ledger, own // Postgres schema, literal SQL, transactional apply — the convention in -// docs/package-migrations.md. -import postgres from "postgres"; +// docs/package-migrations.md. Mechanics (schema/ledger bootstrap, +// transactional apply, the advisory lock across concurrent hub replicas) +// live in @corbits/migration-runner — this file owns only the domain SQL. +import { + applyPackageMigrations, + type ApplyPackageMigrationsReport, + type PackageMigration, +} from "@corbits/migration-runner"; -export interface InferenceCatalogMigration { - name: string; - sql: string; -} +export type InferenceCatalogMigration = PackageMigration; const SCHEMA = "inference_catalog"; @@ -32,63 +35,17 @@ export const inferenceCatalogMigrations: readonly InferenceCatalogMigration[] = const LEDGER_TABLE = "inference_catalog_migrations"; -function quoteIdentifier(name: string): string { - return `"${name.replace(/"/g, '""')}"`; -} - -function quoteQualified(schema: string, name: string): string { - return `${quoteIdentifier(schema)}.${quoteIdentifier(name)}`; -} - -export interface ApplyInferenceCatalogMigrationsReport { - applied: string[]; - alreadyApplied: string[]; -} +export type ApplyInferenceCatalogMigrationsReport = + ApplyPackageMigrationsReport; export async function applyInferenceCatalogMigrations( databaseUrl: string, ): Promise { - const sql = postgres(databaseUrl, { max: 1, onnotice: () => undefined }); - try { - await sql.unsafe(`CREATE SCHEMA IF NOT EXISTS ${quoteIdentifier(SCHEMA)}`); - - await sql.unsafe( - `CREATE TABLE IF NOT EXISTS ${quoteQualified(SCHEMA, LEDGER_TABLE)} (` + - `name text PRIMARY KEY, applied_at timestamptz NOT NULL DEFAULT now())`, - ); - - const applied: string[] = []; - const alreadyApplied: string[] = []; - - for (const migration of inferenceCatalogMigrations) { - const existing = await sql.unsafe( - `SELECT 1 FROM ${quoteQualified(SCHEMA, LEDGER_TABLE)} WHERE name = $1`, - [migration.name], - ); - if (existing.length > 0) { - alreadyApplied.push(migration.name); - continue; - } - try { - await sql.begin(async (tx) => { - await tx.unsafe(migration.sql); - await tx.unsafe( - `INSERT INTO ${quoteQualified(SCHEMA, LEDGER_TABLE)} (name) VALUES ($1)`, - [migration.name], - ); - }); - applied.push(migration.name); - } catch (err) { - const message = err instanceof Error ? err.message : String(err); - throw new Error( - `inference-catalog migration ${migration.name} failed: ${message}`, - { cause: err }, - ); - } - } - - return { applied, alreadyApplied }; - } finally { - await sql.end({ timeout: 5 }); - } + return applyPackageMigrations({ + databaseUrl, + schema: SCHEMA, + ledgerTable: LEDGER_TABLE, + migrations: inferenceCatalogMigrations, + packageLabel: "inference-catalog", + }); } diff --git a/packages/insights/package.json b/packages/insights/package.json index 41b5a0e03..6e461d28c 100644 --- a/packages/insights/package.json +++ b/packages/insights/package.json @@ -15,6 +15,7 @@ "test": "bun test" }, "dependencies": { + "@corbits/migration-runner": "workspace:*", "@corbits/provider-pricing": "workspace:*", "@intx/db": "workspace:*", "@intx/hub-api": "workspace:*", diff --git a/packages/insights/src/migrations.ts b/packages/insights/src/migrations.ts index d3e5d2273..1db8ff36d 100644 --- a/packages/insights/src/migrations.ts +++ b/packages/insights/src/migrations.ts @@ -2,13 +2,17 @@ // own ledger table so the package can be extracted without disentangling // history from the platform drizzle journal. Every table this package // owns — including its ledger — lives in its own `insights` Postgres -// schema, never `public`; see docs/package-migrations.md. -import postgres from "postgres"; +// schema, never `public`; see docs/package-migrations.md. Mechanics +// (schema/ledger bootstrap, transactional apply, the advisory lock +// across concurrent hub replicas) live in @corbits/migration-runner — +// this file owns only the domain SQL. +import { + applyPackageMigrations, + type ApplyPackageMigrationsReport, + type PackageMigration, +} from "@corbits/migration-runner"; -export interface InsightsMigration { - name: string; - sql: string; -} +export type InsightsMigration = PackageMigration; const SCHEMA = "insights"; @@ -83,63 +87,16 @@ export const insightsMigrations: readonly InsightsMigration[] = [ const LEDGER_TABLE = "insights_migrations"; -function quoteIdentifier(name: string): string { - return `"${name.replace(/"/g, '""')}"`; -} - -function quoteQualified(schema: string, name: string): string { - return `${quoteIdentifier(schema)}.${quoteIdentifier(name)}`; -} - -export interface ApplyInsightsMigrationsReport { - applied: string[]; - alreadyApplied: string[]; -} +export type ApplyInsightsMigrationsReport = ApplyPackageMigrationsReport; export async function applyInsightsMigrations( databaseUrl: string, ): Promise { - const sql = postgres(databaseUrl, { max: 1, onnotice: () => undefined }); - try { - await sql.unsafe(`CREATE SCHEMA IF NOT EXISTS ${quoteIdentifier(SCHEMA)}`); - - await sql.unsafe( - `CREATE TABLE IF NOT EXISTS ${quoteQualified(SCHEMA, LEDGER_TABLE)} (` + - `name text PRIMARY KEY, applied_at timestamptz NOT NULL DEFAULT now())`, - ); - - const applied: string[] = []; - const alreadyApplied: string[] = []; - - for (const migration of insightsMigrations) { - const existing = await sql.unsafe( - `SELECT 1 FROM ${quoteQualified(SCHEMA, LEDGER_TABLE)} WHERE name = $1`, - [migration.name], - ); - if (existing.length > 0) { - alreadyApplied.push(migration.name); - continue; - } - try { - await sql.begin(async (tx) => { - await tx.unsafe(migration.sql); - await tx.unsafe( - `INSERT INTO ${quoteQualified(SCHEMA, LEDGER_TABLE)} (name) VALUES ($1)`, - [migration.name], - ); - }); - applied.push(migration.name); - } catch (err) { - const message = err instanceof Error ? err.message : String(err); - throw new Error( - `insights migration ${migration.name} failed: ${message}`, - { cause: err }, - ); - } - } - - return { applied, alreadyApplied }; - } finally { - await sql.end({ timeout: 5 }); - } + return applyPackageMigrations({ + databaseUrl, + schema: SCHEMA, + ledgerTable: LEDGER_TABLE, + migrations: insightsMigrations, + packageLabel: "insights", + }); } diff --git a/packages/migration-runner/test/migration-runner.test.ts b/packages/migration-runner/test/migration-runner.test.ts index 8a9188183..213a7646e 100644 --- a/packages/migration-runner/test/migration-runner.test.ts +++ b/packages/migration-runner/test/migration-runner.test.ts @@ -6,10 +6,7 @@ import { afterAll, beforeAll, describe, expect, test } from "bun:test"; import postgres from "postgres"; import { e2eDatabaseUrl } from "../../../scripts/e2e/harness"; -import { - applyPackageMigrations, - type PackageMigration, -} from "../src/index"; +import { applyPackageMigrations, type PackageMigration } from "../src/index"; function scratchUrlFor(e2eUrl: string): string { const url = new URL(e2eUrl); @@ -66,18 +63,14 @@ describeIfDb("applyPackageMigrations", () => { beforeAll(async () => { await withMaintenanceConnection(async (maintenance) => { - await maintenance.unsafe( - `DROP DATABASE IF EXISTS "${scratchDatabase}"`, - ); + await maintenance.unsafe(`DROP DATABASE IF EXISTS "${scratchDatabase}"`); await maintenance.unsafe(`CREATE DATABASE "${scratchDatabase}"`); }); }, 20000); afterAll(async () => { await withMaintenanceConnection(async (maintenance) => { - await maintenance.unsafe( - `DROP DATABASE IF EXISTS "${scratchDatabase}"`, - ); + await maintenance.unsafe(`DROP DATABASE IF EXISTS "${scratchDatabase}"`); }); }, 20000); @@ -129,9 +122,7 @@ describeIfDb("applyPackageMigrations", () => { expect(new Set(appliedNames).size).toBe(appliedNames.length); expect([...appliedNames, ...alreadyAppliedNames].sort()).toEqual( - [...migrations, ...migrations] - .map((migration) => migration.name) - .sort(), + [...migrations, ...migrations].map((migration) => migration.name).sort(), ); const sql = postgres(scratchUrl, { max: 1, onnotice: () => undefined }); diff --git a/packages/preferences/package.json b/packages/preferences/package.json index e469394c9..c8e8303b6 100644 --- a/packages/preferences/package.json +++ b/packages/preferences/package.json @@ -16,6 +16,7 @@ }, "dependencies": { "@corbits/api-query": "workspace:*", + "@corbits/migration-runner": "workspace:*", "@intx/hub-api": "workspace:*", "arktype": "catalog:", "drizzle-orm": "catalog:", diff --git a/packages/preferences/src/migrations.ts b/packages/preferences/src/migrations.ts index 12bd26359..55a83e7db 100644 --- a/packages/preferences/src/migrations.ts +++ b/packages/preferences/src/migrations.ts @@ -2,13 +2,17 @@ // own ledger table so the package can be extracted without disentangling // history from the platform drizzle journal. Every table this package // owns — including its ledger — lives in its own `preferences` Postgres -// schema, never `public`; see docs/package-migrations.md. -import postgres from "postgres"; - -export interface PreferencesMigration { - name: string; - sql: string; -} +// schema, never `public`; see docs/package-migrations.md. Mechanics +// (schema/ledger bootstrap, transactional apply, the advisory lock +// across concurrent hub replicas) live in @corbits/migration-runner — +// this file owns only the domain SQL. +import { + applyPackageMigrations, + type ApplyPackageMigrationsReport, + type PackageMigration, +} from "@corbits/migration-runner"; + +export type PreferencesMigration = PackageMigration; const SCHEMA = "preferences"; @@ -29,63 +33,16 @@ export const preferencesMigrations: readonly PreferencesMigration[] = [ const LEDGER_TABLE = "preferences_migrations"; -function quoteIdentifier(name: string): string { - return `"${name.replace(/"/g, '""')}"`; -} - -function quoteQualified(schema: string, name: string): string { - return `${quoteIdentifier(schema)}.${quoteIdentifier(name)}`; -} - -export interface ApplyPreferencesMigrationsReport { - applied: string[]; - alreadyApplied: string[]; -} +export type ApplyPreferencesMigrationsReport = ApplyPackageMigrationsReport; export async function applyPreferencesMigrations( databaseUrl: string, ): Promise { - const sql = postgres(databaseUrl, { max: 1, onnotice: () => undefined }); - try { - await sql.unsafe(`CREATE SCHEMA IF NOT EXISTS ${quoteIdentifier(SCHEMA)}`); - - await sql.unsafe( - `CREATE TABLE IF NOT EXISTS ${quoteQualified(SCHEMA, LEDGER_TABLE)} (` + - `name text PRIMARY KEY, applied_at timestamptz NOT NULL DEFAULT now())`, - ); - - const applied: string[] = []; - const alreadyApplied: string[] = []; - - for (const migration of preferencesMigrations) { - const existing = await sql.unsafe( - `SELECT 1 FROM ${quoteQualified(SCHEMA, LEDGER_TABLE)} WHERE name = $1`, - [migration.name], - ); - if (existing.length > 0) { - alreadyApplied.push(migration.name); - continue; - } - try { - await sql.begin(async (tx) => { - await tx.unsafe(migration.sql); - await tx.unsafe( - `INSERT INTO ${quoteQualified(SCHEMA, LEDGER_TABLE)} (name) VALUES ($1)`, - [migration.name], - ); - }); - applied.push(migration.name); - } catch (err) { - const message = err instanceof Error ? err.message : String(err); - throw new Error( - `preferences migration ${migration.name} failed: ${message}`, - { cause: err }, - ); - } - } - - return { applied, alreadyApplied }; - } finally { - await sql.end({ timeout: 5 }); - } + return applyPackageMigrations({ + databaseUrl, + schema: SCHEMA, + ledgerTable: LEDGER_TABLE, + migrations: preferencesMigrations, + packageLabel: "preferences", + }); } From 4720e568e3cb21289a5799bef9a5c67a253fc420 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 04:43:30 -0700 Subject: [PATCH 4/6] Update docs: package migrations use shared migration-runner Document the advisory lock @corbits/migration-runner adds around the bootstrap-and-apply run, and point the self-contained-shape reference implementation at access-policy now that all six packages route through applyPackageMigrations instead of a hand-rolled loop. --- docs/package-migrations.md | 43 ++++++++++++++++++++++++-------------- 1 file changed, 27 insertions(+), 16 deletions(-) diff --git a/docs/package-migrations.md b/docs/package-migrations.md index 3ca975142..5dada978c 100644 --- a/docs/package-migrations.md +++ b/docs/package-migrations.md @@ -41,6 +41,14 @@ tables — how its migrations are written, tracked, and applied. ledger row — is left behind half-applied; the next run retries the same migration from a clean slate instead of getting stuck on a schema that looks partially migrated but has no ledger row to explain why. +- **A session-level advisory lock guards the whole run.** Two hub replicas + can boot at the same time and both call a package's migration runner + concurrently; without a lock, both can see a migration as unapplied and + race the same ledger `INSERT`, crashing the loser on a duplicate key. + `@corbits/migration-runner`'s `applyPackageMigrations` holds + `pg_advisory_lock(hashtext(ledgerTable))` for the whole bootstrap-and-apply + run and releases it in a `finally`, so the second replica simply waits and + then finds every migration already applied. ## Why not `drizzle-kit` codegen @@ -63,16 +71,18 @@ package onto the transactional pattern — see below): 1. **Self-contained, transactional** — `@corbits/chat`, `@corbits/notify`, `@corbits/webhook-triggers`, `@corbits/routines`, `@corbits/insights`, - `@corbits/skills`. - The package's `src/migrations.ts` owns a literal `{ name, sql }[]` array, - opens its own short-lived `postgres` client, creates its ledger table if - absent, and applies each not-yet-applied migration inside - `sql.begin(async (tx) => { ... })` — the migration's SQL and its ledger - insert commit or roll back together. `scripts/db-setup.ts` imports the + `@corbits/skills`, `@corbits/bench`, `@corbits/preferences`, + `@corbits/inference-catalog`, `@corbits/evals`, `@corbits/access-policy`. + The package's `src/migrations.ts` owns only a literal `{ name, sql }[]` + array and a thin `applyXMigrations(databaseUrl)` wrapper; the mechanics — + schema/ledger bootstrap, the transactional apply loop, and the advisory + lock — live once in `@corbits/migration-runner`'s + `applyPackageMigrations`, called with the package's schema name, ledger + table name, and migration array. `scripts/db-setup.ts` imports the package's `applyXMigrations(databaseUrl)` function directly and calls it after the platform's own migrations. See - `packages/insights/src/migrations.ts` for the reference implementation - every other package now matches. + `packages/access-policy/src/migrations.ts` for the reference + implementation every other package now matches. 2. **Delegated to an external package** — `@corbits/mailbox`. The package ships and owns its entire migration story (its own literal SQL, its own ledger, its own `mailbox` schema) behind a single exported runner @@ -108,11 +118,12 @@ the new schema by running `workbench reset` rather than by an in-place ## Which shape to use for a new package **Transactional, self-contained** (shape 1 above), in a Postgres schema -named for the package. Copy `packages/insights/src/migrations.ts`'s shape: -a `pgSchema("")` in `schema.ts`, a literal migration array whose SQL -qualifies every table/index with that schema, a package-named ledger table -living in the same schema, and each migration applied inside `sql.begin`. -Reach for the delegated shape only when the package already ships its own -migration runner as part of a larger, independently-owned engine (its own -schema, its own connection handling) — not as a shortcut to skip writing a -ledger table. +named for the package. Copy `packages/access-policy/src/migrations.ts`'s +shape: a `pgSchema("")` in `schema.ts`, a literal migration array +whose SQL qualifies every table/index with that schema, a package-named +ledger table living in the same schema, and an `applyXMigrations` that +calls `@corbits/migration-runner`'s `applyPackageMigrations` with that +schema, ledger table, and migration array. Reach for the delegated shape +only when the package already ships its own migration runner as part of a +larger, independently-owned engine (its own schema, its own connection +handling) — not as a shortcut to skip writing a ledger table. From 26b2a395989a6b32e17eeb17158b7f5a2da99c76 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 06:12:04 -0700 Subject: [PATCH 5/6] Add lock-release-on-failure test to migration-runner Existing coverage proved the lock releases on success, under concurrent replicas, and after the holding connection dies, but not on the path this PR's own finally-block guards: a migration whose SQL fails. Proves a second run reaches the same lock immediately instead of hanging behind the first run's connection. --- .../test/migration-runner.test.ts | 32 +++++++++++++++++++ 1 file changed, 32 insertions(+) diff --git a/packages/migration-runner/test/migration-runner.test.ts b/packages/migration-runner/test/migration-runner.test.ts index 213a7646e..e361a6b0a 100644 --- a/packages/migration-runner/test/migration-runner.test.ts +++ b/packages/migration-runner/test/migration-runner.test.ts @@ -159,4 +159,36 @@ describeIfDb("applyPackageMigrations", () => { expect(result.applied).toEqual(["0001_widget", "0002_widget_label"]); }, 10000); + + test("releases the lock when a migration fails, so the next run isn't blocked behind it", async () => { + const failSchema = `${SCHEMA}_fail`; + const failLedger = `${LEDGER_TABLE}_fail`; + const brokenMigrations: readonly PackageMigration[] = [ + { name: "0001_broken", sql: "SELECT this_column_does_not_exist" }, + ]; + + await expect( + applyPackageMigrations({ + databaseUrl: scratchUrl, + schema: failSchema, + ledgerTable: failLedger, + migrations: brokenMigrations, + packageLabel: "migration-runner-fail-test", + }), + ).rejects.toThrow( + /migration-runner-fail-test migration 0001_broken failed/, + ); + + // A second run reaching the same lock immediately (rather than hanging + // behind the first run's connection) proves the failed run's `finally` + // released it instead of leaving it held until the connection closed. + const result = await applyPackageMigrations({ + databaseUrl: scratchUrl, + schema: failSchema, + ledgerTable: failLedger, + migrations, + packageLabel: "migration-runner-fail-test", + }); + expect(result.applied).toEqual(["0001_widget", "0002_widget_label"]); + }, 10000); }); From a7bc4a00def050ce654a54ab59d905cfb71e955e Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 06:12:18 -0700 Subject: [PATCH 6/6] Update docs: state the hashtext collision risk plainly in migration-runner README MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The prior wording read as a guarantee ("distinct packages never contend") when what's actually true is narrower: the lock key is a 32-bit hash of the ledger table name, not the name itself, so two distinct names could in principle collide — a liveness cost (one package's boot waits on another's), never a correctness one. Verified against a live hashtext() that today's six ledger table names don't collide, and says so, with the reminder to reverify for a seventh. --- packages/migration-runner/README.md | 13 ++++++++++--- 1 file changed, 10 insertions(+), 3 deletions(-) diff --git a/packages/migration-runner/README.md b/packages/migration-runner/README.md index e28073cc0..2b844a47c 100644 --- a/packages/migration-runner/README.md +++ b/packages/migration-runner/README.md @@ -20,9 +20,16 @@ the ledger and finds every migration already applied. Postgres also releases a session-level advisory lock automatically if the holding connection dies, so a crashed replica never wedges the lock for the next boot. -The lock key is derived from the ledger table name, which is already unique -per package (`bench_migrations`, `insights_migrations`, ...), so distinct -packages never contend on the same lock. +The lock key is `hashtext(ledgerTable)`, a 32-bit hash of the ledger table +name rather than the name itself, so two distinct table names could in +principle hash to the same key — two packages would then serialize their +boot-time migrations against each other instead of running in parallel, a +liveness cost (one waits its turn) and never a correctness one (each still +applies to its own schema and ledger). The six current ledger tables +(`access_policy_migrations`, `bench_migrations`, `evals_migrations`, +`inference_catalog_migrations`, `insights_migrations`, +`preferences_migrations`) do not collide — verified against a live +`hashtext()`. Confirm the same before naming a seventh. ## What it does not change