Compare commits
18
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
9cd461cc08 | ||
|
|
9c8be25ff0 | ||
|
|
83ea07ed6a | ||
|
|
61ebe977d3 | ||
|
|
763595a5a3 | ||
|
|
398617391b | ||
|
|
a0e0e88f07 | ||
|
|
c80a5db5f5 | ||
|
|
4fa67afc3d | ||
|
|
8bd78d2624 | ||
|
|
b3c4b8b463 | ||
|
|
57af5b9dc2 | ||
|
|
44fbf8c587 | ||
|
|
76840d8a57 | ||
|
|
d459ded168 | ||
|
|
43b67d213d | ||
|
|
d42b779644 | ||
|
|
4564ed5bec |
@@ -38,6 +38,8 @@ project: &project
|
||||
- frontend/**
|
||||
- docker/**
|
||||
- scripts/RestartHelper.java
|
||||
- scripts/db-migration/**
|
||||
- .github/workflows/db-migration-test.yml
|
||||
|
||||
frontend: &frontend
|
||||
- frontend/**
|
||||
|
||||
@@ -13,7 +13,7 @@ Usage:
|
||||
"""
|
||||
|
||||
# Sample for Windows:
|
||||
# python .github/scripts/check_language_toml.py --reference-file frontend/public/locales/en-GB/translation.toml --branch "" --files frontend/public/locales/de-DE/translation.toml frontend/public/locales/fr-FR/translation.toml
|
||||
# python .github/scripts/check_language_toml.py --reference-file frontend/editor/public/locales/en-GB/translation.toml --branch "" --files frontend/editor/public/locales/de-DE/translation.toml frontend/editor/public/locales/fr-FR/translation.toml
|
||||
|
||||
import argparse
|
||||
import glob
|
||||
@@ -308,7 +308,7 @@ def check_for_differences(reference_file, file_list, branch, actor):
|
||||
report.append("## ❌ Overall Check Status: **_Failed_**")
|
||||
report.append("")
|
||||
report.append(
|
||||
f"@{actor} please check your translation if it conforms to the standard. Follow the format of [en-GB/translation.toml](https://github.com/Stirling-Tools/Stirling-PDF/blob/main/frontend/public/locales/en-GB/translation.toml)"
|
||||
f"@{actor} please check your translation if it conforms to the standard. Follow the format of [en-GB/translation.toml](https://github.com/Stirling-Tools/Stirling-PDF/blob/main/frontend/editor/public/locales/en-GB/translation.toml)"
|
||||
)
|
||||
else:
|
||||
report.append("## ✅ Overall Check Status: **_Success_**")
|
||||
|
||||
@@ -287,6 +287,7 @@ jobs:
|
||||
- /stirling/V2-PR-${{ needs.check-pr.outputs.pr_number }}/data:/usr/share/tessdata:rw
|
||||
- /stirling/V2-PR-${{ needs.check-pr.outputs.pr_number }}/config:/configs:rw
|
||||
- /stirling/V2-PR-${{ needs.check-pr.outputs.pr_number }}/logs:/logs:rw
|
||||
- /stirling/V2-PR-${{ needs.check-pr.outputs.pr_number }}/storage:/storage:rw
|
||||
environment:
|
||||
DISABLE_ADDITIONAL_FEATURES: "false"
|
||||
SECURITY_ENABLELOGIN: "true"
|
||||
@@ -309,7 +310,7 @@ jobs:
|
||||
|
||||
ssh -i ../private.key -o StrictHostKeyChecking=no -o UserKnownHostsFile=/dev/null -T ${{ secrets.NEW_VPS_USERNAME }}@${{ secrets.NEW_VPS_HOST }} << ENDSSH
|
||||
# Create V2 PR-specific directories
|
||||
mkdir -p /stirling/V2-PR-${{ needs.check-pr.outputs.pr_number }}/{data,config,logs}
|
||||
mkdir -p /stirling/V2-PR-${{ needs.check-pr.outputs.pr_number }}/{data,config,logs,storage}
|
||||
|
||||
# Move docker-compose file to correct location
|
||||
mv /tmp/docker-compose-v2.yml /stirling/V2-PR-${{ needs.check-pr.outputs.pr_number }}/docker-compose.yml
|
||||
|
||||
+124
-27
@@ -1,8 +1,9 @@
|
||||
name: AI Engine CI
|
||||
|
||||
# Validates the Python AI engine: regenerates tool models, runs fixers,
|
||||
# lint, type-check, and tests. Called from build.yml on PRs and merge_group;
|
||||
# also runs directly on push to main as a post-merge safety net.
|
||||
# Validates the Python AI engine: regenerates tool models and runs the
|
||||
# engine quality gate (lint, type-check, format-check, tests). Called from
|
||||
# build.yml on PRs and merge_group; also runs directly on push to main as
|
||||
# a post-merge safety net.
|
||||
on:
|
||||
workflow_call:
|
||||
push:
|
||||
@@ -51,27 +52,95 @@ jobs:
|
||||
run: task engine:tool-models
|
||||
|
||||
- name: Verify tool models are up to date
|
||||
id: tool-models-check
|
||||
continue-on-error: true
|
||||
run: git diff --exit-code engine/src/stirling/models/tool_models.py
|
||||
|
||||
- name: Comment on tool models check failure
|
||||
# Only post a comment on PRs. github-script's PR helpers need an
|
||||
# issue/PR number, which doesn't exist on merge_group runs.
|
||||
if: steps.tool-models-check.outcome == 'failure' && github.event_name == 'pull_request'
|
||||
continue-on-error: true
|
||||
uses: actions/github-script@3a2844b7e9c422d3c10d287c895573f7108da1b3 # v9.0.0
|
||||
with:
|
||||
script: |
|
||||
const marker = '<!-- tool-models-check -->';
|
||||
const body = [
|
||||
marker,
|
||||
'### Tool Models Check Failed',
|
||||
'',
|
||||
'The generated `engine/src/stirling/models/tool_models.py` is out of date with the Java OpenAPI spec and will need to be regenerated before it can be merged in.',
|
||||
'',
|
||||
'Run `task engine:tool-models` to regenerate, then commit the updated file.',
|
||||
].join('\n');
|
||||
const { data: comments } = await github.rest.issues.listComments({
|
||||
owner: context.repo.owner,
|
||||
repo: context.repo.repo,
|
||||
issue_number: context.issue.number,
|
||||
});
|
||||
const existing = comments.find(c => c.body.includes(marker));
|
||||
if (existing) {
|
||||
await github.rest.issues.updateComment({
|
||||
owner: context.repo.owner,
|
||||
repo: context.repo.repo,
|
||||
comment_id: existing.id,
|
||||
body,
|
||||
});
|
||||
} else {
|
||||
await github.rest.issues.createComment({
|
||||
owner: context.repo.owner,
|
||||
repo: context.repo.repo,
|
||||
issue_number: context.issue.number,
|
||||
body,
|
||||
});
|
||||
}
|
||||
|
||||
- name: Fail if tool models check failed
|
||||
if: steps.tool-models-check.outcome == 'failure'
|
||||
run: |
|
||||
if ! git diff --exit-code engine/src/stirling/models/tool_models.py; then
|
||||
echo "tool_models.py is out of date."
|
||||
echo "Run 'task engine:tool-models' locally and commit the updated file."
|
||||
exit 1
|
||||
fi
|
||||
echo "============================================"
|
||||
echo " Tool Models Check Failed"
|
||||
echo "============================================"
|
||||
echo ""
|
||||
echo "The generated engine/src/stirling/models/tool_models.py"
|
||||
echo "is out of date with the Java OpenAPI spec and will"
|
||||
echo "need to be regenerated before it can be merged in."
|
||||
echo ""
|
||||
echo "Run 'task engine:tool-models' to regenerate, then"
|
||||
echo "commit the updated file."
|
||||
echo "============================================"
|
||||
exit 1
|
||||
|
||||
- name: Run fixers
|
||||
run: task engine:fix
|
||||
- name: Remove tool models check comment on success
|
||||
if: steps.tool-models-check.outcome == 'success' && github.event_name == 'pull_request'
|
||||
continue-on-error: true
|
||||
uses: actions/github-script@3a2844b7e9c422d3c10d287c895573f7108da1b3 # v9.0.0
|
||||
with:
|
||||
script: |
|
||||
const marker = '<!-- tool-models-check -->';
|
||||
const { data: comments } = await github.rest.issues.listComments({
|
||||
owner: context.repo.owner,
|
||||
repo: context.repo.repo,
|
||||
issue_number: context.issue.number,
|
||||
});
|
||||
const existing = comments.find(c => c.body.includes(marker));
|
||||
if (existing) {
|
||||
await github.rest.issues.deleteComment({
|
||||
owner: context.repo.owner,
|
||||
repo: context.repo.repo,
|
||||
comment_id: existing.id,
|
||||
});
|
||||
}
|
||||
|
||||
- name: Verify fixes are committed
|
||||
id: fixer_changes
|
||||
run: |
|
||||
if ! git diff --quiet; then
|
||||
git --no-pager diff --stat
|
||||
echo "::error::There are issues with your Python code that will need to be fixed before they can be merged in. Run 'task engine:fix' to auto-fix what can be fixed automatically, then run 'task engine:check' to see what still needs fixing manually."
|
||||
exit 1
|
||||
fi
|
||||
- name: Quality-check engine
|
||||
id: engine-check
|
||||
run: task engine:check
|
||||
continue-on-error: true
|
||||
|
||||
- name: Comment on fixer failures
|
||||
if: steps.fixer_changes.outcome == 'failure' && github.event_name == 'pull_request'
|
||||
- name: Comment on engine check failure
|
||||
# Only post a comment on PRs. github-script's PR helpers need an
|
||||
# issue/PR number, which doesn't exist on merge_group runs.
|
||||
if: steps.engine-check.outcome == 'failure' && github.event_name == 'pull_request'
|
||||
continue-on-error: true
|
||||
uses: actions/github-script@3a2844b7e9c422d3c10d287c895573f7108da1b3 # v9.0.0
|
||||
with:
|
||||
@@ -107,11 +176,39 @@ jobs:
|
||||
});
|
||||
}
|
||||
|
||||
- name: Run linting
|
||||
run: task engine:lint
|
||||
- name: Fail if engine check failed
|
||||
if: steps.engine-check.outcome == 'failure'
|
||||
run: |
|
||||
echo "============================================"
|
||||
echo " Engine Check Failed"
|
||||
echo "============================================"
|
||||
echo ""
|
||||
echo "There are issues with your Python code that"
|
||||
echo "will need to be fixed before they can be merged in."
|
||||
echo ""
|
||||
echo "Run 'task engine:fix' to auto-fix what can be"
|
||||
echo "fixed automatically, then run 'task engine:check'"
|
||||
echo "to see what still needs fixing manually."
|
||||
echo "============================================"
|
||||
exit 1
|
||||
|
||||
- name: Run type checking
|
||||
run: task engine:typecheck
|
||||
|
||||
- name: Run tests
|
||||
run: task engine:test
|
||||
- name: Remove engine check comment on success
|
||||
if: steps.engine-check.outcome == 'success' && github.event_name == 'pull_request'
|
||||
continue-on-error: true
|
||||
uses: actions/github-script@3a2844b7e9c422d3c10d287c895573f7108da1b3 # v9.0.0
|
||||
with:
|
||||
script: |
|
||||
const marker = '<!-- engine-check -->';
|
||||
const { data: comments } = await github.rest.issues.listComments({
|
||||
owner: context.repo.owner,
|
||||
repo: context.repo.repo,
|
||||
issue_number: context.issue.number,
|
||||
});
|
||||
const existing = comments.find(c => c.body.includes(marker));
|
||||
if (existing) {
|
||||
await github.rest.issues.deleteComment({
|
||||
owner: context.repo.owner,
|
||||
repo: context.repo.repo,
|
||||
comment_id: existing.id,
|
||||
});
|
||||
}
|
||||
|
||||
@@ -67,7 +67,7 @@ jobs:
|
||||
MAVEN_PASSWORD: ${{ secrets.MAVEN_PASSWORD }}
|
||||
MAVEN_PUBLIC_URL: ${{ secrets.MAVEN_PUBLIC_URL }}
|
||||
|
||||
- name: Comment on Java formatting failure
|
||||
- name: Comment on backend format check failure
|
||||
# Only post a comment on PRs. github-script's PR helpers need an
|
||||
# issue/PR number, which doesn't exist on merge_group runs.
|
||||
if: steps.spotless-check.outcome == 'failure' && github.event_name == 'pull_request'
|
||||
@@ -78,15 +78,11 @@ jobs:
|
||||
const marker = '<!-- java-formatting-check -->';
|
||||
const body = [
|
||||
marker,
|
||||
'### Java Formatting Check Failed',
|
||||
'### Backend Format Check Failed',
|
||||
'',
|
||||
'Your code has formatting issues. Run the following command to fix them:',
|
||||
'There are formatting issues in your Java code that will need to be fixed before they can be merged in.',
|
||||
'',
|
||||
'```bash',
|
||||
'task backend:format',
|
||||
'```',
|
||||
'',
|
||||
'Then commit and push the changes.',
|
||||
'Run `task backend:format` to auto-fix, then commit and push the changes.',
|
||||
].join('\n');
|
||||
const { data: comments } = await github.rest.issues.listComments({
|
||||
owner: context.repo.owner,
|
||||
@@ -110,22 +106,43 @@ jobs:
|
||||
});
|
||||
}
|
||||
|
||||
- name: Fail if Java formatting issues found
|
||||
- name: Fail if backend format check failed
|
||||
if: steps.spotless-check.outcome == 'failure'
|
||||
run: |
|
||||
echo "============================================"
|
||||
echo " Java Formatting Check Failed"
|
||||
echo " Backend Format Check Failed"
|
||||
echo "============================================"
|
||||
echo ""
|
||||
echo "Your code has formatting issues."
|
||||
echo "Run the following command to fix them:"
|
||||
echo "There are formatting issues in your Java code"
|
||||
echo "that will need to be fixed before they can be"
|
||||
echo "merged in."
|
||||
echo ""
|
||||
echo " task backend:format"
|
||||
echo ""
|
||||
echo "Then commit and push the changes."
|
||||
echo "Run 'task backend:format' to auto-fix, then"
|
||||
echo "commit and push the changes."
|
||||
echo "============================================"
|
||||
exit 1
|
||||
|
||||
- name: Remove backend format check comment on success
|
||||
if: steps.spotless-check.outcome == 'success' && github.event_name == 'pull_request'
|
||||
continue-on-error: true
|
||||
uses: actions/github-script@3a2844b7e9c422d3c10d287c895573f7108da1b3 # v9.0.0
|
||||
with:
|
||||
script: |
|
||||
const marker = '<!-- java-formatting-check -->';
|
||||
const { data: comments } = await github.rest.issues.listComments({
|
||||
owner: context.repo.owner,
|
||||
repo: context.repo.repo,
|
||||
issue_number: context.issue.number,
|
||||
});
|
||||
const existing = comments.find(c => c.body.includes(marker));
|
||||
if (existing) {
|
||||
await github.rest.issues.deleteComment({
|
||||
owner: context.repo.owner,
|
||||
repo: context.repo.repo,
|
||||
comment_id: existing.id,
|
||||
});
|
||||
}
|
||||
|
||||
- name: Build with Gradle and spring security ${{ matrix.spring-security }}
|
||||
run: task backend:build:ci
|
||||
env:
|
||||
|
||||
@@ -68,6 +68,18 @@ jobs:
|
||||
uses: ./.github/workflows/backend-build.yml
|
||||
secrets: inherit
|
||||
|
||||
db-migration-test:
|
||||
# Boots the current bootJar against H2 fixtures captured from past
|
||||
# releases (v2.0.0 / v2.5.0 / v2.10.0) and verifies admin login still
|
||||
# works after Hibernate's ddl-auto=update migrates the schema. Gated on
|
||||
# the `project` filter so doc-only PRs skip this ~5-minute job.
|
||||
if: needs.files-changed.outputs.project == 'true'
|
||||
needs: [files-changed]
|
||||
permissions:
|
||||
contents: read
|
||||
uses: ./.github/workflows/db-migration-test.yml
|
||||
secrets: inherit
|
||||
|
||||
check-generateOpenApiDocs:
|
||||
if: needs.files-changed.outputs.openapi == 'true'
|
||||
needs: [files-changed]
|
||||
@@ -184,6 +196,7 @@ jobs:
|
||||
needs:
|
||||
- files-changed
|
||||
- build
|
||||
- db-migration-test
|
||||
- check-generateOpenApiDocs
|
||||
- frontend-validation
|
||||
- playwright-e2e
|
||||
@@ -208,6 +221,7 @@ jobs:
|
||||
RESULTS: |
|
||||
files-changed=${{ needs.files-changed.result }}
|
||||
build=${{ needs.build.result }}
|
||||
db-migration-test=${{ needs.db-migration-test.result }}
|
||||
check-generateOpenApiDocs=${{ needs.check-generateOpenApiDocs.result }}
|
||||
frontend-validation=${{ needs.frontend-validation.result }}
|
||||
playwright-e2e=${{ needs.playwright-e2e.result }}
|
||||
|
||||
@@ -0,0 +1,93 @@
|
||||
name: DB migration smoke test
|
||||
|
||||
# Boots the current Stirling-PDF JAR against H2 fixtures captured from past
|
||||
# releases (v2.0.0 / v2.5.0 / v2.10.0) and verifies admin login still works.
|
||||
# Catches schema changes that would break existing user databases under
|
||||
# Hibernate's `ddl-auto=update` upgrade path.
|
||||
|
||||
on:
|
||||
workflow_call:
|
||||
|
||||
permissions:
|
||||
contents: read
|
||||
|
||||
jobs:
|
||||
pick:
|
||||
uses: ./.github/workflows/_runner-pick.yml
|
||||
|
||||
migration-test:
|
||||
needs: pick
|
||||
runs-on: ${{ needs.pick.outputs.is_fork == 'true' && 'ubuntu-latest' || 'depot-ubuntu-24.04-8' }}
|
||||
timeout-minutes: 30
|
||||
env:
|
||||
DEPOT_TOKEN: ${{ secrets.DEPOT_TOKEN }}
|
||||
steps:
|
||||
- name: Harden Runner
|
||||
uses: step-security/harden-runner@ab7a9404c0f3da075243ca237b5fac12c98deaa5 # v2.19.3
|
||||
with:
|
||||
egress-policy: audit
|
||||
|
||||
- name: Checkout repository
|
||||
uses: actions/checkout@de0fac2e4500dabe0009e67214ff5f5447ce83dd # v6.0.2
|
||||
|
||||
- name: Set up JDK 25
|
||||
uses: actions/setup-java@be666c2fcd27ec809703dec50e508c2fdc7f6654 # v5.2.0
|
||||
with:
|
||||
java-version: 25
|
||||
distribution: temurin
|
||||
|
||||
- name: Cache Gradle dependency artifacts
|
||||
uses: actions/cache@27d5ce7f107fe9357f9df03efb73ab90386fccae # v5.0.5
|
||||
with:
|
||||
path: |
|
||||
~/.gradle/wrapper
|
||||
~/.gradle/caches/modules-2/files-2.1
|
||||
~/.gradle/caches/modules-2/metadata-2.*
|
||||
key: gradle-deps-${{ runner.os }}-jdk-25-${{ hashFiles('**/gradle/wrapper/gradle-wrapper.properties', '**/*.gradle', '**/*.gradle.kts', 'settings.gradle', 'settings.gradle.kts', 'gradle/libs.versions.toml') }}
|
||||
|
||||
- name: Setup Gradle
|
||||
uses: gradle/actions/setup-gradle@50e97c2cd7a37755bbfafc9c5b7cafaece252f6e # v6.1.0
|
||||
with:
|
||||
gradle-version: 9.3.1
|
||||
cache-disabled: true
|
||||
|
||||
# No `-PnoSpotless` here yet because the upstream cache layer matches the
|
||||
# backend build's; reuse keeps cold-cache cost identical.
|
||||
- name: Build Stirling-PDF JAR
|
||||
env:
|
||||
MAVEN_USER: ${{ secrets.MAVEN_USER }}
|
||||
MAVEN_PASSWORD: ${{ secrets.MAVEN_PASSWORD }}
|
||||
MAVEN_PUBLIC_URL: ${{ secrets.MAVEN_PUBLIC_URL }}
|
||||
run: ./gradlew :stirling-pdf:bootJar -PnoSpotless --no-daemon
|
||||
|
||||
- name: Locate built JAR
|
||||
id: jar
|
||||
run: |
|
||||
jar=$(find app/core/build/libs -maxdepth 1 -name 'Stirling-PDF*.jar' -o -name 'stirling-pdf*.jar' 2>/dev/null \
|
||||
| grep -vE '(-plain|-sources)\.jar$' | head -n 1)
|
||||
if [[ -z "$jar" ]]; then
|
||||
echo "::error::No JAR under app/core/build/libs"
|
||||
ls -lah app/core/build/libs || true
|
||||
exit 1
|
||||
fi
|
||||
# Absolute path - the migration script pushd's into a temp workdir
|
||||
# before invoking java, which would dangle a relative path.
|
||||
jar=$(realpath "$jar")
|
||||
echo "path=$jar" >> "$GITHUB_OUTPUT"
|
||||
echo "Built JAR: $jar"
|
||||
|
||||
- name: Run migration smoke test
|
||||
env:
|
||||
STIRLING_JAR: ${{ steps.jar.outputs.path }}
|
||||
run: bash scripts/db-migration/run-migration-test.sh
|
||||
|
||||
- name: Upload app logs on failure
|
||||
if: failure()
|
||||
uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a # v7.0.1
|
||||
with:
|
||||
name: db-migration-app-logs
|
||||
# Path matches the preserved workdir in run-migration-test.sh -
|
||||
# only failing fixtures leave a directory behind.
|
||||
path: /tmp/stirling-migration-failed-*/app.log
|
||||
retention-days: 7
|
||||
if-no-files-found: warn
|
||||
@@ -586,21 +586,23 @@ jobs:
|
||||
if: always() && steps.digicert-setup.conclusion != 'failure'
|
||||
shell: bash
|
||||
run: |
|
||||
mkdir -p ./dist
|
||||
# Absolute dist path so the cd below can't break the copy targets.
|
||||
DIST="$GITHUB_WORKSPACE/dist"
|
||||
mkdir -p "$DIST"
|
||||
cd ./frontend/editor/src-tauri/target
|
||||
|
||||
# Find and rename artifacts based on platform
|
||||
if [ "${{ matrix.platform }}" = "windows-latest" ]; then
|
||||
# Only ship the MSI installer on Windows. The loose exe and WiX toolset exes
|
||||
# are not the user-facing installer - the MSI contains the signed inner exe.
|
||||
find . -name "*.msi" -exec cp {} "../../../dist/Stirling-PDF-${{ matrix.name }}.msi" \;
|
||||
find . -name "*.msi" -exec cp {} "$DIST/Stirling-PDF-${{ matrix.name }}.msi" \;
|
||||
elif [ "${{ matrix.platform }}" = "macos-15" ]; then
|
||||
find . -name "*.dmg" -exec cp {} "../../../dist/Stirling-PDF-${{ matrix.name }}.dmg" \;
|
||||
find . -name "*.app" -exec cp -r {} "../../../dist/Stirling-PDF-${{ matrix.name }}.app" \;
|
||||
find . -name "*.dmg" -exec cp {} "$DIST/Stirling-PDF-${{ matrix.name }}.dmg" \;
|
||||
find . -name "*.app" -exec cp -r {} "$DIST/Stirling-PDF-${{ matrix.name }}.app" \;
|
||||
else
|
||||
find . -name "*.deb" -exec cp {} "../../../dist/Stirling-PDF-${{ matrix.name }}.deb" \;
|
||||
find . -name "*.rpm" -exec cp {} "../../../dist/Stirling-PDF-${{ matrix.name }}.rpm" \;
|
||||
find . -name "*.AppImage" -exec cp {} "../../../dist/Stirling-PDF-${{ matrix.name }}.AppImage" \;
|
||||
find . -name "*.deb" -exec cp {} "$DIST/Stirling-PDF-${{ matrix.name }}.deb" \;
|
||||
find . -name "*.rpm" -exec cp {} "$DIST/Stirling-PDF-${{ matrix.name }}.rpm" \;
|
||||
find . -name "*.AppImage" -exec cp {} "$DIST/Stirling-PDF-${{ matrix.name }}.AppImage" \;
|
||||
fi
|
||||
|
||||
- name: Upload build artifacts
|
||||
|
||||
@@ -157,6 +157,9 @@ jobs:
|
||||
JPDFIUM_PLATFORMS: ${{ matrix.jpdfium_platforms }}
|
||||
run: task desktop:prepare
|
||||
|
||||
- name: Run Tauri/Cargo tests
|
||||
run: task desktop:test
|
||||
|
||||
# DigiCert KeyLocker Setup (Cloud HSM)
|
||||
- name: Setup DigiCert KeyLocker
|
||||
id: digicert-setup
|
||||
@@ -417,20 +420,22 @@ jobs:
|
||||
- name: Rename artifacts
|
||||
shell: bash
|
||||
run: |
|
||||
mkdir -p ./dist
|
||||
# Absolute dist path so the cd below can't break the copy targets.
|
||||
DIST="$GITHUB_WORKSPACE/dist"
|
||||
mkdir -p "$DIST"
|
||||
cd ./frontend/editor/src-tauri/target
|
||||
|
||||
# Find and rename artifacts based on platform
|
||||
if [ "${{ matrix.platform }}" = "windows-latest" ]; then
|
||||
# Only ship the MSI installer. The loose exe and WiX toolset exes
|
||||
# are not the user-facing installer - the MSI contains the signed inner exe.
|
||||
find . -name "*.msi" -exec cp {} "../../../dist/Stirling-PDF-${{ matrix.name }}.msi" \;
|
||||
find . -name "*.msi" -exec cp {} "$DIST/Stirling-PDF-${{ matrix.name }}.msi" \;
|
||||
elif [ "${{ matrix.platform }}" = "macos-15" ]; then
|
||||
find . -name "*.dmg" -exec cp {} "../../../dist/Stirling-PDF-${{ matrix.name }}.dmg" \;
|
||||
find . -name "*.dmg" -exec cp {} "$DIST/Stirling-PDF-${{ matrix.name }}.dmg" \;
|
||||
else
|
||||
find . -name "*.deb" -exec cp {} "../../../dist/Stirling-PDF-${{ matrix.name }}.deb" \;
|
||||
find . -name "*.rpm" -exec cp {} "../../../dist/Stirling-PDF-${{ matrix.name }}.rpm" \;
|
||||
find . -name "*.AppImage" -exec cp {} "../../../dist/Stirling-PDF-${{ matrix.name }}.AppImage" \;
|
||||
find . -name "*.deb" -exec cp {} "$DIST/Stirling-PDF-${{ matrix.name }}.deb" \;
|
||||
find . -name "*.rpm" -exec cp {} "$DIST/Stirling-PDF-${{ matrix.name }}.rpm" \;
|
||||
find . -name "*.AppImage" -exec cp {} "$DIST/Stirling-PDF-${{ matrix.name }}.AppImage" \;
|
||||
fi
|
||||
|
||||
# Verify the MSI AND the inner exe extracted from it are signed.
|
||||
|
||||
+8
-3
@@ -23,6 +23,10 @@ customFiles/
|
||||
configs/
|
||||
watchedFolders/
|
||||
clientWebUI/
|
||||
# Scratch dir used by local fixture-regeneration runs (see
|
||||
# app/proprietary/src/test/resources/db-migration-fixtures/README.md).
|
||||
# Holds downloaded JARs and disposable workdirs. Never committed.
|
||||
.alpha-local/
|
||||
!cucumber/
|
||||
!cucumber/exampleFiles/
|
||||
!cucumber/exampleFiles/example_html.zip
|
||||
@@ -157,9 +161,7 @@ app/proprietary/build
|
||||
common/build
|
||||
proprietary/build
|
||||
stirling-pdf/build
|
||||
frontend/src-tauri/provisioner/target
|
||||
frontend/src-tauri/target
|
||||
frontend/src-tauri/runtime/
|
||||
frontend/editor/src-tauri/provisioner/target
|
||||
|
||||
# Byte-compiled / optimized / DLL files
|
||||
__pycache__/
|
||||
@@ -275,3 +277,6 @@ docs/type3/signatures/
|
||||
# Playwright MCP screenshots / traces
|
||||
.playwright-mcp/
|
||||
*.playwright-mcp.png
|
||||
|
||||
# Local screenshot artifacts from *-screenshots.spec.ts
|
||||
frontend/screenshots/
|
||||
|
||||
@@ -22,12 +22,13 @@ tasks:
|
||||
vars:
|
||||
PORT: '{{.PORT | default "8080"}}'
|
||||
AIENGINE_URL: '{{.AIENGINE_URL | default ""}}'
|
||||
AIENGINE_TIMEOUTSECONDS: '{{.AIENGINE_TIMEOUTSECONDS | default "120"}}'
|
||||
env:
|
||||
SERVER_PORT: '{{.PORT}}'
|
||||
cmds:
|
||||
- cmd: '{{if .AIENGINE_URL}}AIENGINE_URL={{.AIENGINE_URL}} AIENGINE_ENABLED=true AIENGINE_TIMEOUTSECONDS=120 {{end}}cmd /c ".\gradlew.bat :stirling-pdf:bootRun"'
|
||||
- cmd: '{{if .AIENGINE_URL}}AIENGINE_URL={{.AIENGINE_URL}} AIENGINE_ENABLED=true AIENGINE_TIMEOUTSECONDS={{.AIENGINE_TIMEOUTSECONDS}} {{end}}cmd /c ".\gradlew.bat :stirling-pdf:bootRun"'
|
||||
platforms: [windows]
|
||||
- cmd: '{{if .AIENGINE_URL}}AIENGINE_URL={{.AIENGINE_URL}} AIENGINE_ENABLED=true AIENGINE_TIMEOUTSECONDS=120 {{end}}./gradlew :stirling-pdf:bootRun'
|
||||
- cmd: '{{if .AIENGINE_URL}}AIENGINE_URL={{.AIENGINE_URL}} AIENGINE_ENABLED=true AIENGINE_TIMEOUTSECONDS={{.AIENGINE_TIMEOUTSECONDS}} {{end}}./gradlew :stirling-pdf:bootRun'
|
||||
platforms: [linux, darwin]
|
||||
|
||||
dev:bundled:
|
||||
@@ -40,10 +41,10 @@ tasks:
|
||||
platforms: [linux, darwin]
|
||||
|
||||
dev:saas:
|
||||
desc: "Start backend in SaaS flavor against Supabase (loads .env.saas.local)"
|
||||
desc: "Start backend in SaaS flavor against Supabase"
|
||||
# `dotenv:` reads from the root Taskfile's directory (".") because this
|
||||
# subtaskfile is included with `dir: .`. Drop the file at the repo root.
|
||||
dotenv: ['.env.saas.local']
|
||||
# subtaskfile is included with `dir: .`.
|
||||
dotenv: ['app/.env.saas.local', 'app/.env.saas']
|
||||
ignore_error: true
|
||||
vars:
|
||||
PORT: '{{.PORT | default "8080"}}'
|
||||
|
||||
@@ -78,6 +78,13 @@ tasks:
|
||||
cmds:
|
||||
- npx tauri build --bundles appimage
|
||||
|
||||
test:
|
||||
desc: "Run Tauri/Cargo tests"
|
||||
deps: [prepare]
|
||||
dir: editor/src-tauri
|
||||
cmds:
|
||||
- cargo test
|
||||
|
||||
clean:
|
||||
desc: "Clean Tauri/Cargo build artifacts"
|
||||
dir: editor
|
||||
|
||||
@@ -257,7 +257,6 @@ tasks:
|
||||
- task: lint
|
||||
- task: format:check
|
||||
- task: test
|
||||
- task: theme:check-unused-vars
|
||||
|
||||
check:all:
|
||||
desc: "Full CI quality gate"
|
||||
@@ -299,17 +298,3 @@ tasks:
|
||||
deps: [install]
|
||||
cmds:
|
||||
- node editor/scripts/generate-licenses.js
|
||||
|
||||
theme:clean-unused-vars:
|
||||
desc: "Remove unused CSS variables from theme.css"
|
||||
deps: [install]
|
||||
vars:
|
||||
REMOVE: '{{.REMOVE | default "false"}}'
|
||||
cmds:
|
||||
- node editor/scripts/clean-unused-theme-vars.js{{if eq .REMOVE "true"}} --remove{{end}} {{.CLI_ARGS}}
|
||||
|
||||
theme:check-unused-vars:
|
||||
desc: "Fail if theme.css contains unused CSS variables"
|
||||
deps: [install]
|
||||
cmds:
|
||||
- node editor/scripts/clean-unused-theme-vars.js --check
|
||||
|
||||
@@ -58,6 +58,23 @@ tasks:
|
||||
BACKEND_URL: 'http://localhost:{{.BACKEND_PORT}}'
|
||||
OPEN: "true"
|
||||
|
||||
dev:saas:
|
||||
desc: "Start SaaS backend + frontend concurrently on free ports"
|
||||
vars:
|
||||
PORTS:
|
||||
sh: '{{if eq OS "windows"}}{{.FIND_FREE_PORT_PS}} 8080 5173{{else}}{{.FIND_FREE_PORT_SH}} 8080 5173{{end}}'
|
||||
BACKEND_PORT: '{{index (splitList "\n" .PORTS) 0}}'
|
||||
FRONTEND_PORT: '{{index (splitList "\n" .PORTS) 1}}'
|
||||
deps:
|
||||
- task: backend:dev:saas
|
||||
vars:
|
||||
PORT: '{{.BACKEND_PORT}}'
|
||||
- task: frontend:dev:saas
|
||||
vars:
|
||||
PORT: '{{.FRONTEND_PORT}}'
|
||||
BACKEND_URL: 'http://localhost:{{.BACKEND_PORT}}'
|
||||
OPEN: "true"
|
||||
|
||||
dev:all:
|
||||
desc: "Start backend + frontend + engine concurrently on free ports"
|
||||
vars:
|
||||
|
||||
@@ -1,20 +1,20 @@
|
||||
###############################################################################
|
||||
# Stirling-PDF SaaS local environment template.
|
||||
# Stirling-PDF SaaS environment defaults.
|
||||
#
|
||||
# Copy this file to `.env.saas.local` (gitignored) and fill in real values.
|
||||
# Loaded by `task backend:dev:saas` via Taskfile's `dotenv:` directive, then
|
||||
# read by Spring Boot's `${...}` placeholders in application-saas.properties
|
||||
# and application-dev.properties.
|
||||
# This file is committed and provides non-secret defaults loaded by
|
||||
# `task backend:dev:saas`. Put real values for secrets (passwords, project
|
||||
# refs, edge function secrets) in `.env.saas.local` - any variable set there
|
||||
# takes precedence over what's defined here.
|
||||
#
|
||||
# DO NOT commit `.env.saas.local`. Only `.env.saas.example` is checked in.
|
||||
# DO NOT commit `.env.saas.local`. Only `.env.saas` is checked in.
|
||||
###############################################################################
|
||||
|
||||
# ---------- Supabase project ----------
|
||||
# Project reference (the subdomain part of <ref>.supabase.co). Required.
|
||||
# Example dev project:
|
||||
# Set in .env.saas.local.
|
||||
SAAS_DB_PROJECT_REF=
|
||||
|
||||
# Edge function secret used by billing/license rollup calls.
|
||||
# Edge function secret used by billing/license rollup calls. Set in .env.saas.local.
|
||||
SUPABASE_EDGE_FUNCTION_SECRET=
|
||||
|
||||
# ---------- Database (saas profile) ----------
|
||||
@@ -28,7 +28,7 @@ SAAS_DB_PASSWORD=
|
||||
# ---------- Database (dev profile overrides) ----------
|
||||
# Used when `--spring.profiles.include=dev` is active. The dev profile
|
||||
# defaults the URL/username to the shared dev Supabase project, but the
|
||||
# password must still be provided here.
|
||||
# password must still be provided in .env.saas.local.
|
||||
SAAS_DEV_DB_URL=
|
||||
SAAS_DEV_DB_USERNAME=postgres
|
||||
SAAS_DEV_DB_PASSWORD=
|
||||
@@ -0,0 +1,3 @@
|
||||
# Whitelist committed env defaults. `.env.saas.local` (and any other .env*)
|
||||
# stays ignored via the root .gitignore.
|
||||
!.env.saas
|
||||
@@ -44,6 +44,10 @@
|
||||
"moduleName": ".*",
|
||||
"moduleLicense": "The MIT License"
|
||||
},
|
||||
{
|
||||
"moduleName": ".*",
|
||||
"moduleLicense": "MIT-0"
|
||||
},
|
||||
{
|
||||
"moduleName": "com.github.jai-imageio:jai-imageio-core",
|
||||
"moduleLicense": "LICENSE.txt"
|
||||
|
||||
@@ -78,6 +78,9 @@ dependencies {
|
||||
runtimeOnly "com.stirling:jpdfium-natives-${platform}:1.0.1"
|
||||
}
|
||||
|
||||
// Bucket4j (local in-process token bucket for RateLimitStore default impl)
|
||||
implementation 'com.bucket4j:bucket4j_jdk17-core:8.19.0'
|
||||
|
||||
// ArchUnit: enforces module dependency direction (see ArchitectureTest)
|
||||
testImplementation 'com.tngtech.archunit:archunit-junit5:1.4.2'
|
||||
}
|
||||
|
||||
+5
-1
@@ -77,6 +77,10 @@ public @interface AutoJobPostMapping {
|
||||
/**
|
||||
* Relative resource weight (1-100). See {@link
|
||||
* stirling.software.common.enumeration.ResourceWeight} for the standard tiers.
|
||||
*
|
||||
* <p>The default is a sentinel ({@link Integer#MIN_VALUE}); {@code
|
||||
* AutoJobPostMappingWeightTest} fails the build if any endpoint leaves it unset. Runtime
|
||||
* readers clamp the value into {@code [1, 100]}.
|
||||
*/
|
||||
int resourceWeight() default 1;
|
||||
int resourceWeight() default Integer.MIN_VALUE;
|
||||
}
|
||||
|
||||
@@ -0,0 +1,24 @@
|
||||
package stirling.software.common.cluster;
|
||||
|
||||
/** Health and identity facade for the active cluster backplane. */
|
||||
public interface ClusterBackplane {
|
||||
|
||||
/** Returns {@code true} when the backplane is reachable; used for health endpoints. */
|
||||
boolean isHealthy();
|
||||
|
||||
/** Returns {@code "inprocess"} or {@code "valkey"}. */
|
||||
String backplaneType();
|
||||
|
||||
/** Returns this JVM's stable node id (matches {@code Cluster.resolvedNodeId()}). */
|
||||
String localNodeId();
|
||||
|
||||
/**
|
||||
* Whether this JVM should run the local {@link
|
||||
* stirling.software.common.service.TaskManager#cleanupOldJobs()} loop. Distributed backplanes
|
||||
* own job expiry via their own TTL, so they should override this to return {@code false}.
|
||||
* Defaults to {@code true} so in-process behavior is preserved without an explicit override.
|
||||
*/
|
||||
default boolean shouldRunLocalCleanup() {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,63 @@
|
||||
package stirling.software.common.cluster;
|
||||
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
|
||||
import jakarta.annotation.PostConstruct;
|
||||
|
||||
import lombok.RequiredArgsConstructor;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
|
||||
import stirling.software.common.model.ApplicationProperties;
|
||||
import stirling.software.common.model.ApplicationProperties.Cluster;
|
||||
|
||||
/**
|
||||
* Validates that cluster mode is internally consistent.
|
||||
*
|
||||
* <p>Cluster settings are bound on the central {@link ApplicationProperties} under {@code
|
||||
* cluster.*}; this class reads {@link ApplicationProperties#getCluster()} and runs guards in {@link
|
||||
* PostConstruct}. When {@code cluster.enabled=false} (the default) all checks are skipped so a
|
||||
* single-instance install needs no new config.
|
||||
*/
|
||||
@Slf4j
|
||||
@Configuration
|
||||
@RequiredArgsConstructor
|
||||
public class ClusterConfig {
|
||||
|
||||
private final ApplicationProperties applicationProperties;
|
||||
|
||||
@PostConstruct
|
||||
void validate() {
|
||||
Cluster cluster = applicationProperties.getCluster();
|
||||
if (!cluster.isEnabled()) {
|
||||
return;
|
||||
}
|
||||
String backplane = cluster.getBackplane();
|
||||
if ("valkey".equalsIgnoreCase(backplane)) {
|
||||
String url = cluster.getValkey() == null ? null : cluster.getValkey().getUrl();
|
||||
if (url == null || url.isBlank()) {
|
||||
throw new IllegalStateException(
|
||||
"cluster.enabled=true with backplane=valkey requires"
|
||||
+ " cluster.valkey.url to be set (e.g."
|
||||
+ " redis://valkey:6379).");
|
||||
}
|
||||
} else if ("inprocess".equalsIgnoreCase(backplane)) {
|
||||
// enabled+inprocess only coordinates the local JVM; cross-node lookups will 410.
|
||||
log.warn(
|
||||
"cluster.enabled=true with backplane=inprocess - only the local"
|
||||
+ " JVM is coordinated. Cross-node lookups and the file proxy will fail."
|
||||
+ " Use backplane=valkey for real multi-node deployments.");
|
||||
} else {
|
||||
// Fail fast on typos like "valky" so Spring doesn't later report a cryptic
|
||||
// "no ClusterBackplane bean" - the operator-facing error names the bad value.
|
||||
throw new IllegalStateException(
|
||||
"cluster.enabled=true with unknown backplane '"
|
||||
+ backplane
|
||||
+ "'. Valid values: inprocess | valkey.");
|
||||
}
|
||||
log.info(
|
||||
"Cluster mode enabled (backplane={}, nodeRole={}, nodeId={}).",
|
||||
backplane,
|
||||
cluster.resolvedRole(),
|
||||
cluster.resolvedNodeId());
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,12 @@
|
||||
package stirling.software.common.cluster;
|
||||
|
||||
import java.time.Instant;
|
||||
|
||||
/**
|
||||
* Snapshot of a peer node as recorded in the {@link InstanceRegistry}.
|
||||
*
|
||||
* @param internalAddress {@code host:port} the node listens on for {@code /internal/cluster/**}
|
||||
* @param role one of {@code WEB}, {@code WORKER}, {@code BOTH}
|
||||
*/
|
||||
public record ClusterNode(
|
||||
String nodeId, String internalAddress, Instant lastHeartbeat, String role) {}
|
||||
@@ -0,0 +1,21 @@
|
||||
package stirling.software.common.cluster;
|
||||
|
||||
import java.time.Duration;
|
||||
import java.util.Optional;
|
||||
|
||||
/** Cluster-wide mutual exclusion primitive; non-reentrant by contract. */
|
||||
public interface DistributedLock {
|
||||
|
||||
Optional<LockHandle> tryAcquire(String lockKey, Duration leaseTime);
|
||||
|
||||
interface LockHandle extends AutoCloseable {
|
||||
void release();
|
||||
|
||||
boolean renew(Duration leaseTime);
|
||||
|
||||
@Override
|
||||
default void close() {
|
||||
release();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,45 @@
|
||||
package stirling.software.common.cluster;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.io.InputStream;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
|
||||
/** Low-level storage seam for result/job files. */
|
||||
public interface FileStore {
|
||||
|
||||
/** Stored file record. */
|
||||
record Stored(String fileId, long size) {}
|
||||
|
||||
/** Store the given stream and return a generated file id and total bytes written. */
|
||||
Stored store(InputStream in, String originalName) throws IOException;
|
||||
|
||||
/**
|
||||
* Store the file at {@code source} and return a generated file id and total bytes written.
|
||||
*
|
||||
* <p>Default implementation opens {@code source} as a stream and delegates to {@link
|
||||
* #store(InputStream, String)}. Local-disk implementations should override to use a direct
|
||||
* file-to-file copy ({@code Files.copy(source, dest)} can use {@code sendfile(2)} on Linux),
|
||||
* which avoids the two-memory-copy hit of streaming a disk-backed upload through the JVM heap.
|
||||
*/
|
||||
default Stored store(Path source, String originalName) throws IOException {
|
||||
try (InputStream in = Files.newInputStream(source)) {
|
||||
return store(in, originalName);
|
||||
}
|
||||
}
|
||||
|
||||
/** Open the stored file for streaming reads. Caller closes. */
|
||||
InputStream retrieve(String fileId) throws IOException;
|
||||
|
||||
/** Load the stored file into a byte array. */
|
||||
byte[] retrieveBytes(String fileId) throws IOException;
|
||||
|
||||
/** Size of the stored file in bytes. */
|
||||
long size(String fileId) throws IOException;
|
||||
|
||||
/** Delete the stored file. Returns true if a file was removed. */
|
||||
boolean delete(String fileId);
|
||||
|
||||
/** Whether the file id exists in the store. */
|
||||
boolean exists(String fileId);
|
||||
}
|
||||
@@ -0,0 +1,18 @@
|
||||
package stirling.software.common.cluster;
|
||||
|
||||
import java.time.Duration;
|
||||
import java.util.Collection;
|
||||
import java.util.Optional;
|
||||
|
||||
/** Maps {@code nodeId} to its internal cluster address, with TTL'd heartbeats. */
|
||||
public interface InstanceRegistry {
|
||||
|
||||
/** Register or refresh this node. Idempotent so a wiped backplane self-heals on next tick. */
|
||||
void register(ClusterNode node, Duration heartbeatTtl);
|
||||
|
||||
Optional<ClusterNode> lookup(String nodeId);
|
||||
|
||||
Collection<ClusterNode> activeNodes();
|
||||
|
||||
void deregister(String nodeId);
|
||||
}
|
||||
@@ -0,0 +1,24 @@
|
||||
package stirling.software.common.cluster;
|
||||
|
||||
import java.time.Duration;
|
||||
import java.util.Collection;
|
||||
import java.util.Optional;
|
||||
|
||||
/** Cluster-visible storage for job status and result metadata, with TTL'd entries. */
|
||||
public interface JobStore {
|
||||
|
||||
/** Persist or overwrite a job entry. {@code ttl} sets the lifetime of the entry. */
|
||||
void put(JobStoreEntry entry, Duration ttl);
|
||||
|
||||
Optional<JobStoreEntry> get(String jobId);
|
||||
|
||||
void delete(String jobId);
|
||||
|
||||
boolean exists(String jobId);
|
||||
|
||||
/** Reverse lookup: which job owns this result file id? */
|
||||
Optional<String> findJobIdByFileId(String fileId);
|
||||
|
||||
/** Snapshot of every active entry. Used by admin/stats endpoints; may be O(n). */
|
||||
Collection<JobStoreEntry> all();
|
||||
}
|
||||
@@ -0,0 +1,30 @@
|
||||
package stirling.software.common.cluster;
|
||||
|
||||
import java.time.Instant;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
/**
|
||||
* Cluster-visible projection of a job's status and result metadata, as persisted in {@link
|
||||
* JobStore}.
|
||||
*
|
||||
* @param owningNodeId the node id that originally executed the job
|
||||
*/
|
||||
public record JobStoreEntry(
|
||||
String jobId,
|
||||
JobState state,
|
||||
String owningNodeId,
|
||||
Instant createdAt,
|
||||
Instant completedAt,
|
||||
String error,
|
||||
List<String> fileIds,
|
||||
Map<String, String> resultMeta) {
|
||||
|
||||
/** Lifecycle states for a job as observed by the cluster. */
|
||||
public enum JobState {
|
||||
PENDING,
|
||||
RUNNING,
|
||||
COMPLETE,
|
||||
FAILED
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,16 @@
|
||||
package stirling.software.common.cluster;
|
||||
|
||||
import java.time.Duration;
|
||||
import java.util.Optional;
|
||||
|
||||
/** Short-TTL namespaced key/value cache backed by the cluster backplane. */
|
||||
public interface KeyValueCache {
|
||||
|
||||
void put(String namespace, String key, String value, Duration ttl);
|
||||
|
||||
Optional<String> get(String namespace, String key);
|
||||
|
||||
void evict(String namespace, String key);
|
||||
|
||||
void evictNamespace(String namespace);
|
||||
}
|
||||
@@ -0,0 +1,18 @@
|
||||
package stirling.software.common.cluster;
|
||||
|
||||
import java.time.Duration;
|
||||
|
||||
/** Token-bucket rate limiting backed by the cluster backplane. */
|
||||
public interface RateLimitStore {
|
||||
|
||||
/**
|
||||
* Attempt to consume one token from the bucket identified by {@code bucketKey}.
|
||||
*
|
||||
* @param bucketKey opaque key identifying the bucket (e.g. {@code api:user:123})
|
||||
* @param capacity bucket capacity
|
||||
* @param refillPeriod time window over which {@code capacity} tokens refill
|
||||
*/
|
||||
RateLimitDecision tryConsume(String bucketKey, long capacity, Duration refillPeriod);
|
||||
|
||||
record RateLimitDecision(boolean allowed, long remainingTokens, long nanosToWaitForRefill) {}
|
||||
}
|
||||
@@ -0,0 +1,7 @@
|
||||
package stirling.software.common.cluster;
|
||||
|
||||
/** Records one increment per sticky-session miss (a 410 Gone for a job owned by another node). */
|
||||
@FunctionalInterface
|
||||
public interface StickyMissRecorder {
|
||||
void recordStickyMiss();
|
||||
}
|
||||
+31
@@ -0,0 +1,31 @@
|
||||
package stirling.software.common.cluster.inprocess;
|
||||
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
|
||||
import stirling.software.common.cluster.ClusterBackplane;
|
||||
import stirling.software.common.model.ApplicationProperties;
|
||||
|
||||
@Slf4j
|
||||
public class InProcessClusterBackplane implements ClusterBackplane {
|
||||
|
||||
private final ApplicationProperties applicationProperties;
|
||||
|
||||
public InProcessClusterBackplane(ApplicationProperties applicationProperties) {
|
||||
this.applicationProperties = applicationProperties;
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean isHealthy() {
|
||||
return true;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String backplaneType() {
|
||||
return "inprocess";
|
||||
}
|
||||
|
||||
@Override
|
||||
public String localNodeId() {
|
||||
return applicationProperties.getCluster().resolvedNodeId();
|
||||
}
|
||||
}
|
||||
+65
@@ -0,0 +1,65 @@
|
||||
package stirling.software.common.cluster.inprocess;
|
||||
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
|
||||
import stirling.software.common.cluster.ClusterBackplane;
|
||||
import stirling.software.common.cluster.DistributedLock;
|
||||
import stirling.software.common.cluster.InstanceRegistry;
|
||||
import stirling.software.common.cluster.JobStore;
|
||||
import stirling.software.common.cluster.KeyValueCache;
|
||||
import stirling.software.common.cluster.RateLimitStore;
|
||||
import stirling.software.common.model.ApplicationProperties;
|
||||
|
||||
/**
|
||||
* Default cluster backplane wiring: every interface gets an {@code InProcess*} bean. Active when
|
||||
* cluster mode is off or {@code cluster.backplane=inprocess}.
|
||||
*/
|
||||
@Slf4j
|
||||
@Configuration
|
||||
@ConditionalOnExpression(
|
||||
"!${cluster.enabled:false} ||"
|
||||
+ " '${cluster.backplane:inprocess}'.equalsIgnoreCase('inprocess')")
|
||||
public class InProcessClusterConfiguration {
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean
|
||||
public ClusterBackplane clusterBackplane(ApplicationProperties applicationProperties) {
|
||||
log.info("Cluster backplane: in-process (single node)");
|
||||
return new InProcessClusterBackplane(applicationProperties);
|
||||
}
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean
|
||||
public JobStore jobStore() {
|
||||
return new InProcessJobStore();
|
||||
}
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean
|
||||
public RateLimitStore rateLimitStore() {
|
||||
return new InProcessRateLimitStore();
|
||||
}
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean
|
||||
public DistributedLock distributedLock() {
|
||||
return new InProcessDistributedLock();
|
||||
}
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean
|
||||
public KeyValueCache keyValueCache() {
|
||||
return new InProcessKeyValueCache();
|
||||
}
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean
|
||||
public InstanceRegistry instanceRegistry() {
|
||||
return new InProcessInstanceRegistry();
|
||||
}
|
||||
}
|
||||
+128
@@ -0,0 +1,128 @@
|
||||
package stirling.software.common.cluster.inprocess;
|
||||
|
||||
import java.time.Duration;
|
||||
import java.util.Optional;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
|
||||
import stirling.software.common.cluster.DistributedLock;
|
||||
|
||||
/**
|
||||
* In-process {@link DistributedLock}, non-reentrant per the interface contract, with lease-expiry
|
||||
* semantics that mirror a SET-NX-EX style distributed backend.
|
||||
*
|
||||
* <p>Each lock state carries a per-acquire {@code ownerToken} and an {@code expiryNanos}; another
|
||||
* caller can take over once the lease has elapsed even if the original holder never called {@link
|
||||
* LockHandle#release()}. This matters mostly for parity with the Valkey-backed implementation
|
||||
* (Redis {@code SETEX} auto-expires the key); within a single JVM a crashed holder takes its lock
|
||||
* state with it, but tests and code that rely on the {@code leaseTime} parameter still need it to
|
||||
* be honored.
|
||||
*
|
||||
* <p>Expiry is lazy: an expired lock state lingers in the map until the next acquire attempt for
|
||||
* the same key replaces it. Per-key cleanup also happens on explicit {@link LockHandle#release()},
|
||||
* so a balanced acquire/release workload keeps the map size bounded.
|
||||
*/
|
||||
public class InProcessDistributedLock implements DistributedLock {
|
||||
|
||||
private final ConcurrentHashMap<String, LockState> locks = new ConcurrentHashMap<>();
|
||||
private final AtomicLong tokenSeq = new AtomicLong();
|
||||
|
||||
/**
|
||||
* Lease state for a single acquired lock. {@code ownerToken} prevents a former holder from
|
||||
* releasing or renewing a lock now owned by someone else after lease expiry; {@code
|
||||
* expiryNanos} is read/written only inside {@link ConcurrentHashMap#compute} so the bin lock
|
||||
* provides the necessary happens-before guarantee.
|
||||
*/
|
||||
private static final class LockState {
|
||||
final long ownerToken;
|
||||
long expiryNanos;
|
||||
|
||||
LockState(long ownerToken, long expiryNanos) {
|
||||
this.ownerToken = ownerToken;
|
||||
this.expiryNanos = expiryNanos;
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public Optional<LockHandle> tryAcquire(String lockKey, Duration leaseTime) {
|
||||
long token = tokenSeq.incrementAndGet();
|
||||
long nowNanos = System.nanoTime();
|
||||
long expiryNanos = nowNanos + leaseTime.toNanos();
|
||||
boolean[] acquired = {false};
|
||||
locks.compute(
|
||||
lockKey,
|
||||
(k, existing) -> {
|
||||
if (existing == null || existing.expiryNanos - nowNanos <= 0L) {
|
||||
// No lock, or the previous lease has expired - we take it. Subtraction
|
||||
// form avoids the long-overflow trap that would bite a naive
|
||||
// expiryNanos <= nowNanos comparison around System.nanoTime() rollover.
|
||||
acquired[0] = true;
|
||||
return new LockState(token, expiryNanos);
|
||||
}
|
||||
return existing;
|
||||
});
|
||||
if (!acquired[0]) {
|
||||
return Optional.empty();
|
||||
}
|
||||
return Optional.of(new InProcessHandle(lockKey, token));
|
||||
}
|
||||
|
||||
private void releaseInternal(String lockKey, long token) {
|
||||
locks.compute(
|
||||
lockKey,
|
||||
(k, existing) -> {
|
||||
if (existing == null || existing.ownerToken != token) {
|
||||
// Already removed, expired-and-replaced, or never ours.
|
||||
return existing;
|
||||
}
|
||||
return null;
|
||||
});
|
||||
}
|
||||
|
||||
private boolean renewInternal(String lockKey, long token, Duration leaseTime) {
|
||||
long nowNanos = System.nanoTime();
|
||||
boolean[] renewed = {false};
|
||||
locks.compute(
|
||||
lockKey,
|
||||
(k, existing) -> {
|
||||
if (existing == null
|
||||
|| existing.ownerToken != token
|
||||
|| existing.expiryNanos - nowNanos <= 0L) {
|
||||
// Lock is gone or expired; renewal is a no-op so the caller can detect it.
|
||||
return existing;
|
||||
}
|
||||
existing.expiryNanos = nowNanos + leaseTime.toNanos();
|
||||
renewed[0] = true;
|
||||
return existing;
|
||||
});
|
||||
return renewed[0];
|
||||
}
|
||||
|
||||
private final class InProcessHandle implements LockHandle {
|
||||
private final String lockKey;
|
||||
private final long token;
|
||||
private boolean released;
|
||||
|
||||
InProcessHandle(String lockKey, long token) {
|
||||
this.lockKey = lockKey;
|
||||
this.token = token;
|
||||
}
|
||||
|
||||
@Override
|
||||
public synchronized void release() {
|
||||
if (released) {
|
||||
return;
|
||||
}
|
||||
released = true;
|
||||
releaseInternal(lockKey, token);
|
||||
}
|
||||
|
||||
@Override
|
||||
public synchronized boolean renew(Duration leaseTime) {
|
||||
if (released) {
|
||||
return false;
|
||||
}
|
||||
return renewInternal(lockKey, token, leaseTime);
|
||||
}
|
||||
}
|
||||
}
|
||||
+42
@@ -0,0 +1,42 @@
|
||||
package stirling.software.common.cluster.inprocess;
|
||||
|
||||
import java.time.Duration;
|
||||
import java.util.Collection;
|
||||
import java.util.Collections;
|
||||
import java.util.Optional;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
|
||||
import stirling.software.common.cluster.ClusterNode;
|
||||
import stirling.software.common.cluster.InstanceRegistry;
|
||||
|
||||
public class InProcessInstanceRegistry implements InstanceRegistry {
|
||||
|
||||
private final AtomicReference<ClusterNode> self = new AtomicReference<>();
|
||||
|
||||
@Override
|
||||
public void register(ClusterNode node, Duration heartbeatTtl) {
|
||||
self.set(node);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Optional<ClusterNode> lookup(String nodeId) {
|
||||
ClusterNode current = self.get();
|
||||
return current != null && current.nodeId().equals(nodeId)
|
||||
? Optional.of(current)
|
||||
: Optional.empty();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Collection<ClusterNode> activeNodes() {
|
||||
ClusterNode current = self.get();
|
||||
return current == null ? Collections.emptyList() : Collections.singletonList(current);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void deregister(String nodeId) {
|
||||
ClusterNode current = self.get();
|
||||
if (current != null && current.nodeId().equals(nodeId)) {
|
||||
self.set(null);
|
||||
}
|
||||
}
|
||||
}
|
||||
+96
@@ -0,0 +1,96 @@
|
||||
package stirling.software.common.cluster.inprocess;
|
||||
|
||||
import java.time.Duration;
|
||||
import java.time.Instant;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collection;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
|
||||
import stirling.software.common.cluster.JobStore;
|
||||
import stirling.software.common.cluster.JobStoreEntry;
|
||||
|
||||
@Slf4j
|
||||
public class InProcessJobStore implements JobStore {
|
||||
|
||||
private final ConcurrentHashMap<String, Holder> entries = new ConcurrentHashMap<>();
|
||||
|
||||
@Override
|
||||
public void put(JobStoreEntry entry, Duration ttl) {
|
||||
Instant expiry = ttl == null ? Instant.MAX : Instant.now().plus(ttl);
|
||||
entries.put(entry.jobId(), new Holder(entry, expiry));
|
||||
}
|
||||
|
||||
@Override
|
||||
public Optional<JobStoreEntry> get(String jobId) {
|
||||
Holder h = entries.get(jobId);
|
||||
if (h == null) {
|
||||
return Optional.empty();
|
||||
}
|
||||
if (h.isExpired()) {
|
||||
entries.remove(jobId, h);
|
||||
return Optional.empty();
|
||||
}
|
||||
return Optional.of(h.entry);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void delete(String jobId) {
|
||||
entries.remove(jobId);
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean exists(String jobId) {
|
||||
return get(jobId).isPresent();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Optional<String> findJobIdByFileId(String fileId) {
|
||||
for (Map.Entry<String, Holder> e : entries.entrySet()) {
|
||||
Holder h = e.getValue();
|
||||
if (h.isExpired()) {
|
||||
continue;
|
||||
}
|
||||
List<String> fileIds = h.entry.fileIds();
|
||||
if (fileIds != null && fileIds.contains(fileId)) {
|
||||
return Optional.of(e.getKey());
|
||||
}
|
||||
}
|
||||
return Optional.empty();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Collection<JobStoreEntry> all() {
|
||||
List<JobStoreEntry> result = new ArrayList<>(entries.size());
|
||||
for (Holder h : entries.values()) {
|
||||
if (!h.isExpired()) {
|
||||
result.add(h.entry);
|
||||
}
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
/** Drop entries whose TTL has elapsed. Called by the {@code TaskManager} cleanup scheduler. */
|
||||
public int purgeExpired() {
|
||||
int removed = 0;
|
||||
Instant now = Instant.now();
|
||||
for (Map.Entry<String, Holder> e : entries.entrySet()) {
|
||||
if (!e.getValue().expiry.equals(Instant.MAX) && e.getValue().expiry.isBefore(now)) {
|
||||
if (entries.remove(e.getKey(), e.getValue())) {
|
||||
removed++;
|
||||
}
|
||||
}
|
||||
}
|
||||
return removed;
|
||||
}
|
||||
|
||||
private record Holder(JobStoreEntry entry, Instant expiry) {
|
||||
boolean isExpired() {
|
||||
return !expiry.equals(Instant.MAX) && expiry.isBefore(Instant.now());
|
||||
}
|
||||
}
|
||||
}
|
||||
+55
@@ -0,0 +1,55 @@
|
||||
package stirling.software.common.cluster.inprocess;
|
||||
|
||||
import java.time.Duration;
|
||||
import java.time.Instant;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
|
||||
import stirling.software.common.cluster.KeyValueCache;
|
||||
|
||||
public class InProcessKeyValueCache implements KeyValueCache {
|
||||
|
||||
private final ConcurrentHashMap<String, ConcurrentHashMap<String, Expiring>> namespaces =
|
||||
new ConcurrentHashMap<>();
|
||||
|
||||
@Override
|
||||
public void put(String namespace, String key, String value, Duration ttl) {
|
||||
Instant expiry = ttl == null ? Instant.MAX : Instant.now().plus(ttl);
|
||||
namespaces
|
||||
.computeIfAbsent(namespace, n -> new ConcurrentHashMap<>())
|
||||
.put(key, new Expiring(value, expiry));
|
||||
}
|
||||
|
||||
@Override
|
||||
public Optional<String> get(String namespace, String key) {
|
||||
Map<String, Expiring> ns = namespaces.get(namespace);
|
||||
if (ns == null) {
|
||||
return Optional.empty();
|
||||
}
|
||||
Expiring e = ns.get(key);
|
||||
if (e == null) {
|
||||
return Optional.empty();
|
||||
}
|
||||
if (e.expiry.isBefore(Instant.now())) {
|
||||
ns.remove(key, e);
|
||||
return Optional.empty();
|
||||
}
|
||||
return Optional.of(e.value);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void evict(String namespace, String key) {
|
||||
Map<String, Expiring> ns = namespaces.get(namespace);
|
||||
if (ns != null) {
|
||||
ns.remove(key);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void evictNamespace(String namespace) {
|
||||
namespaces.remove(namespace);
|
||||
}
|
||||
|
||||
private record Expiring(String value, Instant expiry) {}
|
||||
}
|
||||
+49
@@ -0,0 +1,49 @@
|
||||
package stirling.software.common.cluster.inprocess;
|
||||
|
||||
import java.time.Duration;
|
||||
import java.util.Collections;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.Map;
|
||||
|
||||
import io.github.bucket4j.Bandwidth;
|
||||
import io.github.bucket4j.Bucket;
|
||||
import io.github.bucket4j.ConsumptionProbe;
|
||||
import io.github.bucket4j.local.LocalBucketBuilder;
|
||||
|
||||
import stirling.software.common.cluster.RateLimitStore;
|
||||
|
||||
/** Bucket4j-backed token bucket implementation of {@link RateLimitStore}. */
|
||||
public class InProcessRateLimitStore implements RateLimitStore {
|
||||
|
||||
/** Cap to bound memory; oldest accessed buckets are evicted. */
|
||||
private static final int MAX_BUCKETS = 10_000;
|
||||
|
||||
private final Map<String, Bucket> buckets =
|
||||
Collections.synchronizedMap(
|
||||
new LinkedHashMap<String, Bucket>(256, 0.75f, true) {
|
||||
@Override
|
||||
protected boolean removeEldestEntry(Map.Entry<String, Bucket> eldest) {
|
||||
return size() > MAX_BUCKETS;
|
||||
}
|
||||
});
|
||||
|
||||
@Override
|
||||
public RateLimitDecision tryConsume(String bucketKey, long capacity, Duration refillPeriod) {
|
||||
String compositeKey = bucketKey + "|" + capacity + "|" + refillPeriod.toNanos();
|
||||
Bucket bucket =
|
||||
buckets.computeIfAbsent(compositeKey, k -> buildBucket(capacity, refillPeriod));
|
||||
ConsumptionProbe probe = bucket.tryConsumeAndReturnRemaining(1);
|
||||
return new RateLimitDecision(
|
||||
probe.isConsumed(),
|
||||
probe.getRemainingTokens(),
|
||||
probe.isConsumed() ? 0L : probe.getNanosToWaitForRefill());
|
||||
}
|
||||
|
||||
private static Bucket buildBucket(long capacity, Duration refillPeriod) {
|
||||
Bandwidth limit =
|
||||
Bandwidth.builder().capacity(capacity).refillGreedy(capacity, refillPeriod).build();
|
||||
LocalBucketBuilder builder = Bucket.builder();
|
||||
builder.addLimit(limit);
|
||||
return builder.build();
|
||||
}
|
||||
}
|
||||
+128
@@ -0,0 +1,128 @@
|
||||
package stirling.software.common.cluster.inprocess;
|
||||
|
||||
import java.io.BufferedInputStream;
|
||||
import java.io.IOException;
|
||||
import java.io.InputStream;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.UUID;
|
||||
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
|
||||
import stirling.software.common.cluster.FileStore;
|
||||
|
||||
/** Local-disk {@link FileStore} storing files under a base directory keyed by a UUID file id. */
|
||||
@Slf4j
|
||||
public class LocalDiskFileStore implements FileStore {
|
||||
|
||||
private final String baseDirPath;
|
||||
|
||||
public LocalDiskFileStore(String baseDirPath) {
|
||||
this.baseDirPath = baseDirPath;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Stored store(InputStream in, String originalName) throws IOException {
|
||||
String fileId = UUID.randomUUID().toString();
|
||||
Path filePath = resolve(fileId);
|
||||
Files.createDirectories(filePath.getParent());
|
||||
boolean success = false;
|
||||
try {
|
||||
long size = Files.copy(in, filePath);
|
||||
success = true;
|
||||
return new Stored(fileId, size);
|
||||
} finally {
|
||||
if (!success) {
|
||||
try {
|
||||
Files.deleteIfExists(filePath);
|
||||
} catch (IOException cleanupEx) {
|
||||
log.warn(
|
||||
"Failed to clean up partial file {} after store failure",
|
||||
filePath,
|
||||
cleanupEx);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* File-to-file copy. {@link Files#copy(Path, Path, java.nio.file.CopyOption...)} can use {@code
|
||||
* sendfile(2)} on Linux for a zero-copy kernel transfer when source and destination share a
|
||||
* filesystem, avoiding the streaming overhead of pulling the bytes through the JVM heap. Reads
|
||||
* the source size before copying so the post-copy stat is unnecessary.
|
||||
*/
|
||||
@Override
|
||||
public Stored store(Path source, String originalName) throws IOException {
|
||||
String fileId = UUID.randomUUID().toString();
|
||||
Path filePath = resolve(fileId);
|
||||
Files.createDirectories(filePath.getParent());
|
||||
long size = Files.size(source);
|
||||
boolean success = false;
|
||||
try {
|
||||
Files.copy(source, filePath);
|
||||
success = true;
|
||||
return new Stored(fileId, size);
|
||||
} finally {
|
||||
if (!success) {
|
||||
try {
|
||||
Files.deleteIfExists(filePath);
|
||||
} catch (IOException cleanupEx) {
|
||||
log.warn(
|
||||
"Failed to clean up partial file {} after store failure",
|
||||
filePath,
|
||||
cleanupEx);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public InputStream retrieve(String fileId) throws IOException {
|
||||
return new BufferedInputStream(Files.newInputStream(resolve(fileId)));
|
||||
}
|
||||
|
||||
@Override
|
||||
public byte[] retrieveBytes(String fileId) throws IOException {
|
||||
Path filePath = resolve(fileId);
|
||||
if (!Files.exists(filePath)) {
|
||||
throw new IOException("File not found with ID: " + fileId);
|
||||
}
|
||||
return Files.readAllBytes(filePath);
|
||||
}
|
||||
|
||||
@Override
|
||||
public long size(String fileId) throws IOException {
|
||||
Path filePath = resolve(fileId);
|
||||
if (!Files.exists(filePath)) {
|
||||
throw new IOException("File not found with ID: " + fileId);
|
||||
}
|
||||
return Files.size(filePath);
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean delete(String fileId) {
|
||||
try {
|
||||
return Files.deleteIfExists(resolve(fileId));
|
||||
} catch (IOException e) {
|
||||
log.error("Error deleting file with ID: {}", fileId, e);
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean exists(String fileId) {
|
||||
return Files.exists(resolve(fileId));
|
||||
}
|
||||
|
||||
public Path resolve(String fileId) {
|
||||
if (fileId.contains("..") || fileId.contains("/") || fileId.contains("\\")) {
|
||||
throw new IllegalArgumentException("Invalid file ID");
|
||||
}
|
||||
Path basePath = Path.of(baseDirPath).normalize().toAbsolutePath();
|
||||
Path resolvedPath = basePath.resolve(fileId).normalize();
|
||||
if (!resolvedPath.startsWith(basePath)) {
|
||||
throw new IllegalArgumentException("File ID resolves to an invalid path");
|
||||
}
|
||||
return resolvedPath;
|
||||
}
|
||||
}
|
||||
+29
@@ -0,0 +1,29 @@
|
||||
package stirling.software.common.cluster.inprocess;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Value;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
|
||||
import stirling.software.common.cluster.FileStore;
|
||||
|
||||
/**
|
||||
* Always-on wiring for the per-node local-disk {@link FileStore}. Active when {@code
|
||||
* cluster.artifactStore=local} (the default; {@code matchIfMissing=true}). The S3 artifact-store
|
||||
* supplies its own bean when {@code cluster.artifactStore=s3}.
|
||||
*/
|
||||
@Configuration
|
||||
@ConditionalOnProperty(
|
||||
prefix = "cluster",
|
||||
name = "artifactStore",
|
||||
havingValue = "local",
|
||||
matchIfMissing = true)
|
||||
public class LocalDiskFileStoreConfiguration {
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean
|
||||
public FileStore fileStore(@Value("${stirling.tempDir:/tmp/stirling-files}") String tempDir) {
|
||||
return new LocalDiskFileStore(tempDir);
|
||||
}
|
||||
}
|
||||
+27
-5
@@ -80,6 +80,7 @@ public class ConfigInitializer {
|
||||
YamlHelper settingsFile = new YamlHelper(settingTempPath);
|
||||
|
||||
migrateEnterpriseEditionToPremium(settingsFile, settingsTemplateFile);
|
||||
migrateProFeaturesKeyCasing(settingsFile, settingsTemplateFile);
|
||||
|
||||
boolean changesMade =
|
||||
settingsTemplateFile.updateValuesFromYaml(settingsFile, settingsTemplateFile);
|
||||
@@ -116,31 +117,52 @@ public class ConfigInitializer {
|
||||
}
|
||||
if (yaml.getValueByExactKeyPath("enterpriseEdition", "SSOAutoLogin") != null) {
|
||||
template.updateValue(
|
||||
List.of("premium", "proFeatures", "SSOAutoLogin"),
|
||||
List.of("premium", "proFeatures", "ssoAutoLogin"),
|
||||
yaml.getValueByExactKeyPath("enterpriseEdition", "SSOAutoLogin"));
|
||||
}
|
||||
if (yaml.getValueByExactKeyPath("enterpriseEdition", "CustomMetadata", "autoUpdateMetadata")
|
||||
!= null) {
|
||||
template.updateValue(
|
||||
List.of("premium", "proFeatures", "CustomMetadata", "autoUpdateMetadata"),
|
||||
List.of("premium", "proFeatures", "customMetadata", "autoUpdateMetadata"),
|
||||
yaml.getValueByExactKeyPath(
|
||||
"enterpriseEdition", "CustomMetadata", "autoUpdateMetadata"));
|
||||
}
|
||||
if (yaml.getValueByExactKeyPath("enterpriseEdition", "CustomMetadata", "author") != null) {
|
||||
template.updateValue(
|
||||
List.of("premium", "proFeatures", "CustomMetadata", "author"),
|
||||
List.of("premium", "proFeatures", "customMetadata", "author"),
|
||||
yaml.getValueByExactKeyPath("enterpriseEdition", "CustomMetadata", "author"));
|
||||
}
|
||||
if (yaml.getValueByExactKeyPath("enterpriseEdition", "CustomMetadata", "creator") != null) {
|
||||
template.updateValue(
|
||||
List.of("premium", "proFeatures", "CustomMetadata", "creator"),
|
||||
List.of("premium", "proFeatures", "customMetadata", "creator"),
|
||||
yaml.getValueByExactKeyPath("enterpriseEdition", "CustomMetadata", "creator"));
|
||||
}
|
||||
if (yaml.getValueByExactKeyPath("enterpriseEdition", "CustomMetadata", "producer")
|
||||
!= null) {
|
||||
template.updateValue(
|
||||
List.of("premium", "proFeatures", "CustomMetadata", "producer"),
|
||||
List.of("premium", "proFeatures", "customMetadata", "producer"),
|
||||
yaml.getValueByExactKeyPath("enterpriseEdition", "CustomMetadata", "producer"));
|
||||
}
|
||||
}
|
||||
|
||||
// TODO: Remove post migration
|
||||
// settings.yml.template renamed the two non-camelCase proFeatures keys
|
||||
// ("SSOAutoLogin" -> "ssoAutoLogin", "CustomMetadata" -> "customMetadata") so the whole
|
||||
// settings pipeline is consistent camelCase. The save path (YamlHelper.updateValue) matches
|
||||
// keys case-sensitively, so without this carry-forward an existing install's values written
|
||||
// under the old PascalCase keys would be dropped on upgrade and reset to template defaults.
|
||||
void migrateProFeaturesKeyCasing(YamlHelper yaml, YamlHelper template) {
|
||||
Object ssoAutoLogin = yaml.getValueByExactKeyPath("premium", "proFeatures", "SSOAutoLogin");
|
||||
if (ssoAutoLogin != null) {
|
||||
template.updateValue(List.of("premium", "proFeatures", "ssoAutoLogin"), ssoAutoLogin);
|
||||
}
|
||||
for (String field : List.of("autoUpdateMetadata", "author", "creator", "producer")) {
|
||||
Object value =
|
||||
yaml.getValueByExactKeyPath("premium", "proFeatures", "CustomMetadata", field);
|
||||
if (value != null) {
|
||||
template.updateValue(
|
||||
List.of("premium", "proFeatures", "customMetadata", field), value);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -13,6 +13,7 @@ import java.util.Arrays;
|
||||
import java.util.Collection;
|
||||
import java.util.List;
|
||||
import java.util.Locale;
|
||||
import java.util.UUID;
|
||||
|
||||
import org.springframework.boot.context.properties.ConfigurationProperties;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
@@ -77,6 +78,7 @@ public class ApplicationProperties {
|
||||
private PdfEditor pdfEditor = new PdfEditor();
|
||||
private AiEngine aiEngine = new AiEngine();
|
||||
private InternalApi internalApi = new InternalApi();
|
||||
private Cluster cluster = new Cluster();
|
||||
|
||||
@Bean
|
||||
public PropertySource<?> dynamicYamlPropertySource(ConfigurableEnvironment environment)
|
||||
@@ -254,6 +256,106 @@ public class ApplicationProperties {
|
||||
private int longRunningTimeoutSeconds = 600;
|
||||
}
|
||||
|
||||
/**
|
||||
* Cluster backplane configuration. All keys live under the top-level {@code cluster.*} prefix
|
||||
* (e.g. env var {@code CLUSTER_ENABLED}). The master switch is {@link #enabled} and defaults to
|
||||
* off; when off the in-process backplane is wired and no other cluster keys are required.
|
||||
*/
|
||||
@Data
|
||||
public static class Cluster {
|
||||
|
||||
/** Master switch. When {@code false} (default) the in-process backplane is wired. */
|
||||
private boolean enabled = false;
|
||||
|
||||
/** Backplane implementation selector. Valid values: {@code inprocess} | {@code valkey}. */
|
||||
private String backplane = "inprocess";
|
||||
|
||||
/**
|
||||
* Transient cluster job-artifact store selector. Valid values: {@code local} | {@code s3}.
|
||||
*
|
||||
* <p>This is distinct from {@code storage.provider}, which selects the backend for
|
||||
* persistent user-uploaded files. The two switches exist because the user-facing storage
|
||||
* feature is optional ({@code storage.enabled=false} is common) but every multi-node
|
||||
* cluster still needs a shared artifact store to serve cross-node downloads. Both
|
||||
* implementations share credentials from {@code storage.s3.*} when set to {@code s3}.
|
||||
*/
|
||||
private String artifactStore = "local";
|
||||
|
||||
private Valkey valkey = new Valkey();
|
||||
private Node node = new Node();
|
||||
|
||||
private transient String cachedNodeId;
|
||||
|
||||
public NodeRole resolvedRole() {
|
||||
if (node == null || node.getRole() == null) {
|
||||
return NodeRole.BOTH;
|
||||
}
|
||||
String value = node.getRole().trim().toUpperCase(Locale.ROOT);
|
||||
try {
|
||||
return NodeRole.valueOf(value);
|
||||
} catch (IllegalArgumentException ex) {
|
||||
return NodeRole.BOTH;
|
||||
}
|
||||
}
|
||||
|
||||
public synchronized String resolvedNodeId() {
|
||||
if (node != null && node.getId() != null && !node.getId().isBlank()) {
|
||||
return node.getId();
|
||||
}
|
||||
if (cachedNodeId == null) {
|
||||
cachedNodeId = UUID.randomUUID().toString();
|
||||
}
|
||||
return cachedNodeId;
|
||||
}
|
||||
|
||||
public enum NodeRole {
|
||||
WEB,
|
||||
WORKER,
|
||||
BOTH
|
||||
}
|
||||
|
||||
@Data
|
||||
public static class Valkey {
|
||||
/**
|
||||
* {@code redis://host:6379} or {@code rediss://...} for TLS. Required when cluster mode
|
||||
* is on and backplane is valkey.
|
||||
*/
|
||||
private String url = "";
|
||||
|
||||
private Tls tls = new Tls();
|
||||
|
||||
@Data
|
||||
public static class Tls {
|
||||
/**
|
||||
* When {@code true}, skip Valkey/Redis TLS certificate verification (dev/test
|
||||
* only). Leave {@code false} in production.
|
||||
*/
|
||||
private boolean skipCertVerification = false;
|
||||
}
|
||||
}
|
||||
|
||||
@Data
|
||||
public static class Node {
|
||||
/** Optional explicit node id. Blank = auto-generated UUID at startup. */
|
||||
private String id = "";
|
||||
|
||||
/** {@code web} | {@code worker} | {@code both}. */
|
||||
private String role = "both";
|
||||
|
||||
/**
|
||||
* Internal cluster address advertised in the instance registry (host:port). Blank =
|
||||
* derived at startup.
|
||||
*/
|
||||
private String internalAddress = "";
|
||||
|
||||
/** {@code http} | {@code https} - scheme used when peers call this node. */
|
||||
private String scheme = "http";
|
||||
|
||||
/** Heartbeat publish interval for the instance registry, in milliseconds. */
|
||||
private long heartbeatIntervalMs = 5000;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* HTTP timeouts for loopback calls to internal Stirling API endpoints, used by the AI workflow
|
||||
* executor and the pipeline processor. A bounded read timeout prevents a hung tool (e.g. an
|
||||
@@ -426,6 +528,16 @@ public class ApplicationProperties {
|
||||
private String provider;
|
||||
private Client client = new Client();
|
||||
|
||||
/**
|
||||
* When true, the OAuth2/OIDC login flow logs the full set of ID token and UserInfo
|
||||
* claims at INFO level (and again at ERROR level if the username attribute cannot be
|
||||
* resolved). Used to diagnose provider misconfiguration (for example ADFS not returning
|
||||
* an {@code email} claim). WARNING: writes PII (sub, email, name) to application logs.
|
||||
* Leave disabled in production; enable only while actively troubleshooting and disable
|
||||
* again afterwards.
|
||||
*/
|
||||
private Boolean debugLogging = false;
|
||||
|
||||
public void setScopes(String scopes) {
|
||||
List<String> scopesList =
|
||||
Arrays.stream(scopes.split(",")).map(String::trim).toList();
|
||||
@@ -676,6 +788,7 @@ public class ApplicationProperties {
|
||||
private boolean enabled = false;
|
||||
private String provider = "local";
|
||||
private Local local = new Local();
|
||||
private S3 s3 = new S3();
|
||||
private Quotas quotas = new Quotas();
|
||||
private Sharing sharing = new Sharing();
|
||||
private Signing signing = new Signing();
|
||||
@@ -685,6 +798,57 @@ public class ApplicationProperties {
|
||||
private String basePath = InstallationPathConfig.getPath() + "storage";
|
||||
}
|
||||
|
||||
@Data
|
||||
public static class S3 {
|
||||
/**
|
||||
* Optional custom endpoint (e.g. {@code https://<account>.r2.cloudflarestorage.com},
|
||||
* {@code https://<project>.supabase.co/storage/v1/s3}, or {@code http://localhost:9000}
|
||||
* for MinIO). Blank = use AWS regional default.
|
||||
*/
|
||||
private String endpoint = "";
|
||||
|
||||
private String bucket = "";
|
||||
|
||||
private String region = "us-east-1";
|
||||
|
||||
private String accessKey = "";
|
||||
private String secretKey = "";
|
||||
|
||||
/**
|
||||
* When {@code true} use path-style URLs ({@code <endpoint>/<bucket>/<key>}) instead of
|
||||
* virtual-hosted ({@code <bucket>.<endpoint>/<key>}). MinIO and most S3-compatible
|
||||
* gateways require path-style; AWS S3 prefers virtual-hosted.
|
||||
*/
|
||||
private boolean pathStyleAccess = false;
|
||||
|
||||
/**
|
||||
* When {@code false} (default), {@code endpoint} hostnames that resolve to private,
|
||||
* loopback, or link-local addresses are rejected at startup to block SSRF attacks via
|
||||
* the cloud metadata service (e.g. {@code http://169.254.169.254/}). Set to {@code
|
||||
* true} to opt in for MinIO / in-cluster S3 endpoints on private networks.
|
||||
*/
|
||||
private boolean allowPrivateEndpoints = false;
|
||||
|
||||
/**
|
||||
* Controls when the SDK adds an {@code x-amz-checksum-*} header on PUT/UploadPart.
|
||||
* Default {@code WHEN_SUPPORTED} (the SDK default since 2.30) makes the SDK send a
|
||||
* CRC32 checksum on every upload - this works on AWS S3, MinIO, current Supabase,
|
||||
* Backblaze B2 (post-July-2025), and modern R2. Set to {@code WHEN_REQUIRED} to
|
||||
* suppress the auto-checksum on vendors that reject unknown {@code x-amz-checksum-*}
|
||||
* headers (older Backblaze B2, some R2 corner cases, GCS S3 endpoint). Invalid values
|
||||
* fall back to {@code WHEN_SUPPORTED}.
|
||||
*/
|
||||
private String requestChecksumCalculation = "WHEN_SUPPORTED";
|
||||
|
||||
/**
|
||||
* Controls when the SDK validates returned {@code x-amz-checksum-*} headers on GET
|
||||
* responses. Default {@code WHEN_SUPPORTED}. Set to {@code WHEN_REQUIRED} if your
|
||||
* vendor never returns these headers and you see false-positive checksum-mismatch
|
||||
* errors. Invalid values fall back to {@code WHEN_SUPPORTED}.
|
||||
*/
|
||||
private String responseChecksumValidation = "WHEN_SUPPORTED";
|
||||
}
|
||||
|
||||
@Data
|
||||
public static class Sharing {
|
||||
private boolean enabled = false;
|
||||
|
||||
@@ -1,15 +1,13 @@
|
||||
package stirling.software.common.service;
|
||||
|
||||
import java.io.BufferedInputStream;
|
||||
import java.io.BufferedOutputStream;
|
||||
import java.io.ByteArrayInputStream;
|
||||
import java.io.IOException;
|
||||
import java.io.InputStream;
|
||||
import java.io.OutputStream;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.UUID;
|
||||
import java.io.PipedInputStream;
|
||||
import java.io.PipedOutputStream;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Value;
|
||||
import org.springframework.core.io.Resource;
|
||||
import org.springframework.stereotype.Service;
|
||||
import org.springframework.web.multipart.MultipartFile;
|
||||
@@ -18,9 +16,11 @@ import org.springframework.web.servlet.mvc.method.annotation.StreamingResponseBo
|
||||
import lombok.RequiredArgsConstructor;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
|
||||
import stirling.software.common.cluster.FileStore;
|
||||
|
||||
/**
|
||||
* Service for storing and retrieving files with unique file IDs. Used by the AutoJobPostMapping
|
||||
* system to handle file references.
|
||||
* system to handle file references. Disk I/O is delegated to the injected {@link FileStore} bean.
|
||||
*/
|
||||
@Service
|
||||
@RequiredArgsConstructor
|
||||
@@ -30,251 +30,150 @@ public class FileStorage {
|
||||
/** Holds the result of a stream-to-disk store operation: the file ID and the bytes written. */
|
||||
public record StoredFile(String fileId, long size) {}
|
||||
|
||||
@Value("${stirling.tempDir:/tmp/stirling-files}")
|
||||
private String tempDirPath;
|
||||
|
||||
private final FileOrUploadService fileOrUploadService;
|
||||
private final FileStore fileStore;
|
||||
|
||||
/**
|
||||
* Store a file and return its unique ID
|
||||
*
|
||||
* @param file The file to store
|
||||
* @return The unique ID assigned to the file
|
||||
* @throws IOException If there is an error storing the file
|
||||
*/
|
||||
public String storeFile(MultipartFile file) throws IOException {
|
||||
String fileId = generateFileId();
|
||||
Path filePath = getFilePath(fileId);
|
||||
|
||||
// Ensure the directory exists
|
||||
Files.createDirectories(filePath.getParent());
|
||||
|
||||
// Transfer the file to the storage location
|
||||
file.transferTo(filePath.toFile());
|
||||
|
||||
log.debug("Stored file with ID: {}", fileId);
|
||||
return fileId;
|
||||
}
|
||||
|
||||
/**
|
||||
* Store a byte array as a file and return its unique ID
|
||||
*
|
||||
* @param bytes The byte array to store
|
||||
* @param originalName The original name of the file (for extension)
|
||||
* @return The unique ID assigned to the file
|
||||
* @throws IOException If there is an error storing the file
|
||||
*/
|
||||
public String storeBytes(byte[] bytes, String originalName) throws IOException {
|
||||
String fileId = generateFileId();
|
||||
Path filePath = getFilePath(fileId);
|
||||
|
||||
// Ensure the directory exists
|
||||
Files.createDirectories(filePath.getParent());
|
||||
|
||||
// Write the bytes to the file
|
||||
Files.write(filePath, bytes);
|
||||
|
||||
log.debug("Stored byte array with ID: {}", fileId);
|
||||
return fileId;
|
||||
}
|
||||
|
||||
/**
|
||||
* Retrieve a file by its ID as a MultipartFile
|
||||
*
|
||||
* @param fileId The ID of the file to retrieve
|
||||
* @return The file as a MultipartFile
|
||||
* @throws IOException If the file doesn't exist or can't be read
|
||||
*/
|
||||
public MultipartFile retrieveFile(String fileId) throws IOException {
|
||||
Path filePath = getFilePath(fileId);
|
||||
|
||||
if (!Files.exists(filePath)) {
|
||||
throw new IOException("File not found with ID: " + fileId);
|
||||
// Fast path: when Spring buffered the multipart to disk (typical for large uploads), the
|
||||
// backing Resource exposes a real File. Hand the Path to the FileStore so it can do a
|
||||
// file-to-file copy (Linux sendfile, no copy through Java heap) rather than streaming
|
||||
// the bytes through an 8K buffer. Falls back to the InputStream path for in-memory
|
||||
// multiparts, exotic Resource impls, and anything that does not back onto a File.
|
||||
Resource res;
|
||||
try {
|
||||
res = file.getResource();
|
||||
} catch (RuntimeException ignored) {
|
||||
res = null;
|
||||
}
|
||||
if (res != null && res.isFile()) {
|
||||
try {
|
||||
FileStore.Stored stored =
|
||||
fileStore.store(res.getFile().toPath(), file.getOriginalFilename());
|
||||
log.debug("Stored file with ID: {} (fast path)", stored.fileId());
|
||||
return stored.fileId();
|
||||
} catch (IOException ex) {
|
||||
// Some Resource impls advertise isFile()=true but throw on getFile(); fall through.
|
||||
log.debug("Resource fast path failed, falling back to stream copy", ex);
|
||||
}
|
||||
}
|
||||
try (InputStream in = file.getInputStream()) {
|
||||
FileStore.Stored stored = fileStore.store(in, file.getOriginalFilename());
|
||||
log.debug("Stored file with ID: {}", stored.fileId());
|
||||
return stored.fileId();
|
||||
}
|
||||
}
|
||||
|
||||
byte[] fileData = Files.readAllBytes(filePath);
|
||||
public String storeBytes(byte[] bytes, String originalName) throws IOException {
|
||||
FileStore.Stored stored = fileStore.store(new ByteArrayInputStream(bytes), originalName);
|
||||
log.debug("Stored byte array with ID: {}", stored.fileId());
|
||||
return stored.fileId();
|
||||
}
|
||||
|
||||
public MultipartFile retrieveFile(String fileId) throws IOException {
|
||||
byte[] fileData = fileStore.retrieveBytes(fileId);
|
||||
return fileOrUploadService.toMockMultipartFile(fileId, fileData);
|
||||
}
|
||||
|
||||
/**
|
||||
* Retrieve a file by its ID as a byte array
|
||||
*
|
||||
* @param fileId The ID of the file to retrieve
|
||||
* @return The file as a byte array
|
||||
* @throws IOException If the file doesn't exist or can't be read
|
||||
*/
|
||||
public byte[] retrieveBytes(String fileId) throws IOException {
|
||||
Path filePath = getFilePath(fileId);
|
||||
|
||||
if (!Files.exists(filePath)) {
|
||||
throw new IOException("File not found with ID: " + fileId);
|
||||
}
|
||||
|
||||
return Files.readAllBytes(filePath);
|
||||
return fileStore.retrieveBytes(fileId);
|
||||
}
|
||||
|
||||
/**
|
||||
* Retrieve a file by its ID as a streaming InputStream. The caller is responsible for closing
|
||||
* the returned stream.
|
||||
*
|
||||
* @param fileId The ID of the file to retrieve
|
||||
* @return A buffered InputStream for the file
|
||||
* @throws IOException If the file doesn't exist or can't be read
|
||||
*/
|
||||
public InputStream retrieveInputStream(String fileId) throws IOException {
|
||||
Path filePath = getFilePath(fileId);
|
||||
// Let Files.newInputStream throw NoSuchFileException naturally — avoids TOCTOU race
|
||||
// between exists-check and open when another thread may delete concurrently.
|
||||
return new BufferedInputStream(Files.newInputStream(filePath));
|
||||
return fileStore.retrieve(fileId);
|
||||
}
|
||||
|
||||
/**
|
||||
* Store data from an InputStream as a file and return its unique ID and byte count. Streams
|
||||
* directly to disk without buffering the entire content in heap.
|
||||
*
|
||||
* @param inputStream The input stream to read from
|
||||
* @param originalName The original name of the file (unused, kept for API symmetry)
|
||||
* @return A {@link StoredFile} containing the file ID and the number of bytes written
|
||||
* @throws IOException If there is an error storing the file
|
||||
*/
|
||||
public StoredFile storeInputStream(InputStream inputStream, String originalName)
|
||||
throws IOException {
|
||||
String fileId = generateFileId();
|
||||
Path filePath = getFilePath(fileId);
|
||||
Files.createDirectories(filePath.getParent());
|
||||
long size = Files.copy(inputStream, filePath);
|
||||
log.debug("Stored input stream with ID: {}", fileId);
|
||||
return new StoredFile(fileId, size);
|
||||
FileStore.Stored stored = fileStore.store(inputStream, originalName);
|
||||
log.debug("Stored input stream with ID: {}", stored.fileId());
|
||||
return new StoredFile(stored.fileId(), stored.size());
|
||||
}
|
||||
|
||||
public String storeFromStreamingBody(StreamingResponseBody body, String originalName)
|
||||
throws IOException {
|
||||
String fileId = generateFileId();
|
||||
Path filePath = getFilePath(fileId);
|
||||
Files.createDirectories(filePath.getParent());
|
||||
boolean success = false;
|
||||
try (OutputStream os = new BufferedOutputStream(Files.newOutputStream(filePath))) {
|
||||
body.writeTo(os);
|
||||
success = true;
|
||||
} finally {
|
||||
if (!success) {
|
||||
// Hold Throwable not IOException: an unchecked failure (NPE, IllegalState, OOM, etc.)
|
||||
// from the body writer would otherwise close the pipe with EOF and the consumer would
|
||||
// return a truncated file with no error surfaced to the caller.
|
||||
AtomicReference<Throwable> bodyError = new AtomicReference<>();
|
||||
try (PipedOutputStream out = new PipedOutputStream();
|
||||
PipedInputStream in = new PipedInputStream(out, 8192)) {
|
||||
var executor = Executors.newSingleThreadExecutor(Thread.ofVirtual().factory());
|
||||
java.util.concurrent.Future<?> task = null;
|
||||
try {
|
||||
task =
|
||||
executor.submit(
|
||||
() -> {
|
||||
try {
|
||||
body.writeTo(out);
|
||||
} catch (Throwable ex) {
|
||||
bodyError.set(ex);
|
||||
} finally {
|
||||
try {
|
||||
out.close();
|
||||
} catch (IOException ignored) {
|
||||
// closed on the consumer side too
|
||||
}
|
||||
}
|
||||
});
|
||||
FileStore.Stored stored = fileStore.store(in, originalName);
|
||||
Throwable writerErr = bodyError.get();
|
||||
if (writerErr != null) {
|
||||
// Body failed mid-write: the FileStore persisted a truncated entry.
|
||||
// Best-effort delete so we don't leak partial files; never let cleanup
|
||||
// mask the original writer error.
|
||||
try {
|
||||
fileStore.delete(stored.fileId());
|
||||
} catch (RuntimeException cleanupEx) {
|
||||
log.warn(
|
||||
"Failed to delete partial file {} after writer error: {}",
|
||||
stored.fileId(),
|
||||
cleanupEx.getMessage());
|
||||
}
|
||||
if (writerErr instanceof IOException ioe) {
|
||||
throw ioe;
|
||||
}
|
||||
throw new IOException(
|
||||
"StreamingResponseBody writer failed: " + writerErr.getMessage(),
|
||||
writerErr);
|
||||
}
|
||||
log.debug("Stored StreamingResponseBody with ID: {}", stored.fileId());
|
||||
return stored.fileId();
|
||||
} finally {
|
||||
// Interrupt and join the writer task: shutdown() alone returns immediately and a
|
||||
// failed store leaves the writer running, leaking a thread per failed upload.
|
||||
if (task != null && !task.isDone()) {
|
||||
task.cancel(true);
|
||||
}
|
||||
executor.shutdown();
|
||||
try {
|
||||
Files.deleteIfExists(filePath);
|
||||
} catch (IOException cleanupEx) {
|
||||
log.warn(
|
||||
"Failed to clean up partial file {} after store failure",
|
||||
filePath,
|
||||
cleanupEx);
|
||||
if (!executor.awaitTermination(5, java.util.concurrent.TimeUnit.SECONDS)) {
|
||||
executor.shutdownNow();
|
||||
}
|
||||
} catch (InterruptedException ie) {
|
||||
executor.shutdownNow();
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
}
|
||||
}
|
||||
log.debug("Stored StreamingResponseBody with ID: {}", fileId);
|
||||
return fileId;
|
||||
}
|
||||
|
||||
/**
|
||||
* Persist a {@link Resource} body to disk, returning the generated file ID. Used by the async
|
||||
* job pipeline to capture {@code ResponseEntity<Resource>} results produced by controllers.
|
||||
*/
|
||||
public String storeFromResource(Resource resource, String originalName) throws IOException {
|
||||
String fileId = generateFileId();
|
||||
Path filePath = getFilePath(fileId);
|
||||
Files.createDirectories(filePath.getParent());
|
||||
boolean success = false;
|
||||
try (InputStream in = resource.getInputStream()) {
|
||||
Files.copy(in, filePath);
|
||||
success = true;
|
||||
} finally {
|
||||
if (!success) {
|
||||
try {
|
||||
Files.deleteIfExists(filePath);
|
||||
} catch (IOException cleanupEx) {
|
||||
log.warn(
|
||||
"Failed to clean up partial file {} after store failure",
|
||||
filePath,
|
||||
cleanupEx);
|
||||
}
|
||||
}
|
||||
FileStore.Stored stored = fileStore.store(in, originalName);
|
||||
log.debug("Stored Resource with ID: {}", stored.fileId());
|
||||
return stored.fileId();
|
||||
}
|
||||
log.debug("Stored Resource with ID: {}", fileId);
|
||||
return fileId;
|
||||
}
|
||||
|
||||
/**
|
||||
* Delete a file by its ID
|
||||
*
|
||||
* @param fileId The ID of the file to delete
|
||||
* @return true if the file was deleted, false otherwise
|
||||
*/
|
||||
public boolean deleteFile(String fileId) {
|
||||
try {
|
||||
Path filePath = getFilePath(fileId);
|
||||
return Files.deleteIfExists(filePath);
|
||||
} catch (IOException e) {
|
||||
log.error("Error deleting file with ID: {}", fileId, e);
|
||||
return false;
|
||||
}
|
||||
return fileStore.delete(fileId);
|
||||
}
|
||||
|
||||
/**
|
||||
* Check if a file exists by its ID
|
||||
*
|
||||
* @param fileId The ID of the file to check
|
||||
* @return true if the file exists, false otherwise
|
||||
*/
|
||||
public boolean fileExists(String fileId) {
|
||||
Path filePath = getFilePath(fileId);
|
||||
return Files.exists(filePath);
|
||||
return fileStore.exists(fileId);
|
||||
}
|
||||
|
||||
/**
|
||||
* Get the size of a file by its ID without loading the content into memory
|
||||
*
|
||||
* @param fileId The ID of the file
|
||||
* @return The size of the file in bytes
|
||||
* @throws IOException If the file doesn't exist or can't be read
|
||||
*/
|
||||
public long getFileSize(String fileId) throws IOException {
|
||||
Path filePath = getFilePath(fileId);
|
||||
|
||||
if (!Files.exists(filePath)) {
|
||||
throw new IOException("File not found with ID: " + fileId);
|
||||
}
|
||||
|
||||
return Files.size(filePath);
|
||||
}
|
||||
|
||||
/**
|
||||
* Get the path for a file ID
|
||||
*
|
||||
* @param fileId The ID of the file
|
||||
* @return The path to the file
|
||||
* @throws IllegalArgumentException if fileId contains path traversal characters or resolves
|
||||
* outside base directory
|
||||
*/
|
||||
private Path getFilePath(String fileId) {
|
||||
// Validate fileId to prevent path traversal
|
||||
if (fileId.contains("..") || fileId.contains("/") || fileId.contains("\\")) {
|
||||
throw new IllegalArgumentException("Invalid file ID");
|
||||
}
|
||||
|
||||
Path basePath = Path.of(tempDirPath).normalize().toAbsolutePath();
|
||||
Path resolvedPath = basePath.resolve(fileId).normalize();
|
||||
|
||||
// Ensure resolved path is within the base directory
|
||||
if (!resolvedPath.startsWith(basePath)) {
|
||||
throw new IllegalArgumentException("File ID resolves to an invalid path");
|
||||
}
|
||||
|
||||
return resolvedPath;
|
||||
}
|
||||
|
||||
/**
|
||||
* Generate a unique file ID
|
||||
*
|
||||
* @return A unique file ID
|
||||
*/
|
||||
private String generateFileId() {
|
||||
return UUID.randomUUID().toString();
|
||||
return fileStore.size(fileId);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -3,9 +3,13 @@ package stirling.software.common.service;
|
||||
import java.io.BufferedInputStream;
|
||||
import java.io.IOException;
|
||||
import java.io.InputStream;
|
||||
import java.time.Duration;
|
||||
import java.time.Instant;
|
||||
import java.time.LocalDateTime;
|
||||
import java.time.ZoneId;
|
||||
import java.time.temporal.ChronoUnit;
|
||||
import java.util.ArrayList;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Locale;
|
||||
import java.util.Map;
|
||||
@@ -16,6 +20,7 @@ import java.util.concurrent.TimeUnit;
|
||||
import java.util.zip.ZipEntry;
|
||||
import java.util.zip.ZipInputStream;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.beans.factory.annotation.Value;
|
||||
import org.springframework.http.MediaType;
|
||||
import org.springframework.stereotype.Service;
|
||||
@@ -26,6 +31,10 @@ import jakarta.annotation.PreDestroy;
|
||||
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
|
||||
import stirling.software.common.cluster.ClusterBackplane;
|
||||
import stirling.software.common.cluster.JobStore;
|
||||
import stirling.software.common.cluster.JobStoreEntry;
|
||||
import stirling.software.common.cluster.JobStoreEntry.JobState;
|
||||
import stirling.software.common.model.job.JobResult;
|
||||
import stirling.software.common.model.job.JobStats;
|
||||
import stirling.software.common.model.job.ResultFile;
|
||||
@@ -40,20 +49,20 @@ public class TaskManager {
|
||||
private int jobResultExpiryMinutes = 30;
|
||||
|
||||
private final FileStorage fileStorage;
|
||||
private final JobStore jobStore;
|
||||
private final ClusterBackplane clusterBackplane;
|
||||
private final ScheduledExecutorService cleanupExecutor =
|
||||
Executors.newSingleThreadScheduledExecutor(
|
||||
Thread.ofVirtual().name("task-cleanup-", 0).factory());
|
||||
|
||||
/** Initialize the task manager and start the cleanup scheduler */
|
||||
public TaskManager(FileStorage fileStorage) {
|
||||
@Autowired
|
||||
public TaskManager(
|
||||
FileStorage fileStorage, JobStore jobStore, ClusterBackplane clusterBackplane) {
|
||||
this.fileStorage = fileStorage;
|
||||
this.jobStore = jobStore;
|
||||
this.clusterBackplane = clusterBackplane;
|
||||
|
||||
// Schedule periodic cleanup of old job results
|
||||
cleanupExecutor.scheduleAtFixedRate(
|
||||
this::cleanupOldJobs,
|
||||
10, // Initial delay
|
||||
10, // Interval
|
||||
TimeUnit.MINUTES);
|
||||
cleanupExecutor.scheduleAtFixedRate(this::cleanupOldJobs, 10, 10, TimeUnit.MINUTES);
|
||||
|
||||
log.debug(
|
||||
"Task manager initialized with job result expiry of {} minutes",
|
||||
@@ -66,7 +75,9 @@ public class TaskManager {
|
||||
* @param jobId The job ID
|
||||
*/
|
||||
public void createTask(String jobId) {
|
||||
jobResults.put(jobId, JobResult.createNew(jobId));
|
||||
JobResult result = JobResult.createNew(jobId);
|
||||
jobResults.put(jobId, result);
|
||||
writeThrough(jobId, result);
|
||||
log.debug("Created task with job ID: {}", jobId);
|
||||
}
|
||||
|
||||
@@ -79,6 +90,7 @@ public class TaskManager {
|
||||
public void setResult(String jobId, Object result) {
|
||||
JobResult jobResult = getOrCreateJobResult(jobId);
|
||||
jobResult.completeWithResult(result);
|
||||
writeThrough(jobId, jobResult);
|
||||
log.debug("Set result for job ID: {}", jobId);
|
||||
}
|
||||
|
||||
@@ -101,6 +113,7 @@ public class TaskManager {
|
||||
extractZipToIndividualFiles(fileId, originalFileName);
|
||||
if (!extractedFiles.isEmpty()) {
|
||||
jobResult.completeWithFiles(extractedFiles);
|
||||
writeThrough(jobId, jobResult);
|
||||
log.debug(
|
||||
"Set multiple file results for job ID: {} with {} files extracted from"
|
||||
+ " ZIP",
|
||||
@@ -127,6 +140,7 @@ public class TaskManager {
|
||||
"Failed to get file size for job {}: {}. Using size 0.", jobId, e.getMessage());
|
||||
jobResult.completeWithSingleFile(fileId, originalFileName, contentType, 0);
|
||||
}
|
||||
writeThrough(jobId, jobResult);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -138,6 +152,7 @@ public class TaskManager {
|
||||
public void setMultipleFileResults(String jobId, List<ResultFile> resultFiles) {
|
||||
JobResult jobResult = getOrCreateJobResult(jobId);
|
||||
jobResult.completeWithFiles(resultFiles);
|
||||
writeThrough(jobId, jobResult);
|
||||
log.debug(
|
||||
"Set multiple file results for job ID: {} with {} files",
|
||||
jobId,
|
||||
@@ -153,6 +168,7 @@ public class TaskManager {
|
||||
public void setError(String jobId, String error) {
|
||||
JobResult jobResult = getOrCreateJobResult(jobId);
|
||||
jobResult.failWithError(error);
|
||||
writeThrough(jobId, jobResult);
|
||||
log.debug("Set error for job ID: {}: {}", jobId, error);
|
||||
}
|
||||
|
||||
@@ -169,6 +185,7 @@ public class TaskManager {
|
||||
// If no result or error has been set, mark it as complete with an empty result
|
||||
jobResult.completeWithResult("Task completed successfully");
|
||||
}
|
||||
writeThrough(jobId, jobResult);
|
||||
log.debug("Marked job ID: {} as complete", jobId);
|
||||
}
|
||||
|
||||
@@ -205,6 +222,7 @@ public class TaskManager {
|
||||
JobResult jobResult = jobResults.get(jobId);
|
||||
if (jobResult != null) {
|
||||
jobResult.addNote(note);
|
||||
writeThrough(jobId, jobResult);
|
||||
log.debug("Added note to job ID: {}: {}", jobId, note);
|
||||
return true;
|
||||
}
|
||||
@@ -295,8 +313,11 @@ public class TaskManager {
|
||||
return jobResults.computeIfAbsent(jobId, JobResult::createNew);
|
||||
}
|
||||
|
||||
/** Clean up old completed job results */
|
||||
/** Clean up old completed job results. No-op in cluster mode; the backplane TTL owns expiry. */
|
||||
public void cleanupOldJobs() {
|
||||
if (clusterBackplane != null && !clusterBackplane.shouldRunLocalCleanup()) {
|
||||
return;
|
||||
}
|
||||
LocalDateTime expiryThreshold =
|
||||
LocalDateTime.now().minus(jobResultExpiryMinutes, ChronoUnit.MINUTES);
|
||||
int removedCount = 0;
|
||||
@@ -315,6 +336,9 @@ public class TaskManager {
|
||||
|
||||
// Remove the job result
|
||||
jobResults.remove(entry.getKey());
|
||||
if (jobStore != null) {
|
||||
jobStore.delete(entry.getKey());
|
||||
}
|
||||
removedCount++;
|
||||
}
|
||||
}
|
||||
@@ -327,6 +351,53 @@ public class TaskManager {
|
||||
}
|
||||
}
|
||||
|
||||
/** Mirror the in-memory {@code JobResult} into the cluster-visible {@link JobStore}. */
|
||||
private void writeThrough(String jobId, JobResult result) {
|
||||
if (jobStore == null) {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
jobStore.put(toEntry(jobId, result), Duration.ofMinutes(jobResultExpiryMinutes));
|
||||
} catch (RuntimeException ex) {
|
||||
log.warn("JobStore write-through failed for job {}: {}", jobId, ex.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
private JobStoreEntry toEntry(String jobId, JobResult result) {
|
||||
JobState state;
|
||||
if (result.isComplete()) {
|
||||
state = result.getError() != null ? JobState.FAILED : JobState.COMPLETE;
|
||||
} else {
|
||||
state = JobState.PENDING;
|
||||
}
|
||||
Instant createdAt = toInstant(result.getCreatedAt());
|
||||
Instant completedAt = toInstant(result.getCompletedAt());
|
||||
List<String> fileIds = new ArrayList<>();
|
||||
if (result.hasFiles()) {
|
||||
for (ResultFile rf : result.getAllResultFiles()) {
|
||||
fileIds.add(rf.getFileId());
|
||||
}
|
||||
}
|
||||
Map<String, String> meta = new HashMap<>();
|
||||
if (result.getNotes() != null && !result.getNotes().isEmpty()) {
|
||||
meta.put("notesCount", Integer.toString(result.getNotes().size()));
|
||||
}
|
||||
String owningNodeId = clusterBackplane == null ? "local" : clusterBackplane.localNodeId();
|
||||
return new JobStoreEntry(
|
||||
jobId,
|
||||
state,
|
||||
owningNodeId,
|
||||
createdAt,
|
||||
completedAt,
|
||||
result.getError(),
|
||||
fileIds,
|
||||
meta);
|
||||
}
|
||||
|
||||
private Instant toInstant(LocalDateTime ldt) {
|
||||
return ldt == null ? null : ldt.atZone(ZoneId.systemDefault()).toInstant();
|
||||
}
|
||||
|
||||
/** Shutdown the cleanup executor */
|
||||
@PreDestroy
|
||||
public void shutdown() {
|
||||
@@ -370,7 +441,7 @@ public class TaskManager {
|
||||
while ((entry = zipIn.getNextEntry()) != null) {
|
||||
if (!entry.isDirectory()) {
|
||||
String contentType = determineContentType(entry.getName());
|
||||
// storeInputStream returns the fileId and byte count — no extra stat needed
|
||||
// storeInputStream returns the fileId and byte count - no extra stat needed
|
||||
FileStorage.StoredFile stored =
|
||||
fileStorage.storeInputStream(zipIn, entry.getName());
|
||||
|
||||
@@ -458,7 +529,8 @@ public class TaskManager {
|
||||
}
|
||||
|
||||
/**
|
||||
* Find the job key that owns a given file ID.
|
||||
* Find the job key that owns a given file ID. Checks the local in-memory map first, then falls
|
||||
* back to the cluster-visible {@link JobStore}.
|
||||
*
|
||||
* @param fileId file identifier to look up
|
||||
* @return scoped job key if found, otherwise null
|
||||
@@ -474,6 +546,18 @@ public class TaskManager {
|
||||
}
|
||||
}
|
||||
}
|
||||
if (jobStore != null) {
|
||||
// Propagate JobStore failures: returning null on a backplane outage would conflate
|
||||
// "no such file" with "lookup unavailable" and the caller would respond 404 to a
|
||||
// transient blip that should be retried. Let Spring's exception handler surface a
|
||||
// 5xx so clients know to retry.
|
||||
try {
|
||||
return jobStore.findJobIdByFileId(fileId).orElse(null);
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("JobStore findJobIdByFileId failed for {}: {}", fileId, e.getMessage());
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -83,7 +83,16 @@ public class RequestUriUtils {
|
||||
return false;
|
||||
}
|
||||
|
||||
// Blocklist of backend/non-frontend paths that should still go through filters
|
||||
// Blocklist of backend/non-frontend paths that should still go through filters.
|
||||
//
|
||||
// `/files` was historically a backend route; it is now a frontend route
|
||||
// owned by HomePage / FileManagerView. Direct-nav or refresh on /files
|
||||
// (or /files/<folder-uuid>) was returning the Spring auth filter's 401
|
||||
// JSON instead of serving index.html, so the SPA never got a chance to
|
||||
// mount and the user saw a raw error response. There are no `/files`
|
||||
// backend mappings at the servlet root - the real storage endpoints
|
||||
// live under `/api/v1/storage/files`, which is filtered out a few lines
|
||||
// up by the `startsWith("/api/")` guard.
|
||||
String[] backendOnlyPrefixes = {
|
||||
"/register",
|
||||
"/pipeline",
|
||||
@@ -91,7 +100,6 @@ public class RequestUriUtils {
|
||||
"/pdfjs-legacy",
|
||||
"/fonts",
|
||||
"/images",
|
||||
"/files",
|
||||
"/css",
|
||||
"/js",
|
||||
"/swagger",
|
||||
@@ -181,7 +189,7 @@ public class RequestUriUtils {
|
||||
|| trimmedUri.startsWith(
|
||||
"/api/v1/mobile-scanner/") // Mobile scanner endpoints (no auth)
|
||||
|| trimmedUri.startsWith("/v1/api-docs")
|
||||
// Workflow participant endpoints — access controlled by share tokens, not login
|
||||
// Workflow participant endpoints - access controlled by share tokens, not login
|
||||
|| trimmedUri.startsWith("/api/v1/workflow/participant/")
|
||||
// Share-link SPA bootstrap; data APIs remain protected
|
||||
|| trimmedUri.matches("^/share/[^/]+/?$");
|
||||
|
||||
@@ -56,4 +56,17 @@ class ArchitectureTest {
|
||||
.resideInAPackage("stirling.software.saas..");
|
||||
rule.check(commonClasses);
|
||||
}
|
||||
|
||||
@Test
|
||||
void clusterInterfacesHaveNoImplementationDependencies() {
|
||||
ArchRule rule =
|
||||
noClasses()
|
||||
.that()
|
||||
.resideInAPackage("stirling.software.common.cluster..")
|
||||
.should()
|
||||
.dependOnClassesThat()
|
||||
.resideInAnyPackage(
|
||||
"stirling.software.proprietary..", "stirling.software.saas..");
|
||||
rule.check(commonClasses);
|
||||
}
|
||||
}
|
||||
|
||||
+51
@@ -0,0 +1,51 @@
|
||||
package stirling.software.common.cluster;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
import java.time.Instant;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
class BackplaneContractCompilationTest {
|
||||
|
||||
@Test
|
||||
void jobStoreEntryRecordRoundTrips() {
|
||||
Instant now = Instant.now();
|
||||
JobStoreEntry entry =
|
||||
new JobStoreEntry(
|
||||
"job-1",
|
||||
JobStoreEntry.JobState.PENDING,
|
||||
"node-a",
|
||||
now,
|
||||
null,
|
||||
null,
|
||||
List.of("file-1"),
|
||||
Map.of("k", "v"));
|
||||
assertEquals("job-1", entry.jobId());
|
||||
assertEquals(JobStoreEntry.JobState.PENDING, entry.state());
|
||||
assertEquals("node-a", entry.owningNodeId());
|
||||
assertEquals(now, entry.createdAt());
|
||||
assertEquals(List.of("file-1"), entry.fileIds());
|
||||
assertEquals("v", entry.resultMeta().get("k"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void clusterNodeRecordRoundTrips() {
|
||||
Instant heartbeat = Instant.now();
|
||||
ClusterNode node = new ClusterNode("node-a", "10.0.0.1:8080", heartbeat, "BOTH");
|
||||
assertEquals("node-a", node.nodeId());
|
||||
assertEquals("10.0.0.1:8080", node.internalAddress());
|
||||
assertEquals(heartbeat, node.lastHeartbeat());
|
||||
assertEquals("BOTH", node.role());
|
||||
}
|
||||
|
||||
@Test
|
||||
void rateLimitDecisionRecordRoundTrips() {
|
||||
RateLimitStore.RateLimitDecision d = new RateLimitStore.RateLimitDecision(true, 7, 0L);
|
||||
assertEquals(true, d.allowed());
|
||||
assertEquals(7, d.remainingTokens());
|
||||
assertEquals(0L, d.nanosToWaitForRefill());
|
||||
}
|
||||
}
|
||||
+65
@@ -0,0 +1,65 @@
|
||||
package stirling.software.common.cluster;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
|
||||
import java.lang.reflect.Method;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import stirling.software.common.model.ApplicationProperties;
|
||||
import stirling.software.common.model.ApplicationProperties.Cluster;
|
||||
|
||||
class ClusterConfigValidationTest {
|
||||
|
||||
@Test
|
||||
void validationPassesWhenDisabled() {
|
||||
ApplicationProperties props = new ApplicationProperties();
|
||||
ClusterConfig config = new ClusterConfig(props);
|
||||
assertDoesNotThrow(() -> invokeValidate(config));
|
||||
}
|
||||
|
||||
@Test
|
||||
void validationFailsWhenValkeyEnabledWithoutUrl() {
|
||||
ApplicationProperties props = new ApplicationProperties();
|
||||
Cluster cluster = props.getCluster();
|
||||
cluster.setEnabled(true);
|
||||
cluster.setBackplane("valkey");
|
||||
ClusterConfig config = new ClusterConfig(props);
|
||||
assertThrows(IllegalStateException.class, () -> invokeValidate(config));
|
||||
}
|
||||
|
||||
@Test
|
||||
void validationPassesWhenValkeyEnabledWithUrl() {
|
||||
ApplicationProperties props = new ApplicationProperties();
|
||||
Cluster cluster = props.getCluster();
|
||||
cluster.setEnabled(true);
|
||||
cluster.setBackplane("valkey");
|
||||
cluster.getValkey().setUrl("redis://localhost:6379");
|
||||
ClusterConfig config = new ClusterConfig(props);
|
||||
assertDoesNotThrow(() -> invokeValidate(config));
|
||||
}
|
||||
|
||||
@Test
|
||||
void validationPassesWhenInProcessEnabled() {
|
||||
ApplicationProperties props = new ApplicationProperties();
|
||||
Cluster cluster = props.getCluster();
|
||||
cluster.setEnabled(true);
|
||||
cluster.setBackplane("inprocess");
|
||||
ClusterConfig config = new ClusterConfig(props);
|
||||
assertDoesNotThrow(() -> invokeValidate(config));
|
||||
}
|
||||
|
||||
private void invokeValidate(ClusterConfig config) throws Exception {
|
||||
Method m = ClusterConfig.class.getDeclaredMethod("validate");
|
||||
m.setAccessible(true);
|
||||
try {
|
||||
m.invoke(config);
|
||||
} catch (java.lang.reflect.InvocationTargetException ex) {
|
||||
if (ex.getCause() instanceof RuntimeException re) {
|
||||
throw re;
|
||||
}
|
||||
throw ex;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,62 @@
|
||||
package stirling.software.common.cluster;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import stirling.software.common.model.ApplicationProperties;
|
||||
import stirling.software.common.model.ApplicationProperties.Cluster;
|
||||
|
||||
class ClusterPropertiesTest {
|
||||
|
||||
@Test
|
||||
void defaultsAreDisabledAndInprocess() {
|
||||
Cluster props = new ApplicationProperties().getCluster();
|
||||
assertFalse(props.isEnabled());
|
||||
assertEquals("inprocess", props.getBackplane());
|
||||
assertEquals("local", props.getArtifactStore());
|
||||
assertEquals(Cluster.NodeRole.BOTH, props.resolvedRole());
|
||||
assertEquals("", props.getValkey().getUrl());
|
||||
assertFalse(props.getValkey().getTls().isSkipCertVerification());
|
||||
assertEquals("both", props.getNode().getRole());
|
||||
assertEquals("http", props.getNode().getScheme());
|
||||
assertEquals(5000L, props.getNode().getHeartbeatIntervalMs());
|
||||
}
|
||||
|
||||
@Test
|
||||
void resolvedRoleParsesCaseInsensitively() {
|
||||
Cluster props = new ApplicationProperties().getCluster();
|
||||
props.getNode().setRole("WEB");
|
||||
assertEquals(Cluster.NodeRole.WEB, props.resolvedRole());
|
||||
|
||||
props.getNode().setRole("web");
|
||||
assertEquals(Cluster.NodeRole.WEB, props.resolvedRole());
|
||||
|
||||
props.getNode().setRole("Worker");
|
||||
assertEquals(Cluster.NodeRole.WORKER, props.resolvedRole());
|
||||
|
||||
props.getNode().setRole("garbage");
|
||||
assertEquals(Cluster.NodeRole.BOTH, props.resolvedRole());
|
||||
|
||||
props.getNode().setRole(null);
|
||||
assertEquals(Cluster.NodeRole.BOTH, props.resolvedRole());
|
||||
}
|
||||
|
||||
@Test
|
||||
void resolvedNodeIdIsStableAcrossCalls() {
|
||||
Cluster props = new ApplicationProperties().getCluster();
|
||||
String first = props.resolvedNodeId();
|
||||
String second = props.resolvedNodeId();
|
||||
assertNotNull(first);
|
||||
assertEquals(first, second);
|
||||
}
|
||||
|
||||
@Test
|
||||
void resolvedNodeIdHonoursExplicitId() {
|
||||
Cluster props = new ApplicationProperties().getCluster();
|
||||
props.getNode().setId("abc");
|
||||
assertEquals("abc", props.resolvedNodeId());
|
||||
}
|
||||
}
|
||||
+90
@@ -0,0 +1,90 @@
|
||||
package stirling.software.common.cluster;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.springframework.boot.autoconfigure.context.PropertyPlaceholderAutoConfiguration;
|
||||
import org.springframework.boot.test.context.runner.ApplicationContextRunner;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
|
||||
import stirling.software.common.cluster.inprocess.InProcessClusterConfiguration;
|
||||
import stirling.software.common.model.ApplicationProperties;
|
||||
|
||||
/**
|
||||
* Verifies the {@link InProcessClusterConfiguration} conditional wiring: in-process beans wire when
|
||||
* cluster mode is off or {@code backplane=inprocess}, and are skipped when {@code
|
||||
* backplane=valkey}.
|
||||
*/
|
||||
class InProcessConfigurationConditionalTest {
|
||||
|
||||
private final ApplicationContextRunner runner =
|
||||
new ApplicationContextRunner()
|
||||
.withConfiguration(
|
||||
org.springframework.boot.autoconfigure.AutoConfigurations.of(
|
||||
PropertyPlaceholderAutoConfiguration.class))
|
||||
.withUserConfiguration(
|
||||
TestAppPropertiesConfig.class,
|
||||
ClusterConfig.class,
|
||||
InProcessClusterConfiguration.class);
|
||||
|
||||
@Test
|
||||
void inProcessBeansWireWhenClusterDisabled() {
|
||||
runner.run(
|
||||
context ->
|
||||
assertThat(context)
|
||||
.hasNotFailed()
|
||||
.hasSingleBean(ClusterBackplane.class)
|
||||
.hasSingleBean(JobStore.class)
|
||||
.hasSingleBean(RateLimitStore.class)
|
||||
.hasSingleBean(DistributedLock.class)
|
||||
.hasSingleBean(KeyValueCache.class)
|
||||
.hasSingleBean(InstanceRegistry.class));
|
||||
}
|
||||
|
||||
@Test
|
||||
void inProcessBeansWireWhenEnabledWithInProcessBackplane() {
|
||||
runner.withPropertyValues("cluster.enabled=true", "cluster.backplane=inprocess")
|
||||
.run(
|
||||
context ->
|
||||
assertThat(context)
|
||||
.hasNotFailed()
|
||||
.hasSingleBean(ClusterBackplane.class)
|
||||
.hasSingleBean(JobStore.class)
|
||||
.hasSingleBean(RateLimitStore.class)
|
||||
.hasSingleBean(DistributedLock.class)
|
||||
.hasSingleBean(KeyValueCache.class)
|
||||
.hasSingleBean(InstanceRegistry.class));
|
||||
}
|
||||
|
||||
@Test
|
||||
void inProcessBeansSkippedWhenEnabledWithDistributedBackplane() {
|
||||
runner.withPropertyValues(
|
||||
"cluster.enabled=true",
|
||||
"cluster.backplane=valkey",
|
||||
"cluster.valkey.url=redis://localhost:6379")
|
||||
.run(
|
||||
context ->
|
||||
assertThat(context)
|
||||
.hasNotFailed()
|
||||
.doesNotHaveBean(ClusterBackplane.class)
|
||||
.doesNotHaveBean(JobStore.class)
|
||||
.doesNotHaveBean(RateLimitStore.class)
|
||||
.doesNotHaveBean(DistributedLock.class)
|
||||
.doesNotHaveBean(KeyValueCache.class)
|
||||
.doesNotHaveBean(InstanceRegistry.class));
|
||||
}
|
||||
|
||||
/**
|
||||
* Hand-rolled {@link ApplicationProperties} bean: the production class loads YAML at startup
|
||||
* via a {@code @PostConstruct} hook that isn't appropriate for the slice runner, so we wire a
|
||||
* defaults-only instance here.
|
||||
*/
|
||||
@Configuration
|
||||
static class TestAppPropertiesConfig {
|
||||
@Bean
|
||||
ApplicationProperties applicationProperties() {
|
||||
return new ApplicationProperties();
|
||||
}
|
||||
}
|
||||
}
|
||||
+167
@@ -0,0 +1,167 @@
|
||||
package stirling.software.common.cluster.inprocess;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
import java.time.Duration;
|
||||
import java.util.Optional;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import stirling.software.common.cluster.DistributedLock;
|
||||
|
||||
class InProcessDistributedLockTest {
|
||||
|
||||
@Test
|
||||
void acquireReleaseAcquire() {
|
||||
DistributedLock lock = new InProcessDistributedLock();
|
||||
DistributedLock.LockHandle h1 = lock.tryAcquire("k", Duration.ofSeconds(30)).orElseThrow();
|
||||
h1.release();
|
||||
assertTrue(lock.tryAcquire("k", Duration.ofSeconds(30)).isPresent());
|
||||
}
|
||||
|
||||
@Test
|
||||
void reentryFromSameThreadFails() {
|
||||
DistributedLock lock = new InProcessDistributedLock();
|
||||
DistributedLock.LockHandle h1 = lock.tryAcquire("k", Duration.ofSeconds(30)).orElseThrow();
|
||||
Optional<DistributedLock.LockHandle> reentry = lock.tryAcquire("k", Duration.ofSeconds(30));
|
||||
assertFalse(reentry.isPresent(), "in-process lock must be non-reentrant");
|
||||
h1.release();
|
||||
// After release, anyone can acquire again.
|
||||
assertTrue(lock.tryAcquire("k", Duration.ofSeconds(30)).isPresent());
|
||||
}
|
||||
|
||||
@Test
|
||||
void secondAcquireFromAnotherThreadFails() throws InterruptedException {
|
||||
DistributedLock lock = new InProcessDistributedLock();
|
||||
DistributedLock.LockHandle h1 = lock.tryAcquire("k", Duration.ofSeconds(30)).orElseThrow();
|
||||
|
||||
CountDownLatch done = new CountDownLatch(1);
|
||||
AtomicBoolean acquired = new AtomicBoolean(true);
|
||||
Thread t =
|
||||
new Thread(
|
||||
() -> {
|
||||
Optional<DistributedLock.LockHandle> attempt =
|
||||
lock.tryAcquire("k", Duration.ofSeconds(30));
|
||||
acquired.set(attempt.isPresent());
|
||||
attempt.ifPresent(DistributedLock.LockHandle::release);
|
||||
done.countDown();
|
||||
});
|
||||
t.start();
|
||||
assertTrue(done.await(2, TimeUnit.SECONDS));
|
||||
assertFalse(acquired.get());
|
||||
h1.release();
|
||||
}
|
||||
|
||||
@Test
|
||||
void leaseExpiryAllowsTakeoverEvenWithoutRelease() throws InterruptedException {
|
||||
// Acquire with a short lease, never call release, then try to acquire again after the
|
||||
// lease has elapsed. Matches Redis SET-NX-EX semantics - the second caller gets the lock
|
||||
// because the first lease auto-expired. 250ms lease + 350ms wait gives CI generous slack.
|
||||
DistributedLock lock = new InProcessDistributedLock();
|
||||
DistributedLock.LockHandle h1 = lock.tryAcquire("k", Duration.ofMillis(250)).orElseThrow();
|
||||
Thread.sleep(350);
|
||||
Optional<DistributedLock.LockHandle> takeover =
|
||||
lock.tryAcquire("k", Duration.ofSeconds(30));
|
||||
assertTrue(
|
||||
takeover.isPresent(),
|
||||
"expired lease must release the lock so a new caller can take over");
|
||||
// Calling release() on the original handle after takeover must be a no-op (token check).
|
||||
h1.release();
|
||||
// The takeover holder is still the legitimate owner.
|
||||
assertFalse(lock.tryAcquire("k", Duration.ofSeconds(30)).isPresent());
|
||||
takeover.get().release();
|
||||
}
|
||||
|
||||
@Test
|
||||
void renewExtendsLease() throws InterruptedException {
|
||||
// Acquire with a short lease, renew it before it expires, then verify the lock is still
|
||||
// held past the original expiry point. 200ms initial + renew to 2s + wait 350ms.
|
||||
DistributedLock lock = new InProcessDistributedLock();
|
||||
DistributedLock.LockHandle h1 = lock.tryAcquire("k", Duration.ofMillis(200)).orElseThrow();
|
||||
assertTrue(h1.renew(Duration.ofSeconds(2)), "renew on a held lease must succeed");
|
||||
Thread.sleep(350);
|
||||
assertFalse(
|
||||
lock.tryAcquire("k", Duration.ofSeconds(30)).isPresent(),
|
||||
"renew should have pushed expiry well past the original 200ms");
|
||||
h1.release();
|
||||
}
|
||||
|
||||
@Test
|
||||
void renewAfterReleaseFails() {
|
||||
DistributedLock lock = new InProcessDistributedLock();
|
||||
DistributedLock.LockHandle h1 = lock.tryAcquire("k", Duration.ofSeconds(30)).orElseThrow();
|
||||
h1.release();
|
||||
assertFalse(h1.renew(Duration.ofSeconds(30)), "renew on a released handle must fail");
|
||||
}
|
||||
|
||||
/**
|
||||
* Concurrency stress: many threads contending on the same key with each holder respecting the
|
||||
* lease (hold << lease). The lock behaves as a strict mutex in this regime so asserting
|
||||
* mutual exclusion is meaningful. A separate test ({@link
|
||||
* #leaseExpiryAllowsTakeoverEvenWithoutRelease}) covers the takeover-across-expiry branch,
|
||||
* which legitimately allows two holders momentarily and is split-brain behaviour inherent to
|
||||
* any lease-based lock.
|
||||
*/
|
||||
@Test
|
||||
void concurrentContentionPreservesMutualExclusion() throws InterruptedException {
|
||||
DistributedLock lock = new InProcessDistributedLock();
|
||||
int threads = 16;
|
||||
int attemptsPerThread = 200;
|
||||
// Lease far exceeds any plausible hold time, so the takeover branch never triggers in
|
||||
// this test and the lock acts as a strict mutex.
|
||||
Duration lease = Duration.ofSeconds(5);
|
||||
java.util.concurrent.atomic.AtomicInteger concurrentHolders =
|
||||
new java.util.concurrent.atomic.AtomicInteger();
|
||||
java.util.concurrent.atomic.AtomicInteger maxConcurrent =
|
||||
new java.util.concurrent.atomic.AtomicInteger();
|
||||
java.util.concurrent.atomic.AtomicInteger acquires =
|
||||
new java.util.concurrent.atomic.AtomicInteger();
|
||||
java.util.concurrent.atomic.AtomicReference<Throwable> firstFailure =
|
||||
new java.util.concurrent.atomic.AtomicReference<>();
|
||||
CountDownLatch start = new CountDownLatch(1);
|
||||
CountDownLatch done = new CountDownLatch(threads);
|
||||
|
||||
for (int i = 0; i < threads; i++) {
|
||||
new Thread(
|
||||
() -> {
|
||||
try {
|
||||
start.await();
|
||||
for (int j = 0; j < attemptsPerThread; j++) {
|
||||
Optional<DistributedLock.LockHandle> h =
|
||||
lock.tryAcquire("hot", lease);
|
||||
if (h.isPresent()) {
|
||||
int now = concurrentHolders.incrementAndGet();
|
||||
maxConcurrent.accumulateAndGet(now, Math::max);
|
||||
acquires.incrementAndGet();
|
||||
// Trivial critical section; well within lease.
|
||||
concurrentHolders.decrementAndGet();
|
||||
h.get().release();
|
||||
}
|
||||
}
|
||||
} catch (Throwable t) {
|
||||
firstFailure.compareAndSet(null, t);
|
||||
} finally {
|
||||
done.countDown();
|
||||
}
|
||||
},
|
||||
"lock-stress-" + i)
|
||||
.start();
|
||||
}
|
||||
start.countDown();
|
||||
assertTrue(done.await(30, TimeUnit.SECONDS), "stress workers must finish in time");
|
||||
org.junit.jupiter.api.Assertions.assertNull(firstFailure.get(), "no worker may throw");
|
||||
org.junit.jupiter.api.Assertions.assertEquals(
|
||||
1,
|
||||
maxConcurrent.get(),
|
||||
"mutual exclusion violated: more than one holder observed simultaneously");
|
||||
assertTrue(
|
||||
acquires.get() > 0,
|
||||
"at least some acquires must succeed under contention (saw "
|
||||
+ acquires.get()
|
||||
+ ")");
|
||||
}
|
||||
}
|
||||
+28
@@ -0,0 +1,28 @@
|
||||
package stirling.software.common.cluster.inprocess;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
import java.time.Duration;
|
||||
import java.time.Instant;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import stirling.software.common.cluster.ClusterNode;
|
||||
|
||||
class InProcessInstanceRegistryTest {
|
||||
|
||||
@Test
|
||||
void registerThenLookupAndActiveNodes() {
|
||||
InProcessInstanceRegistry registry = new InProcessInstanceRegistry();
|
||||
ClusterNode node = new ClusterNode("node-1", "127.0.0.1:8080", Instant.now(), "BOTH");
|
||||
registry.register(node, Duration.ofSeconds(30));
|
||||
|
||||
assertTrue(registry.lookup("node-1").isPresent());
|
||||
assertEquals("node-1", registry.lookup("node-1").get().nodeId());
|
||||
assertEquals(1, registry.activeNodes().size());
|
||||
|
||||
registry.deregister("node-1");
|
||||
assertTrue(registry.lookup("node-1").isEmpty());
|
||||
}
|
||||
}
|
||||
+91
@@ -0,0 +1,91 @@
|
||||
package stirling.software.common.cluster.inprocess;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
import java.time.Duration;
|
||||
import java.time.Instant;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import stirling.software.common.cluster.JobStoreEntry;
|
||||
|
||||
class InProcessJobStoreTest {
|
||||
|
||||
private final InProcessJobStore store = new InProcessJobStore();
|
||||
|
||||
@Test
|
||||
void putGetDeleteExistsRoundTrip() {
|
||||
JobStoreEntry entry = entry("job-1");
|
||||
store.put(entry, Duration.ofMinutes(30));
|
||||
|
||||
assertTrue(store.exists("job-1"));
|
||||
assertEquals(entry, store.get("job-1").orElseThrow());
|
||||
|
||||
store.delete("job-1");
|
||||
assertFalse(store.exists("job-1"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void ttlExpiry() throws InterruptedException {
|
||||
store.put(entry("job-2"), Duration.ofMillis(50));
|
||||
Thread.sleep(100);
|
||||
assertFalse(store.get("job-2").isPresent());
|
||||
}
|
||||
|
||||
@Test
|
||||
void purgeExpiredRemovesOnlyStaleEntries() throws InterruptedException {
|
||||
store.put(entry("job-fresh"), Duration.ofMinutes(30));
|
||||
store.put(entry("job-stale"), Duration.ofMillis(20));
|
||||
Thread.sleep(80);
|
||||
|
||||
int removed = store.purgeExpired();
|
||||
assertEquals(1, removed);
|
||||
assertTrue(store.exists("job-fresh"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void findJobIdByFileIdReturnsTheRightJob() {
|
||||
store.put(
|
||||
new JobStoreEntry(
|
||||
"job-a",
|
||||
JobStoreEntry.JobState.COMPLETE,
|
||||
"node-1",
|
||||
Instant.now(),
|
||||
Instant.now(),
|
||||
null,
|
||||
List.of("file-1", "file-2"),
|
||||
Map.of()),
|
||||
Duration.ofMinutes(30));
|
||||
store.put(
|
||||
new JobStoreEntry(
|
||||
"job-b",
|
||||
JobStoreEntry.JobState.COMPLETE,
|
||||
"node-1",
|
||||
Instant.now(),
|
||||
Instant.now(),
|
||||
null,
|
||||
List.of("file-3"),
|
||||
Map.of()),
|
||||
Duration.ofMinutes(30));
|
||||
|
||||
assertEquals("job-a", store.findJobIdByFileId("file-1").orElseThrow());
|
||||
assertEquals("job-b", store.findJobIdByFileId("file-3").orElseThrow());
|
||||
assertFalse(store.findJobIdByFileId("missing").isPresent());
|
||||
}
|
||||
|
||||
private JobStoreEntry entry(String id) {
|
||||
return new JobStoreEntry(
|
||||
id,
|
||||
JobStoreEntry.JobState.PENDING,
|
||||
"node-1",
|
||||
Instant.now(),
|
||||
null,
|
||||
null,
|
||||
List.of(),
|
||||
Map.of());
|
||||
}
|
||||
}
|
||||
+41
@@ -0,0 +1,41 @@
|
||||
package stirling.software.common.cluster.inprocess;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
|
||||
import java.time.Duration;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import stirling.software.common.cluster.KeyValueCache;
|
||||
|
||||
class InProcessKeyValueCacheTest {
|
||||
|
||||
@Test
|
||||
void putGetEvict() {
|
||||
KeyValueCache cache = new InProcessKeyValueCache();
|
||||
cache.put("apikey", "a", "userA", Duration.ofMinutes(1));
|
||||
assertEquals("userA", cache.get("apikey", "a").orElseThrow());
|
||||
|
||||
cache.evict("apikey", "a");
|
||||
assertFalse(cache.get("apikey", "a").isPresent());
|
||||
}
|
||||
|
||||
@Test
|
||||
void ttlExpiry() throws InterruptedException {
|
||||
KeyValueCache cache = new InProcessKeyValueCache();
|
||||
cache.put("ns", "k", "v", Duration.ofMillis(40));
|
||||
Thread.sleep(80);
|
||||
assertFalse(cache.get("ns", "k").isPresent());
|
||||
}
|
||||
|
||||
@Test
|
||||
void evictNamespace() {
|
||||
KeyValueCache cache = new InProcessKeyValueCache();
|
||||
cache.put("ns", "a", "1", Duration.ofMinutes(1));
|
||||
cache.put("ns", "b", "2", Duration.ofMinutes(1));
|
||||
cache.evictNamespace("ns");
|
||||
assertFalse(cache.get("ns", "a").isPresent());
|
||||
assertFalse(cache.get("ns", "b").isPresent());
|
||||
}
|
||||
}
|
||||
+56
@@ -0,0 +1,56 @@
|
||||
package stirling.software.common.cluster.inprocess;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
import java.time.Duration;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import stirling.software.common.cluster.RateLimitStore;
|
||||
import stirling.software.common.cluster.RateLimitStore.RateLimitDecision;
|
||||
|
||||
class InProcessRateLimitStoreTest {
|
||||
|
||||
@Test
|
||||
void firstNConsumesAllowed() {
|
||||
RateLimitStore store = new InProcessRateLimitStore();
|
||||
for (int i = 0; i < 5; i++) {
|
||||
assertTrue(store.tryConsume("k", 5, Duration.ofSeconds(60)).allowed(), "i=" + i);
|
||||
}
|
||||
assertFalse(store.tryConsume("k", 5, Duration.ofSeconds(60)).allowed());
|
||||
}
|
||||
|
||||
@Test
|
||||
void remainingTokensDecrements() {
|
||||
RateLimitStore store = new InProcessRateLimitStore();
|
||||
RateLimitDecision d1 = store.tryConsume("k", 5, Duration.ofSeconds(60));
|
||||
RateLimitDecision d2 = store.tryConsume("k", 5, Duration.ofSeconds(60));
|
||||
assertTrue(d1.allowed());
|
||||
assertTrue(d2.allowed());
|
||||
assertEquals(4, d1.remainingTokens());
|
||||
assertEquals(3, d2.remainingTokens());
|
||||
}
|
||||
|
||||
@Test
|
||||
void refillRestoresTokens() throws InterruptedException {
|
||||
RateLimitStore store = new InProcessRateLimitStore();
|
||||
// Capacity 2 with smooth refill over 100 ms -> ~1 token per 50 ms.
|
||||
for (int i = 0; i < 2; i++) {
|
||||
assertTrue(store.tryConsume("k", 2, Duration.ofMillis(100)).allowed());
|
||||
}
|
||||
assertFalse(store.tryConsume("k", 2, Duration.ofMillis(100)).allowed());
|
||||
Thread.sleep(150);
|
||||
assertTrue(store.tryConsume("k", 2, Duration.ofMillis(100)).allowed());
|
||||
}
|
||||
|
||||
@Test
|
||||
void deniedConsumeReportsWaitNanos() {
|
||||
RateLimitStore store = new InProcessRateLimitStore();
|
||||
assertTrue(store.tryConsume("wait", 1, Duration.ofSeconds(10)).allowed());
|
||||
RateLimitDecision denied = store.tryConsume("wait", 1, Duration.ofSeconds(10));
|
||||
assertFalse(denied.allowed());
|
||||
assertTrue(denied.nanosToWaitForRefill() > 0L);
|
||||
}
|
||||
}
|
||||
+42
@@ -0,0 +1,42 @@
|
||||
package stirling.software.common.cluster.inprocess;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertArrayEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
import java.io.ByteArrayInputStream;
|
||||
import java.io.IOException;
|
||||
import java.nio.file.Path;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import stirling.software.common.cluster.FileStore;
|
||||
|
||||
class LocalDiskFileStoreTest {
|
||||
|
||||
@Test
|
||||
void storeRetrieveSizeDeleteExistsRoundTrip(@TempDir Path dir) throws IOException {
|
||||
LocalDiskFileStore store = new LocalDiskFileStore(dir.toString());
|
||||
byte[] payload = "hello-bytes".getBytes();
|
||||
|
||||
FileStore.Stored stored = store.store(new ByteArrayInputStream(payload), "x.txt");
|
||||
assertEquals(payload.length, stored.size());
|
||||
assertTrue(store.exists(stored.fileId()));
|
||||
assertEquals(payload.length, store.size(stored.fileId()));
|
||||
assertArrayEquals(payload, store.retrieveBytes(stored.fileId()));
|
||||
|
||||
assertTrue(store.delete(stored.fileId()));
|
||||
assertFalse(store.exists(stored.fileId()));
|
||||
}
|
||||
|
||||
@Test
|
||||
void traversalIdsAreRejected(@TempDir Path dir) {
|
||||
LocalDiskFileStore store = new LocalDiskFileStore(dir.toString());
|
||||
assertThrows(IllegalArgumentException.class, () -> store.resolve("../foo"));
|
||||
assertThrows(IllegalArgumentException.class, () -> store.resolve("a/b"));
|
||||
assertThrows(IllegalArgumentException.class, () -> store.resolve("a\\b"));
|
||||
}
|
||||
}
|
||||
+95
@@ -0,0 +1,95 @@
|
||||
package stirling.software.common.configuration;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.snakeyaml.engine.v2.api.LoadSettings;
|
||||
|
||||
import stirling.software.common.util.YamlHelper;
|
||||
|
||||
class ConfigInitializerTest {
|
||||
|
||||
private static final LoadSettings LOAD_SETTINGS =
|
||||
LoadSettings.builder()
|
||||
.setUseMarks(true)
|
||||
.setMaxAliasesForCollections(Integer.MAX_VALUE)
|
||||
.setAllowRecursiveKeys(true)
|
||||
.setParseComments(true)
|
||||
.build();
|
||||
|
||||
// Mirrors the proFeatures block of settings.yml.template after the camelCase rename.
|
||||
private static final String CAMEL_CASE_TEMPLATE =
|
||||
"""
|
||||
premium:
|
||||
proFeatures:
|
||||
ssoAutoLogin: false
|
||||
customMetadata:
|
||||
autoUpdateMetadata: false
|
||||
author: username
|
||||
creator: Stirling-PDF
|
||||
producer: Stirling-PDF
|
||||
""";
|
||||
|
||||
@Test
|
||||
void migrateProFeaturesKeyCasing_carriesForwardLegacyPascalCaseValues() {
|
||||
// An existing install whose settings.yml still uses the old PascalCase keys.
|
||||
String legacy =
|
||||
"""
|
||||
premium:
|
||||
proFeatures:
|
||||
SSOAutoLogin: true
|
||||
CustomMetadata:
|
||||
autoUpdateMetadata: true
|
||||
author: alice
|
||||
creator: bob
|
||||
producer: carol
|
||||
""";
|
||||
YamlHelper template = new YamlHelper(LOAD_SETTINGS, CAMEL_CASE_TEMPLATE);
|
||||
YamlHelper existing = new YamlHelper(LOAD_SETTINGS, legacy);
|
||||
|
||||
new ConfigInitializer().migrateProFeaturesKeyCasing(existing, template);
|
||||
|
||||
assertEquals(
|
||||
"true", template.getValueByExactKeyPath("premium", "proFeatures", "ssoAutoLogin"));
|
||||
assertEquals(
|
||||
"true",
|
||||
template.getValueByExactKeyPath(
|
||||
"premium", "proFeatures", "customMetadata", "autoUpdateMetadata"));
|
||||
assertEquals(
|
||||
"alice",
|
||||
template.getValueByExactKeyPath(
|
||||
"premium", "proFeatures", "customMetadata", "author"));
|
||||
assertEquals(
|
||||
"bob",
|
||||
template.getValueByExactKeyPath(
|
||||
"premium", "proFeatures", "customMetadata", "creator"));
|
||||
assertEquals(
|
||||
"carol",
|
||||
template.getValueByExactKeyPath(
|
||||
"premium", "proFeatures", "customMetadata", "producer"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void migrateProFeaturesKeyCasing_withoutLegacyKeys_keepsTemplateDefaults() {
|
||||
// No PascalCase keys present -> this migration step must be a no-op.
|
||||
String alreadyCamel =
|
||||
"""
|
||||
premium:
|
||||
proFeatures:
|
||||
ssoAutoLogin: true
|
||||
customMetadata:
|
||||
author: dave
|
||||
""";
|
||||
YamlHelper template = new YamlHelper(LOAD_SETTINGS, CAMEL_CASE_TEMPLATE);
|
||||
YamlHelper existing = new YamlHelper(LOAD_SETTINGS, alreadyCamel);
|
||||
|
||||
new ConfigInitializer().migrateProFeaturesKeyCasing(existing, template);
|
||||
|
||||
assertEquals(
|
||||
"false", template.getValueByExactKeyPath("premium", "proFeatures", "ssoAutoLogin"));
|
||||
assertEquals(
|
||||
"username",
|
||||
template.getValueByExactKeyPath(
|
||||
"premium", "proFeatures", "customMetadata", "author"));
|
||||
}
|
||||
}
|
||||
+27
@@ -0,0 +1,27 @@
|
||||
package stirling.software.common.service;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertArrayEquals;
|
||||
import static org.mockito.Mockito.mock;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.nio.file.Path;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import stirling.software.common.cluster.inprocess.LocalDiskFileStore;
|
||||
|
||||
class FileStorageDelegationTest {
|
||||
|
||||
@Test
|
||||
void storeBytesThenRetrieveBytesRoundTripsThroughFileStore(@TempDir Path tempDir)
|
||||
throws IOException {
|
||||
FileStorage fs =
|
||||
new FileStorage(
|
||||
mock(FileOrUploadService.class),
|
||||
new LocalDiskFileStore(tempDir.toString()));
|
||||
byte[] payload = "round-trip".getBytes();
|
||||
String id = fs.storeBytes(payload, "x.bin");
|
||||
assertArrayEquals(payload, fs.retrieveBytes(id));
|
||||
}
|
||||
}
|
||||
@@ -3,6 +3,7 @@ package stirling.software.common.service;
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
import static org.mockito.Mockito.*;
|
||||
|
||||
import java.io.ByteArrayInputStream;
|
||||
import java.io.IOException;
|
||||
import java.io.InputStream;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
@@ -13,29 +14,30 @@ import java.util.stream.Stream;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
import org.mockito.InjectMocks;
|
||||
import org.mockito.Mock;
|
||||
import org.mockito.MockitoAnnotations;
|
||||
import org.springframework.core.io.ByteArrayResource;
|
||||
import org.springframework.core.io.Resource;
|
||||
import org.springframework.http.MediaType;
|
||||
import org.springframework.test.util.ReflectionTestUtils;
|
||||
import org.springframework.web.multipart.MultipartFile;
|
||||
|
||||
import stirling.software.common.cluster.inprocess.LocalDiskFileStore;
|
||||
|
||||
class FileStorageTest {
|
||||
|
||||
@TempDir Path tempDir;
|
||||
|
||||
@Mock private FileOrUploadService fileOrUploadService;
|
||||
|
||||
@InjectMocks private FileStorage fileStorage;
|
||||
private FileStorage fileStorage;
|
||||
|
||||
private MultipartFile mockFile;
|
||||
|
||||
@BeforeEach
|
||||
void setUp() {
|
||||
void setUp() throws IOException {
|
||||
MockitoAnnotations.openMocks(this);
|
||||
ReflectionTestUtils.setField(fileStorage, "tempDirPath", tempDir.toString());
|
||||
fileStorage =
|
||||
new FileStorage(fileOrUploadService, new LocalDiskFileStore(tempDir.toString()));
|
||||
|
||||
// Create a mock MultipartFile
|
||||
mockFile = mock(MultipartFile.class);
|
||||
@@ -47,17 +49,7 @@ class FileStorageTest {
|
||||
void testStoreFile() throws IOException {
|
||||
// Arrange
|
||||
byte[] fileContent = "Test PDF content".getBytes();
|
||||
when(mockFile.getBytes()).thenReturn(fileContent);
|
||||
|
||||
// Set up mock to handle transferTo by writing the file
|
||||
doAnswer(
|
||||
invocation -> {
|
||||
java.io.File file = invocation.getArgument(0);
|
||||
Files.write(file.toPath(), fileContent);
|
||||
return null;
|
||||
})
|
||||
.when(mockFile)
|
||||
.transferTo(any(java.io.File.class));
|
||||
when(mockFile.getInputStream()).thenReturn(new ByteArrayInputStream(fileContent));
|
||||
|
||||
// Act
|
||||
String fileId = fileStorage.storeFile(mockFile);
|
||||
@@ -65,7 +57,7 @@ class FileStorageTest {
|
||||
// Assert
|
||||
assertNotNull(fileId);
|
||||
assertTrue(Files.exists(tempDir.resolve(fileId)));
|
||||
verify(mockFile).transferTo(any(java.io.File.class));
|
||||
assertArrayEquals(fileContent, Files.readAllBytes(tempDir.resolve(fileId)));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -247,11 +239,11 @@ class FileStorageTest {
|
||||
filesBefore = s.count();
|
||||
}
|
||||
|
||||
// Act + Assert: IOException must propagate out — not be swallowed.
|
||||
// Act + Assert: IOException must propagate out - not be swallowed.
|
||||
assertThrows(
|
||||
IOException.class, () -> fileStorage.storeFromResource(flakyResource, "n.pdf"));
|
||||
|
||||
// Assert: no partial file lingers under the storage directory — the finally
|
||||
// Assert: no partial file lingers under the storage directory - the finally
|
||||
// branch's deleteIfExists must have cleaned it up.
|
||||
long filesAfter;
|
||||
try (Stream<Path> s = Files.list(tempDir)) {
|
||||
|
||||
+120
@@ -0,0 +1,120 @@
|
||||
package stirling.software.common.service;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.mockito.ArgumentMatchers.any;
|
||||
import static org.mockito.Mockito.never;
|
||||
import static org.mockito.Mockito.spy;
|
||||
import static org.mockito.Mockito.verify;
|
||||
|
||||
import java.time.LocalDateTime;
|
||||
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.mockito.Mock;
|
||||
import org.mockito.MockitoAnnotations;
|
||||
import org.springframework.test.util.ReflectionTestUtils;
|
||||
|
||||
import stirling.software.common.cluster.ClusterBackplane;
|
||||
import stirling.software.common.cluster.JobStoreEntry;
|
||||
import stirling.software.common.cluster.JobStoreEntry.JobState;
|
||||
import stirling.software.common.cluster.inprocess.InProcessClusterBackplane;
|
||||
import stirling.software.common.cluster.inprocess.InProcessJobStore;
|
||||
import stirling.software.common.model.ApplicationProperties;
|
||||
import stirling.software.common.model.job.JobResult;
|
||||
|
||||
class TaskManagerJobStoreDelegationTest {
|
||||
|
||||
@Mock private FileStorage fileStorage;
|
||||
|
||||
private InProcessJobStore jobStore;
|
||||
private ClusterBackplane backplane;
|
||||
private TaskManager taskManager;
|
||||
|
||||
@BeforeEach
|
||||
void setUp() {
|
||||
MockitoAnnotations.openMocks(this);
|
||||
jobStore = spy(new InProcessJobStore());
|
||||
backplane = new InProcessClusterBackplane(new ApplicationProperties());
|
||||
taskManager = new TaskManager(fileStorage, jobStore, backplane);
|
||||
ReflectionTestUtils.setField(taskManager, "jobResultExpiryMinutes", 30);
|
||||
}
|
||||
|
||||
@Test
|
||||
void createTaskWritesPendingEntry() {
|
||||
taskManager.createTask("job-1");
|
||||
JobStoreEntry entry = jobStore.get("job-1").orElseThrow();
|
||||
assertEquals(JobState.PENDING, entry.state());
|
||||
assertEquals(backplane.localNodeId(), entry.owningNodeId());
|
||||
}
|
||||
|
||||
@Test
|
||||
void setCompleteFlipsToComplete() {
|
||||
taskManager.createTask("job-2");
|
||||
taskManager.setResult("job-2", "ok");
|
||||
taskManager.setComplete("job-2");
|
||||
JobStoreEntry entry = jobStore.get("job-2").orElseThrow();
|
||||
assertEquals(JobState.COMPLETE, entry.state());
|
||||
}
|
||||
|
||||
@Test
|
||||
void setErrorFlipsToFailed() {
|
||||
taskManager.createTask("job-3");
|
||||
taskManager.setError("job-3", "boom");
|
||||
JobStoreEntry entry = jobStore.get("job-3").orElseThrow();
|
||||
assertEquals(JobState.FAILED, entry.state());
|
||||
assertEquals("boom", entry.error());
|
||||
}
|
||||
|
||||
@Test
|
||||
void cleanupOldJobsIsNoopWhenBackplaneIsNotInProcess() {
|
||||
ClusterBackplane mockedValkeyBackplane =
|
||||
new ClusterBackplane() {
|
||||
@Override
|
||||
public boolean isHealthy() {
|
||||
return true;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String backplaneType() {
|
||||
return "valkey";
|
||||
}
|
||||
|
||||
@Override
|
||||
public String localNodeId() {
|
||||
return "node-1";
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean shouldRunLocalCleanup() {
|
||||
return false;
|
||||
}
|
||||
};
|
||||
TaskManager tm = new TaskManager(fileStorage, jobStore, mockedValkeyBackplane);
|
||||
ReflectionTestUtils.setField(tm, "jobResultExpiryMinutes", 30);
|
||||
tm.createTask("job-4");
|
||||
tm.setComplete("job-4");
|
||||
ageJobPastExpiry(tm, "job-4");
|
||||
tm.cleanupOldJobs();
|
||||
// cleanup must short-circuit before touching jobStore in cluster mode; the backplane
|
||||
// TTL owns expiry there. If the gate fired correctly, delete is never called.
|
||||
verify(jobStore, never()).delete(any());
|
||||
}
|
||||
|
||||
@Test
|
||||
void cleanupOldJobsDeletesFromJobStoreWhenBackplaneIsInProcess() {
|
||||
taskManager.createTask("job-5");
|
||||
taskManager.setComplete("job-5");
|
||||
ageJobPastExpiry(taskManager, "job-5");
|
||||
taskManager.cleanupOldJobs();
|
||||
verify(jobStore).delete("job-5");
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private static void ageJobPastExpiry(TaskManager tm, String jobId) {
|
||||
var jobResults =
|
||||
(java.util.Map<String, JobResult>) ReflectionTestUtils.getField(tm, "jobResults");
|
||||
JobResult result = jobResults.get(jobId);
|
||||
ReflectionTestUtils.setField(result, "completedAt", LocalDateTime.now().minusHours(2));
|
||||
ReflectionTestUtils.setField(result, "complete", true);
|
||||
}
|
||||
}
|
||||
@@ -1,20 +1,28 @@
|
||||
package stirling.software.common.service;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
import static org.mockito.ArgumentMatchers.any;
|
||||
import static org.mockito.ArgumentMatchers.anyString;
|
||||
import static org.mockito.Mockito.*;
|
||||
|
||||
import java.time.LocalDateTime;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.mockito.ArgumentCaptor;
|
||||
import org.mockito.InjectMocks;
|
||||
import org.mockito.Mock;
|
||||
import org.mockito.MockitoAnnotations;
|
||||
import org.springframework.http.MediaType;
|
||||
import org.springframework.test.util.ReflectionTestUtils;
|
||||
|
||||
import stirling.software.common.cluster.ClusterBackplane;
|
||||
import stirling.software.common.cluster.JobStore;
|
||||
import stirling.software.common.cluster.JobStoreEntry;
|
||||
import stirling.software.common.cluster.JobStoreEntry.JobState;
|
||||
import stirling.software.common.model.job.JobResult;
|
||||
import stirling.software.common.model.job.JobStats;
|
||||
import stirling.software.common.model.job.ResultFile;
|
||||
@@ -22,6 +30,8 @@ import stirling.software.common.model.job.ResultFile;
|
||||
class TaskManagerTest {
|
||||
|
||||
@Mock private FileStorage fileStorage;
|
||||
@Mock private JobStore jobStore;
|
||||
@Mock private ClusterBackplane clusterBackplane;
|
||||
|
||||
@InjectMocks private TaskManager taskManager;
|
||||
|
||||
@@ -30,6 +40,10 @@ class TaskManagerTest {
|
||||
@BeforeEach
|
||||
void setUp() {
|
||||
closeable = MockitoAnnotations.openMocks(this);
|
||||
// Treat the backplane as in-process so cleanupOldJobs is not short-circuited.
|
||||
lenient().when(clusterBackplane.backplaneType()).thenReturn("inprocess");
|
||||
lenient().when(clusterBackplane.localNodeId()).thenReturn("test-node");
|
||||
lenient().when(clusterBackplane.shouldRunLocalCleanup()).thenReturn(true);
|
||||
ReflectionTestUtils.setField(taskManager, "jobResultExpiryMinutes", 30);
|
||||
}
|
||||
|
||||
@@ -270,6 +284,33 @@ class TaskManagerTest {
|
||||
verify(fileStorage).deleteFile("file-id");
|
||||
}
|
||||
|
||||
@Test
|
||||
void testCleanupOldJobs_NoOpWhenBackplaneOwnsExpiry() {
|
||||
// When the backplane reports it should NOT run local cleanup (e.g. a distributed
|
||||
// backplane with its own TTL), the cleanup loop must leave local state untouched.
|
||||
when(clusterBackplane.shouldRunLocalCleanup()).thenReturn(false);
|
||||
|
||||
// Seed an old completed job that would normally be removed.
|
||||
String oldJobId = "old-job-distributed";
|
||||
taskManager.createTask(oldJobId);
|
||||
JobResult oldJob = taskManager.getJobResult(oldJobId);
|
||||
ReflectionTestUtils.setField(oldJob, "completedAt", LocalDateTime.now().minusHours(1));
|
||||
ReflectionTestUtils.setField(oldJob, "complete", true);
|
||||
|
||||
Map<String, JobResult> jobResultsMap =
|
||||
(Map<String, JobResult>) ReflectionTestUtils.getField(taskManager, "jobResults");
|
||||
assertNotNull(jobResultsMap);
|
||||
assertTrue(jobResultsMap.containsKey(oldJobId));
|
||||
|
||||
// Act
|
||||
taskManager.cleanupOldJobs();
|
||||
|
||||
// Assert: nothing was removed locally, and no jobStore.delete was issued.
|
||||
assertTrue(jobResultsMap.containsKey(oldJobId));
|
||||
verify(jobStore, never()).delete(anyString());
|
||||
verify(fileStorage, never()).deleteFile(anyString());
|
||||
}
|
||||
|
||||
@Test
|
||||
void testShutdown() {
|
||||
// This mainly tests that the shutdown method doesn't throw exceptions
|
||||
@@ -310,4 +351,33 @@ class TaskManagerTest {
|
||||
// Assert
|
||||
assertFalse(result);
|
||||
}
|
||||
|
||||
@Test
|
||||
void testWriteThroughOnUpdate() {
|
||||
// Mutating calls must write through to the injected JobStore.
|
||||
String jobId = "write-through-job";
|
||||
taskManager.createTask(jobId);
|
||||
taskManager.setResult(jobId, "done");
|
||||
|
||||
ArgumentCaptor<JobStoreEntry> captor = ArgumentCaptor.forClass(JobStoreEntry.class);
|
||||
verify(jobStore, atLeast(2)).put(captor.capture(), any());
|
||||
|
||||
JobStoreEntry last = captor.getValue();
|
||||
assertEquals(jobId, last.jobId());
|
||||
assertEquals(JobState.COMPLETE, last.state());
|
||||
assertEquals("test-node", last.owningNodeId());
|
||||
}
|
||||
|
||||
@Test
|
||||
void testFindJobKeyByFileId_FallsBackToJobStore() {
|
||||
// When the file id is not in the local map, TaskManager delegates to JobStore.
|
||||
String fileId = "remote-file-id";
|
||||
String expectedJobKey = "remote-job-key";
|
||||
when(jobStore.findJobIdByFileId(fileId)).thenReturn(Optional.of(expectedJobKey));
|
||||
|
||||
String actual = taskManager.findJobKeyByFileId(fileId);
|
||||
|
||||
assertEquals(expectedJobKey, actual);
|
||||
verify(jobStore).findJobIdByFileId(fileId);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2,6 +2,8 @@ package stirling.software.common.util;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collections;
|
||||
import java.util.List;
|
||||
@@ -12,9 +14,47 @@ import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
import org.mockito.MockedStatic;
|
||||
import org.mockito.Mockito;
|
||||
|
||||
import stirling.software.common.configuration.InstallationPathConfig;
|
||||
|
||||
public class GeneralUtilsTest {
|
||||
|
||||
// Regression guard for the SSO auto-login persistence bug: the admin UI writes camelCase
|
||||
// proFeatures keys, so saveKeyToSettings must match (and persist) them against the camelCase
|
||||
// settings.yml.template. A case mismatch makes YamlHelper.updateValue silently no-op.
|
||||
@Test
|
||||
void saveKeyToSettings_persistsCamelCaseProFeatureKeys(@TempDir Path tempDir) throws Exception {
|
||||
Path settings = tempDir.resolve("settings.yml");
|
||||
Files.writeString(
|
||||
settings,
|
||||
"""
|
||||
premium:
|
||||
proFeatures:
|
||||
ssoAutoLogin: false
|
||||
customMetadata:
|
||||
author: username
|
||||
""");
|
||||
|
||||
try (MockedStatic<InstallationPathConfig> mocked =
|
||||
Mockito.mockStatic(InstallationPathConfig.class)) {
|
||||
mocked.when(InstallationPathConfig::getSettingsPath).thenReturn(settings.toString());
|
||||
|
||||
GeneralUtils.saveKeyToSettings("premium.proFeatures.ssoAutoLogin", true);
|
||||
GeneralUtils.saveKeyToSettings("premium.proFeatures.customMetadata.author", "alice");
|
||||
}
|
||||
|
||||
YamlHelper reloaded = new YamlHelper(settings);
|
||||
assertEquals(
|
||||
"true", reloaded.getValueByExactKeyPath("premium", "proFeatures", "ssoAutoLogin"));
|
||||
assertEquals(
|
||||
"alice",
|
||||
reloaded.getValueByExactKeyPath(
|
||||
"premium", "proFeatures", "customMetadata", "author"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void testParsePageListWithAll() {
|
||||
List<Integer> result = GeneralUtils.parsePageList(new String[] {"all"}, 5, false);
|
||||
|
||||
@@ -98,6 +98,17 @@ class RequestUriUtilsTest {
|
||||
assertTrue(RequestUriUtils.isFrontendRoute("", "/split-pdf"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void testIsFrontendRoute_filesRouteOwnedByFrontend() {
|
||||
// /files and /files/<folder-uuid> are FileManagerView routes - they
|
||||
// must fall through to the SPA index.html, not get blocked by the
|
||||
// backend auth filter. Regression test for direct-nav/refresh on
|
||||
// the file manager returning a 401 JSON.
|
||||
assertTrue(RequestUriUtils.isFrontendRoute("", "/files"));
|
||||
assertTrue(
|
||||
RequestUriUtils.isFrontendRoute("", "/files/3331910a-4155-4f71-8111-e38c896bc458"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void testIsFrontendRoute_pathWithExtension() {
|
||||
assertFalse(RequestUriUtils.isFrontendRoute("", "/some/file.pdf"));
|
||||
@@ -183,7 +194,7 @@ class RequestUriUtilsTest {
|
||||
|
||||
@Test
|
||||
void testIsPublicAuthEndpoint_shareRootNotPublic() {
|
||||
// Avoid matching bare "/share" or "/share/" — must have a token segment
|
||||
// Avoid matching bare "/share" or "/share/" - must have a token segment
|
||||
assertFalse(RequestUriUtils.isPublicAuthEndpoint("/share", ""));
|
||||
assertFalse(RequestUriUtils.isPublicAuthEndpoint("/share/", ""));
|
||||
}
|
||||
@@ -197,7 +208,7 @@ class RequestUriUtilsTest {
|
||||
|
||||
@Test
|
||||
void testIsPublicAuthEndpoint_shareApiStillProtected() {
|
||||
// Share-link data APIs must NOT be public — they enforce auth + access checks
|
||||
// Share-link data APIs must NOT be public - they enforce auth + access checks
|
||||
assertFalse(RequestUriUtils.isPublicAuthEndpoint("/api/v1/storage/share-links/abc123", ""));
|
||||
assertFalse(
|
||||
RequestUriUtils.isPublicAuthEndpoint(
|
||||
|
||||
+5
-1
@@ -32,6 +32,7 @@ import stirling.software.SPDF.model.json.PdfJsonTextElement;
|
||||
import stirling.software.SPDF.service.PdfJsonConversionService;
|
||||
import stirling.software.common.annotations.AutoJobPostMapping;
|
||||
import stirling.software.common.annotations.api.GeneralApi;
|
||||
import stirling.software.common.enumeration.ResourceWeight;
|
||||
import stirling.software.common.model.api.general.EditTextOperation;
|
||||
import stirling.software.common.util.ExceptionUtils;
|
||||
import stirling.software.common.util.GeneralUtils;
|
||||
@@ -75,7 +76,10 @@ public class EditTextController {
|
||||
new StringToArrayListPropertyEditor<>(EditTextOperation.class));
|
||||
}
|
||||
|
||||
@AutoJobPostMapping(consumes = "multipart/form-data", value = "/edit-text")
|
||||
@AutoJobPostMapping(
|
||||
consumes = "multipart/form-data",
|
||||
value = "/edit-text",
|
||||
resourceWeight = ResourceWeight.LARGE_WEIGHT)
|
||||
@StandardPdfResponse
|
||||
@Operation(
|
||||
summary = "Edit text in a PDF via find and replace",
|
||||
|
||||
+11
-42
@@ -40,7 +40,8 @@ public class ScalePagesController {
|
||||
private final CustomPDFDocumentFactory pdfDocumentFactory;
|
||||
private final TempFileManager tempFileManager;
|
||||
|
||||
private static PDRectangle getTargetSize(String targetPDRectangle, PDDocument sourceDocument) {
|
||||
private static PDRectangle getTargetSize(
|
||||
String targetPDRectangle, String orientation, PDDocument sourceDocument) {
|
||||
if ("KEEP".equals(targetPDRectangle)) {
|
||||
if (sourceDocument.getNumberOfPages() == 0) {
|
||||
throw ExceptionUtils.createInvalidPageSizeException("KEEP");
|
||||
@@ -57,18 +58,19 @@ public class ScalePagesController {
|
||||
}
|
||||
|
||||
Map<String, PDRectangle> sizeMap = getSizeMap();
|
||||
|
||||
if (sizeMap.containsKey(targetPDRectangle)) {
|
||||
return sizeMap.get(targetPDRectangle);
|
||||
PDRectangle base = sizeMap.get(targetPDRectangle);
|
||||
if (base == null) {
|
||||
throw ExceptionUtils.createInvalidPageSizeException(targetPDRectangle);
|
||||
}
|
||||
|
||||
throw ExceptionUtils.createInvalidPageSizeException(targetPDRectangle);
|
||||
if ("LANDSCAPE".equalsIgnoreCase(orientation)) {
|
||||
return new PDRectangle(base.getHeight(), base.getWidth());
|
||||
}
|
||||
return base;
|
||||
}
|
||||
|
||||
private static Map<String, PDRectangle> getSizeMap() {
|
||||
Map<String, PDRectangle> sizeMap = new HashMap<>();
|
||||
|
||||
// Portrait sizes (A0-A6)
|
||||
sizeMap.put("A0", PDRectangle.A0);
|
||||
sizeMap.put("A1", PDRectangle.A1);
|
||||
sizeMap.put("A2", PDRectangle.A2);
|
||||
@@ -76,42 +78,8 @@ public class ScalePagesController {
|
||||
sizeMap.put("A4", PDRectangle.A4);
|
||||
sizeMap.put("A5", PDRectangle.A5);
|
||||
sizeMap.put("A6", PDRectangle.A6);
|
||||
|
||||
// Landscape sizes (A0-A6)
|
||||
sizeMap.put(
|
||||
"A0_LANDSCAPE",
|
||||
new PDRectangle(PDRectangle.A0.getHeight(), PDRectangle.A0.getWidth()));
|
||||
sizeMap.put(
|
||||
"A1_LANDSCAPE",
|
||||
new PDRectangle(PDRectangle.A1.getHeight(), PDRectangle.A1.getWidth()));
|
||||
sizeMap.put(
|
||||
"A2_LANDSCAPE",
|
||||
new PDRectangle(PDRectangle.A2.getHeight(), PDRectangle.A2.getWidth()));
|
||||
sizeMap.put(
|
||||
"A3_LANDSCAPE",
|
||||
new PDRectangle(PDRectangle.A3.getHeight(), PDRectangle.A3.getWidth()));
|
||||
sizeMap.put(
|
||||
"A4_LANDSCAPE",
|
||||
new PDRectangle(PDRectangle.A4.getHeight(), PDRectangle.A4.getWidth()));
|
||||
sizeMap.put(
|
||||
"A5_LANDSCAPE",
|
||||
new PDRectangle(PDRectangle.A5.getHeight(), PDRectangle.A5.getWidth()));
|
||||
sizeMap.put(
|
||||
"A6_LANDSCAPE",
|
||||
new PDRectangle(PDRectangle.A6.getHeight(), PDRectangle.A6.getWidth()));
|
||||
|
||||
// Portrait US sizes
|
||||
sizeMap.put("LETTER", PDRectangle.LETTER);
|
||||
sizeMap.put("LEGAL", PDRectangle.LEGAL);
|
||||
|
||||
// Landscape US sizes
|
||||
sizeMap.put(
|
||||
"LETTER_LANDSCAPE",
|
||||
new PDRectangle(PDRectangle.LETTER.getHeight(), PDRectangle.LETTER.getWidth()));
|
||||
sizeMap.put(
|
||||
"LEGAL_LANDSCAPE",
|
||||
new PDRectangle(PDRectangle.LEGAL.getHeight(), PDRectangle.LEGAL.getWidth()));
|
||||
|
||||
return sizeMap;
|
||||
}
|
||||
|
||||
@@ -128,13 +96,14 @@ public class ScalePagesController {
|
||||
throws IOException {
|
||||
MultipartFile file = request.getFileInput();
|
||||
String targetPDRectangle = request.getPageSize();
|
||||
String orientation = request.getOrientation();
|
||||
float scaleFactor = request.getScaleFactor();
|
||||
|
||||
try (PDDocument sourceDocument = pdfDocumentFactory.load(file);
|
||||
PDDocument outputDocument =
|
||||
pdfDocumentFactory.createNewDocumentBasedOnOldDocument(sourceDocument)) {
|
||||
|
||||
PDRectangle targetSize = getTargetSize(targetPDRectangle, sourceDocument);
|
||||
PDRectangle targetSize = getTargetSize(targetPDRectangle, orientation, sourceDocument);
|
||||
|
||||
// Create LayerUtility once outside the loop for better performance
|
||||
LayerUtility layerUtility = new LayerUtility(outputDocument);
|
||||
|
||||
+16
-4
@@ -275,7 +275,10 @@ public class ConvertImgPDFController {
|
||||
GeneralUtils.generateFilename(file[0].getOriginalFilename(), "_converted.pdf"));
|
||||
}
|
||||
|
||||
@AutoJobPostMapping(consumes = MediaType.MULTIPART_FORM_DATA_VALUE, value = "/cbz/pdf")
|
||||
@AutoJobPostMapping(
|
||||
consumes = MediaType.MULTIPART_FORM_DATA_VALUE,
|
||||
value = "/cbz/pdf",
|
||||
resourceWeight = ResourceWeight.MEDIUM_WEIGHT)
|
||||
@Operation(
|
||||
summary = "Convert CBZ comic book archive to PDF",
|
||||
description =
|
||||
@@ -301,7 +304,10 @@ public class ConvertImgPDFController {
|
||||
return WebResponseUtils.pdfFileToWebResponse(pdfFile, filename);
|
||||
}
|
||||
|
||||
@AutoJobPostMapping(consumes = MediaType.MULTIPART_FORM_DATA_VALUE, value = "/pdf/cbz")
|
||||
@AutoJobPostMapping(
|
||||
consumes = MediaType.MULTIPART_FORM_DATA_VALUE,
|
||||
value = "/pdf/cbz",
|
||||
resourceWeight = ResourceWeight.LARGE_WEIGHT)
|
||||
@Operation(
|
||||
summary = "Convert PDF to CBZ comic book archive",
|
||||
description =
|
||||
@@ -324,7 +330,10 @@ public class ConvertImgPDFController {
|
||||
return WebResponseUtils.zipFileToWebResponse(cbzFile, filename);
|
||||
}
|
||||
|
||||
@AutoJobPostMapping(consumes = MediaType.MULTIPART_FORM_DATA_VALUE, value = "/cbr/pdf")
|
||||
@AutoJobPostMapping(
|
||||
consumes = MediaType.MULTIPART_FORM_DATA_VALUE,
|
||||
value = "/cbr/pdf",
|
||||
resourceWeight = ResourceWeight.MEDIUM_WEIGHT)
|
||||
@Operation(
|
||||
summary = "Convert CBR comic book archive to PDF",
|
||||
description =
|
||||
@@ -350,7 +359,10 @@ public class ConvertImgPDFController {
|
||||
return WebResponseUtils.bytesToWebResponse(pdfBytes, filename);
|
||||
}
|
||||
|
||||
@AutoJobPostMapping(consumes = MediaType.MULTIPART_FORM_DATA_VALUE, value = "/pdf/cbr")
|
||||
@AutoJobPostMapping(
|
||||
consumes = MediaType.MULTIPART_FORM_DATA_VALUE,
|
||||
value = "/pdf/cbr",
|
||||
resourceWeight = ResourceWeight.LARGE_WEIGHT)
|
||||
@Operation(
|
||||
summary = "Convert PDF to CBR comic book archive",
|
||||
description =
|
||||
|
||||
+10
-4
@@ -141,7 +141,8 @@ public class AttachmentController {
|
||||
|
||||
@AutoJobPostMapping(
|
||||
consumes = MediaType.MULTIPART_FORM_DATA_VALUE,
|
||||
value = "/extract-attachments")
|
||||
value = "/extract-attachments",
|
||||
resourceWeight = ResourceWeight.SMALL_WEIGHT)
|
||||
@Operation(
|
||||
summary = "Extract attachments from PDF",
|
||||
description =
|
||||
@@ -176,7 +177,10 @@ public class AttachmentController {
|
||||
}
|
||||
}
|
||||
|
||||
@AutoJobPostMapping(consumes = MediaType.MULTIPART_FORM_DATA_VALUE, value = "/list-attachments")
|
||||
@AutoJobPostMapping(
|
||||
consumes = MediaType.MULTIPART_FORM_DATA_VALUE,
|
||||
value = "/list-attachments",
|
||||
resourceWeight = ResourceWeight.SMALL_WEIGHT)
|
||||
@Operation(
|
||||
summary = "List attachments in PDF",
|
||||
description =
|
||||
@@ -193,7 +197,8 @@ public class AttachmentController {
|
||||
|
||||
@AutoJobPostMapping(
|
||||
consumes = MediaType.MULTIPART_FORM_DATA_VALUE,
|
||||
value = "/rename-attachment")
|
||||
value = "/rename-attachment",
|
||||
resourceWeight = ResourceWeight.SMALL_WEIGHT)
|
||||
@StandardPdfResponse
|
||||
@Operation(
|
||||
summary = "Rename attachment in PDF",
|
||||
@@ -228,7 +233,8 @@ public class AttachmentController {
|
||||
|
||||
@AutoJobPostMapping(
|
||||
consumes = MediaType.MULTIPART_FORM_DATA_VALUE,
|
||||
value = "/delete-attachment")
|
||||
value = "/delete-attachment",
|
||||
resourceWeight = ResourceWeight.SMALL_WEIGHT)
|
||||
@StandardPdfResponse
|
||||
@Operation(
|
||||
summary = "Delete attachment from PDF",
|
||||
|
||||
+42
-12
@@ -12,6 +12,8 @@ import org.springframework.web.bind.annotation.RequestParam;
|
||||
|
||||
import io.swagger.v3.oas.annotations.Hidden;
|
||||
|
||||
import jakarta.servlet.http.HttpServletRequest;
|
||||
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
|
||||
import stirling.software.SPDF.config.EndpointConfiguration;
|
||||
@@ -91,6 +93,44 @@ public class ConfigController {
|
||||
return null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve the frontend URL the client should advertise to phones / share-link recipients.
|
||||
* Priority: explicit system.frontendUrl, then the Host the user is already using to reach this
|
||||
* server (works for Docker, reverse proxies, and bare-metal LANs), then a detected site-local
|
||||
* IPv4, then empty.
|
||||
*/
|
||||
// visible for testing
|
||||
String resolveFrontendUrl(HttpServletRequest request, AppConfig appConfig) {
|
||||
String configured = applicationProperties.getSystem().getFrontendUrl();
|
||||
if (configured != null && !configured.isBlank()) {
|
||||
return configured;
|
||||
}
|
||||
if (request != null) {
|
||||
String host = request.getServerName();
|
||||
if (host != null && !host.isBlank() && !isLoopbackHost(host)) {
|
||||
String scheme = request.getScheme();
|
||||
int port = request.getServerPort();
|
||||
boolean defaultPort =
|
||||
("http".equals(scheme) && port == 80)
|
||||
|| ("https".equals(scheme) && port == 443);
|
||||
return defaultPort ? scheme + "://" + host : scheme + "://" + host + ":" + port;
|
||||
}
|
||||
}
|
||||
String localIp = GeneralUtils.getLocalNetworkIp();
|
||||
if (localIp != null) {
|
||||
String scheme = appConfig.getBackendUrl().startsWith("https") ? "https" : "http";
|
||||
return scheme + "://" + localIp + ":" + appConfig.getServerPort();
|
||||
}
|
||||
return "";
|
||||
}
|
||||
|
||||
private static boolean isLoopbackHost(String host) {
|
||||
return "localhost".equalsIgnoreCase(host)
|
||||
|| "127.0.0.1".equals(host)
|
||||
|| "::1".equals(host)
|
||||
|| "0:0:0:0:0:0:0:1".equals(host);
|
||||
}
|
||||
|
||||
/** Check if running Enterprise edition dynamically. */
|
||||
private Boolean isRunningEE() {
|
||||
// Use LicenseService for fresh license status if available
|
||||
@@ -107,7 +147,7 @@ public class ConfigController {
|
||||
}
|
||||
|
||||
@GetMapping("/app-config")
|
||||
public ResponseEntity<Map<String, Object>> getAppConfig() {
|
||||
public ResponseEntity<Map<String, Object>> getAppConfig(HttpServletRequest request) {
|
||||
Map<String, Object> configData = new HashMap<>();
|
||||
|
||||
try {
|
||||
@@ -124,17 +164,7 @@ public class ConfigController {
|
||||
configData.put("serverPort", appConfig.getServerPort());
|
||||
|
||||
String frontendUrl = applicationProperties.getSystem().getFrontendUrl();
|
||||
if ((frontendUrl == null || frontendUrl.isBlank())
|
||||
&& Boolean.parseBoolean(
|
||||
System.getProperty("STIRLING_PDF_TAURI_MODE", "false"))) {
|
||||
String localIp = GeneralUtils.getLocalNetworkIp();
|
||||
if (localIp != null) {
|
||||
String scheme =
|
||||
appConfig.getBackendUrl().startsWith("https") ? "https" : "http";
|
||||
frontendUrl = scheme + "://" + localIp + ":" + appConfig.getServerPort();
|
||||
}
|
||||
}
|
||||
configData.put("frontendUrl", frontendUrl != null ? frontendUrl : "");
|
||||
configData.put("frontendUrl", resolveFrontendUrl(request, appConfig));
|
||||
|
||||
// Add mobile scanner settings
|
||||
configData.put(
|
||||
|
||||
+7
-2
@@ -160,14 +160,19 @@ public class ReactRoutingController {
|
||||
return ResponseEntity.ok().contentType(MediaType.TEXT_HTML).body(cachedCallbackHtml);
|
||||
}
|
||||
|
||||
// `files` was historically a backend static-asset directory and was therefore
|
||||
// in the exclusion list - removing it lets /files and /files/<folder-uuid>
|
||||
// forward to the SPA index.html, which is what FileManagerView expects.
|
||||
// (Real storage endpoints live under /api/v1/storage/files, already
|
||||
// excluded by the leading `api` token in the same regex.)
|
||||
@GetMapping(
|
||||
"/{path:^(?!api|static|robots\\.txt|favicon\\.ico|manifest.*\\.json|pipeline|pdfjs|pdfjs-legacy|pdfium|vendor|fonts|images|files|css|js|assets|locales|modern-logo|classic-logo|Login|og_images|samples)[^\\.]*$}")
|
||||
"/{path:^(?!api|static|robots\\.txt|favicon\\.ico|manifest.*\\.json|pipeline|pdfjs|pdfjs-legacy|pdfium|vendor|fonts|images|css|js|assets|locales|modern-logo|classic-logo|Login|og_images|samples)[^\\.]*$}")
|
||||
public ResponseEntity<String> forwardRootPaths(HttpServletRequest request) throws IOException {
|
||||
return serveIndexHtml(request);
|
||||
}
|
||||
|
||||
@GetMapping(
|
||||
"/{path:^(?!api|static|pipeline|pdfjs|pdfjs-legacy|pdfium|vendor|fonts|images|files|css|js|assets|locales|modern-logo|classic-logo|Login|og_images|samples)[^\\.]*}/{subpath:^(?!.*\\.).*$}")
|
||||
"/{path:^(?!api|static|pipeline|pdfjs|pdfjs-legacy|pdfium|vendor|fonts|images|css|js|assets|locales|modern-logo|classic-logo|Login|og_images|samples)[^\\.]*}/{subpath:^(?!.*\\.).*$}")
|
||||
public ResponseEntity<String> forwardNestedPaths(HttpServletRequest request)
|
||||
throws IOException {
|
||||
return serveIndexHtml(request);
|
||||
|
||||
+40
-2
@@ -22,6 +22,7 @@ import org.springframework.web.bind.annotation.ExceptionHandler;
|
||||
import org.springframework.web.bind.annotation.RestControllerAdvice;
|
||||
import org.springframework.web.multipart.MaxUploadSizeExceededException;
|
||||
import org.springframework.web.multipart.support.MissingServletRequestPartException;
|
||||
import org.springframework.web.server.ResponseStatusException;
|
||||
import org.springframework.web.servlet.NoHandlerFoundException;
|
||||
|
||||
import jakarta.servlet.http.HttpServletRequest;
|
||||
@@ -196,12 +197,12 @@ public class GlobalExceptionHandler {
|
||||
/**
|
||||
* Checks whether the given IOException indicates that the client disconnected before the
|
||||
* response could be written (broken pipe, connection reset, etc.). When this happens there is
|
||||
* no point in serialising a {@link ProblemDetail} body because the socket is already closed —
|
||||
* no point in serialising a {@link ProblemDetail} body because the socket is already closed -
|
||||
* and attempting to do so may trigger a secondary {@code HttpMessageNotWritableException} if
|
||||
* the response Content-Type was already committed as a non-JSON type (e.g. image/png).
|
||||
*/
|
||||
private static boolean isClientDisconnectException(IOException ex) {
|
||||
// Walk the causal chain — Jetty/Tomcat may wrap the low-level SocketException
|
||||
// Walk the causal chain - Jetty/Tomcat may wrap the low-level SocketException
|
||||
Throwable current = ex;
|
||||
while (current != null) {
|
||||
String msg = current.getMessage();
|
||||
@@ -1040,6 +1041,43 @@ public class GlobalExceptionHandler {
|
||||
* @param request the HTTP servlet request
|
||||
* @return ProblemDetail with appropriate HTTP status
|
||||
*/
|
||||
/**
|
||||
* Handle ResponseStatusException explicitly so its embedded HTTP status reaches the client
|
||||
* instead of being swallowed by the {@code RuntimeException} catch-all (which would downgrade
|
||||
* every controller-thrown 400/404/409 to a generic 500). Folder/file storage controllers and
|
||||
* any other code that throws {@code ResponseStatusException} relies on this handler taking
|
||||
* precedence.
|
||||
*/
|
||||
@ExceptionHandler(ResponseStatusException.class)
|
||||
public ResponseEntity<ProblemDetail> handleResponseStatusException(
|
||||
ResponseStatusException ex, HttpServletRequest request) {
|
||||
HttpStatus status =
|
||||
HttpStatus.resolve(ex.getStatusCode().value()) != null
|
||||
? HttpStatus.valueOf(ex.getStatusCode().value())
|
||||
: HttpStatus.INTERNAL_SERVER_ERROR;
|
||||
String reason = ex.getReason() != null ? ex.getReason() : status.getReasonPhrase();
|
||||
ProblemDetail problemDetail = createBaseProblemDetail(status, reason, request);
|
||||
problemDetail.setType(URI.create("/errors/" + status.value()));
|
||||
problemDetail.setTitle(status.getReasonPhrase());
|
||||
problemDetail.setProperty("title", status.getReasonPhrase());
|
||||
// 5xx is operator-relevant; 4xx is a normal client-rejection - log at the right level.
|
||||
if (status.is5xxServerError()) {
|
||||
log.error(
|
||||
"ResponseStatusException {} at {}: {}",
|
||||
status.value(),
|
||||
request.getRequestURI(),
|
||||
reason,
|
||||
ex);
|
||||
} else {
|
||||
log.debug(
|
||||
"ResponseStatusException {} at {}: {}",
|
||||
status.value(),
|
||||
request.getRequestURI(),
|
||||
reason);
|
||||
}
|
||||
return ResponseEntity.status(status).contentType(PROBLEM_JSON).body(problemDetail);
|
||||
}
|
||||
|
||||
@ExceptionHandler(RuntimeException.class)
|
||||
public ResponseEntity<ProblemDetail> handleRuntimeException(
|
||||
RuntimeException ex, HttpServletRequest request) {
|
||||
|
||||
@@ -18,4 +18,11 @@ public class PDFWithPageSize extends PDFFile {
|
||||
requiredMode = Schema.RequiredMode.REQUIRED,
|
||||
allowableValues = {"A0", "A1", "A2", "A3", "A4", "A5", "A6", "LETTER", "LEGAL", "KEEP"})
|
||||
private String pageSize;
|
||||
|
||||
@Schema(
|
||||
description =
|
||||
"Orientation to apply to the target page size. Ignored when pageSize is KEEP.",
|
||||
defaultValue = "PORTRAIT",
|
||||
allowableValues = {"PORTRAIT", "LANDSCAPE"})
|
||||
private String orientation = "PORTRAIT";
|
||||
}
|
||||
|
||||
+5
-1
@@ -13,6 +13,7 @@ import lombok.RequiredArgsConstructor;
|
||||
import stirling.software.SPDF.config.swagger.MarkdownConversionResponse;
|
||||
import stirling.software.common.annotations.AutoJobPostMapping;
|
||||
import stirling.software.common.annotations.api.ConvertApi;
|
||||
import stirling.software.common.enumeration.ResourceWeight;
|
||||
import stirling.software.common.model.api.PDFFile;
|
||||
import stirling.software.common.util.PDFToFile;
|
||||
import stirling.software.common.util.TempFileManager;
|
||||
@@ -23,7 +24,10 @@ public class ConvertPDFToMarkdown {
|
||||
|
||||
private final TempFileManager tempFileManager;
|
||||
|
||||
@AutoJobPostMapping(consumes = MediaType.MULTIPART_FORM_DATA_VALUE, value = "/pdf/markdown")
|
||||
@AutoJobPostMapping(
|
||||
consumes = MediaType.MULTIPART_FORM_DATA_VALUE,
|
||||
value = "/pdf/markdown",
|
||||
resourceWeight = ResourceWeight.MEDIUM_WEIGHT)
|
||||
@MarkdownConversionResponse
|
||||
@Operation(
|
||||
summary = "Convert PDF to Markdown",
|
||||
|
||||
@@ -20,6 +20,7 @@ security:
|
||||
password: "" # initial password for the first login
|
||||
oauth2:
|
||||
enabled: false # set to 'true' to enable login (Note: enableLogin must also be 'true' for this to work)
|
||||
debugLogging: false # set to 'true' to log full ID token and UserInfo claims during OAuth2/OIDC login. Use this to diagnose claim issues (e.g. "Attribute value for 'email' cannot be null" with ADFS). WARNING: writes PII (sub, email, name) to logs; disable after troubleshooting.
|
||||
client:
|
||||
keycloak:
|
||||
issuer: "" # URL of the Keycloak realm's OpenID Connect Discovery endpoint
|
||||
@@ -93,8 +94,8 @@ premium:
|
||||
key: 00000000-0000-0000-0000-000000000000
|
||||
enabled: false # Enable license key checks for pro/enterprise features
|
||||
proFeatures:
|
||||
SSOAutoLogin: false
|
||||
CustomMetadata:
|
||||
ssoAutoLogin: false
|
||||
customMetadata:
|
||||
autoUpdateMetadata: false
|
||||
author: username
|
||||
creator: Stirling-PDF
|
||||
@@ -245,6 +246,39 @@ storage:
|
||||
provider: local # storage provider: 'local' for filesystem storage, 'database' for DB-backed storage
|
||||
local:
|
||||
basePath: './storage' # base directory for stored files
|
||||
# ====================================================================================
|
||||
# S3-COMPATIBLE OBJECT STORAGE - PRO / ENTERPRISE LICENSE REQUIRED
|
||||
# storage.provider=s3, storage.provider=database, and cluster.artifactStore=s3 all
|
||||
# require a valid Pro or Enterprise license.
|
||||
# ====================================================================================
|
||||
# Used when provider=s3 (persistent user uploads) and/or cluster.artifactStore=s3
|
||||
# (transient cluster artifacts). The two consumers share this block.
|
||||
# Vendor cheat sheet (set the highlighted flags to taste):
|
||||
# AWS S3 -> endpoint='' region='<your-region>' pathStyleAccess=false
|
||||
# Cloudflare R2 -> endpoint='https://<acct>.r2.cloudflarestorage.com' region='auto'
|
||||
# pathStyleAccess=false; if uploads fail with 'unsupported header
|
||||
# x-amz-checksum-*' set requestChecksumCalculation=WHEN_REQUIRED
|
||||
# Supabase Storage -> endpoint='https://<project>.supabase.co/storage/v1/s3'
|
||||
# region='<project-region>' pathStyleAccess=true
|
||||
# (filenames with non-ASCII display fine - the storage key is opaque)
|
||||
# MinIO (in-cluster) -> endpoint='http://minio:9000' region='us-east-1'
|
||||
# pathStyleAccess=true allowPrivateEndpoints=true
|
||||
# Backblaze B2 -> endpoint='https://s3.<region>.backblazeb2.com'
|
||||
# If on a B2 deployment older than July-2025 and uploads return
|
||||
# 'Unsupported header x-amz-checksum-crc32', set
|
||||
# requestChecksumCalculation=WHEN_REQUIRED
|
||||
# DigitalOcean Spaces -> endpoint='https://<region>.digitaloceanspaces.com'
|
||||
# Note: 5GB per-object cap (regardless of multipart)
|
||||
s3:
|
||||
endpoint: "" # blank = use AWS regional default; otherwise full URL incl. https://
|
||||
bucket: "" # required when provider=s3 or cluster.artifactStore=s3
|
||||
region: us-east-1
|
||||
accessKey: "" # blank = fall back to AWS DefaultCredentialsProvider (env / profile / IMDS)
|
||||
secretKey: ""
|
||||
pathStyleAccess: false # true for MinIO and Supabase; false for AWS/R2/most CDNs
|
||||
allowPrivateEndpoints: false # true required when endpoint resolves to a private/loopback IP (e.g. in-cluster MinIO). SSRF guard - leave false for any internet-facing vendor.
|
||||
requestChecksumCalculation: WHEN_SUPPORTED # WHEN_SUPPORTED|WHEN_REQUIRED|DISABLED. Set WHEN_REQUIRED if your vendor rejects auto-added x-amz-checksum-* headers (older Backblaze B2, some R2 corner cases).
|
||||
responseChecksumValidation: WHEN_SUPPORTED # WHEN_SUPPORTED|WHEN_REQUIRED|DISABLED. Set WHEN_REQUIRED if you see false-positive checksum-mismatch errors on GET from a vendor that never returns checksum headers.
|
||||
quotas:
|
||||
maxStorageMbPerUser: -1 # Max storage per user in MB; -1 disables per-user cap
|
||||
maxStorageMbTotal: -1 # Max storage across all users in MB; -1 disables total cap
|
||||
@@ -330,6 +364,24 @@ aiEngine:
|
||||
url: http://localhost:5001 # URL of the Python AI engine
|
||||
timeoutSeconds: 120 # Timeout in seconds for AI engine requests
|
||||
|
||||
# Cluster configuration. NOT YET ENABLED - scaffolding for later work. Leave at defaults.
|
||||
cluster:
|
||||
enabled: false # Master switch. 'false' (default) wires the in-process backplane and skips all cluster checks. Single-instance installs do not need to change anything here.
|
||||
backplane: inprocess # Backplane implementation: 'inprocess' (single JVM only) or 'valkey' (multi-node via Valkey/Redis)
|
||||
artifactStore: local # Transient cluster job-artifact backend: 'local' (per-node disk; single-node only) or 's3' (shared object store; required for multi-node). Distinct from 'storage.provider' which controls persistent user uploads - when both are 's3' they share the storage.s3.* credentials block. Multi-node deployments MUST set this to 's3'.
|
||||
s3:
|
||||
keyPrefix: transient/ # Bucket key prefix used by the cluster artifact store when artifactStore=s3. Trailing slash recommended. Lets a single bucket host both persistent uploads (storage.s3.*) and transient job artifacts under separate prefixes.
|
||||
valkey:
|
||||
url: "" # Valkey/Redis URL, e.g. 'redis://valkey:6379' or 'rediss://...' for TLS. Required when enabled=true and backplane=valkey.
|
||||
tls:
|
||||
skipCertVerification: false # set to 'true' to skip TLS certificate verification on Valkey connections (dev/test only)
|
||||
node:
|
||||
id: "" # Optional explicit node id. Blank = auto-generated UUID at startup.
|
||||
role: both # 'web' (serves HTTP), 'worker' (runs jobs), or 'both' (default)
|
||||
internalAddress: "" # host:port advertised in the instance registry for peer-to-peer cluster traffic. Blank = derived at startup.
|
||||
scheme: http # 'http' or 'https' - scheme peers use to call this node's /internal/cluster/** endpoints
|
||||
heartbeatIntervalMs: 5000 # Heartbeat publish interval for the instance registry (ms)
|
||||
|
||||
pdfEditor:
|
||||
fallback-font: classpath:/static/fonts/NotoSans-Regular.ttf # Override to point at a custom fallback font
|
||||
cache:
|
||||
|
||||
+132
@@ -0,0 +1,132 @@
|
||||
package stirling.software.SPDF.config;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.lang.reflect.Method;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.springframework.core.io.Resource;
|
||||
import org.springframework.core.io.support.PathMatchingResourcePatternResolver;
|
||||
import org.springframework.core.io.support.ResourcePatternResolver;
|
||||
import org.springframework.core.type.classreading.CachingMetadataReaderFactory;
|
||||
import org.springframework.core.type.classreading.MetadataReader;
|
||||
import org.springframework.core.type.classreading.MetadataReaderFactory;
|
||||
import org.springframework.core.type.filter.TypeFilter;
|
||||
|
||||
import stirling.software.common.annotations.AutoJobPostMapping;
|
||||
|
||||
/**
|
||||
* Build-time guardrail: every {@link AutoJobPostMapping} method must declare an explicit {@code
|
||||
* resourceWeight}.
|
||||
*
|
||||
* <p>The credits interceptor multiplies {@code resourceWeight} into the per-call charge. An
|
||||
* endpoint that falls through to the annotation default produces a charge derived from a value
|
||||
* nobody chose — silently under- or over-billing depending on the endpoint's true cost. Forcing
|
||||
* each method to pick a value from {@link stirling.software.common.enumeration.ResourceWeight}
|
||||
* keeps the choice deliberate.
|
||||
*
|
||||
* <p>The annotation's default is {@link Integer#MIN_VALUE} (a sentinel). Runtime readers clamp the
|
||||
* value into {@code [1, 100]}, so a missed declaration can't crash production — this test is the
|
||||
* contract, the clamp is the safety net.
|
||||
*
|
||||
* <p>Lives in {@code :stirling-pdf} (core) because that's the module whose compile classpath
|
||||
* transitively sees every other module's controllers ({@code :common}, {@code :proprietary}, and
|
||||
* {@code :saas} when enabled).
|
||||
*/
|
||||
class AutoJobPostMappingWeightTest {
|
||||
|
||||
private static final String SCAN_BASE_PACKAGE = "stirling.software";
|
||||
|
||||
@Test
|
||||
void everyAutoJobPostMappingDeclaresExplicitResourceWeight() throws Exception {
|
||||
List<String> offenders = findOffendingMethods();
|
||||
|
||||
assertTrue(
|
||||
offenders.isEmpty(),
|
||||
() ->
|
||||
"The following @AutoJobPostMapping methods do not declare an explicit"
|
||||
+ " resourceWeight. Pick a value from"
|
||||
+ " stirling.software.common.enumeration.ResourceWeight (SMALL,"
|
||||
+ " MEDIUM, LARGE, XLARGE) and add it to the annotation:\n - "
|
||||
+ String.join("\n - ", offenders));
|
||||
}
|
||||
|
||||
private List<String> findOffendingMethods() throws IOException, ClassNotFoundException {
|
||||
List<String> offenders = new ArrayList<>();
|
||||
for (Class<?> candidate : scanForCandidateClasses()) {
|
||||
for (Method method : candidate.getDeclaredMethods()) {
|
||||
AutoJobPostMapping annotation = method.getAnnotation(AutoJobPostMapping.class);
|
||||
if (annotation == null) {
|
||||
continue;
|
||||
}
|
||||
if (annotation.resourceWeight() == Integer.MIN_VALUE) {
|
||||
offenders.add(candidate.getName() + "#" + method.getName());
|
||||
}
|
||||
}
|
||||
}
|
||||
return offenders;
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns every class under {@link #SCAN_BASE_PACKAGE} that has an @AutoJobPostMapping method.
|
||||
*/
|
||||
private List<Class<?>> scanForCandidateClasses() throws IOException, ClassNotFoundException {
|
||||
ResourcePatternResolver resolver = new PathMatchingResourcePatternResolver();
|
||||
MetadataReaderFactory metadataReaderFactory = new CachingMetadataReaderFactory(resolver);
|
||||
|
||||
String pattern = "classpath*:" + SCAN_BASE_PACKAGE.replace('.', '/') + "/**/*.class";
|
||||
Resource[] resources = resolver.getResources(pattern);
|
||||
|
||||
// Pre-filter by reading annotation metadata from the class file so we don't have to load
|
||||
// every class on the test classpath just to find the few that are annotated.
|
||||
TypeFilter mentionsAutoJobPostMapping =
|
||||
(reader, factory) ->
|
||||
reader.getAnnotationMetadata()
|
||||
.getAnnotatedMethods(AutoJobPostMapping.class.getName())
|
||||
.size()
|
||||
> 0;
|
||||
|
||||
List<Class<?>> matches = new ArrayList<>();
|
||||
for (Resource resource : resources) {
|
||||
if (!resource.isReadable()) {
|
||||
continue;
|
||||
}
|
||||
MetadataReader reader = metadataReaderFactory.getMetadataReader(resource);
|
||||
if (!mentionsAutoJobPostMapping.match(reader, metadataReaderFactory)) {
|
||||
continue;
|
||||
}
|
||||
matches.add(Class.forName(reader.getClassMetadata().getClassName()));
|
||||
}
|
||||
return matches;
|
||||
}
|
||||
|
||||
/**
|
||||
* Sanity check that the classpath scan returns non-empty; otherwise the main test passes
|
||||
* vacuously.
|
||||
*/
|
||||
@Test
|
||||
void scannerFindsAtLeastOneAutoJobPostMapping() throws Exception {
|
||||
long count =
|
||||
scanForCandidateClasses().stream()
|
||||
.flatMap(c -> java.util.Arrays.stream(c.getDeclaredMethods()))
|
||||
.filter(m -> m.isAnnotationPresent(AutoJobPostMapping.class))
|
||||
.count();
|
||||
|
||||
assertTrue(
|
||||
count > 10,
|
||||
() ->
|
||||
"Expected the classpath scan to find many @AutoJobPostMapping methods but"
|
||||
+ " found only "
|
||||
+ count
|
||||
+ ". Scanner regression?");
|
||||
}
|
||||
|
||||
@SuppressWarnings("unused")
|
||||
private static String describeCandidates(List<Class<?>> candidates) {
|
||||
return candidates.stream().map(Class::getName).collect(Collectors.joining(", "));
|
||||
}
|
||||
}
|
||||
+2
-1
@@ -237,7 +237,8 @@ class ScalePagesControllerTest {
|
||||
|
||||
ScalePagesRequest request = new ScalePagesRequest();
|
||||
request.setFileInput(file);
|
||||
request.setPageSize("A4_LANDSCAPE");
|
||||
request.setPageSize("A4");
|
||||
request.setOrientation("LANDSCAPE");
|
||||
request.setScaleFactor(1.0f);
|
||||
|
||||
setupFactory();
|
||||
|
||||
+71
@@ -15,10 +15,14 @@ import org.springframework.context.ApplicationContext;
|
||||
import org.springframework.http.HttpStatus;
|
||||
import org.springframework.http.ResponseEntity;
|
||||
|
||||
import jakarta.servlet.http.HttpServletRequest;
|
||||
|
||||
import stirling.software.SPDF.config.EndpointConfiguration;
|
||||
import stirling.software.SPDF.config.EndpointConfiguration.DisableReason;
|
||||
import stirling.software.SPDF.config.EndpointConfiguration.EndpointAvailability;
|
||||
import stirling.software.common.configuration.AppConfig;
|
||||
import stirling.software.common.model.ApplicationProperties;
|
||||
import stirling.software.common.model.ApplicationProperties.System;
|
||||
import stirling.software.common.service.LicenseServiceInterface;
|
||||
import stirling.software.common.service.ServerCertificateServiceInterface;
|
||||
import stirling.software.common.service.UserServiceInterface;
|
||||
@@ -173,4 +177,71 @@ class ConfigControllerTest {
|
||||
assertEquals(HttpStatus.OK, response.getStatusCode());
|
||||
verify(endpointConfiguration).getAllEndpoints();
|
||||
}
|
||||
|
||||
@Test
|
||||
void resolveFrontendUrl_prefersExplicitConfiguredValue() {
|
||||
System sys = mock(System.class);
|
||||
when(applicationProperties.getSystem()).thenReturn(sys);
|
||||
when(sys.getFrontendUrl()).thenReturn("https://pdf.example.com");
|
||||
|
||||
// Request would say something else, but configured wins.
|
||||
HttpServletRequest req = mock(HttpServletRequest.class);
|
||||
AppConfig appConfig = mock(AppConfig.class);
|
||||
|
||||
assertEquals(
|
||||
"https://pdf.example.com", configController.resolveFrontendUrl(req, appConfig));
|
||||
}
|
||||
|
||||
@Test
|
||||
void resolveFrontendUrl_usesRequestHostWhenNotConfigured() {
|
||||
System sys = mock(System.class);
|
||||
when(applicationProperties.getSystem()).thenReturn(sys);
|
||||
when(sys.getFrontendUrl()).thenReturn(null);
|
||||
|
||||
HttpServletRequest req = mock(HttpServletRequest.class);
|
||||
when(req.getServerName()).thenReturn("192.168.1.100");
|
||||
when(req.getScheme()).thenReturn("http");
|
||||
when(req.getServerPort()).thenReturn(8080);
|
||||
|
||||
assertEquals(
|
||||
"http://192.168.1.100:8080",
|
||||
configController.resolveFrontendUrl(req, mock(AppConfig.class)));
|
||||
}
|
||||
|
||||
@Test
|
||||
void resolveFrontendUrl_elidesDefaultHttpsPort() {
|
||||
System sys = mock(System.class);
|
||||
when(applicationProperties.getSystem()).thenReturn(sys);
|
||||
when(sys.getFrontendUrl()).thenReturn("");
|
||||
|
||||
HttpServletRequest req = mock(HttpServletRequest.class);
|
||||
when(req.getServerName()).thenReturn("pdf.example.com");
|
||||
when(req.getScheme()).thenReturn("https");
|
||||
when(req.getServerPort()).thenReturn(443);
|
||||
|
||||
assertEquals(
|
||||
"https://pdf.example.com",
|
||||
configController.resolveFrontendUrl(req, mock(AppConfig.class)));
|
||||
}
|
||||
|
||||
@Test
|
||||
void resolveFrontendUrl_fallsThroughOnLoopbackHost() {
|
||||
System sys = mock(System.class);
|
||||
when(applicationProperties.getSystem()).thenReturn(sys);
|
||||
when(sys.getFrontendUrl()).thenReturn(null);
|
||||
|
||||
HttpServletRequest req = mock(HttpServletRequest.class);
|
||||
when(req.getServerName()).thenReturn("localhost");
|
||||
|
||||
AppConfig appConfig = mock(AppConfig.class);
|
||||
when(appConfig.getBackendUrl()).thenReturn("http://localhost:8080");
|
||||
when(appConfig.getServerPort()).thenReturn("8080");
|
||||
|
||||
// Detected IP (if any) wins over loopback request host. We can't assert the
|
||||
// exact value (depends on the host running the test) but we can assert it
|
||||
// never returns "localhost".
|
||||
String result = configController.resolveFrontendUrl(req, appConfig);
|
||||
assertNotNull(result);
|
||||
assertFalse(result.contains("localhost"));
|
||||
}
|
||||
}
|
||||
|
||||
+93
@@ -0,0 +1,93 @@
|
||||
package stirling.software.common.configuration;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
import static org.mockito.Mockito.mockStatic;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
import org.mockito.MockedStatic;
|
||||
|
||||
import stirling.software.common.util.GeneralUtils;
|
||||
import stirling.software.common.util.YamlHelper;
|
||||
|
||||
/**
|
||||
* End-to-end check of the container-restart path. {@link ConfigInitializer#ensureConfigExists()} is
|
||||
* what runs on every startup, merging the on-disk settings.yml with the bundled
|
||||
* settings.yml.template. These tests exercise it against the real template on the classpath to
|
||||
* prove admin-saved proFeatures values survive a restart - the bug behind "the SSO auto-login
|
||||
* button resets every time the container resets".
|
||||
*/
|
||||
class ConfigInitializerRestartTest {
|
||||
|
||||
private static String read(Path settings, String... keyPath) throws IOException {
|
||||
return String.valueOf(new YamlHelper(settings).getValueByExactKeyPath(keyPath));
|
||||
}
|
||||
|
||||
@Test
|
||||
void ssoAutoLoginAndCustomMetadata_persistAcrossRestart(@TempDir Path tmp) throws Exception {
|
||||
Path settings = tmp.resolve("settings.yml");
|
||||
Path custom = tmp.resolve("custom_settings.yml");
|
||||
|
||||
try (MockedStatic<InstallationPathConfig> paths =
|
||||
mockStatic(InstallationPathConfig.class)) {
|
||||
paths.when(InstallationPathConfig::getSettingsPath).thenReturn(settings.toString());
|
||||
paths.when(InstallationPathConfig::getCustomSettingsPath).thenReturn(custom.toString());
|
||||
|
||||
ConfigInitializer init = new ConfigInitializer();
|
||||
|
||||
// First boot: settings.yml created from the bundled template (camelCase, default off).
|
||||
init.ensureConfigExists();
|
||||
assertEquals("false", read(settings, "premium", "proFeatures", "ssoAutoLogin"));
|
||||
|
||||
// Admin enables SSO auto-login and edits custom metadata via the exact save path the
|
||||
// admin settings controller uses.
|
||||
GeneralUtils.saveKeyToSettings("premium.proFeatures.ssoAutoLogin", true);
|
||||
GeneralUtils.saveKeyToSettings("premium.proFeatures.customMetadata.author", "acme");
|
||||
|
||||
// Container restart: ensureConfigExists merges the saved file with the template again.
|
||||
init.ensureConfigExists();
|
||||
|
||||
assertEquals("true", read(settings, "premium", "proFeatures", "ssoAutoLogin"));
|
||||
assertEquals(
|
||||
"acme", read(settings, "premium", "proFeatures", "customMetadata", "author"));
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void legacyPascalCaseConfig_isMigratedAndPreservedOnRestart(@TempDir Path tmp)
|
||||
throws Exception {
|
||||
Path settings = tmp.resolve("settings.yml");
|
||||
Path custom = tmp.resolve("custom_settings.yml");
|
||||
|
||||
try (MockedStatic<InstallationPathConfig> paths =
|
||||
mockStatic(InstallationPathConfig.class)) {
|
||||
paths.when(InstallationPathConfig::getSettingsPath).thenReturn(settings.toString());
|
||||
paths.when(InstallationPathConfig::getCustomSettingsPath).thenReturn(custom.toString());
|
||||
|
||||
ConfigInitializer init = new ConfigInitializer();
|
||||
|
||||
// Seed a full settings.yml as an OLD install would have written it: PascalCase keys
|
||||
// with
|
||||
// SSO auto-login enabled.
|
||||
init.ensureConfigExists();
|
||||
String legacy =
|
||||
Files.readString(settings)
|
||||
.replace("ssoAutoLogin: false", "SSOAutoLogin: true")
|
||||
.replace("customMetadata:", "CustomMetadata:");
|
||||
Files.writeString(settings, legacy);
|
||||
|
||||
// Upgrade restart.
|
||||
init.ensureConfigExists();
|
||||
|
||||
// Value carried forward onto the new camelCase key; the legacy PascalCase key is gone.
|
||||
assertEquals("true", read(settings, "premium", "proFeatures", "ssoAutoLogin"));
|
||||
assertNull(
|
||||
new YamlHelper(settings)
|
||||
.getValueByExactKeyPath("premium", "proFeatures", "SSOAutoLogin"));
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -123,6 +123,8 @@ SwaggerDoc.json
|
||||
*.tar.gz
|
||||
*.rar
|
||||
*.db
|
||||
# Whitelist the H2 fixtures that feed the version-migration CI smoke test.
|
||||
!src/test/resources/db-migration-fixtures/*.mv.db
|
||||
/build
|
||||
/app/proprietary/build/
|
||||
|
||||
|
||||
@@ -5,6 +5,8 @@ repositories {
|
||||
|
||||
ext {
|
||||
jwtVersion = '0.13.0'
|
||||
awsSdkVersion = '2.44.12'
|
||||
testcontainersMinioVersion = '1.21.4'
|
||||
}
|
||||
|
||||
bootRun {
|
||||
@@ -71,6 +73,13 @@ dependencies {
|
||||
implementation('com.coveo:saml-client:5.0.0') {
|
||||
exclude group: 'org.opensaml', module: 'opensaml-core'
|
||||
}
|
||||
|
||||
implementation "software.amazon.awssdk:s3:$awsSdkVersion"
|
||||
implementation "software.amazon.awssdk:url-connection-client:$awsSdkVersion"
|
||||
|
||||
testImplementation "org.testcontainers:minio:$testcontainersMinioVersion"
|
||||
testImplementation "org.testcontainers:junit-jupiter:$testcontainersMinioVersion"
|
||||
testImplementation "org.testcontainers:localstack:$testcontainersMinioVersion"
|
||||
}
|
||||
|
||||
tasks.register('prepareKotlinBuildScriptModel') {}
|
||||
|
||||
+200
@@ -0,0 +1,200 @@
|
||||
package stirling.software.proprietary.cluster.s3;
|
||||
|
||||
import java.net.InetAddress;
|
||||
import java.net.URI;
|
||||
import java.net.URISyntaxException;
|
||||
import java.net.UnknownHostException;
|
||||
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
|
||||
import stirling.software.common.model.ApplicationProperties;
|
||||
|
||||
import software.amazon.awssdk.auth.credentials.AwsBasicCredentials;
|
||||
import software.amazon.awssdk.auth.credentials.DefaultCredentialsProvider;
|
||||
import software.amazon.awssdk.auth.credentials.StaticCredentialsProvider;
|
||||
import software.amazon.awssdk.core.checksums.RequestChecksumCalculation;
|
||||
import software.amazon.awssdk.core.checksums.ResponseChecksumValidation;
|
||||
import software.amazon.awssdk.http.urlconnection.UrlConnectionHttpClient;
|
||||
import software.amazon.awssdk.regions.Region;
|
||||
import software.amazon.awssdk.services.s3.S3Client;
|
||||
import software.amazon.awssdk.services.s3.S3ClientBuilder;
|
||||
import software.amazon.awssdk.services.s3.S3Configuration;
|
||||
import software.amazon.awssdk.services.s3.presigner.S3Presigner;
|
||||
|
||||
/**
|
||||
* Shared factory for {@link S3Client} and {@link S3Presigner} instances used by both {@code
|
||||
* S3StorageProvider} and {@code S3FileStore}, so endpoint/region/credentials wiring lives in
|
||||
* exactly one place.
|
||||
*/
|
||||
@Slf4j
|
||||
public final class S3Clients {
|
||||
|
||||
private S3Clients() {}
|
||||
|
||||
/** Paired client and presigner with coordinated lifecycle. */
|
||||
public record Bundle(S3Client client, S3Presigner presigner) implements AutoCloseable {
|
||||
@Override
|
||||
public void close() {
|
||||
try {
|
||||
presigner.close();
|
||||
} catch (Exception e) {
|
||||
log.warn("Error closing S3 presigner", e);
|
||||
}
|
||||
try {
|
||||
client.close();
|
||||
} catch (Exception e) {
|
||||
log.warn("Error closing S3 client", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Build a client+presigner pair from the shared S3 config block. */
|
||||
public static Bundle build(ApplicationProperties.Storage.S3 cfg, String usage) {
|
||||
if (cfg == null) {
|
||||
throw new IllegalStateException(
|
||||
usage + " requires storage.s3.* configuration to be set");
|
||||
}
|
||||
if (cfg.getBucket() == null || cfg.getBucket().isBlank()) {
|
||||
throw new IllegalStateException(usage + " requires storage.s3.bucket to be set");
|
||||
}
|
||||
String region =
|
||||
cfg.getRegion() == null || cfg.getRegion().isBlank()
|
||||
? "us-east-1"
|
||||
: cfg.getRegion();
|
||||
|
||||
S3Configuration s3Configuration =
|
||||
S3Configuration.builder().pathStyleAccessEnabled(cfg.isPathStyleAccess()).build();
|
||||
|
||||
RequestChecksumCalculation requestChecksum =
|
||||
parseRequestChecksum(cfg.getRequestChecksumCalculation());
|
||||
ResponseChecksumValidation responseChecksum =
|
||||
parseResponseChecksum(cfg.getResponseChecksumValidation());
|
||||
|
||||
S3ClientBuilder clientBuilder =
|
||||
S3Client.builder()
|
||||
.httpClient(UrlConnectionHttpClient.create())
|
||||
.region(Region.of(region))
|
||||
.serviceConfiguration(s3Configuration)
|
||||
.requestChecksumCalculation(requestChecksum)
|
||||
.responseChecksumValidation(responseChecksum);
|
||||
|
||||
S3Presigner.Builder presignerBuilder =
|
||||
S3Presigner.builder()
|
||||
.region(Region.of(region))
|
||||
.serviceConfiguration(s3Configuration);
|
||||
|
||||
if (cfg.getEndpoint() != null && !cfg.getEndpoint().isBlank()) {
|
||||
URI endpoint;
|
||||
try {
|
||||
endpoint = new URI(cfg.getEndpoint());
|
||||
} catch (URISyntaxException e) {
|
||||
throw new IllegalStateException(
|
||||
"Invalid storage.s3.endpoint: " + cfg.getEndpoint(), e);
|
||||
}
|
||||
validateEndpointHost(endpoint, cfg.isAllowPrivateEndpoints());
|
||||
clientBuilder.endpointOverride(endpoint);
|
||||
presignerBuilder.endpointOverride(endpoint);
|
||||
}
|
||||
|
||||
boolean hasStaticCreds =
|
||||
cfg.getAccessKey() != null
|
||||
&& !cfg.getAccessKey().isBlank()
|
||||
&& cfg.getSecretKey() != null
|
||||
&& !cfg.getSecretKey().isBlank();
|
||||
if (hasStaticCreds) {
|
||||
AwsBasicCredentials credentials =
|
||||
AwsBasicCredentials.create(cfg.getAccessKey(), cfg.getSecretKey());
|
||||
StaticCredentialsProvider provider = StaticCredentialsProvider.create(credentials);
|
||||
clientBuilder.credentialsProvider(provider);
|
||||
presignerBuilder.credentialsProvider(provider);
|
||||
} else {
|
||||
clientBuilder.credentialsProvider(DefaultCredentialsProvider.create());
|
||||
presignerBuilder.credentialsProvider(DefaultCredentialsProvider.create());
|
||||
}
|
||||
|
||||
log.debug(
|
||||
"Configured S3 {}: bucket={}, region={}, endpoint={}, pathStyle={}",
|
||||
usage,
|
||||
cfg.getBucket(),
|
||||
region,
|
||||
cfg.getEndpoint() == null || cfg.getEndpoint().isBlank()
|
||||
? "<aws-default>"
|
||||
: cfg.getEndpoint(),
|
||||
cfg.isPathStyleAccess());
|
||||
|
||||
return new Bundle(clientBuilder.build(), presignerBuilder.build());
|
||||
}
|
||||
|
||||
/**
|
||||
* Block SSRF via the S3 endpoint setting. An admin who can edit config could otherwise point
|
||||
* the SDK at the cloud metadata service (e.g. {@code http://169.254.169.254/}) and exfiltrate
|
||||
* instance-role credentials. Reject any endpoint whose host resolves to a loopback, link-local,
|
||||
* or RFC1918 private address unless the operator has explicitly opted in via {@code
|
||||
* storage.s3.allow-private-endpoints=true}.
|
||||
*/
|
||||
static void validateEndpointHost(URI endpoint, boolean allowPrivate) {
|
||||
if (allowPrivate) {
|
||||
return;
|
||||
}
|
||||
String host = endpoint.getHost();
|
||||
if (host == null || host.isBlank()) {
|
||||
throw new IllegalStateException("storage.s3.endpoint must include a host: " + endpoint);
|
||||
}
|
||||
InetAddress[] addresses;
|
||||
try {
|
||||
addresses = InetAddress.getAllByName(host);
|
||||
} catch (UnknownHostException e) {
|
||||
throw new IllegalStateException(
|
||||
"Unable to resolve storage.s3.endpoint host '" + host + "'", e);
|
||||
}
|
||||
for (InetAddress address : addresses) {
|
||||
if (isPrivateOrLocal(address)) {
|
||||
throw new IllegalStateException(
|
||||
"storage.s3.endpoint host '"
|
||||
+ host
|
||||
+ "' resolves to private/link-local address "
|
||||
+ address.getHostAddress()
|
||||
+ "; set storage.s3.allow-private-endpoints=true to opt in"
|
||||
+ " (e.g. for MinIO or in-cluster S3).");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private static boolean isPrivateOrLocal(InetAddress address) {
|
||||
return address.isLoopbackAddress()
|
||||
|| address.isLinkLocalAddress()
|
||||
|| address.isSiteLocalAddress()
|
||||
|| address.isAnyLocalAddress()
|
||||
|| address.isMulticastAddress();
|
||||
}
|
||||
|
||||
static RequestChecksumCalculation parseRequestChecksum(String value) {
|
||||
if (value == null || value.isBlank()) {
|
||||
return RequestChecksumCalculation.WHEN_SUPPORTED;
|
||||
}
|
||||
try {
|
||||
return RequestChecksumCalculation.valueOf(
|
||||
value.trim().toUpperCase(java.util.Locale.ROOT));
|
||||
} catch (IllegalArgumentException ex) {
|
||||
log.warn(
|
||||
"Unknown storage.s3.request-checksum-calculation value '{}', falling back to WHEN_SUPPORTED",
|
||||
value);
|
||||
return RequestChecksumCalculation.WHEN_SUPPORTED;
|
||||
}
|
||||
}
|
||||
|
||||
static ResponseChecksumValidation parseResponseChecksum(String value) {
|
||||
if (value == null || value.isBlank()) {
|
||||
return ResponseChecksumValidation.WHEN_SUPPORTED;
|
||||
}
|
||||
try {
|
||||
return ResponseChecksumValidation.valueOf(
|
||||
value.trim().toUpperCase(java.util.Locale.ROOT));
|
||||
} catch (IllegalArgumentException ex) {
|
||||
log.warn(
|
||||
"Unknown storage.s3.response-checksum-validation value '{}', falling back to WHEN_SUPPORTED",
|
||||
value);
|
||||
return ResponseChecksumValidation.WHEN_SUPPORTED;
|
||||
}
|
||||
}
|
||||
}
|
||||
+226
@@ -0,0 +1,226 @@
|
||||
package stirling.software.proprietary.cluster.s3;
|
||||
|
||||
import java.io.BufferedInputStream;
|
||||
import java.io.IOException;
|
||||
import java.io.InputStream;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.nio.file.StandardCopyOption;
|
||||
import java.util.Optional;
|
||||
import java.util.UUID;
|
||||
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
|
||||
import stirling.software.common.cluster.FileStore;
|
||||
|
||||
import software.amazon.awssdk.core.ResponseInputStream;
|
||||
import software.amazon.awssdk.core.exception.SdkException;
|
||||
import software.amazon.awssdk.core.sync.RequestBody;
|
||||
import software.amazon.awssdk.services.s3.S3Client;
|
||||
import software.amazon.awssdk.services.s3.model.DeleteObjectRequest;
|
||||
import software.amazon.awssdk.services.s3.model.GetObjectRequest;
|
||||
import software.amazon.awssdk.services.s3.model.GetObjectResponse;
|
||||
import software.amazon.awssdk.services.s3.model.HeadObjectRequest;
|
||||
import software.amazon.awssdk.services.s3.model.HeadObjectResponse;
|
||||
import software.amazon.awssdk.services.s3.model.NoSuchKeyException;
|
||||
import software.amazon.awssdk.services.s3.model.PutObjectRequest;
|
||||
import software.amazon.awssdk.services.s3.model.S3Exception;
|
||||
|
||||
/**
|
||||
* S3-backed {@link FileStore} for transient job-result files. Objects are namespaced under a
|
||||
* configurable key prefix (default {@code transient/}) and can coexist in the same bucket as {@code
|
||||
* S3StorageProvider}.
|
||||
*/
|
||||
@Slf4j
|
||||
public class S3FileStore implements FileStore, AutoCloseable {
|
||||
|
||||
public static final String DEFAULT_KEY_PREFIX = "transient/";
|
||||
|
||||
private final S3Client s3Client;
|
||||
private final String bucket;
|
||||
private final String keyPrefix;
|
||||
private final boolean ownsClient;
|
||||
|
||||
public S3FileStore(S3Client s3Client, String bucket) {
|
||||
this(s3Client, bucket, DEFAULT_KEY_PREFIX, true);
|
||||
}
|
||||
|
||||
public S3FileStore(S3Client s3Client, String bucket, String keyPrefix) {
|
||||
this(s3Client, bucket, keyPrefix, true);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param ownsClient when true, {@link #close()} will close the supplied client. Set to false in
|
||||
* tests that share the client with another consumer.
|
||||
*/
|
||||
public S3FileStore(S3Client s3Client, String bucket, String keyPrefix, boolean ownsClient) {
|
||||
if (bucket == null || bucket.isBlank()) {
|
||||
throw new IllegalArgumentException("S3 bucket must be configured");
|
||||
}
|
||||
this.s3Client = s3Client;
|
||||
this.bucket = bucket;
|
||||
this.keyPrefix = normalizePrefix(keyPrefix);
|
||||
this.ownsClient = ownsClient;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Stored store(InputStream in, String originalName) throws IOException {
|
||||
String fileId = UUID.randomUUID().toString();
|
||||
// S3 PUT requires a known content-length; spool to a temp file first so memory stays
|
||||
// bounded for large payloads, then stream the file to S3 via RequestBody.fromFile.
|
||||
Path tempFile = Files.createTempFile("s3-upload-", ".bin");
|
||||
long size;
|
||||
try {
|
||||
try (InputStream src = in) {
|
||||
Files.copy(src, tempFile, StandardCopyOption.REPLACE_EXISTING);
|
||||
}
|
||||
size = Files.size(tempFile);
|
||||
PutObjectRequest request =
|
||||
PutObjectRequest.builder().bucket(bucket).key(resolveKey(fileId)).build();
|
||||
try {
|
||||
s3Client.putObject(request, RequestBody.fromFile(tempFile));
|
||||
} catch (SdkException e) {
|
||||
throw new IOException("Failed to upload object to S3", e);
|
||||
}
|
||||
} finally {
|
||||
try {
|
||||
Files.deleteIfExists(tempFile);
|
||||
} catch (IOException cleanupError) {
|
||||
log.warn("Failed to delete S3 upload temp file: {}", tempFile, cleanupError);
|
||||
}
|
||||
}
|
||||
return new Stored(fileId, size);
|
||||
}
|
||||
|
||||
@Override
|
||||
public InputStream retrieve(String fileId) throws IOException {
|
||||
validateFileId(fileId);
|
||||
GetObjectRequest request =
|
||||
GetObjectRequest.builder().bucket(bucket).key(resolveKey(fileId)).build();
|
||||
try {
|
||||
ResponseInputStream<GetObjectResponse> stream = s3Client.getObject(request);
|
||||
return new BufferedInputStream(stream);
|
||||
} catch (NoSuchKeyException e) {
|
||||
throw new IOException("File not found with ID: " + fileId, e);
|
||||
} catch (SdkException e) {
|
||||
throw new IOException("Failed to load object from S3", e);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public byte[] retrieveBytes(String fileId) throws IOException {
|
||||
validateFileId(fileId);
|
||||
GetObjectRequest request =
|
||||
GetObjectRequest.builder().bucket(bucket).key(resolveKey(fileId)).build();
|
||||
try (ResponseInputStream<GetObjectResponse> stream = s3Client.getObject(request)) {
|
||||
return stream.readAllBytes();
|
||||
} catch (NoSuchKeyException e) {
|
||||
throw new IOException("File not found with ID: " + fileId, e);
|
||||
} catch (SdkException e) {
|
||||
throw new IOException("Failed to load object from S3", e);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public long size(String fileId) throws IOException {
|
||||
validateFileId(fileId);
|
||||
HeadObjectRequest request =
|
||||
HeadObjectRequest.builder().bucket(bucket).key(resolveKey(fileId)).build();
|
||||
try {
|
||||
HeadObjectResponse response = s3Client.headObject(request);
|
||||
return Optional.ofNullable(response.contentLength()).orElse(0L);
|
||||
} catch (NoSuchKeyException e) {
|
||||
throw new IOException("File not found with ID: " + fileId, e);
|
||||
} catch (S3Exception e) {
|
||||
if (e.statusCode() == 404) {
|
||||
throw new IOException("File not found with ID: " + fileId, e);
|
||||
}
|
||||
throw new IOException("Failed to head object in S3", e);
|
||||
} catch (SdkException e) {
|
||||
throw new IOException("Failed to head object in S3", e);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean delete(String fileId) {
|
||||
try {
|
||||
validateFileId(fileId);
|
||||
} catch (IllegalArgumentException e) {
|
||||
log.warn("Refusing to delete invalid file id: {}", fileId);
|
||||
return false;
|
||||
}
|
||||
try {
|
||||
s3Client.deleteObject(
|
||||
DeleteObjectRequest.builder().bucket(bucket).key(resolveKey(fileId)).build());
|
||||
return true;
|
||||
} catch (SdkException e) {
|
||||
log.error("Error deleting file with ID: {}", fileId, e);
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean exists(String fileId) {
|
||||
try {
|
||||
validateFileId(fileId);
|
||||
} catch (IllegalArgumentException e) {
|
||||
return false;
|
||||
}
|
||||
HeadObjectRequest request =
|
||||
HeadObjectRequest.builder().bucket(bucket).key(resolveKey(fileId)).build();
|
||||
try {
|
||||
s3Client.headObject(request);
|
||||
return true;
|
||||
} catch (NoSuchKeyException e) {
|
||||
return false;
|
||||
} catch (S3Exception e) {
|
||||
if (e.statusCode() == 404) {
|
||||
return false;
|
||||
}
|
||||
log.warn("Error checking existence for file ID: {}", fileId, e);
|
||||
return false;
|
||||
} catch (SdkException e) {
|
||||
log.warn("Error checking existence for file ID: {}", fileId, e);
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
if (!ownsClient) {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
s3Client.close();
|
||||
} catch (Exception e) {
|
||||
log.warn("Error closing S3 client", e);
|
||||
}
|
||||
}
|
||||
|
||||
String resolveKey(String fileId) {
|
||||
return keyPrefix + fileId;
|
||||
}
|
||||
|
||||
private static void validateFileId(String fileId) {
|
||||
if (fileId == null || fileId.isBlank()) {
|
||||
throw new IllegalArgumentException("File ID must not be blank");
|
||||
}
|
||||
if (fileId.contains("..") || fileId.contains("/") || fileId.contains("\\")) {
|
||||
throw new IllegalArgumentException("Invalid file ID");
|
||||
}
|
||||
}
|
||||
|
||||
private static String normalizePrefix(String prefix) {
|
||||
if (prefix == null || prefix.isBlank()) {
|
||||
return "";
|
||||
}
|
||||
String trimmed = prefix.trim();
|
||||
if (trimmed.startsWith("/")) {
|
||||
trimmed = trimmed.substring(1);
|
||||
}
|
||||
if (!trimmed.isEmpty() && !trimmed.endsWith("/")) {
|
||||
trimmed = trimmed + "/";
|
||||
}
|
||||
return trimmed;
|
||||
}
|
||||
}
|
||||
+37
@@ -0,0 +1,37 @@
|
||||
package stirling.software.proprietary.cluster.s3;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Value;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
|
||||
import lombok.RequiredArgsConstructor;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
|
||||
import stirling.software.common.cluster.FileStore;
|
||||
import stirling.software.common.model.ApplicationProperties;
|
||||
|
||||
/** Activates the S3-backed transient {@link FileStore} when {@code cluster.artifactStore=s3}. */
|
||||
@Slf4j
|
||||
@Configuration
|
||||
@RequiredArgsConstructor
|
||||
@ConditionalOnProperty(prefix = "cluster", name = "artifactStore", havingValue = "s3")
|
||||
public class S3FileStoreConfiguration {
|
||||
|
||||
private final ApplicationProperties applicationProperties;
|
||||
|
||||
@Bean(destroyMethod = "close")
|
||||
@ConditionalOnMissingBean
|
||||
public FileStore fileStore(@Value("${cluster.s3.keyPrefix:transient/}") String keyPrefix) {
|
||||
ApplicationProperties.Storage.S3 cfg = applicationProperties.getStorage().getS3();
|
||||
S3Clients.Bundle bundle = S3Clients.build(cfg, "cluster file store");
|
||||
// FileStore has no signed-URL contract; close the unused presigner immediately.
|
||||
try {
|
||||
bundle.presigner().close();
|
||||
} catch (Exception ignored) {
|
||||
}
|
||||
log.info("Cluster FileStore: s3 (bucket={}, keyPrefix={})", cfg.getBucket(), keyPrefix);
|
||||
return new S3FileStore(bundle.client(), cfg.getBucket(), keyPrefix, true);
|
||||
}
|
||||
}
|
||||
+1
-1
@@ -11,7 +11,7 @@ import jakarta.validation.constraints.NotNull;
|
||||
import lombok.Data;
|
||||
|
||||
@Data
|
||||
@Schema(description = "Run an AI workflow against one or more PDF files")
|
||||
@Schema(description = "Run an AI workflow")
|
||||
public class AiWorkflowRequest {
|
||||
|
||||
@NotNull
|
||||
|
||||
+5
-1
@@ -49,7 +49,11 @@ public class EEAppConfig {
|
||||
@Profile("security & !saas")
|
||||
@Bean(name = "SSOAutoLogin")
|
||||
public boolean ssoAutoLogin() {
|
||||
return applicationProperties.getPremium().getProFeatures().isSsoAutoLogin();
|
||||
boolean enabled = applicationProperties.getPremium().getProFeatures().isSsoAutoLogin();
|
||||
if (enabled) {
|
||||
licenseKeyChecker.requireProOrEnterprise("premium.proFeatures.ssoAutoLogin=true");
|
||||
}
|
||||
return enabled;
|
||||
}
|
||||
|
||||
// TODO: Remove post migration
|
||||
|
||||
+16
-1
@@ -32,7 +32,10 @@ public class LicenseKeyChecker {
|
||||
|
||||
private final UserLicenseSettingsService licenseSettingsService;
|
||||
|
||||
private License premiumEnabledResult = License.NORMAL;
|
||||
// volatile: written by evaluateLicense() on the @Scheduled refresh thread, read by request
|
||||
// threads via getPremiumLicenseEnabledResult() / requireProOrEnterprise(). Ensures readers see
|
||||
// the latest tier rather than a stale cached value.
|
||||
private volatile License premiumEnabledResult = License.NORMAL;
|
||||
|
||||
public LicenseKeyChecker(
|
||||
KeygenLicenseVerifier licenseService,
|
||||
@@ -133,4 +136,16 @@ public class LicenseKeyChecker {
|
||||
public License getPremiumLicenseEnabledResult() {
|
||||
return premiumEnabledResult;
|
||||
}
|
||||
|
||||
/**
|
||||
* Throws {@link IllegalStateException} if the current license is not Pro or Enterprise. Used by
|
||||
* boot-time gates to fail fast when an operator enables a premium-only setting without a valid
|
||||
* license. {@code configuredAs} is the human-readable property path (e.g. {@code
|
||||
* "storage.provider=s3"}) and appears in the exception message.
|
||||
*/
|
||||
public void requireProOrEnterprise(String configuredAs) {
|
||||
if (premiumEnabledResult != License.SERVER && premiumEnabledResult != License.ENTERPRISE) {
|
||||
throw new IllegalStateException(configuredAs + " requires a Pro or Enterprise license");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+5
-1
@@ -17,6 +17,7 @@ import lombok.extern.slf4j.Slf4j;
|
||||
|
||||
import stirling.software.common.annotations.AutoJobPostMapping;
|
||||
import stirling.software.common.annotations.api.GeneralApi;
|
||||
import stirling.software.common.enumeration.ResourceWeight;
|
||||
import stirling.software.proprietary.security.model.api.Email;
|
||||
import stirling.software.proprietary.security.service.EmailService;
|
||||
|
||||
@@ -39,7 +40,10 @@ public class EmailController {
|
||||
* attachment.
|
||||
* @return ResponseEntity with success or error message.
|
||||
*/
|
||||
@AutoJobPostMapping(consumes = MediaType.MULTIPART_FORM_DATA_VALUE, value = "/send-email")
|
||||
@AutoJobPostMapping(
|
||||
consumes = MediaType.MULTIPART_FORM_DATA_VALUE,
|
||||
value = "/send-email",
|
||||
resourceWeight = ResourceWeight.SMALL_WEIGHT)
|
||||
@Operation(
|
||||
summary = "Send an email with an attachment",
|
||||
description =
|
||||
|
||||
+171
-4
@@ -1,6 +1,10 @@
|
||||
package stirling.software.proprietary.security.service;
|
||||
|
||||
import java.util.Collections;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.Set;
|
||||
import java.util.TreeSet;
|
||||
|
||||
import org.springframework.security.authentication.LockedException;
|
||||
import org.springframework.security.oauth2.client.oidc.userinfo.OidcUserRequest;
|
||||
@@ -8,6 +12,8 @@ import org.springframework.security.oauth2.client.oidc.userinfo.OidcUserService;
|
||||
import org.springframework.security.oauth2.client.userinfo.OAuth2UserService;
|
||||
import org.springframework.security.oauth2.core.OAuth2AuthenticationException;
|
||||
import org.springframework.security.oauth2.core.OAuth2Error;
|
||||
import org.springframework.security.oauth2.core.oidc.OidcIdToken;
|
||||
import org.springframework.security.oauth2.core.oidc.OidcUserInfo;
|
||||
import org.springframework.security.oauth2.core.oidc.user.DefaultOidcUser;
|
||||
import org.springframework.security.oauth2.core.oidc.user.OidcUser;
|
||||
|
||||
@@ -39,20 +45,37 @@ public class CustomOAuth2UserService implements OAuth2UserService<OidcUserReques
|
||||
|
||||
@Override
|
||||
public OidcUser loadUser(OidcUserRequest userRequest) throws OAuth2AuthenticationException {
|
||||
String registrationId = userRequest.getClientRegistration().getRegistrationId();
|
||||
boolean debugLogging = Boolean.TRUE.equals(oauth2Properties.getDebugLogging());
|
||||
// Resolved inside the try so a bad/null useAsUsername (IllegalArgumentException from
|
||||
// valueOf, or NPE on toUpperCase) is caught and wrapped as OAuth2AuthenticationException
|
||||
// by the existing handlers below, matching the pre-debugLogging behaviour.
|
||||
String usernameAttributeKey = null;
|
||||
|
||||
try {
|
||||
OidcUser user = delegate.loadUser(userRequest);
|
||||
String usernameAttributeKey =
|
||||
usernameAttributeKey =
|
||||
UsernameAttribute.valueOf(oauth2Properties.getUseAsUsername().toUpperCase())
|
||||
.getName();
|
||||
OidcUser user = delegate.loadUser(userRequest);
|
||||
|
||||
if (debugLogging) {
|
||||
logClaimDump(
|
||||
"OAuth2/OIDC login claims received",
|
||||
registrationId,
|
||||
usernameAttributeKey,
|
||||
user.getIdToken(),
|
||||
user.getUserInfo(),
|
||||
user.getAttributes(),
|
||||
false);
|
||||
}
|
||||
|
||||
// Extract SSO provider information
|
||||
String ssoProviderId = user.getSubject(); // Standard OIDC 'sub' claim
|
||||
String ssoProvider = userRequest.getClientRegistration().getRegistrationId();
|
||||
String username = user.getAttribute(usernameAttributeKey);
|
||||
|
||||
log.debug(
|
||||
"OAuth2 login - Provider: {}, ProviderId: {}, Username: {}",
|
||||
ssoProvider,
|
||||
registrationId,
|
||||
ssoProviderId,
|
||||
username);
|
||||
|
||||
@@ -79,10 +102,154 @@ public class CustomOAuth2UserService implements OAuth2UserService<OidcUserReques
|
||||
usernameAttributeKey);
|
||||
} catch (IllegalArgumentException e) {
|
||||
log.error("Error loading OIDC user: {}", e.getMessage());
|
||||
// Only emit the claim dump if we successfully resolved usernameAttributeKey. A null
|
||||
// value here means UsernameAttribute.valueOf rejected the configured useAsUsername
|
||||
// before delegate.loadUser ran — that error message is self-explanatory and a claim
|
||||
// dump would have no resolved-key to compare against.
|
||||
if (debugLogging && usernameAttributeKey != null) {
|
||||
// The DefaultOidcUser constructor (or our own checks) rejected the chosen
|
||||
// username attribute. Dump the claims we DID receive so the operator can pick
|
||||
// a different value for security.oauth2.useAsUsername.
|
||||
logClaimDump(
|
||||
"OAuth2/OIDC login FAILED - dumping received claims",
|
||||
registrationId,
|
||||
usernameAttributeKey,
|
||||
userRequest.getIdToken(),
|
||||
null,
|
||||
userRequest.getIdToken() == null
|
||||
? Collections.emptyMap()
|
||||
: userRequest.getIdToken().getClaims(),
|
||||
true);
|
||||
}
|
||||
throw new OAuth2AuthenticationException(new OAuth2Error(e.getMessage()), e);
|
||||
} catch (Exception e) {
|
||||
log.error("Unexpected error loading OIDC user", e);
|
||||
if (debugLogging && usernameAttributeKey != null && userRequest.getIdToken() != null) {
|
||||
logClaimDump(
|
||||
"OAuth2/OIDC login FAILED (unexpected error) - dumping ID token claims",
|
||||
registrationId,
|
||||
usernameAttributeKey,
|
||||
userRequest.getIdToken(),
|
||||
null,
|
||||
userRequest.getIdToken().getClaims(),
|
||||
true);
|
||||
}
|
||||
throw new OAuth2AuthenticationException("Unexpected error during authentication");
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Emits a multi-line diagnostic dump of the claims returned by the OAuth2/OIDC provider. Only
|
||||
* invoked when {@code security.oauth2.debugLogging=true}.
|
||||
*
|
||||
* @param banner short title for the log block
|
||||
* @param registrationId Spring client registration id (e.g. "demarest", "keycloak")
|
||||
* @param usernameAttributeKey the claim key the application is configured to use as username
|
||||
* @param idToken the decoded ID token, may be null on unexpected failures
|
||||
* @param userInfo the decoded UserInfo response, may be null if the provider returned none
|
||||
* @param mergedAttributes the merged attribute map Spring uses for {@code getAttribute()}
|
||||
* @param failure true if logging in the error path (uses ERROR level), false for INFO
|
||||
*/
|
||||
private void logClaimDump(
|
||||
String banner,
|
||||
String registrationId,
|
||||
String usernameAttributeKey,
|
||||
OidcIdToken idToken,
|
||||
OidcUserInfo userInfo,
|
||||
Map<String, Object> mergedAttributes,
|
||||
boolean failure) {
|
||||
StringBuilder sb = new StringBuilder();
|
||||
sb.append("\n========== [OAUTH2 DEBUG] ").append(banner).append(" ==========\n");
|
||||
sb.append("Provider registrationId : ").append(registrationId).append('\n');
|
||||
sb.append("Configured useAsUsername: ")
|
||||
.append(oauth2Properties.getUseAsUsername())
|
||||
.append(" (looks up claim key '")
|
||||
.append(usernameAttributeKey)
|
||||
.append("')\n");
|
||||
|
||||
if (idToken != null) {
|
||||
Map<String, Object> idClaims = idToken.getClaims();
|
||||
sb.append("\n-- ID token claims (")
|
||||
.append(idClaims == null ? 0 : idClaims.size())
|
||||
.append(") --\n");
|
||||
appendClaims(sb, idClaims);
|
||||
sb.append("ID token issued at : ").append(idToken.getIssuedAt()).append('\n');
|
||||
sb.append("ID token expires at: ").append(idToken.getExpiresAt()).append('\n');
|
||||
} else {
|
||||
sb.append("\n-- ID token: <null> --\n");
|
||||
}
|
||||
|
||||
if (userInfo != null && userInfo.getClaims() != null) {
|
||||
sb.append("\n-- UserInfo endpoint claims (")
|
||||
.append(userInfo.getClaims().size())
|
||||
.append(") --\n");
|
||||
appendClaims(sb, userInfo.getClaims());
|
||||
} else {
|
||||
sb.append("\n-- UserInfo endpoint claims: none returned --\n");
|
||||
}
|
||||
|
||||
if (mergedAttributes != null) {
|
||||
sb.append("\n-- Merged attribute keys available to useAsUsername: ")
|
||||
.append(new TreeSet<>(mergedAttributes.keySet()))
|
||||
.append("\n");
|
||||
Object resolved = mergedAttributes.get(usernameAttributeKey);
|
||||
sb.append("-- Value at '")
|
||||
.append(usernameAttributeKey)
|
||||
.append("' : ")
|
||||
.append(resolved == null ? "<NULL — this is why login fails>" : resolved)
|
||||
.append('\n');
|
||||
|
||||
if (resolved == null) {
|
||||
Set<String> hints = suggestUsernameClaims(mergedAttributes.keySet());
|
||||
if (!hints.isEmpty()) {
|
||||
sb.append(
|
||||
"-- Hint: the following claim(s) are present and map to a"
|
||||
+ " known UsernameAttribute value — try setting"
|
||||
+ " security.oauth2.useAsUsername to one of: ")
|
||||
.append(hints)
|
||||
.append('\n');
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
sb.append(
|
||||
"\nWARNING: this block contains PII. Set security.oauth2.debugLogging=false once"
|
||||
+ " troubleshooting is complete.\n");
|
||||
sb.append("========== [/OAUTH2 DEBUG] ==========");
|
||||
|
||||
if (failure) {
|
||||
log.error(sb.toString());
|
||||
} else {
|
||||
log.info(sb.toString());
|
||||
}
|
||||
}
|
||||
|
||||
private static void appendClaims(StringBuilder sb, Map<String, Object> claims) {
|
||||
if (claims == null || claims.isEmpty()) {
|
||||
sb.append(" (no claims)\n");
|
||||
return;
|
||||
}
|
||||
// Sort for stable, scannable output
|
||||
new TreeSet<>(claims.keySet())
|
||||
.forEach(
|
||||
key -> {
|
||||
Object value = claims.get(key);
|
||||
sb.append(" ").append(key).append(" = ").append(value).append('\n');
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns the intersection of the claim keys the provider actually returned and the keys that
|
||||
* {@link UsernameAttribute} accepts — i.e. valid values the operator could put in {@code
|
||||
* security.oauth2.useAsUsername} to make this login work.
|
||||
*/
|
||||
private static Set<String> suggestUsernameClaims(Set<String> availableClaimKeys) {
|
||||
Set<String> supported = new TreeSet<>();
|
||||
for (UsernameAttribute attr : UsernameAttribute.values()) {
|
||||
if (availableClaimKeys.contains(attr.getName())) {
|
||||
supported.add(attr.getName());
|
||||
}
|
||||
}
|
||||
return supported;
|
||||
}
|
||||
}
|
||||
|
||||
+95
@@ -0,0 +1,95 @@
|
||||
package stirling.software.proprietary.storage.config;
|
||||
|
||||
import java.util.Locale;
|
||||
import java.util.Optional;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Value;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
|
||||
import jakarta.annotation.PostConstruct;
|
||||
|
||||
import lombok.RequiredArgsConstructor;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
|
||||
import stirling.software.common.model.ApplicationProperties;
|
||||
import stirling.software.proprietary.security.configuration.ee.LicenseKeyChecker;
|
||||
|
||||
/**
|
||||
* Fails fast at boot if cluster mode is enabled with node-local storage. Validates both {@code
|
||||
* storage.provider} (persistent uploads) and {@code cluster.artifactStore} (transient job-result
|
||||
* files): neither may be {@code local} when {@code cluster.enabled=true}. Additionally enforces
|
||||
* that any S3-backed configuration ({@code storage.provider=s3} or {@code
|
||||
* cluster.artifactStore=s3}) is accompanied by a valid Pro / Enterprise license.
|
||||
*/
|
||||
@Configuration
|
||||
@RequiredArgsConstructor
|
||||
@Slf4j
|
||||
public class ClusterStorageGate {
|
||||
|
||||
private final ApplicationProperties applicationProperties;
|
||||
private final LicenseKeyChecker licenseKeyChecker;
|
||||
|
||||
@Value("${cluster.enabled:false}")
|
||||
private boolean clusterEnabled;
|
||||
|
||||
@Value("${cluster.artifactStore:local}")
|
||||
private String clusterArtifactStore;
|
||||
|
||||
@PostConstruct
|
||||
void validate() {
|
||||
// License enforcement runs regardless of cluster.enabled: even a single-node setup that
|
||||
// selects a remote backend must hold a Pro or higher license.
|
||||
ApplicationProperties.Storage storage = applicationProperties.getStorage();
|
||||
if (storage != null && storage.isEnabled()) {
|
||||
String provider = normalize(storage.getProvider());
|
||||
if ("s3".equals(provider) || "database".equals(provider)) {
|
||||
licenseKeyChecker.requireProOrEnterprise("storage.provider=" + provider);
|
||||
}
|
||||
}
|
||||
if ("s3".equals(normalize(clusterArtifactStore))) {
|
||||
licenseKeyChecker.requireProOrEnterprise("cluster.artifactStore=s3");
|
||||
}
|
||||
|
||||
if (!clusterEnabled) {
|
||||
return;
|
||||
}
|
||||
if (storage != null && storage.isEnabled()) {
|
||||
validate(
|
||||
"storage.provider",
|
||||
storage.getProvider(),
|
||||
"Local filesystem storage cannot be shared across cluster nodes."
|
||||
+ " Configure storage.provider=s3 (with storage.s3.bucket /"
|
||||
+ " endpoint / credentials) or storage.provider=database before"
|
||||
+ " enabling clustering.");
|
||||
}
|
||||
validate(
|
||||
"cluster.artifactStore",
|
||||
clusterArtifactStore,
|
||||
"Per-node disk cannot back transient job-result files in a multi-node"
|
||||
+ " deployment; downloads would 404 whenever the load balancer routes"
|
||||
+ " a follow-up request to a different node. Configure"
|
||||
+ " cluster.artifactStore=s3 (reuses storage.s3.* config)"
|
||||
+ " before enabling clustering.");
|
||||
}
|
||||
|
||||
private static String normalize(String value) {
|
||||
return Optional.ofNullable(value).orElse("local").trim().toLowerCase(Locale.ROOT);
|
||||
}
|
||||
|
||||
private static void validate(String propertyName, String configuredValue, String remediation) {
|
||||
String normalized =
|
||||
Optional.ofNullable(configuredValue)
|
||||
.orElse("local")
|
||||
.trim()
|
||||
.toLowerCase(Locale.ROOT);
|
||||
if ("local".equals(normalized)) {
|
||||
throw new IllegalStateException(
|
||||
"Cluster mode (cluster.enabled=true) is incompatible with "
|
||||
+ propertyName
|
||||
+ "=local. "
|
||||
+ remediation);
|
||||
}
|
||||
log.info(
|
||||
"Cluster storage gate: clusterEnabled=true, {}={} -> OK", propertyName, normalized);
|
||||
}
|
||||
}
|
||||
+15
-1
@@ -15,8 +15,11 @@ import lombok.extern.slf4j.Slf4j;
|
||||
|
||||
import stirling.software.common.configuration.InstallationPathConfig;
|
||||
import stirling.software.common.model.ApplicationProperties;
|
||||
import stirling.software.proprietary.cluster.s3.S3Clients;
|
||||
import stirling.software.proprietary.security.configuration.ee.LicenseKeyChecker;
|
||||
import stirling.software.proprietary.storage.provider.DatabaseStorageProvider;
|
||||
import stirling.software.proprietary.storage.provider.LocalStorageProvider;
|
||||
import stirling.software.proprietary.storage.provider.S3StorageProvider;
|
||||
import stirling.software.proprietary.storage.provider.StorageProvider;
|
||||
import stirling.software.proprietary.storage.repository.StoredFileBlobRepository;
|
||||
|
||||
@@ -27,8 +30,9 @@ public class StorageProviderConfig {
|
||||
|
||||
private final ApplicationProperties applicationProperties;
|
||||
private final StoredFileBlobRepository storedFileBlobRepository;
|
||||
private final LicenseKeyChecker licenseKeyChecker;
|
||||
|
||||
@Bean
|
||||
@Bean(destroyMethod = "close")
|
||||
public StorageProvider storageProvider() {
|
||||
boolean storageEnabled = applicationProperties.getStorage().isEnabled();
|
||||
String providerName =
|
||||
@@ -37,8 +41,13 @@ public class StorageProviderConfig {
|
||||
.trim()
|
||||
.toLowerCase(Locale.ROOT);
|
||||
if ("database".equals(providerName)) {
|
||||
licenseKeyChecker.requireProOrEnterprise("storage.provider=database");
|
||||
return new DatabaseStorageProvider(storedFileBlobRepository);
|
||||
}
|
||||
if ("s3".equals(providerName)) {
|
||||
licenseKeyChecker.requireProOrEnterprise("storage.provider=s3");
|
||||
return buildS3Provider(applicationProperties.getStorage().getS3());
|
||||
}
|
||||
if (!"local".equals(providerName)) {
|
||||
throw new IllegalStateException("Storage provider not supported: " + providerName);
|
||||
}
|
||||
@@ -71,4 +80,9 @@ public class StorageProviderConfig {
|
||||
}
|
||||
return new LocalStorageProvider(basePath);
|
||||
}
|
||||
|
||||
private S3StorageProvider buildS3Provider(ApplicationProperties.Storage.S3 cfg) {
|
||||
S3Clients.Bundle bundle = S3Clients.build(cfg, "storage provider");
|
||||
return new S3StorageProvider(bundle.client(), bundle.presigner(), cfg.getBucket());
|
||||
}
|
||||
}
|
||||
|
||||
+91
@@ -0,0 +1,91 @@
|
||||
package stirling.software.proprietary.storage.controller;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.UUID;
|
||||
|
||||
import org.springframework.http.HttpStatus;
|
||||
import org.springframework.http.ResponseEntity;
|
||||
import org.springframework.web.bind.annotation.PatchMapping;
|
||||
import org.springframework.web.bind.annotation.PathVariable;
|
||||
import org.springframework.web.bind.annotation.RequestBody;
|
||||
import org.springframework.web.bind.annotation.RequestMapping;
|
||||
import org.springframework.web.bind.annotation.RestController;
|
||||
|
||||
import jakarta.validation.Valid;
|
||||
import jakarta.validation.constraints.NotNull;
|
||||
import jakarta.validation.constraints.Size;
|
||||
|
||||
import lombok.AllArgsConstructor;
|
||||
import lombok.Data;
|
||||
import lombok.NoArgsConstructor;
|
||||
import lombok.RequiredArgsConstructor;
|
||||
|
||||
import stirling.software.proprietary.storage.service.FolderService;
|
||||
|
||||
/**
|
||||
* Folder placement endpoints for existing stored files. Thin adapter: validates the request shape,
|
||||
* delegates the transaction to {@link FolderService}, then maps the result onto the HTTP status.
|
||||
* Authentication, storage-gate, ownership checks, and the bulk cap all live on the service (where
|
||||
* {@code @Transactional} also lives) so the JDBC connection isn't held through JSON serialization.
|
||||
*/
|
||||
@RestController
|
||||
@RequestMapping("/api/v1/storage/files")
|
||||
@RequiredArgsConstructor
|
||||
public class FileFolderPlacementController {
|
||||
|
||||
private static final int BULK_MOVE_MAX_FILES = 1000;
|
||||
|
||||
private final FolderService folderService;
|
||||
|
||||
/** Move a single file to a folder (or to root when folderId is null). */
|
||||
@PatchMapping("/{fileId}/folder")
|
||||
public ResponseEntity<Void> moveFileToFolder(
|
||||
@PathVariable Long fileId, @Valid @RequestBody FolderPlacement body) {
|
||||
folderService.moveFileToFolder(fileId, body.getFolderId());
|
||||
return ResponseEntity.noContent().build();
|
||||
}
|
||||
|
||||
/**
|
||||
* Bulk move - fewer round-trips than calling the single endpoint N times. Returns 200 on full
|
||||
* success, 207 (Multi-Status) when some files were skipped (typically because they don't belong
|
||||
* to the caller).
|
||||
*/
|
||||
@PatchMapping("/folder")
|
||||
public ResponseEntity<BulkMoveResponse> bulkMove(@Valid @RequestBody BulkMoveRequest body) {
|
||||
FolderService.BulkMoveResult result =
|
||||
folderService.bulkMoveFilesToFolder(body.getFolderId(), body.getFileIds());
|
||||
HttpStatus status =
|
||||
result.skippedFileIds().isEmpty() ? HttpStatus.OK : HttpStatus.MULTI_STATUS;
|
||||
return ResponseEntity.status(status)
|
||||
.body(new BulkMoveResponse(result.movedFileIds(), result.skippedFileIds()));
|
||||
}
|
||||
|
||||
@Data
|
||||
@NoArgsConstructor
|
||||
@AllArgsConstructor
|
||||
public static class FolderPlacement {
|
||||
private UUID folderId;
|
||||
}
|
||||
|
||||
@Data
|
||||
@NoArgsConstructor
|
||||
@AllArgsConstructor
|
||||
public static class BulkMoveRequest {
|
||||
private UUID folderId;
|
||||
|
||||
@NotNull
|
||||
@Size(
|
||||
min = 1,
|
||||
max = BULK_MOVE_MAX_FILES,
|
||||
message = "fileIds must contain between 1 and 1000 entries")
|
||||
private List<Long> fileIds;
|
||||
}
|
||||
|
||||
@Data
|
||||
@NoArgsConstructor
|
||||
@AllArgsConstructor
|
||||
public static class BulkMoveResponse {
|
||||
private List<Long> movedFileIds;
|
||||
private List<Long> skippedFileIds;
|
||||
}
|
||||
}
|
||||
+46
-2
@@ -1,7 +1,11 @@
|
||||
package stirling.software.proprietary.storage.controller;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.net.URI;
|
||||
import java.time.Duration;
|
||||
import java.util.List;
|
||||
import java.util.Locale;
|
||||
import java.util.Optional;
|
||||
|
||||
import org.springframework.http.ContentDisposition;
|
||||
import org.springframework.http.HttpHeaders;
|
||||
@@ -25,6 +29,7 @@ import org.springframework.web.server.ResponseStatusException;
|
||||
import io.swagger.v3.oas.annotations.tags.Tag;
|
||||
|
||||
import lombok.RequiredArgsConstructor;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
|
||||
import stirling.software.proprietary.security.model.User;
|
||||
import stirling.software.proprietary.storage.model.FileShare;
|
||||
@@ -35,17 +40,22 @@ import stirling.software.proprietary.storage.model.api.ShareLinkMetadataResponse
|
||||
import stirling.software.proprietary.storage.model.api.ShareLinkResponse;
|
||||
import stirling.software.proprietary.storage.model.api.ShareWithUserRequest;
|
||||
import stirling.software.proprietary.storage.model.api.StoredFileResponse;
|
||||
import stirling.software.proprietary.storage.provider.StorageProvider;
|
||||
import stirling.software.proprietary.storage.service.FileStorageService;
|
||||
|
||||
@RestController
|
||||
@RequestMapping("/api/v1/storage")
|
||||
@RequiredArgsConstructor
|
||||
@Slf4j
|
||||
@Tag(
|
||||
name = "File Storage",
|
||||
description = "Stored file management, sharing, and share link operations")
|
||||
public class FileStorageController {
|
||||
|
||||
private static final Duration SIGNED_URL_TTL = Duration.ofMinutes(5);
|
||||
|
||||
private final FileStorageService fileStorageService;
|
||||
private final StorageProvider storageProvider;
|
||||
|
||||
@PostMapping(
|
||||
value = "/files",
|
||||
@@ -91,7 +101,9 @@ public class FileStorageController {
|
||||
User user = fileStorageService.requireAuthenticatedUser();
|
||||
StoredFile file = fileStorageService.getAccessibleFile(user, fileId);
|
||||
fileStorageService.requireReadAccess(user, file);
|
||||
return buildFileResponse(file, inline);
|
||||
Optional<ResponseEntity<org.springframework.core.io.Resource>> redirect =
|
||||
tryRedirectToSignedUrl(file, inline);
|
||||
return redirect.orElseGet(() -> buildFileResponse(file, inline));
|
||||
}
|
||||
|
||||
@DeleteMapping("/files/{fileId}")
|
||||
@@ -189,7 +201,9 @@ public class FileStorageController {
|
||||
fileStorageService.requireReadAccess(share);
|
||||
fileStorageService.recordShareAccess(share, authentication, inline);
|
||||
StoredFile file = share.getFile();
|
||||
return buildFileResponse(file, inline);
|
||||
Optional<ResponseEntity<org.springframework.core.io.Resource>> redirect =
|
||||
tryRedirectToSignedUrl(file, inline);
|
||||
return redirect.orElseGet(() -> buildFileResponse(file, inline));
|
||||
}
|
||||
|
||||
@GetMapping("/share-links/{token}/metadata")
|
||||
@@ -272,4 +286,34 @@ public class FileStorageController {
|
||||
&& authentication.isAuthenticated()
|
||||
&& !"anonymousUser".equals(authentication.getPrincipal());
|
||||
}
|
||||
|
||||
private Optional<ResponseEntity<org.springframework.core.io.Resource>> tryRedirectToSignedUrl(
|
||||
StoredFile file, boolean inline) {
|
||||
if (file == null || file.getStorageKey() == null || file.getStorageKey().isBlank()) {
|
||||
return Optional.empty();
|
||||
}
|
||||
try {
|
||||
Optional<URI> signed =
|
||||
storageProvider.signedDownloadUrl(
|
||||
file.getStorageKey(),
|
||||
SIGNED_URL_TTL,
|
||||
inline,
|
||||
file.getOriginalFilename());
|
||||
if (signed.isEmpty()) {
|
||||
return Optional.empty();
|
||||
}
|
||||
HttpHeaders headers = new HttpHeaders();
|
||||
headers.setLocation(signed.get());
|
||||
ResponseEntity<org.springframework.core.io.Resource> response =
|
||||
ResponseEntity.status(HttpStatus.FOUND).headers(headers).build();
|
||||
return Optional.of(response);
|
||||
} catch (IOException e) {
|
||||
log.warn(
|
||||
"Failed to create signed download URL for file {} (key: {}), falling back to streaming",
|
||||
file.getId(),
|
||||
file.getStorageKey(),
|
||||
e);
|
||||
return Optional.empty();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+70
@@ -0,0 +1,70 @@
|
||||
package stirling.software.proprietary.storage.controller;
|
||||
|
||||
import java.net.URI;
|
||||
import java.util.List;
|
||||
import java.util.UUID;
|
||||
|
||||
import org.springframework.http.HttpStatus;
|
||||
import org.springframework.http.ResponseEntity;
|
||||
import org.springframework.web.bind.annotation.DeleteMapping;
|
||||
import org.springframework.web.bind.annotation.GetMapping;
|
||||
import org.springframework.web.bind.annotation.PatchMapping;
|
||||
import org.springframework.web.bind.annotation.PathVariable;
|
||||
import org.springframework.web.bind.annotation.PostMapping;
|
||||
import org.springframework.web.bind.annotation.RequestBody;
|
||||
import org.springframework.web.bind.annotation.RequestMapping;
|
||||
import org.springframework.web.bind.annotation.RestController;
|
||||
|
||||
import jakarta.validation.Valid;
|
||||
|
||||
import lombok.RequiredArgsConstructor;
|
||||
|
||||
import stirling.software.proprietary.storage.model.api.CreateFolderRequest;
|
||||
import stirling.software.proprietary.storage.model.api.FolderResponse;
|
||||
import stirling.software.proprietary.storage.model.api.UpdateFolderRequest;
|
||||
import stirling.software.proprietary.storage.service.FolderService;
|
||||
|
||||
/**
|
||||
* REST endpoints for user-owned folders. Phase A - no folder-level sharing yet (Phase 3).
|
||||
*
|
||||
* <p>All operations are scoped to the authenticated user; existing single-file storage endpoints in
|
||||
* {@link FileStorageController} are left alone so the cert-signing and standard upload flows are
|
||||
* unaffected.
|
||||
*/
|
||||
@RestController
|
||||
@RequestMapping("/api/v1/storage/folders")
|
||||
@RequiredArgsConstructor
|
||||
public class FolderController {
|
||||
|
||||
private final FolderService folderService;
|
||||
|
||||
@GetMapping
|
||||
public List<FolderResponse> listFolders() {
|
||||
return folderService.listFolders();
|
||||
}
|
||||
|
||||
@PostMapping
|
||||
public ResponseEntity<FolderResponse> createFolder(
|
||||
@Valid @RequestBody CreateFolderRequest request) {
|
||||
FolderResponse response = folderService.createFolder(request);
|
||||
// 201 Created with Location header - conventional REST. The idempotent re-return path
|
||||
// (same id resubmitted) also lands here; treating it as 201 keeps wire semantics simple.
|
||||
return ResponseEntity.status(HttpStatus.CREATED)
|
||||
.location(URI.create("/api/v1/storage/folders/" + response.id()))
|
||||
.body(response);
|
||||
}
|
||||
|
||||
@PatchMapping("/{folderId}")
|
||||
public ResponseEntity<FolderResponse> updateFolder(
|
||||
@PathVariable UUID folderId, @Valid @RequestBody UpdateFolderRequest request) {
|
||||
return ResponseEntity.ok(folderService.updateFolder(folderId, request));
|
||||
}
|
||||
|
||||
@DeleteMapping("/{folderId}")
|
||||
public ResponseEntity<DeleteFolderResponse> deleteFolder(@PathVariable UUID folderId) {
|
||||
List<UUID> removed = folderService.deleteFolder(folderId);
|
||||
return ResponseEntity.ok(new DeleteFolderResponse(removed));
|
||||
}
|
||||
|
||||
public record DeleteFolderResponse(List<UUID> removedFolderIds) {}
|
||||
}
|
||||
+104
@@ -0,0 +1,104 @@
|
||||
package stirling.software.proprietary.storage.model;
|
||||
|
||||
import java.io.Serializable;
|
||||
import java.time.LocalDateTime;
|
||||
import java.util.UUID;
|
||||
|
||||
import org.hibernate.annotations.CreationTimestamp;
|
||||
import org.hibernate.annotations.OnDelete;
|
||||
import org.hibernate.annotations.OnDeleteAction;
|
||||
import org.hibernate.annotations.UpdateTimestamp;
|
||||
|
||||
import jakarta.persistence.Column;
|
||||
import jakarta.persistence.Entity;
|
||||
import jakarta.persistence.FetchType;
|
||||
import jakarta.persistence.Id;
|
||||
import jakarta.persistence.Index;
|
||||
import jakarta.persistence.JoinColumn;
|
||||
import jakarta.persistence.ManyToOne;
|
||||
import jakarta.persistence.Table;
|
||||
import jakarta.persistence.Version;
|
||||
|
||||
import lombok.Getter;
|
||||
import lombok.NoArgsConstructor;
|
||||
import lombok.Setter;
|
||||
|
||||
import stirling.software.proprietary.security.model.User;
|
||||
|
||||
/**
|
||||
* A user-owned folder used by the file manager UI to organise stored files. Phase A entity - no
|
||||
* folder-level sharing yet (Phase 3).
|
||||
*
|
||||
* <p>The id is a UUID rather than a numeric auto-increment so it round-trips with the
|
||||
* client-generated {@code FolderId} and survives cross-device sync without re-keying.
|
||||
*/
|
||||
@Entity
|
||||
@Table(
|
||||
name = "folders",
|
||||
indexes = {
|
||||
@Index(name = "idx_folders_owner", columnList = "owner_id"),
|
||||
@Index(name = "idx_folders_parent", columnList = "parent_folder_id"),
|
||||
@Index(name = "idx_folders_owner_parent", columnList = "owner_id, parent_folder_id")
|
||||
})
|
||||
@NoArgsConstructor
|
||||
@Getter
|
||||
@Setter
|
||||
public class Folder implements Serializable {
|
||||
|
||||
private static final long serialVersionUID = 1L;
|
||||
|
||||
/**
|
||||
* Dialect-portable UUID column. The previous {@code columnDefinition = "uuid"} was
|
||||
* Postgres-specific and broke on H2/MariaDB. Hibernate's {@code UUID} mapping picks the right
|
||||
* native type per dialect (BINARY(16) on H2/MariaDB, uuid on Postgres) when no explicit
|
||||
* columnDefinition is set.
|
||||
*/
|
||||
@Id
|
||||
@Column(name = "folder_id", nullable = false)
|
||||
private UUID id;
|
||||
|
||||
/**
|
||||
* {@code OnDeleteAction.CASCADE} so deleting the owning {@code User} cascades to this row at
|
||||
* the DB level - UserService.deleteUserRelatedData doesn't enumerate folders today, and leaving
|
||||
* the FK without an action throws a constraint violation on user delete.
|
||||
*/
|
||||
@ManyToOne(fetch = FetchType.LAZY)
|
||||
@JoinColumn(name = "owner_id", nullable = false)
|
||||
@OnDelete(action = OnDeleteAction.CASCADE)
|
||||
private User owner;
|
||||
|
||||
/**
|
||||
* Parent folder; null = root. {@code OnDeleteAction.CASCADE} so a backend-side parent delete
|
||||
* cleans children automatically, matching the service-layer recursive-delete contract.
|
||||
*/
|
||||
@ManyToOne(fetch = FetchType.LAZY)
|
||||
@JoinColumn(name = "parent_folder_id")
|
||||
@OnDelete(action = OnDeleteAction.CASCADE)
|
||||
private Folder parent;
|
||||
|
||||
@Column(name = "name", nullable = false, length = 255)
|
||||
private String name;
|
||||
|
||||
@Column(name = "color", length = 32)
|
||||
private String color;
|
||||
|
||||
@Column(name = "icon", length = 64)
|
||||
private String icon;
|
||||
|
||||
/**
|
||||
* Optimistic-locking version. Cross-PC sync without this lets last-write-win silently. The
|
||||
* column is nullable so existing rows from a pre-version deployment can be backfilled by
|
||||
* Hibernate's update-on-write rather than failing the ddl-auto upgrade.
|
||||
*/
|
||||
@Version
|
||||
@Column(name = "version")
|
||||
private Long version;
|
||||
|
||||
@CreationTimestamp
|
||||
@Column(name = "created_at", updatable = false)
|
||||
private LocalDateTime createdAt;
|
||||
|
||||
@UpdateTimestamp
|
||||
@Column(name = "updated_at")
|
||||
private LocalDateTime updatedAt;
|
||||
}
|
||||
+18
-1
@@ -6,6 +6,8 @@ import java.util.HashSet;
|
||||
import java.util.Set;
|
||||
|
||||
import org.hibernate.annotations.CreationTimestamp;
|
||||
import org.hibernate.annotations.OnDelete;
|
||||
import org.hibernate.annotations.OnDeleteAction;
|
||||
import org.hibernate.annotations.UpdateTimestamp;
|
||||
|
||||
import jakarta.persistence.CascadeType;
|
||||
@@ -35,7 +37,8 @@ import stirling.software.proprietary.workflow.model.WorkflowSession;
|
||||
name = "stored_files",
|
||||
indexes = {
|
||||
@Index(name = "idx_stored_files_owner", columnList = "owner_id"),
|
||||
@Index(name = "idx_stored_files_workflow", columnList = "workflow_session_id")
|
||||
@Index(name = "idx_stored_files_workflow", columnList = "workflow_session_id"),
|
||||
@Index(name = "idx_stored_files_folder", columnList = "folder_id")
|
||||
})
|
||||
@NoArgsConstructor
|
||||
@Getter
|
||||
@@ -106,6 +109,20 @@ public class StoredFile implements Serializable {
|
||||
orphanRemoval = true)
|
||||
private Set<FileShare> shares = new HashSet<>();
|
||||
|
||||
/**
|
||||
* Optional folder placement for the file manager UI. Null = root. Hibernate ddl-auto will add
|
||||
* this as a nullable column on upgrade so existing records continue to work untouched.
|
||||
*
|
||||
* <p>{@code OnDeleteAction.SET_NULL} so any backend that drops a folder row (admin script,
|
||||
* future cleanup job, cascading user delete) cleanly orphans files to root rather than leaving
|
||||
* dangling FK references. The application path ({@code FolderRepository.clearFolderForFiles})
|
||||
* still runs first as a belt-and-braces.
|
||||
*/
|
||||
@ManyToOne(fetch = FetchType.LAZY)
|
||||
@JoinColumn(name = "folder_id")
|
||||
@OnDelete(action = OnDeleteAction.SET_NULL)
|
||||
private Folder folder;
|
||||
|
||||
@CreationTimestamp
|
||||
@Column(name = "created_at", updatable = false)
|
||||
private LocalDateTime createdAt;
|
||||
|
||||
+43
@@ -0,0 +1,43 @@
|
||||
package stirling.software.proprietary.storage.model.api;
|
||||
|
||||
import java.util.UUID;
|
||||
|
||||
import jakarta.validation.constraints.NotBlank;
|
||||
import jakarta.validation.constraints.Pattern;
|
||||
import jakarta.validation.constraints.Size;
|
||||
|
||||
import lombok.AllArgsConstructor;
|
||||
import lombok.Data;
|
||||
import lombok.NoArgsConstructor;
|
||||
|
||||
@Data
|
||||
@NoArgsConstructor
|
||||
@AllArgsConstructor
|
||||
public class CreateFolderRequest {
|
||||
|
||||
/**
|
||||
* Client-generated UUID - lets the caller round-trip the same id it stored locally. Optional;
|
||||
* the server generates one when missing.
|
||||
*/
|
||||
private UUID id;
|
||||
|
||||
@NotBlank
|
||||
@Size(max = 255)
|
||||
private String name;
|
||||
|
||||
private UUID parentFolderId;
|
||||
|
||||
/** Hex colour string (#rrggbb or #rrggbbaa) - matches the frontend palette format. */
|
||||
@Size(max = 32)
|
||||
@Pattern(
|
||||
regexp = "^#[0-9a-fA-F]{6}([0-9a-fA-F]{2})?$",
|
||||
message = "color must be a #RRGGBB or #RRGGBBAA hex value")
|
||||
private String color;
|
||||
|
||||
/** Icon identifier - lowercase alphanumerics, hyphens, underscores only. */
|
||||
@Size(max = 64)
|
||||
@Pattern(
|
||||
regexp = "^[a-z0-9_-]+$",
|
||||
message = "icon must be a lowercase id (a-z, 0-9, '-' or '_')")
|
||||
private String icon;
|
||||
}
|
||||
+38
@@ -0,0 +1,38 @@
|
||||
package stirling.software.proprietary.storage.model.api;
|
||||
|
||||
import java.time.LocalDateTime;
|
||||
import java.util.UUID;
|
||||
|
||||
import stirling.software.proprietary.storage.model.Folder;
|
||||
|
||||
/**
|
||||
* Outbound DTO for folder responses. Records are immutable, value-equality-based, and far less
|
||||
* accident-prone than a {@code @Data} class with public setters.
|
||||
*/
|
||||
public record FolderResponse(
|
||||
UUID id,
|
||||
String name,
|
||||
UUID parentFolderId,
|
||||
String color,
|
||||
String icon,
|
||||
Long version,
|
||||
LocalDateTime createdAt,
|
||||
LocalDateTime updatedAt) {
|
||||
|
||||
public static FolderResponse from(Folder folder) {
|
||||
// {@code folder.getParent().getId()} on a lazy proxy returns the FK value cached at the
|
||||
// join column WITHOUT initialising the proxy under standard Hibernate, so this does
|
||||
// not N+1. If a future Hibernate update changes that, switch the JPQL list query to a
|
||||
// constructor projection.
|
||||
UUID parentId = folder.getParent() == null ? null : folder.getParent().getId();
|
||||
return new FolderResponse(
|
||||
folder.getId(),
|
||||
folder.getName(),
|
||||
parentId,
|
||||
folder.getColor(),
|
||||
folder.getIcon(),
|
||||
folder.getVersion(),
|
||||
folder.getCreatedAt(),
|
||||
folder.getUpdatedAt());
|
||||
}
|
||||
}
|
||||
+7
@@ -2,6 +2,7 @@ package stirling.software.proprietary.storage.model.api;
|
||||
|
||||
import java.time.LocalDateTime;
|
||||
import java.util.List;
|
||||
import java.util.UUID;
|
||||
|
||||
import lombok.Builder;
|
||||
import lombok.Getter;
|
||||
@@ -22,4 +23,10 @@ public class StoredFileResponse {
|
||||
private final List<SharedUserResponse> sharedUsers;
|
||||
private final List<ShareLinkResponse> shareLinks;
|
||||
private final String filePurpose;
|
||||
|
||||
/**
|
||||
* Optional folder placement (Phase A). Null when the file lives at the root or when the server
|
||||
* build doesn't have the folders feature enabled - existing clients should treat null as root.
|
||||
*/
|
||||
private final UUID folderId;
|
||||
}
|
||||
|
||||
+58
@@ -0,0 +1,58 @@
|
||||
package stirling.software.proprietary.storage.model.api;
|
||||
|
||||
import java.util.UUID;
|
||||
|
||||
import jakarta.validation.constraints.Pattern;
|
||||
import jakarta.validation.constraints.Size;
|
||||
|
||||
import lombok.AllArgsConstructor;
|
||||
import lombok.Data;
|
||||
import lombok.NoArgsConstructor;
|
||||
|
||||
/**
|
||||
* PATCH-style update - every field is optional. Send only the fields you want to change.
|
||||
*
|
||||
* <p>The {@code reparent} flag distinguishes "do not change parent" from "move to root" since
|
||||
* {@code parentFolderId == null} alone is ambiguous in a sparse body. We use a boxed {@link
|
||||
* Boolean} so a missing field deserialises to {@code null} (= "do not reparent") rather than to
|
||||
* primitive {@code false}, removing a class of "I PATCHed only the name but the server reset my
|
||||
* parent" footguns.
|
||||
*
|
||||
* <p>When the trimmed name is empty (e.g. {@code " "}) the service rejects the request with HTTP
|
||||
* 400 - silent drops are too easy to mistake for a successful rename.
|
||||
*/
|
||||
@Data
|
||||
@NoArgsConstructor
|
||||
@AllArgsConstructor
|
||||
public class UpdateFolderRequest {
|
||||
|
||||
/** When provided, must contain at least one non-whitespace character. */
|
||||
@Size(max = 255)
|
||||
@Pattern(regexp = "\\S.*", message = "name must not be blank")
|
||||
private String name;
|
||||
|
||||
private Boolean reparent;
|
||||
private UUID parentFolderId;
|
||||
|
||||
@Size(max = 32)
|
||||
@Pattern(
|
||||
regexp = "^(|#[0-9a-fA-F]{6}([0-9a-fA-F]{2})?)$",
|
||||
message = "color must be empty or a #RRGGBB / #RRGGBBAA hex value")
|
||||
private String color;
|
||||
|
||||
@Size(max = 64)
|
||||
@Pattern(
|
||||
regexp = "^([a-z0-9_-]+)?$",
|
||||
message = "icon must be a lowercase id (a-z, 0-9, '-' or '_') or empty")
|
||||
private String icon;
|
||||
|
||||
/**
|
||||
* Convenience accessor - treats null as "do not reparent". Named differently from the
|
||||
* Lombok-generated {@code getReparent()} so callers don't accidentally use one for the other
|
||||
* (the getter is nullable {@code Boolean}; this method collapses to primitive).
|
||||
*/
|
||||
@com.fasterxml.jackson.annotation.JsonIgnore
|
||||
public boolean shouldReparent() {
|
||||
return Boolean.TRUE.equals(reparent);
|
||||
}
|
||||
}
|
||||
+5
-1
@@ -24,6 +24,9 @@ public class LocalStorageProvider implements StorageProvider {
|
||||
|
||||
@Override
|
||||
public StoredObject store(User owner, MultipartFile file) throws IOException {
|
||||
if (owner == null || owner.getId() == null) {
|
||||
throw new IllegalArgumentException("owner.id is required for local storage key");
|
||||
}
|
||||
String originalFilename = sanitizeFilename(file.getOriginalFilename());
|
||||
String storageKey =
|
||||
owner.getId()
|
||||
@@ -77,6 +80,7 @@ public class LocalStorageProvider implements StorageProvider {
|
||||
if (filename == null || filename.isBlank()) {
|
||||
return "file";
|
||||
}
|
||||
return Paths.get(filename).getFileName().toString();
|
||||
String stripped = Paths.get(filename).getFileName().toString().replaceAll("\\p{Cntrl}", "");
|
||||
return stripped.isBlank() ? "file" : stripped;
|
||||
}
|
||||
}
|
||||
|
||||
+190
@@ -0,0 +1,190 @@
|
||||
package stirling.software.proprietary.storage.provider;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.io.InputStream;
|
||||
import java.net.URI;
|
||||
import java.net.URISyntaxException;
|
||||
import java.nio.file.Paths;
|
||||
import java.time.Duration;
|
||||
import java.util.Optional;
|
||||
import java.util.UUID;
|
||||
|
||||
import org.springframework.core.io.InputStreamResource;
|
||||
import org.springframework.core.io.Resource;
|
||||
import org.springframework.web.multipart.MultipartFile;
|
||||
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
|
||||
import stirling.software.proprietary.security.model.User;
|
||||
|
||||
import software.amazon.awssdk.core.ResponseInputStream;
|
||||
import software.amazon.awssdk.core.exception.SdkException;
|
||||
import software.amazon.awssdk.core.sync.RequestBody;
|
||||
import software.amazon.awssdk.services.s3.S3Client;
|
||||
import software.amazon.awssdk.services.s3.model.DeleteObjectRequest;
|
||||
import software.amazon.awssdk.services.s3.model.GetObjectRequest;
|
||||
import software.amazon.awssdk.services.s3.model.GetObjectResponse;
|
||||
import software.amazon.awssdk.services.s3.model.NoSuchKeyException;
|
||||
import software.amazon.awssdk.services.s3.model.PutObjectRequest;
|
||||
import software.amazon.awssdk.services.s3.presigner.S3Presigner;
|
||||
import software.amazon.awssdk.services.s3.presigner.model.GetObjectPresignRequest;
|
||||
import software.amazon.awssdk.services.s3.presigner.model.PresignedGetObjectRequest;
|
||||
|
||||
/** {@link StorageProvider} backed by an S3-compatible object store. */
|
||||
@Slf4j
|
||||
public class S3StorageProvider implements StorageProvider, AutoCloseable {
|
||||
|
||||
private final S3Client s3Client;
|
||||
private final S3Presigner s3Presigner;
|
||||
private final String bucket;
|
||||
|
||||
public S3StorageProvider(S3Client s3Client, S3Presigner s3Presigner, String bucket) {
|
||||
if (bucket == null || bucket.isBlank()) {
|
||||
throw new IllegalArgumentException("S3 bucket must be configured");
|
||||
}
|
||||
this.s3Client = s3Client;
|
||||
this.s3Presigner = s3Presigner;
|
||||
this.bucket = bucket;
|
||||
}
|
||||
|
||||
@Override
|
||||
public StoredObject store(User owner, MultipartFile file) throws IOException {
|
||||
if (owner == null || owner.getId() == null) {
|
||||
throw new IllegalArgumentException("owner.id is required for S3 storage key");
|
||||
}
|
||||
String originalFilename = sanitizeFilename(file.getOriginalFilename());
|
||||
// Key is opaque ({ownerId}/{uuid}) so non-ASCII filenames don't break vendors that
|
||||
// restrict key charset (e.g. Supabase Storage returns 400 Invalid key on unicode).
|
||||
// The display name is preserved in StoredObject.originalFilename and the DB row.
|
||||
String storageKey = owner.getId() + "/" + UUID.randomUUID();
|
||||
|
||||
PutObjectRequest.Builder request =
|
||||
PutObjectRequest.builder().bucket(bucket).key(storageKey);
|
||||
if (file.getContentType() != null && !file.getContentType().isBlank()) {
|
||||
request.contentType(file.getContentType());
|
||||
}
|
||||
try (InputStream inputStream = file.getInputStream()) {
|
||||
s3Client.putObject(
|
||||
request.build(), RequestBody.fromInputStream(inputStream, file.getSize()));
|
||||
} catch (SdkException e) {
|
||||
throw new IOException("Failed to upload object to S3", e);
|
||||
}
|
||||
|
||||
return StoredObject.builder()
|
||||
.storageKey(storageKey)
|
||||
.originalFilename(originalFilename)
|
||||
.contentType(file.getContentType())
|
||||
.sizeBytes(file.getSize())
|
||||
.build();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Resource load(String storageKey) throws IOException {
|
||||
GetObjectRequest request =
|
||||
GetObjectRequest.builder().bucket(bucket).key(storageKey).build();
|
||||
try {
|
||||
ResponseInputStream<GetObjectResponse> stream = s3Client.getObject(request);
|
||||
long contentLength =
|
||||
stream.response().contentLength() != null
|
||||
? stream.response().contentLength()
|
||||
: -1;
|
||||
return new InputStreamResource(stream) {
|
||||
@Override
|
||||
public long contentLength() {
|
||||
return contentLength;
|
||||
}
|
||||
};
|
||||
} catch (NoSuchKeyException e) {
|
||||
throw new IOException("File not found", e);
|
||||
} catch (SdkException e) {
|
||||
throw new IOException("Failed to load object from S3", e);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void delete(String storageKey) throws IOException {
|
||||
try {
|
||||
s3Client.deleteObject(
|
||||
DeleteObjectRequest.builder().bucket(bucket).key(storageKey).build());
|
||||
} catch (SdkException e) {
|
||||
throw new IOException("Failed to delete object from S3", e);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public Optional<URI> signedDownloadUrl(String storageKey, Duration ttl) throws IOException {
|
||||
return signedDownloadUrl(storageKey, ttl, false, null);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Optional<URI> signedDownloadUrl(
|
||||
String storageKey, Duration ttl, boolean inline, String originalFilename)
|
||||
throws IOException {
|
||||
if (storageKey == null || storageKey.isBlank()) {
|
||||
return Optional.empty();
|
||||
}
|
||||
Duration effectiveTtl =
|
||||
ttl == null || ttl.isZero() || ttl.isNegative() ? Duration.ofMinutes(5) : ttl;
|
||||
try {
|
||||
GetObjectRequest.Builder getBuilder =
|
||||
GetObjectRequest.builder().bucket(bucket).key(storageKey);
|
||||
String disposition = buildContentDisposition(inline, originalFilename);
|
||||
if (disposition != null) {
|
||||
getBuilder.responseContentDisposition(disposition);
|
||||
}
|
||||
GetObjectPresignRequest presignRequest =
|
||||
GetObjectPresignRequest.builder()
|
||||
.signatureDuration(effectiveTtl)
|
||||
.getObjectRequest(getBuilder.build())
|
||||
.build();
|
||||
PresignedGetObjectRequest presigned = s3Presigner.presignGetObject(presignRequest);
|
||||
return Optional.of(presigned.url().toURI());
|
||||
} catch (SdkException | URISyntaxException e) {
|
||||
log.warn("Failed to create presigned S3 GET URL for key {}", storageKey, e);
|
||||
return Optional.empty();
|
||||
}
|
||||
}
|
||||
|
||||
// Returns null when originalFilename is blank; S3 falls back to its own default in that case.
|
||||
static String buildContentDisposition(boolean inline, String originalFilename) {
|
||||
if (originalFilename == null || originalFilename.isBlank()) {
|
||||
return null;
|
||||
}
|
||||
// Strip CR/LF and other control chars before path parsing (Paths.get throws on them on
|
||||
// Windows, and they defeat header parsers).
|
||||
String stripped = originalFilename.replaceAll("\\p{Cntrl}", "");
|
||||
// Use only the basename to avoid leaking directory structure into the header.
|
||||
int lastSeparator = Math.max(stripped.lastIndexOf('/'), stripped.lastIndexOf('\\'));
|
||||
if (lastSeparator >= 0) {
|
||||
stripped = stripped.substring(lastSeparator + 1);
|
||||
}
|
||||
if (stripped.isBlank()) {
|
||||
return null;
|
||||
}
|
||||
// Escape per RFC 6266 quoted-string rules.
|
||||
String escaped = stripped.replace("\\", "\\\\").replace("\"", "\\\"");
|
||||
return (inline ? "inline" : "attachment") + "; filename=\"" + escaped + "\"";
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
try {
|
||||
s3Presigner.close();
|
||||
} catch (Exception e) {
|
||||
log.warn("Error closing S3 presigner", e);
|
||||
}
|
||||
try {
|
||||
s3Client.close();
|
||||
} catch (Exception e) {
|
||||
log.warn("Error closing S3 client", e);
|
||||
}
|
||||
}
|
||||
|
||||
private String sanitizeFilename(String filename) {
|
||||
if (filename == null || filename.isBlank()) {
|
||||
return "file";
|
||||
}
|
||||
String stripped = Paths.get(filename).getFileName().toString().replaceAll("\\p{Cntrl}", "");
|
||||
return stripped.isBlank() ? "file" : stripped;
|
||||
}
|
||||
}
|
||||
+30
-1
@@ -1,16 +1,45 @@
|
||||
package stirling.software.proprietary.storage.provider;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.net.URI;
|
||||
import java.time.Duration;
|
||||
import java.util.Optional;
|
||||
|
||||
import org.springframework.core.io.Resource;
|
||||
import org.springframework.web.multipart.MultipartFile;
|
||||
|
||||
import stirling.software.proprietary.security.model.User;
|
||||
|
||||
public interface StorageProvider {
|
||||
public interface StorageProvider extends AutoCloseable {
|
||||
StoredObject store(User owner, MultipartFile file) throws IOException;
|
||||
|
||||
Resource load(String storageKey) throws IOException;
|
||||
|
||||
void delete(String storageKey) throws IOException;
|
||||
|
||||
/**
|
||||
* Releases any backend-specific resources. Default no-op so {@link LocalStorageProvider} and
|
||||
* {@link DatabaseStorageProvider} (which hold no closeable handles) satisfy Spring's
|
||||
* {@code @Bean(destroyMethod = "close")} signature requirement without ceremony. {@code
|
||||
* S3StorageProvider} overrides this to close the underlying SDK client + presigner.
|
||||
*/
|
||||
@Override
|
||||
default void close() {}
|
||||
|
||||
/**
|
||||
* Returns a presigned download URL valid for {@code ttl}, or {@link Optional#empty()} if the
|
||||
* provider does not support signed URLs (callers fall back to {@link #load(String)}).
|
||||
*/
|
||||
default Optional<URI> signedDownloadUrl(String storageKey, Duration ttl) throws IOException {
|
||||
return signedDownloadUrl(storageKey, ttl, false, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Like {@link #signedDownloadUrl(String, Duration)} with explicit Content-Disposition control.
|
||||
*/
|
||||
default Optional<URI> signedDownloadUrl(
|
||||
String storageKey, Duration ttl, boolean inline, String originalFilename)
|
||||
throws IOException {
|
||||
return Optional.empty();
|
||||
}
|
||||
}
|
||||
|
||||
Some files were not shown because too many files have changed in this diff Show More
Reference in New Issue
Block a user