Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions .github/workflows/delphi-characterization.yml
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,8 @@ jobs:
run: |
python -m pip install --quiet pyyaml
python -m unittest tests.test_delphi_storage_codec -v
python -m unittest tests.test_delphi_postgres_results -v
python -m unittest tests.test_delphi_result_resource -v

- name: Codec golden test (Node reads the files Python wrote, writes the cross file)
working-directory: server
Expand Down
63 changes: 63 additions & 0 deletions .github/workflows/dynamo-removal.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
name: Dynamo removal

on:
pull_request:
paths:
- 'delphi/**'
- 'queue-rs/**'
- 'server/**'
- 'client-report/**'
- 'scripts/test-dynamo-removal.sh'
- 'scripts/prove-dynamo-report.sh'
- 'ci/dynamo-removal/**'
- '.github/workflows/dynamo-removal.yml'
workflow_dispatch:

permissions:
contents: read

concurrency:
group: dynamo-removal-${{ github.ref }}
cancel-in-progress: true

jobs:
postgres-report:
runs-on: ubuntu-24.04
timeout-minutes: 150
env:
COMPOSE_PROJECT_NAME: polis-graph-test-dynamo-${{ github.run_id }}-${{ github.run_attempt }}
POLIS_RECOVERY_PG_PORT: '55449'
RECOVERY_PG_PORT: '55449'
steps:
- uses: actions/checkout@11bd71901bbe5b1630ceea73d27597364c9af683 # v4.2.2
with:
persist-credentials: false
- uses: actions/setup-python@a26af69be951a213d495a4c3e4e4022e16d87065 # v5
with:
python-version: '3.12'
- uses: actions/setup-node@49933ea5288caeca8642d1e84afbd3f7d6820020 # v4
with:
node-version: '22'
- name: Queue toolchain
working-directory: queue-rs
run: rustup show active-toolchain
- name: Install uv
run: python -m pip install uv==0.9.2
- name: Queue, Postgres results, importer and report with DynamoDB stopped
run: bash scripts/test-dynamo-removal.sh
- name: Retain proof logs and report screenshots
if: always()
uses: actions/upload-artifact@ea165f8d65b6e75b540449e92b4886f43607fa02 # v4.6.2
with:
name: dynamo-removal-proof
retention-days: 7
include-hidden-files: true
path: |
.dynamo-proof/${{ env.COMPOSE_PROJECT_NAME }}/*.log
.dynamo-proof/${{ env.COMPOSE_PROJECT_NAME }}/demo-result.json
.dynamo-proof/${{ env.COMPOSE_PROJECT_NAME }}/sql-source-sha256.json
.dynamo-proof/${{ env.COMPOSE_PROJECT_NAME }}/install/results.json
.dynamo-proof/${{ env.COMPOSE_PROJECT_NAME }}/results-sql/results.json
.dynamo-proof/${{ env.COMPOSE_PROJECT_NAME }}/report/*.log
.dynamo-proof/${{ env.COMPOSE_PROJECT_NAME }}/report/browser.json
.dynamo-proof/${{ env.COMPOSE_PROJECT_NAME }}/report/*.png
27 changes: 27 additions & 0 deletions .github/workflows/dynamo-writers.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
name: Postgres Delphi writers
on:
pull_request:
paths: ['delphi/**', 'server/**', 'queue-rs/**', 'scripts/test-dynamo-writers.sh', '.github/workflows/dynamo-writers.yml']
workflow_dispatch:
permissions:
contents: read
jobs:
writers:
runs-on: ubuntu-24.04
timeout-minutes: 90
env:
COMPOSE_PROJECT_NAME: polis-graph-test-writers-${{ github.run_id }}-${{ github.run_attempt }}
POLIS_RECOVERY_PG_PORT: '55453'
RECOVERY_PG_PORT: '55453'
steps:
- uses: actions/checkout@11bd71901bbe5b1630ceea73d27597364c9af683
with:
persist-credentials: false
- run: bash scripts/test-dynamo-writers.sh
- uses: actions/upload-artifact@ea165f8d65b6e75b540449e92b4886f43607fa02
if: always()
with:
name: delphi-writers-proof
path: .writer-proof/**/*.log
include-hidden-files: true
retention-days: 7
10 changes: 10 additions & 0 deletions ci/dynamo-removal/report-nginx.conf
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
events {}
http {
include /etc/nginx/mime.types;
server {
listen 8080;
root /report;
location /api/ { proxy_pass http://astra-dynamo1424-server:5000; proxy_read_timeout 180s; }
location / { try_files $uri /index_report.html; }
}
}
78 changes: 78 additions & 0 deletions ci/dynamo-removal/report-proof.cjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,78 @@
/* Read-only browser proof: use generated local reports, never a production URL. */
const fs=require('node:fs');
const path=require('node:path');
const assert=require('node:assert/strict');
const {assertFixtureReports,assertReportCoverage}=require('./report-quality.cjs');
const playwright=process.env.DYNAMO_PROOF_PLAYWRIGHT;
assert(playwright,'DYNAMO_PROOF_PLAYWRIGHT must point to an installed local Playwright package');
const {chromium}=require(playwright);
const base=process.env.DYNAMO_PROOF_REPORT_URL;
assert(/^http:\/\/(127\.0\.0\.1|localhost):\d+$/.test(base),'local browser proof URL required');
const rid=process.env.DYNAMO_PROOF_REPORT_ID;
assert(/^rlocal[a-z0-9]+$/.test(rid),'generated local report ID required');
const output=process.env.DYNAMO_PROOF_OUTPUT;
assert(output,'DYNAMO_PROOF_OUTPUT required');fs.mkdirSync(output,{recursive:true});
(async()=>{
const browser=await chromium.launch({headless:true});const receipts=[];
try {
for(const route of ['report','topicStats','topicReport']){
const page=await browser.newPage({viewport:{width:1440,height:1100}});
const errors=[];const responses=[];let narrativeEvidence;
page.on('pageerror',error=>errors.push(error.message));
page.on('response',async response=>{
if(!response.url().includes('/api/'))return;
let body;try{body=await response.text()}catch{body='unavailable'}
responses.push({url:response.url(),status:response.status(),body});
});
await page.goto(`${base}/${route}/${rid}`,{waitUntil:'networkidle'});
if(route==='report') {
const response=await page.request.get(`${base}/api/v3/delphi/visualizations?report_id=${encodeURIComponent(rid)}`);
assert(response.ok(),'visualization metadata endpoint succeeded');
const metadata=await response.json();assert.equal(metadata.status,'success');
assert(metadata.jobs?.length>0,'actual queue metadata returned with DynamoDB unavailable');
assert(metadata.jobs.every(job=>job.workLive===false),'published completed graph has no live work');
responses.push({url:response.url(),status:response.status(),body:JSON.stringify(metadata)});
}
if(route==='topicStats')await page.getByText('Group Consensus',{exact:false}).first().waitFor({timeout:30000});
if(route==='topicReport'){
const selector=page.locator('select');
await selector.first().waitFor({timeout:30000});
const options=await selector.first().locator('option').evaluateAll(options=>options.map(o=>({value:o.value,text:o.textContent})));
const sourceResponse=await page.request.get(`${base}/api/v3/delphi/reports?report_id=${encodeURIComponent(rid)}`);
assert(sourceResponse.ok(),'narrative source endpoint succeeded');
const source=await sourceResponse.json();assert.equal(source.status,'success');
const checked=assertFixtureReports(source.reports);
const topicResponse=await page.request.get(`${base}/api/v3/delphi?report_id=${encodeURIComponent(rid)}`);
assert(topicResponse.ok(),'topic source endpoint succeeded');
const topics=await topicResponse.json();assert.equal(topics.status,'success');
const checkedTopics=Object.values(topics.runs).flatMap(run=>Object.values(run.topics_by_layer).flatMap(Object.values)).length;
assertReportCoverage(source.reports,topics.runs);
const named=options.find(o=>o.value && source.reports?.[o.value]?.report_data);
assert(named,'at least one available narrative section maps to the rendered selector');
const stored=source.reports[named.value];
const {clauses}=checked[named.value];
await selector.first().selectOption(named.value);
await page.waitForLoadState('networkidle');
await page.locator('.topic-text-content .paragraph').first().waitFor({timeout:30000});
const rendered=await page.locator('.topic-text-content').innerText();
for(const clause of clauses)assert(rendered.includes(clause),'stored narrative clause renders exactly');
narrativeEvidence={section:named.value,model:stored.model,job_id:stored.job_id,clauses:clauses.length,metadata:stored.metadata,checkedSections:Object.keys(checked),checkedTopics};
}
const body=await page.locator('body').innerText();
await page.screenshot({path:path.join(output,route+'.png'),fullPage:true});
await page.screenshot({path:path.join(output,route+'-viewport.png')});
const record={route,body,errors,responses,narrativeEvidence};receipts.push(record);
fs.writeFileSync(path.join(output,'browser.json'),JSON.stringify(receipts,null,2));
assert.equal(errors.length,0,route+': '+errors.join('; '));
assert(responses.length>0,route+': API requests observed');
assert(responses.every(r=>r.status<400),route+': '+responses.filter(r=>r.status>=400).map(r=>r.url));
for(const response of responses){let payload;try{payload=JSON.parse(response.body)}catch{continue}assert.notEqual(payload.status,'error',response.url+': application error')}
// Existing date-display behavior is outside this storage/queue proof.
// Preserve it with the default-backend route and view recordings.
if(route==='report')assert(/people voted/.test(body),route+': participant summary rendered');
if(route==='topicReport')assert(body.length>300,route+': narrative content rendered');
console.log('PASS',route,'API responses',responses.length);
await page.close();
}
} finally {await browser.close()}
})().catch(error=>{console.error(error);process.exitCode=1});
38 changes: 38 additions & 0 deletions ci/dynamo-removal/report-quality.cjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
const assert = require('node:assert/strict');

function assertReportCoverage(reports, runs) {
const topics = Object.values(runs || {}).flatMap(run => Object.values(run.topics_by_layer || {}).flatMap(Object.values));
assert(topics.length > 0, 'topic coverage requires served topics');
const expected = topics.map(topic => {
assert(/^.+#\d+#\d+$/.test(topic.topic_key), 'valid topic key required');
return topic.topic_key.replaceAll('#', '_');
});
assert.equal(new Set(expected).size, expected.length, 'duplicate topic section');
const jobs = new Set(topics.map(topic => topic.topic_key.split('#')[0]));
assert.equal(jobs.size, 1, 'one coherent served narrative job required');
const [job] = jobs;
expected.push(...['groups', 'group_informed_consensus', 'uncertainty'].map(name => `${job}_global_${name}`));
const actual = Object.entries(reports || {}).map(([key, report]) => {
assert.equal(report.section, key, 'duplicate or mismatched report section');
assert.equal(report.job_id, job, 'report belongs to served narrative job');
return key;
});
assert.deepEqual(actual.sort(), expected.sort(), 'complete topic and global section coverage');
return expected.length;
}

function assertFixtureReports(reports) {
assert(reports && Object.keys(reports).length > 0, 'served reports required');
const checked = {};
for (const [section, stored] of Object.entries(reports)) {
const document = typeof stored.report_data === 'string' ? JSON.parse(stored.report_data) : stored.report_data;
assert.equal(stored.model, 'local-narrative-fixture/1');
assert.equal(stored.metadata?.provider_fixture, true);
assert.equal(document.provider_fixture, true);
const clauses = document.paragraphs.flatMap(p => p.sentences.flatMap(s => s.clauses.map(c => c.text)));
assert.deepEqual(clauses, ['Fixed narrative stand-in for queue, Postgres storage and report rendering proof. No LLM provider was called.']);
checked[section] = { document, clauses };
}
return checked;
}
module.exports = { assertFixtureReports, assertReportCoverage };
33 changes: 33 additions & 0 deletions ci/dynamo-removal/report-quality.test.cjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,33 @@
const { test } = require('node:test');
const assert = require('node:assert/strict');
const { assertReportCoverage } = require('./report-quality.cjs');
const coverageFixture = () => {
const runs = { current: { topics_by_layer: { 0: {
0: { topic_key: 'generated#0#0' }, 1: { topic_key: 'generated#0#1' }
} } } };
const sections = ['generated_0_0', 'generated_0_1', 'generated_global_groups',
'generated_global_group_informed_consensus', 'generated_global_uncertainty'];
return { runs, reports: Object.fromEntries(sections.map(section => [section, { section, job_id: 'generated' }])) };
};
test('requires all topic and global sections from one served job', () => {
const { reports, runs } = coverageFixture();
assert.equal(assertReportCoverage(reports, runs), 5);
});
test('rejects a missing topic or global even when remaining reports are valid', () => {
for (const missing of ['generated_0_1', 'generated_global_uncertainty']) {
const { reports, runs } = coverageFixture(); delete reports[missing];
assert.throws(() => assertReportCoverage(reports, runs), /complete topic and global/);
}
});
test('rejects duplicate topic keys or duplicate report identities', () => {
const { reports, runs } = coverageFixture();
runs.current.topics_by_layer[0][1].topic_key = 'generated#0#0';
assert.throws(() => assertReportCoverage(reports, runs), /duplicate topic/);
const fresh = coverageFixture();
fresh.reports.generated_0_1.section = 'generated_0_0';
assert.throws(() => assertReportCoverage(fresh.reports, fresh.runs), /duplicate or mismatched/);
});
test('rejects reports from a different job', () => {
const { reports, runs } = coverageFixture(); reports.generated_0_0.job_id = 'other-job';
assert.throws(() => assertReportCoverage(reports, runs), /served narrative job/);
});
1 change: 1 addition & 0 deletions delphi/Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -137,6 +137,7 @@ RUN install -d /opt/polis-unflip && \
# logs a failed-send error and traceback for every flush, which buried the
# poller's own output; tracing is therefore off unless explicitly enabled.
CMD ["bash", "-c", "\
if [ \"${DELPHI_RESULT_BACKEND}\" = \"postgres\" ]; then export POLIS_JOBS_ENABLED=1; exec polis-jobs; fi; \
echo 'Ensuring DynamoDB tables are set up (runs in all environments)...'; \
python create_dynamodb_tables.py --region ${AWS_REGION} && \
echo 'DynamoDB table setup script finished.'; \
Expand Down
5 changes: 5 additions & 0 deletions delphi/create_dynamodb_tables.py
Original file line number Diff line number Diff line change
Expand Up @@ -530,6 +530,8 @@ def _create_tables(dynamodb, tables, existing_tables):
def create_tables(endpoint_url=None, region_name='us-east-1',
delete_existing=False, evoc_only=False, polismath_only=False,
aws_profile=None):
if os.environ.get('DELPHI_RESULT_BACKEND') == 'postgres':
return []
# Use the environment variable if endpoint_url is not provided
if endpoint_url is None:
endpoint_url = os.environ.get('DYNAMODB_ENDPOINT')
Expand Down Expand Up @@ -601,6 +603,9 @@ def create_tables(endpoint_url=None, region_name='us-east-1',
return created_tables

def main():
if os.environ.get('DELPHI_RESULT_BACKEND') == 'postgres':
print('Postgres result schema is managed by migrations; no Dynamo bootstrap')
return
# Parse arguments
parser = argparse.ArgumentParser(description='Create DynamoDB tables for Delphi system')
parser.add_argument('--endpoint-url', type=str, default=None,
Expand Down
60 changes: 60 additions & 0 deletions delphi/docs/LEGACY_DYNAMO_IMPORT.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
# Import a local DynamoDB export

`scripts/import_dynamo_export.py` moves a bounded, checksummed export into a
normal Postgres graph job. The worker verifies source bytes, codec version,
worker code, importer code and runtime; normal queue finalization stores its
immutable artifact and M28 result rows together. No historical job executes.

The export command requires an explicit loopback endpoint and only uses dummy
local credentials. It never falls back to an AWS endpoint. Specify every family
to export with repeated `--family` arguments. The exporter reads all scan pages
and writes the frozen `delphi-storage-codec/1` files and `SHA256SUMS`.
Keep the source quiescent during export. Consistent scan pages do not provide
a transaction snapshot across a whole table or all families; checksums verify
the completed export but cannot detect changes between scan pages.

```sh
python scripts/import_dynamo_export.py export-local /tmp/generated-export \
--endpoint http://127.0.0.1:8000 \
--family Delphi_CommentEmbeddings --family Delphi_NarrativeReports \
--family Delphi_JobQueue --family Delphi_JobActiveGuard
python scripts/import_dynamo_export.py preview /tmp/generated-export \
--zid "$DEMO_ZID" --report-id "$DEMO_REPORT"
python scripts/import_dynamo_export.py enqueue /tmp/generated-export \
--zid "$DEMO_ZID" --report-id "$DEMO_REPORT" --env local-demo --scope import-demo
```

For enqueue, `QUEUE_DATABASE_URL` must use a queue executor login.
`DATABASE_URL` supplies read access to verify that every explicitly named
report belongs to the selected conversation. Preview is offline; it checks
explicit source bindings but cannot verify the live reports mapping. Neither
command remaps conversation IDs or report IDs. Mixed-conversation exports fail.

The `delphi` worker claims `graph_narrative` and runs the declared
`legacy-dynamo-export/1` model. Repeating the same enqueue returns the same
graph. Source, code or runtime changes produce a new request; an active scope
still prevents overlapping admission. Admission does not publish the result.
After successful execution, publication uses the ordinary explicit
`pd_graph_publish` generation check, so a failed import cannot replace a report.
The approved core contract does not support replacing a failed graph with a
superseding branch; that capability is deferred to #1436. The importer exposes
ordinary admission, status verification and publication only.

Every source row is counted as one of:

- Imported result rows in the 18 M28 result families.
- Archived T15/T16 job and guard rows in `legacy_control_files`, with exact
codec bytes on the immutable artifact. They are never converted to current
jobs, active scopes, leases or provider submissions.
- Quarantined result rows containing NUL strings/keys, which JSONB cannot
represent. Their exact codec bytes and `postgres-jsonb-nul` reason remain in
`quarantine` on the immutable artifact. Binary zero bytes are valid base64
codec data and are imported normally.

Unknown families, duplicate/noncanonical keys, missing/changed/unlisted files,
wrong conversation/report bindings and oversized artifacts fail before enqueue.
The initial importer is bounded to 450,000 serialized output bytes and the
queue's existing input limit. It refuses larger exports without truncation;
large archives need a separate artifact transport. It processes result values
as tagged codec values, preserving decimal numbers, sets, binary values and
JSON stored as strings. It does not read, write or reinterpret stored votes.
Loading
Loading