Skip to content

Commit 2575e4c

Browse files
Merge pull request #53 from SolidLabResearch/fix/source-pod-authentication
Merge protected source-Pod authentication review fixes.
2 parents 26609f9 + 61c8a72 commit 2575e4c

18 files changed

Lines changed: 386 additions & 211 deletions

‎README.md‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -136,13 +136,13 @@ Send a JSON message containing a query, a processing type, and a client identifi
136136

137137
## Credentials
138138

139-
### Normal live operation
139+
Heimdall can access both public and protected Solid streams. Authentication is a property of the source Pod or stream deployment, not of whether a query is `live` or `historical+live`.
140140

141-
Normal `live` operation, including the Sensors-style live evaluation path, does not load source-Pod client credentials.
141+
For a protected source, copy [source-pod-credentials.example.json](./config/source-pod-credentials.example.json) to `config/source-pod-credentials.local.json`, fill it locally, and keep it untracked. Alternatively, set `HEIMDALL_SOURCE_POD_CREDENTIALS_FILE` to a local credential-file path. Each key is a source stream URL or a Pod URL prefix; matching is restricted to the same origin and path boundary, and the most-specific matching entry is selected. Same-origin resources discovered from a configured stream can reuse that source session; cross-origin notification resources require their own matching configuration. The entry contains a CSS client-credential `id`, `secret`, and `idp`.
142142

143-
### historical+live
143+
When configured, Heimdall reuses the resulting authenticated session for relevant source operations: Type Index lookup, LDES metadata and historical retrieval, notification inbox and subscription-server discovery, notification subscription creation, and notification-event retrieval. When no matching source credentials are configured, these operations use ordinary unauthenticated HTTP, so public streams do not need a credentials file. A protected resource with no usable credentials reports the underlying authorization failure; an invalid configured entry reports a configuration error.
144144

145-
`historical+live` requires a client-credential entry for each source stream. Copy [source-pod-credentials.example.json](./config/source-pod-credentials.example.json) to `config/source-pod-credentials.local.json`, fill it locally, and keep it untracked. Alternatively, set `HEIMDALL_SOURCE_POD_CREDENTIALS_FILE` to a local credential-file path.
145+
This does not assert that the Sensors evaluation used authentication; its deployment-specific authorization setup is not established by this repository.
146146

147147
### Legacy aggregation-Pod functionality
148148

‎src/config/privateCredentials.test.ts‎

Lines changed: 15 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
import * as fs from 'fs';
22
import * as os from 'os';
33
import * as path from 'path';
4-
import { aggregationPodAccountFile, loadAggregationPodCredentials, loadSourcePodCredentials } from './privateCredentials';
4+
import { aggregationPodAccountFile, loadAggregationPodCredentials, loadOptionalSourcePodCredentials, loadSourcePodCredentials } from './privateCredentials';
55

66
describe('private credential configuration', () => {
77
const previousSource = process.env.HEIMDALL_SOURCE_POD_CREDENTIALS_FILE;
@@ -36,4 +36,18 @@ describe('private credential configuration', () => {
3636
delete process.env.HEIMDALL_SOURCE_POD_CREDENTIALS_FILE;
3737
expect(() => loadSourcePodCredentials()).toThrow('Missing source-Pod credentials');
3838
});
39+
40+
it('treats absent source credentials as an optional public-stream configuration', () => {
41+
delete process.env.HEIMDALL_SOURCE_POD_CREDENTIALS_FILE;
42+
expect(loadOptionalSourcePodCredentials()).toEqual({});
43+
});
44+
45+
it('fails clearly for an explicitly missing or malformed source credential file', () => {
46+
process.env.HEIMDALL_SOURCE_POD_CREDENTIALS_FILE = path.join(os.tmpdir(), 'does-not-exist-heimdall-source.json');
47+
expect(() => loadOptionalSourcePodCredentials()).toThrow('Missing source-Pod credentials file');
48+
const malformedPath = path.join(fs.mkdtempSync(path.join(os.tmpdir(), 'heimdall-credentials-')), 'malformed.json');
49+
fs.writeFileSync(malformedPath, '{not-json');
50+
process.env.HEIMDALL_SOURCE_POD_CREDENTIALS_FILE = malformedPath;
51+
expect(() => loadOptionalSourcePodCredentials()).toThrow('Unable to load source-Pod credentials');
52+
});
3953
});

‎src/config/privateCredentials.ts‎

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,24 @@ export function loadSourcePodCredentials(): SourcePodCredentials {
1616
return loadJson('HEIMDALL_SOURCE_POD_CREDENTIALS_FILE', 'config/source-pod-credentials.local.json', 'source-Pod credentials');
1717
}
1818

19+
/**
20+
* Source-Pod authentication is optional: public streams must not require a
21+
* local credential file. Invalid configured files still fail loudly.
22+
*/
23+
export function loadOptionalSourcePodCredentials(): SourcePodCredentials {
24+
const configuredPath = process.env.HEIMDALL_SOURCE_POD_CREDENTIALS_FILE;
25+
const filePath = configuredPath || path.resolve(process.cwd(), 'config/source-pod-credentials.local.json');
26+
if (!fs.existsSync(filePath)) {
27+
if (configuredPath) throw new Error(`Missing source-Pod credentials file configured by HEIMDALL_SOURCE_POD_CREDENTIALS_FILE: ${filePath}`);
28+
return {};
29+
}
30+
try {
31+
return JSON.parse(fs.readFileSync(filePath, 'utf8')) as SourcePodCredentials;
32+
} catch (error) {
33+
throw new Error(`Unable to load source-Pod credentials from ${filePath}: ${(error as Error).message}`);
34+
}
35+
}
36+
1937
export function loadAggregationPodCredentials(): AggregationPodCredentials {
2038
return loadJson('HEIMDALL_AGGREGATION_POD_CREDENTIALS_FILE', 'config/aggregation-pod-credentials.local.json', 'aggregation-Pod credentials');
2139
}

‎src/server/WebSocketHandler.test.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -210,7 +210,7 @@ WHERE {
210210

211211
expect(result.ldes_query).toContain(`<${streamUrl}>`);
212212
expect(result.ldes_query).not.toContain('undefined');
213-
expect(discovery).toHaveBeenCalledWith(podUrl, expect.any(Array));
213+
expect(discovery).toHaveBeenCalledWith(podUrl, expect.any(Array), undefined, expect.anything());
214214
const metrics = fs.readFileSync(path.join((handler as any).metric_writer.resultsDir, 'initialization.csv'), 'utf8');
215215
expect(metrics).toContain('stream_discovery');
216216
discovery.mockRestore();

‎src/server/WebSocketHandler.ts‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ import { MetricWriter } from '../evaluation/MetricWriter';
1919
import { createHash } from 'crypto';
2020
import { SharedStreamRegistry } from '../service/heimdall/SharedStreamRegistry';
2121
import { loadAggregationPodCredentials } from '../config/privateCredentials';
22+
import { SourcePodAccess } from '../service/heimdall/SourcePodAccess';
2223

2324
/**
2425
* Class for handling the Websocket server.
@@ -37,6 +38,7 @@ export class WebSocketHandler {
3738
public logger: any;
3839
private query_registry: QueryRegistry;
3940
private readonly metric_writer: MetricWriter;
41+
private readonly sourcePodAccess: SourcePodAccess;
4042
private aggregationPublisherRegistered = false;
4143
/**
4244
* Creates an instance of WebSocketHandler.
@@ -55,6 +57,7 @@ export class WebSocketHandler {
5557
this.connections = new Map<string, WebSocket[]>();
5658
this.parser = new RSPQLParser();
5759
this.metric_writer = metric_writer;
60+
this.sourcePodAccess = sharedStreamRegistry?.sourcePodAccess || new SourcePodAccess();
5861
this.query_registry = new QueryRegistry(metric_writer, sharedStreamRegistry);
5962
this.n3_parser = new Parser({ format: 'N-Triples' });
6063
this.logger.info({}, 'websocket_handler_initialized');
@@ -333,7 +336,7 @@ export class WebSocketHandler {
333336
const interest_metric = new AggregationFocusExtractor(query).extract_focus();
334337
const start_epoch_ms = Date.now();
335338
const start_monotonic_ns = process.hrtime.bigint();
336-
const streams = await find_relevant_streams(pod_url, interest_metric);
339+
const streams = await find_relevant_streams(pod_url, interest_metric, undefined, this.sourcePodAccess);
337340
const ldes_stream = streams[0];
338341
if (ldes_stream === undefined) {
339342
throw new Error(`No relevant LDES stream found for Pod source ${pod_url}`);

‎src/service/heimdall/DecentralizedFileStreamer.ts‎

Lines changed: 14 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -8,13 +8,12 @@ const websocketConnection = require('websocket').connection;
88
const WebSocketClient = require('websocket').client;
99
import { Quad } from "n3";
1010
import { QuadWithID } from "../../utils/Types";
11-
import { session_with_credentials } from "../../utils/authentication/CSSAuthentication";
1211
import { readMembersRateLimited } from "../../utils/ldes-in-ldp/EventSource";
1312
import { RateLimitedLDPCommunication } from "rate-limited-ldp-communication";
1413
import { hash_string_md5 } from "../../utils/Util";
1514
import { TREE } from "@treecg/ldes-snapshot";
16-
import { Session } from "@inrupt/solid-client-authn-node";
1715
import { create_subscription, extract_ldp_inbox, extract_subscription_server } from "../../utils/notifications/Util";
16+
import { SourcePodAccess } from './SourcePodAccess';
1817
import * as HEIMDALL_SETUP from '../../config/heimdall_setup.json';
1918
import { resolveHeimdallSetupConfig, resolveHeimdallWebSocketUrl } from "../../config/heimdallConfig";
2019
import { HEIMDALL_WEBSOCKET_PROTOCOL } from "../../server/websocketProtocols";
@@ -32,7 +31,7 @@ export class DecentralizedFileStreamer {
3231
public ldes!: LDESinLDP;
3332
public comunica_engine: QueryEngine;
3433
public communication: Promise<SolidCommunication | LDPCommunication | RateLimitedLDPCommunication>;
35-
public session: any;
34+
private readonly sourcePodAccess: SourcePodAccess;
3635
public observation_array: any[];
3736
public query: string
3837
public query_hash: string;
@@ -43,17 +42,18 @@ export class DecentralizedFileStreamer {
4342
/**
4443
* Creates an instance of DecentralizedFileStreamer.
4544
* @param {string} ldes_stream - The LDES stream URL.
46-
* @param {session_credentials} session_credentials - The credentials of the Solid Pod.
45+
* @param {SourcePodAccess} sourcePodAccess - Shared optional authenticated access for the source Pod.
4746
* @param {Date} from_date - The start date of the events to be read from the Solid Pod.
4847
* @param {Date} to_date - The end date of the events to be read from the Solid Pod.
4948
* @param {RSPEngine} rsp_engine - The RSP Engine.
5049
* @param {string} query - The query to be executed.
5150
* @param {*} logger - The logger object.
5251
* @memberof DecentralizedFileStreamer
5352
*/
54-
constructor(ldes_stream: string, session_credentials: session_credentials, from_date: Date, to_date: Date, rsp_engine: RSPEngine, query: string, logger: any) {
53+
constructor(ldes_stream: string, sourcePodAccess: SourcePodAccess, from_date: Date, to_date: Date, rsp_engine: RSPEngine, query: string, logger: any) {
5554
this.ldes_stream = ldes_stream;
56-
this.communication = this.get_communication(session_credentials);
55+
this.sourcePodAccess = sourcePodAccess;
56+
this.communication = this.sourcePodAccess.communicationFor(ldes_stream);
5757
this.from_date = from_date;
5858
this.to_date = to_date;
5959
this.query = query;
@@ -63,7 +63,6 @@ export class DecentralizedFileStreamer {
6363
this.stream_name = rsp_engine.getStream(this.ldes_stream);
6464
this.comunica_engine = new QueryEngine();
6565
this.observation_array = [];
66-
this.subscribing_latest_events(this.stream_name);
6766
DecentralizedFileStreamer.connect_with_server(resolveHeimdallWebSocketUrl(HEIMDALL_SETUP)).then(() => {
6867
console.log(`The connection with the websocket server was established.`);
6968
});
@@ -78,16 +77,6 @@ export class DecentralizedFileStreamer {
7877
* @returns {Promise<SolidCommunication | LDPCommunication>} - The communication object with the Solid Pod.
7978
* @memberof DecentralizedFileStreamer
8079
*/
81-
public async get_communication(credentials: session_credentials) {
82-
const session = await this.get_session(credentials);
83-
if (session) {
84-
return new SolidCommunication(session);
85-
}
86-
else {
87-
return new LDPCommunication();
88-
}
89-
}
90-
9180
/**
9281
* Adding the events which might have been added between
9382
* the start of the file streamer and the start of the websocket
@@ -240,12 +229,13 @@ export class DecentralizedFileStreamer {
240229
async subscribing_latest_events(stream_name: RDFStream | undefined) {
241230
if (stream_name !== undefined) {
242231
console.log(`Subscribing to the latest events of the stream`, stream_name);
243-
const inbox = await extract_ldp_inbox(this.ldes_stream);
232+
const sourceFetch = await this.sourcePodAccess.fetchFor(this.ldes_stream);
233+
const inbox = await extract_ldp_inbox(this.ldes_stream, undefined, sourceFetch);
244234
if (inbox !== undefined) {
245-
const subscription_server = await extract_subscription_server(inbox);
235+
const subscription_server = await extract_subscription_server(inbox, undefined, await this.sourcePodAccess.fetchFor(inbox, this.ldes_stream));
246236
if (subscription_server !== undefined) {
247237
const server = subscription_server.location;
248-
const response_subscription = await create_subscription(server, inbox);
238+
const response_subscription = await create_subscription(server, inbox, undefined, await this.sourcePodAccess.fetchFor(server, this.ldes_stream));
249239
if (response_subscription) {
250240
this.logger.info(`The subscription has been succesful.`);
251241
} else {
@@ -274,7 +264,7 @@ export class DecentralizedFileStreamer {
274264
*/
275265
async get_inbox_container(stream: string): Promise<string | undefined> {
276266
console.log(`Getting the inbox container from`, stream);
277-
const ldes_in_ldp: LDESinLDP = new LDESinLDP(stream, new LDPCommunication());
267+
const ldes_in_ldp: LDESinLDP = new LDESinLDP(stream, await this.sourcePodAccess.communicationFor(stream));
278268
const metadata = await ldes_in_ldp.readMetadata();
279269
for (const quad of metadata) {
280270
if (quad.predicate.value === 'http://www.w3.org/ns/ldp#inbox') {
@@ -308,7 +298,7 @@ export class DecentralizedFileStreamer {
308298
"sendTo": `${heimdallSetupConfig.heimdallHttpServerUrl}`
309299
};
310300

311-
const response = await fetch(webhook_notification_server, {
301+
const response = await (await this.sourcePodAccess.fetchFor(webhook_notification_server, ldes_stream))(webhook_notification_server, {
312302
method: 'POST',
313303
headers: {
314304
'Content-Type': 'application/ld+json',
@@ -335,7 +325,7 @@ export class DecentralizedFileStreamer {
335325
"type": "http://www.w3.org/ns/solid/notifications#WebSocketChannel2023",
336326
"topic": `${ldes_stream}`
337327
}
338-
const repsonse = await fetch(notification_server, {
328+
const repsonse = await (await this.sourcePodAccess.fetchFor(notification_server, ldes_stream))(notification_server, {
339329
method: 'POST',
340330
headers: {
341331
'Content-Type': 'application/ld+json',
@@ -379,7 +369,7 @@ export class DecentralizedFileStreamer {
379369
"type": "http://www.w3.org/ns/solid/notifications#WebSocketChannel2023",
380370
"topic": `${inbox_container}`
381371
}
382-
const repsonse = await fetch(notification_server, {
372+
const repsonse = await (await this.sourcePodAccess.fetchFor(notification_server, ldes_stream))(notification_server, {
383373
method: 'POST',
384374
headers: {
385375
'Content-Type': 'application/ld+json',
@@ -408,15 +398,6 @@ export class DecentralizedFileStreamer {
408398
get_file_streamer_start_time() {
409399
return this.file_streamer_start_time;
410400
}
411-
/**
412-
* Get the session with the credentials.
413-
* @param {session_credentials} credentials - The credentials of the solid pod for which you can generate an authenticated session to communicated to the Solid Pod's LDP.
414-
* @returns {Promise<Session>} - The authenticated session.
415-
* @memberof DecentralizedFileStreamer
416-
*/
417-
async get_session(credentials: session_credentials): Promise<Session> {
418-
return await session_with_credentials(credentials);
419-
}
420401
/**
421402
* Send a message to Heimdall's WebSocket server.
422403
* @static

0 commit comments

Comments
 (0)